feat(companion): add runtime event subscription hooks
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 optional `autoConnect` mode for one-shot RPC calls.
|
- `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 optional `autoConnect` mode and `subscribeEvents()` for gateway stream events.
|
||||||
- `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)
|
||||||
|
|||||||
@@ -299,6 +299,20 @@
|
|||||||
],
|
],
|
||||||
"test_status": "pnpm test:run src/companion/heartbeatLoop.test.ts src/companion/platformClients.test.ts + pnpm typecheck passing"
|
"test_status": "pnpm test:run src/companion/heartbeatLoop.test.ts src/companion/platformClients.test.ts + pnpm typecheck passing"
|
||||||
},
|
},
|
||||||
|
"companion-runtime-event-subscriptions": {
|
||||||
|
"status": "completed",
|
||||||
|
"date": "2026-02-17",
|
||||||
|
"updated": "2026-02-17",
|
||||||
|
"summary": "Added event subscription support to `CompanionRuntimeClient` via `subscribeEvents()` so companion runtimes can receive gateway stream events with safe callback isolation and explicit unsubscribe flow.",
|
||||||
|
"files_modified": [
|
||||||
|
"src/companion/runtimeClient.ts",
|
||||||
|
"src/companion/runtimeClient.test.ts",
|
||||||
|
"src/companion/index.ts",
|
||||||
|
"README.md",
|
||||||
|
"docs/plans/state.json"
|
||||||
|
],
|
||||||
|
"test_status": "pnpm test:run src/companion/runtimeClient.test.ts src/companion/platformClients.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",
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ export { CompanionHeartbeatLoop } from './heartbeatLoop.js';
|
|||||||
|
|
||||||
export type {
|
export type {
|
||||||
CompanionRuntimeClientOptions,
|
CompanionRuntimeClientOptions,
|
||||||
|
CompanionEventHandler,
|
||||||
RegisterNodeInput,
|
RegisterNodeInput,
|
||||||
ListNodesInput,
|
ListNodesInput,
|
||||||
SetNodeStatusInput,
|
SetNodeStatusInput,
|
||||||
|
|||||||
@@ -96,6 +96,61 @@ afterAll(async () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
describe('CompanionRuntimeClient', () => {
|
describe('CompanionRuntimeClient', () => {
|
||||||
|
it('dispatches gateway events to subscribed handlers and supports unsubscribe', () => {
|
||||||
|
const client = new CompanionRuntimeClient({
|
||||||
|
url: 'ws://127.0.0.1:1',
|
||||||
|
});
|
||||||
|
const handler = vi.fn();
|
||||||
|
const unsubscribe = client.subscribeEvents(handler);
|
||||||
|
|
||||||
|
(client as unknown as { handleMessage: (raw: string) => void }).handleMessage(
|
||||||
|
JSON.stringify({
|
||||||
|
id: 42,
|
||||||
|
event: 'agent.stream',
|
||||||
|
data: { token: 'hello' },
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(handler).toHaveBeenCalledWith('agent.stream', { token: 'hello' });
|
||||||
|
|
||||||
|
unsubscribe();
|
||||||
|
|
||||||
|
(client as unknown as { handleMessage: (raw: string) => void }).handleMessage(
|
||||||
|
JSON.stringify({
|
||||||
|
id: 43,
|
||||||
|
event: 'agent.stream',
|
||||||
|
data: { token: 'world' },
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(handler).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('isolates subscriber callback failures', () => {
|
||||||
|
const client = new CompanionRuntimeClient({
|
||||||
|
url: 'ws://127.0.0.1:1',
|
||||||
|
});
|
||||||
|
const badHandler = vi.fn(() => {
|
||||||
|
throw new Error('subscriber failed');
|
||||||
|
});
|
||||||
|
const goodHandler = vi.fn();
|
||||||
|
client.subscribeEvents(badHandler);
|
||||||
|
client.subscribeEvents(goodHandler);
|
||||||
|
|
||||||
|
expect(() => {
|
||||||
|
(client as unknown as { handleMessage: (raw: string) => void }).handleMessage(
|
||||||
|
JSON.stringify({
|
||||||
|
id: 44,
|
||||||
|
event: 'agent.stream',
|
||||||
|
data: { token: 'safe' },
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
}).not.toThrow();
|
||||||
|
|
||||||
|
expect(badHandler).toHaveBeenCalledOnce();
|
||||||
|
expect(goodHandler).toHaveBeenCalledWith('agent.stream', { token: 'safe' });
|
||||||
|
});
|
||||||
|
|
||||||
it('connects and performs node registration + capability discovery', async () => {
|
it('connects and performs node registration + capability discovery', async () => {
|
||||||
if (!LISTEN_ALLOWED) {
|
if (!LISTEN_ALLOWED) {
|
||||||
return;
|
return;
|
||||||
|
|||||||
@@ -41,6 +41,8 @@ export interface CompanionRuntimeClientOptions {
|
|||||||
websocketFactory?: (url: string) => WebSocket;
|
websocketFactory?: (url: string) => WebSocket;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type CompanionEventHandler = (event: string, data: unknown) => void;
|
||||||
|
|
||||||
export interface RegisterNodeInput {
|
export interface RegisterNodeInput {
|
||||||
nodeId: string;
|
nodeId: string;
|
||||||
role: string;
|
role: string;
|
||||||
@@ -257,6 +259,7 @@ export class CompanionRuntimeClient {
|
|||||||
private connectPromise: Promise<void> | null = null;
|
private connectPromise: Promise<void> | null = null;
|
||||||
private nextId = 1;
|
private nextId = 1;
|
||||||
private pending = new Map<number, PendingRequest>();
|
private pending = new Map<number, PendingRequest>();
|
||||||
|
private readonly eventHandlers = new Set<CompanionEventHandler>();
|
||||||
|
|
||||||
constructor(options: CompanionRuntimeClientOptions) {
|
constructor(options: CompanionRuntimeClientOptions) {
|
||||||
this.url = options.url;
|
this.url = options.url;
|
||||||
@@ -342,6 +345,13 @@ export class CompanionRuntimeClient {
|
|||||||
ws.close(code, reason);
|
ws.close(code, reason);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
subscribeEvents(handler: CompanionEventHandler): () => void {
|
||||||
|
this.eventHandlers.add(handler);
|
||||||
|
return () => {
|
||||||
|
this.eventHandlers.delete(handler);
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
async call<T>(method: string, params?: Record<string, unknown>): Promise<T> {
|
async call<T>(method: string, params?: Record<string, unknown>): Promise<T> {
|
||||||
if (!this.connected) {
|
if (!this.connected) {
|
||||||
if (!this.autoConnect) {
|
if (!this.autoConnect) {
|
||||||
@@ -498,6 +508,13 @@ export class CompanionRuntimeClient {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if ('event' in parsed) {
|
if ('event' in parsed) {
|
||||||
|
for (const handler of this.eventHandlers) {
|
||||||
|
try {
|
||||||
|
handler(parsed.event, parsed.data);
|
||||||
|
} catch {
|
||||||
|
// Event subscribers are userland callbacks; isolate failures.
|
||||||
|
}
|
||||||
|
}
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user