diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f667766f8..248a2ba97 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -172,9 +172,9 @@ jobs: echo "cross-platform-run-id=$cross_platform_run_id" } >> "$GITHUB_OUTPUT" - # Fast lint + type check (no Rust, no native build needed) + # Fast lint, type check, and browser bundle build (no Rust, no native build needed) check: - name: Lint & type check + name: Lint, type check & web build runs-on: ubuntu-22.04 steps: - uses: actions/checkout@v4 @@ -189,6 +189,8 @@ jobs: - run: bun install --frozen-lockfile - name: Type check workspace run: bun run ci:check:full + - name: Build collab web + run: bun run collab:web:build # Linux x64 baseline + modern: required by `test`, so it runs on every PR # unless native_artifact_lookup found a cached run. Release runs always diff --git a/README.md b/README.md index 10c470db5..036b0bfc8 100644 --- a/README.md +++ b/README.md @@ -394,16 +394,21 @@ The same prompt cards surface over ACP, so editors get the picker without writin Node and TypeScript hosts pull the engine in directly. The package exposes `ModelRegistry`, `SessionManager`, `createAgentSession`, and `discoverAuthStorage`; the session emits typed events you subscribe to. ```ts -import { ModelRegistry, SessionManager, createAgentSession, discoverAuthStorage } from "@oh-my-pi/pi-coding-agent"; +import { + ModelRegistry, + SessionManager, + createAgentSession, + discoverAuthStorage, +} from "@oh-my-pi/pi-coding-agent"; const auth = await discoverAuthStorage(); const models = new ModelRegistry(auth); await models.refresh(); const { session } = await createAgentSession({ - sessionManager: SessionManager.inMemory(), - authStorage: auth, - modelRegistry: models, + sessionManager: SessionManager.inMemory(), + authStorage: auth, + modelRegistry: models, }); await session.prompt("list .ts files"); ``` @@ -481,6 +486,7 @@ For architecture and contribution guidelines, see [packages/coding-agent/DEVELOP | Package | Description | | --------------------------------------------------------- | -------------------------------------------------------------------------- | +| **[@oh-my-pi/collab-web](packages/collab-web)** | Browser guest client, mock host, and local relay for collab live sessions | | **[@oh-my-pi/pi-ai](packages/ai)** | Multi-provider LLM client with streaming and model/provider integration | | **[@oh-my-pi/pi-catalog](packages/catalog)** | Model catalog: bundled model database, provider descriptors, and identity | | **[@oh-my-pi/pi-agent-core](packages/agent)** | Agent runtime with tool calling and state management | @@ -489,6 +495,7 @@ For architecture and contribution guidelines, see [packages/coding-agent/DEVELOP | **[@oh-my-pi/pi-natives](packages/natives)** | N-API bindings for grep, shell, image, text, syntax highlighting, and more | | **[@oh-my-pi/omp-stats](packages/stats)** | Local observability dashboard for AI usage statistics | | **[@oh-my-pi/pi-utils](packages/utils)** | Shared utilities (logging, streams, dirs/env/process helpers) | +| **[@oh-my-pi/pi-wire](packages/wire)** | Shared collab live-session protocol types and relay constants | | **[@oh-my-pi/hashline](packages/hashline)** | Line-anchored patch language and applier behind the `edit` tool | | **[@oh-my-pi/pi-mnemopi](packages/mnemopi)** | Local SQLite memory engine for Oh My Pi agents | | **[@oh-my-pi/snapcompact](packages/snapcompact)** | Bitmap-frame context compression package and SQuAD eval suite | diff --git a/bun.lock b/bun.lock index d16539588..ed3571444 100644 --- a/bun.lock +++ b/bun.lock @@ -77,6 +77,7 @@ "@oh-my-pi/pi-natives": "catalog:", "@oh-my-pi/pi-tui": "catalog:", "@oh-my-pi/pi-utils": "catalog:", + "@oh-my-pi/pi-wire": "catalog:", "@oh-my-pi/snapcompact": "catalog:", "@opentelemetry/api": "catalog:", "@opentelemetry/context-async-hooks": "catalog:", @@ -106,6 +107,22 @@ "@huggingface/transformers": "catalog:", }, }, + "packages/collab-web": { + "name": "@oh-my-pi/collab-web", + "version": "15.11.7", + "dependencies": { + "@oh-my-pi/pi-wire": "catalog:", + "lucide-react": "catalog:", + "marked": "catalog:", + "react": "catalog:", + "react-dom": "catalog:", + }, + "devDependencies": { + "@types/bun": "catalog:", + "@types/react": "catalog:", + "@types/react-dom": "catalog:", + }, + }, "packages/hashline": { "name": "@oh-my-pi/hashline", "version": "15.11.7", @@ -252,6 +269,13 @@ "@types/bun": "catalog:", }, }, + "packages/wire": { + "name": "@oh-my-pi/pi-wire", + "version": "15.11.7", + "devDependencies": { + "@types/bun": "catalog:", + }, + }, "python/robomp/web": { "name": "robomp-web", "version": "0.1.0", @@ -290,6 +314,7 @@ "@oh-my-pi/pi-natives": "15.11.7", "@oh-my-pi/pi-tui": "15.11.7", "@oh-my-pi/pi-utils": "15.11.7", + "@oh-my-pi/pi-wire": "15.11.7", "@oh-my-pi/snapcompact": "15.11.7", "@opentelemetry/api": "^1.9.1", "@opentelemetry/context-async-hooks": "^2.7.1", @@ -667,6 +692,8 @@ "@octokit/types": ["@octokit/types@16.0.0", "", { "dependencies": { "@octokit/openapi-types": "^27.0.0" } }, "sha512-sKq+9r1Mm4efXW1FCk7hFSeJo4QKreL/tTbR0rz/qx/r1Oa2VV83LTA/H/MuCOX7uCIJmQVRKBcbmWoySjAnSg=="], + "@oh-my-pi/collab-web": ["@oh-my-pi/collab-web@workspace:packages/collab-web"], + "@oh-my-pi/hashline": ["@oh-my-pi/hashline@workspace:packages/hashline"], "@oh-my-pi/omp-stats": ["@oh-my-pi/omp-stats@workspace:packages/stats"], @@ -687,6 +714,8 @@ "@oh-my-pi/pi-utils": ["@oh-my-pi/pi-utils@workspace:packages/utils"], + "@oh-my-pi/pi-wire": ["@oh-my-pi/pi-wire@workspace:packages/wire"], + "@oh-my-pi/snapcompact": ["@oh-my-pi/snapcompact@workspace:packages/snapcompact"], "@oh-my-pi/swarm-extension": ["@oh-my-pi/swarm-extension@workspace:packages/swarm-extension"], diff --git a/docs/collab.md b/docs/collab.md index 9122eb136..2b05fcd89 100644 --- a/docs/collab.md +++ b/docs/collab.md @@ -68,6 +68,10 @@ Everything that mutates the host session or machine is host-only: `/model`, `/co Known v1 limit for guests: a turn already streaming when you join becomes visible from its next message boundary. +## Web client + +`packages/collab-web` is a standalone browser client for the same links — no omp install needed on the guest side. It renders the live transcript (streaming text, thinking, tool cards), a subagent panel with on-demand transcripts, and a composer with the same guest powers (prompt, interrupt, hub actions). Run `bun run dev` in the package for a local instance, `bun run mock-host` for an offline scripted host to develop against, and `bun run build` to emit a static `dist/` deployable anywhere (HTTPS required for WebCrypto). The client never talks to anything but the relay, and the key stays in the URL fragment. + ## Settings | Setting | Default | Meaning | diff --git a/package.json b/package.json index 3c91a2ff4..0330dc41a 100644 --- a/package.json +++ b/package.json @@ -31,6 +31,7 @@ "@oh-my-pi/pi-natives": "15.11.7", "@oh-my-pi/pi-tui": "15.11.7", "@oh-my-pi/pi-utils": "15.11.7", + "@oh-my-pi/pi-wire": "15.11.7", "@oh-my-pi/snapcompact": "15.11.7", "@opentelemetry/api": "^1.9.1", "@opentelemetry/context-async-hooks": "^2.7.1", @@ -92,6 +93,10 @@ "dev": "bun --cwd=packages/coding-agent src/cli.ts", "dev:timing": "PI_TIMING=x bun --cwd=packages/coding-agent --preload ../utils/src/module-timer.ts src/cli.ts", "stats": "bun --cwd=packages/coding-agent src/cli.ts stats", + "collab:web:dev": "bun --cwd=packages/collab-web run dev", + "collab:relay": "bun --cwd=packages/collab-web run relay", + "collab:mock-host": "bun --cwd=packages/collab-web run mock-host", + "collab:web:build": "bun --cwd=packages/collab-web run build", "claude:trace": "bun scripts/claude-trace.ts", "build": "bun run --workspaces --if-present build", "build:native": "bun --cwd=packages/natives run build", diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 8b6c981f5..af69f82ce 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -11,6 +11,7 @@ ### Changed +- Moved collab live-session wire contracts into `@oh-my-pi/pi-wire` and stopped broadcasting unsupported session events or entries to collab guests. - 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 diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index 009694006..48c924ebb 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -57,6 +57,7 @@ "@oh-my-pi/pi-natives": "catalog:", "@oh-my-pi/pi-tui": "catalog:", "@oh-my-pi/pi-utils": "catalog:", + "@oh-my-pi/pi-wire": "catalog:", "@oh-my-pi/snapcompact": "catalog:", "@opentelemetry/api": "catalog:", "@opentelemetry/context-async-hooks": "catalog:", diff --git a/packages/coding-agent/src/collab/host.ts b/packages/coding-agent/src/collab/host.ts index e0f6bbf67..041533879 100644 --- a/packages/coding-agent/src/collab/host.ts +++ b/packages/coding-agent/src/collab/host.ts @@ -12,11 +12,13 @@ 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 { BusChannel, AgentEvent as WireAgentEvent, SessionEntry as WireSessionEntry } from "@oh-my-pi/pi-wire"; 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 type { SessionEntry as StoredSessionEntry } from "../session/session-manager"; import { TASK_SUBAGENT_LIFECYCLE_CHANNEL, TASK_SUBAGENT_PROGRESS_CHANNEL } from "../task"; import { generateRoomKey, importRoomKey } from "./crypto"; import { @@ -46,6 +48,45 @@ 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 WIRE_AGENT_EVENT_TYPES: Record = { + agent_start: true, + agent_end: true, + turn_start: true, + turn_end: true, + message_start: true, + message_update: true, + message_end: true, + tool_execution_start: true, + tool_execution_update: true, + tool_execution_end: true, + notice: true, + auto_compaction_start: true, + auto_compaction_end: true, + auto_retry_start: true, + auto_retry_end: true, + thinking_level_changed: true, +}; + +const WIRE_SESSION_ENTRY_TYPES: Record = { + message: true, + custom_message: true, + compaction: true, + branch_summary: true, + model_change: true, + thinking_level_change: true, +}; +const COLLAB_BUS_CHANNELS = [ + TASK_SUBAGENT_LIFECYCLE_CHANNEL, + TASK_SUBAGENT_PROGRESS_CHANNEL, +] as const satisfies readonly BusChannel[]; + +function isWireAgentEvent(event: AgentSessionEvent): event is AgentSessionEvent & WireAgentEvent { + return event.type in WIRE_AGENT_EVENT_TYPES; +} + +function isWireSessionEntry(entry: StoredSessionEntry): entry is StoredSessionEntry & WireSessionEntry { + return entry.type in WIRE_SESSION_ENTRY_TYPES; +} 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; @@ -145,18 +186,18 @@ export class CollabHost { } this.#unsubscribe = this.#ctx.session.subscribe(event => { - this.#broadcast({ t: "event", event }); + if (isWireAgentEvent(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]) { + for (const channel of COLLAB_BUS_CHANNELS) { 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 }); + if (isWireSessionEntry(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(); @@ -249,12 +290,13 @@ export class CollabHost { } logger.info("collab welcome exceeded size threshold; stripped images", { stripped }); } + const entries = snapshot.entries.filter(isWireSessionEntry); this.#socket?.send( { t: "welcome", proto: COLLAB_PROTO, header: snapshot.header, - entries: snapshot.entries, + entries, state: this.#buildState(), agents: this.#snapshotAgents(), }, diff --git a/packages/coding-agent/src/collab/protocol.ts b/packages/coding-agent/src/collab/protocol.ts index fd3e66221..41dd8715a 100644 --- a/packages/coding-agent/src/collab/protocol.ts +++ b/packages/coding-agent/src/collab/protocol.ts @@ -6,68 +6,54 @@ * 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 { + BusChannel, + GuestFrame, + ParsedCollabLink, + Participant, + SessionState, + AgentSnapshot as WireAgentSnapshot, +} from "@oh-my-pi/pi-wire"; +import { DEFAULT_RELAY_URL, ENVELOPE_HEADER_LENGTH, ROOM_ID_BYTES } from "@oh-my-pi/pi-wire"; 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; +export type { + CollabPromptDetails, + ParsedCollabLink, + RelayControlMessage, + RelayControlToGuest, + RelayControlToHost, +} from "@oh-my-pi/pi-wire"; +export { COLLAB_PROMPT_MESSAGE_TYPE, COLLAB_PROTO } from "@oh-my-pi/pi-wire"; +export { DEFAULT_RELAY_URL, ENVELOPE_HEADER_LENGTH, ROOM_ID_BYTES }; -/** 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; -} +export type CollabParticipant = Participant; +export type AgentSnapshot = WireAgentSnapshot; /** 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; +export type CollabSessionState = SessionState & { /** * 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). */ +/** + * Encrypted payload frames (inside AES-GCM, JSON). The wire package pins the + * JSON skeleton (`WireFrame`); host-side frames carry the rich session types + * that serialize into those shapes. + */ export type CollabFrame = - // guest -> host - | { t: "hello"; proto: number; name: string } + // guest -> host (hello/abort/agent-cmd/fetch-transcript are taken verbatim from the wire grammar) + | Exclude | { 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"; @@ -81,7 +67,7 @@ export type CollabFrame = | { t: "event"; event: AgentSessionEvent } | { t: "state"; state: CollabSessionState } /** Mirrored EventBus traffic (task subagent lifecycle/progress channels only). */ - | { t: "bus"; channel: string; data: unknown } + | { t: "bus"; channel: BusChannel; 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. */ @@ -89,20 +75,12 @@ export type CollabFrame = | { 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); @@ -125,23 +103,11 @@ export function rewriteEnvelopePeer(data: Uint8Array, peerId: number): void { // 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); diff --git a/packages/collab-web/CHANGELOG.md b/packages/collab-web/CHANGELOG.md new file mode 100644 index 000000000..741f450bf --- /dev/null +++ b/packages/collab-web/CHANGELOG.md @@ -0,0 +1,19 @@ +# Changelog + +## [Unreleased] +### Added + +- Added deep-link auto-connection support from `##` URLs when opening the web app +- Added subagent-focused UI with a side rail and detail drawer that surfaces each subagent’s lifecycle, running progress, and per-subagent transcript +- Added session status controls in the shell, including connection banners, toast notifications, and rejoin/new-link actions after a session ends +- Added the collab web package with the browser guest client, mock host, local relay, and relay contract tests. + +### Changed + +- Changed relay socket behavior to retry transient disconnections with exponential backoff while treating terminal relay-close conditions and decryption failures as non-retriable +- Changed subagent transcript decoding to handle streamed JSONL payload chunks incrementally by preserving carry-over data across chunks +- Replaced the vendored collab wire type mirror with shared `@oh-my-pi/pi-wire` protocol contracts. + +### Security + +- Hardened transcript Markdown rendering by escaping embedded HTML and allowing only safe link schemes diff --git a/packages/collab-web/README.md b/packages/collab-web/README.md new file mode 100644 index 000000000..60fe965f5 --- /dev/null +++ b/packages/collab-web/README.md @@ -0,0 +1,36 @@ +# @oh-my-pi/collab-web + +Web client for [omp collab sessions](../../docs/collab.md). Paste a `/collab` link into the browser and you get the same live session guests see in the TUI: streaming transcript, tool-call cards, subagent panel with live transcripts, and a composer that prompts (or interrupts) the host agent. + +## Quick start + +```sh +# dev server (Bun HTML dev server with HMR) — http://localhost:3000 +bun run dev + +# offline demo: local relay + scripted mock host; prints a ws://localhost link +bun run mock-host +``` + +Host a session from any omp instance (`/collab`, or `/collab ws://localhost:7466` to use the mock relay), then paste the printed link into the connect screen. Deep links work too: `http://localhost:3000/##` auto-connects on load. + +## Build & deploy + +```sh +bun run build # static site in dist/ +``` + +`dist/` is a fully static SPA — host it anywhere. Two runtime requirements: + +- **Secure context**: room keys are unwrapped with WebCrypto (`crypto.subtle`), which browsers expose only on `https://` or `localhost`. +- **Relay reachability**: the client connects straight to the relay over WebSocket (`wss://` for anything that isn't localhost). The default relay is `wss://relay.omp.sh`; bare `#` links resolve against it. + +The room key never leaves the URL fragment — it is not sent to the relay or any server. + +## Architecture + +- `src/lib/` — vendored wire codec (`codec.ts` AES-256-GCM, `link.ts` envelope + link grammar), `socket.ts` reconnecting relay socket, `client.ts` guest session store (`GuestClient` + immutable snapshots for `useSyncExternalStore`). Shared protocol shapes come from `@oh-my-pi/pi-wire`. +- `src/components/` — `transcript/` (entries, markdown, tool cards), `agents/` (panel + transcript drawer), `shell/` (connect screen, header, composer, banners, toasts). +- `scripts/` — `local-relay.ts` (content-blind relay on `Bun.serve`), `mock-host.ts` + `fixture.ts` (scripted host for offline dev). + +The package is intentionally standalone — no dependency on `@oh-my-pi/pi-coding-agent` at runtime or type level. Wire-shape drift is prevented by consuming the same `@oh-my-pi/pi-wire` contracts as the host, with sealed-frame interop still covered by `test/codec.test.ts`. diff --git a/packages/collab-web/index.html b/packages/collab-web/index.html new file mode 100644 index 000000000..4c5f5d86f --- /dev/null +++ b/packages/collab-web/index.html @@ -0,0 +1,13 @@ + + + + + + + omp collab + + +
+ + + diff --git a/packages/collab-web/package.json b/packages/collab-web/package.json new file mode 100644 index 000000000..e4106532f --- /dev/null +++ b/packages/collab-web/package.json @@ -0,0 +1,61 @@ +{ + "type": "module", + "name": "@oh-my-pi/collab-web", + "version": "15.11.7", + "private": true, + "description": "Browser guest client and local relay tools for omp collab live sessions", + "homepage": "https://omp.sh", + "author": "Can Boluk", + "license": "MIT", + "repository": { + "type": "git", + "url": "git+https://github.com/can1357/oh-my-pi.git", + "directory": "packages/collab-web" + }, + "bugs": { + "url": "https://github.com/can1357/oh-my-pi/issues" + }, + "keywords": [ + "collab", + "relay", + "websocket", + "react", + "agent" + ], + "scripts": { + "dev": "bun ./index.html", + "mock-host": "bun scripts/mock-host.ts", + "relay": "bun scripts/local-relay.ts", + "build": "bun build ./index.html --outdir=dist --minify", + "prepack": "bun run build", + "test": "bun test --parallel", + "check": "biome check . && bun run check:types", + "check:types": "tsgo -p tsconfig.json --noEmit", + "lint": "biome lint .", + "fix": "biome check --write --unsafe .", + "fmt": "biome format --write ." + }, + "dependencies": { + "@oh-my-pi/pi-wire": "catalog:", + "lucide-react": "catalog:", + "marked": "catalog:", + "react": "catalog:", + "react-dom": "catalog:" + }, + "devDependencies": { + "@types/bun": "catalog:", + "@types/react": "catalog:", + "@types/react-dom": "catalog:" + }, + "engines": { + "bun": ">=1.3.14" + }, + "files": [ + "dist", + "src", + "scripts", + "index.html", + "README.md", + "CHANGELOG.md" + ] +} diff --git a/packages/collab-web/scripts/fixture.ts b/packages/collab-web/scripts/fixture.ts new file mode 100644 index 000000000..7c4f51984 --- /dev/null +++ b/packages/collab-web/scripts/fixture.ts @@ -0,0 +1,656 @@ +/** + * Canned session data for the offline collab harness. + * + * One realistic host session (header + ~14 entries covering every renderer + * branch), an agent registry (main + a running sub with ticking progress + a + * parked sub with a transcript), a subagent transcript JSONL blob, and a + * scripted streaming turn the mock host replays on every guest prompt. + */ +import type { + AgentEvent, + AgentSnapshot, + AssistantMessage, + SessionEntry, + SessionHeader, + SubagentProgressPayload, + ToolCallContent, + ToolResultMessage, + WireModel, + WireUsage, +} from "@oh-my-pi/pi-wire"; + +export const HOST_DISPLAY_NAME = "kai"; + +export const fixtureModel: WireModel = { + id: "claude-haiku-4-5", + name: "Claude Haiku 4.5", + provider: "anthropic", + contextWindow: 200_000, +}; + +const NOW = Date.now(); +const MIN = 60_000; + +function iso(tsMs: number): string { + return new Date(tsMs).toISOString(); +} + +function mkUsage(input: number, output: number, cacheRead: number, cost: number): WireUsage { + return { input, output, cacheRead, cacheWrite: 0, totalTokens: input + output + cacheRead, cost: { total: cost } }; +} + +export const fixtureHeader: SessionHeader = { + type: "session", + id: "mock-collab-session", + title: "relay reconnect audit", + timestamp: iso(NOW - 32 * MIN), + cwd: "/Users/kai/Projects/pi", +}; + +const ASSISTANT_AUDIT_TEXT = `## Reconnect audit + +Three things to verify before touching anything: + +- backoff window (1s → 30s, exponential) +- fatal close codes that must *never* retry +- the resync \`hello\` → \`welcome\` handshake on reopen + +\`\`\`ts +const FATAL = new Set([4001, 4004, 4009, 4029]); +const delay = Math.min(1000 * 2 ** attempt, 30_000); +\`\`\` + +Checking the actual implementation now.`; + +export const fixtureEntries: SessionEntry[] = [ + { + id: "e01", + parentId: null, + timestamp: iso(NOW - 31 * MIN), + type: "message", + message: { + role: "user", + content: + "the guest socket sometimes stays dead after a relay redeploy — can you audit the reconnect path in relay-client.ts?", + timestamp: NOW - 31 * MIN, + }, + }, + { + id: "e02", + parentId: "e01", + timestamp: iso(NOW - 30 * MIN), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "thinking", + thinking: + "Reconnect semantics live in two places: the backoff loop in relay-client.ts and the fatal close codes in protocol.ts. I should read both before claiming anything — a redeploy closes with 1001, which must be treated as transient.", + }, + { type: "text", text: ASSISTANT_AUDIT_TEXT }, + { + type: "toolCall", + id: "call-bash-01", + name: "bash", + arguments: { command: 'rg -n "BACKOFF|MAX_PENDING" packages/coding-agent/src/collab/relay-client.ts' }, + intent: "Checking backoff constants", + }, + { + type: "toolCall", + id: "call-read-01", + name: "read", + arguments: { path: "packages/coding-agent/src/collab/relay-clinet.ts", offset: 160, limit: 60 }, + intent: "Reading the close handler", + }, + ], + model: fixtureModel.id, + usage: mkUsage(2_410, 386, 18_200, 0.0119), + stopReason: "toolUse", + timestamp: NOW - 30 * MIN, + }, + }, + { + id: "e03", + parentId: "e02", + timestamp: iso(NOW - 30 * MIN + 4_000), + type: "message", + message: { + role: "toolResult", + toolCallId: "call-bash-01", + toolName: "bash", + content: [ + { + type: "text", + text: "20:const BACKOFF_BASE_MS = 1_000;\n21:const BACKOFF_MAX_MS = 30_000;\n23:const MAX_PENDING_SENDS = 256;", + }, + ], + isError: false, + timestamp: NOW - 30 * MIN + 4_000, + }, + }, + { + id: "e04", + parentId: "e03", + timestamp: iso(NOW - 30 * MIN + 5_000), + type: "message", + message: { + role: "toolResult", + toolCallId: "call-read-01", + toolName: "read", + content: [ + { + type: "text", + text: "ENOENT: no such file or directory\n open 'packages/coding-agent/src/collab/relay-clinet.ts'\n at open (node:internal/fs/promises:642:23)\n at readRange (src/tools/read.ts:88:9)", + }, + ], + isError: true, + timestamp: NOW - 30 * MIN + 5_000, + }, + }, + { + id: "e05", + parentId: "e04", + timestamp: iso(NOW - 29 * MIN), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "text", + text: "Typo on my end — the file is `relay-client.ts`. The backoff constants match the spec (1s base, 30s cap, 256-frame reconnect buffer); reading the close handler next.", + }, + ], + model: fixtureModel.id, + usage: mkUsage(3_080, 64, 20_540, 0.0124), + stopReason: "stop", + timestamp: NOW - 29 * MIN, + }, + }, + { + id: "e06", + parentId: "e05", + timestamp: iso(NOW - 26 * MIN), + type: "message", + message: { + role: "user", + content: "what about half-open sockets after a NAT rebind? do we ping at all?", + timestamp: NOW - 26 * MIN, + }, + }, + { + id: "e07", + parentId: "e06", + timestamp: iso(NOW - 25 * MIN), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "text", + text: "The relay sends protocol-level pings, so a half-open socket dies within the idle window and lands in the normal close path. From there the client splits three ways:\n\n- transient drop → exponential retry, fresh `hello` on reopen\n- fatal close code (4001/4004/4009/4029) → surface the reason, stay down\n- decrypt failure → treated as fatal, since the key cannot heal", + }, + ], + model: fixtureModel.id, + usage: mkUsage(3_420, 142, 21_900, 0.0131), + stopReason: "stop", + timestamp: NOW - 25 * MIN, + }, + }, + { + id: "e08", + parentId: "e07", + timestamp: iso(NOW - 21 * MIN), + type: "custom_message", + customType: "collab-prompt", + content: + "also verify the relay drops the send buffer past 256 frames — we saw memory creep on a 6h session last week", + details: { from: "ada" }, + display: true, + }, + { + id: "e09", + parentId: "e08", + timestamp: iso(NOW - 20 * MIN), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "text", + text: "Good catch ada — `MAX_PENDING_SENDS` caps the reconnect buffer at 256 sealed envelopes and overflow is dropped (debug-logged), so no unbounded growth. The cost is silent frame loss, which the `welcome` resync on reopen covers.", + }, + ], + model: fixtureModel.id, + usage: mkUsage(3_900, 98, 23_400, 0.0138), + stopReason: "stop", + timestamp: NOW - 20 * MIN, + }, + }, + { + id: "e10", + parentId: "e09", + timestamp: iso(NOW - 14 * MIN), + type: "compaction", + summary: + "Audited the collab reconnect path: exponential backoff 1s→30s confirmed; fatal close codes 4001/4004/4009/4029 never retry; decrypt failures are terminal; the reconnect buffer is capped at 256 sealed frames with drop-on-overflow, recovered by the welcome resync. Open question: ping cadence on idle relays behind aggressive NATs.", + shortSummary: "reconnect audit findings", + firstKeptEntryId: "e08", + tokensBefore: 48_213, + }, + { + id: "e11", + parentId: "e10", + timestamp: iso(NOW - 13 * MIN), + type: "model_change", + model: fixtureModel.id, + }, + { + id: "e12", + parentId: "e11", + timestamp: iso(NOW - 13 * MIN + 5_000), + type: "thinking_level_change", + thinkingLevel: "medium", + }, + { + id: "e13", + parentId: "e12", + timestamp: iso(NOW - 9 * MIN), + type: "message", + message: { + role: "user", + content: "summarize what changed and what's left", + timestamp: NOW - 9 * MIN, + }, + }, + { + id: "e14", + parentId: "e13", + timestamp: iso(NOW - 8 * MIN), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "text", + text: "Done: backoff window verified, fatal-code table cross-checked against the relay, buffer cap documented. Left: an integration test that kills the relay mid-stream and asserts the guest resyncs from `welcome` without duplicated entries — RelayProbe is running that now, DocSweep already swept the docs.", + }, + ], + model: fixtureModel.id, + usage: mkUsage(12_640, 187, 31_780, 0.0212), + stopReason: "stop", + timestamp: NOW - 8 * MIN, + }, + }, +]; + +// ─── agents ────────────────────────────────────────────────────────────────── + +export const fixtureAgents: AgentSnapshot[] = [ + { + id: "main", + displayName: "Main", + kind: "main", + status: "running", + hasSessionFile: true, + createdAt: NOW - 32 * MIN, + lastActivity: NOW - 5_000, + }, + { + id: "RelayProbe", + displayName: "RelayProbe", + kind: "sub", + parentId: "main", + status: "running", + hasSessionFile: true, + createdAt: NOW - 6 * MIN, + lastActivity: NOW - 2_000, + }, + { + id: "DocSweep", + displayName: "DocSweep", + kind: "sub", + parentId: "main", + status: "parked", + hasSessionFile: true, + createdAt: NOW - 25 * MIN, + lastActivity: NOW - 11 * MIN, + }, +]; + +const PROBE_TOOLS = ["bash", "read", "search", "edit"] as const; +const PROBE_TOOL_ARGS: Record<(typeof PROBE_TOOLS)[number], string> = { + bash: "bun test packages/coding-agent/test/collab --filter reconnect", + read: "packages/coding-agent/src/collab/relay-client.ts:168-197", + search: "scheduleRetry|failFatal", + edit: "packages/coding-agent/test/collab/reconnect.test.ts", +}; + +/** Progress payload for the running sub; `tick` advances the counters. */ +export function makeProbeProgress(tick: number): SubagentProgressPayload { + const tool = PROBE_TOOLS[tick % PROBE_TOOLS.length]!; + const recentTools = [1, 2, 3].map(back => { + const prior = PROBE_TOOLS[(tick + PROBE_TOOLS.length * back - back) % PROBE_TOOLS.length]!; + return { tool: prior, args: PROBE_TOOL_ARGS[prior], endMs: Date.now() - back * 2_000 }; + }); + return { + index: 0, + agent: "task", + task: "probe relay reconnect under packet loss", + parentToolCallId: "call-task-01", + assignment: "Kill the relay mid-stream and assert the guest resyncs from welcome without duplicate entries.", + sessionFile: "/tmp/omp/agents/RelayProbe.jsonl", + progress: { + index: 0, + id: "RelayProbe", + agent: "task", + status: "running", + task: "probe relay reconnect under packet loss", + description: "relay reconnect probe", + lastIntent: "Replaying drop scenario", + currentTool: tool, + currentToolArgs: PROBE_TOOL_ARGS[tool], + recentTools, + recentOutput: ["3 sockets reconnected in 1.2s", "0 duplicate entries after resync"], + toolCount: 9 + tick, + requests: 4 + Math.floor(tick / 3), + tokens: 18_400 + tick * 450, + contextTokens: 22_300 + tick * 510, + contextWindow: fixtureModel.contextWindow, + cost: 0.041 + tick * 0.0012, + durationMs: 95_000 + tick * 2_000, + resolvedModel: fixtureModel.id, + }, + }; +} + +// ─── subagent transcript ───────────────────────────────────────────────────── + +const SUB_T0 = NOW - 25 * MIN; + +const subagentTranscriptLines: unknown[] = [ + { type: "session", id: "mock-docsweep", timestamp: iso(SUB_T0), cwd: "/Users/kai/Projects/pi" }, + // Unknown entry type — guests must skip it (tolerant default branch). + { type: "session_init", id: "s00", parentId: null, timestamp: iso(SUB_T0), version: 3 }, + { + id: "s01", + parentId: null, + timestamp: iso(SUB_T0 + 2_000), + type: "message", + message: { + role: "user", + content: "Sweep docs/collab.md for stale close-code references and report mismatches.", + timestamp: SUB_T0 + 2_000, + }, + }, + { + id: "s02", + parentId: "s01", + timestamp: iso(SUB_T0 + 9_000), + type: "message", + message: { + role: "assistant", + content: [ + { type: "thinking", thinking: "Grep the doc for 4xxx codes, then diff against protocol.ts." }, + { type: "text", text: "Scanning `docs/collab.md` for close-code mentions." }, + { + type: "toolCall", + id: "sub-call-01", + name: "search", + arguments: { pattern: "40\\d\\d", path: "docs/collab.md" }, + intent: "Finding close codes", + }, + ], + model: fixtureModel.id, + usage: mkUsage(1_180, 96, 0, 0.0021), + stopReason: "toolUse", + timestamp: SUB_T0 + 9_000, + }, + }, + { + id: "s03", + parentId: "s02", + timestamp: iso(SUB_T0 + 11_000), + type: "message", + message: { + role: "toolResult", + toolCallId: "sub-call-01", + toolName: "search", + content: [ + { + type: "text", + text: "docs/collab.md:41: 4001 room closed\ndocs/collab.md:42: 4004 no such room\ndocs/collab.md:43: 4009 host conflict\ndocs/collab.md:44: 4029 room full", + }, + ], + isError: false, + timestamp: SUB_T0 + 11_000, + }, + }, + { + id: "s04", + parentId: "s03", + timestamp: iso(SUB_T0 + 20_000), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "toolCall", + id: "sub-call-02", + name: "read", + arguments: { path: "packages/coding-agent/src/collab/relay-client.ts", offset: 13, limit: 6 }, + intent: "Cross-checking the fatal table", + }, + ], + model: fixtureModel.id, + usage: mkUsage(1_460, 41, 980, 0.0024), + stopReason: "toolUse", + timestamp: SUB_T0 + 20_000, + }, + }, + { + id: "s05", + parentId: "s04", + timestamp: iso(SUB_T0 + 22_000), + type: "message", + message: { + role: "toolResult", + toolCallId: "sub-call-02", + toolName: "read", + content: [ + { + type: "text", + text: '13:const FATAL_CLOSE_REASONS: Record = {\n14:\t4001: "room closed",\n15:\t4004: "no such room",\n16:\t4009: "a host is already connected for this room",\n17:\t4029: "room is full",\n18:};', + }, + ], + isError: false, + timestamp: SUB_T0 + 22_000, + }, + }, + { + id: "s06", + parentId: "s05", + timestamp: iso(SUB_T0 + 30_000), + type: "message", + message: { + role: "assistant", + content: [ + { + type: "text", + text: "All four close codes in `docs/collab.md` match `FATAL_CLOSE_REASONS` — no stale references. The doc could additionally mention that decrypt failures are fatal; flagged as a suggestion, not a mismatch.", + }, + ], + model: fixtureModel.id, + usage: mkUsage(1_720, 88, 1_240, 0.0029), + stopReason: "stop", + timestamp: SUB_T0 + 30_000, + }, + }, +]; + +/** DocSweep's session file, served by the mock host's fetch-transcript handler. */ +export const subagentTranscriptJsonl: string = `${subagentTranscriptLines + .map(line => JSON.stringify(line)) + .join("\n")}\n`; + +// ─── scripted streaming turn ───────────────────────────────────────────────── + +export type ScriptedStep = + | { kind: "event"; event: AgentEvent } + | { kind: "entry"; entry: SessionEntry } + | { kind: "state"; streaming: boolean }; + +const TURN_THINKING_1 = "Guest wants a live check. "; +const TURN_THINKING_2 = "Guest wants a live check. I'll run the reconnect suite once and summarize the result."; +const TURN_TEXT_1 = "Kicking off a live reconnect probe "; +const TURN_TEXT_2 = "Kicking off a live reconnect probe — one suite run, then a verdict."; +const TURN_CLOSE_1 = "Probe passed: 3 reconnects, "; +const TURN_CLOSE_2 = "Probe passed: 3 reconnects, 0 duplicate entries after resync. The reconnect path holds."; + +/** + * One scripted streaming turn, replayed by the mock host at ~40ms cadence. + * + * `seq` keeps ids unique across replays; `parentId` chains the appended + * entries onto the current transcript tail. + */ +export function makeScriptedTurn(seq: number, parentId: string | null): ScriptedStep[] { + const ts = Date.now(); + const a1Id = `turn${seq}-a1`; + const callId = `turn${seq}-call1`; + const r1Id = `turn${seq}-r1`; + const a2Id = `turn${seq}-a2`; + const command = "bun test packages/coding-agent/test/collab --filter reconnect"; + const toolResultText = + "3 tests passed (reconnect.test.ts)\n3 sockets reconnected in 1.2s\n0 duplicate entries after resync"; + + const partial = (content: AssistantMessage["content"]): AssistantMessage => ({ + role: "assistant", + content, + model: fixtureModel.id, + usage: mkUsage(0, 0, 0, 0), + stopReason: "stop", + timestamp: ts, + }); + + const toolCall: ToolCallContent = { + type: "toolCall", + id: callId, + name: "bash", + arguments: { command }, + intent: "Running the reconnect suite", + }; + + const a1Final: AssistantMessage = { + role: "assistant", + content: [{ type: "thinking", thinking: TURN_THINKING_2 }, { type: "text", text: TURN_TEXT_2 }, toolCall], + model: fixtureModel.id, + usage: mkUsage(4_310, 164, 24_800, 0.0147), + stopReason: "toolUse", + timestamp: ts, + }; + + const r1Message: ToolResultMessage = { + role: "toolResult", + toolCallId: callId, + toolName: "bash", + content: [{ type: "text", text: toolResultText }], + isError: false, + timestamp: ts, + }; + + const a2Final: AssistantMessage = { + role: "assistant", + content: [{ type: "text", text: TURN_CLOSE_2 }], + model: fixtureModel.id, + usage: mkUsage(4_690, 52, 25_400, 0.0153), + stopReason: "stop", + timestamp: ts, + }; + + return [ + { kind: "event", event: { type: "agent_start" } }, + { kind: "state", streaming: true }, + { kind: "event", event: { type: "turn_start" } }, + { kind: "event", event: { type: "message_start", message: partial([]) } }, + { + kind: "event", + event: { type: "message_update", message: partial([{ type: "thinking", thinking: TURN_THINKING_1 }]) }, + }, + { + kind: "event", + event: { type: "message_update", message: partial([{ type: "thinking", thinking: TURN_THINKING_2 }]) }, + }, + { + kind: "event", + event: { + type: "message_update", + message: partial([ + { type: "thinking", thinking: TURN_THINKING_2 }, + { type: "text", text: TURN_TEXT_1 }, + ]), + }, + }, + { + kind: "event", + event: { + type: "message_update", + message: partial([ + { type: "thinking", thinking: TURN_THINKING_2 }, + { type: "text", text: TURN_TEXT_2 }, + ]), + }, + }, + { + kind: "event", + event: { + type: "message_update", + message: partial([ + { type: "thinking", thinking: TURN_THINKING_2 }, + { type: "text", text: TURN_TEXT_2 }, + toolCall, + ]), + }, + }, + { kind: "event", event: { type: "message_end", message: a1Final } }, + { kind: "entry", entry: { id: a1Id, parentId, timestamp: iso(ts), type: "message", message: a1Final } }, + { + kind: "event", + event: { + type: "tool_execution_start", + toolCallId: callId, + toolName: "bash", + args: { command }, + intent: "Running the reconnect suite", + }, + }, + { + kind: "event", + event: { + type: "tool_execution_update", + toolCallId: callId, + toolName: "bash", + args: { command }, + partialResult: "3 sockets reconnected in 1.2s", + }, + }, + { + kind: "event", + event: { + type: "tool_execution_end", + toolCallId: callId, + toolName: "bash", + result: toolResultText, + isError: false, + }, + }, + { kind: "entry", entry: { id: r1Id, parentId: a1Id, timestamp: iso(ts), type: "message", message: r1Message } }, + { kind: "event", event: { type: "message_start", message: partial([]) } }, + { kind: "event", event: { type: "message_update", message: partial([{ type: "text", text: TURN_CLOSE_1 }]) } }, + { kind: "event", event: { type: "message_update", message: partial([{ type: "text", text: TURN_CLOSE_2 }]) } }, + { kind: "event", event: { type: "message_end", message: a2Final } }, + { kind: "entry", entry: { id: a2Id, parentId: r1Id, timestamp: iso(ts), type: "message", message: a2Final } }, + { kind: "event", event: { type: "turn_end" } }, + { kind: "event", event: { type: "agent_end" } }, + { kind: "state", streaming: false }, + ]; +} diff --git a/packages/collab-web/scripts/local-relay.ts b/packages/collab-web/scripts/local-relay.ts new file mode 100644 index 000000000..cfa1eeae2 --- /dev/null +++ b/packages/collab-web/scripts/local-relay.ts @@ -0,0 +1,169 @@ +/** + * Offline stand-in for the public collab relay (`wss://relay.omp.sh`). + * + * Speaks the exact relay contract the real clients expect: + * - `GET /r/?role=host|guest` upgrades to a WebSocket. + * - The host creates the room; a second host is rejected with close 4009 and + * a guest joining a missing room with close 4004. + * - Host binary frames: envelope peerId 0 broadcasts to every guest, peerId N + * targets that guest only — forwarded unchanged either way. + * - Guest binary frames: the first 4 envelope bytes are rewritten to the + * sender's peerId, then forwarded to the host. + * - TEXT control to the host: `{"t":"peer-joined","peer":N}` / `{"t":"peer-left","peer":N}`. + * - Host disconnect: TEXT `{"t":"room-closed"}` to every guest, then close 4001 + * and the room is garbage-collected. + * + * The relay never sees plaintext: payloads stay sealed end to end. + */ +import { rewriteEnvelopePeer, unpackEnvelope } from "../src/lib/link"; + +const ROOM_PATH_RE = /^\/r\/([A-Za-z0-9_-]{10,64})$/; + +const DEFAULT_PORT = 7466; + +interface SocketData { + roomId: string; + role: "host" | "guest"; + /** Assigned on open for guests; the host stays 0. */ + peerId: number; +} + +type RelaySocket = Bun.ServerWebSocket; + +interface Room { + host: RelaySocket; + guests: Map; + nextPeerId: number; +} + +export interface LocalRelay { + /** ws://localhost: — append `/r/?role=…` to connect. */ + url: string; + /** Closes every room and stops the server. Idempotent. */ + stop(): void; +} + +export function startLocalRelay(port = 0): LocalRelay { + const rooms = new Map(); + + const server = Bun.serve({ + port, + fetch(req, srv): Response | undefined { + const url = new URL(req.url); + const match = ROOM_PATH_RE.exec(url.pathname); + const role = url.searchParams.get("role"); + if (!match || (role !== "host" && role !== "guest")) { + return new Response("not found", { status: 404 }); + } + const data: SocketData = { roomId: match[1]!, role, peerId: 0 }; + if (srv.upgrade(req, { data })) return undefined; + return new Response("websocket upgrade required", { status: 426 }); + }, + websocket: { + open(ws: RelaySocket): void { + const { roomId, role } = ws.data; + if (role === "host") { + if (rooms.has(roomId)) { + ws.close(4009, "a host is already connected for this room"); + return; + } + rooms.set(roomId, { host: ws, guests: new Map(), nextPeerId: 1 }); + return; + } + const room = rooms.get(roomId); + if (!room) { + ws.close(4004, "no such room"); + return; + } + const peerId = room.nextPeerId++; + ws.data.peerId = peerId; + room.guests.set(peerId, ws); + room.host.send(JSON.stringify({ t: "peer-joined", peer: peerId })); + }, + message(ws: RelaySocket, message: string | Buffer): void { + if (typeof message === "string") return; // clients never send TEXT + const room = rooms.get(ws.data.roomId); + if (!room) return; + if (ws.data.role === "host") { + const envelope = unpackEnvelope(message); + if (!envelope) return; + if (envelope.peerId === 0) { + for (const guest of room.guests.values()) guest.send(message); + } else { + room.guests.get(envelope.peerId)?.send(message); + } + return; + } + if (message.byteLength < 4) return; + rewriteEnvelopePeer(message, ws.data.peerId); + room.host.send(message); + }, + close(ws: RelaySocket): void { + const { roomId, role, peerId } = ws.data; + const room = rooms.get(roomId); + if (!room) return; + if (role === "host") { + // Rejected second host: the live room is not ours to tear down. + if (room.host !== ws) return; + rooms.delete(roomId); + const closure = JSON.stringify({ t: "room-closed" }); + for (const guest of room.guests.values()) { + guest.send(closure); + guest.close(4001, "room closed"); + } + room.guests.clear(); + return; + } + if (room.guests.delete(peerId)) { + room.host.send(JSON.stringify({ t: "peer-left", peer: peerId })); + } + }, + }, + }); + + return { + url: `ws://localhost:${server.port}`, + stop(): void { + for (const room of rooms.values()) { + const closure = JSON.stringify({ t: "room-closed" }); + for (const guest of room.guests.values()) { + guest.send(closure); + guest.close(4001, "room closed"); + } + room.host.close(1001, "relay shutting down"); + } + rooms.clear(); + server.stop(true); + }, + }; +} +function parsePort(argv: readonly string[]): number { + let raw: string | undefined; + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]!; + if (arg === "--port") raw = argv[i + 1]; + else if (arg.startsWith("--port=")) raw = arg.slice("--port=".length); + } + if (raw === undefined) return DEFAULT_PORT; + const port = Number(raw); + if (!Number.isInteger(port) || port < 0 || port > 65_535) { + console.error(`local-relay: invalid --port ${raw}`); + process.exit(1); + } + return port; +} + +if (import.meta.main) { + const relay = startLocalRelay(parsePort(Bun.argv.slice(2))); + let stopping = false; + const shutdown = (): void => { + if (stopping) return; + stopping = true; + relay.stop(); + process.exit(0); + }; + console.log(`local collab relay listening on ${relay.url}`); + console.log("connect with /r/?role=host|guest; Ctrl+C stops the relay"); + process.on("SIGINT", shutdown); + process.on("SIGTERM", shutdown); +} diff --git a/packages/collab-web/scripts/mock-host.ts b/packages/collab-web/scripts/mock-host.ts new file mode 100644 index 000000000..01cf06e49 --- /dev/null +++ b/packages/collab-web/scripts/mock-host.ts @@ -0,0 +1,357 @@ +/** + * Offline mock collab host: starts the local relay, opens a room as host, and + * serves the canned fixture session to any collab-web guest that joins. + * + * bun scripts/mock-host.ts [--port 7466] + * + * Replays a scripted streaming turn on every guest prompt, ticks subagent + * progress on the bus every 2s, and answers fetch-transcript with byte slices + * of the fixture JSONL — exactly the frames a real `omp /collab` host emits. + */ + +import type { AgentSnapshot, HostFrame, SessionEntry, SessionState, WireFrame } from "@oh-my-pi/pi-wire"; +import { generateRoomKey, importRoomKey, open, seal } from "../src/lib/codec"; +import { COLLAB_PROTO, formatCollabLink, generateRoomId, packEnvelope, unpackEnvelope } from "../src/lib/link"; +import { + fixtureAgents, + fixtureEntries, + fixtureHeader, + fixtureModel, + HOST_DISPLAY_NAME, + makeProbeProgress, + makeScriptedTurn, + type ScriptedStep, + subagentTranscriptJsonl, +} from "./fixture"; +import { startLocalRelay } from "./local-relay"; + +const DEFAULT_PORT = 7466; +const STEP_INTERVAL_MS = 40; +const TICK_INTERVAL_MS = 2_000; +const AGENTS_SNAPSHOT_EVERY = 5; + +function parsePort(argv: string[]): number { + let raw: string | undefined; + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]!; + if (arg === "--port") raw = argv[i + 1]; + else if (arg.startsWith("--port=")) raw = arg.slice("--port=".length); + } + if (raw === undefined) return DEFAULT_PORT; + const port = Number(raw); + if (!Number.isInteger(port) || port < 1 || port > 65_535) { + console.error(`mock-host: invalid --port ${raw}`); + process.exit(1); + } + return port; +} + +const port = parsePort(Bun.argv.slice(2)); +const relay = startLocalRelay(port); +const roomId = generateRoomId(); +const rawKey = generateRoomKey(); +const key = await importRoomKey(rawKey); +const link = formatCollabLink(relay.url, roomId, rawKey); + +// ── mutable session state ──────────────────────────────────────────────────── + +const entries: SessionEntry[] = [...fixtureEntries]; +const agents: AgentSnapshot[] = fixtureAgents.map(agent => ({ ...agent })); +const peers = new Map(); +const transcriptBytes = new TextEncoder().encode(subagentTranscriptJsonl); +const transcriptDecoder = new TextDecoder(); + +let lastEntryId: string | null = entries[entries.length - 1]?.id ?? null; +let streaming = false; +let queuedPrompts = 0; +let turnSeq = 0; +let liveEntrySeq = 0; +let replayQueue: ScriptedStep[] = []; +let replayTimer: Timer | null = null; +let tick = 0; +let shuttingDown = false; + +// ── sealed transport (order-preserving, mirrors relay-client) ──────────────── + +const ws = new WebSocket(`${relay.url}/r/${roomId}?role=host`); +ws.binaryType = "arraybuffer"; + +let sendChain: Promise = Promise.resolve(); +let recvChain: Promise = Promise.resolve(); + +/** Seal and send a frame; peerId 0 broadcasts, N targets that guest. */ +function sendFrame(frame: HostFrame, targetPeer = 0): void { + sendChain = sendChain + .then(async () => { + if (ws.readyState !== WebSocket.OPEN) return; + const sealed = await seal(key, frame); + ws.send(packEnvelope(targetPeer, sealed)); + }) + .catch((err: unknown) => { + console.error("mock-host: send failed:", err); + }); +} + +function buildState(): SessionState { + const participants: SessionState["participants"] = [{ name: HOST_DISPLAY_NAME, role: "host" }]; + for (const name of peers.values()) participants.push({ name, role: "guest" }); + const tokens = 51_000 + entries.length * 120; + return { + isStreaming: streaming, + queuedMessageCount: queuedPrompts, + sessionName: fixtureHeader.title, + cwd: fixtureHeader.cwd, + model: fixtureModel, + thinkingLevel: "medium", + contextUsage: { + tokens, + contextWindow: fixtureModel.contextWindow, + percent: (tokens / fixtureModel.contextWindow) * 100, + }, + participants, + }; +} + +function broadcastState(): void { + sendFrame({ t: "state", state: buildState() }); +} + +function appendEntry(entry: SessionEntry): void { + entries.push(entry); + lastEntryId = entry.id; + sendFrame({ t: "entry", entry }); +} + +function notice(level: "info" | "warning" | "error", message: string): void { + sendFrame({ t: "event", event: { type: "notice", level, message, source: "collab" } }); +} + +// ── scripted turn replay ───────────────────────────────────────────────────── + +function startReplay(): void { + turnSeq++; + replayQueue = makeScriptedTurn(turnSeq, lastEntryId); + scheduleStep(); +} + +function scheduleStep(): void { + replayTimer = setTimeout(() => { + replayTimer = null; + const step = replayQueue.shift(); + if (step) applyStep(step); + if (replayQueue.length > 0) { + scheduleStep(); + return; + } + if (queuedPrompts > 0) { + queuedPrompts--; + startReplay(); + } + }, STEP_INTERVAL_MS); +} + +function applyStep(step: ScriptedStep): void { + switch (step.kind) { + case "event": + sendFrame({ t: "event", event: step.event }); + break; + case "entry": + appendEntry(step.entry); + break; + case "state": + streaming = step.streaming; + broadcastState(); + break; + } +} + +function cancelReplay(): void { + if (replayTimer !== null) { + clearTimeout(replayTimer); + replayTimer = null; + } + replayQueue = []; +} + +// ── guest frame handling ───────────────────────────────────────────────────── + +function peerName(fromPeer: number): string { + return peers.get(fromPeer) ?? `guest-${fromPeer}`; +} + +function handleHello(name: string, proto: number, fromPeer: number): void { + if (proto !== COLLAB_PROTO) { + sendFrame( + { 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}`; + peers.set(fromPeer, cleanName); + sendFrame( + { + t: "welcome", + proto: COLLAB_PROTO, + header: fixtureHeader, + entries: [...entries], + state: buildState(), + agents: agents.map(agent => ({ ...agent })), + }, + fromPeer, + ); + console.log(`mock-host: ${cleanName} joined (peer ${fromPeer})`); + broadcastState(); +} + +function handlePrompt(text: string, fromPeer: number): void { + liveEntrySeq++; + appendEntry({ + id: `live-${liveEntrySeq}`, + parentId: lastEntryId, + timestamp: new Date().toISOString(), + type: "custom_message", + customType: "collab-prompt", + content: text, + details: { from: peerName(fromPeer) }, + display: true, + }); + if (replayTimer !== null || replayQueue.length > 0) { + queuedPrompts++; + broadcastState(); + return; + } + startReplay(); +} + +function handleAbort(fromPeer: number): void { + const wasReplaying = replayTimer !== null || replayQueue.length > 0; + cancelReplay(); + queuedPrompts = 0; + notice("info", `${peerName(fromPeer)} interrupted`); + if (wasReplaying) sendFrame({ t: "event", event: { type: "agent_end" } }); + streaming = false; + broadcastState(); +} + +function handleAgentCmd(cmd: string, agentId: string, fromPeer: number): void { + notice("info", `${peerName(fromPeer)} sent agent-cmd ${cmd} → ${agentId}`); +} + +function handleFetchTranscript(reqId: number, fromByte: number, fromPeer: number): void { + const total = transcriptBytes.byteLength; + const start = Math.max(0, Math.min(fromByte, total)); + const text = start >= total ? "" : transcriptDecoder.decode(transcriptBytes.subarray(start)); + // We always serve to EOF, so the next offset base is the full size. + sendFrame({ t: "transcript", reqId, text, newSize: total }, fromPeer); +} + +function handleFrame(frame: WireFrame, fromPeer: number): void { + switch (frame.t) { + case "hello": + handleHello(frame.name, frame.proto, fromPeer); + break; + case "prompt": + handlePrompt(frame.text, fromPeer); + break; + case "abort": + handleAbort(fromPeer); + break; + case "agent-cmd": + handleAgentCmd(frame.cmd, frame.agentId, fromPeer); + break; + case "fetch-transcript": + handleFetchTranscript(frame.reqId, frame.fromByte, fromPeer); + break; + default: + // Host-frame echoes or unknown types: ignore. + break; + } +} + +function handleControl(text: string): void { + let msg: unknown; + try { + msg = JSON.parse(text); + } catch { + return; + } + if (typeof msg !== "object" || msg === null) return; + const control = msg as { t?: unknown; peer?: unknown }; + if (control.t === "peer-left" && typeof control.peer === "number") { + const name = peers.get(control.peer); + peers.delete(control.peer); + if (name) console.log(`mock-host: ${name} left (peer ${control.peer})`); + broadcastState(); + } +} + +ws.onopen = () => { + console.log("mock collab host ready"); + console.log(`join link: ${link}`); + console.log("paste the link into the collab-web connect screen (bun ./index.html), Ctrl+C stops the host"); +}; + +ws.onmessage = event => { + const data: unknown = event.data; + if (typeof data === "string") { + handleControl(data); + return; + } + if (!(data instanceof ArrayBuffer)) return; + const envelope = unpackEnvelope(new Uint8Array(data)); + if (!envelope) return; + recvChain = recvChain + .then(async () => { + const frame = await open(key, envelope.payload); + handleFrame(frame, envelope.peerId); + }) + .catch((err: unknown) => { + console.error("mock-host: dropping undecryptable frame:", err); + }); +}; + +ws.onclose = event => { + if (shuttingDown) return; + console.error(`mock-host: relay socket closed (${event.code} ${event.reason || "no reason"})`); + shutdown(1); +}; + +// ── progress ticker ────────────────────────────────────────────────────────── + +const tickInterval: Timer = setInterval(() => { + tick++; + sendFrame({ t: "bus", channel: "task:subagent:progress", data: makeProbeProgress(tick) }); + const now = Date.now(); + for (const agent of agents) { + if (agent.status === "running") agent.lastActivity = now; + } + if (tick % AGENTS_SNAPSHOT_EVERY === 0) { + sendFrame({ t: "agents", agents: agents.map(agent => ({ ...agent })) }); + } +}, TICK_INTERVAL_MS); + +// ── shutdown ───────────────────────────────────────────────────────────────── + +function shutdown(code: number): void { + if (shuttingDown) return; + shuttingDown = true; + cancelReplay(); + clearInterval(tickInterval); + sendFrame({ t: "bye", reason: "mock host shutting down" }); + // Let the bye flush through the send chain before tearing the room down. + void sendChain.finally(() => { + try { + ws.close(1000); + } catch { + // already closing + } + relay.stop(); + process.exit(code); + }); +} + +process.on("SIGINT", () => { + console.log("\nmock-host: shutting down"); + shutdown(0); +}); diff --git a/packages/collab-web/src/app.tsx b/packages/collab-web/src/app.tsx new file mode 100644 index 000000000..916e318d7 --- /dev/null +++ b/packages/collab-web/src/app.tsx @@ -0,0 +1,170 @@ +import type { ReactNode } from "react"; +import { useCallback, useEffect, useMemo, useRef, useState } from "react"; +import { AgentDrawer } from "./components/agents/AgentDrawer"; +import { AgentsPanel } from "./components/agents/AgentsPanel"; +import { Banners } from "./components/shell/Banners"; +import { Composer } from "./components/shell/Composer"; +import { ConnectScreen } from "./components/shell/ConnectScreen"; +import { HeaderBar } from "./components/shell/HeaderBar"; +import { Toasts } from "./components/shell/Toasts"; +import { Transcript } from "./components/transcript/Transcript"; +import { GuestClient } from "./lib/client"; +import { useGuestSnapshot } from "./lib/use-guest"; +import "./components/shell/shell.css"; + +const NAME_KEY = "omp.collab.name"; + +interface Creds { + link: string; + name: string; +} + +function storedName(): string { + try { + return localStorage.getItem(NAME_KEY) ?? "guest"; + } catch { + return "guest"; + } +} + +/** Deep link = everything after the FIRST `#` (the link itself contains another `#`). */ +function hashLink(): string | null { + const href = window.location.href; + const i = href.indexOf("#"); + if (i < 0 || i + 1 >= href.length) return null; + return href.slice(i + 1); +} + +export function App(): ReactNode { + const [client, setClient] = useState(null); + const [connectError, setConnectError] = useState(null); + const credsRef = useRef(null); + + const connect = useCallback((link: string, name: string): void => { + let next: GuestClient; + try { + next = new GuestClient(link, name); + } catch (err) { + setConnectError(err instanceof Error ? err.message : String(err)); + return; + } + next.connect(); + try { + localStorage.setItem(NAME_KEY, name); + } catch { + // storage unavailable (private mode) — non-fatal + } + credsRef.current = { link, name }; + window.location.hash = link; + setConnectError(null); + setClient(prev => { + prev?.close(); + return next; + }); + }, []); + + const leave = useCallback((): void => { + setClient(prev => { + prev?.close(); + return null; + }); + history.replaceState(null, "", window.location.pathname + window.location.search); + }, []); + + const rejoin = useCallback((): void => { + const creds = credsRef.current; + if (creds) connect(creds.link, creds.name); + }, [connect]); + + // Deep link: a page load with a hash auto-connects. + useEffect(() => { + const link = hashLink(); + if (link) connect(link, storedName()); + }, [connect]); + + useEffect(() => { + if (!client) document.title = "omp collab"; + }, [client]); + + if (!client) { + return ; + } + return ; +} + +interface SessionProps { + client: GuestClient; + onLeave(): void; + onRejoin(): void; +} + +function Session({ client, onLeave, onRejoin }: SessionProps): ReactNode { + const snap = useGuestSnapshot(client); + const [railOpen, setRailOpen] = useState(false); + const [selectedId, setSelectedId] = useState(null); + const autoOpenedRef = useRef(false); + + const subCount = useMemo(() => snap.agents.filter(a => a.kind === "sub").length, [snap.agents]); + + // Auto-open the rail the first time a subagent appears. + useEffect(() => { + if (subCount > 0 && !autoOpenedRef.current) { + autoOpenedRef.current = true; + setRailOpen(true); + } + }, [subCount]); + + const title = snap.header?.title ?? snap.state?.sessionName ?? "session"; + useEffect(() => { + document.title = `${title} · omp collab`; + }, [title]); + + const drawerAgent = selectedId != null ? snap.agents.find(a => a.id === selectedId) : undefined; + + return ( +
+ setRailOpen(open => !open)} + onLeave={onLeave} + /> +
+
+
+ +
+
+ {railOpen && ( + + )} +
+ + {drawerAgent && ( + setSelectedId(null)} + /> + )} + + +
+ ); +} diff --git a/packages/collab-web/src/components/agents/AgentDrawer.tsx b/packages/collab-web/src/components/agents/AgentDrawer.tsx new file mode 100644 index 000000000..6078ace60 --- /dev/null +++ b/packages/collab-web/src/components/agents/AgentDrawer.tsx @@ -0,0 +1,184 @@ +import type { AgentSnapshot, SessionEntry, SubagentProgressPayload } from "@oh-my-pi/pi-wire"; +import { OctagonX, RotateCcw, SendHorizontal, X } from "lucide-react"; +import type { ReactNode } from "react"; +import { useEffect, useState } from "react"; +import type { GuestClient } from "../../lib/client"; +import { fmtCost, fmtDuration, fmtTokens } from "../../lib/format"; +import { parseJsonl } from "../../lib/jsonl"; +import type { TranscriptProps } from "../transcript/Transcript"; +import { Transcript } from "../transcript/Transcript"; + +const EMPTY_TOOLS: TranscriptProps["activeTools"] = new Map(); +const POLL_MS = 1200; + +export function AgentDrawer(props: { + agent: AgentSnapshot; + progress?: SubagentProgressPayload; + client: GuestClient; + onClose(): void; +}): ReactNode { + const { agent, progress, client, onClose } = props; + const [entries, setEntries] = useState([]); + const [draft, setDraft] = useState(""); + + useEffect(() => { + const onKey = (e: KeyboardEvent) => { + if (e.key === "Escape") onClose(); + }; + window.addEventListener("keydown", onKey); + return () => window.removeEventListener("keydown", onKey); + }, [onClose]); + + // Live transcript: poll the host-side session file while the drawer is + // open, appending parsed JSONL entries. State resets when the agent + // changes; the interval and any in-flight reply are dropped on cleanup. + useEffect(() => { + setEntries([]); + if (!agent.hasSessionFile) return; + let disposed = false; + let inFlight = false; + let cursor = 0; + let carry = ""; + let acc: readonly SessionEntry[] = []; + const poll = async (): Promise => { + if (disposed || inFlight) return; + inFlight = true; + try { + const reply = await client.fetchTranscript(agent.id, cursor); + if (disposed || reply === null) return; // timeout/error → keep polling + cursor = reply.newSize; + if (!reply.text) return; + const parsed = parseJsonl(reply.text, carry); + carry = parsed.carry; + const fresh: SessionEntry[] = []; + for (const item of parsed.items) { + if (typeof item !== "object" || item === null) continue; + if ((item as { type?: unknown }).type === "session") continue; + fresh.push(item as SessionEntry); + } + if (fresh.length > 0) { + acc = [...acc, ...fresh]; + setEntries(acc); + } + } finally { + inFlight = false; + } + }; + void poll(); + const timer = setInterval(() => { + void poll(); + }, POLL_MS); + return () => { + disposed = true; + clearInterval(timer); + }; + }, [agent.id, agent.hasSessionFile, client]); + + const sendChat = () => { + const text = draft.trim(); + if (!text) return; + client.sendAgentCmd("chat", agent.id, text); + setDraft(""); + }; + + const p = progress?.progress; + const model = p?.resolvedModel; + const ctxPct = + p?.contextTokens !== undefined && p.contextWindow + ? Math.min(100, (p.contextTokens / p.contextWindow) * 100) + : null; + + return ( + + ); +} diff --git a/packages/collab-web/src/components/agents/AgentsPanel.tsx b/packages/collab-web/src/components/agents/AgentsPanel.tsx new file mode 100644 index 000000000..92bbed098 --- /dev/null +++ b/packages/collab-web/src/components/agents/AgentsPanel.tsx @@ -0,0 +1,131 @@ +import type { + AgentProgress, + AgentSnapshot, + SubagentLifecyclePayload, + SubagentProgressPayload, +} from "@oh-my-pi/pi-wire"; +import type { ReactNode } from "react"; +import { useEffect, useMemo, useState } from "react"; +import { fmtCost, fmtDuration, fmtTokens, relTime } from "../../lib/format"; +import "./agents.css"; + +/** Re-render tick so running-tool durations and relative times stay live. */ +function useNow(intervalMs: number): number { + const [now, setNow] = useState(() => Date.now()); + useEffect(() => { + const timer = setInterval(() => setNow(Date.now()), intervalMs); + return () => clearInterval(timer); + }, [intervalMs]); + return now; +} + +/** + * Best-effort start timestamp for the in-flight tool. The host serializes the + * full AgentProgress (which carries `currentToolStartMs`); the wire mirror + * omits it, so read it tolerantly and fall back to the last tool's end time. + */ +function toolStartMs(p: AgentProgress): number | null { + const start = (p as { currentToolStartMs?: unknown }).currentToolStartMs; + if (typeof start === "number") return start; + const lastEnd = p.recentTools[0]?.endMs; + return typeof lastEnd === "number" ? lastEnd : null; +} + +function activityLine( + agent: AgentSnapshot, + p: AgentProgress | undefined, + lc: SubagentLifecyclePayload | undefined, + now: number, +): string { + if (p?.currentTool) { + const start = toolStartMs(p); + if (start !== null) return `${p.currentTool} · ${fmtDuration(Math.max(0, now - start))}`; + return p.currentTool; + } + if (p?.lastIntent) return p.lastIntent; + if (lc) return lc.status; + return agent.status; +} + +function AgentRow(props: { + agent: AgentSnapshot; + payload: SubagentProgressPayload | undefined; + lifecycle: SubagentLifecyclePayload | undefined; + selected: boolean; + now: number; + onSelect(id: string | null): void; +}): ReactNode { + const { agent, payload, lifecycle, selected, now, onSelect } = props; + const p = payload?.progress; + return ( + + ); +} + +export function AgentsPanel(props: { + agents: readonly AgentSnapshot[]; + progress: ReadonlyMap; + lifecycle: ReadonlyMap; + selectedId: string | null; + onSelect(id: string | null): void; +}): ReactNode { + const { agents, progress, lifecycle, selectedId, onSelect } = props; + const now = useNow(1000); + + const sorted = useMemo(() => { + const mains: AgentSnapshot[] = []; + const subs: AgentSnapshot[] = []; + for (const agent of agents) (agent.kind === "main" ? mains : subs).push(agent); + subs.sort((a, b) => { + const ar = a.status === "running" ? 0 : 1; + const br = b.status === "running" ? 0 : 1; + if (ar !== br) return ar - br; + return b.lastActivity - a.lastActivity; + }); + return { mains, subs }; + }, [agents]); + + return ( +
+ {sorted.mains.map(agent => ( + + ))} + {sorted.subs.map(agent => ( + + ))} + {sorted.subs.length === 0 ?
no subagents
: null} +
+ ); +} diff --git a/packages/collab-web/src/components/agents/agents.css b/packages/collab-web/src/components/agents/agents.css new file mode 100644 index 000000000..2810af1e4 --- /dev/null +++ b/packages/collab-web/src/components/agents/agents.css @@ -0,0 +1,384 @@ +/* Agents rail + drawer (T3). `ag-*` prefix; design-token vars only. */ + +/* ── Panel ────────────────────────────────────────────────────────── */ + +.ag-panel { + display: flex; + flex-direction: column; + gap: 2px; + min-height: 0; + padding: 8px; + overflow-y: auto; + font-size: 13px; +} + +.ag-row { + display: flex; + flex-direction: column; + gap: 3px; + width: 100%; + padding: 7px 9px; + border: 1px solid transparent; + border-radius: var(--radius); + background: none; + color: var(--fg); + font: inherit; + text-align: left; + cursor: pointer; + transition: background 150ms ease-out; +} + +.ag-row:hover { + background: var(--bg-raised); +} + +.ag-row--selected { + border-color: var(--border); + background: var(--bg-raised); +} + +.ag-row:focus-visible { + outline: 2px solid var(--ring); + outline-offset: 2px; +} + +.ag-row-head { + display: flex; + align-items: center; + gap: 7px; + min-width: 0; +} + +.ag-row-name { + flex: 1; + min-width: 0; + overflow: hidden; + font-weight: 500; + text-overflow: ellipsis; + white-space: nowrap; +} + +.ag-row-activity { + overflow: hidden; + padding-left: 14px; + color: var(--fg-muted); + font-family: var(--font-mono); + font-size: 11px; + text-overflow: ellipsis; + white-space: nowrap; +} + +.ag-row-meta { + display: flex; + gap: 10px; + padding-left: 14px; + color: var(--fg-faint); + font-family: var(--font-mono); + font-size: 11px; +} + +.ag-row-meta-when { + margin-left: auto; +} + +.ag-empty { + padding: 14px 10px; + color: var(--fg-faint); + font-size: 12px; + font-style: italic; + text-align: center; +} + +/* ── Status dot ───────────────────────────────────────────────────── */ + +.ag-dot { + flex: none; + width: 7px; + height: 7px; + border-radius: 50%; + background: var(--fg-faint); +} + +.ag-dot--running { + background: var(--ok); + animation: ag-pulse 1.6s ease-out infinite; +} + +.ag-dot--parked { + background: var(--warn); +} + +.ag-dot--aborted { + background: var(--err); +} + +@keyframes ag-pulse { + 0% { + box-shadow: 0 0 0 0 color-mix(in srgb, var(--ok) 45%, transparent); + } + 70%, + 100% { + box-shadow: 0 0 0 6px transparent; + } +} + +/* ── Chips ────────────────────────────────────────────────────────── */ + +.ag-chip { + flex: none; + padding: 1px 5px; + border: 1px solid var(--border); + border-radius: var(--radius-sm); + color: var(--fg-muted); + font-family: var(--font-mono); + font-size: 10px; + line-height: 1.5; + white-space: nowrap; +} + +.ag-chip--running { + border-color: color-mix(in srgb, var(--ok) 35%, transparent); + color: var(--ok); +} + +.ag-chip--parked { + border-color: color-mix(in srgb, var(--warn) 35%, transparent); + color: var(--warn); +} + +.ag-chip--aborted { + border-color: color-mix(in srgb, var(--err) 35%, transparent); + color: var(--err); +} + +.ag-chip--model { + overflow: hidden; + max-width: 160px; + color: var(--accent-muted); + text-overflow: ellipsis; +} + +/* ── Drawer ───────────────────────────────────────────────────────── */ + +.ag-drawer { + position: fixed; + top: 0; + right: 0; + bottom: 0; + z-index: 40; + display: flex; + flex-direction: column; + width: min(440px, 92vw); + border-left: 1px solid var(--border-strong); + background: var(--bg-overlay); + box-shadow: + -12px 0 32px color-mix(in srgb, var(--bg-inset) 55%, transparent), + -32px 0 80px color-mix(in srgb, var(--bg-inset) 35%, transparent); + animation: ag-drawer-in 150ms ease-out; +} + +@keyframes ag-drawer-in { + from { + transform: translateX(100%); + } + to { + transform: translateX(0); + } +} + +.ag-drawer-head { + display: flex; + align-items: center; + gap: 8px; + padding: 10px 14px; + border-bottom: 1px solid var(--border); +} + +.ag-drawer-title { + display: flex; + flex: 1; + align-items: center; + gap: 7px; + min-width: 0; +} + +.ag-drawer-name { + overflow: hidden; + font-size: 13px; + font-weight: 600; + text-overflow: ellipsis; + white-space: nowrap; +} + +.ag-drawer-actions { + display: flex; + flex: none; + align-items: center; + gap: 6px; +} + +/* ── Buttons ──────────────────────────────────────────────────────── */ + +.ag-btn { + display: inline-flex; + align-items: center; + gap: 5px; + padding: 3px 8px; + border: 1px solid var(--border); + border-radius: var(--radius); + background: var(--bg-raised); + color: var(--fg-muted); + font: inherit; + font-size: 12px; + cursor: pointer; + transition: + background 150ms ease-out, + color 150ms ease-out; +} + +.ag-btn:hover { + background: color-mix(in srgb, var(--bg-raised) 94%, var(--fg)); + color: var(--fg); +} + +.ag-btn:active { + transform: scale(0.97); +} + +.ag-btn:focus-visible { + outline: 2px solid var(--ring); + outline-offset: 2px; +} + +.ag-btn--danger { + border-color: color-mix(in srgb, var(--err) 35%, transparent); + color: var(--err); +} + +.ag-btn--danger:hover { + color: var(--err); +} + +.ag-iconbtn { + display: inline-flex; + align-items: center; + justify-content: center; + padding: 3px; + border: none; + border-radius: var(--radius-sm); + background: none; + color: var(--fg-muted); + cursor: pointer; + transition: color 150ms ease-out; +} + +.ag-iconbtn:hover { + color: var(--fg); +} + +.ag-iconbtn:active { + transform: scale(0.97); +} + +.ag-iconbtn:focus-visible { + outline: 2px solid var(--ring); + outline-offset: 2px; +} + +/* ── Stats strip ──────────────────────────────────────────────────── */ + +.ag-stats { + display: flex; + flex-wrap: wrap; + align-items: center; + gap: 12px; + padding: 8px 14px; + border-bottom: 1px solid var(--border); + color: var(--fg-muted); + font-family: var(--font-mono); + font-size: 11px; +} + +.ag-stat { + display: flex; + align-items: baseline; + gap: 5px; +} + +.ag-stat-label { + color: var(--fg-faint); +} + +.ag-stat-value { + color: var(--fg); +} + +.ag-gauge { + display: inline-block; + overflow: hidden; + align-self: center; + width: 72px; + height: 4px; + border-radius: var(--radius-sm); + background: var(--bg-inset); +} + +.ag-gauge-fill { + display: block; + height: 100%; + border-radius: inherit; + background: var(--accent); +} + +.ag-gauge-fill--warn { + background: var(--warn); +} + +/* ── Drawer body & chat ───────────────────────────────────────────── */ + +.ag-drawer-body { + display: flex; + flex: 1; + flex-direction: column; + min-height: 0; + padding: 10px 14px; + overflow-y: auto; +} + +.ag-chat { + display: flex; + flex: none; + gap: 6px; + padding: 10px 14px; + border-top: 1px solid var(--border); +} + +.ag-chat-input { + flex: 1; + min-width: 0; + padding: 6px 9px; + border: 1px solid var(--border); + border-radius: var(--radius); + background: var(--bg-inset); + color: var(--fg); + font: inherit; + font-size: 13px; +} + +.ag-chat-input::placeholder { + color: var(--fg-faint); +} + +.ag-chat-input:focus-visible { + outline: 2px solid var(--ring); + outline-offset: 2px; +} + +@media (prefers-reduced-motion: reduce) { + .ag-drawer { + animation: none; + } + .ag-dot--running { + animation: none; + } +} diff --git a/packages/collab-web/src/components/shell/Banners.tsx b/packages/collab-web/src/components/shell/Banners.tsx new file mode 100644 index 000000000..40b836a49 --- /dev/null +++ b/packages/collab-web/src/components/shell/Banners.tsx @@ -0,0 +1,47 @@ +import type { ReactNode } from "react"; +import type { ConnectionPhase } from "../../lib/client"; + +export interface BannersProps { + phase: ConnectionPhase; + endedReason: string | null; + onRejoin(): void; + onNewLink(): void; +} + +export function Banners({ phase, endedReason, onRejoin, onNewLink }: BannersProps): ReactNode { + if (phase === "connecting" || phase === "waiting") { + return ( +
+ + {phase === "connecting" ? "connecting to relay…" : "joining session…"} +
+ ); + } + if (phase === "reconnecting") { + return ( +
+ + reconnecting… +
+ ); + } + if (phase === "ended") { + return ( +
+
+
session ended
+ {endedReason &&
{endedReason}
} +
+ + +
+
+
+ ); + } + return null; +} diff --git a/packages/collab-web/src/components/shell/Composer.tsx b/packages/collab-web/src/components/shell/Composer.tsx new file mode 100644 index 000000000..d76bd33ed --- /dev/null +++ b/packages/collab-web/src/components/shell/Composer.tsx @@ -0,0 +1,88 @@ +import { SendHorizontal, Square } from "lucide-react"; +import type { KeyboardEvent, ReactNode } from "react"; +import { useCallback, useLayoutEffect, useRef, useState } from "react"; +import type { GuestClient, GuestSnapshot } from "../../lib/client"; + +export interface ComposerProps { + client: GuestClient; + snapshot: GuestSnapshot; +} + +/** Textarea metrics: line-height 20px + 8px vertical padding × 2 (kept in sync with shell.css). */ +const LINE_PX = 20; +const PAD_Y = 16; +const MAX_ROWS = 8; + +export function Composer({ client, snapshot }: ComposerProps): ReactNode { + const [text, setText] = useState(""); + const taRef = useRef(null); + + const live = snapshot.phase === "live"; + const busy = snapshot.working || (snapshot.state?.isStreaming ?? false); + const queued = snapshot.state?.queuedMessageCount ?? 0; + const canSend = live && text.trim().length > 0; + + useLayoutEffect(() => { + const el = taRef.current; + if (!el) return; + el.style.height = "0px"; + const max = MAX_ROWS * LINE_PX + PAD_Y; + el.style.height = `${Math.max(LINE_PX + PAD_Y, Math.min(el.scrollHeight, max))}px`; + el.style.overflowY = el.scrollHeight > max ? "auto" : "hidden"; + }, [text]); + + const send = useCallback((): void => { + const trimmed = text.trim(); + if (!trimmed || !live) return; + client.sendPrompt(trimmed); + setText(""); + }, [client, live, text]); + + const onKeyDown = (e: KeyboardEvent): void => { + if (e.key === "Enter" && !e.shiftKey) { + e.preventDefault(); + send(); + } + }; + + return ( +
+
+