diff --git a/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs b/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs index 23cc01e..3177095 100644 --- a/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs +++ b/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs @@ -134,21 +134,20 @@ public void OnTradingPhaseChanged(in TradingPhaseChangedEvent e) } /// - /// Issue #322: B3-aligned best-effort wire markers on - /// securityTradingEvent for halt and resume. The downstream - /// consumer just needs them to be distinct and non-NULL; the - /// SecurityStatus_3 frame's securityTradingStatus still - /// carries the engine's preserved so the - /// post-resume phase is unambiguous. + /// Issue #583: encodes administrative halt/resume using only official + /// UMDF 2.2.0 SecurityStatus_3 semantics. On halt, + /// securityTradingStatus reports the official FORBIDDEN status + /// (not the pre-halt phase); on resume, the engine's preserved + /// is restored as securityTradingStatus. + /// securityTradingEvent uses the official + /// SECURITY_STATUS_CHANGE / SECURITY_REJOINS_SECURITY_GROUP_STATUS + /// values to mark the divergence/rejoin with the security group's phase. /// public void OnInstrumentHalted(in InstrumentHaltedEvent e) { AssertOnLoopThread(); - byte phaseByte = _phaseSnapshot.TryGetValue(e.SecurityId, out var phase) - ? (byte)phase - : (byte)TradingPhase.Open; - UmdfFrameBuilder.WriteInstrumentHalted(FrameSink, e.SecurityId, phaseByte, e.RptSeq, e.TransactTimeNanos); + UmdfFrameBuilder.WriteInstrumentHalted(FrameSink, e.SecurityId, e.RptSeq, e.TransactTimeNanos); } public void OnInstrumentResumed(in InstrumentResumedEvent e) diff --git a/src/B3.Exchange.Core/ISnapshotBookSource.cs b/src/B3.Exchange.Core/ISnapshotBookSource.cs index e780d79..6194478 100644 --- a/src/B3.Exchange.Core/ISnapshotBookSource.cs +++ b/src/B3.Exchange.Core/ISnapshotBookSource.cs @@ -26,6 +26,16 @@ public interface ISnapshotBookSource /// uint GetCurrentRptSeq(long securityId); + /// + /// The current official UMDF 2.2.0 SecurityTradingStatus value + /// for (issue #583): FORBIDDEN while the + /// instrument is administratively halted, otherwise the instrument's + /// current trading-group phase. stamps + /// this into the trailing SecurityStatus_3 packet of every + /// published snapshot. + /// + byte GetSecurityTradingStatus(long securityId); + /// /// Iterates the resting orders for on the /// requested side in price-time priority. Caller must fully drain the @@ -55,6 +65,10 @@ public MatchingEngineSnapshotSource(MatchingEngine engine, IReadOnlyList s public uint GetCurrentRptSeq(long securityId) => _engine.GetCurrentRptSeq(securityId); + public byte GetSecurityTradingStatus(long securityId) => _engine.IsHalted(securityId, out _) + ? B3.Umdf.WireEncoder.UmdfFrameBuilder.SecurityTradingStatusForbidden + : (byte)_engine.GetTradingPhase(securityId); + public IEnumerable EnumerateBook(long securityId, Side side) => _engine.EnumerateBook(securityId, side); } diff --git a/src/B3.Exchange.Core/SnapshotRotator.cs b/src/B3.Exchange.Core/SnapshotRotator.cs index 7d0e693..00325c4 100644 --- a/src/B3.Exchange.Core/SnapshotRotator.cs +++ b/src/B3.Exchange.Core/SnapshotRotator.cs @@ -11,7 +11,10 @@ namespace B3.Exchange.Core; /// SnapshotFullRefresh_Header_30 + as many /// SnapshotFullRefresh_Orders_MBO_71 chunks as the book requires (each /// chunk capped at or -/// the caller-supplied per-chunk cap, whichever is smaller). +/// the caller-supplied per-chunk cap, whichever is smaller), followed by a +/// trailing SecurityStatus_3 packet reporting the instrument's current +/// official trading status (issue #583) via +/// . /// /// /// Threading: reads the matching engine's resting @@ -177,6 +180,7 @@ public int PublishFor(long securityId, ushort incrementalSequenceVersion) firstSequenceNumber: firstSeq, sendingTimeNanos: nowNanos, securityId: securityId, + securityTradingStatus: _source.GetSecurityTradingStatus(securityId), lastRptSeq: lastRptSeq, incrementalSequenceVersion: incrementalSequenceVersion, bids: System.Runtime.InteropServices.CollectionsMarshal.AsSpan(bidsList), diff --git a/src/B3.Umdf.WireEncoder/SnapshotPacketBuilder.cs b/src/B3.Umdf.WireEncoder/SnapshotPacketBuilder.cs index 2b12685..6c3cd55 100644 --- a/src/B3.Umdf.WireEncoder/SnapshotPacketBuilder.cs +++ b/src/B3.Umdf.WireEncoder/SnapshotPacketBuilder.cs @@ -5,6 +5,10 @@ namespace B3.Umdf.WireEncoder; /// the first packet carries the SnapshotFullRefresh_Header_30 frame /// followed by as many SnapshotFullRefresh_Orders_MBO_71 frames as fit /// in the destination buffer; subsequent packets carry only Orders_71 frames. +/// The final packet carries a single SecurityStatus_3 frame reporting +/// the instrument's current official trading status — e.g. FORBIDDEN while +/// administratively halted (issue #583) — so late subscribers can bootstrap +/// without waiting for the next incremental status transition. /// Each packet is written to a caller-supplied via the /// delegate so the caller can pool the buffer and /// invoke its multicast publisher synchronously. @@ -63,7 +67,12 @@ public static int GetPacketCount( bufferSize - WireOffsets.PacketHeaderSize, maxEntriesPerChunk); } - return packetCount; + + // Issue #583: every snapshot ends with one trailing packet carrying + // the instrument's current SecurityStatus_3 state, so late + // subscribers can bootstrap FORBIDDEN/halted instruments without + // waiting for the next incremental status transition. + return checked(packetCount + 1); } /// @@ -84,6 +93,12 @@ public static int GetPacketCount( /// stamped on every packet of the snapshot so consumers see it as one /// logical refresh. /// Instrument identifier (Trade/Order semantic). + /// Current official UMDF 2.2.0 + /// SecurityTradingStatus value for this instrument (issue #583): + /// e.g. FORBIDDEN=18 while administratively halted, or the current + /// trading-group phase otherwise. Written as a trailing + /// SecurityStatus_3 packet so late subscribers can bootstrap + /// without waiting for the next incremental status transition. /// Last incremental RptSeq published before this /// snapshot was taken; consumers gate snapshot acceptance on this value. /// Pass null for an empty illiquid snapshot (per B3 §7.4). @@ -103,6 +118,7 @@ public static int WriteSnapshot( uint firstSequenceNumber, ulong sendingTimeNanos, long securityId, + byte securityTradingStatus, uint? lastRptSeq, ushort incrementalSequenceVersion, ReadOnlySpan bids, @@ -164,6 +180,27 @@ public static int WriteSnapshot( packetCount++; } + // -------- Trailing packet: current SecurityStatus_3 (issue #583) ---- + // Written as its own packet (not merged into the header packet) so + // the header/orders packing logic above is unaffected by the + // instrument's status and consumers can process it independently. + // securityTradingEvent is NULL (255): this is a point-in-time status + // report, not a transition event. + seqNum++; + p = UmdfWireEncoder.WritePacketHeader(buffer, channelNumber, snapshotSequenceVersion, seqNum, sendingTimeNanos); + p += UmdfWireEncoder.WriteSecurityStatusFrame( + buffer.Slice(p), + securityId: securityId, + tradingSessionId: 0, + securityTradingStatus: securityTradingStatus, + securityTradingEvent: 255, + tradeDate: 0, + tradSesOpenTimeNanos: 0, + transactTimeNanos: sendingTimeNanos, + rptSeq: 0); // schema: "Sequence number per instrument update. (Zeroed in snapshot feed)" + onPacket(buffer.Slice(0, p)); + packetCount++; + return packetCount; } @@ -178,6 +215,11 @@ private static int SnapshotHeaderFrameSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SnapHeaderBlockLength; + private static int SecurityStatusFrameSize + => WireOffsets.FramingHeaderSize + + WireOffsets.SbeMessageHeaderSize + + WireOffsets.SecurityStatusBlockLength; + private static void ValidateLayout( int bufferSize, int bidCount, @@ -199,6 +241,14 @@ private static void ValidateLayout( $"buffer too small ({bufferSize} bytes) for a single Orders_71 entry frame ({OrdersFrameSize(1)} bytes); increase buffer size.", nameof(bufferSize)); } + // Issue #583: the trailing SecurityStatus_3 packet reuses the same + // scratch buffer, starting a fresh PacketHeader at offset 0. + if (bufferSize < WireOffsets.PacketHeaderSize + SecurityStatusFrameSize) + { + throw new ArgumentException( + $"buffer too small ({bufferSize} bytes) for the trailing SecurityStatus_3 packet ({WireOffsets.PacketHeaderSize + SecurityStatusFrameSize} bytes); increase buffer size.", + nameof(bufferSize)); + } } private static void ConsumeEntriesThatFit( diff --git a/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs b/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs index d58aa25..4f63353 100644 --- a/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs +++ b/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs @@ -10,18 +10,27 @@ namespace B3.Umdf.WireEncoder; public static class UmdfFrameBuilder { /// - /// SecurityTradingEvent byte written for an instrument halt - /// (SecurityStatus_3). Matches the B3-aligned value documented - /// in issue #322. + /// Official UMDF 2.2.0 SecurityTradingStatus value meaning the + /// instrument is not available for trading (administrative halt). + /// Issue #583: replaces the previous out-of-domain marker approach. /// - public const byte SecurityTradingEventHalt = 1; + public const byte SecurityTradingStatusForbidden = 18; /// - /// SecurityTradingEvent byte written for an instrument resume - /// (SecurityStatus_3). Matches the B3-aligned value documented - /// in issue #322. + /// Official UMDF 2.2.0 SecurityTradingEvent value meaning the + /// instrument's status is being maintained separately from its + /// security group's phase. Written when an instrument enters + /// administrative halt (issue #583). /// - public const byte SecurityTradingEventResume = 2; + public const byte SecurityTradingEventSecurityStatusChange = 101; + + /// + /// Official UMDF 2.2.0 SecurityTradingEvent value meaning the + /// instrument's status now follows its security group's phase again. + /// Written when an instrument resumes from administrative halt, + /// restoring the preserved pre-halt phase (issue #583). + /// + public const byte SecurityTradingEventSecurityRejoinsGroupStatus = 102; /// /// Writes an Order_MBO_50 NEW frame (action=NEW). @@ -164,14 +173,18 @@ public static void WriteTradingPhaseChanged( } /// - /// Writes a SecurityStatus_3 frame for an instrument halt - /// (securityTradingEvent = ). - /// Used for OnInstrumentHalted. + /// Writes a SecurityStatus_3 frame for an instrument halt. + /// securityTradingStatus is always + /// (issue #583: the + /// official "not available for trading" status, not the pre-halt + /// phase) and securityTradingEvent is + /// , marking + /// that this instrument's status now diverges from its security + /// group's phase. Used for OnInstrumentHalted. /// public static void WriteInstrumentHalted( IUmdfFrameSink sink, long securityId, - byte securityTradingStatus, uint rptSeq, ulong transactTimeNanos) { @@ -183,8 +196,8 @@ public static void WriteInstrumentHalted( dst, securityId: securityId, tradingSessionId: 0, - securityTradingStatus: securityTradingStatus, - securityTradingEvent: SecurityTradingEventHalt, + securityTradingStatus: SecurityTradingStatusForbidden, + securityTradingEvent: SecurityTradingEventSecurityStatusChange, tradeDate: 0, tradSesOpenTimeNanos: 0, transactTimeNanos: transactTimeNanos, @@ -193,14 +206,19 @@ public static void WriteInstrumentHalted( } /// - /// Writes a SecurityStatus_3 frame for an instrument resume - /// (securityTradingEvent = ). - /// Used for OnInstrumentResumed. + /// Writes a SecurityStatus_3 frame for an instrument resume. + /// is the pre-halt + /// phase the engine preserved while the instrument was halted, and + /// securityTradingEvent is + /// , + /// marking that this instrument's status now follows its security + /// group's phase again (issue #583). Used for + /// OnInstrumentResumed. /// public static void WriteInstrumentResumed( IUmdfFrameSink sink, long securityId, - byte securityTradingStatus, + byte restoredSecurityTradingStatus, uint rptSeq, ulong transactTimeNanos) { @@ -212,8 +230,8 @@ public static void WriteInstrumentResumed( dst, securityId: securityId, tradingSessionId: 0, - securityTradingStatus: securityTradingStatus, - securityTradingEvent: SecurityTradingEventResume, + securityTradingStatus: restoredSecurityTradingStatus, + securityTradingEvent: SecurityTradingEventSecurityRejoinsGroupStatus, tradeDate: 0, tradSesOpenTimeNanos: 0, transactTimeNanos: transactTimeNanos, diff --git a/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs b/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs index 9b19f59..eb10f70 100644 --- a/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs +++ b/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs @@ -33,9 +33,14 @@ private sealed class FakeSource : ISnapshotBookSource public Dictionary RptSeqBySecurity { get; } = new(); public Dictionary<(long, Side), List> Books { get; } = new(); + /// Default OPEN=17; tests override per issue #583 scenarios. + public byte SecurityTradingStatus { get; set; } = 17; + public uint GetCurrentRptSeq(long securityId) => RptSeqBySecurity.TryGetValue(securityId, out uint rptSeq) ? rptSeq : 0; + public byte GetSecurityTradingStatus(long securityId) => SecurityTradingStatus; + public IEnumerable EnumerateBook(long securityId, Side side) => Books.TryGetValue((securityId, side), out var l) ? l : Enumerable.Empty(); } @@ -58,9 +63,9 @@ public void EmptyBook_NoPriorIncrementals_EmitsIlliquidHeaderOnly() int packets = rot.PublishNext(incrementalSequenceVersion: 3); - Assert.Equal(1, packets); - Assert.Single(sink.Packets); - Assert.Equal(1u, rot.SequenceNumber); + Assert.Equal(2, packets); // header (illiquid) + trailing SecurityStatus_3 (issue #583) + Assert.Equal(2, sink.Packets.Count); + Assert.Equal(2u, rot.SequenceNumber); var pkt = sink.Packets[0]; ref readonly var hdr = ref MemoryMarshal.AsRef(pkt.AsSpan(0, PacketHeaderSize)); @@ -75,6 +80,10 @@ public void EmptyBook_NoPriorIncrementals_EmitsIlliquidHeaderOnly() Assert.Equal(0u, snapHdr.Data.TotNumReports); Assert.Null(snapHdr.Data.LastRptSeq); // illiquid: §7.4 Assert.Equal((ushort)3, snapHdr.Data.LastSequenceVersion); + + // Issue #583: trailing packet reports the current SecurityTradingStatus. + int statusOffset = FrameOffset + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + Assert.Equal(src.SecurityTradingStatus, sink.Packets[1][statusOffset]); } [Fact] @@ -96,8 +105,8 @@ public void NonEmptyBook_BuildsHeaderPlusOrdersFrame_WithLastRptSeq() var rot = new SnapshotRotator(channelNumber: 84, source: src, sink: sink); int packets = rot.PublishNext(incrementalSequenceVersion: 4); - Assert.Equal(1, packets); - Assert.Single(sink.Packets); + Assert.Equal(2, packets); // header+orders packet + trailing status packet + Assert.Equal(2, sink.Packets.Count); var pkt = sink.Packets[0]; Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( @@ -114,6 +123,38 @@ public void NonEmptyBook_BuildsHeaderPlusOrdersFrame_WithLastRptSeq() + WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SnapOrdersHeaderBlockLength + 2; Assert.Equal((byte)3, pkt[groupNumInGroupOff]); + + // Issue #583: SecurityStatus_3.rptSeq must be zeroed in the snapshot + // feed per schema, even though LastRptSeq (17) is non-null above. + int trailingRptSeqOffset = FrameOffset + WireOffsets.SecurityStatusBodyRptSeqOffset; + Assert.Equal(0u, BitConverter.ToUInt32(sink.Packets[1], trailingRptSeqOffset)); + } + + [Fact] + public void HaltedInstrument_TrailingStatusPacket_ReportsForbidden() + { + // Issue #583: a late subscriber joining while an instrument is halted + // must learn the FORBIDDEN status from the snapshot's trailing + // SecurityStatus_3 packet, not just from a future incremental event. + const byte securityTradingStatusForbidden = 18; + var src = new FakeSource + { + SecurityIds = new[] { 42L }, + SecurityTradingStatus = securityTradingStatusForbidden, + }; + src.RptSeqBySecurity[42L] = 9; + src.Books[(42L, Side.Buy)] = new() { Order(1, Side.Buy, 100_0000, 500) }; + + var sink = new CapturingSink(); + var rot = new SnapshotRotator(channelNumber: 84, source: src, sink: sink); + + int packets = rot.PublishNext(incrementalSequenceVersion: 1); + + Assert.Equal(2, packets); + Assert.Equal(2, sink.Packets.Count); + + int statusOffset = FrameOffset + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + Assert.Equal(securityTradingStatusForbidden, sink.Packets[1][statusOffset]); } [Fact] @@ -125,20 +166,22 @@ public void RoundRobinsThroughInstruments_AndAdvancesSnapSequence() for (int i = 0; i < 4; i++) rot.PublishNext(incrementalSequenceVersion: 1); // wraps after 3 - Assert.Equal(4, sink.Packets.Count); - // Each tick → 1 packet (empty book), so SequenceNumber == 4. - Assert.Equal(4u, rot.SequenceNumber); + // Each tick → 2 packets (empty-book header + trailing status), issue #583. + Assert.Equal(8, sink.Packets.Count); + Assert.Equal(8u, rot.SequenceNumber); // SecurityIDs read out of each header packet must follow the rotation - // order 11, 22, 33, 11. + // order 11, 22, 33, 11. Header packets are the first of each 2-packet + // publish, at indices 0, 2, 4, 6. long[] expected = new[] { 11L, 22L, 33L, 11L }; for (int i = 0; i < expected.Length; i++) { + int headerIdx = i * 2; Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( - sink.Packets[i].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var hdr)); + sink.Packets[headerIdx].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var hdr)); Assert.Equal(expected[i], (long)(ulong)hdr.Data.SecurityID); - ref readonly var packetHdr = ref MemoryMarshal.AsRef(sink.Packets[i].AsSpan(0, PacketHeaderSize)); - Assert.Equal((uint)(i + 1), packetHdr.SequenceNumber); + ref readonly var packetHdr = ref MemoryMarshal.AsRef(sink.Packets[headerIdx].AsSpan(0, PacketHeaderSize)); + Assert.Equal((uint)(headerIdx + 1), packetHdr.SequenceNumber); } } @@ -156,11 +199,13 @@ public void PublishFor_UsesExactTargetSecurityRptSeq() rot.PublishFor(11, incrementalSequenceVersion: 4); rot.PublishFor(33, incrementalSequenceVersion: 4); + // Each PublishFor emits 2 packets (empty-book header + trailing + // status, issue #583); header packets sit at indices 0, 2, 4. uint[] expected = [2, 7, 11]; for (int i = 0; i < expected.Length; i++) { Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( - sink.Packets[i].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), + sink.Packets[i * 2].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var header)); Assert.Equal(expected[i], header.Data.LastRptSeq); } @@ -207,7 +252,7 @@ public void LargeBookAboveLegacyBufferLimit_PublishesCompleteSnapshotAndStepsSeq Assert.Equal((ushort)17, header.Data.LastSequenceVersion); int totalEntries = 0; - for (int i = 0; i < packets; i++) + for (int i = 0; i < packets - 1; i++) { var pkt = sink.Packets[i]; int p = PacketHeaderSize; @@ -308,7 +353,8 @@ public void BumpSequenceVersion_DoesNotChangeIncrementalLastSequenceVersion() rot.PublishNext(incrementalSequenceVersion: 9); rot.PublishNext(incrementalSequenceVersion: 9); - Assert.Equal(2u, rot.SequenceNumber); + // Each publish → 2 packets (empty-book header + trailing status). + Assert.Equal(4u, rot.SequenceNumber); Assert.Equal((ushort)1, rot.SequenceVersion); rot.BumpSequenceVersion(); @@ -317,12 +363,14 @@ public void BumpSequenceVersion_DoesNotChangeIncrementalLastSequenceVersion() Assert.Equal(0u, rot.SequenceNumber); rot.PublishNext(incrementalSequenceVersion: 9); - // Next packet: SequenceVersion=2, SequenceNumber=1 - ref readonly var hdr = ref MemoryMarshal.AsRef(sink.Packets[^1].AsSpan(0, PacketHeaderSize)); + // Next publish: SequenceVersion=2, header packet SequenceNumber=1, + // trailing status packet SequenceNumber=2. + var headerPkt = sink.Packets[^2]; + ref readonly var hdr = ref MemoryMarshal.AsRef(headerPkt.AsSpan(0, PacketHeaderSize)); Assert.Equal((ushort)2, hdr.SequenceVersion); Assert.Equal(1u, hdr.SequenceNumber); Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( - sink.Packets[^1].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var snapHdr)); + headerPkt.AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var snapHdr)); Assert.Equal((ushort)9, snapHdr.Data.LastSequenceVersion); } @@ -404,9 +452,9 @@ public void SnapshotTick_PublishedOnDispatcherThread_AndIsolatedFromIncrementalS Assert.True(disp.EnqueueSnapshotTick()); Drain(disp); - Assert.Single(snapSink.Packets); // snapshot went to snap-sink only + Assert.Equal(2, snapSink.Packets.Count); // header + trailing status (issue #583) Assert.Empty(incSink.Packets); // incremental sink untouched - Assert.Equal(1u, rotator.SequenceNumber); // snap-channel seq advanced + Assert.Equal(2u, rotator.SequenceNumber); // snap-channel seq advanced Assert.Equal(0u, disp.SequenceNumber); // inc-channel seq untouched Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( snapSink.Packets[0].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var snapshotHeader)); @@ -449,7 +497,7 @@ public void SnapshotTick_ReadsLiveRestingOrders_FromMatchingEngine() Assert.True(disp.EnqueueSnapshotTick()); Drain(disp); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); // header(+orders) + trailing status (issue #583) var pkt = snapSink.Packets[0]; int frameOff = WireOffsets.PacketHeaderSize + WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize; Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( @@ -488,7 +536,7 @@ public void SnapshotTick_ReflectsIncrementalVersionBump() Assert.True(disp.EnqueueSnapshotTick()); Drain(disp); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); // header(+orders) + trailing status (issue #583) ref readonly var packetHeader = ref MemoryMarshal.AsRef( snapSink.Packets[0].AsSpan(0, PacketHeaderSize)); Assert.Equal((ushort)2, packetHeader.SequenceVersion); @@ -538,7 +586,7 @@ public void SnapshotTick_ReflectsRestoredIncrementalVersion() Assert.True(restoredDispatcher.EnqueueSnapshotTick()); Drain(restoredDispatcher); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); // header(+orders) + trailing status (issue #583) ref readonly var packetHeader = ref MemoryMarshal.AsRef( snapSink.Packets[0].AsSpan(0, PacketHeaderSize)); Assert.Equal((ushort)1, packetHeader.SequenceVersion); diff --git a/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs b/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs index f57d5d0..b5ee88a 100644 --- a/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs +++ b/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs @@ -329,7 +329,8 @@ public async Task RestoreReplay_StartsDurableEpoch_PreservesState_AndOrdersReset Assert.Equal(HaltReason.RegulatoryHalt, restoredHalt.Reason); Assert.True(dispatcher.EnqueueSnapshotTick()); - Assert.True(WaitFor(() => snapshotSink.Packets.Count == 1)); + // Header(+orders) packet + trailing SecurityStatus_3 packet (issue #583). + Assert.True(WaitFor(() => snapshotSink.Packets.Count == 2)); Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( snapshotSink.Packets[0].AsSpan( SnapshotFrameOffset, WireOffsets.SnapHeaderBlockLength), @@ -340,6 +341,11 @@ public async Task RestoreReplay_StartsDurableEpoch_PreservesState_AndOrdersReset Assert.Equal(30u, snapshotHeader.Data.LastRptSeq); Assert.Equal((ushort)10, snapshotHeader.Data.LastSequenceVersion); + // Rotation publishes Petr first; it isn't halted, so the trailing + // status packet reports its restored OPEN phase, not FORBIDDEN. + int statusOffset = SnapshotFrameOffset + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + Assert.Equal((byte)TradingPhase.Open, snapshotSink.Packets[1][statusOffset]); + Assert.True(dispatcher.EnqueueNewOrder( new NewOrderCommand( "LIVE-BID", Petr, Side.Buy, OrderType.Limit, TimeInForce.Day, diff --git a/tests/B3.Umdf.WireEncoder.Tests/SnapshotPacketBuilderTests.cs b/tests/B3.Umdf.WireEncoder.Tests/SnapshotPacketBuilderTests.cs index 6d1ed17..614ded1 100644 --- a/tests/B3.Umdf.WireEncoder.Tests/SnapshotPacketBuilderTests.cs +++ b/tests/B3.Umdf.WireEncoder.Tests/SnapshotPacketBuilderTests.cs @@ -28,21 +28,22 @@ private sealed class CapturingHandler } [Fact] - public void EmptyBook_EmitsSingleHeaderOnlyPacket() + public void EmptyBook_EmitsHeaderAndTrailingStatusPacket() { var sink = new CapturingHandler(); var buf = new byte[SnapshotPacketBuilder.DefaultPacketBufferSize]; int packets = SnapshotPacketBuilder.WriteSnapshot(buf, channelNumber: 84, snapshotSequenceVersion: 0, firstSequenceNumber: 100, - sendingTimeNanos: 1_700_000_000_000_000_000UL, securityId: 7L, lastRptSeq: null, + sendingTimeNanos: 1_700_000_000_000_000_000UL, securityId: 7L, + securityTradingStatus: 17, lastRptSeq: null, incrementalSequenceVersion: 7, bids: ReadOnlySpan.Empty, asks: ReadOnlySpan.Empty, onPacket: sink.OnPacket); - Assert.Equal(1, packets); - Assert.Single(sink.Packets); + Assert.Equal(2, packets); + Assert.Equal(2, sink.Packets.Count); // PacketHeader.SequenceNumber == 100 ref readonly var hdr = ref MemoryMarshal.AsRef(sink.Packets[0].AsSpan(0, 16)); @@ -56,6 +57,17 @@ public void EmptyBook_EmitsSingleHeaderOnlyPacket() Assert.Equal(0u, rdr.Data.TotNumOffers); Assert.Null(rdr.Data.LastRptSeq); Assert.Equal((ushort)7, rdr.Data.LastSequenceVersion); + + // Issue #583: trailing packet carries a SecurityStatus_3 frame + // stamped with the caller-supplied trading status and RptSeq NULL + // (255 sentinel event, rptSeq 0 for the illiquid/no-history case). + var statusPkt = sink.Packets[1]; + ref readonly var statusHdr = ref MemoryMarshal.AsRef(statusPkt.AsSpan(0, 16)); + Assert.Equal(101u, statusHdr.SequenceNumber); + int statusOffset = FrameOffset + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + int eventOffset = FrameOffset + WireOffsets.SecurityStatusBodySecurityTradingEventOffset; + Assert.Equal((byte)17, statusPkt[statusOffset]); + Assert.Equal((byte)255, statusPkt[eventOffset]); } [Fact] @@ -68,10 +80,10 @@ public void SmallBook_SinglePacketCarriesHeaderAndOrders() var asks = new[] { Ask(101_0000L, 300L, 3) }; int packets = SnapshotPacketBuilder.WriteSnapshot( - buf, 84, 0, 50, 1UL, 7L, lastRptSeq: 99u, incrementalSequenceVersion: 8, - bids, asks, sink.OnPacket); + buf, 84, 0, 50, 1UL, 7L, securityTradingStatus: 17, lastRptSeq: 99u, + incrementalSequenceVersion: 8, bids, asks, sink.OnPacket); - Assert.Equal(1, packets); + Assert.Equal(2, packets); var pkt = sink.Packets[0]; // Header counts derived from input spans. @@ -113,7 +125,8 @@ public void BookLargerThanPacketBuffer_ChunksAcrossPacketsWithSequentialSeqNums( for (int i = 0; i < asks.Length; i++) asks[i] = Ask(101_0000L + i, 100L, 10_000 + i); int packets = SnapshotPacketBuilder.WriteSnapshot(buf, 84, 0, - firstSequenceNumber: 1000, sendingTimeNanos: 42UL, securityId: 7L, lastRptSeq: 12345u, + firstSequenceNumber: 1000, sendingTimeNanos: 42UL, securityId: 7L, + securityTradingStatus: 17, lastRptSeq: 12345u, incrementalSequenceVersion: 11, bids, asks, sink.OnPacket); @@ -137,9 +150,10 @@ public void BookLargerThanPacketBuffer_ChunksAcrossPacketsWithSequentialSeqNums( Assert.Equal(700u, hdr.Data.TotNumBids); Assert.Equal(300u, hdr.Data.TotNumOffers); - // Walk both packets summing NumInGroup across every Orders_71 frame. + // Walk every packet except the trailing SecurityStatus_3 packet, + // summing NumInGroup across every Orders_71 frame. int totalEntries = 0; - for (int i = 0; i < packets; i++) + for (int i = 0; i < packets - 1; i++) { var pkt = sink.Packets[i]; int p = WireOffsets.PacketHeaderSize; @@ -156,6 +170,10 @@ public void BookLargerThanPacketBuffer_ChunksAcrossPacketsWithSequentialSeqNums( } Assert.Equal(1_000, totalEntries); Assert.All(sink.Packets, packet => Assert.InRange(packet.Length, 1, buf.Length)); + + // Trailing packet carries the SecurityStatus_3 frame. + int statusOffset = FrameOffset + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + Assert.Equal((byte)17, sink.Packets[^1][statusOffset]); } [Fact] @@ -171,7 +189,7 @@ public void BufferTooSmallForSingleEntry_Throws() var bids = new[] { Bid(1L, 1L, 1) }; Assert.Throws(() => - SnapshotPacketBuilder.WriteSnapshot(buf, 84, 0, 0, 0, 7L, null, 1, + SnapshotPacketBuilder.WriteSnapshot(buf, 84, 0, 0, 0, 7L, 17, null, 1, bids, ReadOnlySpan.Empty, sink.OnPacket)); } @@ -185,13 +203,15 @@ public void RespectsCustomMaxEntriesPerChunk() for (int i = 0; i < bids.Length; i++) bids[i] = Bid(100L - i, 1L, i + 1); // cap=3 → expect 4 chunks (3,3,3,1). They all fit in one 16 KB buffer - // alongside the snapshot header → 1 packet. + // alongside the snapshot header → 1 book packet + 1 trailing status + // packet = 2 packets total. int packets = SnapshotPacketBuilder.WriteSnapshot( - buf, 84, 0, 0, 0, 7L, lastRptSeq: 1u, incrementalSequenceVersion: 1, + buf, 84, 0, 0, 0, 7L, securityTradingStatus: 17, lastRptSeq: 1u, + incrementalSequenceVersion: 1, bids, ReadOnlySpan.Empty, sink.OnPacket, maxEntriesPerChunk: 3); - Assert.Equal(1, packets); + Assert.Equal(2, packets); // Walk the packet and count Orders_71 frames + their NumInGroup. var pkt = sink.Packets[0]; int p = WireOffsets.PacketHeaderSize; @@ -220,12 +240,12 @@ public void RejectsInvalidMaxEntriesPerChunk() var sink = new CapturingHandler(); var buf = new byte[1024]; Assert.Throws(() => - SnapshotPacketBuilder.WriteSnapshot(buf, 0, 0, 0, 0, 0, null, 1, + SnapshotPacketBuilder.WriteSnapshot(buf, 0, 0, 0, 0, 0, 17, null, 1, ReadOnlySpan.Empty, ReadOnlySpan.Empty, sink.OnPacket, maxEntriesPerChunk: 0)); Assert.Throws(() => - SnapshotPacketBuilder.WriteSnapshot(buf, 0, 0, 0, 0, 0, null, 1, + SnapshotPacketBuilder.WriteSnapshot(buf, 0, 0, 0, 0, 0, 17, null, 1, ReadOnlySpan.Empty, ReadOnlySpan.Empty, sink.OnPacket, maxEntriesPerChunk: 256)); diff --git a/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs b/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs index 38c1534..7618095 100644 --- a/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs +++ b/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs @@ -151,7 +151,7 @@ public void WriteInstrumentHalted_ReservesAndCommitsCorrectSize() { var sink = MakeSink(); UmdfFrameBuilder.WriteInstrumentHalted(sink, - securityId: 1L, securityTradingStatus: 2, rptSeq: 6u, transactTimeNanos: 0ul); + securityId: 1L, rptSeq: 6u, transactTimeNanos: 0ul); int expected = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBlockLength; @@ -160,15 +160,18 @@ public void WriteInstrumentHalted_ReservesAndCommitsCorrectSize() } [Fact] - public void WriteInstrumentHalted_WritesHaltEventByte() + public void WriteInstrumentHalted_WritesForbiddenStatusAndSecurityStatusChangeEvent() { var sink = MakeSink(); UmdfFrameBuilder.WriteInstrumentHalted(sink, - securityId: 1L, securityTradingStatus: 2, rptSeq: 6u, transactTimeNanos: 0ul); + securityId: 1L, rptSeq: 6u, transactTimeNanos: 0ul); + int statusOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; int eventOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBodySecurityTradingEventOffset; - Assert.Equal(UmdfFrameBuilder.SecurityTradingEventHalt, sink.Buffer[eventOffset]); + Assert.Equal(UmdfFrameBuilder.SecurityTradingStatusForbidden, sink.Buffer[statusOffset]); + Assert.Equal(UmdfFrameBuilder.SecurityTradingEventSecurityStatusChange, sink.Buffer[eventOffset]); } [Fact] @@ -176,7 +179,7 @@ public void WriteInstrumentResumed_ReservesAndCommitsCorrectSize() { var sink = MakeSink(); UmdfFrameBuilder.WriteInstrumentResumed(sink, - securityId: 1L, securityTradingStatus: 2, rptSeq: 7u, transactTimeNanos: 0ul); + securityId: 1L, restoredSecurityTradingStatus: 17, rptSeq: 7u, transactTimeNanos: 0ul); int expected = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBlockLength; @@ -185,15 +188,18 @@ public void WriteInstrumentResumed_ReservesAndCommitsCorrectSize() } [Fact] - public void WriteInstrumentResumed_WritesResumeEventByte() + public void WriteInstrumentResumed_WritesRestoredStatusAndRejoinsGroupEvent() { var sink = MakeSink(); UmdfFrameBuilder.WriteInstrumentResumed(sink, - securityId: 1L, securityTradingStatus: 2, rptSeq: 7u, transactTimeNanos: 0ul); + securityId: 1L, restoredSecurityTradingStatus: 17, rptSeq: 7u, transactTimeNanos: 0ul); + int statusOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; int eventOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBodySecurityTradingEventOffset; - Assert.Equal(UmdfFrameBuilder.SecurityTradingEventResume, sink.Buffer[eventOffset]); + Assert.Equal((byte)17, sink.Buffer[statusOffset]); + Assert.Equal(UmdfFrameBuilder.SecurityTradingEventSecurityRejoinsGroupStatus, sink.Buffer[eventOffset]); } [Fact]