Packages
tai
0.0.13
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/markets/order_book.ex
defmodule Tai.Markets.OrderBook do
@moduledoc """
Manage and query the state for an order book for a symbol on a feed
"""
use GenServer
alias Tai.{Markets, PubSub}
@type t :: %Markets.OrderBook{
venue_id: atom,
product_symbol: atom,
bids: map,
asks: map
}
@enforce_keys [
:venue_id,
:product_symbol,
:bids,
:asks
]
defstruct [
:venue_id,
:product_symbol,
:bids,
:asks
]
def start_link(feed_id: venue_id, symbol: product_symbol) do
name = to_name(venue_id, product_symbol)
order_book = %Markets.OrderBook{
venue_id: venue_id,
product_symbol: product_symbol,
bids: %{},
asks: %{}
}
GenServer.start_link(__MODULE__, order_book, name: name)
end
def init(state), do: {:ok, state}
def handle_call({:quotes, depth: depth}, _from, state) do
order_book = %Markets.OrderBook{
venue_id: state.venue_id,
product_symbol: state.product_symbol,
bids: state |> ordered_bids |> take(depth),
asks: state |> ordered_asks |> take(depth)
}
{:reply, {:ok, order_book}, state}
end
def handle_call({:replace, snapshot}, _from, _state) do
PubSub.broadcast(
{:order_book_snapshot, snapshot.venue_id, snapshot.product_symbol},
{:order_book_snapshot, snapshot.venue_id, snapshot.product_symbol, snapshot}
)
Tai.Events.broadcast(%Tai.Events.OrderBookSnapshot{
venue_id: snapshot.venue_id,
symbol: snapshot.product_symbol,
snapshot: snapshot
})
{:reply, :ok, snapshot}
end
def handle_call({:update, %Markets.OrderBook{bids: bids, asks: asks} = changes}, _from, state) do
PubSub.broadcast(
{:order_book_changes, state.venue_id, state.product_symbol},
{:order_book_changes, state.venue_id, state.product_symbol, changes}
)
new_state =
state
|> update_side(:bids, bids)
|> update_side(:asks, asks)
{:reply, :ok, new_state}
end
@doc """
Return bid/asks up to the given depth. If depth is not provided it returns
the full order book.
"""
def quotes(name, depth \\ :all) do
GenServer.call(name, {:quotes, depth: depth})
end
@doc """
Return the bid/ask at the top of the book
"""
@spec inside_quote(atom, atom) :: {:ok, Markets.Quote.t()}
def inside_quote(venue_id, product_symbol) do
name = to_name(venue_id, product_symbol)
name
|> quotes(1)
|> case do
{:ok, %{bids: bids, asks: asks}} ->
inside_bid = List.first(bids)
inside_ask = List.first(asks)
q = %Markets.Quote{
venue_id: venue_id,
product_symbol: product_symbol,
bid: inside_bid,
ask: inside_ask
}
{:ok, q}
end
end
@spec replace(t) :: :ok
def replace(%Markets.OrderBook{} = replacement) do
replacement.venue_id
|> Markets.OrderBook.to_name(replacement.product_symbol)
|> GenServer.call({:replace, replacement})
end
@deprecated "use Tai.Markets.OrderBook.update/1 instead"
def update(name, %Markets.OrderBook{} = changes) do
GenServer.call(name, {:update, changes})
end
@spec update(t) :: :ok
def update(%Markets.OrderBook{} = changes) do
changes.venue_id
|> Markets.OrderBook.to_name(changes.product_symbol)
|> GenServer.call({:update, changes})
end
@spec to_name(atom, atom) :: atom
def to_name(venue_id, product_symbol) do
:"#{__MODULE__}_#{venue_id}_#{product_symbol}"
end
defp ordered_bids(state) do
state.bids
|> Map.keys()
|> Enum.sort()
|> Enum.reverse()
|> with_price_levels(state.bids)
end
defp ordered_asks(state) do
state.asks
|> Map.keys()
|> Enum.sort()
|> with_price_levels(state.asks)
end
defp with_price_levels(prices, level_details) do
prices
|> Enum.map(fn price ->
{size, processed_at, server_changed_at} = level_details[price]
%Markets.PriceLevel{
price: price,
size: size,
processed_at: processed_at,
server_changed_at: server_changed_at
}
end)
end
defp take(list, :all), do: list
defp take(list, depth) do
list
|> Enum.take(depth)
end
defp update_side(state, side, price_levels) do
new_side =
state
|> Map.get(side)
|> Map.merge(price_levels)
|> Map.drop(price_levels |> drop_prices)
state
|> Map.put(side, new_side)
end
defp drop_prices(price_levels) do
price_levels
|> Enum.filter(fn {_price, {size, _processed_at, _server_changed_at}} -> size == 0 end)
|> Enum.map(fn {price, _} -> price end)
end
end