Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/clustering.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
4 changes: 3 additions & 1 deletion scripts/docker_ci_fix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
[
Expand Down
4 changes: 3 additions & 1 deletion scripts/docker_cluster_ha.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
18 changes: 17 additions & 1 deletion src/cluster/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(())
Expand Down
4 changes: 3 additions & 1 deletion src/cluster/media/peer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -546,7 +546,9 @@ pub async fn accept_auth_negotiated<S: AsyncRead + AsyncWrite + Unpin>(
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,
Expand Down
4 changes: 4 additions & 0 deletions src/media_output.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
127 changes: 122 additions & 5 deletions src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64, (String, String)>,
}

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)> {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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
Expand All @@ -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
}
}
Expand Down Expand Up @@ -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<u64> = 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()));
Expand Down Expand Up @@ -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;
Expand Down
Loading