Packages
tai
0.0.2
0.0.75
0.0.74
0.0.73
0.0.72
0.0.71
0.0.70
0.0.69
0.0.68
0.0.67
0.0.66
0.0.65
0.0.64
0.0.63
0.0.62
0.0.61
0.0.60
0.0.59
0.0.58
0.0.57
0.0.56
0.0.55
0.0.54
0.0.53
0.0.52
0.0.51
0.0.50
0.0.49
0.0.48
0.0.47
0.0.46
0.0.45
0.0.44
0.0.43
0.0.42
0.0.41
0.0.40
0.0.39
0.0.38
0.0.37
0.0.36
0.0.35
0.0.34
0.0.33
0.0.32
0.0.31
0.0.30
0.0.29
0.0.28
0.0.27
0.0.26
0.0.25
0.0.24
0.0.23
0.0.22
0.0.21
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14
0.0.13
0.0.12
0.0.11
0.0.10
0.0.9
0.0.8
0.0.7
0.0.6
0.0.5
0.0.4
0.0.3
0.0.2
0.0.1
A composable, real time, market data and trade execution toolkit
Current section
Files
Jump to
Current section
Files
lib/tai/advisor.ex
defmodule Tai.Advisor do
@moduledoc """
A behavior for implementing a server that receives order book changes.
It can be used to monitor multiple quote streams and create, update or cancel orders.
"""
alias Tai.{Advisor, PubSub, Trading.Order}
alias Tai.Markets.{OrderBook, Quote}
@typedoc """
The state of the running advisor
"""
@type t :: Advisor
@enforce_keys [:advisor_id, :exchanges, :order_books, :inside_quotes, :store]
defstruct advisor_id: nil, exchanges: [], order_books: %{}, inside_quotes: %{}, store: %{}
@doc """
Callback when order book has bid or ask changes
"""
@callback handle_order_book_changes(order_book_feed_id :: Atom.t, symbol :: Atom.t, changes :: term, state :: Advisor.t) :: :ok
@doc """
Callback when the highest bid or lowest ask changes price or size
"""
@callback handle_inside_quote(order_book_feed_id :: Atom.t, symbol :: Atom.t, inside_quote :: Quote.t, changes :: Map.t | List.t, state :: Advisor.t) :: :ok | {:ok, actions :: Map.t}
@doc """
Callback when an order is enqueued
"""
@callback handle_order_enqueued(order :: Order.t, state :: Advisor.t) :: :ok
@doc """
Callback when an order is created on the server
"""
@callback handle_order_create_ok(order :: Order.t, state :: Advisor.t) :: :ok
@doc """
Callback when an order creation fails
"""
@callback handle_order_create_error(reason :: term, order :: Order.t, state :: Advisor.t) :: :ok
@doc """
Callback when an order has been cancelled in the outbox but before the
request has been sent to the exchange.
"""
@callback handle_order_cancelling(order :: Order.t, state :: Advisor.t) :: :ok
@doc """
Callback when an order has been cancelled on the exchange
"""
@callback handle_order_cancelled(order :: Order.t, state :: Advisor.t) :: :ok
@doc """
Returns an atom that will identify the process
## Examples
iex> Tai.Advisor.to_name(:my_test_advisor)
:advisor_my_test_advisor
"""
def to_name(advisor_id), do: :"advisor_#{advisor_id}"
defmacro __using__(_) do
quote location: :keep do
use GenServer
require Logger
require Tai.TimeFrame
alias Tai.{Advisor, Markets.OrderBook, Trading.OrderOutbox}
@behaviour Advisor
def start_link(advisor_id: advisor_id, order_books: order_books, exchanges: exchanges) do
GenServer.start_link(
__MODULE__,
%Advisor{
advisor_id: advisor_id,
order_books: order_books,
exchanges: exchanges,
inside_quotes: %{},
store: %{}
},
name: advisor_id |> Advisor.to_name
)
end
@doc false
def init(%Advisor{order_books: order_books, exchanges: exchanges} = state) do
subscribe_to_order_book_channels(order_books)
subscribe_to_exchange_channels(exchanges)
{:ok, state}
end
@doc false
def handle_info({:order_book_snapshot, order_book_feed_id, symbol, snapshot}, state) do
new_state = state
|> cache_inside_quote(order_book_feed_id, symbol)
|> execute_handle_inside_quote(order_book_feed_id, symbol, snapshot)
{:noreply, new_state}
end
@doc false
def handle_info({:order_book_changes, order_book_feed_id, symbol, changes}, state) do
new_state = Tai.TimeFrame.debug "[#{state.advisor_id |> Advisor.to_name}] handle_info({:order_book_changes...})" do
handle_order_book_changes(order_book_feed_id, symbol, changes, state)
previous_inside_quote = state |> cached_inside_quote(order_book_feed_id, symbol)
if inside_quote_is_stale?(previous_inside_quote, changes) do
state
|> cache_inside_quote(order_book_feed_id, symbol)
|> execute_handle_inside_quote(order_book_feed_id, symbol, changes, previous_inside_quote)
else
state
end
end
{:noreply, new_state}
end
@doc false
def handle_info({:order_enqueued, order}, state) do
handle_order_enqueued(order, state)
{:noreply, state}
end
@doc false
def handle_info({:order_create_ok, order}, state) do
handle_order_create_ok(order, state)
{:noreply, state}
end
@doc false
def handle_info({:order_create_error, reason, order}, state) do
handle_order_create_error(reason, order, state)
{:noreply, state}
end
@doc false
def handle_info({:order_cancelling, order}, state) do
handle_order_cancelling(order, state)
{:noreply, state}
end
@doc false
def handle_info({:order_cancelled, order}, state) do
handle_order_cancelled(order, state)
{:noreply, state}
end
@doc """
Returns the current state of the order book up to the given depth
## Examples
iex> Tai.Advisor.quotes(feed_id: :test_feed_a, symbol: :btcusd, depth: 1)
{:ok, %Tai.Markets.OrderBook{bids: [], asks: []}
"""
def quotes(feed_id: order_book_feed_id, symbol: symbol, depth: depth) do
[feed_id: order_book_feed_id, symbol: symbol]
|> OrderBook.to_name
|> OrderBook.quotes(depth)
end
@doc """
Returns the inside quote stored before the last 'handle_inside_quote' callback
## Examples
iex> Tai.Advisor.cached_inside_quote(state, :test_feed_a, :btcusd)
%Tai.Markets.Quote{
bid: %Tai.Markets.PriceLevel{price: 101.1, size: 1.1, processed_at: nil, server_changed_at: nil},
ask: %Tai.Markets.PriceLevel{price: 101.2, size: 0.1, processed_at: nil, server_changed_at: nil}
}
"""
def cached_inside_quote(%{inside_quotes: inside_quotes}, order_book_feed_id, symbol) do
inside_quotes
|> Map.get([feed_id: order_book_feed_id, symbol: symbol] |> OrderBook.to_name)
end
def handle_order_book_changes(order_book_feed_id, symbol, changes, state), do: :ok
def handle_inside_quote(order_book_feed_id, symbol, inside_quote, changes, state), do: :ok
def handle_order_enqueued(order, state), do: :ok
def handle_order_create_ok(order, state), do: :ok
def handle_order_create_error(reason, order, state), do: :ok
def handle_order_cancelling(order, state), do: :ok
def handle_order_cancelled(order, state), do: :ok
defp subscribe_to_order_book_channels(order_books) do
order_books
|> Enum.each(fn {order_book_feed_id, symbols} ->
symbols
|> Enum.each(fn symbol ->
PubSub.subscribe([
{:order_book_snapshot, order_book_feed_id, symbol},
{:order_book_changes, order_book_feed_id, symbol}
])
end)
end)
end
defp subscribe_to_exchange_channels(exchanges) do
PubSub.subscribe(exchanges)
end
defp fetch_inside_quote(order_book_feed_id, symbol) do
[feed_id: order_book_feed_id, symbol: symbol, depth: 1]
|> quotes
|> case do
{:ok, %OrderBook{bids: bids, asks: asks}} ->
%Quote{bid: bids |> List.first, ask: asks |> List.first}
end
end
defp cache_inside_quote(state, order_book_feed_id, symbol) do
current_inside_quote = fetch_inside_quote(order_book_feed_id, symbol)
order_book_key = [feed_id: order_book_feed_id, symbol: symbol] |> OrderBook.to_name
new_inside_quotes = state.inside_quotes |> Map.put(order_book_key, current_inside_quote)
state |> Map.put(:inside_quotes, new_inside_quotes)
end
defp inside_quote_is_stale?(previous_inside_quote, %OrderBook{bids: bids, asks: asks} = changes) do
(bids |> Enum.any? && bids |> inside_bid_is_stale?(previous_inside_quote)) || (asks |> Enum.any? && asks |> inside_ask_is_stale?(previous_inside_quote))
end
defp inside_bid_is_stale?(_bids, nil), do: false
defp inside_bid_is_stale?(bids, %Quote{} = prev_quote) do
bids
|> Enum.any?(fn {price, {size, _processed_at, _server_changed_at}} ->
price >= prev_quote.bid.price || (price == prev_quote.bid.price && size != prev_quote.bid.size)
end)
end
defp inside_ask_is_stale?(asks, nil), do: false
defp inside_ask_is_stale?(asks, %Quote{} = prev_quote) do
asks
|> Enum.any?(fn {price, {size, _processed_at, _server_changed_at}} ->
price <= prev_quote.ask.price || (price == prev_quote.ask.price && size != prev_quote.ask.size)
end)
end
defp execute_handle_inside_quote(state, order_book_feed_id, symbol, changes, previous_inside_quote \\ nil) do
current_inside_quote = cached_inside_quote(state, order_book_feed_id, symbol)
if current_inside_quote == previous_inside_quote do
state
else
handle_inside_quote(order_book_feed_id, symbol, current_inside_quote, changes, state)
|> normalize_handler_response
|> cancel_orders
|> submit_orders
|> update_store(state)
end
end
defp normalize_handler_response(:ok) do
{:ok, %{}}
|> normalize_handler_response
end
defp normalize_handler_response({:ok, actions}) do
default_actions = %{cancel_orders: [], limit_orders: []}
{:ok, default_actions |> Map.merge(actions)}
end
defp cancel_orders({:ok, %{cancel_orders: cancel_orders}} = handler_response) do
cancel_orders
|> OrderOutbox.cancel
handler_response
end
defp submit_orders({:ok, %{limit_orders: limit_orders}} = handler_response) do
limit_orders
|> OrderOutbox.add
handler_response
end
defp update_store({:ok, %{store: store}}, state), do: state |> Map.put(:store, store)
defp update_store({:ok, %{}}, state), do: state
defoverridable [
handle_order_book_changes: 4,
handle_inside_quote: 5,
handle_order_enqueued: 2,
handle_order_create_ok: 2,
handle_order_create_error: 3,
handle_order_cancelling: 2,
handle_order_cancelled: 2
]
end
end
end