diff --git a/.changeset/session-delete-action.md b/.changeset/session-delete-action.md new file mode 100644 index 00000000000..8b69403b4f4 --- /dev/null +++ b/.changeset/session-delete-action.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": minor +--- + +Expose permanent session deletion via `POST /api/v1/sessions/{session_id}:delete` and broadcast `event.session.deleted` over WebSocket. diff --git a/apps/kimi-inspect/src/activity/store.test.ts b/apps/kimi-inspect/src/activity/store.test.ts index 64f191f3a15..3e600fe8122 100644 --- a/apps/kimi-inspect/src/activity/store.test.ts +++ b/apps/kimi-inspect/src/activity/store.test.ts @@ -199,6 +199,21 @@ describe('SessionActivityHub', () => { expect(hub.store.get('s1')).toBeUndefined(); expect(onListChanged).toHaveBeenCalledTimes(1); + instances[0]!.emitFrame({ + type: 'event.session.work_changed', + session_id: 's2', + payload: { type: 'event.session.work_changed', busy: true }, + }); + expect(hub.store.get('s2')).toBeDefined(); + + instances[0]!.emitFrame({ + type: 'event.session.deleted', + session_id: '__global__', + payload: { type: 'event.session.deleted', sessionId: 's2', workspace_id: 'wd_1' }, + }); + expect(hub.store.get('s2')).toBeUndefined(); + expect(onListChanged).toHaveBeenCalledTimes(2); + for (const type of [ 'event.workspace.created', 'event.workspace.updated', @@ -206,7 +221,7 @@ describe('SessionActivityHub', () => { ]) { instances[0]!.emitFrame({ type, session_id: '__global__', payload: {} }); } - expect(onListChanged).toHaveBeenCalledTimes(4); + expect(onListChanged).toHaveBeenCalledTimes(5); hub.close(); }); }); diff --git a/apps/kimi-inspect/src/activity/store.ts b/apps/kimi-inspect/src/activity/store.ts index addadaca64f..af77ad6c062 100644 --- a/apps/kimi-inspect/src/activity/store.ts +++ b/apps/kimi-inspect/src/activity/store.ts @@ -106,6 +106,10 @@ export class SessionActivityHub { this.store.remove(sessionId); opts.onListChanged(); }, + onSessionDeleted: (sessionId) => { + this.store.remove(sessionId); + opts.onListChanged(); + }, onWorkspaceChanged: () => opts.onListChanged(), onReconnected: () => void this.seed(), }, diff --git a/apps/kimi-inspect/src/activity/ws.ts b/apps/kimi-inspect/src/activity/ws.ts index f4ecbe2149b..6a1887f05c4 100644 --- a/apps/kimi-inspect/src/activity/ws.ts +++ b/apps/kimi-inspect/src/activity/ws.ts @@ -64,6 +64,10 @@ export interface GlobalEventsWsHandlers { * carries the `__global__` watermark; the real session id rides in the * payload. */ onSessionArchived?: ((sessionId: string) => void) | undefined; + /** A session was permanently deleted (list-level signal). Same envelope + * shape as `event.session.archived`: the real session id rides in the + * payload. */ + onSessionDeleted?: (sessionId: string) => void; /** A workspace was created / updated / deleted (list-level signal). */ onWorkspaceChanged?: (() => void) | undefined; /** A DI unit of the engine's scope tree changed state (debug feed). */ @@ -195,6 +199,14 @@ export class GlobalEventsWs { } return; } + case 'event.session.deleted': { + const payload = frame.payload as { sessionId?: unknown } | undefined; + const deletedId = payload?.sessionId; + if (typeof deletedId === 'string' && deletedId !== '') { + this.handlers.onSessionDeleted?.(deletedId); + } + return; + } case 'event.workspace.created': case 'event.workspace.updated': case 'event.workspace.deleted': { diff --git a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleEvents.ts b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleEvents.ts index d762ddd641b..155e250c6eb 100644 --- a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleEvents.ts +++ b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleEvents.ts @@ -13,6 +13,18 @@ export interface SessionArchived { readonly payload: SessionArchivedPayload; } +export interface SessionDeletedPayload { + readonly sessionId: string; + readonly workspaceId: string; +} + +export class SessionDeleted extends Event2<{ readonly payload: SessionDeletedPayload }> { + static override readonly type = 'event.session.deleted'; +} +export interface SessionDeleted { + readonly payload: SessionDeletedPayload; +} + export interface SessionCreatedPayload { readonly agentId: string; readonly sessionId: string; diff --git a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts index a6b2ac436e3..abfc5cd0567 100644 --- a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts @@ -82,7 +82,7 @@ import { IWorkspaceMcpService } from '#/workspace/workspaceMcp/workspaceMcp'; import { PLUGIN_SKILL_SOURCE_ID } from '#/features/skill/catalog/skillSource'; import { agentScopeOf, sessionDirOf, sessionScopeOf } from './internal/addressing'; -import { SessionArchived } from './sessionLifecycleEvents'; +import { SessionArchived, SessionDeleted } from './sessionLifecycleEvents'; import { assertForkTurnIndex, sliceMainRecordsAtTurn, @@ -449,6 +449,11 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.index.remove(sessionId); this.appendLogStore.append('', 'session_index.jsonl', { sessionId, deleted: true }); await this.appendLogStore.flush(); + this.event.publish( + new SessionDeleted({ + payload: { sessionId, workspaceId: this.workspaceContext.workspaceId }, + }), + ); } private async announceWillClose(event: SessionWillCloseEvent): Promise { diff --git a/packages/kap-server/src/openapi/transforms.ts b/packages/kap-server/src/openapi/transforms.ts index b2c8b860385..2713faa6dcb 100644 --- a/packages/kap-server/src/openapi/transforms.ts +++ b/packages/kap-server/src/openapi/transforms.ts @@ -41,7 +41,10 @@ import { questionResolveRequestSchema, questionResolveResultSchema, } from '../protocol/rest-question'; -import { archiveSessionResponseSchema } from '../protocol/rest-session'; +import { + archiveSessionResponseSchema, + deleteSessionResponseSchema, +} from '../protocol/rest-session'; const binarySchema = { type: 'string', @@ -183,18 +186,32 @@ function patchSessionAction(paths: Record): void { const operation = asRecord(pathItem?.['post']); if (pathItem === undefined || operation === undefined) return; + projectSessionAction(paths, pathItem, 'archive', 'runSessionArchiveAction', { + description: 'Session archive response', + content: jsonContent(openApiDocumentEnvelopeJsonSchema(archiveSessionResponseSchema)), + }); + projectSessionAction(paths, pathItem, 'delete', 'runSessionDeleteAction', { + description: 'Session delete response', + content: jsonContent(openApiDocumentEnvelopeJsonSchema(deleteSessionResponseSchema)), + }); + delete paths[internalPath]; +} + +function projectSessionAction( + paths: Record, + pathItem: Record, + action: string, + operationId: string, + okResponse: Record, +): void { const cloned = cloneRecord(pathItem); replacePathParamName(cloned, 'tail', 'session_id'); const clonedOperation = asRecord(cloned['post']); if (clonedOperation !== undefined) { - clonedOperation['operationId'] = 'runSessionArchiveAction'; - setResponse(clonedOperation, '200', { - description: 'Session archive response', - content: jsonContent(openApiDocumentEnvelopeJsonSchema(archiveSessionResponseSchema)), - }); + clonedOperation['operationId'] = operationId; + setResponse(clonedOperation, '200', okResponse); } - paths['/api/v1/sessions/{session_id}:archive'] = cloned; - delete paths[internalPath]; + paths[`/api/v1/sessions/{session_id}:${action}`] = cloned; } function patchFsAction(paths: Record): void { diff --git a/packages/kap-server/src/protocol/events-zod.ts b/packages/kap-server/src/protocol/events-zod.ts index 720d59d7bea..a72736449cc 100644 --- a/packages/kap-server/src/protocol/events-zod.ts +++ b/packages/kap-server/src/protocol/events-zod.ts @@ -584,6 +584,11 @@ export const sessionArchivedEventSchema = z.object({ workspace_id: z.string().min(1), }); +export const sessionDeletedEventSchema = z.object({ + type: z.literal('event.session.deleted'), + workspace_id: z.string().min(1), +}); + export const workspaceCreatedEventSchema = z.object({ type: z.literal('event.workspace.created'), workspace: workspaceSchema, @@ -1055,6 +1060,7 @@ export const agentEventSchema = z.discriminatedUnion('type', [ sessionMetaUpdatedEventSchema, sessionCreatedEventSchema, sessionArchivedEventSchema, + sessionDeletedEventSchema, workspaceCreatedEventSchema, workspaceUpdatedEventSchema, workspaceDeletedEventSchema, diff --git a/packages/kap-server/src/protocol/rest-session.ts b/packages/kap-server/src/protocol/rest-session.ts index 1abc32e6b9f..41d037b82aa 100644 --- a/packages/kap-server/src/protocol/rest-session.ts +++ b/packages/kap-server/src/protocol/rest-session.ts @@ -162,8 +162,10 @@ export type ArchiveSessionResponse = z.infer; -export const deleteSessionResponseSchema = archiveSessionResponseSchema; -export type DeleteSessionResponse = ArchiveSessionResponse; +export const deleteSessionResponseSchema = z.object({ + deleted: z.literal(true), +}); +export type DeleteSessionResponse = z.infer; export const sessionAbortResponseSchema = z.object({ aborted: z.boolean(), diff --git a/packages/kap-server/src/routes/sessions.ts b/packages/kap-server/src/routes/sessions.ts index 7230c2cdfbc..a0eff3061b9 100644 --- a/packages/kap-server/src/routes/sessions.ts +++ b/packages/kap-server/src/routes/sessions.ts @@ -32,6 +32,7 @@ import { type SessionSummary, } from '@moonshot-ai/agent-core-v2'; import { SessionMetaUpdated } from '@moonshot-ai/agent-core-v2/session/sessionMetadata/sessionMetaEvents'; +import { IGlobalSearchService } from '../search/searchService'; import { ErrorCode } from '../protocol/error-codes'; import { pageResponseSchema } from '../protocol/pagination'; import { toProtocolMessage } from '../services/messages/messageProjection'; @@ -41,6 +42,7 @@ import { compactSessionResponseSchema, createSessionChildRequestSchema, createSessionRequestSchema, + deleteSessionResponseSchema, forkSessionRequestSchema, getSessionGoalResponseSchema, listSessionChildrenResponseSchema, @@ -594,6 +596,7 @@ export function registerSessionsRoutes(app: SessionRouteHost, core: Scope): void sessionAbortResponseSchema, startBtwSessionResponseSchema, archiveSessionResponseSchema, + deleteSessionResponseSchema, ]), }, errors: { @@ -844,7 +847,15 @@ export function registerSessionsRoutes(app: SessionRouteHost, core: Scope): void ); } -type SessionAction = 'fork' | 'compact' | 'undo' | 'abort' | 'btw' | 'restore' | 'archive'; +type SessionAction = + | 'fork' + | 'compact' + | 'undo' + | 'abort' + | 'btw' + | 'restore' + | 'archive' + | 'delete'; interface SessionActionExtra { readonly core: Scope; @@ -865,6 +876,7 @@ const sessionActions: ActionTable = { btw: { handle: btwSessionAction }, restore: { handle: restoreSessionAction }, archive: { handle: archiveSessionAction }, + delete: { handle: deleteSessionAction }, }; async function forkSessionAction( @@ -982,6 +994,20 @@ async function archiveSessionAction(ctx: SessionActionCtx): Promise { reply.send(okEnvelope({ archived: true }, req.id)); } +async function deleteSessionAction(ctx: SessionActionCtx): Promise { + const { core, req, reply, id } = ctx; + try { + await core.accessor.get(ISessionManager).delete(id); + } catch (error) { + if (!isError2(error) || error.code !== ErrorCodes.SESSION_NOT_FOUND) throw error; + await core.accessor.get(IGlobalSearchService).deleteSession(id); + throw error; + } + await core.accessor.get(IGlobalSearchService).deleteSession(id); + requestLog(req)?.info({ session_id: id, action: 'delete' }, 'session action completed'); + reply.send(okEnvelope({ deleted: true }, req.id)); +} + export interface SessionWireFields { readonly id: string; readonly workspaceId: string; diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index 70750caeebf..800eb1f3a4e 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -1,6 +1,6 @@ import { createHash } from 'node:crypto'; -import { open, readFile, readdir, rm, stat } from 'node:fs/promises'; -import { join, relative } from 'node:path'; +import { appendFile, mkdir, open, readFile, readdir, rm, stat, writeFile } from 'node:fs/promises'; +import { dirname, join, relative } from 'node:path'; import { LockError, @@ -34,6 +34,7 @@ import { import { analyzeWireLine, type StepEffect, type TurnEffect } from './wireExtract.ts'; const TEXT_INDEX_NAME = 'body'; +const DELETED_LEDGER_HORIZON_MS = 24 * 60 * 60 * 1000; const TRI_INDEX_NAME = 'tri'; const WIRE_FILENAME = 'wire.jsonl'; @@ -239,6 +240,92 @@ export class SearchIndexCore { return this.options.log; } + private get deletedLedgerPath(): string { + return join(this.indexDir, 'deleted-sessions.jsonl'); + } + + private get deletedLedgerSnapshotPath(): string { + return join(this.indexDir, 'deleted-sessions.snapshot.json'); + } + + private parseDeletedLedgerLines(raw: string, into: Map): void { + for (const line of raw.split('\n')) { + if (line.length === 0) continue; + if (line.startsWith('-')) { + into.delete(line.slice(1).split('\t')[0]!); + continue; + } + const [id, at] = line.split('\t'); + if (id !== undefined && id.length > 0) into.set(id, Number(at ?? 0)); + } + } + + private async readDeletedLedger(): Promise> { + const ledger = new Map(); + let watermark = 0; + try { + const snapshot: unknown = JSON.parse(await readFile(this.deletedLedgerSnapshotPath, 'utf8')); + if (typeof snapshot === 'object' && snapshot !== null) { + const w = (snapshot as { watermark?: unknown }).watermark; + if (typeof w === 'number') watermark = w; + const entries = (snapshot as { entries?: unknown }).entries; + if (Array.isArray(entries)) { + for (const entry of entries) { + if (Array.isArray(entry) && typeof entry[0] === 'string') { + ledger.set(entry[0], typeof entry[1] === 'number' ? entry[1] : 0); + } + } + } + } + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + const stats = await stat(this.deletedLedgerPath).catch((error: unknown) => { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return null; + throw error; + }); + if (stats !== null && stats.size > watermark) { + const handle = await open(this.deletedLedgerPath, 'r'); + try { + const length = stats.size - watermark; + const buffer = Buffer.alloc(length); + await handle.read(buffer, 0, length, watermark); + this.parseDeletedLedgerLines(buffer.toString('utf8'), ledger); + } finally { + await handle.close(); + } + } + return ledger; + } + + private async writeDeletedLedgerSnapshot( + keepIds: Set, + ledger: Map, + ): Promise { + const stats = await stat(this.deletedLedgerPath).catch((error: unknown) => { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return null; + throw error; + }); + const horizon = Date.now() - DELETED_LEDGER_HORIZON_MS; + const entries = [...ledger] + .filter(([id, at]) => keepIds.has(id) || at >= horizon) + .map(([id, at]) => [id, at]); + await writeFile( + this.deletedLedgerSnapshotPath, + JSON.stringify({ watermark: stats?.size ?? 0, entries }), + 'utf8', + ); + } + + private async appendDeletedLedger(sessionId: string): Promise { + await mkdir(dirname(this.deletedLedgerPath), { recursive: true }); + await appendFile(this.deletedLedgerPath, `${sessionId}\t${Date.now()}\n`, 'utf8'); + } + + private async retractDeletedLedger(sessionId: string): Promise { + await appendFile(this.deletedLedgerPath, `-${sessionId}\t${Date.now()}\n`, 'utf8'); + } + ensureOpen(): Promise { this.openPromise ??= this.openDb().then( () => { @@ -438,6 +525,15 @@ export class SearchIndexCore { return { ...outcome, lockToken: this.lockToken, lifecycle: this.lifecycleState() }; } + async deleteSession(sessionId: string): Promise { + if (this.disposed) return; + await this.appendDeletedLedger(sessionId); + await this.ensureOpen(); + const db = this.db; + if (!db || db.readOnly || this.disposed) return; + await this.deleteSessionDocs(db, sessionId); + } + private async runSync(sessions: readonly SyncSessionInput[]): Promise { if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; this.syncReplaced = false; @@ -455,6 +551,40 @@ export class SearchIndexCore { if (!currentIds.has(sessionId)) await this.deleteSessionDocs(db, sessionId); } + let deletedLedger: Map; + try { + deletedLedger = await this.readDeletedLedger(); + } catch (error) { + this.log.warn('global search: failed to read the session deletion ledger; skipping its purge this pass', { + error: errorMessage(error), + }); + deletedLedger = new Map(); + } + const summaryById = new Map(sessions.map((s) => [s.id, s])); + const retractedIds = new Set(); + for (const [sessionId, deletedAt] of deletedLedger) { + if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; + const incarnation = summaryById.get(sessionId); + if (incarnation !== undefined && incarnation.updatedAt > deletedAt) { + await this.retractDeletedLedger(sessionId); + retractedIds.add(sessionId); + continue; + } + await this.deleteSessionDocs(db, sessionId); + } + if (deletedLedger.size > 0) { + try { + const keepIds = new Set( + [...deletedLedger.keys()].filter((id) => currentIds.has(id) && !retractedIds.has(id)), + ); + await this.writeDeletedLedgerSnapshot(keepIds, deletedLedger); + } catch (error) { + this.log.warn('global search: failed to compact the session deletion ledger', { + error: errorMessage(error), + }); + } + } + let indexed = 0; for (const summary of sessions) { if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; @@ -865,7 +995,23 @@ export class SearchIndexCore { const boundary = page.kind === 'keyset' ? page.boundary : undefined; const matched = matchDocs(q, candidates, boundary, budget); incomplete ??= matched.incomplete; - const { pageRows, hasMore } = paginateRows(q, page, matched.rows); + let deletedLedger: Map; + try { + deletedLedger = await this.readDeletedLedger(); + } catch (error) { + throw new GlobalSearchError( + 'index_unavailable', + `failed to read the session deletion ledger: ${errorMessage(error)}`, + ); + } + const visibleRows = + deletedLedger.size === 0 + ? matched.rows + : matched.rows.filter((row) => { + const sep = row.key.indexOf('/'); + return sep <= 0 || !deletedLedger.has(row.key.slice(0, sep)); + }); + const { pageRows, hasMore } = paginateRows(q, page, visibleRows); return { kind: 'page', rows: pageRows, diff --git a/packages/kap-server/src/search/searchService.ts b/packages/kap-server/src/search/searchService.ts index 11dfe134b64..3ea00a4da7b 100644 --- a/packages/kap-server/src/search/searchService.ts +++ b/packages/kap-server/src/search/searchService.ts @@ -101,6 +101,7 @@ export interface IGlobalSearchService { readonly _serviceBrand: undefined; search(query: GlobalSearchQuery): Promise; reindex(): Promise<{ sessions: number; documents: number }>; + deleteSession(sessionId: string): Promise; status(): Promise<{ sessions: number; documents: number; @@ -163,6 +164,7 @@ export interface SearchBackend { refresh(): Promise; reindex(): Promise; status(): Promise; + deleteSession(sessionId: string): Promise; dispose(): Promise; } @@ -209,6 +211,10 @@ export class InlineSearchBackend implements SearchBackend { return this.core.status(); } + deleteSession(sessionId: string): Promise { + return this.core.deleteSession(sessionId); + } + dispose(): Promise { dropLiveLockToken(this.core.lockTokenView); return this.core.close(); @@ -262,6 +268,22 @@ export class GlobalSearchService implements IGlobalSearchService { this.liveSource = source; } + async deleteSession(sessionId: string): Promise { + if (this.disposed) return; + this.summaries.delete(sessionId); + try { + await this.syncPromise?.catch(() => {}); + await this.backend.deleteSession(sessionId); + this.requestSync(); + } catch (error) { + this.log.warn('global search: failed to purge a deleted session from the index', { + sessionId, + error: error instanceof Error ? error.message : String(error), + }); + throw error; + } + } + private get indexDir(): string { return join(this.bootstrap.homeDir, INDEX_DIR_NAME); } diff --git a/packages/kap-server/src/search/worker/entry.ts b/packages/kap-server/src/search/worker/entry.ts index 8be6b16ce28..448e235fd89 100644 --- a/packages/kap-server/src/search/worker/entry.ts +++ b/packages/kap-server/src/search/worker/entry.ts @@ -94,6 +94,11 @@ async function dispatch(request: SearchWorkerCall): Promise { } case 'status': return core.status(); + case 'deleteSession': + await core.deleteSession( + (request.params as { sessionId: string }).sessionId, + ); + return null; case 'close': return null; } diff --git a/packages/kap-server/src/search/worker/host.ts b/packages/kap-server/src/search/worker/host.ts index f8d1288f5d9..0b67149f536 100644 --- a/packages/kap-server/src/search/worker/host.ts +++ b/packages/kap-server/src/search/worker/host.ts @@ -165,6 +165,10 @@ export class SearchWorkerHost { return this.call('status'); } + async deleteSession(sessionId: string): Promise { + await this.call('deleteSession', { sessionId }); + } + async killWorkerForTest(): Promise { const worker = this.worker; if (worker === null) return; diff --git a/packages/kap-server/src/search/worker/protocol.ts b/packages/kap-server/src/search/worker/protocol.ts index 5989697246d..c901be26e14 100644 --- a/packages/kap-server/src/search/worker/protocol.ts +++ b/packages/kap-server/src/search/worker/protocol.ts @@ -23,6 +23,7 @@ export type SearchWorkerCall = | { readonly id: number; readonly v: number; readonly type: 'refresh' } | { readonly id: number; readonly v: number; readonly type: 'reindex' } | { readonly id: number; readonly v: number; readonly type: 'status' } + | { readonly id: number; readonly v: number; readonly type: 'deleteSession'; readonly params: { readonly sessionId: string } } | { readonly id: number; readonly v: number; readonly type: 'close' }; export type SearchWorkerCallType = SearchWorkerCall['type']; @@ -47,6 +48,7 @@ export interface SearchWorkerResultMap { readonly refresh: SearchWorkerOpenResult; readonly reindex: SearchWorkerOpenResult; readonly status: CoreStatus; + readonly deleteSession: null; readonly close: null; } diff --git a/packages/kap-server/src/transport/ws/v1/events.ts b/packages/kap-server/src/transport/ws/v1/events.ts index 9353b80ef79..4d03735a0b6 100644 --- a/packages/kap-server/src/transport/ws/v1/events.ts +++ b/packages/kap-server/src/transport/ws/v1/events.ts @@ -48,6 +48,11 @@ export interface SessionArchivedEvent { readonly workspace_id: string; } +export interface SessionDeletedEvent { + readonly type: 'event.session.deleted'; + readonly workspace_id: string; +} + export interface WorkspaceCreatedEvent { readonly type: 'event.workspace.created'; readonly workspace: Workspace; @@ -218,6 +223,7 @@ export type AgentEvent = | SessionMetaUpdatedEvent | SessionCreatedEvent | SessionArchivedEvent + | SessionDeletedEvent | WorkspaceCreatedEvent | WorkspaceUpdatedEvent | WorkspaceDeletedEvent diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts index 919803d96f9..0444fa09a9b 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -1,3 +1,4 @@ +import { open, rm } from 'node:fs/promises'; import type { AgentActivityState, ApprovalResponse, @@ -138,10 +139,65 @@ export class SessionEventBroadcaster { private readonly globalTargets = new Set(); private readonly diEventTargets = new Set(); private readonly pendingStates = new Map>(); + private readonly journalRemovalAttempts = new Map(); private readonly maxBufferSize: number; private readonly coreEventSubscription: IDisposable; private closed = false; + private async journalEpoch(sessionId: string): Promise { + try { + const handle = await open(sessionJournalPath(this.opts.eventsDir, sessionId), 'r'); + try { + const buffer = Buffer.alloc(4096); + const { bytesRead } = await handle.read(buffer, 0, 4096, 0); + const firstLine = buffer.toString('utf8', 0, bytesRead).split('\n', 1)[0] ?? ''; + if (firstLine.length === 0) return undefined; + const epoch = (JSON.parse(firstLine) as { epoch?: unknown }).epoch; + return typeof epoch === 'string' ? epoch : undefined; + } finally { + await handle.close(); + } + } catch { + return undefined; + } + } + + private async removeSessionJournal(sessionId: string): Promise { + try { + await rm(sessionJournalPath(this.opts.eventsDir, sessionId), { force: true }); + this.journalRemovalAttempts.delete(sessionId); + } catch (error: unknown) { + if (this.closed) return; + const tracked = this.journalRemovalAttempts.get(sessionId); + const epoch = tracked?.epoch ?? (await this.journalEpoch(sessionId)); + if (tracked?.epoch !== undefined) { + const current = await this.journalEpoch(sessionId); + if (current !== tracked.epoch) { + this.journalRemovalAttempts.delete(sessionId); + return; + } + } + const attempt = (tracked?.attempt ?? 0) + 1; + this.journalRemovalAttempts.set(sessionId, { attempt, epoch }); + if (attempt >= 4) { + this.opts.logger?.error?.( + { sessionId, err: String(error) }, + 'session journal could not be removed after repeated attempts', + ); + this.journalRemovalAttempts.delete(sessionId); + return; + } + this.opts.logger?.warn( + { sessionId, attempt, err: String(error) }, + 'session journal removal failed; retrying', + ); + const timer = setTimeout(() => { + void this.removeSessionJournal(sessionId); + }, attempt * 5_000); + timer.unref?.(); + } + } + constructor( private readonly opts: { readonly eventsDir: string; @@ -641,6 +697,37 @@ export class SessionEventBroadcaster { ); return; } + if (event.type === 'event.session.deleted') { + const payload = sessionDeletedPayload(corePayload); + if (payload === undefined) return; + void (async () => { + try { + const pending = this.pendingStates.get(payload.sessionId); + if (pending !== undefined) await pending.catch(() => undefined); + const state = this.sessions.get(payload.sessionId); + if (state !== undefined) { + this.sessions.delete(payload.sessionId); + await disposeSessionState(state); + } + this.opts.transcriptService?.dropSession(payload.sessionId); + await this.removeSessionJournal(payload.sessionId); + } catch (error: unknown) { + this.opts.logger?.warn( + { sessionId: payload.sessionId, err: String(error) }, + 'session deletion cleanup failed; dispatching event.session.deleted anyway', + ); + } + await this.dispatchGlobal({ + type: 'event.session.deleted', + workspace_id: payload.workspaceId, + agentId: 'main', + sessionId: payload.sessionId, + } as Event); + })().catch((error: unknown) => + this.logDispatchError(GLOBAL_SESSION_ID, 'event.session.deleted', error), + ); + return; + } if (event.type === 'event.workspace.created' || event.type === 'event.workspace.updated') { const workspace = workspaceLifecyclePayload(corePayload); if (workspace === undefined) return; @@ -1365,6 +1452,20 @@ function sessionArchivedPayload( return { sessionId: candidate.sessionId, workspaceId: candidate.workspaceId }; } +function sessionDeletedPayload( + payload: unknown, +): { sessionId: string; workspaceId: string } | undefined { + if (typeof payload !== 'object' || payload === null) return undefined; + const candidate = payload as { sessionId?: unknown; workspaceId?: unknown }; + if (typeof candidate.sessionId !== 'string' || candidate.sessionId.length === 0) { + return undefined; + } + if (typeof candidate.workspaceId !== 'string' || candidate.workspaceId.length === 0) { + return undefined; + } + return { sessionId: candidate.sessionId, workspaceId: candidate.workspaceId }; +} + function workspaceLifecyclePayload(payload: unknown): Workspace | undefined { if (typeof payload !== 'object' || payload === null) return undefined; const candidate = (payload as { workspace?: unknown }).workspace; diff --git a/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap b/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap index a2a5f4b78c6..7a90a626599 100644 --- a/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap +++ b/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap @@ -408,6 +408,10 @@ exports[`API surface snapshot > matches the documented v2 route table and meta e "POST", "/api/v1/sessions/{session_id}:archive", ], + [ + "POST", + "/api/v1/sessions/{session_id}:delete", + ], [ "POST", "/api/v1/sessions/{session_id}/{tail}", diff --git a/packages/kap-server/test/openapi.test.ts b/packages/kap-server/test/openapi.test.ts index 20c9a13388e..b32ad7887e2 100644 --- a/packages/kap-server/test/openapi.test.ts +++ b/packages/kap-server/test/openapi.test.ts @@ -60,12 +60,13 @@ describe('server-v2 OpenAPI', () => { expect(paths['/api/v1/sessions/{session_id}/fs/{*}']).toBeDefined(); }); - it('projects the session-action dispatcher into archive only', async () => { + it('projects the session-action dispatcher into archive and delete only', async () => { const doc = await fetchOpenApi(); const paths = asRecord(doc['paths']); expect(paths['/api/v1/sessions/{tail}']).toBeUndefined(); expect(paths['/api/v1/sessions/{session_id}:archive']).toBeDefined(); + expect(paths['/api/v1/sessions/{session_id}:delete']).toBeDefined(); expect(paths['/api/v1/sessions/{session_id}:fork']).toBeUndefined(); expect(paths['/api/v1/sessions/{session_id}:undo']).toBeUndefined(); @@ -74,6 +75,12 @@ describe('server-v2 OpenAPI', () => { const params = archiveOp['parameters'] as Array>; expect(params.some((p) => p['in'] === 'path' && p['name'] === 'session_id')).toBe(true); expect(params.some((p) => p['name'] === 'tail')).toBe(false); + + const deleteOp = operation(doc, '/api/v1/sessions/{session_id}:delete', 'post'); + expect(deleteOp['operationId']).toBe('runSessionDeleteAction'); + const deleteParams = deleteOp['parameters'] as Array>; + expect(deleteParams.some((p) => p['in'] === 'path' && p['name'] === 'session_id')).toBe(true); + expect(deleteParams.some((p) => p['name'] === 'tail')).toBe(false); }); it('describes the file upload as multipart/form-data', async () => { diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index dc9eb305ed6..89149323a55 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -298,6 +298,77 @@ describe('GlobalSearchService', () => { expect(injected.items).toEqual([]); }); + it('purges a deleted session from the index immediately via deleteSession', async () => { + const s1 = summary('s1', '删除测试', T1); + await writeWire(home!, 's1', 'main', [userLine('即将被删除的苹果', T1)]); + const service = track(makeService(home!, staticIndex([s1]))); + await service.reindex(); + expect((await service.search({ query: '苹果' })).items.length).toBeGreaterThan(0); + + await service.deleteSession('s1'); + + expect((await service.search({ query: '苹果' })).items).toEqual([]); + }); + + it('keeps a deleted session out of search results even if a stale sync re-indexes it', async () => { + const s1 = summary('s1', '删除测试', T1); + await writeWire(home!, 's1', 'main', [userLine('即将被删除的苹果', T1)]); + const service = track(makeService(home!, staticIndex([s1]))); + await service.reindex(); + expect((await service.search({ query: '苹果' })).items.length).toBeGreaterThan(0); + + await service.deleteSession('s1'); + await settleSync(service); + + expect((await service.search({ query: '苹果' })).items).toEqual([]); + }); + + it('records concurrent deletions without losing tombstones', async () => { + const s1 = summary('s1', '删除测试一', T1); + const s2 = summary('s2', '删除测试二', T1); + await writeWire(home!, 's1', 'main', [userLine('第一个苹果', T1)]); + await writeWire(home!, 's2', 'main', [userLine('第二个苹果', T1)]); + const service = track(makeService(home!, staticIndex([s1, s2]))); + await service.reindex(); + expect((await service.search({ query: '苹果' })).items.length).toBe(2); + + await Promise.all([service.deleteSession('s1'), service.deleteSession('s2')]); + + expect((await service.search({ query: '苹果' })).items).toEqual([]); + }); + + it('fails the search instead of serving deleted content when the ledger is unreadable', async () => { + const s1 = summary('s1', '删除测试', T1); + await writeWire(home!, 's1', 'main', [userLine('即将被删除的苹果', T1)]); + const service = track(makeService(home!, staticIndex([s1]))); + await service.reindex(); + await service.deleteSession('s1'); + await rm(join(home!, 'search-index', 'deleted-sessions.jsonl'), { force: true }); + await mkdir(join(home!, 'search-index', 'deleted-sessions.jsonl')); + + await expect(service.search({ query: '苹果' })).rejects.toThrow(); + }); + + it('lifts the tombstone when the session id is recreated with a newer incarnation', async () => { + const s1 = summary('s1', '旧会话', T1); + await writeWire(home!, 's1', 'main', [userLine('旧苹果', T1)]); + const first = track(makeService(home!, staticIndex([s1]))); + await first.reindex(); + await first.dispose(); + + await appendFile(join(home!, 'search-index', 'deleted-sessions.jsonl'), 's1\t1000\n', 'utf8'); + const s1v2 = summary('s1', '新会话', T2); + await writeWire(home!, 's1', 'main', [userLine('新香蕉', T2)]); + const second = track(makeService(home!, staticIndex([s1v2]))); + await settleSync(second); + + expect((await second.search({ query: '香蕉' })).items.length).toBe(1); + const snapshot = JSON.parse( + await readFile(join(home!, 'search-index', 'deleted-sessions.snapshot.json'), 'utf8'), + ) as { entries: Array<[string, number]> }; + expect(snapshot.entries.some(([id]) => id === 's1')).toBe(false); + }); + it('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); diff --git a/packages/kap-server/test/sessionEventBroadcaster.test.ts b/packages/kap-server/test/sessionEventBroadcaster.test.ts index d86dc69447f..d2915568795 100644 --- a/packages/kap-server/test/sessionEventBroadcaster.test.ts +++ b/packages/kap-server/test/sessionEventBroadcaster.test.ts @@ -1,4 +1,4 @@ -import { mkdtemp, rm } from 'node:fs/promises'; +import { access, mkdir, mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; @@ -1260,6 +1260,77 @@ describe('SessionEventBroadcaster', () => { expect(globalView.deliveries).toEqual(['immediate']); }); + it('fans out event.session.deleted to every connection, including for cold sessions', async () => { + const globalView = collectingTarget(); + bc.addGlobalTarget(globalView.target); + + eventBus.emit({ + type: 'event.session.deleted', + payload: { sessionId: 'cold-1', workspaceId: 'wd_cold' }, + }); + + await vi.waitFor(() => expect(globalView.envelopes).toHaveLength(1)); + expect(globalView.envelopes[0]).toMatchObject({ + type: 'event.session.deleted', + session_id: '__global__', + payload: { + type: 'event.session.deleted', + agentId: 'main', + sessionId: 'cold-1', + workspace_id: 'wd_cold', + }, + }); + expect(globalView.deliveries).toEqual(['immediate']); + }); + + it('purges cached state and the journal when a materialized session is deleted', async () => { + const lc = new FakeLifecycle(); + const main = lc.addAgent('main'); + sessions.set('s1', lc); + const view = collectingTarget(); + expect(await bc.subscribe('s1', view.target)).toBe(true); + + main.bus.emit(agentEvent('turn.started', { turnId: 1 })); + await bc.getCursor('s1'); + await vi.waitFor(async () => { + await access(join(dir, 's1.jsonl')); + }); + + sessions.delete('s1'); + eventBus.emit({ + type: 'event.session.deleted', + payload: { sessionId: 's1', workspaceId: 'wd_1' }, + }); + + await vi.waitFor(async () => { + await expect(access(join(dir, 's1.jsonl'))).rejects.toThrow(); + }); + expect(await bc.subscribe('s1', collectingTarget().target)).toBe(false); + }); + + it('still fans out event.session.deleted when the journal cleanup fails', async () => { + const globalView = collectingTarget(); + bc.addGlobalTarget(globalView.target); + await mkdir(join(dir, 'cold-2.jsonl')); + + eventBus.emit({ + type: 'event.session.deleted', + payload: { sessionId: 'cold-2', workspaceId: 'wd_cold' }, + }); + + await vi.waitFor(() => expect(globalView.envelopes).toHaveLength(1)); + expect(globalView.envelopes[0]).toMatchObject({ + type: 'event.session.deleted', + session_id: '__global__', + payload: { + type: 'event.session.deleted', + agentId: 'main', + sessionId: 'cold-2', + workspace_id: 'wd_cold', + }, + }); + }); + it('fans out event.workspace.created/updated with the wire workspace shape', async () => { const globalView = collectingTarget(); bc.addGlobalTarget(globalView.target); diff --git a/packages/kap-server/test/sessions.test.ts b/packages/kap-server/test/sessions.test.ts index 405615f7398..b421c0e58ec 100644 --- a/packages/kap-server/test/sessions.test.ts +++ b/packages/kap-server/test/sessions.test.ts @@ -905,6 +905,66 @@ describe('server-v2 /api/v1/sessions', () => { expect(body.code).toBe(40401); }); + it('deletes a session via :delete and publishes event.session.deleted', async () => { + const cwd = home as string; + const created = await postJson('/api/v1/sessions', { metadata: { cwd } }); + const id = created.body.data.id; + const workspaceId = created.body.data.workspace_id; + + const events: Event2[] = []; + const sub = (server as RunningServer).core.accessor + .get(IEventService) + .subscribe((event) => events.push(event)); + try { + const deleted = await postJson<{ deleted: boolean }>(`/api/v1/sessions/${id}:delete`); + expect(deleted.body.code).toBe(0); + expect(deleted.body.data).toEqual({ deleted: true }); + + const got = await getJson(`/api/v1/sessions/${id}`); + expect(got.body.code).toBe(40401); + + expect( + events + .filter((event) => event.type === 'event.session.deleted') + .map((event) => (event as { readonly payload?: unknown }).payload), + ).toEqual([{ sessionId: id, workspaceId }]); + } finally { + sub.dispose(); + } + }); + + it('deletes a cold session via :delete and publishes event.session.deleted', async () => { + const cwd = home as string; + const created = await postJson('/api/v1/sessions', { metadata: { cwd } }); + const id = created.body.data.id; + const workspaceId = created.body.data.workspace_id; + await closeSessionById((server as RunningServer).core.accessor, id); + expect(getLiveSessionById((server as RunningServer).core.accessor, id)).toBeUndefined(); + + const events: Event2[] = []; + const sub = (server as RunningServer).core.accessor + .get(IEventService) + .subscribe((event) => events.push(event)); + try { + const deleted = await postJson<{ deleted: boolean }>(`/api/v1/sessions/${id}:delete`); + expect(deleted.body.code).toBe(0); + expect(deleted.body.data).toEqual({ deleted: true }); + + expect( + events + .filter((event) => event.type === 'event.session.deleted') + .map((event) => (event as { readonly payload?: unknown }).payload), + ).toEqual([{ sessionId: id, workspaceId }]); + } finally { + sub.dispose(); + } + }); + + it('returns 40401 when deleting a missing session', async () => { + const { body } = await postJson('/api/v1/sessions/sess_missing:delete'); + expect(body.code).toBe(40401); + }); + it('cold-loads a persisted session on :undo instead of 40401', async () => { const cwd = home as string; const created = await postJson('/api/v1/sessions', { metadata: { cwd } });