Packages
tai
0.0.43
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/process_quote.ex
defmodule Tai.Markets.ProcessQuote do
use GenServer
alias Tai.Markets.{Quote, PricePoint}
defmodule State do
@type market_quote :: Tai.Markets.Quote.t()
@type t :: %State{
market_quote: market_quote | nil,
depth: pos_integer
}
@enforce_keys ~w(market_quote depth)a
defstruct ~w(market_quote depth)a
end
@type venue_id :: Tai.Venues.Adapter.venue_id()
@type product :: Tai.Venues.Product.t()
@type product_symbol :: Tai.Venues.Product.symbol()
@spec start_link(product: product, depth: pos_integer) :: GenServer.on_start()
def start_link(product: product, depth: depth) when depth > 0 do
state = %State{market_quote: nil, depth: depth}
name = product.venue_id |> to_name(product.symbol)
GenServer.start_link(__MODULE__, state, name: name)
end
@spec to_name(venue_id, product_symbol) :: atom
def to_name(venue, symbol), do: :"#{__MODULE__}_#{venue}_#{symbol}"
def init(state), do: {:ok, state}
def handle_cast({:order_book_snapshot, order_book, change_set}, state) do
new_market_quote = build_market_quote(order_book, change_set, state.depth)
new_state = state |> Map.put(:market_quote, new_market_quote)
{:noreply, new_state, {:continue, :broadcast_market_quote}}
end
def handle_cast({:order_book_apply, order_book, change_set}, state) do
new_market_quote = build_market_quote(order_book, change_set, state.depth)
if market_quote_changed?(state.market_quote, new_market_quote) do
new_state = state |> Map.put(:market_quote, new_market_quote)
{:noreply, new_state, {:continue, :broadcast_market_quote}}
else
{:noreply, state}
end
end
def handle_continue(:broadcast_market_quote, state) do
msg = {:tai, state.market_quote}
{:market_quote, state.market_quote.venue_id, state.market_quote.product_symbol}
|> Tai.PubSub.broadcast(msg)
:market_quote
|> Tai.PubSub.broadcast(msg)
{:noreply, state}
end
defp build_market_quote(order_book, change_set, depth) do
bids = price_points(order_book.bids, depth, &(&1 > &2))
asks = price_points(order_book.asks, depth, &(&1 < &2))
%Quote{
venue_id: order_book.venue_id,
product_symbol: order_book.product_symbol,
last_venue_timestamp: change_set.last_venue_timestamp,
last_received_at: change_set.last_received_at,
bids: bids,
asks: asks
}
end
defp price_points(side, depth, sort_by) do
side
|> Map.keys()
|> Enum.sort(sort_by)
|> Enum.take(depth)
|> Enum.map(&%PricePoint{price: &1, size: side |> Map.fetch!(&1)})
end
defp market_quote_changed?(nil, %Quote{}), do: true
defp market_quote_changed?(current_market_quote, new_market_quote) do
inside_price_point_changed?(current_market_quote.bids, new_market_quote.bids) ||
inside_price_point_changed?(current_market_quote.asks, new_market_quote.asks)
end
defp inside_price_point_changed?(current, new) do
current |> List.first() != new |> List.first()
end
end