From 4d60498642e8d1ea25a83a0395190e74a56fe351 Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 07:00:22 +0000 Subject: [PATCH 01/10] feat(kap-server): expose permanent session deletion via :delete Add a delete session action to the v1 session-action route that permanently removes a session through ISessionManager.delete, returning { deleted: true } and 40401 for unknown sessions. Publish a new event.session.deleted Event2 from the session lifecycle controller and fan it out to global WS subscribers with the same wire shape as event.session.archived. Project the :delete operation into the OpenAPI document alongside :archive. Co-authored-by: qer --- .changeset/session-delete-action.md | 5 ++ .../sessionLifecycleEvents.ts | 12 ++++ .../sessionLifecycleService.ts | 7 ++- packages/kap-server/src/openapi/transforms.ts | 33 +++++++--- .../kap-server/src/protocol/events-zod.ts | 6 ++ .../kap-server/src/protocol/rest-session.ts | 6 +- packages/kap-server/src/routes/sessions.ts | 20 ++++++- .../kap-server/src/transport/ws/v1/events.ts | 6 ++ .../ws/v1/sessionEventBroadcaster.ts | 27 +++++++++ .../apiSurface.snapshot.test.ts.snap | 4 ++ packages/kap-server/test/openapi.test.ts | 9 ++- .../test/sessionEventBroadcaster.test.ts | 23 +++++++ packages/kap-server/test/sessions.test.ts | 60 +++++++++++++++++++ 13 files changed, 205 insertions(+), 13 deletions(-) create mode 100644 .changeset/session-delete-action.md diff --git a/.changeset/session-delete-action.md b/.changeset/session-delete-action.md new file mode 100644 index 00000000000..99c4c6bab55 --- /dev/null +++ b/.changeset/session-delete-action.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Expose permanent session deletion via `POST /api/v1/sessions/{session_id}:delete` and broadcast `event.session.deleted` over WebSocket. 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..1a2e0798c7c 100644 --- a/packages/kap-server/src/routes/sessions.ts +++ b/packages/kap-server/src/routes/sessions.ts @@ -41,6 +41,7 @@ import { compactSessionResponseSchema, createSessionChildRequestSchema, createSessionRequestSchema, + deleteSessionResponseSchema, forkSessionRequestSchema, getSessionGoalResponseSchema, listSessionChildrenResponseSchema, @@ -594,6 +595,7 @@ export function registerSessionsRoutes(app: SessionRouteHost, core: Scope): void sessionAbortResponseSchema, startBtwSessionResponseSchema, archiveSessionResponseSchema, + deleteSessionResponseSchema, ]), }, errors: { @@ -844,7 +846,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 +875,7 @@ const sessionActions: ActionTable = { btw: { handle: btwSessionAction }, restore: { handle: restoreSessionAction }, archive: { handle: archiveSessionAction }, + delete: { handle: deleteSessionAction }, }; async function forkSessionAction( @@ -982,6 +993,13 @@ 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; + await core.accessor.get(ISessionManager).delete(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/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..54568d447a3 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -641,6 +641,19 @@ export class SessionEventBroadcaster { ); return; } + if (event.type === 'event.session.deleted') { + const payload = sessionDeletedPayload(corePayload); + if (payload === undefined) return; + void 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 +1378,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/sessionEventBroadcaster.test.ts b/packages/kap-server/test/sessionEventBroadcaster.test.ts index d86dc69447f..6d1271336a6 100644 --- a/packages/kap-server/test/sessionEventBroadcaster.test.ts +++ b/packages/kap-server/test/sessionEventBroadcaster.test.ts @@ -1260,6 +1260,29 @@ 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('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 } }); From 04e8ae1e459ba9893a0dd006ff7ff8fdb265aaea Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 09:48:00 +0000 Subject: [PATCH 02/10] fix(kap-server): purge broadcaster state and journal on session deletion The broadcaster kept a deleted session's SessionState cached, so a later subscription or cursor request for that id could replay the retained event journal after a supposedly permanent deletion. The deletion branch now awaits any in-flight state creation, evicts and disposes the cached state (closing its journal), drops the transcript live store, and removes the journal file before fanning out event.session.deleted. Co-authored-by: qer --- .../ws/v1/sessionEventBroadcaster.ts | 24 ++++++++++++----- .../test/sessionEventBroadcaster.test.ts | 27 ++++++++++++++++++- 2 files changed, 44 insertions(+), 7 deletions(-) diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts index 54568d447a3..616693151f3 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 { rm } from 'node:fs/promises'; import type { AgentActivityState, ApprovalResponse, @@ -644,12 +645,23 @@ export class SessionEventBroadcaster { if (event.type === 'event.session.deleted') { const payload = sessionDeletedPayload(corePayload); if (payload === undefined) return; - void this.dispatchGlobal({ - type: 'event.session.deleted', - workspace_id: payload.workspaceId, - agentId: 'main', - sessionId: payload.sessionId, - } as Event).catch((error: unknown) => + void (async () => { + 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 rm(sessionJournalPath(this.opts.eventsDir, payload.sessionId), { force: true }); + 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; diff --git a/packages/kap-server/test/sessionEventBroadcaster.test.ts b/packages/kap-server/test/sessionEventBroadcaster.test.ts index 6d1271336a6..90340755f2c 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, mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; @@ -1283,6 +1283,31 @@ describe('SessionEventBroadcaster', () => { 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('fans out event.workspace.created/updated with the wire workspace shape', async () => { const globalView = collectingTarget(); bc.addGlobalTarget(globalView.target); From f35f8bb2c5873ce17055d33a8ffc11697647787e Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 10:25:41 +0000 Subject: [PATCH 03/10] fix(kap-server): fan out session deletion even when journal cleanup fails A failing journal rm (EACCES, EBUSY, an unexpected directory at the journal path) rejected the deletion branch before dispatchGlobal ran, so connected clients never learned that the session was permanently deleted even though the core deletion and REST response had already succeeded. The purge now logs cleanup failures and the event is dispatched regardless. Co-authored-by: qer --- .../ws/v1/sessionEventBroadcaster.ts | 23 +++++++++++------ .../test/sessionEventBroadcaster.test.ts | 25 ++++++++++++++++++- 2 files changed, 39 insertions(+), 9 deletions(-) diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts index 616693151f3..d0770ab6ecc 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -646,15 +646,22 @@ export class SessionEventBroadcaster { const payload = sessionDeletedPayload(corePayload); if (payload === undefined) return; void (async () => { - 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); + 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 rm(sessionJournalPath(this.opts.eventsDir, payload.sessionId), { force: true }); + } catch (error: unknown) { + this.opts.logger?.warn( + { sessionId: payload.sessionId, err: String(error) }, + 'session deletion cleanup failed; dispatching event.session.deleted anyway', + ); } - this.opts.transcriptService?.dropSession(payload.sessionId); - await rm(sessionJournalPath(this.opts.eventsDir, payload.sessionId), { force: true }); await this.dispatchGlobal({ type: 'event.session.deleted', workspace_id: payload.workspaceId, diff --git a/packages/kap-server/test/sessionEventBroadcaster.test.ts b/packages/kap-server/test/sessionEventBroadcaster.test.ts index 90340755f2c..d2915568795 100644 --- a/packages/kap-server/test/sessionEventBroadcaster.test.ts +++ b/packages/kap-server/test/sessionEventBroadcaster.test.ts @@ -1,4 +1,4 @@ -import { access, 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'; @@ -1308,6 +1308,29 @@ describe('SessionEventBroadcaster', () => { 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); From a3a10a8c83f0751f36521a0da590d907ec63d5c6 Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 11:24:50 +0000 Subject: [PATCH 04/10] feat(kimi-inspect): handle event.session.deleted in the activity feed The global-events WebSocket consumer handled event.session.archived but ignored the new deletion event, so a session deleted by another client stayed visible in Kimi Inspect until the next polling tick. The switch now maps the payload session id to a new onSessionDeleted handler that removes the session from the activity store and invalidates the list, mirroring the archived path. Co-authored-by: qer --- apps/kimi-inspect/src/activity/store.test.ts | 17 ++++++++++++++++- apps/kimi-inspect/src/activity/store.ts | 4 ++++ apps/kimi-inspect/src/activity/ws.ts | 12 ++++++++++++ 3 files changed, 32 insertions(+), 1 deletion(-) 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..9138b88fbac 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) | undefined; /** 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': { From de53f98a4ecfa5492287fa8912d47ddb3e10dc3f Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 12:03:35 +0000 Subject: [PATCH 05/10] feat(kap-server): purge deleted sessions from the global search index Deleting a session removed the core session and its session-index entry but left its prompt and assistant text in the search MiniDb until a later writer sync happened to re-scan, so a search right after a permanent deletion could still return the deleted content. The delete route now purges the session's documents immediately through a new IGlobalSearchService.deleteSession, wired through the inline backend and the search worker protocol, with failures logged rather than failing the delete. Also classify the changeset as minor (new user-visible capability) and drop a redundant undefined union from the kimi-inspect callback type. Co-authored-by: qer --- .changeset/session-delete-action.md | 2 +- apps/kimi-inspect/src/activity/ws.ts | 2 +- packages/kap-server/src/routes/sessions.ts | 2 ++ packages/kap-server/src/search/indexCore.ts | 8 ++++++++ .../kap-server/src/search/searchService.ts | 19 +++++++++++++++++++ .../kap-server/src/search/worker/entry.ts | 5 +++++ packages/kap-server/src/search/worker/host.ts | 4 ++++ .../kap-server/src/search/worker/protocol.ts | 2 ++ .../test/search/searchService.test.ts | 12 ++++++++++++ 9 files changed, 54 insertions(+), 2 deletions(-) diff --git a/.changeset/session-delete-action.md b/.changeset/session-delete-action.md index 99c4c6bab55..8b69403b4f4 100644 --- a/.changeset/session-delete-action.md +++ b/.changeset/session-delete-action.md @@ -1,5 +1,5 @@ --- -"@moonshot-ai/kimi-code": patch +"@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/ws.ts b/apps/kimi-inspect/src/activity/ws.ts index 9138b88fbac..6a1887f05c4 100644 --- a/apps/kimi-inspect/src/activity/ws.ts +++ b/apps/kimi-inspect/src/activity/ws.ts @@ -67,7 +67,7 @@ export interface GlobalEventsWsHandlers { /** 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) | undefined; + 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). */ diff --git a/packages/kap-server/src/routes/sessions.ts b/packages/kap-server/src/routes/sessions.ts index 1a2e0798c7c..3ae222461fb 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'; @@ -996,6 +997,7 @@ async function archiveSessionAction(ctx: SessionActionCtx): Promise { async function deleteSessionAction(ctx: SessionActionCtx): Promise { const { core, req, reply, id } = ctx; await core.accessor.get(ISessionManager).delete(id); + 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)); } diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index 70750caeebf..802b55511bd 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -438,6 +438,14 @@ export class SearchIndexCore { return { ...outcome, lockToken: this.lockToken, lifecycle: this.lifecycleState() }; } + async deleteSession(sessionId: string): Promise { + if (this.disposed) return; + 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; diff --git a/packages/kap-server/src/search/searchService.ts b/packages/kap-server/src/search/searchService.ts index 11dfe134b64..819f33c0d91 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,19 @@ export class GlobalSearchService implements IGlobalSearchService { this.liveSource = source; } + async deleteSession(sessionId: string): Promise { + if (this.disposed) return; + this.summaries.delete(sessionId); + try { + await this.backend.deleteSession(sessionId); + } 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), + }); + } + } + 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/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index dc9eb305ed6..5ca3e824d8c 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -298,6 +298,18 @@ 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('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); From b628094846aa947465144ec58804f3a79ec65c4f Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 12:45:27 +0000 Subject: [PATCH 06/10] fix(kap-server): make the search-index purge durable with a deletion ledger MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The immediate per-document purge left three holes: it was not serialized against an in-flight sync that had snapshotted the session (the stale sync could recreate documents after the purge), it silently no-oped when another process held the MiniDb write lock, and a failed purge was logged and forgotten with no retry. Deletion now appends the session id to a bounded ledger file in the index directory first — durable even when the database itself is read-only — and every sync pass in any process purges the ledger's sessions idempotently, while search results filter ledger members outright so deleted content is never served in the meantime. The service also serializes the purge behind the in-flight sync and schedules a fresh pass, so the best-effort immediate deletion and the guaranteed ledger cleanup reinforce each other. Co-authored-by: qer --- packages/kap-server/src/search/indexCore.ts | 44 +++++++++++++++++-- .../kap-server/src/search/searchService.ts | 2 + .../test/search/searchService.test.ts | 13 ++++++ 3 files changed, 56 insertions(+), 3 deletions(-) diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index 802b55511bd..9f8342e9db1 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 { mkdir, open, readFile, readdir, rm, stat, writeFile } from 'node:fs/promises'; +import { dirname, join, relative } from 'node:path'; import { LockError, @@ -239,6 +239,29 @@ export class SearchIndexCore { return this.options.log; } + private get deletedLedgerPath(): string { + return join(this.indexDir, 'deleted-sessions.json'); + } + + private async readDeletedLedger(): Promise> { + try { + const parsed: unknown = JSON.parse(await readFile(this.deletedLedgerPath, 'utf8')); + if (!Array.isArray(parsed)) return new Set(); + return new Set(parsed.filter((entry): entry is string => typeof entry === 'string')); + } catch { + return new Set(); + } + } + + private async appendDeletedLedger(sessionId: string): Promise { + const ledger = await this.readDeletedLedger(); + if (ledger.has(sessionId)) return; + ledger.add(sessionId); + const capped = [...ledger].slice(-500); + await mkdir(dirname(this.deletedLedgerPath), { recursive: true }); + await writeFile(this.deletedLedgerPath, JSON.stringify(capped), 'utf8'); + } + ensureOpen(): Promise { this.openPromise ??= this.openDb().then( () => { @@ -440,6 +463,7 @@ export class SearchIndexCore { 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; @@ -463,6 +487,12 @@ export class SearchIndexCore { if (!currentIds.has(sessionId)) await this.deleteSessionDocs(db, sessionId); } + const deletedLedger = await this.readDeletedLedger(); + for (const sessionId of deletedLedger) { + if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; + await this.deleteSessionDocs(db, sessionId); + } + let indexed = 0; for (const summary of sessions) { if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; @@ -873,7 +903,15 @@ 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); + const deletedLedger = await this.readDeletedLedger(); + 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 819f33c0d91..88d222fe6bd 100644 --- a/packages/kap-server/src/search/searchService.ts +++ b/packages/kap-server/src/search/searchService.ts @@ -272,7 +272,9 @@ export class GlobalSearchService implements IGlobalSearchService { 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, diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index 5ca3e824d8c..fe93f74ef14 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -310,6 +310,19 @@ describe('GlobalSearchService', () => { 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('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); From 64310f8aff352430256227c6e1abb84de8029211 Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 13:10:30 +0000 Subject: [PATCH 07/10] fix(kap-server): make the deletion ledger append-only and propagate its failure The read-modify-write ledger raced with itself: two concurrent deletions (or two read-only server processes) could read the same snapshot and overwrite each other's tombstone, and the 500-entry cap could evict intents the writer had not consumed yet. The ledger is now an append-only JSONL file, so concurrent records never clobber each other and no entry is ever evicted; every sync pass keeps purging the recorded sessions idempotently. When even the ledger append fails (disk full, unwritable directory), the error now propagates through the service instead of being logged and swallowed, so the delete route no longer reports success while nothing durable records the purge intent. Co-authored-by: qer --- packages/kap-server/src/search/indexCore.ts | 15 +++++---------- packages/kap-server/src/search/searchService.ts | 1 + .../kap-server/test/search/searchService.test.ts | 14 ++++++++++++++ 3 files changed, 20 insertions(+), 10 deletions(-) diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index 9f8342e9db1..b085d15db72 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -1,5 +1,5 @@ import { createHash } from 'node:crypto'; -import { mkdir, open, readFile, readdir, rm, stat, writeFile } from 'node:fs/promises'; +import { appendFile, mkdir, open, readFile, readdir, rm, stat } from 'node:fs/promises'; import { dirname, join, relative } from 'node:path'; import { @@ -240,26 +240,21 @@ export class SearchIndexCore { } private get deletedLedgerPath(): string { - return join(this.indexDir, 'deleted-sessions.json'); + return join(this.indexDir, 'deleted-sessions.jsonl'); } private async readDeletedLedger(): Promise> { try { - const parsed: unknown = JSON.parse(await readFile(this.deletedLedgerPath, 'utf8')); - if (!Array.isArray(parsed)) return new Set(); - return new Set(parsed.filter((entry): entry is string => typeof entry === 'string')); + const raw = await readFile(this.deletedLedgerPath, 'utf8'); + return new Set(raw.split('\n').filter((line) => line.length > 0)); } catch { return new Set(); } } private async appendDeletedLedger(sessionId: string): Promise { - const ledger = await this.readDeletedLedger(); - if (ledger.has(sessionId)) return; - ledger.add(sessionId); - const capped = [...ledger].slice(-500); await mkdir(dirname(this.deletedLedgerPath), { recursive: true }); - await writeFile(this.deletedLedgerPath, JSON.stringify(capped), 'utf8'); + await appendFile(this.deletedLedgerPath, `${sessionId}\n`, 'utf8'); } ensureOpen(): Promise { diff --git a/packages/kap-server/src/search/searchService.ts b/packages/kap-server/src/search/searchService.ts index 88d222fe6bd..3ea00a4da7b 100644 --- a/packages/kap-server/src/search/searchService.ts +++ b/packages/kap-server/src/search/searchService.ts @@ -280,6 +280,7 @@ export class GlobalSearchService implements IGlobalSearchService { sessionId, error: error instanceof Error ? error.message : String(error), }); + throw error; } } diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index fe93f74ef14..baa9e7428dc 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -323,6 +323,20 @@ describe('GlobalSearchService', () => { 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('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); From 84375a8875baabaa9908052ab41f63d4d7ec02bc Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 13:38:46 +0000 Subject: [PATCH 08/10] fix(kap-server): harden ledger reads and keep failed cleanups retriable A transient ledger read failure (EACCES, EMFILE, I/O) was folded into an empty ledger, silently disabling the query-time filter and letting permanently deleted text surface; only ENOENT means empty now, search fails closed with index_unavailable when the ledger cannot be read, and the sync pass just skips its consumption round on a read error. And because the route now purges even when the session is already absent (still returning SESSION_NOT_FOUND to the caller), a deletion whose ledger write failed after the core session was removed stays retriable: retrying reaches the cleanup instead of dead-ending on the first line. Co-authored-by: qer --- packages/kap-server/src/routes/sessions.ts | 8 +++++- packages/kap-server/src/search/indexCore.ts | 25 ++++++++++++++++--- .../test/search/searchService.test.ts | 12 +++++++++ 3 files changed, 40 insertions(+), 5 deletions(-) diff --git a/packages/kap-server/src/routes/sessions.ts b/packages/kap-server/src/routes/sessions.ts index 3ae222461fb..a0eff3061b9 100644 --- a/packages/kap-server/src/routes/sessions.ts +++ b/packages/kap-server/src/routes/sessions.ts @@ -996,7 +996,13 @@ async function archiveSessionAction(ctx: SessionActionCtx): Promise { async function deleteSessionAction(ctx: SessionActionCtx): Promise { const { core, req, reply, id } = ctx; - await core.accessor.get(ISessionManager).delete(id); + 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)); diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index b085d15db72..d748328cbb2 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -247,8 +247,9 @@ export class SearchIndexCore { try { const raw = await readFile(this.deletedLedgerPath, 'utf8'); return new Set(raw.split('\n').filter((line) => line.length > 0)); - } catch { - return new Set(); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return new Set(); + throw error; } } @@ -482,7 +483,15 @@ export class SearchIndexCore { if (!currentIds.has(sessionId)) await this.deleteSessionDocs(db, sessionId); } - const deletedLedger = await this.readDeletedLedger(); + let deletedLedger: Set; + 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 Set(); + } for (const sessionId of deletedLedger) { if (this.disposed) return { noop: true, sessions: 0, documents: 0 }; await this.deleteSessionDocs(db, sessionId); @@ -898,7 +907,15 @@ export class SearchIndexCore { const boundary = page.kind === 'keyset' ? page.boundary : undefined; const matched = matchDocs(q, candidates, boundary, budget); incomplete ??= matched.incomplete; - const deletedLedger = await this.readDeletedLedger(); + let deletedLedger: Set; + 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 diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index baa9e7428dc..45ffed2d969 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -337,6 +337,18 @@ describe('GlobalSearchService', () => { 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('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); From 3c24522b88650f443bbdbbe088dbdd852465b338 Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 14:23:37 +0000 Subject: [PATCH 09/10] fix(kap-server): make journal removal retriable and the ledger incarnation-aware MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The broadcaster's journal rm was caught and only logged, so a supposedly permanent deletion could leave the per-session event JSONL on disk forever; removal now retries with bounded backoff and escalates to an error log when it keeps failing. Ledger lines now carry the deletion timestamp, and a sync pass that meets a newer incarnation of the same id retracts the tombstone instead of hiding the recreated session forever; passes also compact the ledger — retaining only entries appended mid-pass and tombstones whose session is still in the authoritative list (so a stale snapshot re-indexing it stays filtered) — keeping per-query ledger reads and per-pass purge work bounded over the installation's lifetime. Co-authored-by: qer --- packages/kap-server/src/search/indexCore.ts | 69 ++++++++++++++++--- .../ws/v1/sessionEventBroadcaster.ts | 30 +++++++- .../test/search/searchService.test.ts | 18 +++++ 3 files changed, 106 insertions(+), 11 deletions(-) diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index d748328cbb2..8df4fe108e2 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -1,5 +1,5 @@ import { createHash } from 'node:crypto'; -import { appendFile, mkdir, open, readFile, readdir, rm, stat } from 'node:fs/promises'; +import { appendFile, mkdir, open, readFile, readdir, rm, stat, writeFile } from 'node:fs/promises'; import { dirname, join, relative } from 'node:path'; import { @@ -243,19 +243,37 @@ export class SearchIndexCore { return join(this.indexDir, 'deleted-sessions.jsonl'); } - private async readDeletedLedger(): Promise> { + private async readDeletedLedgerRaw(): Promise { try { - const raw = await readFile(this.deletedLedgerPath, 'utf8'); - return new Set(raw.split('\n').filter((line) => line.length > 0)); + return await readFile(this.deletedLedgerPath, 'utf8'); } catch (error) { - if ((error as NodeJS.ErrnoException).code === 'ENOENT') return new Set(); + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return ''; throw error; } } + private async readDeletedLedger(): Promise> { + const raw = await this.readDeletedLedgerRaw(); + const ledger = new Map(); + for (const line of raw.split('\n')) { + if (line.length === 0) continue; + if (line.startsWith('-')) { + ledger.delete(line.slice(1).split('\t')[0]!); + continue; + } + const [id, at] = line.split('\t'); + if (id !== undefined && id.length > 0) ledger.set(id, Number(at ?? 0)); + } + return ledger; + } + private async appendDeletedLedger(sessionId: string): Promise { await mkdir(dirname(this.deletedLedgerPath), { recursive: true }); - await appendFile(this.deletedLedgerPath, `${sessionId}\n`, 'utf8'); + 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 { @@ -483,19 +501,50 @@ export class SearchIndexCore { if (!currentIds.has(sessionId)) await this.deleteSessionDocs(db, sessionId); } - let deletedLedger: Set; + let deletedLedger: Map; + let deletedLedgerRawBefore = ''; try { + deletedLedgerRawBefore = await this.readDeletedLedgerRaw(); 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 Set(); + deletedLedger = new Map(); } - for (const sessionId of deletedLedger) { + 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 initialLines = new Set(deletedLedgerRawBefore.split('\n').filter((l) => l.length > 0)); + const appended = (await this.readDeletedLedgerRaw()) + .split('\n') + .filter((line) => line.length > 0 && !initialLines.has(line)); + const retained = [...deletedLedger] + .filter(([id]) => currentIds.has(id) && !retractedIds.has(id)) + .map(([id, at]) => `${id}\t${at}`); + const kept = [...retained, ...appended]; + await writeFile( + this.deletedLedgerPath, + kept.length > 0 ? `${kept.join('\n')}\n` : '', + 'utf8', + ); + } 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) { @@ -907,7 +956,7 @@ export class SearchIndexCore { const boundary = page.kind === 'keyset' ? page.boundary : undefined; const matched = matchDocs(q, candidates, boundary, budget); incomplete ??= matched.incomplete; - let deletedLedger: Set; + let deletedLedger: Map; try { deletedLedger = await this.readDeletedLedger(); } catch (error) { diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts index d0770ab6ecc..bed85b35669 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -139,10 +139,38 @@ 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 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 attempt = (this.journalRemovalAttempts.get(sessionId) ?? 0) + 1; + this.journalRemovalAttempts.set(sessionId, attempt); + 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; @@ -655,7 +683,7 @@ export class SessionEventBroadcaster { await disposeSessionState(state); } this.opts.transcriptService?.dropSession(payload.sessionId); - await rm(sessionJournalPath(this.opts.eventsDir, payload.sessionId), { force: true }); + await this.removeSessionJournal(payload.sessionId); } catch (error: unknown) { this.opts.logger?.warn( { sessionId: payload.sessionId, err: String(error) }, diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index 45ffed2d969..1c0008b08a6 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -349,6 +349,24 @@ describe('GlobalSearchService', () => { 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 ledgerRaw = await readFile(join(home!, 'search-index', 'deleted-sessions.jsonl'), 'utf8'); + expect(ledgerRaw).not.toContain('s1\t1000'); + }); + it('hits session titles as title docs', async () => { const s1 = summary('s1', '季度总结报告', T1); await writeWire(home!, 's1', 'main', [userLine('随便说点什么', T1)]); From a04254558a36c0700ddc03945c6678b6210e3951 Mon Sep 17 00:00:00 2001 From: kimi-agent-bot Date: Fri, 4 Sep 2026 14:49:44 +0000 Subject: [PATCH 10/10] fix(kap-server): race-free ledger compaction and incarnation-safe journal retries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The truncating ledger rewrite could overwrite a tombstone appended by another process mid-compaction, and dropping an acknowledged entry let a read-only peer serve the deleted content from its still-stale MiniDb view. The append log is now never rewritten: the writer alone periodically snapshots the effective tombstones (24h horizon plus entries still present in the authoritative list) into a sidecar file stamped with the append log's byte watermark, and every reader folds the snapshot with just the tail beyond that watermark — so a concurrent append is always included, readers behind a refresh keep filtering within the horizon, and per-query work stays bounded by the snapshot plus the unconsumed tail. Journal removal retries are now bound to the deleted incarnation: the retry tracks the journal's header epoch and cancels instead of deleting when the file at the path belongs to a recreated session with the same id. Co-authored-by: qer --- packages/kap-server/src/search/indexCore.ts | 91 +++++++++++++------ .../ws/v1/sessionEventBroadcaster.ts | 35 ++++++- .../test/search/searchService.test.ts | 6 +- 3 files changed, 100 insertions(+), 32 deletions(-) diff --git a/packages/kap-server/src/search/indexCore.ts b/packages/kap-server/src/search/indexCore.ts index 8df4fe108e2..800eb1f3a4e 100644 --- a/packages/kap-server/src/search/indexCore.ts +++ b/packages/kap-server/src/search/indexCore.ts @@ -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'; @@ -243,30 +244,79 @@ export class SearchIndexCore { return join(this.indexDir, 'deleted-sessions.jsonl'); } - private async readDeletedLedgerRaw(): Promise { - try { - return await readFile(this.deletedLedgerPath, 'utf8'); - } catch (error) { - if ((error as NodeJS.ErrnoException).code === 'ENOENT') return ''; - throw error; - } + private get deletedLedgerSnapshotPath(): string { + return join(this.indexDir, 'deleted-sessions.snapshot.json'); } - private async readDeletedLedger(): Promise> { - const raw = await this.readDeletedLedgerRaw(); - const ledger = new Map(); + private parseDeletedLedgerLines(raw: string, into: Map): void { for (const line of raw.split('\n')) { if (line.length === 0) continue; if (line.startsWith('-')) { - ledger.delete(line.slice(1).split('\t')[0]!); + into.delete(line.slice(1).split('\t')[0]!); continue; } const [id, at] = line.split('\t'); - if (id !== undefined && id.length > 0) ledger.set(id, Number(at ?? 0)); + 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'); @@ -502,9 +552,7 @@ export class SearchIndexCore { } let deletedLedger: Map; - let deletedLedgerRawBefore = ''; try { - deletedLedgerRawBefore = await this.readDeletedLedgerRaw(); deletedLedger = await this.readDeletedLedger(); } catch (error) { this.log.warn('global search: failed to read the session deletion ledger; skipping its purge this pass', { @@ -526,19 +574,10 @@ export class SearchIndexCore { } if (deletedLedger.size > 0) { try { - const initialLines = new Set(deletedLedgerRawBefore.split('\n').filter((l) => l.length > 0)); - const appended = (await this.readDeletedLedgerRaw()) - .split('\n') - .filter((line) => line.length > 0 && !initialLines.has(line)); - const retained = [...deletedLedger] - .filter(([id]) => currentIds.has(id) && !retractedIds.has(id)) - .map(([id, at]) => `${id}\t${at}`); - const kept = [...retained, ...appended]; - await writeFile( - this.deletedLedgerPath, - kept.length > 0 ? `${kept.join('\n')}\n` : '', - 'utf8', + 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), diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts index bed85b35669..0444fa09a9b 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -1,4 +1,4 @@ -import { rm } from 'node:fs/promises'; +import { open, rm } from 'node:fs/promises'; import type { AgentActivityState, ApprovalResponse, @@ -139,19 +139,46 @@ 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 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 attempt = (this.journalRemovalAttempts.get(sessionId) ?? 0) + 1; - this.journalRemovalAttempts.set(sessionId, attempt); + 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) }, diff --git a/packages/kap-server/test/search/searchService.test.ts b/packages/kap-server/test/search/searchService.test.ts index 1c0008b08a6..89149323a55 100644 --- a/packages/kap-server/test/search/searchService.test.ts +++ b/packages/kap-server/test/search/searchService.test.ts @@ -363,8 +363,10 @@ describe('GlobalSearchService', () => { await settleSync(second); expect((await second.search({ query: '香蕉' })).items.length).toBe(1); - const ledgerRaw = await readFile(join(home!, 'search-index', 'deleted-sessions.jsonl'), 'utf8'); - expect(ledgerRaw).not.toContain('s1\t1000'); + 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 () => {