refactor!: align with the DSH plugin conventions and drop the npm publish workflow
Follow the cordis-plugin-development references instead of the ad-hoc choices the first version made. Client half: - register the row button under this package's own id instead of shadowing the shipped `archive` action at a lower priority, so the official hover buttons keep their cells and archiving stays where the harness put it - drop the synthetic `pointerout` dispatched at a `[data-row-key]` ancestor: a plugin does not read or drive another package's DOM, so the tooltip is now positioned from its own button alone - keep the self-rendered primitives, the `--dsw-*` token-only styling, and the modal focus/Escape behavior, and document why Host half: - resolve a conversation's descendants with `sessionQuery.traceSession()` instead of listing every stored header and re-deriving `parentSession` edges - decide liveness from the traced `SessionRecord.live` flag, which removes the `agents` and `sessions` dependencies from `inject` - build the orphan-sweep corpus from `sessionQuery.listSessions()`, which already merges live and persisted sessions - document that `sessionPersistence.locate()` is a JSONL-backend diagnostic hook, not part of the seam, and keep probing it explicitly - record only tree roots in the deferred-deletion ledger: the activation sweep removes a root with everything under it, so a separately recorded child was deleted twice or needed a pass that never runs - return the activation's sweep promise from `apply()` so a test can await it Manifest and docs: - version 2.0.0, `private`, `dsh.manifestVersion`, and `engines` - delete the Gitea npm publish workflow and every npm-publishing task: the package is distributed only through the Git repository and tags - rewrite the README around the current install paths and the plugin's limits
This commit is contained in:
@@ -8,14 +8,21 @@
|
||||
* 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 `<root>/session-<hash>/…`, 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.
|
||||
* 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
|
||||
@@ -23,9 +30,10 @@
|
||||
* 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.
|
||||
* 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
|
||||
*/
|
||||
@@ -47,17 +55,19 @@ 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.
|
||||
* 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()
|
||||
@@ -76,15 +86,18 @@ export function apply(ctx) {
|
||||
fetch: () => handleOrphanSweep(ctx),
|
||||
}), 'session-delete: orphan sweep route')
|
||||
|
||||
let swept
|
||||
ctx.effect(() => {
|
||||
let disposed = false
|
||||
sweepPending(ctx, () => disposed).catch((error) => {
|
||||
swept = sweepPending(ctx, () => disposed).catch((error) => {
|
||||
warn(ctx, `could not finish deferred deletions: ${message(error)}`)
|
||||
})
|
||||
return () => {
|
||||
disposed = true
|
||||
}
|
||||
}, 'session-delete: deferred deletions')
|
||||
|
||||
return swept
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -112,7 +125,7 @@ async function handleRequest(ctx, request, inFlight) {
|
||||
inFlight.add(sessionId)
|
||||
try {
|
||||
const value = await deleteConversation(ctx, sessionId)
|
||||
return Response.json(value, { headers: { 'cache-control': 'no-store' } })
|
||||
return json(value)
|
||||
} catch (error) {
|
||||
warn(ctx, `deleting "${sessionId}" failed: ${message(error)}`)
|
||||
return failure(500, 'delete-failed', message(error))
|
||||
@@ -126,28 +139,21 @@ async function handleRequest(ctx, request, inFlight) {
|
||||
* @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.
|
||||
* session forced the removal to the next activation.
|
||||
*/
|
||||
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 lineage = await resolveLineage(ctx, sessionId)
|
||||
const ids = lineage.ids
|
||||
|
||||
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 }
|
||||
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 = headers.get(id)
|
||||
const header = lineage.headers.get(id)
|
||||
if (header === undefined) continue
|
||||
if (await removeArtifacts(ctx, header)) removed.push(id)
|
||||
}
|
||||
@@ -164,6 +170,65 @@ async function deleteConversation(ctx, sessionId) {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
@@ -246,7 +311,7 @@ async function removeSpillArtifacts(ctx, sessionIds) {
|
||||
async function handleOrphanSweep(ctx) {
|
||||
try {
|
||||
const report = await sweepOrphanSpill(ctx)
|
||||
return Response.json({ ok: true, ...report }, { headers: { 'cache-control': 'no-store' } })
|
||||
return json({ ok: true, ...report })
|
||||
} catch (error) {
|
||||
warn(ctx, `orphan sweep failed: ${message(error)}`)
|
||||
return failure(500, 'sweep-failed', message(error))
|
||||
@@ -259,16 +324,22 @@ async function handleOrphanSweep(ctx) {
|
||||
*
|
||||
* A spill directory is orphaned only when its name has exactly the
|
||||
* `session-<sha256(session id)[:12]>` 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.
|
||||
* 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 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 }
|
||||
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
|
||||
@@ -325,71 +396,24 @@ async function measureDirectory(directory) {
|
||||
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.
|
||||
*
|
||||
* 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.sessionPersistence
|
||||
if (typeof persistence.locate !== 'function') {
|
||||
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)
|
||||
@@ -433,7 +457,9 @@ async function syncDirectory(directory) {
|
||||
|
||||
/**
|
||||
* Drop every durable reference to a session whose artifacts are gone:
|
||||
* Workspace accounting, the archive set, and the pin set.
|
||||
* 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.
|
||||
*/
|
||||
@@ -468,9 +494,15 @@ function announceRemoval(ctx, sessionId) {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*/
|
||||
@@ -478,24 +510,39 @@ 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 headers = await collectHeaders(ctx)
|
||||
const archivedSessionIds = registry.archivedSessionIds
|
||||
const remaining = []
|
||||
for (const sessionId of pending) {
|
||||
if (isDisposed()) {
|
||||
remaining.push(sessionId)
|
||||
continue
|
||||
}
|
||||
if (isLive(ctx, sessionId)) {
|
||||
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 (!registry.archivedSessionIds.includes(sessionId)) continue
|
||||
if (!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)
|
||||
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)}`)
|
||||
@@ -536,9 +583,21 @@ async function writeLedger(ctx, sessionIds) {
|
||||
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]))
|
||||
/**
|
||||
* 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])
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -548,19 +607,21 @@ function rememberPending(ctx, sessionIds) {
|
||||
* @returns absolute ledger path.
|
||||
*/
|
||||
function ledgerPath(ctx) {
|
||||
const root = ctx.sessionPersistence?.root
|
||||
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 Response.json(
|
||||
{ ok: false, error: { code, message: text } },
|
||||
{ status, headers: { 'cache-control': 'no-store' } },
|
||||
)
|
||||
return json({ ok: false, error: { code, message: text } }, status)
|
||||
}
|
||||
|
||||
/** Log one contained diagnostic. */
|
||||
|
||||
Reference in New Issue
Block a user