/** * 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. * * The harness deliberately has no deletion API of its own — the shipped JSONL * backend's own documentation says session files accumulate under `root` until * something removes them externally — so removing them is this bundle's whole * job. Two artifacts are deliberately left alone: * * - content-addressed attachments, because one blob can be referenced by * several sessions and nothing in the harness counts those references; * - the projection cache, which is derived and carries a lifecycle identity * check on read, so it can never resolve to a different conversation. * * 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. * * 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. * * Cross-half contract: the browser half reaches both operations through the * authenticated exact Fetch routes registered below (`ctx.connection.fetch`), * which is how a static client bundle talks to its Host half — the `host.call` * channel belongs to the dynamic client runner, not to a package bundle. * * @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 = [ 'connection', 'sessionPersistence', 'sessionQuery', 'workspaceRegistry', ] /** * Install the deletion routes and finish deletions deferred to this activation. * * @param ctx - Host context carrying the persistence, query, and Workspace services. * @returns the activation's deferred-deletion promise. The runtime ignores it * (the effect owns the work and catches its own failures); awaiting it is how * a test observes that the startup sweep has settled. */ 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') let swept ctx.effect(() => { let disposed = false swept = sweepPending(ctx, () => disposed).catch((error) => { warn(ctx, `could not finish deferred deletions: ${message(error)}`) }) return () => { disposed = true } }, 'session-delete: deferred deletions') return swept } /** * 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 json(value) } 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 the next activation. */ async function deleteConversation(ctx, sessionId) { const lineage = await resolveLineage(ctx, sessionId) const ids = lineage.ids if (lineage.liveIds.length > 0) { for (const id of lineage.ids) await ctx.workspaceRegistry.archiveSession(id, { stopActivity: true }) await rememberPending(ctx, sessionId) return { ok: true, status: 'scheduled', sessionIds: ids, live: lineage.liveIds } } const removed = [] for (const id of ids) { const header = lineage.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, } } /** * Resolve one conversation and its subagent descendants from the query * service's lineage trace. * * `traceSession` is the harness's own corpus observation: one call answers the * whole parent/child structure, so this bundle never has to list every stored * header and re-derive `parentSession` edges itself. Each copied * `SessionRecord` also carries the header the deletion needs and a `live` flag, * which is what keeps the `agents` and `sessions` services out of `inject`. * * @param ctx - Host context. * @param sessionId - conversation the user asked to delete. * @returns the target and descendant ids (parents first), their headers, and * which of them this process currently holds live. */ async function resolveLineage(ctx, sessionId) { let trace try { trace = await ctx.sessionQuery.traceSession(sessionId) } catch (error) { throw new Error(`conversation "${sessionId}" was not found: ${message(error)}`) } const target = trace?.target const header = target?.header if (target === undefined || header === undefined) { throw new Error(`conversation "${sessionId}" was not found`) } if (header.origin === 'subagent') { throw new Error('a subagent conversation is removed together with the conversation that owns it') } const headers = new Map([[String(header.id), header]]) const ids = [String(header.id)] const liveIds = target.live === true ? [String(header.id)] : [] collectDescendants(trace.descendants, headers, ids, liveIds) return { headers, ids, liveIds } } /** * Walk one already-traced descendant tree into flat id order. * @param nodes - lineage nodes, nearest generation first. * @param headers - accumulator keyed by session id. * @param ids - accumulator in parents-first order. * @param liveIds - accumulator of ids this process holds live. */ function collectDescendants(nodes, headers, ids, liveIds) { for (const node of Array.isArray(nodes) ? nodes : []) { const header = node?.session?.header if (header === undefined) continue const id = String(header.id) if (headers.has(id)) continue headers.set(id, header) ids.push(id) if (node.session.live === true) liveIds.push(id) collectDescendants(node.descendants, headers, ids, liveIds) } } /** * 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 json({ ok: true, ...report }) } 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 hashes to it. `listSessions()` already merges * live and persisted sessions into one logical corpus, so "known" needs no * second source. 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 records = await ctx.sessionQuery.listSessions() const known = new Set() for (const record of records) { const id = record?.header?.id if (id !== undefined) known.add(spillSessionDirectory(id)) } const report = { directories: 0, files: 0, bytes: 0, roots: 0, sessions: known.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 } } /** * Remove every stored generation of one session, keeping its write lock file. * * The absolute artifact path comes from the backend's `locate()` hook. That * hook is not part of the abstract `sessionPersistence` contract — the seam * itself declares only `create`, `open`, `flush`, `stat`, and `list` — so it is * probed rather than assumed, and a backend without it fails this one * conversation instead of silently deleting nothing. The returned path names * the highest canonical generation *file*, so its parent directory is the * session-owned directory whose contents go away. * * @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.get('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. Each call is * idempotent for an id it does not hold, so a failure is reported and the * remaining references are still cleared. * @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. * * Every ledger entry is a root conversation, so each iteration resolves its own * subtree and removes it whole; a child is never recorded separately, because * the sweep that removes a root may not re-run for an entry a sibling root * already took with it. An entry is deleted only while it is still archived, so * restoring a conversation in the sidebar cancels its pending deletion: the * durable archive set is what the sidebar's own unarchive drops. * * @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 // The registry-global archive set; membership is the user's own "still // archived" decision and changes only through the sidebar or this plugin. const registry = ctx.workspaceRegistry const archivedSessionIds = registry.archivedSessionIds const remaining = [] for (const sessionId of pending) { if (isDisposed()) { remaining.push(sessionId) continue } let lineage try { lineage = await resolveLineage(ctx, sessionId) } catch (error) { // The conversation is not in the corpus any more: what is left to do is // drop the registry references a previous run did not reach, and stop // retrying it on every activation. warn(ctx, `deferred deletion of "${sessionId}" resolved to nothing: ${message(error)}`) await forgetReferences(ctx, sessionId) continue } if (lineage.liveIds.length > 0) { remaining.push(sessionId) continue } if (!archivedSessionIds.includes(sessionId)) continue try { for (const id of lineage.ids) { const header = lineage.headers.get(id) if (header !== undefined) await removeArtifacts(ctx, header) } await removeSpillArtifacts(ctx, lineage.ids) for (const id of lineage.ids) await forgetReferences(ctx, id) 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 root conversation whose artifacts are removed on the next * activation. * * Only a tree root is recorded: the activation sweep removes a root together * with everything under it, so a separately recorded child would either be * deleted twice or, once its root had taken it, need its entry dropped on a * pass that never runs. * * @param ctx - Host context. * @param sessionId - the root conversation of one deferred deletion. */ async function rememberPending(ctx, sessionId) { const pending = await readLedger(ctx) await writeLedger(ctx, [...pending, sessionId]) } /** * 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.get('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) } /** Answer one JSON value the Client half reads, never cached. */ function json(value, status = 200) { return Response.json(value, { status, headers: { 'cache-control': 'no-store' } }) } /** Build one stable JSON failure the Client half surfaces verbatim. */ function failure(status, code, text) { return json({ ok: false, error: { code, message: text } }, status) } /** 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) }