Packages
tai
0.0.39
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/venue_adapters/bitmex/stream/process_auth.ex
defmodule Tai.VenueAdapters.Bitmex.Stream.ProcessAuth do
use GenServer
alias __MODULE__
defmodule State do
@type venue_id :: Tai.Venues.Adapter.venue_id()
@type t :: %State{venue_id: venue_id, tasks: map}
@enforce_keys ~w(venue_id tasks)a
defstruct ~w(venue_id tasks)a
end
@type venue_id :: Tai.Venues.Adapter.venue_id()
def start_link(venue_id: venue_id) do
state = %State{venue_id: venue_id, tasks: %{}}
name = venue_id |> to_name()
GenServer.start_link(__MODULE__, state, name: name)
end
def init(state), do: {:ok, state}
def handle_cast({venue_msg, received_at}, state) do
{:ok, new_state} =
venue_msg
|> extract()
|> process(received_at, state)
{:noreply, new_state}
end
def handle_info({ref, :ok}, state) when is_reference(ref) do
new_tasks = Map.delete(state.tasks, ref)
new_state = Map.put(state, :tasks, new_tasks)
{:noreply, new_state}
end
def handle_info(_msg, state), do: {:noreply, state}
@spec to_name(venue_id) :: atom
def to_name(venue_id), do: :"#{__MODULE__}_#{venue_id}"
defdelegate extract(msg), to: ProcessAuth.VenueMessage
defp process(messages, received_at, state) do
message_tasks =
messages
|> Enum.reduce(
%{},
fn msg, tasks ->
t = Task.async(fn -> ProcessAuth.Message.process(msg, received_at, state) end)
Map.put(tasks, t.ref, t)
end
)
new_tasks = Map.merge(state.tasks, message_tasks)
new_state = Map.put(state, :tasks, new_tasks)
{:ok, new_state}
end
end