diff --git a/crates/spt-daemon/src/diag294.rs b/crates/spt-daemon/src/diag294.rs new file mode 100644 index 00000000..a37bf9ae --- /dev/null +++ b/crates/spt-daemon/src/diag294.rs @@ -0,0 +1,32 @@ +//! Unlanded releases#294 diagnostic arm. Never requirement evidence. +//! All timestamps share this process's Instant epoch; stdout recording has observer cost. +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::OnceLock; +use std::time::Instant; + +static CLOCK: OnceLock<(Instant, String)> = OnceLock::new(); +static SEQUENCE: AtomicU64 = AtomicU64::new(0); + +pub fn enable(test: &str) { + CLOCK.get_or_init(|| (Instant::now(), test.to_owned())); +} + +pub fn enabled() -> bool { + CLOCK.get().is_some() +} + +pub fn emit(stage: &str, fields: serde_json::Value) { + let Some((epoch, test)) = CLOCK.get() else { + return; + }; + let elapsed_us = epoch.elapsed().as_micros(); + let sequence = SEQUENCE.fetch_add(1, Ordering::Relaxed); + println!( + "IR294 {}", + serde_json::json!({ + "stage": stage, "elapsed_us": elapsed_us, "sequence": sequence, + "clock": "process_instant_since_diag_enable", "pid": std::process::id(), + "test": test, "fields": fields, + }) + ); +} diff --git a/crates/spt-daemon/src/lib.rs b/crates/spt-daemon/src/lib.rs index 359cc206..8a9cf5f7 100644 --- a/crates/spt-daemon/src/lib.rs +++ b/crates/spt-daemon/src/lib.rs @@ -111,6 +111,7 @@ pub mod attach; pub mod attachment; pub mod autostart; pub mod brain; +pub mod diag294; pub mod brainproc; pub mod broker; pub mod codec; diff --git a/crates/spt-daemon/src/nethost.rs b/crates/spt-daemon/src/nethost.rs index cc10a235..f4484fd4 100644 --- a/crates/spt-daemon/src/nethost.rs +++ b/crates/spt-daemon/src/nethost.rs @@ -44,12 +44,13 @@ use spt_proto::identity::{Identity, PublicKey}; use crate::broker::SharedSend; use crate::effect::Minter; -use crate::seedproofx::{prove_membership, MembershipSource, RosterExchange}; use crate::frame::Envelope; use crate::msg::{ net_presence_event_envelope, net_stream_data_envelope, net_stream_eof_envelope, - NetPresenceEvent, NetStreamInfo, PRESENCE_CONNECTED, PRESENCE_DIAL_FAILED, PRESENCE_DISCONNECTED, + NetPresenceEvent, NetStreamInfo, PRESENCE_CONNECTED, PRESENCE_DIAL_FAILED, + PRESENCE_DISCONNECTED, }; +use crate::seedproofx::{prove_membership, MembershipSource, RosterExchange}; /// The reserved [`crate::effect::EffectKey`] session namespace for net-scoped /// effects (a dial has no PTY session). Broker session ids are minted from 1 @@ -674,7 +675,12 @@ impl PresenceLog { /// reply used to carry — now it rides the presence stream so the pump seeds /// `peer-addrs.json` from a non-blocking dial the same way. // [impl->REQ-CONV-1] - fn append_connected(&mut self, conn_id: u64, remote_id_hex: &str, remote_addr: serde_json::Value) { + fn append_connected( + &mut self, + conn_id: u64, + remote_id_hex: &str, + remote_addr: serde_json::Value, + ) { self.push(NetPresenceEvent { seq: 0, kind: PRESENCE_CONNECTED.to_string(), @@ -863,6 +869,8 @@ struct StreamEntry { /// Tables shared between the host's sync surface and its runtime tasks /// (accept loop, per-conn stream acceptors, per-stream read pumps). struct NetShared { + /// Endpoint identity for the opt-in #294 transport observations. + local_id: PublicKey, next_conn_id: AtomicU64, conns: Mutex>, next_stream_id: AtomicU64, @@ -873,8 +881,9 @@ struct NetShared { } impl NetShared { - fn new() -> Self { + fn new(local_id: PublicKey) -> Self { NetShared { + local_id, next_conn_id: AtomicU64::new(1), conns: Mutex::new(HashMap::new()), next_stream_id: AtomicU64::new(1), @@ -994,9 +1003,7 @@ impl DialPlan { stage, ) .await - .ok_or_else(|| { - io::Error::other("seed-proof failed: peer is not a subnet member") - })? + .ok_or_else(|| io::Error::other("seed-proof failed: peer is not a subnet member"))? } None => HashSet::new(), }; @@ -1107,6 +1114,32 @@ fn register_stream( lifetime: crate::msg::StreamLifetime, ) -> u64 { let id = shared.next_stream_id.fetch_add(1, Ordering::Relaxed); + // Observe open/accept before registration, with both local and wire IDs. + // These timestamps are host observations, not measurements of wire latency. + let diag_wire_id = if crate::diag294::enabled() { + match &send { + SendHalf::Quic(send) => Some(u64::from(send.id())), + SendHalf::Loopback(_) => None, + } + } else { + None + }; + if let Some(wire_id) = diag_wire_id { + crate::diag294::emit( + if initiated_locally { + "transport_outbound_opened" + } else { + "transport_inbound_accepted" + }, + serde_json::json!({ + "local_endpoint_id": shared.local_id.to_hex(), + "remote_endpoint_id": remote_id_hex, + "conn_id": conn_id, + "stream_id": id, + "quic_wire_stream_id": wire_id, + }), + ); + } let log = Arc::new(Mutex::new(StreamLog::new(id, shared.ring_cap))); let room = Arc::new(tokio::sync::Notify::new()); shared.streams.lock().unwrap().insert( @@ -1122,6 +1155,22 @@ fn register_stream( lifetime, }), ); + if let Some(wire_id) = diag_wire_id { + crate::diag294::emit( + if initiated_locally { + "transport_outbound_registered" + } else { + "transport_inbound_registered" + }, + serde_json::json!({ + "local_endpoint_id": shared.local_id.to_hex(), + "remote_endpoint_id": remote_id_hex, + "conn_id": conn_id, + "stream_id": id, + "quic_wire_stream_id": wire_id, + }), + ); + } // The read pump — the stream's single producer. A clean end (peer finish) // and a torn end (conn lost) both finish the log; the buffered chunks stay @@ -1301,8 +1350,8 @@ impl NetHost { )) .map_err(|e| io::Error::other(e.to_string()))?; let endpoint = Arc::new(endpoint); - let shared = Arc::new(NetShared::new()); let local_pub = cfg.identity.public_key(); + let shared = Arc::new(NetShared::new(local_pub)); let pair_rate = Arc::new(tokio::sync::Mutex::new( spt_net::net::pairing::ratelimit::PairingRateLimiter::new(), )); @@ -1674,7 +1723,10 @@ impl NetHost { /// the `apply_once` closure, so a concurrent deduped replay always finds it). /// The `minter` matches the journal key the broker built (ADR-0034 namespacing). pub fn record_dial_op(&self, minter: Minter, op_id: u64, conn_id: u64) { - self.dial_ops.lock().unwrap().insert((minter, op_id), conn_id); + self.dial_ops + .lock() + .unwrap() + .insert((minter, op_id), conn_id); } /// The connection a journaled dial `(minter, op_id)` opened, if this broker @@ -2185,7 +2237,9 @@ impl NetHost { // A replay write failure poisons the seat writer-side; the next // producer enqueue halt-and-removes it (the T1 discipline). // [impl->REQ-STREAMLOG-SUBSCRIBER-DISCIPLINE] - log.lock().unwrap().begin_attach(Arc::clone(&sub), from_seq)?; + log.lock() + .unwrap() + .begin_attach(Arc::clone(&sub), from_seq)?; Ok(()) } @@ -2312,7 +2366,9 @@ mod tests { // identity equivalence spt-net's endpoint module is built on. let pk = Identity::from_seed(&[7u8; 32]).public_key(); let id = EndpointId::from_bytes(&pk.to_bytes()).expect("valid ed25519 point"); - let relay: RelayUrl = "https://relay.example.invalid./".parse().expect("relay url"); + let relay: RelayUrl = "https://relay.example.invalid./" + .parse() + .expect("relay url"); let addr = EndpointAddr::new(id) .with_relay_url(relay.clone()) .with_ip_addr("10.0.0.7:4711".parse().unwrap()); @@ -2323,7 +2379,11 @@ mod tests { Some(relay.to_string().as_str()), "the relay path is found in what iroh ACTUALLY emits: {json}" ); - assert_eq!(addr_relay_urls(&json).len(), 1, "the direct path is not a relay: {json}"); + assert_eq!( + addr_relay_urls(&json).len(), + 1, + "the direct path is not a relay: {json}" + ); // A relay-less address reports None — the honest diagnosis, not an error. let direct_only = serde_json::to_value( @@ -2334,7 +2394,10 @@ mod tests { // Junk / absent degrades to None, never a panic (status surface). assert_eq!(addr_home_relay(&serde_json::Value::Null), None); - assert_eq!(addr_home_relay(&serde_json::json!({"addrs": "not-an-array"})), None); + assert_eq!( + addr_home_relay(&serde_json::json!({"addrs": "not-an-array"})), + None + ); } fn hermetic(identity: &Identity) -> NetConfig { @@ -2471,7 +2534,9 @@ mod tests { assert_eq!(host.conn_count(), 1, "one loopback conn held"); // open_stream mints the cross-wired pair: operator row + peer row. - let op_stream = host.open_stream(c1, crate::msg::StreamLifetime::Durable).expect("open loopback stream"); + let op_stream = host + .open_stream(c1, crate::msg::StreamLifetime::Durable) + .expect("open loopback stream"); let infos = host.stream_infos(); assert_eq!(infos.len(), 2, "operator + peer rows"); let op = infos @@ -2574,7 +2639,11 @@ mod tests { ret.append(&[i]); } let (bytes, finished) = ret.drain(); - assert_eq!(bytes, (0..8u8).collect::>(), "retentive loses nothing, in order"); + assert_eq!( + bytes, + (0..8u8).collect::>(), + "retentive loses nothing, in order" + ); assert!(!finished); // …and a second drain is empty (the cursor consumed them exactly once). assert_eq!(ret.drain().0, Vec::::new()); @@ -2585,7 +2654,11 @@ mod tests { for i in 0..8u8 { ord.append(&[i]); } - assert_eq!(ord.drain().0, vec![5, 6, 7], "ordinary keeps only the last cap chunks"); + assert_eq!( + ord.drain().0, + vec![5, 6, 7], + "ordinary keeps only the last cap chunks" + ); } // [unit->REQ-SHELL-4] the loopback tunnel pair under BACKPRESSURE (M11-W3, @@ -2715,11 +2788,18 @@ mod tests { assert_eq!(log.opener_line(), None, "no newline yet: not pinned"); log.append(b"est\",\"session_id\":3}\n{\"kind\":\"input\"}\n"); let want = &b"{\"kind\":\"request\",\"session_id\":3}"[..]; - assert_eq!(log.opener_line().as_deref(), Some(want), "split line reassembled + pinned"); + assert_eq!( + log.opener_line().as_deref(), + Some(want), + "split line reassembled + pinned" + ); for i in 0..64u32 { log.append(format!("{{\"kind\":\"input\",\"n\":{i}}}\n").as_bytes()); } - assert!(log.floor_seq() > 0, "the bounded ring rolled: seq 0 evicted"); + assert!( + log.floor_seq() > 0, + "the bounded ring rolled: seq 0 evicted" + ); assert_eq!( log.opener_line().as_deref(), Some(want), @@ -2735,7 +2815,11 @@ mod tests { fn opener_capture_gives_up_bounded_on_a_newline_less_stream() { let mut log = StreamLog::new(8, 4); log.append(&vec![b'x'; OPENER_PIN_MAX + 1]); - assert_eq!(log.opener_line(), None, "over the cap without a newline: gave up"); + assert_eq!( + log.opener_line(), + None, + "over the cap without a newline: gave up" + ); log.append(b"late-line\n"); assert_eq!(log.opener_line(), None, "Oversize is terminal"); } @@ -2822,7 +2906,10 @@ mod tests { String::from_utf8_lossy(&got).contains("\"outcome\":\"edge\""), "the reply flushed through the retired row (got {got:?})" ); - assert!(finished, "the reply's FIN reached the requester side (clean end, not torn)"); + assert!( + finished, + "the reply's FIN reached the requester side (clean end, not torn)" + ); } // [unit->REQ-REDISPATCH-FINISHED-RETIRE] a dead CONNECTION retires its @@ -2834,7 +2921,9 @@ mod tests { let a = NetHost::start(hermetic(&Identity::generate())).expect("host a"); let b = NetHost::start(hermetic(&Identity::generate())).expect("host b"); let (conn_id, _) = a.dial(b.addr()).expect("dial"); - let sid = a.open_stream(conn_id, crate::msg::StreamLifetime::Durable).expect("open stream"); + let sid = a + .open_stream(conn_id, crate::msg::StreamLifetime::Durable) + .expect("open stream"); a.send_stream(sid, b"{\"hello\":1}\n", false).expect("send"); // B's acceptor registers the peer row (async — poll). let mut saw = false; @@ -2857,7 +2946,10 @@ mod tests { } std::thread::sleep(Duration::from_millis(10)); } - assert!(swept, "the closed-watcher swept the dead conn's stream rows"); + assert!( + swept, + "the closed-watcher swept the dead conn's stream rows" + ); } // ── ADR-0038 Amendment fixes 2+3+4 — the subscriber-seat discipline at @@ -2933,8 +3025,14 @@ mod tests { // The next producer append sees the poisoned seat: halt-and-remove + // lease cancel — never another write attempt against the dead conn. log.append(b"two"); - assert!(log.subscriber.is_none(), "halt-and-remove at the append site"); - assert!(lease.is_canceled(), "poison cancels the serve lease (fix 4)"); + assert!( + log.subscriber.is_none(), + "halt-and-remove at the append site" + ); + assert!( + lease.is_canceled(), + "poison cancels the serve lease (fix 4)" + ); // Later appends stay seatless (no reinstall, no panic, no re-feed). log.append(b"three"); @@ -2944,7 +3042,8 @@ mod tests { // RENEWS the lease — the old worker's handle stays canceled, fresh // sends serve again. One poison never permanently dead-ends a stream. let (fresh, _cf, _rf) = seat_socket_pair(Duration::from_secs(5)); - log.begin_attach(Arc::clone(&fresh), 0).expect("new generation attaches"); + log.begin_attach(Arc::clone(&fresh), 0) + .expect("new generation attaches"); assert!( !log.lease.is_canceled(), "a new subscriber generation renews the serve lease" @@ -2976,8 +3075,14 @@ mod tests { log.finish(); assert!(log.finished, "the read side still records its clean end"); - assert!(log.subscriber.is_none(), "halt-and-remove at the finish site"); - assert!(lease.is_canceled(), "poison at finish cancels the lease too"); + assert!( + log.subscriber.is_none(), + "halt-and-remove at the finish site" + ); + assert!( + lease.is_canceled(), + "poison at finish cancels the lease too" + ); } // [unit->REQ-STREAMLOG-SUBSCRIBER-DISCIPLINE] @@ -2998,7 +3103,8 @@ mod tests { log.begin_attach(Arc::clone(&a), 0).expect("first attach"); // Healthy prior → displaced (the legit brain-swap path). - log.begin_attach(Arc::clone(&b), 0).expect("healthy displacement"); + log.begin_attach(Arc::clone(&b), 0) + .expect("healthy displacement"); assert!(log.subscriber.as_ref().unwrap().is(&b)); // Poisoned but NOT gone (writer still draining) → refuse WouldBlock. @@ -3011,8 +3117,13 @@ mod tests { assert!(err.to_string().contains("subscriber busy"), "{err}"); // Fully gone → the replacement installs. - log.subscriber.as_ref().unwrap().done.store(true, Ordering::Release); - log.begin_attach(Arc::clone(&a), 0).expect("gone prior admits the replacement"); + log.subscriber + .as_ref() + .unwrap() + .done + .store(true, Ordering::Release); + log.begin_attach(Arc::clone(&a), 0) + .expect("gone prior admits the replacement"); assert!(log.subscriber.as_ref().unwrap().is(&a)); } @@ -3084,7 +3195,8 @@ mod tests { // Pin the wire: nothing the writer does lands until we release. let pin = sub.pin_gate_for_test(); - log.begin_attach(Arc::clone(&sub), 0).expect("attach with a pending replay"); + log.begin_attach(Arc::clone(&sub), 0) + .expect("attach with a pending replay"); // Live appends land while the replay is still entirely undrained. for i in 6..10u8 { log.append(&[i]); // live seqs 6..=9 @@ -3096,7 +3208,11 @@ mod tests { for want in 0u64..10 { let env = crate::codec::read_frame(&mut client).expect("frame on the wire"); assert_eq!(env.kind, crate::msg::KIND_NET_STREAM_DATA); - let seq = env.payload.get("seq").and_then(|v| v.as_u64()).expect("seq"); + let seq = env + .payload + .get("seq") + .and_then(|v| v.as_u64()) + .expect("seq"); assert_eq!( seq, want, "replay-then-live seq order must be structural; a live frame \ @@ -3129,7 +3245,8 @@ mod tests { log.append(b"queued-for-the-displaced-writer"); // HEALTHY displacement: B takes the seat (the brain-swap contract). - log.begin_attach(Arc::clone(&b), 0).expect("healthy displacement"); + log.begin_attach(Arc::clone(&b), 0) + .expect("healthy displacement"); let lease_b = Arc::clone(&log.lease); assert!( !Arc::ptr_eq(&lease_a, &lease_b), @@ -3210,7 +3327,8 @@ mod tests { let (owner, _shell) = host.open_loopback_pair().expect("open pair"); let lease = host.stream_lease(owner).expect("lease handle"); assert!(!lease.is_canceled()); - host.send_stream(owner, b"ok", false).expect("a live lease serves"); + host.send_stream(owner, b"ok", false) + .expect("a live lease serves"); lease.cancel(); let err = host @@ -3271,7 +3389,8 @@ mod tests { let _ = shell; let (a, _ca, _ra) = seat_socket_pair(Duration::from_secs(5)); - host.subscribe_stream(owner, Arc::clone(&a), 0).expect("A subscribes"); + host.subscribe_stream(owner, Arc::clone(&a), 0) + .expect("A subscribes"); let (_, seats) = host.stream_counts(); assert_eq!(seats, 1, "A's seat installed"); @@ -3291,8 +3410,10 @@ mod tests { // Identity rule: B displaces A; A's late unsubscribe must NOT evict B. let (a2, _ca2, _ra2) = seat_socket_pair(Duration::from_secs(5)); let (b, _cb, _rb) = seat_socket_pair(Duration::from_secs(5)); - host.subscribe_stream(owner, Arc::clone(&a2), 0).expect("A2 subscribes"); - host.subscribe_stream(owner, Arc::clone(&b), 0).expect("B displaces A2"); + host.subscribe_stream(owner, Arc::clone(&a2), 0) + .expect("A2 subscribes"); + host.subscribe_stream(owner, Arc::clone(&b), 0) + .expect("B displaces A2"); assert!( !host.unsubscribe_stream(owner, &a2), "a displaced caller's release is a no-op" @@ -3315,7 +3436,8 @@ mod tests { // the retired_row_still_flushes_a_late_reply contract above). let (owner, shell) = host.open_loopback_pair().expect("pair 1"); let (sub, _c, _r) = seat_socket_pair(Duration::from_secs(5)); - host.subscribe_stream(shell, Arc::clone(&sub), 0).expect("subscribe"); + host.subscribe_stream(shell, Arc::clone(&sub), 0) + .expect("subscribe"); let (rows_before, seats_before) = host.stream_counts(); assert_eq!(seats_before, 1); assert!(host.retire_stream(shell)); diff --git a/crates/spt-daemon/src/sync.rs b/crates/spt-daemon/src/sync.rs index b5067403..f0641e09 100644 --- a/crates/spt-daemon/src/sync.rs +++ b/crates/spt-daemon/src/sync.rs @@ -342,7 +342,30 @@ pub fn request_sync( cs: &ContextStore, scratch_dir: &Path, ) -> io::Result { - let opened = brain.net_open_stream(conn_id, Some(open_op))?; + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_open_stream.enter", + serde_json::json!({ + "open_op_seq": open_op.seq, + "conn_id": conn_id, + "stream_id": null, + }), + ); + } + let opened = brain.net_open_stream(conn_id, Some(open_op)); + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_open_stream.exit", + serde_json::json!({ + "open_op_seq": open_op.seq, + "conn_id": conn_id, + "stream_id": opened.as_ref().ok().map(|opened| opened.stream_id), + "outcome": if opened.is_ok() { "ok" } else { "error" }, + "error": opened.as_ref().err().map(|error| error.to_string()), + }), + ); + } + let opened = opened?; let stream_id = opened.stream_id; let out = request_sync_on(brain, stream_id, refs, open_op, cs, scratch_dir); // Requester-side lifetime bound (ADR-0040 decisions 1+5): this pull's @@ -373,12 +396,60 @@ fn request_sync_on( cs: &ContextStore, scratch_dir: &Path, ) -> io::Result { - brain.net_stream_subscribe(stream_id, 0)?; + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_stream_subscribe.enter", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + }), + ); + } + let subscribed = brain.net_stream_subscribe(stream_id, 0); + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_stream_subscribe.exit", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + "outcome": if subscribed.is_ok() { "ok" } else { "error" }, + "error": subscribed.as_ref().err().map(|error| error.to_string()), + }), + ); + } + subscribed?; let sync_id = format!("pull-{}", open_op.seq); let mut have_tips = BTreeMap::new(); for r in refs { - if let Some(tip) = cs.branch_store().tip(r)? { + if crate::diag294::enabled() { + crate::diag294::emit( + "request.branch_tip.enter", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + "ref_name": r, + }), + ); + } + let tip = cs.branch_store().tip(r); + if crate::diag294::enabled() { + crate::diag294::emit( + "request.branch_tip.exit", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + "ref_name": r, + "outcome": match &tip { + Ok(Some(_)) => "found", + Ok(None) => "missing", + Err(_) => "error", + }, + "error": tip.as_ref().err().map(|error| error.to_string()), + }), + ); + } + if let Some(tip) = tip? { have_tips.insert(r.clone(), tip); } } @@ -388,7 +459,28 @@ fn request_sync_on( have_tips, } .encode_line(); - brain.net_stream_send(stream_id, &line, None, true)?; + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_stream_send.enter", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + }), + ); + } + let sent = brain.net_stream_send(stream_id, &line, None, true); + if crate::diag294::enabled() { + crate::diag294::emit( + "request.net_stream_send.exit", + serde_json::json!({ + "open_op_seq": open_op.seq, + "stream_id": stream_id, + "outcome": if sent.is_ok() { "ok" } else { "error" }, + "error": sent.as_ref().err().map(|error| error.to_string()), + }), + ); + } + sent?; std::fs::create_dir_all(scratch_dir)?; let bundle: PathBuf = scratch_dir.join(format!("recv-{}.bundle", safe_stem(&sync_id))); @@ -557,7 +649,9 @@ pub fn reconcile_after_sync( // literal, not merely a matter of role conflicts being rare (REQ-EP-7). // [impl->REQ-EP-7] if std::path::Path::new(&file).file_name() - == Some(std::ffi::OsStr::new(spt_store::contextstore::LIVE_ROLE_FILE)) + == Some(std::ffi::OsStr::new( + spt_store::contextstore::LIVE_ROLE_FILE, + )) { outcomes.push((branch, file, ReconcileOutcome::RoleExcluded)); continue; @@ -694,7 +788,10 @@ mod tests { b"\nREMOTE role", ) .unwrap(); - assert_eq!(cs.list_conflicts(&wt, Some("live-role.md")).unwrap().len(), 1); + assert_eq!( + cs.list_conflicts(&wt, Some("live-role.md")).unwrap().len(), + 1 + ); let report = SyncPullReport { applied: vec![ApplyReport { @@ -718,9 +815,8 @@ mod tests { "[adapter]\nname=\"mock\"\nversion=\"1\"\nmin_spt_core_version=\"1\"\n\n\ [session.psyche_resume]\ncommand='{cmd}'\n" ); - let rt = spt_runtime::ManifestRuntime::new( - spt_runtime::Manifest::from_toml_str(&toml).unwrap(), - ); + let rt = + spt_runtime::ManifestRuntime::new(spt_runtime::Manifest::from_toml_str(&toml).unwrap()); let outcomes = reconcile_after_sync( &rt, diff --git a/crates/spt-daemon/tests/sync.rs b/crates/spt-daemon/tests/sync.rs index 07450ed3..97aa034c 100644 --- a/crates/spt-daemon/tests/sync.rs +++ b/crates/spt-daemon/tests/sync.rs @@ -28,6 +28,7 @@ use std::thread; use std::time::Duration; use spt_daemon::brain::{Brain, BrokerEvent}; +use spt_daemon::diag294; use spt_daemon::effect::{MintedOp, Minter}; use spt_daemon::nethost::{NetConfig, NetHost}; use spt_daemon::sync::{ @@ -48,6 +49,7 @@ use spt_test_support::TestHome; /// Hold a scoped temp home for each test (the reconcile turn's merged write /// stamps the node identity under SPT_HOME). fn init_home() -> TestHome { + diag294::enable(thread::current().name().unwrap_or("sync")); TestHome::new() } @@ -76,6 +78,12 @@ fn hermetic() -> NetConfig { fn net_broker(name: &str, dir: &std::path::Path) -> Arc { let host = NetHost::start(hermetic()).expect("net host start"); + diag294::emit( + "broker.identity", + serde_json::json!({ + "broker": name, "endpoint_id": host.node_id_hex(), + }), + ); let broker = Broker::bind_in_with_net(name, dir.join("effects.log"), Some(host)).expect("bind broker"); let serve = Arc::clone(&broker); @@ -95,17 +103,166 @@ fn connect_retry(name: &str) -> Brain { panic!("brain could not connect"); } -fn wait_for_stream_except(brain: &mut Brain, skip: &[u64]) -> (u64, String) { - // 10s budget: polls exit early when the stream appears, so the long bound - // only pays on a loaded runner (2s flaked on gravity under a parallel - // full-workspace run — 2026-06-04). - for _ in 0..400 { +static DIAG_EXPIRY_SIGNAL: std::sync::OnceLock> = + std::sync::OnceLock::new(); + +// Diagnostic-arm control only: prove late visibility cannot rescue acceptance. +#[test] +fn diag_continuation_preserves_original_failure() { + let _home = init_home(); + let dir = tempfile::tempdir().unwrap(); + let (a, b) = (unique_name(), unique_name()); + let _broker_a = net_broker(&a, &dir.path().join("a")); + let _broker_b = net_broker(&b, &dir.path().join("b")); + let mut responder = connect_retry(&a); + let mut requester = connect_retry(&b); + let addr = responder.net_status().unwrap().addr; + let conn = requester + .net_dial(addr, Some(MintedOp::new(Minter::Cli, op()))) + .unwrap(); + let stream = requester + .net_open_stream(conn.conn_id, Some(MintedOp::new(Minter::Pump, op()))) + .unwrap(); + let (tx, rx) = std::sync::mpsc::channel(); + DIAG_EXPIRY_SIGNAL.set(tx).unwrap(); + let writer = thread::spawn(move || { + rx.recv_timeout(Duration::from_secs(90)) + .expect("acceptance expiry signal"); + requester + .net_stream_send(stream.stream_id, b"diagnostic-control", None, true) + .unwrap(); + }); + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + wait_for_stream_except(responder, &[], 0) + })); + writer.join().unwrap(); + let error = result + .err() + .expect("late visibility must not rescue the acceptance failure"); + assert_eq!( + error.downcast_ref::<&str>().copied(), + Some("sync stream never appeared at the responder's broker") + ); +} + +fn wait_for_stream_except(mut brain: Brain, skip: &[u64], exchange: u64) -> (Brain, (u64, String)) { + // Acceptance remains 400 observations, each unsuccessful observation followed + // by 25ms sleep. RPC and scheduling time are additional, not a 10s wall bound. + let started = std::time::Instant::now(); + let mut ipc = Duration::ZERO; + let mut sleeps = Duration::ZERO; + diag294::emit( + "poll.enter", + serde_json::json!({ + "exchange": exchange, "iteration_budget": 400, "skip": skip, + }), + ); + for iteration in 1..=400 { + let call = std::time::Instant::now(); let reply = brain.net_streams().expect("net-streams"); + ipc += call.elapsed(); if let Some(s) = reply.streams.iter().find(|s| !skip.contains(&s.stream_id)) { - return (s.stream_id, s.remote_id_hex.clone()); + diag294::emit( + "poll.success", + serde_json::json!({ + "exchange": exchange, "iterations": iteration, "iteration_budget": 400, + "poll_elapsed_us": started.elapsed().as_micros(), + "ipc_us": ipc.as_micros(), "sleep_us": sleeps.as_micros(), + "stream_id": s.stream_id, "remote_id": s.remote_id_hex, + "last_filtered_count": reply.streams.len(), + }), + ); + let found = (s.stream_id, s.remote_id_hex.clone()); + return (brain, found); } + let sleeping = std::time::Instant::now(); thread::sleep(Duration::from_millis(25)); + sleeps += sleeping.elapsed(); + } + diag294::emit( + "poll.expiry", + serde_json::json!({ + "exchange": exchange, "iterations": 400, "iteration_budget": 400, + "poll_elapsed_us": started.elapsed().as_micros(), + "ipc_us": ipc.as_micros(), "sleep_us": sleeps.as_micros(), + "acceptance": "failed", + }), + ); + if let Some(signal) = DIAG_EXPIRY_SIGNAL.get() { + let _ = signal.send(()); } + // Diagnostic ONLY: preserve the original failure even if later observed. + // The existing blocking carrier is moved, not replaced. An independent + // receiver deadline bounds how long the test waits even if an IPC wedges. + // On that watchdog path the observer thread may remain blocked until the + // test process exits; this is explicitly not proof of continued absence. + let continuation_started = std::time::Instant::now(); + let continuation_deadline = continuation_started + Duration::from_secs(30); + let skip = skip.to_vec(); + let (tx, rx) = std::sync::mpsc::sync_channel(1); + thread::spawn(move || { + let mut iterations = 0u64; + let mut ipc = Duration::ZERO; + let mut sleeps = Duration::ZERO; + while std::time::Instant::now() < continuation_deadline { + iterations += 1; + let call = std::time::Instant::now(); + let reply = brain.net_streams(); + ipc += call.elapsed(); + if std::time::Instant::now() >= continuation_deadline { + break; + } + match reply { + Ok(reply) => { + if let Some(s) = reply.streams.iter().find(|s| !skip.contains(&s.stream_id)) { + let _ = tx.send(serde_json::json!({ + "outcome": "late_eligible_observation", "stream_id": s.stream_id, + "remote_id": s.remote_id_hex, + "first_observed_poll_elapsed_us": started.elapsed().as_micros(), + "iterations": iterations, "ipc_us": ipc.as_micros(), + "sleep_us": sleeps.as_micros(), + })); + return; + } + } + Err(error) => { + let _ = tx.send(serde_json::json!({ + "outcome": "observer_error", "error": error.to_string(), + "iterations": iterations, "ipc_us": ipc.as_micros(), + "sleep_us": sleeps.as_micros(), + })); + return; + } + } + let sleeping = std::time::Instant::now(); + thread::sleep( + Duration::from_millis(25) + .min(continuation_deadline.saturating_duration_since(sleeping)), + ); + sleeps += sleeping.elapsed(); + } + let _ = tx.send(serde_json::json!({ + "outcome": "continuation_bound_reached", "iterations": iterations, + "ipc_us": ipc.as_micros(), "sleep_us": sleeps.as_micros(), + })); + }); + let observation = rx + .recv_timeout(continuation_deadline.saturating_duration_since(std::time::Instant::now())) + .unwrap_or_else(|error| { + serde_json::json!({ + "outcome": "observer_did_not_return_by_bound", "error": error.to_string(), + "continued_non_observation": "not_established", + }) + }); + diag294::emit( + "poll.continuation.exit", + serde_json::json!({ + "exchange": exchange, "continuation_budget_ms": 30000, + "continuation_elapsed_us": continuation_started.elapsed().as_micros(), + "poll_elapsed_us": started.elapsed().as_micros(), + "acceptance": "failed", "observation": observation, + }), + ); panic!("sync stream never appeared at the responder's broker"); } @@ -185,8 +342,34 @@ impl PullLink { let req_root = self.requester_root.clone(); let scratch_req = self.scratch.join("req"); let open_op = MintedOp::new(Minter::Pump, op()); + let exchange = open_op.seq; + static PULL_ORDINAL: AtomicU64 = AtomicU64::new(0); + let pull_ordinal = PULL_ORDINAL.fetch_add(1, Ordering::Relaxed) + 1; + diag294::emit( + "pull.identity", + serde_json::json!({ + "exchange": exchange, "pull_ordinal": pull_ordinal, + "requester": self.requester, "responder": self.responder, + "requester_conn_id": conn.conn_id, "refs": refs, + }), + ); let puller = thread::spawn(move || { - let cs = ContextStore::open_or_init_in(&req_root).unwrap(); + diag294::emit( + "requester.thread.enter", + serde_json::json!({"exchange": exchange}), + ); + diag294::emit( + "requester.store_init.enter", + serde_json::json!({"exchange": exchange}), + ); + let store = ContextStore::open_or_init_in(&req_root); + diag294::emit( + "requester.store_init.exit", + serde_json::json!({ + "exchange": exchange, "outcome": if store.is_ok() { "ok" } else { "error" }, + }), + ); + let cs = store.unwrap(); request_sync( &mut req_brain, conn.conn_id, @@ -198,7 +381,8 @@ impl PullLink { .expect("pull") }); - let (stream, origin) = wait_for_stream_except(&mut serve_brain, &skip); + let (mut serve_brain, (stream, origin)) = + wait_for_stream_except(serve_brain, &skip, exchange); let cs = ContextStore::open_or_init_in(&self.responder_root).unwrap(); let outcome = serve_sync( &mut serve_brain, @@ -551,8 +735,32 @@ fn torn_pull_recovers_by_repulling() { let root_b_clone = root_b.clone(); let scratch_req = scratch.join("req"); let open_op = MintedOp::new(Minter::Pump, op()); + let exchange = open_op.seq; + diag294::emit( + "pull.identity", + serde_json::json!({ + "exchange": exchange, "pull_ordinal": 0, + "requester": name_b, "responder": name_a, + "requester_conn_id": conn.conn_id, "refs": ["a-doyle"], + }), + ); let puller = thread::spawn(move || { - let cs = ContextStore::open_or_init_in(&root_b_clone).unwrap(); + diag294::emit( + "requester.thread.enter", + serde_json::json!({"exchange": exchange}), + ); + diag294::emit( + "requester.store_init.enter", + serde_json::json!({"exchange": exchange}), + ); + let store = ContextStore::open_or_init_in(&root_b_clone); + diag294::emit( + "requester.store_init.exit", + serde_json::json!({ + "exchange": exchange, "outcome": if store.is_ok() { "ok" } else { "error" }, + }), + ); + let cs = store.unwrap(); request_sync( &mut req_brain, conn.conn_id, @@ -563,7 +771,7 @@ fn torn_pull_recovers_by_repulling() { ) .expect_err("torn pull must error") }); - let (stream, _origin) = wait_for_stream_except(&mut serve_brain, &[]); + let (mut serve_brain, (stream, _origin)) = wait_for_stream_except(serve_brain, &[], exchange); serve_brain.net_stream_subscribe(stream, 0).unwrap(); let mut decoder = SyncDecoder::new(); let deadline = serve_brain.call_deadline();