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
65 changes: 65 additions & 0 deletions crates/bridge-core/tests/sessions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -456,6 +456,71 @@ fn close_session_ends_it_tombstones_it_and_forgets_its_transcript() {
assert!(rig.transcripts.entries(&s).is_empty());
}

#[test]
fn closing_sessions_in_a_row_ends_each_and_late_host_events_bring_none_back() {
let mut rig = Rig::new();
rig.host_up();
let ready: Vec<String> = ["alpha", "beta", "alpha"].iter().map(|a| rig.ready_session(a)).collect();
for s in &ready {
rig.say(s, "x");
}
// One more that the host has not answered yet: closed while starting.
rig.send(json!({"type":"create-session","agent":"alpha"}));
let (start_id, starting) = rig.start_request();
rig.take();

let mut all = ready.clone();
all.push(starting.session_id.clone());
for s in &all {
rig.send(json!({"type":"close-session","sessionId":s}));
}

let ended: Vec<String> = rig
.host_frames()
.into_iter()
.filter_map(|f| match f.message {
BridgeMessage::EndSession { session_id } => Some(session_id),
_ => None,
})
.collect();
assert_eq!(ended, all, "every session is ended in the host, the starting one included");
rig.advance(5_000);
let msgs = rig.messages();
let acked: Vec<&str> = msgs
.iter()
.filter_map(|m| match m {
BridgeToPhone::CloseSessionAck(a) if a.success => Some(a.session_id.as_str()),
_ => None,
})
.collect();
assert_eq!(acked, all.iter().map(String::as_str).collect::<Vec<_>>());
let hb = last_heartbeat(&msgs);
assert!(hb.sessions.is_empty());
let mut removed = hb.removed_sessions.clone().unwrap_or_default();
removed.sort();
let mut expected = all.clone();
expected.sort();
assert_eq!(removed, expected);

// What the host still had in flight arrives after the closes.
rig.host_reply(&start_id, HostMessage::Ack);
rig.host_event(&starting.session_id, SessionEvent::Ready {});
rig.say(&ready[0], "late");
rig.host_event(&ready[1], SessionEvent::Ended { error: None, resume_lost: false });
rig.advance(5_000);
let msgs = rig.messages();
assert!(outputs(&msgs).is_empty(), "no output of a closed session reaches the phone");
if let Some(hb) = msgs.iter().rev().find_map(|m| match m {
BridgeToPhone::Sessions(h) => Some(h),
_ => None,
}) {
assert!(hb.sessions.is_empty(), "a closed session came back: {:?}", hb.sessions);
}
for s in &all {
assert!(rig.transcripts.entries(s).is_empty());
}
}

