fix: re-polled steering after yield and drained queued follow-ups after turns
- Re-polled steering at the loop yield boundary and included it in the pre-stop pending batch so late messages are processed immediately. - Added session-side draining for stranded queued messages, scheduling an auto-continue when a prompt settles and follow-ups or steers remain. - Added a regression test for late steering injection at yield and updated mid-turn collab prompt handling to keep steering messages in the pending display queue until consumed.
This commit is contained in:
+3
-17
@@ -91,32 +91,18 @@ Known v1 limit for guests: a turn already streaming when you join becomes visibl
|
||||
|---|---|---|
|
||||
| `collab.relayUrl` | `wss://my.omp.sh` | Relay used by `/collab` when no relay is passed inline |
|
||||
| `collab.displayName` | OS username | Name shown to other participants |
|
||||
| `share.serverUrl` | `https://my.omp.sh/s` | Share viewer/upload base used by `/share` (same Go service; links are `<base>/<id>#<key>`) |
|
||||
| `share.serverUrl` | `https://my.omp.sh/s` | Share viewer/upload base used by `/share` (links are `<base>/<id>#<key>`) |
|
||||
| `share.redactSecrets` | `true` | Run the secret obfuscator over `/share` snapshots before upload |
|
||||
|
||||
## Self-hosting the relay
|
||||
|
||||
The relay is a small content-blind Go service (`omp-collab-relay`, in the pi-www repo under `relay/`). It keeps no state beyond live connections and exposes:
|
||||
The relay is a small content-blind Go service. It keeps no state beyond live connections and exposes:
|
||||
|
||||
- `GET /` — the static collab-web guest client (target of the `/collab` deep link),
|
||||
- `GET /r/<roomId>?role=host|guest` — WebSocket upgrade,
|
||||
- `POST /s` / `GET /s/<id>` / `GET /s/<id>/raw` — `/share` blob upload, viewer page, and blob fetch (see the relay README),
|
||||
- `POST /s` / `GET /s/<id>` / `GET /s/<id>/raw` — `/share` blob upload, viewer page, and blob fetch,
|
||||
- `GET /healthz` — liveness.
|
||||
|
||||
Run it:
|
||||
|
||||
```sh
|
||||
go build -o omp-collab-relay .
|
||||
RELAY_BIND=0.0.0.0:7475 ./omp-collab-relay
|
||||
```
|
||||
|
||||
`RELAY_BIND` accepts `host:port`, a bare port (binds localhost), or a unix socket path (front it with a TLS-terminating reverse proxy — guests other than localhost require `wss://`). Then:
|
||||
|
||||
```
|
||||
/collab my-relay.example.com
|
||||
```
|
||||
|
||||
or set `collab.relayUrl` in `/settings`.
|
||||
|
||||
## Architecture notes
|
||||
|
||||
|
||||
+1
-1
@@ -71,7 +71,7 @@
|
||||
- Optional: `timeout`.
|
||||
|
||||
**Execution**
|
||||
- `file: "*"`: `runWorkspaceDiagnostics()` detects project type from root markers and runs one subprocess command: Rust `cargo check --message-format=short`, TypeScript `npx tsc --noEmit`, Go `go build ./...`, Python `pyright`.
|
||||
- `file: "*"`: `runWorkspaceDiagnostics()` detects project type from root markers and runs one subprocess command: Rust `cargo check --message-format=short`, TypeScript `npx tsc --noEmit`, Python `pyright`.
|
||||
- Concrete file or glob: `resolveDiagnosticTargets()` treats non-globs as one target, otherwise expands a `Bun.Glob` up to `MAX_GLOB_DIAGNOSTIC_TARGETS`.
|
||||
- Per file, every matching server runs: custom clients call `lint(file)`; real LSP servers optionally wait for project load, capture `diagnosticsVersion`, `refreshFile()`, then `waitForDiagnostics()` for fresh `publishDiagnostics` (settles on the latest publish; exact-version match accepted immediately).
|
||||
- Results are deduplicated by range+message and severity-sorted.
|
||||
|
||||
@@ -815,11 +815,15 @@ async function runLoopBody(
|
||||
// Agent would stop here. Drain non-interrupting asides + follow-up messages.
|
||||
await config.onBeforeYield?.();
|
||||
// Skip queue drains when externally aborted (same stranding hazard as above).
|
||||
// Re-poll steering too: a steer can land between the stop-boundary dequeue
|
||||
// above and this yield point (e.g. queued while onBeforeYield ran). Without
|
||||
// this poll it would strand in the queue until the next manual prompt.
|
||||
const lateSteering = signal?.aborted ? [] : (await config.getSteeringMessages?.()) || [];
|
||||
const asideMessages = signal?.aborted ? [] : resolveAsides(await config.getAsideMessages?.());
|
||||
const followUpMessages = signal?.aborted ? [] : (await config.getFollowUpMessages?.()) || [];
|
||||
if (asideMessages.length > 0 || followUpMessages.length > 0) {
|
||||
if (lateSteering.length > 0 || asideMessages.length > 0 || followUpMessages.length > 0) {
|
||||
// Set as pending so the inner loop processes them before stopping.
|
||||
pendingMessages = [...asideMessages, ...followUpMessages];
|
||||
pendingMessages = [...lateSteering, ...asideMessages, ...followUpMessages];
|
||||
continue;
|
||||
}
|
||||
|
||||
|
||||
@@ -80,6 +80,38 @@ describe("Agent", () => {
|
||||
expect(mock.calls.length).toBe(2);
|
||||
});
|
||||
|
||||
it("delivers a steer that lands at the yield boundary instead of stranding it", async () => {
|
||||
// Regression: a steering message queued after the stop-boundary dequeue
|
||||
// (e.g. while onBeforeYield runs) was silently stranded in the queue until
|
||||
// the next manual prompt. The outer yield drain must re-poll steering.
|
||||
const mock = createMockModel({ responses: [{ content: ["First answer"] }, { content: ["Steer answer"] }] });
|
||||
const agent = new Agent({ streamFn: mock.stream });
|
||||
let injected = false;
|
||||
agent.setOnBeforeYield(() => {
|
||||
if (injected) return;
|
||||
injected = true;
|
||||
agent.steer({
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "Late steer" }],
|
||||
steering: true,
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
});
|
||||
|
||||
await agent.prompt("Initial");
|
||||
|
||||
expect(mock.calls.length).toBe(2);
|
||||
expect(agent.hasQueuedMessages()).toBe(false);
|
||||
const steerDelivered = agent.state.messages.some(
|
||||
message =>
|
||||
message.role === "user" &&
|
||||
Array.isArray(message.content) &&
|
||||
message.content.some(part => part.type === "text" && part.text === "Late steer"),
|
||||
);
|
||||
expect(steerDelivered).toBe(true);
|
||||
expect(agent.state.messages[agent.state.messages.length - 1].role).toBe("assistant");
|
||||
});
|
||||
|
||||
it("prompt() emits assistant error lifecycle for Anthropic output-blocked stream errors before assistant start", async () => {
|
||||
const mock = createMockModel({ responses: [] });
|
||||
const errorText = "Output blocked by content filtering policy";
|
||||
|
||||
@@ -56,7 +56,7 @@
|
||||
|
||||
- Updated collab link handling to accept compact `roomId#key` links and relay hosts without explicit scheme when joining or starting sessions
|
||||
- Added `/collab stop` and `/collab status` options to control and inspect active shared sessions
|
||||
- Added `/collab`, `/join`, and `/leave` for live session sharing: the host shares an end-to-end encrypted link (AES-256-GCM key only in the link fragment; the relay sees opaque bytes) and guests render the session natively in their own TUI — streaming text, tool cards, footer state, ctrl+o expansion, `/dump` — and can prompt or interrupt the host's agent. Guest prompts render with an author badge; session-mutating commands stay host-only. Defaults to the public `my.omp.sh` relay (`collab.relayUrl`); a self-hostable Go relay lives in the pi-www repo as `omp-collab-relay`
|
||||
- Added `/collab`, `/join`, and `/leave` for live session sharing: the host shares an end-to-end encrypted link (AES-256-GCM key only in the link fragment; the relay sees opaque bytes) and guests render the session natively in their own TUI — streaming text, tool cards, footer state, ctrl+o expansion, `/dump` — and can prompt or interrupt the host's agent. Guest prompts render with an author badge; session-mutating commands stay host-only. Defaults to the public `my.omp.sh` relay (`collab.relayUrl`)
|
||||
- Added `collab.relayUrl` and `collab.displayName` settings plus a `collab` status-line segment showing the participant count (host) or guest role. Guests mirror the host's real model/thinking state into their replica agent and the host's subagent ecosystem end-to-end: the live subagent HUD, the Agent Hub table with live progress, hub chat/kill/revive (routed to the host), and on-demand subagent transcript viewing
|
||||
- Added `omp join <link>` subcommand that launches the interactive TUI and immediately runs `/join <link>`
|
||||
|
||||
|
||||
@@ -29,6 +29,7 @@ import {
|
||||
COLLAB_PROTO,
|
||||
type CollabFrame,
|
||||
type CollabParticipant,
|
||||
type CollabPromptDetails,
|
||||
type CollabSessionState,
|
||||
formatCollabLink,
|
||||
formatCollabWebLink,
|
||||
@@ -364,13 +365,24 @@ export class CollabHost {
|
||||
const name = peer.name;
|
||||
const content: string | (TextContent | ImageContent)[] =
|
||||
images && images.length > 0 ? [{ type: "text", text }, ...images] : text;
|
||||
const details: CollabPromptDetails & { __pendingDisplayTag?: string } = { from: name };
|
||||
if (this.#ctx.session.isStreaming) {
|
||||
// Mid-turn guest prompts are steered: register the pending-display twin
|
||||
// so queuedMessageCount reflects the queued steer (host pending bar +
|
||||
// guests' "queued ×N" badge). The tag dequeues the entry when the agent
|
||||
// consumes the message (mirrors the skill-prompt path).
|
||||
details.__pendingDisplayTag = this.#ctx.session.enqueueCustomMessageDisplay(text, "steer");
|
||||
this.#ctx.updatePendingMessagesDisplay();
|
||||
this.#ctx.ui.requestRender();
|
||||
this.#scheduleStateBroadcast();
|
||||
}
|
||||
this.#ctx.session
|
||||
.promptCustomMessage(
|
||||
{
|
||||
customType: COLLAB_PROMPT_MESSAGE_TYPE,
|
||||
content,
|
||||
display: true,
|
||||
details: { from: name },
|
||||
details,
|
||||
attribution: "user",
|
||||
},
|
||||
{ streamingBehavior: "steer" },
|
||||
|
||||
@@ -1339,7 +1339,7 @@ export const SETTINGS_SCHEMA = {
|
||||
tab: "interaction",
|
||||
group: "Collab",
|
||||
label: "Relay URL",
|
||||
description: "Relay used by /collab (wss://host[:port]; self-host with the omp-collab-relay service)",
|
||||
description: "Relay used by /collab (wss://host[:port])",
|
||||
},
|
||||
},
|
||||
|
||||
@@ -3948,9 +3948,6 @@ export const SETTINGS_SCHEMA = {
|
||||
|
||||
"dev.autoqaPush.endpoint": {
|
||||
type: "string",
|
||||
// Bundled QA collector — runs `/work/pi-www/autoqa` behind qa.omp.sh.
|
||||
// Override via `PI_AUTO_QA_PUSH_URL` or `dev.autoqaPush.endpoint`
|
||||
// in `config.yml` to point at a self-hosted instance.
|
||||
default: "https://qa.omp.sh/v1/grievances" as const,
|
||||
ui: {
|
||||
tab: "tools",
|
||||
|
||||
@@ -1139,9 +1139,24 @@ export class AgentSession {
|
||||
if (this.#promptInFlightCount === 0) {
|
||||
this.#releasePowerAssertion();
|
||||
this.#flushPendingAgentEnd();
|
||||
this.#drainStrandedQueuedMessages();
|
||||
}
|
||||
}
|
||||
|
||||
/** A steer/follow-up can land after the agent loop's final queue poll but
|
||||
* before the prompt unwinds: #promptInFlightCount keeps isStreaming true
|
||||
* through post-prompt recovery, so senders (collab guests, skills) still
|
||||
* queue via agent.steer()/followUp() instead of starting a fresh prompt.
|
||||
* Without a drain those messages strand invisibly until the next manual
|
||||
* prompt. Runs when the session settles; the guard makes it a no-op when
|
||||
* the queue was consumed normally or a new turn already started. */
|
||||
#drainStrandedQueuedMessages(): void {
|
||||
if (!this.agent.hasQueuedMessages()) return;
|
||||
this.#scheduleAgentContinue({
|
||||
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
|
||||
});
|
||||
}
|
||||
|
||||
#resetInFlight(): void {
|
||||
this.#promptInFlightCount = 0;
|
||||
this.#releasePowerAssertion();
|
||||
@@ -5255,6 +5270,14 @@ export class AgentSession {
|
||||
} else {
|
||||
this.agent.steer(normalizedAppMessage);
|
||||
}
|
||||
// The isStreaming check above can be stale: image normalization is
|
||||
// awaited, so the turn may have ended in between, leaving the message
|
||||
// queued on an idle agent. Mirror #queueSteer's idle-path delivery.
|
||||
if (this.#canAutoContinueForFollowUp()) {
|
||||
this.#scheduleAgentContinue({
|
||||
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
|
||||
});
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
/**
|
||||
* Contract: a custom message steered into a streaming session (the collab-host
|
||||
* and skill-prompt path: `promptCustomMessage(..., { streamingBehavior: "steer" })`)
|
||||
* is always delivered — never silently stranded in the agent's steering queue.
|
||||
*
|
||||
* Two regression seams, both observed as "guest messages just disappear" in
|
||||
* collab sessions:
|
||||
* 1. A steer landing at the run's yield boundary (after the stop-boundary
|
||||
* dequeue) must force another turn instead of stranding.
|
||||
* 2. A steer landing while the prompt unwinds (isStreaming stays true through
|
||||
* post-prompt recovery, but the loop is already done) must be drained when
|
||||
* the session settles.
|
||||
*/
|
||||
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
import * as os from "node:os";
|
||||
import * as path from "node:path";
|
||||
import { Agent } from "@oh-my-pi/pi-agent-core";
|
||||
import { createMockModel, type MockModel } from "@oh-my-pi/pi-ai/providers/mock";
|
||||
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
|
||||
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
import { Snowflake } from "@oh-my-pi/pi-utils";
|
||||
|
||||
const COLLAB_PROMPT_TYPE = "collab-prompt";
|
||||
|
||||
interface SteerHarness {
|
||||
session: AgentSession;
|
||||
sessionManager: SessionManager;
|
||||
mock: MockModel;
|
||||
}
|
||||
|
||||
describe("AgentSession queued steer delivery", () => {
|
||||
let tempDir: string;
|
||||
let session: AgentSession;
|
||||
const authStorages: AuthStorage[] = [];
|
||||
|
||||
beforeEach(() => {
|
||||
tempDir = path.join(os.tmpdir(), `pi-steer-strand-${Snowflake.next()}`);
|
||||
fs.mkdirSync(tempDir, { recursive: true });
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
await session?.dispose();
|
||||
for (const authStorage of authStorages.splice(0)) {
|
||||
authStorage.close();
|
||||
}
|
||||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
async function createSession(responses: { content: string[] }[]): Promise<SteerHarness> {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
|
||||
const mock = createMockModel({ responses });
|
||||
const agent = new Agent({
|
||||
getApiKey: () => "test-key",
|
||||
initialState: { model, systemPrompt: ["Test"], tools: [] },
|
||||
streamFn: mock.stream,
|
||||
});
|
||||
const sessionManager = SessionManager.inMemory();
|
||||
const settings = Settings.isolated({ "compaction.enabled": false });
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir, `auth-${Snowflake.next()}.db`));
|
||||
authStorages.push(authStorage);
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
|
||||
session = new AgentSession({ agent, sessionManager, settings, modelRegistry });
|
||||
return { session, sessionManager, mock };
|
||||
}
|
||||
|
||||
function steerCollabPrompt(target: AgentSession, text: string): Promise<void> {
|
||||
return target.promptCustomMessage(
|
||||
{
|
||||
customType: COLLAB_PROMPT_TYPE,
|
||||
content: text,
|
||||
display: true,
|
||||
details: { from: "guest" },
|
||||
attribution: "user",
|
||||
},
|
||||
{ streamingBehavior: "steer" },
|
||||
);
|
||||
}
|
||||
|
||||
/** Resolves with the entry text when a collab-prompt entry is persisted. */
|
||||
function nextCollabEntry(sessionManager: SessionManager): Promise<string> {
|
||||
const { promise, resolve } = Promise.withResolvers<string>();
|
||||
sessionManager.onEntryAppended = entry => {
|
||||
if (entry.type === "custom_message" && entry.customType === COLLAB_PROMPT_TYPE) {
|
||||
resolve(typeof entry.content === "string" ? entry.content : JSON.stringify(entry.content));
|
||||
}
|
||||
};
|
||||
return promise;
|
||||
}
|
||||
|
||||
it("delivers a collab steer that lands at the run's yield boundary", async () => {
|
||||
const { session, sessionManager, mock } = await createSession([
|
||||
{ content: ["host answer"] },
|
||||
{ content: ["ack guest"] },
|
||||
]);
|
||||
const entryAppended = nextCollabEntry(sessionManager);
|
||||
|
||||
let streamingAtInject: boolean | undefined;
|
||||
let injected = false;
|
||||
session.agent.setOnBeforeYield(async () => {
|
||||
if (injected) return;
|
||||
injected = true;
|
||||
// The session is still mid-prompt here, so this takes the steer path.
|
||||
streamingAtInject = session.isStreaming;
|
||||
await steerCollabPrompt(session, "guest steer at yield");
|
||||
});
|
||||
|
||||
await session.prompt("hello");
|
||||
|
||||
expect(streamingAtInject).toBe(true);
|
||||
expect(await entryAppended).toBe("guest steer at yield");
|
||||
expect(mock.calls.length).toBe(2);
|
||||
expect(session.agent.hasQueuedMessages()).toBe(false);
|
||||
});
|
||||
|
||||
it("drains a steer stranded in the agent queue when the session settles", async () => {
|
||||
const { session, sessionManager, mock } = await createSession([
|
||||
{ content: ["host answer"] },
|
||||
{ content: ["ack guest"] },
|
||||
]);
|
||||
const entryAppended = nextCollabEntry(sessionManager);
|
||||
|
||||
// Inject from the wire agent_end subscriber: it fires synchronously while
|
||||
// the session settles (#promptInFlightCount just hit 0), after the agent
|
||||
// loop's final queue poll — a message queued here is invisible to the run
|
||||
// and must be picked up by the settle-time drain.
|
||||
const secondRunDone = Promise.withResolvers<void>();
|
||||
let agentEnds = 0;
|
||||
session.subscribe(event => {
|
||||
if (event.type !== "agent_end") return;
|
||||
agentEnds++;
|
||||
if (agentEnds === 1) {
|
||||
session.agent.steer({
|
||||
role: "custom",
|
||||
customType: COLLAB_PROMPT_TYPE,
|
||||
content: "guest steer at settle",
|
||||
display: true,
|
||||
details: { from: "guest" },
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
} else if (agentEnds === 2) {
|
||||
secondRunDone.resolve();
|
||||
}
|
||||
});
|
||||
|
||||
await session.prompt("hello");
|
||||
expect(await entryAppended).toBe("guest steer at settle");
|
||||
await secondRunDone.promise;
|
||||
|
||||
expect(mock.calls.length).toBe(2);
|
||||
expect(session.agent.hasQueuedMessages()).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -25,7 +25,7 @@ interface RelayData {
|
||||
|
||||
type RelaySocket = Bun.ServerWebSocket<RelayData>;
|
||||
|
||||
/** Single-room relay mirroring the omp-collab-relay forwarding contract. */
|
||||
/** Single-room test relay mirroring the production forwarding contract. */
|
||||
function startTestRelay(): { url: string; stop(): void } {
|
||||
let host: RelaySocket | null = null;
|
||||
const guests = new Map<number, RelaySocket>();
|
||||
|
||||
@@ -0,0 +1,214 @@
|
||||
/**
|
||||
* Contract: a guest prompt that arrives while the host agent is streaming is
|
||||
* steered AND becomes visible as a queued message — the host registers the
|
||||
* pending-display twin (so `queuedMessageCount` covers it, feeding the web
|
||||
* composer's "queued ×N" badge and the TUI pending bar) and tags the custom
|
||||
* message so the entry is dequeued when the agent consumes it.
|
||||
*/
|
||||
import { afterEach, describe, expect, it } from "bun:test";
|
||||
import { importRoomKey } from "@oh-my-pi/pi-coding-agent/collab/crypto";
|
||||
import { CollabHost } from "@oh-my-pi/pi-coding-agent/collab/host";
|
||||
import {
|
||||
COLLAB_PROTO,
|
||||
type CollabFrame,
|
||||
parseCollabLink,
|
||||
rewriteEnvelopePeer,
|
||||
unpackEnvelope,
|
||||
} from "@oh-my-pi/pi-coding-agent/collab/protocol";
|
||||
import { CollabSocket } from "@oh-my-pi/pi-coding-agent/collab/relay-client";
|
||||
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
|
||||
|
||||
interface RelayData {
|
||||
role: "host" | "guest";
|
||||
peerId: number;
|
||||
}
|
||||
|
||||
type RelaySocket = Bun.ServerWebSocket<RelayData>;
|
||||
|
||||
/** Single-room relay mirroring the omp-collab-relay forwarding contract. */
|
||||
function startTestRelay(): { url: string; stop(): void } {
|
||||
let host: RelaySocket | null = null;
|
||||
const guests = new Map<number, RelaySocket>();
|
||||
let nextPeerId = 1;
|
||||
const server = Bun.serve({
|
||||
port: 0,
|
||||
fetch(req, srv): Response | undefined {
|
||||
const role = new URL(req.url).searchParams.get("role") === "host" ? "host" : "guest";
|
||||
const data: RelayData = { role, peerId: 0 };
|
||||
if (srv.upgrade(req, { data })) return undefined;
|
||||
return new Response("upgrade failed", { status: 400 });
|
||||
},
|
||||
websocket: {
|
||||
open(ws: RelaySocket): void {
|
||||
if (ws.data.role === "host") {
|
||||
host = ws;
|
||||
return;
|
||||
}
|
||||
ws.data.peerId = nextPeerId++;
|
||||
guests.set(ws.data.peerId, ws);
|
||||
host?.send(JSON.stringify({ t: "peer-joined", peer: ws.data.peerId }));
|
||||
},
|
||||
message(ws: RelaySocket, message: string | Buffer): void {
|
||||
if (typeof message === "string") return;
|
||||
const bytes = new Uint8Array(message);
|
||||
if (ws.data.role === "host") {
|
||||
const envelope = unpackEnvelope(bytes);
|
||||
if (!envelope) return;
|
||||
if (envelope.peerId === 0) {
|
||||
for (const guest of guests.values()) guest.send(bytes);
|
||||
} else {
|
||||
guests.get(envelope.peerId)?.send(bytes);
|
||||
}
|
||||
return;
|
||||
}
|
||||
rewriteEnvelopePeer(bytes, ws.data.peerId);
|
||||
host?.send(bytes);
|
||||
},
|
||||
close(ws: RelaySocket): void {
|
||||
if (ws.data.role === "guest") {
|
||||
guests.delete(ws.data.peerId);
|
||||
host?.send(JSON.stringify({ t: "peer-left", peer: ws.data.peerId }));
|
||||
}
|
||||
},
|
||||
},
|
||||
});
|
||||
return { url: `ws://localhost:${server.port}`, stop: () => server.stop(true) };
|
||||
}
|
||||
|
||||
interface CapturedPrompt {
|
||||
details?: { from?: string; __pendingDisplayTag?: string };
|
||||
}
|
||||
|
||||
interface StreamingHostHarness {
|
||||
ctx: InteractiveModeContext;
|
||||
prompts: CapturedPrompt[];
|
||||
enqueued: { text: string; mode: string; tag: string }[];
|
||||
nextPrompt(): Promise<CapturedPrompt>;
|
||||
}
|
||||
|
||||
/** Context double for a host whose agent is mid-turn (isStreaming === true). */
|
||||
function makeStreamingHostContext(): StreamingHostHarness {
|
||||
const prompts: CapturedPrompt[] = [];
|
||||
const enqueued: { text: string; mode: string; tag: string }[] = [];
|
||||
const promptWaiters: ((prompt: CapturedPrompt) => void)[] = [];
|
||||
const ctx = {
|
||||
settings: { get: () => "" },
|
||||
sessionManager: {
|
||||
getSessionId: () => "sess-1",
|
||||
getCwd: () => "/tmp",
|
||||
snapshotForReplication: () => ({
|
||||
header: { type: "session", id: "sess-1", timestamp: new Date().toISOString(), cwd: "/tmp" },
|
||||
entries: [],
|
||||
}),
|
||||
onEntryAppended: undefined,
|
||||
},
|
||||
session: {
|
||||
isStreaming: true,
|
||||
get queuedMessageCount(): number {
|
||||
return enqueued.length;
|
||||
},
|
||||
sessionName: "test",
|
||||
model: undefined,
|
||||
thinkingLevel: undefined,
|
||||
subscribe: () => () => {},
|
||||
emitNotice: () => {},
|
||||
enqueueCustomMessageDisplay: (text: string, mode: string) => {
|
||||
const tag = `tag-${enqueued.length + 1}`;
|
||||
enqueued.push({ text, mode, tag });
|
||||
return tag;
|
||||
},
|
||||
promptCustomMessage: (message: CapturedPrompt) => {
|
||||
const captured: CapturedPrompt = { details: message.details };
|
||||
prompts.push(captured);
|
||||
for (const waiter of promptWaiters.splice(0)) waiter(captured);
|
||||
return Promise.resolve();
|
||||
},
|
||||
},
|
||||
eventBus: undefined,
|
||||
statusLine: {
|
||||
setCollabStatus: () => {},
|
||||
invalidate: () => {},
|
||||
getCachedContextBreakdown: () => ({ usedTokens: 0, contextWindow: 0 }),
|
||||
},
|
||||
ui: { requestRender: () => {} },
|
||||
updatePendingMessagesDisplay: () => {},
|
||||
showStatus: () => {},
|
||||
collabHost: undefined,
|
||||
} as unknown as InteractiveModeContext;
|
||||
const nextPrompt = (): Promise<CapturedPrompt> => {
|
||||
const { promise, resolve } = Promise.withResolvers<CapturedPrompt>();
|
||||
promptWaiters.push(resolve);
|
||||
return promise;
|
||||
};
|
||||
return { ctx, prompts, enqueued, nextPrompt };
|
||||
}
|
||||
|
||||
interface TestGuest {
|
||||
socket: CollabSocket;
|
||||
nextFrame(): Promise<CollabFrame>;
|
||||
}
|
||||
|
||||
async function joinAsGuest(link: string, name: string): Promise<TestGuest> {
|
||||
const parsed = parseCollabLink(link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
const writeToken = parsed.writeToken ? Buffer.from(parsed.writeToken).toString("base64url") : undefined;
|
||||
const key = await importRoomKey(parsed.key);
|
||||
const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "guest", key });
|
||||
const queue: CollabFrame[] = [];
|
||||
const waiters: ((frame: CollabFrame) => void)[] = [];
|
||||
socket.onFrame = frame => {
|
||||
const waiter = waiters.shift();
|
||||
if (waiter) waiter(frame);
|
||||
else queue.push(frame);
|
||||
};
|
||||
socket.onOpen = () => socket.send({ t: "hello", proto: COLLAB_PROTO, name, writeToken });
|
||||
socket.connect();
|
||||
const nextFrame = (): Promise<CollabFrame> => {
|
||||
const queued = queue.shift();
|
||||
if (queued) return Promise.resolve(queued);
|
||||
const { promise, resolve } = Promise.withResolvers<CollabFrame>();
|
||||
waiters.push(resolve);
|
||||
return promise;
|
||||
};
|
||||
return { socket, nextFrame };
|
||||
}
|
||||
|
||||
const cleanups: (() => void | Promise<void>)[] = [];
|
||||
|
||||
afterEach(async () => {
|
||||
for (const cleanup of cleanups.splice(0).reverse()) await cleanup();
|
||||
});
|
||||
|
||||
describe("collab mid-turn guest prompts", () => {
|
||||
it("registers the steer as a queued message and reports it to guests via state", async () => {
|
||||
const relay = startTestRelay();
|
||||
cleanups.push(relay.stop);
|
||||
const harness = makeStreamingHostContext();
|
||||
const host = new CollabHost(harness.ctx);
|
||||
await host.start(relay.url);
|
||||
cleanups.push(() => host.stop("test done"));
|
||||
|
||||
const guest = await joinAsGuest(host.link, "writer");
|
||||
cleanups.push(() => guest.socket.close());
|
||||
const welcome = await guest.nextFrame();
|
||||
if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`);
|
||||
|
||||
const prompted = harness.nextPrompt();
|
||||
guest.socket.send({ t: "prompt", text: "steer the host" });
|
||||
const prompt = await prompted;
|
||||
|
||||
// Display twin registered before dispatch, tag forwarded in details so
|
||||
// the session dequeues the entry when the agent consumes the message.
|
||||
expect(harness.enqueued).toEqual([{ text: "steer the host", mode: "steer", tag: "tag-1" }]);
|
||||
expect(prompt.details).toEqual({ from: "writer", __pendingDisplayTag: "tag-1" });
|
||||
|
||||
// The queued steer must reach guests through state.queuedMessageCount —
|
||||
// that field drives the web composer's "queued ×N" badge.
|
||||
let sawQueuedCount = false;
|
||||
for (let i = 0; i < 10 && !sawQueuedCount; i++) {
|
||||
const frame = await guest.nextFrame();
|
||||
if (frame.t === "state" && frame.state.queuedMessageCount === 1) sawQueuedCount = true;
|
||||
}
|
||||
expect(sawQueuedCount).toBe(true);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user