diff --git a/atra-gateway/lib/atra_gateway/matching_engine.ex b/atra-gateway/lib/atra_gateway/matching_engine.ex index 731892c..0cef41d 100644 --- a/atra-gateway/lib/atra_gateway/matching_engine.ex +++ b/atra-gateway/lib/atra_gateway/matching_engine.ex @@ -2,77 +2,118 @@ defmodule AtraGateway.MatchingEngine do use GenServer require Logger alias AtraGateway.Orders - + def start_link(_) do GenServer.start_link(__MODULE__, nil, name: __MODULE__) end - + def init(_) do - :ets.new(:order_buffer, [:named_table, :public, write_concurrency: true]) - schedule_batch_processing() - {:ok, %{batch_size: 1000, batch_interval_ms: 100}} + lane_count = lane_count() + {:ok, %{next_seq: %{}, lane_count: lane_count}} end def place_order(order_params) do - # put order in ets buffer - true = :ets.insert(:order_buffer, {System.monotonic_time(), order_params}) - {:ok, %{status: :accepted}} + GenServer.call(__MODULE__, {:place_order, order_params}, 5_000) end - defp schedule_batch_processing do - Process.send_after(self(), :process_batch, 100) + def place_orders(order_params_list) do + GenServer.call(__MODULE__, {:place_orders, order_params_list}, 10_000) end - def handle_info(:process_batch, %{batch_size: batch_size} = state) do - # take <= batch_size orders from the buffer - orders = :ets.take(:order_buffer, batch_size) - - if orders != [] do - orders - |> Enum.sort_by(&elem(&1, 0)) # sort by timestamp - |> Enum.map(&elem(&1, 1)) # take just the order params - |> Enum.chunk_every(50) # process in smaller chunks - |> Task.async_stream(&process_order_chunk/1, - max_concurrency: 10, - timeout: 5000) - |> Stream.run() + def handle_call({:place_order, order_params}, _from, state) do + with {:ok, sequenced, next_state} <- sequence_order(order_params, state), + {:ok, response} <- grpc_place_order(sequenced) do + {:reply, {:ok, response}, next_state} + else + {:error, reason} -> + Logger.error("order placement failed: #{inspect(reason)}") + {:reply, {:error, reason}, state} end + end + + def handle_call({:place_orders, order_params_list}, _from, state) do + {responses, new_state} = + Enum.reduce(order_params_list, {[], state}, fn order_params, {acc, st} -> + case sequence_order(order_params, st) do + {:ok, sequenced, next_state} -> + case grpc_place_order(sequenced) do + {:ok, response} -> {[response | acc], next_state} + {:error, _} -> {acc, st} + end + + {:error, _} -> + {acc, st} + end + end) + + {:reply, {:ok, Enum.reverse(responses)}, new_state} + end + + defp sequence_order(order_params, state) do + instrument_id = Map.fetch!(order_params, :instrument_id) + next_seq = Map.get(state.next_seq, instrument_id, 1) + provided_seq = Map.get(order_params, :sequence_number, nil) + strict? = strict_sequence_validation?() + + sequence_number = + cond do + is_nil(provided_seq) -> next_seq + strict? and provided_seq != next_seq -> :out_of_order + true -> provided_seq + end - schedule_batch_processing() - {:noreply, state} + if sequence_number == :out_of_order do + {:error, :out_of_order_sequence} + else + ingress_timestamp_ns = System.os_time(:nanosecond) + order = + order_params + |> Map.put(:sequence_number, sequence_number) + |> Map.put(:ingress_timestamp_ns, ingress_timestamp_ns) + |> Map.put(:lane_id, lane_for_instrument(instrument_id, state.lane_count)) + + next_expected = max(next_seq, sequence_number + 1) + next_state = put_in(state.next_seq[instrument_id], next_expected) + {:ok, order, next_state} + end end - defp process_order_chunk(orders) do + defp grpc_place_order(params) do :poolboy.transaction(:grpc_pool, fn pid -> channel = AtraGateway.GrpcConnection.get_channel(pid) - - # Convert orders to proto format - proto_requests = Enum.map(orders, fn params -> - %Orderbook.OrderRequest{ + + request = %Orderbook.OrderRequest{ id: params.id, price: to_string(params.price), quantity: to_string(params.quantity), side: proto_side(params.side), - order_type: proto_order_type(params.type)} end) - - request = %Orderbook.OrderBatchRequest{ - orders: proto_requests + order_type: proto_order_type(params.type), + instrument_id: params.instrument_id, + sequence_number: params.sequence_number, + ingress_timestamp_ns: params.ingress_timestamp_ns, + idempotency_key: Map.get(params, :idempotency_key) } - - case Orderbook.OrderBookService.Stub.place_orders(channel, request) do - {:ok, %Orderbook.OrderBatchResponse{orders: responses}} -> - Enum.map(responses, &Orders.from_proto/1) - {:error, reason} -> - Logger.error("Batch order placement failed: #{inspect(reason)}") - {:error, reason} + + case Orderbook.OrderBookService.Stub.place_order(channel, request) do + {:ok, response} -> {:ok, Orders.from_proto(response)} + {:error, reason} -> {:error, reason} end end) end + defp lane_count do + System.get_env("ATRA_GATEWAY_LANE_COUNT", "4") |> String.to_integer() + end + + defp strict_sequence_validation? do + System.get_env("ATRA_GATEWAY_STRICT_SEQUENCE_VALIDATION", "false") in ["1", "true", "TRUE", "yes", "YES"] + end + + defp lane_for_instrument(instrument_id, lane_count), do: rem(instrument_id, lane_count) + defp proto_side(:bid), do: :BID defp proto_side(:ask), do: :ASK defp proto_order_type(:limit), do: :LIMIT defp proto_order_type(:market), do: :MARKET - end diff --git a/atra-gateway/lib/atra_gateway/orders.ex b/atra-gateway/lib/atra_gateway/orders.ex index 9bdd83a..f544982 100644 --- a/atra-gateway/lib/atra_gateway/orders.ex +++ b/atra-gateway/lib/atra_gateway/orders.ex @@ -6,7 +6,7 @@ defmodule AtraGateway.Orders do @doc """ Creates a new order map with required fields. """ - def new(price, quantity, side, type \\ :limit) do + def new(price, quantity, side, type \\ :limit, instrument_id \\ 1) do # we will want a better ID generation strategy id = System.unique_integer([:positive, :monotonic]) @@ -14,6 +14,7 @@ defmodule AtraGateway.Orders do id: id, price: price, quantity: quantity, + instrument_id: instrument_id, side: side, type: type } @@ -31,7 +32,11 @@ defmodule AtraGateway.Orders do side: atom_from_proto_side(response.side), type: atom_from_proto_order_type(response.order_type), status: atom_from_proto_status(response.status), - timestamp: proto_timestamp_to_datetime(response.timestamp) + timestamp: proto_timestamp_to_datetime(response.timestamp), + instrument_id: response.instrument_id, + sequence_number: response.sequence_number, + ingress_timestamp_ns: response.ingress_timestamp_ns, + idempotency_key: response.idempotency_key } end diff --git a/atra-gateway/lib/atra_gateway/server.ex b/atra-gateway/lib/atra_gateway/server.ex index 6f36f5e..911e45d 100644 --- a/atra-gateway/lib/atra_gateway/server.ex +++ b/atra-gateway/lib/atra_gateway/server.ex @@ -11,23 +11,16 @@ defmodule AtraGateway.Server do def place_order(request, _stream) do Logger.info("atra.gateway.request.place_order id:#{inspect(request.id)}") + instrument_id = if request.instrument_id == 0, do: 1, else: request.instrument_id case AtraGateway.MatchingEngine.place_order(AtraGateway.Orders.new( parse_numeric(request.price), parse_numeric(request.quantity), atom_from_proto_side(request.side), - atom_from_proto_order_type(request.order_type) + atom_from_proto_order_type(request.order_type), + instrument_id )) do - {:ok, _} -> - # forwarding messages directly to the matcher isn't good practice. we'll use INQ later. - :poolboy.transaction(:grpc_pool, fn pid -> - channel = AtraGateway.GrpcConnection.get_channel(pid) - case Orderbook.OrderBookService.Stub.place_order(channel, request) do - {:ok, response} -> response - {:error, reason} -> - Logger.error("atra.gateway.error.request request:place_order note:'#{inspect(reason)}'") - raise GRPC.RPCError, status: :internal, message: "Internal error" - end - end) + {:ok, response} -> + to_proto_response(response) {:error, reason} -> Logger.error("atra.gateway.error.request request:place_order note:'#{inspect(reason)}'") raise GRPC.RPCError, status: :internal, message: "Internal error" @@ -88,15 +81,29 @@ defmodule AtraGateway.Server do def place_orders(request, _stream) do Logger.info("atra.gateway.request.place_order batch:#{length(request.orders)}") - :poolboy.transaction(:grpc_pool, fn pid -> - channel = AtraGateway.GrpcConnection.get_channel(pid) - case Orderbook.OrderBookService.Stub.place_orders(channel, request) do - {:ok, response} -> response - {:error, reason} -> - Logger.error("atra.gateway.error.request request:place_orders note:'#{inspect(reason)}'") - raise GRPC.RPCError, status: :internal, message: "Internal error" - end - end) + orders = + request.orders + |> Enum.map(fn req -> + instrument_id = if req.instrument_id == 0, do: 1, else: req.instrument_id + AtraGateway.Orders.new( + parse_numeric(req.price), + parse_numeric(req.quantity), + atom_from_proto_side(req.side), + atom_from_proto_order_type(req.order_type), + instrument_id + ) + end) + + case AtraGateway.MatchingEngine.place_orders(orders) do + {:ok, responses} -> + %Orderbook.OrderBatchResponse{ + orders: Enum.map(responses, &to_proto_response/1) + } + + {:error, reason} -> + Logger.error("atra.gateway.error.request request:place_orders note:'#{inspect(reason)}'") + raise GRPC.RPCError, status: :internal, message: "Internal error" + end end # convert proto enums to atoms @@ -105,5 +112,32 @@ defmodule AtraGateway.Server do defp atom_from_proto_order_type(:LIMIT), do: :limit defp atom_from_proto_order_type(:MARKET), do: :market + + defp to_proto_response(response) do + %Orderbook.OrderResponse{ + id: response.id, + price: to_string(response.price), + quantity: to_string(response.quantity), + remaining_quantity: to_string(response.remaining_quantity), + side: proto_side(response.side), + order_type: proto_order_type(response.type), + status: proto_status(response.status), + instrument_id: response.instrument_id, + sequence_number: response.sequence_number, + ingress_timestamp_ns: response.ingress_timestamp_ns, + idempotency_key: response.idempotency_key + } + end + + defp proto_side(:bid), do: :BID + defp proto_side(:ask), do: :ASK + + defp proto_order_type(:limit), do: :LIMIT + defp proto_order_type(:market), do: :MARKET + + defp proto_status(:pending), do: :PENDING + defp proto_status(:partially_filled), do: :PARTIALLY_FILLED + defp proto_status(:filled), do: :FILLED + defp proto_status(:cancelled), do: :CANCELLED end diff --git a/atra-gateway/lib/proto/orderbook.pb.ex b/atra-gateway/lib/proto/orderbook.pb.ex index 58504fa..2507fe8 100644 --- a/atra-gateway/lib/proto/orderbook.pb.ex +++ b/atra-gateway/lib/proto/orderbook.pb.ex @@ -37,6 +37,10 @@ defmodule Orderbook.OrderRequest do field :quantity, 3, type: :string field :side, 4, type: Orderbook.Side, enum: true field :order_type, 5, type: Orderbook.OrderType, json_name: "orderType", enum: true + field :instrument_id, 6, type: :uint32, json_name: "instrumentId" + field :sequence_number, 7, type: :uint64, json_name: "sequenceNumber" + field :ingress_timestamp_ns, 8, type: :uint64, json_name: "ingressTimestampNs", optional: true + field :idempotency_key, 9, type: :string, json_name: "idempotencyKey", optional: true end defmodule Orderbook.OrderResponse do @@ -52,6 +56,10 @@ defmodule Orderbook.OrderResponse do field :order_type, 6, type: Orderbook.OrderType, json_name: "orderType", enum: true field :status, 7, type: Orderbook.OrderStatus, enum: true field :timestamp, 8, type: Google.Protobuf.Timestamp + field :instrument_id, 9, type: :uint32, json_name: "instrumentId" + field :sequence_number, 10, type: :uint64, json_name: "sequenceNumber" + field :ingress_timestamp_ns, 11, type: :uint64, json_name: "ingressTimestampNs", optional: true + field :idempotency_key, 12, type: :string, json_name: "idempotencyKey", optional: true end defmodule Orderbook.CancelOrderRequest do @@ -115,6 +123,9 @@ defmodule Orderbook.Trade do field :quantity, 4, type: :string field :side, 5, type: Orderbook.Side, enum: true field :timestamp, 6, type: Google.Protobuf.Timestamp + field :maker_sequence_number, 7, type: :uint64, json_name: "makerSequenceNumber" + field :taker_sequence_number, 8, type: :uint64, json_name: "takerSequenceNumber" + field :ingress_timestamp_ns, 9, type: :uint64, json_name: "ingressTimestampNs", optional: true end defmodule Orderbook.TradeHistoryResponse do diff --git a/atra-proto/proto/orderbook.proto b/atra-proto/proto/orderbook.proto index c789aca..61caac3 100644 --- a/atra-proto/proto/orderbook.proto +++ b/atra-proto/proto/orderbook.proto @@ -18,6 +18,10 @@ message OrderRequest { string quantity = 3; Side side = 4; OrderType order_type = 5; + uint32 instrument_id = 6; + uint64 sequence_number = 7; + optional uint64 ingress_timestamp_ns = 8; + optional string idempotency_key = 9; } message OrderResponse { @@ -29,6 +33,10 @@ message OrderResponse { OrderType order_type = 6; OrderStatus status = 7; google.protobuf.Timestamp timestamp = 8; + uint32 instrument_id = 9; + uint64 sequence_number = 10; + optional uint64 ingress_timestamp_ns = 11; + optional string idempotency_key = 12; } message CancelOrderRequest { @@ -64,6 +72,9 @@ message Trade { string quantity = 4; Side side = 5; google.protobuf.Timestamp timestamp = 6; + uint64 maker_sequence_number = 7; + uint64 taker_sequence_number = 8; + optional uint64 ingress_timestamp_ns = 9; } message TradeHistoryResponse {