diff --git a/docs/collab.md b/docs/collab.md new file mode 100644 index 000000000..9122eb136 --- /dev/null +++ b/docs/collab.md @@ -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 ` | Start sharing through a specific relay (`relay.example.com`, `ws://localhost:7475`) | +| `/collab status` | Show link + participants | +| `/collab stop` | Stop sharing | +| `/join ` | Join a shared session as a guest | +| `/leave` | Leave (guest) or stop sharing (host) | + +## Link format + +``` +# → default relay (relay.omp.sh) +host[:port]/r/# → custom relay, wss:// inferred +ws://localhost:7475/r/# → plain ws, allowed for localhost only +``` + +The fragment (`#`) 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/?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/.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. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 741a4a5bd..3b1264c6c 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 ` subcommand that launches the interactive TUI and immediately runs `/join ` ### 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 diff --git a/packages/coding-agent/src/cli-commands.ts b/packages/coding-agent/src/cli-commands.ts index cc4e3a668..6faa4cad3 100644 --- a/packages/coding-agent/src/cli-commands.ts +++ b/packages/coding-agent/src/cli-commands.ts @@ -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) }, diff --git a/packages/coding-agent/src/cli/args.ts b/packages/coding-agent/src/cli/args.ts index f5737d2a3..5dcd39d7e 100644 --- a/packages/coding-agent/src/cli/args.ts +++ b/packages/coding-agent/src/cli/args.ts @@ -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; diff --git a/packages/coding-agent/src/collab/crypto.ts b/packages/coding-agent/src/collab/crypto.ts new file mode 100644 index 000000000..e8b7ab158 --- /dev/null +++ b/packages/coding-agent/src/collab/crypto.ts @@ -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 { + 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 { + 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 { + 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 { + if (bytes.buffer instanceof ArrayBuffer && bytes.byteOffset === 0 && bytes.byteLength === bytes.buffer.byteLength) { + return bytes as Uint8Array; + } + const copy = new Uint8Array(bytes.byteLength); + copy.set(bytes); + return copy; +} diff --git a/packages/coding-agent/src/collab/guest.ts b/packages/coding-agent/src/collab/guest.ts new file mode 100644 index 000000000..520bb51c2 --- /dev/null +++ b/packages/coding-agent/src/collab/guest.ts @@ -0,0 +1,421 @@ +/** + * Guest side of a collab live session. + * + * `/join ` 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 = { + 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; + +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 = 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(); + #pendingTranscripts = new Map 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 { + 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(); + 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 { + 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 { + 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(); + 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 { + 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, + }); + } +} diff --git a/packages/coding-agent/src/collab/host.ts b/packages/coding-agent/src/collab/host.ts new file mode 100644 index 000000000..e0f6bbf67 --- /dev/null +++ b/packages/coding-agent/src/collab/host.ts @@ -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 = { + 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(); + #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 { + 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(); + 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 { + if (this.#stopped) return; + this.#socket?.send({ t: "bye", reason }); + await this.#teardown(); + } + + async #teardown(): Promise { + 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 { + 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(); + } +} diff --git a/packages/coding-agent/src/collab/protocol.ts b/packages/coding-agent/src/collab/protocol.ts new file mode 100644 index 000000000..fd3e66221 --- /dev/null +++ b/packages/coding-agent/src/collab/protocol.ts @@ -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:///r/# +// ═══════════════════════════════════════════════════════════════════════════ + +export const ROOM_ID_BYTES = 16; + +/** Default public relay; bare `#` 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 = { localhost: true, "127.0.0.1": true, "::1": true, "[::1]": true }; + +export interface ParsedCollabLink { + /** wss://host[:port]/r/ — 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 + * `#`, 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 `#` → 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/ 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 # 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 }; +} diff --git a/packages/coding-agent/src/collab/relay-client.ts b/packages/coding-agent/src/collab/relay-client.ts new file mode 100644 index 000000000..130e92dec --- /dev/null +++ b/packages/coding-agent/src/collab/relay-client.ts @@ -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 = { + 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/ — 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 = Promise.resolve(); + /** Serializes open() so frames are delivered in arrival order. */ + #recvChain: Promise = 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; + } + } +} diff --git a/packages/coding-agent/src/commands/join.ts b/packages/coding-agent/src/commands/join.ts new file mode 100644 index 000000000..c93cc4fb9 --- /dev/null +++ b/packages/coding-agent/src/commands/join.ts @@ -0,0 +1,39 @@ +/** + * Join a shared collab session from the CLI: launches the interactive TUI and + * immediately runs `/join `. + */ +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 { + const { args } = await this.parse(Join); + const link = args.link?.trim(); + if (!link) { + process.stderr.write(`Usage: ${APP_NAME} join \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, []); + } +} diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index e5dcf73f8..16ad64edc 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -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 = { "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 = { @@ -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", diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 225af3ee1..7025a9317 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -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 { const mode = new InteractiveMode( session, @@ -414,6 +416,12 @@ async function runInteractiveMode( } } + // `omp join `: 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. diff --git a/packages/coding-agent/src/modes/components/agent-hub.ts b/packages/coding-agent/src/modes/components/agent-hub.ts index 9dedb5f14..6dadb32ba 100644 --- a/packages/coding-agent/src/modes/components/agent-hub.ts +++ b/packages/coding-agent/src/modes/components/agent-hub.ts @@ -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) }); + }); } } diff --git a/packages/coding-agent/src/modes/components/collab-prompt-message.ts b/packages/coding-agent/src/modes/components/collab-prompt-message.ts new file mode 100644 index 000000000..657259bff --- /dev/null +++ b/packages/coding-agent/src/modes/components/collab-prompt-message.ts @@ -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) { + 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), + }), + ); + } +} diff --git a/packages/coding-agent/src/modes/components/status-line/component.ts b/packages/coding-agent/src/modes/components/status-line/component.ts index 91d4d4096..74ba39b25 100644 --- a/packages/coding-agent/src/modes/components/status-line/component.ts +++ b/packages/coding-agent/src/modes/components/status-line/component.ts @@ -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, diff --git a/packages/coding-agent/src/modes/components/status-line/presets.ts b/packages/coding-agent/src/modes/components/status-line/presets.ts index 210906b8a..b7628d97c 100644 --- a/packages/coding-agent/src/modes/components/status-line/presets.ts +++ b/packages/coding-agent/src/modes/components/status-line/presets.ts @@ -2,7 +2,7 @@ import type { PresetDef, StatusLinePreset } from "./types"; export const STATUS_LINE_PRESETS: Record = { 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: { diff --git a/packages/coding-agent/src/modes/components/status-line/segments.ts b/packages/coding-agent/src/modes/components/status-line/segments.ts index c0c57835d..708e646b7 100644 --- a/packages/coding-agent/src/modes/components/status-line/segments.ts +++ b/packages/coding-agent/src/modes/components/status-line/segments.ts @@ -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 = { cache_hit: cacheHitSegment, session_name: sessionNameSegment, usage: usageSegment, + collab: collabSegment, }; export function renderSegment(id: StatusLineSegmentId, ctx: SegmentContext): RenderedSegment { diff --git a/packages/coding-agent/src/modes/components/status-line/types.ts b/packages/coding-agent/src/modes/components/status-line/types.ts index 933072ab2..953fbf571 100644 --- a/packages/coding-agent/src/modes/components/status-line/types.ts +++ b/packages/coding-agent/src/modes/components/status-line/types.ts @@ -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; diff --git a/packages/coding-agent/src/modes/components/tips.txt b/packages/coding-agent/src/modes/components/tips.txt index f606541c8..8458cfa95 100644 --- a/packages/coding-agent/src/modes/components/tips.txt +++ b/packages/coding-agent/src/modes/components/tips.txt @@ -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 \ No newline at end of file +`/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 ` to watch tool calls stream and prompt the agent from their own omp \ No newline at end of file diff --git a/packages/coding-agent/src/modes/controllers/input-controller.ts b/packages/coding-agent/src/modes/controllers/input-controller.ts index ccbb0e958..183c2961e 100644 --- a/packages/coding-agent/src/modes/controllers/input-controller.ts +++ b/packages/coding-agent/src/modes/controllers/input-controller.ts @@ -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 diff --git a/packages/coding-agent/src/modes/controllers/selector-controller.ts b/packages/coding-agent/src/modes/controllers/selector-controller.ts index 9c6c1822b..5a9270fed 100644 --- a/packages/coding-agent/src/modes/controllers/selector-controller.ts +++ b/packages/coding-agent/src/modes/controllers/selector-controller.ts @@ -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, { diff --git a/packages/coding-agent/src/modes/interactive-mode.ts b/packages/coding-agent/src/modes/interactive-mode.ts index 9c59cf97b..83b0fceb8 100644 --- a/packages/coding-agent/src/modes/interactive-mode.ts +++ b/packages/coding-agent/src/modes/interactive-mode.ts @@ -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 = new Set(); skillCommands: Map = 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; diff --git a/packages/coding-agent/src/modes/types.ts b/packages/coding-agent/src/modes/types.ts index ac580702c..2f25f9e6e 100644 --- a/packages/coding-agent/src/modes/types.ts +++ b/packages/coding-agent/src/modes/types.ts @@ -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; diff --git a/packages/coding-agent/src/modes/utils/ui-helpers.ts b/packages/coding-agent/src/modes/utils/ui-helpers.ts index 5a5176a31..3bbeeb870 100644 --- a/packages/coding-agent/src/modes/utils/ui-helpers.ts +++ b/packages/coding-agent/src/modes/utils/ui-helpers.ts @@ -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); + this.ctx.chatContainer.addChild(component); + break; + } if (message.customType === SKILL_PROMPT_MESSAGE_TYPE) { const component = new SkillMessageComponent(message as CustomMessage); component.setExpanded(this.ctx.toolOutputExpanded); diff --git a/packages/coding-agent/src/session/session-manager.ts b/packages/coding-agent/src/session/session-manager.ts index 81649ff86..7c3ef3f91 100644 --- a/packages/coding-agent/src/session/session-manager.ts +++ b/packages/coding-agent/src/session/session-manager.ts @@ -1994,6 +1994,12 @@ export class SessionManager { #byId: Map = new Map(); #labelsById: Map = 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. diff --git a/packages/coding-agent/src/slash-commands/builtin-registry.ts b/packages/coding-agent/src/slash-commands/builtin-registry.ts index b0de73542..01187dfef 100644 --- a/packages/coding-agent/src/slash-commands/builtin-registry.ts +++ b/packages/coding-agent/src/slash-commands/builtin-registry.ts @@ -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 = [ 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: "", + allowArgs: true, + handleTui: async (command, runtime) => { + const ctx = runtime.ctx; + ctx.editor.setText(""); + const link = command.args.trim(); + if (!link) { + ctx.showError("Usage: /join "); + 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; diff --git a/packages/coding-agent/test/collab/crypto.test.ts b/packages/coding-agent/test/collab/crypto.test.ts new file mode 100644 index 000000000..6a1ea5587 --- /dev/null +++ b/packages/coding-agent/test/collab/crypto.test.ts @@ -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(); + }); +}); diff --git a/packages/coding-agent/test/collab/session-replication.test.ts b/packages/coding-agent/test/collab/session-replication.test.ts new file mode 100644 index 000000000..c60fd8172 --- /dev/null +++ b/packages/coding-agent/test/collab/session-replication.test.ts @@ -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"); + }); +}); diff --git a/packages/coding-agent/test/issue-953-repro.test.ts b/packages/coding-agent/test/issue-953-repro.test.ts index 4c8aeb064..8f3901f93 100644 --- a/packages/coding-agent/test/issue-953-repro.test.ts +++ b/packages/coding-agent/test/issue-953-repro.test.ts @@ -20,6 +20,7 @@ function createCtx(usage: Partial): SegmentContext planMode: null, loopMode: null, goalMode: null, + collab: null, usageStats: { input: 0, output: 0, diff --git a/packages/coding-agent/test/join-command.test.ts b/packages/coding-agent/test/join-command.test.ts new file mode 100644 index 000000000..16246c0f2 --- /dev/null +++ b/packages/coding-agent/test/join-command.test.ts @@ -0,0 +1,16 @@ +/** + * `omp join ` must route to the registered `join` subcommand instead of + * being rewritten to `launch join ` 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 ` 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"], + }); + }); +}); diff --git a/packages/coding-agent/test/status-line-overflow.test.ts b/packages/coding-agent/test/status-line-overflow.test.ts index 6efc1207a..f4fe3a1e1 100644 --- a/packages/coding-agent/test/status-line-overflow.test.ts +++ b/packages/coding-agent/test/status-line-overflow.test.ts @@ -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, diff --git a/packages/coding-agent/test/status-line-path.test.ts b/packages/coding-agent/test/status-line-path.test.ts index 08ef1bb9e..65c9ffac3 100644 --- a/packages/coding-agent/test/status-line-path.test.ts +++ b/packages/coding-agent/test/status-line-path.test.ts @@ -31,6 +31,7 @@ function createPathContext(): SegmentContext { planMode: null, loopMode: null, goalMode: null, + collab: null, usageStats: { input: 0, output: 0,