diff --git a/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs b/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs index 23cc01e..0645a18 100644 --- a/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs +++ b/src/B3.Exchange.Core/ChannelDispatcher.Sinks.cs @@ -133,22 +133,10 @@ public void OnTradingPhaseChanged(in TradingPhaseChangedEvent e) UmdfFrameBuilder.WriteTradingPhaseChanged(FrameSink, e.SecurityId, (byte)e.Phase, e.RptSeq, e.TransactTimeNanos); } - /// - /// 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. - /// - 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..663580e 100644 --- a/src/B3.Exchange.Core/ISnapshotBookSource.cs +++ b/src/B3.Exchange.Core/ISnapshotBookSource.cs @@ -4,9 +4,10 @@ namespace B3.Exchange.Core; /// /// Minimal read-only view that needs over a -/// matching engine to materialise a per-symbol UMDF snapshot. Abstracted so -/// the rotator can be unit-tested with a hand-rolled book without spinning -/// a real . +/// matching engine to materialise a per-symbol UMDF snapshot plus the +/// current SecurityStatus_3 recovery frame. Abstracted so the +/// rotator can be unit-tested with a hand-rolled book without spinning a +/// real . /// /// All members are invoked on the owning 's /// dispatch thread — implementations need not be thread-safe. @@ -32,6 +33,13 @@ public interface ISnapshotBookSource /// enumerator before mutating the book. /// IEnumerable EnumerateBook(long securityId, Side side); + + /// + /// Returns the current SecurityStatus_3 pair for + /// without mutating engine state or + /// advancing the incremental RptSeq. + /// + (byte SecurityTradingStatus, byte SecurityTradingEvent) GetCurrentSecurityStatus(long securityId); } /// @@ -57,4 +65,16 @@ public MatchingEngineSnapshotSource(MatchingEngine engine, IReadOnlyList s public IEnumerable EnumerateBook(long securityId, Side side) => _engine.EnumerateBook(securityId, side); + + public (byte SecurityTradingStatus, byte SecurityTradingEvent) GetCurrentSecurityStatus(long securityId) + { + if (_engine.IsHalted(securityId, out _)) + { + return ( + (byte)TradingPhase.Forbidden, + (byte)B3.Umdf.Mbo.Sbe.V16.SecurityTradingEvent.SECURITY_STATUS_CHANGE); + } + + return ((byte)_engine.GetTradingPhase(securityId), byte.MaxValue); + } } diff --git a/src/B3.Exchange.Core/SnapshotRotator.cs b/src/B3.Exchange.Core/SnapshotRotator.cs index 7d0e693..0b98500 100644 --- a/src/B3.Exchange.Core/SnapshotRotator.cs +++ b/src/B3.Exchange.Core/SnapshotRotator.cs @@ -11,7 +11,9 @@ 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 one +/// standalone SecurityStatus_3 recovery packet for that instrument's +/// current status. /// /// /// Threading: reads the matching engine's resting @@ -45,6 +47,10 @@ public sealed class SnapshotRotator private readonly INanosTimeSource _timeSource; private readonly int _maxEntriesPerChunk; private readonly byte[] _packetBuf; + private const int SecurityStatusPacketSize = WireOffsets.PacketHeaderSize + + WireOffsets.FramingHeaderSize + + WireOffsets.SbeMessageHeaderSize + + WireOffsets.SecurityStatusBlockLength; private int _rotationIndex; @@ -109,8 +115,9 @@ public void BumpSequenceVersion() /// Publishes a complete snapshot for the next instrument in the rotation. /// Returns the number of UDP packets emitted: 0 when /// is empty; otherwise at - /// least 1, since even an empty book emits the header packet. Caller MUST - /// invoke from the dispatcher thread so that + /// least 2, since even an empty book emits the header packet plus a + /// SecurityStatus_3 packet. Caller MUST invoke from the dispatcher + /// thread so that /// observes a stable book. /// must be captured from the /// dispatcher in the same turn as the book and RptSeq watermark. @@ -157,11 +164,13 @@ public int PublishFor(long securityId, ushort incrementalSequenceVersion) bool isBookEmpty = bidsList.Count == 0 && asksList.Count == 0; uint? lastRptSeq = isBookEmpty && currentRptSeq == 0 ? null : currentRptSeq; - int requiredPackets = SnapshotPacketBuilder.GetPacketCount( + int requiredSnapshotPackets = SnapshotPacketBuilder.GetPacketCount( _packetBuf.Length, bidsList.Count, asksList.Count, _maxEntriesPerChunk); + const int statusPackets = 1; + int requiredPackets = checked(requiredSnapshotPackets + statusPackets); EnsurePacketSequenceCapacity(requiredPackets); ushort version = SequenceVersion; @@ -169,8 +178,9 @@ public int PublishFor(long securityId, ushort incrementalSequenceVersion) ulong nowNanos = _timeSource.NowNanos(); byte channel = _channelNumber; var sink = _sink; + var (securityTradingStatus, securityTradingEvent) = _source.GetCurrentSecurityStatus(securityId); - int packetsEmitted = SnapshotPacketBuilder.WriteSnapshot( + int snapshotPacketsEmitted = SnapshotPacketBuilder.WriteSnapshot( buffer: _packetBuf, channelNumber: channel, snapshotSequenceVersion: version, @@ -184,13 +194,27 @@ public int PublishFor(long securityId, ushort incrementalSequenceVersion) onPacket: pkt => sink.Publish(channel, pkt), maxEntriesPerChunk: _maxEntriesPerChunk); - if (packetsEmitted != requiredPackets) + if (snapshotPacketsEmitted != requiredSnapshotPackets) { throw new InvalidOperationException( - $"snapshot packet count changed during encode: reserved {requiredPackets}, emitted {packetsEmitted}"); + $"snapshot packet count changed during encode: reserved {requiredSnapshotPackets}, emitted {snapshotPacketsEmitted}"); } - SequenceNumber = (uint)((ulong)SequenceNumber + (uint)packetsEmitted); - return packetsEmitted; + + int statusPacketLength = WriteSecurityStatusPacket( + dst: _packetBuf, + channelNumber: channel, + sequenceVersion: version, + sequenceNumber: firstSeq + (uint)snapshotPacketsEmitted, + sendingTimeNanos: nowNanos, + securityId: securityId, + securityTradingStatus: securityTradingStatus, + securityTradingEvent: securityTradingEvent, + transactTimeNanos: nowNanos, + rptSeq: currentRptSeq); + sink.Publish(channel, _packetBuf.AsSpan(0, statusPacketLength)); + + SequenceNumber = (uint)((ulong)SequenceNumber + (uint)requiredPackets); + return requiredPackets; } private void EnsurePacketSequenceCapacity(int requiredPackets) @@ -205,4 +229,38 @@ private void EnsurePacketSequenceCapacity(int requiredPackets) $"snapshot channel {_channelNumber} automatic epoch rollover"); SequenceNumber = 0; } + + private static int WriteSecurityStatusPacket( + Span dst, + byte channelNumber, + ushort sequenceVersion, + uint sequenceNumber, + ulong sendingTimeNanos, + long securityId, + byte securityTradingStatus, + byte securityTradingEvent, + ulong transactTimeNanos, + uint rptSeq) + { + if (dst.Length < SecurityStatusPacketSize) + throw new ArgumentException("buffer too small for SecurityStatus_3 packet", nameof(dst)); + + int packetHeaderLength = UmdfWireEncoder.WritePacketHeader( + dst, + channelNumber: channelNumber, + sequenceVersion: sequenceVersion, + sequenceNumber: sequenceNumber, + sendingTimeNanos: sendingTimeNanos); + int frameLength = UmdfWireEncoder.WriteSecurityStatusFrame( + dst.Slice(packetHeaderLength), + securityId: securityId, + tradingSessionId: 0, + securityTradingStatus: securityTradingStatus, + securityTradingEvent: securityTradingEvent, + tradeDate: 0, + tradSesOpenTimeNanos: 0, + transactTimeNanos: transactTimeNanos, + rptSeq: rptSeq); + return packetHeaderLength + frameLength; + } } diff --git a/src/B3.Exchange.Matching/Commands.cs b/src/B3.Exchange.Matching/Commands.cs index e957285..9e47a3c 100644 --- a/src/B3.Exchange.Matching/Commands.cs +++ b/src/B3.Exchange.Matching/Commands.cs @@ -115,11 +115,10 @@ public enum RejectReason : byte /// /// Reason an instrument was placed into administrative halt by the -/// operator halt/resume API (issue #322). Carried on the outbound -/// SecurityStatus_3 frame's securityTradingEvent byte -/// (best-effort mapping to FIX SecurityTradingEvent values) and on -/// the engine-level InstrumentHaltedEvent for downstream -/// observability. +/// operator halt/resume API (issue #322). Carried on the engine-level +/// InstrumentHaltedEvent for downstream observability. The vendored +/// UMDF 2.2.0 schema does not expose a field for this detail, so the wire +/// encoder publishes only the official halt/resume status semantics. /// public enum HaltReason : byte { diff --git a/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs b/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs index d58aa25..6d59706 100644 --- a/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs +++ b/src/B3.Umdf.WireEncoder/UmdfFrameBuilder.cs @@ -11,17 +11,23 @@ public static class UmdfFrameBuilder { /// /// SecurityTradingEvent byte written for an instrument halt - /// (SecurityStatus_3). Matches the B3-aligned value documented - /// in issue #322. + /// (SecurityStatus_3). The earlier issue #322 implementation + /// incorrectly emitted a proprietary out-of-domain value; this constant + /// now maps to the vendored UMDF 2.2.0 + /// SECURITY_STATUS_CHANGE enum member. /// - public const byte SecurityTradingEventHalt = 1; + public const byte SecurityTradingEventHalt = + (byte)B3.Umdf.Mbo.Sbe.V16.SecurityTradingEvent.SECURITY_STATUS_CHANGE; /// /// SecurityTradingEvent byte written for an instrument resume - /// (SecurityStatus_3). Matches the B3-aligned value documented - /// in issue #322. + /// (SecurityStatus_3). The earlier issue #322 implementation + /// incorrectly emitted a proprietary out-of-domain value; this constant + /// now maps to the vendored UMDF 2.2.0 + /// SECURITY_REJOINS_SECURITY_GROUP_STATUS enum member. /// - public const byte SecurityTradingEventResume = 2; + public const byte SecurityTradingEventResume = + (byte)B3.Umdf.Mbo.Sbe.V16.SecurityTradingEvent.SECURITY_REJOINS_SECURITY_GROUP_STATUS; /// /// Writes an Order_MBO_50 NEW frame (action=NEW). @@ -164,14 +170,16 @@ public static void WriteTradingPhaseChanged( } /// - /// Writes a SecurityStatus_3 frame for an instrument halt - /// (securityTradingEvent = ). + /// Writes a SecurityStatus_3 frame for an instrument halt. + /// Administrative halts always encode + /// securityTradingStatus=FORBIDDEN and + /// securityTradingEvent = + /// . /// Used for OnInstrumentHalted. /// public static void WriteInstrumentHalted( IUmdfFrameSink sink, long securityId, - byte securityTradingStatus, uint rptSeq, ulong transactTimeNanos) { @@ -183,7 +191,7 @@ public static void WriteInstrumentHalted( dst, securityId: securityId, tradingSessionId: 0, - securityTradingStatus: securityTradingStatus, + securityTradingStatus: (byte)B3.Umdf.Mbo.Sbe.V16.SecurityTradingStatus.FORBIDDEN, securityTradingEvent: SecurityTradingEventHalt, tradeDate: 0, tradSesOpenTimeNanos: 0, diff --git a/tests/B3.Exchange.Core.Tests/ChannelDispatcherTests.cs b/tests/B3.Exchange.Core.Tests/ChannelDispatcherTests.cs index 0313757..98e6422 100644 --- a/tests/B3.Exchange.Core.Tests/ChannelDispatcherTests.cs +++ b/tests/B3.Exchange.Core.Tests/ChannelDispatcherTests.cs @@ -9,6 +9,7 @@ using B3.Exchange.Core; using B3.Exchange.TestSupport; using B3.Exchange.Matching; +using B3.Umdf.Mbo.Sbe.V16; using B3.Umdf.WireEncoder; using Microsoft.Extensions.Logging.Abstractions; @@ -51,6 +52,15 @@ public partial class ChannelDispatcherTests private static long Px(decimal p) => (long)(p * 10_000m); + private static SecurityStatus_3DataReader ReadSecurityStatus(byte[] packet) + { + Assert.True(SecurityStatus_3Data.TryParse( + packet.AsSpan(WireOffsets.PacketHeaderSize + WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize, + WireOffsets.SecurityStatusBlockLength), + out var status)); + return status; + } + private sealed class RecordingPacketSink : IUmdfPacketSink { public List Packets { get; } = new(); @@ -2261,6 +2271,44 @@ public async Task Issue322_OperatorResumeInstrument_RestoresOrderAcceptanceAndCl Assert.Single(session.News); } + [Fact] + public void Issue583_OperatorHalt_EmitsForbiddenWithOfficialSecurityStatusChangeEvent() + { + var (disp, pkt, _) = NewDispatcher(); + + Assert.True(disp.EnqueueOperatorSetTradingPhase(Petr, TradingPhase.Reserved)); + DrainInbound(disp); + pkt.Packets.Clear(); + + Assert.True(disp.EnqueueOperatorHalt(Petr, HaltReason.RegulatoryHalt, null)); + DrainInbound(disp); + + var status = ReadSecurityStatus(Assert.Single(pkt.Packets)); + Assert.Equal(SecurityTradingStatus.FORBIDDEN, status.Data.SecurityTradingStatus); + Assert.Equal(SecurityTradingEvent.SECURITY_STATUS_CHANGE, status.Data.SecurityTradingEvent); + Assert.Equal(2u, status.Data.RptSeq); + } + + [Fact] + public void Issue583_OperatorResume_EmitsRestoredPhaseWithOfficialRejoinEvent() + { + var (disp, pkt, _) = NewDispatcher(); + + Assert.True(disp.EnqueueOperatorSetTradingPhase(Petr, TradingPhase.Reserved)); + DrainInbound(disp); + Assert.True(disp.EnqueueOperatorHalt(Petr, HaltReason.NewsHold, null)); + DrainInbound(disp); + pkt.Packets.Clear(); + + Assert.True(disp.EnqueueOperatorResume(Petr)); + DrainInbound(disp); + + var status = ReadSecurityStatus(Assert.Single(pkt.Packets)); + Assert.Equal(SecurityTradingStatus.RESERVED, status.Data.SecurityTradingStatus); + Assert.Equal(SecurityTradingEvent.SECURITY_REJOINS_SECURITY_GROUP_STATUS, status.Data.SecurityTradingEvent); + Assert.Equal(3u, status.Data.RptSeq); + } + [Fact] public void Issue322_OperatorPhaseChange_WhileHalted_FaultsCompletionWithInvalidOperation() { diff --git a/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs b/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs index 9b19f59..e23b072 100644 --- a/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs +++ b/tests/B3.Exchange.Core.Tests/SnapshotRotatorTests.cs @@ -31,6 +31,7 @@ private sealed class FakeSource : ISnapshotBookSource { public IReadOnlyList SecurityIds { get; init; } = Array.Empty(); public Dictionary RptSeqBySecurity { get; } = new(); + public Dictionary SecurityStatusBySecurity { get; } = new(); public Dictionary<(long, Side), List> Books { get; } = new(); public uint GetCurrentRptSeq(long securityId) @@ -38,6 +39,11 @@ public uint GetCurrentRptSeq(long securityId) public IEnumerable EnumerateBook(long securityId, Side side) => Books.TryGetValue((securityId, side), out var l) ? l : Enumerable.Empty(); + + public (byte SecurityTradingStatus, byte SecurityTradingEvent) GetCurrentSecurityStatus(long securityId) + => SecurityStatusBySecurity.TryGetValue(securityId, out var status) + ? status + : ((byte)SecurityTradingStatus.OPEN, byte.MaxValue); } private sealed class CapturingSink : IUmdfPacketSink @@ -49,6 +55,14 @@ private sealed class CapturingSink : IUmdfPacketSink private static RestingOrderView Order(long oid, Side side, long px, long qty, ulong nanos = 1_000UL, uint firm = 7) => new(oid, side, px, qty, firm, nanos); + private static SecurityStatus_3DataReader ReadSecurityStatusPacket(byte[] packet) + { + Assert.True(SecurityStatus_3Data.TryParse( + packet.AsSpan(FrameOffset, WireOffsets.SecurityStatusBlockLength), + out var status)); + return status; + } + [Fact] public void EmptyBook_NoPriorIncrementals_EmitsIlliquidHeaderOnly() { @@ -58,9 +72,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); + 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 +89,11 @@ 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); + + var status = ReadSecurityStatusPacket(sink.Packets[1]); + Assert.Equal(SecurityTradingStatus.OPEN, status.Data.SecurityTradingStatus); + Assert.Null(status.Data.SecurityTradingEvent); + Assert.Null(status.Data.RptSeq); } [Fact] @@ -96,8 +115,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); + Assert.Equal(2, sink.Packets.Count); var pkt = sink.Packets[0]; Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( @@ -114,6 +133,11 @@ public void NonEmptyBook_BuildsHeaderPlusOrdersFrame_WithLastRptSeq() + WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SnapOrdersHeaderBlockLength + 2; Assert.Equal((byte)3, pkt[groupNumInGroupOff]); + + var status = ReadSecurityStatusPacket(sink.Packets[1]); + Assert.Equal(SecurityTradingStatus.OPEN, status.Data.SecurityTradingStatus); + Assert.Null(status.Data.SecurityTradingEvent); + Assert.Equal(17u, status.Data.RptSeq); } [Fact] @@ -125,9 +149,8 @@ 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); + 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. @@ -135,10 +158,10 @@ public void RoundRobinsThroughInstruments_AndAdvancesSnapSequence() for (int i = 0; i < expected.Length; i++) { Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( - sink.Packets[i].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var hdr)); + sink.Packets[i * 2].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[i * 2].AsSpan(0, PacketHeaderSize)); + Assert.Equal((uint)(2 * i + 1), packetHdr.SequenceNumber); } } @@ -160,12 +183,37 @@ public void PublishFor_UsesExactTargetSecurityRptSeq() 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); } } + [Fact] + public void PublishFor_EmitsCurrentSecurityStatusPacketWithoutAdvancingRptSeq() + { + var src = new FakeSource { SecurityIds = new[] { 42L, 43L } }; + src.RptSeqBySecurity[42L] = 17; + src.RptSeqBySecurity[43L] = 21; + src.SecurityStatusBySecurity[42L] = ((byte)SecurityTradingStatus.FORBIDDEN, (byte)SecurityTradingEvent.SECURITY_STATUS_CHANGE); + src.SecurityStatusBySecurity[43L] = ((byte)SecurityTradingStatus.RESERVED, byte.MaxValue); + var sink = new CapturingSink(); + var rot = new SnapshotRotator(channelNumber: 5, source: src, sink: sink); + + rot.PublishFor(42L, incrementalSequenceVersion: 4); + rot.PublishFor(43L, incrementalSequenceVersion: 4); + + var halted = ReadSecurityStatusPacket(sink.Packets[1]); + Assert.Equal(SecurityTradingStatus.FORBIDDEN, halted.Data.SecurityTradingStatus); + Assert.Equal(SecurityTradingEvent.SECURITY_STATUS_CHANGE, halted.Data.SecurityTradingEvent); + Assert.Equal(17u, halted.Data.RptSeq); + + var active = ReadSecurityStatusPacket(sink.Packets[3]); + Assert.Equal(SecurityTradingStatus.RESERVED, active.Data.SecurityTradingStatus); + Assert.Null(active.Data.SecurityTradingEvent); + Assert.Equal(21u, active.Data.RptSeq); + } + [Fact] public void LargeBookAboveLegacyBufferLimit_PublishesCompleteSnapshotAndStepsSequence() { @@ -185,11 +233,13 @@ public void LargeBookAboveLegacyBufferLimit_PublishesCompleteSnapshotAndStepsSeq sink: sink, timeSource: new FakeNanosTimeSource(123_456UL)); + int snapshotPackets = SnapshotPacketBuilder.GetPacketCount( + SnapshotPacketBuilder.DefaultPacketBufferSize, 20_000, 5_000); int packets = rot.PublishNext(incrementalSequenceVersion: 17); - Assert.True(packets >= 2, $"expected multi-packet snapshot, got {packets}"); + Assert.Equal(snapshotPackets + 1, packets); Assert.Equal((uint)packets, rot.SequenceNumber); - for (int i = 0; i < packets; i++) + for (int i = 0; i < snapshotPackets; i++) { ref readonly var pkthdr = ref MemoryMarshal.AsRef(sink.Packets[i].AsSpan(0, PacketHeaderSize)); Assert.Equal((uint)(i + 1), pkthdr.SequenceNumber); @@ -198,6 +248,9 @@ public void LargeBookAboveLegacyBufferLimit_PublishesCompleteSnapshotAndStepsSeq Assert.InRange(sink.Packets[i].Length, 1, SnapshotPacketBuilder.DefaultPacketBufferSize); } + ref readonly var statusPacketHdr = ref MemoryMarshal.AsRef(sink.Packets[^1].AsSpan(0, PacketHeaderSize)); + Assert.Equal((uint)(snapshotPackets + 1), statusPacketHdr.SequenceNumber); + Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( sink.Packets[0].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var header)); Assert.Equal(25_000u, header.Data.TotNumReports); @@ -207,7 +260,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 < snapshotPackets; i++) { var pkt = sink.Packets[i]; int p = PacketHeaderSize; @@ -223,6 +276,9 @@ public void LargeBookAboveLegacyBufferLimit_PublishesCompleteSnapshotAndStepsSeq } } Assert.Equal(25_000, totalEntries); + + var status = ReadSecurityStatusPacket(sink.Packets[^1]); + Assert.Equal(99u, status.Data.RptSeq); } [Fact] @@ -238,9 +294,10 @@ public void MultiPacketAtSequenceBoundary_UsesMaxThenBumpsEpochWithoutWrapping() .ToList(); var sink = new CapturingSink(); var rot = new SnapshotRotator(channelNumber: 84, source: src, sink: sink); - int requiredPackets = SnapshotPacketBuilder.GetPacketCount( + int requiredSnapshotPackets = SnapshotPacketBuilder.GetPacketCount( SnapshotPacketBuilder.DefaultPacketBufferSize, 700, 300); - Assert.True(requiredPackets > 1); + int requiredPackets = requiredSnapshotPackets + 1; + Assert.True(requiredSnapshotPackets > 1); rot.CreateTestProbe().SetSequence( version: 7, number: uint.MaxValue - (uint)requiredPackets); @@ -250,7 +307,7 @@ public void MultiPacketAtSequenceBoundary_UsesMaxThenBumpsEpochWithoutWrapping() Assert.Equal(requiredPackets, firstPublishPackets); Assert.Equal((ushort)7, rot.SequenceVersion); Assert.Equal(uint.MaxValue, rot.SequenceNumber); - for (int i = 0; i < requiredPackets; i++) + for (int i = 0; i < requiredSnapshotPackets; i++) { ref readonly var header = ref MemoryMarshal.AsRef( sink.Packets[i].AsSpan(0, PacketHeaderSize)); @@ -258,13 +315,17 @@ public void MultiPacketAtSequenceBoundary_UsesMaxThenBumpsEpochWithoutWrapping() Assert.Equal(uint.MaxValue - (uint)requiredPackets + (uint)i + 1, header.SequenceNumber); Assert.NotEqual(0u, header.SequenceNumber); } + ref readonly var firstStatusHeader = ref MemoryMarshal.AsRef( + sink.Packets[requiredSnapshotPackets].AsSpan(0, PacketHeaderSize)); + Assert.Equal((ushort)7, firstStatusHeader.SequenceVersion); + Assert.Equal(uint.MaxValue, firstStatusHeader.SequenceNumber); int secondPublishPackets = rot.PublishNext(incrementalSequenceVersion: 55); Assert.Equal(requiredPackets, secondPublishPackets); Assert.Equal((ushort)8, rot.SequenceVersion); Assert.Equal((uint)requiredPackets, rot.SequenceNumber); - for (int i = 0; i < requiredPackets; i++) + for (int i = 0; i < requiredSnapshotPackets; i++) { ref readonly var header = ref MemoryMarshal.AsRef( sink.Packets[requiredPackets + i].AsSpan(0, PacketHeaderSize)); @@ -272,6 +333,10 @@ public void MultiPacketAtSequenceBoundary_UsesMaxThenBumpsEpochWithoutWrapping() Assert.Equal((uint)(i + 1), header.SequenceNumber); Assert.NotEqual(0u, header.SequenceNumber); } + ref readonly var secondStatusHeader = ref MemoryMarshal.AsRef( + sink.Packets[(requiredPackets * 2) - 1].AsSpan(0, PacketHeaderSize)); + Assert.Equal((ushort)8, secondStatusHeader.SequenceVersion); + Assert.Equal((uint)requiredPackets, secondStatusHeader.SequenceNumber); Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( sink.Packets[requiredPackets].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var snapshotHeader)); @@ -308,7 +373,7 @@ public void BumpSequenceVersion_DoesNotChangeIncrementalLastSequenceVersion() rot.PublishNext(incrementalSequenceVersion: 9); rot.PublishNext(incrementalSequenceVersion: 9); - Assert.Equal(2u, rot.SequenceNumber); + Assert.Equal(4u, rot.SequenceNumber); Assert.Equal((ushort)1, rot.SequenceVersion); rot.BumpSequenceVersion(); @@ -317,12 +382,11 @@ 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)); + ref readonly var hdr = ref MemoryMarshal.AsRef(sink.Packets[^2].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)); + sink.Packets[^2].AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var snapHdr)); Assert.Equal((ushort)9, snapHdr.Data.LastSequenceVersion); } @@ -345,6 +409,9 @@ public void BookWithRptSeqButNoOrders_PublishesHeaderWithLastRptSeq() Assert.Equal(0u, hdr.Data.TotNumReports); Assert.Equal(5u, hdr.Data.LastRptSeq); Assert.Equal((ushort)1, hdr.Data.LastSequenceVersion); + + var status = ReadSecurityStatusPacket(sink.Packets[1]); + Assert.Equal(5u, status.Data.RptSeq); } } @@ -379,6 +446,14 @@ private sealed class CapturingSink : IUmdfPacketSink public void Publish(byte channelNumber, ReadOnlySpan packet) => Packets.Add(packet.ToArray()); } + private static SecurityStatus_3DataReader ReadSecurityStatusPacket(byte[] packet) + { + Assert.True(SecurityStatus_3Data.TryParse( + packet.AsSpan(FrameOffset, WireOffsets.SecurityStatusBlockLength), + out var status)); + return status; + } + [Fact] public void SnapshotTick_PublishedOnDispatcherThread_AndIsolatedFromIncrementalSink() { @@ -404,9 +479,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); // snapshot header + security status 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 +524,7 @@ public void SnapshotTick_ReadsLiveRestingOrders_FromMatchingEngine() Assert.True(disp.EnqueueSnapshotTick()); Drain(disp); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); var pkt = snapSink.Packets[0]; int frameOff = WireOffsets.PacketHeaderSize + WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize; Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( @@ -462,6 +537,61 @@ public void SnapshotTick_ReadsLiveRestingOrders_FromMatchingEngine() Assert.Equal(disp.SequenceVersion, hdr.Data.LastSequenceVersion); } + [Fact] + public void SnapshotTick_PublishesCurrentSecurityStatusWithoutAdvancingIncrementalRptSeq() + { + var incSink = new CapturingSink(); + var snapSink = new CapturingSink(); + MatchingEngine? engine = null; + var disp = new ChannelDispatcher(channelNumber: 1, + engineFactory: s => { engine = new MatchingEngine(new[] { Petr4 }, s, NullLogger.Instance); return engine; }, + options: new ChannelDispatcherOptions + { + PacketSink = incSink, + Outbound = new ChannelDispatcherTests_RecordingOutbound(), + Logger = NullLogger.Instance, + TimeSource = new FakeNanosTimeSource(1UL), + TradeDate = 1, + }); + var rotator = new SnapshotRotator(channelNumber: 1, + source: new MatchingEngineSnapshotSource(engine!, new[] { Petr }), + sink: snapSink, timeSource: new FakeNanosTimeSource(1UL)); + disp.AttachSnapshotRotator(rotator); + + Assert.True(disp.EnqueueOperatorSetTradingPhase(Petr, TradingPhase.Reserved)); + Drain(disp); + uint activeRptSeq = engine!.GetCurrentRptSeq(Petr); + snapSink.Packets.Clear(); + + Assert.True(disp.EnqueueSnapshotTick()); + Drain(disp); + + var activeStatus = ReadSecurityStatusPacket(snapSink.Packets[1]); + Assert.Equal(SecurityTradingStatus.RESERVED, activeStatus.Data.SecurityTradingStatus); + Assert.Null(activeStatus.Data.SecurityTradingEvent); + Assert.Equal(activeRptSeq, activeStatus.Data.RptSeq); + Assert.Equal(activeRptSeq, engine.GetCurrentRptSeq(Petr)); + + Assert.True(disp.EnqueueOperatorHalt(Petr, HaltReason.RegulatoryHalt, null)); + Drain(disp); + uint haltedRptSeq = engine.GetCurrentRptSeq(Petr); + snapSink.Packets.Clear(); + + Assert.True(disp.EnqueueSnapshotTick()); + Drain(disp); + + var snapshotHeader = snapSink.Packets[0]; + Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( + snapshotHeader.AsSpan(FrameOffset, WireOffsets.SnapHeaderBlockLength), out var hdr)); + Assert.Equal(haltedRptSeq, hdr.Data.LastRptSeq); + + var haltedStatus = ReadSecurityStatusPacket(snapSink.Packets[1]); + Assert.Equal(SecurityTradingStatus.FORBIDDEN, haltedStatus.Data.SecurityTradingStatus); + Assert.Equal(SecurityTradingEvent.SECURITY_STATUS_CHANGE, haltedStatus.Data.SecurityTradingEvent); + Assert.Equal(haltedRptSeq, haltedStatus.Data.RptSeq); + Assert.Equal(haltedRptSeq, engine.GetCurrentRptSeq(Petr)); + } + [Fact] public void SnapshotTick_ReflectsIncrementalVersionBump() { @@ -488,7 +618,7 @@ public void SnapshotTick_ReflectsIncrementalVersionBump() Assert.True(disp.EnqueueSnapshotTick()); Drain(disp); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); ref readonly var packetHeader = ref MemoryMarshal.AsRef( snapSink.Packets[0].AsSpan(0, PacketHeaderSize)); Assert.Equal((ushort)2, packetHeader.SequenceVersion); @@ -538,7 +668,7 @@ public void SnapshotTick_ReflectsRestoredIncrementalVersion() Assert.True(restoredDispatcher.EnqueueSnapshotTick()); Drain(restoredDispatcher); - Assert.Single(snapSink.Packets); + Assert.Equal(2, snapSink.Packets.Count); 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..d04574e 100644 --- a/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs +++ b/tests/B3.Exchange.Persistence.Tests/ChannelDispatcherStartupEpochTests.cs @@ -329,7 +329,10 @@ public async Task RestoreReplay_StartsDurableEpoch_PreservesState_AndOrdersReset Assert.Equal(HaltReason.RegulatoryHalt, restoredHalt.Reason); Assert.True(dispatcher.EnqueueSnapshotTick()); - Assert.True(WaitFor(() => snapshotSink.Packets.Count == 1)); + // Snapshot rotation now emits the book snapshot packet followed by a + // standalone SecurityStatus_3 recovery packet for the instrument + // (issue #583), so a single tick yields 2 packets, not 1. + Assert.True(WaitFor(() => snapshotSink.Packets.Count == 2)); Assert.True(SnapshotFullRefresh_Header_30Data.TryParse( snapshotSink.Packets[0].AsSpan( SnapshotFrameOffset, WireOffsets.SnapHeaderBlockLength), diff --git a/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs b/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs index 38c1534..3765e09 100644 --- a/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs +++ b/tests/B3.Umdf.WireEncoder.Tests/UmdfFrameBuilderTests.cs @@ -1,3 +1,4 @@ +using B3.Umdf.Mbo.Sbe.V16; using B3.Umdf.WireEncoder; namespace B3.Umdf.WireEncoder.Tests; @@ -151,7 +152,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,14 +161,18 @@ public void WriteInstrumentHalted_ReservesAndCommitsCorrectSize() } [Fact] - public void WriteInstrumentHalted_WritesHaltEventByte() + public void WriteInstrumentHalted_WritesForbiddenStatusAndOfficialEventByte() { var sink = MakeSink(); UmdfFrameBuilder.WriteInstrumentHalted(sink, - securityId: 1L, securityTradingStatus: 2, rptSeq: 6u, transactTimeNanos: 0ul); + securityId: 1L, rptSeq: 6u, transactTimeNanos: 0ul); int eventOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBodySecurityTradingEventOffset; + int statusOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + + WireOffsets.SecurityStatusBodySecurityTradingStatusOffset; + Assert.Equal((byte)SecurityTradingStatus.FORBIDDEN, sink.Buffer[statusOffset]); + Assert.Equal((byte)SecurityTradingEvent.SECURITY_STATUS_CHANGE, UmdfFrameBuilder.SecurityTradingEventHalt); Assert.Equal(UmdfFrameBuilder.SecurityTradingEventHalt, sink.Buffer[eventOffset]); } @@ -193,6 +198,7 @@ public void WriteInstrumentResumed_WritesResumeEventByte() int eventOffset = WireOffsets.FramingHeaderSize + WireOffsets.SbeMessageHeaderSize + WireOffsets.SecurityStatusBodySecurityTradingEventOffset; + Assert.Equal((byte)SecurityTradingEvent.SECURITY_REJOINS_SECURITY_GROUP_STATUS, UmdfFrameBuilder.SecurityTradingEventResume); Assert.Equal(UmdfFrameBuilder.SecurityTradingEventResume, sink.Buffer[eventOffset]); }