diff --git a/.github/workflows/release-cross.yml b/.github/workflows/release-cross.yml new file mode 100644 index 0000000000..3389877be5 --- /dev/null +++ b/.github/workflows/release-cross.yml @@ -0,0 +1,138 @@ +name: Cross release (fork) + +on: + workflow_dispatch: + push: + tags: + - "v*-opencrabs.*" + +permissions: + contents: read + +jobs: + linux: + name: Linux packages + runs-on: ubuntu-22.04 + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-node@v4 + with: + node-version: 20 + cache: npm + + - uses: dtolnay/rust-toolchain@stable + + - name: Install Linux system dependencies + run: bash scripts/install-linux-deps-debian.sh + + - run: npm ci + + - name: Build Linux packages + run: | + set -euo pipefail + printf '%s\n' '{"bundle":{"createUpdaterArtifacts":false}}' > "$RUNNER_TEMP/tauri.fork.conf.json" + npm run build:linux -- --config "$RUNNER_TEMP/tauri.fork.conf.json" + + - name: Stage Linux packages + run: | + set -euo pipefail + shopt -s nullglob + mkdir -p release-artifacts + cp target/release/bundle/deb/*.deb release-artifacts/ + cp target/release/bundle/appimage/*.AppImage release-artifacts/ + + - name: Upload Linux packages + uses: actions/upload-artifact@v4 + with: + name: linux-packages + path: release-artifacts/ + if-no-files-found: error + + linux-rpm: + name: Linux rpm package + runs-on: ubuntu-latest + # Same recipe as upstream: build inside EL 10 so the rpm's glibc + # requirement matches the oldest advertised target; an Ubuntu-linked + # rpm would install on EL 10 and then fail to load. + container: almalinux:10 + steps: + - name: Install git + run: dnf install -y git + + - uses: actions/checkout@v4 + + - uses: actions/setup-node@v4 + with: + node-version: 20 + cache: npm + + - uses: dtolnay/rust-toolchain@stable + + - name: Install Enterprise Linux system dependencies + run: bash scripts/install-linux-deps-fedora.sh + + - run: npm ci + + - name: Build rpm package + run: npm run build:fedora + + - name: Verify and stage the rpm + run: | + set -euo pipefail + shopt -s nullglob + bundle=(target/release/bundle/rpm/*.rpm) + if (( ${#bundle[@]} != 1 )); then + echo "Expected exactly one .rpm" >&2 + exit 1 + fi + requires="$(rpm -qp --requires "${bundle[0]}")" + for dep in webkit2gtk4.1 gtk3 libappindicator-gtk3 librsvg2 openssl-libs; do + if ! grep -qx "$dep" <<<"$requires"; then + echo "rpm is missing Requires: $dep" >&2 + exit 1 + fi + done + mkdir -p release-artifacts + cp "${bundle[0]}" release-artifacts/ + + - name: Upload rpm package + uses: actions/upload-artifact@v4 + with: + name: linux-rpm-package + path: release-artifacts/ + if-no-files-found: error + + windows: + name: Windows packages + runs-on: windows-latest + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-node@v4 + with: + node-version: 20 + cache: npm + + - uses: dtolnay/rust-toolchain@stable + + - run: npm ci + + - name: Build Windows installer + shell: pwsh + run: | + '{"bundle":{"createUpdaterArtifacts":false}}' | Set-Content -Encoding utf8 "$env:RUNNER_TEMP\tauri.fork.conf.json" + npx tauri build --bundles nsis --config "$env:RUNNER_TEMP\tauri.fork.conf.json" + + - name: Stage Windows packages + shell: pwsh + run: | + New-Item -ItemType Directory -Force "release-artifacts" | Out-Null + Copy-Item "target/release/bundle/nsis/*.exe" "release-artifacts/" + + - name: Upload Windows packages + uses: actions/upload-artifact@v4 + with: + name: windows-packages + path: release-artifacts/ + if-no-files-found: error diff --git a/src-tauri/src/harness.rs b/src-tauri/src/harness.rs index ec717665fa..3018b2d162 100644 --- a/src-tauri/src/harness.rs +++ b/src-tauri/src/harness.rs @@ -820,6 +820,19 @@ pub fn harness_resolve_antigravity() -> Result { }) } +/// Resolve the OpenCrabs CLI (`opencrabs`). +#[tauri::command(async)] +pub fn harness_resolve_opencrabs() -> Result { + resolve_opencrabs() + .map(|path| CursorBinary { + path: path.to_string_lossy().into_owned(), + }) + .ok_or_else(|| { + "OpenCrabs CLI not found. Install OpenCrabs from https://github.com/adolfousier/opencrabs, then retry." + .into() + }) +} + /// Bind an ephemeral loopback port for `opencode serve`. #[tauri::command] pub fn harness_free_port() -> Result { @@ -1994,6 +2007,7 @@ fn resolve_harness_binary_default(provider: &str) -> Option { "hermes" => resolve_hermes(), "antigravity" => resolve_antigravity(), "devin" => resolve_devin(), + "opencrabs" => resolve_opencrabs(), _ => None, } } @@ -2451,6 +2465,35 @@ fn resolve_antigravity() -> Option { first_binary(candidates) } +fn resolve_opencrabs() -> Option { + let home = dirs_home().map(PathBuf::from); + let mut candidates: Vec = Vec::new(); + + if let Some(home) = &home { + candidates.push(home.join(".opencrabs/bin/opencrabs")); + candidates.push(home.join(".local/bin/opencrabs")); + candidates.push(home.join(".cargo/bin/opencrabs")); + candidates.push(home.join(".npm-global/bin/opencrabs")); + candidates.push(home.join("n/bin/opencrabs")); + } + #[cfg(target_os = "macos")] + candidates.push(PathBuf::from("/opt/homebrew/bin/opencrabs")); + candidates.push(PathBuf::from("/usr/local/bin/opencrabs")); + candidates.push(PathBuf::from("/usr/bin/opencrabs")); + candidates.push(PathBuf::from("/snap/bin/opencrabs")); + if let Some(from_shell) = which_via_login_shell("opencrabs") { + candidates.push(from_shell); + } + + first_binary_matching(candidates, is_opencrabs_binary) +} + +/// No same-named collision is known for `opencrabs`; an executable with the +/// right name is the agent CLI. +fn is_opencrabs_binary(path: &Path) -> bool { + binary_name_eq(path, "opencrabs") +} + fn is_pi_coding_agent(path: &Path) -> bool { if !path.is_file() { return false; diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 829dc13f5e..8564443d91 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -457,6 +457,7 @@ pub fn run() { harness::harness_resolve_hermes, harness::harness_resolve_devin, harness::harness_resolve_antigravity, + harness::harness_resolve_opencrabs, harness::harness_free_port, harness::harness_spawn, codex_mono_store::codex_mono_store_prepare, diff --git a/src/assets/providers/opencrabs.svg b/src/assets/providers/opencrabs.svg new file mode 100644 index 0000000000..d4177a568a --- /dev/null +++ b/src/assets/providers/opencrabs.svg @@ -0,0 +1,20 @@ + + + + + + + + + + + + + + + + + + + diff --git a/src/features/sessions/model/attachments.ts b/src/features/sessions/model/attachments.ts index 007221c9b2..799492688b 100644 --- a/src/features/sessions/model/attachments.ts +++ b/src/features/sessions/model/attachments.ts @@ -522,7 +522,7 @@ function fallbackName(mimeType: string): string { return "attachment"; } -function fileUri(path: string): string { +export function fileUri(path: string): string { const normalized = path.replace(/\\/g, "/"); const abs = normalized.startsWith("/") ? normalized : `/${normalized}`; return `file://${abs.split("/").map(encodeURIComponent).join("/")}`; diff --git a/src/features/sessions/model/models.ts b/src/features/sessions/model/models.ts index 3f5469b76c..27b6137880 100644 --- a/src/features/sessions/model/models.ts +++ b/src/features/sessions/model/models.ts @@ -204,6 +204,12 @@ export const MODELS: AgentModel[] = [ name: "Adaptive", nativeId: "adaptive", }, + { + id: "opencrabs:default", + harness: "opencrabs", + name: "Configured model", + nativeId: "", + }, ]; export const DEFAULT_MODEL_ID: Record = { @@ -218,6 +224,7 @@ export const DEFAULT_MODEL_ID: Record = { hermes: "hermes:default", antigravity: "antigravity:gemini-3.8-flash-high", devin: "devin:adaptive", + opencrabs: "opencrabs:default", }; const FAVORITES_KEY = "monocode.favoriteModels"; @@ -248,6 +255,7 @@ const HARNESS_ORDER: HarnessId[] = [ "hermes", "antigravity", "devin", + "opencrabs", ]; const EMPTY_MODELS: AgentModel[] = []; diff --git a/src/features/sessions/model/session.ts b/src/features/sessions/model/session.ts index c707542468..7b0b3efb07 100644 --- a/src/features/sessions/model/session.ts +++ b/src/features/sessions/model/session.ts @@ -26,7 +26,8 @@ export type HarnessId = | "fx" | "hermes" | "antigravity" - | "devin"; + | "devin" + | "opencrabs"; export const HARNESSES: HarnessId[] = [ "claude", @@ -40,6 +41,7 @@ export const HARNESSES: HarnessId[] = [ "hermes", "antigravity", "devin", + "opencrabs", ]; export type BlockRole = @@ -539,6 +541,7 @@ export const HARNESS_LABEL: Record = { hermes: "hermes", antigravity: "antigravity", devin: "devin", + opencrabs: "opencrabs", }; export const HARNESS_TITLE: Record = { @@ -553,6 +556,7 @@ export const HARNESS_TITLE: Record = { hermes: "Hermes Agent", antigravity: "Antigravity", devin: "Devin", + opencrabs: "OpenCrabs", }; /** fx ACP rejects attachment prompt blocks. */ diff --git a/src/features/sessions/ui/HarnessIcon.tsx b/src/features/sessions/ui/HarnessIcon.tsx index 560ead8b2f..c715937aee 100644 --- a/src/features/sessions/ui/HarnessIcon.tsx +++ b/src/features/sessions/ui/HarnessIcon.tsx @@ -10,6 +10,7 @@ import omp from "../../../assets/providers/omp.svg"; import opencode from "../../../assets/providers/opencode.svg"; import pi from "../../../assets/providers/pi.svg"; import antigravity from "../../../assets/providers/antigravity.svg"; +import opencrabs from "../../../assets/providers/opencrabs.svg"; import type { HarnessId } from "../model/session"; export const HARNESS_ICONS: Record = { @@ -24,6 +25,7 @@ export const HARNESS_ICONS: Record = { hermes, antigravity, devin, + opencrabs, }; /** White marks that must follow `currentColor` so they stay visible in light mode. */ diff --git a/src/integrations/harness/core/availability.ts b/src/integrations/harness/core/availability.ts index 1ed8f5b7c1..882a9af2af 100644 --- a/src/integrations/harness/core/availability.ts +++ b/src/integrations/harness/core/availability.ts @@ -11,6 +11,7 @@ import { resolveHermesBinary, resolveOmpBinary, resolveOpenCodeBinary, + resolveOpenCrabsBinary, resolvePiBinary, } from "./child"; import { isLiveHarness } from "./registry"; @@ -59,6 +60,10 @@ const CLI: Record = { ? "irm https://static.devin.ai/cli/setup.ps1 | iex" : "curl -fsSL https://cli.devin.ai/install.sh | bash", }, + opencrabs: { + name: "OpenCrabs CLI", + install: "cargo install --git https://github.com/adolfousier/opencrabs", + }, }; let inflight: Promise | null = null; @@ -176,6 +181,14 @@ export function probeHarnessAvailability( return [id, false] as const; } } + if (id === "opencrabs") { + try { + await resolveOpenCrabsBinary(); + return [id, true] as const; + } catch { + return [id, false] as const; + } + } return [id, false] as const; }), ) diff --git a/src/integrations/harness/core/availabilityState.ts b/src/integrations/harness/core/availabilityState.ts index 732d04bc40..ee60ac667a 100644 --- a/src/integrations/harness/core/availabilityState.ts +++ b/src/integrations/harness/core/availabilityState.ts @@ -19,6 +19,7 @@ let availability: HarnessAvailability = { hermes: false, antigravity: false, devin: false, + opencrabs: false, }; let version = 0; let probedAt = 0; diff --git a/src/integrations/harness/core/child.ts b/src/integrations/harness/core/child.ts index d0239a1625..da611bc13e 100644 --- a/src/integrations/harness/core/child.ts +++ b/src/integrations/harness/core/child.ts @@ -414,6 +414,7 @@ async function resolveHarnessBinary( hermes: "harness_resolve_hermes", antigravity: "harness_resolve_antigravity", devin: "harness_resolve_devin", + opencrabs: "harness_resolve_opencrabs", }; return invoke(command[provider]); } @@ -487,6 +488,10 @@ export function resolveAntigravityBinary( }>; } +export function resolveOpenCrabsBinary(): Promise<{ path: string }> { + return invoke("harness_resolve_opencrabs"); +} + export function freeHarnessPort(): Promise { return invoke("harness_free_port"); } diff --git a/src/integrations/harness/core/register.ts b/src/integrations/harness/core/register.ts index eaca92d4bd..1eccece993 100644 --- a/src/integrations/harness/core/register.ts +++ b/src/integrations/harness/core/register.ts @@ -9,6 +9,7 @@ import { ensureOmpRegistered } from "../providers/omp/ompAdapter"; import { ensurePiRegistered } from "../providers/pi/piAdapter"; import { ensureAntigravityRegistered } from "../providers/antigravity/antigravityAdapter"; import { ensureDevinRegistered } from "../providers/devin/devinAdapter"; +import { ensureOpenCrabsRegistered } from "../providers/opencrabs/opencrabsAdapter"; /** Register all known live harness adapters. Idempotent. */ export function registerBuiltinHarnesses(): void { @@ -23,4 +24,5 @@ export function registerBuiltinHarnesses(): void { ensureHermesRegistered(); ensureAntigravityRegistered(); ensureDevinRegistered(); + ensureOpenCrabsRegistered(); } diff --git a/src/integrations/harness/core/textHarness.ts b/src/integrations/harness/core/textHarness.ts index 2da606ed15..79054dede5 100644 --- a/src/integrations/harness/core/textHarness.ts +++ b/src/integrations/harness/core/textHarness.ts @@ -13,6 +13,7 @@ const TEXT_HARNESSES: HarnessId[] = [ "codex", "grok", "opencode", + "opencrabs", ]; /** Pick the harness used for titles, commit messages, and PR text. */ diff --git a/src/integrations/harness/index.ts b/src/integrations/harness/index.ts index bc8b1e5c76..1b5ab22aa8 100644 --- a/src/integrations/harness/index.ts +++ b/src/integrations/harness/index.ts @@ -108,6 +108,16 @@ export { respondAntigravityApproval, bindAntigravitySession, } from "./providers/antigravity/antigravity"; +export { + sendOpenCrabsTurn, + steerOpenCrabsTurn, + cancelOpenCrabsTurn, + respondOpenCrabsApproval, + stopOpenCrabsSession, + forgetOpenCrabsSession, + bindOpenCrabsSession, + compactOpenCrabsContext, +} from "./providers/opencrabs/opencrabs"; export { generateCursorSessionTitle } from "./providers/cursor/cursorTitle"; export { generateCodexSessionTitle } from "./providers/codex/codexTitle"; export { generateOpenCodeSessionTitle } from "./providers/opencode/opencodeTitle"; @@ -117,6 +127,7 @@ export { generateOmpSessionTitle, } from "./providers/pi/piTitle"; export { generateGrokSessionTitle } from "./providers/grok/grokTitle"; +export { generateOpenCrabsSessionTitle } from "./providers/opencrabs/opencrabsTitle"; export { generateCursorCommitMessage, generateCursorPrContent, @@ -138,6 +149,10 @@ export { generateGrokCommitMessage, generateGrokPrContent, } from "./providers/grok/grokGit"; +export { + generateOpenCrabsCommitMessage, + generateOpenCrabsPrContent, +} from "./providers/opencrabs/opencrabsGit"; export { generateCommitMessage, generatePrContent, diff --git a/src/integrations/harness/providers/opencrabs/opencrabs.ts b/src/integrations/harness/providers/opencrabs/opencrabs.ts new file mode 100644 index 0000000000..73264df634 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabs.ts @@ -0,0 +1,244 @@ +import { nativeModelId } from "../../../../features/sessions/model/models"; +import { openCrabsPromptBlocks } from "./opencrabsPrompt"; +import type { + ApprovalDecision, + CompactContextInput, + SendTurnInput, + SteerTurnInput, +} from "../../core/types"; +import { + cancelledThreads, + CONTROL_TIMEOUT_MS, + ensureLive, + liveByThread, + type Live, + PROMPT_TIMEOUT_MS, + SERVER_HELP, + stopOpenCrabsSession, +} from "./opencrabsLive"; + +// Lifecycle (spawn, handshake, resume, teardown) lives in opencrabsLive; +// re-exported here so existing import sites keep working. +export { + bindOpenCrabsSession, + forgetOpenCrabsSession, + spawnArgs, + stopOpenCrabsSession, +} from "./opencrabsLive"; + +/** + * Live OpenCrabs adapter turn orchestration. Spawns `opencrabs acp` (via + * opencrabsLive) and talks Agent Client Protocol over stdio. Permission + * requests surface in the UI unless the runtime mode auto-answers them. + */ +export async function sendOpenCrabsTurn(input: SendTurnInput): Promise { + let live: Live; + try { + live = await ensureLive(input); + } catch (error) { + cancelledThreads.delete(input.sessionId); + throw error; + } + if (cancelledThreads.delete(input.sessionId)) return; + + live.onEvent = input.onEvent; + live.runtimeMode = input.runtimeMode; + live.planning = input.intent === "plan"; + live.turns = live.turns + .catch(() => undefined) + .then(async () => { + live.cancelled = false; + live.muteUpdates = false; + try { + await applyModelSelection(live, input); + await applyRuntimeMode(live, input); + if (live.cancelled) return; + await prompt(live, input); + } catch (error) { + if (live.cancelled) return; + throw error; + } + }); + try { + await live.turns; + } catch (error) { + // A timed-out or failed turn leaves the child's protocol state unknowable. + // Keep the provider session id, but recycle the process so the next turn + // resumes on a fresh transport instead of a wedged one. + if (liveByThread.get(input.sessionId) === live) { + await stopOpenCrabsSession(input.sessionId); + } + throw error; + } +} + +export async function steerOpenCrabsTurn(input: SteerTurnInput): Promise { + const live = liveByThread.get(input.sessionId); + if (!live) throw new Error("No active OpenCrabs session"); + const blocks = await openCrabsPromptBlocks(input.text, input.attachments); + if (blocks.length === 0) return; + await live.acp + // Plain method name, not the "_session/steer" ext-prefix: released + // opencrabs (v0.5.2) only registers the plain name, and JSON-RPC drops + // unknown notifications silently — an ext-prefixed notify is a no-op on + // every released binary. The parity server accepts both spellings. + .notify("session/steer", { + sessionId: live.acpSessionId, + prompt: blocks, + }) + .catch(() => undefined); +} + +export function respondOpenCrabsApproval( + sessionId: string, + requestId: number, + decision: ApprovalDecision, +): void { + liveByThread.get(sessionId)?.approvals.get(requestId)?.(decision); +} + +/** Abort the in-flight prompt without tearing down the ACP session. */ +export async function cancelOpenCrabsTurn(sessionId: string): Promise { + const live = liveByThread.get(sessionId); + if (!live) { + cancelledThreads.add(sessionId); + return; + } + live.cancelled = true; + live.muteUpdates = true; + for (const [, resolve] of live.approvals) resolve("deny"); + live.approvals.clear(); + await live.acp + .notify("session/cancel", { sessionId: live.acpSessionId }) + .catch(() => undefined); + live.acp.rejectPending(new Error("cancelled")); +} + +/** + * Compact the session's context window via `session/compact`. The server + * runs its native summarization turn; the request resolves when it ends. + * Compaction of a long conversation is a full turn, so it rides the turn + * queue and the prompt timeout rather than the control timeout. + */ +export async function compactOpenCrabsContext( + input: CompactContextInput, +): Promise { + const live = liveByThread.get(input.sessionId) ?? + (await ensureLive({ ...input, text: "" })); + live.onEvent = input.onEvent; + live.turns = live.turns + .catch(() => undefined) + .then(async () => { + if (live.cancelled) return; + await live.acp.request( + "session/compact", + { sessionId: live.acpSessionId }, + PROMPT_TIMEOUT_MS, + ); + }); + try { + await live.turns; + } catch (error) { + // Same contract as a failed prompt turn: a timed-out compact leaves the + // child's protocol state unknowable. Keep the ACP session id, recycle + // the process so the next turn resumes on a fresh transport instead of + // a wedged one. + if (liveByThread.get(input.sessionId) === live) { + await stopOpenCrabsSession(input.sessionId); + } + throw error; + } +} + +/** + * Model selection is best-effort: the static catalog ships only `default` + * (empty native id, skipped here), while live catalog entries carry + * `provider/model` pairs the server routes through `session/set_model`. + */ +async function applyModelSelection( + live: Live, + input: SendTurnInput, +): Promise { + const base = nativeModelId(input.model).trim(); + if (!base) return; + try { + await live.acp.request( + "session/set_model", + { sessionId: live.acpSessionId, modelId: base }, + CONTROL_TIMEOUT_MS, + ); + // The badge only hears about model switches through configChanged — + // without it the picker and the turn can quietly disagree. + live.onEvent({ type: "session.configChanged", model: input.model }); + } catch (error) { + // The switch failed but the turn proceeds on the previous model — say + // so instead of letting the picker silently disagree with the turn. + live.onEvent({ + type: "session.error", + message: `Model switch to ${base} failed — continuing with the previous model (${String(error)})`, + }); + } +} + +/** + * Push the runtime/plan mode server-side so the approval policy lives where + * the tools run. Client-side gating in handlePermission stays as backstop, + * and an older binary without set_mode support degrades to it. + */ +async function applyRuntimeMode( + live: Live, + input: SendTurnInput, +): Promise { + const modeId = input.intent === "plan" ? "plan" : input.runtimeMode; + await live.acp + .request( + "session/set_mode", + { sessionId: live.acpSessionId, modeId }, + CONTROL_TIMEOUT_MS, + ) + .catch((error: unknown) => { + const detail = error instanceof Error ? error.message : String(error); + console.debug("[monocode] opencrabs set_mode failed", detail); + if (/timed out|not running|exited|closed|pipe/i.test(detail)) throw error; + // Old binaries without set_mode answer "method not found" — that is + // the designed degradation to client-side gating, not a failure worth + // a transcript line. Anything else the user should hear: the next + // turn then runs a different approval policy than the chip promises. + if (!/method not found/i.test(detail)) { + live.onEvent({ + type: "session.error", + message: `Mode switch to ${modeId} failed — the next turn keeps the previous tool-approval policy (${detail})`, + }); + } + }); +} + +async function prompt(live: Live, input: SendTurnInput): Promise { + try { + const blocks = await openCrabsPromptBlocks(input.text, input.attachments); + if (blocks.length === 0) return; + await live.acp.request( + "session/prompt", + { + sessionId: live.acpSessionId, + prompt: blocks, + }, + PROMPT_TIMEOUT_MS, + ); + if (live.cancelled) return; + live.onEvent({ type: "message.completed" }); + live.onEvent({ type: "reasoning.completed" }); + } catch (error) { + if (live.cancelled) return; + const detail = error instanceof Error ? error.message : String(error); + live.onEvent({ + type: "session.error", + message: /timed out|not running|exited|closed|pipe|method not found/i.test( + detail, + ) + ? `${detail.trim()}\n\n${SERVER_HELP}` + : detail, + }); + throw error; + } +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsAdapter.ts b/src/integrations/harness/providers/opencrabs/opencrabsAdapter.ts new file mode 100644 index 0000000000..58749843e2 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsAdapter.ts @@ -0,0 +1,44 @@ +import { + bindOpenCrabsSession, + cancelOpenCrabsTurn, + compactOpenCrabsContext, + forgetOpenCrabsSession, + respondOpenCrabsApproval, + sendOpenCrabsTurn, + steerOpenCrabsTurn, + stopOpenCrabsSession, +} from "./opencrabs"; +import { openCrabsCommands } from "./opencrabsCommands"; +import { + generateOpenCrabsBranchName, + generateOpenCrabsCommitMessage, + generateOpenCrabsPrContent, +} from "./opencrabsGit"; +import { generateOpenCrabsSessionTitle } from "./opencrabsTitle"; +import { registerHarness, type HarnessAdapter } from "../../core/registry"; + +export const openCrabsAdapter: HarnessAdapter = { + id: "opencrabs", + live: true, + sendTurn: sendOpenCrabsTurn, + steerTurn: steerOpenCrabsTurn, + cancelTurn: cancelOpenCrabsTurn, + respondApproval: respondOpenCrabsApproval, + stopSession: stopOpenCrabsSession, + forgetSession: forgetOpenCrabsSession, + bindSession: bindOpenCrabsSession, + compactContext: compactOpenCrabsContext, + commands: openCrabsCommands, + generateTitle: generateOpenCrabsSessionTitle, + generateCommitMessage: generateOpenCrabsCommitMessage, + generatePrContent: generateOpenCrabsPrContent, + generateBranchName: generateOpenCrabsBranchName, +}; + +let registered = false; + +export function ensureOpenCrabsRegistered(): void { + if (registered) return; + registerHarness(openCrabsAdapter); + registered = true; +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsApproval.test.ts b/src/integrations/harness/providers/opencrabs/opencrabsApproval.test.ts new file mode 100644 index 0000000000..6948c613fe --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsApproval.test.ts @@ -0,0 +1,90 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { handlePermission, type ApprovalTarget } from "./opencrabsApproval"; +import type { AcpClient } from "../../core/acp"; +import type { HarnessEvent } from "../../core/types"; +import type { RuntimeMode } from "../../../../features/sessions/model/session"; + +const permissionParams = { + sessionId: "s-1", + toolCall: { + toolCallId: "call-1", + title: "bash", + kind: "execute", + rawInput: { command: "ls" }, + }, + options: [{ optionId: "allow-once" }, { optionId: "reject-once" }], +}; + +function makeTarget(): { + target: ApprovalTarget; + respond: ReturnType; + events: HarnessEvent[]; +} { + const respond = vi.fn().mockResolvedValue(undefined); + const events: HarnessEvent[] = []; + const target: ApprovalTarget = { + acp: { respond } as unknown as AcpClient, + planning: false, + runtimeMode: "supervised" as RuntimeMode, + onEvent: (event) => events.push(event), + approvals: new Map(), + }; + return { target, respond, events }; +} + +describe("handlePermission dialog expiry", () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + afterEach(() => { + vi.useRealTimers(); + }); + + it("expires the dialog to deny just under the server's 300s deadline", async () => { + const { target, respond, events } = makeTarget(); + const pending = handlePermission(target, 7, permissionParams); + await vi.advanceTimersByTimeAsync(0); + expect(events.some((e) => e.type === "approval.requested")).toBe(true); + + await vi.advanceTimersByTimeAsync(290_000); + await pending; + + const resolved = events[events.length - 1]; + expect(resolved).toMatchObject({ + type: "approval.resolved", + requestId: 7, + decision: "deny", + }); + expect(respond).toHaveBeenCalledTimes(1); + expect(respond).toHaveBeenCalledWith(7, { + outcome: { outcome: "selected", optionId: "reject-once" }, + }); + expect(target.approvals.size).toBe(0); + + // The settled timer is cleared: no second resolution long after expiry. + await vi.advanceTimersByTimeAsync(600_000); + expect(respond).toHaveBeenCalledTimes(1); + expect(events.filter((e) => e.type === "approval.resolved")).toHaveLength(1); + }); + + it("a user answer before expiry cancels the timer", async () => { + const { target, respond, events } = makeTarget(); + const pending = handlePermission(target, 8, permissionParams); + await vi.advanceTimersByTimeAsync(0); + + target.approvals.get(8)?.("allow"); + await pending; + + expect(respond).toHaveBeenCalledTimes(1); + expect(respond).toHaveBeenCalledWith(8, { + outcome: { outcome: "selected", optionId: "allow-once" }, + }); + + // Long past the deadline the cleared timer never fires a stale deny. + await vi.advanceTimersByTimeAsync(600_000); + expect(respond).toHaveBeenCalledTimes(1); + const resolved = events.filter((e) => e.type === "approval.resolved"); + expect(resolved).toHaveLength(1); + expect(resolved[0]).toMatchObject({ decision: "allow" }); + }); +}); diff --git a/src/integrations/harness/providers/opencrabs/opencrabsApproval.ts b/src/integrations/harness/providers/opencrabs/opencrabsApproval.ts new file mode 100644 index 0000000000..9df7cc80f2 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsApproval.ts @@ -0,0 +1,97 @@ +import type { ApprovalDecision, HarnessEvent } from "../../core/types"; +import type { AcpClient } from "../../core/acp"; +import type { RuntimeMode } from "../../../../features/sessions/model/session"; +import { + autoPermissionOption, + permissionOptionId, + permissionRequestFromAcp, +} from "./opencrabsProtocol"; + +export type ApprovalTarget = { + acp: AcpClient; + planning: boolean; + runtimeMode: RuntimeMode; + onEvent: (event: HarnessEvent) => void; + approvals: Map void>; +}; + +/** The server abandons an unanswered permission after 300s and proceeds as + * denied (PERMISSION_TIMEOUT in its turn bridge). The dialog must not + * outlive that deadline — answering after the server moved on sends a + * response nobody waits for, and an eternally open dialog lies about the + * session's state. Expire just under the server deadline, resolving deny, + * so UI and server stay in agreement. */ +const APPROVAL_EXPIRY_MS = 290_000; + +export async function handlePermission( + target: ApprovalTarget, + id: number, + params: unknown, +): Promise { + const request = permissionRequestFromAcp(params); + if (request.callId) { + target.onEvent({ + type: "tool.updated", + callId: request.callId, + title: request.title, + kind: request.kind, + preview: request.preview, + }); + } + + if (target.planning) { + const normalized = (request.kind ?? "").toLowerCase(); + const readOnly = normalized === "read" || normalized === "search"; + await target.acp.respond(id, { + outcome: { + outcome: "selected", + optionId: permissionOptionId(readOnly ? "allow" : "deny", request.optionIds), + }, + }); + return; + } + + const auto = autoPermissionOption( + target.runtimeMode, + request.kind, + request.optionIds, + ); + if (auto) { + await target.acp.respond(id, { + outcome: { outcome: "selected", optionId: auto }, + }); + return; + } + + target.onEvent({ + type: "approval.requested", + requestId: id, + title: request.title, + kind: request.kind, + callId: request.callId, + preview: request.preview, + }); + + const decision = await new Promise((resolve) => { + // Settle exactly once: user answer, cancel sweep, or expiry — whichever + // comes first clears the timer and the map slot, so a late timer or a + // late click can never double-resolve. + const settle = (outcome: ApprovalDecision) => { + clearTimeout(timer); + target.approvals.delete(id); + resolve(outcome); + }; + const timer = setTimeout(() => { + settle("deny"); + }, APPROVAL_EXPIRY_MS); + target.approvals.set(id, settle); + }); + target.onEvent({ type: "approval.resolved", requestId: id, decision }); + + await target.acp.respond(id, { + outcome: { + outcome: "selected", + optionId: permissionOptionId(decision, request.optionIds), + }, + }); +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsCommands.ts b/src/integrations/harness/providers/opencrabs/opencrabsCommands.ts new file mode 100644 index 0000000000..a310965226 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsCommands.ts @@ -0,0 +1,50 @@ +import { + nativeCommandInvocation, + type NativeCommand, + type NativeCommandProvider, +} from "../../core/nativeCommands"; + +/** Slash commands pushed by the server, cached per MonoCode thread. */ +const commandsByThread = new Map(); +const commandSubscribers = new Map void>>(); + +export function cacheNativeCommands( + threadId: string, + rows: { name: string; description: string }[], +): void { + const commands: NativeCommand[] = rows.map((row) => ({ + name: row.name, + description: row.description, + invocation: nativeCommandInvocation("opencrabs", row.name), + source: "opencrabs" as const, + })); + commandsByThread.set(threadId, commands); + commandSubscribers.get(threadId)?.forEach((cb) => cb(commands)); +} + +export function clearNativeCommands(threadId: string): void { + commandsByThread.delete(threadId); + commandSubscribers.delete(threadId); +} + +/** + * The server's `available_commands_update` push, surfaced as a command + * provider: built-ins, skills, and the user's own commands.toml entries are + * slash-able from the picker once a session is live. + */ +export const openCrabsCommands: NativeCommandProvider = { + discover: async (context) => + commandsByThread.get(context.sessionId ?? "") ?? [], + subscribe: (context, onCommands) => { + const key = context.sessionId ?? ""; + const set = commandSubscribers.get(key) ?? new Set(); + set.add(onCommands); + commandSubscribers.set(key, set); + const cached = commandsByThread.get(key); + if (cached) onCommands(cached); + return () => { + set.delete(onCommands); + }; + }, + rawSlashCommands: true, +}; diff --git a/src/integrations/harness/providers/opencrabs/opencrabsCompact.test.ts b/src/integrations/harness/providers/opencrabs/opencrabsCompact.test.ts new file mode 100644 index 0000000000..8d95a4d8f1 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsCompact.test.ts @@ -0,0 +1,98 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { Attachment } from "../../../../features/sessions/model/session"; + +const acpState = vi.hoisted(() => ({ failCompact: false })); + +vi.mock("@tauri-apps/api/core", () => ({ + invoke: (command: string, args: unknown) => invokeMock(command, args), +})); + +vi.mock("../../core/child", () => ({ + resolveOpenCrabsBinary: vi + .fn() + .mockResolvedValue({ path: "/fake/opencrabs" }), + spawnChild: vi.fn().mockResolvedValue(undefined), + killChild: vi.fn().mockResolvedValue(undefined), + watchChild: vi.fn(), + unwatchChild: vi.fn(), +})); + +vi.mock("../../core/acp", () => { + class FakeAcp { + constructor(_sessionId: string, _handlers: Record) {} + request(method: string) { + if (method === "session/compact" && acpState.failCompact) { + return Promise.reject(new Error("Request timed out")); + } + switch (method) { + case "initialize": + return Promise.resolve({ protocolVersion: 1 }); + case "session/new": + return Promise.resolve({ + sessionId: "acp-1", + models: { currentModelId: null, availableModels: [] }, + }); + default: + return Promise.resolve({}); + } + } + notify() { + return Promise.resolve(); + } + respond() { + return Promise.resolve(); + } + respondError() { + return Promise.resolve(); + } + rejectPending() {} + pushLine() {} + close() {} + } + return { AcpClient: FakeAcp }; +}); + +const invokeMock = vi.fn(); + +import { compactOpenCrabsContext } from "./opencrabs"; +import { killChild, spawnChild } from "../../core/child"; + +function compactInput() { + return { + sessionId: "t-compact", + cwd: "/workspace", + model: "opencrabs:", + runtimeMode: "supervised", + intent: "ask", + onEvent: () => undefined, + attachments: [] as Attachment[], + }; +} + +describe("compactOpenCrabsContext", () => { + beforeEach(() => { + acpState.failCompact = false; + vi.clearAllMocks(); + }); + + it("recycles the transport when compact fails, so the next call respawns", async () => { + // 1. First compact on a cold session: spawns, session/new, compact OK. + await expect(compactOpenCrabsContext(compactInput())).resolves.toBeUndefined(); + + // 2. Same session + cwd reuses the live transport; its compact times + // out — the failure must surface, not swallow. + acpState.failCompact = true; + await expect(compactOpenCrabsContext(compactInput())).rejects.toThrow( + "Request timed out", + ); + + // 3. The wedged transport was recycled: child killed, and the next + // compact spawns a fresh child instead of reusing the dead one. + expect(killChild).toHaveBeenCalledWith("t-compact"); + const spawnsBefore = vi.mocked(spawnChild).mock.calls.length; + + acpState.failCompact = false; + await expect(compactOpenCrabsContext(compactInput())).resolves.toBeUndefined(); + expect(vi.mocked(spawnChild).mock.calls.length).toBe(spawnsBefore + 1); + }); +}); diff --git a/src/integrations/harness/providers/opencrabs/opencrabsGit.ts b/src/integrations/harness/providers/opencrabs/opencrabsGit.ts new file mode 100644 index 0000000000..634a7bdf73 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsGit.ts @@ -0,0 +1,87 @@ +import { gitRangeContext, gitStagedContext } from "../../../../platform/tauri/fs"; +import { + buildBranchNamePrompt, + buildCommitMessagePrompt, + buildPrContentPrompt, + formatCommitMessage, + parseBranchName, + parseCommitMessage, + parsePrContent, + type PrContent, +} from "../../../../features/source-control/model/gitText"; +import { runOpenCrabsTextPrompt } from "./opencrabsText"; + +const GIT_TIMEOUT_MS = 90_000; + +export async function generateOpenCrabsCommitMessage( + cwd: string, +): Promise { + const context = await gitStagedContext(cwd); + const output = await runOpenCrabsTextPrompt({ + cwd, + prompt: buildCommitMessagePrompt({ + branch: context.branch, + stagedSummary: context.summary, + stagedPatch: context.patch, + }), + timeoutMs: GIT_TIMEOUT_MS, + }); + const parsed = parseCommitMessage(output); + if (parsed) return formatCommitMessage(parsed); + const snippet = output.trim().replace(/\s+/g, " ").slice(0, 240); + throw new Error( + snippet + ? `Could not generate a commit message. Model replied: ${snippet}` + : "Could not generate a commit message. OpenCrabs returned no text.", + ); +} + +export async function generateOpenCrabsPrContent( + cwd: string, +): Promise<(PrContent & { base: string; head: string }) | null> { + const range = await gitRangeContext(cwd); + let parsed: PrContent | null = null; + try { + const output = await runOpenCrabsTextPrompt({ + cwd, + prompt: buildPrContentPrompt({ + baseBranch: range.base, + headBranch: range.head, + commitSummary: range.commitSummary, + diffSummary: range.diffSummary, + diffPatch: range.diffPatch, + }), + timeoutMs: GIT_TIMEOUT_MS, + }); + parsed = parsePrContent(output); + } catch (error) { + console.debug("[monocode] opencrabs pr content", error); + } + const title = + parsed?.title || + range.commitSummary.split(/\r?\n/)[0]?.trim() || + `Update ${range.head}`; + return { + title, + body: parsed?.body || range.commitSummary.trim(), + base: range.base, + head: range.head, + }; +} + +export async function generateOpenCrabsBranchName( + cwd: string, + message: string, +): Promise { + try { + const output = await runOpenCrabsTextPrompt({ + cwd, + prompt: buildBranchNamePrompt(message), + timeoutMs: GIT_TIMEOUT_MS, + }); + return parseBranchName(output); + } catch (error) { + console.debug("[monocode] opencrabs branch name", error); + return null; + } +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsLive.ts b/src/integrations/harness/providers/opencrabs/opencrabsLive.ts new file mode 100644 index 0000000000..c6d3c1d04c --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsLive.ts @@ -0,0 +1,364 @@ +import { + nativeModelId, + setHarnessModels, +} from "../../../../features/sessions/model/models"; +import type { RuntimeMode } from "../../../../features/sessions/model/session"; +import { AcpClient, type AcpHandlers } from "../../core/acp"; +import { + killChild, + resolveOpenCrabsBinary, + spawnChild, + unwatchChild, + watchChild, +} from "../../core/child"; +import { + eventsFromAcpUpdate, + modelsFromSessionNew, + nativeCommandsFromUpdate, + sessionIdFromResult, +} from "./opencrabsProtocol"; +import { AcpSubagents } from "../../core/acpSubagents"; +import { handlePermission } from "./opencrabsApproval"; +import { + cacheNativeCommands, + clearNativeCommands, +} from "./opencrabsCommands"; +import type { + ApprovalDecision, + HarnessEvent, + SendTurnInput, +} from "../../core/types"; + +/** + * Session lifecycle for the OpenCrabs adapter: process spawn, ACP + * handshake, resume, and teardown. Turn orchestration lives in + * `opencrabs.ts`; this module owns the machinery that keeps a live child + * and its ACP session id alive across turns. + */ + +export type Live = { + acp: AcpClient; + acpSessionId: string; + threadId: string; + cwd: string; + muteUpdates: boolean; + cancelled: boolean; + runtimeMode: RuntimeMode; + planning: boolean; + subagents: AcpSubagents; + onEvent: (event: HarnessEvent) => void; + approvals: Map void>; + turns: Promise; +}; + +export type Resume = { + acpSessionId: string; + cwd: string; +}; + +// OpenCrabs boots a full runtime (config, brain files, provider handshake) +// before it can answer `initialize`, so give it more room than a thin CLI. +export const INIT_TIMEOUT_MS = 30_000; +export const SESSION_TIMEOUT_MS = 45_000; +export const CONTROL_TIMEOUT_MS = 15_000; +export const PROMPT_TIMEOUT_MS = 30 * 60_000; + +export const SERVER_HELP = + "The ACP server mode ships in opencrabs v0.5.2 and later. " + + "Check `opencrabs --version`, then upgrade (or point the resolver at a " + + "newer binary) and retry."; + +const CLIENT_CAPABILITIES = { + fs: { readTextFile: false, writeTextFile: false }, + terminal: false, +}; + +export const liveByThread = new Map(); +export const resumeByThread = new Map(); +export const cancelledThreads = new Set(); + +export async function ensureLive(input: SendTurnInput): Promise { + const existing = liveByThread.get(input.sessionId); + if (existing && existing.cwd === input.cwd) { + existing.onEvent = input.onEvent; + existing.runtimeMode = input.runtimeMode; + existing.planning = input.intent === "plan"; + return existing; + } + if (existing) { + resumeByThread.delete(input.sessionId); + await stopOpenCrabsSession(input.sessionId); + } + + const resume = resumeByThread.get(input.sessionId); + const canLoad = resume != null && resume.cwd === input.cwd; + if (resume && resume.cwd !== input.cwd) { + resumeByThread.delete(input.sessionId); + } + + const { path } = await resolveOpenCrabsBinary(); + const handlers: AcpHandlers = {}; + const acp = new AcpClient(input.sessionId, handlers); + const liveRef: { current: Live | null } = { current: null }; + const muteGate = { current: false }; + + handlers.onNotification = (method, params) => { + // Control-plane carve-out: the commands push can land during the + // session/load window — before `live` exists and while transcript replay + // is muted. Muting exists to keep replay out of the transcript; it must + // not drop command discovery, or every restarted session loses the + // autocomplete catalog. + const pushedCommands = nativeCommandsFromUpdate(params); + if (pushedCommands) { + cacheNativeCommands(input.sessionId, pushedCommands); + return; + } + if (muteGate.current) return; + const live = liveRef.current; + if (!live || live.muteUpdates) return; + handleNotification(live, method, params); + }; + handlers.onRequest = (id, method, params) => { + const live = liveRef.current; + if (!live) { + void acp + .respondError(id, { + code: -32601, + message: `Method not found: ${method}`, + }) + .catch(() => undefined); + return; + } + void handleRequest(live, id, method, params); + }; + + // ensureLive runs once per session, so these handlers outlive the turn that + // created them. Route through liveRef so events after turn 1 reach the + // current turn's listener instead of a finished one. + const emit = (event: HarnessEvent) => { + (liveRef.current?.onEvent ?? input.onEvent)(event); + }; + + watchChild( + input.sessionId, + (line) => acp.pushLine(line), + (code) => { + acp.close(new Error("opencrabs exited")); + liveByThread.delete(input.sessionId); + emit({ type: "session.ended", code }); + }, + (line) => { + console.debug("[monocode] opencrabs stderr", line); + }, + ); + + await spawnChild( + input.sessionId, + path, + spawnArgs(input.model), + input.cwd, + undefined, + "opencrabs", + ); + + try { + await acp.request( + "initialize", + { + protocolVersion: 1, + clientCapabilities: CLIENT_CAPABILITIES, + clientInfo: { name: "monocode", version: "0.1.0" }, + }, + INIT_TIMEOUT_MS, + ); + + let setup: unknown; + let acpSessionId: string | undefined; + let resumeFailureReason: string | undefined; + + if (canLoad && resume) { + muteGate.current = true; + try { + setup = await acp.request( + "session/load", + { + sessionId: resume.acpSessionId, + cwd: input.cwd, + mcpServers: [], + }, + SESSION_TIMEOUT_MS, + ); + acpSessionId = sessionIdFromResult(setup) ?? resume.acpSessionId; + } catch (error) { + // Context loss must be visible: MonoCode renders the old transcript + // locally, but a fresh opencrabs session has no memory of it. The + // notice is deferred until the replacement session is confirmed — + // saying "started a fresh session" before session/new succeeds + // would claim a fallback that may never exist. + resumeFailureReason = error instanceof Error ? error.message : String(error); + setup = undefined; + acpSessionId = undefined; + } finally { + muteGate.current = false; + } + } + + if (!acpSessionId) { + setup = await acp.request( + "session/new", + { cwd: input.cwd, mcpServers: [] }, + SESSION_TIMEOUT_MS, + ); + acpSessionId = sessionIdFromResult(setup); + } + if (!acpSessionId) throw new Error("opencrabs did not return a session id"); + if (resumeFailureReason) { + // The fallback session exists: now the boundary notice is true. + emit({ + type: "interjection", + customType: "custom", + text: `OpenCrabs session could not be resumed (${resumeFailureReason}) — started a fresh session. Earlier messages in this transcript are no longer in the agent's context.`, + }); + } + + // Live catalog: replace the static "Configured model" picker entry with + // the server's configured provider/model pairs. + const catalog = modelsFromSessionNew(setup); + if (catalog.available.length > 0) { + setHarnessModels( + "opencrabs", + catalog.available.map((entry) => ({ + id: `opencrabs:${entry.modelId}`, + harness: "opencrabs" as const, + name: entry.name, + nativeId: entry.modelId, + })), + ); + } + + const live: Live = { + acp, + acpSessionId, + threadId: input.sessionId, + cwd: input.cwd, + // Never mute a resumed session: the load-window replay is already + // dropped (live is null and muteGate is closed until the response + // resolves; NDJSON ordering guarantees no stragglers). Unsolicited + // pushes after load are the cross-surface mirror — the whole point. + muteUpdates: false, + cancelled: false, + runtimeMode: input.runtimeMode, + planning: input.intent === "plan", + subagents: new AcpSubagents(), + onEvent: input.onEvent, + approvals: new Map(), + turns: Promise.resolve(), + }; + liveRef.current = live; + liveByThread.set(input.sessionId, live); + resumeByThread.set(input.sessionId, { + acpSessionId, + cwd: input.cwd, + }); + live.onEvent({ + type: "session.providerBound", + providerSessionId: acpSessionId, + }); + // Reflect the server's current model in the thread badge — on load this + // is the restored per-session pick, so the picker survives restarts. + if (catalog.current) { + live.onEvent({ + type: "session.configChanged", + model: `opencrabs:${catalog.current}`, + }); + } + live.onEvent({ type: "session.started" }); + return live; + } catch (error) { + acp.close(error instanceof Error ? error : new Error(String(error))); + await stopOpenCrabsSession(input.sessionId); + // A bare "initialize timed out" names no harness and suggests nothing. + // The connect path is the one place the transcript error line cannot + // say who dropped the ball — name it and hand over the recovery lever. + const detail = error instanceof Error ? error.message : String(error); + throw new Error( + `OpenCrabs failed to connect: ${detail.trim()}\n\n${SERVER_HELP}`, + ); + } +} + +export function spawnArgs(model: string): string[] { + const native = nativeModelId(model).trim(); + // "default" is the placeholder id from the static catalog, not a model the + // server knows — spawning `--model default` only works by fallback luck. + // Omit the flag so the server uses its configured default model. + return native && native.toLowerCase() !== "default" + ? ["acp", "--model", native] + : ["acp"]; +} + +/** Kill the process but keep the ACP session id so we can session/load. */ +export async function stopOpenCrabsSession(sessionId: string): Promise { + cancelledThreads.delete(sessionId); + const live = liveByThread.get(sessionId); + liveByThread.delete(sessionId); + if (live) { + live.muteUpdates = true; + for (const [, resolve] of live.approvals) resolve("deny"); + live.approvals.clear(); + } + live?.acp.close(); + unwatchChild(sessionId); + await killChild(sessionId).catch(() => undefined); +} + +/** Delete or idle detach — drop the OpenCrabs conversation too. */ +export async function forgetOpenCrabsSession(sessionId: string): Promise { + clearNativeCommands(sessionId); + resumeByThread.delete(sessionId); + await stopOpenCrabsSession(sessionId); +} + +/** Seed ACP resume state for a restored MonoCode session. */ +export function bindOpenCrabsSession( + threadId: string, + acpSessionId: string, + cwd: string, +): void { + const sessionId = acpSessionId.trim(); + if (!threadId || !sessionId || !cwd.trim()) return; + resumeByThread.set(threadId, { acpSessionId: sessionId, cwd }); +} + +function handleNotification(live: Live, method: string, params: unknown) { + if (method !== "session/update") return; + const commands = nativeCommandsFromUpdate(params); + if (commands) { + cacheNativeCommands(live.threadId, commands); + return; + } + for (const event of live.subagents.route( + params, + eventsFromAcpUpdate(params), + )) { + live.onEvent(event); + } +} + +async function handleRequest( + live: Live, + id: number, + method: string, + params: unknown, +) { + if (method === "session/request_permission") { + await handlePermission(live, id, params); + return; + } + await live.acp + .respondError(id, { + code: -32601, + message: `Method not found: ${method}`, + }) + .catch(() => undefined); +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsPrompt.test.ts b/src/integrations/harness/providers/opencrabs/opencrabsPrompt.test.ts new file mode 100644 index 0000000000..6f01c3c0b9 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsPrompt.test.ts @@ -0,0 +1,169 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { Attachment } from "../../../../features/sessions/model/session"; + +const invokeMock = vi.fn(); + +vi.mock("@tauri-apps/api/core", () => ({ + invoke: (command: string, args: unknown) => invokeMock(command, args), +})); + +import { openCrabsPromptBlocks } from "./opencrabsPrompt"; + +function imageAttachment(over: Partial): Attachment { + return { + id: "att-1", + name: "shot.png", + mimeType: "image/png", + kind: "image", + size: 4, + ...over, + } as Attachment; +} + +describe("openCrabsPromptBlocks", () => { + beforeEach(() => { + invokeMock.mockReset(); + }); + + it("passes plain text through unchanged", async () => { + const blocks = await openCrabsPromptBlocks("hello", []); + expect(blocks).toEqual([{ type: "text", text: "hello" }]); + expect(invokeMock).not.toHaveBeenCalled(); + }); + + it("rewrites a disk-backed image as a resource_link without persisting", async () => { + const file = imageAttachment({ path: "/tmp/shot.png", data: "aW1hZ2U=" }); + const blocks = await openCrabsPromptBlocks("look", [file]); + expect(blocks[1]).toEqual({ + type: "resource_link", + uri: "file:///tmp/shot.png", + name: "shot.png", + mimeType: "image/png", + size: 4, + }); + expect(invokeMock).not.toHaveBeenCalled(); + }); + + it("persists a pasted blob and links the temp path", async () => { + invokeMock.mockResolvedValue("/tmp/monocode-attachments/1-2-pasted-image.png"); + const file = imageAttachment({ data: "aW1hZ2U=" }); + const blocks = await openCrabsPromptBlocks("", [file]); + expect(invokeMock).toHaveBeenCalledWith("write_attachment", { + name: "pasted-image.png", + data: "aW1hZ2U=", + }); + expect(blocks[0]).toEqual({ + type: "resource_link", + uri: "file:///tmp/monocode-attachments/1-2-pasted-image.png", + name: "pasted-image.png", + mimeType: "image/png", + size: 4, + }); + }); + + it("fails loud when a pasted blob cannot be persisted", async () => { + invokeMock.mockRejectedValue(new Error("disk full")); + const file = imageAttachment({ data: "aW1hZ2U=" }); + await expect(openCrabsPromptBlocks("", [file])).rejects.toThrow( + "disk full", + ); + }); + + it("leaves non-vision resource_link attachments untouched", async () => { + const file: Attachment = { + id: "att-2", + name: "notes.md", + mimeType: "text/markdown", + kind: "attachment", + size: 10, + path: "/tmp/notes.md", + } as Attachment; + const blocks = await openCrabsPromptBlocks("", [file]); + expect(blocks[0]).toEqual({ + type: "resource_link", + uri: "file:///tmp/notes.md", + name: "notes.md", + mimeType: "text/markdown", + size: 10, + }); + expect(invokeMock).not.toHaveBeenCalled(); + }); + + it("links a disk-backed audio attachment without persisting", async () => { + const file: Attachment = { + id: "att-3", + name: "voice.mp3", + mimeType: "audio/mpeg", + kind: "audio", + size: 42, + path: "/tmp/voice.mp3", + } as Attachment; + const blocks = await openCrabsPromptBlocks("transcribe this", [file]); + expect(blocks[1]).toEqual({ + type: "resource_link", + uri: "file:///tmp/voice.mp3", + name: "voice.mp3", + mimeType: "audio/mpeg", + size: 42, + }); + expect(invokeMock).not.toHaveBeenCalled(); + }); + + it("persists a pasted audio blob with a kind-specific name", async () => { + invokeMock.mockResolvedValue("/tmp/monocode-attachments/3-pasted-audio.mp3"); + const file: Attachment = { + id: "att-4", + name: "voice.mp3", + mimeType: "audio/mpeg", + kind: "audio", + size: 42, + data: "YXVkaW8=", + } as Attachment; + const blocks = await openCrabsPromptBlocks("", [file]); + expect(invokeMock).toHaveBeenCalledWith("write_attachment", { + name: "pasted-audio.mp3", + data: "YXVkaW8=", + }); + expect(blocks[0]).toEqual({ + type: "resource_link", + uri: "file:///tmp/monocode-attachments/3-pasted-audio.mp3", + name: "pasted-audio.mp3", + mimeType: "audio/mpeg", + size: 42, + }); + }); + + it("fails loud when an attachment has neither path nor data", async () => { + const file: Attachment = { + id: "att-5", + name: "ghost.pdf", + mimeType: "application/pdf", + kind: "file", + size: 1, + } as Attachment; + await expect(openCrabsPromptBlocks("", [file])).rejects.toThrow( + "no local file path or data", + ); + expect(invokeMock).not.toHaveBeenCalled(); + }); +}); + +describe("spawnArgs (ACP child spawn)", () => { + // The wire-proven dogfood leak: the "Default" catalog placeholder spawned + // `--model default`, which only worked by server fallback luck. + it("omits --model when the placeholder id is selected", async () => { + const { spawnArgs } = await import("./opencrabs"); + expect(spawnArgs("default")).toEqual(["acp"]); + expect(spawnArgs("Default")).toEqual(["acp"]); + }); + + it("passes a real model pair through", async () => { + const { spawnArgs } = await import("./opencrabs"); + // Pairs pass through whole — the server's set_model is pair-aware. + expect(spawnArgs("infer/ali/glm-5.3")).toEqual([ + "acp", + "--model", + "infer/ali/glm-5.3", + ]); + }); +}); diff --git a/src/integrations/harness/providers/opencrabs/opencrabsPrompt.ts b/src/integrations/harness/providers/opencrabs/opencrabsPrompt.ts new file mode 100644 index 0000000000..27a62db3d8 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsPrompt.ts @@ -0,0 +1,81 @@ +import { invoke } from "@tauri-apps/api/core"; +import { fileUri, type PromptContentBlock } from "../../../../features/sessions/model/attachments"; +import type { Attachment } from "../../../../features/sessions/model/session"; + +const EXT_BY_MIME: Record = { + "image/png": "png", + "image/jpeg": "jpg", + "image/jpg": "jpg", + "image/gif": "gif", + "image/webp": "webp", + "audio/mpeg": "mp3", + "audio/mp3": "mp3", + "audio/wav": "wav", + "audio/ogg": "ogg", + "audio/m4a": "m4a", + "audio/webm": "webm", + "application/pdf": "pdf", + "text/plain": "txt", + "text/markdown": "md", +}; + +/** + * Prompt blocks shaped for `opencrabs acp`. + * + * The server declares `promptCapabilities.image: false` and reads + * `resource_link` blocks as on-disk path references, so every attachment + * kind — image, audio, file — is rewritten as a link. The rule is uniform: + * a disk path links as-is; a pasted blob is persisted through the shared + * `write_attachment` command first; an attachment with neither fails the + * send with a visible error. Nothing is silently dropped. + */ +export async function openCrabsPromptBlocks( + text: string, + attachments: Attachment[] = [], +): Promise { + const blocks: PromptContentBlock[] = []; + const trimmed = text.trim(); + if (trimmed) blocks.push({ type: "text", text: trimmed }); + for (const file of attachments) { + blocks.push(await attachmentBlock(file)); + } + return blocks; +} + +async function attachmentBlock(file: Attachment): Promise { + if (file.path?.trim()) { + return { + type: "resource_link", + uri: fileUri(file.path), + name: file.name, + mimeType: file.mimeType, + size: file.size, + }; + } + if (!file.data) { + throw new Error( + `Cannot attach ${JSON.stringify(file.name)}: no local file path or data is available. Attach the file again.`, + ); + } + const name = persistName(file); + const path = await invoke("write_attachment", { + name, + data: file.data, + }); + return { + type: "resource_link", + uri: fileUri(path), + name, + mimeType: file.mimeType, + size: file.size, + }; +} + +/** Pasted blobs get a generated name: kind plus the best extension guess. */ +function persistName(file: Attachment): string { + const ext = + EXT_BY_MIME[file.mimeType] ?? + (file.name.includes(".") ? file.name.split(".").pop() : undefined) ?? + "bin"; + return `pasted-${file.kind}.${ext}`; +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsProtocol.test.ts b/src/integrations/harness/providers/opencrabs/opencrabsProtocol.test.ts new file mode 100644 index 0000000000..c9c498285d --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsProtocol.test.ts @@ -0,0 +1,214 @@ +import { describe, expect, it } from "vitest"; + +import { eventsFromAcpUpdate } from "./opencrabsProtocol"; + +/** + * Usage metering end-to-end: the opencrabs ACP server emits + * `{ sessionUpdate: "usage", usage: { used, size } }` (size is the ACP + * spec's field for the context window total). The meter must translate + * `size` into MonoCode's `window` — with `used` alone, `contextRatio` + * returns null and the context meter renders nothing. + * + * Regression: the adapter originally read only window/contextWindow/ + * context_window, while the server sent `size` — the meter was silently + * dead for every opencrabs session. + */ +describe("eventsFromAcpUpdate usage metering", () => { + it("maps the opencrabs wire shape (used + size) to a context event", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "usage", + usage: { used: 1234, size: 200000 }, + }, + }); + expect(events).toEqual([ + { type: "context", used: 1234, window: 200000 }, + ]); + }); + + it("still maps the legacy window field names", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "usage", + usage: { used: 500, context_window: 128000 }, + }, + }); + expect(events).toEqual([ + { type: "context", used: 500, window: 128000 }, + ]); + }); + + it("emits nothing when the usage block has no readable numbers", () => { + const events = eventsFromAcpUpdate({ + update: { sessionUpdate: "usage", usage: {} }, + }); + expect(events).toEqual([]); + }); +}); + +describe("acpSizeField validation (CodeRabbit: spec is unsigned integer)", () => { + it("accepts a non-negative integer size", () => { + const events = eventsFromAcpUpdate({ + update: { sessionUpdate: "usage", usage: { used: 10, size: 0 } }, + }); + expect(events).toEqual([{ type: "context", used: 10, window: 0 }]); + }); + + it("rejects a negative size instead of poisoning the meter window", () => { + const events = eventsFromAcpUpdate({ + update: { sessionUpdate: "usage", usage: { used: 10, size: -5 } }, + }); + // used survives, window stays undefined — the meter keeps its previous + // window instead of adopting garbage. + expect(events).toEqual([{ type: "context", used: 10, window: undefined }]); + }); + + it("rejects a fractional size", () => { + const events = eventsFromAcpUpdate({ + update: { sessionUpdate: "usage", usage: { used: 10, size: 200000.5 } }, + }); + expect(events).toEqual([{ type: "context", used: 10, window: undefined }]); + }); +}); + +describe("eventsFromAcpUpdate image resource_links", () => { + /** + * The renderer has no filesystem (vite externalizes node:fs — a statSync + * import here fails the production build), so the server stats image + * files as it emits and sends mimeType + size on the block. The adapter + * is pure JSON in, events out. + */ + it("emits image.generated carrying the server's wire metadata", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "tool_call_update", + toolCallId: "call-img-1", + status: "completed", + rawOutput: "chart saved", + content: [ + { type: "text", text: "chart saved" }, + { + type: "resource_link", + uri: "file:///tmp/acp-adapter-img-test.png", + name: "chart.png", + mimeType: "image/png", + size: 1234, + }, + ], + }, + }); + const img = events.find((e) => e.type === "image.generated"); + expect(img).toMatchObject({ + type: "image.generated", + itemId: "call-img-1:1", + path: "/tmp/acp-adapter-img-test.png", + name: "chart.png", + mimeType: "image/png", + size: 1234, + }); + expect(events.some((e) => e.type === "tool.updated")).toBe(true); + }); + + it("falls back to the extension map for blocks without metadata", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "tool_call_update", + toolCallId: "call-img-2", + content: [ + { type: "resource_link", uri: "file:///x/y.jpg", name: "y.jpg" }, + ], + }, + }); + expect(events.find((e) => e.type === "image.generated")).toMatchObject({ + mimeType: "image/jpeg", + size: 0, + }); + }); + + it("ignores non-image extensions and non-file uris", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "tool_call_update", + toolCallId: "call-img-3", + content: [ + { type: "resource_link", uri: "file:///etc/hosts" }, + { type: "resource_link", uri: "https://example.com/x.png" }, + ], + }, + }); + expect(events.some((e) => e.type === "image.generated")).toBe(false); + }); +}); + +describe("eventsFromAcpUpdate subagent classification", () => { + it("classifies the spawn title prefix as an agent card", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "tool_call", + toolCallId: "call-agent-1", + kind: "other", + title: "subagent: fix-auth", + status: "in_progress", + rawInput: { label: "fix-auth" }, + }, + }); + const tool = events[0]; + expect(tool.type).toBe("tool.updated"); + expect((tool as { kind?: string }).kind).toBe("agent"); + }); + + it("leaves ordinary tools unclassified", () => { + const events = eventsFromAcpUpdate({ + update: { + sessionUpdate: "tool_call", + toolCallId: "call-bash-1", + kind: "execute", + title: "bash", + status: "in_progress", + }, + }); + expect((events[0] as { kind?: string }).kind).toBe("execute"); + }); +}); + +/** + * Cross-surface mirror: the opencrabs server pushes rows written by OTHER + * surfaces (Telegram, TUI, cron) as standard session/update frames — + * `{ sessionUpdate, content: { type: "text", text } }`. User rows must + * render user-side (interjection), and each turn's user row must arrive + * BEFORE its agent rows so the interjection block bounds the assistant + * streaming block per turn (no cross-turn blob merging). + */ +describe("eventsFromAcpUpdate cross-surface mirror", () => { + const chunk = (kind: string, text: string) => ({ + update: { sessionUpdate: kind, content: { type: "text", text } }, + }); + + it("renders user_message_chunk as a user-side interjection", () => { + const events = eventsFromAcpUpdate(chunk("user_message_chunk", "from telegram")); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("interjection"); + expect((events[0] as { text: string }).text).toBe("from telegram"); + expect((events[0] as { customType: string }).customType).toBe("custom"); + }); + + it("maps a mirrored turn in order: interjection, thought, message", () => { + const stream = [ + ...eventsFromAcpUpdate(chunk("user_message_chunk", "do the thing")), + ...eventsFromAcpUpdate(chunk("agent_thought_chunk", "pondering")), + ...eventsFromAcpUpdate(chunk("agent_message_chunk", "done")), + ]; + expect(stream.map((e) => e.type)).toEqual([ + "interjection", + "reasoning.delta", + "message.delta", + ]); + }); + + it("maps the batched user_message variant too", () => { + const events = eventsFromAcpUpdate(chunk("user_message", "whole row")); + expect(events).toHaveLength(1); + expect(events[0].type).toBe("interjection"); + expect((events[0] as { text: string }).text).toBe("whole row"); + }); +}); diff --git a/src/integrations/harness/providers/opencrabs/opencrabsProtocol.ts b/src/integrations/harness/providers/opencrabs/opencrabsProtocol.ts new file mode 100644 index 0000000000..2889f056ae --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsProtocol.ts @@ -0,0 +1,471 @@ +import type { RuntimeMode, ToolPreview } from "../../../../features/sessions/model/session"; +import { normalizeTaskListStatus } from "../../../../features/sessions/model/taskList"; +import type { ApprovalDecision, HarnessEvent } from "../../core/types"; +import { + composeToolTitle, + extractSearchQuery, + extractShellCommand, + extractSkillName, + extractToolPreview, +} from "../../core/preview"; +import { acpAgentInfo } from "../../core/acpSubagents"; + +export type OpenCrabsPermissionRequest = { + title: string; + kind?: string; + callId?: string; + preview?: ToolPreview; + optionIds: string[]; +}; + +/** + * ACP `session/update` params -> MonoCode harness events. + * + * Unlike fx, the OpenCrabs ACP server sends structured tool fields + * (kind/title/rawInput/locations), so the shared extraction in ./preview + * does the work and no harness-specific result mining is needed. + */ +export function eventsFromAcpUpdate(params: unknown): HarnessEvent[] { + const rec = asRecord(params); + const update = asRecord(rec?.update) ?? rec; + if (!update) return []; + const kind = String( + update.sessionUpdate ?? update.session_update ?? update.type ?? "", + ); + + if (kind === "user_message_chunk" || kind === "user_message") { + // Mirrored surfaces: user rows here were sent from ANOTHER surface + // (Telegram, TUI, cron) and pushed by the server's cross-surface + // watcher. Render them user-side as interjection blocks — which also + // bounds each turn's assistant streaming block naturally. + const text = textFromContent( + update.content ?? update.text, + kind === "user_message" ? "\n" : "", + ); + return text ? [{ type: "interjection", customType: "custom", text }] : []; + } + + if (kind === "agent_message_chunk" || kind === "agent_message") { + const text = textFromContent( + update.content ?? update.text, + kind === "agent_message" ? "\n" : "", + ); + return text ? [{ type: "message.delta", text }] : []; + } + + if (kind === "agent_thought_chunk" || kind === "agent_thought") { + const text = textFromContent( + update.content ?? update.text, + kind === "agent_thought" ? "\n" : "", + ); + return text ? [{ type: "reasoning.delta", text }] : []; + } + + if ( + kind === "tool_call" || + kind === "tool_call_update" || + kind === "tool_call_content_chunk" + ) { + const tool = + asRecord(update.toolCall) ?? asRecord(update.tool_call) ?? update; + const callId = String( + tool.toolCallId ?? + tool.tool_call_id ?? + update.toolCallId ?? + update.tool_call_id ?? + "", + ); + if (!callId) return []; + const toolKind = stringField(update, "kind") ?? stringField(tool, "kind"); + const status = stringField(update, "status") ?? stringField(tool, "status"); + const preview = extractToolPreview(update, tool); + const title = composeToolTitle({ + kind: toolKind, + title: toolLabel(update) ?? toolLabel(tool), + command: extractShellCommand( + update.rawInput, + tool.rawInput, + update.raw_input, + tool.raw_input, + update.input, + tool.input, + ), + skill: extractSkillName( + update.rawInput, + tool.rawInput, + update.raw_input, + tool.raw_input, + update.input, + tool.input, + ), + path: preview?.path, + query: + preview?.query ?? + extractSearchQuery( + update.rawInput ?? + tool.rawInput ?? + update.raw_input ?? + tool.raw_input ?? + update.input ?? + tool.input, + ), + previewKind: preview?.kind, + }); + return [ + ...imagesFromContent(update, tool, callId), + { + type: "tool.updated", + callId, + kind: toolKind, + status, + detail: toolDetail(update, tool), + preview, + // Classify on the RAW label — composeToolTitle strips the + // "subagent:" prefix acpAgentInfo's detection needs. + ...acpAgentInfo(update, tool, toolKind, toolLabel(update) ?? toolLabel(tool)), + title: title || toolLabel(update) || toolLabel(tool), + }, + ]; + } + + if (kind === "plan" || kind === "current_plan") { + const event = planEvent(update); + return event ? [event] : []; + } + + const usage = usageFromUpdate(update); + return usage ? [usage] : []; +} + +export function permissionRequestFromAcp( + params: unknown, +): OpenCrabsPermissionRequest { + const rec = asRecord(params); + const subject = asRecord(rec?.subject); + const tool = + asRecord(rec?.toolCall) ?? + asRecord(rec?.tool_call) ?? + asRecord(subject?.toolCall) ?? + asRecord(subject) ?? + rec ?? + {}; + const command = stringField(subject ?? {}, "command"); + const kind = stringField(tool, "kind") ?? stringField(subject ?? {}, "kind"); + const preview = mergePreview( + extractToolPreview(tool, tool), + subject ? extractToolPreview(subject, subject) : undefined, + ); + const title = + composeToolTitle({ + kind, + title: toolLabel(tool) ?? toolLabel(subject ?? {}) ?? command, + command: command ?? extractShellCommand(tool, subject), + skill: extractSkillName(tool, subject), + path: preview?.path, + query: preview?.query ?? extractSearchQuery(tool), + previewKind: preview?.kind, + }) || "Permission"; + const options = Array.isArray(rec?.options) ? rec.options : []; + const optionIds = options + .map((item) => asRecord(item)?.optionId ?? asRecord(item)?.option_id) + .filter((value): value is string => typeof value === "string"); + + return { + title, + kind, + callId: + stringField(tool, "toolCallId") ?? + stringField(tool, "tool_call_id") ?? + stringField(rec ?? {}, "toolCallId") ?? + stringField(subject ?? {}, "toolCallId"), + preview, + optionIds, + }; +} + +/** + * Runtime-mode auto-answer for permission requests. `supervised` always asks; + * `auto-accept-edits` still asks for execute/other; looser modes auto-allow. + */ +export function autoPermissionOption( + runtimeMode: RuntimeMode, + kind: string | undefined, + optionIds: string[], +): string | null { + if (optionIds.length === 0) return null; + const tool = (kind ?? "").toLowerCase(); + if (runtimeMode === "supervised") return null; + if ( + runtimeMode === "auto-accept-edits" && + (tool === "execute" || tool === "other") + ) { + return null; + } + if (runtimeMode === "full-access") { + return pickOption(optionIds, [ + "allow-always", + "allow_always", + "allow-once", + "allow_once", + ]); + } + return pickOption(optionIds, [ + "allow-once", + "allow_once", + "allow-always", + "allow_always", + ]); +} + +export function permissionOptionId( + decision: ApprovalDecision, + optionIds: string[], +): string { + if (decision === "allow") { + return ( + pickOption(optionIds, [ + "allow-once", + "allow_once", + "allow-always", + "allow_always", + "allow", + ]) ?? "allow-once" + ); + } + return ( + pickOption(optionIds, [ + "reject-once", + "reject_once", + "reject-always", + "reject_always", + "reject", + "deny", + ]) ?? "reject-once" + ); +} + +export function sessionIdFromResult(result: unknown): string | undefined { + const rec = asRecord(result); + const id = rec?.sessionId ?? rec?.session_id ?? rec?.id; + return typeof id === "string" && id.trim() ? id.trim() : undefined; +} + +/** `session/new` carries the live catalog: models.availableModels + currentModelId. */ +export function modelsFromSessionNew(result: unknown): { + available: { modelId: string; name: string }[]; + current: string | undefined; +} { + const rec = asRecord(result); + const models = asRecord(rec?.models); + const available = Array.isArray(models?.availableModels) + ? models.availableModels.flatMap((item) => { + const entry = asRecord(item); + const modelId = String(entry?.modelId ?? "").trim(); + if (!modelId) return []; + return [{ modelId, name: String(entry?.name ?? modelId).trim() || modelId }]; + }) + : []; + const currentRaw = models?.currentModelId; + return { + available, + current: typeof currentRaw === "string" && currentRaw.trim() ? currentRaw.trim() : undefined, + }; +} + +/** + * The server's `available_commands_update` push: slash commands usable in + * prompts. Returns null for every other update so the caller's regular + * session/update routing is untouched. + */ +export function nativeCommandsFromUpdate( + params: unknown, +): { name: string; description: string }[] | null { + const update = asRecord(asRecord(params)?.update); + if (update?.sessionUpdate !== "available_commands_update") return null; + const list = Array.isArray(update.availableCommands) + ? update.availableCommands + : []; + return list.flatMap((value) => { + const row = asRecord(value); + const name = typeof row?.name === "string" ? row.name.trim() : ""; + if (!name || /[\s/\\]/.test(name)) return []; + return [ + { + name, + description: + typeof row?.description === "string" ? row.description : "", + }, + ]; + }); +} + +/** Image resource_links the tool's content references on disk. */ +const IMAGE_MIME: Record = { + png: "image/png", + jpg: "image/jpeg", + jpeg: "image/jpeg", + gif: "image/gif", + webp: "image/webp", +}; + +function imagesFromContent( + update: Record, + tool: Record, + callId: string, +): HarnessEvent[] { + const content = update.content ?? tool.content; + if (!Array.isArray(content)) return []; + const events: HarnessEvent[] = []; + for (const [i, block] of content.entries()) { + const rec = asRecord(block); + if (!rec || rec.type !== "resource_link") continue; + const uri = typeof rec.uri === "string" ? rec.uri : ""; + if (!uri.startsWith("file://")) continue; + const path = decodeURIComponent(uri.slice("file://".length)); + const ext = path.split(".").pop()?.toLowerCase() ?? ""; + // The server stats the file before emitting the block (the renderer + // has no filesystem), so mimeType and size ride on the wire; the + // extension map is the fallback for older servers. + const mimeType = + typeof rec.mimeType === "string" && rec.mimeType + ? rec.mimeType + : IMAGE_MIME[ext]; + if (!mimeType) continue; + events.push({ + type: "image.generated", + itemId: `${callId}:${i}`, + path, + name: String(rec.name ?? path.split("/").pop() ?? "image"), + mimeType, + size: Number(rec.size ?? 0), + }); + } + return events; +} + +function planEvent(update: Record): HarnessEvent | null { + const entries = update.entries ?? update.plan; + if (Array.isArray(entries)) { + const items = entries.flatMap((item) => { + const rec = asRecord(item); + if (!rec) return []; + const content = String(rec.content ?? rec.text ?? rec.title ?? "").trim(); + if (!content) return []; + return [{ text: content, status: normalizeTaskListStatus(rec.status) }]; + }); + return { type: "tasks.updated", items }; + } + if (typeof update.text === "string" && update.text.trim()) { + return { type: "plan", text: update.text }; + } + return null; +} + +function usageFromUpdate(update: Record): HarnessEvent | null { + const usage = + asRecord(update.usage) ?? + asRecord(update.tokenUsage) ?? + asRecord(update.token_usage); + if (!usage) return null; + const used = + numberField(usage, "used") ?? + numberField(usage, "usedTokens") ?? + numberField(usage, "used_tokens"); + const window = + numberField(usage, "window") ?? + numberField(usage, "contextWindow") ?? + numberField(usage, "context_window") ?? + acpSizeField(usage); + if (used == null && window == null) return null; + return { type: "context", used: used ?? undefined, window: window ?? undefined }; +} + +// ACP `usage.size` is an unsigned integer by spec: a negative or fractional +// value is protocol garbage, not a window, and must not reach +// `mergeContextUsage` (a bogus window makes the meter render nothing). +function acpSizeField(rec: Record): number | undefined { + const value = numberField(rec, "size"); + return value != null && Number.isInteger(value) && value >= 0 ? value : undefined; +} + +function mergePreview( + a: ToolPreview | undefined, + b: ToolPreview | undefined, +): ToolPreview | undefined { + if (!a) return b; + if (!b) return a; + return { ...b, ...a, path: a.path ?? b.path, query: a.query ?? b.query }; +} + +function toolLabel(rec: Record): string | undefined { + return ( + stringField(rec, "title") ?? + stringField(rec, "name") ?? + stringField(rec, "toolName") ?? + stringField(rec, "tool_name") + ); +} + +function toolDetail( + update: Record, + tool: Record, +): string | undefined { + const content = + textFromContent(update.content, "\n") || + textFromContent(tool.content, "\n"); + if (content.trim()) return cap(content); + const output = update.rawOutput ?? tool.rawOutput; + if (typeof output === "string" && output.trim()) return cap(output); + const outputText = textFromContent(output); + return outputText.trim() ? cap(outputText) : undefined; +} + +function cap(value: string, max = 8_000): string { + const text = value.trim(); + if (text.length <= max) return text; + return `${text.slice(0, max)}\n…`; +} + +function pickOption(optionIds: string[], preferred: string[]): string | null { + for (const id of preferred) { + if (optionIds.includes(id)) return id; + } + return null; +} + +function textFromContent(content: unknown, separator = ""): string { + if (typeof content === "string") return content; + const rec = asRecord(content); + if (rec && typeof rec.text === "string") return rec.text; + if (rec && rec.content != null) return textFromContent(rec.content, separator); + if (Array.isArray(content)) { + return content + .map((item) => textFromContent(item, separator)) + .filter(Boolean) + .join(separator); + } + return ""; +} + +export function asRecord(value: unknown): Record | null { + if (value && typeof value === "object" && !Array.isArray(value)) { + return value as Record; + } + return null; +} + +export function stringField( + rec: Record, + key: string, +): string | undefined { + const value = rec[key]; + return typeof value === "string" && value.trim() ? value : undefined; +} + +function numberField( + rec: Record, + key: string, +): number | undefined { + const value = rec[key]; + return typeof value === "number" && Number.isFinite(value) ? value : undefined; +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsText.test.ts b/src/integrations/harness/providers/opencrabs/opencrabsText.test.ts new file mode 100644 index 0000000000..2049749fa5 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsText.test.ts @@ -0,0 +1,125 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const spawned: Array<{ command: string; args: string[]; cwd: string }> = []; +let onLine: ((line: string) => void) | undefined; +let onClose: (() => void) | undefined; +let killCalls = 0; + +vi.mock("../../core/child", () => ({ + resolveOpenCrabsBinary: async () => ({ path: "/fake/opencrabs" }), + spawnChild: async (_id: string, command: string, args: string[], cwd: string) => { + spawned.push({ command, args, cwd }); + }, + killChild: async () => { + killCalls += 1; + }, + unwatchChild: () => undefined, + watchChild: ( + _id: string, + line: (l: string) => void, + close: () => void, + ) => { + onLine = line; + onClose = close; + }, +})); + +import { parseRunSummary, runOpenCrabsTextPrompt } from "./opencrabsText"; + +describe("parseRunSummary", () => { + it("reads the content field of the final JSON object", () => { + const stdout = [ + "🤔 Processing...", + "TITLE-OK", + "{", + ' "content": "TITLE-OK",', + ' "cost": 0.01,', + ' "session_id": "abc"', + "}", + ].join("\n"); + expect(parseRunSummary(stdout)).toBe("TITLE-OK"); + }); + + it("skips brace lines that are not the summary object", () => { + const stdout = '{"noise": true}\nplain answer\n{"content": "real"}'; + expect(parseRunSummary(stdout)).toBe("real"); + }); + + it("returns null when no summary object is present", () => { + expect(parseRunSummary("just plain text\nmore text")).toBeNull(); + }); +}); + +describe("runOpenCrabsTextPrompt", () => { + beforeEach(() => { + spawned.length = 0; + onLine = undefined; + onClose = undefined; + killCalls = 0; + }); + + it("spawns opencrabs run with json format and returns the content", async () => { + const pending = runOpenCrabsTextPrompt({ + cwd: "/repo", + prompt: "title this", + }); + await vi.waitFor(() => { + expect(spawned).toHaveLength(1); + }); + expect(spawned[0]).toEqual({ + command: "/fake/opencrabs", + args: ["run", "--quiet", "--format", "json", "title this"], + cwd: "/repo", + }); + onLine?.("🤔 Processing..."); + onLine?.('{"content": "My Title"}'); + onClose?.(); + await expect(pending).resolves.toBe("My Title"); + }); + + it("rejects when the run produces no summary object", async () => { + const pending = runOpenCrabsTextPrompt({ cwd: "/repo", prompt: "hi" }); + await vi.waitFor(() => { + expect(spawned).toHaveLength(1); + }); + onLine?.("everything broke"); + onClose?.(); + await expect(pending).rejects.toThrow(/everything broke/); + }); + + it("rejects with a timeout error, kills once, and keeps the queue usable", async () => { + vi.useFakeTimers(); + try { + const timed = runOpenCrabsTextPrompt({ + cwd: "/repo", + prompt: "slow", + timeoutMs: 1000, + }); + // Let the spawn settle; the child never emits an exit event. + await vi.advanceTimersByTimeAsync(0); + onLine?.("partial output"); + + // Attach the rejection handler before the timer fires, so the rejection + // never has an unhandled window (Node PromiseRejectionHandledWarning). + const rejection = expect(timed).rejects.toThrow(/timed out after 1000ms/); + await vi.advanceTimersByTimeAsync(1000); + await rejection; + // Timeout kill only — the finally must not fire a second kill. + expect(killCalls).toBe(1); + + // The serialized queue must not be wedged by the timed-out turn. + const next = runOpenCrabsTextPrompt({ + cwd: "/repo", + prompt: "recover", + timeoutMs: 1000, + }); + await vi.advanceTimersByTimeAsync(0); + expect(spawned).toHaveLength(2); + onLine?.('{"content": "recovered"}'); + onClose?.(); + await expect(next).resolves.toBe("recovered"); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/src/integrations/harness/providers/opencrabs/opencrabsText.ts b/src/integrations/harness/providers/opencrabs/opencrabsText.ts new file mode 100644 index 0000000000..539cb34d33 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsText.ts @@ -0,0 +1,132 @@ +import { + killChild, + resolveOpenCrabsBinary, + spawnChild, + unwatchChild, + watchChild, +} from "../../core/child"; + +const TEXT_CHILD_ID = "monocode-opencrabs-text"; +const REQUEST_TIMEOUT_MS = 45_000; + +/** + * One-shot text generation through `opencrabs run --quiet --format json`. + * + * Unlike the REPL/server harnesses there is nothing to keep warm: each + * prompt spawns a short-lived child, captures stdout until exit, and reads + * the structured summary object. `--quiet` (parity branch and later) keeps + * stdout payload-only; the parser below still scans backwards for the final + * parseable object with a string `content` field, so progress noise from a + * binary without the flag degrades gracefully instead of corrupting output. + * Calls are serialized because each one carries the full headless context. + */ +let turns: Promise = Promise.resolve(); + +export async function runOpenCrabsTextPrompt(input: { + cwd: string; + prompt: string; + timeoutMs?: number; +}): Promise { + const run = turns.catch(() => undefined).then(() => promptOnce(input)); + turns = run.then( + () => undefined, + () => undefined, + ); + return run; +} + +async function promptOnce(input: { + cwd: string; + prompt: string; + timeoutMs?: number; +}): Promise { + const { path } = await resolveOpenCrabsBinary(); + const timeoutMs = input.timeoutMs ?? REQUEST_TIMEOUT_MS; + let stdout = ""; + let exited = false; + let notifyExit: () => void = () => undefined; + const exitPromise = new Promise((resolve) => { + notifyExit = resolve; + }); + + watchChild( + TEXT_CHILD_ID, + (line) => { + stdout += `${line}\n`; + }, + () => { + exited = true; + notifyExit(); + }, + ); + + let timedOut = false; + let timeoutKill: Promise | undefined; + const timer = setTimeout(() => { + if (exited) return; + timedOut = true; + // killChild() drops the exit watcher, so resolve the exit promise first — + // otherwise `await exitPromise` hangs forever and blocks the turns queue. + notifyExit(); + timeoutKill = killChild(TEXT_CHILD_ID).catch(() => undefined); + }, timeoutMs); + + try { + await spawnChild( + TEXT_CHILD_ID, + path, + ["run", "--quiet", "--format", "json", input.prompt], + input.cwd, + undefined, + "opencrabs", + ); + await exitPromise; + } finally { + clearTimeout(timer); + unwatchChild(TEXT_CHILD_ID); + if (timeoutKill) { + await timeoutKill; + } else if (!exited) { + await killChild(TEXT_CHILD_ID).catch(() => undefined); + } + } + + if (timedOut) { + throw new Error(`OpenCrabs text generation timed out after ${timeoutMs}ms.`); + } + + const text = parseRunSummary(stdout); + if (!text) { + const nonEmpty = stdout.trim().split("\n").filter(Boolean); + const tail = nonEmpty[nonEmpty.length - 1] ?? ""; + throw new Error( + tail + ? `OpenCrabs text generation failed: ${tail.slice(0, 200)}` + : "OpenCrabs returned empty output.", + ); + } + return text; +} + +/** Last JSON object on stdout with a string `content` field, if any. */ +export function parseRunSummary(stdout: string): string | null { + const lines = stdout.split("\n"); + for (let i = lines.length - 1; i >= 0; i--) { + if (!lines[i].startsWith("{")) continue; + try { + const parsed: unknown = JSON.parse(lines.slice(i).join("\n")); + if ( + parsed && + typeof parsed === "object" && + "content" in parsed && + typeof (parsed as { content: unknown }).content === "string" + ) { + const content = (parsed as { content: string }).content.trim(); + if (content) return content; + } + } catch { + // Not the summary object; keep scanning upwards. + } + } + return null; +} diff --git a/src/integrations/harness/providers/opencrabs/opencrabsTitle.ts b/src/integrations/harness/providers/opencrabs/opencrabsTitle.ts new file mode 100644 index 0000000000..18a4137da1 --- /dev/null +++ b/src/integrations/harness/providers/opencrabs/opencrabsTitle.ts @@ -0,0 +1,26 @@ +import { + buildThreadTitlePrompt, + parseGeneratedSessionTitle, + type GeneratedSessionTitle, +} from "../../../../features/sessions/model/sessionTitle"; +import { runOpenCrabsTextPrompt } from "./opencrabsText"; + +const TITLE_TIMEOUT_MS = 45_000; + +export async function generateOpenCrabsSessionTitle(input: { + sessionId: string; + cwd: string; + message: string; +}): Promise { + try { + const output = await runOpenCrabsTextPrompt({ + cwd: input.cwd, + prompt: buildThreadTitlePrompt(input.message), + timeoutMs: TITLE_TIMEOUT_MS, + }); + return parseGeneratedSessionTitle(output, input.message); + } catch (error) { + console.debug("[monocode] opencrabs session title", error); + return null; + } +}