fix(companion): reject pending event waits on teardown
This commit is contained in:
@@ -1190,7 +1190,7 @@ Methods:
|
|||||||
- `system.capabilities` returns gateway protocol and node policy snapshot.
|
- `system.capabilities` returns gateway protocol and node policy snapshot.
|
||||||
|
|
||||||
Companion runtime helper:
|
Companion runtime helper:
|
||||||
- `src/companion/runtimeClient.ts` provides a typed Node/WebSocket client for companion runtimes (macOS/iOS/Android workers) with wrappers for `node.register`, `node.capabilities.get`, `node.location.set/get`, `node.status.set`, `node.push_token.set`, `system.capabilities`, `system.nodes`, and canvas artifact RPCs (`canvas.put/get/list/delete/clear`), plus convenience helpers (`bootstrapNode`, optional `autoConnect`, `dispose()`) and event helpers (`subscribeEvents()`, `subscribeEvent()`, `subscribeAgentStream()`, `subscribeAgentTyping()`, `subscribeContextWarning()`, `waitForEvent()` with timeout/predicate/abort support, `waitForAgentStream()`, `waitForAgentTyping()`, `waitForContextWarning()`, `clearEventSubscriptions()`).
|
- `src/companion/runtimeClient.ts` provides a typed Node/WebSocket client for companion runtimes (macOS/iOS/Android workers) with wrappers for `node.register`, `node.capabilities.get`, `node.location.set/get`, `node.status.set`, `node.push_token.set`, `system.capabilities`, `system.nodes`, and canvas artifact RPCs (`canvas.put/get/list/delete/clear`), plus convenience helpers (`bootstrapNode`, optional `autoConnect`, `dispose()`) and event helpers (`subscribeEvents()`, `subscribeEvent()`, `subscribeAgentStream()`, `subscribeAgentTyping()`, `subscribeContextWarning()`, `waitForEvent()` with timeout/predicate/abort support and deterministic teardown cancellation, `waitForAgentStream()`, `waitForAgentTyping()`, `waitForContextWarning()`, `clearEventSubscriptions()`).
|
||||||
- `src/companion/platformClients.ts` provides platform-focused wrappers:
|
- `src/companion/platformClients.ts` provides platform-focused wrappers:
|
||||||
- `MacOSCompanionClient` (`platform: "macos"`, APNs push registration)
|
- `MacOSCompanionClient` (`platform: "macos"`, APNs push registration)
|
||||||
- `IOSCompanionClient` (`platform: "ios"`, APNs push registration)
|
- `IOSCompanionClient` (`platform: "ios"`, APNs push registration)
|
||||||
|
|||||||
@@ -580,6 +580,19 @@
|
|||||||
],
|
],
|
||||||
"test_status": "pnpm test:run src/companion/platformClients.test.ts src/companion/runtimeClient.test.ts src/companion/heartbeatLoop.test.ts src/companion/platformClients.integration.test.ts + pnpm typecheck passing"
|
"test_status": "pnpm test:run src/companion/platformClients.test.ts src/companion/runtimeClient.test.ts src/companion/heartbeatLoop.test.ts src/companion/platformClients.integration.test.ts + pnpm typecheck passing"
|
||||||
},
|
},
|
||||||
|
"companion-runtime-waiter-teardown-rejection": {
|
||||||
|
"status": "completed",
|
||||||
|
"date": "2026-02-17",
|
||||||
|
"updated": "2026-02-17",
|
||||||
|
"summary": "Hardened `waitForEvent()` lifecycle semantics by rejecting pending waiters immediately on teardown paths (`disconnect`, `dispose`, `clearEventSubscriptions`) instead of waiting for timeout.",
|
||||||
|
"files_modified": [
|
||||||
|
"src/companion/runtimeClient.ts",
|
||||||
|
"src/companion/runtimeClient.test.ts",
|
||||||
|
"README.md",
|
||||||
|
"docs/plans/state.json"
|
||||||
|
],
|
||||||
|
"test_status": "pnpm test:run src/companion/runtimeClient.test.ts src/companion/platformClients.test.ts src/companion/heartbeatLoop.test.ts src/companion/platformClients.integration.test.ts + pnpm typecheck passing"
|
||||||
|
},
|
||||||
"browser-tools-activation-clarity": {
|
"browser-tools-activation-clarity": {
|
||||||
"status": "completed",
|
"status": "completed",
|
||||||
"date": "2026-02-17",
|
"date": "2026-02-17",
|
||||||
|
|||||||
@@ -339,6 +339,30 @@ describe('CompanionRuntimeClient', () => {
|
|||||||
await awaited;
|
await awaited;
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('waitForEvent rejects immediately when event subscriptions are cleared', async () => {
|
||||||
|
const client = new CompanionRuntimeClient({
|
||||||
|
url: 'ws://127.0.0.1:1',
|
||||||
|
});
|
||||||
|
|
||||||
|
const awaited = expect(
|
||||||
|
client.waitForEvent('agent.stream', { timeoutMs: 10_000 }),
|
||||||
|
).rejects.toThrow('Event subscriptions cleared');
|
||||||
|
client.clearEventSubscriptions();
|
||||||
|
await awaited;
|
||||||
|
});
|
||||||
|
|
||||||
|
it('waitForEvent rejects immediately on disconnect', async () => {
|
||||||
|
const client = new CompanionRuntimeClient({
|
||||||
|
url: 'ws://127.0.0.1:1',
|
||||||
|
});
|
||||||
|
|
||||||
|
const awaited = expect(
|
||||||
|
client.waitForEvent('agent.stream', { timeoutMs: 10_000 }),
|
||||||
|
).rejects.toThrow('Disconnected');
|
||||||
|
client.disconnect();
|
||||||
|
await awaited;
|
||||||
|
});
|
||||||
|
|
||||||
it('waitForAgentStream resolves on agent.stream events', async () => {
|
it('waitForAgentStream resolves on agent.stream events', async () => {
|
||||||
const client = new CompanionRuntimeClient({
|
const client = new CompanionRuntimeClient({
|
||||||
url: 'ws://127.0.0.1:1',
|
url: 'ws://127.0.0.1:1',
|
||||||
|
|||||||
@@ -274,6 +274,7 @@ export class CompanionRuntimeClient {
|
|||||||
private nextId = 1;
|
private nextId = 1;
|
||||||
private pending = new Map<number, PendingRequest>();
|
private pending = new Map<number, PendingRequest>();
|
||||||
private readonly eventHandlers = new Set<CompanionEventHandler>();
|
private readonly eventHandlers = new Set<CompanionEventHandler>();
|
||||||
|
private readonly pendingEventWaits = new Set<(error: Error) => void>();
|
||||||
|
|
||||||
constructor(options: CompanionRuntimeClientOptions) {
|
constructor(options: CompanionRuntimeClientOptions) {
|
||||||
const requestTimeoutMs = options.requestTimeoutMs ?? 15_000;
|
const requestTimeoutMs = options.requestTimeoutMs ?? 15_000;
|
||||||
@@ -354,12 +355,14 @@ export class CompanionRuntimeClient {
|
|||||||
|
|
||||||
disconnect(code?: number, reason?: string): void {
|
disconnect(code?: number, reason?: string): void {
|
||||||
if (!this.ws) {
|
if (!this.ws) {
|
||||||
|
this.rejectEventWaits(new Error('Disconnected'));
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
const ws = this.ws;
|
const ws = this.ws;
|
||||||
this.ws = null;
|
this.ws = null;
|
||||||
this.rejectAllPending(new Error('Disconnected'));
|
this.rejectAllPending(new Error('Disconnected'));
|
||||||
|
this.rejectEventWaits(new Error('Disconnected'));
|
||||||
ws.close(code, reason);
|
ws.close(code, reason);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -377,6 +380,7 @@ export class CompanionRuntimeClient {
|
|||||||
|
|
||||||
clearEventSubscriptions(): void {
|
clearEventSubscriptions(): void {
|
||||||
this.eventHandlers.clear();
|
this.eventHandlers.clear();
|
||||||
|
this.rejectEventWaits(new Error('Event subscriptions cleared'));
|
||||||
}
|
}
|
||||||
|
|
||||||
subscribeEvent<TData = unknown>(
|
subscribeEvent<TData = unknown>(
|
||||||
@@ -436,9 +440,15 @@ export class CompanionRuntimeClient {
|
|||||||
abortCleanup();
|
abortCleanup();
|
||||||
abortCleanup = null;
|
abortCleanup = null;
|
||||||
}
|
}
|
||||||
|
this.pendingEventWaits.delete(cancelWait);
|
||||||
fn();
|
fn();
|
||||||
};
|
};
|
||||||
|
|
||||||
|
const cancelWait = (error: Error) => {
|
||||||
|
finish(() => reject(error));
|
||||||
|
};
|
||||||
|
this.pendingEventWaits.add(cancelWait);
|
||||||
|
|
||||||
const unsubscribe = this.subscribeEvent<TData>(eventName, (data) => {
|
const unsubscribe = this.subscribeEvent<TData>(eventName, (data) => {
|
||||||
if (predicate && !predicate(data)) {
|
if (predicate && !predicate(data)) {
|
||||||
return;
|
return;
|
||||||
@@ -452,7 +462,7 @@ export class CompanionRuntimeClient {
|
|||||||
|
|
||||||
if (signal) {
|
if (signal) {
|
||||||
const onAbort = () => {
|
const onAbort = () => {
|
||||||
finish(() => reject(new Error(`Aborted while waiting for event ${eventName}`)));
|
cancelWait(new Error(`Aborted while waiting for event ${eventName}`));
|
||||||
};
|
};
|
||||||
signal.addEventListener('abort', onAbort, { once: true });
|
signal.addEventListener('abort', onAbort, { once: true });
|
||||||
abortCleanup = () => {
|
abortCleanup = () => {
|
||||||
@@ -697,6 +707,13 @@ export class CompanionRuntimeClient {
|
|||||||
}
|
}
|
||||||
this.pending.clear();
|
this.pending.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private rejectEventWaits(error: Error): void {
|
||||||
|
for (const cancel of this.pendingEventWaits) {
|
||||||
|
cancel(error);
|
||||||
|
}
|
||||||
|
this.pendingEventWaits.clear();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function withToken(url: string, token?: string): string {
|
function withToken(url: string, token?: string): string {
|
||||||
|
|||||||
Reference in New Issue
Block a user