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/venues/boot.ex
defmodule Tai.Venues.Boot do
@moduledoc """
Coordinates the asynchronous hydration of a venue:
- products
- asset balances
- fees
"""
alias Tai.Venues.Boot
@type adapter :: Tai.Venues.Adapter.t()
@spec run(adapter :: adapter) :: {:ok, adapter} | {:error, {adapter, [reason :: term]}}
def run(%Tai.Venues.Adapter{} = adapter) do
adapter
|> hydrate_products_and_balances
|> wait_for_products
|> hydrate_fees_and_positions_and_start_streams
|> wait_for_balances_and_fees
end
defp hydrate_products_and_balances(adapter) do
t_products = Task.async(Boot.Products, :hydrate, [adapter])
t_balances = Task.async(Boot.AssetBalances, :hydrate, [adapter])
{adapter, t_products, t_balances}
end
defp wait_for_products({adapter, t_products, t_balances}) do
working_tasks = [asset_balances: t_balances]
case Task.await(t_products, adapter.timeout) do
{:ok, products} ->
{:ok, adapter, working_tasks, products}
{:error, reason} ->
err_reasons = [products: reason]
{:error, adapter, working_tasks, err_reasons}
end
end
defp hydrate_fees_and_positions_and_start_streams({:ok, adapter, working_tasks, products}) do
t_fees = Task.async(Boot.Fees, :hydrate, [adapter, products])
t_positions = Task.async(Boot.Positions, :hydrate, [adapter])
t_stream = Task.async(Boot.Stream, :start, [adapter, products])
t_order_books = Task.async(Boot.OrderBooks, :start, [adapter, products])
new_working_tasks = [{:fees, t_fees} | working_tasks]
new_working_tasks = [{:positions, t_positions} | new_working_tasks]
new_working_tasks = [{:order_books, t_stream} | new_working_tasks]
new_working_tasks = [{:order_books, t_order_books} | new_working_tasks]
{:ok, adapter, new_working_tasks}
end
defp hydrate_fees_and_positions_and_start_streams({:error, _, _, _} = error), do: error
defp wait_for_balances_and_fees({:ok, adapter, working_tasks}) do
adapter
|> collect_remaining_errors(working_tasks, [])
end
defp wait_for_balances_and_fees({:error, adapter, working_tasks, err_reasons}) do
adapter
|> collect_remaining_errors(working_tasks, err_reasons)
end
defp collect_remaining_errors(adapter, [], err_reasons) do
if Enum.empty?(err_reasons) do
{:ok, adapter}
else
{:error, {adapter, err_reasons}}
end
end
defp collect_remaining_errors(adapter, [{name, working} | tasks], err_reasons) do
case Task.await(working, adapter.timeout) do
{:error, reason} ->
adapter |> collect_remaining_errors(tasks, [{name, reason} | err_reasons])
_ ->
adapter |> collect_remaining_errors(tasks, err_reasons)
end
end
end