Skip to content

Commit 2ac8abc

Browse files
committed
fix(agent-core-v2): keep late task settlement silent after agent teardown
A background task can settle after its agent is gone. Gate the task-started, task-terminated and notification paths on the agent still being active, and stop tasks after the loop reaches quiescence rather than before it drains, so a late settle cannot reach the wire.
1 parent ee29130 commit 2ac8abc

7 files changed

Lines changed: 48 additions & 11 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@pymodel/pythinker-code": patch
3+
---
4+
5+
A background task that finishes after its agent is closed no longer emits stray task events.

‎packages/agent-core-v2/src/agent/task/taskService.ts‎

Lines changed: 19 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import {
1414
} from '#/_base/utils/abort';
1515
import { setClampedTimeout } from '#/_base/utils/timer';
1616
import { escapeXml, escapeXmlAttr } from '#/_base/utils/xml-escape';
17-
import { IEventBus } from '#/app/event/eventBus';
17+
import { IEventBus, ISessionEventBus } from '#/app/event/eventBus';
1818
import { Error2, ErrorCodes } from '#/errors';
1919
import { z } from 'zod';
2020
import {
@@ -241,6 +241,7 @@ export class AgentTaskService extends Disposable implements IAgentTaskService {
241241
@IAgentScopeContext private readonly scopeContext: IAgentScopeContext,
242242
@ITaskService private readonly taskService: ITaskService,
243243
@IEventBus private readonly eventBus: IEventBus,
244+
@ISessionEventBus private readonly sessionEventBus: ISessionEventBus,
244245
@IEventDispatcher private readonly dispatcher: IEventDispatcher,
245246
@IAgentLifecycleService private readonly agentLifecycle: IAgentLifecycleService,
246247
@IAgentLoopService private readonly loop: IAgentLoopService,
@@ -853,6 +854,10 @@ export class AgentTaskService extends Disposable implements IAgentTaskService {
853854
return resolveAgentTaskConfig(this.config)?.keepAliveOnExit === true;
854855
}
855856

857+
private lifecycleActive(): boolean {
858+
return this.sessionEventBus.isAgentActive(this.scopeContext.agentContext);
859+
}
860+
856861
async wait(
857862
taskId: string,
858863
timeoutMs = 30_000,
@@ -1085,19 +1090,23 @@ export class AgentTaskService extends Disposable implements IAgentTaskService {
10851090
}
10861091

10871092
private recordTaskStarted(info: AgentTaskInfo): void {
1088-
void this.dispatcher.dispatch(
1089-
new TaskStarted({ agentId: this.scopeContext.agentId, info }),
1090-
);
1093+
if (this.lifecycleActive()) {
1094+
void this.dispatcher.dispatch(
1095+
new TaskStarted({ agentId: this.scopeContext.agentId, info }),
1096+
);
1097+
}
10911098
this.telemetry.track2('background_task_created', {
10921099
task_id: info.taskId,
10931100
kind: info.kind === 'process' ? 'bash' : info.kind,
10941101
});
10951102
}
10961103

10971104
private recordTaskTerminated(info: AgentTaskInfo, outputTail?: string): void {
1098-
void this.dispatcher.dispatch(
1099-
new TaskTerminated({ agentId: this.scopeContext.agentId, info, outputTail }),
1100-
);
1105+
if (this.lifecycleActive()) {
1106+
void this.dispatcher.dispatch(
1107+
new TaskTerminated({ agentId: this.scopeContext.agentId, info, outputTail }),
1108+
);
1109+
}
11011110
this.telemetry.track2('background_task_completed', {
11021111
task_id: info.taskId,
11031112
kind: info.kind,
@@ -1107,8 +1116,10 @@ export class AgentTaskService extends Disposable implements IAgentTaskService {
11071116
}
11081117

11091118
private async notifyAgentTask(info: AgentTaskInfo): Promise<void> {
1119+
if (!this.lifecycleActive()) return;
11101120
const context = await this.buildAgentTaskNotificationContext(info);
11111121
if (context === undefined) return;
1122+
if (!this.lifecycleActive()) return;
11121123
const key = notificationKey(context.origin);
11131124
if (this.deliveredNotificationKeys.has(key)) return;
11141125
const request = new TaskNotificationStepRequest(
@@ -1307,6 +1318,7 @@ export class AgentTaskService extends Disposable implements IAgentTaskService {
13071318
}
13081319

13091320
private fireNotificationHook(notification: AgentTaskNotification): void {
1321+
if (!this.lifecycleActive()) return;
13101322
void this.dispatcher.dispatch(
13111323
new TaskNotified({
13121324
agentId: this.scopeContext.agentId,

‎packages/agent-core-v2/src/app/event/eventBus.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ export const IEventBus: ServiceIdentifier<IEventBus> = createDecorator<IEventBus
1818
export interface ISessionEventBus extends IEventBus {
1919
activateAgent(agent: AgentContext): void;
2020
deactivateAgent(agent: AgentContext): void;
21+
isAgentActive(agent: AgentContext): boolean;
2122
sourceOf(event: Event2<any>): AgentContext | undefined;
2223
onAgent<P extends AgentDomainTrait, E extends Event2<P>>(
2324
agent: AgentContext,

‎packages/agent-core-v2/src/app/event/eventBusService.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,10 @@ export class EventBusService extends Service implements ISessionEventBus {
2525
if (this.agents.get(agent.agentId) === agent) this.agents.delete(agent.agentId);
2626
}
2727

28+
isAgentActive(agent: AgentContext): boolean {
29+
return this.agents.get(agent.agentId) === agent;
30+
}
31+
2832
publish(event: Event2<any>, agent?: AgentContext): void {
2933
const cls = event.constructor as Event2Class;
3034
if (cls.agentDomain) {

‎packages/agent-core-v2/src/session/agentLifecycle/agentLifecycleService.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -468,7 +468,6 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle
468468
const reason = abortError('Agent removed');
469469
let quiescence: IDisposable | undefined;
470470
try {
471-
await phase(() => handle.accessor.get(IAgentTaskService).stopAllOnExit('Session closed'));
472471
const compaction = handle.accessor.get(IAgentFullCompactionService).compacting;
473472
const compactionSettled = compaction?.promise.catch(() => undefined) ?? Promise.resolve();
474473
const prompt = handle.accessor.get(IAgentPromptService);
@@ -486,6 +485,7 @@ export class AgentLifecycleService extends Disposable implements IAgentLifecycle
486485
await phase(() => {
487486
quiescence = loop.tryAcquireQuiescence();
488487
});
488+
await phase(() => handle.accessor.get(IAgentTaskService).stopAllOnExit('Session closed'));
489489
await phase(() => handle.accessor.get(IEventDispatcher).flush());
490490
await phase(() => managed.runtimeSet.close());
491491
managed.killSpace();

‎packages/agent-core-v2/test/agent/task/taskService.test.ts‎

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -637,7 +637,10 @@ describe('AgentTaskService', () => {
637637
expect(forceStop).not.toHaveBeenCalled();
638638
});
639639

640-
it('scope disposal leaves a process running when keepAliveOnExit is set', async () => {
640+
it('scope disposal leaves a process running when keepAliveOnExit is set, and its late settle stays silent after deactivation', async () => {
641+
const { records } = capturingWire();
642+
const track2 = vi.fn();
643+
ix.stub(ITelemetryService, { track2 });
641644
stubTaskConfig({ keepAliveOnExit: true });
642645
const stdout = new Readable({ read() {} });
643646
const stderr = new Readable({ read() {} });
@@ -656,7 +659,8 @@ describe('AgentTaskService', () => {
656659
dispose: vi.fn().mockResolvedValue(undefined),
657660
} as unknown as IHostProcess;
658661
const svc = ix.get(IAgentTaskService);
659-
svc.registerTask(new ProcessTask(proc, 'keep-running', 'long-running process'));
662+
const taskId = svc.registerTask(new ProcessTask(proc, 'keep-running', 'long-running process'));
663+
const agentContext = ix.get(IAgentScopeContext).agentContext;
660664
await Promise.resolve();
661665

662666
disposables.dispose();
@@ -665,10 +669,18 @@ describe('AgentTaskService', () => {
665669
expect(proc.kill).not.toHaveBeenCalled();
666670
expect(proc.dispose).not.toHaveBeenCalled();
667671

672+
eventBus.deactivateAgent(agentContext);
668673
stdout.push(null);
669674
stderr.push(null);
670675
resolveWait(0);
671-
await Promise.resolve();
676+
await waitForCondition(() => svc.getTask(taskId)?.status === 'completed');
677+
678+
expect(svc.getTask(taskId)?.status).toBe('completed');
679+
expect(records.filter((record) => record['type'] === 'task.terminated')).toHaveLength(0);
680+
expect(track2.mock.calls.map(([event]) => event)).toEqual([
681+
'background_task_created',
682+
'background_task_completed',
683+
]);
672684
});
673685

674686
it('stop requests force-stop when killGracePeriodMs is zero', async () => {

‎packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -783,6 +783,9 @@ describe('AgentLifecycleService', () => {
783783

784784
expect(stopAllOnExit).toHaveBeenCalledWith('Session closed');
785785
expect(promptDrain).toHaveBeenCalledOnce();
786+
expect(stopAllOnExit.mock.invocationCallOrder[0]).toBeGreaterThan(
787+
promptDrain.mock.invocationCallOrder[0]!,
788+
);
786789
});
787790

788791
it('remove waits for prompt intake to drain before disposing the agent scope', async () => {

0 commit comments

Comments
 (0)