Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
125 changes: 83 additions & 42 deletions atra-gateway/lib/atra_gateway/matching_engine.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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
9 changes: 7 additions & 2 deletions atra-gateway/lib/atra_gateway/orders.ex
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,15 @@ 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])

%{
id: id,
price: price,
quantity: quantity,
instrument_id: instrument_id,
side: side,
type: type
}
Expand All @@ -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

Expand Down
76 changes: 55 additions & 21 deletions atra-gateway/lib/atra_gateway/server.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand All @@ -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

11 changes: 11 additions & 0 deletions atra-gateway/lib/proto/orderbook.pb.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
11 changes: 11 additions & 0 deletions atra-proto/proto/orderbook.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down
Loading