mirror of
https://github.com/davidkaya/aryx.git
synced 2026-08-29 14:07:13 +02:00
fix: harden activity event handling
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
@@ -0,0 +1,28 @@
|
|||||||
|
import type { AgentActivityEvent, TurnDeltaEvent } from '@shared/contracts/sidecar';
|
||||||
|
import type { ChatMessageRecord } from '@shared/domain/session';
|
||||||
|
|
||||||
|
export interface RunTurnPendingCommand {
|
||||||
|
kind: 'run-turn';
|
||||||
|
resolve: (messages: ChatMessageRecord[]) => void;
|
||||||
|
reject: (error: Error) => void;
|
||||||
|
onDelta: (event: TurnDeltaEvent) => void | Promise<void>;
|
||||||
|
onActivity: (event: AgentActivityEvent) => void | Promise<void>;
|
||||||
|
errored: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function markRunTurnPendingErrored(
|
||||||
|
pending: RunTurnPendingCommand,
|
||||||
|
error: unknown,
|
||||||
|
): Error {
|
||||||
|
const normalized = error instanceof Error ? error : new Error(String(error));
|
||||||
|
if (!pending.errored) {
|
||||||
|
pending.errored = true;
|
||||||
|
pending.reject(normalized);
|
||||||
|
}
|
||||||
|
|
||||||
|
return normalized;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function shouldHandleRunTurnEvent(pending: RunTurnPendingCommand): boolean {
|
||||||
|
return !pending.errored;
|
||||||
|
}
|
||||||
@@ -12,6 +12,11 @@ import type {
|
|||||||
} from '@shared/contracts/sidecar';
|
} from '@shared/contracts/sidecar';
|
||||||
import type { ChatMessageRecord } from '@shared/domain/session';
|
import type { ChatMessageRecord } from '@shared/domain/session';
|
||||||
import { createSidecarEnvironment } from '@main/sidecar/sidecarEnvironment';
|
import { createSidecarEnvironment } from '@main/sidecar/sidecarEnvironment';
|
||||||
|
import {
|
||||||
|
markRunTurnPendingErrored,
|
||||||
|
shouldHandleRunTurnEvent,
|
||||||
|
type RunTurnPendingCommand,
|
||||||
|
} from '@main/sidecar/runTurnPending';
|
||||||
import { resolveSidecarProcess } from '@main/sidecar/sidecarRuntime';
|
import { resolveSidecarProcess } from '@main/sidecar/sidecarRuntime';
|
||||||
|
|
||||||
type PendingCommand =
|
type PendingCommand =
|
||||||
@@ -25,13 +30,7 @@ type PendingCommand =
|
|||||||
resolve: (issues: ValidatePatternCommand['pattern'] extends never ? never : unknown) => void;
|
resolve: (issues: ValidatePatternCommand['pattern'] extends never ? never : unknown) => void;
|
||||||
reject: (error: Error) => void;
|
reject: (error: Error) => void;
|
||||||
}
|
}
|
||||||
| {
|
| RunTurnPendingCommand;
|
||||||
kind: 'run-turn';
|
|
||||||
resolve: (messages: ChatMessageRecord[]) => void;
|
|
||||||
reject: (error: Error) => void;
|
|
||||||
onDelta: (event: TurnDeltaEvent) => void | Promise<void>;
|
|
||||||
onActivity: (event: AgentActivityEvent) => void | Promise<void>;
|
|
||||||
};
|
|
||||||
|
|
||||||
export class SidecarClient {
|
export class SidecarClient {
|
||||||
private process?: ChildProcessWithoutNullStreams;
|
private process?: ChildProcessWithoutNullStreams;
|
||||||
@@ -129,6 +128,7 @@ export class SidecarClient {
|
|||||||
reject,
|
reject,
|
||||||
onDelta: onDelta ?? (() => undefined),
|
onDelta: onDelta ?? (() => undefined),
|
||||||
onActivity: onActivity ?? (() => undefined),
|
onActivity: onActivity ?? (() => undefined),
|
||||||
|
errored: false,
|
||||||
});
|
});
|
||||||
} else if (command.type === 'validate-pattern') {
|
} else if (command.type === 'validate-pattern') {
|
||||||
this.pending.set(command.requestId, {
|
this.pending.set(command.requestId, {
|
||||||
@@ -183,27 +183,33 @@ export class SidecarClient {
|
|||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
case 'turn-delta':
|
case 'turn-delta':
|
||||||
if (pending.kind === 'run-turn') {
|
if (pending.kind === 'run-turn' && shouldHandleRunTurnEvent(pending)) {
|
||||||
this.invokeRunTurnHandler(event.requestId, pending, () => pending.onDelta(event));
|
this.invokeRunTurnHandler(event.requestId, pending, () => pending.onDelta(event));
|
||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
case 'agent-activity':
|
case 'agent-activity':
|
||||||
if (pending.kind === 'run-turn') {
|
if (pending.kind === 'run-turn' && shouldHandleRunTurnEvent(pending)) {
|
||||||
this.invokeRunTurnHandler(event.requestId, pending, () => pending.onActivity(event));
|
this.invokeRunTurnHandler(event.requestId, pending, () => pending.onActivity(event));
|
||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
case 'turn-complete':
|
case 'turn-complete':
|
||||||
if (pending.kind === 'run-turn') {
|
if (pending.kind === 'run-turn') {
|
||||||
pending.resolve(event.messages);
|
if (shouldHandleRunTurnEvent(pending)) {
|
||||||
|
pending.resolve(event.messages);
|
||||||
|
}
|
||||||
this.pending.delete(event.requestId);
|
this.pending.delete(event.requestId);
|
||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
case 'command-error':
|
case 'command-error':
|
||||||
pending.reject(new Error(event.message));
|
if (pending.kind === 'run-turn') {
|
||||||
|
markRunTurnPendingErrored(pending, new Error(event.message));
|
||||||
|
} else {
|
||||||
|
pending.reject(new Error(event.message));
|
||||||
|
}
|
||||||
this.pending.delete(event.requestId);
|
this.pending.delete(event.requestId);
|
||||||
return;
|
return;
|
||||||
case 'command-complete':
|
case 'command-complete':
|
||||||
if (pending.kind !== 'run-turn') {
|
if (pending.kind !== 'run-turn' || pending.errored) {
|
||||||
this.pending.delete(event.requestId);
|
this.pending.delete(event.requestId);
|
||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
@@ -212,12 +218,14 @@ export class SidecarClient {
|
|||||||
|
|
||||||
private invokeRunTurnHandler(
|
private invokeRunTurnHandler(
|
||||||
requestId: string,
|
requestId: string,
|
||||||
pending: Extract<PendingCommand, { kind: 'run-turn' }>,
|
pending: RunTurnPendingCommand,
|
||||||
callback: () => void | Promise<void>,
|
callback: () => void | Promise<void>,
|
||||||
): void {
|
): void {
|
||||||
void Promise.resolve(callback()).catch((error: unknown) => {
|
void Promise.resolve(callback()).catch((error: unknown) => {
|
||||||
this.pending.delete(requestId);
|
markRunTurnPendingErrored(pending, error);
|
||||||
pending.reject(error instanceof Error ? error : new Error(String(error)));
|
if (this.pending.get(requestId) !== pending) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ export function applySessionEventActivity(
|
|||||||
if (event.kind === 'agent-activity') {
|
if (event.kind === 'agent-activity') {
|
||||||
const agentKey = resolveAgentKey(event);
|
const agentKey = resolveAgentKey(event);
|
||||||
if (!agentKey) {
|
if (!agentKey) {
|
||||||
|
console.warn('[kopaya activity] Dropping agent-activity event without agentId/agentName.', event);
|
||||||
return current;
|
return current;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,48 @@
|
|||||||
|
import { describe, expect, test } from 'bun:test';
|
||||||
|
|
||||||
|
import {
|
||||||
|
markRunTurnPendingErrored,
|
||||||
|
shouldHandleRunTurnEvent,
|
||||||
|
type RunTurnPendingCommand,
|
||||||
|
} from '@main/sidecar/runTurnPending';
|
||||||
|
|
||||||
|
describe('run turn pending helpers', () => {
|
||||||
|
test('marks a run-turn pending command as errored and rejects it once', () => {
|
||||||
|
const rejected: Error[] = [];
|
||||||
|
const pending: RunTurnPendingCommand = {
|
||||||
|
kind: 'run-turn',
|
||||||
|
resolve: () => undefined,
|
||||||
|
reject: (error) => rejected.push(error),
|
||||||
|
onDelta: () => undefined,
|
||||||
|
onActivity: () => undefined,
|
||||||
|
errored: false,
|
||||||
|
};
|
||||||
|
|
||||||
|
const first = markRunTurnPendingErrored(pending, 'boom');
|
||||||
|
const second = markRunTurnPendingErrored(pending, new Error('later'));
|
||||||
|
|
||||||
|
expect(first).toBeInstanceOf(Error);
|
||||||
|
expect(first.message).toBe('boom');
|
||||||
|
expect(second.message).toBe('later');
|
||||||
|
expect(pending.errored).toBe(true);
|
||||||
|
expect(rejected).toHaveLength(1);
|
||||||
|
expect(rejected[0].message).toBe('boom');
|
||||||
|
});
|
||||||
|
|
||||||
|
test('stops handling turn events after the pending command has errored', () => {
|
||||||
|
const pending: RunTurnPendingCommand = {
|
||||||
|
kind: 'run-turn',
|
||||||
|
resolve: () => undefined,
|
||||||
|
reject: () => undefined,
|
||||||
|
onDelta: () => undefined,
|
||||||
|
onActivity: () => undefined,
|
||||||
|
errored: false,
|
||||||
|
};
|
||||||
|
|
||||||
|
expect(shouldHandleRunTurnEvent(pending)).toBe(true);
|
||||||
|
|
||||||
|
markRunTurnPendingErrored(pending, new Error('boom'));
|
||||||
|
|
||||||
|
expect(shouldHandleRunTurnEvent(pending)).toBe(false);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -83,6 +83,42 @@ describe('session activity helpers', () => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('warns when an agent-activity event is missing identifiers', () => {
|
||||||
|
const originalWarn = console.warn;
|
||||||
|
const warnings: unknown[][] = [];
|
||||||
|
console.warn = (...args: unknown[]) => {
|
||||||
|
warnings.push(args);
|
||||||
|
};
|
||||||
|
|
||||||
|
try {
|
||||||
|
const current: SessionActivityMap = {
|
||||||
|
'session-1': {
|
||||||
|
architect: {
|
||||||
|
agentId: 'architect',
|
||||||
|
agentName: 'Architect',
|
||||||
|
activityType: 'thinking',
|
||||||
|
},
|
||||||
|
},
|
||||||
|
};
|
||||||
|
|
||||||
|
expect(
|
||||||
|
applySessionEventActivity(current, {
|
||||||
|
sessionId: 'session-1',
|
||||||
|
kind: 'agent-activity',
|
||||||
|
occurredAt: '2026-03-23T00:00:00.000Z',
|
||||||
|
activityType: 'thinking',
|
||||||
|
agentId: ' ',
|
||||||
|
agentName: ' ',
|
||||||
|
}),
|
||||||
|
).toBe(current);
|
||||||
|
|
||||||
|
expect(warnings).toHaveLength(1);
|
||||||
|
expect(String(warnings[0][0])).toContain('Dropping agent-activity event');
|
||||||
|
} finally {
|
||||||
|
console.warn = originalWarn;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
test('clears stale activity when a session restarts', () => {
|
test('clears stale activity when a session restarts', () => {
|
||||||
const current: SessionActivityMap = {
|
const current: SessionActivityMap = {
|
||||||
'session-1': {
|
'session-1': {
|
||||||
|
|||||||
Reference in New Issue
Block a user