-
Notifications
You must be signed in to change notification settings - Fork 1
Async Engine
The Go SDK's asyncengine subpackage and the C++ SDK's
openpit::asyncengine::TypedAsyncEngine turn an AccountSync engine into a
concurrent facade with per-account ordering guarantees. These helpers are one
possible integration pattern - applications may provide their own sharded
queues, actor model, or equivalent dispatcher.
Both SDKs provide sharded and dynamic strategies, futures, observer callbacks,
graceful stop, and hard stop. The Go examples below use goroutines and
context.Context; the C++ surface uses worker threads, Future<T>, and
std::chrono deadlines with the same account-pinning contract.
Use a bundled async engine when the application:
- builds the engine with
AccountSyncto allow concurrent access across many accounts; - wants the SDK to guarantee that no two operations for the same account ever run inside the engine concurrently;
- prefers
Future-style return values over manual channel/wait-group bookkeeping; - needs graceful and hard stop modes with deadlines.
If FullSync or NoSync is sufficient, or if you already have a
custom per-account dispatcher you are happy with, you do not need this
package.
asyncengine ships two dispatch strategies. Each chosen at build time so
the hot path stays branch-free in steady state.
- Go:
asyncengine.NewBuilder(engine).Sharded(workers) - C++:
openpit::asyncengine::MakeTypedAsyncEngine(engine, workers)orTypedBuilder<Driver>(driver).Sharded(workers)
The two C++ forms differ only in who owns the engine adapter. The typed builder
borrows a driver that the caller keeps alive, which is what a custom driver
needs. MakeTypedAsyncEngine is the shortcut for the common case of wrapping a
plain engine: it owns the default adapter itself and exposes the typed async
methods - start pre-trade, execute pre-trade, apply execution report, apply
account adjustment, and apply drop copy - without any driver lifetime plumbing.
The source engine must still outlive the async wrapper.
- A fixed number of worker queues created up front. Go runs one goroutine per shard; C++ runs one worker thread per shard.
- The account ID is hashed (Fibonacci mix) into one of the shards.
- Routing is one multiply-shift per submit. There is no per-account map lookup; synchronization is limited to ordering the submit against queue capacity and stop.
Strengths:
- The cheapest hot path: ~constant per submit.
- O(1) memory regardless of how many distinct accounts are active.
- Predictable worker count.
Trade-offs:
- A hot account saturates a single shard while the others stay idle.
- No per-account observer signals (queue created/removed never fire).
- Different accounts that hash to the same shard interleave through the same channel.
Pick Sharded when the active account set is broad and roughly balanced
or when raw throughput per submit matters more than per-account
isolation.
-
Go:
asyncengine.NewBuilder(engine).Dynamic().MaxQueues(n).IdleCleanupAfter(d) -
C++:
TypedBuilder<Driver>(driver).Dynamic().MaxQueues(n).IdleCleanupAfter(d) -
Per-account queue created on first submit.
-
A background cleanup worker retires queues that have been empty and untouched for
IdleCleanupAfter. -
MaxQueuescaps the live queue count; submits for unknown accounts returnErrQueueLimitin Go orErrorCode::QueueLimitin C++ past the cap.MaxQueues(0)removes the cap.
Strengths:
- Full per-account isolation: a slow account does not starve others that happen to share a shard.
- Per-account observer signals (
OnQueueCreated,OnQueueRemoved, per-account latency). - Memory scales with the active set rather than the total population.
Trade-offs:
- A synchronized map lookup on every submit.
- A periodic cleanup worker.
- Slightly higher per-account memory because each queue holds its own bounded buffer and worker.
Pick Dynamic when account activity is skewed, when you want per-account
metrics, or when the population is large enough that statically
allocating shards would be wasteful.
MaxQueues defaults to runtime.NumCPU() * 32, a value designed to be
effectively non-restrictive on typical hosts while still bounding
pathological growth (e.g. ephemeral per-request accounts). The limit counts
account queues only: MaxQueues(n) permits n usable account queues, and
MaxQueues(1) permits one. Override with MaxQueues(0) for unbounded growth.
Every queued operation returns a future (pkg/future in Go,
openpit::asyncengine::Future<T> in C++), resolved exactly once by a worker or
synchronously when submission fails before queueing.
Go operations that mirror one return value use future.Future[T]; tuple-shaped
operations use future.Future2[A, B]. C++ typed operations return
Future<StartOutcome>, Future<ExecuteOutcome>, or the concrete result type.
AsyncRequest::Execute uses PairFuture for its reservation-or-rejects pair.
Drop copy carries the same accepted-or-rejected pair as the main stage:
future.Future2[*AsyncDropCopyOperation, []reject.Reject] in Go and
Future<asyncengine::DropCopyOutcome<Driver>> in C++, whose outcome holds a
std::shared_ptr<AsyncDropCopyOperation<Driver>> and the rejects.
As with reservations, the accepted value is an async wrapper, not the
direct-engine handle: AsyncDropCopyOperation wraps
pretrade.DropCopyOperation in Go and pretrade::DropCopyOperation in C++ so
that finalization re-enters the producing account queue.
The wrapper also forwards the operation snapshots: Lock,
AccountAdjustments, AccountBlock, and IsAccountBlocked. These reads are
synchronous and serialize with finalization on the same wrapper; they do not
enqueue a separate engine call.
Go future operations:
-
f.Await(ctx)- block until resolved or ctx fires.Future[T]yields(T, error);Future2[A, B]yields(A, B, error). -
f.Done()- non-blocking check. -
f.TryGet()- non-blocking read. -
f.Wait()- channel that closes on resolution (forselect).
The future and the wrapped engine objects are decoupled: cancelling the
context passed to Await does not stop the underlying engine call; it
only stops the caller from waiting. In Go the future still owns any eventual
drop-copy operation after cancellation; await it again or use TryGet, then
commit or roll it back and close it. Abandoning the future is not a commit.
In C++, Await() returns the value or throws the carried async error on the
waiting thread; Await(timeout) returns std::nullopt when the timeout expires.
The underlying task continues in both cases.
AsyncEngine requires the wrapped engine to be built with AccountSync.
Per-account correctness holds in both strategies: no two operations for
the same account ever run inside the engine concurrently, and within one
account every queued operation runs in submit order.
Parallelism across different accounts depends on the strategy.
Dynamic gives full per-account isolation, so distinct accounts are
always processed in parallel. Sharded gives per-shard serialization:
distinct accounts that hash to the same shard share one worker and are
processed one after another (a hot account can block others on its
shard). Neither relaxes the per-account invariant above.
The Threading Contract of the engine itself is respected end-to-end.
AsyncRequest, AsyncReservation, and AsyncDropCopyOperation keep the same
per-account queue across the boundary. Executing a started request, or
committing, rolling back, or closing a reservation or drop-copy operation
through its wrapper, re-enters the same per-account chain.
That routing covers explicit calls only. Neither wrapper installs a destructor
hook: Go has no finalizer at all, and the C++ wrapper's implicit release runs
when the last shared_ptr owner goes away, on whatever thread drops it. So
releasing a reservation or drop-copy operation without an explicit
Close-flavored call performs its implicit rollback off the account queue.
Finish every wrapper with CommitAndClose, RollbackAndClose, or Close when
that lane matters.
Finalizing through a wrapper does not change what a mutation finalizer owes the
engine: it has no right to fail, and a failure arms the engine kill switch
instead of failing the void call. Because every mutation a Go or C++ policy
registers is a custom-policy mutation, that block covers every account, so the
next queued pre-trade operation for any account resolves with a
SystemUnavailable reject. Go can clear it in the lane with
asyncEngine.Accounts().UnblockAll(ctx); in C++ clear it on the wrapped
engine's own Accounts() handle. See
Account Blocking.
ApplyDropCopy uses the same account queue as every other operation for that
account. A readable account ID is required: an order that does not expose one
is refused up front, without being queued, and the future resolves immediately
with ErrMissingAccountID in Go or ErrorCode::MissingAccountId in C++ -
exactly as StartPreTrade and ExecutePreTrade already behave. The core would
reject such an order with MissingRequiredField before any policy callback
runs, so there is nothing to gain by queueing it. Existing account and
account-group blocks do not prevent a queued drop copy from reaching the
engine.
The C++ typed facade is non-copyable but move-constructible. Its dispatch state keeps a stable address, so already-issued request and reservation wrappers stay valid after a move. Move assignment is disabled because replacing a live target could invalidate wrappers created by that target. Keep the resulting engine alive until all wrappers are released.
- Go
StopGraceful(ctx)refuses new submissions and waits for every already-queued task to run to completion. Returnsctx.Err()if ctx fires before workers drain; the engine is then partially stopped andStopHardmay be invoked to complete the shutdown. - Go
StopHard(ctx)refuses new submissions, aborts every task that has not yet started withErrStopped, and waits for the currently-running task in each worker to finish. Returnsctx.Err()if ctx fires before the in-flight task in each worker finishes. - C++
StopGraceful(timeout)andStopHard(timeout)have the same queue behavior and returntruewhen workers drained before the deadline. A hard stop resolves queued tasks withErrorCode::Stopped.
When the builder was wired via AccountSyncReadyEngineBuilder.BuildAsync,
the underlying *Engine is released when the stop completes successfully
(returns nil). If StopGraceful or StopHard returns ctx.Err()
(timeout), the engine is not yet released; the caller must complete
the shutdown (call StopHard) before the release happens. When the
caller used asyncengine.NewBuilder(engine) and did not wire
WithStopUnderlying, the engine handle remains owned by the caller.
asyncengine.Observer in Go and openpit::asyncengine::Observer in C++ are
optional diagnostic interfaces. The default is NoopObserver. Wire only the
methods you need - the package itself
has zero external observability dependencies (no OpenTelemetry,
Prometheus, or logging in the import graph), so you decide where the
signals go.
Available callbacks: OnEnqueue, OnDequeue, OnComplete,
OnSlowSubmit, OnQueueFullBlocked, OnQueueCreated, OnQueueRemoved,
OnSubmitCancelled.
Observer callbacks are diagnostic only. In C++, exceptions thrown by an observer are ignored so observability cannot change queue or task semantics.
Every callback carries the account ID the task was routed by, and that ID is
always a real account: the dispatcher never queues work it cannot key and never
substitutes a sentinel. AccountID(0) in a callback therefore means account
zero, nothing else.
package main
import (
"context"
"log"
"time"
"go.openpit.dev/openpit"
"go.openpit.dev/openpit/model"
"go.openpit.dev/openpit/param"
"go.openpit.dev/openpit/pretrade/policies"
)
func main() {
// Build an AccountSync engine and wrap it into an async facade in one
// chain. BuildAsync is only available on the AccountSync builder; for
// FullSync or NoSync engines the bundled async helper is not used.
asyncBuilder, err := openpit.NewEngineBuilder().
AccountSync().
Builtin(policies.BuildOrderValidation()).
Builtin(
policies.BuildRateLimit().BrokerBarrier(
policies.RateLimitBrokerBarrier{
Limit: policies.RateLimit{
MaxOrders: 100,
Window: time.Second,
},
},
),
).
BuildAsync()
if err != nil {
log.Fatal(err)
}
// Pick a dispatch strategy. Use Sharded for cheap routing across a
// balanced account population; use Dynamic for per-account isolation
// and per-account metrics.
async, err := asyncBuilder.Dynamic().
MaxQueues(0).
IdleCleanupAfter(5 * time.Minute).
Build()
if err != nil {
log.Fatal(err)
}
defer func() {
if err := async.StopGraceful(context.Background()); err != nil {
log.Printf("StopGraceful: %v", err)
}
}()
// Submit a start-stage call. The future resolves once the worker has
// executed the call. AsyncRequest.Execute and Close are queued in the
// same per-account chain so AccountSync is never violated.
usd, err := param.NewAsset("USD")
if err != nil {
log.Fatal(err)
}
aapl, err := param.NewAsset("AAPL")
if err != nil {
log.Fatal(err)
}
order := model.NewOrder()
op := order.EnsureOperationView()
op.SetInstrument(param.NewInstrument(aapl, usd))
op.SetAccountID(param.NewAccountIDFromUint64(99224416))
op.SetSide(param.SideBuy)
price, err := param.NewPriceFromString("185")
if err != nil {
log.Fatal(err)
}
qty, err := param.NewQuantityFromString("100")
if err != nil {
log.Fatal(err)
}
op.SetTradeAmount(param.NewQuantityTradeAmount(qty))
op.SetPrice(price)
request, rejects, err := async.StartPreTrade(
context.Background(),
order,
).Await(context.Background())
if err != nil {
log.Fatal(err)
}
if request == nil {
// Rejected at the start stage; inspect rejects.
_ = rejects
return
}
// The async request preserves AccountSync across the Start - Execute
// boundary by routing Execute through the same per-account queue. The
// future yields the same (reservation, rejects, error) tuple the
// synchronous main stage returns.
reservation, rejects, err := request.Execute(
context.Background(),
).Await(context.Background())
if err != nil {
log.Fatal(err)
}
if reservation == nil {
_ = rejects
return
}
if _, err := reservation.CommitAndClose(
context.Background(),
).Await(context.Background()); err != nil {
log.Fatal(err)
}
}Go Submit(ctx, accountID, fn) and C++
Submit(accountId, std::function<void()>, timeout) enqueue caller-owned work
into the same per-account queue. Use this to run client-side work atomically
with respect to engine calls on the same account - for example, to persist an
order before Execute, or update a strategy book after Commit.
Splitting a logical transaction into two Submit calls "surfaces"
between them: tasks for the same account from elsewhere in the system
can interleave between the two halves. To keep that boundary closed,
bundle the work into a single Submit whose fn does both halves.
- Aborted tasks resolve their future with
ErrStoppedin Go orErrorCode::Stoppedin C++. For theAsyncReservationandAsyncDropCopyOperationmethods that "and close" (Close,CommitAndClose,RollbackAndClose) the underlying handle is still released as a safety net even on abort. An aborted plainCommitorRollbackdoes not release it, so the caller must still close the wrapper. - An aborted
AsyncRequest.Executereleases the underlying request before resolving the future. The caller does not need a separate cleanup path. - The Go
ctxor C++timeoutpassed to a submit method controls how long the producer waits for queue space. Once queued, the worker runs or aborts the task.
- Threading Contract: per-mode sync contract that AsyncEngine respects.
- Pre-trade Pipeline: the Request/Reservation lifecycle that AsyncRequest/AsyncReservation wraps.
- Getting Started: the synchronous flow that AsyncEngine layers on top of.
- Account Blocking: the engine-wide block and the mutation finalizer contract that queued finalization is subject to.