diff --git a/docs/clustering.md b/docs/clustering.md index aa93626..8a06d2d 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 | diff --git a/scripts/docker_ci_fix.py b/scripts/docker_ci_fix.py index 9ef3d91..e0a097d 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( [ diff --git a/scripts/docker_cluster_ha.py b/scripts/docker_cluster_ha.py index 793c168..d56207e 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", diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index fdeaec9..b4d557d 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(()) diff --git a/src/cluster/media/peer.rs b/src/cluster/media/peer.rs index a5ee7ff..2c95fb1 100644 --- a/src/cluster/media/peer.rs +++ b/src/cluster/media/peer.rs @@ -546,7 +546,9 @@ 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, diff --git a/src/media_output.rs b/src/media_output.rs index c49afe6..1d1a82e 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 diff --git a/src/server.rs b/src/server.rs index 87dfbd0..6decd44 100644 --- a/src/server.rs +++ b/src/server.rs @@ -289,20 +289,50 @@ 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); + // 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()) + && previous != key + && self.routes.get(&previous) == Some(&conn_id) + { + self.routes.remove(&previous); + ended.push(previous); + } + if self.routes.get(&key) != Some(&conn_id) { + 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 } /// Removes and returns the routes whose publisher is no longer among @@ -316,6 +346,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 +2776,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())); @@ -3150,6 +3186,87 @@ 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 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 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;