From f6fdc16d23b3d57896e068d4f4c07f640055e401 Mon Sep 17 00:00:00 2001 From: Zaki Date: Mon, 23 Mar 2026 00:04:02 -0700 Subject: [PATCH] fix(security): block cross-channel approval thread hijacking (#1485) Add source_channel field to Thread struct and verify channel authorization before allowing approval messages to target threads. The web gateway channel is allowed as a trusted approval UI, and threads without a source_channel (e.g. deserialized from older DB records) are permitted for backward compatibility. Closes #1485 Co-Authored-By: Claude Opus 4.6 (1M context) --- .claude/skills/pr-review-batch/SKILL.md | 259 ++++++++ .superset/config.json | 4 + .../2026-02-21-docker-detection-design.md | 101 +++ docs/plans/2026-02-21-docker-detection.md | 450 ++++++++++++++ docs/plans/2026-02-21-skills-tab-design.md | 66 ++ docs/plans/2026-02-21-skills-tab.md | 574 ++++++++++++++++++ .../2026-03-09-routine-silent-failure.md | 480 +++++++++++++++ docs/plans/2026-03-11-security-merge-train.md | 82 +++ ...026-03-22-engine-v2-acceptance-criteria.md | 184 ++++++ scripts/monitor-prs.sh | 139 +++++ src/agent/agent_loop.rs | 2 +- src/agent/compaction.rs | 10 +- src/agent/dispatcher.rs | 4 +- src/agent/session.rs | 178 +++--- src/agent/session_manager.rs | 37 +- src/agent/thread_ops.rs | 37 +- src/channels/web/handlers/chat.rs | 2 +- src/channels/web/server.rs | 2 +- 18 files changed, 2507 insertions(+), 104 deletions(-) create mode 100644 .claude/skills/pr-review-batch/SKILL.md create mode 100644 .superset/config.json create mode 100644 docs/plans/2026-02-21-docker-detection-design.md create mode 100644 docs/plans/2026-02-21-docker-detection.md create mode 100644 docs/plans/2026-02-21-skills-tab-design.md create mode 100644 docs/plans/2026-02-21-skills-tab.md create mode 100644 docs/plans/2026-03-09-routine-silent-failure.md create mode 100644 docs/plans/2026-03-11-security-merge-train.md create mode 100644 docs/plans/2026-03-22-engine-v2-acceptance-criteria.md create mode 100755 scripts/monitor-prs.sh diff --git a/.claude/skills/pr-review-batch/SKILL.md b/.claude/skills/pr-review-batch/SKILL.md new file mode 100644 index 00000000..450bb7d6 --- /dev/null +++ b/.claude/skills/pr-review-batch/SKILL.md @@ -0,0 +1,259 @@ +--- +name: pr-review-batch +description: IronClaw maintainer PR review -- batch review open PRs against ironclaw project standards (Rust, WASM tools, dual-backend DB, security-first) +triggers: + - review PR + - review PRs + - review open PRs + - batch review + - "review #" + - check PRs +--- + +# IronClaw PR Review Workflow + +Maintainer review workflow for the **nearai/ironclaw** repository. Optimized for batch review with parallel data fetching, security-first evaluation against IronClaw's Rust/WASM architecture, and structured GitHub review comments. + +- **Repository:** nearai/ironclaw +- **Maintainer GitHub:** zmanian +- **Primary language:** Rust (async tokio, wasmtime, axum) +- **Key subsystems:** WASM tool sandbox, dual-backend DB (postgres + libsql), LLM provider decorator chain, multi-channel system, SKILL.md skills, registry/installer +- **CI jobs that matter:** Formatting, Clippy (default, all-features, libsql-only, Windows), Regression test enforcement +- **CI jobs that DON'T prove much:** classify, scope (these always pass, even on fork PRs with no secrets) + +## Review Modes + +The user controls how interactive the review is. Detect the mode from their message: + +| User Says | Mode | Behavior | +|-----------|------|----------| +| "Review 938, 933" | **Autonomous** | Fetch, evaluate, post reviews without stopping | +| "Review PRs. Interview me" | **Interactive** | Present findings, ask for input before posting | +| "Check on open PRs" | **Triage** | Summarize state of each PR, ask what to review in depth | +| "Approve 683 and 687" | **Direct verdict** | Post the specified verdict without full analysis | + +**Default is autonomous** unless the user says "interview", "ask me", "discuss", "check with me", or similar. + +## Step 1: Parse PR Numbers + +Extract PR numbers from the user's message. Accept formats: +- "Review 938, 933, 918" +- "Review #834 and #922" +- "Review all open PRs" (use `gh pr list --state open --limit 30`) + +## Step 2: Fetch Data (Parallel) + +For EACH PR, fetch all of these in parallel: + +```bash +# Metadata: title, author, base/head branch, size +gh pr view --json title,author,state,headRefName,baseRefName,additions,deletions,changedFiles \ + --jq '{title, author: .author.login, state, base: .baseRefName, head: .headRefName, additions, deletions, changedFiles}' + +# Full diff +gh pr diff --patch + +# CI status +gh pr checks + +# Previous reviews (for re-reviews) +gh pr view --json reviews --jq '.reviews[] | {author: .author.login, state: .state, body: .body[:200]}' +``` + +For large diffs (>1000 lines), use `gh pr diff --patch | head -500` first, then fetch remaining sections as needed. Note Cargo.lock churn separately -- don't count it as meaningful diff. + +## Step 3: Evaluate Each PR + +Check in this priority order: + +### 3a. CI Status +- All checks must pass -- not just classify/scope. Must have: Formatting, Clippy (all 3 feature combos), Regression test enforcement. +- **Fork PRs (critical gotcha):** Only classify/scope run because GitHub Actions secrets aren't available for fork PRs. The PR will APPEAR to have passing checks. Never trust this. Flag it -- local CI verification or maintainer-triggered re-run required before merge. + +### 3b. Previous Reviews +- Check if zmanian already reviewed -- if so, this is a re-review +- For re-reviews: verify each previous feedback item was addressed, referencing specific commit hashes +- Note reviews from Gemini, Copilot -- cross-reference their findings but don't trust blindly + +### 3c. Security (Highest Priority -- IronClaw-Specific) +- **Identity file write protection:** PROTECTED_IDENTITY_FILES (AGENTS.md, SOUL.md, USER.md, IDENTITY.md) must not become LLM-writable +- **Tool approval requirements:** ApprovalRequirement changes (Never vs UnlessAutoApproved vs Always) -- especially for tools that cross trust boundaries (tool_install, tool_auth, build_tool, shell) +- **WASM sandbox boundaries:** fuel limits, memory limits, network allowlists must not be weakened +- **Credential handling:** no secrets in logs/errors/SSE events; use `redact_params()` before broadcast +- **SSRF vectors:** URL validation must resolve DNS before checking for private/loopback IPs +- **Prompt injection defense:** sanitizer/validator/policy changes in `src/safety/` + +### 3d. Correctness (IronClaw-Specific) +- **No `.unwrap()/.expect()` in production code** (tests are fine) +- **String safety:** no byte-index slicing (`&s[..n]`) on user/external strings -- use `is_char_boundary()` or `char_indices()` +- **Dual-backend DB:** new persistence features must support BOTH postgres AND libsql. Check for missing trait implementations. +- **Feature flags:** changes must compile under `--no-default-features --features libsql`, default, and `--all-features` +- **Transaction safety:** multi-step DB operations wrapped in transactions (both backends) +- **LLM provider decorator chain:** new `LlmProvider` trait methods must be delegated in ALL wrapper types (grep `impl LlmProvider for`) + +### 3e. Architecture & Conventions +- `crate::` for cross-module imports (not `super::` except tests and intra-module) +- `thiserror` for error types in `error.rs`; map errors with context via `.map_err()` +- Strong types over strings (enums, newtypes) +- Module specs followed -- if a module has a CLAUDE.md (agent, web, db, llm, setup, tools, workspace), check it +- Module-owned initialization: init logic lives in owning module as public factory fn, not in main.rs/app.rs +- No unnecessary dependencies (check `~/.claude/approved-dependencies.md` list) + +### 3f. Tests +- Bug fixes MUST have regression tests (enforced by CI regression-check job and commit-msg hook) +- Tests use `tempfile` crate, not hardcoded `/tmp/` paths +- No real network requests in tests (use mocks or RFC 5737 TEST-NET IPs like 192.0.2.1) +- Test names and comments match actual test behavior and assertions +- `[skip-regression-check]` in commit message or PR label only if genuinely not feasible + +## Step 4: Interview (Interactive Mode) + +In interactive mode, present findings and ask for the maintainer's judgment before posting. **Do NOT post reviews until the maintainer confirms.** + +### When to Interview (Even in Autonomous Mode) + +Always pause and ask the maintainer when you encounter: + +1. **Judgment calls on architecture direction** -- "This PR adds a named provider for Z.AI. Should we prefer named providers or push contributors toward openai_compatible for niche providers?" +2. **Security tradeoffs with usability** -- "Removing approval from tool_install reduces friction but weakens the trust boundary. What's your stance?" +3. **Scope creep concerns** -- "This PR started as a bug fix but adds 300 lines of new feature. Accept as-is or ask to split?" +4. **Dependency additions** -- "This adds `datafusion` (heavy dep). Worth it for the use case?" +5. **Contradictory signals** -- "Gemini approved but Copilot flagged a real issue. The code works but the pattern is fragile." +6. **Taking over vs requesting changes** -- "This PR has 5+ issues. Want me to take it over or send detailed feedback?" +7. **Merge ordering for conflicting PRs** -- "PRs #933 and #918 both modify cli/mod.rs. Which should land first?" + +### Interview Format + +Present findings concisely, then ask a specific question: + +``` +**PR #922: Relax tool approval requirements** + +The HTTP GET change is clean (tiered: credentials->Always, GET->Never, other->UnlessAutoApproved). + +But it also removes approval from: +- build_tool (can execute shell commands) +- tool_install (downloads WASM modules) +- tool_auth (grants credentials to tools) + +These cross the trust boundary. Options: +1. Approve as-is (maximum convenience) +2. Request changes: keep build_tool + extension tools gated, accept the rest +3. Request changes: revert everything except HTTP GET and list_dir + +Which direction? +``` + +Wait for the maintainer's response before posting. + +### Triage Mode + +In triage mode, present a dashboard first: + +``` +| PR | Author | Title | CI | Reviews | Age | Risk | +|----|--------|-------|----|---------|-----|------| +| #938 | reidliu41 | Z.AI provider | green | none | 1d | low | +| #922 | ilblackdragon | relax approvals | green | copilot:concern | 2d | medium | +| #927 | ilblackdragon | chat onboarding | green | zmanian:changes | 3d | high | +``` + +Then ask: "Which ones should I review in depth? Or should I go through all of them?" + +## Step 5: Determine Verdict + +| Verdict | Criteria | +|---------|----------| +| **APPROVE** | Clean, follows IronClaw patterns, full CI green, no security issues, tests present | +| **REQUEST CHANGES** | Security regressions, functional bugs, .expect() in production, trust boundary violations, missing dual-backend support, missing error handling | +| **COMMENT** | Good direction but needs discussion, or already approved with observations | + +In interactive mode, confirm the verdict with the maintainer before posting. In autonomous mode, post directly. + +## Step 6: Post Reviews + +Post reviews via `gh pr review` using HEREDOC for body formatting. + +### New Review Format + +``` +## Review: + +<1-2 sentence assessment> + +Positives: +- +- + +### : + + +### : + + +Minor notes: +- + + +``` + +Severity levels: Critical, Concerning, Minor (non-blocking) + +### Re-Review Format + +``` +## Re-review: + +All/N items from my previous review have been resolved: + +1. **** -- Fixed in commit . . +2. **** -- Fixed.
. + + + +LGTM. +``` + +## Step 7: Handle GitHub API Errors + +GitHub 502s are common during batch posting. Retry with `sleep 5` between attempts. Post reviews sequentially (not in parallel) to avoid rate limits. + +## Step 8: Summary + +After all reviews are posted, provide a summary table: + +``` +| PR | Title | Verdict | +|----|-------|---------| +| #938 | Z.AI provider | Approved | +| #933 | channels list CLI | Approved | +| #918 | skills CLI | Approved | +``` + +Note cross-PR conflicts (e.g., PRs that both modify `src/cli/mod.rs` and snapshot files). + +## Special Cases + +### Fork PRs +Only classify/scope CI jobs run. **Never merge with only these passing.** Either: +- Run local CI: `cargo check --all-features && cargo clippy --all && cargo test` +- Or trigger full CI by pushing a maintainer commit to the PR branch + +### Registry/WASM PRs +- Verify artifact URLs match the naming convention: `---wasm32-wasip2.tar.gz` +- Check SHA256 checksums against actual release assets +- Ensure `name` field in manifest matches crate_name in source config +- Cross-reference with `.github/workflows/release.yml` for automated patching + +### Taking Over a PR +When a contributor PR has too many issues: +1. Create new branch from staging +2. Cherry-pick or apply the contributor's changes +3. Fix the issues +4. Create superseding PR referencing the original + +### Cross-PR Context +When PRs are related (e.g., all touch registry manifests, or both modify cli/mod.rs), post context comments on each explaining how they fit together and merge ordering. + +### Batch Merge +When the user says "merge" after reviews, use `gh pr merge --squash` for each approved PR. Verify CI is still green before each merge. diff --git a/.superset/config.json b/.superset/config.json new file mode 100644 index 00000000..9d9aba5a --- /dev/null +++ b/.superset/config.json @@ -0,0 +1,4 @@ +{ + "setup": [], + "teardown": [] +} diff --git a/docs/plans/2026-02-21-docker-detection-design.md b/docs/plans/2026-02-21-docker-detection-design.md new file mode 100644 index 00000000..a760b314 --- /dev/null +++ b/docs/plans/2026-02-21-docker-detection-design.md @@ -0,0 +1,101 @@ +# Proactive Docker Detection + +Date: 2026-02-21 + +## Problem + +IronClaw's sandbox system requires Docker but provides no proactive guidance. Docker availability is only checked at runtime when a sandbox job is attempted, resulting in a confusing error. Users have no way to know during setup or startup whether Docker is properly configured. + +## Goals + +1. Detect Docker installation AND daemon running status at two points: setup wizard and every startup +2. Provide platform-specific installation guidance (macOS, Linux, Windows) +3. Surface Docker status clearly in the boot screen +4. Allow users to skip/continue without Docker (sandbox is optional) + +## Non-Goals + +- Auto-installing Docker +- Changing the default sandbox setting (stays `enabled: false`) +- Modifying the existing `connect_docker()` function + +## Design + +### Docker Status Model + +New file `src/sandbox/detect.rs` with centralized detection: + +```rust +pub enum DockerStatus { + Available, // Binary on PATH + daemon responding to ping + NotInstalled, // `docker` binary not found on PATH + NotRunning, // Binary found but daemon not responding + Disabled, // Sandbox not enabled (no check performed) +} + +pub enum Platform { MacOS, Linux, Windows } + +pub struct DockerDetection { + pub status: DockerStatus, + pub platform: Platform, +} +``` + +Detection logic: +1. Check if `docker` binary exists on PATH (reuse `which`/`where` pattern from `skills/gating.rs`) +2. If found, attempt `connect_docker()` to ping the daemon +3. Return `Available`, `NotInstalled`, or `NotRunning` + +### Platform-Specific Guidance + +| Platform | Not Installed | Not Running | +|----------|--------------|-------------| +| macOS | "Install Docker Desktop: https://docs.docker.com/desktop/install/mac-install/" | "Start Docker Desktop from Applications, or run: open -a Docker" | +| Linux | "Install Docker Engine: https://docs.docker.com/engine/install/" | "Start the Docker daemon: sudo systemctl start docker" | +| Windows | "Install Docker Desktop: https://docs.docker.com/desktop/install/windows-install/" | "Start Docker Desktop from the Start menu" | + +### Wizard Step (First-Run) + +Add Step 8 "Docker Sandbox" (current steps 8 becomes 9, total becomes 9): + +1. Ask "Do you want to enable Docker sandbox for isolated code execution?" +2. If yes, run Docker detection +3. Based on status: + - **Available**: Enable sandbox, confirm + - **Not Installed**: Show install guidance, offer to skip or retry after installing + - **Not Running**: Show start guidance, offer to skip or retry +4. If user skips, sandbox stays disabled + +### Startup Check (Every Launch) + +In `main.rs`, when `config.sandbox.enabled == true`, before creating `ContainerJobManager`: + +1. Run `DockerDetection::check()` +2. If **Available**: proceed normally +3. If **NotInstalled** or **NotRunning**: log warning, disable sandbox for this session, continue startup + +### Boot Screen Changes + +`BootInfo` gains `docker_status: DockerStatus` field. + +Features line rendering: +- `Available` + enabled: `sandbox` (as today) +- `NotInstalled` + enabled in config: `sandbox (docker not installed)` +- `NotRunning` + enabled in config: `sandbox (docker not running)` +- `Disabled`: no sandbox shown (as today) + +Warning lines shown in yellow when Docker is configured but unavailable. + +## Files + +| Action | File | Change | +|--------|------|--------| +| Create | `src/sandbox/detect.rs` | Detection logic, platform hints | +| Modify | `src/sandbox/mod.rs` | Export `detect` module | +| Modify | `src/setup/wizard.rs` | Add Docker/Sandbox wizard step | +| Modify | `src/main.rs` | Startup check before ContainerJobManager | +| Modify | `src/boot_screen.rs` | Show Docker status | + +## Dependencies + +No new crate dependencies. Uses existing `bollard` (via `connect_docker()`), `std::process::Command` (for binary detection), and `std::env::consts::OS` (for platform detection). diff --git a/docs/plans/2026-02-21-docker-detection.md b/docs/plans/2026-02-21-docker-detection.md new file mode 100644 index 00000000..aec71f2c --- /dev/null +++ b/docs/plans/2026-02-21-docker-detection.md @@ -0,0 +1,450 @@ +# Docker Detection Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** Add proactive Docker detection at startup and in the setup wizard, with platform-specific installation guidance. + +**Architecture:** New `src/sandbox/detect.rs` module for centralized Docker detection. Wizard gets a new step. Startup check in `main.rs` warns and disables sandbox if Docker unavailable. Boot screen shows Docker status. + +**Tech Stack:** Rust, bollard (existing), std::process::Command + +--- + +### Task 1: Create `src/sandbox/detect.rs` -- Docker Detection Module + +**Files:** +- Create: `src/sandbox/detect.rs` +- Modify: `src/sandbox/mod.rs` + +**Step 1: Write the failing test** + +```rust +// In src/sandbox/detect.rs + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_detect_platform() { + let platform = Platform::current(); + // Should return a valid platform on any CI/dev machine + match platform { + Platform::MacOS | Platform::Linux | Platform::Windows => {} + } + } + + #[test] + fn test_install_hint_not_empty() { + for platform in [Platform::MacOS, Platform::Linux, Platform::Windows] { + assert!(!platform.install_hint().is_empty()); + assert!(!platform.start_hint().is_empty()); + } + } + + #[test] + fn test_docker_status_display() { + assert_eq!(DockerStatus::Available.as_str(), "available"); + assert_eq!(DockerStatus::NotInstalled.as_str(), "not installed"); + assert_eq!(DockerStatus::NotRunning.as_str(), "not running"); + assert_eq!(DockerStatus::Disabled.as_str(), "disabled"); + } + + #[test] + fn test_docker_status_is_ok() { + assert!(DockerStatus::Available.is_ok()); + assert!(!DockerStatus::NotInstalled.is_ok()); + assert!(!DockerStatus::NotRunning.is_ok()); + assert!(!DockerStatus::Disabled.is_ok()); + } + + #[tokio::test] + async fn test_check_docker_returns_valid_status() { + let result = check_docker().await; + // On CI without Docker, should be NotInstalled or NotRunning + // On dev with Docker, should be Available + // Either way, should not panic + match result.status { + DockerStatus::Available + | DockerStatus::NotInstalled + | DockerStatus::NotRunning => {} + DockerStatus::Disabled => panic!("check_docker should never return Disabled"), + } + } +} +``` + +**Step 2: Write the implementation** + +```rust +//! Proactive Docker detection with platform-specific guidance. + +/// Docker daemon availability status. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DockerStatus { + /// Docker binary found on PATH and daemon responding to ping. + Available, + /// `docker` binary not found on PATH. + NotInstalled, + /// Binary found but daemon not responding. + NotRunning, + /// Sandbox feature not enabled (no check performed). + Disabled, +} + +impl DockerStatus { + pub fn is_ok(&self) -> bool { + matches!(self, DockerStatus::Available) + } + + pub fn as_str(&self) -> &'static str { + match self { + DockerStatus::Available => "available", + DockerStatus::NotInstalled => "not installed", + DockerStatus::NotRunning => "not running", + DockerStatus::Disabled => "disabled", + } + } +} + +/// Host platform for install guidance. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Platform { + MacOS, + Linux, + Windows, +} + +impl Platform { + pub fn current() -> Self { + match std::env::consts::OS { + "macos" => Platform::MacOS, + "windows" => Platform::Windows, + _ => Platform::Linux, + } + } + + pub fn install_hint(&self) -> &'static str { + match self { + Platform::MacOS => "Install Docker Desktop: https://docs.docker.com/desktop/install/mac-install/", + Platform::Linux => "Install Docker Engine: https://docs.docker.com/engine/install/", + Platform::Windows => "Install Docker Desktop: https://docs.docker.com/desktop/install/windows-install/", + } + } + + pub fn start_hint(&self) -> &'static str { + match self { + Platform::MacOS => "Start Docker Desktop from Applications, or run: open -a Docker", + Platform::Linux => "Start the Docker daemon: sudo systemctl start docker", + Platform::Windows => "Start Docker Desktop from the Start menu", + } + } +} + +/// Result of a Docker detection check. +pub struct DockerDetection { + pub status: DockerStatus, + pub platform: Platform, +} + +/// Check whether Docker is installed and running. +/// +/// 1. Checks if `docker` binary exists on PATH +/// 2. If found, tries to connect and ping the Docker daemon +/// 3. Returns `Available`, `NotInstalled`, or `NotRunning` +pub async fn check_docker() -> DockerDetection { + let platform = Platform::current(); + + // Step 1: Check if docker binary is on PATH + if !docker_binary_exists() { + return DockerDetection { + status: DockerStatus::NotInstalled, + platform, + }; + } + + // Step 2: Try to connect to the daemon + match crate::sandbox::connect_docker().await { + Ok(_) => DockerDetection { + status: DockerStatus::Available, + platform, + }, + Err(_) => DockerDetection { + status: DockerStatus::NotRunning, + platform, + }, + } +} + +/// Check if the `docker` binary exists on PATH. +fn docker_binary_exists() -> bool { + #[cfg(unix)] + { + std::process::Command::new("which") + .arg("docker") + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .is_ok_and(|s| s.success()) + } + #[cfg(windows)] + { + std::process::Command::new("where") + .arg("docker") + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .is_ok_and(|s| s.success()) + } +} +``` + +**Step 3: Export from `src/sandbox/mod.rs`** + +Add `pub mod detect;` and re-export key types. + +**Step 4: Run tests** + +Run: `cargo test sandbox::detect::tests -- --nocapture` +Expected: All pass + +**Step 5: Clippy** + +Run: `cargo clippy --all --all-features` +Expected: Zero warnings on new code + +**Step 6: Commit** + +```bash +git add src/sandbox/detect.rs src/sandbox/mod.rs +git commit -m "feat: add Docker detection module with platform guidance" +``` + +--- + +### Task 2: Update Boot Screen to Show Docker Status + +**Files:** +- Modify: `src/boot_screen.rs` + +**Step 1: Add `docker_status` to `BootInfo`** + +Add field: `pub docker_status: DockerStatus` (import from `crate::sandbox::detect::DockerStatus`). + +**Step 2: Update `print_boot_screen` features rendering** + +When sandbox is enabled in config but Docker isn't available, show a warning: +- `DockerStatus::Available`: "sandbox" (as today) +- `DockerStatus::NotInstalled`: "sandbox (docker not installed)" in yellow +- `DockerStatus::NotRunning`: "sandbox (docker not running)" in yellow +- `DockerStatus::Disabled`: don't show sandbox (as today) + +**Step 3: Update tests** + +Update all 3 existing `BootInfo` test structs to include `docker_status` field. + +**Step 4: Run tests** + +Run: `cargo test boot_screen::tests` +Expected: All pass + +**Step 5: Commit** + +```bash +git add src/boot_screen.rs +git commit -m "feat: show Docker status in boot screen" +``` + +--- + +### Task 3: Add Startup Docker Check in `main.rs` + +**Files:** +- Modify: `src/main.rs` + +**Step 1: Add Docker check before `ContainerJobManager` creation** + +Before line ~989 (`let container_job_manager = if config.sandbox.enabled`), insert: + +```rust +// Proactive Docker detection +let docker_status = if config.sandbox.enabled { + let detection = ironclaw::sandbox::detect::check_docker().await; + match detection.status { + ironclaw::sandbox::detect::DockerStatus::Available => { + tracing::info!("Docker is available"); + detection.status + } + ironclaw::sandbox::detect::DockerStatus::NotInstalled => { + tracing::warn!( + "Docker is not installed. Sandbox disabled for this session. {}", + detection.platform.install_hint() + ); + detection.status + } + ironclaw::sandbox::detect::DockerStatus::NotRunning => { + tracing::warn!( + "Docker is installed but not running. Sandbox disabled for this session. {}", + detection.platform.start_hint() + ); + detection.status + } + ironclaw::sandbox::detect::DockerStatus::Disabled => detection.status, + } +} else { + ironclaw::sandbox::detect::DockerStatus::Disabled +}; +``` + +Then gate the `ContainerJobManager` creation on `docker_status.is_ok()`: +```rust +let container_job_manager = if config.sandbox.enabled && docker_status.is_ok() { + // ... existing code ... +``` + +**Step 2: Pass `docker_status` to `BootInfo`** + +In the boot screen construction, add the `docker_status` field. + +**Step 3: Run full test suite** + +Run: `cargo test` +Expected: All pass + +**Step 4: Commit** + +```bash +git add src/main.rs +git commit -m "feat: check Docker availability at startup" +``` + +--- + +### Task 4: Add Docker/Sandbox Wizard Step + +**Files:** +- Modify: `src/setup/wizard.rs` + +**Step 1: Increment `total_steps` from 8 to 9** + +**Step 2: Add `step_docker_sandbox()` method** + +Insert after Extensions (step 7), before Heartbeat: + +```rust +/// Step 8: Docker Sandbox +async fn step_docker_sandbox(&mut self) -> Result<(), SetupError> { + print_info("The Docker sandbox provides isolated execution for code generation,"); + print_info("builds, and untrusted commands. It requires Docker to be installed."); + println!(); + + if !confirm("Enable Docker sandbox?", false).map_err(SetupError::Io)? { + self.settings.sandbox.enabled = false; + print_info("Sandbox disabled. You can enable it later with SANDBOX_ENABLED=true."); + return Ok(()); + } + + // Check Docker availability + let detection = crate::sandbox::detect::check_docker().await; + + match detection.status { + crate::sandbox::detect::DockerStatus::Available => { + self.settings.sandbox.enabled = true; + print_success("Docker is installed and running. Sandbox enabled."); + } + crate::sandbox::detect::DockerStatus::NotInstalled => { + println!(); + print_error("Docker is not installed."); + print_info(detection.platform.install_hint()); + println!(); + // Offer retry or skip + if confirm("Retry after installing Docker?", false).map_err(SetupError::Io)? { + let retry = crate::sandbox::detect::check_docker().await; + if retry.status.is_ok() { + self.settings.sandbox.enabled = true; + print_success("Docker is now available. Sandbox enabled."); + } else { + self.settings.sandbox.enabled = false; + print_info("Docker still not available. Sandbox disabled for now."); + } + } else { + self.settings.sandbox.enabled = false; + print_info("Sandbox disabled. Install Docker and set SANDBOX_ENABLED=true later."); + } + } + crate::sandbox::detect::DockerStatus::NotRunning => { + println!(); + print_error("Docker is installed but not running."); + print_info(detection.platform.start_hint()); + println!(); + if confirm("Retry after starting Docker?", false).map_err(SetupError::Io)? { + let retry = crate::sandbox::detect::check_docker().await; + if retry.status.is_ok() { + self.settings.sandbox.enabled = true; + print_success("Docker is now running. Sandbox enabled."); + } else { + self.settings.sandbox.enabled = false; + print_info("Docker still not responding. Sandbox disabled for now."); + } + } else { + self.settings.sandbox.enabled = false; + print_info("Sandbox disabled. Start Docker and set SANDBOX_ENABLED=true later."); + } + } + _ => { + self.settings.sandbox.enabled = false; + } + } + + Ok(()) +} +``` + +**Step 3: Wire into `run()` method** + +```rust +// Step 8: Docker Sandbox +print_step(8, total_steps, "Docker Sandbox"); +self.step_docker_sandbox().await?; +self.persist_after_step().await; + +// Step 9: Heartbeat (was Step 8) +print_step(9, total_steps, "Background Tasks"); +self.step_heartbeat()?; +self.persist_after_step().await; +``` + +**Step 4: Run tests** + +Run: `cargo test setup` +Expected: All pass + +**Step 5: Commit** + +```bash +git add src/setup/wizard.rs +git commit -m "feat: add Docker sandbox step to setup wizard" +``` + +--- + +### Task 5: Final Verification + +**Step 1: Run full test suite** + +Run: `cargo test` +Expected: All pass + +**Step 2: Run clippy** + +Run: `cargo clippy --all --all-features --benches --tests --examples` +Expected: Zero warnings + +**Step 3: Check for unwrap/expect in production code** + +Grep changed files for `.unwrap()` and `.expect(` -- should have none in production code. + +**Step 4: Verify both feature flags compile** + +Run: `cargo check` and `cargo check --no-default-features --features libsql` +Expected: Both clean diff --git a/docs/plans/2026-02-21-skills-tab-design.md b/docs/plans/2026-02-21-skills-tab-design.md new file mode 100644 index 00000000..ab4a5e81 --- /dev/null +++ b/docs/plans/2026-02-21-skills-tab-design.md @@ -0,0 +1,66 @@ +# Skills Tab - Web UI Design + +## Goal + +Add a Skills tab to the IronClaw web gateway that lets users browse installed skills, search ClawHub for new skills, and install/remove skills -- all from the browser. + +## Scope + +**Frontend only.** The REST API endpoints already exist: + +| Method | Endpoint | Purpose | +|--------|----------|---------| +| GET | `/api/skills` | List installed skills | +| POST | `/api/skills/search` | Search ClawHub + local | +| POST | `/api/skills/install` | Install (requires `X-Confirm-Action: true`) | +| DELETE | `/api/skills/{name}` | Remove (requires `X-Confirm-Action: true`) | + +No Rust changes needed. + +## Layout + +Three sections inside the tab panel: + +### 1. Search ClawHub + +A search input at the top. On submit, calls `POST /api/skills/search` and renders catalog results as dashed-border cards (matching the "available extension" pattern). Cards that match an already-installed skill show "Installed" instead of an Install button. + +Staggered fade-in animation on search results for polish. + +### 2. Installed Skills + +Grid of cards for all locally loaded skills. Each card shows: +- **Name** (bold, `.ext-name` style) +- **Trust badge**: "Trusted" (green) or "Installed" (blue) -- small pill +- **Version** (small, secondary text) +- **Description** (`.ext-desc` style) +- **Activation keywords** as small tags (`.ext-keywords` style) +- **Remove button** -- only for registry-installed skills (trust=Installed), not user-placed trusted skills + +### 3. Install by URL + +A small form matching the WASM install form pattern: +- Name input +- URL input (HTTPS) +- Install button + +## Visual Design + +Reuses existing `.ext-card`, `.extensions-list`, `.extensions-section`, `.btn-ext` classes. New CSS limited to: +- `.skill-trust` badge pill (green for Trusted, blue for Installed) +- `.skill-version` small version label +- Staggered `@keyframes skillFadeIn` for search results +- `.skill-search-box` for the search input styling + +## Files Modified + +- `src/channels/web/static/index.html` -- Add Skills tab button + panel markup +- `src/channels/web/static/app.js` -- Add `loadSkills()`, `searchClawHub()`, `installSkill()`, `removeSkill()`, render functions, wire into `switchTab()` +- `src/channels/web/static/style.css` -- Trust badge styles, search box, fade-in animation + +## Decisions + +- Reuse ext-card classes rather than creating a parallel card system +- Trust badge differentiates skills from extensions visually +- Confirmation uses `window.confirm()` dialog matching `removeExtension()` pattern +- Search is manual (button/enter) not live-as-you-type to avoid hammering ClawHub diff --git a/docs/plans/2026-02-21-skills-tab.md b/docs/plans/2026-02-21-skills-tab.md new file mode 100644 index 00000000..1eca1c72 --- /dev/null +++ b/docs/plans/2026-02-21-skills-tab.md @@ -0,0 +1,574 @@ +# Skills Tab Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** Add a Skills tab to the IronClaw web UI for browsing installed skills, searching ClawHub, and installing/removing skills. + +**Architecture:** Frontend-only changes to three static files (HTML, CSS, JS). The REST API (`/api/skills/*`) already exists and needs no modification. Follows the existing Extensions tab pattern: card grid layout, `apiFetch()` helper, `showToast()` for feedback. + +**Tech Stack:** Vanilla HTML/CSS/JS (no frameworks), existing design system (CSS variables, `ext-card` family classes). + +--- + +### Task 1: Add Skills tab button and panel markup to index.html + +**Files:** +- Modify: `src/channels/web/static/index.html:39-44` (tab bar) and `188-232` (before extensions panel) + +**Step 1: Add the Skills tab button** + +In `index.html`, inside the `.tab-bar` div, add a Skills button between Extensions and the spacer. Change lines 43-44 from: + +```html + +
+``` + +to: + +```html + + +
+``` + +**Step 2: Add the Skills tab panel** + +Add the Skills panel markup after the Extensions panel closing `` (after line 232) and before the toasts div: + +```html + +
+
+
+

Search ClawHub

+ +
+
+
+

Installed Skills

+
+
Loading skills...
+
+
+
+

Install Skill by URL

+
+ + + +
+
+
+
+``` + +**Step 3: Verify the HTML is well-formed** + +Open the file and confirm the new panel is between the Extensions panel closing tag and `
`. + +**Step 4: Commit** + +```bash +git add src/channels/web/static/index.html +git commit -m "feat(web): add Skills tab markup to index.html" +``` + +--- + +### Task 2: Add Skills CSS (trust badges, search box, fade-in animation) + +**Files:** +- Modify: `src/channels/web/static/style.css` (append before the `@media` responsive block at line 2810) + +**Step 1: Add skill-specific CSS** + +Insert the following CSS before the `/* --- Activity toolbar --- */` comment (before line 2810): + +```css +/* --- Skills tab --- */ + +.skill-search-box { + display: flex; + gap: 8px; + align-items: center; + margin-bottom: 12px; +} + +.skill-search-box input { + flex: 1; + padding: 8px 12px; + background: var(--bg); + border: 1px solid var(--border); + border-radius: var(--radius); + color: var(--text); + font-size: 13px; +} + +.skill-search-box input:focus { + outline: none; + border-color: var(--accent); + box-shadow: 0 0 0 3px rgba(52, 211, 153, 0.1); +} + +.skill-search-box button { + padding: 8px 20px; + background: var(--accent); + color: #09090b; + border: none; + border-radius: var(--radius); + cursor: pointer; + font-size: 13px; + font-weight: 600; + transition: background 0.2s, transform 0.2s; +} + +.skill-search-box button:hover { + background: var(--accent-hover); + transform: translateY(-1px); +} + +.skill-trust { + font-size: 10px; + padding: 2px 6px; + border-radius: 8px; + font-weight: 500; + text-transform: uppercase; + letter-spacing: 0.3px; +} + +.skill-trust.trust-trusted { + background: rgba(52, 211, 153, 0.15); + color: var(--success); +} + +.skill-trust.trust-installed { + background: rgba(96, 165, 250, 0.15); + color: #60a5fa; +} + +.skill-version { + font-size: 11px; + color: var(--text-secondary); + font-family: var(--font-mono); +} + +@keyframes skillFadeIn { + from { opacity: 0; transform: translateY(8px); } + to { opacity: 1; transform: translateY(0); } +} + +.skill-search-result { + animation: skillFadeIn 0.3s ease-out both; +} +``` + +**Step 2: Commit** + +```bash +git add src/channels/web/static/style.css +git commit -m "feat(web): add Skills tab CSS styles" +``` + +--- + +### Task 3: Wire Skills tab into switchTab() and keyboard shortcuts + +**Files:** +- Modify: `src/channels/web/static/app.js:823-827` (switchTab function) and `2700-2704` (keyboard shortcuts) + +**Step 1: Add skills tab loading to switchTab()** + +In the `switchTab()` function, after line 827 (`if (tab === 'extensions') loadExtensions();`), add: + +```javascript + if (tab === 'skills') loadSkills(); +``` + +**Step 2: Update keyboard shortcut tab array** + +At line 2702, change: + +```javascript + const tabs = ['chat', 'memory', 'jobs', 'routines', 'extensions']; +``` + +to: + +```javascript + const tabs = ['chat', 'memory', 'jobs', 'routines', 'extensions', 'skills']; +``` + +And update the key range check at line 2700 from `'5'` to `'6'`: + +```javascript + if (mod && e.key >= '1' && e.key <= '6') { +``` + +**Step 3: Commit** + +```bash +git add src/channels/web/static/app.js +git commit -m "feat(web): wire Skills tab into switchTab and keyboard shortcuts" +``` + +--- + +### Task 4: Implement loadSkills() -- render installed skills + +**Files:** +- Modify: `src/channels/web/static/app.js` (add new section after the Extensions section, before keyboard shortcuts) + +**Step 1: Add the loadSkills function** + +Add this code block before the `// --- Keyboard shortcuts ---` comment (before line 2692): + +```javascript +// --- Skills --- + +function loadSkills() { + var skillsList = document.getElementById('skills-list'); + apiFetch('/api/skills').then(function(data) { + if (!data.skills || data.skills.length === 0) { + skillsList.innerHTML = '
No skills installed
'; + return; + } + skillsList.innerHTML = ''; + for (var i = 0; i < data.skills.length; i++) { + skillsList.appendChild(renderSkillCard(data.skills[i])); + } + }).catch(function(err) { + skillsList.innerHTML = '
Failed to load skills: ' + escapeHtml(err.message) + '
'; + }); +} + +function renderSkillCard(skill) { + var card = document.createElement('div'); + card.className = 'ext-card'; + + var header = document.createElement('div'); + header.className = 'ext-header'; + + var name = document.createElement('span'); + name.className = 'ext-name'; + name.textContent = skill.name; + header.appendChild(name); + + var trust = document.createElement('span'); + var trustClass = skill.trust.toLowerCase() === 'trusted' ? 'trust-trusted' : 'trust-installed'; + trust.className = 'skill-trust ' + trustClass; + trust.textContent = skill.trust; + header.appendChild(trust); + + var version = document.createElement('span'); + version.className = 'skill-version'; + version.textContent = 'v' + skill.version; + header.appendChild(version); + + card.appendChild(header); + + var desc = document.createElement('div'); + desc.className = 'ext-desc'; + desc.textContent = skill.description; + card.appendChild(desc); + + if (skill.keywords && skill.keywords.length > 0) { + var kw = document.createElement('div'); + kw.className = 'ext-keywords'; + kw.textContent = 'Activates on: ' + skill.keywords.join(', '); + card.appendChild(kw); + } + + var actions = document.createElement('div'); + actions.className = 'ext-actions'; + + // Only show Remove for registry-installed skills, not user-placed trusted skills + if (skill.trust.toLowerCase() !== 'trusted') { + var removeBtn = document.createElement('button'); + removeBtn.className = 'btn-ext remove'; + removeBtn.textContent = 'Remove'; + removeBtn.addEventListener('click', function() { removeSkill(skill.name); }); + actions.appendChild(removeBtn); + } + + card.appendChild(actions); + return card; +} +``` + +**Step 2: Commit** + +```bash +git add src/channels/web/static/app.js +git commit -m "feat(web): implement loadSkills and renderSkillCard" +``` + +--- + +### Task 5: Implement searchClawHub() -- search and render catalog results + +**Files:** +- Modify: `src/channels/web/static/app.js` (add after `renderSkillCard`, before keyboard shortcuts) + +**Step 1: Add search and catalog card rendering** + +Add this code after the `renderSkillCard` function: + +```javascript +function searchClawHub() { + var input = document.getElementById('skill-search-input'); + var query = input.value.trim(); + if (!query) return; + + var resultsDiv = document.getElementById('skill-search-results'); + resultsDiv.innerHTML = '
Searching...
'; + + apiFetch('/api/skills/search', { + method: 'POST', + body: { query: query }, + }).then(function(data) { + resultsDiv.innerHTML = ''; + + // Show catalog results + if (data.catalog && data.catalog.length > 0) { + // Build a set of installed skill names for quick lookup + var installedNames = {}; + if (data.installed) { + for (var j = 0; j < data.installed.length; j++) { + installedNames[data.installed[j].name] = true; + } + } + + for (var i = 0; i < data.catalog.length; i++) { + var card = renderCatalogSkillCard(data.catalog[i], installedNames); + card.style.animationDelay = (i * 0.06) + 's'; + resultsDiv.appendChild(card); + } + } + + // Show matching installed skills too + if (data.installed && data.installed.length > 0) { + for (var k = 0; k < data.installed.length; k++) { + var installedCard = renderSkillCard(data.installed[k]); + installedCard.style.animationDelay = ((data.catalog ? data.catalog.length : 0) + k) * 0.06 + 's'; + installedCard.classList.add('skill-search-result'); + resultsDiv.appendChild(installedCard); + } + } + + if (resultsDiv.children.length === 0) { + resultsDiv.innerHTML = '
No skills found for "' + escapeHtml(query) + '"
'; + } + }).catch(function(err) { + resultsDiv.innerHTML = '
Search failed: ' + escapeHtml(err.message) + '
'; + }); +} + +function renderCatalogSkillCard(entry, installedNames) { + var card = document.createElement('div'); + card.className = 'ext-card ext-available skill-search-result'; + + var header = document.createElement('div'); + header.className = 'ext-header'; + + var name = document.createElement('span'); + name.className = 'ext-name'; + name.textContent = entry.name || entry.slug; + header.appendChild(name); + + if (entry.version) { + var version = document.createElement('span'); + version.className = 'skill-version'; + version.textContent = 'v' + entry.version; + header.appendChild(version); + } + + card.appendChild(header); + + if (entry.description) { + var desc = document.createElement('div'); + desc.className = 'ext-desc'; + desc.textContent = entry.description; + card.appendChild(desc); + } + + var actions = document.createElement('div'); + actions.className = 'ext-actions'; + + var slug = entry.slug || entry.name; + var isInstalled = installedNames[entry.name] || installedNames[slug]; + + if (isInstalled) { + var label = document.createElement('span'); + label.className = 'ext-active-label'; + label.textContent = 'Installed'; + actions.appendChild(label); + } else { + var installBtn = document.createElement('button'); + installBtn.className = 'btn-ext install'; + installBtn.textContent = 'Install'; + installBtn.addEventListener('click', (function(s, btn) { + return function() { + if (!confirm('Install skill "' + s + '" from ClawHub?')) return; + btn.disabled = true; + btn.textContent = 'Installing...'; + installSkill(s, null, btn); + }; + })(slug, installBtn)); + actions.appendChild(installBtn); + } + + card.appendChild(actions); + return card; +} + +// Wire up Enter key on search input +document.getElementById('skill-search-input').addEventListener('keydown', function(e) { + if (e.key === 'Enter') searchClawHub(); +}); +``` + +**Step 2: Commit** + +```bash +git add src/channels/web/static/app.js +git commit -m "feat(web): implement ClawHub search with staggered card animation" +``` + +--- + +### Task 6: Implement installSkill() and removeSkill() + +**Files:** +- Modify: `src/channels/web/static/app.js` (add after search functions, before keyboard shortcuts) + +**Step 1: Add install and remove functions** + +Add this code after the search event listener: + +```javascript +function installSkill(nameOrSlug, url, btn) { + var body = { name: nameOrSlug }; + if (url) body.url = url; + + apiFetch('/api/skills/install', { + method: 'POST', + headers: { 'X-Confirm-Action': 'true' }, + body: body, + }).then(function(res) { + if (res.success) { + showToast('Installed skill "' + nameOrSlug + '"', 'success'); + } else { + showToast('Install failed: ' + (res.message || 'unknown error'), 'error'); + } + loadSkills(); + if (btn) { btn.disabled = false; btn.textContent = 'Install'; } + }).catch(function(err) { + showToast('Install failed: ' + err.message, 'error'); + if (btn) { btn.disabled = false; btn.textContent = 'Install'; } + }); +} + +function removeSkill(name) { + if (!confirm('Remove skill "' + name + '"?')) return; + apiFetch('/api/skills/' + encodeURIComponent(name), { + method: 'DELETE', + headers: { 'X-Confirm-Action': 'true' }, + }).then(function(res) { + if (res.success) { + showToast('Removed skill "' + name + '"', 'success'); + } else { + showToast('Remove failed: ' + (res.message || 'unknown error'), 'error'); + } + loadSkills(); + }).catch(function(err) { + showToast('Remove failed: ' + err.message, 'error'); + }); +} + +function installSkillFromForm() { + var name = document.getElementById('skill-install-name').value.trim(); + if (!name) { showToast('Skill name is required', 'error'); return; } + var url = document.getElementById('skill-install-url').value.trim() || null; + if (url && !url.startsWith('https://')) { + showToast('URL must use HTTPS', 'error'); + return; + } + if (!confirm('Install skill "' + name + '"?')) return; + installSkill(name, url, null); + document.getElementById('skill-install-name').value = ''; + document.getElementById('skill-install-url').value = ''; +} +``` + +**Step 2: Commit** + +```bash +git add src/channels/web/static/app.js +git commit -m "feat(web): implement installSkill, removeSkill, and form handler" +``` + +--- + +### Task 7: Fix apiFetch to merge extra headers properly + +**Files:** +- Modify: `src/channels/web/static/app.js:86-98` (apiFetch function) + +**Context:** The current `apiFetch` function sets `opts.headers` as an object and always overwrites with `Authorization`. When we pass `headers: { 'X-Confirm-Action': 'true' }` in options, the current code does `opts.headers = opts.headers || {}` which preserves our custom headers, then adds Authorization. However, `fetch()` expects headers as a `Headers` object or plain object -- the plain object approach works fine. Verify this works by reading the function carefully. + +**Step 1: Verify apiFetch handles extra headers** + +Read `app.js:86-98`. The current code: +```javascript +function apiFetch(path, options) { + const opts = options || {}; + opts.headers = opts.headers || {}; + opts.headers['Authorization'] = 'Bearer ' + token; + ... +} +``` + +This correctly merges: if we pass `{ headers: { 'X-Confirm-Action': 'true' } }`, it keeps our header and adds Authorization. **No change needed.** Move on. + +**Step 2: Commit (skip -- no changes)** + +--- + +### Task 8: Manual testing and final commit + +**Step 1: Verify the HTML is valid** + +Open `src/channels/web/static/index.html` and confirm: +- The Skills tab button appears in the tab bar +- The `tab-skills` panel has the correct structure +- No unclosed tags + +**Step 2: Verify the JS doesn't have syntax errors** + +Run a quick syntax check (if node is available): +```bash +node -c src/channels/web/static/app.js +``` + +**Step 3: Test the tab appears and loads** + +Start the app and open the web gateway. Verify: +1. Skills tab appears in the tab bar between Extensions and the spacer +2. Clicking it shows the three sections +3. Installed skills load and display with trust badges and keywords +4. ClawHub search returns results with staggered animation +5. Install from search works (with confirm dialog) +6. Remove works for registry-installed skills +7. Install by URL form works +8. Cmd+6 keyboard shortcut switches to Skills tab + +**Step 4: Final commit if any fixes were needed** + +```bash +git add src/channels/web/static/index.html src/channels/web/static/app.js src/channels/web/static/style.css +git commit -m "feat(web): complete Skills tab with ClawHub search, install, and remove" +``` diff --git a/docs/plans/2026-03-09-routine-silent-failure.md b/docs/plans/2026-03-09-routine-silent-failure.md new file mode 100644 index 00000000..439becfd --- /dev/null +++ b/docs/plans/2026-03-09-routine-silent-failure.md @@ -0,0 +1,480 @@ +# Fix Routine Silent Failures (#697) Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** When full_job routines fail due to missing sandbox/Docker infrastructure, surface loud, clear errors to the user instead of failing silently. + +**Architecture:** Three layers of improvement: (1) incorporate PR #711's sync mechanism so dispatched job completions/failures propagate back to routine runs, (2) fail fast at dispatch time when sandbox is configured but Docker is unavailable by threading sandbox availability into RoutineEngine, (3) send a user-visible notification at startup when sandbox is disabled due to missing Docker. + +**Tech Stack:** Rust, tokio, thiserror + +--- + +## Prerequisites + +- Branch from `main` (not from the existing `fix/697-routine-silent-failure` branch) +- We will incorporate PR #711's changes as part of this PR, making #711 superseded + +--- + +### Task 1: Add `list_dispatched_routine_runs` to Database trait and implementations + +PR #711 adds this method. We incorporate it here. + +**Files:** +- Modify: `src/db/mod.rs` (RoutineStore trait) +- Modify: `src/db/postgres.rs` +- Modify: `src/db/libsql/routines.rs` +- Modify: `src/history/store.rs` + +**Step 1: Add trait method to RoutineStore** + +In `src/db/mod.rs`, add to the `RoutineStore` trait (after `link_routine_run_to_job`): + +```rust +/// List routine runs that were dispatched as full_job (status = 'running' +/// with a linked job_id). Used by the routine engine to sync completion +/// status from the background job. +async fn list_dispatched_routine_runs(&self) -> Result, DatabaseError>; +``` + +**Step 2: Implement for PostgreSQL** + +In `src/db/postgres.rs`, add the implementation (delegating to `Store`): + +```rust +async fn list_dispatched_routine_runs(&self) -> Result, DatabaseError> { + self.inner.list_dispatched_routine_runs().await +} +``` + +**Step 3: Implement for libSQL** + +In `src/db/libsql/routines.rs`, add: + +```rust +pub async fn list_dispatched_routine_runs( + &self, +) -> Result, DatabaseError> { + let conn = self.pool.connection().await.map_err(|e| { + DatabaseError::Query(format!("failed to get connection: {e}")) + })?; + let mut rows = conn + .query( + "SELECT id, routine_id, trigger_type, trigger_detail, started_at, \ + completed_at, status, result_summary, tokens_used, job_id, created_at \ + FROM routine_runs WHERE status = 'running' AND job_id IS NOT NULL", + (), + ) + .await + .map_err(|e| DatabaseError::Query(e.to_string()))?; + + let mut runs = Vec::new(); + while let Some(row) = rows.next().await.map_err(|e| DatabaseError::Query(e.to_string()))? { + runs.push(parse_routine_run_row(&row)?); + } + Ok(runs) +} +``` + +**Step 4: Implement for Store wrapper** + +In `src/history/store.rs`, add: + +```rust +pub async fn list_dispatched_routine_runs(&self) -> Result, DatabaseError> { + sqlx::query_as::<_, RoutineRunRow>( + "SELECT id, routine_id, trigger_type, trigger_detail, started_at, \ + completed_at, status, result_summary, tokens_used, job_id, created_at \ + FROM routine_runs WHERE status = 'running' AND job_id IS NOT NULL" + ) + .fetch_all(&self.pool) + .await + .map(|rows| rows.into_iter().map(Into::into).collect()) + .map_err(|e| DatabaseError::Query(e.to_string())) +} +``` + +**Step 5: Verify compilation** + +```bash +cargo check +cargo check --no-default-features --features libsql +``` + +**Step 6: Commit** + +```bash +git add src/db/mod.rs src/db/postgres.rs src/db/libsql/routines.rs src/history/store.rs +git commit -m "feat(db): add list_dispatched_routine_runs for routine-job sync (#697)" +``` + +--- + +### Task 2: Add sync_dispatched_runs and fix dispatch status in routine_engine + +Incorporates PR #711's core fix: change `execute_full_job` to return `RunStatus::Running` instead of `Ok`, and add the periodic sync mechanism. + +**Files:** +- Modify: `src/agent/routine_engine.rs` + +**Step 1: Write tests for job-state-to-run-status mapping and Running notification gating** + +Add to the `mod tests` block at the bottom of `routine_engine.rs`: + +```rust +#[test] +fn test_running_status_does_not_notify() { + let config = NotifyConfig { + on_success: true, + on_failure: true, + on_attention: true, + ..Default::default() + }; + + let should_notify = match RunStatus::Running { + RunStatus::Ok => config.on_success, + RunStatus::Attention => config.on_attention, + RunStatus::Failed => config.on_failure, + RunStatus::Running => false, + }; + assert!(!should_notify); +} + +#[test] +fn test_full_job_dispatch_returns_running_status() { + assert_eq!(RunStatus::Running.to_string(), "running"); +} + +/// Regression test for #697: full_job routines were immediately marked Ok +/// on dispatch, so failures/completions were never synced back. +#[test] +fn test_job_state_to_run_status_mapping() { + use crate::context::JobState; + + let map_state = |state: JobState, reason: Option<&str>| -> Option<(RunStatus, String)> { + let last_reason = reason.map(|s| s.to_string()); + match state { + JobState::Completed | JobState::Submitted | JobState::Accepted => { + let summary = + last_reason.unwrap_or_else(|| "Job completed successfully".to_string()); + Some((RunStatus::Ok, summary)) + } + JobState::Failed => { + let summary = last_reason + .unwrap_or_else(|| "Job failed (no error message recorded)".to_string()); + Some((RunStatus::Failed, summary)) + } + JobState::Cancelled => Some((RunStatus::Failed, "Job was cancelled".to_string())), + JobState::Pending | JobState::InProgress | JobState::Stuck => None, + } + }; + + let (status, _) = map_state(JobState::Completed, None).unwrap(); + assert_eq!(status, RunStatus::Ok); + + let (status, _) = map_state(JobState::Failed, Some("OOM killed")).unwrap(); + assert_eq!(status, RunStatus::Failed); + assert_eq!(summary, "OOM killed"); + + let (status, summary) = map_state(JobState::Failed, None).unwrap(); + assert_eq!(status, RunStatus::Failed); + assert!(summary.contains("no error message")); + + assert!(map_state(JobState::Pending, None).is_none()); + assert!(map_state(JobState::InProgress, None).is_none()); + assert!(map_state(JobState::Stuck, None).is_none()); +} +``` + +**Step 2: Run tests to verify they fail** + +```bash +cargo test routine_engine::tests --all-features +``` + +Expected: compilation error since `sync_dispatched_runs` doesn't exist yet. + +**Step 3: Add import and sync methods** + +Add `use crate::context::JobState;` to the imports. + +Add `sync_dispatched_runs` and `complete_dispatched_run` methods to `impl RoutineEngine` (after `check_cron_triggers`). See PR #711 diff for exact implementation. + +Change `execute_full_job` return from: +```rust +Ok((RunStatus::Ok, Some(summary), None)) +``` +to: +```rust +Ok((RunStatus::Running, Some(summary), None)) +``` + +Update the summary message to include "Status will be updated when the job completes." + +Add `engine.sync_dispatched_runs().await;` to the cron ticker loop in `spawn_cron_ticker`, after `check_cron_triggers`. + +**Step 4: Run tests** + +```bash +cargo test routine_engine::tests --all-features +``` + +Expected: PASS + +**Step 5: Commit** + +```bash +git add src/agent/routine_engine.rs +git commit -m "fix(routines): sync dispatched full_job runs with job completion (#697)" +``` + +--- + +### Task 3: Fail fast when sandbox is unavailable at dispatch time + +This is the new work beyond PR #711. Thread sandbox availability into `RoutineEngine` so `execute_full_job` can fail immediately with a clear error instead of dispatching a doomed job. + +**Files:** +- Modify: `src/agent/routine_engine.rs` +- Modify: `src/agent/agent_loop.rs` + +**Step 1: Write the failing test** + +Add to `mod tests` in `routine_engine.rs`: + +```rust +#[test] +fn test_sandbox_unavailable_error_message() { + let err = RoutineError::JobDispatchFailed { + reason: "Sandbox is enabled but Docker is not available. \ + Install Docker or set SANDBOX_ENABLED=false to run full_job routines." + .to_string(), + }; + let msg = err.to_string(); + assert!(msg.contains("Docker is not available")); + assert!(msg.contains("SANDBOX_ENABLED")); +} +``` + +**Step 2: Run test to verify it passes (this one is a unit test for the error variant)** + +```bash +cargo test routine_engine::tests::test_sandbox_unavailable_error_message --all-features +``` + +Expected: PASS (error variant already exists, we're just testing the message). + +**Step 3: Add `sandbox_available` field to `RoutineEngine`** + +In `src/agent/routine_engine.rs`, add a field to the `RoutineEngine` struct: + +```rust +/// Whether sandbox/Docker infrastructure is available for full_job execution. +sandbox_available: bool, +``` + +Update `RoutineEngine::new` to accept and store it: + +```rust +pub fn new( + config: RoutineConfig, + store: Arc, + llm: Arc, + workspace: Arc, + notify_tx: mpsc::Sender, + scheduler: Option>, + sandbox_available: bool, +) -> Self { + Self { + config, + store, + llm, + workspace, + notify_tx, + running_count: Arc::new(AtomicUsize::new(0)), + event_cache: Arc::new(RwLock::new(Vec::new())), + scheduler, + sandbox_available, + } +} +``` + +**Step 4: Add sandbox check in `execute_full_job`** + +At the top of `execute_full_job`, before the scheduler check, add a sandbox availability check. This requires passing `sandbox_available` through `EngineContext`. + +Add `sandbox_available: bool` to `EngineContext`. + +Update `spawn_fire` and `fire_manual` to pass `self.sandbox_available` into `EngineContext`. + +In `execute_full_job`, add before the scheduler check: + +```rust +if !ctx.sandbox_available { + return Err(RoutineError::JobDispatchFailed { + reason: "Sandbox is enabled but Docker is not available. \ + Install Docker or set SANDBOX_ENABLED=false to run full_job routines." + .to_string(), + }); +} +``` + +**Step 5: Update call site in `agent_loop.rs`** + +In `src/agent/agent_loop.rs`, where `RoutineEngine::new` is called (~line 442), pass the sandbox availability. The `Agent` struct needs to know Docker status. The simplest approach: + +Add a `sandbox_available: bool` field to `Agent` (or to `AgentDeps`). Set it during construction based on the `docker_status` from `main.rs`. The value flows: `main.rs` detects Docker -> passes `sandbox_available` bool through `AppComponents` or `AgentDeps` -> `Agent` passes it to `RoutineEngine::new`. + +Look at how `main.rs` passes config to `Agent`. The `docker_status` is computed in `main.rs`. The cleanest path: +- Add `sandbox_available: bool` to `AppComponents` (set in `main.rs`) +- Thread it through to `AgentDeps` -> `Agent` -> `RoutineEngine::new` + +Alternatively, since `config.sandbox.enabled` is already available in the agent, just add one more bool. Check the existing flow and pick the minimal path. + +**Step 6: Verify compilation** + +```bash +cargo check --all-features +cargo check --no-default-features --features libsql +``` + +**Step 7: Run tests** + +```bash +cargo test routine_engine::tests --all-features +``` + +**Step 8: Commit** + +```bash +git add src/agent/routine_engine.rs src/agent/agent_loop.rs src/main.rs src/app.rs +git commit -m "fix(routines): fail fast when sandbox unavailable at dispatch time (#697)" +``` + +--- + +### Task 4: Surface sandbox unavailability to user via notification channel + +Currently the Docker detection warning only goes to `tracing::warn` (logs). Users on TUI/web never see it. Send a user-visible notification after channels are set up. + +**Files:** +- Modify: `src/main.rs` + +**Step 1: Write the test** + +This is a startup behavior change, so the test is an integration-level assertion. Add a unit test for the notification message formatting: + +In `src/agent/routine_engine.rs` tests (or a new test in main.rs tests if they exist): + +```rust +#[test] +fn test_sandbox_warning_message_format() { + let msg = format!( + "Sandbox is enabled but Docker is not available -- full_job routines will fail. {}", + "Install Docker Desktop from https://docker.com/get-started" + ); + assert!(msg.contains("full_job routines will fail")); + assert!(msg.contains("Docker")); +} +``` + +**Step 2: Add startup notification in `main.rs`** + +After the channel manager is set up and the agent is running, if `config.sandbox.enabled && !docker_status.is_ok()`, send a warning message through the channel manager. The pattern already exists for heartbeat/routine notifications. + +The exact location: after `channels` is fully initialized (after all channels are added), but before the agent run loop. Find where `channels.broadcast_all` is accessible. + +The simplest approach: after the agent starts (`agent.run()` is typically the last call), but since that blocks, the notification should be sent *before* `agent.run()` is called, using a spawned task or inline send. + +Look at where heartbeat startup notifications go. Mirror that pattern: + +```rust +if config.sandbox.enabled && !docker_status.is_ok() { + let warning = format!( + "Warning: Sandbox is enabled but Docker is not available -- \ + full_job routines will fail until Docker is running. {}", + docker_status_detection.platform.install_hint() + ); + let response = OutgoingResponse { + content: warning, + thread_id: None, + attachments: Vec::new(), + metadata: serde_json::json!({ + "source": "system", + "type": "warning", + }), + }; + let channels_clone = channels.clone(); + tokio::spawn(async move { + // Small delay to let channels finish connecting + tokio::time::sleep(std::time::Duration::from_secs(2)).await; + let _ = channels_clone.broadcast_all("default", response).await; + }); +} +``` + +Note: we need to preserve the `detection` struct (not just `docker_status`) to access `platform.install_hint()`. Adjust the variable binding in the Docker detection block to keep it available. + +**Step 3: Verify compilation** + +```bash +cargo check --all-features +``` + +**Step 4: Commit** + +```bash +git add src/main.rs +git commit -m "feat(startup): notify user when sandbox unavailable (#697)" +``` + +--- + +### Task 5: Final verification and cleanup + +**Step 1: Run full test suite** + +```bash +cargo fmt +cargo clippy --all --benches --tests --examples --all-features +cargo test --all-features +``` + +**Step 2: Verify both feature configurations compile** + +```bash +cargo check --no-default-features --features libsql +cargo check +``` + +**Step 3: Run pre-commit safety checks** + +```bash +grep -rnE '\.unwrap\(|\.expect\(' src/agent/routine_engine.rs src/main.rs +``` + +Expect: no hits in production code (test code is fine). + +**Step 4: Create final commit if any formatting/clippy fixes needed** + +```bash +git add -A +git commit -m "style: formatting and clippy fixes (#697)" +``` + +--- + +## Summary of Changes + +| What | Where | Why | +|------|-------|-----| +| `list_dispatched_routine_runs` DB method | `db/mod.rs`, postgres, libsql, store | Query for running routine runs with linked jobs | +| `sync_dispatched_runs()` engine method | `routine_engine.rs` | Periodically sync job completion back to routine runs | +| `RunStatus::Running` on dispatch | `routine_engine.rs` | Don't mark as Ok before job actually completes | +| `sandbox_available` flag | `RoutineEngine`, `EngineContext` | Fail fast at dispatch when Docker missing | +| Startup notification | `main.rs` | Warn user visibly when sandbox is disabled | + +## PR Scope + +This PR supersedes PR #711 by incorporating its changes plus the additional fail-fast and startup notification work. PR #711 can be closed after this merges. diff --git a/docs/plans/2026-03-11-security-merge-train.md b/docs/plans/2026-03-11-security-merge-train.md new file mode 100644 index 00000000..58974546 --- /dev/null +++ b/docs/plans/2026-03-11-security-merge-train.md @@ -0,0 +1,82 @@ +# Security Merge Train Status Board + +Date opened: 2026-03-11 +Last updated: 2026-03-12 +Base branch: `staging` +Current `staging` head: `acea1143cf70f7fa593c077620c979d5aa260de9` + +This board started as the security merge train plan and now tracks the live status of the approved-PR merge effort. + +## Current Branch Health + +- Full staging batch for `acea1143cf70f7fa593c077620c979d5aa260de9` completed green. +- E2E, Linux tests, Windows builds, Docker build, WASM WIT compatibility, staging gate, and summary all passed. +- Current gating problem is no longer branch regressions. It is fresh review requirements on replacement PRs. + +## Merged Into `staging` + +| PR | Title | Outcome | +|---|---|---| +| #510 | fix(security): add DOMPurify and sanitize rendered markdown | Merged | +| #518 | fix(security): resolve DNS once and reuse for SSRF validation | Merged | +| #520 | fix(security): harden auth token env overlay usage / WASM metadata loading hardening | Merged | +| #949 | fix(setup): drain residual events and filter key kind in onboard prompts | Merged | +| #935 | fix(mcp): stdio/unix transports skip initialize handshake | Merged | +| #760 | fix(agent): block thread_id-based context pollution across users | Merged | +| #752 | fix(mcp): header safety validation and Authorization conflict bug from #704 | Merged | +| #735 | fix: drain tunnel pipes to prevent zombie process | Merged | +| #684 | fix(setup): validate channel credentials during setup | Merged | +| #850 | docs: add Russian localization (README.ru.md) | Merged | +| #851 | feat(setup): display ASCII art banner during onboarding | Merged | +| #964 | fix(ci): disambiguate WASM bundle filenames to prevent tool/channel collision | Merged | +| #839 | fix(test): stabilize openai compat oversized-body regression | Merged | +| #472 | Fix systemctl unit | Merged | + +## Security Replacement Queue + +These supersede the originally approved but dirty security PRs. + +| Replacement PR | Supersedes | CI | Auto-merge | Merge blocker | Notes | +|---|---|---|---|---|---| +| #966 | #514 | Green | Enabled | `REVIEW_REQUIRED` | CSP replacement; includes E2E coverage | +| #967 | #516 | Green | Enabled | `REVIEW_REQUIRED` | FullAccess policy guard | +| #968 | #522 | Green | Enabled | `REVIEW_REQUIRED` | Safe env overlay / set_var invariants | +| #970 | #513 | Green | Enabled | `REVIEW_REQUIRED` | Webhook HMAC migration | + +## General Replacement Queue + +These supersede other approved dirty PRs that were still worth carrying forward. + +| Replacement PR | Supersedes | CI | Auto-merge | Merge blocker | Notes | +|---|---|---|---|---|---| +| #986 | #793 | Green | Enabled | `REVIEW_REQUIRED` | Non-OAuth HTTP MCP clients now carry session manager | +| #987 | #679 | In progress / early checks green | Enabled | `REVIEW_REQUIRED` | Preserves `selected_model` when re-running setup on the same backend | + +## Approved Originals Still Open + +| PR | Title | Current state | Recommended action | Notes | +|---|---|---|---|---| +| #514 | fix(security): add Content-Security-Policy header to web gateway | Dirty | Ignore in favor of #966 | Replacement path is active | +| #516 | fix(security): require explicit `SANDBOX_ALLOW_FULL_ACCESS` to enable FullAccess policy | Dirty | Ignore in favor of #967 | Replacement path is active | +| #522 | fix(security): make unsafe `env::set_var` calls safe with explicit invariants | Dirty | Ignore in favor of #968 | Replacement path is active | +| #513 | fix(security): migrate webhook auth to HMAC-SHA256 signature header | Dirty | Ignore in favor of #970 | Replacement path is active | +| #793 | fix(mcp): set session manager on non-OAuth HTTP MCP clients | Dirty | Ignore in favor of #986 | Replacement path is active | +| #679 | fix(setup): preserve model selection on provider re-run | Dirty | Ignore in favor of #987 | Replacement path is active | +| #737 | 汉化v0.1.0 | Dirty | Do not open a faithful replacement | `staging` already has a divergent i18n implementation | +| #831 | refactor(orchestrator/api): use `test_secrets_store()` helper in credentials test | Dirty + draft | Do not rescue as-is | Current diff has drifted far beyond the title / intended scope | +| #934 | fix(memory): reject absolute filesystem paths with corrective routing | Unstable | Do not merge as-is | Default-branch workflow change would break staging promotion in this repo | +| #616 | feat: adds context-llm tool support | Unstable | Separate review pass needed | Too large for the safe merge train | + +## Practical Merge Order From Here + +1. Get fresh approval on `#966`, `#967`, `#968`, `#970`, `#986`, `#987`. +2. Let auto-merge land them as checks clear. +3. Re-run full staging CI after each actual merge to `staging`. +4. Treat `#934`, `#737`, `#831`, and `#616` as separate workstreams, not part of the current safe merge train. + +## Key Findings + +- The repo token used here cannot bypass the required-review ruleset, even with `gh pr merge --admin`. +- Direct pushes to `staging` are blocked by repo rules (`GH013`). +- The replacement-PR path is the workable route for approved dirty PRs. +- `#934` is not merely stale. Its workflow change is unsafe here because this repository's default branch is `staging`, not `main`. diff --git a/docs/plans/2026-03-22-engine-v2-acceptance-criteria.md b/docs/plans/2026-03-22-engine-v2-acceptance-criteria.md new file mode 100644 index 00000000..e6c02024 --- /dev/null +++ b/docs/plans/2026-03-22-engine-v2-acceptance-criteria.md @@ -0,0 +1,184 @@ +# Engine v2 Acceptance Criteria + +**Date:** 2026-03-22 +**Status:** Active +**Author:** Zaki Manian +**Goal:** Define the merge bar for replacing the v1 agent loop with the v2 engine (`crates/ironclaw_engine/`). Phase 6 is not done until every criterion below is met. + +--- + +## Overview + +The v2 engine replaces ~10 v1 abstractions (Session, Job, Routine, Channel, Tool, Skill, Hook, Observer, Extension, LoopDelegate) with 5 primitives (Thread, Step, Capability, MemoryDoc, Project). Phases 1-5 are complete: types, execution loop, CodeAct/Monty, budget controls, and conversation surface. + +Phase 6 delivers the bridge adapters (`LlmBridgeAdapter`, `StoreBridgeAdapter`, `EffectBridgeAdapter`) that connect the engine to existing IronClaw infrastructure. The acceptance criteria below define what "ready to replace v1" means. Nothing merges to `staging` until all pass. + +--- + +## Acceptance Criteria + +### 1. Behavioral Equivalence + +Every observable behavior of v1 must be reproduced by v2 running through bridge adapters. + +| # | Criterion | Verification | +|---|-----------|-------------| +| 1.1 | All existing E2E test fixtures pass through `EngineV2Delegate` | `cargo test --features integration -p ironclaw -- engine_v2` and `cd tests/e2e && pytest` with `ENGINE_V2=true` | +| 1.2 | Tool dispatch produces identical outputs for identical inputs | Add property test: for each built-in tool, run same `(name, params)` through v1 `execute_tool_with_safety()` and v2 `EffectBridgeAdapter::execute_action()`, assert outputs match | +| 1.3 | Error handling is equivalent: no silent failures where v1 errors, no errors where v1 succeeds | Diff test: run full E2E trace fixtures through both paths, compare `LoopOutcome` variants. Specifically test: invalid tool name, malformed params, timeout, policy deny | +| 1.4 | Approval flows work identically | Test sequence: tool with `requires_approval` -> pause -> user approves -> resume -> completion. Must produce same SSE events (`approval_needed`, `approval_resolved`) | +| 1.5 | System commands (`/help`, `/model`, `/status`, `/skills`, `/job`) produce equivalent responses | Command parity test: submit each system command through v2 conversation surface, compare output structure | +| 1.6 | Compaction produces equivalent context reduction | Run a 50-turn conversation through both engines, trigger compaction, compare resulting context window token count (must be within 5%) | + +**Blocking:** 1.1, 1.2, 1.3, 1.4 are hard blockers. 1.5 and 1.6 may be deferred to Phase 7 with written justification. + +### 2. Performance + +No performance regressions. Improvements expected from context-as-variables but not required. + +| # | Criterion | Target | Verification | +|---|-----------|--------|-------------| +| 2.1 | P50 step latency | Within +10% of v1 | Benchmark harness: `cargo bench -p ironclaw --bench step_latency` with mock LLM (fixed 50ms response). Run 1000 steps, compare distributions. Harness must test both engines in the same binary. | +| 2.2 | P95 step latency | Within +10% of v1 | Same harness as 2.1 | +| 2.3 | P99 step latency | Within +15% of v1 | Same harness as 2.1 (wider margin for tail latency) | +| 2.4 | Monty VM startup | < 1ms (verify the 0.06ms claim) | Dedicated microbenchmark: `cargo bench -p ironclaw_engine --bench monty_startup`. Time `MontyVm::new()` over 10,000 iterations, report P50/P99. Must include independent measurement, not self-reported. | +| 2.5 | Token efficiency | Neutral or improved | Measure total tokens (prompt + completion) for the same 10-turn conversation fixture through both engines. v2 must not use more tokens than v1. Context-as-variables should reduce prompt tokens by 10-30% on conversations with tool output > 4KB. | +| 2.6 | Memory per thread | No regression | Measure RSS delta when spawning 100 threads with mock LLM. v2 must not exceed v1 by more than 10%. | + +**Blocking:** 2.1, 2.2, 2.3 are hard blockers. 2.4, 2.5, 2.6 are soft blockers (documented regressions acceptable with mitigation plan). + +### 3. Safety and Security + +The engine itself contains no safety logic by design. Safety is enforced at the bridge boundary (`EffectBridgeAdapter`). This must be airtight. + +| # | Criterion | Verification | +|---|-----------|-------------| +| 3.1 | `SafetyLayer` (prompt injection, leak detection, content validation) is applied on every action execution through `EffectBridgeAdapter` | Unit test: mock `EffectExecutor` that logs calls, verify `SafetyLayer::validate_tool_input()` and `SafetyLayer::sanitize_tool_output()` are called for every `execute_action()` invocation. No code path bypasses this. | +| 3.2 | Policy engine enforces `Deny > RequireApproval > Allow` with zero bypasses | Test matrix: for each `EffectType` variant (ReadLocal, ReadExternal, WriteLocal, WriteExternal, CredentialedNetwork, Compute, Financial), create conflicting rules and verify Deny always wins, RequireApproval wins over Allow. Cover the case where a single action triggers multiple effect types. | +| 3.3 | Thread tree is acyclic with bounded depth | `ThreadTree::attach()` must reject cycles (test: A->B->C->A). `ThreadConfig::max_depth` must be enforced (test: exceed depth limit, verify `ThreadError::DepthExceeded`). Default max depth: 8. | +| 3.4 | Capability leases are checked before every action execution | Audit `ExecutionLoop::run()` and `execute_action_calls()`: no path from LLM response to `EffectExecutor::execute_action()` that skips `LeaseManager::check_lease()`. Verify with test: expired lease -> action denied, revoked lease -> action denied, exhausted `max_uses` -> action denied. | +| 3.5 | Monty VM panics cannot crash the host | Test: inject Python code that triggers a Monty panic (e.g., stack overflow, infinite allocation). Verify the step completes with `StepStatus::Failed`, thread continues or fails gracefully, no process abort. Specifically test all resource limits: 30s timeout, 64MB memory, 1M allocations. | +| 3.6 | No new attack surfaces | Review checklist (manual, documented in PR): (a) lease forgery: `LeaseId` cannot be guessed or constructed outside `LeaseManager::grant()`, (b) policy bypass: no public method on `ExecutionLoop` that executes actions without policy check, (c) effect escalation: action's declared `EffectType` cannot be changed after capability registration, (d) cross-thread lease usage: lease bound to `thread_id` is enforced. | + +**Blocking:** All items are hard blockers. 3.6 is a manual review checklist that must be signed off in the merge PR. + +### 4. Persistence and Migration + +Production requires durable state. `InMemoryStore` is for tests only. + +| # | Criterion | Verification | +|---|-----------|-------------| +| 4.1 | `StoreBridgeAdapter` implements the full `Store` trait (18 methods) for both PostgreSQL and libSQL | Integration test per backend: create thread -> add steps -> append events -> save leases -> restart process -> load thread -> verify all data intact. Run with `cargo test --features integration` (postgres) and default (libSQL). | +| 4.2 | Database migrations create all required tables | Migration V14+ creates: `engine_threads`, `engine_steps`, `engine_events`, `engine_projects`, `engine_memory_docs`, `engine_capability_leases`. Test: run migrations on empty database, verify tables exist with correct schemas. Both backends. | +| 4.3 | Thread state survives process restart | Integration test: start thread -> execute 3 steps -> kill process -> restart -> resume thread -> verify step count is 3, thread state is correct, events are intact. | +| 4.4 | In-flight v1 sessions continue working when v2 is enabled | Test: create v1 session with active thread -> enable `ENGINE_V2=true` -> new messages on the existing session use v1 path (not v2). Only new threads use v2. Verify with assertion on delegate type. | +| 4.5 | Data migration path is documented | `docs/plans/` must contain a migration guide covering: (a) which v1 tables map to which v2 tables, (b) whether historical data is migrated or v2 starts fresh, (c) rollback procedure if migration fails. | + +**Blocking:** 4.1, 4.2, 4.3, 4.4 are hard blockers. 4.5 is required documentation but may ship as a separate document in the same milestone. + +### 5. Observability + +The engine must emit enough telemetry to debug production issues without attaching a debugger. + +| # | Criterion | Verification | +|---|-----------|-------------| +| 5.1 | Step execution duration is recorded | Each `Step` must have `started_at` and `completed_at` timestamps. Verify via unit test: execute a step, assert both fields are set and `completed_at > started_at`. | +| 5.2 | Token usage is tracked per step and per thread | `Step::token_usage` must be populated from `LlmOutput`. Thread-level aggregation: `thread.steps.iter().map(|s| s.token_usage).sum()`. Verify: run 5 steps with known token counts from mock LLM, assert thread total matches. | +| 5.3 | Policy decision counters | `PolicyEngine` must expose counts of `Allow`, `Deny`, and `RequireApproval` decisions. Verify: run 10 actions with mixed policies, assert counters match expected values. These must be queryable (not just logged). | +| 5.4 | Active lease gauge | `LeaseManager` must expose current active lease count. Verify: grant 5 leases, revoke 2, expire 1, assert gauge reads 2. | +| 5.5 | Event sourcing query performance | `Store::load_events(thread_id)` must return within 100ms for a thread with 1000 events. Benchmark test with both backends. | +| 5.6 | Structured logging for execution loop | Each step must emit `tracing` spans with: `thread_id`, `step_index`, `execution_tier`, `duration_ms`, `token_count`. Verify by capturing tracing output in test and asserting field presence. | + +**Blocking:** 5.1, 5.2, 5.6 are hard blockers. 5.3, 5.4, 5.5 are soft blockers (must be filed as issues if deferred). + +### 6. Rollout Strategy + +No big-bang cutover. Gradual rollout with rollback capability. + +| # | Criterion | Verification | +|---|-----------|-------------| +| 6.1 | Feature flag `ENGINE_V2` controls engine selection | When `ENGINE_V2=true`: new threads use `EngineV2Delegate`. When `ENGINE_V2=false` (default): all threads use v1. Verify: start with flag off, create thread (v1), set flag on, create thread (v2), both work. | +| 6.2 | Existing threads continue on their original engine | A thread started on v1 must remain on v1 even when `ENGINE_V2=true`. Thread metadata must record which engine version created it. Verify: create v1 thread, enable v2, send message to v1 thread, assert v1 delegate is used. | +| 6.3 | Rollback path: disable flag, no data loss | Enable v2, create threads, disable v2. v2 threads become read-only (no new messages accepted) but their data persists. New threads use v1. No data corruption in either direction. | +| 6.4 | Percentage-based rollout support | `ENGINE_V2_ROLLOUT_PERCENT=10` routes 10% of new threads to v2 (hash of thread_id mod 100). This enables canary deployment. Verify: create 100 threads with rollout at 10%, assert approximately 10 use v2. | +| 6.5 | Canary validation period | Before full rollout, v2 must run on >= 10% of new threads for at least 1 week with no P0/P1 incidents. This is a process gate, not a code test. Document the canary checklist in the rollout runbook. | + +**Blocking:** 6.1, 6.2, 6.3 are hard blockers. 6.4 is a soft blocker. 6.5 is a process requirement. + +--- + +## Non-Goals for Phase 6 + +These are explicitly out of scope. Do not implement them as part of Phase 6 acceptance. + +- **Full reflection pipeline** (Phase 7) -- thread post-mortem analysis and lesson extraction +- **WASM/Docker thread isolation** (Phase 8) -- running threads in sandboxed containers +- **Performance optimization beyond parity** -- v2 should match v1, not beat it (improvements are welcome but not required) +- **Mission system** -- `Mission` type is defined but not wired up +- **Provenance tracking / taint analysis** -- structs exist but enforcement is Phase 7 +- **Two-phase commit for Financial effects** -- design is documented in Phase 6 spec, but implementation may defer to Phase 7 if no Financial-effect tools exist yet +- **Dual model routing** -- `LlmBridgeAdapter` should support it structurally but it is not a Phase 6 acceptance criterion + +--- + +## Verification Plan + +### Automated Tests (CI-blocking) + +```bash +# 1. Engine unit tests (existing) +cargo test -p ironclaw_engine + +# 2. Bridge adapter tests (new) +cargo test -p ironclaw -- bridge + +# 3. Integration tests with both backends +cargo test --features integration -- engine_v2 +cargo test -- engine_v2 # libSQL path + +# 4. E2E tests with v2 engine +cd tests/e2e && ENGINE_V2=true pytest + +# 5. Behavioral equivalence diff tests +cargo test -- behavioral_equivalence + +# 6. Performance benchmarks (CI-reported, not CI-blocking) +cargo bench -p ironclaw --bench step_latency +cargo bench -p ironclaw_engine --bench monty_startup +``` + +### Manual Review (PR-blocking) + +- [ ] Security audit checklist (criterion 3.6) signed off by reviewer +- [ ] Migration documentation (criterion 4.5) exists and reviewed +- [ ] Canary runbook (criterion 6.5) exists + +### Test Fixtures Required + +| Fixture | Purpose | Location | +|---------|---------|----------| +| `trace_basic_conversation.json` | Multi-turn chat with tool calls | `tests/fixtures/engine_v2/` | +| `trace_approval_flow.json` | Tool requiring approval -> approve -> complete | `tests/fixtures/engine_v2/` | +| `trace_error_handling.json` | Invalid tool, malformed params, timeout | `tests/fixtures/engine_v2/` | +| `trace_compaction.json` | 50-turn conversation triggering compaction | `tests/fixtures/engine_v2/` | +| `trace_codeact.json` | CodeAct/Monty execution with tool dispatch | `tests/fixtures/engine_v2/` | + +### Benchmark Harness Requirements + +The step latency benchmark must: +1. Use the same mock LLM (fixed response, configurable latency) for both engines +2. Run in the same binary to eliminate process-level variance +3. Report P50/P95/P99 with confidence intervals +4. Run at least 1000 iterations per engine +5. Warm up with 100 iterations before measurement +6. Be added to CI as a reporting job (not a gate) with regression alerts at +15% + +### Definition of Done + +Phase 6 is complete when: +1. All hard-blocker criteria pass in CI +2. All soft-blocker criteria either pass or have filed issues with mitigation plans +3. Security review checklist is signed off +4. Migration documentation exists +5. Canary runbook exists +6. PR is approved by at least one reviewer who has read this document diff --git a/scripts/monitor-prs.sh b/scripts/monitor-prs.sh new file mode 100755 index 00000000..f57d8d5f --- /dev/null +++ b/scripts/monitor-prs.sh @@ -0,0 +1,139 @@ +#!/usr/bin/env bash +set -euo pipefail + +usage() { + cat <<'EOF' +Usage: scripts/monitor-prs.sh [--repo owner/name] [--author login] + +Shows open PRs for the author with: +- review decision +- latest review summary +- failing or pending checks + +Defaults: +- repo: current gitHub repo from `gh repo view` +- author: currently authenticated GitHub user from `gh api user` +EOF +} + +repo="" +author="" + +while [ $# -gt 0 ]; do + case "$1" in + --repo) + repo="${2:-}" + shift 2 + ;; + --author) + author="${2:-}" + shift 2 + ;; + -h|--help) + usage + exit 0 + ;; + *) + echo "Unknown argument: $1" >&2 + usage >&2 + exit 1 + ;; + esac +done + +if ! command -v gh >/dev/null 2>&1; then + echo "gh CLI is required" >&2 + exit 1 +fi + +if ! command -v jq >/dev/null 2>&1; then + echo "jq is required" >&2 + exit 1 +fi + +if [ -z "$repo" ]; then + repo="$(gh repo view --json nameWithOwner -q .nameWithOwner)" +fi + +if [ -z "$author" ]; then + author="$(gh api user -q .login)" +fi + +json_fields="number,title,url,headRefName,reviewDecision,latestReviews,statusCheckRollup" +prs="$(gh pr list --repo "$repo" --author "$author" --state open --limit 100 --json "$json_fields")" + +count="$(printf '%s' "$prs" | jq 'length')" +echo "Open PRs for $author in $repo: $count" +echo + +if [ "$count" -eq 0 ]; then + exit 0 +fi + +printf '%s' "$prs" | jq -r ' + def check_name: + .name // .context // .workflowName // "unknown-check"; + + def failing_checks: + [.statusCheckRollup[]? + | select(.status == "COMPLETED" and (.conclusion // .state // "") != "SUCCESS") + | { + name: check_name, + workflow: (.workflowName // ""), + url: (.detailsUrl // "") + }]; + + def pending_checks: + [.statusCheckRollup[]? + | select(.status != "COMPLETED") + | { + name: check_name, + workflow: (.workflowName // ""), + url: (.detailsUrl // "") + }]; + + .[] + | . as $pr + | failing_checks as $failing + | pending_checks as $pending + | [ + ("#" + (.number | tostring) + " " + .title), + (" Branch: " + .headRefName), + (" URL: " + .url), + (" Review: " + (.reviewDecision // "UNKNOWN")), + ( + if (.latestReviews | length) > 0 then + " Latest review: " + + .latestReviews[0].state + + " by " + + .latestReviews[0].author.login + + " at " + + .latestReviews[0].submittedAt + else + " Latest review: none" + end + ), + (" Checks: " + ($failing | length | tostring) + " failing, " + + ($pending | length | tostring) + " pending"), + ( + if ($failing | length) > 0 then + ($failing[] | " FAIL: " + .name + + (if .workflow != "" then " [" + .workflow + "]" else "" end) + + (if .url != "" then " -> " + .url else "" end)) + else + " FAIL: none" + end + ), + ( + if ($pending | length) > 0 then + ($pending[] | " PENDING: " + .name + + (if .workflow != "" then " [" + .workflow + "]" else "" end) + + (if .url != "" then " -> " + .url else "" end)) + else + " PENDING: none" + end + ) + ] + | .[] + , "" +' diff --git a/src/agent/agent_loop.rs b/src/agent/agent_loop.rs index 5cbd8166..5da7320b 100644 --- a/src/agent/agent_loop.rs +++ b/src/agent/agent_loop.rs @@ -731,7 +731,7 @@ impl Agent { { use crate::agent::session::Thread; let mut sess = session.lock().await; - let thread = Thread::with_id(id, sess.id); + let thread = Thread::with_id(id, sess.id, Some("gateway")); sess.active_thread = Some(id); sess.threads.entry(id).or_insert(thread); } diff --git a/src/agent/compaction.rs b/src/agent/compaction.rs index 30bb2b6c..c69f9608 100644 --- a/src/agent/compaction.rs +++ b/src/agent/compaction.rs @@ -319,7 +319,7 @@ mod tests { #[test] fn test_format_turns() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Hello"); thread.complete_turn("Hi there"); thread.start_turn("How are you?"); @@ -351,7 +351,7 @@ mod tests { /// Helper: build a thread with `n` completed turns. /// Turn `i` has user_input "msg-{i}" and response "resp-{i}". fn make_thread(n: usize) -> Thread { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); for i in 0..n { thread.start_turn(format!("msg-{}", i)); thread.complete_turn(format!("resp-{}", i)); @@ -457,7 +457,7 @@ mod tests { async fn test_compact_truncate_empty_turns() { let llm = Arc::new(StubLlm::new("unused")); let compactor = make_compactor(llm); - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); assert!(thread.turns.is_empty()); let result = compactor @@ -698,7 +698,7 @@ mod tests { #[test] fn test_format_turns_for_storage_with_tool_calls() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Search for X"); // Record a tool call on the current turn if let Some(turn) = thread.turns.last_mut() { @@ -719,7 +719,7 @@ mod tests { #[test] fn test_format_turns_for_storage_incomplete_turn() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("In progress message"); // Don't complete the turn diff --git a/src/agent/dispatcher.rs b/src/agent/dispatcher.rs index 03548219..9d85238c 100644 --- a/src/agent/dispatcher.rs +++ b/src/agent/dispatcher.rs @@ -2140,7 +2140,7 @@ mod tests { // Initialize a thread in the session so the loop can record tool calls. let thread_id = { let mut sess = session.lock().await; - sess.create_thread().id + sess.create_thread("test").id }; let message = IncomingMessage::new("test", "test-user", "do something"); @@ -2245,7 +2245,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("test-user"))); let thread_id = { let mut sess = session.lock().await; - sess.create_thread().id + sess.create_thread("test").id }; let message = IncomingMessage::new("test", "test-user", "keep calling tools"); diff --git a/src/agent/session.rs b/src/agent/session.rs index 45594922..0eda753b 100644 --- a/src/agent/session.rs +++ b/src/agent/session.rs @@ -68,8 +68,8 @@ impl Session { } /// Create a new thread in this session. - pub fn create_thread(&mut self) -> &mut Thread { - let thread = Thread::new(self.id); + pub fn create_thread(&mut self, channel: &str) -> &mut Thread { + let thread = Thread::new(self.id, Some(channel)); let thread_id = thread.id; self.active_thread = Some(thread_id); self.last_active_at = Utc::now(); @@ -87,9 +87,9 @@ impl Session { } /// Get or create the active thread. - pub fn get_or_create_thread(&mut self) -> &mut Thread { + pub fn get_or_create_thread(&mut self, channel: &str) -> &mut Thread { match self.active_thread { - None => self.create_thread(), + None => self.create_thread(channel), Some(id) => { if self.threads.contains_key(&id) { // Entry existence confirmed by contains_key above. @@ -100,7 +100,7 @@ impl Session { } else { // Stale active_thread ID: create a new thread, which // updates self.active_thread to the new thread's ID. - self.create_thread() + self.create_thread(channel) } } } @@ -225,6 +225,9 @@ pub struct Thread { /// Messages queued while the thread was processing a turn. #[serde(default, skip_serializing_if = "VecDeque::is_empty")] pub pending_messages: VecDeque, + /// Channel that created this thread (for approval authorization). + #[serde(default)] + pub source_channel: Option, } /// Maximum number of messages that can be queued while a thread is processing. @@ -235,7 +238,7 @@ pub const MAX_PENDING_MESSAGES: usize = 10; impl Thread { /// Create a new thread. - pub fn new(session_id: Uuid) -> Self { + pub fn new(session_id: Uuid, source_channel: Option<&str>) -> Self { let now = Utc::now(); Self { id: Uuid::new_v4(), @@ -248,11 +251,12 @@ impl Thread { pending_approval: None, pending_auth: None, pending_messages: VecDeque::new(), + source_channel: source_channel.map(String::from), } } /// Create a thread with a specific ID (for DB hydration). - pub fn with_id(id: Uuid, session_id: Uuid) -> Self { + pub fn with_id(id: Uuid, session_id: Uuid, source_channel: Option<&str>) -> Self { let now = Utc::now(); Self { id, @@ -265,6 +269,7 @@ impl Thread { pending_approval: None, pending_auth: None, pending_messages: VecDeque::new(), + source_channel: source_channel.map(String::from), } } @@ -696,13 +701,13 @@ mod tests { let mut session = Session::new("user-123"); assert!(session.active_thread.is_none()); - session.create_thread(); + session.create_thread("test"); assert!(session.active_thread.is_some()); } #[test] fn test_thread_turns() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Hello"); assert_eq!(thread.state, ThreadState::Processing); @@ -715,7 +720,7 @@ mod tests { #[test] fn test_thread_messages() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("First message"); thread.complete_turn("First response"); @@ -738,7 +743,7 @@ mod tests { #[test] fn test_restore_from_messages() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // First add some turns thread.start_turn("Original message"); @@ -764,7 +769,7 @@ mod tests { #[test] fn test_restore_from_messages_incomplete_turn() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Messages with incomplete last turn (no assistant response) let messages = vec![ @@ -783,7 +788,7 @@ mod tests { #[test] fn test_enter_auth_mode() { let before = Utc::now(); - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); assert!(thread.pending_auth.is_none()); thread.enter_auth_mode("telegram".to_string()); @@ -796,7 +801,7 @@ mod tests { #[test] fn test_take_pending_auth() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.enter_auth_mode("notion".to_string()); let pending = thread.take_pending_auth(); @@ -811,7 +816,7 @@ mod tests { #[test] fn test_pending_auth_serialization() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.enter_auth_mode("openai".to_string()); let json = serde_json::to_string(&thread).expect("should serialize"); @@ -841,7 +846,7 @@ mod tests { #[test] fn test_pending_auth_default_none() { // Deserialization of old data without pending_auth should default to None - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.pending_auth = None; let json = serde_json::to_string(&thread).expect("serialize"); @@ -855,7 +860,7 @@ mod tests { fn test_thread_with_id() { let specific_id = Uuid::new_v4(); let session_id = Uuid::new_v4(); - let thread = Thread::with_id(specific_id, session_id); + let thread = Thread::with_id(specific_id, session_id, None); assert_eq!(thread.id, specific_id); assert_eq!(thread.session_id, session_id); @@ -867,7 +872,7 @@ mod tests { fn test_thread_with_id_restore_messages() { let thread_id = Uuid::new_v4(); let session_id = Uuid::new_v4(); - let mut thread = Thread::with_id(thread_id, session_id); + let mut thread = Thread::with_id(thread_id, session_id, None); let messages = vec![ ChatMessage::user("Hello from DB"), @@ -886,7 +891,7 @@ mod tests { #[test] fn test_restore_from_messages_empty() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Add a turn first, then restore with empty vec thread.start_turn("hello"); @@ -902,7 +907,7 @@ mod tests { #[test] fn test_restore_from_messages_only_assistant_messages() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Only assistant messages (no user messages to anchor turns) let messages = vec![ @@ -919,7 +924,7 @@ mod tests { #[test] fn test_restore_from_messages_multiple_user_messages_in_a_row() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Two user messages with no assistant response between them let messages = vec![ @@ -946,8 +951,8 @@ mod tests { fn test_thread_switch() { let mut session = Session::new("user-1"); - let t1_id = session.create_thread().id; - let t2_id = session.create_thread().id; + let t1_id = session.create_thread("test").id; + let t2_id = session.create_thread("test").id; // After creating two threads, active should be the last one assert_eq!(session.active_thread, Some(t2_id)); @@ -967,8 +972,8 @@ mod tests { fn test_get_or_create_thread_idempotent() { let mut session = Session::new("user-1"); - let tid1 = session.get_or_create_thread().id; - let tid2 = session.get_or_create_thread().id; + let tid1 = session.get_or_create_thread("test").id; + let tid2 = session.get_or_create_thread("test").id; // Should return the same thread (not create a new one each time) assert_eq!(tid1, tid2); @@ -977,7 +982,7 @@ mod tests { #[test] fn test_truncate_turns() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); for i in 0..5 { thread.start_turn(format!("msg-{}", i)); @@ -1001,7 +1006,7 @@ mod tests { #[test] fn test_truncate_turns_noop_when_fewer() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("only one"); thread.complete_turn("response"); @@ -1013,7 +1018,7 @@ mod tests { #[test] fn test_thread_interrupt_and_resume() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("do something"); assert_eq!(thread.state, ThreadState::Processing); @@ -1031,7 +1036,7 @@ mod tests { #[test] fn test_resume_only_from_interrupted() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Idle thread: resume should be a no-op assert_eq!(thread.state, ThreadState::Idle); @@ -1047,7 +1052,7 @@ mod tests { #[test] fn test_turn_fail() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("risky operation"); thread.fail_turn("connection timed out"); @@ -1063,7 +1068,7 @@ mod tests { #[test] fn test_messages_with_incomplete_last_turn() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("first"); thread.complete_turn("first reply"); @@ -1079,7 +1084,7 @@ mod tests { #[test] fn test_thread_serialization_round_trip() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("hello"); thread.complete_turn("world"); @@ -1097,7 +1102,7 @@ mod tests { #[test] fn test_session_serialization_round_trip() { let mut session = Session::new("user-ser"); - session.create_thread(); + session.create_thread("test"); session.auto_approve_tool("echo"); let json = serde_json::to_string(&session).unwrap(); @@ -1135,7 +1140,7 @@ mod tests { #[test] fn test_turn_number_increments() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Before any turns, turn_number() is 1 (1-indexed for display) assert_eq!(thread.turn_number(), 1); @@ -1150,7 +1155,7 @@ mod tests { #[test] fn test_complete_turn_on_empty_thread() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Completing a turn when there are no turns should be a safe no-op thread.complete_turn("phantom response"); @@ -1160,7 +1165,7 @@ mod tests { #[test] fn test_fail_turn_on_empty_thread() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Failing a turn when there are no turns should be a safe no-op thread.fail_turn("phantom error"); @@ -1170,7 +1175,7 @@ mod tests { #[test] fn test_pending_approval_flow() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); let approval = PendingApproval { request_id: Uuid::new_v4(), @@ -1197,7 +1202,7 @@ mod tests { #[test] fn test_clear_pending_approval() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); let approval = PendingApproval { request_id: Uuid::new_v4(), @@ -1226,7 +1231,7 @@ mod tests { assert!(session.active_thread().is_none()); assert!(session.active_thread_mut().is_none()); - let tid = session.create_thread().id; + let tid = session.create_thread("test").id; assert!(session.active_thread().is_some()); assert_eq!(session.active_thread().unwrap().id, tid); @@ -1243,7 +1248,7 @@ mod tests { #[test] fn test_messages_includes_tool_calls() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Search for X"); { @@ -1275,7 +1280,7 @@ mod tests { #[test] fn test_messages_multiple_tool_calls_per_turn() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Do two things"); { @@ -1302,7 +1307,7 @@ mod tests { #[test] fn test_restore_from_messages_with_tool_calls() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); // Build a message sequence with tool calls let tc = ToolCall { @@ -1333,7 +1338,7 @@ mod tests { #[test] fn test_restore_from_messages_with_tool_error() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); let tc = ToolCall { id: "call_0".to_string(), @@ -1363,7 +1368,7 @@ mod tests { fn test_messages_round_trip_with_tools() { // Build a thread with tool calls, get messages(), restore, get messages() again // The two message sequences should be equivalent. - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Do search"); { @@ -1376,7 +1381,7 @@ mod tests { let messages_original = thread.messages(); // Restore into a new thread - let mut thread2 = Thread::new(Uuid::new_v4()); + let mut thread2 = Thread::new(Uuid::new_v4(), None); thread2.restore_from_messages(messages_original.clone()); let messages_restored = thread2.messages(); @@ -1398,7 +1403,7 @@ mod tests { #[test] fn test_restore_multi_stage_tool_calls() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); let tc1 = ToolCall { id: "call_a".to_string(), @@ -1439,7 +1444,7 @@ mod tests { #[test] fn test_messages_truncates_large_tool_results() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("Read big file"); { @@ -1462,13 +1467,11 @@ mod tests { #[test] fn test_thread_message_queue() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Queue is initially empty assert!(thread.pending_messages.is_empty()); assert!(thread.take_pending_message().is_none()); - // Queue messages and verify FIFO ordering assert!(thread.queue_message("first".to_string())); assert!(thread.queue_message("second".to_string())); assert!(thread.queue_message("third".to_string())); @@ -1479,17 +1482,14 @@ mod tests { assert_eq!(thread.take_pending_message(), Some("third".to_string())); assert!(thread.take_pending_message().is_none()); - // Fill to capacity — all 10 should succeed for i in 0..MAX_PENDING_MESSAGES { assert!(thread.queue_message(format!("msg-{}", i))); } assert_eq!(thread.pending_messages.len(), MAX_PENDING_MESSAGES); - // 11th message rejected by queue_message itself assert!(!thread.queue_message("overflow".to_string())); assert_eq!(thread.pending_messages.len(), MAX_PENDING_MESSAGES); - // Drain and verify order for i in 0..MAX_PENDING_MESSAGES { assert_eq!(thread.take_pending_message(), Some(format!("msg-{}", i))); } @@ -1498,13 +1498,11 @@ mod tests { #[test] fn test_thread_message_queue_serialization() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Empty queue should not appear in serialization (skip_serializing_if) let json = serde_json::to_string(&thread).unwrap(); assert!(!json.contains("pending_messages")); - // Non-empty queue should serialize and deserialize thread.queue_message("queued msg".to_string()); let json = serde_json::to_string(&thread).unwrap(); assert!(json.contains("pending_messages")); @@ -1517,11 +1515,9 @@ mod tests { #[test] fn test_thread_message_queue_default_on_old_data() { - // Deserialization of old data without pending_messages should default to empty - let thread = Thread::new(Uuid::new_v4()); + let thread = Thread::new(Uuid::new_v4(), None); let json = serde_json::to_string(&thread).unwrap(); - // The field is absent (skip_serializing_if), simulating old data assert!(!json.contains("pending_messages")); let restored: Thread = serde_json::from_str(&json).unwrap(); assert!(restored.pending_messages.is_empty()); @@ -1529,18 +1525,15 @@ mod tests { #[test] fn test_interrupt_clears_pending_messages() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Start a turn so there's something to interrupt thread.start_turn("initial input"); - // Queue several messages while "processing" thread.queue_message("queued-1".to_string()); thread.queue_message("queued-2".to_string()); thread.queue_message("queued-3".to_string()); assert_eq!(thread.pending_messages.len(), 3); - // Interrupt should clear the queue thread.interrupt(); assert!(thread.pending_messages.is_empty()); assert_eq!(thread.state, ThreadState::Interrupted); @@ -1548,27 +1541,22 @@ mod tests { #[test] fn test_thread_state_idle_after_full_drain() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Simulate a full drain cycle: start turn, queue messages, complete turn, - // then drain all queued messages as a single merged turn (#259). thread.start_turn("turn 1"); assert_eq!(thread.state, ThreadState::Processing); thread.queue_message("queued-a".to_string()); thread.queue_message("queued-b".to_string()); - // Complete the turn (simulates process_user_input finishing) thread.complete_turn("response 1"); assert_eq!(thread.state, ThreadState::Idle); - // Drain: merge all queued messages and process as a single turn let merged = thread.drain_pending_messages().unwrap(); assert_eq!(merged, "queued-a\nqueued-b"); thread.start_turn(&merged); thread.complete_turn("response for merged"); - // Queue is fully drained, thread is idle assert!(thread.drain_pending_messages().is_none()); assert!(thread.pending_messages.is_empty()); assert_eq!(thread.state, ThreadState::Idle); @@ -1576,12 +1564,10 @@ mod tests { #[test] fn test_drain_pending_messages_merges_with_newlines() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Empty queue returns None assert!(thread.drain_pending_messages().is_none()); - // Single message returned as-is (no trailing newline) thread.queue_message("only one".to_string()); assert_eq!( thread.drain_pending_messages(), @@ -1589,7 +1575,6 @@ mod tests { ); assert!(thread.pending_messages.is_empty()); - // Multiple messages joined with newlines thread.queue_message("hey".to_string()); thread.queue_message("can you check the server".to_string()); thread.queue_message("it started 10 min ago".to_string()); @@ -1599,25 +1584,62 @@ mod tests { ); assert!(thread.pending_messages.is_empty()); - // Queue is empty after drain assert!(thread.drain_pending_messages().is_none()); } #[test] fn test_requeue_drained_preserves_content_at_front() { - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); - // Re-queue into empty queue thread.requeue_drained("failed batch".to_string()); assert_eq!(thread.pending_messages.len(), 1); assert_eq!(thread.pending_messages[0], "failed batch"); - // New messages go behind the re-queued content thread.queue_message("new msg".to_string()); assert_eq!(thread.pending_messages.len(), 2); - // Drain should return re-queued content first (front of queue) let merged = thread.drain_pending_messages().unwrap(); assert_eq!(merged, "failed batch\nnew msg"); } + + #[test] + fn test_thread_new_stores_source_channel() { + let thread = Thread::new(Uuid::new_v4(), Some("gateway")); + assert_eq!(thread.source_channel.as_deref(), Some("gateway")); + } + + #[test] + fn test_thread_new_none_channel() { + let thread = Thread::new(Uuid::new_v4(), None); + assert!(thread.source_channel.is_none()); + } + + #[test] + fn test_thread_with_id_stores_source_channel() { + let thread = Thread::with_id(Uuid::new_v4(), Uuid::new_v4(), Some("http")); + assert_eq!(thread.source_channel.as_deref(), Some("http")); + } + + #[test] + fn test_create_thread_sets_source_channel() { + let mut session = Session::new("user-chan"); + let thread_id = session.create_thread("gateway").id; + let thread = session.threads.get(&thread_id).unwrap(); + assert_eq!(thread.source_channel.as_deref(), Some("gateway")); + } + + #[test] + fn test_source_channel_serde_backcompat() { + let json = r#"{ + "id": "00000000-0000-0000-0000-000000000001", + "session_id": "00000000-0000-0000-0000-000000000002", + "state": "Idle", + "turns": [], + "created_at": "2025-01-01T00:00:00Z", + "updated_at": "2025-01-01T00:00:00Z", + "metadata": null + }"#; + let thread: Thread = serde_json::from_str(json).unwrap(); + assert!(thread.source_channel.is_none()); + } } diff --git a/src/agent/session_manager.rs b/src/agent/session_manager.rs index 3bf20697..0d04b5f6 100644 --- a/src/agent/session_manager.rs +++ b/src/agent/session_manager.rs @@ -167,7 +167,7 @@ impl SessionManager { // Create new thread (always create a new one for a new key) let thread_id = { let mut sess = session.lock().await; - let thread = sess.create_thread(); + let thread = sess.create_thread(channel); thread.id }; @@ -443,7 +443,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-hydrate"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(thread_id, sess.id); + let thread = Thread::with_id(thread_id, sess.id, None); sess.threads.insert(thread_id, thread); sess.active_thread = Some(thread_id); } @@ -567,7 +567,7 @@ mod tests { // Simulate hydration: create thread with a known UUID { let mut sess = session.lock().await; - let thread = Thread::with_id(known_uuid, session_id); + let thread = Thread::with_id(known_uuid, session_id, None); sess.threads.insert(known_uuid, thread); } @@ -594,7 +594,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-idem"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } @@ -623,7 +623,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-undo"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } @@ -647,7 +647,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-new"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } @@ -755,7 +755,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-cross"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } @@ -782,7 +782,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-cross"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } @@ -921,7 +921,7 @@ mod tests { let session = Arc::new(Mutex::new(Session::new("user-direct"))); { let mut sess = session.lock().await; - let thread = Thread::with_id(tid, sess.id); + let thread = Thread::with_id(tid, sess.id, None); sess.threads.insert(tid, thread); } { @@ -947,4 +947,23 @@ mod tests { "should have exactly 1 thread, not a duplicate" ); } + + #[tokio::test] + async fn test_thread_stores_source_channel() { + let manager = SessionManager::new(); + let (session, thread_id) = manager + .resolve_thread("user1", "gateway", Some("ext-1")) + .await; + let sess = session.lock().await; + let thread = sess.threads.get(&thread_id).unwrap(); + assert_eq!(thread.source_channel.as_deref(), Some("gateway")); + } + + #[tokio::test] + async fn test_different_channels_get_different_threads() { + let manager = SessionManager::new(); + let (_, tid1) = manager.resolve_thread("user1", "gateway", None).await; + let (_, tid2) = manager.resolve_thread("user1", "web", None).await; + assert_ne!(tid1, tid2); + } } diff --git a/src/agent/thread_ops.rs b/src/agent/thread_ops.rs index eec29099..c5983ce7 100644 --- a/src/agent/thread_ops.rs +++ b/src/agent/thread_ops.rs @@ -141,7 +141,8 @@ impl Agent { sess.id }; - let mut thread = crate::agent::session::Thread::with_id(thread_uuid, session_id); + let mut thread = + crate::agent::session::Thread::with_id(thread_uuid, session_id, Some(&message.channel)); if !chat_messages.is_empty() { thread.restore_from_messages(chat_messages); } @@ -954,6 +955,28 @@ impl Agent { approved: bool, always: bool, ) -> Result { + // Verify channel authorization: the approving channel must match the + // thread's source channel, OR be the web gateway (trusted approval UI). + { + let sess = session.lock().await; + if let Some(thread) = sess.threads.get(&thread_id) { + let authorized = thread.source_channel.as_ref().is_none_or(|src| { + src == &message.channel || message.channel == "web" + }); + if !authorized { + tracing::warn!( + %thread_id, + source_channel = ?thread.source_channel, + approval_channel = %message.channel, + "Blocked cross-channel approval attempt" + ); + return Ok(SubmissionResult::error( + "approval not authorized for this channel", + )); + } + } + } + // Get pending approval for this thread let pending = { let mut sess = session.lock().await; @@ -1744,7 +1767,7 @@ impl Agent { .get_or_create_session(&message.user_id) .await; let mut sess = session.lock().await; - let thread = sess.create_thread(); + let thread = sess.create_thread(&message.channel); let thread_id = thread.id; Ok(SubmissionResult::ok_with_message(format!( "New thread: {}", @@ -2035,7 +2058,7 @@ mod tests { let session_id = Uuid::new_v4(); let thread_id = Uuid::new_v4(); - let mut thread = Thread::with_id(thread_id, session_id); + let mut thread = Thread::with_id(thread_id, session_id, None); // Set thread to AwaitingApproval with a pending tool approval let pending = PendingApproval { @@ -2103,7 +2126,7 @@ mod tests { use crate::agent::session::{MAX_PENDING_MESSAGES, Thread, ThreadState}; use uuid::Uuid; - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("processing something"); assert_eq!(thread.state, ThreadState::Processing); @@ -2129,7 +2152,7 @@ mod tests { use crate::agent::session::{Thread, ThreadState}; use uuid::Uuid; - let mut thread = Thread::new(Uuid::new_v4()); + let mut thread = Thread::new(Uuid::new_v4(), None); thread.start_turn("processing"); thread.queue_message("pending-1".to_string()); @@ -2159,7 +2182,7 @@ mod tests { let thread_id = Uuid::new_v4(); let session_id = Uuid::new_v4(); - let mut thread = Thread::with_id(thread_id, session_id); + let mut thread = Thread::with_id(thread_id, session_id, None); thread.start_turn("working"); assert_eq!(thread.state, ThreadState::Processing); @@ -2186,7 +2209,7 @@ mod tests { let thread_id = Uuid::new_v4(); let session_id = Uuid::new_v4(); - let mut thread = Thread::with_id(thread_id, session_id); + let mut thread = Thread::with_id(thread_id, session_id, None); thread.start_turn("working"); assert_eq!(thread.state, ThreadState::Processing); diff --git a/src/channels/web/handlers/chat.rs b/src/channels/web/handlers/chat.rs index 5cb2b9ea..f91eaaf2 100644 --- a/src/channels/web/handlers/chat.rs +++ b/src/channels/web/handlers/chat.rs @@ -543,7 +543,7 @@ pub async fn chat_new_thread_handler( let session = session_manager.get_or_create_session(&state.user_id).await; let (thread_id, info) = { let mut sess = session.lock().await; - let thread = sess.create_thread(); + let thread = sess.create_thread("web"); let id = thread.id; let info = ThreadInfo { id: thread.id, diff --git a/src/channels/web/server.rs b/src/channels/web/server.rs index 7edaad67..81e54ae5 100644 --- a/src/channels/web/server.rs +++ b/src/channels/web/server.rs @@ -1658,7 +1658,7 @@ async fn chat_new_thread_handler( let session = session_manager.get_or_create_session(&state.user_id).await; let (thread_id, info) = { let mut sess = session.lock().await; - let thread = sess.create_thread(); + let thread = sess.create_thread("web"); let id = thread.id; let info = ThreadInfo { id: thread.id,