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
3 changes: 3 additions & 0 deletions crates/agent-protocol/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,9 @@ mod tests {
assert_eq!(s.credentials["anthropic_api_key"].expose(), "sk");
bridge_rt(json!({"v":1,"id":"3","kind":"start-session","payload":{"sessionId":"s","agent":"fake","cwd":"/"}}));
bridge_rt(json!({"v":1,"id":"4","kind":"end-session","payload":{"sessionId":"s"}}));
bridge_rt(json!({"v":1,"id":"4b","kind":"delete-conversation","payload":{
"sessionId":"s","agent":"claude-code","cwd":"/w","conversationId":"native-1"
}}));
bridge_rt(json!({"v":1,"id":"5","kind":"prompt","payload":{"sessionId":"s","text":"hi"}}));
bridge_rt(json!({"v":1,"id":"6","kind":"interrupt","payload":{"sessionId":"s"}}));
bridge_rt(json!({"v":1,"id":"6b","kind":"stop-task","payload":{"sessionId":"s","taskId":"b1"}}));
Expand Down
13 changes: 13 additions & 0 deletions crates/agent-protocol/src/messages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -316,6 +316,19 @@ pub enum BridgeMessage {
StartSession(Box<StartSession>),
/// Stop the session and forget it. Reply: `ack`. No `ended` follows.
EndSession { session_id: String },
/// Delete the agent's own record of a conversation — its transcript
/// files, or its session on the agent's server — once the user deleted
/// the bridge session `session_id` that ran it. Sent only for
/// conversations the bridge itself started, never for one the user
/// began elsewhere. The host first waits for `session_id` to finish
/// ending. Reply: `ack`, also when there was nothing to delete, or
/// `error`.
DeleteConversation {
session_id: String,
agent: String,
cwd: String,
conversation_id: String,
},
/// Hand user input to the agent. Reply: `ack`.
Prompt { session_id: String, text: String },
/// Stop the running turn. Reply: `ack`.
Expand Down
41 changes: 41 additions & 0 deletions crates/bridge-core/src/engine/host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,16 @@ pub(crate) enum HostCall {
/// A session's MCP status, asked for or answering a toggle.
SessionMcp { session_id: String },
CheckCredential { ticket: u64, agent: String, credential: String, value: agent_protocol::Secret },
DeleteConversation(ConversationDelete),
}

/// A conversation of a deleted session, for its agent to delete.
#[derive(Debug, Clone)]
pub(crate) struct ConversationDelete {
pub session_id: String,
pub agent: String,
pub cwd: String,
pub conversation_id: String,
}

fn cancelled(kind: &CardKind, reason: &str) -> BridgeMessage {
Expand Down Expand Up @@ -115,6 +125,10 @@ impl Engine {
self.host.up = false;
self.host.initialized = false;
for call in std::mem::take(&mut self.host.calls).into_values() {
if let HostCall::DeleteConversation(delete) = call {
self.host.deletes.push(delete);
continue;
}
// Session-bound requests die with the sessions, handled below.
if !matches!(
call,
Expand Down Expand Up @@ -175,6 +189,9 @@ impl Engine {
for id in waiting {
self.spawn(&id);
}
for delete in std::mem::take(&mut self.host.deletes) {
self.delete_conversation(delete);
}
}
Ok(_) => log::error!("[Engine] The agent host answered initialize with the wrong reply"),
Err(err) => log::error!("[Engine] The agent host failed to initialize: {err}"),
Expand Down Expand Up @@ -216,9 +233,33 @@ impl Engine {
};
self.on_credential_checked(ticket, &agent, &credential, &value, valid);
}
HostCall::DeleteConversation(delete) => match result {
Ok(_) => log::info!("[Engine] Conversation {} of {} deleted", delete.conversation_id, delete.session_id),
Err(err) => log::warn!(
"[Engine] Deleting conversation {} of {} failed: {err}",
delete.conversation_id,
delete.session_id
),
},
}
}

/// Ask the host to delete a deleted session's conversation, or keep the
/// request until the host is back.
pub(super) fn delete_conversation(&mut self, delete: ConversationDelete) {
if !self.host.initialized {
self.host.deletes.push(delete);
return;
}
let message = BridgeMessage::DeleteConversation {
session_id: delete.session_id.clone(),
agent: delete.agent.clone(),
cwd: delete.cwd.clone(),
conversation_id: delete.conversation_id.clone(),
};
self.call(HostCall::DeleteConversation(delete), message);
}

fn on_option_reply(&mut self, session_id: &str, option: SessionOption, value: String, result: Result<(), String>) {
let Some(session) = self.sessions.get_mut(session_id) else { return };
let confirmed = match result {
Expand Down
6 changes: 5 additions & 1 deletion crates/bridge-core/src/engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ use crate::settings::{load_or_default, ProviderProfile, StoredCredentials, Store
use crate::sync::{SyncConfig, SyncServer, SyncTimer};
use crate::time::iso;

use host::HostCall;
use host::{ConversationDelete, HostCall};

pub const DEFAULT_HEARTBEAT_INTERVAL_MS: u64 = 60_000;
/// How often a session's `committed` flag is checked against git (catches
Expand Down Expand Up @@ -111,6 +111,10 @@ struct HostLink {
initialized: bool,
next_id: u64,
calls: BTreeMap<String, HostCall>,
/// Conversations of deleted sessions the host has yet to delete: asked
/// while it was down, or in flight when it went down. Sent once it is
/// initialized again (deleting twice is harmless).
deletes: Vec<ConversationDelete>,
}

struct PairingWindow {
Expand Down
16 changes: 14 additions & 2 deletions crates/bridge-core/src/engine/phone.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use protocol::events::{
CommandsMsg, InputFailedReason, ModelsMsg, PluginAckMsg, SessionFailedMsg, SessionPendingMsg,
};

use super::{Engine, HostCall};
use super::{ConversationDelete, Engine, HostCall};
use crate::catalog::{is_effort, is_mode};
use crate::io::{Effect, InboundEvent, Via};
use crate::registry::{project_of, SessionRecord};
Expand Down Expand Up @@ -366,7 +366,19 @@ impl Engine {
fn on_close_session(&mut self, session_id: &str) {
let existed = self.sessions.contains_key(session_id);
self.close_runner(session_id, "Session closed");
self.sessions.remove(session_id);
// Every conversation a record points to was started by the bridge (a
// session can only be created fresh), so the agent's copy goes too —
// including one dropped as missing, which may exist after all.
if let Some(Session { rec, .. }) = self.sessions.remove(session_id) {
for conversation_id in [rec.native_session_id, rec.previous_native_session_id].into_iter().flatten() {
self.delete_conversation(ConversationDelete {
session_id: session_id.to_string(),
agent: rec.agent.clone(),
cwd: rec.cwd.clone(),
conversation_id,
});
}
}
self.tombstones.add(session_id);
self.registry_dirty = true;
if let Err(err) = self.transcripts.remove(session_id) {
Expand Down
60 changes: 60 additions & 0 deletions crates/bridge-core/tests/sessions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -456,6 +456,66 @@ fn close_session_ends_it_tombstones_it_and_forgets_its_transcript() {
assert!(rig.transcripts.entries(&s).is_empty());
}

fn conversation_deletes(frames: &[agent_protocol::BridgeFrame]) -> Vec<(String, String, String)> {
frames
.iter()
.filter_map(|f| match &f.message {
BridgeMessage::DeleteConversation { session_id, agent, conversation_id, .. } => {
Some((session_id.clone(), agent.clone(), conversation_id.clone()))
}
_ => None,
})
.collect()
}

#[test]
fn closing_a_session_deletes_its_conversation_after_ending_it() {
let mut rig = Rig::new();
rig.host_up();
let s = rig.ready_session("alpha");
info_native(&mut rig, &s, "n1");
rig.take();
rig.send(json!({"type":"close-session","sessionId":s}));
let frames = rig.host_frames();
let end = frames.iter().position(|f| matches!(f.message, BridgeMessage::EndSession { .. }));
let delete = frames.iter().position(|f| matches!(f.message, BridgeMessage::DeleteConversation { .. }));
assert!(end.is_some() && end < delete, "ended first: {frames:?}");
assert_eq!(conversation_deletes(&frames), vec![(s.clone(), "alpha".into(), "n1".into())]);
}

#[test]
fn closing_a_session_with_no_conversation_deletes_nothing() {
let mut rig = Rig::new();
rig.host_up();
let s = rig.ready_session("alpha");
rig.take();
rig.send(json!({"type":"close-session","sessionId":s}));
assert!(conversation_deletes(&rig.host_frames()).is_empty());
}

#[test]
fn a_conversation_delete_waits_for_the_host_and_survives_its_death() {
let mut rig = Rig::new();
rig.host_up();
let a = rig.ready_session("alpha");
info_native(&mut rig, &a, "na");
let b = rig.ready_session("alpha");
info_native(&mut rig, &b, "nb");
rig.take();

// In flight when the host dies: sent again once it is back.
rig.send(json!({"type":"close-session","sessionId":a}));
assert_eq!(conversation_deletes(&rig.host_frames()).len(), 1);
rig.input(Input::HostDown { reason: "exit 1".into() });
// Asked while it is down: held until then.
rig.send(json!({"type":"close-session","sessionId":b}));
assert!(conversation_deletes(&rig.host_frames()).is_empty());

rig.host_up();
let deleted: Vec<String> = conversation_deletes(&rig.host_frames()).into_iter().map(|(_, _, c)| c).collect();
assert_eq!(deleted, vec!["na".to_string(), "nb".to_string()]);
}

#[test]
fn closing_sessions_in_a_row_ends_each_and_late_host_events_bring_none_back() {
let mut rig = Rig::new();
Expand Down
1 change: 1 addition & 0 deletions docs/PROTOCOL.md
Original file line number Diff line number Diff line change
Expand Up @@ -473,6 +473,7 @@ The bridge's ids are `b1, b2, …`; the host's are `h1, h2, …`.
| `initialize {bridgeVersion}` | `initialized {hostVersion, agents: AgentInfo[]}` |
| `start-session {sessionId, agent, cwd, mode?, effort?, model?, resume?, credentials, env, provider?}` | `ack` once starting (progress follows as events), or `error` |
| `end-session {sessionId}` | `ack`; no `ended` follows |
| `delete-conversation {sessionId, agent, cwd, conversationId}` | `ack` once the agent's own record of a deleted session's conversation is gone (also when there was none), after `sessionId` finished ending; or `error` |
| `prompt {sessionId, text}` | `ack` |
| `interrupt {sessionId}` | `ack` |
| `stop-task {sessionId, taskId}` | `ack` once asked (the task's next `background_task` entry says it stopped), or `error` |
Expand Down
50 changes: 50 additions & 0 deletions packages/agent-host/src/__tests__/host.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,10 @@
* bridge, and session lifetimes.
*/
import { describe, it, expect } from 'vitest';
import type { Driver } from '../driver';
import { FakeDriver } from '../drivers/fake';
import { AgentHost, parseBridgeFrame } from '../host';
import type { AgentInfo } from '../types';

interface Frame {
v: number;
Expand Down Expand Up @@ -126,6 +128,54 @@ describe('AgentHost', () => {
expect(h.reply('z')).toMatchObject({ kind: 'ack' });
});

it('deleting a conversation waits for its session to finish ending', async () => {
const log: string[] = [];
let release!: () => void;
const stopped = new Promise<void>((r) => (release = r));
const driver: Driver = {
info: () => ({ id: 'slow', displayName: 'Slow', modes: [], efforts: [], supports: {} as AgentInfo['supports'], credentials: [] }),
startSession: () => ({
prompt: () => {},
interrupt: async () => {},
setOption: async () => {},
getUsage: async () => null,
end: async () => {
await stopped;
log.push('ended');
},
}),
listModels: async () => ({ models: [] }),
deleteConversation: async (id, cwd) => {
log.push(`deleted ${id} in ${cwd}`);
},
};
const out: Frame[] = [];
const host = new AgentHost([driver], { write: (l) => out.push(JSON.parse(l) as Frame), log: () => {} }, '9.9.9');
const send = (id: string, kind: string, payload: Record<string, unknown>) =>
host.handleLine(JSON.stringify({ v: 1, id, kind, payload }));
await send('st', 'start-session', { sessionId: 's1', agent: 'slow', cwd: '/w' });

// A delete while the session runs is refused: the agent still writes.
await send('d0', 'delete-conversation', { sessionId: 's1', agent: 'slow', cwd: '/w', conversationId: 'n1' });
expect(out.find((f) => f.id === 'd0')).toMatchObject({ kind: 'error', payload: { message: 'session s1 is still running' } });

// Handled concurrently, as main.ts does: the delete waits for the end.
const ending = send('e', 'end-session', { sessionId: 's1' });
const deleting = send('d1', 'delete-conversation', { sessionId: 's1', agent: 'slow', cwd: '/w', conversationId: 'n1' });
await new Promise((r) => setTimeout(r, 10));
expect(log).toEqual([]);
release();
await Promise.all([ending, deleting]);
expect(log).toEqual(['ended', 'deleted n1 in /w']);
expect(out.find((f) => f.id === 'd1')).toEqual({ v: 1, id: 'd1', kind: 'ack' });
});

it('deleting a conversation of an agent that keeps none is acknowledged', async () => {
const h = harness();
await h.send('delete-conversation', { sessionId: 's1', agent: 'fake', cwd: '/w', conversationId: 'n1' }, 'd');
expect(h.reply('d')).toMatchObject({ kind: 'ack' });
});

it('a lost resume target ends the session at once, saying so', async () => {
const h = harness();
await h.send('start-session', { sessionId: 's1', agent: 'fake', cwd: '/w', resume: 'lost' });
Expand Down
10 changes: 10 additions & 0 deletions packages/agent-host/src/driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,16 @@ export interface Driver {
/** Check a credential value with its provider: true/false, or undefined
* when it could not be checked. */
checkCredential?(credential: string, value: string): Promise<boolean | undefined>;
/**
* Delete the agent's own record of a conversation (an `info`
* `nativeSessionId` it reported, run in `cwd`): its transcript files, or
* its session on the agent's server, with any subagent conversations it
* spawned. Called only once the session that ran it has ended. Resolves
* when nothing of it is left — also when there was nothing to begin with —
* and rejects with the reason it could not. Absent: the agent keeps
* nothing the host can remove.
*/
deleteConversation?(conversationId: string, cwd: string): Promise<void>;
/** Manages the agent's plugins, when it has any. */
readonly plugins?: PluginManager;
/** Manages the agent's MCP servers, when it supports them. */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,11 @@ class ScriptedFacade implements SdkFacade {
return [{ id: 'claude-sonnet-5', label: 'Sonnet 5' }];
}

readonly deleted: Array<[string, string]> = [];
async deleteSession(sessionId: string, cwd: string): Promise<void> {
this.deleted.push([sessionId, cwd]);
}

get last(): { opts: SdkSessionOptions; handle: ScriptedHandle } {
return this.sessions[this.sessions.length - 1]!;
}
Expand Down Expand Up @@ -179,6 +184,27 @@ describe('Claude session lifecycle', () => {
await ctx.waitFor((e) => e.type === 'info' && e.mode === 'default');
});

it('ending waits for the CLI to close its stream, so nothing writes after a delete', async () => {
const { ctx, session, handle } = start();
await ctx.waitFor((e) => e.type === 'ready');
let closed = false;
handle.end = async () => {
handle.ended = true;
setTimeout(() => {
closed = true;
handle.close();
}, 20);
};
await session.end();
expect(closed).toBe(true);
});

it('deleting a conversation deletes its transcript in its directory', async () => {
const facade = new ScriptedFacade();
await new ClaudeDriver({ facade }).deleteConversation('native-1', '/w');
expect(facade.deleted).toEqual([['native-1', '/w']]);
});

it('a clean stream end ends the session normally; before confirmation it is a failure', async () => {
const ok = start();
await ok.ctx.waitFor((e) => e.type === 'ready');
Expand Down
Loading
Loading