#[test]
fn shutdown_publishes_every_session_offline_and_stops() {
let mut rig = Rig::new();
Expand Down
8 changes: 8 additions & 0 deletions crates/client-runtime/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -948,6 +948,14 @@ impl Loop {
}
Msg::Pause => {
// Backgrounded: the OS may kill the process from here on.
// A delete still in its undo window is committed now —
// its toast is out of sight, and the window lives only
// in memory: a process killed before it closed would
// never send the `close-session`, and the session would
// come back once its dismissal expired.
if self.stores.delete_controller.is_pending() {
self.on_undo_timer().await;
}
self.flush_acks();
self.flush_writes().await;
self.ws.set_ping_interval(BACKGROUND_PING_EVERY);
Expand Down
102 changes: 102 additions & 0 deletions crates/client-runtime/src/runtime/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1758,6 +1758,108 @@ async fn the_undo_toast_clears_itself_when_the_window_expires_without_a_tap() {
.await;
}

/// A heartbeat listing `ids` for the machine of the session-delete test.
fn heartbeat_listing(ids: &[&str]) -> BridgeToPhone {
let sessions: Vec<serde_json::Value> = ids
.iter()
.map(|id| {
serde_json::json!({"id": id, "agent": "claude-code", "slug": id, "cwd": "/w",
"lastActivity": "t", "lineCount": 0, "title": null, "project": "p"})
})
.collect();
protocol::codec::decode_bridge_to_phone(
&serde_json::json!({
"type": "sessions",
"machine": "laptop",
"sessions": sessions,
"agents": [],
"protocolVersion": protocol::capabilities::PROTOCOL_VERSION,
})
.to_string(),
)
.unwrap()
}

/// Deleting sessions one after another: each delete commits the one before
/// it at once (one `close-session` each), a heartbeat sent before the
/// bridge has closed the later ones does not bring them back, and the
/// last one — still in its undo window — is committed as soon as the app
/// is backgrounded rather than lost with a killed process.
#[tokio::test]
async fn deleting_sessions_in_a_row_closes_each_one_exactly_once() {
LocalSet::new()
.run_until(async {
let mut mock = mock_relay().await;
let phone = keypair_from_secret_hex(SEC_PHONE).unwrap();
let machine = generate_keypair();
let mut state = client_core::stores::machines::MachinesState::default();
state.register_machine(&machine.pubkey_hex, "bridge", None, None, &[]);
let kv = MemoryKv::seeded([(
crate::stores::MACHINES_KEY,
client_core::stores::machines::serialize_machines(&state.machines),
)]);
let ports = CorePorts { kv: Rc::new(kv), ..CorePorts::default() };
let core = core_for_ports(&mock, &phone, Rc::new(Spy::default()), ports).await;
core.start();
eose_all(&mut mock).await;
let refresh = next_command_via(&mut mock, &phone, &phone.pubkey_hex, &machine).await;
assert!(matches!(refresh, Some(PhoneToBridge::RefreshSessions(_))), "{refresh:?}");

let listed = heartbeat_listing(&["s1", "s2", "s3"]);
push_bridge_to_phone_event(&mock, &machine, &phone.pubkey_hex, "cd-1", &listed);
settle().await;
let sessions = |view: MachinesView| -> Vec<String> {
view.machines[&machine.pubkey_hex].sessions.keys().cloned().collect()
};
assert_eq!(sessions(core.machines_view().await), ["s1", "s2", "s3"]);

for id in ["s1", "s2", "s3"] {
core.dispatch(Intent::DeleteSession {
machine: machine.pubkey_hex.clone(),
session_id: id.into(),
label: None,
})
.await;
}
assert!(sessions(core.machines_view().await).is_empty());
let toast = core.ui_view().await.undo_toast.expect("the last delete is undoable");
assert_eq!(toast.session_id, "s3");

// s1 and s2 were committed by the delete that followed each.
let closed = |cmd: Option<PhoneToBridge>| match cmd {
Some(PhoneToBridge::CloseSession(m)) => m.session_id,
other => panic!("expected a close-session, got {other:?}"),
};
for id in ["s1", "s2"] {
assert_eq!(closed(next_command_via(&mut mock, &phone, &phone.pubkey_hex, &machine).await), id);
}

// The bridge has closed s1 only; its next list still has the
// other two, which must stay deleted on the phone.
let after_first = heartbeat_listing(&["s2", "s3"]);
push_bridge_to_phone_event(&mock, &machine, &phone.pubkey_hex, "cd-1", &after_first);
settle().await;
assert!(sessions(core.machines_view().await).is_empty());

// Backgrounded inside s3's undo window: committed now, well
// before the window would have run out.
core.pause();
let last = tokio::time::timeout(
std::time::Duration::from_millis(client_core::delete_controller::UNDO_DELAY_MS / 2),
next_command_via(&mut mock, &phone, &phone.pubkey_hex, &machine),
)
.await
.expect("the pending delete is sent when the app is backgrounded");
assert_eq!(closed(last), "s3");
assert!(core.ui_view().await.undo_toast.is_none());

// Nothing is left to undo, and nothing is sent twice.
core.dispatch(Intent::UndoDelete).await;
assert!(sessions(core.machines_view().await).is_empty());
})
.await;
}

/// The next command EVENT the identity `phone` signed for `machine`,
/// with its payload decrypted as from `payload_key` (the phone's
/// identity or its session key). Every EVENT is ACKed on the way.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -479,6 +479,24 @@ describe('DeepSeekSession lifecycle', () => {
await vi.waitFor(() => expect(ready.harness.child.killed).toContain('SIGTERM'));
});

it('closes sessions ended together one by one, and the shared process only after the last', async () => {
const ready = withDriver();
const sessions: DriverSession[] = [];
for (const sessionId of ['b1', 'b2', 'b3']) {
await started(ready, { sessionId });
sessions.push(ready.session);
}
expect(ready.spawns).toHaveLength(1);

await Promise.all([sessions[0]!.end(), sessions[1]!.end()]);
expect([...ready.harness.closed].sort()).toEqual(['s1', 's2']);
expect(ready.harness.child.killed).toEqual([]);

await Promise.all([sessions[2]!.end(), sessions[2]!.end()]);
expect([...ready.harness.closed].sort()).toEqual(['s1', 's2', 's3']);
await vi.waitFor(() => expect(ready.harness.child.killed).toContain('SIGTERM'));
});

it('ends every session of a harness process that goes away', async () => {
const ready = withDriver();
const first = await started(ready, { sessionId: 'b1' });
Expand Down
Loading