feat(coding-agent): added collaborative sessions with host and guest command support
- Added AES-GCM room-key crypto, relay link parsing, and invalid-link validation. - Added `/collab`, `/join`, and `/leave` command handling for collaborative sessions. - Added startup `join` argument wiring to execute `/join` during interactive launch. - Added status-line, prompt, and command-routing updates for guest/host collaboration UX.
This commit is contained in:
+110
@@ -0,0 +1,110 @@
|
||||
# Collab: Live Session Sharing
|
||||
|
||||
`/collab` shares your running session with other omp instances in real time. Guests render the **same session natively in their own TUI** — streaming assistant text, tool-call cards, footer state (cwd, model, context %, cost), ctrl+o expansion, `/dump` — no terminal mirroring. Guests can prompt and interrupt the agent; the host machine runs the agent and all tools.
|
||||
|
||||
## Quick start
|
||||
|
||||
Host:
|
||||
|
||||
```
|
||||
/collab
|
||||
```
|
||||
|
||||
prints a link like
|
||||
|
||||
```
|
||||
Collab link: mgAYTZwEnpRQtca0CTgn-Q#gdJUbTovD94ofDaa8YvhY0-ty16w4fn8PgB6PLnoA30
|
||||
```
|
||||
|
||||
Guest (any directory, any machine):
|
||||
|
||||
```
|
||||
/join mgAYTZwEnpRQtca0CTgn-Q#gdJU…
|
||||
```
|
||||
|
||||
The guest's previous session is restored on `/leave` (or when the host stops).
|
||||
|
||||
### Commands
|
||||
|
||||
| Command | Effect |
|
||||
|---|---|
|
||||
| `/collab` | Start sharing (or re-print the link when already hosting) |
|
||||
| `/collab <relay>` | Start sharing through a specific relay (`relay.example.com`, `ws://localhost:7475`) |
|
||||
| `/collab status` | Show link + participants |
|
||||
| `/collab stop` | Stop sharing |
|
||||
| `/join <link>` | Join a shared session as a guest |
|
||||
| `/leave` | Leave (guest) or stop sharing (host) |
|
||||
|
||||
## Link format
|
||||
|
||||
```
|
||||
<roomId>#<key> → default relay (relay.omp.sh)
|
||||
host[:port]/r/<roomId>#<key> → custom relay, wss:// inferred
|
||||
ws://localhost:7475/r/<roomId>#<key> → plain ws, allowed for localhost only
|
||||
```
|
||||
|
||||
The fragment (`#<key>`) is the 32-byte AES-256-GCM room key, base64url-encoded. Fragments never appear in HTTP requests, and the key is never sent to the relay.
|
||||
|
||||
## End-to-end encryption
|
||||
|
||||
Every session payload (entries, events, state, prompts) is sealed with AES-256-GCM before it touches the socket. The relay sees only:
|
||||
|
||||
- room ids and connection counts,
|
||||
- opaque ciphertext frames and their sizes,
|
||||
- a 4-byte routing prefix (which guest a frame targets).
|
||||
|
||||
Possession of the link is the trust boundary: anyone with the full link can read the session and prompt the agent. Share it like a secret.
|
||||
|
||||
## Guest permission model
|
||||
|
||||
Single trust level. Guests can:
|
||||
|
||||
- read the entire session (including the back-transcript at join time),
|
||||
- prompt the agent (rendered with their name badge on every participant's transcript; the LLM sees the prompt text verbatim — names are display-only),
|
||||
- interrupt the agent (Esc),
|
||||
- use the Agent Hub against the host's subagents: live table and progress, chat (steers the host's subagent), kill, revive, and transcript viewing (fetched from the host on demand).
|
||||
|
||||
Everything that mutates the host session or machine is host-only: `/model`, `/compact`, `/resume`, `/branch`, bash (`!`), python (`$`), skills, etc. Guests keep a small local allowlist (`/dump`, `/export`, `/copy`, `/help`, `/hotkeys`, `/theme`, `/settings`, `/leave`, `/collab`, `/exit`).
|
||||
|
||||
Known v1 limit for guests: a turn already streaming when you join becomes visible from its next message boundary.
|
||||
|
||||
## Settings
|
||||
|
||||
| Setting | Default | Meaning |
|
||||
|---|---|---|
|
||||
| `collab.relayUrl` | `wss://relay.omp.sh` | Relay used by `/collab` when no relay is passed inline |
|
||||
| `collab.displayName` | OS username | Name shown to other participants |
|
||||
|
||||
## 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:
|
||||
|
||||
- `GET /r/<roomId>?role=host|guest` — WebSocket upgrade,
|
||||
- `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
|
||||
|
||||
Hub topology — the host is authoritative, guests never peer:
|
||||
|
||||
1. `entry` frames — durable session entries, broadcast pre-blob-externalization so images stay inline (guests cannot resolve host blob refs). Guests append them verbatim (ids preserved) to a replica session file under `~/.omp/collab/<roomId>.jsonl` and into the agent's message array, which is why `/dump` and context estimates work.
|
||||
2. `event` frames — live agent events, fed straight into the guest's normal event controller; rendering is events-only to prevent double-render.
|
||||
3. `state` frames — debounced footer snapshots: streaming flag, the host's full model object and thinking level (applied to the guest's replica agent state, so model display and context-window math are native), host context numbers, and participants.
|
||||
4. `bus` frames — mirrored task-subagent lifecycle/progress EventBus traffic, republished on the guest's local bus so the subagent HUD and status-line count work natively.
|
||||
5. `agents` frames — agent-registry snapshots feeding a guest-local registry, so the Agent Hub table renders host subagents.
|
||||
|
||||
Guest→host: `hello`, `prompt`, `abort`, `agent-cmd` (hub chat/kill/revive), and `fetch-transcript` (incremental subagent-transcript reads answered by targeted `transcript` frames). The replica loads through the regular `/resume` machinery, so theming, ctrl+o, and transcript behavior are native by construction; the guest process never chdirs to host paths.
|
||||
@@ -1,11 +1,22 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Added
|
||||
|
||||
- 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 `relay.omp.sh` relay (`collab.relayUrl`); a self-hostable Go relay lives in the pi-www repo as `omp-collab-relay`
|
||||
- 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>`
|
||||
|
||||
### Changed
|
||||
|
||||
- The Ctrl+P role-cycle track and the plan-approval model slider now color segments by track position from the theme's own palette (accent/success/warning/error + markdown/syntax hues, deduplicated per theme since many themes alias them) instead of role-keyed colors with a gray fallback for custom roles
|
||||
|
||||
### Security
|
||||
|
||||
- Rejected non-local `ws://` relay URLs and invalid room keys when parsing collab links to prevent insecure or malformed session joins
|
||||
|
||||
## [15.11.7] - 2026-06-12
|
||||
### Added
|
||||
|
||||
|
||||
@@ -26,6 +26,7 @@ export const commands: CommandEntry[] = [
|
||||
{ name: "gallery", load: () => import("./commands/gallery").then(m => m.default) },
|
||||
{ name: "grievances", load: () => import("./commands/grievances").then(m => m.default) },
|
||||
{ name: "install", load: () => import("./commands/install").then(m => m.default) },
|
||||
{ name: "join", load: () => import("./commands/join").then(m => m.default) },
|
||||
{ name: "plugin", load: () => import("./commands/plugin").then(m => m.default) },
|
||||
{ name: "setup", load: () => import("./commands/setup").then(m => m.default) },
|
||||
{ name: "shell", load: () => import("./commands/shell").then(m => m.default) },
|
||||
|
||||
@@ -32,6 +32,8 @@ export interface Args {
|
||||
sessionDir?: string;
|
||||
providerSessionId?: string;
|
||||
fork?: string;
|
||||
/** Collab link to join at startup (set by the `join` subcommand; no CLI flag). */
|
||||
join?: string;
|
||||
models?: string[];
|
||||
tools?: string[];
|
||||
noTools?: boolean;
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
/**
|
||||
* AES-256-GCM sealing for collab frames.
|
||||
*
|
||||
* The room key lives only in the link fragment; the relay sees opaque bytes.
|
||||
* Sealed layout: `[12B IV][ciphertext+tag]`.
|
||||
*/
|
||||
import type { CollabFrame } from "./protocol";
|
||||
|
||||
const AES_ALGORITHM = "AES-GCM";
|
||||
const IV_LENGTH = 12;
|
||||
const KEY_LENGTH = 32;
|
||||
const TEXT_ENCODER = new TextEncoder();
|
||||
const TEXT_DECODER = new TextDecoder();
|
||||
|
||||
export function generateRoomKey(): Uint8Array {
|
||||
const key = new Uint8Array(KEY_LENGTH);
|
||||
crypto.getRandomValues(key);
|
||||
return key;
|
||||
}
|
||||
|
||||
export function importRoomKey(raw: Uint8Array): Promise<CryptoKey> {
|
||||
if (raw.byteLength !== KEY_LENGTH) {
|
||||
throw new Error(`Room key must be ${KEY_LENGTH} bytes, got ${raw.byteLength}`);
|
||||
}
|
||||
return crypto.subtle.importKey("raw", asStrict(raw), AES_ALGORITHM, false, ["encrypt", "decrypt"]);
|
||||
}
|
||||
|
||||
export async function seal(key: CryptoKey, frame: CollabFrame): Promise<Uint8Array> {
|
||||
const iv = new Uint8Array(IV_LENGTH);
|
||||
crypto.getRandomValues(iv);
|
||||
const plaintext = TEXT_ENCODER.encode(JSON.stringify(frame));
|
||||
const ciphertext = new Uint8Array(await crypto.subtle.encrypt({ name: AES_ALGORITHM, iv }, key, plaintext));
|
||||
const out = new Uint8Array(IV_LENGTH + ciphertext.byteLength);
|
||||
out.set(iv, 0);
|
||||
out.set(ciphertext, IV_LENGTH);
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Inverse of {@link seal}. Throws on auth failure or malformed input. */
|
||||
export async function open(key: CryptoKey, data: Uint8Array): Promise<CollabFrame> {
|
||||
if (data.byteLength <= IV_LENGTH) {
|
||||
throw new Error("Sealed frame too short");
|
||||
}
|
||||
const iv = asStrict(data.subarray(0, IV_LENGTH));
|
||||
const ciphertext = asStrict(data.subarray(IV_LENGTH));
|
||||
const plaintext = new Uint8Array(await crypto.subtle.decrypt({ name: AES_ALGORITHM, iv }, key, ciphertext));
|
||||
return JSON.parse(TEXT_DECODER.decode(plaintext)) as CollabFrame;
|
||||
}
|
||||
|
||||
function asStrict(bytes: Uint8Array): Uint8Array<ArrayBuffer> {
|
||||
if (bytes.buffer instanceof ArrayBuffer && bytes.byteOffset === 0 && bytes.byteLength === bytes.buffer.byteLength) {
|
||||
return bytes as Uint8Array<ArrayBuffer>;
|
||||
}
|
||||
const copy = new Uint8Array(bytes.byteLength);
|
||||
copy.set(bytes);
|
||||
return copy;
|
||||
}
|
||||
@@ -0,0 +1,421 @@
|
||||
/**
|
||||
* Guest side of a collab live session.
|
||||
*
|
||||
* `/join <link>` writes the host's snapshot to a replica session file and
|
||||
* drives it through the normal `/resume` machinery, then applies live frames:
|
||||
* entries → SessionManager + agent.replaceMessages, events →
|
||||
* EventController.handleEvent, state → status-line overrides plus real
|
||||
* model/thinking state applied to the replica agent. The host's subagent
|
||||
* ecosystem is mirrored too: agent snapshots populate a local AgentRegistry
|
||||
* (Agent Hub), EventBus traffic (observer HUD) is republished, and hub
|
||||
* actions (chat/kill/revive/transcript reads) round-trip over the wire.
|
||||
* Everything renders through the same components, so ctrl+o, theming, and
|
||||
* transcript behavior are native by construction.
|
||||
*/
|
||||
import * as path from "node:path";
|
||||
import type { ThinkingLevel } from "@oh-my-pi/pi-agent-core";
|
||||
import type { ImageContent } from "@oh-my-pi/pi-ai";
|
||||
import { getConfigRootDir, logger } from "@oh-my-pi/pi-utils";
|
||||
import type { AgentHubRemote } from "../modes/components/agent-hub";
|
||||
import type { InteractiveModeContext } from "../modes/types";
|
||||
import { AgentRegistry } from "../registry/agent-registry";
|
||||
import type { AgentSessionEvent } from "../session/agent-session";
|
||||
import { shouldDisableReasoning, toReasoningEffort } from "../thinking";
|
||||
import { setSessionTerminalTitle } from "../utils/title-generator";
|
||||
import { importRoomKey } from "./crypto";
|
||||
import { collabDisplayName } from "./host";
|
||||
import {
|
||||
type AgentSnapshot,
|
||||
COLLAB_PROTO,
|
||||
type CollabFrame,
|
||||
type CollabSessionState,
|
||||
parseCollabLink,
|
||||
} from "./protocol";
|
||||
import { CollabSocket } from "./relay-client";
|
||||
|
||||
/** Commands a guest may run locally; everything else is host-only. */
|
||||
export const COLLAB_GUEST_ALLOWED_COMMANDS: Record<string, true> = {
|
||||
dump: true,
|
||||
export: true,
|
||||
copy: true,
|
||||
help: true,
|
||||
hotkeys: true,
|
||||
theme: true,
|
||||
settings: true,
|
||||
leave: true,
|
||||
collab: true,
|
||||
exit: true,
|
||||
quit: true,
|
||||
};
|
||||
const WELCOME_TIMEOUT_MS = 30_000;
|
||||
const TRANSCRIPT_TIMEOUT_MS = 20_000;
|
||||
|
||||
type WelcomeFrame = Extract<CollabFrame, { t: "welcome" }>;
|
||||
|
||||
export class CollabGuestLink {
|
||||
#ctx: InteractiveModeContext;
|
||||
#socket: CollabSocket | null = null;
|
||||
#roomId = "";
|
||||
/** Previous session file to restore on leave; null = previous session was unsaved. */
|
||||
#returnSessionFile: string | null = null;
|
||||
/** Frames apply strictly in arrival order through this chain. */
|
||||
#applyChain: Promise<void> = Promise.resolve();
|
||||
#welcomed = false;
|
||||
#left = false;
|
||||
/** False until the first assistant message_start (real or synthesized) since (re)sync. */
|
||||
#assistantStreamSynced = false;
|
||||
state: CollabSessionState | null = null;
|
||||
/** Local mirror of the host's agent ecosystem (refs carry `session: null`). */
|
||||
readonly agentRegistry = new AgentRegistry();
|
||||
/** Per-agent `hasSessionFile` from the last snapshot; gates remote transcript fetches. */
|
||||
#agentHasTranscript = new Map<string, boolean>();
|
||||
#pendingTranscripts = new Map<number, (r: { text: string; newSize: number } | null) => void>();
|
||||
#nextReqId = 1;
|
||||
readonly #hubRemote: AgentHubRemote = {
|
||||
chat: (id, text) => {
|
||||
this.#socket?.send({ t: "agent-cmd", cmd: "chat", agentId: id, text });
|
||||
},
|
||||
kill: id => {
|
||||
this.#socket?.send({ t: "agent-cmd", cmd: "kill", agentId: id });
|
||||
},
|
||||
revive: id => {
|
||||
this.#socket?.send({ t: "agent-cmd", cmd: "revive", agentId: id });
|
||||
},
|
||||
readTranscript: (id, fromByte) => {
|
||||
const socket = this.#socket;
|
||||
if (!socket || this.#agentHasTranscript.get(id) === false) {
|
||||
return Promise.resolve(null);
|
||||
}
|
||||
const reqId = this.#nextReqId++;
|
||||
const { promise, resolve } = Promise.withResolvers<{ text: string; newSize: number } | null>();
|
||||
const timer = setTimeout(() => {
|
||||
this.#pendingTranscripts.delete(reqId);
|
||||
resolve(null);
|
||||
}, TRANSCRIPT_TIMEOUT_MS);
|
||||
this.#pendingTranscripts.set(reqId, result => {
|
||||
clearTimeout(timer);
|
||||
resolve(result);
|
||||
});
|
||||
socket.send({ t: "fetch-transcript", reqId, agentId: id, fromByte });
|
||||
return promise;
|
||||
},
|
||||
};
|
||||
|
||||
/** Agent Hub actions routed to the host over the wire. */
|
||||
get hubRemote(): AgentHubRemote {
|
||||
return this.#hubRemote;
|
||||
}
|
||||
|
||||
constructor(ctx: InteractiveModeContext) {
|
||||
this.#ctx = ctx;
|
||||
}
|
||||
|
||||
async join(link: string): Promise<void> {
|
||||
const parsed = parseCollabLink(link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
this.#roomId = parsed.roomId;
|
||||
const key = await importRoomKey(parsed.key);
|
||||
|
||||
this.#returnSessionFile = this.#ctx.sessionManager.getSessionFile() ?? null;
|
||||
|
||||
const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "guest", key });
|
||||
this.#socket = socket;
|
||||
|
||||
const firstWelcome = Promise.withResolvers<void>();
|
||||
let joined = false;
|
||||
|
||||
socket.onOpen = () => {
|
||||
// (Re)connect: re-introduce ourselves; the host answers with a fresh
|
||||
// welcome which (re)syncs the replica.
|
||||
this.#welcomed = false;
|
||||
socket.send({ t: "hello", proto: COLLAB_PROTO, name: collabDisplayName(this.#ctx) });
|
||||
};
|
||||
socket.onFrame = frame => {
|
||||
this.#applyChain = this.#applyChain
|
||||
.then(async () => {
|
||||
if (frame.t === "welcome") {
|
||||
await this.#applyWelcome(frame, joined);
|
||||
if (!joined) {
|
||||
joined = true;
|
||||
firstWelcome.resolve();
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (!this.#welcomed || this.#left) return;
|
||||
this.#applyFrame(frame);
|
||||
})
|
||||
.catch(err => logger.warn("collab guest frame apply failed", { type: frame.t, error: String(err) }));
|
||||
};
|
||||
socket.onClose = (reason, willReconnect) => {
|
||||
this.#flushPendingTranscripts();
|
||||
if (this.#left) return;
|
||||
if (!joined) {
|
||||
firstWelcome.reject(new Error(reason));
|
||||
return;
|
||||
}
|
||||
if (willReconnect) {
|
||||
this.#ctx.showStatus(`Collab connection lost (${reason}), reconnecting…`, { dim: true });
|
||||
return;
|
||||
}
|
||||
this.#ctx.showStatus(`Collab session ended (${reason})`);
|
||||
void this.#restoreLocalSession();
|
||||
};
|
||||
socket.connect();
|
||||
|
||||
const timeout = setTimeout(
|
||||
() => firstWelcome.reject(new Error("timed out waiting for the host's welcome")),
|
||||
WELCOME_TIMEOUT_MS,
|
||||
);
|
||||
try {
|
||||
await firstWelcome.promise;
|
||||
} catch (err) {
|
||||
this.#left = true;
|
||||
socket.close();
|
||||
this.#socket = null;
|
||||
throw err;
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
|
||||
this.#ctx.collabGuest = this;
|
||||
}
|
||||
|
||||
/** User-initiated leave (or post-disconnect cleanup): restore the previous session. */
|
||||
async leave(_reason: string): Promise<void> {
|
||||
if (this.#left) return;
|
||||
this.#socket?.close();
|
||||
await this.#restoreLocalSession();
|
||||
}
|
||||
|
||||
sendPrompt(text: string, images?: ImageContent[]): void {
|
||||
this.#socket?.send({ t: "prompt", text, images: images && images.length > 0 ? images : undefined });
|
||||
}
|
||||
|
||||
sendAbort(): void {
|
||||
this.#socket?.send({ t: "abort" });
|
||||
}
|
||||
|
||||
/** Write the welcome snapshot to the replica file and (re)load it through the resume machinery. */
|
||||
async #applyWelcome(frame: WelcomeFrame, isResync: boolean): Promise<void> {
|
||||
if (this.#left) return;
|
||||
const replicaPath = path.join(getConfigRootDir(), "collab", `${this.#roomId}.jsonl`);
|
||||
const lines = [frame.header, ...frame.entries].map(entry => JSON.stringify(entry)).join("\n");
|
||||
await Bun.write(replicaPath, `${lines}\n`);
|
||||
|
||||
// Resume sequence (selector-controller.handleResumeSession) minus
|
||||
// applyCwdChange: the guest process never chdirs to a host path. The
|
||||
// SessionManager still adopts the header cwd for display/relativization.
|
||||
this.#clearTransientUi();
|
||||
this.#clearAgentMirror();
|
||||
await this.#ctx.session.switchSession(replicaPath);
|
||||
this.state = frame.state;
|
||||
this.#applyHostState(frame.state);
|
||||
this.#ctx.resetObserverRegistry();
|
||||
this.#applyAgentSnapshots(frame.agents);
|
||||
this.#assistantStreamSynced = false;
|
||||
setSessionTerminalTitle(frame.state.sessionName ?? frame.header.title, frame.state.cwd);
|
||||
this.#ctx.chatContainer.clear();
|
||||
this.#ctx.renderInitialMessages({ clearTerminalHistory: true });
|
||||
await this.#ctx.reloadTodos();
|
||||
this.#updateStatusSegment();
|
||||
this.#welcomed = true;
|
||||
this.#ctx.showStatus(isResync ? "Reconnected to collab session" : "Joined collab session");
|
||||
}
|
||||
|
||||
#applyFrame(frame: CollabFrame): void {
|
||||
switch (frame.t) {
|
||||
case "entry": {
|
||||
// Entries are never rendered directly — rendering is events-only
|
||||
// (prevents double-render). They keep the replica file, the agent's
|
||||
// message array (/dump, context estimates), and todos current.
|
||||
this.#ctx.sessionManager.ingestReplicatedEntry(frame.entry);
|
||||
if (frame.entry.type === "message") {
|
||||
this.#ctx.session.agent.replaceMessages([...this.#ctx.session.messages, frame.entry.message]);
|
||||
}
|
||||
break;
|
||||
}
|
||||
case "event":
|
||||
this.#applyEvent(frame.event);
|
||||
break;
|
||||
case "state": {
|
||||
this.state = frame.state;
|
||||
this.#applyHostState(frame.state);
|
||||
setSessionTerminalTitle(frame.state.sessionName, frame.state.cwd);
|
||||
this.#updateStatusSegment();
|
||||
// Reconciler: events normally drive the loader; clear a stale one if
|
||||
// the host reports idle (e.g. events lost across a reconnect).
|
||||
if (!frame.state.isStreaming && this.#ctx.loadingAnimation) {
|
||||
this.#ctx.loadingAnimation.stop();
|
||||
this.#ctx.loadingAnimation = undefined;
|
||||
}
|
||||
this.#ctx.statusLine.invalidate();
|
||||
this.#ctx.ui.requestRender();
|
||||
break;
|
||||
}
|
||||
case "bus":
|
||||
// Mirrored host EventBus traffic (task subagent lifecycle/progress)
|
||||
// feeding the observer HUD and Agent Hub progress columns.
|
||||
this.#ctx.eventBus?.emit(frame.channel, frame.data);
|
||||
break;
|
||||
case "agents":
|
||||
this.#applyAgentSnapshots(frame.agents);
|
||||
break;
|
||||
case "transcript": {
|
||||
const resolve = this.#pendingTranscripts.get(frame.reqId);
|
||||
if (resolve) {
|
||||
this.#pendingTranscripts.delete(frame.reqId);
|
||||
resolve(frame.error ? null : { text: frame.text, newSize: frame.newSize });
|
||||
}
|
||||
break;
|
||||
}
|
||||
case "bye": {
|
||||
this.#ctx.showStatus(`Collab session ended (${frame.reason})`);
|
||||
this.#socket?.close();
|
||||
void this.#restoreLocalSession();
|
||||
break;
|
||||
}
|
||||
case "error":
|
||||
this.#ctx.showError(`Collab host: ${frame.message}`);
|
||||
break;
|
||||
default:
|
||||
logger.debug("collab guest ignoring unexpected frame", { type: frame.t });
|
||||
}
|
||||
}
|
||||
|
||||
#applyEvent(event: AgentSessionEvent): void {
|
||||
// Orphan-delta guard: when joining mid-turn the message_start for the
|
||||
// in-flight assistant message predates the snapshot. message_update
|
||||
// carries the full accumulating message, so synthesize the missing start
|
||||
// before the first orphaned update; every other handler is tolerant of
|
||||
// unknown anchors (guarded by streamingComponent/pendingTools lookups).
|
||||
if (event.type === "message_start" && event.message.role === "assistant") {
|
||||
this.#assistantStreamSynced = true;
|
||||
} else if (
|
||||
event.type === "message_update" &&
|
||||
event.message.role === "assistant" &&
|
||||
!this.#assistantStreamSynced
|
||||
) {
|
||||
this.#assistantStreamSynced = true;
|
||||
void this.#ctx.eventController.handleEvent({ type: "message_start", message: event.message });
|
||||
}
|
||||
void this.#ctx.eventController.handleEvent(event);
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply the host's real model/thinking state to the replica agent so model
|
||||
* display and context-window math are native (no display-string overrides).
|
||||
* Pure agent-state mutation: session.setModel/setThinkingLevel would
|
||||
* persist entries and clamp to local credentials.
|
||||
*/
|
||||
#applyHostState(state: CollabSessionState): void {
|
||||
const session = this.#ctx.session;
|
||||
if (
|
||||
state.model &&
|
||||
(session.agent.state.model?.id !== state.model.id ||
|
||||
session.agent.state.model?.provider !== state.model.provider)
|
||||
) {
|
||||
session.agent.setModel(state.model);
|
||||
}
|
||||
const level = state.thinkingLevel as ThinkingLevel | undefined;
|
||||
session.agent.setThinkingLevel(toReasoningEffort(level));
|
||||
session.agent.setDisableReasoning(shouldDisableReasoning(level));
|
||||
}
|
||||
|
||||
/** Diff a host agent snapshot into the local registry (refs keep `session: null`). */
|
||||
#applyAgentSnapshots(agents: AgentSnapshot[]): void {
|
||||
const seen = new Set<string>();
|
||||
for (const snap of agents) seen.add(snap.id);
|
||||
for (const ref of this.agentRegistry.list()) {
|
||||
if (!seen.has(ref.id)) {
|
||||
this.agentRegistry.unregister(ref.id);
|
||||
this.#agentHasTranscript.delete(ref.id);
|
||||
}
|
||||
}
|
||||
for (const snap of agents) {
|
||||
if (this.agentRegistry.get(snap.id)) {
|
||||
this.agentRegistry.setStatus(snap.id, snap.status);
|
||||
} else {
|
||||
this.agentRegistry.register({
|
||||
id: snap.id,
|
||||
displayName: snap.displayName,
|
||||
kind: snap.kind,
|
||||
parentId: snap.parentId,
|
||||
session: null,
|
||||
status: snap.status,
|
||||
});
|
||||
}
|
||||
// Refs are returned by reference: patch host timestamps directly so
|
||||
// hub age/activity columns reflect the host, not local registration.
|
||||
const ref = this.agentRegistry.get(snap.id);
|
||||
if (ref) {
|
||||
ref.createdAt = snap.createdAt;
|
||||
ref.lastActivity = snap.lastActivity;
|
||||
ref.displayName = snap.displayName;
|
||||
}
|
||||
this.#agentHasTranscript.set(snap.id, snap.hasSessionFile);
|
||||
}
|
||||
}
|
||||
|
||||
#clearAgentMirror(): void {
|
||||
for (const ref of this.agentRegistry.list()) {
|
||||
this.agentRegistry.unregister(ref.id);
|
||||
}
|
||||
this.#agentHasTranscript.clear();
|
||||
}
|
||||
|
||||
/** Resolve every in-flight transcript request with null (resolvers clear their own timers). */
|
||||
#flushPendingTranscripts(): void {
|
||||
for (const resolve of this.#pendingTranscripts.values()) {
|
||||
resolve(null);
|
||||
}
|
||||
this.#pendingTranscripts.clear();
|
||||
}
|
||||
|
||||
#clearTransientUi(): void {
|
||||
this.#ctx.statusContainer.clear();
|
||||
this.#ctx.pendingMessagesContainer.clear();
|
||||
this.#ctx.compactionQueuedMessages = [];
|
||||
this.#ctx.streamingComponent = undefined;
|
||||
this.#ctx.streamingMessage = undefined;
|
||||
this.#ctx.pendingTools.clear();
|
||||
if (this.#ctx.loadingAnimation) {
|
||||
this.#ctx.loadingAnimation.stop();
|
||||
this.#ctx.loadingAnimation = undefined;
|
||||
}
|
||||
}
|
||||
|
||||
async #restoreLocalSession(): Promise<void> {
|
||||
if (this.#left) return;
|
||||
this.#left = true;
|
||||
this.#socket = null;
|
||||
this.#ctx.collabGuest = undefined;
|
||||
this.#ctx.statusLine.setCollabStatus(null);
|
||||
this.#flushPendingTranscripts();
|
||||
this.#clearAgentMirror();
|
||||
this.#ctx.resetObserverRegistry();
|
||||
this.#clearTransientUi();
|
||||
// Replica file stays on disk: it is a valid session file outside the
|
||||
// sessions dir, so it never shows up in /resume but remains readable.
|
||||
if (this.#returnSessionFile) {
|
||||
await this.#ctx.handleResumeSession(this.#returnSessionFile);
|
||||
return;
|
||||
}
|
||||
await this.#ctx.session.newSession();
|
||||
setSessionTerminalTitle(this.#ctx.sessionManager.getSessionName(), this.#ctx.sessionManager.getCwd());
|
||||
this.#ctx.statusLine.invalidate();
|
||||
this.#ctx.statusLine.setSessionStartTime(Date.now());
|
||||
this.#ctx.updateEditorTopBorder();
|
||||
this.#ctx.updateEditorBorderColor();
|
||||
this.#ctx.renderInitialMessages({ clearTerminalHistory: true });
|
||||
await this.#ctx.reloadTodos();
|
||||
this.#ctx.ui.requestRender(true, { clearScrollback: true });
|
||||
}
|
||||
|
||||
#updateStatusSegment(): void {
|
||||
this.#ctx.statusLine.setCollabStatus({
|
||||
role: "guest",
|
||||
participantCount: this.state?.participants.length ?? 1,
|
||||
stateOverride: this.state,
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,452 @@
|
||||
/**
|
||||
* Host side of a collab live session.
|
||||
*
|
||||
* Taps the host session's event stream and SessionManager append chokepoint,
|
||||
* broadcasting entries/events/state to guests through the relay. Guests prompt
|
||||
* and abort through us; the host machine runs the agent and tools. The host's
|
||||
* subagent ecosystem is mirrored too: task EventBus traffic (observer HUD),
|
||||
* agent-registry snapshots (Agent Hub table), hub chat/kill/revive commands,
|
||||
* and incremental subagent-transcript reads.
|
||||
*/
|
||||
import * as fs from "node:fs/promises";
|
||||
import * as os from "node:os";
|
||||
import type { ImageContent, TextContent } from "@oh-my-pi/pi-ai";
|
||||
import { logger } from "@oh-my-pi/pi-utils";
|
||||
import type { InteractiveModeContext } from "../modes/types";
|
||||
import { AgentLifecycleManager } from "../registry/agent-lifecycle";
|
||||
import { AgentRegistry } from "../registry/agent-registry";
|
||||
import type { AgentSessionEvent } from "../session/agent-session";
|
||||
import { stripImagesFromMessage, USER_INTERRUPT_LABEL } from "../session/messages";
|
||||
import { TASK_SUBAGENT_LIFECYCLE_CHANNEL, TASK_SUBAGENT_PROGRESS_CHANNEL } from "../task";
|
||||
import { generateRoomKey, importRoomKey } from "./crypto";
|
||||
import {
|
||||
type AgentSnapshot,
|
||||
COLLAB_PROMPT_MESSAGE_TYPE,
|
||||
COLLAB_PROTO,
|
||||
type CollabFrame,
|
||||
type CollabParticipant,
|
||||
type CollabSessionState,
|
||||
formatCollabLink,
|
||||
generateRoomId,
|
||||
parseCollabLink,
|
||||
} from "./protocol";
|
||||
import { CollabSocket } from "./relay-client";
|
||||
|
||||
/** Events that change the footer state guests render. */
|
||||
const STATE_TRIGGER_EVENTS: Record<string, true> = {
|
||||
agent_start: true,
|
||||
agent_end: true,
|
||||
message_end: true,
|
||||
tool_execution_end: true,
|
||||
thinking_level_changed: true,
|
||||
auto_compaction_end: true,
|
||||
};
|
||||
|
||||
const STATE_DEBOUNCE_MS = 100;
|
||||
const AGENTS_DEBOUNCE_MS = 100;
|
||||
const STREAMING_STATE_INTERVAL_MS = 2000;
|
||||
const WELCOME_IMAGE_STRIP_THRESHOLD = 24 * 1024 * 1024;
|
||||
const CONNECT_TIMEOUT_MS = 15_000;
|
||||
/** Max bytes served per fetch-transcript reply (guest re-requests from `newSize`). */
|
||||
const TRANSCRIPT_READ_CAP = 4 * 1024 * 1024;
|
||||
|
||||
/** Display name for this process's user in collab sessions. */
|
||||
export function collabDisplayName(ctx: InteractiveModeContext): string {
|
||||
const configured = (ctx.settings.get("collab.displayName") ?? "").trim();
|
||||
if (configured) return configured;
|
||||
try {
|
||||
return os.userInfo().username;
|
||||
} catch {
|
||||
return "anonymous";
|
||||
}
|
||||
}
|
||||
|
||||
export class CollabHost {
|
||||
#ctx: InteractiveModeContext;
|
||||
#socket: CollabSocket | null = null;
|
||||
#link = "";
|
||||
#sessionId = "";
|
||||
#unsubscribe?: () => void;
|
||||
#peers = new Map<number, string>();
|
||||
#lastStateJson = "";
|
||||
#stateDebounce: Timer | null = null;
|
||||
#streamingInterval: Timer | null = null;
|
||||
#agentsDebounce: Timer | null = null;
|
||||
#busUnsubscribers: (() => void)[] = [];
|
||||
#registryUnsubscribe?: () => void;
|
||||
#stopped = false;
|
||||
|
||||
constructor(ctx: InteractiveModeContext) {
|
||||
this.#ctx = ctx;
|
||||
}
|
||||
|
||||
get link(): string {
|
||||
return this.#link;
|
||||
}
|
||||
|
||||
get participants(): CollabParticipant[] {
|
||||
const list: CollabParticipant[] = [{ name: collabDisplayName(this.#ctx), role: "host" }];
|
||||
for (const name of this.#peers.values()) list.push({ name, role: "guest" });
|
||||
return list;
|
||||
}
|
||||
|
||||
async start(relayUrl: string): Promise<void> {
|
||||
const rawKey = generateRoomKey();
|
||||
const roomId = generateRoomId();
|
||||
this.#link = formatCollabLink(relayUrl, roomId, rawKey);
|
||||
const parsed = parseCollabLink(this.#link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
const key = await importRoomKey(rawKey);
|
||||
|
||||
const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "host", key });
|
||||
this.#socket = socket;
|
||||
this.#sessionId = this.#ctx.sessionManager.getSessionId();
|
||||
|
||||
const firstOpen = Promise.withResolvers<void>();
|
||||
let opened = false;
|
||||
socket.onOpen = () => {
|
||||
if (!opened) {
|
||||
opened = true;
|
||||
firstOpen.resolve();
|
||||
}
|
||||
};
|
||||
socket.onFrame = (frame, fromPeer) => this.#handleFrame(frame, fromPeer);
|
||||
socket.onControl = msg => {
|
||||
if (msg.t === "peer-left") this.#handlePeerLeft(msg.peer);
|
||||
};
|
||||
socket.onClose = (reason, willReconnect) => {
|
||||
if (this.#stopped) return;
|
||||
if (!opened) {
|
||||
firstOpen.reject(new Error(reason));
|
||||
return;
|
||||
}
|
||||
if (willReconnect) {
|
||||
this.#ctx.showStatus(`Collab relay connection lost (${reason}), reconnecting…`, { dim: true });
|
||||
} else {
|
||||
void this.#teardown();
|
||||
this.#ctx.session.emitNotice("warning", `Collab ended: ${reason}`, "collab");
|
||||
}
|
||||
};
|
||||
socket.connect();
|
||||
|
||||
const timeout = setTimeout(
|
||||
() => firstOpen.reject(new Error("timed out connecting to relay")),
|
||||
CONNECT_TIMEOUT_MS,
|
||||
);
|
||||
try {
|
||||
await firstOpen.promise;
|
||||
} catch (err) {
|
||||
this.#stopped = true;
|
||||
socket.close();
|
||||
this.#socket = null;
|
||||
throw err;
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
|
||||
this.#unsubscribe = this.#ctx.session.subscribe(event => {
|
||||
this.#broadcast({ t: "event", event });
|
||||
this.#onEventForState(event);
|
||||
});
|
||||
const bus = this.#ctx.eventBus;
|
||||
if (bus) {
|
||||
for (const channel of [TASK_SUBAGENT_LIFECYCLE_CHANNEL, TASK_SUBAGENT_PROGRESS_CHANNEL]) {
|
||||
this.#busUnsubscribers.push(bus.on(channel, data => this.#broadcast({ t: "bus", channel, data })));
|
||||
}
|
||||
}
|
||||
this.#registryUnsubscribe = AgentRegistry.global().onChange(() => this.#scheduleAgentsBroadcast());
|
||||
this.#ctx.sessionManager.onEntryAppended = entry => {
|
||||
this.#broadcast({ t: "entry", entry });
|
||||
// Model/thinking/title changes land as entries while idle; refresh
|
||||
// guest state promptly (debounce + JSON diff dedupe).
|
||||
this.#scheduleStateBroadcast();
|
||||
};
|
||||
this.#updateStatusSegment();
|
||||
}
|
||||
|
||||
/** Broadcast a goodbye, detach all taps, and close the socket. */
|
||||
async stop(reason: string): Promise<void> {
|
||||
if (this.#stopped) return;
|
||||
this.#socket?.send({ t: "bye", reason });
|
||||
await this.#teardown();
|
||||
}
|
||||
|
||||
async #teardown(): Promise<void> {
|
||||
if (this.#stopped) return;
|
||||
this.#stopped = true;
|
||||
this.#ctx.sessionManager.onEntryAppended = undefined;
|
||||
this.#unsubscribe?.();
|
||||
this.#unsubscribe = undefined;
|
||||
for (const unsubscribe of this.#busUnsubscribers) unsubscribe();
|
||||
this.#busUnsubscribers = [];
|
||||
this.#registryUnsubscribe?.();
|
||||
this.#registryUnsubscribe = undefined;
|
||||
clearTimeout(this.#stateDebounce ?? undefined);
|
||||
this.#stateDebounce = null;
|
||||
clearTimeout(this.#agentsDebounce ?? undefined);
|
||||
this.#agentsDebounce = null;
|
||||
clearInterval(this.#streamingInterval ?? undefined);
|
||||
this.#streamingInterval = null;
|
||||
this.#peers.clear();
|
||||
this.#socket?.close();
|
||||
this.#socket = null;
|
||||
this.#ctx.collabHost = undefined;
|
||||
this.#ctx.statusLine.setCollabStatus(null);
|
||||
this.#ctx.ui.requestRender();
|
||||
}
|
||||
|
||||
#broadcast(frame: CollabFrame): void {
|
||||
if (this.#stopped || !this.#socket) return;
|
||||
if (this.#ctx.sessionManager.getSessionId() !== this.#sessionId) {
|
||||
void this.stop("session switched");
|
||||
this.#ctx.session.emitNotice("warning", "Collab ended: session switched", "collab");
|
||||
return;
|
||||
}
|
||||
this.#socket.send(frame);
|
||||
}
|
||||
|
||||
#handleFrame(frame: CollabFrame, fromPeer: number): void {
|
||||
switch (frame.t) {
|
||||
case "hello":
|
||||
this.#handleHello(frame.name, frame.proto, fromPeer);
|
||||
break;
|
||||
case "prompt":
|
||||
this.#handlePrompt(frame.text, frame.images, fromPeer);
|
||||
break;
|
||||
case "abort":
|
||||
this.#handleAbort(fromPeer);
|
||||
break;
|
||||
case "agent-cmd":
|
||||
this.#handleAgentCmd(frame.cmd, frame.agentId, frame.text, fromPeer);
|
||||
break;
|
||||
case "fetch-transcript":
|
||||
void this.#handleFetchTranscript(frame.reqId, frame.agentId, frame.fromByte, fromPeer);
|
||||
break;
|
||||
default:
|
||||
logger.debug("collab host ignoring unexpected frame", { type: frame.t, fromPeer });
|
||||
}
|
||||
}
|
||||
|
||||
#handleHello(name: string, proto: number, fromPeer: number): void {
|
||||
if (proto !== COLLAB_PROTO) {
|
||||
this.#socket?.send(
|
||||
{ t: "error", message: `protocol mismatch: host speaks v${COLLAB_PROTO}, guest sent v${proto}` },
|
||||
fromPeer,
|
||||
);
|
||||
return;
|
||||
}
|
||||
const cleanName = name.trim().slice(0, 64) || `guest-${fromPeer}`;
|
||||
this.#peers.set(fromPeer, cleanName);
|
||||
|
||||
// Snapshot and send synchronously: no awaits between snapshot and send, so
|
||||
// later entries/events queue behind the welcome on the same socket and the
|
||||
// guest never sees a gap.
|
||||
const snapshot = this.#ctx.sessionManager.snapshotForReplication();
|
||||
if (JSON.stringify(snapshot).length > WELCOME_IMAGE_STRIP_THRESHOLD) {
|
||||
let stripped = 0;
|
||||
for (const entry of snapshot.entries) {
|
||||
if (entry.type === "message") stripped += stripImagesFromMessage(entry.message);
|
||||
}
|
||||
logger.info("collab welcome exceeded size threshold; stripped images", { stripped });
|
||||
}
|
||||
this.#socket?.send(
|
||||
{
|
||||
t: "welcome",
|
||||
proto: COLLAB_PROTO,
|
||||
header: snapshot.header,
|
||||
entries: snapshot.entries,
|
||||
state: this.#buildState(),
|
||||
agents: this.#snapshotAgents(),
|
||||
},
|
||||
fromPeer,
|
||||
);
|
||||
this.#ctx.session.emitNotice("info", `${cleanName} joined the collab session`, "collab");
|
||||
this.#updateStatusSegment();
|
||||
this.#scheduleStateBroadcast();
|
||||
}
|
||||
|
||||
#handlePrompt(text: string, images: ImageContent[] | undefined, fromPeer: number): void {
|
||||
const name = this.#peers.get(fromPeer) ?? `guest-${fromPeer}`;
|
||||
const content: string | (TextContent | ImageContent)[] =
|
||||
images && images.length > 0 ? [{ type: "text", text }, ...images] : text;
|
||||
this.#ctx.session
|
||||
.promptCustomMessage(
|
||||
{
|
||||
customType: COLLAB_PROMPT_MESSAGE_TYPE,
|
||||
content,
|
||||
display: true,
|
||||
details: { from: name },
|
||||
attribution: "user",
|
||||
},
|
||||
{ streamingBehavior: "steer" },
|
||||
)
|
||||
.catch(err => {
|
||||
logger.warn("collab guest prompt failed", { error: String(err) });
|
||||
this.#socket?.send({ t: "error", message: `prompt failed: ${String(err)}` }, fromPeer);
|
||||
});
|
||||
}
|
||||
|
||||
#handleAbort(fromPeer: number): void {
|
||||
const name = this.#peers.get(fromPeer) ?? `guest-${fromPeer}`;
|
||||
void this.#ctx.session
|
||||
.abort()
|
||||
.then(() => this.#ctx.session.emitNotice("info", `${name} interrupted`, "collab"))
|
||||
.catch(err => logger.warn("collab guest abort failed", { error: String(err) }));
|
||||
}
|
||||
|
||||
#handlePeerLeft(peer: number): void {
|
||||
const name = this.#peers.get(peer);
|
||||
this.#peers.delete(peer);
|
||||
if (name) this.#ctx.session.emitNotice("info", `${name} left the collab session`, "collab");
|
||||
this.#updateStatusSegment();
|
||||
this.#scheduleStateBroadcast();
|
||||
}
|
||||
|
||||
#buildState(): CollabSessionState {
|
||||
const session = this.#ctx.session;
|
||||
// Context numbers come from the status line's breakdown — not
|
||||
// session.getContextUsage() — so guests render exactly what the host's
|
||||
// own footer shows.
|
||||
const breakdown = this.#ctx.statusLine.getCachedContextBreakdown();
|
||||
return {
|
||||
isStreaming: session.isStreaming,
|
||||
queuedMessageCount: session.queuedMessageCount,
|
||||
sessionName: session.sessionName,
|
||||
cwd: this.#ctx.sessionManager.getCwd(),
|
||||
model: session.model,
|
||||
thinkingLevel: session.thinkingLevel,
|
||||
contextUsage: {
|
||||
tokens: breakdown.usedTokens,
|
||||
contextWindow: breakdown.contextWindow,
|
||||
percent: breakdown.contextWindow > 0 ? (breakdown.usedTokens / breakdown.contextWindow) * 100 : null,
|
||||
},
|
||||
participants: this.participants,
|
||||
};
|
||||
}
|
||||
|
||||
#onEventForState(event: AgentSessionEvent): void {
|
||||
if (!STATE_TRIGGER_EVENTS[event.type]) return;
|
||||
this.#scheduleStateBroadcast();
|
||||
if (event.type === "agent_start" && !this.#streamingInterval) {
|
||||
this.#streamingInterval = setInterval(() => this.#scheduleStateBroadcast(), STREAMING_STATE_INTERVAL_MS);
|
||||
} else if (event.type === "agent_end" && this.#streamingInterval) {
|
||||
clearInterval(this.#streamingInterval);
|
||||
this.#streamingInterval = null;
|
||||
}
|
||||
}
|
||||
|
||||
#snapshotAgents(): AgentSnapshot[] {
|
||||
return AgentRegistry.global()
|
||||
.list()
|
||||
.map(ref => ({
|
||||
id: ref.id,
|
||||
displayName: ref.displayName,
|
||||
kind: ref.kind,
|
||||
parentId: ref.parentId,
|
||||
status: ref.status,
|
||||
hasSessionFile: !!ref.sessionFile,
|
||||
createdAt: ref.createdAt,
|
||||
lastActivity: ref.lastActivity,
|
||||
}));
|
||||
}
|
||||
|
||||
#scheduleAgentsBroadcast(): void {
|
||||
if (this.#stopped || this.#agentsDebounce) return;
|
||||
this.#agentsDebounce = setTimeout(() => {
|
||||
this.#agentsDebounce = null;
|
||||
this.#broadcast({ t: "agents", agents: this.#snapshotAgents() });
|
||||
}, AGENTS_DEBOUNCE_MS);
|
||||
}
|
||||
|
||||
#handleAgentCmd(cmd: "chat" | "kill" | "revive", agentId: string, text: string | undefined, fromPeer: number): void {
|
||||
const fail = (err: unknown) => {
|
||||
logger.warn("collab agent-cmd failed", { cmd, agentId, error: String(err) });
|
||||
this.#socket?.send({ t: "error", message: `agent ${agentId}: ${String(err)}` }, fromPeer);
|
||||
};
|
||||
switch (cmd) {
|
||||
case "chat": {
|
||||
const trimmed = text?.trim();
|
||||
if (!trimmed) {
|
||||
this.#socket?.send({ t: "error", message: `agent ${agentId}: empty chat message` }, fromPeer);
|
||||
return;
|
||||
}
|
||||
// Mirrors the hub's #submitChatMessage: revive if parked, steer if mid-turn.
|
||||
AgentLifecycleManager.global()
|
||||
.ensureLive(agentId)
|
||||
.then(session => session.prompt(trimmed, { streamingBehavior: "steer" }))
|
||||
.catch(fail);
|
||||
break;
|
||||
}
|
||||
case "kill": {
|
||||
const kill = async () => {
|
||||
const ref = AgentRegistry.global().get(agentId);
|
||||
if (ref && ref.status === "running" && ref.session) {
|
||||
await ref.session.abort({ reason: USER_INTERRUPT_LABEL });
|
||||
}
|
||||
await AgentLifecycleManager.global().release(agentId);
|
||||
};
|
||||
kill().catch(fail);
|
||||
break;
|
||||
}
|
||||
case "revive":
|
||||
AgentLifecycleManager.global().ensureLive(agentId).catch(fail);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/** Incremental transcript read mirroring the hub's readFileIncremental contract. */
|
||||
async #handleFetchTranscript(reqId: number, agentId: string, fromByte: number, fromPeer: number): Promise<void> {
|
||||
const reply = (text: string, newSize: number, error?: string) =>
|
||||
this.#socket?.send({ t: "transcript", reqId, text, newSize, error }, fromPeer);
|
||||
const file = AgentRegistry.global().get(agentId)?.sessionFile;
|
||||
if (!file) {
|
||||
reply("", fromByte, "no transcript available");
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const stat = await fs.stat(file);
|
||||
if (stat.size <= fromByte) {
|
||||
reply("", stat.size);
|
||||
return;
|
||||
}
|
||||
const want = Math.min(stat.size - fromByte, TRANSCRIPT_READ_CAP);
|
||||
const handle = await fs.open(file, "r");
|
||||
let bytesRead: number;
|
||||
const buf = Buffer.allocUnsafe(want);
|
||||
try {
|
||||
({ bytesRead } = await handle.read(buf, 0, want, fromByte));
|
||||
} finally {
|
||||
await handle.close();
|
||||
}
|
||||
let slice = buf.subarray(0, bytesRead);
|
||||
const reachedEof = fromByte + bytesRead >= stat.size;
|
||||
if (!reachedEof) {
|
||||
// Trim to the last complete JSONL line so no line or UTF-8 char is split.
|
||||
const lastNewline = slice.lastIndexOf(0x0a);
|
||||
slice = slice.subarray(0, lastNewline >= 0 ? lastNewline + 1 : 0);
|
||||
}
|
||||
reply(slice.toString("utf-8"), reachedEof ? stat.size : fromByte + slice.byteLength);
|
||||
} catch (err) {
|
||||
logger.debug("collab transcript read failed", { agentId, error: String(err) });
|
||||
reply("", fromByte, String(err));
|
||||
}
|
||||
}
|
||||
|
||||
#scheduleStateBroadcast(): void {
|
||||
if (this.#stopped || this.#stateDebounce) return;
|
||||
this.#stateDebounce = setTimeout(() => {
|
||||
this.#stateDebounce = null;
|
||||
const state = this.#buildState();
|
||||
const json = JSON.stringify(state);
|
||||
if (json === this.#lastStateJson) return;
|
||||
this.#lastStateJson = json;
|
||||
this.#broadcast({ t: "state", state });
|
||||
}, STATE_DEBOUNCE_MS);
|
||||
}
|
||||
|
||||
#updateStatusSegment(): void {
|
||||
this.#ctx.statusLine.setCollabStatus({ role: "host", participantCount: this.#peers.size + 1 });
|
||||
this.#ctx.statusLine.invalidate();
|
||||
this.#ctx.ui.requestRender();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,225 @@
|
||||
/**
|
||||
* Collab live-session wire protocol.
|
||||
*
|
||||
* Hub topology: the host is authoritative, guests never peer. All session
|
||||
* payloads (`CollabFrame`) travel AES-256-GCM sealed; the relay only sees the
|
||||
* plaintext envelope (`[4B uint32 BE peerId][sealed payload]`) plus TEXT JSON
|
||||
* control messages that carry no session data.
|
||||
*/
|
||||
import type { ImageContent, Model } from "@oh-my-pi/pi-ai";
|
||||
import type { ContextUsage } from "../extensibility/extensions/types";
|
||||
import type { AgentSessionEvent } from "../session/agent-session";
|
||||
import type { SessionEntry, SessionHeader } from "../session/session-manager";
|
||||
|
||||
export const COLLAB_PROTO = 1;
|
||||
|
||||
/** customType of guest prompts injected on the host (rendered with an author badge). */
|
||||
export const COLLAB_PROMPT_MESSAGE_TYPE = "collab-prompt";
|
||||
|
||||
/** Display metadata attached to collab guest prompts. */
|
||||
export interface CollabPromptDetails {
|
||||
from?: string;
|
||||
}
|
||||
|
||||
export interface CollabParticipant {
|
||||
name: string;
|
||||
role: "host" | "guest";
|
||||
}
|
||||
|
||||
/** Serializable mirror of an {@link AgentRef} (live session handle stripped). */
|
||||
export interface AgentSnapshot {
|
||||
id: string;
|
||||
displayName: string;
|
||||
kind: "main" | "sub";
|
||||
parentId?: string;
|
||||
status: "running" | "idle" | "parked" | "aborted";
|
||||
/** Whether the host has a transcript file for this agent (gates remote transcript fetch). */
|
||||
hasSessionFile: boolean;
|
||||
createdAt: number;
|
||||
lastActivity: number;
|
||||
}
|
||||
|
||||
/** Debounced footer snapshot broadcast by the host. */
|
||||
export interface CollabSessionState {
|
||||
isStreaming: boolean;
|
||||
queuedMessageCount: number;
|
||||
sessionName?: string;
|
||||
/** Host cwd — display/title/relativization only; guest never chdirs. */
|
||||
cwd: string;
|
||||
/**
|
||||
* Host model (full catalog object). Guests apply it to their replica
|
||||
* agent state so model display and context-window math are native.
|
||||
*/
|
||||
model?: Model;
|
||||
/** Host effective thinking level (ThinkingLevel value). */
|
||||
thinkingLevel?: string;
|
||||
/** Host status-line context numbers (guest system prompt/tools differ, so local estimates drift). */
|
||||
contextUsage?: ContextUsage;
|
||||
participants: CollabParticipant[];
|
||||
}
|
||||
|
||||
/** Encrypted payload frames (inside AES-GCM, JSON). */
|
||||
export type CollabFrame =
|
||||
// guest -> host
|
||||
| { t: "hello"; proto: number; name: string }
|
||||
| { t: "prompt"; text: string; images?: ImageContent[] }
|
||||
| { t: "abort" }
|
||||
/** Agent Hub action routed to the host (chat requires `text`). */
|
||||
| { t: "agent-cmd"; cmd: "chat" | "kill" | "revive"; agentId: string; text?: string }
|
||||
/** Incremental subagent-transcript read (mirrors the hub's readFileIncremental contract). */
|
||||
| { t: "fetch-transcript"; reqId: number; agentId: string; fromByte: number }
|
||||
// host -> guest
|
||||
| {
|
||||
t: "welcome";
|
||||
proto: number;
|
||||
header: SessionHeader;
|
||||
entries: SessionEntry[];
|
||||
state: CollabSessionState;
|
||||
agents: AgentSnapshot[];
|
||||
}
|
||||
| { t: "entry"; entry: SessionEntry }
|
||||
| { t: "event"; event: AgentSessionEvent }
|
||||
| { t: "state"; state: CollabSessionState }
|
||||
/** Mirrored EventBus traffic (task subagent lifecycle/progress channels only). */
|
||||
| { t: "bus"; channel: string; data: unknown }
|
||||
/** Full agent-registry snapshot (debounced on registry change). */
|
||||
| { t: "agents"; agents: AgentSnapshot[] }
|
||||
/** Targeted reply to fetch-transcript; `text` is decoded JSONL from `fromByte`, `newSize` the next offset base. */
|
||||
| { t: "transcript"; reqId: number; text: string; newSize: number; error?: string }
|
||||
| { t: "bye"; reason: string }
|
||||
| { t: "error"; message: string };
|
||||
|
||||
/** Relay → host control message (TEXT JSON, unencrypted, no session data). */
|
||||
export type RelayControlToHost = { t: "peer-joined" | "peer-left"; peer: number };
|
||||
/** Relay → guest control message (TEXT JSON, unencrypted). */
|
||||
export type RelayControlToGuest = { t: "room-closed" };
|
||||
export type RelayControlMessage = RelayControlToHost | RelayControlToGuest;
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// Wire envelope: [4B uint32 BE peerId][sealed payload]
|
||||
// Host→relay: peerId 0 broadcasts to all guests; peerId N targets guest N.
|
||||
// Guest→relay: always 0; the relay rewrites it to the sender's id.
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
export const ENVELOPE_HEADER_LENGTH = 4;
|
||||
|
||||
export function packEnvelope(peerId: number, sealed: Uint8Array): Uint8Array {
|
||||
const out = new Uint8Array(ENVELOPE_HEADER_LENGTH + sealed.byteLength);
|
||||
new DataView(out.buffer).setUint32(0, peerId, false);
|
||||
out.set(sealed, ENVELOPE_HEADER_LENGTH);
|
||||
return out;
|
||||
}
|
||||
|
||||
export function unpackEnvelope(data: Uint8Array): { peerId: number; payload: Uint8Array } | null {
|
||||
if (data.byteLength < ENVELOPE_HEADER_LENGTH) return null;
|
||||
const peerId = new DataView(data.buffer, data.byteOffset, ENVELOPE_HEADER_LENGTH).getUint32(0, false);
|
||||
return { peerId, payload: data.subarray(ENVELOPE_HEADER_LENGTH) };
|
||||
}
|
||||
|
||||
/** Rewrite the peerId in place without copying the payload. */
|
||||
export function rewriteEnvelopePeer(data: Uint8Array, peerId: number): void {
|
||||
new DataView(data.buffer, data.byteOffset, ENVELOPE_HEADER_LENGTH).setUint32(0, peerId, false);
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// Link format: wss://<host[:port]>/r/<roomId>#<base64url-32-byte-key>
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
export const ROOM_ID_BYTES = 16;
|
||||
|
||||
/** Default public relay; bare `<roomId>#<key>` links resolve against it. */
|
||||
export const DEFAULT_RELAY_URL = "wss://relay.omp.sh";
|
||||
|
||||
const ROOM_PATH_RE = /^\/r\/([A-Za-z0-9_-]{10,64})$/;
|
||||
const BARE_LINK_RE = /^([A-Za-z0-9_-]{10,64})#([A-Za-z0-9_-]+)$/;
|
||||
const B64URL_RE = /^[A-Za-z0-9_-]+$/;
|
||||
const LOCAL_HOSTNAMES: Record<string, true> = { localhost: true, "127.0.0.1": true, "::1": true, "[::1]": true };
|
||||
|
||||
export interface ParsedCollabLink {
|
||||
/** wss://host[:port]/r/<roomId> — no query, no fragment. */
|
||||
wsUrl: string;
|
||||
roomId: string;
|
||||
key: Uint8Array;
|
||||
}
|
||||
|
||||
export function generateRoomId(): string {
|
||||
const bytes = new Uint8Array(ROOM_ID_BYTES);
|
||||
crypto.getRandomValues(bytes);
|
||||
return Buffer.from(bytes).toString("base64url");
|
||||
}
|
||||
|
||||
/** Normalize a relay base URL (ws/wss/http/https) into a ws/wss origin, or an error. */
|
||||
function normalizeRelayOrigin(relayUrl: string): { origin: string } | { error: string } {
|
||||
let url: URL;
|
||||
try {
|
||||
url = new URL(relayUrl);
|
||||
} catch {
|
||||
return { error: `Invalid relay URL: ${relayUrl}` };
|
||||
}
|
||||
let scheme: string;
|
||||
switch (url.protocol) {
|
||||
case "wss:":
|
||||
case "https:":
|
||||
scheme = "wss:";
|
||||
break;
|
||||
case "ws:":
|
||||
case "http:":
|
||||
scheme = "ws:";
|
||||
break;
|
||||
default:
|
||||
return { error: `Unsupported relay URL scheme: ${url.protocol}` };
|
||||
}
|
||||
if (scheme === "ws:" && !LOCAL_HOSTNAMES[url.hostname]) {
|
||||
return { error: "relay link must be wss:// (plain ws:// is only allowed for localhost)" };
|
||||
}
|
||||
const port = url.port ? `:${url.port}` : "";
|
||||
return { origin: `${scheme}//${url.hostname}${port}` };
|
||||
}
|
||||
|
||||
/**
|
||||
* Render the shareable link. Compact forms: the default relay collapses to
|
||||
* `<roomId>#<key>`, other wss relays drop the scheme (`host[:port]/r/…`);
|
||||
* only localhost ws:// links keep their full URL so parsing cannot
|
||||
* mis-infer wss.
|
||||
*/
|
||||
export function formatCollabLink(relayUrl: string, roomId: string, key: Uint8Array): string {
|
||||
const normalized = normalizeRelayOrigin(relayUrl);
|
||||
if ("error" in normalized) throw new Error(normalized.error);
|
||||
const keyText = Buffer.from(key).toString("base64url");
|
||||
if (normalized.origin === DEFAULT_RELAY_URL) return `${roomId}#${keyText}`;
|
||||
const compact = normalized.origin.startsWith("wss://")
|
||||
? normalized.origin.slice("wss://".length)
|
||||
: normalized.origin;
|
||||
return `${compact}/r/${roomId}#${keyText}`;
|
||||
}
|
||||
|
||||
export function parseCollabLink(link: string): ParsedCollabLink | { error: string } {
|
||||
let text = link.trim();
|
||||
// Bare `<roomId>#<key>` → default relay.
|
||||
const bare = BARE_LINK_RE.exec(text);
|
||||
if (bare) text = `${DEFAULT_RELAY_URL}/r/${bare[1]}#${bare[2]}`;
|
||||
// Scheme-less `host[:port]/r/…` → wss.
|
||||
else if (!text.includes("://")) text = `wss://${text}`;
|
||||
let url: URL;
|
||||
try {
|
||||
url = new URL(text);
|
||||
} catch {
|
||||
return { error: `Invalid collab link: ${link}` };
|
||||
}
|
||||
const normalized = normalizeRelayOrigin(url.origin);
|
||||
if ("error" in normalized) return normalized;
|
||||
const match = ROOM_PATH_RE.exec(url.pathname);
|
||||
if (!match) {
|
||||
return { error: "Collab link must contain a /r/<roomId> path" };
|
||||
}
|
||||
const roomId = match[1]!;
|
||||
const fragment = url.hash.startsWith("#") ? url.hash.slice(1) : url.hash;
|
||||
if (!fragment) {
|
||||
return { error: "Collab link is missing the #<key> fragment" };
|
||||
}
|
||||
const key = B64URL_RE.test(fragment) ? new Uint8Array(Buffer.from(fragment, "base64url")) : null;
|
||||
if (key?.byteLength !== 32) {
|
||||
return { error: "Collab link key must be 32 base64url bytes" };
|
||||
}
|
||||
return { wsUrl: `${normalized.origin}/r/${roomId}`, roomId, key };
|
||||
}
|
||||
@@ -0,0 +1,216 @@
|
||||
/**
|
||||
* Client-side WebSocket wrapper for collab live-session sharing.
|
||||
*
|
||||
* Connects to a relay room, seals/opens AES-GCM frames, and reconnects with
|
||||
* exponential backoff on transient drops. Fatal relay close codes (room gone,
|
||||
* host conflict, room full) and decryption failures never reconnect.
|
||||
*/
|
||||
import { logger } from "@oh-my-pi/pi-utils";
|
||||
import { open, seal } from "./crypto";
|
||||
import type { CollabFrame, RelayControlMessage } from "./protocol";
|
||||
import { packEnvelope, unpackEnvelope } from "./protocol";
|
||||
|
||||
const FATAL_CLOSE_REASONS: Record<number, string> = {
|
||||
4001: "room closed",
|
||||
4004: "no such room",
|
||||
4009: "a host is already connected for this room",
|
||||
4029: "room is full",
|
||||
};
|
||||
|
||||
const BACKOFF_BASE_MS = 1_000;
|
||||
const BACKOFF_MAX_MS = 30_000;
|
||||
/** Max enveloped frames buffered while a reconnect is pending; overflow is dropped. */
|
||||
const MAX_PENDING_SENDS = 256;
|
||||
|
||||
export interface CollabSocketOptions {
|
||||
/** wss://host[:port]/r/<roomId> — no query string. */
|
||||
wsUrl: string;
|
||||
role: "host" | "guest";
|
||||
key: CryptoKey;
|
||||
}
|
||||
|
||||
export class CollabSocket {
|
||||
/** Fires after every successful (re)connect. */
|
||||
onOpen?: () => void;
|
||||
onFrame?: (frame: CollabFrame, fromPeer: number) => void;
|
||||
onControl?: (msg: RelayControlMessage) => void;
|
||||
/** Fires once per terminal close (intentional, fatal code, or bad key). willReconnect=true for transient drops that will retry. */
|
||||
onClose?: (reason: string, willReconnect: boolean) => void;
|
||||
|
||||
readonly #opts: CollabSocketOptions;
|
||||
#ws: WebSocket | null = null;
|
||||
#retryTimer: NodeJS.Timeout | undefined;
|
||||
#attempt = 0;
|
||||
/** Terminal state: intentional close or fatal failure. Cleared by connect(). */
|
||||
#closed = false;
|
||||
/** Serializes seal() so frames hit the wire in send() order. */
|
||||
#sendChain: Promise<void> = Promise.resolve();
|
||||
/** Serializes open() so frames are delivered in arrival order. */
|
||||
#recvChain: Promise<void> = Promise.resolve();
|
||||
/** Envelopes sealed while disconnected, flushed on the next open. */
|
||||
#pendingSends: Uint8Array[] = [];
|
||||
|
||||
constructor(opts: CollabSocketOptions) {
|
||||
this.#opts = opts;
|
||||
}
|
||||
|
||||
get isOpen(): boolean {
|
||||
return this.#ws?.readyState === WebSocket.OPEN;
|
||||
}
|
||||
|
||||
connect(): void {
|
||||
if (this.#ws || this.#retryTimer) return;
|
||||
this.#closed = false;
|
||||
this.#attempt = 0;
|
||||
this.#openSocket();
|
||||
}
|
||||
|
||||
send(frame: CollabFrame, targetPeer = 0): void {
|
||||
this.#sendChain = this.#sendChain
|
||||
.then(async () => {
|
||||
if (this.#closed) {
|
||||
logger.debug("collab: dropping frame, socket closed", { t: frame.t });
|
||||
return;
|
||||
}
|
||||
const sealed = await seal(this.#opts.key, frame);
|
||||
const envelope = packEnvelope(targetPeer, sealed);
|
||||
const ws = this.#ws;
|
||||
if (ws && ws.readyState === WebSocket.OPEN) {
|
||||
ws.send(envelope);
|
||||
return;
|
||||
}
|
||||
if (this.#pendingSends.length >= MAX_PENDING_SENDS) {
|
||||
logger.debug("collab: dropping frame, reconnect buffer full", { t: frame.t });
|
||||
return;
|
||||
}
|
||||
this.#pendingSends.push(envelope);
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
logger.debug("collab: send failed", { error: String(err) });
|
||||
});
|
||||
}
|
||||
|
||||
/** Intentional close: clears any retry timer, suppresses reconnect. A later connect() starts fresh. */
|
||||
close(): void {
|
||||
const hadActivity = this.#ws !== null || this.#retryTimer !== undefined;
|
||||
this.#clearRetry();
|
||||
const wasClosed = this.#closed;
|
||||
this.#closed = true;
|
||||
this.#pendingSends.length = 0;
|
||||
const ws = this.#ws;
|
||||
this.#ws = null;
|
||||
if (ws) {
|
||||
try {
|
||||
ws.close(1000);
|
||||
} catch {
|
||||
// already closing/closed
|
||||
}
|
||||
}
|
||||
if (hadActivity && !wasClosed) this.onClose?.("closed", false);
|
||||
}
|
||||
|
||||
#openSocket(): void {
|
||||
const ws = new WebSocket(`${this.#opts.wsUrl}?role=${this.#opts.role}`);
|
||||
ws.binaryType = "arraybuffer";
|
||||
this.#ws = ws;
|
||||
ws.onopen = () => {
|
||||
if (this.#ws !== ws) return;
|
||||
this.#attempt = 0;
|
||||
for (const envelope of this.#pendingSends) ws.send(envelope);
|
||||
this.#pendingSends.length = 0;
|
||||
this.onOpen?.();
|
||||
};
|
||||
ws.onmessage = (event: MessageEvent) => {
|
||||
if (this.#ws !== ws) return;
|
||||
this.#handleMessage(ws, event.data);
|
||||
};
|
||||
ws.onerror = () => {
|
||||
// The paired close event carries the actionable state; nothing to do here.
|
||||
};
|
||||
ws.onclose = (event: CloseEvent) => {
|
||||
if (this.#ws !== ws) return;
|
||||
this.#ws = null;
|
||||
this.#handleClose(event.code, event.reason);
|
||||
};
|
||||
}
|
||||
|
||||
#handleMessage(ws: WebSocket, data: unknown): void {
|
||||
if (typeof data === "string") {
|
||||
try {
|
||||
this.onControl?.(JSON.parse(data) as RelayControlMessage);
|
||||
} catch {
|
||||
logger.debug("collab: ignoring malformed control message");
|
||||
}
|
||||
return;
|
||||
}
|
||||
const bytes = data instanceof ArrayBuffer ? new Uint8Array(data) : data instanceof Uint8Array ? data : null;
|
||||
if (!bytes) return;
|
||||
const envelope = unpackEnvelope(bytes);
|
||||
if (!envelope) return;
|
||||
this.#recvChain = this.#recvChain
|
||||
.then(async () => {
|
||||
if (this.#ws !== ws) return;
|
||||
let frame: CollabFrame;
|
||||
try {
|
||||
frame = await open(this.#opts.key, envelope.payload);
|
||||
} catch {
|
||||
this.#failFatal("bad key or corrupted frame");
|
||||
return;
|
||||
}
|
||||
if (this.#ws !== ws) return;
|
||||
this.onFrame?.(frame, envelope.peerId);
|
||||
})
|
||||
.catch((err: unknown) => {
|
||||
logger.debug("collab: frame handler failed", { error: String(err) });
|
||||
});
|
||||
}
|
||||
|
||||
#handleClose(code: number, reason: string): void {
|
||||
if (this.#closed) return;
|
||||
const fatalReason = FATAL_CLOSE_REASONS[code];
|
||||
if (fatalReason !== undefined) {
|
||||
this.#closed = true;
|
||||
this.#pendingSends.length = 0;
|
||||
this.onClose?.(fatalReason, false);
|
||||
return;
|
||||
}
|
||||
this.onClose?.(reason || `connection lost (code ${code})`, true);
|
||||
this.#scheduleRetry();
|
||||
}
|
||||
|
||||
/** Decryption failure: wrong key or corrupted frame. Never reconnect. */
|
||||
#failFatal(reason: string): void {
|
||||
if (this.#closed) return;
|
||||
this.#closed = true;
|
||||
this.#clearRetry();
|
||||
this.#pendingSends.length = 0;
|
||||
const ws = this.#ws;
|
||||
this.#ws = null;
|
||||
if (ws) {
|
||||
try {
|
||||
ws.close(1000);
|
||||
} catch {
|
||||
// already closing/closed
|
||||
}
|
||||
}
|
||||
this.onClose?.(reason, false);
|
||||
}
|
||||
|
||||
#scheduleRetry(): void {
|
||||
const base = Math.min(BACKOFF_BASE_MS * 2 ** this.#attempt, BACKOFF_MAX_MS);
|
||||
this.#attempt++;
|
||||
const delay = base * (0.75 + Math.random() * 0.5);
|
||||
this.#retryTimer = setTimeout(() => {
|
||||
this.#retryTimer = undefined;
|
||||
if (this.#closed) return;
|
||||
this.#openSocket();
|
||||
}, delay);
|
||||
}
|
||||
|
||||
#clearRetry(): void {
|
||||
if (this.#retryTimer !== undefined) {
|
||||
clearTimeout(this.#retryTimer);
|
||||
this.#retryTimer = undefined;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
/**
|
||||
* Join a shared collab session from the CLI: launches the interactive TUI and
|
||||
* immediately runs `/join <link>`.
|
||||
*/
|
||||
import { APP_NAME } from "@oh-my-pi/pi-utils";
|
||||
import { Args, Command } from "@oh-my-pi/pi-utils/cli";
|
||||
import { parseArgs } from "../cli/args";
|
||||
import { runRootCommand } from "../main";
|
||||
|
||||
export default class Join extends Command {
|
||||
static description = "Join a shared collab session (same as /join)";
|
||||
|
||||
static args = {
|
||||
link: Args.string({
|
||||
description: "Collab link shared by the host (/collab)",
|
||||
required: true,
|
||||
}),
|
||||
};
|
||||
|
||||
static examples = [`${APP_NAME} join wss://relay.omp.sh/s/abc123#key`];
|
||||
|
||||
async run(): Promise<void> {
|
||||
const { args } = await this.parse(Join);
|
||||
const link = args.link?.trim();
|
||||
if (!link) {
|
||||
process.stderr.write(`Usage: ${APP_NAME} join <link>\n`);
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
if (!process.stdin.isTTY || !process.stdout.isTTY) {
|
||||
process.stderr.write(`${APP_NAME} join requires an interactive terminal\n`);
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
const parsed = parseArgs([]);
|
||||
parsed.join = link;
|
||||
await runRootCommand(parsed, []);
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
import { THINKING_EFFORTS } from "@oh-my-pi/pi-ai";
|
||||
import { SHAPE_VARIANT_NAMES } from "@oh-my-pi/snapcompact";
|
||||
import { DEFAULT_RELAY_URL } from "../collab/protocol";
|
||||
import { AUTO_THINKING, getConfiguredThinkingLevelMetadata, getThinkingLevelMetadata } from "../thinking";
|
||||
import {
|
||||
TINY_MODEL_DEVICE_DEFAULT,
|
||||
@@ -101,6 +102,7 @@ export const TAB_GROUPS: Record<SettingTab, readonly string[]> = {
|
||||
"Approvals",
|
||||
"Notifications",
|
||||
"Speech",
|
||||
"Collab",
|
||||
"Magic Keywords",
|
||||
"Startup & Updates",
|
||||
"Power (macOS)",
|
||||
@@ -147,7 +149,8 @@ export type StatusLineSegmentId =
|
||||
| "cache_write"
|
||||
| "cache_hit"
|
||||
| "session_name"
|
||||
| "usage";
|
||||
| "usage"
|
||||
| "collab";
|
||||
|
||||
/** Submenu choice metadata. */
|
||||
export type SubmenuOption<V extends string = string> = {
|
||||
@@ -1327,6 +1330,29 @@ export const SETTINGS_SCHEMA = {
|
||||
},
|
||||
},
|
||||
|
||||
// Collab
|
||||
"collab.relayUrl": {
|
||||
type: "string",
|
||||
default: DEFAULT_RELAY_URL,
|
||||
ui: {
|
||||
tab: "interaction",
|
||||
group: "Collab",
|
||||
label: "Relay URL",
|
||||
description: "Relay used by /collab (wss://host[:port]; self-host with the omp-collab-relay service)",
|
||||
},
|
||||
},
|
||||
|
||||
"collab.displayName": {
|
||||
type: "string",
|
||||
default: "",
|
||||
ui: {
|
||||
tab: "interaction",
|
||||
group: "Collab",
|
||||
label: "Display Name",
|
||||
description: "Name shown to other collab participants (default: OS username)",
|
||||
},
|
||||
},
|
||||
|
||||
// Speech-to-text
|
||||
"stt.enabled": {
|
||||
type: "boolean",
|
||||
|
||||
@@ -66,6 +66,7 @@ import {
|
||||
import type { AgentSession } from "./session/agent-session";
|
||||
import type { AuthStorage } from "./session/auth-storage";
|
||||
import { resolveResumableSession, type SessionInfo, SessionManager } from "./session/session-manager";
|
||||
import { executeBuiltinSlashCommand } from "./slash-commands/builtin-registry";
|
||||
import { discoverTitleSystemPromptFile, resolvePromptInput } from "./system-prompt";
|
||||
import { initTelemetryExport, isTelemetryExportEnabled } from "./telemetry-export";
|
||||
import { AUTO_THINKING } from "./thinking";
|
||||
@@ -346,6 +347,7 @@ async function runInteractiveMode(
|
||||
initialMessage?: string,
|
||||
initialImages?: ImageContent[],
|
||||
titleSystemPrompt?: string,
|
||||
joinLink?: string,
|
||||
): Promise<void> {
|
||||
const mode = new InteractiveMode(
|
||||
session,
|
||||
@@ -414,6 +416,12 @@ async function runInteractiveMode(
|
||||
}
|
||||
}
|
||||
|
||||
// `omp join <link>`: dispatch through the same builtin path as a typed
|
||||
// `/join` so collab guards and error rendering stay in one place.
|
||||
if (joinLink !== undefined) {
|
||||
await executeBuiltinSlashCommand(`/join ${joinLink}`, { ctx: mode });
|
||||
}
|
||||
|
||||
if (initialMessage !== undefined) {
|
||||
try {
|
||||
using _keepalive = new EventLoopKeepalive();
|
||||
@@ -1298,6 +1306,7 @@ export async function runRootCommand(
|
||||
initialMessage,
|
||||
initialImages,
|
||||
titleSystemPrompt,
|
||||
parsedArgs.join,
|
||||
);
|
||||
} else {
|
||||
// Branch-only single-shot runner: keep print-mode code out of normal interactive startup.
|
||||
|
||||
@@ -81,6 +81,15 @@ function statusBadge(status: AgentStatus): string {
|
||||
}
|
||||
}
|
||||
|
||||
/** Guest-side proxy for hub actions executed on the collab host. */
|
||||
export interface AgentHubRemote {
|
||||
chat(id: string, text: string): void;
|
||||
kill(id: string): void;
|
||||
revive(id: string): void;
|
||||
/** Mirrors readFileIncremental: text from fromByte (complete JSONL lines), newSize = next fromByte base; null = unavailable. */
|
||||
readTranscript(id: string, fromByte: number): Promise<{ text: string; newSize: number } | null>;
|
||||
}
|
||||
|
||||
export interface AgentHubDeps {
|
||||
/** Progress/status snapshot source (task lifecycle + progress channels). */
|
||||
observers: SessionObserverRegistry;
|
||||
@@ -94,6 +103,8 @@ export interface AgentHubDeps {
|
||||
lifecycle?: AgentLifecycleManager;
|
||||
/** Injectable for tests; defaults to the process-global bus. */
|
||||
irc?: IrcBus;
|
||||
/** Collab guest: route actions/transcripts to the host instead of local sessions. */
|
||||
remote?: AgentHubRemote;
|
||||
}
|
||||
|
||||
export class AgentHubOverlayComponent extends Container {
|
||||
@@ -106,6 +117,11 @@ export class AgentHubOverlayComponent extends Container {
|
||||
#hubKeys: KeyId[];
|
||||
#unsubscribers: Array<() => void> = [];
|
||||
#ageTimer: NodeJS.Timeout | undefined;
|
||||
#remote: AgentHubRemote | undefined;
|
||||
#remoteFetchInFlight = false;
|
||||
/** Invalidates stale in-flight fetch callbacks after openChat resets the cache. */
|
||||
#remoteFetchToken = 0;
|
||||
#remoteTranscriptUnavailable = false;
|
||||
|
||||
// Table state
|
||||
#view: "table" | "chat" = "table";
|
||||
@@ -143,6 +159,7 @@ export class AgentHubOverlayComponent extends Container {
|
||||
this.#onDone = deps.onDone;
|
||||
this.#requestRender = deps.requestRender;
|
||||
this.#hubKeys = deps.hubKeys;
|
||||
this.#remote = deps.remote;
|
||||
|
||||
this.#editor = new Editor(getEditorTheme());
|
||||
this.#editor.setMaxHeight(4);
|
||||
@@ -196,6 +213,9 @@ export class AgentHubOverlayComponent extends Container {
|
||||
this.#chatAgentId = id;
|
||||
this.#notice = undefined;
|
||||
this.#transcriptCache = undefined;
|
||||
this.#remoteTranscriptUnavailable = false;
|
||||
this.#remoteFetchInFlight = false;
|
||||
this.#remoteFetchToken++;
|
||||
this.#scrollOffset = 0;
|
||||
this.#selectedEntryIndex = 0;
|
||||
this.#expandedEntries.clear();
|
||||
@@ -238,6 +258,8 @@ export class AgentHubOverlayComponent extends Container {
|
||||
|
||||
/** Subscribe to the chat agent's live session (if any) for transcript refreshes. Idempotent per session. */
|
||||
#attachLiveSession(): void {
|
||||
// Remote refs carry no live session handle; refreshes come from observer onChange.
|
||||
if (this.#remote) return;
|
||||
const session = this.#chatAgentId ? (this.#registry.get(this.#chatAgentId)?.session ?? undefined) : undefined;
|
||||
if (session === this.#attachedSession) return;
|
||||
this.#detachLiveSession();
|
||||
@@ -391,6 +413,11 @@ export class AgentHubOverlayComponent extends Container {
|
||||
return;
|
||||
}
|
||||
this.#notice = undefined;
|
||||
if (this.#remote) {
|
||||
this.#remote.revive(ref.id);
|
||||
this.#requestRender();
|
||||
return;
|
||||
}
|
||||
// Fire-and-forget; failures surface as an inline notice
|
||||
this.#lifecycle()
|
||||
.ensureLive(ref.id)
|
||||
@@ -405,6 +432,12 @@ export class AgentHubOverlayComponent extends Container {
|
||||
const ref = this.#rows[this.#selectedRow];
|
||||
if (!ref) return;
|
||||
this.#notice = undefined;
|
||||
if (this.#remote) {
|
||||
this.#remote.kill(ref.id);
|
||||
this.#refreshRows();
|
||||
this.#requestRender();
|
||||
return;
|
||||
}
|
||||
void (async () => {
|
||||
try {
|
||||
if (ref.status === "running" && ref.session) {
|
||||
@@ -512,7 +545,10 @@ export class AgentHubOverlayComponent extends Container {
|
||||
|
||||
// Load transcript first so model info is available for the header
|
||||
let messageEntries: SessionMessageEntry[] | null = null;
|
||||
if (ref?.sessionFile) {
|
||||
if (this.#remote) {
|
||||
if (id) this.#fetchRemoteTranscript(id);
|
||||
messageEntries = this.#transcriptCache?.entries ?? [];
|
||||
} else if (ref?.sessionFile) {
|
||||
messageEntries = this.#loadTranscript(ref.sessionFile);
|
||||
}
|
||||
|
||||
@@ -530,12 +566,18 @@ export class AgentHubOverlayComponent extends Container {
|
||||
this.#viewerEntries = [];
|
||||
if (!ref) {
|
||||
contentLines.push(theme.fg("dim", "Agent no longer registered."));
|
||||
} else if (!ref.sessionFile) {
|
||||
} else if (!this.#remote && !ref.sessionFile) {
|
||||
contentLines.push(theme.fg("dim", "No session file available yet."));
|
||||
} else if (!messageEntries) {
|
||||
contentLines.push(theme.fg("dim", "Unable to read session file."));
|
||||
} else if (messageEntries.length === 0) {
|
||||
contentLines.push(theme.fg("dim", "No messages yet."));
|
||||
if (this.#remote && this.#remoteTranscriptUnavailable) {
|
||||
contentLines.push(theme.fg("dim", "Transcript lives on the host — not available."));
|
||||
} else if (this.#remote && !this.#transcriptCache) {
|
||||
contentLines.push(theme.fg("dim", "Loading transcript from host…"));
|
||||
} else {
|
||||
contentLines.push(theme.fg("dim", "No messages yet."));
|
||||
}
|
||||
} else {
|
||||
this.#buildTranscriptLines(messageEntries, contentLines);
|
||||
}
|
||||
@@ -580,6 +622,12 @@ export class AgentHubOverlayComponent extends Container {
|
||||
if (!id || !trimmed) return;
|
||||
this.#editor.setText("");
|
||||
this.#notice = undefined;
|
||||
if (this.#remote) {
|
||||
this.#remote.chat(id, trimmed);
|
||||
this.#scheduleChatRefresh();
|
||||
this.#requestRender();
|
||||
return;
|
||||
}
|
||||
void (async () => {
|
||||
try {
|
||||
// Revives a parked agent; returns the live session for running/idle.
|
||||
@@ -1024,31 +1072,80 @@ export class AgentHubOverlayComponent extends Container {
|
||||
return this.#loadTranscript(sessionFile);
|
||||
}
|
||||
|
||||
if (!this.#transcriptCache) {
|
||||
this.#transcriptCache = { path: sessionFile, bytesRead: 0, entries: [] };
|
||||
}
|
||||
this.#ingestTranscriptChunk(sessionFile, result.text, fromByte);
|
||||
return this.#transcriptCache?.entries ?? null;
|
||||
}
|
||||
|
||||
if (result.text.length > 0) {
|
||||
const lastNewline = result.text.lastIndexOf("\n");
|
||||
if (lastNewline >= 0) {
|
||||
const completeChunk = result.text.slice(0, lastNewline + 1);
|
||||
const newEntries = parseSessionEntries(completeChunk);
|
||||
for (const entry of newEntries) {
|
||||
if (entry.type === "message") {
|
||||
this.#transcriptCache.entries.push(entry);
|
||||
// Extract model from first assistant message
|
||||
const msg = entry.message;
|
||||
if (!this.#transcriptCache.model && msg.role === "assistant") {
|
||||
this.#transcriptCache.model = msg.model;
|
||||
}
|
||||
} else if (entry.type === "model_change") {
|
||||
this.#transcriptCache.model = entry.model;
|
||||
}
|
||||
/** Parse a complete-line JSONL chunk into the transcript cache and advance bytesRead. Shared by the local file and remote paths. */
|
||||
#ingestTranscriptChunk(cacheKey: string, text: string, fromByte: number): void {
|
||||
if (!this.#transcriptCache) {
|
||||
this.#transcriptCache = { path: cacheKey, bytesRead: 0, entries: [] };
|
||||
}
|
||||
if (text.length === 0) return;
|
||||
const lastNewline = text.lastIndexOf("\n");
|
||||
if (lastNewline < 0) return;
|
||||
const completeChunk = text.slice(0, lastNewline + 1);
|
||||
const newEntries = parseSessionEntries(completeChunk);
|
||||
for (const entry of newEntries) {
|
||||
if (entry.type === "message") {
|
||||
this.#transcriptCache.entries.push(entry);
|
||||
// Extract model from first assistant message
|
||||
const msg = entry.message;
|
||||
if (!this.#transcriptCache.model && msg.role === "assistant") {
|
||||
this.#transcriptCache.model = msg.model;
|
||||
}
|
||||
this.#transcriptCache.bytesRead = fromByte + Buffer.byteLength(completeChunk, "utf-8");
|
||||
} else if (entry.type === "model_change") {
|
||||
this.#transcriptCache.model = entry.model;
|
||||
}
|
||||
}
|
||||
return this.#transcriptCache.entries;
|
||||
this.#transcriptCache.bytesRead = fromByte + Buffer.byteLength(completeChunk, "utf-8");
|
||||
}
|
||||
|
||||
/** Kick an incremental transcript fetch from the collab host (single-flight). */
|
||||
#fetchRemoteTranscript(id: string): void {
|
||||
const remote = this.#remote;
|
||||
if (!remote || this.#remoteFetchInFlight) return;
|
||||
const cacheKey = `remote:${id}`;
|
||||
if (this.#transcriptCache && this.#transcriptCache.path !== cacheKey) {
|
||||
this.#transcriptCache = undefined;
|
||||
}
|
||||
const fromByte = this.#transcriptCache?.bytesRead ?? 0;
|
||||
this.#remoteFetchInFlight = true;
|
||||
const token = ++this.#remoteFetchToken;
|
||||
void remote
|
||||
.readTranscript(id, fromByte)
|
||||
.then(result => {
|
||||
if (token !== this.#remoteFetchToken) return;
|
||||
this.#remoteFetchInFlight = false;
|
||||
if (this.#chatAgentId !== id) return;
|
||||
if (!result) {
|
||||
if (!this.#transcriptCache || this.#transcriptCache.entries.length === 0) {
|
||||
if (!this.#remoteTranscriptUnavailable) {
|
||||
this.#remoteTranscriptUnavailable = true;
|
||||
this.#scheduleChatRefresh();
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (result.newSize < fromByte) {
|
||||
// Host transcript truncated/rotated — restart from 0.
|
||||
this.#transcriptCache = undefined;
|
||||
this.#fetchRemoteTranscript(id);
|
||||
return;
|
||||
}
|
||||
this.#remoteTranscriptUnavailable = false;
|
||||
const hadCache = this.#transcriptCache !== undefined;
|
||||
const before = this.#transcriptCache?.entries.length ?? 0;
|
||||
this.#ingestTranscriptChunk(cacheKey, result.text, fromByte);
|
||||
const after = this.#transcriptCache?.entries.length ?? 0;
|
||||
// Only refresh on new content (or first completed fetch) — an
|
||||
// unconditional rebuild would re-kick the fetch in a tight loop.
|
||||
if (after > before || !hadCache) this.#scheduleChatRefresh();
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
if (token === this.#remoteFetchToken) this.#remoteFetchInFlight = false;
|
||||
logger.warn("Agent hub: remote transcript fetch failed", { id, error: String(error) });
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import type { TextContent } from "@oh-my-pi/pi-ai";
|
||||
import { Container, Markdown, Text } from "@oh-my-pi/pi-tui";
|
||||
import type { CollabPromptDetails } from "../../collab/protocol";
|
||||
import type { CustomMessage } from "../../session/messages";
|
||||
import { getMarkdownTheme, theme } from "../theme/theme";
|
||||
|
||||
/**
|
||||
* Renders a collab guest prompt on every participant's transcript: a
|
||||
* user-message-styled bubble prefixed with the author's name.
|
||||
*/
|
||||
export class CollabPromptMessageComponent extends Container {
|
||||
constructor(message: CustomMessage<CollabPromptDetails>) {
|
||||
super();
|
||||
const from = message.details?.from?.trim() || "guest";
|
||||
this.addChild(new Text(theme.fg("accent", `\x1b[1m«${from}»\x1b[22m ›`), 1, 0));
|
||||
const text =
|
||||
typeof message.content === "string"
|
||||
? message.content
|
||||
: message.content
|
||||
.filter((content): content is TextContent => content.type === "text")
|
||||
.map(content => content.text)
|
||||
.join("");
|
||||
this.addChild(
|
||||
new Markdown(text, 1, 1, getMarkdownTheme(), {
|
||||
bgColor: (value: string) => theme.bg("userMessageBg", value),
|
||||
color: (value: string) => theme.fg("userMessageText", value),
|
||||
}),
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -18,6 +18,7 @@ import { renderSegment, type SegmentContext } from "./segments";
|
||||
import { getSeparator } from "./separators";
|
||||
import { calculateTokensPerSecond } from "./token-rate";
|
||||
import type {
|
||||
CollabStatus,
|
||||
EffectiveStatusLineSettings,
|
||||
StatusLineSegmentId,
|
||||
StatusLineSegmentOptions,
|
||||
@@ -152,6 +153,7 @@ export class StatusLineComponent implements Component {
|
||||
#planModeStatus: { enabled: boolean; paused: boolean } | null = null;
|
||||
#loopModeStatus: { enabled: boolean } | null = null;
|
||||
#goalModeStatus: { enabled: boolean; paused: boolean } | null = null;
|
||||
#collabStatus: CollabStatus | null = null;
|
||||
|
||||
// Git status caching (1s TTL)
|
||||
#cachedGitStatus: { staged: number; unstaged: number; untracked: number } | null = null;
|
||||
@@ -217,6 +219,11 @@ export class StatusLineComponent implements Component {
|
||||
this.#subagentCount = count;
|
||||
}
|
||||
|
||||
/** Active subagent count as currently displayed (collab state mirroring). */
|
||||
get subagentCount(): number {
|
||||
return this.#subagentCount;
|
||||
}
|
||||
|
||||
setSessionStartTime(time: number): void {
|
||||
this.#sessionStartTime = time;
|
||||
}
|
||||
@@ -233,6 +240,10 @@ export class StatusLineComponent implements Component {
|
||||
this.#goalModeStatus = status ?? null;
|
||||
}
|
||||
|
||||
setCollabStatus(status: CollabStatus | null): void {
|
||||
this.#collabStatus = status;
|
||||
}
|
||||
|
||||
setHookStatus(key: string, text: string | undefined): void {
|
||||
if (text === undefined) {
|
||||
this.#hookStatuses.delete(key);
|
||||
@@ -642,7 +653,15 @@ export class StatusLineComponent implements Component {
|
||||
contextTokens = breakdown.usedTokens;
|
||||
contextWindow = breakdown.contextWindow || contextWindow;
|
||||
}
|
||||
const contextPercent = contextWindow > 0 ? (contextTokens / contextWindow) * 100 : 0;
|
||||
let contextPercent = contextWindow > 0 ? (contextTokens / contextWindow) * 100 : 0;
|
||||
|
||||
// Collab guest: context comes from the host's state frames — the local
|
||||
// replica does no accounting of its own.
|
||||
const collabState = this.#collabStatus?.stateOverride;
|
||||
if (collabState?.contextUsage) {
|
||||
contextWindow = collabState.contextUsage.contextWindow || contextWindow;
|
||||
contextPercent = collabState.contextUsage.percent ?? contextPercent;
|
||||
}
|
||||
|
||||
return {
|
||||
session: this.session,
|
||||
@@ -651,6 +670,7 @@ export class StatusLineComponent implements Component {
|
||||
planMode: this.#planModeStatus,
|
||||
loopMode: this.#loopModeStatus,
|
||||
goalMode: this.#goalModeStatus,
|
||||
collab: this.#collabStatus,
|
||||
usageStats,
|
||||
contextPercent,
|
||||
contextWindow,
|
||||
|
||||
@@ -2,7 +2,7 @@ import type { PresetDef, StatusLinePreset } from "./types";
|
||||
|
||||
export const STATUS_LINE_PRESETS: Record<StatusLinePreset, PresetDef> = {
|
||||
default: {
|
||||
leftSegments: ["pi", "model", "mode", "path", "git", "pr", "context_pct", "cost"],
|
||||
leftSegments: ["pi", "model", "mode", "collab", "path", "git", "pr", "context_pct", "cost"],
|
||||
rightSegments: ["session_name"],
|
||||
separator: "powerline-thin",
|
||||
segmentOptions: {
|
||||
|
||||
@@ -493,6 +493,18 @@ const sessionNameSegment: StatusLineSegment = {
|
||||
},
|
||||
};
|
||||
|
||||
const collabSegment: StatusLineSegment = {
|
||||
id: "collab",
|
||||
render(ctx) {
|
||||
if (!ctx.collab) return { content: "", visible: false };
|
||||
const label =
|
||||
ctx.collab.role === "host"
|
||||
? `⇄ collab:${ctx.collab.participantCount}`
|
||||
: `⇄ collab guest:${ctx.collab.participantCount}`;
|
||||
return { content: theme.fg("accent", label), visible: true };
|
||||
},
|
||||
};
|
||||
|
||||
function pickUsageColor(percent: number): "muted" | "warning" | "error" {
|
||||
if (percent >= 80) return "error";
|
||||
if (percent >= 50) return "warning";
|
||||
@@ -573,6 +585,7 @@ export const SEGMENTS: Record<StatusLineSegmentId, StatusLineSegment> = {
|
||||
cache_hit: cacheHitSegment,
|
||||
session_name: sessionNameSegment,
|
||||
usage: usageSegment,
|
||||
collab: collabSegment,
|
||||
};
|
||||
|
||||
export function renderSegment(id: StatusLineSegmentId, ctx: SegmentContext): RenderedSegment {
|
||||
|
||||
@@ -1,8 +1,17 @@
|
||||
import type { CollabSessionState } from "../../../collab/protocol";
|
||||
import type { StatusLinePreset, StatusLineSegmentId, StatusLineSeparatorStyle } from "../../../config/settings-schema";
|
||||
import type { AgentSession } from "../../../session/agent-session";
|
||||
|
||||
export type { StatusLinePreset, StatusLineSegmentId, StatusLineSeparatorStyle };
|
||||
|
||||
/** Collab session indicator + (guest-only) host-state override for segments. */
|
||||
export interface CollabStatus {
|
||||
role: "host" | "guest";
|
||||
participantCount: number;
|
||||
/** Guest only: host footer snapshot that overrides locally computed values. */
|
||||
stateOverride?: CollabSessionState | null;
|
||||
}
|
||||
|
||||
export interface StatusLineSegmentOptions {
|
||||
model?: { showThinkingLevel?: boolean };
|
||||
path?: { abbreviate?: boolean; maxLength?: number; stripWorkPrefix?: boolean };
|
||||
@@ -49,6 +58,7 @@ export interface SegmentContext {
|
||||
enabled: boolean;
|
||||
paused: boolean;
|
||||
} | null;
|
||||
collab: CollabStatus | null;
|
||||
// Cached values for performance (computed once per render)
|
||||
usageStats: {
|
||||
input: number;
|
||||
|
||||
@@ -16,4 +16,5 @@ Press alt+p (or /switch) to switch provider, and ctrl+p to cycle role models smo
|
||||
Press ctrl+r to search your prompt history and reuse a past message
|
||||
`/force read` pins the next turn to one specific tool when the model keeps reaching for the wrong one
|
||||
`/copy code` grabs the last code block to your clipboard — `/copy cmd` grabs the last shell/python command
|
||||
`/shake` rips heavy tool results out of context to reclaim tokens without a full /compact — `/shake images` drops just images
|
||||
`/shake` rips heavy tool results out of context to reclaim tokens without a full /compact — `/shake images` drops just images
|
||||
Pair up live: `/collab` shares your session through an end-to-end encrypted relay link — a teammate runs `/join <link>` to watch tool calls stream and prompt the agent from their own omp
|
||||
@@ -124,6 +124,16 @@ export class InputController {
|
||||
if (this.ctx.hasActiveOmfg() && this.ctx.handleOmfgEscape()) {
|
||||
return;
|
||||
}
|
||||
if (this.ctx.collabGuest) {
|
||||
// Guest Esc: ask the host to interrupt its agent; the local replica
|
||||
// session is never streaming, so the native abort path below would
|
||||
// no-op.
|
||||
if (this.ctx.collabGuest.state?.isStreaming || this.ctx.loadingAnimation) {
|
||||
this.ctx.notifyInterrupting();
|
||||
this.ctx.collabGuest.sendAbort();
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (this.ctx.loadingAnimation) {
|
||||
if (this.ctx.cancelPendingSubmission()) {
|
||||
return;
|
||||
@@ -391,6 +401,32 @@ export class InputController {
|
||||
text = slashResult;
|
||||
}
|
||||
|
||||
// Collab guest: prompts execute on the host; local slash/skill/bash/
|
||||
// python execution is host-only (builtins are gated inside
|
||||
// executeBuiltinSlashCommand, which already consumed allowed ones).
|
||||
if (this.ctx.collabGuest) {
|
||||
if (text.startsWith("/")) {
|
||||
this.ctx.showStatus(`${text.split(/\s+/, 1)[0]} is host-only during a collab session`);
|
||||
this.ctx.editor.setText("");
|
||||
return;
|
||||
}
|
||||
if (text.startsWith("!") || text.startsWith("$")) {
|
||||
this.ctx.showStatus("Local execution is host-only during a collab session");
|
||||
this.ctx.editor.setText("");
|
||||
return;
|
||||
}
|
||||
this.ctx.editor.addToHistory(text);
|
||||
this.ctx.editor.setText("");
|
||||
this.ctx.editor.imageLinks = undefined;
|
||||
const images = inputImages && inputImages.length > 0 ? [...inputImages] : undefined;
|
||||
this.ctx.pendingImages = [];
|
||||
this.ctx.pendingImageLinks = [];
|
||||
// No local render: the prompt comes back from the host as a
|
||||
// collab-prompt event/entry and renders with the author badge.
|
||||
this.ctx.collabGuest.sendPrompt(text, images);
|
||||
return;
|
||||
}
|
||||
|
||||
// Handle skill commands (/skill:name [args]). Enter ⇒ steer (matches the
|
||||
// free-text Enter semantics applied a few lines below at the streaming
|
||||
// branch). Ctrl+Enter routes through `handleFollowUp` and dispatches the
|
||||
|
||||
@@ -1188,6 +1188,8 @@ export class SelectorController {
|
||||
hubKeys,
|
||||
onDone: done,
|
||||
requestRender: () => this.ctx.ui.requestRender(),
|
||||
registry: this.ctx.collabGuest?.agentRegistry,
|
||||
remote: this.ctx.collabGuest?.hubRemote,
|
||||
});
|
||||
|
||||
overlayHandle = this.ctx.ui.showOverlay(hub, {
|
||||
|
||||
@@ -49,6 +49,8 @@ import {
|
||||
} from "@oh-my-pi/pi-utils";
|
||||
import chalk from "chalk";
|
||||
import { reset as resetCapabilities } from "../capability";
|
||||
import type { CollabGuestLink } from "../collab/guest";
|
||||
import type { CollabHost } from "../collab/host";
|
||||
import { KeybindingsManager } from "../config/keybindings";
|
||||
import { isSettingsInitialized, onStatusLineSessionAccentChanged, Settings, settings } from "../config/settings";
|
||||
import { clearClaudePluginRootsCache } from "../discovery/helpers";
|
||||
@@ -396,6 +398,8 @@ export class InteractiveMode implements InteractiveModeContext {
|
||||
fileSlashCommands: Set<string> = new Set();
|
||||
skillCommands: Map<string, string> = new Map();
|
||||
oauthManualInput: OAuthManualInputManager = new OAuthManualInputManager();
|
||||
collabHost?: CollabHost;
|
||||
collabGuest?: CollabGuestLink;
|
||||
|
||||
#pendingSlashCommands: SlashCommand[] = [];
|
||||
#cleanupUnsubscribe?: () => void;
|
||||
@@ -422,6 +426,12 @@ export class InteractiveMode implements InteractiveModeContext {
|
||||
readonly #commandController: CommandController;
|
||||
readonly #todoCommandController: TodoCommandController;
|
||||
readonly #eventController: EventController;
|
||||
get eventController(): EventController {
|
||||
return this.#eventController;
|
||||
}
|
||||
get eventBus(): EventBus | undefined {
|
||||
return this.#eventBus;
|
||||
}
|
||||
readonly #extensionUiController: ExtensionUiController;
|
||||
readonly #inputController: InputController;
|
||||
readonly #selectorController: SelectorController;
|
||||
|
||||
@@ -2,6 +2,8 @@ import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import type { CompactionOutcome } from "@oh-my-pi/pi-agent-core/compaction";
|
||||
import type { AssistantMessage, ImageContent, Message, UsageReport } from "@oh-my-pi/pi-ai";
|
||||
import type { Component, Container, EditorTheme, Loader, Spacer, Text, TUI } from "@oh-my-pi/pi-tui";
|
||||
import type { CollabGuestLink } from "../collab/guest";
|
||||
import type { CollabHost } from "../collab/host";
|
||||
import type { KeybindingsManager } from "../config/keybindings";
|
||||
import type { Settings } from "../config/settings";
|
||||
import type {
|
||||
@@ -19,6 +21,7 @@ import type { HistoryStorage } from "../session/history-storage";
|
||||
import type { SessionContext, SessionManager } from "../session/session-manager";
|
||||
import type { ShakeMode } from "../session/shake-types";
|
||||
import type { LspStartupServerInfo } from "../tools";
|
||||
import type { EventBus } from "../utils/event-bus";
|
||||
import type { AssistantMessageComponent } from "./components/assistant-message";
|
||||
import type { BashExecutionComponent } from "./components/bash-execution";
|
||||
import type { CustomEditor } from "./components/custom-editor";
|
||||
@@ -29,6 +32,7 @@ import type { HookSelectorComponent, HookSelectorOptions } from "./components/ho
|
||||
import type { StatusLineComponent } from "./components/status-line";
|
||||
import type { ToolExecutionHandle } from "./components/tool-execution";
|
||||
import type { TranscriptContainer } from "./components/transcript-container";
|
||||
import type { EventController } from "./controllers/event-controller";
|
||||
import type { LoopLimitRuntime } from "./loop-limit";
|
||||
import type { OAuthManualInputManager } from "./oauth-manual-input";
|
||||
import type { Theme } from "./theme/theme";
|
||||
@@ -101,6 +105,10 @@ export interface InteractiveModeContext {
|
||||
mcpManager?: MCPManager;
|
||||
lspServers?: LspStartupServerInfo[];
|
||||
titleSystemPrompt?: string;
|
||||
collabHost?: CollabHost;
|
||||
collabGuest?: CollabGuestLink;
|
||||
eventController: EventController;
|
||||
eventBus?: EventBus;
|
||||
|
||||
// State
|
||||
isInitialized: boolean;
|
||||
|
||||
@@ -1,11 +1,13 @@
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import type { AssistantMessage, ImageContent, Message } from "@oh-my-pi/pi-ai";
|
||||
import { type Component, Spacer, Text, TruncatedText } from "@oh-my-pi/pi-tui";
|
||||
import { COLLAB_PROMPT_MESSAGE_TYPE, type CollabPromptDetails } from "../../collab/protocol";
|
||||
import { settings } from "../../config/settings";
|
||||
import { getFileSnapshotStore } from "../../edit/file-snapshot-store";
|
||||
import { AssistantMessageComponent } from "../../modes/components/assistant-message";
|
||||
import { BashExecutionComponent } from "../../modes/components/bash-execution";
|
||||
import { BranchSummaryMessageComponent } from "../../modes/components/branch-summary-message";
|
||||
import { CollabPromptMessageComponent } from "../../modes/components/collab-prompt-message";
|
||||
import { CompactionSummaryMessageComponent } from "../../modes/components/compaction-summary-message";
|
||||
import { CustomMessageComponent } from "../../modes/components/custom-message";
|
||||
import { DynamicBorder } from "../../modes/components/dynamic-border";
|
||||
@@ -185,6 +187,11 @@ export class UiHelpers {
|
||||
this.ctx.chatContainer.addChild(component);
|
||||
break;
|
||||
}
|
||||
if (message.customType === COLLAB_PROMPT_MESSAGE_TYPE) {
|
||||
const component = new CollabPromptMessageComponent(message as CustomMessage<CollabPromptDetails>);
|
||||
this.ctx.chatContainer.addChild(component);
|
||||
break;
|
||||
}
|
||||
if (message.customType === SKILL_PROMPT_MESSAGE_TYPE) {
|
||||
const component = new SkillMessageComponent(message as CustomMessage<SkillPromptDetails>);
|
||||
component.setExpanded(this.ctx.toolOutputExpanded);
|
||||
|
||||
@@ -1994,6 +1994,12 @@ export class SessionManager {
|
||||
#byId: Map<string, SessionEntry> = new Map();
|
||||
#labelsById: Map<string, string> = new Map();
|
||||
#leafId: string | null = null;
|
||||
/**
|
||||
* Collab replication tap: invoked for every appended entry with the
|
||||
* in-memory (pre-blob-externalization) entry, so inline images survive.
|
||||
* Failures are swallowed — a broadcast error must never break persistence.
|
||||
*/
|
||||
onEntryAppended?: (entry: SessionEntry) => void;
|
||||
#usageStatistics = {
|
||||
input: 0,
|
||||
output: 0,
|
||||
@@ -2927,6 +2933,44 @@ export class SessionManager {
|
||||
this.#usageStatistics.cost += usage.cost.total;
|
||||
}
|
||||
}
|
||||
if (this.onEntryAppended) {
|
||||
try {
|
||||
this.onEntryAppended(entry);
|
||||
} catch (err) {
|
||||
logger.warn("collab entry hook failed", { error: String(err) });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Append a foreign (host-authored) entry verbatim, preserving its
|
||||
* `id`/`parentId` — no id minting. Used by collab guests to mirror the
|
||||
* host session into the local replica file.
|
||||
*/
|
||||
ingestReplicatedEntry(entry: SessionEntry): void {
|
||||
this.#appendEntry(entry);
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the session for collab replication: the live header plus a deep
|
||||
* copy of every entry (the host mutates entries in place on
|
||||
* truncation/rewrite paths, so guests must not share references).
|
||||
*/
|
||||
snapshotForReplication(): { header: SessionHeader; entries: SessionEntry[] } {
|
||||
const live = this.getHeader();
|
||||
const header: SessionHeader = live
|
||||
? structuredClone(live)
|
||||
: {
|
||||
type: "session",
|
||||
version: CURRENT_SESSION_VERSION,
|
||||
id: this.#sessionId,
|
||||
title: this.#sessionName,
|
||||
titleSource: this.#titleSource,
|
||||
timestamp: new Date().toISOString(),
|
||||
cwd: this.cwd,
|
||||
};
|
||||
const entries = structuredClone(this.#fileEntries.filter(e => e.type !== "session")) as SessionEntry[];
|
||||
return { header, entries };
|
||||
}
|
||||
|
||||
/** Append a message as child of current leaf, then advance leaf. Returns entry id.
|
||||
|
||||
@@ -5,6 +5,8 @@ import { getOAuthProviders } from "@oh-my-pi/pi-ai/oauth";
|
||||
import { setNextRequestDebugPath } from "@oh-my-pi/pi-ai/utils/request-debug";
|
||||
import { Snowflake, setProjectDir } from "@oh-my-pi/pi-utils";
|
||||
import { $ } from "bun";
|
||||
import { COLLAB_GUEST_ALLOWED_COMMANDS, CollabGuestLink } from "../collab/guest";
|
||||
import { CollabHost } from "../collab/host";
|
||||
import type { SettingPath, SettingValue } from "../config/settings";
|
||||
import { settings } from "../config/settings";
|
||||
import {
|
||||
@@ -442,6 +444,115 @@ const BUILTIN_SLASH_COMMAND_REGISTRY: ReadonlyArray<SlashCommandSpec> = [
|
||||
runtime.ctx.editor.setText("");
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "collab",
|
||||
description: "Share this session live via a relay",
|
||||
inlineHint: "[start|stop|status] [relayUrl]",
|
||||
allowArgs: true,
|
||||
handleTui: async (command, runtime) => {
|
||||
const ctx = runtime.ctx;
|
||||
ctx.editor.setText("");
|
||||
const args = command.args.trim();
|
||||
const [first = ""] = args.split(/\s+/, 1);
|
||||
if (first === "stop") {
|
||||
if (!ctx.collabHost) {
|
||||
ctx.showStatus("Not hosting a collab session");
|
||||
return;
|
||||
}
|
||||
await ctx.collabHost.stop("host stopped");
|
||||
ctx.showStatus("Collab stopped");
|
||||
return;
|
||||
}
|
||||
if (first === "status") {
|
||||
if (ctx.collabHost) {
|
||||
const names = ctx.collabHost.participants.map(p => (p.role === "host" ? `${p.name} (host)` : p.name));
|
||||
ctx.showStatus(`Collab: ${names.join(", ")} — ${ctx.collabHost.link}`);
|
||||
} else if (ctx.collabGuest) {
|
||||
ctx.showStatus("In a collab session as a guest (/leave to exit)");
|
||||
} else {
|
||||
ctx.showStatus("Not in a collab session");
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (ctx.collabGuest) {
|
||||
ctx.showError("Already in a collab session as a guest (/leave first)");
|
||||
return;
|
||||
}
|
||||
if (ctx.collabHost) {
|
||||
ctx.session.emitNotice("info", `Collab link: ${ctx.collabHost.link}`, "collab");
|
||||
return;
|
||||
}
|
||||
const explicitUrl = first === "start" ? args.slice("start".length).trim() : args;
|
||||
const relayInput = explicitUrl || ctx.settings.get("collab.relayUrl") || "";
|
||||
if (!relayInput) {
|
||||
ctx.showError(
|
||||
"No relay configured. Set collab.relayUrl in /settings or pass one: /collab relay.example.com",
|
||||
);
|
||||
return;
|
||||
}
|
||||
// Scheme-less relay args default to wss (ws:// must be spelled out for localhost).
|
||||
const relayUrl = relayInput.includes("://") ? relayInput : `wss://${relayInput}`;
|
||||
const host = new CollabHost(ctx);
|
||||
try {
|
||||
await host.start(relayUrl);
|
||||
} catch (err) {
|
||||
ctx.showError(`Failed to start collab session: ${errorMessage(err)}`);
|
||||
return;
|
||||
}
|
||||
ctx.collabHost = host;
|
||||
ctx.session.emitNotice(
|
||||
"info",
|
||||
`Collab link: ${host.link}\nAnyone with this link can read the session and prompt the agent.`,
|
||||
"collab",
|
||||
);
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "join",
|
||||
description: "Join a shared collab session",
|
||||
inlineHint: "<link>",
|
||||
allowArgs: true,
|
||||
handleTui: async (command, runtime) => {
|
||||
const ctx = runtime.ctx;
|
||||
ctx.editor.setText("");
|
||||
const link = command.args.trim();
|
||||
if (!link) {
|
||||
ctx.showError("Usage: /join <link>");
|
||||
return;
|
||||
}
|
||||
if (ctx.collabHost) {
|
||||
ctx.showError("Stop hosting first (/collab stop)");
|
||||
return;
|
||||
}
|
||||
if (ctx.collabGuest) {
|
||||
ctx.showError("Already in a collab session (/leave first)");
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await new CollabGuestLink(ctx).join(link);
|
||||
} catch (err) {
|
||||
ctx.showError(`Failed to join collab session: ${errorMessage(err)}`);
|
||||
}
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "leave",
|
||||
description: "Leave the collab session",
|
||||
handleTui: async (_command, runtime) => {
|
||||
const ctx = runtime.ctx;
|
||||
ctx.editor.setText("");
|
||||
if (ctx.collabGuest) {
|
||||
await ctx.collabGuest.leave("left");
|
||||
return;
|
||||
}
|
||||
if (ctx.collabHost) {
|
||||
await ctx.collabHost.stop("host stopped");
|
||||
ctx.showStatus("Collab stopped");
|
||||
return;
|
||||
}
|
||||
ctx.showStatus("Not in a collab session");
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "browser",
|
||||
description: "Toggle browser headless vs visible mode",
|
||||
@@ -1890,6 +2001,13 @@ export async function executeBuiltinSlashCommand(
|
||||
if (parsed.args.length > 0 && !command.allowArgs) {
|
||||
return false;
|
||||
}
|
||||
// Collab guests run a read-mostly replica: session-mutating builtins are
|
||||
// host-only; the allowlist covers purely local/read-only commands.
|
||||
if (runtime.ctx.collabGuest && !COLLAB_GUEST_ALLOWED_COMMANDS[command.name]) {
|
||||
runtime.ctx.showStatus(`/${command.name} is host-only during a collab session`);
|
||||
runtime.ctx.editor.setText("");
|
||||
return true;
|
||||
}
|
||||
if (command.handleTui) {
|
||||
const result = await command.handleTui(parsed, runtime);
|
||||
if (result && typeof result === "object" && "prompt" in result) return result.prompt;
|
||||
|
||||
@@ -0,0 +1,103 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { generateRoomKey, importRoomKey, open, seal } from "@oh-my-pi/pi-coding-agent/collab/crypto";
|
||||
import {
|
||||
type CollabFrame,
|
||||
DEFAULT_RELAY_URL,
|
||||
formatCollabLink,
|
||||
generateRoomId,
|
||||
packEnvelope,
|
||||
parseCollabLink,
|
||||
rewriteEnvelopePeer,
|
||||
unpackEnvelope,
|
||||
} from "@oh-my-pi/pi-coding-agent/collab/protocol";
|
||||
|
||||
describe("collab crypto", () => {
|
||||
it("round-trips a frame through seal/open", async () => {
|
||||
const key = await importRoomKey(generateRoomKey());
|
||||
const frame: CollabFrame = { t: "prompt", text: "check bun.lock — and ünïcode 🚀" };
|
||||
const sealed = await seal(key, frame);
|
||||
expect(await open(key, sealed)).toEqual(frame);
|
||||
});
|
||||
|
||||
it("rejects tampered ciphertext", async () => {
|
||||
const key = await importRoomKey(generateRoomKey());
|
||||
const sealed = await seal(key, { t: "abort" });
|
||||
sealed[sealed.length - 1]! ^= 0xff;
|
||||
expect(open(key, sealed)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it("rejects frames sealed with a different key", async () => {
|
||||
const sealed = await seal(await importRoomKey(generateRoomKey()), { t: "abort" });
|
||||
const otherKey = await importRoomKey(generateRoomKey());
|
||||
expect(open(otherKey, sealed)).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
|
||||
describe("collab link format", () => {
|
||||
const key = generateRoomKey();
|
||||
const roomId = generateRoomId();
|
||||
|
||||
it("collapses the default relay to a bare roomId#key link", () => {
|
||||
const link = formatCollabLink(DEFAULT_RELAY_URL, roomId, key);
|
||||
expect(link).toBe(`${roomId}#${Buffer.from(key).toString("base64url")}`);
|
||||
const parsed = parseCollabLink(link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
expect(parsed.wsUrl).toBe(`${DEFAULT_RELAY_URL}/r/${roomId}`);
|
||||
expect(parsed.roomId).toBe(roomId);
|
||||
expect(parsed.key).toEqual(key);
|
||||
});
|
||||
|
||||
it("drops the wss scheme for custom relays and infers it on parse", () => {
|
||||
const link = formatCollabLink("wss://relay.example.com:8443", roomId, key);
|
||||
expect(link.startsWith("relay.example.com:8443/r/")).toBe(true);
|
||||
const parsed = parseCollabLink(link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
expect(parsed.wsUrl).toBe(`wss://relay.example.com:8443/r/${roomId}`);
|
||||
});
|
||||
|
||||
it("keeps full ws:// URLs for localhost relays", () => {
|
||||
const link = formatCollabLink("ws://localhost:7475", roomId, key);
|
||||
expect(link.startsWith("ws://localhost:7475/r/")).toBe(true);
|
||||
const parsed = parseCollabLink(link);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
expect(parsed.wsUrl).toBe(`ws://localhost:7475/r/${roomId}`);
|
||||
});
|
||||
|
||||
it("rewrites https relay URLs to wss", () => {
|
||||
const parsed = parseCollabLink(`https://relay.example.com/r/${roomId}#${Buffer.from(key).toString("base64url")}`);
|
||||
if ("error" in parsed) throw new Error(parsed.error);
|
||||
expect(parsed.wsUrl).toBe(`wss://relay.example.com/r/${roomId}`);
|
||||
});
|
||||
|
||||
it("rejects plain ws:// for non-localhost hosts", () => {
|
||||
const parsed = parseCollabLink(`ws://relay.example.com/r/${roomId}#${Buffer.from(key).toString("base64url")}`);
|
||||
expect("error" in parsed && parsed.error.includes("wss://")).toBe(true);
|
||||
});
|
||||
|
||||
it("rejects keys that are not 32 base64url bytes", () => {
|
||||
expect("error" in parseCollabLink(`${roomId}#dG9vc2hvcnQ`)).toBe(true);
|
||||
expect("error" in parseCollabLink(`${roomId}#not+base64url/`)).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("collab wire envelope", () => {
|
||||
it("round-trips peer id and payload", () => {
|
||||
const payload = new Uint8Array([1, 2, 3, 250]);
|
||||
const packed = packEnvelope(0xdeadbeef, payload);
|
||||
const unpacked = unpackEnvelope(packed);
|
||||
expect(unpacked?.peerId).toBe(0xdeadbeef);
|
||||
expect(unpacked?.payload).toEqual(payload);
|
||||
});
|
||||
|
||||
it("rewrites the peer id in place without touching the payload", () => {
|
||||
const packed = packEnvelope(0, new Uint8Array([9, 8, 7]));
|
||||
rewriteEnvelopePeer(packed, 42);
|
||||
const unpacked = unpackEnvelope(packed);
|
||||
expect(unpacked?.peerId).toBe(42);
|
||||
expect(unpacked?.payload).toEqual(new Uint8Array([9, 8, 7]));
|
||||
});
|
||||
|
||||
it("returns null for frames shorter than the header", () => {
|
||||
expect(unpackEnvelope(new Uint8Array([0, 0]))).toBeNull();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,110 @@
|
||||
import { afterEach, describe, expect, it } from "bun:test";
|
||||
import * as path from "node:path";
|
||||
import { isBlobRef } from "@oh-my-pi/pi-coding-agent/session/blob-store";
|
||||
import { type SessionEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
const tempDirs: TempDir[] = [];
|
||||
|
||||
function makeManager(): { manager: SessionManager; cwd: string } {
|
||||
const dir = TempDir.createSync("@pi-collab-repl-");
|
||||
tempDirs.push(dir);
|
||||
const cwd = dir.path();
|
||||
return { manager: SessionManager.create(cwd, path.join(cwd, "sessions")), cwd };
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(tempDirs.splice(0).map(dir => dir.remove()));
|
||||
});
|
||||
|
||||
// Comfortably above BLOB_EXTERNALIZE_THRESHOLD (1024 base64 chars).
|
||||
const BIG_IMAGE_B64 = Buffer.alloc(4096, 7).toString("base64");
|
||||
|
||||
describe("SessionManager collab replication", () => {
|
||||
it("onEntryAppended receives the in-memory entry with inline image data while the persisted line externalizes it", async () => {
|
||||
const { manager } = makeManager();
|
||||
const captured: SessionEntry[] = [];
|
||||
manager.onEntryAppended = entry => captured.push(entry);
|
||||
|
||||
manager.appendMessage({
|
||||
role: "user",
|
||||
content: [
|
||||
{ type: "text", text: "look at this" },
|
||||
{ type: "image", data: BIG_IMAGE_B64, mimeType: "image/png" },
|
||||
],
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
// Persistence is deferred until an assistant message exists; force it.
|
||||
await manager.rewriteEntries();
|
||||
|
||||
expect(captured).toHaveLength(1);
|
||||
const hooked = captured[0]!;
|
||||
if (hooked.type !== "message" || hooked.message.role !== "user" || typeof hooked.message.content === "string") {
|
||||
throw new Error("unexpected hook entry shape");
|
||||
}
|
||||
const hookedImage = hooked.message.content.find(c => c.type === "image");
|
||||
expect(hookedImage?.data).toBe(BIG_IMAGE_B64);
|
||||
|
||||
const file = manager.getSessionFile();
|
||||
if (!file) throw new Error("expected a persisted session file");
|
||||
const lines = (await Bun.file(file).text()).split("\n").filter(Boolean);
|
||||
const persisted = lines.map(line => JSON.parse(line)).find(e => e.type === "message");
|
||||
const persistedImage = persisted.message.content.find((c: { type: string; data?: string }) => c.type === "image");
|
||||
expect(isBlobRef(persistedImage.data)).toBe(true);
|
||||
});
|
||||
|
||||
it("swallows hook failures so persistence is never broken by a broadcast error", () => {
|
||||
const { manager } = makeManager();
|
||||
manager.onEntryAppended = () => {
|
||||
throw new Error("socket exploded");
|
||||
};
|
||||
const id = manager.appendMessage({ role: "user", content: "still works", timestamp: Date.now() });
|
||||
expect(manager.getEntry(id)?.id).toBe(id);
|
||||
});
|
||||
|
||||
it("ingestReplicatedEntry preserves foreign ids and advances the leaf", async () => {
|
||||
const { manager } = makeManager();
|
||||
const rootId = manager.appendMessage({ role: "user", content: "root", timestamp: Date.now() });
|
||||
|
||||
const foreign: SessionEntry = {
|
||||
type: "message",
|
||||
id: "feed0001",
|
||||
parentId: rootId,
|
||||
timestamp: new Date().toISOString(),
|
||||
message: { role: "user", content: "from the host", timestamp: Date.now() },
|
||||
};
|
||||
manager.ingestReplicatedEntry(foreign);
|
||||
|
||||
expect(manager.getEntry("feed0001")?.parentId).toBe(rootId);
|
||||
// Leaf advanced: the next locally appended entry chains off the ingested one.
|
||||
const nextId = manager.appendMessage({ role: "user", content: "after", timestamp: Date.now() });
|
||||
expect(manager.getEntry(nextId)?.parentId).toBe("feed0001");
|
||||
|
||||
// Round-trip: a fresh manager loading the file sees the foreign ids verbatim.
|
||||
await manager.rewriteEntries();
|
||||
const file = manager.getSessionFile();
|
||||
if (!file) throw new Error("expected a persisted session file");
|
||||
const { manager: loaded } = makeManager();
|
||||
await loaded.setSessionFile(file);
|
||||
expect(loaded.getEntry("feed0001")?.parentId).toBe(rootId);
|
||||
expect(loaded.getEntry(nextId)?.parentId).toBe("feed0001");
|
||||
});
|
||||
|
||||
it("snapshotForReplication deep-copies entries and preserves the header identity", () => {
|
||||
const { manager, cwd } = makeManager();
|
||||
manager.appendMessage({ role: "user", content: "snapshot me", timestamp: Date.now() });
|
||||
|
||||
const snapshot = manager.snapshotForReplication();
|
||||
expect(snapshot.header.id).toBe(manager.getSessionId());
|
||||
expect(snapshot.header.cwd).toBe(path.resolve(cwd));
|
||||
expect(snapshot.entries).toHaveLength(1);
|
||||
|
||||
// Deep copy: mutating the snapshot must not leak into the live session.
|
||||
const entry = snapshot.entries[0]!;
|
||||
if (entry.type !== "message") throw new Error("unexpected entry type");
|
||||
entry.message = { role: "user", content: "mutated", timestamp: 0 };
|
||||
const live = manager.getEntry(entry.id);
|
||||
if (live?.type !== "message" || live.message.role !== "user") throw new Error("unexpected live entry");
|
||||
expect(live.message.content).toBe("snapshot me");
|
||||
});
|
||||
});
|
||||
@@ -20,6 +20,7 @@ function createCtx(usage: Partial<SegmentContext["usageStats"]>): SegmentContext
|
||||
planMode: null,
|
||||
loopMode: null,
|
||||
goalMode: null,
|
||||
collab: null,
|
||||
usageStats: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
/**
|
||||
* `omp join <link>` must route to the registered `join` subcommand instead of
|
||||
* being rewritten to `launch join <link>` and forwarded to the LLM as an
|
||||
* initial prompt (same failure mode as #1496).
|
||||
*/
|
||||
import { describe, expect, test } from "bun:test";
|
||||
import { isSubcommand, resolveCliArgv } from "@oh-my-pi/pi-coding-agent/cli-commands";
|
||||
|
||||
describe("join command is registered as a top-level subcommand", () => {
|
||||
test("CLI runner routes `join <link>` to the join command, not launch", () => {
|
||||
expect(isSubcommand("join")).toBe(true);
|
||||
expect(resolveCliArgv(["join", "wss://relay.omp.sh/s/abc#key"])).toEqual({
|
||||
argv: ["join", "wss://relay.omp.sh/s/abc#key"],
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -45,6 +45,7 @@ function createCtx(overrides?: { pathMaxLength?: number; branch?: string | null
|
||||
planMode: null,
|
||||
loopMode: null,
|
||||
goalMode: null,
|
||||
collab: null,
|
||||
usageStats: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
|
||||
@@ -31,6 +31,7 @@ function createPathContext(): SegmentContext {
|
||||
planMode: null,
|
||||
loopMode: null,
|
||||
goalMode: null,
|
||||
collab: null,
|
||||
usageStats: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
|
||||
Reference in New Issue
Block a user