diff --git a/.changeset/fix-restore-crash-loop.md b/.changeset/fix-restore-crash-loop.md new file mode 100644 index 00000000000..8707e68e06c --- /dev/null +++ b/.changeset/fix-restore-crash-loop.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Fix repeated server crashes when resuming a session that was interrupted in the middle of a turn. diff --git a/packages/agent-core-v2/src/_base/di/instantiation.ts b/packages/agent-core-v2/src/_base/di/instantiation.ts index 4553e8d76bb..7adc73615ea 100644 --- a/packages/agent-core-v2/src/_base/di/instantiation.ts +++ b/packages/agent-core-v2/src/_base/di/instantiation.ts @@ -196,6 +196,7 @@ export interface IInstantiationService { provideAll(entries: ReadonlyArray): void; unprovide(id: ServiceIdentifier): void; dispose(): void; + disposeAsync(): Promise; } export const IInstantiationService: ServiceIdentifier = diff --git a/packages/agent-core-v2/src/_base/di/instantiationService.ts b/packages/agent-core-v2/src/_base/di/instantiationService.ts index 32bab205248..cd3ca5272a3 100644 --- a/packages/agent-core-v2/src/_base/di/instantiationService.ts +++ b/packages/agent-core-v2/src/_base/di/instantiationService.ts @@ -487,6 +487,10 @@ export class InstantiationService implements IInstantiationService { return this._ledger.register(disposer, label); } + anchorKernelFinalizer(disposer: Disposer, label: string): LedgerEntry { + return this._ledger.registerFinalizer(disposer, label); + } + private _getFiberHost(): FiberHost { this._fiberHost ??= { mintUid: () => ++this._root()._nextUnitUid, @@ -648,18 +652,31 @@ export class InstantiationService implements IInstantiationService { return new InstantiationService(services, this._strict, this, this._enableTracing); } + private _disposePromise: Promise | undefined; + dispose(): void { + void this.disposeAsync(); + } + + disposeAsync(): Promise { + this._disposePromise ??= this.disposeCore(); + return this._disposePromise; + } + + private disposeCore(): Promise { if (this._disposed) { - return; + return Promise.resolve(); } this._disposed = true; + const childTeardowns: Promise[] = []; + let teardown: void | Promise = undefined; try { for (const child of Array.from(this._children)) { - child.dispose(); + childTeardowns.push(child.disposeAsync()); } this._children.clear(); - void this._ledger.teardown('scope-close'); + teardown = this._ledger.teardown('scope-close'); this._services.dispose(); this.cascade.dispose(); for (const view of this._collectionViews.values()) { @@ -674,6 +691,7 @@ export class InstantiationService implements IInstantiationService { this._parent._children.delete(this); } } + return Promise.all([...childTeardowns, Promise.resolve(teardown)]).then(() => undefined); } private _createInstance(ctor: any, args: unknown[], _trace: Trace, unit?: { diff --git a/packages/agent-core-v2/src/_base/di/scope.ts b/packages/agent-core-v2/src/_base/di/scope.ts index ca66120984f..4db1ae68387 100644 --- a/packages/agent-core-v2/src/_base/di/scope.ts +++ b/packages/agent-core-v2/src/_base/di/scope.ts @@ -112,7 +112,7 @@ export interface IScopeHandle { readonly id: string; readonly kind: K; readonly accessor: ServicesAccessor; - dispose(): void; + dispose(): void | Promise; } export type IAppScopeHandle = IScopeHandle<'app'>; @@ -171,7 +171,7 @@ export function createScopedChildHandle( get: (serviceId: ServiceIdentifier): T => child.invokeFunction((a) => a.get(serviceId)), }; - return { id, kind, accessor, dispose: () => child.dispose() }; + return { id, kind, accessor, dispose: () => child.disposeAsync() }; } export class Scope implements IDisposable { diff --git a/packages/agent-core-v2/src/_base/di/scopeUnits.ts b/packages/agent-core-v2/src/_base/di/scopeUnits.ts index 7219222c3b8..db02af0757b 100644 --- a/packages/agent-core-v2/src/_base/di/scopeUnits.ts +++ b/packages/agent-core-v2/src/_base/di/scopeUnits.ts @@ -22,7 +22,7 @@ export function watchScopeUnits(container: InstantiationService, kind: ScopeKind const foldLedger = new Ledger(`scope-units:${kind}`); container.anchorKernelEntry((reason) => foldLedger.teardown(reason), `scope-units:${kind}`); - const materialized = new Map void>(); + const materialized = new Map void | Promise>(); const materialize = (record: StoredRecord): void => { const recipe = record.value as ServiceRecipe; @@ -32,7 +32,7 @@ export function watchScopeUnits(container: InstantiationService, kind: ScopeKind if (isClassRecipe(recipe)) { const instance = host.constructService(recipe, undefined) as Partial; unitLedger.register(() => { - instance.dispose?.(); + return instance.dispose?.(); }, `unit:${name}`); } else { const facade = new FiberRuntime( @@ -57,23 +57,23 @@ export function watchScopeUnits(container: InstantiationService, kind: ScopeKind } let retracted = false; - const retract = (): void => { + const retract = (): void | Promise => { if (retracted) { - return; + return undefined; } retracted = true; materialized.delete(record.id); - void unitLedger.teardown('unload'); + return unitLedger.teardown('unload'); }; if (!record.providerBook.isActive) { - retract(); + void retract(); return; } record.providerBook.register(() => { - retract(); + void retract(); }, `scope-units:${kind}`); foldLedger.register(() => { - retract(); + return retract(); }, `record:${name}`); materialized.set(record.id, retract); }; @@ -92,7 +92,7 @@ export function watchScopeUnits(container: InstantiationService, kind: ScopeKind } for (const [id, retract] of Array.from(materialized)) { if (!seen.has(id)) { - retract(); + void retract(); } } }; diff --git a/packages/agent-core-v2/src/_base/di/test.ts b/packages/agent-core-v2/src/_base/di/test.ts index 5bdfc3e25cc..5b41a4e68fe 100644 --- a/packages/agent-core-v2/src/_base/di/test.ts +++ b/packages/agent-core-v2/src/_base/di/test.ts @@ -29,7 +29,9 @@ export function createScopedTestHost(appStubs: ScopeSeed = []): ScopedTestHost { id: handle.id, kind: handle.kind, accessor: handle.accessor, - dispose: () => handle.dispose(), + dispose: () => { + void handle.dispose(); + }, } as Scope; } return app.createChild(kind, id, { seeds: stubs }); diff --git a/packages/agent-core-v2/src/_base/di/testInstantiationService.ts b/packages/agent-core-v2/src/_base/di/testInstantiationService.ts index 7537f4de2cc..aaa07a17588 100644 --- a/packages/agent-core-v2/src/_base/di/testInstantiationService.ts +++ b/packages/agent-core-v2/src/_base/di/testInstantiationService.ts @@ -262,6 +262,14 @@ export class TestInstantiationService extends InstantiationService implements ID super.dispose(); } } + + public override disposeAsync(): Promise { + sinon.restore(); + if (this._properDispose) { + return super.disposeAsync(); + } + return Promise.resolve(); + } } interface SinonOptions { diff --git a/packages/agent-core-v2/src/_base/lifecycle/ledger.ts b/packages/agent-core-v2/src/_base/lifecycle/ledger.ts index f39b15c177c..d49c9ce0786 100644 --- a/packages/agent-core-v2/src/_base/lifecycle/ledger.ts +++ b/packages/agent-core-v2/src/_base/lifecycle/ledger.ts @@ -65,6 +65,11 @@ export class Ledger { return this._push({ label, kind: 'disposer', active: true, run: disposer }); } + registerFinalizer(disposer: Disposer, label: string = 'finalizer'): LedgerEntry { + this._assertActive('registerFinalizer'); + return this._push({ label, kind: 'disposer', active: true, run: disposer }, true); + } + effect(body: EffectBody, label: string = 'effect'): LedgerEntry { this._assertActive('effect'); const out = body(); @@ -151,11 +156,15 @@ export class Ledger { return infos; } - private _push(record: EntryRecord): LedgerEntry { + private _push(record: EntryRecord, front = false): LedgerEntry { if (Ledger.captureStacks) { record.stack = new Error('Ledger registration').stack; } - this._records.push(record); + if (front) { + this._records.unshift(record); + } else { + this._records.push(record); + } return { label: record.label, get disposed() { diff --git a/packages/agent-core-v2/src/session/agentLifecycle/agentLifecycleService.ts b/packages/agent-core-v2/src/session/agentLifecycle/agentLifecycleService.ts index be90c46eb84..95842611cdd 100644 --- a/packages/agent-core-v2/src/session/agentLifecycle/agentLifecycleService.ts +++ b/packages/agent-core-v2/src/session/agentLifecycle/agentLifecycleService.ts @@ -222,6 +222,7 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle eventBus?.activateAgent(agent); let managed: ManagedAgent | undefined; let didCreate = false; + let finalizerArmed = false; try { const handle = createScopedChildHandle( this.instantiation, @@ -237,13 +238,17 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle }], ], configureContainer: (container) => { + container.anchorKernelFinalizer(() => { + eventBus?.deactivateAgent(agent); + }, 'agent-event-bus-deactivate'); + finalizerArmed = true; this.adopt({ id: agentId, kind: LifecycleScope.Agent, accessor: { get: (id) => container.invokeFunction((accessor) => accessor.get(id)), }, - dispose: () => { container.dispose(); }, + dispose: () => container.disposeAsync(), }); managed = this.roster.get(agentId); }, @@ -274,10 +279,10 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle await managed.runtimeSet.close().catch(() => undefined); managed.killSpace(); try { - managed.handle.dispose(); + await managed.handle.dispose(); } catch { } } - eventBus?.deactivateAgent(agent); + if (!finalizerArmed) eventBus?.deactivateAgent(agent); if (didCreate) this.onDidCloseEmitter.fire(agent); throw error; } @@ -446,10 +451,7 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle await Promise.all([loop.settled(), compactionSettled, prompt.drain(reason)]); await managed.runtimeSet.close(); managed.killSpace(); - handle.dispose(); - this.instantiation.invokeFunction((accessor) => - (accessor.get(ISessionEventBus) as ISessionEventBus | undefined)?.deactivateAgent(agent), - ); + await handle.dispose(); if (this.roster.get(agent.agentId) === managed) this.roster.delete(agent.agentId); this.onDidCloseEmitter.fire(agent); } diff --git a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts index 0f90357e167..b5ed990d891 100644 --- a/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts +++ b/packages/agent-core-v2/src/workspace/sessionLifecycle/sessionLifecycleService.ts @@ -208,7 +208,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec const sessionDir = handle.accessor.get(ISessionContext).sessionDir; this.sessions.delete(sessionId); await this.drainAgents(handle).catch(() => {}); - handle.dispose(); + void handle.dispose(); await this.hostFs.remove(sessionDir).catch(() => {}); throw error; } @@ -282,7 +282,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.pluginAgentProfileLoader.ready, ]); } catch (error) { - handle.dispose(); + void handle.dispose(); void this.explicitAgentProfileLoader.reload().catch(() => undefined); throw error; } @@ -364,7 +364,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.announceCreated({ sessionId, handle, source: 'resume' }); } catch (error) { this.sessions.delete(sessionId); - handle.dispose(); + void handle.dispose(); throw error; } return handle; @@ -387,7 +387,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec await this.appendLogStore.drainRetirements(); await drainSessionMetadataWrites(); await this.indexMirror.drain(); - handle.dispose(); + void handle.dispose(); await drainLogCloses(); this._onDidCloseSession.fire({ sessionId }); } @@ -408,7 +408,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec this.sessions.delete(sessionId); await drainSessionMetadataWrites(); await this.indexMirror.drain(); - handle.dispose(); + void handle.dispose(); await drainLogCloses(); this._onDidArchiveSession.fire({ sessionId }); } @@ -589,7 +589,7 @@ export class SessionLifecycleService extends Disposable implements ISessionLifec } if (target !== undefined) { try { - target.dispose(); + void target.dispose(); } catch { } } diff --git a/packages/agent-core-v2/test/_base/di/child.test.ts b/packages/agent-core-v2/test/_base/di/child.test.ts index efdb838eb7d..5a13c768959 100644 --- a/packages/agent-core-v2/test/_base/di/child.test.ts +++ b/packages/agent-core-v2/test/_base/di/child.test.ts @@ -214,6 +214,64 @@ describe('InstantiationService.createChild', () => { expect(events).toEqual(['disposed']); }); + it('repeated disposeAsync returns the in-flight teardown promise', async () => { + const events: string[] = []; + let releaseGate!: () => void; + const ix = new InstantiationService(new ServiceCollection()); + ix.anchorKernelEntry(() => { + events.push('finalizer'); + }, 'finalizer'); + ix.anchorKernelEntry(() => { + events.push('gate-entered'); + return new Promise((resolve) => { + releaseGate = resolve; + }); + }, 'gate'); + + const first = ix.disposeAsync(); + const second = ix.disposeAsync(); + let secondSettled = false; + void second.then(() => { + secondSettled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(events).toEqual(['gate-entered']); + expect(secondSettled).toBe(false); + releaseGate(); + await Promise.all([first, second]); + expect(events).toEqual(['gate-entered', 'finalizer']); + }); + + it('disposeAsync awaits asynchronous child container teardown', async () => { + const events: string[] = []; + let releaseChildGate!: () => void; + const parent = new InstantiationService(new ServiceCollection()); + const child = parent.createChild(new ServiceCollection()) as InstantiationService; + child.anchorKernelEntry(() => { + events.push('child-finalizer'); + }, 'child-finalizer'); + child.anchorKernelEntry(() => { + events.push('child-gate-entered'); + return new Promise((resolve) => { + releaseChildGate = resolve; + }); + }, 'child-gate'); + parent.anchorKernelEntry(() => { + events.push('parent-finalizer'); + }, 'parent-finalizer'); + + let settled = false; + const disposal = parent.disposeAsync().then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(events).toEqual(['child-gate-entered', 'parent-finalizer']); + expect(settled).toBe(false); + releaseChildGate(); + await disposal; + expect(events).toEqual(['child-gate-entered', 'parent-finalizer', 'child-finalizer']); + }); + it('parent dispose propagates to children', () => { const events: string[] = []; interface IParentSvc { diff --git a/packages/agent-core-v2/test/features/skill/workspace/skillCatalog.test.ts b/packages/agent-core-v2/test/features/skill/workspace/skillCatalog.test.ts index dfaf12126d2..bb9d20e9ac7 100644 --- a/packages/agent-core-v2/test/features/skill/workspace/skillCatalog.test.ts +++ b/packages/agent-core-v2/test/features/skill/workspace/skillCatalog.test.ts @@ -185,7 +185,7 @@ function makeHost( }); const disposeHost = host.dispose.bind(host); host.dispose = () => { - workspaceHandle.dispose(); + void workspaceHandle.dispose(); disposeHost(); }; return { host, workspace: workspaceHandle, config }; diff --git a/packages/agent-core-v2/test/features/todo/sessionTodo.test.ts b/packages/agent-core-v2/test/features/todo/sessionTodo.test.ts index 2fe14e2939a..83bd077eb17 100644 --- a/packages/agent-core-v2/test/features/todo/sessionTodo.test.ts +++ b/packages/agent-core-v2/test/features/todo/sessionTodo.test.ts @@ -166,7 +166,7 @@ function makeRuntimeAgent( registry.untrack(managed); await managed.runtimeSet.close(); managed.killSpace(); - handle.dispose(); + await handle.dispose(); }, }; } @@ -248,7 +248,7 @@ describe('TodoAgentRuntime', () => { expect(managed.runtimeSet.resolve(AgentTodo)).toBe(todo); expect(reminders).toBe(1); await managed.runtimeSet.close(); - handle.dispose(); + await handle.dispose(); }); it('rejects resolve and lease tracking once the runtime set is closed', async () => { @@ -417,7 +417,7 @@ describe('TodoAgentRuntime', () => { expect(creates).toBe(1); expect(managed.runtimeSet.inspect()[0]).toMatchObject({ status: 'materialized' }); await managed.runtimeSet.close(); - handle.dispose(); + await handle.dispose(); }); it('reports registered, materialized, retired, and definition generations', async () => { @@ -467,7 +467,7 @@ describe('TodoAgentRuntime', () => { managed.attachDurableRuntimes(); expect(managed.runtimeSet.inspect()[0]).toMatchObject({ status: 'materialized', state: [] }); await managed.runtimeSet.close(); - handle.dispose(); + await handle.dispose(); }); it('retains actor failure status and inspection diagnostics', async () => { @@ -518,7 +518,7 @@ describe('TodoAgentRuntime', () => { error: 'actor failed', }); await managed.runtimeSet.close(); - handle.dispose(); + await handle.dispose(); }); }); diff --git a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts index e4dafc94d3f..59461bb0028 100644 --- a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts +++ b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts @@ -2,6 +2,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { SyncDescriptor } from '#/_base/di/descriptors'; import { Disposable, DisposableStore } from '#/_base/di/lifecycle'; +import { IInstantiationService } from '#/_base/di/instantiation'; +import { InstantiationService } from '#/_base/di/instantiationService'; import { LifecycleScope } from '#/app/scopes'; import { type ISessionScopeHandle } from '#/_base/di/scope'; import { TestInstantiationService } from '#/_base/di/test'; @@ -29,7 +31,7 @@ import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory' import '#/agent/contextMemory/contextMemoryService'; import { INHERITED_IN_FLIGHT_TOOL_OUTPUT } from '#/agent/contextMemory/openToolExchange'; import type { ContextMessage } from '#/agent/contextMemory/types'; -import { agentContextOf } from '#/agent/scopeContext/scopeContext'; +import { agentContextOf, IAgentScopeContext } from '#/agent/scopeContext/scopeContext'; import { IAgentIdentity } from '#/app/agentIdentity/agentIdentity'; import { IBuiltinAgentProfileLoader } from '#/app/agentProfileCatalog/builtinAgentProfileLoader'; import { IModelCatalog } from '#/kosong/model/catalog'; @@ -70,6 +72,7 @@ import { IConfigService } from '#/app/config/config'; import { ISessionEventBus } from '#/app/event/eventBus'; import { EventBusService } from '#/app/event/eventBusService'; import '#/app/event/eventBusService'; +import { AgentActivityUpdated } from '#/agent/activityView/activityView'; import { IAgentBlobService } from '#/agent/blob/agentBlobService'; import { IAgentPluginService } from '#/agent/plugin/agentPlugin'; import { ILogService } from '#/_base/log/log'; @@ -97,6 +100,7 @@ import '#/agent/toolActivation/toolActivationService'; import { IAgentMediaToolsRegistrar } from '#/agent/media/mediaTools'; import { ISessionWorkspaceContext } from '#/session/workspaceContext/workspaceContext'; import { FakeRuntime } from '#/runtime/fakeRuntime'; +import { ScopeUnits } from '#/_base/di/fiber'; import { IRuntimeResolver, IWorkspaceInstanceManager, @@ -503,6 +507,179 @@ describe('AgentLifecycleService', () => { expect(svc.handleOf('main')).toBeUndefined(); }); + it('remove keeps the lifecycle context active through async scope teardown', async () => { + const svc = ix.get(IAgentLifecycleService); + const bus = ix.get(ISessionEventBus); + const main = await svc.create({ agentId: 'main' }); + const seen: string[] = []; + disposables.add(bus.subscribe(AgentActivityUpdated, (event) => seen.push(event.lifecycle))); + const agentScope = ix.children.find((child) => child.debugLabel === 'main'); + expect(agentScope).toBeDefined(); + let releaseDrain!: () => void; + let gateEntered!: () => void; + const entered = new Promise((resolve) => { + gateEntered = resolve; + }); + agentScope!.anchorKernelEntry(() => { + gateEntered(); + return new Promise((resolve) => { + releaseDrain = resolve; + }); + }, 'test-async-disposer'); + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown): void => { + unhandled.push(reason); + }; + process.on('unhandledRejection', onUnhandled); + try { + const removal = svc.remove(main); + await entered; + bus.publish( + new AgentActivityUpdated({ lifecycle: 'disposed', background: [], agentId: 'main' }), + main, + ); + expect(seen).toEqual(['disposed']); + releaseDrain(); + await removal; + expect(() => + bus.publish( + new AgentActivityUpdated({ lifecycle: 'disposed', background: [], agentId: 'main' }), + main, + ), + ).toThrow("Agent event 'agent.activity.updated' has no active lifecycle context"); + expect(unhandled).toEqual([]); + } finally { + process.off('unhandledRejection', onUnhandled); + } + }); + + function contributeDisposeBeacon( + dispose: (eventBus: ISessionEventBus, scope: IAgentScopeContext) => void | Promise, + ): void { + class DisposeBeacon { + constructor( + @ISessionEventBus private readonly eventBus: ISessionEventBus, + @IAgentScopeContext private readonly scope: IAgentScopeContext, + ) {} + + dispose(): void | Promise { + return dispose(this.eventBus, this.scope); + } + } + + ix.fiberHost.addCollectionRecord( + ScopeUnits(LifecycleScope.Agent), + 'test', + new Ledger('test'), + DisposeBeacon, + ); + } + + function publishDisposed(eventBus: ISessionEventBus, scope: IAgentScopeContext): void { + eventBus.publish( + new AgentActivityUpdated({ + lifecycle: 'disposed', + background: [], + agentId: scope.agentId, + }), + scope.agentContext, + ); + } + + it('remove deactivates after scope-units contributed units are torn down', async () => { + const svc = ix.get(IAgentLifecycleService); + const bus = ix.get(ISessionEventBus); + const seen: string[] = []; + disposables.add(bus.subscribe(AgentActivityUpdated, (event) => seen.push(event.lifecycle))); + + contributeDisposeBeacon(publishDisposed); + + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown): void => { + unhandled.push(reason); + }; + process.on('unhandledRejection', onUnhandled); + try { + const main = await svc.create({ agentId: 'main' }); + await svc.remove(main); + expect(seen).toEqual(['disposed']); + expect(unhandled).toEqual([]); + } finally { + process.off('unhandledRejection', onUnhandled); + } + }); + + it('create failure after scope creation keeps the context active through async teardown', async () => { + registerAgent.mockRejectedValueOnce(new Error('boom')); + const svc = ix.get(IAgentLifecycleService); + const bus = ix.get(ISessionEventBus); + const seen: string[] = []; + disposables.add(bus.subscribe(AgentActivityUpdated, (event) => seen.push(event.lifecycle))); + + class GatedBeacon { + constructor( + @ISessionEventBus private readonly eventBus: ISessionEventBus, + @IAgentScopeContext private readonly scope: IAgentScopeContext, + @IInstantiationService instantiation: IInstantiationService, + ) { + (instantiation as InstantiationService).anchorKernelEntry( + () => new Promise((resolve) => setTimeout(resolve, 20)), + 'beacon-gate', + ); + } + + dispose(): void { + publishDisposed(this.eventBus, this.scope); + } + } + + ix.fiberHost.addCollectionRecord( + ScopeUnits(LifecycleScope.Agent), + 'test', + new Ledger('test'), + GatedBeacon, + ); + + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown): void => { + unhandled.push(reason); + }; + process.on('unhandledRejection', onUnhandled); + try { + await expect(svc.create({ agentId: 'main' })).rejects.toThrow('boom'); + expect(seen).toEqual(['disposed']); + expect(unhandled).toEqual([]); + } finally { + process.off('unhandledRejection', onUnhandled); + } + }); + + it('remove awaits asynchronous contributed-unit teardown before deactivating', async () => { + const svc = ix.get(IAgentLifecycleService); + const bus = ix.get(ISessionEventBus); + const seen: string[] = []; + disposables.add(bus.subscribe(AgentActivityUpdated, (event) => seen.push(event.lifecycle))); + + contributeDisposeBeacon(async (eventBus, scope) => { + await new Promise((resolve) => setTimeout(resolve, 0)); + publishDisposed(eventBus, scope); + }); + + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown): void => { + unhandled.push(reason); + }; + process.on('unhandledRejection', onUnhandled); + try { + const main = await svc.create({ agentId: 'main' }); + await svc.remove(main); + expect(seen).toEqual(['disposed']); + expect(unhandled).toEqual([]); + } finally { + process.off('unhandledRejection', onUnhandled); + } + }); + it('remove stops the agent background tasks before disposal', async () => { const svc = ix.get(IAgentLifecycleService); const main = await svc.create({ agentId: 'main' }); @@ -1179,7 +1356,7 @@ describe('AgentLifecycleService', () => { const originalDispose = handle.dispose.bind(handle); handle.dispose = () => { order.push('scope-disposed'); - originalDispose(); + return originalDispose(); }; svc.resolve(main, AgentTodo).get(); expect(reminders).toBe(1); diff --git a/packages/agent-core-v2/test/wire/resume.test.ts b/packages/agent-core-v2/test/wire/resume.test.ts index 6db7748e5f3..5f1d72a50cf 100644 --- a/packages/agent-core-v2/test/wire/resume.test.ts +++ b/packages/agent-core-v2/test/wire/resume.test.ts @@ -130,6 +130,40 @@ describe('Agent resume', () => { } }); + it('resumes a journal that stops mid-turn without any turn-closing record', async () => { + const persistence = new RecordingAgentPersistence([ + resumeConfigRecord(), + contextAppendRecord(0, [{ role: 'user', text: 'dangling turn', origin: { kind: 'user' } }]), + turnPromptRecord(0, { kind: 'user' }), + { + type: 'context.append_loop_event', + event: { type: 'step.begin', uuid: 'step-0', turnId: '0', step: 1 }, + }, + ] as unknown as WireRecord[]); + const ctx = testAgent({ persistence, autoConfigure: false }); + + try { + await ctx.restorePersisted(); + + expect(ctx.llmCalls).toHaveLength(0); + expect(turnCurrentId(ctx)).toBe(0); + + ctx.mockNextResponse({ type: 'text', text: 'Fresh response after resume.' }); + await ctx.rpc.prompt({ input: [{ type: 'text', text: 'Fresh prompt after resume' }] }); + await ctx.untilTurnEnd(); + + expect(findRpcEvent(ctx.allEvents, 'turn.started')?.args).toMatchObject({ turnId: 1 }); + expect(findRpcEvent(ctx.allEvents, 'turn.ended')?.args).toMatchObject({ + turnId: 1, + reason: 'completed', + }); + expect(findRpcEvent(ctx.allEvents, 'error')).toBeUndefined(); + await ctx.expectResumeMatches(); + } finally { + await ctx.dispose(); + } + }); + it('does not reconcile a legacy interruption whose delivery was recorded', async () => { const persistence = new RecordingAgentPersistence([ resumeConfigRecord(), diff --git a/packages/kap-server/src/start.ts b/packages/kap-server/src/start.ts index 586e7fb5fd2..9a97059c90c 100644 --- a/packages/kap-server/src/start.ts +++ b/packages/kap-server/src/start.ts @@ -228,6 +228,19 @@ export async function startServer(opts: ServerStartOptions): Promise { + logger.error( + { err: reason instanceof Error ? reason : new Error(String(reason)) }, + 'unhandledRejection', + ); + }; + const onUncaughtException = (err: unknown): void => { + logger.fatal( + { err: err instanceof Error ? err : new Error(String(err)) }, + 'uncaughtException', + ); + process.exit(1); + }; const authFailureLimiter = exposureClass === 'loopback' ? undefined : createAuthFailureLimiter({ logger }); @@ -371,7 +384,12 @@ export async function startServer(opts: ServerStartOptions): Promise { expect(() => core.accessor.get(IBootstrapService)).toThrow(); expect(await listLiveServerInstances(home)).toEqual([]); }); + + it('installs process-level rejection handlers while running and removes them on close', async () => { + home = await mkdtemp(join(tmpdir(), 'kimi-server-v2-')); + const rejectionBefore = process.listenerCount('unhandledRejection'); + const exceptionBefore = process.listenerCount('uncaughtException'); + server = await startServer({ + hostIdentity: TEST_HOST_IDENTITY, + host: '127.0.0.1', + port: 0, + homeDir: home, + logLevel: 'silent', + }); + + expect(process.listenerCount('unhandledRejection')).toBe(rejectionBefore + 1); + expect(process.listenerCount('uncaughtException')).toBe(exceptionBefore + 1); + + await server.close(); + server = undefined; + + expect(process.listenerCount('unhandledRejection')).toBe(rejectionBefore); + expect(process.listenerCount('uncaughtException')).toBe(exceptionBefore); + }); + + it('does not leave process handlers installed when startup fails', async () => { + home = await mkdtemp(join(tmpdir(), 'kimi-server-v2-')); + const emptyAssets = await mkdtemp(join(tmpdir(), 'kimi-server-v2-assets-')); + const rejectionBefore = process.listenerCount('unhandledRejection'); + const exceptionBefore = process.listenerCount('uncaughtException'); + try { + await expect( + startServer({ + hostIdentity: TEST_HOST_IDENTITY, + host: '127.0.0.1', + port: 0, + homeDir: home, + logLevel: 'silent', + webAssetsDir: emptyAssets, + }), + ).rejects.toThrow('web assets'); + expect(process.listenerCount('unhandledRejection')).toBe(rejectionBefore); + expect(process.listenerCount('uncaughtException')).toBe(exceptionBefore); + } finally { + await rm(emptyAssets, { recursive: true, force: true }); + } + }); }); function silentLogger() {