From 05fbd103e47def36c42a0160ad66eef847e9fc3f Mon Sep 17 00:00:00 2001 From: deymosh <232244450+deymosh@users.noreply.github.com> Date: Thu, 8 Oct 2026 14:18:03 +0200 Subject: [PATCH] feat: deleting a session also deletes the agent's conversation Deleting a session from the phone removed the bridge's record and transcript but left the agent's own copy behind: Claude Code's transcript files, OpenCode's server-side session, the DeepSeek Harness's session directory. They piled up unseen and could still be resumed by hand. A new driver-protocol request, delete-conversation, asks the agent host to delete one conversation once its session has ended. The bridge sends it on close-session for each conversation the record points to (the current one, and one dropped as missing), all of them started by the bridge itself. A request the host could not take (down, or it died with the request in flight) is kept and sent once it is initialized again; deleting twice is harmless. The host handles lines concurrently, so the delete waits for that session's end-session to finish and is refused while it still runs. Each driver deletes what its agent keeps: - Claude Code: the SDK's deleteSession (transcript plus subagent transcripts). Ending a session now also waits, bounded, for the CLI to exit, so it cannot write the transcript back after the delete. - OpenCode: session.delete, which removes child sessions too; a 404 is already deleted. - DeepSeek Harness: no API, so the session directory is removed, with every subagent session whose header names it as parent. Verified live on the three agents: after the delete, resuming the conversation finds nothing, and its files are gone. Co-Authored-By: Claude Code --- crates/agent-protocol/src/lib.rs | 3 + crates/agent-protocol/src/messages.rs | 13 +++ crates/bridge-core/src/engine/host.rs | 41 +++++++ crates/bridge-core/src/engine/mod.rs | 6 +- crates/bridge-core/src/engine/phone.rs | 16 ++- crates/bridge-core/tests/sessions.rs | 60 ++++++++++ docs/PROTOCOL.md | 1 + .../agent-host/src/__tests__/host.test.ts | 50 +++++++++ packages/agent-host/src/driver.ts | 10 ++ .../claude/__tests__/claudeDriver.test.ts | 26 +++++ .../agent-host/src/drivers/claude/driver.ts | 22 +++- .../agent-host/src/drivers/claude/facade.ts | 17 ++- .../src/drivers/claude/testModeFacade.ts | 1 + .../__tests__/deepseekConversations.test.ts | 66 +++++++++++ .../src/drivers/deepseek/conversations.ts | 103 ++++++++++++++++++ .../agent-host/src/drivers/deepseek/driver.ts | 5 + .../opencode/__tests__/opencodeDriver.test.ts | 21 ++++ .../agent-host/src/drivers/opencode/driver.ts | 9 ++ packages/agent-host/src/generated/protocol.ts | 30 +++++ packages/agent-host/src/host.ts | 24 +++- 20 files changed, 516 insertions(+), 8 deletions(-) create mode 100644 packages/agent-host/src/drivers/deepseek/__tests__/deepseekConversations.test.ts create mode 100644 packages/agent-host/src/drivers/deepseek/conversations.ts diff --git a/crates/agent-protocol/src/lib.rs b/crates/agent-protocol/src/lib.rs index 28040e52..b76b6297 100644 --- a/crates/agent-protocol/src/lib.rs +++ b/crates/agent-protocol/src/lib.rs @@ -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"}})); diff --git a/crates/agent-protocol/src/messages.rs b/crates/agent-protocol/src/messages.rs index ae95de46..b28403d1 100644 --- a/crates/agent-protocol/src/messages.rs +++ b/crates/agent-protocol/src/messages.rs @@ -316,6 +316,19 @@ pub enum BridgeMessage { StartSession(Box), /// 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`. diff --git a/crates/bridge-core/src/engine/host.rs b/crates/bridge-core/src/engine/host.rs index 2d08262b..31923e39 100644 --- a/crates/bridge-core/src/engine/host.rs +++ b/crates/bridge-core/src/engine/host.rs @@ -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 { @@ -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, @@ -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}"), @@ -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 { diff --git a/crates/bridge-core/src/engine/mod.rs b/crates/bridge-core/src/engine/mod.rs index dcf64744..f4774d10 100644 --- a/crates/bridge-core/src/engine/mod.rs +++ b/crates/bridge-core/src/engine/mod.rs @@ -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 @@ -111,6 +111,10 @@ struct HostLink { initialized: bool, next_id: u64, calls: BTreeMap, + /// 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, } struct PairingWindow { diff --git a/crates/bridge-core/src/engine/phone.rs b/crates/bridge-core/src/engine/phone.rs index 7e2b45d4..7f3c86e9 100644 --- a/crates/bridge-core/src/engine/phone.rs +++ b/crates/bridge-core/src/engine/phone.rs @@ -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}; @@ -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) { diff --git a/crates/bridge-core/tests/sessions.rs b/crates/bridge-core/tests/sessions.rs index 1b90f401..b0c13a89 100644 --- a/crates/bridge-core/tests/sessions.rs +++ b/crates/bridge-core/tests/sessions.rs @@ -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 = 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(); diff --git a/docs/PROTOCOL.md b/docs/PROTOCOL.md index 5f7b1fbe..46738d2f 100644 --- a/docs/PROTOCOL.md +++ b/docs/PROTOCOL.md @@ -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` | diff --git a/packages/agent-host/src/__tests__/host.test.ts b/packages/agent-host/src/__tests__/host.test.ts index 59abb04e..514e7f42 100644 --- a/packages/agent-host/src/__tests__/host.test.ts +++ b/packages/agent-host/src/__tests__/host.test.ts @@ -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; @@ -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((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) => + 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' }); diff --git a/packages/agent-host/src/driver.ts b/packages/agent-host/src/driver.ts index b50012be..3bc67d1d 100644 --- a/packages/agent-host/src/driver.ts +++ b/packages/agent-host/src/driver.ts @@ -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; + /** + * 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; /** Manages the agent's plugins, when it has any. */ readonly plugins?: PluginManager; /** Manages the agent's MCP servers, when it supports them. */ diff --git a/packages/agent-host/src/drivers/claude/__tests__/claudeDriver.test.ts b/packages/agent-host/src/drivers/claude/__tests__/claudeDriver.test.ts index 6488c6ef..736132e8 100644 --- a/packages/agent-host/src/drivers/claude/__tests__/claudeDriver.test.ts +++ b/packages/agent-host/src/drivers/claude/__tests__/claudeDriver.test.ts @@ -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 { + this.deleted.push([sessionId, cwd]); + } + get last(): { opts: SdkSessionOptions; handle: ScriptedHandle } { return this.sessions[this.sessions.length - 1]!; } @@ -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'); diff --git a/packages/agent-host/src/drivers/claude/driver.ts b/packages/agent-host/src/drivers/claude/driver.ts index 8e340d25..956c26a3 100644 --- a/packages/agent-host/src/drivers/claude/driver.ts +++ b/packages/agent-host/src/drivers/claude/driver.ts @@ -173,6 +173,9 @@ export class ClaudeSession implements DriverSession { private queuedInput: string[] = []; private ready = false; private ended = false; + /** The loop reading the SDK's messages; it finishes once the CLI has + * closed its stream, that is, has exited. */ + private consuming: Promise | null = null; /** Load the plugins as they are on disk now. */ async reloadPlugins(): Promise { if (!this.ended) await this.handle?.reloadPlugins(); @@ -273,7 +276,7 @@ export class ClaudeSession implements DriverSession { }, ); } - this.consume(handle).catch((err) => this.finish(`SDK stream consumer failed: ${err}`)); + this.consuming = this.consume(handle).catch((err) => this.finish(`SDK stream consumer failed: ${err}`)); } private markReady(): void { @@ -695,9 +698,22 @@ export class ClaudeSession implements DriverSession { async end(): Promise { this.ended = true; await this.handle?.end(); + // Aborting does not wait for the CLI to exit, and until it has it may + // still write to the conversation's transcript — which would bring a + // deleted conversation back. Bounded: a wedged CLI must not hold up the + // reply. + if (this.consuming) { + let timer: NodeJS.Timeout | undefined; + const timeout = new Promise((resolve) => (timer = setTimeout(resolve, CLI_EXIT_WAIT_MS))); + await Promise.race([this.consuming, timeout]); + clearTimeout(timer); + } } } +/** How long ending a session waits for its CLI to exit. */ +const CLI_EXIT_WAIT_MS = 5_000; + export class ClaudeDriver implements Driver { /** The on-demand install in progress or done; dropped when it fails, so * the next session tries again (a bridge started offline recovers). */ @@ -829,6 +845,10 @@ export class ClaudeDriver implements Driver { return undefined; } } + + deleteConversation(conversationId: string, cwd: string): Promise { + return this.options.facade.deleteSession(conversationId, cwd); + } } interface HookAsk { diff --git a/packages/agent-host/src/drivers/claude/facade.ts b/packages/agent-host/src/drivers/claude/facade.ts index ebbb4026..ce6c155e 100644 --- a/packages/agent-host/src/drivers/claude/facade.ts +++ b/packages/agent-host/src/drivers/claude/facade.ts @@ -13,7 +13,7 @@ * adapter, the driver) import them from HERE, type-only, never from the SDK * package directly. */ -import { query } from '@anthropic-ai/claude-agent-sdk'; +import { deleteSession, query } from '@anthropic-ai/claude-agent-sdk'; import { randomUUID } from 'node:crypto'; import * as fs from 'node:fs'; import * as os from 'node:os'; @@ -227,6 +227,9 @@ export interface SdkFacade { * With `discovery` set and no live session to ask, a throwaway session is * spawned to answer (never given a prompt) and closed again. */ supportedModels(discovery?: ModelDiscoveryOptions): Promise; + /** Delete a conversation's transcript and its subagent transcripts. + * Resolves also when it has none (it never had a turn). */ + deleteSession(sessionId: string, cwd: string): Promise; } // --- Model-list aggregation (CDX-022) --- @@ -1065,4 +1068,16 @@ export class RealSdkFacade implements SdkFacade { // configured for an LLM gateway.) return firstSupportedModels([...this.handles].filter((h) => !h.customProvider)); } + + async deleteSession(sessionId: string, cwd: string): Promise { + try { + await deleteSession(sessionId, { dir: cwd }); + } catch (err) { + // The SDK also throws when it finds no transcript, which a conversation + // that never had a turn does not have: only a transcript still there + // is a failure. + const left = claudeProjectDirs(cwd).some((dir) => fs.existsSync(path.join(dir, `${sessionId}.jsonl`))); + if (left) throw err; + } + } } diff --git a/packages/agent-host/src/drivers/claude/testModeFacade.ts b/packages/agent-host/src/drivers/claude/testModeFacade.ts index 199fbae7..e76e23c5 100644 --- a/packages/agent-host/src/drivers/claude/testModeFacade.ts +++ b/packages/agent-host/src/drivers/claude/testModeFacade.ts @@ -210,4 +210,5 @@ export class TestModeSdkFacade implements SdkFacade { async supportedModels(): Promise { return [{ id: 'test-mode', label: 'Test Mode' }]; } + async deleteSession(): Promise {} } diff --git a/packages/agent-host/src/drivers/deepseek/__tests__/deepseekConversations.test.ts b/packages/agent-host/src/drivers/deepseek/__tests__/deepseekConversations.test.ts new file mode 100644 index 00000000..fdc14976 --- /dev/null +++ b/packages/agent-host/src/drivers/deepseek/__tests__/deepseekConversations.test.ts @@ -0,0 +1,66 @@ +/** + * Deleting a harness conversation from its on-disk layout: the conversation, + * its subagents (found through their headers, compressed or not), and + * nothing else. + */ +import { existsSync } from 'node:fs'; +import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import * as path from 'node:path'; +import { zstdCompressSync } from 'node:zlib'; +import { afterEach, describe, expect, it } from 'vitest'; +import { deleteDshConversation } from '../conversations'; + +let home: string; + +afterEach(async () => { + if (home) await rm(home, { recursive: true, force: true }); +}); + +/** A session directory whose log starts with `header`, then one event. */ +async function session(project: string, id: string, header: Record = {}, zstd = true): Promise { + const dir = path.join(home, 'sessions', project, id); + await mkdir(dir, { recursive: true }); + const head = `${JSON.stringify({ type: 'session', version: 4, id, ...header })}\n`; + const body = `${JSON.stringify({ type: 'event' })}\n`; + const content = zstd + ? Buffer.concat([zstdCompressSync(Buffer.from(head)), zstdCompressSync(Buffer.from(body))]) + : Buffer.from(head + body); + await writeFile(path.join(dir, zstd ? 'session.v4.jsonl.zstd' : 'session.v4.jsonl'), content); + await writeFile(path.join(dir, 'session.lock'), ''); + return dir; +} + +const subagentOf = (parent: string) => ({ origin: 'subagent', parentSession: parent }); + +describe('deleting a DeepSeek Harness conversation', () => { + it('removes it with its subagents, theirs, and nothing else', async () => { + home = await mkdtemp(path.join(tmpdir(), 'dsh-home-')); + const target = await session('--w--', 'a'); + const child = await session('--w--', 'a-child', subagentOf('a')); + const grandchild = await session('--w--', 'a-grandchild', subagentOf('a-child'), false); + const sibling = await session('--w--', 'b'); + const othersChild = await session('--w--', 'b-child', subagentOf('b')); + const elsewhere = await session('--x--', 'c'); + + await deleteDshConversation(home, 'a'); + + for (const gone of [target, child, grandchild]) expect(existsSync(gone)).toBe(false); + for (const kept of [sibling, othersChild, elsewhere]) expect(existsSync(kept)).toBe(true); + }); + + it('finds it in whichever project it was run in', async () => { + home = await mkdtemp(path.join(tmpdir(), 'dsh-home-')); + const target = await session('--some-project~0020dir--', 'a'); + await deleteDshConversation(home, 'a'); + expect(existsSync(target)).toBe(false); + }); + + it('nothing to delete is not an error; an id that is not a path segment is', async () => { + home = await mkdtemp(path.join(tmpdir(), 'dsh-home-')); + await expect(deleteDshConversation(home, 'missing')).resolves.toBeUndefined(); + await session('--w--', 'a'); + await expect(deleteDshConversation(home, '../--w--')).rejects.toThrow(/not a DeepSeek Harness conversation id/); + expect(existsSync(path.join(home, 'sessions', '--w--', 'a'))).toBe(true); + }); +}); diff --git a/packages/agent-host/src/drivers/deepseek/conversations.ts b/packages/agent-host/src/drivers/deepseek/conversations.ts new file mode 100644 index 00000000..00570c12 --- /dev/null +++ b/packages/agent-host/src/drivers/deepseek/conversations.ts @@ -0,0 +1,103 @@ +/** + * Deleting a harness conversation. The harness has no API for it, so this + * follows its on-disk layout (dsh-session-persistence-jsonl): + * + * $DSH_HOME/sessions///session.v.jsonl[.zstd] + * + * where `` is a lossy, truncated encoding of the session's cwd — + * not recomputed here: a conversation id is a UUID, so it is looked up in + * every project directory instead. Each subagent conversation is a sibling + * directory of its own, whose header (the first line of its log) names its + * parent as `parentSession`, with `origin: "subagent"`; those go too, and + * their own subagents with them. + */ +import { createReadStream } from 'node:fs'; +import { readdir, rm } from 'node:fs/promises'; +import * as path from 'node:path'; +import { StringDecoder } from 'node:string_decoder'; +import { createZstdDecompress } from 'node:zlib'; + +const LOG_FILE = /^session\.v(\d+)\.jsonl(\.zstd)?$/; +/** A header is a short JSON object; anything longer is not one. */ +const HEADER_MAX = 64 * 1024; + +/** Delete conversation `id` and its subagent conversations from `home`. + * Resolves also when there is none. */ +export async function deleteDshConversation(home: string, id: string): Promise { + if (!isSegment(id)) throw new Error(`'${id}' is not a DeepSeek Harness conversation id`); + const root = path.join(home, 'sessions'); + for (const project of await subdirs(root)) { + const dir = path.join(root, project); + const names = await subdirs(dir); + if (!names.includes(id)) continue; + const doomed = [id]; + const parentOf = new Map(); + for (const name of names) { + const header = await readHeader(path.join(dir, name)); + if (header?.origin === 'subagent' && typeof header.parentSession === 'string') parentOf.set(name, header.parentSession); + } + for (let i = 0; i < doomed.length; i++) { + for (const [child, parent] of parentOf) if (parent === doomed[i] && !doomed.includes(child)) doomed.push(child); + } + for (const name of doomed.reverse()) await rm(path.join(dir, name), { recursive: true, force: true }); + } +} + +function isSegment(id: string): boolean { + return id !== '' && id !== '.' && id !== '..' && !/[/\\]/.test(id); +} + +async function subdirs(dir: string): Promise { + try { + return (await readdir(dir, { withFileTypes: true })).filter((e) => e.isDirectory()).map((e) => e.name); + } catch { + return []; + } +} + +/** The header of the session in `dir`, from its newest log; undefined when + * it has none that reads as one. */ +async function readHeader(dir: string): Promise<{ origin?: unknown; parentSession?: unknown } | undefined> { + let files: string[]; + try { + files = await readdir(dir); + } catch { + return undefined; + } + const version = (name: string) => Number(LOG_FILE.exec(name)?.[1] ?? -1); + const log = files.filter((f) => LOG_FILE.test(f)).sort((a, b) => version(b) - version(a))[0]; + if (!log) return undefined; + const line = await firstLine(path.join(dir, log)); + if (line === undefined) return undefined; + try { + const value: unknown = JSON.parse(line); + return typeof value === 'object' && value !== null ? value : undefined; + } catch { + return undefined; + } +} + +/** The first line of a log, decompressing a `.zstd` one only as far as it. + * A log still being written may end in an incomplete frame; the header is + * in the first one. */ +async function firstLine(file: string): Promise { + const source = createReadStream(file); + const stream = file.endsWith('.zstd') ? source.pipe(createZstdDecompress()) : source; + source.on('error', (err) => stream.destroy(err)); + const decoder = new StringDecoder('utf8'); + let text = ''; + try { + for await (const chunk of stream as AsyncIterable) { + text += decoder.write(chunk); + const end = text.indexOf('\n'); + if (end >= 0) return text.slice(0, end); + if (text.length > HEADER_MAX) return undefined; + } + return text === '' ? undefined : text; + } catch { + return undefined; + } finally { + source.destroy(); + stream.destroy(); + } +} diff --git a/packages/agent-host/src/drivers/deepseek/driver.ts b/packages/agent-host/src/drivers/deepseek/driver.ts index b0be197b..53d76219 100644 --- a/packages/agent-host/src/drivers/deepseek/driver.ts +++ b/packages/agent-host/src/drivers/deepseek/driver.ts @@ -45,6 +45,7 @@ import { askPlugin, listSessionCommands, runSessionCommand, steerSession } from import { HARNESS_PLUGIN, installHarnessPlugin, QUESTION_MARKER } from './plugin'; import { ASK_USER_TOOL, installProfileTools } from './profileTools'; import { parseQuestionLine, planReviewOf, toAnswerItems, toQuestionSpecs, type PlanReview, type PushedQuestionLine } from './questions'; +import { deleteDshConversation } from './conversations'; import { gatewayModelsUrl, syncGatewayCatalog } from './gateway'; import { DSH_LABEL } from './install'; import { DeepSeekMcp } from './mcp'; @@ -1006,6 +1007,10 @@ export class DeepSeekDriver implements Driver { } } + deleteConversation(conversationId: string): Promise { + return deleteDshConversation(this.options.home, conversationId); + } + /** * One `dsh plugin` invocation, through the same CLI the sessions run. It * runs pnpm in the profile, whose `node_modules` also holds CodeDeck's own diff --git a/packages/agent-host/src/drivers/opencode/__tests__/opencodeDriver.test.ts b/packages/agent-host/src/drivers/opencode/__tests__/opencodeDriver.test.ts index 68745245..07a88e2a 100644 --- a/packages/agent-host/src/drivers/opencode/__tests__/opencodeDriver.test.ts +++ b/packages/agent-host/src/drivers/opencode/__tests__/opencodeDriver.test.ts @@ -45,6 +45,27 @@ function start(client: FakeClient, overrides: Partial = {}, handle return ctx; } +describe('OpenCode conversation delete', () => { + const deleting = (result: unknown) => { + const client = clientWith([]) as FakeClient & { session: { delete: ReturnType } }; + client.session.delete = vi.fn().mockResolvedValue(result); + return { client, driver: OpenCodeDriver.withClient(client) }; + }; + + it('deletes the session on the server, in its directory', async () => { + const { client, driver } = deleting({ data: true, error: undefined, response: { status: 200 } }); + await driver.deleteConversation('ses_1', '/w'); + expect(client.session.delete).toHaveBeenCalledWith({ sessionID: 'ses_1', directory: '/w' }); + }); + + it('a session the server no longer knows is already deleted; another error is not', async () => { + const gone = deleting({ data: undefined, error: { name: 'NotFoundError' }, response: { status: 404 } }); + await expect(gone.driver.deleteConversation('ses_1', '/w')).resolves.toBeUndefined(); + const failed = deleting({ data: undefined, error: { name: 'BadRequest' }, response: { status: 400 } }); + await expect(failed.driver.deleteConversation('ses_1', '/w')).rejects.toThrow(/could not delete session ses_1/); + }); +}); + describe('OpenCode session lifecycle', () => { it('reports its conversation id, becomes ready, and ends normally when the stream closes', async () => { const ctx = start(clientWith([])); diff --git a/packages/agent-host/src/drivers/opencode/driver.ts b/packages/agent-host/src/drivers/opencode/driver.ts index 191199a2..90e0cd17 100644 --- a/packages/agent-host/src/drivers/opencode/driver.ts +++ b/packages/agent-host/src/drivers/opencode/driver.ts @@ -1090,6 +1090,15 @@ export class OpenCodeDriver implements Driver { } } + /** OpenCode deletes a session's subagent sessions with it. */ + async deleteConversation(conversationId: string, cwd: string): Promise { + const client = await this.client(); + const { error, response } = await client.session.delete({ sessionID: conversationId, directory: cwd }); + if (error && response?.status !== 404) { + throw new Error(`OpenCode could not delete session ${conversationId}: ${JSON.stringify(error)}`); + } + } + async shutdown(): Promise { this.stopped = true; await this.server?.close(); diff --git a/packages/agent-host/src/generated/protocol.ts b/packages/agent-host/src/generated/protocol.ts index f3e30e7f..ea23c0e1 100644 --- a/packages/agent-host/src/generated/protocol.ts +++ b/packages/agent-host/src/generated/protocol.ts @@ -131,6 +131,21 @@ export type BridgeMessage_Deserialize = { kind: "end-session"; payload: { sessionId: 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`. + */ +{ kind: "delete-conversation"; payload: { + sessionId: string, + agent: string, + cwd: string, + conversationId: string, +} } | /** Hand user input to the agent. Reply: `ack`. */ { kind: "prompt"; payload: { sessionId: string, @@ -248,6 +263,21 @@ export type BridgeMessage_Serialize = { kind: "end-session"; payload: { sessionId: 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`. + */ +{ kind: "delete-conversation"; payload: { + sessionId: string, + agent: string, + cwd: string, + conversationId: string, +} } | /** Hand user input to the agent. Reply: `ack`. */ { kind: "prompt"; payload: { sessionId: string, diff --git a/packages/agent-host/src/host.ts b/packages/agent-host/src/host.ts index 7654f579..8f890c08 100644 --- a/packages/agent-host/src/host.ts +++ b/packages/agent-host/src/host.ts @@ -63,6 +63,9 @@ function errorText(err: unknown): string { export class AgentHost { private readonly drivers: Map; private readonly sessions = new Map(); + /** Sessions an `end-session` is still stopping: what deleting their + * conversation waits for, since lines are handled concurrently. */ + private readonly ending = new Map>(); private readonly pending = new Map void>(); private nextRequestId = 0; private shuttingDown = false; @@ -144,14 +147,29 @@ export class AgentHost { payload: { hostVersion: this.hostVersion, agents: [...this.drivers.values()].map((d) => d.info()) }, }; case 'end-session': { - const slot = this.sessions.get(message.payload.sessionId); - this.sessions.delete(message.payload.sessionId); + const { sessionId } = message.payload; + const slot = this.sessions.get(sessionId); + this.sessions.delete(sessionId); if (slot) { slot.closed = true; - await slot.session?.end(); + const ending = Promise.resolve(slot.session?.end()); + const settled = ending.catch(() => {}); + this.ending.set(sessionId, settled); + void settled.finally(() => { + if (this.ending.get(sessionId) === settled) this.ending.delete(sessionId); + }); + await ending; } return ack(); } + case 'delete-conversation': { + const { sessionId, agent, cwd, conversationId } = message.payload; + // The agent must have stopped writing to the conversation first. + if (this.sessions.has(sessionId)) throw new Error(`session ${sessionId} is still running`); + await this.ending.get(sessionId); + await this.driver(agent).deleteConversation?.(conversationId, cwd); + return ack(); + } case 'prompt': this.session(message.payload.sessionId).prompt(message.payload.text); return ack();