diff --git a/docs/stream_socket.rst b/docs/stream_socket.rst index 54f6aae5f..9afbd3a24 100644 --- a/docs/stream_socket.rst +++ b/docs/stream_socket.rst @@ -150,7 +150,10 @@ Message types: produced so far, i.e. the sequence the next broadcast EVENT will carry; the snapshot is a state message, not an event in that series, so a consumer seeds its gap tracking from it and must not treat the following EVENT with - the same sequence as a duplicate. The capture-fault edges are emitted by + the same sequence as a duplicate. Both that sequence and the header + ``generation`` are stamped when the consumer connects, so they describe + the moment of connection rather than the last status change. The + capture-fault edges are emitted by zmc once per transition; because the socket survives camera reconnects, ``connection_failed`` is observable exactly when media has stopped. diff --git a/src/zm_stream_socket.cpp b/src/zm_stream_socket.cpp index a0b96db77..98f7e9269 100644 --- a/src/zm_stream_socket.cpp +++ b/src/zm_stream_socket.cpp @@ -365,10 +365,11 @@ void StreamSocket::SendMonitorEvent(std::vector payload) { void StreamSocket::SetSnapshotEvent(std::vector payload) { std::lock_guard lock(mutex_); - // The snapshot is the consumer's authoritative current status on connect; it - // is tagged with the current event sequence as a baseline and not broadcast. - snapshot_ = MakeMessage(MessageType::Event, StreamId::Monitor, 0, - event_sequence_, 0, std::move(payload), true); + // The snapshot is the consumer's authoritative current status on connect. It + // is not broadcast, and only the body is kept: AcceptClient frames it, so + // the header carries the generation and the event-sequence baseline that are + // current when the consumer connects rather than when the status last moved. + snapshot_payload_ = std::move(payload); } void StreamSocket::InvalidateKeyframe() { @@ -549,8 +550,11 @@ void StreamSocket::AcceptClient() { EnqueueLocked(*client, hello_audio_); if (hello_video_) EnqueueLocked(*client, hello_video_); - if (snapshot_) - EnqueueLocked(*client, snapshot_); + if (!snapshot_payload_.empty()) { + EnqueueLocked(*client, MakeMessage(MessageType::Event, StreamId::Monitor, 0, + event_sequence_, 0, + std::vector(snapshot_payload_), true)); + } if (keyframe_) EnqueueLocked(*client, keyframe_); diff --git a/src/zm_stream_socket.h b/src/zm_stream_socket.h index a8f127d4e..1da825197 100644 --- a/src/zm_stream_socket.h +++ b/src/zm_stream_socket.h @@ -106,7 +106,9 @@ class StreamSocket { // Cache the current-status snapshot replayed to each new consumer on connect // (the events analogue of the cached keyframe). payload is a pre-built EVENT - // body of code kEventSnapshot. Caching only; does not broadcast. + // body of code kEventSnapshot. Caching only; does not broadcast. The header + // is framed when a consumer connects, so it always carries the generation + // and event sequence in effect at that moment. void SetSnapshotEvent(std::vector payload); // Drop the cached keyframe. Called when the capture source closes so a @@ -195,7 +197,7 @@ class StreamSocket { MessagePtr hello_video_; MessagePtr hello_audio_; MessagePtr keyframe_; - MessagePtr snapshot_; // current-status EVENT, replayed on connect + std::vector snapshot_payload_; // current-status EVENT body, framed per connect std::vector hello_video_payload_; std::vector hello_audio_payload_; uint32_t generation_ = 0; diff --git a/tests/zm_stream_socket.cpp b/tests/zm_stream_socket.cpp index 7146a9e22..4fa904bfb 100644 --- a/tests/zm_stream_socket.cpp +++ b/tests/zm_stream_socket.cpp @@ -809,3 +809,33 @@ TEST_CASE("StreamSocket::SetStreams drops a video stream the source no longer ha server.Stop(); } + +TEST_CASE("StreamSocket frames the snapshot with the generation and event sequence at connect", "[stream_socket]") { + StreamSocket server(1, kSockPath); + REQUIRE(server.Start()); + + // Snapshot cached while generation 0 and no events produced yet + server.SetSnapshotEvent(make_state_event(kEventSnapshot, 0, 0, "IDLE")); + + // Then the stream is reconfigured and two events go by with nobody listening + codec_parameters_ptr video = make_h264_parameters(); + server.SetVideoParams(video.get(), {0, 0}); + video->width = 1280; + server.SetVideoParams(video.get(), {0, 0}); + server.SendMonitorEvent(make_state_event(kEventConnectionFailed, 0, 0, "")); + server.SendMonitorEvent(make_state_event(kEventConnectionRestored, 0, 0, "")); + + TestClient client; + REQUIRE(client.Connect()); + ReceivedMessage message; + REQUIRE(client.ReadMessage(message)); // video HELLO + REQUIRE(message.header.generation == 1); + + REQUIRE(client.ReadMessage(message)); + REQUIRE(message.header.type == static_cast(MessageType::Event)); + // Header reflects the moment of connect, not the moment the status last moved + REQUIRE(message.header.generation == 1); + REQUIRE(message.header.sequence == 2); + + server.Stop(); +}