From 83eec4108fdd037b3aafc0d081bbde66b0b5d20f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 17 Aug 2026 02:09:02 +0000 Subject: [PATCH 1/8] fix(security): harden cluster admin_proof replay and plaintext heartbeat counts - Track recently accepted admin_proof digests per node so captured ClientWrite/membership proofs cannot be replayed while the API token is unchanged. - Ignore heartbeat-supplied session counts on plaintext clusters; only mTLS-bound peers may update remote play/publish load used for admission. Co-authored-by: Alexander Wagner --- src/cluster/manager.rs | 74 +++++++++++++++++++++++++++++++-------- tests/cluster_security.rs | 12 +++++++ 2 files changed, 71 insertions(+), 15 deletions(-) diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 21fcad4..4b3e3a2 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -1,7 +1,8 @@ //! ClusterManager: Raft lifecycle, durable mutations, media hub. +use std::collections::HashMap; use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, Instant}; use openraft::BasicNode; use openraft::ChangeMembers; @@ -45,6 +46,8 @@ const CLUSTER_ID_SETTING: &str = "cluster_id"; /// seed actually landing. const BOOTSTRAP_SEEDED_SETTING: &str = "bootstrap_seeded"; const RAFT_WRITE_TIMEOUT_MSG: &str = "raft write timed out"; +/// Reject re-sent `admin_proof` values captured from the control plane. +const ADMIN_PROOF_REPLAY_TTL: Duration = Duration::from_secs(600); /// Resolve control/media addresses from a heartbeat payload. /// @@ -127,6 +130,9 @@ pub struct ClusterManager { /// reported as replicas with no data behind them. standby_subs: Mutex>, state_machine: SqliteStateMachine, + /// Recently accepted `admin_proof` digests — blocks captured ClientWrite / + /// membership proofs from being replayed while the API token is unchanged. + used_admin_proofs: Mutex>, } /// Shared with AppState so Raft applies can mark deleted streams / revoked viewers @@ -291,6 +297,7 @@ impl ClusterManager { pending_node_cleanups: Mutex::new(std::collections::HashSet::new()), standby_subs: Mutex::new(std::collections::HashMap::new()), state_machine: sm_handle.clone(), + used_admin_proofs: Mutex::new(HashMap::new()), }); // Control + media listeners @@ -341,18 +348,24 @@ impl ClusterManager { }); } } - counts_hb - .peer_session_counts - .lock() - .insert(info.node_id, (info.publishers, info.players)); - counts_hb - .peer_stream_players - .lock() - .insert(info.node_id, info.stream_players.into_iter().collect()); - counts_hb - .peer_viewer_players - .lock() - .insert(info.node_id, info.viewer_players.into_iter().collect()); + // Session counts are only trustworthy when mTLS binds the + // authenticated peer to a certificate identity. On plaintext + // clusters any `CLUSTER_SECRET` holder can authenticate as + // another member and lie about remote play/publish load. + if counts_hb.config.tls_enabled { + counts_hb + .peer_session_counts + .lock() + .insert(info.node_id, (info.publishers, info.players)); + counts_hb + .peer_stream_players + .lock() + .insert(info.node_id, info.stream_players.into_iter().collect()); + counts_hb + .peer_viewer_players + .lock() + .insert(info.node_id, info.viewer_players.into_iter().collect()); + } }); let mgr_admin = Arc::clone(&mgr); let on_admin: Arc< @@ -908,7 +921,17 @@ impl ClusterManager { Ok(crate::cluster::security::admin_proof(&token, payload)) } + fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) { + guard.retain(|_, seen_at| { + now.checked_duration_since(*seen_at) + .is_none_or(|age| age < ADMIN_PROOF_REPLAY_TTL) + }); + } + fn verify_admin_proof(&self, proof: &str, payload: &str) -> bool { + if proof.is_empty() { + return false; + } let binding = self.session_hooks.lock(); let Some(hooks) = binding.as_ref() else { return false; @@ -917,10 +940,20 @@ impl ClusterManager { if token.is_empty() { return false; } - crate::cluster::security::secrets_equal( + if !crate::cluster::security::secrets_equal( &crate::cluster::security::admin_proof(&token, payload), proof, - ) + ) { + return false; + } + let mut used = self.used_admin_proofs.lock(); + let now = Instant::now(); + Self::purge_expired_admin_proofs(&mut used, now); + if used.contains_key(proof) { + return false; + } + used.insert(proof.to_string(), now); + true } async fn forward_admin( @@ -2378,4 +2411,15 @@ mod heartbeat_routing_tests { assert_eq!(ctrl.as_deref(), Some("new-ctrl:1940")); assert_eq!(media.as_deref(), Some("new-media:1941")); } + + #[test] + fn purge_expired_admin_proofs_drops_stale_entries() { + let mut map = HashMap::new(); + map.insert( + "stale-proof".to_string(), + Instant::now() - ADMIN_PROOF_REPLAY_TTL - Duration::from_secs(1), + ); + ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); + assert!(map.is_empty()); + } } diff --git a/tests/cluster_security.rs b/tests/cluster_security.rs index 52459b3..c74f88d 100644 --- a/tests/cluster_security.rs +++ b/tests/cluster_security.rs @@ -78,6 +78,18 @@ fn secrets_equal_rejects_length_pairs_that_overflow_u8() { assert!(!secrets_equal("", &a)); } +#[test] +fn admin_proof_is_deterministic_so_replay_cache_is_required() { + let token = "api-token-for-tests-only"; + let payload = r#"{"SetApiToken":{"token":"stolen"}}"#; + let a = admin_proof(token, payload); + let b = admin_proof(token, payload); + assert_eq!( + a, b, + "identical payloads must yield identical proofs so nodes need replay tracking" + ); +} + #[test] fn tls_identity_requires_cert_marker_when_tls_on() { let mut der = Vec::new(); From 93532f5f5f471703b6821e70c5e5c602b5d6993a Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 06:53:56 +0200 Subject: [PATCH 2/8] fix(cluster): preserve plaintext session accounting --- src/cluster/manager.rs | 35 +++++++++++++++++------------------ 1 file changed, 17 insertions(+), 18 deletions(-) diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 4b3e3a2..726fbf7 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -348,24 +348,23 @@ impl ClusterManager { }); } } - // Session counts are only trustworthy when mTLS binds the - // authenticated peer to a certificate identity. On plaintext - // clusters any `CLUSTER_SECRET` holder can authenticate as - // another member and lie about remote play/publish load. - if counts_hb.config.tls_enabled { - counts_hb - .peer_session_counts - .lock() - .insert(info.node_id, (info.publishers, info.players)); - counts_hb - .peer_stream_players - .lock() - .insert(info.node_id, info.stream_players.into_iter().collect()); - counts_hb - .peer_viewer_players - .lock() - .insert(info.node_id, info.viewer_players.into_iter().collect()); - } + // Plaintext clustering explicitly uses possession of CLUSTER_SECRET + // plus Raft membership as its trust boundary; mTLS strengthens that + // boundary with per-node certificate identity. Keep session counts in + // both modes because cluster-wide stats, viewer limits and drain/delete + // accounting depend on these heartbeat caches. + counts_hb + .peer_session_counts + .lock() + .insert(info.node_id, (info.publishers, info.players)); + counts_hb + .peer_stream_players + .lock() + .insert(info.node_id, info.stream_players.into_iter().collect()); + counts_hb + .peer_viewer_players + .lock() + .insert(info.node_id, info.viewer_players.into_iter().collect()); }); let mgr_admin = Arc::clone(&mgr); let on_admin: Arc< From af27deb3efd9d2570dd762955656f3d691728a6a Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Wed, 19 Aug 2026 04:54:08 +0000 Subject: [PATCH 3/8] Update Cargo.lock --- Cargo.lock | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 5777f7f..f781254 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1769,9 +1769,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.16" +version = "0.11.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" +checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" dependencies = [ "aws-lc-rs", "bytes", @@ -1912,18 +1912,18 @@ dependencies = [ [[package]] name = "ref-cast" -version = "1.0.26" +version = "1.0.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "216e8f773d7923bcba9ceb86a86c93cabb3903a11872fc3f138c49630e50b96d" +checksum = "7e440fb4e4b4147295338efb76001ab9e4efc0e5839df2c47fc5ac2381d365c3" dependencies = [ "ref-cast-impl", ] [[package]] name = "ref-cast-impl" -version = "1.0.26" +version = "1.0.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2c9283685feec7d69af75fb0e858d5e7378f33fe4fc699383b2916ab9273e03c" +checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" dependencies = [ "proc-macro2", "quote", @@ -3324,9 +3324,9 @@ dependencies = [ [[package]] name = "zerovec-derive" -version = "0.11.4" +version = "0.11.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "47402523226a02bfe5230160dc3ccc089aa6f6f19e7fcbb4e6f824bbb1b4aa62" +checksum = "9f212a141d820099d57ffafb9569be9617a6f27d3dc881fbee8fb56642f917a9" dependencies = [ "proc-macro2", "quote", From 3bcaeeb318cc1d47036e00025ca354bcd62f5f33 Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 07:06:15 +0200 Subject: [PATCH 4/8] chore: apply PR 169 Cursor finding fix --- .github/workflows/fix-pr169-cursor.yml | 251 +++++++++++++++++++++++++ 1 file changed, 251 insertions(+) create mode 100644 .github/workflows/fix-pr169-cursor.yml diff --git a/.github/workflows/fix-pr169-cursor.yml b/.github/workflows/fix-pr169-cursor.yml new file mode 100644 index 0000000..0c759b4 --- /dev/null +++ b/.github/workflows/fix-pr169-cursor.yml @@ -0,0 +1,251 @@ +name: Apply PR 169 Cursor finding fix + +on: + push: + branches: + - cursor/application-security-review-2fc5 + paths: + - .github/workflows/fix-pr169-cursor.yml + +permissions: + contents: write + +jobs: + apply-fix: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + with: + ref: cursor/application-security-review-2fc5 + fetch-depth: 0 + + - name: Apply retry-safe admin proof fix + shell: bash + run: | + python3 - <<'PY' + from pathlib import Path + + path = Path('src/cluster/manager.rs') + text = path.read_text() + + replacements = [] + + replacements.append(( + ''' let proof = crate::cluster::security::admin_proof(&token, &req_str);''', + ''' let proof = ClusterManager::mint_fresh_admin_proof(&token, &req_str);''' + )) + + replacements.append(( + ''' Ok(crate::cluster::security::admin_proof(&token, payload)) + } + + fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) {''', + ''' Ok(Self::mint_fresh_admin_proof(&token, payload)) + } + + /// Mint a unique proof for each admin attempt. The random nonce is signed + /// together with the canonical payload, so repeated legitimate operations + /// do not reuse the same replay-cache key. + fn mint_fresh_admin_proof(api_token: &str, payload: &str) -> String { + let nonce = hex::encode(crate::cluster::security::auth_nonce()); + let signed_payload = format!("{nonce}:{payload}"); + let mac = crate::cluster::security::admin_proof(api_token, &signed_payload); + format!("{nonce}.{mac}") + } + + fn admin_proof_signature_valid(api_token: &str, payload: &str, proof: &str) -> bool { + let Some((nonce, mac)) = proof.split_once('.') else { + return false; + }; + // auth_nonce() is 16 bytes => exactly 32 lowercase/uppercase hex chars. + if nonce.len() != 32 + || hex::decode(nonce) + .map(|decoded| decoded.len() != 16) + .unwrap_or(true) + { + return false; + } + let signed_payload = format!("{nonce}:{payload}"); + crate::cluster::security::secrets_equal( + &crate::cluster::security::admin_proof(api_token, &signed_payload), + mac, + ) + } + + fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) {''' + )) + + replacements.append(( + ''' if !crate::cluster::security::secrets_equal( + &crate::cluster::security::admin_proof(&token, payload), + proof, + ) { + return false; + } + let mut used = self.used_admin_proofs.lock(); + let now = Instant::now(); + Self::purge_expired_admin_proofs(&mut used, now); + if used.contains_key(proof) { + return false; + } + used.insert(proof.to_string(), now); + true + } + + async fn forward_admin(''', + ''' if !Self::admin_proof_signature_valid(&token, payload, proof) { + return false; + } + let mut used = self.used_admin_proofs.lock(); + let now = Instant::now(); + Self::purge_expired_admin_proofs(&mut used, now); + if used.contains_key(proof) { + return false; + } + // Join proofs must remain retryable if add_learner/forwarding fails. + // accept_join() consumes them only after the join actually succeeds. + if !payload.starts_with("Join:") { + used.insert(proof.to_string(), now); + } + true + } + + fn consume_admin_proof(&self, proof: &str) { + let mut used = self.used_admin_proofs.lock(); + let now = Instant::now(); + Self::purge_expired_admin_proofs(&mut used, now); + used.insert(proof.to_string(), now); + } + + async fn forward_admin(''' + )) + + replacements.append(( + ''' Err(RaftError::APIError(ClientWriteError::ForwardToLeader(ftl))) => { + let leader_addr = ftl + .leader_node + .map(|n| n.addr) + .or_else(|| { + ftl.leader_id + .and_then(|id| self.meta.get(id).map(|(ctrl, _)| ctrl)) + }) + .ok_or_else(|| "forward_to_leader: no leader address available".to_string())?; + return network::forward_join( + &leader_addr, + &self.config.secret, + self.config.node_id, + node_id, + control_addr, + media_addr, + proof, + self.tls_client.clone(), + ) + .await; + }''', + ''' Err(RaftError::APIError(ClientWriteError::ForwardToLeader(ftl))) => { + let leader_addr = ftl + .leader_node + .map(|n| n.addr) + .or_else(|| { + ftl.leader_id + .and_then(|id| self.meta.get(id).map(|(ctrl, _)| ctrl)) + }) + .ok_or_else(|| "forward_to_leader: no leader address available".to_string())?; + let proof_for_cache = proof.clone(); + let result = network::forward_join( + &leader_addr, + &self.config.secret, + self.config.node_id, + node_id, + control_addr, + media_addr, + proof, + self.tls_client.clone(), + ) + .await; + if result.is_ok() { + self.consume_admin_proof(&proof_for_cache); + } + return result; + }''' + )) + + replacements.append(( + ''' crate::log_info!("Cluster: node {node_id} joined as learner"); + Ok((self.cluster_id(), peers)) + }''', + ''' crate::log_info!("Cluster: node {node_id} joined as learner"); + self.consume_admin_proof(&proof); + Ok((self.cluster_id(), peers)) + }''' + )) + + replacements.append(( + ''' fn purge_expired_admin_proofs_drops_stale_entries() { + let mut map = HashMap::new(); + map.insert( + "stale-proof".to_string(), + Instant::now() - ADMIN_PROOF_REPLAY_TTL - Duration::from_secs(1), + ); + ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); + assert!(map.is_empty()); + } +}''', + ''' fn purge_expired_admin_proofs_drops_stale_entries() { + let mut map = HashMap::new(); + map.insert( + "stale-proof".to_string(), + Instant::now() - ADMIN_PROOF_REPLAY_TTL - Duration::from_secs(1), + ); + ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); + assert!(map.is_empty()); + } + + #[test] + fn fresh_admin_proofs_are_unique_and_payload_bound() { + let token: String = (0u8..24).map(|i| char::from(b'a' + (i % 26))).collect(); + let payload = "AdminDrain:2"; + let first = ClusterManager::mint_fresh_admin_proof(&token, payload); + let second = ClusterManager::mint_fresh_admin_proof(&token, payload); + assert_ne!(first, second); + assert!(ClusterManager::admin_proof_signature_valid( + &token, payload, &first + )); + assert!(ClusterManager::admin_proof_signature_valid( + &token, payload, &second + )); + assert!(!ClusterManager::admin_proof_signature_valid( + &token, + "AdminResume:2", + &first + )); + } +}''' + )) + + for old, new in replacements: + count = text.count(old) + if count != 1: + raise SystemExit(f'expected exactly one match, got {count}: {old[:100]!r}') + text = text.replace(old, new, 1) + + path.write_text(text) + + # This workflow is a one-shot transport for the patch. Remove it from + # the final PR diff before committing the actual fix. + Path('.github/workflows/fix-pr169-cursor.yml').unlink() + PY + + cargo fmt --all + + - name: Validate cluster tests + run: cargo test --features cluster + + - name: Commit fix back to PR branch + shell: bash + run: | + git config user.name "github-actions[bot]" + git config user.email "41898282+github-actions[bot]@users.noreply.github.com" + git add src/cluster/manager.rs .github/workflows/fix-pr169-cursor.yml + git commit -m "fix(security): make admin proof replay protection retry-safe" + git push origin HEAD:cursor/application-security-review-2fc5 From 0a80d6a2c56be9f95fd410dbcff34e81f5c23bc2 Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 07:10:36 +0200 Subject: [PATCH 5/8] chore: trigger PR 169 Cursor fix workflow --- .github/workflows/fix-pr169-cursor.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/fix-pr169-cursor.yml b/.github/workflows/fix-pr169-cursor.yml index 0c759b4..56ef02f 100644 --- a/.github/workflows/fix-pr169-cursor.yml +++ b/.github/workflows/fix-pr169-cursor.yml @@ -1,3 +1,4 @@ +# Re-trigger now that this workflow already exists on the branch. name: Apply PR 169 Cursor finding fix on: From 1aacd0354987a9822d2da875c22a88257679908f Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 07:24:15 +0200 Subject: [PATCH 6/8] fix(security): make admin proof replay protection retry-safe --- src/cluster/manager.rs | 86 ++++++++++++++++++++++++++++++++++++------ 1 file changed, 75 insertions(+), 11 deletions(-) diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 726fbf7..4e7fbcc 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -130,7 +130,7 @@ pub struct ClusterManager { /// reported as replicas with no data behind them. standby_subs: Mutex>, state_machine: SqliteStateMachine, - /// Recently accepted `admin_proof` digests — blocks captured ClientWrite / + /// Recently accepted `admin_proof` values — blocks captured ClientWrite / /// membership proofs from being replayed while the API token is unchanged. used_admin_proofs: Mutex>, } @@ -747,7 +747,7 @@ impl ClusterManager { .map_err(|e| CoordError::Cluster(e.to_string()))?; let req_str = serde_json::to_string(&req) .map_err(|e| CoordError::Cluster(e.to_string()))?; - let proof = crate::cluster::security::admin_proof(&token, &req_str); + let proof = Self::mint_fresh_admin_proof(&token, &req_str); let addr = ftl .leader_node .map(|n| n.addr) @@ -917,7 +917,35 @@ impl ClusterManager { .map(|h| h.api_token.read().clone()) .filter(|t| !t.is_empty()) .ok_or_else(|| "API token unavailable for cluster admin action".to_string())?; - Ok(crate::cluster::security::admin_proof(&token, payload)) + Ok(Self::mint_fresh_admin_proof(&token, payload)) + } + + /// Mint a unique proof for each admin attempt. The nonce is bound into the + /// API-token proof so an exact captured proof cannot be re-randomized by an + /// attacker, while a legitimate retry gets a different replay-cache key. + fn mint_fresh_admin_proof(api_token: &str, payload: &str) -> String { + let nonce = hex::encode(crate::cluster::security::auth_nonce()); + let signed_payload = format!("{nonce}:{payload}"); + let mac = crate::cluster::security::admin_proof(api_token, &signed_payload); + format!("{nonce}.{mac}") + } + + fn admin_proof_signature_valid(api_token: &str, payload: &str, proof: &str) -> bool { + let Some((nonce, mac)) = proof.split_once('.') else { + return false; + }; + if nonce.len() != 32 + || hex::decode(nonce) + .map(|decoded| decoded.len() != 16) + .unwrap_or(true) + { + return false; + } + let signed_payload = format!("{nonce}:{payload}"); + crate::cluster::security::secrets_equal( + &crate::cluster::security::admin_proof(api_token, &signed_payload), + mac, + ) } fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) { @@ -939,10 +967,7 @@ impl ClusterManager { if token.is_empty() { return false; } - if !crate::cluster::security::secrets_equal( - &crate::cluster::security::admin_proof(&token, payload), - proof, - ) { + if !Self::admin_proof_signature_valid(&token, payload, proof) { return false; } let mut used = self.used_admin_proofs.lock(); @@ -951,10 +976,23 @@ impl ClusterManager { if used.contains_key(proof) { return false; } - used.insert(proof.to_string(), now); + // A join proof is a configured one-time capability. Keep it retryable + // until add_learner/forwarding has actually succeeded; accept_join() + // records it after success. Other admin attempts mint a fresh nonce on + // each retry, so they can be consumed immediately here. + if !payload.starts_with("Join:") { + used.insert(proof.to_string(), now); + } true } + fn consume_admin_proof(&self, proof: &str) { + let mut used = self.used_admin_proofs.lock(); + let now = Instant::now(); + Self::purge_expired_admin_proofs(&mut used, now); + used.insert(proof.to_string(), now); + } + async fn forward_admin( &self, node_id: NodeId, @@ -1015,7 +1053,8 @@ impl ClusterManager { .and_then(|id| self.meta.get(id).map(|(ctrl, _)| ctrl)) }) .ok_or_else(|| "forward_to_leader: no leader address available".to_string())?; - return network::forward_join( + let proof_for_cache = proof.clone(); + let result = network::forward_join( &leader_addr, &self.config.secret, self.config.node_id, @@ -1026,6 +1065,10 @@ impl ClusterManager { self.tls_client.clone(), ) .await; + if result.is_ok() { + self.consume_admin_proof(&proof_for_cache); + } + return result; } Err(e) => { let msg = e.to_string(); @@ -1084,6 +1127,7 @@ impl ClusterManager { }) .collect(); crate::log_info!("Cluster: node {node_id} joined as learner"); + self.consume_admin_proof(&proof); Ok((self.cluster_id(), peers)) } @@ -1253,7 +1297,7 @@ impl ClusterManager { if self.local_live_sessions(&stream_id) == 0 { self.pending_drain_clears.lock().remove(&stream_id); if let Some(hooks) = self.session_hooks.lock().as_ref() { - hooks.deleted_streams.lock().remove(&stream_id); + hooks.deleted_streams.lock().remove(stream_id); } } } @@ -2385,7 +2429,7 @@ fn _cluster_manager_markers() {} #[cfg(test)] mod heartbeat_routing_tests { - use super::heartbeat_routing_addrs; + use super::*; #[test] fn plaintext_heartbeat_ignores_peer_supplied_routing_addrs() { @@ -2421,4 +2465,24 @@ mod heartbeat_routing_tests { ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); assert!(map.is_empty()); } + + #[test] + fn fresh_admin_proofs_are_unique_and_payload_bound() { + let token: String = (0u8..24).map(|i| char::from(b'a' + (i % 26))).collect(); + let payload = "AdminDrain:2"; + let first = ClusterManager::mint_fresh_admin_proof(&token, payload); + let second = ClusterManager::mint_fresh_admin_proof(&token, payload); + assert_ne!(first, second); + assert!(ClusterManager::admin_proof_signature_valid( + &token, payload, &first + )); + assert!(ClusterManager::admin_proof_signature_valid( + &token, payload, &second + )); + assert!(!ClusterManager::admin_proof_signature_valid( + &token, + "AdminResume:2", + &first + )); + } } From fddcab8bab4d7d6fc0d4c888b93b6b79cfc5beb8 Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 07:24:52 +0200 Subject: [PATCH 7/8] chore: remove temporary PR fix workflow --- .github/workflows/fix-pr169-cursor.yml | 252 ------------------------- 1 file changed, 252 deletions(-) delete mode 100644 .github/workflows/fix-pr169-cursor.yml diff --git a/.github/workflows/fix-pr169-cursor.yml b/.github/workflows/fix-pr169-cursor.yml deleted file mode 100644 index 56ef02f..0000000 --- a/.github/workflows/fix-pr169-cursor.yml +++ /dev/null @@ -1,252 +0,0 @@ -# Re-trigger now that this workflow already exists on the branch. -name: Apply PR 169 Cursor finding fix - -on: - push: - branches: - - cursor/application-security-review-2fc5 - paths: - - .github/workflows/fix-pr169-cursor.yml - -permissions: - contents: write - -jobs: - apply-fix: - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@v4 - with: - ref: cursor/application-security-review-2fc5 - fetch-depth: 0 - - - name: Apply retry-safe admin proof fix - shell: bash - run: | - python3 - <<'PY' - from pathlib import Path - - path = Path('src/cluster/manager.rs') - text = path.read_text() - - replacements = [] - - replacements.append(( - ''' let proof = crate::cluster::security::admin_proof(&token, &req_str);''', - ''' let proof = ClusterManager::mint_fresh_admin_proof(&token, &req_str);''' - )) - - replacements.append(( - ''' Ok(crate::cluster::security::admin_proof(&token, payload)) - } - - fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) {''', - ''' Ok(Self::mint_fresh_admin_proof(&token, payload)) - } - - /// Mint a unique proof for each admin attempt. The random nonce is signed - /// together with the canonical payload, so repeated legitimate operations - /// do not reuse the same replay-cache key. - fn mint_fresh_admin_proof(api_token: &str, payload: &str) -> String { - let nonce = hex::encode(crate::cluster::security::auth_nonce()); - let signed_payload = format!("{nonce}:{payload}"); - let mac = crate::cluster::security::admin_proof(api_token, &signed_payload); - format!("{nonce}.{mac}") - } - - fn admin_proof_signature_valid(api_token: &str, payload: &str, proof: &str) -> bool { - let Some((nonce, mac)) = proof.split_once('.') else { - return false; - }; - // auth_nonce() is 16 bytes => exactly 32 lowercase/uppercase hex chars. - if nonce.len() != 32 - || hex::decode(nonce) - .map(|decoded| decoded.len() != 16) - .unwrap_or(true) - { - return false; - } - let signed_payload = format!("{nonce}:{payload}"); - crate::cluster::security::secrets_equal( - &crate::cluster::security::admin_proof(api_token, &signed_payload), - mac, - ) - } - - fn purge_expired_admin_proofs(guard: &mut HashMap, now: Instant) {''' - )) - - replacements.append(( - ''' if !crate::cluster::security::secrets_equal( - &crate::cluster::security::admin_proof(&token, payload), - proof, - ) { - return false; - } - let mut used = self.used_admin_proofs.lock(); - let now = Instant::now(); - Self::purge_expired_admin_proofs(&mut used, now); - if used.contains_key(proof) { - return false; - } - used.insert(proof.to_string(), now); - true - } - - async fn forward_admin(''', - ''' if !Self::admin_proof_signature_valid(&token, payload, proof) { - return false; - } - let mut used = self.used_admin_proofs.lock(); - let now = Instant::now(); - Self::purge_expired_admin_proofs(&mut used, now); - if used.contains_key(proof) { - return false; - } - // Join proofs must remain retryable if add_learner/forwarding fails. - // accept_join() consumes them only after the join actually succeeds. - if !payload.starts_with("Join:") { - used.insert(proof.to_string(), now); - } - true - } - - fn consume_admin_proof(&self, proof: &str) { - let mut used = self.used_admin_proofs.lock(); - let now = Instant::now(); - Self::purge_expired_admin_proofs(&mut used, now); - used.insert(proof.to_string(), now); - } - - async fn forward_admin(''' - )) - - replacements.append(( - ''' Err(RaftError::APIError(ClientWriteError::ForwardToLeader(ftl))) => { - let leader_addr = ftl - .leader_node - .map(|n| n.addr) - .or_else(|| { - ftl.leader_id - .and_then(|id| self.meta.get(id).map(|(ctrl, _)| ctrl)) - }) - .ok_or_else(|| "forward_to_leader: no leader address available".to_string())?; - return network::forward_join( - &leader_addr, - &self.config.secret, - self.config.node_id, - node_id, - control_addr, - media_addr, - proof, - self.tls_client.clone(), - ) - .await; - }''', - ''' Err(RaftError::APIError(ClientWriteError::ForwardToLeader(ftl))) => { - let leader_addr = ftl - .leader_node - .map(|n| n.addr) - .or_else(|| { - ftl.leader_id - .and_then(|id| self.meta.get(id).map(|(ctrl, _)| ctrl)) - }) - .ok_or_else(|| "forward_to_leader: no leader address available".to_string())?; - let proof_for_cache = proof.clone(); - let result = network::forward_join( - &leader_addr, - &self.config.secret, - self.config.node_id, - node_id, - control_addr, - media_addr, - proof, - self.tls_client.clone(), - ) - .await; - if result.is_ok() { - self.consume_admin_proof(&proof_for_cache); - } - return result; - }''' - )) - - replacements.append(( - ''' crate::log_info!("Cluster: node {node_id} joined as learner"); - Ok((self.cluster_id(), peers)) - }''', - ''' crate::log_info!("Cluster: node {node_id} joined as learner"); - self.consume_admin_proof(&proof); - Ok((self.cluster_id(), peers)) - }''' - )) - - replacements.append(( - ''' fn purge_expired_admin_proofs_drops_stale_entries() { - let mut map = HashMap::new(); - map.insert( - "stale-proof".to_string(), - Instant::now() - ADMIN_PROOF_REPLAY_TTL - Duration::from_secs(1), - ); - ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); - assert!(map.is_empty()); - } -}''', - ''' fn purge_expired_admin_proofs_drops_stale_entries() { - let mut map = HashMap::new(); - map.insert( - "stale-proof".to_string(), - Instant::now() - ADMIN_PROOF_REPLAY_TTL - Duration::from_secs(1), - ); - ClusterManager::purge_expired_admin_proofs(&mut map, Instant::now()); - assert!(map.is_empty()); - } - - #[test] - fn fresh_admin_proofs_are_unique_and_payload_bound() { - let token: String = (0u8..24).map(|i| char::from(b'a' + (i % 26))).collect(); - let payload = "AdminDrain:2"; - let first = ClusterManager::mint_fresh_admin_proof(&token, payload); - let second = ClusterManager::mint_fresh_admin_proof(&token, payload); - assert_ne!(first, second); - assert!(ClusterManager::admin_proof_signature_valid( - &token, payload, &first - )); - assert!(ClusterManager::admin_proof_signature_valid( - &token, payload, &second - )); - assert!(!ClusterManager::admin_proof_signature_valid( - &token, - "AdminResume:2", - &first - )); - } -}''' - )) - - for old, new in replacements: - count = text.count(old) - if count != 1: - raise SystemExit(f'expected exactly one match, got {count}: {old[:100]!r}') - text = text.replace(old, new, 1) - - path.write_text(text) - - # This workflow is a one-shot transport for the patch. Remove it from - # the final PR diff before committing the actual fix. - Path('.github/workflows/fix-pr169-cursor.yml').unlink() - PY - - cargo fmt --all - - - name: Validate cluster tests - run: cargo test --features cluster - - - name: Commit fix back to PR branch - shell: bash - run: | - git config user.name "github-actions[bot]" - git config user.email "41898282+github-actions[bot]@users.noreply.github.com" - git add src/cluster/manager.rs .github/workflows/fix-pr169-cursor.yml - git commit -m "fix(security): make admin proof replay protection retry-safe" - git push origin HEAD:cursor/application-security-review-2fc5 From 40922e3f79a840743717dfd3064ba4b9dee598de Mon Sep 17 00:00:00 2001 From: Alexander Wagner Date: Wed, 19 Aug 2026 07:28:00 +0200 Subject: [PATCH 8/8] fix: restore stream-id borrow in drain cleanup --- src/cluster/manager.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 4e7fbcc..94eb754 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -1297,7 +1297,7 @@ impl ClusterManager { if self.local_live_sessions(&stream_id) == 0 { self.pending_drain_clears.lock().remove(&stream_id); if let Some(hooks) = self.session_hooks.lock().as_ref() { - hooks.deleted_streams.lock().remove(stream_id); + hooks.deleted_streams.lock().remove(&stream_id); } } }