From ba47df37f0c89bff4f0c3dd7d0bc2ef6c8062b92 Mon Sep 17 00:00:00 2001 From: Adam Malczewski Date: Sun, 21 Jun 2026 02:08:44 +0900 Subject: feat(message-queue): per-conversation queue + steering injection A per-conversation message queue (new message-queue extension) holds user messages enqueued while a turn generates; delivered mid-turn as steering at the tool-result boundary (or carried to a new turn if no tool call fires). - kernel: RunTurnInput.drainSteering callback (generic; kernel stays pure) - wire 0.7.0->0.8.0: QueuedMessage, QueuePayload, TurnSteeringEvent (additive) - transport-contract 0.11.0->0.12.0: POST /conversations/:id/queue + chat.queue WS op - message-queue ext: queue state + per-conversation custom surface (rendererId message-queue) - session-orchestrator: enqueue facade + drainSteering wiring + post-seal carry - transport-http/ws: queue endpoint + chat.queue op (fixes WsClientMessage exhaustive switch) - host-bin: register message-queue 1043 vitest + 199 transport bun pass; tsc/biome clean; boot smoke clean. FE courier: frontend-message-queue-handoff.md. --- packages/message-queue/src/extension.ts | 82 +++++++++++++++++++++++++++++++++ 1 file changed, 82 insertions(+) create mode 100644 packages/message-queue/src/extension.ts (limited to 'packages/message-queue/src/extension.ts') diff --git a/packages/message-queue/src/extension.ts b/packages/message-queue/src/extension.ts new file mode 100644 index 0000000..26d19ca --- /dev/null +++ b/packages/message-queue/src/extension.ts @@ -0,0 +1,82 @@ +/** + * message-queue extension — owns the per-conversation steering message queue + * (state + surface + drain-and-clear). Plugs into the kernel host via a + * manifest + activate(host); the session-orchestrator drains it via the + * `messageQueueHandle` service. Does NOT touch the turn loop. + */ +import type { Extension, HostAPI, Manifest } from "@dispatch/kernel"; +import type { SurfaceContext, SurfaceProvider } from "@dispatch/surface-registry"; +import { surfaceRegistryHandle } from "@dispatch/surface-registry"; +import type { SurfaceSpec } from "@dispatch/ui-contract"; +import { buildQueueSpec, MESSAGE_QUEUE_SURFACE_ID } from "./pure.js"; +import { createMessageQueueService, messageQueueHandle } from "./service.js"; + +export const manifest: Manifest = { + id: "message-queue", + name: "Message Queue", + version: "0.0.0", + apiVersion: "^0.1.0", + trust: "bundled", + activation: "eager", + dependsOn: ["surface-registry"], + capabilities: {}, + contributes: { + services: ["message-queue"], + }, +}; + +export function activate(host: HostAPI): void { + const registry = host.getService(surfaceRegistryHandle); + + const subscribers = new Set<() => void>(); + + const service = createMessageQueueService({ + id: () => crypto.randomUUID(), + now: () => Date.now(), + notify: () => { + for (const sub of subscribers) { + sub(); + } + }, + logger: host.logger, + }); + + host.provideService(messageQueueHandle, service); + + function getSpec(context?: SurfaceContext): SurfaceSpec { + const convId = context?.conversationId; + const messages = convId === undefined ? [] : service.getQueue(convId); + return buildQueueSpec(messages); + } + + function invoke(_actionId: string, _payload?: unknown, _context?: SurfaceContext): void { + // The message-queue surface is read-only: a client renders the queue and + // the session-orchestrator drains it. No client-facing actions. + } + + const provider: SurfaceProvider = { + catalogEntry: { + id: MESSAGE_QUEUE_SURFACE_ID, + region: "side", + title: "Message Queue", + scope: "conversation", + }, + getSpec, + invoke, + subscribe(onChange) { + subscribers.add(onChange); + return () => { + subscribers.delete(onChange); + }; + }, + }; + + registry.register(provider); + + host.logger.info("message-queue: registered"); +} + +export const extension: Extension = { + manifest, + activate, +}; -- cgit v1.2.3