diff --git a/README.md b/README.md index e82d1e5..af1be6b 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,7 @@ Hosted clients should treat `bifrost-rs` as the signer authority. - `runtime_status()` is the canonical aggregated read model. - `readiness()` is the narrower capability view. +- `peer_status()` reports peer capability, policy, latency, and nonce-inventory telemetry for operator UIs. - `drain_runtime_events()` is incremental and lossy-safe; clients must recover truth from `runtime_status()`. - `prepare_sign()` and `prepare_ecdh()` are the normal operation-prep APIs. - `wipe_state()` is the canonical signer-side reset path. diff --git a/crates/bifrost-app/src/host/client.rs b/crates/bifrost-app/src/host/client.rs index 90df69b..9f76246 100644 --- a/crates/bifrost-app/src/host/client.rs +++ b/crates/bifrost-app/src/host/client.rs @@ -192,9 +192,10 @@ mod tests { #[cfg(unix)] fn test_socket_path(name: &str) -> PathBuf { + let short_name = name.chars().next().unwrap_or('x'); let unique = format!( - "bifrost-app-client-{}-{}-{}.sock", - name, + "bac-{}-{}-{}.sock", + short_name, std::process::id(), std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) diff --git a/crates/bifrost-app/tests/daemon_control_errors.rs b/crates/bifrost-app/tests/daemon_control_errors.rs index 84df41d..07732c4 100644 --- a/crates/bifrost-app/tests/daemon_control_errors.rs +++ b/crates/bifrost-app/tests/daemon_control_errors.rs @@ -15,12 +15,13 @@ fn token_with_byte(b: u8) -> DaemonToken { } fn temp_path(name: &str, suffix: &str) -> PathBuf { + let short_name = name.chars().next().unwrap_or('x'); let nonce = SystemTime::now() .duration_since(UNIX_EPOCH) .expect("clock") .as_nanos(); std::env::temp_dir().join(format!( - "bifrost-daemon-errors-{name}-{}-{nonce}.{suffix}", + "bde-{short_name}-{}-{nonce}.{suffix}", std::process::id() )) } diff --git a/crates/bifrost-app/tests/daemon_lifecycle.rs b/crates/bifrost-app/tests/daemon_lifecycle.rs index 4822e05..3822e10 100644 --- a/crates/bifrost-app/tests/daemon_lifecycle.rs +++ b/crates/bifrost-app/tests/daemon_lifecycle.rs @@ -19,12 +19,13 @@ fn token_with_byte(b: u8) -> DaemonToken { } fn temp_path(name: &str, suffix: &str) -> PathBuf { + let short_name = name.chars().next().unwrap_or('x'); let nonce = SystemTime::now() .duration_since(UNIX_EPOCH) .expect("system clock before unix epoch") .as_nanos(); std::env::temp_dir().join(format!( - "bifrost-daemon-{name}-{}-{nonce}.{suffix}", + "bd-{short_name}-{}-{nonce}.{suffix}", std::process::id() )) } diff --git a/crates/bifrost-bridge-tokio/src/lib.rs b/crates/bifrost-bridge-tokio/src/lib.rs index aa4c89a..0c9a19d 100644 --- a/crates/bifrost-bridge-tokio/src/lib.rs +++ b/crates/bifrost-bridge-tokio/src/lib.rs @@ -847,6 +847,9 @@ fn completed_operation_kind(operation: &CompletedOperation) -> &'static str { CompletedOperation::Sign { .. } => "sign", CompletedOperation::Ecdh { .. } => "ecdh", CompletedOperation::Ping { .. } => "ping", + CompletedOperation::PingServed { .. } => "ping_served", + CompletedOperation::EcdhServed { .. } => "ecdh_served", + CompletedOperation::SignServed { .. } => "sign_served", CompletedOperation::Onboard { .. } => "onboard", CompletedOperation::OnboardServed { .. } => "onboard_served", } diff --git a/crates/bifrost-bridge-wasm/src/lib.rs b/crates/bifrost-bridge-wasm/src/lib.rs index ad89d65..2ce6ba4 100644 --- a/crates/bifrost-bridge-wasm/src/lib.rs +++ b/crates/bifrost-bridge-wasm/src/lib.rs @@ -286,6 +286,18 @@ enum CompletedOperationJson { request_id: String, peer: String, }, + PingServed { + request_id: String, + peer_pubkey32_hex: String, + }, + EcdhServed { + request_id: String, + peer_pubkey32_hex: String, + }, + SignServed { + request_id: String, + peer_pubkey32_hex: String, + }, Onboard { request_id: String, group_member_count: usize, @@ -1587,6 +1599,92 @@ mod tests { assert!(err.to_string().contains("message_hex_32")); } + #[test] + fn inbound_sign_nonce_miss_surfaces_sign_failure() { + let bundle = + create_keyset(CreateKeysetConfig::new("Test Group", 2, 2)).expect("create keyset"); + let group = bundle.group.clone(); + let alice_share = bundle.shares[0].clone(); + let bob_share = bundle.shares[1].clone(); + let bob_peer = hex::encode(&group.members[1].pubkey[1..]); + let mut bob_seed_state = DeviceState::new(bob_share.idx, *bob_share.seckey.expose_bytes()); + let stale_bob_nonces = bob_seed_state + .nonce_pool + .generate_for_peer( + alice_share.idx, + 10, + &bob_seed_state.secrets.nonce_pool_secret, + ) + .expect("generate stale bob nonces"); + + let alice_bootstrap = RuntimeBootstrapInput { + group: GroupPackageWire::from(group.clone()), + share: SharePackageWire::from(alice_share), + peers: vec![bob_peer.clone()], + initial_peer_nonces: vec![BootstrapPeerNoncesInput { + peer: bob_peer.clone(), + nonces: stale_bob_nonces.into_iter().map(Into::into).collect(), + }], + }; + let bob_bootstrap = RuntimeBootstrapInput { + group: GroupPackageWire::from(group), + share: SharePackageWire::from(bob_share), + peers: vec![hex::encode(&bundle.group.members[0].pubkey[1..])], + initial_peer_nonces: Vec::new(), + }; + let now = 1_700_000_000_000u64; + + let mut alice = WasmBridgeRuntime::new(); + alice + .init_runtime( + "{}".to_string(), + serde_json::to_string(&alice_bootstrap).expect("alice bootstrap"), + ) + .expect("init alice"); + let mut bob = WasmBridgeRuntime::new(); + bob.init_runtime( + "{}".to_string(), + serde_json::to_string(&bob_bootstrap).expect("bob bootstrap"), + ) + .expect("init bob"); + + alice + .handle_command( + serde_json::json!({ + "type": "sign", + "message_hex_32": hex::encode([0x77; 32]) + }) + .to_string(), + ) + .expect("queue sign"); + alice.tick(now).expect("tick alice"); + let outbound: Vec = + serde_json::from_str(&alice.drain_outbound_events().expect("alice outbound")) + .expect("decode alice outbound"); + assert_eq!(outbound.len(), 1); + + bob.handle_inbound_event( + serde_json::to_string(&outbound[0]).expect("encode inbound event"), + ) + .expect("bob handle inbound"); + bob.tick(now + 1).expect("tick bob"); + + let failures: serde_json::Value = + serde_json::from_str(&bob.drain_failures().expect("bob failures")) + .expect("decode failures"); + let failure_array = failures.as_array().expect("failure array"); + assert_eq!(failure_array.len(), 1); + assert_eq!(failure_array[0]["op_type"], "sign"); + assert_eq!(failure_array[0]["code"], "peer_rejected"); + assert!( + failure_array[0]["message"] + .as_str() + .expect("failure message") + .to_ascii_lowercase() + .contains("nonce unavailable") + ); + } + #[test] fn restore_runtime_round_trip_preserves_runtime_metadata_and_status() { let bundle = @@ -1991,6 +2089,15 @@ mod tests { serde_json::from_str(&bob.drain_outbound_events().expect("bob outbound")) .expect("decode bob outbound"); assert_eq!(bob_outbound.len(), 1); + let bob_completions: serde_json::Value = + serde_json::from_str(&bob.drain_completions().expect("bob completions")) + .expect("decode bob completions"); + let bob_completion_array = bob_completions.as_array().expect("completion array"); + assert_eq!(bob_completion_array.len(), 1); + assert_eq!( + bob_completion_array[0]["PingServed"]["peer_pubkey32_hex"], + bob_bootstrap.peers[0] + ); alice .handle_inbound_event(serde_json::to_string(&bob_outbound[0]).expect("encode reply")) @@ -2134,6 +2241,18 @@ impl From for CompletedOperationJson { shared_secret_hex32: hex::encode(shared_secret), }, CompletedOperation::Ping { request_id, peer } => Self::Ping { request_id, peer }, + CompletedOperation::PingServed { request_id, peer } => Self::PingServed { + request_id, + peer_pubkey32_hex: peer, + }, + CompletedOperation::EcdhServed { request_id, peer } => Self::EcdhServed { + request_id, + peer_pubkey32_hex: peer, + }, + CompletedOperation::SignServed { request_id, peer } => Self::SignServed { + request_id, + peer_pubkey32_hex: peer, + }, CompletedOperation::Onboard { request_id, group_member_count, diff --git a/crates/bifrost-codec/src/wire.rs b/crates/bifrost-codec/src/wire.rs index 8a56390..ab008fd 100644 --- a/crates/bifrost-codec/src/wire.rs +++ b/crates/bifrost-codec/src/wire.rs @@ -144,6 +144,8 @@ pub struct PingPayloadWire { pub version: u16, pub advertised_nonces: Vec, pub held_peer_nonce_codes: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub recognized_peer_nonce_codes: Option>, pub policy_profile: Option, /// Hex of the sender's nonce-pool generation. Empty/absent on legacy peers. #[serde(default)] @@ -637,6 +639,15 @@ impl TryFrom for PingPayload { "ping held peer nonce codes exceed max size", )); } + if value + .recognized_peer_nonce_codes + .as_ref() + .is_some_and(|codes| codes.len() > MAX_NONCE_PACKAGE) + { + return Err(crate::error::CodecError::InvalidPayload( + "ping recognized peer nonce codes exceed max size", + )); + } let advertised_nonces = value .advertised_nonces .into_iter() @@ -647,6 +658,15 @@ impl TryFrom for PingPayload { .into_iter() .map(|code| hexbytes::decode(&code)) .collect::, _>>()?; + let recognized_peer_nonce_codes = value + .recognized_peer_nonce_codes + .map(|codes| { + codes + .into_iter() + .map(|code| hexbytes::decode(&code)) + .collect::, _>>() + }) + .transpose()?; let nonce_pool_generation = if value.nonce_pool_generation.is_empty() { bifrost_core::nonce::UNKNOWN_POOL_GENERATION @@ -658,6 +678,7 @@ impl TryFrom for PingPayload { version: value.version, advertised_nonces, held_peer_nonce_codes, + recognized_peer_nonce_codes, policy_profile: value.policy_profile.map(TryInto::try_into).transpose()?, nonce_pool_generation, }) @@ -678,6 +699,12 @@ impl From for PingPayloadWire { .into_iter() .map(|code| hexbytes::encode(&code)) .collect(), + recognized_peer_nonce_codes: value.recognized_peer_nonce_codes.map(|codes| { + codes + .into_iter() + .map(|code| hexbytes::encode(&code)) + .collect() + }), policy_profile: value.policy_profile.map(Into::into), nonce_pool_generation: if value.nonce_pool_generation == bifrost_core::nonce::UNKNOWN_POOL_GENERATION @@ -1033,6 +1060,7 @@ mod tests { code: [3u8; 32], }], held_peer_nonce_codes: vec![[4u8; 32], [5u8; 32]], + recognized_peer_nonce_codes: Some(vec![[6u8; 32]]), policy_profile: None, nonce_pool_generation: [9u8; 32], }; @@ -1047,6 +1075,7 @@ mod tests { version: 1, advertised_nonces: Vec::new(), held_peer_nonce_codes: Vec::new(), + recognized_peer_nonce_codes: None, policy_profile: None, nonce_pool_generation: String::new(), }; diff --git a/crates/bifrost-codec/tests/wire_fuzz.rs b/crates/bifrost-codec/tests/wire_fuzz.rs index da40685..e4b4859 100644 --- a/crates/bifrost-codec/tests/wire_fuzz.rs +++ b/crates/bifrost-codec/tests/wire_fuzz.rs @@ -235,6 +235,11 @@ fn gen_ping(rng: &mut Rng) -> PingPayloadWire { version: if rng.bool() { 2 } else { rng.u16() }, advertised_nonces: (0..list_len(rng).min(6)).map(|_| gen_nonce(rng)).collect(), held_peer_nonce_codes: (0..list_len(rng).min(6)).map(|_| field(rng, 32)).collect(), + recognized_peer_nonce_codes: if rng.bool() { + None + } else { + Some((0..list_len(rng).min(6)).map(|_| field(rng, 32)).collect()) + }, policy_profile: if rng.bool() { None } else { diff --git a/crates/bifrost-core/src/nonce.rs b/crates/bifrost-core/src/nonce.rs index 5fbb288..487ace5 100644 --- a/crates/bifrost-core/src/nonce.rs +++ b/crates/bifrost-core/src/nonce.rs @@ -236,6 +236,17 @@ impl NoncePool { } } + pub fn retain_incoming_codes(&mut self, peer_idx: u16, codes: &[Bytes32]) { + let Some(map) = self.incoming.get_mut(&peer_idx) else { + return; + }; + let retain = codes.iter().copied().collect::>(); + map.retain(|code, _| retain.contains(code)); + if let Some(order) = self.incoming_order.get_mut(&peer_idx) { + order.retain(|code| map.contains_key(code)); + } + } + pub fn consume_incoming(&mut self, peer_idx: u16) -> Option { let map = self.incoming.get_mut(&peer_idx)?; let order = self.incoming_order.get_mut(&peer_idx)?; @@ -252,6 +263,22 @@ impl NoncePool { None } + pub fn consume_latest_incoming(&mut self, peer_idx: u16) -> Option { + let map = self.incoming.get_mut(&peer_idx)?; + let order = self.incoming_order.get_mut(&peer_idx)?; + while let Some(code) = order.pop_back() { + if let Some(nonce) = map.remove(&code) { + return Some(MemberPublicNonce { + idx: peer_idx, + binder_pn: nonce.binder_pn, + hidden_pn: nonce.hidden_pn, + code: nonce.code, + }); + } + } + None + } + pub fn take_outgoing_signing_nonces( &mut self, peer_idx: u16, @@ -532,4 +559,51 @@ mod tests { assert_eq!(consumed_first.code, first.code); assert_eq!(consumed_second.code, second.code); } + + #[test] + fn retain_incoming_codes_prunes_stale_entries_and_order() { + let mut pool = NoncePool::new(1, NoncePoolConfig::default()); + pool.init_peer(2); + + let stale = DerivedPublicNonce { + binder_pn: [2u8; 33], + hidden_pn: [3u8; 33], + code: [10u8; 32], + }; + let current = DerivedPublicNonce { + binder_pn: [4u8; 33], + hidden_pn: [5u8; 33], + code: [11u8; 32], + }; + pool.store_incoming(2, vec![stale, current.clone()]); + pool.retain_incoming_codes(2, &[current.code]); + + let consumed = pool.consume_incoming(2).expect("retained nonce"); + assert_eq!(consumed.code, current.code); + assert!(pool.consume_incoming(2).is_none()); + } + + #[test] + fn consume_latest_incoming_prefers_newly_stored_nonces() { + let mut pool = NoncePool::new(1, NoncePoolConfig::default()); + pool.init_peer(2); + + let old = DerivedPublicNonce { + binder_pn: [2u8; 33], + hidden_pn: [3u8; 33], + code: [10u8; 32], + }; + let new = DerivedPublicNonce { + binder_pn: [4u8; 33], + hidden_pn: [5u8; 33], + code: [11u8; 32], + }; + pool.store_incoming(2, vec![old.clone()]); + pool.store_incoming(2, vec![new.clone()]); + + let consumed_new = pool.consume_latest_incoming(2).expect("latest nonce"); + let consumed_old = pool.consume_latest_incoming(2).expect("old nonce"); + assert_eq!(consumed_new.code, new.code); + assert_eq!(consumed_old.code, old.code); + } } diff --git a/crates/bifrost-core/src/types.rs b/crates/bifrost-core/src/types.rs index 2ddf051..0b2965c 100644 --- a/crates/bifrost-core/src/types.rs +++ b/crates/bifrost-core/src/types.rs @@ -469,6 +469,8 @@ pub struct PingPayload { pub advertised_nonces: Vec, #[serde(with = "serde_fixed_array::vec_bytes32")] pub held_peer_nonce_codes: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub recognized_peer_nonce_codes: Option>, pub policy_profile: Option, /// The sender's nonce-pool generation; a change signals the sender reset its /// outgoing pool, so the receiver discards the nonces it holds from it. Absent @@ -527,6 +529,7 @@ mod tests { code: [7u8; 32], }], held_peer_nonce_codes: vec![[11u8; 32]], + recognized_peer_nonce_codes: Some(vec![[12u8; 32]]), policy_profile: Some(PeerScopedPolicyProfile { for_peer: [8u8; 32], revision: 9, @@ -565,6 +568,10 @@ mod tests { assert_eq!(onboard.group.members.len(), 2); assert_eq!(ping.version, 2); assert_eq!(ping.held_peer_nonce_codes.len(), 1); + assert_eq!( + ping.recognized_peer_nonce_codes.as_deref(), + Some(&[[12u8; 32]][..]) + ); assert_eq!( ping.policy_profile .as_ref() diff --git a/crates/bifrost-devtools/src/e2e.rs b/crates/bifrost-devtools/src/e2e.rs index a563a67..ab05376 100644 --- a/crates/bifrost-devtools/src/e2e.rs +++ b/crates/bifrost-devtools/src/e2e.rs @@ -1,8 +1,8 @@ use std::env; use std::fs::{self, OpenOptions}; -use std::io::Write; +use std::io::{self, Write}; use std::path::{Path, PathBuf}; -use std::process::{Child, Command, Stdio}; +use std::process::{Child, Command, ExitStatus, Output, Stdio}; use std::thread; use std::time::{Duration, Instant}; @@ -11,6 +11,8 @@ use serde_json::Value; const DEFAULT_RELAY: &str = "ws://127.0.0.1:8194"; const DEFAULT_VAULT_PASSPHRASE: &str = "igloo-shell-e2e-passphrase"; +const EXEC_BUSY_RETRIES: usize = 5; +const EXEC_BUSY_RETRY_DELAY: Duration = Duration::from_millis(50); pub fn run_e2e_node_command(args: &[String]) -> Result<()> { let mut out_dir: Option = None; @@ -737,8 +739,7 @@ fn run_command( args: &[&str], name: &str, ) -> Result<()> { - let status = build_command(exe, shell_env, args) - .status() + let status = command_status_with_busy_retry(exe, shell_env, args) .with_context(|| format!("run {name}"))?; if !status.success() { bail!("{name} failed: {status}"); @@ -771,8 +772,7 @@ fn run_shell_json( shell_env: Option<&ManagedShellEnv>, args: &[&str], ) -> Result { - let output = build_command(shell_exe, shell_env, args) - .output() + let output = command_output_with_busy_retry(shell_exe, shell_env, args) .context("capture igloo-shell output")?; if !output.status.success() { bail!( @@ -793,6 +793,39 @@ fn build_command(exe: &Path, shell_env: Option<&ManagedShellEnv>, args: &[&str]) command } +fn command_status_with_busy_retry( + exe: &Path, + shell_env: Option<&ManagedShellEnv>, + args: &[&str], +) -> io::Result { + retry_if_executable_busy(|| build_command(exe, shell_env, args).status()) +} + +fn command_output_with_busy_retry( + exe: &Path, + shell_env: Option<&ManagedShellEnv>, + args: &[&str], +) -> io::Result { + retry_if_executable_busy(|| build_command(exe, shell_env, args).output()) +} + +fn retry_if_executable_busy(mut run: impl FnMut() -> io::Result) -> io::Result { + for attempt in 0..EXEC_BUSY_RETRIES { + match run() { + Ok(value) => return Ok(value), + Err(err) if is_executable_busy(&err) && attempt + 1 < EXEC_BUSY_RETRIES => { + thread::sleep(EXEC_BUSY_RETRY_DELAY); + } + Err(err) => return Err(err), + } + } + unreachable!("retry loop always returns before exhausting attempts") +} + +fn is_executable_busy(err: &io::Error) -> bool { + cfg!(unix) && err.raw_os_error() == Some(26) +} + fn append_output(path: &Path, label: &str, output: &str) -> Result<()> { let mut file = OpenOptions::new() .append(true) @@ -973,8 +1006,18 @@ mod tests { unsafe { env::remove_var("IGLOO_SHELL_BIN"); } - let resolved = resolve_shell_exe(None).expect("default resolution"); - assert!(resolved.ends_with("repos/igloo-shell/target/debug/igloo-shell")); + let default_path = infra_root() + .expect("infra root") + .join("repos/igloo-shell/target/debug/igloo-shell"); + match resolve_shell_exe(None) { + Ok(resolved) => assert_eq!(resolved, default_path), + Err(err) => { + assert!(!default_path.is_file()); + let message = err.to_string(); + assert!(message.contains("missing igloo-shell binary")); + assert!(message.contains(&default_path.display().to_string())); + } + } } #[test] diff --git a/crates/bifrost-signer/src/lib.rs b/crates/bifrost-signer/src/lib.rs index 036b6ef..adb2374 100644 --- a/crates/bifrost-signer/src/lib.rs +++ b/crates/bifrost-signer/src/lib.rs @@ -716,6 +716,21 @@ pub enum CompletedOperation { request_id: String, peer: String, }, + /// Recorded by the responder after it answers a ping request. + PingServed { + request_id: String, + peer: String, + }, + /// Recorded by the responder after it returns an ECDH share. + EcdhServed { + request_id: String, + peer: String, + }, + /// Recorded by the responder after it returns a signing share. + SignServed { + request_id: String, + peer: String, + }, Onboard { request_id: String, group_member_count: usize, @@ -737,6 +752,9 @@ impl CompletedOperation { CompletedOperation::Sign { request_id, .. } | CompletedOperation::Ecdh { request_id, .. } | CompletedOperation::Ping { request_id, .. } + | CompletedOperation::PingServed { request_id, .. } + | CompletedOperation::EcdhServed { request_id, .. } + | CompletedOperation::SignServed { request_id, .. } | CompletedOperation::Onboard { request_id, .. } | CompletedOperation::OnboardServed { request_id, .. } => request_id, } @@ -1648,7 +1666,7 @@ impl SigningDevice { let nonce = self .state .nonce_pool - .consume_incoming(idx) + .consume_latest_incoming(idx) .ok_or(SignerError::NonceUnavailable)?; member_nonce_sets.push(bifrost_core::types::MemberNonceCommitmentSet { idx, @@ -1821,7 +1839,7 @@ impl SigningDevice { .member_idx_by_pubkey .get(peer) .ok_or_else(|| SignerError::UnknownPeer(peer.to_string()))?; - let payload = self.ping_payload(peer, peer_idx)?; + let payload = self.ping_payload(peer, peer_idx, None)?; let request_id = self.next_request_id(); self.latest_request_id = Some(request_id.clone()); @@ -2012,14 +2030,23 @@ impl SigningDevice { pub fn expire_stale(&mut self, now: u64) -> Vec { let mut stale = Vec::new(); + let mut ping_failed_peers = Vec::new(); self.state.pending_operations.retain(|id, op| { if op.timeout_at <= now { + let failed_peer = if matches!(op.op_type, PendingOpType::Ping) { + op.target_peers.first().cloned() + } else { + None + }; + if let Some(peer) = failed_peer.clone() { + ping_failed_peers.push(peer); + } stale.push(OperationFailure { request_id: id.clone(), op_type: op.op_type.clone(), code: OperationFailureCode::Timeout, message: "locked peer response timeout".to_string(), - failed_peer: None, + failed_peer, }); false } else { @@ -2029,6 +2056,9 @@ impl SigningDevice { for failure in &stale { self.state.op_started_ms.remove(&failure.request_id); } + for peer in ping_failed_peers { + self.state.peer_last_seen.remove(&peer); + } // Drop parked approvals the operator never resolved. Silent: the // requester's own operation has already timed out by now, so there is no // peer left waiting for a response. @@ -2265,12 +2295,13 @@ impl SigningDevice { SignerError::InvalidRequest(e.to_string()) })?; let peer_generation = ping.nonce_pool_generation; + let recognized_peer_nonce_codes = self + .recognized_peer_nonce_codes_for_reported_inventory( + sender_idx, + &ping.held_peer_nonce_codes, + ); self.reconcile_peer_generation(&sender, sender_idx, peer_generation); - if !ping.advertised_nonces.is_empty() { - self.state - .nonce_pool - .store_incoming(sender_idx, ping.advertised_nonces); - } + self.apply_ping_nonce_payload(sender_idx, &ping); self.store_remote_nonce_inventory_observation( &sender, sender_idx, @@ -2282,14 +2313,21 @@ impl SigningDevice { self.store_remote_scoped_policy(&sender, profile)?; } + let served_request_id = envelope.request_id; + let served_peer = sender.clone(); let response = BridgeEnvelope { - request_id: envelope.request_id, + request_id: served_request_id.clone(), sent_at: now, payload: BridgePayload::PingResponse(PingPayloadWire::from( - self.ping_payload(&sender, sender_idx)?, + self.ping_payload(&sender, sender_idx, Some(recognized_peer_nonce_codes))?, )), }; - self.encrypt_for_peers(&[sender], &response) + let outbound = self.encrypt_for_peers(&[sender], &response)?; + self.completions.push_back(CompletedOperation::PingServed { + request_id: served_request_id, + peer: served_peer, + }); + Ok(outbound) } BridgePayload::OnboardRequest(wire) => { let request: bifrost_core::types::OnboardRequest = @@ -2464,12 +2502,19 @@ impl SigningDevice { partial.replenish = Some(replenish); } + let served_request_id = envelope.request_id; + let served_peer = sender.clone(); let response = BridgeEnvelope { - request_id: envelope.request_id, + request_id: served_request_id.clone(), sent_at: now, payload: BridgePayload::SignResponse(PartialSigPackageWire::from(partial)), }; - self.encrypt_for_peers(&[sender], &response) + let outbound = self.encrypt_for_peers(&[sender], &response)?; + self.completions.push_back(CompletedOperation::SignServed { + request_id: served_request_id, + peer: served_peer, + }); + Ok(outbound) } BridgePayload::EcdhRequest(wire) => { let req: EcdhPackage = @@ -2486,12 +2531,19 @@ impl SigningDevice { let response = ecdh_create_from_share(&req.members, &self.share, &targets) .map_err(|e| SignerError::InvalidRequest(e.to_string()))?; + let served_request_id = envelope.request_id; + let served_peer = sender.clone(); let envelope = BridgeEnvelope { - request_id: envelope.request_id, + request_id: served_request_id.clone(), sent_at: now, payload: BridgePayload::EcdhResponse(EcdhPackageWire::from(response)), }; - self.encrypt_for_peers(&[sender], &envelope) + let outbound = self.encrypt_for_peers(&[sender], &envelope)?; + self.completions.push_back(CompletedOperation::EcdhServed { + request_id: served_request_id, + peer: served_peer, + }); + Ok(outbound) } BridgePayload::Error(_) => Ok(Vec::new()), BridgePayload::PingResponse(_) => { @@ -2547,11 +2599,7 @@ impl SigningDevice { .ok_or_else(|| SignerError::UnknownPeer(sender.to_string()))?; let peer_generation = ping.nonce_pool_generation; self.reconcile_peer_generation(sender, sender_idx, peer_generation); - if !ping.advertised_nonces.is_empty() { - self.state - .nonce_pool - .store_incoming(sender_idx, ping.advertised_nonces); - } + self.apply_ping_nonce_payload(sender_idx, &ping); self.store_remote_nonce_inventory_observation( sender, sender_idx, @@ -2845,6 +2893,26 @@ impl SigningDevice { held_codes } + fn recognized_peer_nonce_codes_for_reported_inventory( + &self, + peer_idx: u16, + held_codes: &[Bytes32], + ) -> Vec { + let current_codes = self.state.nonce_pool.outgoing_public_nonce_codes(peer_idx); + if current_codes.is_empty() || held_codes.is_empty() { + return Vec::new(); + } + let current_code_set = current_codes.into_iter().collect::>(); + let mut recognized_codes = held_codes + .iter() + .copied() + .filter(|code| current_code_set.contains(code)) + .collect::>(); + recognized_codes.sort_unstable(); + recognized_codes.dedup(); + recognized_codes + } + fn peer_needs_nonce_refill(&self, peer: &str, peer_idx: u16) -> bool { self.normalized_remote_held_nonce_codes(peer, peer_idx) .len() @@ -2946,16 +3014,42 @@ impl SigningDevice { .collect()) } - fn ping_payload(&mut self, peer: &str, peer_idx: u16) -> Result { + fn ping_payload( + &mut self, + peer: &str, + peer_idx: u16, + recognized_peer_nonce_codes: Option>, + ) -> Result { Ok(PingPayload { version: 2, advertised_nonces: self.advertised_nonces_for_peer(peer, peer_idx)?, held_peer_nonce_codes: self.state.nonce_pool.incoming_nonce_codes(peer_idx), + recognized_peer_nonce_codes: Some( + recognized_peer_nonce_codes + .unwrap_or_else(|| self.normalized_remote_held_nonce_codes(peer, peer_idx)), + ), policy_profile: Some(self.local_policy_profile_for(peer)?), nonce_pool_generation: self.state.nonce_pool.generation(), }) } + fn apply_ping_nonce_payload(&mut self, peer_idx: u16, payload: &PingPayload) { + if let Some(recognized_codes) = payload.recognized_peer_nonce_codes.as_ref() { + let mut retained = recognized_codes.clone(); + retained.extend(payload.advertised_nonces.iter().map(|nonce| nonce.code)); + retained.sort_unstable(); + retained.dedup(); + self.state + .nonce_pool + .retain_incoming_codes(peer_idx, &retained); + } + if !payload.advertised_nonces.is_empty() { + self.state + .nonce_pool + .store_incoming(peer_idx, payload.advertised_nonces.clone()); + } + } + fn reject_request( &self, peer: &str, @@ -3457,6 +3551,143 @@ mod tests { let ready_idx = decode_member_index(&fixture.group.members, &ready_peer).expect("ready idx"); assert_eq!(session.members, vec![fixture.local_share.idx, ready_idx]); + + let responder_peers = fixture + .group + .members + .iter() + .filter(|member| member.idx != ready_share.idx) + .map(|member| hex::encode(&member.pubkey[1..])) + .collect::>(); + let mut responder = SigningDevice::new( + fixture.group.clone(), + ready_share, + responder_peers, + peer_state, + DeviceConfig::default(), + ) + .expect("ready peer signer"); + let responder_effects = responder + .apply(SignerInput::ProcessEvent { + event: outbound[0].clone(), + }) + .expect("process sign request"); + assert_eq!(responder_effects.outbound.len(), 1); + let sign_served_peer = responder_effects + .completions + .iter() + .find_map(|completion| match completion { + CompletedOperation::SignServed { peer, .. } => Some(peer.clone()), + _ => None, + }) + .expect("expected sign-served completion"); + assert_eq!(sign_served_peer, fixture.signer.local_pubkey32()); + } + + #[test] + fn ping_refill_prunes_stale_peer_nonces_before_signing() { + let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); + let peer = fixture.signer.peers[0].clone(); + let peer_share = share_for_peer(&fixture.group, &fixture.shares, &peer); + + let mut stale_peer_state = + DeviceState::new(peer_share.idx, *peer_share.seckey.expose_bytes()); + let stale_nonces = stale_peer_state + .nonce_pool + .generate_for_peer( + fixture.local_share.idx, + 10, + &stale_peer_state.secrets.nonce_pool_secret, + ) + .expect("generate stale peer nonces"); + fixture + .signer + .state + .nonce_pool + .store_incoming(peer_share.idx, stale_nonces); + assert!(fixture.signer.state.nonce_pool.can_sign(peer_share.idx)); + + let mut current_peer = build_peer_signer(&fixture.group, &peer_share); + let ping = fixture + .signer + .apply(SignerInput::BeginPing { peer: peer.clone() }) + .expect("begin sync ping") + .outbound; + assert_eq!(ping.len(), 1); + + let response = current_peer + .process_event(&ping[0]) + .expect("peer processes sync ping"); + assert_eq!(response.len(), 1); + fixture + .signer + .process_event(&response[0]) + .expect("requester processes sync response"); + + let request = fixture + .signer + .initiate_sign([0x66; 32]) + .expect("initiate sign after sync"); + assert_eq!(request.len(), 1); + + let effects = current_peer + .apply(SignerInput::ProcessEvent { + event: request[0].clone(), + }) + .expect("current peer processes sign request"); + assert_eq!(effects.outbound.len(), 1); + } + + #[test] + fn initiate_sign_prefers_newly_advertised_peer_nonces() { + let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); + let peer = fixture.signer.peers[0].clone(); + let peer_share = share_for_peer(&fixture.group, &fixture.shares, &peer); + + let mut stale_peer_state = + DeviceState::new(peer_share.idx, *peer_share.seckey.expose_bytes()); + let stale_nonces = stale_peer_state + .nonce_pool + .generate_for_peer( + fixture.local_share.idx, + 10, + &stale_peer_state.secrets.nonce_pool_secret, + ) + .expect("generate stale peer nonces"); + fixture + .signer + .state + .nonce_pool + .store_incoming(peer_share.idx, stale_nonces); + + let mut current_peer = build_peer_signer(&fixture.group, &peer_share); + let fresh_nonces = current_peer + .state + .nonce_pool + .generate_for_peer( + fixture.local_share.idx, + 10, + ¤t_peer.state.secrets.nonce_pool_secret, + ) + .expect("generate fresh peer nonces"); + fixture + .signer + .state + .nonce_pool + .store_incoming(peer_share.idx, fresh_nonces); + + let request = fixture + .signer + .initiate_sign([0x77; 32]) + .expect("initiate sign with mixed stale and fresh nonces"); + assert_eq!(request.len(), 1); + + let effects = current_peer + .apply(SignerInput::ProcessEvent { + event: request[0].clone(), + }) + .expect("current peer processes sign request"); + assert_eq!(effects.outbound.len(), 1); } #[test] @@ -3563,6 +3794,66 @@ mod tests { assert_eq!(status.nonce_history[1].held, 7); } + #[test] + fn peer_status_reports_latest_response_latency() { + let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); + let request_id = "req-latency".to_string(); + let sender = fixture.signer.peers[0].clone(); + let now = now_unix_secs(); + fixture + .signer + .state + .peer_last_seen + .insert(sender.clone(), now); + fixture.signer.state.pending_operations.insert( + request_id.clone(), + PendingOperation { + op_type: PendingOpType::Ping, + request_id: request_id.clone(), + started_at: now.saturating_sub(2), + timeout_at: now + 30, + target_peers: vec![sender.clone()], + threshold: 1, + collected_responses: vec![], + context: PendingOpContext::PingRequest, + }, + ); + fixture + .signer + .state + .op_started_ms + .insert(request_id.clone(), now_unix_millis().saturating_sub(2_000)); + + let response = BridgeEnvelope { + request_id, + sent_at: now, + payload: BridgePayload::PingResponse(PingPayloadWire::from(PingPayload { + version: 2, + advertised_nonces: Vec::new(), + held_peer_nonce_codes: Vec::new(), + recognized_peer_nonce_codes: None, + policy_profile: None, + nonce_pool_generation: fixture.signer.state.nonce_pool.generation(), + })), + }; + + fixture + .signer + .match_pending_response(&response, &sender) + .expect("match ping response"); + + let status = fixture + .signer + .peer_status() + .into_iter() + .find(|entry| entry.pubkey == sender) + .expect("peer status"); + let latency_ms = status.last_response_latency_ms.expect("latency"); + + assert!(latency_ms >= 2_000); + assert!(latency_ms < 3_500); + } + #[test] fn telemetry_rings_are_bounded() { let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); @@ -3709,6 +4000,8 @@ mod tests { DeviceConfig::default(), ) .expect("requester signer"); + let requester_pubkey = + decode_member_pubkey(&inviter.group, requester_share.idx).expect("requester pubkey"); let ping = requester .apply(SignerInput::BeginPing { @@ -3717,11 +4010,22 @@ mod tests { .expect("begin ping") .outbound; assert_eq!(ping.len(), 1); - let ping_response = inviter + let ping_effects = inviter .signer - .process_event(&ping[0]) + .apply(SignerInput::ProcessEvent { + event: ping[0].clone(), + }) .expect("process ping"); - assert_eq!(ping_response.len(), 1); + assert_eq!(ping_effects.outbound.len(), 1); + let ping_served_peer = ping_effects + .completions + .iter() + .find_map(|completion| match completion { + CompletedOperation::PingServed { peer, .. } => Some(peer.clone()), + _ => None, + }) + .expect("expected ping-served completion"); + assert_eq!(ping_served_peer, requester_pubkey); let onboard = requester .apply(SignerInput::BeginOnboard { @@ -3739,8 +4043,6 @@ mod tests { assert_eq!(inviter_effects.outbound.len(), 1); // The responder records an OnboardServed completion carrying the // requester's x-only pubkey so hosts can mark its share onboarded. - let requester_pubkey = - decode_member_pubkey(&inviter.group, requester_share.idx).expect("requester pubkey"); let served_peer = inviter_effects .completions .iter() @@ -3773,6 +4075,11 @@ mod tests { let peer = fixture.signer.peers[0].clone(); let peer_share = share_for_peer(&fixture.group, &fixture.shares, &peer); let mut peer_signer = build_peer_signer(&fixture.group, &peer_share); + fixture + .signer + .state + .peer_last_seen + .insert(peer.clone(), now_unix_secs()); let first_ping = fixture .signer @@ -3788,6 +4095,11 @@ mod tests { }) .expect("expire timed-out ping"); assert_eq!(expired.failures.len(), 1); + assert_eq!( + expired.failures[0].failed_peer.as_deref(), + Some(peer.as_str()) + ); + assert!(!fixture.signer.state.peer_last_seen.contains_key(&peer)); let second_ping = fixture .signer @@ -3993,6 +4305,78 @@ mod tests { assert_eq!(failure.failed_peer, Some(locked_peer)); } + #[test] + fn inbound_sign_request_failure_reports_sign_op_type() { + let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); + let locked_peer = fixture.signer.peers[0].clone(); + let peer_share = share_for_peer(&fixture.group, &fixture.shares, &locked_peer); + let mut peer_signer = build_peer_signer(&fixture.group, &peer_share); + + let advertised = peer_signer + .state + .nonce_pool + .generate_for_peer( + fixture.local_share.idx, + 10, + &peer_signer.state.secrets.nonce_pool_secret, + ) + .expect("generate peer nonces"); + fixture + .signer + .state + .nonce_pool + .store_incoming(peer_share.idx, advertised.clone()); + + let request = fixture + .signer + .initiate_sign([0x55; 32]) + .expect("initiate sign"); + assert_eq!(request.len(), 1); + let envelope = + decode_envelope_for_local(&peer_share, fixture.signer.local_pubkey32(), &request[0]); + let BridgePayload::SignRequest(wire) = envelope.payload.clone() else { + panic!("expected sign request"); + }; + let session = SignSessionPackage::try_from(wire).expect("decode sign request"); + let peer_nonce_set = session + .nonces + .as_ref() + .and_then(|sets| sets.iter().find(|entry| entry.idx == peer_share.idx)) + .expect("peer nonce set"); + let mut codes_by_hash: Vec = vec![[0u8; 32]; session.hashes.len()]; + for entry in &peer_nonce_set.entries { + codes_by_hash[entry.hash_index as usize] = entry.code; + } + + peer_signer + .state + .nonce_pool + .take_outgoing_signing_nonces_many(fixture.local_share.idx, &codes_by_hash) + .expect("spend referenced peer nonce"); + + let outbound = peer_signer + .process_event(&request[0]) + .expect("inbound failure is surfaced through failures"); + assert!(outbound.is_empty()); + + let failures = peer_signer.take_failures(); + assert_eq!(failures.len(), 1); + let failure = &failures[0]; + assert_eq!(failure.request_id, envelope.request_id); + assert!(matches!(failure.op_type, PendingOpType::Sign)); + assert_eq!(failure.code, OperationFailureCode::PeerRejected); + assert!( + failure + .message + .to_ascii_lowercase() + .contains("nonce unavailable") + ); + assert_eq!( + failure.failed_peer.as_deref(), + Some(fixture.signer.local_pubkey32()) + ); + } + #[test] fn inbound_onboard_request_rejects_unsupported_version() { let mut fixture = fixture(PeerSelectionStrategy::DeterministicSorted); @@ -4094,6 +4478,7 @@ mod tests { version: 2, advertised_nonces: Vec::new(), held_peer_nonce_codes: Vec::new(), + recognized_peer_nonce_codes: None, policy_profile: None, nonce_pool_generation: bifrost_core::nonce::UNKNOWN_POOL_GENERATION, })), @@ -4479,6 +4864,7 @@ mod tests { version: 2, advertised_nonces: Vec::new(), held_peer_nonce_codes: Vec::new(), + recognized_peer_nonce_codes: None, policy_profile: None, nonce_pool_generation: bifrost_core::nonce::UNKNOWN_POOL_GENERATION, })), @@ -4624,10 +5010,22 @@ mod tests { let peer_share = share_for_peer(&fixture.group, &fixture.shares, &peer); let mut peer_signer = build_peer_signer(&fixture.group, &peer_share); - let responses = peer_signer - .process_event(&outbound[0]) + let peer_effects = peer_signer + .apply(SignerInput::ProcessEvent { + event: outbound[0].clone(), + }) .expect("peer processes ecdh"); + let responses = peer_effects.outbound; assert_eq!(responses.len(), 1); + let ecdh_served_peer = peer_effects + .completions + .iter() + .find_map(|completion| match completion { + CompletedOperation::EcdhServed { peer, .. } => Some(peer.clone()), + _ => None, + }) + .expect("expected ecdh-served completion"); + assert_eq!(ecdh_served_peer, fixture.signer.local_pubkey32()); let response = decode_envelope_for_local(&fixture.local_share, &peer, &responses[0]); let matched = fixture