From 9e3000f1db17381aa1199dba3f759b92db829144 Mon Sep 17 00:00:00 2001 From: Rusty Eddy Date: Mon, 28 Sep 2026 21:34:27 -0700 Subject: [PATCH 1/2] feat: sim snapshot per-position marks and margin-aware buying power Implements #412 (ADR-066): - account.Snapshot gains per-position marks: PositionMark carries the listing key, price, and the time it was observed. It is validated against the open positions (it must name an open position, appear at most once, have a positive price, and be observed no later than the snapshot). Marks are optional per position; Marks() and Mark(key) expose them. - The simulator reports its marks with the clock time each was set (fill, Advance, or ObserveMark), and replaces its private positionKey with account.ListingKey. - sim.AccountConfig.InitialMarginRatio (optional, positive) configures the account's margin model. When set, MarginUsed = required margin on gross notional, MarginAvailable = equity - MarginUsed (negative when over the limit, with no forced action), and BuyingPower = MarginAvailable / ratio, floored at zero. All are computed with the shared internal/account/margin calculation. When nil, the legacy fields are unchanged. - AccountConfig.MarginModelInfo describes the choice ("none", or initial-margin-ratio v1 with ratio=) for the run manifest (#414). Closes #412 Co-Authored-By: Claude Opus 5.5 --- internal/account/mark.go | 70 +++++++ internal/account/mark_test.go | 100 +++++++++ internal/account/snapshot.go | 30 +++ internal/adapters/broker/sim/account.go | 125 ++++++++--- internal/adapters/broker/sim/advance.go | 20 +- internal/adapters/broker/sim/broker.go | 29 +-- internal/adapters/broker/sim/config.go | 56 ++++- internal/adapters/broker/sim/margin_test.go | 219 ++++++++++++++++++++ 8 files changed, 589 insertions(+), 60 deletions(-) create mode 100644 internal/account/mark.go create mode 100644 internal/account/mark_test.go create mode 100644 internal/adapters/broker/sim/margin_test.go diff --git a/internal/account/mark.go b/internal/account/mark.go new file mode 100644 index 0000000..d64538f --- /dev/null +++ b/internal/account/mark.go @@ -0,0 +1,70 @@ +package account + +import ( + "fmt" + "time" + + runtimeorder "github.com/rustyeddy/trader/internal/order" + "github.com/rustyeddy/trader/num" +) + +// PositionMark is the current valuation price of one open position's +// listing, and when it was observed (ADR-066). +// +// A mark is "as of the reporter's last observation of that listing" — +// for the simulator, the last fill or bar close it saw — not a live +// price. In a multi-instrument account, different listings' marks can +// have different AsOf times; each is the latest observation at or +// before the snapshot's own AsOf. +type PositionMark struct { + // Listing identifies the open position this mark values. + Listing ListingKey + // Price is the valuation price. It must be positive. + Price num.Price + // AsOf is when Price was observed. It must be set and must not be + // after the snapshot's AsOf. + AsOf time.Time +} + +// checkMarks validates marks against the snapshot's positions and +// returns them in positions order. Every mark must name an open +// position, at most once. A position may have no mark: marks are +// optional for a reporter that has none, and a consumer that needs one +// treats its absence as an error, never as zero or AvgPrice. +func checkMarks(positions []runtimeorder.Position, marks []PositionMark, asOf time.Time) ([]PositionMark, error) { + if len(marks) == 0 { + return nil, nil + } + open := make(map[ListingKey]struct{}, len(positions)) + for _, p := range positions { + open[KeyOf(p.Listing)] = struct{}{} + } + byKey := make(map[ListingKey]PositionMark, len(marks)) + for i, m := range marks { + if _, ok := open[m.Listing]; !ok { + return nil, fmt.Errorf("entry %d: no open position in listing %s/%s/%s", + i, m.Listing.InstrumentID, m.Listing.Provider, m.Listing.Venue) + } + if _, dup := byKey[m.Listing]; dup { + return nil, fmt.Errorf("entry %d: duplicate mark for listing %s/%s/%s", + i, m.Listing.InstrumentID, m.Listing.Provider, m.Listing.Venue) + } + if m.Price.IsZero() { + return nil, fmt.Errorf("entry %d: price must be positive", i) + } + if m.AsOf.IsZero() { + return nil, fmt.Errorf("entry %d: as-of time must be set", i) + } + if m.AsOf.After(asOf) { + return nil, fmt.Errorf("entry %d: as-of %s is after snapshot as-of %s", i, m.AsOf, asOf) + } + byKey[m.Listing] = m + } + ordered := make([]PositionMark, 0, len(byKey)) + for _, p := range positions { + if m, ok := byKey[KeyOf(p.Listing)]; ok { + ordered = append(ordered, m) + } + } + return ordered, nil +} diff --git a/internal/account/mark_test.go b/internal/account/mark_test.go new file mode 100644 index 0000000..0333dab --- /dev/null +++ b/internal/account/mark_test.go @@ -0,0 +1,100 @@ +package account + +import ( + "testing" + "time" + + runtimeorder "github.com/rustyeddy/trader/internal/order" + "github.com/rustyeddy/trader/num" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// twoPositionParams is validParams with open EUR/USD and GBP/USD +// positions, in that order. +func twoPositionParams(t *testing.T) SnapshotParams { + t.Helper() + p := validParams(t) + gbp := mustListing(t, "GBP", "USD", "OANDA", "GBP_USD") + p.Positions = append(p.Positions, mustPosition(t, p.AccountID, gbp)) + return p +} + +func markFor(p runtimeorder.Position, price string, at time.Time) PositionMark { + return PositionMark{Listing: KeyOf(p.Listing), Price: num.MustParsePrice(price), AsOf: at} +} + +func TestSnapshotMarks(t *testing.T) { + p := twoPositionParams(t) + earlier := p.AsOf.Add(-time.Hour) + eurMark := markFor(p.Positions[0], "1.1", p.AsOf) + gbpMark := markFor(p.Positions[1], "1.3", earlier) + // Supplied out of Positions order; returned in Positions order. + p.Marks = []PositionMark{gbpMark, eurMark} + + s, err := NewSnapshot(p) + require.NoError(t, err) + assert.Equal(t, []PositionMark{eurMark, gbpMark}, s.Marks()) + + got, ok := s.Mark(KeyOf(p.Positions[1].Listing)) + require.True(t, ok) + assert.Equal(t, gbpMark, got, "each mark keeps its own AsOf") + + // Marks returns a copy. + s.Marks()[0].Price = num.MustParsePrice("99") + assert.Equal(t, eurMark, s.Marks()[0]) +} + +func TestSnapshotMarksOptional(t *testing.T) { + p := twoPositionParams(t) + p.Marks = []PositionMark{markFor(p.Positions[1], "1.3", p.AsOf)} + s, err := NewSnapshot(p) + require.NoError(t, err) + _, ok := s.Mark(KeyOf(p.Positions[0].Listing)) + assert.False(t, ok, "a position may have no mark") + + s, err = NewSnapshot(validParams(t)) + require.NoError(t, err) + assert.Empty(t, s.Marks()) +} + +func TestSnapshotRejectsInvalidMarks(t *testing.T) { + other := mustListing(t, "AUD", "USD", "OANDA", "AUD_USD") + cases := []struct { + name string + mutate func(p *SnapshotParams) + }{ + {"no open position in listing", func(p *SnapshotParams) { + p.Marks = []PositionMark{{Listing: KeyOf(other), Price: num.MustParsePrice("1"), AsOf: p.AsOf}} + }}, + {"duplicate", func(p *SnapshotParams) { + m := markFor(p.Positions[0], "1", p.AsOf) + p.Marks = []PositionMark{m, m} + }}, + {"zero price", func(p *SnapshotParams) { + p.Marks = []PositionMark{{Listing: KeyOf(p.Positions[0].Listing), AsOf: p.AsOf}} + }}, + {"zero as-of", func(p *SnapshotParams) { + p.Marks = []PositionMark{{Listing: KeyOf(p.Positions[0].Listing), Price: num.MustParsePrice("1")}} + }}, + {"as-of after snapshot", func(p *SnapshotParams) { + p.Marks = []PositionMark{markFor(p.Positions[0], "1", p.AsOf.Add(time.Second))} + }}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + p := validParams(t) + tc.mutate(&p) + _, err := NewSnapshot(p) + assert.ErrorIs(t, err, ErrInvalidSnapshot) + }) + } +} + +func TestKeyOf(t *testing.T) { + l := mustEurUsdListing(t) + k := KeyOf(l) + assert.Equal(t, l.InstrumentID(), k.InstrumentID) + assert.Equal(t, l.Provider(), k.Provider) + assert.Equal(t, l.Venue(), k.Venue) +} diff --git a/internal/account/snapshot.go b/internal/account/snapshot.go index 227d457..cc3acca 100644 --- a/internal/account/snapshot.go +++ b/internal/account/snapshot.go @@ -57,6 +57,7 @@ type Snapshot struct { financing num.Money positions []runtimeorder.Position + marks []PositionMark openOrders []runtimeorder.Order } @@ -109,6 +110,11 @@ type SnapshotParams struct { // must case-insensitively equal Broker, and no two entries may name // the same (instrument, provider, venue) listing. Positions []runtimeorder.Position + // Marks is the current valuation price of open positions, at most + // one per position (see PositionMark). Every entry must name an + // entry of Positions by its ListingKey. It may omit positions, or be + // empty, when the reporter has no current price for them. + Marks []PositionMark // OpenOrders is this account's outstanding orders. Every entry's // Request.AccountID must equal AccountID, every entry's // Request.Listing.Provider must case-insensitively equal Broker, no @@ -162,6 +168,11 @@ func NewSnapshot(params SnapshotParams) (Snapshot, error) { return Snapshot{}, fmt.Errorf("%w: positions: %v", ErrInvalidSnapshot, err) } + marks, err := checkMarks(positions, params.Marks, params.AsOf) + if err != nil { + return Snapshot{}, fmt.Errorf("%w: marks: %v", ErrInvalidSnapshot, err) + } + openOrders, err := checkOpenOrders(params.AccountID, params.Broker, params.OpenOrders) if err != nil { return Snapshot{}, fmt.Errorf("%w: open orders: %v", ErrInvalidSnapshot, err) @@ -183,6 +194,7 @@ func NewSnapshot(params SnapshotParams) (Snapshot, error) { fees: params.Fees, financing: params.Financing, positions: positions, + marks: marks, openOrders: openOrders, }, nil } @@ -373,6 +385,24 @@ func (s Snapshot) Positions() []runtimeorder.Position { return cloned } +// Marks returns a copy of the open positions' current valuation +// prices, in Positions order. A position with no current price has no +// entry. +func (s Snapshot) Marks() []PositionMark { + return append([]PositionMark(nil), s.marks...) +} + +// Mark returns the current valuation price of the open position in +// key's listing, if the snapshot has one. +func (s Snapshot) Mark(key ListingKey) (PositionMark, bool) { + for _, m := range s.marks { + if m.Listing == key { + return m, true + } + } + return PositionMark{}, false +} + // OpenOrders returns a deep copy of the account's outstanding orders. // Mutating the returned slice, or any pointer/slice field it reaches, // does not affect s. diff --git a/internal/adapters/broker/sim/account.go b/internal/adapters/broker/sim/account.go index 59be574..6ef1099 100644 --- a/internal/adapters/broker/sim/account.go +++ b/internal/adapters/broker/sim/account.go @@ -10,6 +10,7 @@ import ( "github.com/rustyeddy/trader/instrument" "github.com/rustyeddy/trader/internal/account" + "github.com/rustyeddy/trader/internal/account/margin" brokerpkg "github.com/rustyeddy/trader/internal/broker" "github.com/rustyeddy/trader/internal/id" runtimeorder "github.com/rustyeddy/trader/internal/order" @@ -17,17 +18,11 @@ import ( "github.com/rustyeddy/trader/order" ) -// positionKey identifies one (instrument, provider, venue) listing for -// accountState.positions, matching the uniqueness rule -// account.NewSnapshot itself enforces on Snapshot.Positions. -type positionKey struct { - instrumentID instrument.ID - provider string - venue string -} - -func keyForListing(l instrument.Listing) positionKey { - return positionKey{instrumentID: l.InstrumentID(), provider: l.Provider(), venue: l.Venue()} +// mark is the last known price of one listing and the simulator time +// it was recorded — exposed on account.Snapshot as a PositionMark. +type mark struct { + price num.Price + at time.Time } // reducibleQuantity reports how much of position a ReduceOnly order @@ -97,7 +92,7 @@ type accountState struct { // is simply empty" convention for a flat account. Use // commitPosition, never a direct map write, so this invariant // cannot be violated by accident. - positions map[positionKey]runtimeorder.Position + positions map[account.ListingKey]runtimeorder.Position // marks holds the last known price per listing (issue #152, // M3-09): set from a market order's fill price (Submit), a // triggered limit/stop fill's price, or — even when no order @@ -107,8 +102,9 @@ type accountState struct { // known market observation," not live/real-time mark-to-market — // see snapshotLocked's doc comment. Entries are never deleted, even // after a position closes, so the last traded price remains - // available for history/display. - marks map[positionKey]num.Price + // available for history/display. Each mark records the simulator + // clock time it was set, which Snapshot reports as the mark's AsOf. + marks map[account.ListingKey]mark // realizedPnL is this account's cumulative realized profit and // loss, denominated in currency. It moves only when a fill reduces, // closes, or reverses a position (see position.go); opening or @@ -123,6 +119,10 @@ type accountState struct { // as account.Snapshot.Fees. fees num.Money + // margin is the account's initial-margin policy (ADR-066), or nil + // when none is configured — see AccountConfig.InitialMarginRatio. + marginPolicy *margin.Ratio + events []brokerpkg.Event nextSequence uint64 @@ -151,14 +151,18 @@ func zeroMoney(currency num.Currency) (num.Money, error) { // UnrealizedPnL is computed from s.marks against each open Position's // AvgPrice (see unrealizedPnLForPosition) — explicitly "as of the // simulator's last known market observation" (whatever last touched -// s.marks for that listing: a fill, or a Broker.Advance revaluation), -// not live/real-time mark-to-market; this package has no ongoing price -// feed to mark against between those events. Equity is s.cash plus -// that UnrealizedPnL. BuyingPower and MarginAvailable mirror s.cash -// directly and MarginUsed is always zero: this package still models an -// unleveraged, fully funded account with no margin policy of its own -// (that is M4's job) — these fields are a deliberate M3 placeholder, -// not a claim of real margin/leverage semantics. +// s.marks for that listing: a fill, or a Broker.Advance or ObserveMark +// revaluation), not live/real-time mark-to-market; this package has no +// ongoing price feed to mark against between those events. Equity is +// s.cash plus that UnrealizedPnL. The same marks are reported as +// Snapshot.Marks, each with the simulator time it was recorded +// (ADR-066), in Positions order. +// +// With no margin model configured (AccountConfig.InitialMarginRatio +// nil), BuyingPower and MarginAvailable mirror s.cash and MarginUsed is +// zero — the original M3 behavior, unchanged. With one configured, the +// three are derived from the shared margin calculation (ADR-066); see +// marginFieldsLocked. func (s *accountState) snapshotLocked() (account.Snapshot, error) { openOrders := make([]runtimeorder.Order, 0, len(s.orders)) for _, o := range s.orders { @@ -186,12 +190,14 @@ func (s *accountState) snapshotLocked() (account.Snapshot, error) { }) unrealizedPnL := s.zero - for key, p := range s.positions { - mark, ok := s.marks[key] + marks := make([]account.PositionMark, 0, len(positions)) + for _, p := range positions { + key := account.KeyOf(p.Listing) + m, ok := s.marks[key] if !ok { continue // a position always has a mark from its opening fill } - delta, err := unrealizedPnLForPosition(p, mark, p.Listing.Spec().SettlementCurrency()) + delta, err := unrealizedPnLForPosition(p, m.price, p.Listing.Spec().SettlementCurrency()) if err != nil { return account.Snapshot{}, err } @@ -199,6 +205,7 @@ func (s *accountState) snapshotLocked() (account.Snapshot, error) { if err != nil { return account.Snapshot{}, err } + marks = append(marks, account.PositionMark{Listing: key, Price: m.price, AsOf: m.at}) } equity, err := s.cash.Add(unrealizedPnL) @@ -206,6 +213,14 @@ func (s *accountState) snapshotLocked() (account.Snapshot, error) { return account.Snapshot{}, err } + buyingPower, marginUsed, marginAvailable := s.cash, s.zero, s.cash + if s.marginPolicy != nil { + buyingPower, marginUsed, marginAvailable, err = s.marginFieldsLocked(positions, equity) + if err != nil { + return account.Snapshot{}, err + } + } + return account.NewSnapshot(account.SnapshotParams{ AccountID: s.ref.AccountID, Broker: s.ref.Broker, @@ -213,24 +228,66 @@ func (s *accountState) snapshotLocked() (account.Snapshot, error) { AsOf: s.asOf, CashBalances: []num.Money{s.cash}, Equity: equity, - BuyingPower: s.cash, - MarginUsed: s.zero, - MarginAvailable: s.cash, + BuyingPower: buyingPower, + MarginUsed: marginUsed, + MarginAvailable: marginAvailable, RealizedPnL: s.realizedPnL, UnrealizedPnL: unrealizedPnL, Fees: s.fees, Financing: s.zero, Positions: positions, + Marks: marks, OpenOrders: openOrders, }) } +// marginFieldsLocked derives the snapshot's margin fields from the +// configured initial-margin ratio (ADR-066), using the shared margin +// calculation so the simulator and the risk rule agree on gross +// exposure: +// +// MarginUsed = required margin on current gross notional +// MarginAvailable = equity − MarginUsed (negative when over the limit) +// BuyingPower = MarginAvailable ÷ ratio, floored at zero +// +// Every open position is valued at its mark. MarginAvailable may go +// negative after an adverse move: v1 has no maintenance margin, so the +// simulator reports the breach and takes no action. BuyingPower is +// floored at zero because it is funds available to open new positions, +// and an over-limit account has none. The caller must hold s.mu. +func (s *accountState) marginFieldsLocked(positions []runtimeorder.Position, equity num.Money) (buyingPower, used, available num.Money, err error) { + marks := make(margin.Marks, len(s.marks)) + for key, m := range s.marks { + marks[key] = m.price + } + req, err := margin.Account(positions, marks, *s.marginPolicy, s.currency) + if err != nil { + return num.Money{}, num.Money{}, num.Money{}, fmt.Errorf("sim: computing margin: %w", err) + } + available, err = equity.Sub(req.Required) + if err != nil { + return num.Money{}, num.Money{}, num.Money{}, fmt.Errorf("sim: computing margin available: %w", err) + } + buyingPower, err = available.DivRate(s.marginPolicy.Value()) + if err != nil { + return num.Money{}, num.Money{}, num.Money{}, fmt.Errorf("sim: computing buying power: %w", err) + } + negative, err := buyingPower.Cmp(s.zero) + if err != nil { + return num.Money{}, num.Money{}, num.Money{}, fmt.Errorf("sim: comparing buying power: %w", err) + } + if negative < 0 { + buyingPower = s.zero + } + return buyingPower, req.Required, available, nil +} + // commitPosition stores pos as the account's current position for // key, or removes any stored entry when pos is Flat — the only way // s.positions is ever written, so its "only non-Flat entries" invariant // (see accountState's doc comment) cannot be violated by a direct map // write at a call site. The caller must already hold s.mu. -func (s *accountState) commitPosition(key positionKey, pos runtimeorder.Position) { +func (s *accountState) commitPosition(key account.ListingKey, pos runtimeorder.Position) { if pos.Side == order.Flat { delete(s.positions, key) return @@ -503,7 +560,7 @@ func (h *accountHandle) Submit(ctx context.Context, req runtimeorder.Request) (r return runtimeorder.Order{}, err } - h.state.commitFill(req.Listing, outcome) + h.state.commitFill(req.Listing, outcome, now) h.state.asOf = now h.state.commitEvents(append([]brokerpkg.Event{acceptEvent, outcome.fillEvent, outcome.filledEvent}, outcome.extraEvents...)...) return outcome.order, nil @@ -538,11 +595,11 @@ type fillOutcome struct { // commitPosition), records the new mark, and updates cash/realizedPnL/ // fees. The caller must already hold s.mu and must call commitEvents // separately (see Submit and accountState.advance). -func (s *accountState) commitFill(listing instrument.Listing, outcome fillOutcome) { +func (s *accountState) commitFill(listing instrument.Listing, outcome fillOutcome, at time.Time) { s.orders[outcome.order.Request.OrderID] = cloneOrder(outcome.order) - key := keyForListing(listing) + key := account.KeyOf(listing) s.commitPosition(key, outcome.position) - s.marks[key] = outcome.mark + s.marks[key] = mark{price: outcome.mark, at: at} s.cash = outcome.cash s.realizedPnL = outcome.realizedPnL s.fees = outcome.fees @@ -601,7 +658,7 @@ func roundFillPriceToTick(price num.Price, side order.Side, tick num.Price) (num func (s *accountState) buildFill(deps Deps, o runtimeorder.Order, price num.Price, causationID id.EventID, sequence uint64) (fillOutcome, error) { req := o.Request - key := keyForListing(req.Listing) + key := account.KeyOf(req.Listing) currency := req.Listing.Spec().SettlementCurrency() // Checked first, before any other part of the fill is built: diff --git a/internal/adapters/broker/sim/advance.go b/internal/adapters/broker/sim/advance.go index 919032f..6799e81 100644 --- a/internal/adapters/broker/sim/advance.go +++ b/internal/adapters/broker/sim/advance.go @@ -8,6 +8,7 @@ import ( "time" "github.com/rustyeddy/trader/instrument" + "github.com/rustyeddy/trader/internal/account" brokerpkg "github.com/rustyeddy/trader/internal/broker" "github.com/rustyeddy/trader/internal/id" runtimeorder "github.com/rustyeddy/trader/internal/order" @@ -129,16 +130,17 @@ func (s *accountState) observeMark(instrumentID instrument.ID, close num.Price, return } + now := deps.Clock.Now() changed := false for key := range s.positions { - if key.instrumentID != instrumentID { + if key.InstrumentID != instrumentID { continue } - s.marks[key] = close + s.marks[key] = mark{price: close, at: now} changed = true } if changed { - s.asOf = deps.Clock.Now() + s.asOf = now } } @@ -189,10 +191,11 @@ func (s *accountState) advance(ctx context.Context, deps Deps, obs Observation) // (issue #152, M3-09) must not go stale merely because there was // nothing to fill this bar. Any fill processed below overwrites // this with its own, more specific fill price. - key := keyForListing(obs.Listing) + key := account.KeyOf(obs.Listing) if _, hasPosition := s.positions[key]; hasPosition { - s.marks[key] = obs.Close - s.asOf = deps.Clock.Now() + now := deps.Clock.Now() + s.marks[key] = mark{price: obs.Close, at: now} + s.asOf = now } var atOpen, withinBar []triggeredOrder @@ -280,8 +283,9 @@ func (s *accountState) advance(ctx context.Context, deps Deps, obs Observation) continue } - s.commitFill(t.order.Request.Listing, outcome) - s.asOf = deps.Clock.Now() + now := deps.Clock.Now() + s.commitFill(t.order.Request.Listing, outcome, now) + s.asOf = now s.commitEvents(append([]brokerpkg.Event{outcome.fillEvent, outcome.filledEvent}, outcome.extraEvents...)...) } diff --git a/internal/adapters/broker/sim/broker.go b/internal/adapters/broker/sim/broker.go index 9e5a187..d67fc7c 100644 --- a/internal/adapters/broker/sim/broker.go +++ b/internal/adapters/broker/sim/broker.go @@ -9,7 +9,6 @@ import ( brokerpkg "github.com/rustyeddy/trader/internal/broker" "github.com/rustyeddy/trader/internal/id" runtimeorder "github.com/rustyeddy/trader/internal/order" - "github.com/rustyeddy/trader/num" ) // Broker is a deterministic, in-memory implementation of broker.Broker @@ -70,18 +69,24 @@ func NewBroker(name string, deps Deps, configs ...AccountConfig) (*Broker, error return nil, fmt.Errorf("account config %d: %w", i, err) } + policy, err := cfg.marginPolicy() + if err != nil { + return nil, fmt.Errorf("account config %d: %w", i, err) + } + accounts[cfg.AccountID] = &accountState{ - ref: ref, - currency: cfg.StartingCash.Currency(), - cash: cfg.StartingCash, - zero: zero, - asOf: now, - orders: make(map[id.OrderID]runtimeorder.Order), - positions: make(map[positionKey]runtimeorder.Position), - marks: make(map[positionKey]num.Price), - realizedPnL: zero, - fees: zero, - changed: make(chan struct{}), + ref: ref, + currency: cfg.StartingCash.Currency(), + cash: cfg.StartingCash, + zero: zero, + asOf: now, + orders: make(map[id.OrderID]runtimeorder.Order), + positions: make(map[account.ListingKey]runtimeorder.Position), + marks: make(map[account.ListingKey]mark), + realizedPnL: zero, + fees: zero, + marginPolicy: policy, + changed: make(chan struct{}), } } diff --git a/internal/adapters/broker/sim/config.go b/internal/adapters/broker/sim/config.go index 8c43857..682052a 100644 --- a/internal/adapters/broker/sim/config.go +++ b/internal/adapters/broker/sim/config.go @@ -4,6 +4,7 @@ import ( "fmt" "github.com/rustyeddy/trader/instrument" + "github.com/rustyeddy/trader/internal/account/margin" "github.com/rustyeddy/trader/internal/clock" "github.com/rustyeddy/trader/internal/id" "github.com/rustyeddy/trader/num" @@ -101,15 +102,42 @@ func (d Deps) validate() error { return nil } -// AccountConfig describes one simulated account's identity and -// deterministic starting capital. StartingCash's Currency becomes the -// account's home Currency; equity, buying power, and margin available -// all start equal to StartingCash, with margin used and PnL starting at -// zero — this package models no leverage or margin policy of its own -// (that is risk's concern, M4). +// AccountConfig describes one simulated account's identity, +// deterministic starting capital, and optional margin model. +// StartingCash's Currency becomes the account's home Currency. type AccountConfig struct { AccountID id.AccountID StartingCash num.Money + + // InitialMarginRatio configures the account's initial-margin model + // (ADR-066): the minimum equity required per unit of gross position + // notional — 1.0 is unlevered, 0.5 permits 2× gross exposure, 0.25 + // permits 4×. It must be positive when set. + // + // When nil, the account has no margin model: BuyingPower and + // MarginAvailable mirror cash and MarginUsed is zero, the original + // M3 behavior. When set, Snapshot derives all three from the shared + // margin calculation (see accountState.marginFieldsLocked). Either + // way, MarginModelInfo describes the choice for a run manifest. + InitialMarginRatio *num.Rate +} + +// marginModelName and marginModelVersion identify the initial-margin +// ratio model in MarginModelInfo. +const ( + marginModelName = "initial-margin-ratio" + marginModelVersion = "v1" +) + +// MarginModelInfo identifies c's margin model for reproducibility +// records (ADR-028, ADR-066): Name "none" when InitialMarginRatio is +// nil, otherwise the ratio model with its ratio in Config, so two runs +// with different ratios are distinguishable. +func (c AccountConfig) MarginModelInfo() ModelInfo { + if c.InitialMarginRatio == nil { + return ModelInfo{Name: "none"} + } + return ModelInfo{Name: marginModelName, Version: marginModelVersion, Config: "ratio=" + c.InitialMarginRatio.String()} } func (c AccountConfig) validate() error { @@ -119,5 +147,21 @@ func (c AccountConfig) validate() error { if !c.StartingCash.IsValid() { return fmt.Errorf("%w: starting cash must be valid money", ErrInvalidConfig) } + if _, err := c.marginPolicy(); err != nil { + return err + } return nil } + +// marginPolicy returns c's initial-margin policy, or nil when none is +// configured. +func (c AccountConfig) marginPolicy() (*margin.Ratio, error) { + if c.InitialMarginRatio == nil { + return nil, nil + } + r, err := margin.NewRatio(*c.InitialMarginRatio) + if err != nil { + return nil, fmt.Errorf("%w: %v", ErrInvalidConfig, err) + } + return &r, nil +} diff --git a/internal/adapters/broker/sim/margin_test.go b/internal/adapters/broker/sim/margin_test.go new file mode 100644 index 0000000..3ca96ec --- /dev/null +++ b/internal/adapters/broker/sim/margin_test.go @@ -0,0 +1,219 @@ +package sim + +import ( + "context" + "testing" + "time" + + "github.com/rustyeddy/trader/internal/account" + "github.com/rustyeddy/trader/internal/clock" + "github.com/rustyeddy/trader/internal/risk" + "github.com/rustyeddy/trader/num" + "github.com/rustyeddy/trader/order" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// marginAccount opens a $10,000 simulated account with the given +// initial-margin ratio ("" for no margin model). +func marginAccount(t *testing.T, ratio string) (*Broker, Deps, *accountHandle) { + t.Helper() + deps := testDeps() + accountID := mustAccountID(t, deps.IDs) + cfg := AccountConfig{AccountID: accountID, StartingCash: usd("10000")} + if ratio != "" { + r := num.MustParseRate(ratio) + cfg.InitialMarginRatio = &r + } + b, err := NewBroker("sim", deps, cfg) + require.NoError(t, err) + acc, err := b.OpenAccount(context.Background(), accountID) + require.NoError(t, err) + return b, deps, acc.(*accountHandle) +} + +func submitMarket(t *testing.T, deps Deps, h *accountHandle, symbol string, side order.Side, qty string) { + t.Helper() + listing := mustEurUsdListing(t) + if symbol == "GBP_USD" { + listing = mustGbpUsdListing(t) + } + _, err := h.Submit(context.Background(), mustMarketRequestFor(t, deps.IDs, h.Reference().AccountID, listing, side, qty)) + require.NoError(t, err) +} + +func snapshot(t *testing.T, h *accountHandle) account.Snapshot { + t.Helper() + s, err := h.Snapshot(context.Background()) + require.NoError(t, err) + return s +} + +// assertMargin checks equity, margin used, margin available, and +// buying power, in that order. +func assertMargin(t *testing.T, s account.Snapshot, equity, used, available, buyingPower string) { + t.Helper() + assert.True(t, s.Equity().Equal(usd(equity)), "equity: want %s, got %s", equity, s.Equity()) + assert.True(t, s.MarginUsed().Equal(usd(used)), "margin used: want %s, got %s", used, s.MarginUsed()) + assert.True(t, s.MarginAvailable().Equal(usd(available)), "margin available: want %s, got %s", available, s.MarginAvailable()) + assert.True(t, s.BuyingPower().Equal(usd(buyingPower)), "buying power: want %s, got %s", buyingPower, s.BuyingPower()) +} + +func TestAccountConfig_InitialMarginRatioValidation(t *testing.T) { + deps := testDeps() + for _, bad := range []string{"0", "-0.5"} { + r := num.MustParseRate(bad) + _, err := NewBroker("sim", deps, AccountConfig{AccountID: mustAccountID(t, deps.IDs), StartingCash: usd("1"), InitialMarginRatio: &r}) + assert.ErrorIs(t, err, ErrInvalidConfig, bad) + } +} + +func TestAccountConfig_MarginModelInfo(t *testing.T) { + assert.Equal(t, ModelInfo{Name: "none"}, AccountConfig{}.MarginModelInfo()) + + half := num.MustParseRate("0.5") + one := num.MustParseRate("1") + infoHalf := AccountConfig{InitialMarginRatio: &half}.MarginModelInfo() + assert.Equal(t, ModelInfo{Name: "initial-margin-ratio", Version: "v1", Config: "ratio=0.5"}, infoHalf) + assert.NotEqual(t, infoHalf, AccountConfig{InitialMarginRatio: &one}.MarginModelInfo(), "different ratios are distinguishable") +} + +func TestSnapshotMarginFields(t *testing.T) { + t.Run("no margin model keeps legacy fields", func(t *testing.T) { + _, deps, h := marginAccount(t, "") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "40000") // 44000 notional, 4.4× equity + // Legacy: buying power and margin available mirror cash. The + // simulator's cash moves only by realized PnL and fees, so it + // still reports the full 10000 — the gap behind #409. + assertMargin(t, snapshot(t, h), "10000", "0", "10000", "10000") + }) + t.Run("flat at 1.0", func(t *testing.T) { + _, _, h := marginAccount(t, "1") + assertMargin(t, snapshot(t, h), "10000", "0", "10000", "10000") + }) + t.Run("flat at 0.5", func(t *testing.T) { + _, _, h := marginAccount(t, "0.5") + assertMargin(t, snapshot(t, h), "10000", "0", "10000", "20000") + }) + t.Run("long at 1.0", func(t *testing.T) { + _, deps, h := marginAccount(t, "1") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "5000") // 5500 notional + assertMargin(t, snapshot(t, h), "10000", "5500", "4500", "4500") + }) + t.Run("long at 0.5", func(t *testing.T) { + _, deps, h := marginAccount(t, "0.5") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "5000") + assertMargin(t, snapshot(t, h), "10000", "2750", "7250", "14500") + }) + t.Run("short at 1.0", func(t *testing.T) { + _, deps, h := marginAccount(t, "1") + submitMarket(t, deps, h, "GBP_USD", order.Sell, "4000") // 5000 notional + assertMargin(t, snapshot(t, h), "10000", "5000", "5000", "5000") + }) + t.Run("over the limit at entry", func(t *testing.T) { + // #412 only reports; refusing the fill is #415. 44000 notional + // at 1.0 leaves margin available at -34000 and no buying power. + _, deps, h := marginAccount(t, "1") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "40000") + assertMargin(t, snapshot(t, h), "10000", "44000", "-34000", "0") + }) + t.Run("multi-instrument long and short is gross", func(t *testing.T) { + _, deps, h := marginAccount(t, "1") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "5000") // 5500 + submitMarket(t, deps, h, "GBP_USD", order.Sell, "2000") // 2500 + assertMargin(t, snapshot(t, h), "10000", "8000", "2000", "2000") + }) + t.Run("multi-instrument at 0.5", func(t *testing.T) { + _, deps, h := marginAccount(t, "0.5") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "5000") // 5500 + submitMarket(t, deps, h, "GBP_USD", order.Sell, "2000") // 2500 + assertMargin(t, snapshot(t, h), "10000", "4000", "6000", "12000") + }) +} + +func TestSnapshotOverLimitAfterAdverseMoveTakesNoAction(t *testing.T) { + b, deps, h := marginAccount(t, "0.5") + // 15000 EUR at 1.1 = 16500 notional, 8250 margin: within 2× of 10000. + submitMarket(t, deps, h, "EUR_USD", order.Buy, "15000") + assertMargin(t, snapshot(t, h), "10000", "8250", "1750", "3500") + + // Close at 0.6: equity 10000 − 7500 = 2500, gross 9000, margin 4500. + require.NoError(t, deps.Clock.(*clock.Simulated).Advance(time.Hour)) + require.NoError(t, b.Advance(context.Background(), mustObservation(t, mustEurUsdListing(t), "0.60000", "0.60000", "0.60000", "0.60000", barTime))) + + s := snapshot(t, h) + assertMargin(t, s, "2500", "4500", "-2000", "0") + require.Len(t, s.Positions(), 1, "no forced liquidation in v1") + assert.True(t, s.Positions()[0].Quantity.Equal(num.MustParseQuantity("15000"))) + assert.Empty(t, s.OpenOrders()) +} + +func TestSnapshotMarks(t *testing.T) { + b, deps, h := marginAccount(t, "") + eur, gbp := mustEurUsdListing(t), mustGbpUsdListing(t) + + submitMarket(t, deps, h, "GBP_USD", order.Buy, "100") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "100") + + t.Run("fill sets the mark at the fill price and time", func(t *testing.T) { + s := snapshot(t, h) + m, ok := s.Mark(account.KeyOf(eur)) + require.True(t, ok) + assert.True(t, m.Price.Equal(num.MustParsePrice("1.10000"))) + assert.True(t, m.AsOf.Equal(testStart)) + }) + t.Run("marks follow positions order and are deterministic", func(t *testing.T) { + s := snapshot(t, h) + marks := s.Marks() + require.Len(t, marks, 2) + for i, p := range s.Positions() { + assert.Equal(t, account.KeyOf(p.Listing), marks[i].Listing) + } + assert.Equal(t, marks, snapshot(t, h).Marks()) + }) + + sim := deps.Clock.(*clock.Simulated) + t.Run("Advance updates only the observed listing", func(t *testing.T) { + require.NoError(t, sim.Advance(time.Hour)) + require.NoError(t, b.Advance(context.Background(), mustObservation(t, eur, "1.11000", "1.12000", "1.10000", "1.11500", barTime))) + s := snapshot(t, h) + m, _ := s.Mark(account.KeyOf(eur)) + assert.True(t, m.Price.Equal(num.MustParsePrice("1.11500"))) + assert.True(t, m.AsOf.Equal(testStart.Add(time.Hour))) + g, _ := s.Mark(account.KeyOf(gbp)) + assert.True(t, g.AsOf.Equal(testStart), "GBP/USD was not observed, so its mark is older") + }) + t.Run("ObserveMark updates the mark and its time", func(t *testing.T) { + require.NoError(t, sim.Advance(time.Hour)) + require.NoError(t, h.ObserveMark(context.Background(), gbp.InstrumentID(), num.MustParsePrice("1.30000"), testStart.Add(2*time.Hour))) + g, _ := snapshot(t, h).Mark(account.KeyOf(gbp)) + assert.True(t, g.Price.Equal(num.MustParsePrice("1.30000"))) + assert.True(t, g.AsOf.Equal(testStart.Add(2*time.Hour))) + }) +} + +// TestFullNotionalSizerRespectsMarginRatio verifies that the +// full-notional sizer (ADR-061), which sizes from min(Equity, +// BuyingPower), respects the configured ratio through the simulator's +// margin-aware BuyingPower. +func TestFullNotionalSizerRespectsMarginRatio(t *testing.T) { + ref := num.MustParsePrice("1.10000") + size := func(t *testing.T, ratio string) num.Quantity { + t.Helper() + _, _, h := marginAccount(t, ratio) + q, err := risk.NewFullNotionalSizer().Size(context.Background(), risk.SizeInput{ + Account: snapshot(t, h), + Listing: mustEurUsdListing(t), + ReferencePrice: &ref, + }) + require.NoError(t, err) + return q + } + // Ratio 1.0: buying power = equity = 10000 → 9090 units (9999 USD). + assert.True(t, size(t, "1").Equal(num.MustParseQuantity("9090"))) + // Ratio 0.5 permits more buying power, but the sizer still caps at + // equity: fully invested, unlevered. + assert.True(t, size(t, "0.5").Equal(num.MustParseQuantity("9090"))) + // Ratio 2.0 halves buying power to 5000 → 4545 units. + assert.True(t, size(t, "2").Equal(num.MustParseQuantity("4545"))) +} From c67b0c1b318f8d38179fe7298c4ee5a950b22846 Mon Sep 17 00:00:00 2001 From: Rusty Eddy Date: Mon, 28 Sep 2026 22:04:35 -0700 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20address=20sim=20marks=20review=20?= =?UTF-8?q?=E2=80=94=20zero=20prices,=20JSONL=20marks,=20current-only=20ma?= =?UTF-8?q?rgin=20marks?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Reject a zero final fill price in buildFill, after tick rounding and slippage and before anything is committed, so a zero price from Deps.Prices or a SlippageModel can't poison later snapshots. A regression test shows the account is unchanged and the same request fills later. - Reject a zero ObserveMark price, and any zero Observation price (Low must be positive), with ErrInvalidObservation before mutation. Covers ObserveMark, Advance, and AdvanceBar. - The JSONL journal now round-trips account marks (markWire). Each mark's listing is resolved against the entry's own positions, and a mark naming no position is ErrCorruptEntry. - marginFieldsLocked builds its mark map from open positions only, so cost doesn't grow with every listing ever traded. - Keep clock-derived mark times (sim.Deps.Clock's documented rule), and pin the Scheduler invariant that the clock equals the bar time when marks are observed. Refs #412 Co-Authored-By: Claude Opus 5.5 --- internal/adapters/broker/sim/account.go | 26 ++++++- internal/adapters/broker/sim/advance.go | 6 ++ internal/adapters/broker/sim/margin_test.go | 60 +++++++++++++++ internal/adapters/broker/sim/observation.go | 9 ++- .../journal/jsonl/marks_internal_test.go | 18 +++++ internal/adapters/journal/jsonl/marks_test.go | 73 +++++++++++++++++++ internal/adapters/journal/jsonl/unwire.go | 30 ++++++++ internal/adapters/journal/jsonl/wire.go | 26 +++++++ internal/backtest/scheduler_test.go | 46 ++++++++++++ 9 files changed, 290 insertions(+), 4 deletions(-) create mode 100644 internal/adapters/journal/jsonl/marks_internal_test.go create mode 100644 internal/adapters/journal/jsonl/marks_test.go diff --git a/internal/adapters/broker/sim/account.go b/internal/adapters/broker/sim/account.go index 6ef1099..d911f71 100644 --- a/internal/adapters/broker/sim/account.go +++ b/internal/adapters/broker/sim/account.go @@ -20,6 +20,12 @@ import ( // mark is the last known price of one listing and the simulator time // it was recorded — exposed on account.Snapshot as a PositionMark. +// +// at comes from Deps.Clock, like every other timestamp this package +// produces, not from an Observation's own time (see Deps.Clock). The +// backtest Scheduler advances the clock to each bar's time before +// observing it, so a bar-derived mark's time equals the bar's time +// (pinned by backtest's TestScheduler_ClockEqualsBarTimeWhenObservingMarks). type mark struct { price num.Price at time.Time @@ -256,9 +262,14 @@ func (s *accountState) snapshotLocked() (account.Snapshot, error) { // floored at zero because it is funds available to open new positions, // and an over-limit account has none. The caller must hold s.mu. func (s *accountState) marginFieldsLocked(positions []runtimeorder.Position, equity num.Money) (buyingPower, used, available num.Money, err error) { - marks := make(margin.Marks, len(s.marks)) - for key, m := range s.marks { - marks[key] = m.price + // Only open positions' marks: s.marks also keeps closed listings' + // last prices for history, which are not inputs to current margin. + marks := make(margin.Marks, len(positions)) + for _, p := range positions { + key := account.KeyOf(p.Listing) + if m, ok := s.marks[key]; ok { + marks[key] = m.price + } } req, err := margin.Account(positions, marks, *s.marginPolicy, s.currency) if err != nil { @@ -766,6 +777,15 @@ func (s *accountState) buildFill(deps Deps, o runtimeorder.Order, price num.Pric price = adjusted } + // The final execution price becomes the position's mark, and + // account.Snapshot rejects a non-positive mark (ADR-066). Checked + // here, while nothing has been committed, so a zero price from + // Deps.Prices or a SlippageModel fails the fill and leaves the + // account unchanged instead of poisoning every later Snapshot. + if price.IsZero() { + return fillOutcome{}, fmt.Errorf("%w: fill price must be positive", runtimeorder.ErrInvalidFill) + } + var commission *num.Money if deps.Commission != nil { c, err := deps.Commission.Commission(req.Listing, req.Side, fillQty, price) diff --git a/internal/adapters/broker/sim/advance.go b/internal/adapters/broker/sim/advance.go index 6799e81..9575a5a 100644 --- a/internal/adapters/broker/sim/advance.go +++ b/internal/adapters/broker/sim/advance.go @@ -90,6 +90,12 @@ func (h *accountHandle) ObserveMark(ctx context.Context, instrumentID instrument if err := ctx.Err(); err != nil { return err } + if close.IsZero() { + // Rejected before any state changes: account.Snapshot rejects + // a non-positive mark (ADR-066), so storing one would break + // every later Snapshot. + return fmt.Errorf("%w: mark price must be positive", ErrInvalidObservation) + } h.state.observeMark(instrumentID, close, h.broker.deps) return nil } diff --git a/internal/adapters/broker/sim/margin_test.go b/internal/adapters/broker/sim/margin_test.go index 3ca96ec..111bfc4 100644 --- a/internal/adapters/broker/sim/margin_test.go +++ b/internal/adapters/broker/sim/margin_test.go @@ -7,6 +7,7 @@ import ( "github.com/rustyeddy/trader/internal/account" "github.com/rustyeddy/trader/internal/clock" + runtimeorder "github.com/rustyeddy/trader/internal/order" "github.com/rustyeddy/trader/internal/risk" "github.com/rustyeddy/trader/num" "github.com/rustyeddy/trader/order" @@ -217,3 +218,62 @@ func TestFullNotionalSizerRespectsMarginRatio(t *testing.T) { // Ratio 2.0 halves buying power to 5000 → 4545 units. assert.True(t, size(t, "2").Equal(num.MustParseQuantity("4545"))) } + +// TestZeroFillPriceLeavesAccountUnchanged: a zero fill price would +// become a mark that account.Snapshot rejects, so the fill must fail +// before any state is committed. +func TestZeroFillPriceLeavesAccountUnchanged(t *testing.T) { + ctx := context.Background() + deps := testDeps() + prices := &mutablePriceSource{prices: map[string]num.Price{"EUR_USD": num.MustParsePrice("0")}} + deps.Prices = prices + accountID := mustAccountID(t, deps.IDs) + b, err := NewBroker("sim", deps, AccountConfig{AccountID: accountID, StartingCash: usd("10000")}) + require.NoError(t, err) + acc, err := b.OpenAccount(ctx, accountID) + require.NoError(t, err) + h := acc.(*accountHandle) + before := snapshot(t, h) + + req := mustMarketRequest(t, deps.IDs, accountID, order.Buy, "100") + _, err = acc.Submit(ctx, req) + require.ErrorIs(t, err, runtimeorder.ErrInvalidFill) + + after := snapshot(t, h) + assert.Empty(t, after.Positions()) + assert.Empty(t, after.Marks()) + assert.Empty(t, after.OpenOrders()) + assert.True(t, before.AsOf().Equal(after.AsOf())) + assert.Empty(t, h.state.orders, "the order was never stored") + + // The same request can still fill once a valid price exists. + prices.set("EUR_USD", num.MustParsePrice("1.10000")) + _, err = acc.Submit(ctx, req) + require.NoError(t, err) + require.Len(t, snapshot(t, h).Positions(), 1) +} + +func TestZeroObservationPricesRejectedBeforeMutation(t *testing.T) { + ctx := context.Background() + b, deps, h := marginAccount(t, "") + submitMarket(t, deps, h, "EUR_USD", order.Buy, "100") + before := snapshot(t, h) + eur := mustEurUsdListing(t) + + t.Run("ObserveMark", func(t *testing.T) { + err := h.ObserveMark(ctx, eur.InstrumentID(), num.MustParsePrice("0"), testStart) + assert.ErrorIs(t, err, ErrInvalidObservation) + }) + t.Run("Advance", func(t *testing.T) { + obs := Observation{Listing: eur, Open: num.MustParsePrice("1"), High: num.MustParsePrice("1"), Low: num.MustParsePrice("0"), Close: num.MustParsePrice("0"), Time: barTime} + assert.ErrorIs(t, b.Advance(ctx, obs), ErrInvalidObservation) + }) + t.Run("AdvanceBar", func(t *testing.T) { + zero, one := num.MustParsePrice("0"), num.MustParsePrice("1") + assert.ErrorIs(t, h.AdvanceBar(ctx, eur, one, one, zero, zero, barTime), ErrInvalidObservation) + }) + + after := snapshot(t, h) + assert.Equal(t, before.Marks(), after.Marks(), "marks unchanged") + assert.True(t, before.AsOf().Equal(after.AsOf()), "as-of unchanged") +} diff --git a/internal/adapters/broker/sim/observation.go b/internal/adapters/broker/sim/observation.go index a3f7293..10b5d6f 100644 --- a/internal/adapters/broker/sim/observation.go +++ b/internal/adapters/broker/sim/observation.go @@ -33,7 +33,8 @@ type Observation struct { } // validate reports whether o is well-formed: Listing must be -// constructed, Time must be non-zero, and Low <= Open, Close <= High. +// constructed, Time must be non-zero, prices must be positive, and +// Low <= Open, Close <= High. func (o Observation) validate() error { if o.Listing.InstrumentID().IsZero() { return fmt.Errorf("%w: listing must be constructed", ErrInvalidObservation) @@ -41,6 +42,12 @@ func (o Observation) validate() error { if o.Time.IsZero() { return fmt.Errorf("%w: time must be set", ErrInvalidObservation) } + // Low is the bar's smallest price, so a positive Low makes every + // price positive. Close becomes a position's mark, which + // account.Snapshot requires to be positive (ADR-066). + if o.Low.IsZero() { + return fmt.Errorf("%w: prices must be positive", ErrInvalidObservation) + } if o.High.Cmp(o.Low) < 0 { return fmt.Errorf("%w: high %s is below low %s", ErrInvalidObservation, o.High, o.Low) } diff --git a/internal/adapters/journal/jsonl/marks_internal_test.go b/internal/adapters/journal/jsonl/marks_internal_test.go new file mode 100644 index 0000000..5e82bfa --- /dev/null +++ b/internal/adapters/journal/jsonl/marks_internal_test.go @@ -0,0 +1,18 @@ +package jsonl + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestFromMarkWiresRejectsMarkWithNoPosition(t *testing.T) { + _, err := fromMarkWires([]markWire{{InstrumentID: "fx:EUR/USD", Provider: "sim"}}, nil) + assert.ErrorIs(t, err, ErrCorruptEntry) +} + +func TestFromMarkWiresEmpty(t *testing.T) { + marks, err := fromMarkWires(nil, nil) + assert.NoError(t, err) + assert.Nil(t, marks) +} diff --git a/internal/adapters/journal/jsonl/marks_test.go b/internal/adapters/journal/jsonl/marks_test.go new file mode 100644 index 0000000..ab00bf6 --- /dev/null +++ b/internal/adapters/journal/jsonl/marks_test.go @@ -0,0 +1,73 @@ +package jsonl_test + +import ( + "context" + "testing" + "time" + + "github.com/rustyeddy/trader/instrument" + "github.com/rustyeddy/trader/internal/account" + "github.com/rustyeddy/trader/internal/id" + "github.com/rustyeddy/trader/internal/journal" + runtimeorder "github.com/rustyeddy/trader/internal/order" + "github.com/rustyeddy/trader/num" + "github.com/rustyeddy/trader/order" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func mustGbpUsdVenueListing(t *testing.T) instrument.Listing { + t.Helper() + inst, err := instrument.NewCurrencyPair(num.MustParseCurrency("GBP"), num.MustParseCurrency("USD")) + require.NoError(t, err) + spec, err := instrument.NewSpec(num.MustParsePrice("0.00001"), num.MustParseQuantity("1"), num.MustParseRate("1"), num.MustParseCurrency("USD")) + require.NoError(t, err) + l, err := instrument.NewListing(instrument.ListingParams{Instrument: inst, Provider: "sim", Venue: "ECN", Symbol: "GBP_USD", Spec: spec, Tradable: true}) + require.NoError(t, err) + return l +} + +// TestWriterReaderRoundTripsAccountMarks proves a snapshot's +// per-position marks, each with its own AsOf, survive a JSONL write and +// read exactly (ADR-066). +func TestWriterReaderRoundTripsAccountMarks(t *testing.T) { + usd := num.MustParseCurrency("USD") + accountID := mustAccountID(t) + asOf := time.Date(2026, 3, 2, 21, 0, 0, 0, time.UTC) + listings := []instrument.Listing{mustEurUsdListing(t), mustGbpUsdVenueListing(t)} + + var positions []runtimeorder.Position + for _, l := range listings { + avg := num.MustParsePrice("1.10000") + p, err := runtimeorder.NewPosition(runtimeorder.Position{AccountID: accountID, Listing: l, Side: order.Long, Quantity: num.MustParseQuantity("100"), AvgPrice: &avg}) + require.NoError(t, err) + positions = append(positions, p) + } + marks := []account.PositionMark{ + {Listing: account.KeyOf(listings[0]), Price: num.MustParsePrice("1.12345"), AsOf: asOf}, + {Listing: account.KeyOf(listings[1]), Price: num.MustParsePrice("1.25000"), AsOf: asOf.Add(-time.Hour)}, + } + zero := num.MustParseMoney("0", usd) + snap, err := account.NewSnapshot(account.SnapshotParams{ + AccountID: accountID, Broker: "sim", Currency: usd, AsOf: asOf, + CashBalances: []num.Money{num.MustParseMoney("10000", usd)}, + Equity: num.MustParseMoney("10000", usd), BuyingPower: zero, MarginUsed: zero, MarginAvailable: zero, + RealizedPnL: zero, UnrealizedPnL: zero, Fees: zero, Financing: zero, + Positions: positions, Marks: marks, + }) + require.NoError(t, err) + + w, path := mustWriter(t) + require.NoError(t, w.Record(context.Background(), journal.Record{RunID: mustRunID(t), Metadata: id.Metadata{Timestamp: asOf}, Kind: journal.KindAccount, Account: &snap})) + require.NoError(t, w.Close()) + + entries := readAll(t, path) + require.Len(t, entries, 1) + got := entries[0].Account.Marks() + require.Len(t, got, 2) + for i, want := range snap.Marks() { + assert.Equal(t, want.Listing, got[i].Listing) + assert.True(t, want.Price.Equal(got[i].Price)) + assert.True(t, want.AsOf.Equal(got[i].AsOf), "mark %d keeps its own AsOf", i) + } +} diff --git a/internal/adapters/journal/jsonl/unwire.go b/internal/adapters/journal/jsonl/unwire.go index 59f33ff..92abc4f 100644 --- a/internal/adapters/journal/jsonl/unwire.go +++ b/internal/adapters/journal/jsonl/unwire.go @@ -255,6 +255,10 @@ func fromAccountWire(w accountWire) (account.Snapshot, error) { } positions = append(positions, p) } + marks, err := fromMarkWires(w.Marks, positions) + if err != nil { + return account.Snapshot{}, err + } openOrders := make([]runtimeorder.Order, 0, len(w.OpenOrders)) for _, ow := range w.OpenOrders { o, err := fromOrderWire(ow) @@ -279,6 +283,7 @@ func fromAccountWire(w accountWire) (account.Snapshot, error) { Fees: w.Fees, Financing: w.Financing, Positions: positions, + Marks: marks, OpenOrders: openOrders, }) if err != nil { @@ -287,6 +292,31 @@ func fromAccountWire(w accountWire) (account.Snapshot, error) { return snap, nil } +// fromMarkWires resolves each mark's listing against positions, the +// listings the same snapshot entry already carries. A mark naming no +// decoded position is corrupt. +func fromMarkWires(ws []markWire, positions []runtimeorder.Position) ([]account.PositionMark, error) { + if len(ws) == 0 { + return nil, nil + } + type wireKey struct{ instrumentID, provider, venue string } + keys := make(map[wireKey]account.ListingKey, len(positions)) + for _, p := range positions { + k := account.KeyOf(p.Listing) + keys[wireKey{k.InstrumentID.String(), k.Provider, k.Venue}] = k + } + marks := make([]account.PositionMark, 0, len(ws)) + for i, mw := range ws { + key, ok := keys[wireKey{mw.InstrumentID, mw.Provider, mw.Venue}] + if !ok { + return nil, fmt.Errorf("%w: mark %d names no position in listing %s/%s/%s", + ErrCorruptEntry, i, mw.InstrumentID, mw.Provider, mw.Venue) + } + marks = append(marks, account.PositionMark{Listing: key, Price: mw.Price, AsOf: mw.AsOf}) + } + return marks, nil +} + func fromStatusWire(w statusWire) (broker.Status, error) { state, err := parseAccountStatus(w.State) if err != nil { diff --git a/internal/adapters/journal/jsonl/wire.go b/internal/adapters/journal/jsonl/wire.go index 70e9d89..53c295e 100644 --- a/internal/adapters/journal/jsonl/wire.go +++ b/internal/adapters/journal/jsonl/wire.go @@ -282,9 +282,32 @@ type accountWire struct { Fees num.Money `json:"fees"` Financing num.Money `json:"financing"` Positions []positionWire `json:"positions,omitempty"` + Marks []markWire `json:"marks,omitempty"` OpenOrders []orderWire `json:"open_orders,omitempty"` } +// markWire is one account.PositionMark. Its listing is identified by +// the same instrument/provider/venue triple as a positionWire's +// listing, and is resolved against the snapshot's own positions on +// decode (fromAccountWire). +type markWire struct { + InstrumentID string `json:"instrument_id"` + Provider string `json:"provider"` + Venue string `json:"venue,omitempty"` + Price num.Price `json:"price"` + AsOf time.Time `json:"as_of"` +} + +func toMarkWire(m account.PositionMark) markWire { + return markWire{ + InstrumentID: m.Listing.InstrumentID.String(), + Provider: m.Listing.Provider, + Venue: m.Listing.Venue, + Price: m.Price, + AsOf: m.AsOf, + } +} + func toAccountWire(s account.Snapshot) accountWire { w := accountWire{ AccountID: s.AccountID(), @@ -305,6 +328,9 @@ func toAccountWire(s account.Snapshot) accountWire { for _, p := range s.Positions() { w.Positions = append(w.Positions, toPositionWire(p)) } + for _, m := range s.Marks() { + w.Marks = append(w.Marks, toMarkWire(m)) + } for _, o := range s.OpenOrders() { w.OpenOrders = append(w.OpenOrders, toOrderWire(o)) } diff --git a/internal/backtest/scheduler_test.go b/internal/backtest/scheduler_test.go index b5dc6f4..b15b435 100644 --- a/internal/backtest/scheduler_test.go +++ b/internal/backtest/scheduler_test.go @@ -1715,3 +1715,49 @@ func TestScheduler_BracketStopLegRejectedAbortsRunAfterJournalingEntryFill(t *te assert.Equal(t, 1, intentKinds[order.IntentAdjustStop], "the stop leg's own sub-intent") assert.True(t, foundRejectedStopDecision, "the stop leg's own risk rejection must be journaled") } + +// clockCheckingObserver records, for every ObserveMark/AdvanceBar call, +// the simulator clock alongside the bar time the Scheduler supplied. +type clockCheckingObserver struct { + clk interface{ Now() time.Time } + observer backtest.MarketObserver + advancer backtest.IntrabarAdvancer + calls []struct{ now, at time.Time } +} + +func (c *clockCheckingObserver) ObserveMark(ctx context.Context, instrumentID instrument.ID, close num.Price, at time.Time) error { + c.calls = append(c.calls, struct{ now, at time.Time }{c.clk.Now(), at}) + return c.observer.ObserveMark(ctx, instrumentID, close, at) +} + +func (c *clockCheckingObserver) AdvanceBar(ctx context.Context, listing instrument.Listing, open, high, low, close num.Price, at time.Time) error { + c.calls = append(c.calls, struct{ now, at time.Time }{c.clk.Now(), at}) + return c.advancer.AdvanceBar(ctx, listing, open, high, low, close, at) +} + +// TestScheduler_ClockEqualsBarTimeWhenObservingMarks pins the invariant +// the simulator relies on when it stamps a mark's AsOf from its own +// clock rather than the supplied bar time (sim.Deps.Clock, ADR-026, +// ADR-066): the Scheduler advances the clock to each bar's time before +// calling ObserveMark or AdvanceBar with it. +func TestScheduler_ClockEqualsBarTimeWhenObservingMarks(t *testing.T) { + mgr := newSchedulerTestManager(t) + replay := newTwoInstrumentReplay(t, mgr) + t.Cleanup(func() { _ = replay.Close() }) + + h := newSchedulerHarness(t, schedulerSpan(t).Start()) + strat := &intrabarStopTriggerStrategy{requirements: bothInstrumentsRequirements(t), instID: eurusdID(t), stopPrice: "1.10065"} + deps := newSchedulerDeps(t, replay, strat, h) + rec := &clockCheckingObserver{clk: h.clockObj, observer: deps.MarketObserver, advancer: deps.IntrabarAdvancer} + deps.MarketObserver = rec + deps.IntrabarAdvancer = rec + + sched, err := backtest.NewScheduler(deps) + require.NoError(t, err) + require.NoError(t, sched.Run(context.Background())) + + require.NotEmpty(t, rec.calls) + for i, c := range rec.calls { + assert.True(t, c.now.Equal(c.at), "call %d: clock %s, bar time %s", i, c.now, c.at) + } +}