From 08483f007f5a4becab5a04751e23534448712615 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:00:30 +0200 Subject: [PATCH 01/12] fix(bug-hunter): S4-1 - docker_ci_fix.py propagates cargo check exit status --- scripts/docker_ci_fix.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/scripts/docker_ci_fix.py b/scripts/docker_ci_fix.py index 9ef3d913..e0a097da 100644 --- a/scripts/docker_ci_fix.py +++ b/scripts/docker_ci_fix.py @@ -39,7 +39,9 @@ def main() -> int: cargo fmt cargo generate-lockfile cargo check --features test-support 2>&1 | tee /src/ci-fix.log | tail -n 60 -echo EXIT=$? >> /src/ci-fix.log +EXIT=${PIPESTATUS[0]} +echo EXIT=$EXIT >> /src/ci-fix.log +exit $EXIT """ r = subprocess.run( [ From 2ecab91dc58951b8b28461bd11095fa54c779342 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:00:30 +0200 Subject: [PATCH 02/12] fix(bug-hunter): S4-2 - docker_cluster_ha.py propagates cargo test exit status --- scripts/docker_cluster_ha.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/scripts/docker_cluster_ha.py b/scripts/docker_cluster_ha.py index 793c1682..d56207e0 100644 --- a/scripts/docker_cluster_ha.py +++ b/scripts/docker_cluster_ha.py @@ -58,9 +58,11 @@ def main() -> int: apt-get update -qq apt-get install -y -qq pkg-config libssl-dev >/dev/null cargo test --features cluster,test-support --test cluster_ha -- --nocapture > /src/cargo-cluster-ha.log 2>&1 -echo EXIT=$? >> /src/cargo-cluster-ha.log +EXIT=$? +echo EXIT=$EXIT >> /src/cargo-cluster-ha.log grep -E '^(test |failures:|error|EXIT=|thread )' /src/cargo-cluster-ha.log | tail -n 80 tail -n 30 /src/cargo-cluster-ha.log +exit $EXIT """ cmd = [ "docker", From 22d36b77caad4a05ce6b22fd99eba96fdb737c3f Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:05:56 +0200 Subject: [PATCH 03/12] fix(bug-hunter): S1-1/S1-2 - drop removed node from health tracker; keep learner publish gate on resume --- src/cluster/manager.rs | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index fdeaec9d..b4d557db 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -531,7 +531,22 @@ impl ClusterManager { match join_action { JoinReseedAction::ResumeExisting => { mgr.refresh_topology_from_any(join_addr).await?; - health.set_local(NodeHealthState::Ready); + // A resumed learner must stay publish-gated; only promotion + // to voter clears the gate (reconciled in + // `prune_topology_to_membership`). The persisted membership + // is the only source available before the metrics loop + // starts; a missing/empty one keeps the learner gate. + let is_voter = mgr + .state_machine + .last_membership() + .membership() + .voter_ids() + .any(|id| id == config.node_id); + health.set_local(if is_voter { + NodeHealthState::Ready + } else { + NodeHealthState::Learner + }); crate::log_info!( "Cluster: resuming existing member node {} (CLUSTER_JOIN set but local raft state present)", config.node_id @@ -1004,6 +1019,7 @@ impl ClusterManager { } self.media.disconnect_peer(node_id); self.meta.remove(node_id); + self.health.remove(node_id); self.network.nodes.write().remove(&node_id); crate::log_info!("Cluster: removed node {node_id}"); Ok(()) From f56eb5f2e1b85890e00782f9c3fd128a723fe759 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:05:56 +0200 Subject: [PATCH 04/12] fix(bug-hunter): S2-1 - bound post-auth media hello read with AUTH_TIMEOUT --- src/cluster/media/peer.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/cluster/media/peer.rs b/src/cluster/media/peer.rs index a5ee7ff0..7a847888 100644 --- a/src/cluster/media/peer.rs +++ b/src/cluster/media/peer.rs @@ -546,7 +546,11 @@ pub async fn accept_auth_negotiated( clear_cluster_auth_failures(peer); write_media_frame(stream, &MediaMessage::AuthOk).await?; - let hello = read_media_frame(stream).await?; + let hello = tokio::time::timeout(AUTH_TIMEOUT, read_media_frame(stream)) + .await + .map_err(|_| { + std::io::Error::new(std::io::ErrorKind::TimedOut, "media hello timeout") + })??; let MediaMessage::Hello { version, node_id: hello_id, From a5be5a6b4873d7f12de5169d6a94bb770e67f827 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:15:50 +0200 Subject: [PATCH 05/12] style(bug-hunter): apply rustfmt to S2-1 timeout closure (CI fmt gate) --- src/cluster/media/peer.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/cluster/media/peer.rs b/src/cluster/media/peer.rs index 7a847888..2c95fb1a 100644 --- a/src/cluster/media/peer.rs +++ b/src/cluster/media/peer.rs @@ -548,9 +548,7 @@ pub async fn accept_auth_negotiated( let hello = tokio::time::timeout(AUTH_TIMEOUT, read_media_frame(stream)) .await - .map_err(|_| { - std::io::Error::new(std::io::ErrorKind::TimedOut, "media hello timeout") - })??; + .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "media hello timeout"))??; let MediaMessage::Hello { version, node_id: hello_id, From baace497835b1c8282240aeb992368caa583f084 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:22:17 +0200 Subject: [PATCH 06/12] fix(bug-hunter): S3-2 - end previous route on same-connection publish rename --- src/server.rs | 33 ++++++++++++++++++++++++++++----- 1 file changed, 28 insertions(+), 5 deletions(-) diff --git a/src/server.rs b/src/server.rs index 87dfbd03..5bcf0ff8 100644 --- a/src/server.rs +++ b/src/server.rs @@ -289,20 +289,37 @@ enum ShardRelayMsg { #[derive(Default)] struct ExportedRoutes { routes: HashMap<(String, String), u64>, + /// The route each connection currently exports, so a publish that + /// switches route on the same connection ends its previous route + /// instead of leaving the other shards' inject claim on it until the + /// stale timeout. + conn_routes: HashMap, } impl ExportedRoutes { - fn record(&mut self, frames: &[librtmp2::RelayFrame]) { + /// Records `frames` and returns the routes that ended because their + /// connection switched to a different route (publish rename). + fn record(&mut self, frames: &[librtmp2::RelayFrame]) -> Vec<(String, String)> { + let mut ended = Vec::new(); for frame in frames { // Frames injected from elsewhere are never re-broadcast. if librtmp2::server::is_external_publisher_id(frame.publisher_conn_id) { continue; } + let conn_id = frame.publisher_conn_id; let key = (frame.app.clone(), frame.stream_name.clone()); - if self.routes.get(&key) != Some(&frame.publisher_conn_id) { - self.routes.insert(key, frame.publisher_conn_id); + if self.conn_routes.get(&conn_id) == Some(&key) { + continue; + } + if let Some(previous) = self.conn_routes.insert(conn_id, key.clone()) { + self.routes.remove(&previous); + ended.push(previous); + } + if self.routes.get(&key) != Some(&conn_id) { + self.routes.insert(key, conn_id); } } + ended } /// Removes and returns the routes whose publisher is no longer among @@ -316,6 +333,11 @@ impl ExportedRoutes { } live }); + // Drop per-connection entries with their route so a later republish + // on the same connection is recorded (and ended) again. + self.conn_routes.retain(|conn_id, route| { + publishing.contains(conn_id) && self.routes.get(route) == Some(conn_id) + }); ended } } @@ -2741,14 +2763,15 @@ impl ShardLoop { let Some(txs) = self.relay_txs.clone() else { return; }; - self.exported_routes.record(frames); + let mut ended = self.exported_routes.record(frames); let publishing: HashSet = server .connections .iter() .filter(|c| c.state == librtmp2::types::ConnState::Publishing) .map(|c| c.conn_id) .collect(); - for (app, stream_name) in self.exported_routes.take_ended(&publishing) { + ended.extend(self.exported_routes.take_ended(&publishing)); + for (app, stream_name) in ended { for i in (0..txs.len()).filter(|&i| i != self.index) { self.pending_route_ends .push((i, app.clone(), stream_name.clone())); From 85125cef6844a2877c7b3a020ac3c593ce612b9b Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:22:18 +0200 Subject: [PATCH 07/12] fix(bug-hunter): S3-4 - mark stalled sink failed so the monitor kills the child --- src/media_output.rs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/media_output.rs b/src/media_output.rs index c49afe66..1d1a82e2 100644 --- a/src/media_output.rs +++ b/src/media_output.rs @@ -754,6 +754,10 @@ impl SinkSender { ); } } else { + // The monitor only kills the FFmpeg child once `failed` is set; + // without it a worker blocked on a stalled child would leak both + // the child and the monitor thread. + self.failed.store(true, Ordering::Release); crate::log_warn!( "Media output '{}' worker did not stop within 5s; continuing shutdown", self.label From 52b2af5ea5765cb7257b14b0603f31c4d084461a Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:33:21 +0200 Subject: [PATCH 08/12] docs(bug-hunter): CLUSTER_SECRET is 32-256 ASCII chars in clustering.md (CodeRabbit review finding) --- docs/clustering.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/clustering.md b/docs/clustering.md index aa936266..8a06d2d4 100644 --- a/docs/clustering.md +++ b/docs/clustering.md @@ -62,7 +62,7 @@ Set in `.env` or via `LRTMP2_CLUSTER_*` process overrides: | `CLUSTER_BOOTSTRAP` | `false` | First voter; mutually exclusive with JOIN | | `CLUSTER_JOIN` | — | Address of an existing control peer | | `CLUSTER_JOIN_PROOF` | — | Required for a fresh join; mint via authenticated `POST /api/v1/cluster/join-proof` | -| `CLUSTER_SECRET` | — | Shared secret (≥16 chars); never logged | +| `CLUSTER_SECRET` | — | Shared secret (32–256 ASCII letters, digits, `-` or `_`); never logged | | `CLUSTER_TLS_ENABLED` | `false` | mTLS for control/media when true | | `CLUSTER_TLS_CERT_FILE` / `KEY` / `CA` | — | Required if TLS enabled | | `CLUSTER_HEARTBEAT_MS` / `CLUSTER_HEARTBEAT_INTERVAL_MS` | `500` | Peer heartbeat interval | From 530c17470e3ae94bb3983820d297a0af49a208fa Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:43:23 +0200 Subject: [PATCH 09/12] fix(bug-hunter): S-RESCAN-1 - repair stale route overwrites without spurious RouteEnded (+regression test) --- src/server.rs | 43 ++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 40 insertions(+), 3 deletions(-) diff --git a/src/server.rs b/src/server.rs index 5bcf0ff8..765c5261 100644 --- a/src/server.rs +++ b/src/server.rs @@ -308,12 +308,21 @@ impl ExportedRoutes { } let conn_id = frame.publisher_conn_id; let key = (frame.app.clone(), frame.stream_name.clone()); - if self.conn_routes.get(&conn_id) == Some(&key) { + // Healthy fast path: this connection already owns this route. + // Anything else must repair the map like the base code did — a + // stale frame from a dead publisher may have overwritten the + // owner, and skipping the repair would emit a spurious RouteEnded + // for a route that is still being published. + if self.conn_routes.get(&conn_id) == Some(&key) + && self.routes.get(&key) == Some(&conn_id) + { continue; } if let Some(previous) = self.conn_routes.insert(conn_id, key.clone()) { - self.routes.remove(&previous); - ended.push(previous); + if previous != key { + self.routes.remove(&previous); + ended.push(previous); + } } if self.routes.get(&key) != Some(&conn_id) { self.routes.insert(key, conn_id); @@ -3173,6 +3182,34 @@ mod tests { assert_eq!(ended, vec![("live".to_string(), "a".to_string())]); } + #[test] + fn exported_routes_repair_stale_overwrites_without_spurious_end() { + use super::ExportedRoutes; + let frame = |conn_id: u64, stream: &str| librtmp2::RelayFrame { + frame_type: librtmp2::types::FrameType::Video, + timestamp: 0, + payload: Vec::new(), + cache_payload: None, + app: "live".to_string(), + stream_name: stream.to_string(), + publisher_conn_id: conn_id, + }; + let mut routes = ExportedRoutes::default(); + routes.record(&[frame(1, "a")]); + // A stale frame from a dead publisher overwrites the route owner. + routes.record(&[frame(2, "a")]); + // The live owner's next frame must repair the map without ending the + // route (regression: the fast path used to skip the repair). + assert!(routes.record(&[frame(1, "a")]).is_empty()); + assert!(routes.take_ended(&HashSet::from([1])).is_empty()); + // Renaming the same connection still ends the previous route. + assert_eq!( + routes.record(&[frame(1, "b")]), + vec![("live".to_string(), "a".to_string())] + ); + assert!(routes.take_ended(&HashSet::from([1])).is_empty()); + } + #[test] fn tls_shards_split_the_connection_cap_statically() { use super::shard_connection_cap; From f001f9952ca1c69e04ae6f35aada3ebe16fd9bd4 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 21:52:40 +0200 Subject: [PATCH 10/12] fix(bug-hunter): S-Codex-1 - reconcile route ends against the final batch map (+regression test) --- src/server.rs | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/src/server.rs b/src/server.rs index 765c5261..2f20eaa3 100644 --- a/src/server.rs +++ b/src/server.rs @@ -328,6 +328,9 @@ impl ExportedRoutes { self.routes.insert(key, conn_id); } } + // A route that was renamed away and then re-claimed later in the same + // batch is still live: do not announce its end. + ended.retain(|key| !self.routes.contains_key(key)); ended } @@ -3210,6 +3213,34 @@ mod tests { assert!(routes.take_ended(&HashSet::from([1])).is_empty()); } + #[test] + fn exported_routes_do_not_end_routes_reclaimed_in_same_batch() { + use super::ExportedRoutes; + let frame = |conn_id: u64, stream: &str| librtmp2::RelayFrame { + frame_type: librtmp2::types::FrameType::Video, + timestamp: 0, + payload: Vec::new(), + cache_payload: None, + app: "live".to_string(), + stream_name: stream.to_string(), + publisher_conn_id: conn_id, + }; + let mut routes = ExportedRoutes::default(); + // Conn 1 renames a -> b, then conn 2 claims a again, all in one batch: + // a must not be announced as ended (it is live under conn 2). + assert!( + routes + .record(&[frame(1, "a"), frame(1, "b"), frame(2, "a")]) + .is_empty() + ); + assert!(routes.take_ended(&HashSet::from([1, 2])).is_empty()); + // Once conn 1 stops publishing, only its current route b ends. + assert_eq!( + routes.take_ended(&HashSet::from([2])), + vec![("live".to_string(), "b".to_string())] + ); + } + #[test] fn tls_shards_split_the_connection_cap_statically() { use super::shard_connection_cap; From ac46e4e633970dbeae02716eb24a99d756216b54 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 22:07:48 +0200 Subject: [PATCH 11/12] style(bug-hunter): collapse if-let in ExportedRoutes::record (CI clippy 1.99 collapsible_if) --- src/server.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/src/server.rs b/src/server.rs index 2f20eaa3..c33bb623 100644 --- a/src/server.rs +++ b/src/server.rs @@ -318,11 +318,11 @@ impl ExportedRoutes { { continue; } - if let Some(previous) = self.conn_routes.insert(conn_id, key.clone()) { - if previous != key { - self.routes.remove(&previous); - ended.push(previous); - } + if let Some(previous) = self.conn_routes.insert(conn_id, key.clone()) + && previous != key + { + self.routes.remove(&previous); + ended.push(previous); } if self.routes.get(&key) != Some(&conn_id) { self.routes.insert(key, conn_id); From e42f4f36bf1f1108173d145f591223621d58d186 Mon Sep 17 00:00:00 2001 From: AlexanderWagnerDev Date: Sun, 4 Oct 2026 22:17:57 +0200 Subject: [PATCH 12/12] fix(bug-hunter): S-CR-2 - only end the previous route if this connection still owns it (+regression test) --- src/server.rs | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) diff --git a/src/server.rs b/src/server.rs index c33bb623..6decd446 100644 --- a/src/server.rs +++ b/src/server.rs @@ -320,6 +320,7 @@ impl ExportedRoutes { } if let Some(previous) = self.conn_routes.insert(conn_id, key.clone()) && previous != key + && self.routes.get(&previous) == Some(&conn_id) { self.routes.remove(&previous); ended.push(previous); @@ -3241,6 +3242,31 @@ mod tests { ); } + #[test] + fn exported_routes_keep_reclaimed_route_when_old_owner_moves_away() { + use super::ExportedRoutes; + let frame = |conn_id: u64, stream: &str| librtmp2::RelayFrame { + frame_type: librtmp2::types::FrameType::Video, + timestamp: 0, + payload: Vec::new(), + cache_payload: None, + app: "live".to_string(), + stream_name: stream.to_string(), + publisher_conn_id: conn_id, + }; + let mut routes = ExportedRoutes::default(); + routes.record(&[frame(1, "a")]); + // Conn 2 reclaims a, then conn 1 moves to b in the same batch: conn 1 + // no longer owns a, so its move must not end a or drop conn 2's entry. + assert!(routes.record(&[frame(2, "a"), frame(1, "b")]).is_empty()); + assert!(routes.take_ended(&HashSet::from([1, 2])).is_empty()); + // Conn 2 stopping ends a; conn 1's b ended separately above. + assert_eq!( + routes.take_ended(&HashSet::from([1])), + vec![("live".to_string(), "a".to_string())] + ); + } + #[test] fn tls_shards_split_the_connection_cap_statically() { use super::shard_connection_cap;