summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/session-orchestrator/src')
-rw-r--r--packages/session-orchestrator/src/extension.ts1
-rw-r--r--packages/session-orchestrator/src/index.ts2
-rw-r--r--packages/session-orchestrator/src/orchestrator.test.ts37
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts28
-rw-r--r--packages/session-orchestrator/src/queue.test.ts4
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() {},
};
}