mirror of
https://github.com/MoonshotAI/kimi-code.git
synced 2026-08-11 01:37:31 +00:00
fix: harden the session close/archive and create/fork lifecycle
The closing registry now records the operation kind: an archive arriving during a plain close waits it out and lands the archived flag on the persisted document instead of riding the close to success, and a failing close hook no longer strands a half-closed session — the teardown always completes while the hook error still reaches the caller. create and fork reserve their target id synchronously with the existence check, so a concurrent create/fork of the same id loses up front and can never tear down the winner's scope or directory.
This commit is contained in:
parent
9a95e1c02f
commit
e7c397a7c8
2 changed files with 230 additions and 22 deletions
|
|
@ -11,8 +11,12 @@
|
||||||
* session work can cancel itself and drop its write-back. A close/archive
|
* session work can cancel itself and drop its write-back. A close/archive
|
||||||
* is tracked in a closing registry from that first synchronous step until
|
* is tracked in a closing registry from that first synchronous step until
|
||||||
* the scope is disposed: `get`/`list` hide the session while it is
|
* the scope is disposed: `get`/`list` hide the session while it is
|
||||||
* closing, and `resume` waits the close out and re-materializes instead
|
* closing, `resume` waits the close out and re-materializes instead of
|
||||||
* of returning the doomed handle,
|
* returning the doomed handle, and a failing close hook never strands the
|
||||||
|
* teardown. A duplicate close joins the in-flight one, while an archive
|
||||||
|
* arriving during a plain close waits it out and then lands the archived
|
||||||
|
* flag directly on the persisted document (it never rides the close to
|
||||||
|
* success without archiving),
|
||||||
* tearing sessions down on close/archive — archiving flags the session's
|
* tearing sessions down on close/archive — archiving flags the session's
|
||||||
* `sessionMetadata`, removes its `agentLifecycle` agents, restoring clears
|
* `sessionMetadata`, removes its `agentLifecycle` agents, restoring clears
|
||||||
* the archived flag, and broadcasts through `event`; session start and
|
* the archived flag, and broadcasts through `event`; session start and
|
||||||
|
|
@ -161,6 +165,11 @@ type MaterializeSessionOptions = Omit<CreateSessionOptions, 'sessionId'> & {
|
||||||
readonly sessionId: string;
|
readonly sessionId: string;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
type ClosingEntry = {
|
||||||
|
readonly kind: 'close' | 'archive';
|
||||||
|
readonly run: Promise<void>;
|
||||||
|
};
|
||||||
|
|
||||||
export class WorkspaceHandlerService extends Disposable implements IWorkspaceHandlerService {
|
export class WorkspaceHandlerService extends Disposable implements IWorkspaceHandlerService {
|
||||||
declare readonly _serviceBrand: undefined;
|
declare readonly _serviceBrand: undefined;
|
||||||
private readonly sessions = new Map<string, ISessionScopeHandle>();
|
private readonly sessions = new Map<string, ISessionScopeHandle>();
|
||||||
|
|
@ -173,8 +182,9 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
private readonly _onDidForkSession = this._register(new Emitter<SessionForkedEvent>());
|
private readonly _onDidForkSession = this._register(new Emitter<SessionForkedEvent>());
|
||||||
readonly onDidForkSession: Event<SessionForkedEvent> = this._onDidForkSession.event;
|
readonly onDidForkSession: Event<SessionForkedEvent> = this._onDidForkSession.event;
|
||||||
private readonly resuming = new Map<string, Promise<ISessionScopeHandle | undefined>>();
|
private readonly resuming = new Map<string, Promise<ISessionScopeHandle | undefined>>();
|
||||||
private readonly closing = new Map<string, Promise<void>>();
|
private readonly closing = new Map<string, ClosingEntry>();
|
||||||
private readonly sessionLifetimes = new Map<string, AbortController>();
|
private readonly sessionLifetimes = new Map<string, AbortController>();
|
||||||
|
private readonly reservedTargets = new Set<string>();
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
@IInstantiationService private readonly instantiation: IInstantiationService,
|
@IInstantiationService private readonly instantiation: IInstantiationService,
|
||||||
|
|
@ -219,6 +229,24 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
|
|
||||||
async create(opts: CreateSessionOptions): Promise<ISessionScopeHandle> {
|
async create(opts: CreateSessionOptions): Promise<ISessionScopeHandle> {
|
||||||
const sessionId = opts.sessionId ?? createSessionId();
|
const sessionId = opts.sessionId ?? createSessionId();
|
||||||
|
if (this.sessions.has(sessionId) || this.reservedTargets.has(sessionId)) {
|
||||||
|
throw new Error2(
|
||||||
|
ErrorCodes.SESSION_ALREADY_EXISTS,
|
||||||
|
`Session "${sessionId}" already exists`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
this.reservedTargets.add(sessionId);
|
||||||
|
try {
|
||||||
|
return await this.createReserved(opts, sessionId);
|
||||||
|
} finally {
|
||||||
|
this.reservedTargets.delete(sessionId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async createReserved(
|
||||||
|
opts: CreateSessionOptions,
|
||||||
|
sessionId: string,
|
||||||
|
): Promise<ISessionScopeHandle> {
|
||||||
const handle = await this.materializeSession({ ...opts, sessionId });
|
const handle = await this.materializeSession({ ...opts, sessionId });
|
||||||
try {
|
try {
|
||||||
const main =
|
const main =
|
||||||
|
|
@ -376,7 +404,7 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
if (closing !== undefined) {
|
if (closing !== undefined) {
|
||||||
// A close is already in flight — wait it out and re-materialize
|
// A close is already in flight — wait it out and re-materialize
|
||||||
// instead of handing back the doomed handle.
|
// instead of handing back the doomed handle.
|
||||||
return closing.then(() => this.resume(sessionId, opts));
|
return closing.run.then(() => this.resume(sessionId, opts));
|
||||||
}
|
}
|
||||||
const live = this.sessions.get(sessionId);
|
const live = this.sessions.get(sessionId);
|
||||||
if (live !== undefined) return Promise.resolve(live);
|
if (live !== undefined) return Promise.resolve(live);
|
||||||
|
|
@ -429,27 +457,43 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
|
|
||||||
async close(sessionId: string): Promise<void> {
|
async close(sessionId: string): Promise<void> {
|
||||||
const inflight = this.closing.get(sessionId);
|
const inflight = this.closing.get(sessionId);
|
||||||
if (inflight !== undefined) return inflight;
|
if (inflight !== undefined) return inflight.run;
|
||||||
const handle = this.sessions.get(sessionId);
|
const handle = this.sessions.get(sessionId);
|
||||||
if (handle === undefined) return;
|
if (handle === undefined) return;
|
||||||
return this.trackClosing(sessionId, this.doClose(sessionId, handle));
|
return this.trackClosing(sessionId, 'close', this.doClose(sessionId, handle));
|
||||||
}
|
}
|
||||||
|
|
||||||
private async doClose(sessionId: string, handle: ISessionScopeHandle): Promise<void> {
|
private async doClose(sessionId: string, handle: ISessionScopeHandle): Promise<void> {
|
||||||
this.invalidateSessionLifetime(sessionId);
|
this.invalidateSessionLifetime(sessionId);
|
||||||
await this.announceWillClose({ sessionId, handle, reason: 'exit' });
|
try {
|
||||||
this.sessions.delete(sessionId);
|
await this.announceWillClose({ sessionId, handle, reason: 'exit' });
|
||||||
await this.drainAgents(handle);
|
} finally {
|
||||||
handle.dispose();
|
// A failing close hook must not strand the session half-open: the
|
||||||
this._onDidCloseSession.fire({ sessionId });
|
// teardown always completes; the hook error still reaches the caller.
|
||||||
|
this.sessions.delete(sessionId);
|
||||||
|
await this.drainAgents(handle);
|
||||||
|
handle.dispose();
|
||||||
|
this._onDidCloseSession.fire({ sessionId });
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async archive(sessionId: string): Promise<void> {
|
async archive(sessionId: string): Promise<void> {
|
||||||
const inflight = this.closing.get(sessionId);
|
const inflight = this.closing.get(sessionId);
|
||||||
if (inflight !== undefined) return inflight;
|
if (inflight !== undefined) {
|
||||||
|
if (inflight.kind === 'archive') return inflight.run;
|
||||||
|
// A plain close is already tearing the session down — the archive
|
||||||
|
// must not ride it to success. Wait it out (tracked, so a resume
|
||||||
|
// queues behind the marker write), then land the archived flag
|
||||||
|
// directly on the persisted document of the now-closed session.
|
||||||
|
return this.trackClosing(
|
||||||
|
sessionId,
|
||||||
|
'archive',
|
||||||
|
inflight.run.then(() => this.markArchivedOnDisk(sessionId)),
|
||||||
|
);
|
||||||
|
}
|
||||||
const handle = this.sessions.get(sessionId);
|
const handle = this.sessions.get(sessionId);
|
||||||
if (handle === undefined) return;
|
if (handle === undefined) return;
|
||||||
return this.trackClosing(sessionId, this.doArchive(sessionId, handle));
|
return this.trackClosing(sessionId, 'archive', this.doArchive(sessionId, handle));
|
||||||
}
|
}
|
||||||
|
|
||||||
private async doArchive(sessionId: string, handle: ISessionScopeHandle): Promise<void> {
|
private async doArchive(sessionId: string, handle: ISessionScopeHandle): Promise<void> {
|
||||||
|
|
@ -461,18 +505,32 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
type: 'event.session.archived',
|
type: 'event.session.archived',
|
||||||
payload: { sessionId },
|
payload: { sessionId },
|
||||||
});
|
});
|
||||||
await this.announceWillClose({ sessionId, handle, reason: 'exit' });
|
try {
|
||||||
this.sessions.delete(sessionId);
|
await this.announceWillClose({ sessionId, handle, reason: 'exit' });
|
||||||
handle.dispose();
|
} finally {
|
||||||
this._onDidArchiveSession.fire({ sessionId });
|
this.sessions.delete(sessionId);
|
||||||
|
handle.dispose();
|
||||||
|
this._onDidArchiveSession.fire({ sessionId });
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async trackClosing(sessionId: string, run: Promise<void>): Promise<void> {
|
private async markArchivedOnDisk(sessionId: string): Promise<void> {
|
||||||
this.closing.set(sessionId, run);
|
const scope = sessionScopeOf(this.handlerScope, sessionId);
|
||||||
|
const meta = await this.docs.get<SessionMeta>(scope, 'state.json');
|
||||||
|
if (meta === undefined) return;
|
||||||
|
await this.docs.set(scope, 'state.json', { ...meta, archived: true, updatedAt: Date.now() });
|
||||||
|
}
|
||||||
|
|
||||||
|
private async trackClosing(
|
||||||
|
sessionId: string,
|
||||||
|
kind: ClosingEntry['kind'],
|
||||||
|
run: Promise<void>,
|
||||||
|
): Promise<void> {
|
||||||
|
this.closing.set(sessionId, { kind, run });
|
||||||
try {
|
try {
|
||||||
await run;
|
await run;
|
||||||
} finally {
|
} finally {
|
||||||
if (this.closing.get(sessionId) === run) this.closing.delete(sessionId);
|
if (this.closing.get(sessionId)?.run === run) this.closing.delete(sessionId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -509,7 +567,7 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
// Wait out an in-flight close of the source instead of capturing its
|
// Wait out an in-flight close of the source instead of capturing its
|
||||||
// doomed handle (a close starting mid-fork remains an accepted,
|
// doomed handle (a close starting mid-fork remains an accepted,
|
||||||
// pre-existing window — fork never quiesces the source).
|
// pre-existing window — fork never quiesces the source).
|
||||||
await this.closing.get(sourceId);
|
await this.closing.get(sourceId)?.run;
|
||||||
const sourceHandle = this.sessions.get(sourceId);
|
const sourceHandle = this.sessions.get(sourceId);
|
||||||
const indexSummary = await this.index.get(sourceId);
|
const indexSummary = await this.index.get(sourceId);
|
||||||
if (
|
if (
|
||||||
|
|
@ -530,6 +588,7 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
let targetId: string | undefined;
|
let targetId: string | undefined;
|
||||||
let target: ISessionScopeHandle | undefined;
|
let target: ISessionScopeHandle | undefined;
|
||||||
let targetSessionDir: string | undefined;
|
let targetSessionDir: string | undefined;
|
||||||
|
let targetReserved = false;
|
||||||
try {
|
try {
|
||||||
const sourceMeta =
|
const sourceMeta =
|
||||||
sourceHandle !== undefined
|
sourceHandle !== undefined
|
||||||
|
|
@ -537,7 +596,18 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
: await this.readMetaFromDisk(sourceId);
|
: await this.readMetaFromDisk(sourceId);
|
||||||
|
|
||||||
targetId = opts.newSessionId ?? createSessionId();
|
targetId = opts.newSessionId ?? createSessionId();
|
||||||
if (this.sessions.has(targetId) || (await this.index.get(targetId)) !== undefined) {
|
if (this.sessions.has(targetId) || this.reservedTargets.has(targetId)) {
|
||||||
|
throw new Error2(
|
||||||
|
ErrorCodes.SESSION_ALREADY_EXISTS,
|
||||||
|
`Session "${targetId}" already exists`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
// The check and the reservation are one synchronous step: a concurrent
|
||||||
|
// create/fork of the same target id loses before either materializes,
|
||||||
|
// so the loser can never tear down the winner's scope or directory.
|
||||||
|
this.reservedTargets.add(targetId);
|
||||||
|
targetReserved = true;
|
||||||
|
if ((await this.index.get(targetId)) !== undefined) {
|
||||||
throw new Error2(
|
throw new Error2(
|
||||||
ErrorCodes.SESSION_ALREADY_EXISTS,
|
ErrorCodes.SESSION_ALREADY_EXISTS,
|
||||||
`Session "${targetId}" already exists`,
|
`Session "${targetId}" already exists`,
|
||||||
|
|
@ -612,6 +682,10 @@ export class WorkspaceHandlerService extends Disposable implements IWorkspaceHan
|
||||||
await this.hostFs.remove(targetSessionDir).catch(() => {});
|
await this.hostFs.remove(targetSessionDir).catch(() => {});
|
||||||
}
|
}
|
||||||
throw error;
|
throw error;
|
||||||
|
} finally {
|
||||||
|
if (targetReserved && targetId !== undefined) {
|
||||||
|
this.reservedTargets.delete(targetId);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -602,6 +602,140 @@ describe('WorkspaceHandlerService', () => {
|
||||||
expect(svc.get('s1')).toBeUndefined();
|
expect(svc.get('s1')).toBeUndefined();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('resume waits out an in-flight close and returns a fresh scope', async () => {
|
||||||
|
const svc = await build([
|
||||||
|
stubPair(IWorkspaceService, persistentWorkspaceStub()),
|
||||||
|
stubPair(ISessionIndex, sessionIndexWithSummary('s1', '/tmp/proj')),
|
||||||
|
stubPair(IAgentLifecycleService, agentLifecycleWithMainStub()),
|
||||||
|
]);
|
||||||
|
const original = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' });
|
||||||
|
|
||||||
|
let closeStarted!: () => void;
|
||||||
|
let releaseClose!: () => void;
|
||||||
|
const closeBegan = new Promise<void>((resolve) => {
|
||||||
|
closeStarted = resolve;
|
||||||
|
});
|
||||||
|
const closeGate = new Promise<void>((resolve) => {
|
||||||
|
releaseClose = resolve;
|
||||||
|
});
|
||||||
|
original.accessor
|
||||||
|
.get(ISessionLifecycleHooks)
|
||||||
|
.onWillCloseSession.register('test-block', async (_event, next) => {
|
||||||
|
closeStarted();
|
||||||
|
await closeGate;
|
||||||
|
await next();
|
||||||
|
});
|
||||||
|
|
||||||
|
const closing = svc.close('s1');
|
||||||
|
await closeBegan;
|
||||||
|
// The closing session is hidden from get/list immediately.
|
||||||
|
expect(svc.get('s1')).toBeUndefined();
|
||||||
|
expect(svc.list()).toEqual([]);
|
||||||
|
|
||||||
|
let settled = false;
|
||||||
|
const resumed = svc.resume('s1').then((handle) => {
|
||||||
|
settled = true;
|
||||||
|
return handle;
|
||||||
|
});
|
||||||
|
await tick();
|
||||||
|
expect(settled).toBe(false);
|
||||||
|
|
||||||
|
releaseClose();
|
||||||
|
await closing;
|
||||||
|
const fresh = await resumed;
|
||||||
|
expect(fresh).toBeDefined();
|
||||||
|
expect(fresh).not.toBe(original);
|
||||||
|
expect(svc.get('s1')).toBe(fresh);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('archive during a plain close waits it out and lands the archived flag on disk', async () => {
|
||||||
|
const docs = new Map<string, unknown>();
|
||||||
|
const docStore = {
|
||||||
|
_serviceBrand: undefined,
|
||||||
|
get: (scope: string, key: string) => Promise.resolve(docs.get(`${scope}/${key}`)),
|
||||||
|
set: (scope: string, key: string, value: unknown) => {
|
||||||
|
docs.set(`${scope}/${key}`, value);
|
||||||
|
return Promise.resolve();
|
||||||
|
},
|
||||||
|
delete: () => Promise.resolve(),
|
||||||
|
list: () => Promise.resolve([]),
|
||||||
|
watch: () => (_listener: unknown) => ({ dispose: () => {} }),
|
||||||
|
acquire: () => ({ dispose: () => {} }),
|
||||||
|
} as unknown as IAtomicDocumentStore;
|
||||||
|
const svc = await build([stubPair(IAtomicDocumentStore, docStore)]);
|
||||||
|
const original = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' });
|
||||||
|
const metaKey = `${original.accessor.get(ISessionContext).metaScope}/state.json`;
|
||||||
|
docs.set(metaKey, { id: 's1', version: 2, createdAt: 1, updatedAt: 1, archived: false });
|
||||||
|
|
||||||
|
let releaseClose!: () => void;
|
||||||
|
const closeGate = new Promise<void>((resolve) => {
|
||||||
|
releaseClose = resolve;
|
||||||
|
});
|
||||||
|
original.accessor
|
||||||
|
.get(ISessionLifecycleHooks)
|
||||||
|
.onWillCloseSession.register('test-block', async (_event, next) => {
|
||||||
|
await closeGate;
|
||||||
|
await next();
|
||||||
|
});
|
||||||
|
|
||||||
|
const closing = svc.close('s1');
|
||||||
|
await tick();
|
||||||
|
const archived = svc.archive('s1');
|
||||||
|
releaseClose();
|
||||||
|
await closing;
|
||||||
|
await archived;
|
||||||
|
|
||||||
|
// The archive must not ride the close: the marker lands on the persisted
|
||||||
|
// document after the close completes.
|
||||||
|
expect(docs.get(metaKey)).toMatchObject({ archived: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
it('still completes the teardown when a close hook fails', async () => {
|
||||||
|
const svc = await build();
|
||||||
|
const original = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' });
|
||||||
|
original.accessor
|
||||||
|
.get(ISessionLifecycleHooks)
|
||||||
|
.onWillCloseSession.register('test-boom', async () => {
|
||||||
|
throw new Error('hook boom');
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(svc.close('s1')).rejects.toThrow('hook boom');
|
||||||
|
// The hook error reaches the caller, but the session is gone, not
|
||||||
|
// stranded as a half-closed zombie.
|
||||||
|
expect(svc.get('s1')).toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects a concurrent create with the same reserved session id', async () => {
|
||||||
|
const svc = await build();
|
||||||
|
|
||||||
|
const [first, second] = await Promise.allSettled([
|
||||||
|
svc.create({ sessionId: 's1', workDir: '/tmp/proj' }),
|
||||||
|
svc.create({ sessionId: 's1', workDir: '/tmp/proj' }),
|
||||||
|
]);
|
||||||
|
|
||||||
|
expect([first.status, second.status].sort()).toEqual(['fulfilled', 'rejected']);
|
||||||
|
const rejection = [first, second].find((result) => result.status === 'rejected');
|
||||||
|
expect((rejection as PromiseRejectedResult).reason).toMatchObject({
|
||||||
|
code: ErrorCodes.SESSION_ALREADY_EXISTS,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it('rejects a concurrent fork racing for the same reserved target id', async () => {
|
||||||
|
const svc = await build();
|
||||||
|
await svc.create({ sessionId: 'src', workDir: '/tmp/proj' });
|
||||||
|
|
||||||
|
const [first, second] = await Promise.allSettled([
|
||||||
|
svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }),
|
||||||
|
svc.fork({ sourceSessionId: 'src', newSessionId: 'dst' }),
|
||||||
|
]);
|
||||||
|
|
||||||
|
expect([first.status, second.status].sort()).toEqual(['fulfilled', 'rejected']);
|
||||||
|
const rejection = [first, second].find((result) => result.status === 'rejected');
|
||||||
|
expect((rejection as PromiseRejectedResult).reason).toMatchObject({
|
||||||
|
code: ErrorCodes.SESSION_ALREADY_EXISTS,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
it('create seeds identity and materializes metadata', async () => {
|
it('create seeds identity and materializes metadata', async () => {
|
||||||
const svc = await build();
|
const svc = await build();
|
||||||
const h = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' });
|
const h = await svc.create({ sessionId: 's1', workDir: '/tmp/proj' });
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue