diff --git a/crates/bridge-core/src/engine/host.rs b/crates/bridge-core/src/engine/host.rs index 48be334..e62fb52 100644 --- a/crates/bridge-core/src/engine/host.rs +++ b/crates/bridge-core/src/engine/host.rs @@ -437,10 +437,14 @@ impl Engine { } SessionEvent::Entries { entries } => self.on_entries(session_id, entries), SessionEvent::Turn { state } => { + // The agent may repeat the state it is in; only a change is + // worth a new session list. if let Some(run) = self.run_mut(session_id) { - run.turn = state; + if run.turn != state { + run.turn = state; + self.list_dirty = true; + } } - self.list_dirty = true; } SessionEvent::Ended { error, resume_lost } => { if let Some(run) = self.run_mut(session_id) { diff --git a/crates/bridge-core/src/engine/mod.rs b/crates/bridge-core/src/engine/mod.rs index c5d7d3b..dcf6474 100644 --- a/crates/bridge-core/src/engine/mod.rs +++ b/crates/bridge-core/src/engine/mod.rs @@ -153,6 +153,8 @@ pub struct Engine { stopped: bool, /// Publish the heartbeat once the current input is handled. list_dirty: bool, + /// The last session list's `rev`. + list_rev: u64, /// Store the registry once the current input is handled. registry_dirty: bool, /// Only `last_activity` changed since the registry was stored: written @@ -193,6 +195,7 @@ impl Engine { started: false, stopped: false, list_dirty: false, + list_rev: 0, registry_dirty: false, activity_dirty: false, } @@ -274,6 +277,7 @@ impl Engine { let last_seen = self.store.get(store_keys::LAST_SEEN).and_then(|v| v.trim().parse().ok()).unwrap_or(0); let processed: Vec = load_or_default(self.store.get(store_keys::PROCESSED_IDS), "processed event ids"); self.ingest = Ingest::new(last_seen, processed); + self.list_rev = self.store.get(store_keys::LIST_REV).and_then(|v| v.trim().parse().ok()).unwrap_or(0); let doc = self.store.get(store_keys::REGISTRY).map(|raw| RegistryDoc::load(&raw)).unwrap_or_default(); self.tombstones = Tombstones::new(doc.removed_sessions); @@ -532,6 +536,7 @@ impl Engine { let agents = self.catalog.descriptors(|a| self.agent_credentials(a)); let credentials = self.bridge_credentials(); let removed = self.tombstones.ids(); + let rev = self.next_list_rev(); let message = BridgeToPhone::Sessions(SessionListMsg { machine: self.config.machine.clone(), host: self.config.host_kind, @@ -545,7 +550,19 @@ impl Engine { removed_sessions: (!removed.is_empty()).then_some(removed), machine_offline: offline.then_some(true), direct: self.config.direct.clone(), + rev: Some(rev), }); self.publish_all(message); } + + /// The next session list's `rev`: the clock in ms, or one past the last + /// one when the clock has not moved past it (two lists in one ms, or a + /// clock set back). Stored, so a restart keeps it increasing. + fn next_list_rev(&mut self) -> u64 { + self.list_rev = self.now().max(self.list_rev + 1); + if let Err(err) = self.store.set(store_keys::LIST_REV, &self.list_rev.to_string()) { + log::warn!("[Engine] Could not store the session list revision: {err}"); + } + self.list_rev + } } diff --git a/crates/bridge-core/src/io.rs b/crates/bridge-core/src/io.rs index 194fd73..3151f0c 100644 --- a/crates/bridge-core/src/io.rs +++ b/crates/bridge-core/src/io.rs @@ -176,6 +176,9 @@ pub mod store_keys { pub const REGISTRY: &str = "registry"; pub const LAST_SEEN: &str = "lastSeenTimestamp"; pub const PROCESSED_IDS: &str = "processedEventIds"; + /// The last session list's `rev`, so the next one is greater even if + /// the clock went back across a restart. + pub const LIST_REV: &str = "listRev"; /// Keys whose values contain secrets. pub const SECRET_KEYS: [&str; 2] = [CREDENTIALS, PROVIDER_PROFILES]; } diff --git a/crates/bridge-core/tests/sessions.rs b/crates/bridge-core/tests/sessions.rs index ea83377..89669c3 100644 --- a/crates/bridge-core/tests/sessions.rs +++ b/crates/bridge-core/tests/sessions.rs @@ -56,6 +56,24 @@ fn the_heartbeat_advertises_the_direct_link() { assert_eq!(direct.endpoints, ["wss://10.0.0.2:7447"]); } +#[test] +fn every_session_list_has_a_greater_rev_than_the_last_even_across_a_restart() { + let mut rig = Rig::new(); + let first = last_heartbeat(&rig.messages()).rev.expect("a rev"); + // Two lists in the same millisecond still differ. + rig.host_up(); + let second = last_heartbeat(&rig.messages()).rev.unwrap(); + assert!(second > first, "{second} > {first}"); + rig.advance(60_000); + rig.send(serde_json::json!({ "type": "refresh-sessions" })); + let third = last_heartbeat(&rig.messages()).rev.unwrap(); + assert!(third > second, "{third} > {second}"); + // A restarted bridge whose clock starts behind keeps counting up. + let mut rig = rig.restart(); + let after = last_heartbeat(&rig.messages()).rev.unwrap(); + assert!(after > third, "{after} > {third}"); +} + #[test] fn without_a_paired_phone_nothing_is_published() { let mut rig = Rig::with(RigOptions { paired: false, ..Default::default() }); diff --git a/crates/bridge-runtime/src/relay.rs b/crates/bridge-runtime/src/relay.rs index 38f8666..3f4fe85 100644 --- a/crates/bridge-runtime/src/relay.rs +++ b/crates/bridge-runtime/src/relay.rs @@ -544,6 +544,7 @@ mod tests { removed_sessions: None, machine_offline: None, direct: None, + rev: None, }) } diff --git a/crates/client-core/src/stores/machines.rs b/crates/client-core/src/stores/machines.rs index ee2de24..f5cb0a1 100644 --- a/crates/client-core/src/stores/machines.rs +++ b/crates/client-core/src/stores/machines.rs @@ -301,6 +301,11 @@ pub struct MachineView { /// A grant sent to this bridge and not confirmed yet. #[serde(default, skip_serializing_if = "Option::is_none")] pub session_grant_sent: Option, + /// The `rev` of the newest session list applied: an older one, arriving + /// late over another path or replayed, is not applied over it. + #[serde(default, skip_serializing_if = "Option::is_none")] + #[specta(type = Option)] + pub list_rev: Option, /// Where the bridge says it can be reached directly (its heartbeat's /// `direct`), as last heard. #[serde(default, skip_serializing_if = "Option::is_none")] @@ -462,6 +467,7 @@ impl MachineView { mcp: BTreeMap::new(), session_grant: None, session_grant_sent: None, + list_rev: None, direct: None, direct_endpoints: Vec::new(), relays: Vec::new(), @@ -662,6 +668,13 @@ impl MachinesState { /// stored value; a field the wire carries always wins. The resurrection /// shield filters non-expired user-dismissed session ids out of `msg` /// BEFORE the merge. + /// Whether `msg` is no newer than the newest list applied for + /// `machine_pubkey`, so applying it would step the machine back. + pub fn is_outdated_list(&self, machine_pubkey: &str, msg: &SessionListMsg) -> bool { + let applied = self.machine(machine_pubkey).and_then(|m| m.list_rev); + matches!((msg.rev, applied), (Some(rev), Some(applied)) if rev <= applied) + } + pub fn apply_session_list(&mut self, machine_pubkey: &str, msg: &SessionListMsg, at: u64) { self.dismissed_sessions = prune_dismissed(&self.dismissed_sessions, at); let dismissed = &self.dismissed_sessions; @@ -703,6 +716,9 @@ impl MachinesState { entry.protocol_version = Some(msg.protocol_version); entry.machine_offline = msg.machine_offline.unwrap_or(false); entry.direct = msg.direct.clone(); + if msg.rev.is_some() { + entry.list_rev = msg.rev; + } entry.last_heartbeat_at = Some(at); entry.sessions = sessions; } diff --git a/crates/client-core/src/stores/pending_sessions.rs b/crates/client-core/src/stores/pending_sessions.rs index 62a20a6..d233c08 100644 --- a/crates/client-core/src/stores/pending_sessions.rs +++ b/crates/client-core/src/stores/pending_sessions.rs @@ -84,12 +84,16 @@ impl PendingSessionsState { /// `session-failed` — flip the placeholder to a visible error. A failure for /// a pending we never saw still surfaces (an invisible failure is the old - /// bug), with empty machine fields. - pub fn apply_failed(&mut self, pending_id: &str, reason: &str, now: u64) { + /// bug), with empty machine fields. Returns whether it is news: false + /// when the placeholder had already failed (the same failure again, over + /// another relay or replayed), so it is told once. + pub fn apply_failed(&mut self, pending_id: &str, reason: &str, now: u64) -> bool { match self.pending.get_mut(pending_id) { Some(existing) => { + let news = existing.state != PendingSessionState::Failed; existing.state = PendingSessionState::Failed; existing.reason = Some(reason.to_string()); + news } None => { self.pending.insert( @@ -104,6 +108,7 @@ impl PendingSessionsState { seen_at: now, }, ); + true } } } diff --git a/crates/client-runtime/src/dispatch.rs b/crates/client-runtime/src/dispatch.rs index c99b582..28cfa63 100644 --- a/crates/client-runtime/src/dispatch.rs +++ b/crates/client-runtime/src/dispatch.rs @@ -202,16 +202,19 @@ impl<'a> Router<'a> { r.pending_sessions_changed = true; } BridgeToPhone::SessionFailed(m) => { - self.stores + let news = self + .stores .pending_sessions .apply_failed(&m.pending_id, &m.reason, self.now); r.pending_sessions_changed = true; - let fx = self.emit_notify(&NotifyEvent::SessionFailed { - machine: machine.to_string(), - session_id: m.pending_id.clone(), - reason: Some(m.reason.clone()).filter(|s| !s.is_empty()), - }); - r.notifies.extend(fx); + if news { + let fx = self.emit_notify(&NotifyEvent::SessionFailed { + machine: machine.to_string(), + session_id: m.pending_id.clone(), + reason: Some(m.reason.clone()).filter(|s| !s.is_empty()), + }); + r.notifies.extend(fx); + } } BridgeToPhone::SessionReady(m) => { if self.stores.pending_sessions.contains(&m.pending_id) { @@ -521,6 +524,13 @@ impl<'a> Router<'a> { if self.stores.machines.machine(machine).is_none() { return; } + // A list no newer than the one applied (late over another relay or + // the direct link, re-published, or replayed after a restart) would + // step every session back and re-announce transitions already told. + if self.stores.machines.is_outdated_list(machine, m) { + log::debug!("session list rev {:?} from {} is outdated; ignored", m.rev, machine.get(..8).unwrap_or(machine)); + return; + } // CDX-026b: capture the pre-merge session states — a backgrounded phone // catching up over sync never sees live cards, so the heartbeat @@ -839,6 +849,7 @@ mod tests { removed_sessions: None, machine_offline: None, direct: None, + rev: None, } } @@ -1165,6 +1176,58 @@ mod tests { assert!(out.persist.contains(&StoreId::Machines)); } + #[tokio::test] + async fn a_session_list_older_than_the_one_applied_changes_nothing_and_tells_nothing() { + let (mut s, ts, kp) = stores().await; + s.machines.register_machine(MACHINE, "laptop", None, None, &[]); + let list = |state, rev| { + let mut m = sessions_msg(vec![info("s1", Some(state), None)]); + m.rev = Some(rev); + BridgeToPhone::Sessions(m) + }; + let route = async |s: &mut CoreStores, msg: BridgeToPhone, now| { + let mut r = Router::new(s, &ts, &kp, now); + r.visible = false; + r.route(MACHINE, &msg).await + }; + assert!(route(&mut s, list(SessionState::Running, 10), 1_000).await.notifies.is_empty()); + let finished = route(&mut s, list(SessionState::Idle, 20), 2_000).await; + assert_eq!(finished.notifies.len(), 1, "the turn's end is told once"); + + // A Running list from before it, late over another relay or the + // direct link, and the Idle one replayed after a restart: nothing. + for (late, now) in [(list(SessionState::Running, 15), 3_000), (list(SessionState::Idle, 20), 4_000)] { + let out = route(&mut s, late, now).await; + assert_eq!(out, RouteResult::default()); + } + assert_eq!(s.machines.session(MACHINE, "s1").unwrap().info.state, Some(SessionState::Idle)); + + // The rev survives the store, so a restarted phone still ignores them. + let stored = serde_json::to_string(&s.machines.machines).unwrap(); + let back: std::collections::BTreeMap = + serde_json::from_str(&stored).unwrap(); + assert_eq!(back[MACHINE].list_rev, Some(20)); + + // A newer turn is told again (past the notification cooldown). + route(&mut s, list(SessionState::Running, 30), 15_000).await; + assert_eq!(route(&mut s, list(SessionState::Idle, 40), 16_000).await.notifies.len(), 1); + } + + #[tokio::test] + async fn the_same_session_failure_is_told_once() { + let (mut s, ts, kp) = stores().await; + let failed = BridgeToPhone::SessionFailed(protocol::events::SessionFailedMsg { + pending_id: "p1".into(), + reason: "no such folder".into(), + }); + let mut r = Router::new(&mut s, &ts, &kp, 1_000); + r.visible = false; + assert_eq!(r.route(MACHINE, &failed).await.notifies.len(), 1); + let mut r = Router::new(&mut s, &ts, &kp, 2_000); + r.visible = false; + assert!(r.route(MACHINE, &failed).await.notifies.is_empty()); + } + #[tokio::test] async fn the_foreground_watched_session_is_never_marked_or_notified() { let (mut s, ts, kp) = stores().await; diff --git a/crates/protocol/fixtures/corpus.json b/crates/protocol/fixtures/corpus.json index bc2afd1..868325e 100644 --- a/crates/protocol/fixtures/corpus.json +++ b/crates/protocol/fixtures/corpus.json @@ -84,6 +84,7 @@ "bridgeToPhone": { "valid": [ { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11 }, + { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11, "rev": 1790000000000 }, { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11, "direct": { "endpoints": ["wss://192.168.1.20:7447", "ws://abc.onion:7448"], "certSha256": "00ff" } }, { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11, "direct": { "endpoints": ["ws://abc.onion:7448"] } }, { "type": "sessions", "machine": "m", "host": "service", "sessions": [{ "id": "s", "agent": "opencode", "slug": "sl", "cwd": "/w", "lastActivity": "t", "lineCount": 1, "title": null, "project": "p", "mode": "build", "effort": "high", "state": "running", "seqHigh": 10 }], "agents": [{ "id": "claude-code", "displayName": "Claude Code", "modes": [{ "id": "default", "label": "Default" }, { "id": "plan", "label": "Plan", "description": "Read-only planning" }], "efforts": [{ "id": "high", "label": "High" }], "defaultMode": "default", "defaultEffort": "high", "supports": { "models": true, "usage": true, "providers": true, "gsd": true, "interrupt": true, "commands": true, "plugins": true }, "credentials": [{ "id": "anthropic_api_key", "label": "Anthropic API key", "present": true, "fromEnv": true }] }, { "id": "opencode", "displayName": "OpenCode" }], "credentials": [{ "id": "github_pat", "label": "GitHub token", "present": false }], "protocolVersion": 11, "capabilities": ["sync/1", "files", "chunked"], "folders": ["a", "b"], "roots": ["/w"], "removedSessions": ["old"], "machineOffline": true }, @@ -157,6 +158,7 @@ { "type": "output", "sessionId": "s", "seq": 1, "entries": [{ "timestamp": "t", "entryType": "tool_call", "callId": "c", "toolName": "X", "kind": "teleport", "title": "t" }] }, { "type": "sessions", "machine": "m", "sessions": [], "protocolVersion": 11 }, { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11, "direct": { "certSha256": "00ff" } }, + { "type": "sessions", "machine": "m", "sessions": [], "agents": [], "protocolVersion": 11, "rev": -1 }, { "type": "sessions", "machine": "m", "sessions": [{ "id": "s", "slug": "sl", "cwd": "/w", "lastActivity": "t", "lineCount": 1, "title": null, "project": "p" }], "agents": [], "protocolVersion": 11 }, { "type": "option-confirmed", "sessionId": "s", "option": "temperature", "value": "1" }, { "type": "input-failed", "sessionId": "s", "reason": "unknown-reason" }, diff --git a/crates/protocol/src/events.rs b/crates/protocol/src/events.rs index 8236d20..30ade0e 100644 --- a/crates/protocol/src/events.rs +++ b/crates/protocol/src/events.rs @@ -44,6 +44,14 @@ pub struct SessionListMsg { /// [`crate::direct`]). Absent when it serves no direct link. #[serde(default, skip_serializing_if = "Option::is_none")] pub direct: Option, + /// The list's revision: strictly greater in each list the bridge + /// publishes (its clock in ms, kept ahead of the last one). The same + /// list reaches a phone over several relays and the direct link, and + /// may be published again, replayed or delayed, so lists arrive out of + /// order; a phone applies only one newer than the newest it applied. + #[serde(default, skip_serializing_if = "Option::is_none")] + #[specta(type = Option)] + pub rev: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize, specta::Type)] diff --git a/docs/PROTOCOL.md b/docs/PROTOCOL.md index 597ee3a..2440851 100644 --- a/docs/PROTOCOL.md +++ b/docs/PROTOCOL.md @@ -84,6 +84,13 @@ advertised; creating a session on it fails with the reason. `session-failed {pendingId, reason}`. The bridge uses the future sessionId as the pendingId; a pending placeholder is also resolved by the session simply appearing in a heartbeat's `sessions[]`. +- **Ordering:** each `sessions` list carries `rev`, strictly greater than the + last one the bridge published (its clock in ms, kept ahead of the stored + last one, so a restart or a clock set back cannot lower it). The same list + arrives over every relay and the direct link, and may be re-published or + replayed, so a phone applies a list only if its `rev` is greater than that + of the newest one it applied for that machine (kept across restarts); an + older one changes nothing and announces nothing. - **Options:** `set-option {sessionId, option: mode|effort|model, value}` → `option-confirmed {sessionId, option, value}`. A refused value publishes nothing (the phone keeps what it knew); a refused effort or model confirms