diff options
Diffstat (limited to 'packages/session-orchestrator/src')
| -rw-r--r-- | packages/session-orchestrator/src/extension.ts | 1 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/index.ts | 2 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.test.ts | 37 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.ts | 28 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/queue.test.ts | 4 |
5 files changed, 66 insertions, 6 deletions
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts index 6e56c2b..cbb3def 100644 --- a/packages/session-orchestrator/src/extension.ts +++ b/packages/session-orchestrator/src/extension.ts @@ -27,6 +27,7 @@ export const manifest: Manifest = { "session-orchestrator/turn-settled", "session-orchestrator/warm-completed", "session-orchestrator/conversation-closed", + "session-orchestrator/conversation-status-changed", ], }, }; diff --git a/packages/session-orchestrator/src/index.ts b/packages/session-orchestrator/src/index.ts index 3d8ad18..d2aacd9 100644 --- a/packages/session-orchestrator/src/index.ts +++ b/packages/session-orchestrator/src/index.ts @@ -2,9 +2,11 @@ export { extension, manifest } from "./extension.js"; export { type ConversationClosedPayload, type ConversationOpenedPayload, + type ConversationStatusChangedPayload, cacheWarmHandle, conversationClosed, conversationOpened, + conversationStatusChanged, createSessionOrchestrator, createWarmService, type EnqueueInput, diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts index 9ee9704..53c1ce7 100644 --- a/packages/session-orchestrator/src/orchestrator.test.ts +++ b/packages/session-orchestrator/src/orchestrator.test.ts @@ -86,6 +86,10 @@ function createInMemoryStore(): ConversationStore & { return null; }, async setConversationTitle() {}, + async getConversationStatus() { + return null; + }, + async setConversationStatus() {}, }; } @@ -532,6 +536,10 @@ describe("turn-sealed event", () => { return null; }, async setConversationTitle() {}, + async getConversationStatus() { + return null; + }, + async setConversationStatus() {}, }; const { orchestrator } = createSessionOrchestrator({ @@ -594,6 +602,10 @@ describe("turn-sealed event", () => { return null; }, async setConversationTitle() {}, + async getConversationStatus() { + return null; + }, + async setConversationStatus() {}, }; const { orchestrator } = createSessionOrchestrator({ @@ -945,6 +957,10 @@ describe("turn metrics persistence", () => { return null; }, async setConversationTitle() {}, + async getConversationStatus() { + return null; + }, + async setConversationStatus() {}, }; const { orchestrator } = createSessionOrchestrator({ @@ -1112,18 +1128,23 @@ describe("lifecycle event hooks", () => { modelName: "mymodel", }); - expect(emitted).toHaveLength(2); + expect(emitted).toHaveLength(4); expect(emitted[0]?.hook).toBe("session-orchestrator/turn-started"); expect(emitted[0]?.payload.conversationId).toBe("conv-lifecycle"); expect(emitted[0]?.payload.cwd).toBe("/work"); expect(emitted[0]?.payload.modelName).toBe("mymodel"); expect(emitted[0]?.order).toBe(0); - expect(emitted[1]?.hook).toBe("session-orchestrator/turn-settled"); - expect(emitted[1]?.payload.conversationId).toBe("conv-lifecycle"); - expect(emitted[1]?.payload.cwd).toBe("/work"); - expect(emitted[1]?.payload.modelName).toBe("mymodel"); - expect(emitted[1]?.order).toBe(1); + expect(emitted[1]?.hook).toBe("session-orchestrator/conversation-status-changed"); + expect((emitted[1]?.payload as unknown as { status: string }).status).toBe("active"); + + expect(emitted[2]?.hook).toBe("session-orchestrator/turn-settled"); + expect(emitted[2]?.payload.conversationId).toBe("conv-lifecycle"); + expect(emitted[2]?.payload.cwd).toBe("/work"); + expect(emitted[2]?.payload.modelName).toBe("mymodel"); + + expect(emitted[3]?.hook).toBe("session-orchestrator/conversation-status-changed"); + expect((emitted[3]?.payload as unknown as { status: string }).status).toBe("idle"); }); }); @@ -2302,6 +2323,10 @@ describe("closeConversation (CR-4c)", () => { hook: "session-orchestrator/conversation-closed", payload: { conversationId: "conv-never-seen" }, }, + { + hook: "session-orchestrator/conversation-status-changed", + payload: { conversationId: "conv-never-seen", status: "closed" }, + }, ]); // Closing again is still safe. diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts index 82ca59e..4ce83ca 100644 --- a/packages/session-orchestrator/src/orchestrator.ts +++ b/packages/session-orchestrator/src/orchestrator.ts @@ -2,6 +2,7 @@ import type { ConversationStore } from "@dispatch/conversation-store"; import type { AgentEvent, ChatMessage, + ConversationStatus, EventHookDescriptor, Logger, ProviderContract, @@ -111,6 +112,22 @@ export interface ConversationOpenedPayload { export const conversationOpened: EventHookDescriptor<ConversationOpenedPayload> = defineEventHook<ConversationOpenedPayload>("session-orchestrator/conversation-opened"); +/** Payload for the conversationStatusChanged bus event. */ +export interface ConversationStatusChangedPayload { + readonly conversationId: string; + readonly status: ConversationStatus; +} + +/** + * Fired when a conversation's lifecycle status changes (active/idle/closed). + * Transport-ws subscribes and broadcasts a `conversation.statusChanged` WS + * message to all connected frontend clients so tabs sync across devices. + */ +export const conversationStatusChanged: EventHookDescriptor<ConversationStatusChangedPayload> = + defineEventHook<ConversationStatusChangedPayload>( + "session-orchestrator/conversation-status-changed", + ); + /** Payload for the warmCompleted bus event. */ export interface WarmCompletedPayload { readonly conversationId: string; @@ -288,6 +305,8 @@ export function createSessionOrchestrator( payloadPromise.then((payload) => { deps.emit?.(turnStarted, payload); + deps.emit?.(conversationStatusChanged, { conversationId, status: "active" }); + void deps.conversationStore.setConversationStatus(conversationId, "active"); }); void (async () => { @@ -419,6 +438,13 @@ export function createSessionOrchestrator( } void payloadPromise.then((payload) => { deps.emit?.(turnSettled, payload); + if (!carried) { + deps.emit?.(conversationStatusChanged, { + conversationId, + status: "idle", + }); + void deps.conversationStore.setConversationStatus(conversationId, "idle"); + } }); } })(); @@ -486,6 +512,8 @@ export function createSessionOrchestrator( turn.controller.abort(); } deps.emit?.(conversationClosed, { conversationId }); + deps.emit?.(conversationStatusChanged, { conversationId, status: "closed" }); + void deps.conversationStore.setConversationStatus(conversationId, "closed"); return { abortedTurn }; }, diff --git a/packages/session-orchestrator/src/queue.test.ts b/packages/session-orchestrator/src/queue.test.ts index 2aba3d4..745ac27 100644 --- a/packages/session-orchestrator/src/queue.test.ts +++ b/packages/session-orchestrator/src/queue.test.ts @@ -82,6 +82,10 @@ function createInMemoryStore(): ConversationStore & { return null; }, async setConversationTitle() {}, + async getConversationStatus() { + return null; + }, + async setConversationStatus() {}, }; } |
