From 63000ccaf2c519161ec6877e9d8fadbde6669da0 Mon Sep 17 00:00:00 2001 From: qer Date: Tue, 8 Sep 2026 01:46:44 +0800 Subject: [PATCH] feat(kap-server): expose session deletion with serialized cleanup --- .changeset/session-delete-action.md | 5 + 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 ++ .../src/app/sessionManager/sessionManager.ts | 1 + .../sessionManager/sessionManagerService.ts | 23 ++- .../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 | 51 ++++++- .../apiSurface.snapshot.test.ts.snap | 4 + packages/kap-server/test/openapi.test.ts | 9 +- .../test/sessionEventBroadcaster.test.ts | 59 +++++++- packages/kap-server/test/sessions.test.ts | 132 ++++++++++++++++++ 18 files changed, 387 insertions(+), 20 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..b286a998e53 --- /dev/null +++ b/.changeset/session-delete-action.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": minor +--- + +Add session deletion to the server API via `POST /api/v1/sessions/{id}:delete`. 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/app/sessionManager/sessionManager.ts b/packages/agent-core-v2/src/app/sessionManager/sessionManager.ts index eff2d413046..c2852e4c785 100644 --- a/packages/agent-core-v2/src/app/sessionManager/sessionManager.ts +++ b/packages/agent-core-v2/src/app/sessionManager/sessionManager.ts @@ -31,6 +31,7 @@ export interface ISessionManager { readonly onDidCreateSession?: Event; readonly onWillCloseSession?: Event; readonly onDidCloseSession?: Event; + readonly onWillDeleteSession?: Event<{ readonly sessionId: string } & IWaitUntil>; readonly onDidArchiveSession?: Event; readonly onDidForkSession?: Event; create(options: CreateManagedSessionOptions): Promise; diff --git a/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts b/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts index dd61cebe17d..d74b74c6098 100644 --- a/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts +++ b/packages/agent-core-v2/src/app/sessionManager/sessionManagerService.ts @@ -50,6 +50,8 @@ export class SessionManager implements ISessionManager { readonly onWillCloseSession = this.willCloseEmitter.event; private readonly didCloseEmitter = new Emitter(); readonly onDidCloseSession = this.didCloseEmitter.event; + private readonly willDeleteEmitter = new Emitter<{ readonly sessionId: string } & IWaitUntil>(); + readonly onWillDeleteSession = this.willDeleteEmitter.event; private readonly didArchiveEmitter = new Emitter(); readonly onDidArchiveSession = this.didArchiveEmitter.event; private readonly didForkEmitter = new Emitter(); @@ -66,9 +68,9 @@ export class SessionManager implements ISessionManager { ? { root: options.workDir } : { workspaceId: options.workspaceId, root: options.workDir }, ); - const controller = this.controllerForWorkspace(workspace.id); - if (options.sessionId === undefined) return controller.create(options); - return this.serializeLifecycle(options.sessionId, () => controller.create(options)); + const create = () => this.controllerForWorkspace(workspace.id).create(options); + if (options.sessionId === undefined) return create(); + return this.serializeLifecycle(options.sessionId, create); } async resume(sessionId: string, options?: ResumeSessionOptions): Promise { @@ -169,6 +171,20 @@ export class SessionManager implements ISessionManager { if (controller === undefined) { throw new Error2(ErrorCodes.SESSION_NOT_FOUND, `session ${sessionId} does not exist`); } + await controller.close(sessionId); + const cleanups: Promise[] = []; + this.willDeleteEmitter.fire({ + sessionId, + signal: new AbortController().signal, + waitUntil: (cleanup) => { + if (Object.isFrozen(cleanups)) throw new Error('waitUntil must be called synchronously'); + cleanups.push(cleanup); + }, + }); + void Object.freeze(cleanups); + const settled = await Promise.allSettled(cleanups); + const failed = settled.find((result) => result.status === 'rejected'); + if (failed?.status === 'rejected') throw failed.reason; await controller.delete(sessionId); }); } @@ -218,6 +234,7 @@ export class SessionManager implements ISessionManager { this.didCreateEmitter.dispose(); this.willCloseEmitter.dispose(); this.didCloseEmitter.dispose(); + this.willDeleteEmitter.dispose(); this.didArchiveEmitter.dispose(); this.didForkEmitter.dispose(); } 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 5789bd9216c..b2f8fc5c84c 100644 --- a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts @@ -97,7 +97,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, @@ -475,6 +475,11 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await dropFileHistorySession({ docs: this.docs, workspaceId: this.workspaceId, 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 f4d7fd89bcb..961e38a8adf 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, @@ -1056,6 +1061,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 1e4fdd8657e..d4799d6fd20 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, @@ -606,6 +607,7 @@ export function registerSessionsRoutes( sessionAbortResponseSchema, startBtwSessionResponseSchema, archiveSessionResponseSchema, + deleteSessionResponseSchema, ]), }, errors: { @@ -854,7 +856,15 @@ export function registerSessionsRoutes( ); } -type SessionAction = 'fork' | 'compact' | 'undo' | 'abort' | 'btw' | 'restore' | 'archive'; +type SessionAction = + | 'fork' + | 'compact' + | 'undo' + | 'abort' + | 'btw' + | 'restore' + | 'archive' + | 'delete'; interface SessionActionExtra { readonly core: Scope; @@ -875,6 +885,7 @@ const sessionActions: ActionTable = { btw: { handle: btwSessionAction }, restore: { handle: restoreSessionAction }, archive: { handle: archiveSessionAction }, + delete: { handle: deleteSessionAction }, }; async function forkSessionAction( @@ -990,6 +1001,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 bb3f6447473..6507e617788 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventBroadcaster.ts @@ -1,3 +1,5 @@ +import { rm } from 'node:fs/promises'; + import type { AgentActivityState, ApprovalResponse, @@ -17,6 +19,7 @@ import { IEventService, ISessionActivityView, ISessionIndex, + ISessionManager, MAIN_AGENT_ID, getLiveSessionById, listSessionPendingInteractions, @@ -140,6 +143,7 @@ export class SessionEventBroadcaster { private readonly pendingStates = new Map>(); private readonly maxBufferSize: number; private readonly coreEventSubscription: IDisposable; + private readonly deletionSubscription: IDisposable | undefined; private closed = false; constructor( @@ -152,6 +156,11 @@ export class SessionEventBroadcaster { }, ) { this.maxBufferSize = opts.maxBufferSize ?? DEFAULT_MAX_BUFFER_SIZE; + this.deletionSubscription = opts.core.accessor.get(ISessionManager).onWillDeleteSession?.( + (event) => { + event.waitUntil(this.purgeSession(event.sessionId)); + }, + ); this.coreEventSubscription = opts.core.accessor .get(IEventService) .subscribe((event) => this.onCoreEvent(event)); @@ -510,6 +519,7 @@ export class SessionEventBroadcaster { if (this.closed) return; this.closed = true; this.coreEventSubscription.dispose(); + this.deletionSubscription?.dispose(); await Promise.all( [...this.pendingStates.values()].map((pending) => pending.catch(() => undefined)), ); @@ -520,6 +530,18 @@ export class SessionEventBroadcaster { this.sessions.clear(); } + private async purgeSession(sessionId: string): Promise { + await this.pendingStates.get(sessionId); + const state = this.sessions.get(sessionId); + if (state !== undefined) { + this.sessions.delete(sessionId); + state.targets.clear(); + await disposeSessionState(state); + } + this.opts.transcriptService?.dropSession(sessionId); + await rm(sessionJournalPath(this.opts.eventsDir, sessionId), { force: true }); + } + private ensureState(sessionId: string): Promise { if (this.closed) return Promise.resolve(undefined); const existing = this.sessions.get(sessionId); @@ -546,7 +568,7 @@ export class SessionEventBroadcaster { sessionJournalPath(this.opts.eventsDir, sessionId), this.opts.logger, ); - if (this.closed) { + if (this.closed || getLiveSessionById(this.opts.core.accessor, sessionId) !== session) { await journal.close(); return undefined; } @@ -649,6 +671,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; @@ -1374,6 +1409,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 9f68a116eee..7f90a88bef4 100644 --- a/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap +++ b/packages/kap-server/test/__snapshots__/apiSurface.snapshot.test.ts.snap @@ -424,6 +424,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 7bbc4593f58..ad265f8a801 100644 --- a/packages/kap-server/test/openapi.test.ts +++ b/packages/kap-server/test/openapi.test.ts @@ -32,12 +32,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(); @@ -46,6 +47,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 de26e0b527d..c555076b75c 100644 --- a/packages/kap-server/test/sessionEventBroadcaster.test.ts +++ b/packages/kap-server/test/sessionEventBroadcaster.test.ts @@ -49,7 +49,7 @@ import { type BroadcastTarget, SessionEventBroadcaster, } from '../src/transport/ws/v1/sessionEventBroadcaster'; -import type { EventEnvelope } from '../src/transport/ws/v1/sessionEventJournal'; +import { SessionEventJournal, type EventEnvelope } from '../src/transport/ws/v1/sessionEventJournal'; import { TranscriptService } from '../src/services/transcript/transcriptService'; type FakeBusEvent = { type: string }; @@ -459,9 +459,12 @@ function makeCore( eventBus = new FakeEventBus(), metaAgents: Record = {}, ): Scope { + const handles = new WeakMap(); const sessionFor = (sid: string) => { const lifecycle = sessions.get(sid); if (lifecycle === undefined) return undefined; + const existing = handles.get(lifecycle); + if (existing !== undefined) return existing; const sessionAccessor = { get: (t: unknown) => { if (t === IAgentLifecycleService) return lifecycle; @@ -470,7 +473,9 @@ function makeCore( return undefined; }, }; - return { id: sid, kind: LifecycleScope.Session, accessor: sessionAccessor, dispose: () => {} }; + const handle = { id: sid, kind: LifecycleScope.Session, accessor: sessionAccessor, dispose: () => {} } as unknown as IScopeHandle; + handles.set(lifecycle, handle); + return handle; }; const sessionLifecycle = { onDidCloseSession: () => ({ dispose: () => {} }), @@ -554,6 +559,33 @@ describe('SessionEventBroadcaster', () => { await rm(dir, { recursive: true, force: true }); }); + it('does not install a pending state after its session disappears', async () => { + const lifecycle = new FakeLifecycle(); + lifecycle.addAgent('main'); + sessions.set('s1', lifecycle); + const journal = await SessionEventJournal.open(join(dir, 's1.jsonl')); + let enter!: () => void; + let release!: () => void; + const entered = new Promise((resolve) => { enter = resolve; }); + const gate = new Promise((resolve) => { release = resolve; }); + const opening = vi.spyOn(SessionEventJournal, 'open').mockImplementation(async () => { + enter(); + await gate; + return journal; + }); + try { + const subscription = bc.subscribe('s1', collectingTarget().target); + await entered; + sessions.delete('s1'); + release(); + expect(await subscription).toBe(false); + expect(await bc.subscribe('s1', collectingTarget().target)).toBe(false); + } finally { + release(); + opening.mockRestore(); + } + }); + it('preserves a real Event2 time in payload and derives the envelope timestamp from it', async () => { const lc = new FakeLifecycle(); const main = lc.addAgent('main'); @@ -1252,6 +1284,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 cfb842f887a..22363feefe0 100644 --- a/packages/kap-server/test/sessions.test.ts +++ b/packages/kap-server/test/sessions.test.ts @@ -28,6 +28,7 @@ import { sessionDirOf, type ScopeSeed, } from '@moonshot-ai/agent-core-v2'; +import { SessionMetaUpdated } from '@moonshot-ai/agent-core-v2/session/sessionMetadata/sessionMetaEvents'; import { TurnStarted } from '@moonshot-ai/agent-core-v2/agent/loop/turnEvents'; import { sessionWarningsResponseSchema } from '@moonshot-ai/agent-core-v2/app/sessionLegacy/sessionProtocol'; import { encodeWorkDirKey } from '@moonshot-ai/agent-core-v2/_base/utils/workdir-slug'; @@ -953,6 +954,137 @@ 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); + await expect(readFile(join(home!, 'server', 'events', `${id}.jsonl`))).rejects.toMatchObject({ code: 'ENOENT' }); + + 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('keeps failed journal cleanup retriable without publishing deletion', async () => { + const created = await postJson('/api/v1/sessions', { metadata: { cwd: home } }); + const id = created.body.data.id; + const journalPath = join(home!, 'server', 'events', `${id}.jsonl`); + await vi.waitFor(async () => expect(await readFile(journalPath, 'utf8')).toContain('journal_header')); + await closeSessionById(server!.core.accessor, id); + await rm(journalPath); + await mkdir(journalPath); + const events: Event2[] = []; + const sub = server!.core.accessor.get(IEventService).subscribe((event) => events.push(event)); + try { + const failed = await postJson(`/api/v1/sessions/${id}:delete`); + expect(failed.body.code).not.toBe(0); + expect(await server!.core.accessor.get(ISessionManager).status(id)).toBeDefined(); + expect(events.filter((event) => event.type === 'event.session.deleted')).toEqual([]); + await rm(journalPath, { recursive: true }); + const retried = await postJson<{ deleted: boolean }>(`/api/v1/sessions/${id}:delete`); + expect(retried.body.data).toEqual({ deleted: true }); + expect(events.filter((event) => event.type === 'event.session.deleted')).toHaveLength(1); + await expect(readFile(journalPath)).rejects.toMatchObject({ code: 'ENOENT' }); + } finally { + sub.dispose(); + } + }); + + it.each(['closing', 'cleanup'] as const)('waits for %s before recreating an explicit session id', async (phase) => { + const manager = server!.core.accessor.get(ISessionManager); + const created = await postJson('/api/v1/sessions', { metadata: { cwd: home } }); + const id = created.body.data.id; + const journalPath = join(home!, 'server', 'events', `${id}.jsonl`); + await vi.waitFor(async () => expect(await readFile(journalPath, 'utf8')).toContain('journal_header')); + const oldJournal = await readFile(journalPath, 'utf8'); + let enter!: () => void; + let release!: () => void; + const entered = new Promise((resolve) => { enter = resolve; }); + const gate = new Promise((resolve) => { release = resolve; }); + const event = phase === 'closing' ? manager.onWillCloseSession! : manager.onWillDeleteSession!; + const sub = event((event) => { + if (event.sessionId !== id) return; + event.waitUntil(gate); + enter(); + }); + try { + const deletion = manager.delete(id); + await entered; + let recreated = false; + const creation = manager.create({ sessionId: id, workDir: home! }).then((handle) => { + recreated = true; + return handle; + }); + await new Promise((resolve) => setImmediate(resolve)); + expect(recreated).toBe(false); + release(); + await deletion; + await creation; + server!.core.accessor.get(IEventService).publish(new SessionMetaUpdated({ + payload: { sessionId: id, agentId: 'main', patch: { title: 'Recreated session' } }, + })); + await vi.waitFor(async () => { + const journal = await readFile(journalPath, 'utf8'); + expect(journal).toContain('journal_header'); + expect(JSON.parse(journal.split('\n')[0]!).epoch).not.toBe(JSON.parse(oldJournal.split('\n')[0]!).epoch); + }); + expect(manager.get(id)).toBeDefined(); + } finally { + release(); + 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 } });