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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions crates/bridge-core/src/engine/host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
17 changes: 17 additions & 0 deletions crates/bridge-core/src/engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -193,6 +195,7 @@ impl Engine {
started: false,
stopped: false,
list_dirty: false,
list_rev: 0,
registry_dirty: false,
activity_dirty: false,
}
Expand Down Expand Up @@ -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<String> = 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);
Expand Down Expand Up @@ -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,
Expand All @@ -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
}
}
3 changes: 3 additions & 0 deletions crates/bridge-core/src/io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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];
}
18 changes: 18 additions & 0 deletions crates/bridge-core/tests/sessions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() });
Expand Down
1 change: 1 addition & 0 deletions crates/bridge-runtime/src/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,7 @@ mod tests {
removed_sessions: None,
machine_offline: None,
direct: None,
rev: None,
})
}

Expand Down
16 changes: 16 additions & 0 deletions crates/client-core/src/stores/machines.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SessionGrant>,
/// 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<specta_typescript::Number>)]
pub list_rev: Option<u64>,
/// Where the bridge says it can be reached directly (its heartbeat's
/// `direct`), as last heard.
#[serde(default, skip_serializing_if = "Option::is_none")]
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down
9 changes: 7 additions & 2 deletions crates/client-core/src/stores/pending_sessions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -104,6 +108,7 @@ impl PendingSessionsState {
seen_at: now,
},
);
true
}
}
}
Expand Down
77 changes: 70 additions & 7 deletions crates/client-runtime/src/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -839,6 +849,7 @@ mod tests {
removed_sessions: None,
machine_offline: None,
direct: None,
rev: None,
}
}

Expand Down Expand Up @@ -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<String, client_core::stores::machines::MachineView> =
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;
Expand Down
2 changes: 2 additions & 0 deletions crates/protocol/fixtures/corpus.json
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand Down Expand Up @@ -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" },
Expand Down
8 changes: 8 additions & 0 deletions crates/protocol/src/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<crate::direct::DirectInfo>,
/// 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<specta_typescript::Number>)]
pub rev: Option<u64>,
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, specta::Type)]
Expand Down
7 changes: 7 additions & 0 deletions docs/PROTOCOL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading