/** * Host half of the session-delete bundle. * * A conversation is durable in exactly two places: its session log (one * directory per session under the persistence root) and the Workspace registry * that references it. Permanent deletion therefore means: stop the session and * gate it against new work, remove its stored artifacts together with its * subagent descendants' logs and their spilled tool output, and drop the * registry references to it. * * Spilled tool output is the session's one other session-scoped artifact: the * spill backend stores oversized results at `/session-/…`, where * the directory is `sha256(session id)` truncated to 12 hex characters and the * root is either the configured one or a `dsh-spill-*` directory under the OS * temp directory. Those files are otherwise reclaimed only by the backend's own * age sweep, so deletion removes them here. Content-addressed attachments are * deliberately left alone: one blob can be referenced by several sessions, and * nothing in the harness counts those references. * * A session that is still live in this process cannot be removed from the * in-memory Session and Agent stores by another plugin, and deleting a log * under a live writer would let the next append recreate a header-less * artifact. Such a session is archived (durably gated and stopped) now and its * artifacts are removed on the next activation, when nothing is live yet. * * The browser half reaches this operation through the authenticated exact * Fetch route registered below; it is the only cross-half entry point this * package needs. * * @module @dsh-plugin/session-delete */ import { createHash } from 'node:crypto' import { lstat, open, readdir, readFile, rm, stat, writeFile } from 'node:fs/promises' import { dirname, join } from 'node:path' import { homedir, tmpdir } from 'node:os' /** Authenticated exact route the Client half posts one deletion to. */ const ROUTE_PATH = '/api/plugin/session-delete/delete' /** Authenticated exact route the General-settings row posts its sweep to. */ const ORPHAN_ROUTE_PATH = '/api/plugin/session-delete/orphans' /** Durable record of sessions whose artifacts still have to be removed. */ const LEDGER_FILENAME = 'session-delete-pending.json' /** The backend's per-session write lock; it survives its session on POSIX. */ const LEASE_FILENAME = 'session.lock' /** Exactly the shape the local spill backend derives from a session id. */ const SPILL_SESSION_DIRECTORY = /^session-[0-9a-f]{12}$/ /** Services this Host half needs before it may run at all. */ export const inject = [ 'agents', 'connection', 'sessionPersistence', 'sessionQuery', 'sessions', 'workspaceRegistry', ] /** * Install the deletion route and finish deletions deferred to this activation. * @param ctx - Host context carrying the Session, Agent, storage, and Workspace services. */ export function apply(ctx) { const inFlight = new Set() ctx.effect(() => ctx.connection.fetch.register({ path: ROUTE_PATH, methods: ['POST'], requestBody: 'buffered', fetch: (request) => handleRequest(ctx, request, inFlight), }), 'session-delete: delete route') ctx.effect(() => ctx.connection.fetch.register({ path: ORPHAN_ROUTE_PATH, methods: ['POST'], requestBody: 'buffered', fetch: () => handleOrphanSweep(ctx), }), 'session-delete: orphan sweep route') ctx.effect(() => { let disposed = false sweepPending(ctx, () => disposed).catch((error) => { warn(ctx, `could not finish deferred deletions: ${message(error)}`) }) return () => { disposed = true } }, 'session-delete: deferred deletions') } /** * Answer one authenticated deletion request. * @param ctx - Host context. * @param request - buffered Fetch request carrying `{ sessionId }`. * @param inFlight - identities already being deleted in this process. * @returns the JSON result or a stable failure the Client half reports. */ async function handleRequest(ctx, request, inFlight) { let body try { body = await request.json() } catch { return failure(400, 'bad-request', 'the request body must be JSON') } const sessionId = body?.sessionId if (typeof sessionId !== 'string' || sessionId.trim() === '') { return failure(400, 'bad-request', 'sessionId must be a non-empty string') } if (inFlight.has(sessionId)) { return failure(409, 'busy', `conversation "${sessionId}" is already being deleted`) } inFlight.add(sessionId) try { const value = await deleteConversation(ctx, sessionId) return Response.json(value, { headers: { 'cache-control': 'no-store' } }) } catch (error) { warn(ctx, `deleting "${sessionId}" failed: ${message(error)}`) return failure(500, 'delete-failed', message(error)) } finally { inFlight.delete(sessionId) } } /** * Delete one conversation and every subagent conversation under it. * @param ctx - Host context. * @param sessionId - conversation the user asked to delete. * @returns `deleted` when the artifacts are gone, `scheduled` when a live * session forced the removal to this run's end. */ async function deleteConversation(ctx, sessionId) { const headers = await collectHeaders(ctx) const target = headers.get(sessionId) if (target === undefined) throw new Error(`conversation "${sessionId}" was not found`) if (target.origin === 'subagent') { throw new Error('a subagent conversation is removed together with the conversation that owns it') } const ids = [sessionId, ...subagentDescendants(headers, sessionId)] const live = ids.filter((id) => isLive(ctx, id)) if (live.length > 0) { for (const id of ids) await ctx.workspaceRegistry.archiveSession(id, { stopActivity: true }) await rememberPending(ctx, ids) return { ok: true, status: 'scheduled', sessionIds: ids, live } } const removed = [] for (const id of ids) { const header = headers.get(id) if (header === undefined) continue if (await removeArtifacts(ctx, header)) removed.push(id) } const spilled = await removeSpillArtifacts(ctx, ids) for (const id of ids) await forgetReferences(ctx, id) for (const id of removed) announceRemoval(ctx, id) return { ok: true, status: 'deleted', sessionIds: ids, removed, spilledDirectories: spilled.directories, spilledFiles: spilled.files, } } /** * Derive one session's spill directory name exactly as the local spill backend * does: `sha256(session id)` truncated to 12 hex characters. * @param sessionId - owning session identity. * @returns the directory name shared by every root. */ function spillSessionDirectory(sessionId) { return `session-${createHash('sha256').update(String(sessionId)).digest('hex').slice(0, 12)}` } /** * Every spill root this process may have written under: the backend's active * root plus the `dsh-spill-*` default roots the backend itself sweeps. * @param ctx - Host context. * @returns absolute candidate roots. */ async function spillRoots(ctx) { const roots = new Set() const configured = ctx.get('spillStore')?.root if (typeof configured === 'string' && configured !== '') roots.add(configured) const base = tmpdir() let entries = [] try { entries = await readdir(base, { withFileTypes: true }) } catch (error) { warn(ctx, `spill roots under ${base} could not be listed: ${message(error)}`) } for (const entry of entries) { if (entry.isDirectory() && entry.name.startsWith('dsh-spill')) roots.add(join(base, entry.name)) } return [...roots] } /** * Remove the spilled tool output of the given sessions. * * Only a real, session-named directory under a spill root is removed; a symlink, * a file, or any other entry is left untouched, and a failure is reported * without failing the deletion that already happened. * * @param ctx - Host context. * @param sessionIds - identities whose spilled output goes away. * @returns how many directories and files were removed. */ async function removeSpillArtifacts(ctx, sessionIds) { const names = sessionIds.map(spillSessionDirectory) let directories = 0 let files = 0 for (const root of await spillRoots(ctx)) { for (const name of names) { const directory = join(root, name) let info try { info = await lstat(directory) } catch { continue } if (info.isSymbolicLink() || !info.isDirectory()) continue try { files += (await readdir(directory)).length } catch { /* a listing failure never blocks the removal */ } try { await rm(directory, { recursive: true, force: true }) directories += 1 } catch (error) { warn(ctx, `spilled output at ${directory} was kept: ${message(error)}`) } } } return { directories, files } } /** * Answer the General-settings cleanup request. * @param ctx - Host context. * @returns the JSON report the settings row shows. */ async function handleOrphanSweep(ctx) { try { const report = await sweepOrphanSpill(ctx) return Response.json({ ok: true, ...report }, { headers: { 'cache-control': 'no-store' } }) } catch (error) { warn(ctx, `orphan sweep failed: ${message(error)}`) return failure(500, 'sweep-failed', message(error)) } } /** * Remove spilled tool output that no session known to this home directory can * own. * * A spill directory is orphaned only when its name has exactly the * `session-` shape the spill backend derives AND no * session known to this process — stored or live — hashes to it. Anything else, * including a symlink or an entry the backend would never create, is left alone. * * @param ctx - Host context. * @returns counts for the settings row. */ async function sweepOrphanSpill(ctx) { const headers = await collectHeaders(ctx) const known = new Set([...headers.keys()].map(spillSessionDirectory)) const report = { directories: 0, files: 0, bytes: 0, roots: 0, sessions: headers.size } for (const root of await spillRoots(ctx)) { report.roots += 1 let entries try { entries = await readdir(root, { withFileTypes: true }) } catch (error) { warn(ctx, `spill root ${root} could not be listed: ${message(error)}`) continue } for (const entry of entries) { if (!entry.isDirectory() || !SPILL_SESSION_DIRECTORY.test(entry.name)) continue if (known.has(entry.name)) continue const directory = join(root, entry.name) try { const info = await lstat(directory) if (info.isSymbolicLink() || !info.isDirectory()) continue const measured = await measureDirectory(directory) await rm(directory, { recursive: true, force: true }) report.directories += 1 report.files += measured.files report.bytes += measured.bytes } catch (error) { warn(ctx, `orphaned spill directory ${directory} was kept: ${message(error)}`) } } } return report } /** * Measure one spill session directory. The backend writes every artifact * directly into it, so leaf files are the whole content. * @param directory - spill session directory. * @returns file count and total bytes. */ async function measureDirectory(directory) { let files = 0 let bytes = 0 let entries try { entries = await readdir(directory, { withFileTypes: true }) } catch { return { files, bytes } } for (const entry of entries) { if (!entry.isFile()) continue try { bytes += (await stat(join(directory, entry.name))).size files += 1 } catch { /* a file that vanished mid-measure simply does not count */ } } return { files, bytes } } /** * Read every session header this process knows: stored ones plus live ones. * @param ctx - Host context. * @returns headers keyed by session id. */ async function collectHeaders(ctx) { const headers = new Map() for (const record of await ctx.sessionQuery.listSessions()) { const header = record?.header if (header !== undefined) headers.set(String(header.id), header) } for (const session of ctx.get('sessions')?.list() ?? []) { headers.set(String(session.id), session.header) } return headers } /** * Collect the subagent conversations stored underneath one conversation. * @param headers - every known session header. * @param rootId - the conversation being deleted. * @returns descendant ids, nearest first. */ function subagentDescendants(headers, rootId) { const children = new Map() for (const header of headers.values()) { if (header.origin !== 'subagent' || header.parentSession === undefined) continue const parent = String(header.parentSession) const rows = children.get(parent) if (rows === undefined) children.set(parent, [String(header.id)]) else rows.push(String(header.id)) } const found = [] const seen = new Set([rootId]) const queue = [rootId] while (queue.length > 0) { for (const child of children.get(queue.shift()) ?? []) { if (seen.has(child)) continue seen.add(child) found.push(child) queue.push(child) } } return found } /** * Whether this process currently holds the session or its agent in memory. * @param ctx - Host context. * @param sessionId - candidate identity. * @returns true while the session is live. */ function isLive(ctx, sessionId) { return ctx.agents.get(sessionId) !== undefined || ctx.get('sessions')?.get(sessionId) !== undefined } /** * Remove every stored generation of one session, keeping its write lock file. * @param ctx - Host context. * @param header - the session's stored header, which names its artifact. * @returns whether anything was removed. */ async function removeArtifacts(ctx, header) { const persistence = ctx.sessionPersistence if (typeof persistence.locate !== 'function') { throw new Error('this session storage backend cannot locate stored artifacts, so nothing was deleted') } const located = persistence.locate(header) const artifactPath = located?.path if (typeof artifactPath !== 'string' || artifactPath === '') { throw new Error(`conversation "${String(header.id)}" has no stored artifact to delete`) } const directory = dirname(artifactPath) let entries try { entries = await readdir(directory, { withFileTypes: true }) } catch (error) { if (error?.code === 'ENOENT') return false throw error } let removed = false for (const entry of entries) { if (entry.name === LEASE_FILENAME) continue await rm(join(directory, entry.name), { recursive: true, force: true }) removed = true } if (removed && process.platform !== 'win32') await syncDirectory(directory) return removed } /** * Flush one directory entry set on POSIX, where removal is not durable until * the parent directory is synced. * @param directory - directory to sync. */ async function syncDirectory(directory) { const handle = await open(directory, 'r') try { await handle.sync() } finally { await handle.close() } } /** * Drop every durable reference to a session whose artifacts are gone: * Workspace accounting, the archive set, and the pin set. * @param ctx - Host context. * @param sessionId - removed identity. */ async function forgetReferences(ctx, sessionId) { const registry = ctx.workspaceRegistry for (const workspace of registry.list()) { try { await workspace.detachSession(sessionId) } catch (error) { warn(ctx, `workspace "${workspace.id}" kept a reference to "${sessionId}": ${message(error)}`) } } for (const [operation, run] of [ ['unpin', () => registry.unpinSession(sessionId)], ['unarchive', () => registry.unarchiveSession(sessionId)], ]) { try { await run() } catch (error) { warn(ctx, `${operation} of "${sessionId}" failed: ${message(error)}`) } } } /** Announce one removed conversation so connected pages drop its row. */ function announceRemoval(ctx, sessionId) { try { ctx.emit('api-session/removed', sessionId) } catch (error) { warn(ctx, `could not announce removal of "${sessionId}": ${message(error)}`) } } /** * Finish deletions deferred by a live session. Each id is deleted only while * it is still archived, so restoring a conversation in the sidebar cancels its * pending deletion. * @param ctx - Host context. * @param isDisposed - whether this plugin already unloaded. */ async function sweepPending(ctx, isDisposed) { const pending = await readLedger(ctx) if (pending.length === 0) return const registry = ctx.workspaceRegistry const headers = await collectHeaders(ctx) const remaining = [] for (const sessionId of pending) { if (isDisposed()) { remaining.push(sessionId) continue } if (isLive(ctx, sessionId)) { remaining.push(sessionId) continue } if (!registry.archivedSessionIds.includes(sessionId)) continue try { const header = headers.get(sessionId) if (header !== undefined) await removeArtifacts(ctx, header) await removeSpillArtifacts(ctx, [sessionId]) await forgetReferences(ctx, sessionId) announceRemoval(ctx, sessionId) } catch (error) { warn(ctx, `could not finish deleting "${sessionId}": ${message(error)}`) remaining.push(sessionId) } } await writeLedger(ctx, remaining) } /** * Read the deferred-deletion record. * @param ctx - Host context. * @returns session ids still waiting for removal. */ async function readLedger(ctx) { try { const parsed = JSON.parse(await readFile(ledgerPath(ctx), 'utf8')) return Array.isArray(parsed?.sessions) ? parsed.sessions.filter((id) => typeof id === 'string' && id !== '') : [] } catch { return [] } } /** * Replace the deferred-deletion record. * @param ctx - Host context. * @param sessionIds - ids still waiting for removal. */ async function writeLedger(ctx, sessionIds) { const unique = [...new Set(sessionIds)] const path = ledgerPath(ctx) if (unique.length === 0) { await rm(path, { force: true }) return } await writeFile(path, `${JSON.stringify({ version: 1, sessions: unique }, null, 2)}\n`, 'utf8') } /** Record the ids whose artifacts are removed on the next activation. */ function rememberPending(ctx, sessionIds) { return readLedger(ctx).then((pending) => writeLedger(ctx, [...pending, ...sessionIds])) } /** * The deferred-deletion record lives beside the session store it describes, so * it follows `$DSH_HOME` without this plugin reading harness configuration. * @param ctx - Host context. * @returns absolute ledger path. */ function ledgerPath(ctx) { const root = ctx.sessionPersistence?.root if (typeof root === 'string' && root !== '') return join(dirname(root), LEDGER_FILENAME) const home = process.env.DSH_HOME const base = typeof home === 'string' && home !== '' ? home : join(homedir(), '.dsh') return join(base, LEDGER_FILENAME) } /** Build one stable JSON failure the Client half surfaces verbatim. */ function failure(status, code, text) { return Response.json( { ok: false, error: { code, message: text } }, { status, headers: { 'cache-control': 'no-store' } }, ) } /** Log one contained diagnostic. */ function warn(ctx, text) { try { ctx.logger?.warn?.(`session-delete: ${text}`) } catch { /* a diagnostic never fails the operation it describes */ } } /** Render one unknown failure as text. */ function message(error) { return error instanceof Error ? error.message : String(error) }