diff --git a/.env.example b/.env.example index 8fd44c5a..873931d7 100644 --- a/.env.example +++ b/.env.example @@ -4,7 +4,7 @@ DATABASE_POOL_SIZE=10 # LLM Provider # LLM_BACKEND=nearai # default -# Possible values: nearai, ollama, openai_compatible, openai, anthropic, tinfoil +# Possible values: nearai, ollama, openai_compatible, openai, anthropic, github_copilot, tinfoil, openai_codex # LLM_REQUEST_TIMEOUT_SECS=120 # Increase for local LLMs (Ollama, vLLM, LM Studio) # === Anthropic Direct === @@ -24,6 +24,17 @@ DATABASE_POOL_SIZE=10 # LLM_USE_CODEX_AUTH=true # CODEX_AUTH_PATH=~/.codex/auth.json +# === GitHub Copilot === +# Uses the OAuth token from your Copilot IDE sign-in (for example +# ~/.config/github-copilot/apps.json on Linux/macOS), or run `ironclaw onboard` +# and choose the GitHub device login flow. +# LLM_BACKEND=github_copilot +# GITHUB_COPILOT_TOKEN=gho_... +# GITHUB_COPILOT_MODEL=gpt-4o +# IronClaw injects standard VS Code Copilot headers automatically. +# Optional advanced headers for custom overrides: +# GITHUB_COPILOT_EXTRA_HEADERS=Copilot-Integration-Id:vscode-chat + # === NEAR AI (Chat Completions API) === # Two auth modes: # 1. Session token (default): Uses browser OAuth (GitHub/Google) on first run. @@ -31,7 +42,7 @@ DATABASE_POOL_SIZE=10 # Base URL defaults to https://private.near.ai # 2. API key: Set NEARAI_API_KEY to use API key auth from cloud.near.ai. # Base URL defaults to https://cloud-api.near.ai -NEARAI_MODEL=zai-org/GLM-5-FP8 +NEARAI_MODEL=Qwen/Qwen3.5-122B-A10B NEARAI_BASE_URL=https://private.near.ai NEARAI_AUTH_URL=https://private.near.ai # NEARAI_SESSION_TOKEN=sess_... # hosting providers: set this @@ -92,6 +103,13 @@ NEARAI_AUTH_URL=https://private.near.ai # long = 1-hour TTL, 2.0ร— (200%) write surcharge # ANTHROPIC_CACHE_RETENTION=short +# === OpenAI Codex (ChatGPT subscription, OAuth) === +# LLM_BACKEND=openai_codex +# OPENAI_CODEX_MODEL=gpt-5.3-codex # default +# OPENAI_CODEX_CLIENT_ID=app_EMoamEEZ73f0CkXaXp7hrann # override (rare) +# OPENAI_CODEX_AUTH_URL=https://auth.openai.com # override (rare) +# OPENAI_CODEX_API_URL=https://chatgpt.com/backend-api/codex # override (rare) + # For full provider setup guide see docs/LLM_PROVIDERS.md # Channel Configuration diff --git a/AGENTS.md b/AGENTS.md index 7be35afb..cc5e7cff 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,6 +1,94 @@ # Agent Rules -## Feature Parity Update Policy +## Purpose and Precedence +- `AGENTS.md` is the quick-start contract for coding agents. It is not the full architecture spec. +- Read the relevant subsystem spec before changing a complex area. When a repo spec exists, treat it as authoritative. +Start with these deeper docs as needed: +- `CLAUDE.md` +- `src/agent/CLAUDE.md` +- `src/channels/web/CLAUDE.md` +- `src/db/CLAUDE.md` +- `src/llm/CLAUDE.md` +- `src/setup/README.md` +- `src/tools/README.md` +- `src/workspace/README.md` +- `src/NETWORK_SECURITY.md` +- `tests/e2e/CLAUDE.md` + +## Architecture Mental Model + +- Channels normalize external input into `IncomingMessage`; `ChannelManager` merges all active channel streams. +- `Agent` owns session/thread/turn handling, submission parsing, the LLM/tool loop, approvals, routines, and background runtime behavior. +- `AppBuilder` is the composition root that wires database, secrets, LLMs, tools, workspace, extensions, skills, hooks, and cost controls before the agent starts. +- The web gateway is a browser-facing API/UI layered on top of the same agent/session/tool systems, not a separate product path. + +## Where to Work + +- Agent/runtime behavior: `src/agent/` +- Web gateway/API/SSE/WebSocket: `src/channels/web/` +- Persistence and DB abstractions: `src/db/` +- Setup/onboarding/configuration flow: `src/setup/` +- LLM providers and routing: `src/llm/` +- Workspace, memory, embeddings, search: `src/workspace/` +- Extensions, tools, channels, MCP, WASM: `src/extensions/`, `src/tools/`, `src/channels/` + +## Ownership and Composition Rules + +- Keep `src/main.rs` and `src/app.rs` orchestration-focused. Do not move module-owned logic into entrypoints. +- Module-specific initialization should live in the owning module behind a public factory/helper, not be reimplemented ad hoc. +- Keep feature-flag branching inside the module that owns the abstraction whenever possible. +- Prefer extending existing traits and registries over hardcoding one-off integration paths. + +## Repo-Wide Coding Rules + +- Avoid `.unwrap()` and `.expect()` in production; prefer proper error handling. They are fine in tests, and in production only for truly infallible invariants (e.g., literals/regexes) with a safety comment. +- Keep clippy clean with zero warnings. +- Prefer `crate::` imports for cross-module references. +- Use strong types and enums over stringly-typed control flow when the shape is known. + +## Database, Setup, and Config Rules + +- New persistence behavior must support both PostgreSQL and libSQL. +- Add new DB operations to the shared DB trait first, then implement both backends. +- Treat bootstrap config, DB-backed settings, and encrypted secrets as distinct layers; do not collapse them casually. +- If onboarding or setup behavior changes, update `src/setup/README.md` in the same branch. +- Do not break config precedence, bootstrap env loading, DB-backed config reload, or post-secrets LLM re-resolution. + +## Security and Runtime Invariants + +- Review any change touching listeners, routes, auth, secrets, sandboxing, approvals, or outbound HTTP with a security mindset. +- Do not weaken bearer-token auth, webhook auth, CORS/origin checks, body limits, rate limits, allowlists, or secret-handling guarantees. +- Treat Docker containers and external services as untrusted. +- Session/thread/turn state matters. Submission parsing happens before normal chat handling. +- Skills are selected deterministically. Tool approval and auth flows are special paths and must not be mixed into normal chat history carelessly. +- Persistent memory is the workspace system, not just transcript storage; preserve file-like semantics, chunking/search behavior, and identity/system-prompt loading. + +## Tools, Channels, and Extensions + +- Use a built-in Rust tool for core internal capabilities tightly coupled to the runtime. +- Use WASM tools or WASM channels for sandboxed extensions and plugin-style integrations. +- Use MCP for external server integrations when the capability belongs outside the main binary. +- Preserve extension lifecycle expectations: install, authenticate/configure, activate, remove. + +## Docs, Parity, and Testing + +- If behavior changes, update the relevant docs/specs in the same branch. - If you change implementation status for any feature tracked in `FEATURE_PARITY.md`, update that file in the same branch. - Do not open a PR that changes feature behavior without checking `FEATURE_PARITY.md` for needed status updates (`โŒ`, `๐Ÿšง`, `โœ…`, notes, and priorities). +- Add the narrowest tests that validate the change: unit tests for local logic, integration tests for runtime/DB/routing behavior, and E2E or trace coverage for gateway, approvals, extensions, or other user-visible flows. + +## Risk and Change Discipline + +- Keep changes scoped; avoid broad refactors unless the task truly requires them. +- Security, database schema, runtime, worker, CI, and secrets changes are high-risk. Call out rollback risks, compatibility concerns, and hidden side effects. +- Preserve existing defaults unless the task explicitly changes them. +- Avoid unrelated file churn and generated-file edits unless required. +- Respect a dirty worktree and never revert user changes you did not make. + +## Before Finishing + +- Confirm whether behavior changes require updates to `FEATURE_PARITY.md`, specs, API docs, or `CHANGELOG.md`. +- Run the most targeted tests/checks that cover the change. +- Re-check security-sensitive paths when touching auth, secrets, network listeners, sandboxing, or approvals. +- Keep the final diff scoped to the task. diff --git a/CLAUDE.md b/CLAUDE.md index d47292e1..e2d84c1e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -158,6 +158,8 @@ src/ โ”‚ โ”œโ”€โ”€ secrets/ # Secrets management (AES-256-GCM, OS keychain for master key) โ”‚ +โ”œโ”€โ”€ profile.rs # Psychographic profile types, 9-dimension analysis framework +โ”‚ โ”œโ”€โ”€ setup/ # 7-step onboarding wizard โ€” see src/setup/README.md โ”‚ โ”œโ”€โ”€ skills/ # SKILL.md prompt extension system โ€” see .claude/rules/skills.md diff --git a/Cargo.lock b/Cargo.lock index 2c5547e0..151edd69 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2339,7 +2339,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -5575,7 +5575,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -5624,7 +5624,7 @@ dependencies = [ "once_cell", "ring", "rustls-pki-types", - "rustls-webpki 0.103.9", + "rustls-webpki 0.103.10", "subtle", "zeroize", ] @@ -5696,9 +5696,9 @@ dependencies = [ [[package]] name = "rustls-webpki" -version = "0.103.9" +version = "0.103.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7df23109aa6c1567d1c575b9952556388da57401e4ace1d15f79eedad0d8f53" +checksum = "df33b2b81ac578cabaf06b89b0631153a3f416b0a886e8a7a1707fb51abbd1ef" dependencies = [ "aws-lc-rs", "ring", @@ -6479,10 +6479,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] diff --git a/FEATURE_PARITY.md b/FEATURE_PARITY.md index d932e9ad..b7cb9103 100644 --- a/FEATURE_PARITY.md +++ b/FEATURE_PARITY.md @@ -247,6 +247,7 @@ This document tracks feature parity between IronClaw (Rust implementation) and O | OpenRouter | โœ… | โœ… | - | Via OpenAI-compatible provider (RigAdapter) | | Tinfoil | โŒ | โœ… | - | Private inference provider (IronClaw-only) | | OpenAI-compatible | โŒ | โœ… | - | Generic OpenAI-compatible endpoint (RigAdapter) | +| GitHub Copilot | โœ… | โœ… | - | Dedicated provider with OAuth token exchange (`GithubCopilotProvider`) | | Ollama (local) | โœ… | โœ… | - | via `rig::providers::ollama` (full support) | | Perplexity | โœ… | โŒ | P3 | Freshness parameter for web_search | | MiniMax | โœ… | โŒ | P3 | Regional endpoint selection | @@ -470,7 +471,7 @@ This document tracks feature parity between IronClaw (Rust implementation) and O | Device pairing | โœ… | โŒ | | | Tailscale identity | โœ… | โŒ | | | Trusted-proxy auth | โœ… | โŒ | Header-based reverse proxy auth | -| OAuth flows | โœ… | ๐Ÿšง | NEAR AI OAuth + Gemini OAuth (PKCE, S256, loopback redirect, offline access) | +| OAuth flows | โœ… | ๐Ÿšง | NEAR AI OAuth + Gemini OAuth (PKCE, S256) + hosted extension/MCP OAuth broker; external auth-proxy rollout still pending | | DM pairing verification | โœ… | โœ… | ironclaw pairing approve, host APIs | | Allowlist/blocklist | โœ… | ๐Ÿšง | allow_from + pairing store | | Per-group tool policies | โœ… | โŒ | | diff --git a/README.md b/README.md index fa73dc45..6e14d9ea 100644 --- a/README.md +++ b/README.md @@ -168,7 +168,7 @@ written to `~/.ironclaw/.env` so they are available before the database connects ### Alternative LLM Providers IronClaw defaults to NEAR AI but supports many LLM providers out of the box. -Built-in providers include **Anthropic**, **OpenAI**, **Google Gemini**, **MiniMax**, +Built-in providers include **Anthropic**, **OpenAI**, **GitHub Copilot**, **Google Gemini**, **MiniMax**, **Mistral**, and **Ollama** (local). OpenAI-compatible services like **OpenRouter** (300+ models), **Together AI**, **Fireworks AI**, and self-hosted servers (**vLLM**, **LiteLLM**) are also supported. diff --git a/README.zh-CN.md b/README.zh-CN.md index a337d713..d818872a 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -165,7 +165,7 @@ ironclaw onboard ### ๆ›ฟไปฃ LLM ๆไพ›ๅ•† IronClaw ้ป˜่ฎคไฝฟ็”จ NEAR AI๏ผŒไฝ†ๅผ€็ฎฑๅณ็”จๅœฐๆ”ฏๆŒๅคš็ง LLM ๆไพ›ๅ•†ใ€‚ -ๅ†…็ฝฎๆไพ›ๅ•†ๅŒ…ๆ‹ฌ **Anthropic**ใ€**OpenAI**ใ€**Google Gemini**ใ€**MiniMax**ใ€**Mistral** ๅ’Œ **Ollama**๏ผˆๆœฌๅœฐ้ƒจ็ฝฒ๏ผ‰ใ€‚ๅŒๆ—ถไนŸๆ”ฏๆŒ OpenAI ๅ…ผๅฎนๆœๅŠก๏ผŒๅฆ‚ **OpenRouter**๏ผˆ300+ ๆจกๅž‹๏ผ‰ใ€**Together AI**ใ€**Fireworks AI** ไปฅๅŠ่‡ชๆ‰˜็ฎกๆœๅŠกๅ™จ๏ผˆ**vLLM**ใ€**LiteLLM**๏ผ‰ใ€‚ +ๅ†…็ฝฎๆไพ›ๅ•†ๅŒ…ๆ‹ฌ **Anthropic**ใ€**OpenAI**ใ€**GitHub Copilot**ใ€**Google Gemini**ใ€**MiniMax**ใ€**Mistral** ๅ’Œ **Ollama**๏ผˆๆœฌๅœฐ้ƒจ็ฝฒ๏ผ‰ใ€‚ๅŒๆ—ถไนŸๆ”ฏๆŒ OpenAI ๅ…ผๅฎนๆœๅŠก๏ผŒๅฆ‚ **OpenRouter**๏ผˆ300+ ๆจกๅž‹๏ผ‰ใ€**Together AI**ใ€**Fireworks AI** ไปฅๅŠ่‡ชๆ‰˜็ฎกๆœๅŠกๅ™จ๏ผˆ**vLLM**ใ€**LiteLLM**๏ผ‰ใ€‚ ๅœจๅ‘ๅฏผไธญ้€‰ๆ‹ฉไฝ ็š„ๆไพ›ๅ•†๏ผŒๆˆ–็›ดๆŽฅ่ฎพ็ฝฎ็Žฏๅขƒๅ˜้‡๏ผš diff --git a/crates/ironclaw_safety/src/lib.rs b/crates/ironclaw_safety/src/lib.rs index 3e9a48ba..d0c3f783 100644 --- a/crates/ironclaw_safety/src/lib.rs +++ b/crates/ironclaw_safety/src/lib.rs @@ -243,6 +243,18 @@ mod tests { assert!(wrapped.contains("Hello ")); } + #[test] + fn test_wrap_for_llm_escapes_attr_chars() { + let config = SafetyConfig { + max_output_length: 100_000, + injection_check_enabled: true, + }; + let safety = SafetyLayer::new(&config); + + let wrapped = safety.wrap_for_llm("bad&\"<>name", "ok", false); + assert!(wrapped.contains("name=\"bad&"<>name\"")); // safety: test assertion in #[cfg(test)] module + } + #[test] fn test_sanitize_action_forces_sanitization_when_injection_check_disabled() { let config = SafetyConfig { diff --git a/deny.toml b/deny.toml index 80aa2215..fddb3d43 100644 --- a/deny.toml +++ b/deny.toml @@ -15,6 +15,8 @@ ignore = [ "RUSTSEC-2026-0020", # wasmtime wasi:http/types.fields panic โ€” mitigated by fuel limits "RUSTSEC-2026-0021", + # rustls-webpki CRL distributionPoint matching โ€” 0.102.8 pinned by libsql transitive dep + "RUSTSEC-2026-0049", ] [licenses] diff --git a/docs/LLM_PROVIDERS.md b/docs/LLM_PROVIDERS.md index 793f0682..765ce8ea 100644 --- a/docs/LLM_PROVIDERS.md +++ b/docs/LLM_PROVIDERS.md @@ -17,6 +17,7 @@ the most common configurations. | Yandex AI Studio | `yandex` | `YANDEX_API_KEY` | YandexGPT models | | MiniMax | `minimax` | `MINIMAX_API_KEY` | MiniMax-M2.7 models | | Cloudflare Workers AI | `cloudflare` | `CLOUDFLARE_API_KEY` | Access to Workers AI | +| GitHub Copilot | `github_copilot` | `GITHUB_COPILOT_TOKEN` | Multi-models | | Ollama | `ollama` | No | Local inference | | AWS Bedrock | `bedrock` | AWS credentials | Native Converse API | | OpenRouter | `openai_compatible` | `LLM_API_KEY` | 300+ models | @@ -78,7 +79,7 @@ GEMINI_MODEL=gemini-2.5-flash |---|---|---| | Function calling | โœ… | `functionDeclarations` / `functionCall` / `functionResponse` | | `generationConfig` | โœ… | `temperature`, `maxOutputTokens` passed from request | -| `thinkingConfig` | โœ… | `thinkingBudget`/`thinkingLevel` for thinking-capable models | +| `thinkingConfig` | โœ… | `thinkingBudget`/`thinkingLevel` for thinking-capable models (does NOT set `includeThoughts`) | | `toolConfig` | โœ… | `functionCallingConfig.mode`: `AUTO`/`ANY`/`NONE` | | SSE streaming | โœ… | Cloud Code API with `streamGenerateContent?alt=sse` | | Token refresh | โœ… | Automatic via refresh token | @@ -106,6 +107,34 @@ API (`generativelanguage.googleapis.com`). --- +## GitHub Copilot + +GitHub Copilot exposes chat endpoint at +`https://api.githubcopilot.com`. IronClaw uses that endpoint directly through the +built-in `github_copilot` provider. + +```env +LLM_BACKEND=github_copilot +GITHUB_COPILOT_TOKEN=gho_... +GITHUB_COPILOT_MODEL=gpt-4o +# Optional advanced headers if your setup needs them: +# GITHUB_COPILOT_EXTRA_HEADERS=Copilot-Integration-Id:vscode-chat +``` + +`ironclaw onboard` can acquire this token for you using GitHub device login. If you +already signed into Copilot through VS Code or a JetBrains IDE, you can also reuse +the `oauth_token` stored in `~/.config/github-copilot/apps.json`. If you prefer, +`LLM_BACKEND=github-copilot` also works as an alias. + +Popular models vary by subscription, but `gpt-4o` is a safe default. IronClaw keeps +model entry manual for this provider because GitHub Copilot model listing may require +extra integration headers on some clients. IronClaw automatically injects the standard +VS Code identity headers (`User-Agent`, `Editor-Version`, `Editor-Plugin-Version`, +`Copilot-Integration-Id`) and lets you override them with +`GITHUB_COPILOT_EXTRA_HEADERS`. + +--- + ## Ollama (local) Install Ollama from [ollama.com](https://ollama.com), pull a model, then: diff --git a/providers.json b/providers.json index 550edd64..517e2a26 100644 --- a/providers.json +++ b/providers.json @@ -77,6 +77,29 @@ "can_list_models": false } }, + { + "id": "github_copilot", + "aliases": [ + "github-copilot", + "githubcopilot", + "copilot" + ], + "protocol": "github_copilot", + "default_base_url": "https://api.githubcopilot.com", + "api_key_env": "GITHUB_COPILOT_TOKEN", + "api_key_required": true, + "model_env": "GITHUB_COPILOT_MODEL", + "default_model": "gpt-4o", + "extra_headers_env": "GITHUB_COPILOT_EXTRA_HEADERS", + "description": "GitHub Copilot Chat API (OAuth token from IDE sign-in)", + "setup": { + "kind": "api_key", + "secret_name": "llm_github_copilot_token", + "key_url": "https://docs.github.com/en/copilot", + "display_name": "GitHub Copilot", + "can_list_models": false + } + }, { "id": "tinfoil", "aliases": [], diff --git a/skills/delegation/SKILL.md b/skills/delegation/SKILL.md new file mode 100644 index 00000000..0163dd32 --- /dev/null +++ b/skills/delegation/SKILL.md @@ -0,0 +1,75 @@ +--- +name: delegation +version: 0.1.0 +description: Helps users delegate tasks, break them into steps, set deadlines, and track progress via routines and memory. +activation: + keywords: + - delegate + - hand off + - assign task + - help me with + - take care of + - remind me to + - schedule + - plan my + - manage my + - track this + patterns: + - "can you.*handle" + - "I need (help|someone) to" + - "take over" + - "set up a reminder" + - "follow up on" + tags: + - personal-assistant + - task-management + - delegation + max_context_tokens: 1500 +--- + +# Task Delegation Assistant + +When the user wants to delegate a task or get help managing something, follow this process: + +## 1. Clarify the Task + +Ask what needs to be done, by when, and any constraints. Get enough detail to act independently but don't over-interrogate. If the request is clear, skip straight to planning. + +## 2. Break It Down + +Decompose the task into concrete, actionable steps. Use `memory_write` to persist the task plan to a path like `tasks/{task-name}.md` with: +- Clear description +- Steps with checkboxes +- Due date (if any) +- Status: pending/in-progress/done + +## 3. Set Up Tracking + +If the task is recurring or has a deadline: +- Create a routine using `routine_create` for scheduled check-ins +- Add a heartbeat item if it needs daily monitoring +- Set up an event-triggered routine if it depends on external input + +## 4. Use Profile Context + +Check `USER.md` for the user's preferences: +- **Proactivity level**: High = check in frequently. Low = only report on completion. +- **Communication style**: Match their preferred tone and detail level. +- **Focus areas**: Prioritize tasks that align with their stated goals. + +## 5. Execute or Queue + +- If you can do it now (search, draft, organize, calculate), do it immediately. +- If it requires waiting, external action, or follow-up, create a reminder routine. +- If it requires tools you don't have, explain what's needed and suggest alternatives. + +## 6. Report Back + +Always confirm the plan with the user before starting execution. After completing, update the task file in memory and notify the user with a concise summary. + +## Communication Guidelines + +- Be direct and action-oriented +- Confirm understanding before acting on ambiguous requests +- When in doubt about autonomy level, ask once then remember the answer +- Use `memory_write` to track delegation preferences for future reference diff --git a/skills/routine-advisor/SKILL.md b/skills/routine-advisor/SKILL.md new file mode 100644 index 00000000..3bb10c72 --- /dev/null +++ b/skills/routine-advisor/SKILL.md @@ -0,0 +1,118 @@ +--- +name: routine-advisor +version: 0.1.0 +description: Suggests relevant cron routines based on user context, goals, and observed patterns +activation: + keywords: + - every day + - every morning + - every week + - routine + - automate + - remind me + - check daily + - monitor + - recurring + - schedule + - habit + - workflow + - keep forgetting + - always have to + - repetitive + - notifications + - digest + - summary + - review daily + - weekly review + patterns: + - "I (always|usually|often|regularly) (check|do|look at|review)" + - "every (morning|evening|week|day|monday|friday)" + - "I (wish|want) (I|it) (could|would) (automatically|auto)" + - "is there a way to (auto|schedule|set up)" + - "can you (check|monitor|watch|track).*for me" + - "I keep (forgetting|missing|having to)" + tags: + - automation + - scheduling + - personal-assistant + - productivity + max_context_tokens: 1500 +--- + +# Routine Advisor + +When the conversation suggests the user has a repeatable task or could benefit from automation, consider suggesting a routine. + +## When to Suggest + +Suggest a routine when you notice: +- The user describes doing something repeatedly ("I check my PRs every morning") +- The user mentions forgetting recurring tasks ("I keep forgetting to...") +- The user asks you to do something that sounds periodic +- You've learned enough about the user to propose a relevant automation +- The user has installed extensions that enable new monitoring capabilities + +## How to Suggest + +Be specific and concrete. Not "Want me to set up a routine?" but rather: "I noticed you review PRs every morning. Want me to create a daily 9am routine that checks your open PRs and sends you a summary?" + +Always include: +1. What the routine would do (specific action) +2. When it would run (specific schedule in plain language) +3. How it would notify them (which channel they're on) + +Wait for the user to confirm before creating. + +## Pacing + +- First 1-3 conversations: Do NOT suggest routines. Focus on helping and learning. +- After learning 2-3 user patterns: Suggest your first routine. Keep it simple. +- After 5+ conversations: Suggest more routines as patterns emerge. +- Never suggest more than 1 routine per conversation unless the user is clearly interested. +- If the user declines, wait at least 3 conversations before suggesting again. + +## Creating Routines + +Use the `routine_create` tool. Before creating, check `routine_list` to avoid duplicates. + +Parameters: +- `trigger_type`: Usually "cron" for scheduled tasks +- `schedule`: Standard cron format. Common schedules: + - Daily 9am: `0 9 * * *` + - Weekday mornings: `0 9 * * MON-FRI` + - Weekly Monday: `0 9 * * MON` + - Every 2 hours during work: `0 9-17/2 * * MON-FRI` + - Sunday evening: `0 18 * * SUN` +- `action_type`: "lightweight" for simple checks, "full_job" for multi-step tasks +- `prompt`: Clear, specific instruction for what the routine should do +- `context_paths`: Workspace files to load as context (e.g., `["context/profile.json", "MEMORY.md"]`) + +## Routine Ideas by User Type + +**Developer:** +- Daily PR review digest (check open PRs, summarize what needs attention) +- CI/CD failure alerts (monitor build status) +- Weekly dependency update check +- Daily standup prep (summarize yesterday's work from daily logs) + +**Professional:** +- Morning briefing (today's priorities from memory + any pending tasks) +- End-of-day summary (what was accomplished, what's pending) +- Weekly goal review (check progress against stated goals) +- Meeting prep reminders + +**Health/Personal:** +- Daily exercise or habit check-in +- Weekly meal planning prompt +- Monthly budget review reminder + +**General:** +- Daily news digest on topics of interest +- Weekly reflection prompt (what went well, what to improve) +- Periodic task/reminder check-in +- Regular cleanup of stale tasks or notes +- Weekly profile evolution (if the user has a profile in `context/profile.json`, suggest a Monday routine that reads the profile via `memory_read`, searches recent conversations for new patterns with `memory_search`, and updates the profile via `memory_write` if any fields should change with confidence > 0.6 โ€” be conservative, only update with clear evidence) + +## Awareness + +Before suggesting, consider what tools and extensions are currently available. Only suggest routines the agent can actually execute. If a routine would need a tool that isn't installed, mention that too: "If you connect your calendar, I could also send you a morning briefing with today's meetings." diff --git a/src/agent/CLAUDE.md b/src/agent/CLAUDE.md index e55c9591..686753de 100644 --- a/src/agent/CLAUDE.md +++ b/src/agent/CLAUDE.md @@ -113,7 +113,7 @@ Check-insert is done under a single write lock to prevent TOCTOU races. A cleanu 4. Detects broken tools via `store.get_broken_tools(5)` (threshold: 5 failures). Requires `with_store()` to be called; returns empty without a store. 5. Attempts to rebuild broken tools via `SoftwareBuilder`. Requires `with_builder()` to be called; returns `ManualRequired` without a builder. -Note: the `stuck_threshold` duration is stored but currently unused (marked `#[allow(dead_code)]`). Stuck detection relies on `JobState::Stuck` being set by the state machine, not wall-clock time comparison. +The `stuck_threshold` duration is used for time-based detection of `InProgress` jobs that have been running longer than the threshold. When `detect_stuck_jobs()` finds such jobs, it transitions them to `Stuck` before returning them, enabling the normal `attempt_recovery()` path. Repair results: `Success`, `Retry`, `Failed`, `ManualRequired`. `Retry` does NOT notify the user (to avoid spam). diff --git a/src/agent/agent_loop.rs b/src/agent/agent_loop.rs index 1780ba9d..a0e8278f 100644 --- a/src/agent/agent_loop.rs +++ b/src/agent/agent_loop.rs @@ -10,6 +10,7 @@ use std::sync::Arc; use futures::StreamExt; +use uuid::Uuid; use crate::agent::context_monitor::ContextMonitor; use crate::agent::heartbeat::spawn_heartbeat; @@ -17,7 +18,7 @@ use crate::agent::routine_engine::{RoutineEngine, spawn_cron_ticker}; use crate::agent::self_repair::{DefaultSelfRepair, RepairResult, SelfRepair}; use crate::agent::session_manager::SessionManager; use crate::agent::submission::{Submission, SubmissionParser, SubmissionResult}; -use crate::agent::{HeartbeatConfig as AgentHeartbeatConfig, Router, Scheduler}; +use crate::agent::{HeartbeatConfig as AgentHeartbeatConfig, Router, Scheduler, SchedulerDeps}; use crate::channels::{ChannelManager, IncomingMessage, OutgoingResponse}; use crate::config::{AgentConfig, HeartbeatConfig, RoutineConfig, SkillsConfig}; use crate::context::ContextManager; @@ -31,6 +32,13 @@ use crate::skills::SkillRegistry; use crate::tools::ToolRegistry; use crate::workspace::Workspace; +/// Static greeting persisted to DB and broadcast on first launch. +/// +/// Sent before the LLM is involved so the user sees something immediately. +/// The conversational onboarding (profile building, channel setup) happens +/// organically in the subsequent turns driven by BOOTSTRAP.md. +const BOOTSTRAP_GREETING: &str = include_str!("../workspace/seeds/GREETING.md"); + /// Collapse a tool output string into a single-line preview for display. pub(crate) fn truncate_for_preview(output: &str, max_chars: usize) -> String { let collapsed: String = output @@ -113,6 +121,17 @@ async fn resolve_routine_notification_target( .await } +pub(crate) fn chat_tool_execution_metadata(message: &IncomingMessage) -> serde_json::Value { + serde_json::json!({ + "notify_channel": message.channel, + "notify_user": message + .routing_target() + .unwrap_or_else(|| message.user_id.clone()), + "notify_thread_id": message.thread_id, + "notify_metadata": message.metadata, + }) +} + fn should_fallback_routine_notification(error: &ChannelError) -> bool { !matches!(error, ChannelError::MissingRoutingTarget { .. }) } @@ -146,6 +165,8 @@ pub struct AgentDeps { pub transcription: Option>, /// Document text extraction middleware for PDF, DOCX, PPTX, etc. pub document_extraction: Option>, + /// Sandbox readiness state for full-job routine dispatch. + pub sandbox_readiness: crate::agent::routine_engine::SandboxReadiness, /// Software builder for self-repair tool rebuilding. pub builder: Option>, } @@ -207,9 +228,12 @@ impl Agent { context_manager.clone(), deps.llm.clone(), deps.safety.clone(), - deps.tools.clone(), - deps.store.clone(), - deps.hooks.clone(), + SchedulerDeps { + tools: deps.tools.clone(), + extension_manager: deps.extension_manager.clone(), + store: deps.store.clone(), + hooks: deps.hooks.clone(), + }, ); if let Some(ref tx) = deps.sse_tx { scheduler.set_sse_sender(tx.clone()); @@ -338,6 +362,32 @@ impl Agent { /// Run the agent main loop. pub async fn run(self) -> Result<(), Error> { + // Proactive bootstrap: persist the static greeting to DB *before* + // starting channels so the first web client sees it via history. + let bootstrap_thread_id = if self + .workspace() + .is_some_and(|ws| ws.take_bootstrap_pending()) + { + tracing::debug!( + "Fresh workspace detected โ€” persisting static bootstrap greeting to DB" + ); + if let Some(store) = self.store() { + let thread_id = store + .get_or_create_assistant_conversation("default", "gateway") + .await + .ok(); + if let Some(id) = thread_id { + self.persist_assistant_response(id, "gateway", "default", BOOTSTRAP_GREETING) + .await; + } + thread_id + } else { + None + } + } else { + None + }; + // Start channels let mut message_stream = self.channels.start_all().await?; @@ -554,8 +604,10 @@ impl Agent { Arc::clone(workspace), notify_tx, Some(self.scheduler.clone()), + self.deps.extension_manager.clone(), self.tools().clone(), self.safety().clone(), + self.deps.sandbox_readiness, )); // Register routine tools @@ -668,6 +720,30 @@ impl Agent { None }; + // Bootstrap phase 2: register the thread in session manager and + // broadcast the greeting via SSE for any clients already connected. + // The greeting was already persisted to DB before start_all(), so + // clients that connect after this point will see it via history. + if let Some(id) = bootstrap_thread_id { + // Use get_or_create_session (not resolve_thread) to avoid creating + // an orphan thread. Then insert the DB-sourced thread directly. + let session = self.session_manager.get_or_create_session("default").await; + { + use crate::agent::session::Thread; + let mut sess = session.lock().await; + let thread = Thread::with_id(id, sess.id); + sess.active_thread = Some(id); + sess.threads.entry(id).or_insert(thread); + } + self.session_manager + .register_thread("default", "gateway", id, session) + .await; + + let mut out = OutgoingResponse::text(BOOTSTRAP_GREETING.to_string()); + out.thread_id = Some(id.to_string()); + let _ = self.channels.broadcast("gateway", "default", out).await; + } + // Main message loop tracing::debug!("Agent {} ready and listening", self.config.name); @@ -861,9 +937,6 @@ impl Agent { } async fn handle_message(&self, message: &IncomingMessage) -> Result, Error> { - // Log at info level only for tracking without exposing PII (user_id can be a phone number) - tracing::info!(message_id = %message.id, "Processing message"); - // Log sensitive details at debug level for troubleshooting tracing::debug!( message_id = %message.id, @@ -942,19 +1015,59 @@ impl Agent { } } - // Resolve session and thread - tracing::debug!( - message_id = %message.id, - "Resolving session and thread" - ); - let (session, thread_id) = self - .session_manager - .resolve_thread( - &message.user_id, - &message.channel, - message.conversation_scope(), - ) - .await; + // Resolve session and thread. Approval submissions are allowed to + // target an already-loaded owned thread by UUID across channels so the + // web approval UI can approve work that originated from HTTP/other + // owner-scoped channels. + let approval_thread_uuid = if matches!( + submission, + Submission::ExecApproval { .. } | Submission::ApprovalResponse { .. } + ) { + message + .conversation_scope() + .and_then(|thread_id| Uuid::parse_str(thread_id).ok()) + } else { + None + }; + + let (session, thread_id) = if let Some(target_thread_id) = approval_thread_uuid { + let session = self + .session_manager + .get_or_create_session(&message.user_id) + .await; + let mut sess = session.lock().await; + if sess.threads.contains_key(&target_thread_id) { + sess.active_thread = Some(target_thread_id); + sess.last_active_at = chrono::Utc::now(); + drop(sess); + self.session_manager + .register_thread( + &message.user_id, + &message.channel, + target_thread_id, + Arc::clone(&session), + ) + .await; + (session, target_thread_id) + } else { + drop(sess); + self.session_manager + .resolve_thread( + &message.user_id, + &message.channel, + message.conversation_scope(), + ) + .await + } + } else { + self.session_manager + .resolve_thread( + &message.user_id, + &message.channel, + message.conversation_scope(), + ) + .await + }; tracing::debug!( message_id = %message.id, thread_id = %thread_id, @@ -1124,9 +1237,10 @@ impl Agent { #[cfg(test)] mod tests { use super::{ - resolve_routine_notification_user, should_fallback_routine_notification, - truncate_for_preview, + chat_tool_execution_metadata, resolve_routine_notification_user, + should_fallback_routine_notification, truncate_for_preview, }; + use crate::channels::IncomingMessage; use crate::error::ChannelError; #[test] @@ -1222,6 +1336,50 @@ mod tests { assert_eq!(resolve_routine_notification_user(&metadata), None); // safety: test-only assertion } + #[test] + fn chat_tool_execution_metadata_prefers_message_routing_target() { + let message = IncomingMessage::new("telegram", "owner-scope", "hello") + .with_sender_id("telegram-user") + .with_thread("thread-7") + .with_metadata(serde_json::json!({ + "chat_id": 424242, + "chat_type": "private", + })); + + let metadata = chat_tool_execution_metadata(&message); + assert_eq!( + metadata.get("notify_channel").and_then(|v| v.as_str()), + Some("telegram") + ); // safety: test-only assertion + assert_eq!( + metadata.get("notify_user").and_then(|v| v.as_str()), + Some("424242") + ); // safety: test-only assertion + assert_eq!( + metadata.get("notify_thread_id").and_then(|v| v.as_str()), + Some("thread-7") + ); // safety: test-only assertion + } + + #[test] + fn chat_tool_execution_metadata_falls_back_to_user_scope_without_route() { + let message = IncomingMessage::new("gateway", "owner-scope", "hello").with_sender_id(""); + + let metadata = chat_tool_execution_metadata(&message); + assert_eq!( + metadata.get("notify_channel").and_then(|v| v.as_str()), + Some("gateway") + ); // safety: test-only assertion + assert_eq!( + metadata.get("notify_user").and_then(|v| v.as_str()), + Some("owner-scope") + ); // safety: test-only assertion + assert_eq!( + metadata.get("notify_thread_id"), + Some(&serde_json::Value::Null) + ); // safety: test-only assertion + } + #[test] fn targeted_routine_notifications_do_not_fallback_without_owner_route() { let error = ChannelError::MissingRoutingTarget { diff --git a/src/agent/dispatcher.rs b/src/agent/dispatcher.rs index d3825b2f..fc3da61b 100644 --- a/src/agent/dispatcher.rs +++ b/src/agent/dispatcher.rs @@ -144,12 +144,7 @@ impl Agent { .with_requester_id(&message.sender_id); job_ctx.http_interceptor = self.deps.http_interceptor.clone(); job_ctx.user_timezone = user_tz.name().to_string(); - job_ctx.metadata = serde_json::json!({ - "notify_channel": message.channel, - "notify_user": message.user_id, - "notify_thread_id": message.thread_id, - "notify_metadata": message.metadata, - }); + job_ctx.metadata = crate::agent::agent_loop::chat_tool_execution_metadata(message); // Build system prompts once for this turn. Two variants: with tools // (normal iterations) and without (force_text final iteration). @@ -1199,6 +1194,7 @@ mod tests { http_interceptor: None, transcription: None, document_extraction: None, + sandbox_readiness: crate::agent::routine_engine::SandboxReadiness::DisabledByConfig, builder: None, }; @@ -2070,6 +2066,7 @@ mod tests { http_interceptor: None, transcription: None, document_extraction: None, + sandbox_readiness: crate::agent::routine_engine::SandboxReadiness::DisabledByConfig, builder: None, }; @@ -2189,6 +2186,7 @@ mod tests { http_interceptor: None, transcription: None, document_extraction: None, + sandbox_readiness: crate::agent::routine_engine::SandboxReadiness::DisabledByConfig, builder: None, }; diff --git a/src/agent/job_monitor.rs b/src/agent/job_monitor.rs index 714caeac..675d0426 100644 --- a/src/agent/job_monitor.rs +++ b/src/agent/job_monitor.rs @@ -14,12 +14,15 @@ //! Agent Loop //! ``` +use std::sync::Arc; + use tokio::sync::{broadcast, mpsc}; use tokio::task::JoinHandle; use uuid::Uuid; use crate::channels::IncomingMessage; use crate::channels::web::types::SseEvent; +use crate::context::{ContextManager, JobState}; /// Route context for forwarding job monitor events back to the user's channel. #[derive(Debug, Clone)] @@ -40,10 +43,23 @@ pub struct JobMonitorRoute { /// Tool use/result and status events are intentionally skipped (too noisy for /// the main agent's context window). pub fn spawn_job_monitor( + job_id: Uuid, + event_rx: broadcast::Receiver<(Uuid, SseEvent)>, + inject_tx: mpsc::Sender, + route: JobMonitorRoute, +) -> JoinHandle<()> { + spawn_job_monitor_with_context(job_id, event_rx, inject_tx, route, None) +} + +/// Like `spawn_job_monitor`, but also transitions the job's in-memory state +/// when it receives a `JobResult` event. This ensures fire-and-forget sandbox +/// jobs don't stay `InProgress` forever in the `ContextManager`. +pub fn spawn_job_monitor_with_context( job_id: Uuid, mut event_rx: broadcast::Receiver<(Uuid, SseEvent)>, inject_tx: mpsc::Sender, route: JobMonitorRoute, + context_manager: Option>, ) -> JoinHandle<()> { let short_id = job_id.to_string()[..8].to_string(); @@ -77,6 +93,26 @@ pub fn spawn_job_monitor( } } SseEvent::JobResult { status, .. } => { + // Transition in-memory state so the job frees its + // max_jobs slot and query tools show the final state. + if let Some(ref cm) = context_manager { + let target = if status == "completed" { + JobState::Completed + } else { + JobState::Failed + }; + let reason = if status != "completed" { + Some(format!("Container finished: {}", status)) + } else { + None + }; + let _ = cm + .update_context(job_id, |ctx| { + let _ = ctx.transition_to(target, reason); + }) + .await; + } + let mut msg = IncomingMessage::new( route.channel.clone(), route.user_id.clone(), @@ -121,6 +157,62 @@ pub fn spawn_job_monitor( }) } +/// Lightweight watcher that only transitions ContextManager state on job +/// completion. Used when monitor routing metadata is absent (no channel to +/// inject messages into) but we still need to free the `max_jobs` slot. +pub fn spawn_completion_watcher( + job_id: Uuid, + mut event_rx: broadcast::Receiver<(Uuid, SseEvent)>, + context_manager: Arc, +) -> JoinHandle<()> { + let short_id = job_id.to_string()[..8].to_string(); + + tokio::spawn(async move { + loop { + match event_rx.recv().await { + Ok((ev_job_id, SseEvent::JobResult { status, .. })) if ev_job_id == job_id => { + let target = if status == "completed" { + JobState::Completed + } else { + JobState::Failed + }; + let reason = if status != "completed" { + Some(format!("Container finished: {}", status)) + } else { + None + }; + let _ = context_manager + .update_context(job_id, |ctx| { + let _ = ctx.transition_to(target, reason); + }) + .await; + tracing::debug!( + job_id = %short_id, + status = %status, + "Completion watcher exiting (job finished)" + ); + break; + } + Ok(_) => {} + Err(broadcast::error::RecvError::Lagged(n)) => { + tracing::warn!( + job_id = %short_id, + skipped = n, + "Completion watcher lagged" + ); + } + Err(broadcast::error::RecvError::Closed) => { + tracing::debug!( + job_id = %short_id, + "Broadcast channel closed, stopping completion watcher" + ); + break; + } + } + } + }) +} + #[cfg(test)] mod tests { use super::*; @@ -211,6 +303,7 @@ mod tests { job_id: job_id.to_string(), status: "completed".to_string(), session_id: None, + fallback_deliverable: None, }, )) .unwrap(); @@ -293,4 +386,139 @@ mod tests { let msg = IncomingMessage::new("monitor", "system", "test").into_internal(); assert!(msg.is_internal); } + + // === Regression: fire-and-forget sandbox jobs must transition out of InProgress === + // Before this fix, spawn_job_monitor only forwarded SSE messages but never + // updated ContextManager. Background sandbox jobs stayed InProgress forever, + // permanently consuming a max_jobs slot. + + #[tokio::test] + async fn test_monitor_transitions_context_on_completion() { + use crate::context::{ContextManager, JobState}; + + let cm = Arc::new(ContextManager::new(5)); + let job_id = Uuid::new_v4(); + cm.register_sandbox_job(job_id, "user-1", "Build app", "desc") + .await + .unwrap(); + + let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16); + let (inject_tx, mut inject_rx) = mpsc::channel::(16); + + let handle = spawn_job_monitor_with_context( + job_id, + event_tx.subscribe(), + inject_tx, + test_route(), + Some(Arc::clone(&cm)), + ); + + // Send completion event + event_tx + .send(( + job_id, + SseEvent::JobResult { + job_id: job_id.to_string(), + status: "completed".to_string(), + session_id: None, + fallback_deliverable: None, + }, + )) + .unwrap(); + + // Drain the injected message + let _ = tokio::time::timeout(std::time::Duration::from_secs(1), inject_rx.recv()).await; + + // Wait for monitor to exit + tokio::time::timeout(std::time::Duration::from_secs(1), handle) + .await + .expect("monitor should exit") + .expect("monitor should not panic"); + + // Job should now be Completed, not InProgress + let ctx = cm.get_context(job_id).await.unwrap(); + assert_eq!(ctx.state, JobState::Completed); + } + + #[tokio::test] + async fn test_monitor_transitions_context_on_failure() { + use crate::context::{ContextManager, JobState}; + + let cm = Arc::new(ContextManager::new(5)); + let job_id = Uuid::new_v4(); + cm.register_sandbox_job(job_id, "user-1", "Build app", "desc") + .await + .unwrap(); + + let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16); + let (inject_tx, mut inject_rx) = mpsc::channel::(16); + + let handle = spawn_job_monitor_with_context( + job_id, + event_tx.subscribe(), + inject_tx, + test_route(), + Some(Arc::clone(&cm)), + ); + + // Send failure event + event_tx + .send(( + job_id, + SseEvent::JobResult { + job_id: job_id.to_string(), + status: "failed".to_string(), + session_id: None, + fallback_deliverable: None, + }, + )) + .unwrap(); + + let _ = tokio::time::timeout(std::time::Duration::from_secs(1), inject_rx.recv()).await; + tokio::time::timeout(std::time::Duration::from_secs(1), handle) + .await + .expect("monitor should exit") + .expect("monitor should not panic"); + + let ctx = cm.get_context(job_id).await.unwrap(); + assert_eq!(ctx.state, JobState::Failed); + } + + // === Regression: completion watcher (no route metadata) === + // When monitor_route_from_ctx() returns None, spawn_completion_watcher + // must still transition the job so the max_jobs slot is freed. + + #[tokio::test] + async fn test_completion_watcher_transitions_on_result() { + use crate::context::{ContextManager, JobState}; + + let cm = Arc::new(ContextManager::new(5)); + let job_id = Uuid::new_v4(); + cm.register_sandbox_job(job_id, "user-1", "Build app", "desc") + .await + .unwrap(); + + let (event_tx, _) = broadcast::channel::<(Uuid, SseEvent)>(16); + let handle = spawn_completion_watcher(job_id, event_tx.subscribe(), Arc::clone(&cm)); + + event_tx + .send(( + job_id, + SseEvent::JobResult { + job_id: job_id.to_string(), + status: "completed".to_string(), + session_id: None, + fallback_deliverable: None, + }, + )) + .unwrap(); + + tokio::time::timeout(std::time::Duration::from_secs(1), handle) + .await + .expect("watcher should exit") + .expect("watcher should not panic"); + + let ctx = cm.get_context(job_id).await.unwrap(); + assert_eq!(ctx.state, JobState::Completed); + } } diff --git a/src/agent/mod.rs b/src/agent/mod.rs index ee980233..84155666 100644 --- a/src/agent/mod.rs +++ b/src/agent/mod.rs @@ -39,8 +39,8 @@ pub use context_monitor::{CompactionStrategy, ContextBreakdown, ContextMonitor}; pub use heartbeat::{HeartbeatConfig, HeartbeatResult, HeartbeatRunner, spawn_heartbeat}; pub use router::{MessageIntent, Router}; pub use routine::{Routine, RoutineAction, RoutineRun, Trigger}; -pub use routine_engine::RoutineEngine; -pub use scheduler::Scheduler; +pub use routine_engine::{RoutineEngine, SandboxReadiness}; +pub use scheduler::{Scheduler, SchedulerDeps}; pub use self_repair::{BrokenTool, RepairResult, RepairTask, SelfRepair, StuckJob}; pub use session::{PendingApproval, PendingAuth, Session, Thread, ThreadState, Turn, TurnState}; pub use session_manager::SessionManager; diff --git a/src/agent/routine.rs b/src/agent/routine.rs index f3850fa0..1b8ca96a 100644 --- a/src/agent/routine.rs +++ b/src/agent/routine.rs @@ -79,6 +79,13 @@ pub enum Trigger { #[serde(default)] filters: std::collections::HashMap, }, + /// Fire on incoming webhook POST to /api/webhooks/{path}. + Webhook { + /// Optional webhook path suffix (defaults to routine id). + path: Option, + /// Optional shared secret for HMAC validation. + secret: Option, + }, /// Only fires via tool call or CLI. Manual, } @@ -90,6 +97,7 @@ impl Trigger { Trigger::Cron { .. } => "cron", Trigger::Event { .. } => "event", Trigger::SystemEvent { .. } => "system_event", + Trigger::Webhook { .. } => "webhook", Trigger::Manual => "manual", } } @@ -171,6 +179,17 @@ impl Trigger { filters, }) } + "webhook" => { + let path = config + .get("path") + .and_then(|v| v.as_str()) + .map(String::from); + let secret = config + .get("secret") + .and_then(|v| v.as_str()) + .map(String::from); + Ok(Trigger::Webhook { path, secret }) + } "manual" => Ok(Trigger::Manual), other => Err(RoutineError::UnknownTriggerType { trigger_type: other.to_string(), @@ -198,6 +217,10 @@ impl Trigger { "event_type": event_type, "filters": filters, }), + Trigger::Webhook { path, secret } => serde_json::json!({ + "path": path, + "secret": secret, + }), Trigger::Manual => serde_json::json!({}), } } @@ -235,11 +258,6 @@ pub enum RoutineAction { /// Max reasoning iterations (default: 10). #[serde(default = "default_max_iterations")] max_iterations: u32, - /// Tool names pre-authorized for `Always`-approval tools (e.g. destructive - /// shell commands, cross-channel messaging). `UnlessAutoApproved` tools are - /// automatically permitted in routine jobs without listing them here. - #[serde(default)] - tool_permissions: Vec, }, } @@ -264,19 +282,6 @@ fn clamp_max_tool_rounds(value: u64) -> u32 { value.clamp(1, MAX_TOOL_ROUNDS_LIMIT as u64) as u32 } -/// Parse a `tool_permissions` JSON array into a `Vec`. -pub fn parse_tool_permissions(value: &serde_json::Value) -> Vec { - value - .get("tool_permissions") - .and_then(|v| v.as_array()) - .map(|arr| { - arr.iter() - .filter_map(|v| v.as_str().map(String::from)) - .collect() - }) - .unwrap_or_default() -} - impl RoutineAction { /// The string tag stored in the DB action_type column. pub fn type_tag(&self) -> &'static str { @@ -351,12 +356,10 @@ impl RoutineAction { .and_then(|v| v.as_u64()) .unwrap_or(default_max_iterations() as u64) as u32; - let tool_permissions = parse_tool_permissions(&config); Ok(RoutineAction::FullJob { title, description, max_iterations, - tool_permissions, }) } other => Err(RoutineError::UnknownActionType { @@ -385,12 +388,10 @@ impl RoutineAction { title, description, max_iterations, - tool_permissions, } => serde_json::json!({ "title": title, "description": description, "max_iterations": max_iterations, - "tool_permissions": tool_permissions, }), } } @@ -516,16 +517,36 @@ pub fn content_hash(content: &str) -> u64 { hasher.finish() } +/// Normalize a cron expression to the 7-field format expected by the `cron` crate. +/// +/// The `cron` crate requires: `sec min hour day-of-month month day-of-week year`. +/// Standard cron uses 5 fields: `min hour day-of-month month day-of-week`. +/// This function auto-expands: +/// - 5-field โ†’ prepend `0` (seconds) and append `*` (year) +/// - 6-field โ†’ append `*` (year) +/// - 7-field โ†’ pass through unchanged +pub fn normalize_cron_expression(schedule: &str) -> String { + let trimmed = schedule.trim(); + let fields: Vec<&str> = trimmed.split_whitespace().collect(); + match fields.len() { + 5 => format!("0 {} *", trimmed), + 6 => format!("{} *", trimmed), + _ => trimmed.to_string(), + } +} + /// Parse a cron expression and compute the next fire time from now. /// +/// Accepts standard 5-field, 6-field, or 7-field cron expressions (auto-normalized). /// When `timezone` is provided and valid, the schedule is evaluated in that /// timezone and the result is converted back to UTC. Otherwise UTC is used. pub fn next_cron_fire( schedule: &str, timezone: Option<&str>, ) -> Result>, RoutineError> { + let normalized = normalize_cron_expression(schedule); let cron_schedule = - cron::Schedule::from_str(schedule).map_err(|e| RoutineError::InvalidCron { + cron::Schedule::from_str(&normalized).map_err(|e| RoutineError::InvalidCron { reason: e.to_string(), })?; if let Some(tz) = timezone.and_then(crate::timezone::parse_timezone) { @@ -705,7 +726,7 @@ pub fn describe_cron(schedule: &str, timezone: Option<&str>) -> String { mod tests { use crate::agent::routine::{ MAX_TOOL_ROUNDS_LIMIT, RoutineAction, RoutineGuardrails, RunStatus, Trigger, content_hash, - describe_cron, next_cron_fire, + describe_cron, next_cron_fire, normalize_cron_expression, }; #[test] @@ -772,13 +793,47 @@ mod tests { title: "Deploy review".to_string(), description: "Review and deploy pending changes".to_string(), max_iterations: 5, - tool_permissions: vec!["shell".to_string()], }; let json = action.to_config_json(); let parsed = RoutineAction::from_db("full_job", json).expect("parse full_job"); assert!( - matches!(parsed, RoutineAction::FullJob { title, max_iterations, tool_permissions, .. } - if title == "Deploy review" && max_iterations == 5 && tool_permissions == vec!["shell".to_string()]) + matches!(parsed, RoutineAction::FullJob { title, max_iterations, .. } + if title == "Deploy review" + && max_iterations == 5) + ); + } + + #[test] + fn test_action_full_job_ignores_legacy_permission_fields() { + let parsed = RoutineAction::from_db( + "full_job", + serde_json::json!({ + "title": "Deploy review", + "description": "Review and deploy pending changes", + "max_iterations": 5, + "tool_permissions": ["shell"], + "permission_mode": "inherit_owner" + }), + ) + .expect("parse full_job"); + assert!(matches!( + parsed, + RoutineAction::FullJob { + ref title, + ref description, + max_iterations, + .. + } if title == "Deploy review" + && description == "Review and deploy pending changes" + && max_iterations == 5 + )); + assert_eq!( + parsed.to_config_json(), + serde_json::json!({ + "title": "Deploy review", + "description": "Review and deploy pending changes", + "max_iterations": 5, + }) ); } @@ -930,9 +985,66 @@ mod tests { .type_tag(), "system_event" ); + assert_eq!( + Trigger::Webhook { + path: None, + secret: None, + } + .type_tag(), + "webhook" + ); assert_eq!(Trigger::Manual.type_tag(), "manual"); } + #[test] + fn test_normalize_cron_5_field() { + // Standard cron: min hour dom month dow + assert_eq!(normalize_cron_expression("0 9 * * 1"), "0 0 9 * * 1 *"); + assert_eq!( + normalize_cron_expression("0 9 * * MON-FRI"), + "0 0 9 * * MON-FRI *" + ); + } + + #[test] + fn test_normalize_cron_6_field() { + // 6-field: sec min hour dom month dow + assert_eq!( + normalize_cron_expression("0 0 9 * * MON-FRI"), + "0 0 9 * * MON-FRI *" + ); + } + + #[test] + fn test_normalize_cron_7_field_passthrough() { + // Already 7-field: no change + assert_eq!( + normalize_cron_expression("0 0 9 * * MON-FRI *"), + "0 0 9 * * MON-FRI *" + ); + } + + #[test] + fn test_next_cron_fire_5_field_accepted() { + // Standard 5-field cron should now work through normalization + let result = next_cron_fire("0 9 * * 1", None); + assert!( + result.is_ok(), + "5-field cron should be accepted: {result:?}" + ); + assert!(result.unwrap().is_some()); + } + + #[test] + fn test_next_cron_fire_5_field_with_timezone() { + let result = next_cron_fire("0 9 * * MON-FRI", Some("America/New_York")); + assert!( + result.is_ok(), + "5-field cron with timezone should be accepted: {result:?}" + ); + assert!(result.unwrap().is_some()); + } + #[test] fn test_action_lightweight_backward_compat_no_use_tools() { // Simulate old DB record without use_tools field diff --git a/src/agent/routine_engine.rs b/src/agent/routine_engine.rs index 2487ac05..2a5f4474 100644 --- a/src/agent/routine_engine.rs +++ b/src/agent/routine_engine.rs @@ -29,11 +29,13 @@ use crate::config::RoutineConfig; use crate::context::{JobContext, JobState}; use crate::db::Database; use crate::error::RoutineError; +use crate::extensions::ExtensionManager; use crate::llm::{ ChatMessage, CompletionRequest, FinishReason, LlmProvider, ToolCall, ToolCompletionRequest, }; use crate::tools::{ - ApprovalContext, ApprovalRequirement, ToolError, ToolRegistry, prepare_tool_params, + ToolError, ToolRegistry, autonomous_allowed_tool_names, autonomous_unavailable_message, + prepare_tool_params, }; use crate::workspace::Workspace; use ironclaw_safety::SafetyLayer; @@ -43,6 +45,17 @@ enum EventMatcher { System { routine: Routine }, } +/// Distinguishes why sandbox is unavailable so error messages are accurate. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SandboxReadiness { + /// Docker is available and sandbox is enabled. + Available, + /// User explicitly disabled sandboxing (SANDBOX_ENABLED=false). + DisabledByConfig, + /// Sandbox is enabled but Docker is not running or not installed. + DockerUnavailable, +} + /// The routine execution engine. pub struct RoutineEngine { config: RoutineConfig, @@ -57,10 +70,14 @@ pub struct RoutineEngine { event_cache: Arc>>, /// Scheduler for dispatching jobs (FullJob mode). scheduler: Option>, + /// Owner-scoped extension activation state for autonomous tool resolution. + extension_manager: Option>, /// Tool registry for lightweight routine tool execution. tools: Arc, /// Safety layer for tool output sanitization. safety: Arc, + /// Sandbox readiness state for full-job dispatch. + sandbox_readiness: SandboxReadiness, /// Timestamp when this engine instance was created. Used by /// `sync_dispatched_runs` to distinguish orphaned runs (from a previous /// process) from actively-watched runs (from this process). @@ -76,8 +93,10 @@ impl RoutineEngine { workspace: Arc, notify_tx: mpsc::Sender, scheduler: Option>, + extension_manager: Option>, tools: Arc, safety: Arc, + sandbox_readiness: SandboxReadiness, ) -> Self { Self { config, @@ -88,8 +107,10 @@ impl RoutineEngine { running_count: Arc::new(AtomicUsize::new(0)), event_cache: Arc::new(RwLock::new(Vec::new())), scheduler, + extension_manager, tools, safety, + sandbox_readiness, boot_time: Utc::now(), } } @@ -686,8 +707,95 @@ impl RoutineEngine { notify_tx: self.notify_tx.clone(), running_count: self.running_count.clone(), scheduler: self.scheduler.clone(), + extension_manager: self.extension_manager.clone(), tools: self.tools.clone(), safety: self.safety.clone(), + sandbox_readiness: self.sandbox_readiness, + }; + + tokio::spawn(async move { + execute_routine(engine, routine, run).await; + }); + + Ok(run_id) + } + + /// Fire a routine from a webhook trigger. + /// + /// Similar to `fire_manual` but records the trigger as `"webhook"` with the + /// webhook path as detail. Skips ownership check (auth is via webhook secret). + /// Enforces enabled check, cooldown, and concurrent run limit. + pub async fn fire_webhook( + &self, + routine_id: Uuid, + webhook_path: &str, + ) -> Result { + let routine = self + .store + .get_routine(routine_id) + .await + .map_err(|e| RoutineError::Database { + reason: e.to_string(), + })? + .ok_or(RoutineError::NotFound { id: routine_id })?; + + if !routine.enabled { + return Err(RoutineError::Disabled { + name: routine.name.clone(), + }); + } + + if !self.check_cooldown(&routine) { + return Err(RoutineError::Cooldown { + name: routine.name.clone(), + }); + } + + if !self.check_concurrent(&routine).await { + return Err(RoutineError::MaxConcurrent { + name: routine.name.clone(), + }); + } + + if self.running_count.load(Ordering::Relaxed) >= self.config.max_concurrent_routines { + return Err(RoutineError::MaxConcurrent { + name: routine.name.clone(), + }); + } + + let run_id = Uuid::new_v4(); + let run = RoutineRun { + id: run_id, + routine_id: routine.id, + trigger_type: "webhook".to_string(), + trigger_detail: Some(webhook_path.to_string()), + started_at: Utc::now(), + completed_at: None, + status: RunStatus::Running, + result_summary: None, + tokens_used: None, + job_id: None, + created_at: Utc::now(), + }; + + if let Err(e) = self.store.create_routine_run(&run).await { + return Err(RoutineError::Database { + reason: format!("failed to create run record: {e}"), + }); + } + + let engine = EngineContext { + config: self.config.clone(), + store: self.store.clone(), + llm: self.llm.clone(), + workspace: self.workspace.clone(), + notify_tx: self.notify_tx.clone(), + running_count: self.running_count.clone(), + scheduler: self.scheduler.clone(), + extension_manager: self.extension_manager.clone(), + tools: self.tools.clone(), + safety: self.safety.clone(), + sandbox_readiness: self.sandbox_readiness, }; tokio::spawn(async move { @@ -721,8 +829,10 @@ impl RoutineEngine { notify_tx: self.notify_tx.clone(), running_count: self.running_count.clone(), scheduler: self.scheduler.clone(), + extension_manager: self.extension_manager.clone(), tools: self.tools.clone(), safety: self.safety.clone(), + sandbox_readiness: self.sandbox_readiness, }; // Record the run in DB, then spawn execution @@ -857,8 +967,10 @@ struct EngineContext { notify_tx: mpsc::Sender, running_count: Arc, scheduler: Option>, + extension_manager: Option>, tools: Arc, safety: Arc, + sandbox_readiness: SandboxReadiness, } /// Execute a routine run. Handles both lightweight and full_job modes. @@ -889,18 +1001,13 @@ async fn execute_routine(ctx: EngineContext, routine: Routine, run: RoutineRun) title, description, max_iterations, - tool_permissions, } => { - execute_full_job( - &ctx, - &routine, - &run, + let execution = FullJobExecutionConfig { title, description, - *max_iterations, - tool_permissions, - ) - .await + max_iterations: *max_iterations, + }; + execute_full_job(&ctx, &routine, &run, &execution).await } }; @@ -1026,15 +1133,36 @@ fn sanitize_routine_name(name: &str) -> String { /// non-active state (not Pending/InProgress/Stuck). Returns the final /// `RunStatus` mapped from the job outcome. This keeps the routine run /// active for the full job lifetime so concurrency guardrails apply. +struct FullJobExecutionConfig<'a> { + title: &'a str, + description: &'a str, + max_iterations: u32, +} + async fn execute_full_job( ctx: &EngineContext, routine: &Routine, run: &RoutineRun, - title: &str, - description: &str, - max_iterations: u32, - tool_permissions: &[String], + execution: &FullJobExecutionConfig<'_>, ) -> Result<(RunStatus, Option, Option), RoutineError> { + match ctx.sandbox_readiness { + SandboxReadiness::Available => {} + SandboxReadiness::DisabledByConfig => { + return Err(RoutineError::JobDispatchFailed { + reason: "Sandboxing is disabled (SANDBOX_ENABLED=false). \ + Full-job routines require sandbox." + .to_string(), + }); + } + SandboxReadiness::DockerUnavailable => { + return Err(RoutineError::JobDispatchFailed { + reason: "Sandbox is enabled but Docker is not available. \ + Install Docker or set SANDBOX_ENABLED=false." + .to_string(), + }); + } + } + let scheduler = ctx .scheduler .as_ref() @@ -1042,8 +1170,10 @@ async fn execute_full_job( reason: "scheduler not available".to_string(), })?; - let mut metadata = - serde_json::json!({ "max_iterations": max_iterations, "owner_id": routine.user_id }); + let mut metadata = serde_json::json!({ + "max_iterations": execution.max_iterations, + "owner_id": routine.user_id + }); // Carry the routine's notify config in job metadata so the message tool // can resolve channel/target per-job without global state mutation. if let Some(channel) = &routine.notify.channel { @@ -1051,17 +1181,12 @@ async fn execute_full_job( } metadata["notify_user"] = serde_json::json!(&routine.notify.user); - // Build approval context: UnlessAutoApproved tools are auto-approved for routines; - // Always tools require explicit listing in tool_permissions. - let approval_context = ApprovalContext::autonomous_with_tools(tool_permissions.iter().cloned()); - let job_id = scheduler - .dispatch_job_with_context( + .dispatch_job( &routine.user_id, - title, - description, + execution.title, + execution.description, Some(metadata), - approval_context, ) .await .map_err(|e| RoutineError::JobDispatchFailed { @@ -1082,7 +1207,7 @@ async fn execute_full_job( tracing::info!( routine = %routine.name, job_id = %job_id, - max_iterations = max_iterations, + max_iterations = execution.max_iterations, "Dispatched full job for routine, watching for completion" ); @@ -1350,6 +1475,9 @@ async fn execute_lightweight_with_tools( description: routine.name.clone(), ..Default::default() }; + let allowed_tools = + autonomous_allowed_tool_names(&ctx.tools, ctx.extension_manager.as_ref(), &routine.user_id) + .await; loop { iteration += 1; @@ -1384,8 +1512,11 @@ async fn execute_lightweight_with_tools( // Tool-enabled iteration let tool_defs = ctx .tools - .tool_definitions_excluding(ROUTINE_TOOL_DENYLIST) - .await; + .tool_definitions() + .await + .into_iter() + .filter(|tool| allowed_tools.contains(&tool.name)) + .collect(); let request_messages = snapshot_messages_for_tool_iteration(&messages); let request = ToolCompletionRequest::new(request_messages, tool_defs) @@ -1420,7 +1551,7 @@ async fn execute_lightweight_with_tools( // Execute tools sequentially for tc in response.tool_calls { - let result = execute_routine_tool(ctx, &job_ctx, &tc).await; + let result = execute_routine_tool(ctx, &job_ctx, &allowed_tools, &tc).await; // Sanitize and wrap result (including errors) let result_content = match result { @@ -1489,31 +1620,16 @@ fn snapshot_messages_for_tool_iteration(messages: &[ChatMessage]) -> Vec, tc: &ToolCall, ) -> Result> { - // Block tools that pose autonomy-escalation risks - if ROUTINE_TOOL_DENYLIST.contains(&tc.name.as_str()) { - return Err(format!( - "Tool '{}' is not available in lightweight routines", - tc.name - ) - .into()); + if !allowed_tools.contains(&tc.name) { + let message = autonomous_unavailable_message(&tc.name, &job_ctx.user_id); + return Err(message.into()); } // Check if tool exists @@ -1524,22 +1640,6 @@ async fn execute_routine_tool( .ok_or_else(|| format!("Tool '{}' not found", tc.name))?; let normalized_params = prepare_tool_params(tool.as_ref(), &tc.arguments); - // Check approval requirement: only allow Never tools in lightweight routines. - // UnlessAutoApproved and Always tools are blocked to prevent prompt injection attacks. - // Lightweight routines can be triggered by external events and may process untrusted data, - // making them vulnerable to prompt injection that could trick the LLM into calling - // sensitive tools. Blocking these tools entirely is the safest approach. - match tool.requires_approval(&normalized_params) { - ApprovalRequirement::Never => {} - ApprovalRequirement::UnlessAutoApproved | ApprovalRequirement::Always => { - return Err(format!( - "Tool '{}' requires manual approval and cannot be used in lightweight routines", - tc.name - ) - .into()); - } - } - // Validate tool parameters let validation = ctx .safety @@ -1680,6 +1780,7 @@ pub fn spawn_cron_ticker( // never races with FullJobWatcher instances from this process. engine.sync_dispatched_runs().await; engine.check_cron_triggers().await; + engine.sync_dispatched_runs().await; } }) } @@ -1693,6 +1794,56 @@ fn truncate(s: &str, max: usize) -> String { } } +/// Sanitize a summary string from job transitions before using in notifications. +/// +/// `last_reason` comes from untrusted container code, so we: +/// 1. Strip control characters (except newline) to prevent terminal injection +/// 2. Strip HTML tags to prevent injection in web-rendered notifications +/// 3. Collapse multiple whitespace/newlines to single spaces for cleaner output +/// 4. Truncate to 500 chars to prevent oversized notifications +#[cfg(test)] +fn sanitize_summary(s: &str) -> String { + // Strip control characters (keep newline for now, collapse later) + let no_control: String = s + .chars() + .filter(|c| !c.is_control() || *c == '\n') + .collect(); + + // Strip HTML tags (e.g. world"), + "Hello alert('xss') world" + ); + assert_eq!( + sanitize_summary("bold and link"), + "bold and link" + ); + assert_eq!(sanitize_summary(""), ""); + } + + #[test] + fn test_sanitize_summary_multibyte_truncation() { + use super::sanitize_summary; + + // Ensure truncation doesn't panic on multi-byte chars near the boundary + let s = "a".repeat(498) + "\u{1F600}\u{1F600}"; // 498 + two 4-byte emoji + let result = sanitize_summary(&s); + assert!(result.len() <= 503); + assert!(result.ends_with("...")); + } } diff --git a/src/agent/scheduler.rs b/src/agent/scheduler.rs index fa7364a4..2e23b35f 100644 --- a/src/agent/scheduler.rs +++ b/src/agent/scheduler.rs @@ -14,10 +14,14 @@ use crate::config::AgentConfig; use crate::context::{ContextManager, JobContext, JobState}; use crate::db::Database; use crate::error::{Error, JobError}; +use crate::extensions::ExtensionManager; use crate::hooks::HookRegistry; use crate::llm::LlmProvider; use crate::safety::SafetyLayer; -use crate::tools::{ApprovalContext, ToolRegistry, prepare_tool_params}; +use crate::tools::{ + ApprovalContext, ToolRegistry, autonomous_allowed_tool_names, autonomous_unavailable_error, + prepare_tool_params, +}; use crate::worker::job::{Worker, WorkerDeps}; /// Message to send to a worker. @@ -45,6 +49,14 @@ struct ScheduledSubtask { handle: JoinHandle>, } +/// Shared scheduler-owned dependencies that are forwarded into autonomous runs. +pub struct SchedulerDeps { + pub tools: Arc, + pub extension_manager: Option>, + pub store: Option>, + pub hooks: Arc, +} + /// Schedules and manages parallel job execution. pub struct Scheduler { config: AgentConfig, @@ -52,6 +64,7 @@ pub struct Scheduler { llm: Arc, safety: Arc, tools: Arc, + extension_manager: Option>, store: Option>, hooks: Arc, /// SSE broadcast sender for live job event streaming. @@ -71,18 +84,17 @@ impl Scheduler { context_manager: Arc, llm: Arc, safety: Arc, - tools: Arc, - store: Option>, - hooks: Arc, + deps: SchedulerDeps, ) -> Self { Self { config, context_manager, llm, safety, - tools, - store, - hooks, + tools: deps.tools, + extension_manager: deps.extension_manager, + store: deps.store, + hooks: deps.hooks, sse_tx: None, http_interceptor: None, jobs: Arc::new(RwLock::new(HashMap::new())), @@ -120,14 +132,21 @@ impl Scheduler { description: &str, metadata: Option, ) -> Result { - self.dispatch_job_inner(user_id, title, description, metadata, None) - .await + let approval_context = self.autonomous_approval_context(user_id).await; + self.dispatch_job_inner( + user_id, + title, + description, + metadata, + Some(approval_context), + ) + .await } /// Dispatch a job with an explicit approval context for autonomous execution. /// /// Same as `dispatch_job`, but the worker will use the given `ApprovalContext` - /// to determine which tools are pre-approved (instead of blocking all non-`Never` tools). + /// to determine the explicit autonomous allowlist for that job. pub async fn dispatch_job_with_context( &self, user_id: &str, @@ -216,6 +235,13 @@ impl Scheduler { Ok(job_id) } + async fn autonomous_approval_context(&self, user_id: &str) -> ApprovalContext { + ApprovalContext::autonomous_with_tools( + autonomous_allowed_tool_names(&self.tools, self.extension_manager.as_ref(), user_id) + .await, + ) + } + /// Schedule a job for execution. pub async fn schedule(&self, job_id: Uuid) -> Result<(), JobError> { self.schedule_with_context(job_id, None).await @@ -518,10 +544,7 @@ impl Scheduler { let blocked = ApprovalContext::is_blocked_or_default(&approval_context, tool_name, requirement); if blocked { - return Err(crate::error::ToolError::AuthRequired { - name: tool_name.to_string(), - } - .into()); + return Err(autonomous_unavailable_error(tool_name, &job_ctx.user_id).into()); } // Delegate to shared tool execution pipeline @@ -776,7 +799,18 @@ mod tests { let tools = Arc::new(ToolRegistry::new()); let hooks = Arc::new(HookRegistry::default()); - Scheduler::new(config, cm, llm, safety, tools, None, hooks) + Scheduler::new( + config, + cm, + llm, + safety, + SchedulerDeps { + tools, + extension_manager: None, + store: None, + hooks, + }, + ) } #[tokio::test] @@ -1003,12 +1037,14 @@ mod tests { async fn test_execute_tool_task_autonomous_unblocks_soft() { let (tools, cm, safety, job_id) = setup_tools_and_job().await; - // Autonomous context auto-approves UnlessAutoApproved + // Autonomous execution only allows tools explicitly in scope. let result = Scheduler::execute_tool_task( tools.clone(), cm.clone(), safety.clone(), - Some(ApprovalContext::autonomous()), + Some(ApprovalContext::autonomous_with_tools([ + "soft_gate".to_string() + ])), job_id, "soft_gate", serde_json::json!({}), @@ -1040,8 +1076,11 @@ mod tests { async fn test_execute_tool_task_autonomous_with_permissions() { let (tools, cm, safety, job_id) = setup_tools_and_job().await; - // Autonomous context with explicit permission for hard_gate - let ctx = ApprovalContext::autonomous_with_tools(["hard_gate".to_string()]); + // Autonomous context with explicit permission for both tools. + let ctx = ApprovalContext::autonomous_with_tools([ + "soft_gate".to_string(), + "hard_gate".to_string(), + ]); let result = Scheduler::execute_tool_task( tools.clone(), diff --git a/src/agent/self_repair.rs b/src/agent/self_repair.rs index db491194..4e58cb15 100644 --- a/src/agent/self_repair.rs +++ b/src/agent/self_repair.rs @@ -66,6 +66,7 @@ pub trait SelfRepair: Send + Sync { /// Default self-repair implementation. pub struct DefaultSelfRepair { context_manager: Arc, + /// Jobs in `InProgress` longer than this are treated as stuck. stuck_threshold: Duration, max_repair_attempts: u32, store: Option>, @@ -111,15 +112,58 @@ impl DefaultSelfRepair { #[async_trait] impl SelfRepair for DefaultSelfRepair { async fn detect_stuck_jobs(&self) -> Vec { - let stuck_ids = self.context_manager.find_stuck_jobs().await; + let stuck_ids = self + .context_manager + .find_stuck_jobs_with_threshold(Some(self.stuck_threshold)) + .await; let mut stuck_jobs = Vec::new(); for job_id in stuck_ids { if let Ok(ctx) = self.context_manager.get_context(job_id).await - && ctx.state == JobState::Stuck + && matches!(ctx.state, JobState::Stuck | JobState::InProgress) { - // Measure stuck_duration from the most recent Stuck transition, - // not from started_at (which reflects when the job first ran). + // InProgress jobs detected by threshold need to be transitioned + // to Stuck before they can be repaired (attempt_recovery requires + // Stuck state). These jobs already passed the threshold check in + // find_stuck_jobs_with_threshold, so skip the duration filter below. + let just_transitioned = ctx.state == JobState::InProgress; + if just_transitioned { + let reason = "exceeded stuck_threshold"; + let transition = self + .context_manager + .update_context(job_id, |ctx| ctx.mark_stuck(reason)) + .await; + match transition { + Ok(Ok(())) => {} + Ok(Err(e)) => { + tracing::warn!( + job = %job_id, + "Failed to mark InProgress job as Stuck: {}", + e + ); + continue; + } + Err(e) => { + tracing::warn!( + job = %job_id, + "Failed to transition InProgress job to Stuck: {}", + e + ); + continue; + } + } + } + + // Re-fetch context after potential InProgress->Stuck transition + // so that stuck_since picks up the new transition timestamp. + let ctx = match self.context_manager.get_context(job_id).await { + Ok(c) => c, + Err(_) => continue, + }; + + // Use the timestamp of the most recent Stuck transition, not started_at. + // A job that ran for hours before becoming stuck should not immediately + // exceed the threshold โ€” we measure from when it actually became stuck. let stuck_since = ctx .transitions .iter() @@ -134,8 +178,10 @@ impl SelfRepair for DefaultSelfRepair { }) .unwrap_or_default(); - // Only report jobs that have been stuck long enough - if stuck_duration < self.stuck_threshold { + // Only report already-Stuck jobs that have been stuck long enough. + // Jobs just transitioned from InProgress skip this check โ€” they + // were already vetted by find_stuck_jobs_with_threshold. + if !just_transitioned && stuck_duration < self.stuck_threshold { continue; } @@ -163,10 +209,17 @@ impl SelfRepair for DefaultSelfRepair { }); } - // Try to recover the job + // Try to recover the job. + // If the job is still InProgress (detected via stuck_threshold), transition + // it to Stuck first so that attempt_recovery() can move it back to InProgress. let result = self .context_manager - .update_context(job.job_id, |ctx| ctx.attempt_recovery()) + .update_context(job.job_id, |ctx| { + if ctx.state == JobState::InProgress { + ctx.transition_to(JobState::Stuck, Some("exceeded stuck_threshold".into()))?; + } + ctx.attempt_recovery() + }) .await; match result { @@ -489,6 +542,82 @@ mod tests { ); } + #[tokio::test] + async fn detect_and_repair_in_progress_job_via_threshold() { + let cm = Arc::new(ContextManager::new(10)); + let job_id = cm.create_job("Long running", "desc").await.unwrap(); + + // Transition to InProgress. + cm.update_context(job_id, |ctx| ctx.transition_to(JobState::InProgress, None)) + .await + .unwrap() + .unwrap(); + + // Backdate started_at to simulate a job running for 10 minutes. + cm.update_context(job_id, |ctx| { + ctx.started_at = Some(Utc::now() - chrono::Duration::seconds(600)); + }) + .await + .unwrap(); + + // Use a 5-minute threshold so the 10-minute job is detected. + let repair = DefaultSelfRepair::new(Arc::clone(&cm), Duration::from_secs(300), 3); + + // detect_stuck_jobs should find it and transition InProgress -> Stuck. + let stuck = repair.detect_stuck_jobs().await; + assert_eq!(stuck.len(), 1); + assert_eq!(stuck[0].job_id, job_id); + + // After detection the job should now be in Stuck state. + let ctx = cm.get_context(job_id).await.unwrap(); + assert_eq!(ctx.state, JobState::Stuck); + + // Repair should recover it: Stuck -> InProgress. + let result = repair.repair_stuck_job(&stuck[0]).await.unwrap(); + assert!( + matches!(result, RepairResult::Success { .. }), + "Expected Success, got: {:?}", + result + ); + + // Job should be back to InProgress after recovery. + let ctx = cm.get_context(job_id).await.unwrap(); + assert_eq!(ctx.state, JobState::InProgress); + } + + #[tokio::test] + async fn detect_broken_tools_returns_empty_without_store() { + let cm = Arc::new(ContextManager::new(10)); + let repair = DefaultSelfRepair::new(cm, Duration::from_secs(60), 3); + + // No store configured, should return empty. + let broken = repair.detect_broken_tools().await; + assert!(broken.is_empty()); + } + + #[tokio::test] + async fn repair_broken_tool_returns_manual_without_builder() { + let cm = Arc::new(ContextManager::new(10)); + let repair = DefaultSelfRepair::new(cm, Duration::from_secs(60), 3); + + let broken = BrokenTool { + name: "test-tool".to_string(), + failure_count: 10, + last_error: Some("crash".to_string()), + first_failure: Utc::now(), + last_failure: Utc::now(), + last_build_result: None, + repair_attempts: 0, + }; + + let result = repair.repair_broken_tool(&broken).await.unwrap(); + assert!( + matches!(result, RepairResult::ManualRequired { .. }), + "Expected ManualRequired without builder, got: {:?}", + result + ); + } + #[tokio::test] async fn detect_stuck_jobs_filters_by_threshold() { let cm = Arc::new(ContextManager::new(10)); @@ -581,39 +710,6 @@ mod tests { ); } - #[tokio::test] - async fn detect_broken_tools_returns_empty_without_store() { - let cm = Arc::new(ContextManager::new(10)); - let repair = DefaultSelfRepair::new(cm, Duration::from_secs(60), 3); - - // No store configured, should return empty. - let broken = repair.detect_broken_tools().await; - assert!(broken.is_empty()); - } - - #[tokio::test] - async fn repair_broken_tool_returns_manual_without_builder() { - let cm = Arc::new(ContextManager::new(10)); - let repair = DefaultSelfRepair::new(cm, Duration::from_secs(60), 3); - - let broken = BrokenTool { - name: "test-tool".to_string(), - failure_count: 10, - last_error: Some("crash".to_string()), - first_failure: Utc::now(), - last_failure: Utc::now(), - last_build_result: None, - repair_attempts: 0, - }; - - let result = repair.repair_broken_tool(&broken).await.unwrap(); - assert!( - matches!(result, RepairResult::ManualRequired { .. }), - "Expected ManualRequired without builder, got: {:?}", - result - ); - } - /// Mock SoftwareBuilder that returns a successful build result. struct MockBuilder { build_count: std::sync::atomic::AtomicU32, diff --git a/src/agent/session_manager.rs b/src/agent/session_manager.rs index 3db275cc..3bf20697 100644 --- a/src/agent/session_manager.rs +++ b/src/agent/session_manager.rs @@ -772,6 +772,33 @@ mod tests { assert_ne!(resolved, tid); } + #[tokio::test] + async fn test_register_then_resolve_same_uuid_on_second_channel_reuses_thread() { + use crate::agent::session::{Session, Thread}; + + let manager = SessionManager::new(); + let tid = Uuid::new_v4(); + + let session = Arc::new(Mutex::new(Session::new("user-cross"))); + { + let mut sess = session.lock().await; + let thread = Thread::with_id(tid, sess.id); + sess.threads.insert(tid, thread); + } + + manager + .register_thread("user-cross", "http", tid, Arc::clone(&session)) + .await; + manager + .register_thread("user-cross", "gateway", tid, Arc::clone(&session)) + .await; + + let (_, resolved) = manager + .resolve_thread("user-cross", "gateway", Some(&tid.to_string())) + .await; + assert_eq!(resolved, tid); + } + // === QA Plan P3 - 4.2: Concurrent session stress tests === #[tokio::test] diff --git a/src/agent/thread_ops.rs b/src/agent/thread_ops.rs index e7809068..0fb968f1 100644 --- a/src/agent/thread_ops.rs +++ b/src/agent/thread_ops.rs @@ -939,6 +939,7 @@ impl Agent { JobContext::with_user(&message.user_id, "chat", "Interactive chat session") .with_requester_id(&message.sender_id); job_ctx.http_interceptor = self.deps.http_interceptor.clone(); + job_ctx.metadata = crate::agent::agent_loop::chat_tool_execution_metadata(message); // Prefer a valid timezone from the approval message, fall back to the // resolved timezone stored when the approval was originally requested. let tz_candidate = message @@ -1961,7 +1962,7 @@ mod tests { context_messages: vec![], deferred_tool_calls: vec![], user_timezone: None, - allow_always: true, + allow_always: false, }; thread.await_approval(pending); diff --git a/src/app.rs b/src/app.rs index 5cbd6a35..d50cefb3 100644 --- a/src/app.rs +++ b/src/app.rs @@ -25,7 +25,7 @@ use crate::tools::ToolRegistry; use crate::tools::mcp::{McpProcessManager, McpSessionManager}; use crate::tools::wasm::SharedCredentialRegistry; use crate::tools::wasm::WasmToolRuntime; -use crate::workspace::{EmbeddingProvider, Workspace}; +use crate::workspace::{EmbeddingCacheConfig, EmbeddingProvider, Workspace}; /// Fully initialized application components, ready for channel wiring /// and agent construction. @@ -312,12 +312,23 @@ impl AppBuilder { .create_provider(&self.config.llm.nearai.base_url, self.session.clone()); // Register memory tools if database is available + let workspace_user_id = self + .config + .channels + .gateway + .as_ref() + .map(|gw| gw.user_id.as_str()) + .unwrap_or("default"); let workspace = if let Some(ref db) = self.db { - let mut ws = Workspace::new_with_db(&self.config.owner_id, db.clone()) + let emb_cache_config = EmbeddingCacheConfig { + max_entries: self.config.embeddings.cache_size, + }; + let mut ws = Workspace::new_with_db(workspace_user_id, db.clone()) .with_search_config(&self.config.search); if let Some(ref emb) = embeddings { - ws = ws.with_embeddings(emb.clone()); + ws = ws.with_embeddings_cached(emb.clone(), emb_cache_config); } + ws = ws.with_memory_layers(self.config.workspace.memory_layers.clone()); let ws = Arc::new(ws); tools.register_memory_tools(Arc::clone(&ws)); Some(ws) @@ -525,7 +536,7 @@ impl AppBuilder { server_name, e ); - return; + return None; } }; @@ -542,6 +553,10 @@ impl AppBuilder { tool_count, server_name ); + return Some(( + server_name, + Arc::new(client), + )); } Err(e) => { tracing::warn!( @@ -572,14 +587,27 @@ impl AppBuilder { } } } + None }); } + let mut startup_clients = Vec::new(); while let Some(result) = join_set.join_next().await { - if let Err(e) = result { - tracing::warn!("MCP server loading task panicked: {}", e); + match result { + Ok(Some(client_pair)) => { + startup_clients.push(client_pair); + } + Ok(None) => {} + Err(e) => { + if e.is_panic() { + tracing::error!("MCP server loading task panicked: {}", e); + } else { + tracing::warn!("MCP server loading task failed: {}", e); + } + } } } + return startup_clients; } Err(e) => { if matches!( @@ -597,10 +625,12 @@ impl AppBuilder { } } } + Vec::new() } }; - let (dev_loaded_tool_names, _) = tokio::join!(wasm_tools_future, mcp_servers_future); + let (dev_loaded_tool_names, startup_mcp_clients) = + tokio::join!(wasm_tools_future, mcp_servers_future); // Load registry catalog entries for extension discovery let mut catalog_entries = match crate::registry::RegistryCatalog::load_or_embedded() { @@ -662,6 +692,17 @@ impl AppBuilder { )); tools.register_extension_tools(Arc::clone(&manager)); tracing::debug!("Extension manager initialized with in-chat discovery tools"); + + if !startup_mcp_clients.is_empty() { + tracing::info!( + count = startup_mcp_clients.len(), + "Injecting startup MCP clients into extension manager" + ); + for (name, client) in startup_mcp_clients { + manager.inject_mcp_client(name, client).await; + } + } + Some(manager) }; @@ -689,11 +730,11 @@ impl AppBuilder { self.init_secrets().await?; // Post-init validation: backends with dedicated config (nearai, gemini_oauth, - // bedrock) handle their own credential resolution. For registry-based backends, - // fail early if no provider config was resolved. + // bedrock, openai_codex) handle their own credential resolution. For registry-based + // backends, fail early if no provider config was resolved. if !matches!( self.config.llm.backend.as_str(), - "nearai" | "gemini_oauth" | "bedrock" + "nearai" | "gemini_oauth" | "bedrock" | "openai_codex" ) && self.config.llm.provider.is_none() { let backend = &self.config.llm.backend; @@ -724,6 +765,17 @@ impl AppBuilder { dev_loaded_tool_names, ) = self.init_extensions(&tools, &hooks).await?; + // Load bootstrap-completed flag from settings so that existing users + // who already completed onboarding don't re-get bootstrap injection. + if let Some(ref ws) = workspace { + let toml_path = crate::settings::Settings::default_toml_path(); + if let Ok(Some(settings)) = crate::settings::Settings::load_toml(&toml_path) + && settings.profile_onboarding_completed + { + ws.mark_bootstrap_completed(); + } + } + // Seed workspace and backfill embeddings if let Some(ref ws) = workspace { // Import workspace files from disk FIRST if WORKSPACE_IMPORT_DIR is set. diff --git a/src/channels/wasm/wrapper.rs b/src/channels/wasm/wrapper.rs index 8f0c9db4..be7768d0 100644 --- a/src/channels/wasm/wrapper.rs +++ b/src/channels/wasm/wrapper.rs @@ -3314,6 +3314,7 @@ mod tests { use std::sync::Arc; use crate::channels::Channel; + use crate::channels::OutgoingResponse; use crate::channels::wasm::capabilities::ChannelCapabilities; use crate::channels::wasm::runtime::{ PreparedChannelModule, WasmChannelRuntime, WasmChannelRuntimeConfig, @@ -3401,6 +3402,16 @@ mod tests { assert!(channel.health_check().await.is_err()); } + #[tokio::test] + async fn test_broadcast_delegates_to_call_on_broadcast() { + let channel = create_test_channel(); + // With `component: None`, call_on_broadcast short-circuits to Ok(()). + let result = channel + .broadcast("146032821", OutgoingResponse::text("hello")) + .await; + assert!(result.is_ok()); + } + #[tokio::test] async fn test_execute_poll_no_wasm_returns_empty() { // When there's no WASM module (None component), execute_poll diff --git a/src/channels/web/handlers/memory.rs b/src/channels/web/handlers/memory.rs index 8e50f25e..fc0e1fe4 100644 --- a/src/channels/web/handlers/memory.rs +++ b/src/channels/web/handlers/memory.rs @@ -123,25 +123,8 @@ pub async fn memory_read_handler( })) } -pub async fn memory_write_handler( - State(state): State>, - Json(req): Json, -) -> Result, (StatusCode, String)> { - let workspace = state.workspace.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Workspace not available".to_string(), - ))?; - - workspace - .write(&req.path, &req.content) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - Ok(Json(MemoryWriteResponse { - path: req.path, - status: "written", - })) -} +// memory_write_handler lives in server.rs (layer-aware version with append, +// privacy redirect, and proper error status codes). pub async fn memory_search_handler( State(state): State>, diff --git a/src/channels/web/handlers/mod.rs b/src/channels/web/handlers/mod.rs index 0573a067..2f942058 100644 --- a/src/channels/web/handlers/mod.rs +++ b/src/channels/web/handlers/mod.rs @@ -26,3 +26,4 @@ pub mod routines; pub mod settings; #[allow(dead_code)] pub mod static_files; +pub mod webhooks; diff --git a/src/channels/web/handlers/routines.rs b/src/channels/web/handlers/routines.rs index 41bfee5a..368a28ae 100644 --- a/src/channels/web/handlers/routines.rs +++ b/src/channels/web/handlers/routines.rs @@ -303,7 +303,9 @@ fn routine_error_status(err: &RoutineError) -> StatusCode { match err { RoutineError::NotFound { .. } => StatusCode::NOT_FOUND, RoutineError::NotAuthorized { .. } => StatusCode::FORBIDDEN, - RoutineError::Disabled { .. } | RoutineError::MaxConcurrent { .. } => StatusCode::CONFLICT, + RoutineError::Disabled { .. } + | RoutineError::Cooldown { .. } + | RoutineError::MaxConcurrent { .. } => StatusCode::CONFLICT, _ => StatusCode::INTERNAL_SERVER_ERROR, } } diff --git a/src/channels/web/handlers/webhooks.rs b/src/channels/web/handlers/webhooks.rs new file mode 100644 index 00000000..7b041a06 --- /dev/null +++ b/src/channels/web/handlers/webhooks.rs @@ -0,0 +1,197 @@ +//! Public webhook trigger endpoint for routine webhook triggers. +//! +//! `POST /api/webhooks/{path}` โ€” matches the path against routines with +//! `Trigger::Webhook { path, secret }`, validates the secret via constant-time +//! comparison, and fires the matching routine through the `RoutineEngine`. + +use std::sync::Arc; + +use axum::{ + Json, + extract::{Path, State}, + http::{HeaderMap, StatusCode}, +}; +use subtle::ConstantTimeEq; + +use crate::agent::routine::Trigger; +use crate::channels::web::server::GatewayState; + +/// Validate the webhook secret for a routine. +/// +/// Returns `Ok(())` if the routine has a configured secret and the provided +/// secret matches via constant-time comparison. Returns an appropriate HTTP +/// error if the secret is missing (403) or invalid (401). +fn validate_webhook_secret( + trigger: &Trigger, + provided_secret: &str, +) -> Result<(), (StatusCode, String)> { + // Require webhook secret โ€” routines without a secret cannot be triggered via webhook + let expected_secret = match trigger { + Trigger::Webhook { + secret: Some(s), .. + } => s, + _ => { + return Err(( + StatusCode::FORBIDDEN, + "Webhook secret not configured for this routine. \ + Set a secret with: ironclaw routine update --webhook-secret " + .to_string(), + )); + } + }; + + if !bool::from(provided_secret.as_bytes().ct_eq(expected_secret.as_bytes())) { + return Err(( + StatusCode::UNAUTHORIZED, + "Invalid webhook secret".to_string(), + )); + } + + Ok(()) +} + +/// Handle incoming webhook POST to `/api/webhooks/{path}`. +/// +/// This endpoint is **public** (no gateway auth token required) but protected +/// by the per-routine webhook secret sent via the `X-Webhook-Secret` header. +pub async fn webhook_trigger_handler( + State(state): State>, + Path(path): Path, + headers: HeaderMap, +) -> Result, (StatusCode, String)> { + // Rate limit check + if !state.webhook_rate_limiter.check() { + return Err(( + StatusCode::TOO_MANY_REQUESTS, + "Rate limit exceeded. Try again shortly.".to_string(), + )); + } + + let store = state.store.as_ref().ok_or(( + StatusCode::SERVICE_UNAVAILABLE, + "Database not available".to_string(), + ))?; + + // Targeted query instead of loading all routines + let routine = store + .get_webhook_routine_by_path(&path) + .await + .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? + .ok_or(( + StatusCode::NOT_FOUND, + "No routine matches this webhook path".to_string(), + ))?; + + let provided_secret = headers + .get("x-webhook-secret") + .and_then(|v| v.to_str().ok()) + .unwrap_or(""); + + validate_webhook_secret(&routine.trigger, provided_secret)?; + + // Fire through the RoutineEngine so guardrails, run tracking, + // notifications, and FullJob dispatch all work correctly. + let engine = { + let guard = state.routine_engine.read().await; + guard.as_ref().cloned().ok_or(( + StatusCode::SERVICE_UNAVAILABLE, + "Routine engine not available".to_string(), + ))? + }; + + let run_id = engine.fire_webhook(routine.id, &path).await.map_err(|e| { + let status = match &e { + crate::error::RoutineError::NotFound { .. } => StatusCode::NOT_FOUND, + crate::error::RoutineError::Disabled { .. } + | crate::error::RoutineError::Cooldown { .. } + | crate::error::RoutineError::MaxConcurrent { .. } => StatusCode::CONFLICT, + _ => StatusCode::INTERNAL_SERVER_ERROR, + }; + (status, e.to_string()) + })?; + + Ok(Json(serde_json::json!({ + "status": "triggered", + "routine_id": routine.id, + "routine_name": routine.name, + "run_id": run_id, + }))) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Routines with `secret: None` must be rejected with 403. + #[test] + fn test_validate_rejects_missing_secret() { + let trigger = Trigger::Webhook { + path: Some("my-hook".to_string()), + secret: None, + }; + let result = validate_webhook_secret(&trigger, "any-secret"); + let (status, msg) = result.unwrap_err(); + assert_eq!(status, StatusCode::FORBIDDEN); + assert!( + msg.contains("not configured"), + "Error should tell user to configure a secret, got: {msg}" + ); + } + + /// Non-webhook triggers must be rejected with 403. + #[test] + fn test_validate_rejects_non_webhook_trigger() { + let trigger = Trigger::Manual; + let result = validate_webhook_secret(&trigger, "any-secret"); + let (status, _) = result.unwrap_err(); + assert_eq!(status, StatusCode::FORBIDDEN); + } + + /// Correct secret passes validation. + #[test] + fn test_validate_accepts_correct_secret() { + let trigger = Trigger::Webhook { + path: Some("my-hook".to_string()), + secret: Some("s3cret-token".to_string()), + }; + assert!(validate_webhook_secret(&trigger, "s3cret-token").is_ok()); + } + + /// Wrong secret returns 401. + #[test] + fn test_validate_rejects_wrong_secret() { + let trigger = Trigger::Webhook { + path: Some("my-hook".to_string()), + secret: Some("correct-secret".to_string()), + }; + let result = validate_webhook_secret(&trigger, "wrong-secret"); + let (status, msg) = result.unwrap_err(); + assert_eq!(status, StatusCode::UNAUTHORIZED); + assert!(msg.contains("Invalid"), "Expected 'Invalid' in: {msg}"); + } + + /// Empty provided secret returns 401 (not a false positive). + #[test] + fn test_validate_rejects_empty_provided_secret() { + let trigger = Trigger::Webhook { + path: Some("my-hook".to_string()), + secret: Some("real-secret".to_string()), + }; + let result = validate_webhook_secret(&trigger, ""); + let (status, _) = result.unwrap_err(); + assert_eq!(status, StatusCode::UNAUTHORIZED); + } + + /// Constant-time comparison: secrets of different lengths are still rejected + /// (not short-circuited in a way that leaks length info). + #[test] + fn test_validate_rejects_different_length_secret() { + let trigger = Trigger::Webhook { + path: None, + secret: Some("short".to_string()), + }; + let result = validate_webhook_secret(&trigger, "a-much-longer-secret-value"); + let (status, _) = result.unwrap_err(); + assert_eq!(status, StatusCode::UNAUTHORIZED); + } +} diff --git a/src/channels/web/mod.rs b/src/channels/web/mod.rs index bfefc5c4..1fdb4455 100644 --- a/src/channels/web/mod.rs +++ b/src/channels/web/mod.rs @@ -98,6 +98,7 @@ impl GatewayChannel { skill_catalog: None, chat_rate_limiter: server::RateLimiter::new(30, 60), oauth_rate_limiter: server::RateLimiter::new(10, 60), + webhook_rate_limiter: server::RateLimiter::new(10, 60), registry_entries: Vec::new(), cost_guard: None, routine_engine: Arc::new(tokio::sync::RwLock::new(None)), @@ -136,6 +137,7 @@ impl GatewayChannel { skill_catalog: self.state.skill_catalog.clone(), chat_rate_limiter: server::RateLimiter::new(30, 60), oauth_rate_limiter: server::RateLimiter::new(10, 60), + webhook_rate_limiter: server::RateLimiter::new(10, 60), registry_entries: self.state.registry_entries.clone(), cost_guard: self.state.cost_guard.clone(), routine_engine: Arc::clone(&self.state.routine_engine), diff --git a/src/channels/web/server.rs b/src/channels/web/server.rs index ab697951..7b24805c 100644 --- a/src/channels/web/server.rs +++ b/src/channels/web/server.rs @@ -19,6 +19,7 @@ use axum::{ routing::{get, post}, }; use serde::Deserialize; +use sha2::{Digest, Sha256}; use tokio::sync::{mpsc, oneshot}; use tokio_stream::StreamExt; use tower_http::cors::{AllowHeaders, CorsLayer}; @@ -35,7 +36,10 @@ use crate::channels::web::handlers::jobs::{ jobs_events_handler, jobs_list_handler, jobs_prompt_handler, jobs_restart_handler, jobs_summary_handler, }; -use crate::channels::web::handlers::routines::{routines_delete_handler, routines_toggle_handler}; +use crate::channels::web::handlers::routines::{ + routines_delete_handler, routines_detail_handler, routines_list_handler, + routines_summary_handler, routines_toggle_handler, routines_trigger_handler, +}; use crate::channels::web::handlers::skills::{ skills_install_handler, skills_list_handler, skills_remove_handler, skills_search_handler, }; @@ -63,6 +67,16 @@ pub type PromptQueue = Arc< pub type RoutineEngineSlot = Arc>>>; +fn redact_oauth_state_for_logs(state: &str) -> String { + let digest = Sha256::digest(state.as_bytes()); + let mut short_hash = String::with_capacity(12); + for byte in &digest[..6] { + use std::fmt::Write as _; + let _ = write!(&mut short_hash, "{byte:02x}"); + } + format!("sha256:{short_hash}:len={}", state.len()) +} + /// Simple sliding-window rate limiter. /// /// Tracks the number of requests in the current window. Resets when the window expires. @@ -176,6 +190,8 @@ pub struct GatewayState { pub chat_rate_limiter: RateLimiter, /// Rate limiter for OAuth callback endpoints (10 requests per 60 seconds). pub oauth_rate_limiter: RateLimiter, + /// Rate limiter for webhook trigger endpoints (10 requests per 60 seconds). + pub webhook_rate_limiter: RateLimiter, /// Registry catalog entries for the available extensions API. /// Populated at startup from `registry/` manifests, independent of extension manager. pub registry_entries: Vec, @@ -219,7 +235,11 @@ pub async fn start_server( "/oauth/slack/callback", get(slack_relay_oauth_callback_handler), ) - .route("/relay/events", post(relay_events_handler)); + .route("/relay/events", post(relay_events_handler)) + .route( + "/api/webhooks/{path}", + post(crate::channels::web::handlers::webhooks::webhook_trigger_handler), + ); // Protected routes (require auth) let auth_state = AuthState { token: auth_token }; @@ -330,6 +350,7 @@ pub async fn start_server( .route("/", get(index_handler)) .route("/style.css", get(css_handler)) .route("/app.js", get(js_handler)) + .route("/theme-init.js", get(theme_init_handler)) .route("/favicon.ico", get(favicon_handler)) .route("/i18n/index.js", get(i18n_index_handler)) .route("/i18n/en.js", get(i18n_en_handler)) @@ -451,6 +472,16 @@ async fn js_handler() -> impl IntoResponse { ) } +async fn theme_init_handler() -> impl IntoResponse { + ( + [ + (header::CONTENT_TYPE, "application/javascript"), + (header::CACHE_CONTROL, "no-cache"), + ], + include_str!("static/theme-init.js"), + ) +} + async fn favicon_handler() -> impl IntoResponse { ( [ @@ -566,22 +597,35 @@ async fn oauth_callback_handler( } }; - // Strip instance prefix from state for registry lookup. - // Platform nginx sends `state=instance:nonce` but flows are keyed by nonce only. - let lookup_key = oauth_defaults::strip_instance_prefix(&state_param); + let decoded_state = match oauth_defaults::decode_hosted_oauth_state(&state_param) { + Ok(decoded) => decoded, + Err(error) => { + let redacted_state = redact_oauth_state_for_logs(&state_param); + tracing::warn!( + state = %redacted_state, + error = %error, + "OAuth callback received with malformed state" + ); + clear_auth_mode(&state).await; + return oauth_error_page("IronClaw"); + } + }; + let lookup_key = decoded_state.flow_id.clone(); let flow = ext_mgr .pending_oauth_flows() .write() .await - .remove(lookup_key); + .remove(&lookup_key); let flow = match flow { Some(f) => f, None => { + let redacted_state = redact_oauth_state_for_logs(&state_param); + let redacted_lookup_key = redact_oauth_state_for_logs(&lookup_key); tracing::warn!( - state = %state_param, - lookup_key = %lookup_key, + state = %redacted_state, + lookup_key = %redacted_lookup_key, "OAuth callback received with unknown or expired state" ); clear_auth_mode(&state).await; @@ -608,33 +652,29 @@ async fn oauth_callback_handler( } // Exchange the authorization code for tokens. - // Use the platform exchange proxy when configured (keeps client_secret off container), - // otherwise call the provider's token URL directly. - let exchange_proxy_url = std::env::var("IRONCLAW_OAUTH_EXCHANGE_URL").ok(); + // Use the platform exchange proxy when configured, otherwise call the + // provider's token URL directly. + let exchange_proxy_url = oauth_defaults::exchange_proxy_url(); let result: Result<(), String> = async { - let token_response = if let (Some(proxy_url), None) = (&exchange_proxy_url, &flow.resource) - { - // Use the platform exchange proxy when configured and no resource - // parameter is needed. The proxy holds client_secret server-side so - // the container never sees it. MCP flows (resource.is_some()) bypass - // the proxy because it doesn't forward the RFC 8707 resource param. + let token_response = if let Some(proxy_url) = &exchange_proxy_url { let gateway_token = flow.gateway_token.as_deref().unwrap_or_default(); - oauth_defaults::exchange_via_proxy( + oauth_defaults::exchange_via_proxy(oauth_defaults::ProxyTokenExchangeRequest { proxy_url, gateway_token, - &code, - &flow.redirect_uri, - flow.code_verifier.as_deref(), - &flow.access_token_field, - ) + token_url: &flow.token_url, + client_id: &flow.client_id, + client_secret: flow.client_secret.as_deref(), + code: &code, + redirect_uri: &flow.redirect_uri, + code_verifier: flow.code_verifier.as_deref(), + access_token_field: &flow.access_token_field, + extra_token_params: &flow.token_exchange_extra_params, + }) .await .map_err(|e| e.to_string())? } else { - // Direct token exchange: uses exchange_oauth_code_with_resource so MCP - // flows can include the RFC 8707 `resource` parameter to scope the - // issued token to the specific MCP server. - oauth_defaults::exchange_oauth_code_with_resource( + oauth_defaults::exchange_oauth_code_with_params( &flow.token_url, &flow.client_id, flow.client_secret.as_deref(), @@ -642,7 +682,7 @@ async fn oauth_callback_handler( &flow.redirect_uri, flow.code_verifier.as_deref(), &flow.access_token_field, - flow.resource.as_deref(), + &flow.token_exchange_extra_params, ) .await .map_err(|e| e.to_string())? @@ -669,10 +709,8 @@ async fn oauth_callback_handler( .await .map_err(|e| e.to_string())?; - // For MCP OAuth flows (identified by resource field), persist the - // client_id so token refresh works without re-authentication. - // The CLI flow stores this in authorize_mcp_server(); the gateway - // callback must do the same. + // Persist the client_id for flows that need it after the session ends + // (for example DCR-based MCP refresh). if let Some(ref client_id_secret) = flow.client_id_secret_name { let params = crate::secrets::CreateSecretParams::new(client_id_secret, &flow.client_id) .with_provider(flow.provider.as_ref().cloned().unwrap_or_default()); @@ -1784,14 +1822,53 @@ async fn memory_write_handler( "Workspace not available".to_string(), ))?; - workspace - .write(&req.path, &req.content) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + // Route through layer-aware methods when a layer is specified + if let Some(ref layer_name) = req.layer { + let result = if req.append { + workspace + .append_to_layer(layer_name, &req.path, &req.content, req.force) + .await + } else { + workspace + .write_to_layer(layer_name, &req.path, &req.content, req.force) + .await + } + .map_err(|e| { + use crate::error::WorkspaceError; + let status = match &e { + WorkspaceError::LayerNotFound { .. } => StatusCode::BAD_REQUEST, + WorkspaceError::LayerReadOnly { .. } => StatusCode::FORBIDDEN, + WorkspaceError::PrivacyRedirectFailed => StatusCode::UNPROCESSABLE_ENTITY, + _ => StatusCode::INTERNAL_SERVER_ERROR, + }; + (status, e.to_string()) + })?; + return Ok(Json(MemoryWriteResponse { + path: req.path, + status: "written", + redirected: Some(result.redirected), + actual_layer: Some(result.actual_layer), + })); + } + + // Non-layer path: honor the append field + if req.append { + workspace + .append(&req.path, &req.content) + .await + .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + } else { + workspace + .write(&req.path, &req.content) + .await + .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + } Ok(Json(MemoryWriteResponse { path: req.path, status: "written", + redirected: None, + actual_layer: None, })) } @@ -2266,7 +2343,7 @@ async fn extensions_setup_handler( "Extension manager not available (secrets store required)".to_string(), ))?; - let secrets = ext_mgr + let setup = ext_mgr .get_setup_schema(&name) .await .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; @@ -2282,7 +2359,8 @@ async fn extensions_setup_handler( Ok(Json(ExtensionSetupResponse { name, kind, - secrets, + secrets: setup.secrets, + fields: setup.fields, })) } @@ -2300,7 +2378,7 @@ async fn extensions_setup_submit_handler( // through to the LLM instead of being intercepted as a token. clear_auth_mode(&state).await; - match ext_mgr.configure(&name, &req.secrets).await { + match ext_mgr.configure(&name, &req.secrets, &req.fields).await { Ok(result) => { let mut resp = if result.verification.is_some() || result.activated { ActionResponse::ok(result.message) @@ -2308,6 +2386,9 @@ async fn extensions_setup_submit_handler( ActionResponse::fail(result.message) }; resp.activated = Some(result.activated); + if result.restart_required || !result.activated { + resp.needs_restart = Some(true); + } resp.auth_url = result.auth_url.clone(); resp.verification = result.verification.clone(); resp.instructions = result.verification.as_ref().map(|v| v.instructions.clone()); @@ -2373,164 +2454,6 @@ async fn pairing_approve_handler( } } -// --- Routines handlers --- - -async fn routines_list_handler( - State(state): State>, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let routines = store - .list_all_routines() - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - let items: Vec = routines.iter().map(RoutineInfo::from_routine).collect(); - - Ok(Json(RoutineListResponse { routines: items })) -} - -async fn routines_summary_handler( - State(state): State>, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let routines = store - .list_all_routines() - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - let total = routines.len() as u64; - let enabled = routines.iter().filter(|r| r.enabled).count() as u64; - let disabled = total - enabled; - let failing = routines - .iter() - .filter(|r| r.consecutive_failures > 0) - .count() as u64; - - let today_start = chrono::Utc::now() - .date_naive() - .and_hms_opt(0, 0, 0) - .map(|dt| dt.and_utc()); - let runs_today = if let Some(start) = today_start { - routines - .iter() - .filter(|r| r.last_run_at.is_some_and(|ts| ts >= start)) - .count() as u64 - } else { - 0 - }; - - Ok(Json(RoutineSummaryResponse { - total, - enabled, - disabled, - failing, - runs_today, - })) -} - -async fn routines_detail_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let store = state.store.as_ref().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Database not available".to_string(), - ))?; - - let routine_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; - - let routine = store - .get_routine(routine_id) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? - .ok_or((StatusCode::NOT_FOUND, "Routine not found".to_string()))?; - - let runs = store - .list_routine_runs(routine_id, 20) - .await - .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - - let recent_runs: Vec = runs - .iter() - .map(|run| RoutineRunInfo { - id: run.id, - trigger_type: run.trigger_type.clone(), - started_at: run.started_at.to_rfc3339(), - completed_at: run.completed_at.map(|dt| dt.to_rfc3339()), - status: format!("{:?}", run.status), - result_summary: run.result_summary.clone(), - tokens_used: run.tokens_used, - job_id: run.job_id, - }) - .collect(); - let routine_info = RoutineInfo::from_routine(&routine); - - Ok(Json(RoutineDetailResponse { - id: routine.id, - name: routine.name.clone(), - description: routine.description.clone(), - enabled: routine.enabled, - trigger_type: routine_info.trigger_type, - trigger_raw: routine_info.trigger_raw, - trigger_summary: routine_info.trigger_summary, - trigger: serde_json::to_value(&routine.trigger).unwrap_or_default(), - action: serde_json::to_value(&routine.action).unwrap_or_default(), - guardrails: serde_json::to_value(&routine.guardrails).unwrap_or_default(), - notify: serde_json::to_value(&routine.notify).unwrap_or_default(), - last_run_at: routine.last_run_at.map(|dt| dt.to_rfc3339()), - next_fire_at: routine.next_fire_at.map(|dt| dt.to_rfc3339()), - run_count: routine.run_count, - consecutive_failures: routine.consecutive_failures, - created_at: routine.created_at.to_rfc3339(), - recent_runs, - })) -} - -async fn routines_trigger_handler( - State(state): State>, - Path(id): Path, -) -> Result, (StatusCode, String)> { - let engine = { - let guard = state.routine_engine.read().await; - guard.as_ref().cloned().ok_or(( - StatusCode::SERVICE_UNAVAILABLE, - "Routine engine not available".to_string(), - ))? - }; - - let routine_id = Uuid::parse_str(&id) - .map_err(|_| (StatusCode::BAD_REQUEST, "Invalid routine ID".to_string()))?; - - let run_id = engine - .fire_manual(routine_id, Some(&state.user_id)) - .await - .map_err(|e| { - let status = match &e { - crate::error::RoutineError::NotFound { .. } => StatusCode::NOT_FOUND, - crate::error::RoutineError::NotAuthorized { .. } => StatusCode::FORBIDDEN, - crate::error::RoutineError::Disabled { .. } - | crate::error::RoutineError::MaxConcurrent { .. } => StatusCode::CONFLICT, - _ => StatusCode::INTERNAL_SERVER_ERROR, - }; - (status, e.to_string()) - })?; - - Ok(Json(serde_json::json!({ - "status": "triggered", - "routine_id": routine_id, - "run_id": run_id, - }))) -} - async fn routines_runs_handler( State(state): State>, Path(id): Path, @@ -2960,6 +2883,7 @@ mod tests { scheduler: None, chat_rate_limiter: RateLimiter::new(30, 60), oauth_rate_limiter: RateLimiter::new(10, 60), + webhook_rate_limiter: RateLimiter::new(10, 60), registry_entries: vec![], cost_guard: None, routine_engine: Arc::new(tokio::sync::RwLock::new(None)), @@ -3311,7 +3235,7 @@ mod tests { secrets, sse_sender: None, gateway_token: None, - resource: None, + token_exchange_extra_params: std::collections::HashMap::new(), client_id_secret_name: None, created_at, }; @@ -3379,7 +3303,7 @@ mod tests { secrets, sse_sender: Some(sender), gateway_token: None, - resource: None, + token_exchange_extra_params: std::collections::HashMap::new(), client_id_secret_name: None, created_at, }; @@ -3482,7 +3406,7 @@ mod tests { secrets, sse_sender: None, gateway_token: None, - resource: None, + token_exchange_extra_params: std::collections::HashMap::new(), client_id_secret_name: None, // Expired โ€” handler will reject after lookup (no network I/O) created_at, @@ -3534,6 +3458,85 @@ mod tests { ); } + #[tokio::test] + async fn test_oauth_callback_accepts_versioned_hosted_state() { + use axum::body::Body; + use tower::ServiceExt; + + let secrets: Arc = + Arc::new(crate::secrets::InMemorySecretsStore::new(Arc::new( + crate::secrets::SecretsCrypto::new(secrecy::SecretString::from( + TEST_GATEWAY_CRYPTO_KEY.to_string(), + )) + .expect("crypto"), + ))); + let (ext_mgr, _wasm_tools_dir, _wasm_channels_dir) = test_ext_mgr(secrets.clone()); + + let Some(created_at) = expired_flow_created_at() else { + eprintln!("Skipping versioned OAuth state test: monotonic uptime below expiry window"); + return; + }; + let flow = crate::cli::oauth_defaults::PendingOAuthFlow { + extension_name: "test_tool".to_string(), + display_name: "Test Tool".to_string(), + token_url: "https://example.com/token".to_string(), + client_id: "client123".to_string(), + client_secret: None, + redirect_uri: "https://example.com/oauth/callback".to_string(), + code_verifier: None, + access_token_field: "access_token".to_string(), + secret_name: "test_token".to_string(), + provider: None, + validation_endpoint: None, + scopes: vec![], + user_id: "test".to_string(), + secrets, + sse_sender: None, + gateway_token: None, + token_exchange_extra_params: std::collections::HashMap::new(), + client_id_secret_name: None, + created_at, + }; + + ext_mgr + .pending_oauth_flows() + .write() + .await + .insert("test_nonce".to_string(), flow); + + let state = test_gateway_state(Some(ext_mgr.clone())); + let app = test_oauth_router(state); + let versioned_state = + crate::cli::oauth_defaults::encode_hosted_oauth_state("test_nonce", Some("myinstance")); + + let req = axum::http::Request::builder() + .uri(format!( + "/oauth/callback?code=fake_code&state={}", + urlencoding::encode(&versioned_state) + )) + .body(Body::empty()) + .expect("request"); + + let resp = ServiceExt::>::oneshot(app, req) + .await + .expect("response"); + assert_eq!(resp.status(), StatusCode::OK); + + let body = axum::body::to_bytes(resp.into_body(), 1024 * 64) + .await + .expect("body"); + let html = String::from_utf8_lossy(&body); + assert!(html.contains("Authorization Failed")); + assert!( + ext_mgr + .pending_oauth_flows() + .read() + .await + .get("test_nonce") + .is_none() + ); + } + // --- Slack relay OAuth CSRF tests --- fn test_relay_oauth_router(state: Arc) -> Router { diff --git a/src/channels/web/static/app.js b/src/channels/web/static/app.js index bc23d68c..075aa7cc 100644 --- a/src/channels/web/static/app.js +++ b/src/channels/web/static/app.js @@ -1,5 +1,69 @@ // IronClaw Web Gateway - Client +// --- Theme Management (dark / light / system) --- +// Icon switching is handled by pure CSS via data-theme-mode on . + +function getSystemTheme() { + return window.matchMedia('(prefers-color-scheme: light)').matches ? 'light' : 'dark'; +} + +const VALID_THEME_MODES = { dark: true, light: true, system: true }; + +function getThemeMode() { + const stored = localStorage.getItem('ironclaw-theme'); + return (stored && VALID_THEME_MODES[stored]) ? stored : 'system'; +} + +function resolveTheme(mode) { + return mode === 'system' ? getSystemTheme() : mode; +} + +function applyTheme(mode) { + const resolved = resolveTheme(mode); + document.documentElement.setAttribute('data-theme', resolved); + document.documentElement.setAttribute('data-theme-mode', mode); + const titleKeys = { dark: 'theme.tooltipDark', light: 'theme.tooltipLight', system: 'theme.tooltipSystem' }; + const btn = document.getElementById('theme-toggle'); + if (btn) btn.title = (typeof I18n !== 'undefined' && titleKeys[mode]) ? I18n.t(titleKeys[mode]) : ('Theme: ' + mode); + const announce = document.getElementById('theme-announce'); + if (announce) announce.textContent = (typeof I18n !== 'undefined') ? I18n.t('theme.announce', { mode: mode }) : ('Theme: ' + mode); +} + +function toggleTheme() { + const cycle = { dark: 'light', light: 'system', system: 'dark' }; + const current = getThemeMode(); + const next = cycle[current] || 'dark'; + localStorage.setItem('ironclaw-theme', next); + applyTheme(next); +} + +// Apply theme immediately (FOUC prevention is done via inline script in , +// but we call again here to ensure tooltip is set after DOM is ready). +applyTheme(getThemeMode()); + +// Delay enabling theme transition to avoid flash on initial load. +requestAnimationFrame(function() { + requestAnimationFrame(function() { + document.body.classList.add('theme-transition'); + }); +}); + +// Listen for OS theme changes โ€” only re-apply when in 'system' mode. +const mql = window.matchMedia('(prefers-color-scheme: light)'); +const onSchemeChange = function() { + if (getThemeMode() === 'system') { + applyTheme('system'); + } +}; +if (mql.addEventListener) { + mql.addEventListener('change', onSchemeChange); +} else if (mql.addListener) { + mql.addListener(onSchemeChange); +} + +// Bind theme toggle button (CSP-compliant โ€” no inline onclick). +document.getElementById('theme-toggle').addEventListener('click', toggleTheme); + let token = ''; let eventSource = null; let logEventSource = null; @@ -100,6 +164,30 @@ document.getElementById('token-input').addEventListener('keydown', (e) => { if (e.key === 'Enter') authenticate(); }); +// --- Static element event bindings (CSP-compliant, no inline handlers) --- +document.getElementById('auth-connect-btn').addEventListener('click', () => authenticate()); +document.getElementById('restart-overlay').addEventListener('click', () => cancelRestart()); +document.getElementById('restart-close-btn').addEventListener('click', () => cancelRestart()); +document.getElementById('restart-cancel-btn').addEventListener('click', () => cancelRestart()); +document.getElementById('restart-confirm-btn').addEventListener('click', () => confirmRestart()); +document.getElementById('language-btn').addEventListener('click', () => toggleLanguageMenu()); +// Language option clicks handled by delegated data-action="switch-language" handler. +document.getElementById('restart-btn').addEventListener('click', () => triggerRestart()); +document.getElementById('thread-new-btn').addEventListener('click', () => createNewThread()); +document.getElementById('thread-toggle-btn').addEventListener('click', () => toggleThreadSidebar()); +document.getElementById('assistant-thread').addEventListener('click', () => switchToAssistant()); +document.getElementById('send-btn').addEventListener('click', () => sendMessage()); +document.getElementById('memory-edit-btn').addEventListener('click', () => startMemoryEdit()); +document.getElementById('memory-save-btn').addEventListener('click', () => saveMemoryEdit()); +document.getElementById('memory-cancel-btn').addEventListener('click', () => cancelMemoryEdit()); +document.getElementById('logs-server-level').addEventListener('change', function() { setServerLogLevel(this.value); }); +document.getElementById('logs-pause-btn').addEventListener('click', () => toggleLogsPause()); +document.getElementById('logs-clear-btn').addEventListener('click', () => clearLogs()); +document.getElementById('wasm-install-btn').addEventListener('click', () => installWasmExtension()); +document.getElementById('mcp-add-btn').addEventListener('click', () => addMcpServer()); +document.getElementById('skill-search-btn').addEventListener('click', () => searchClawHub()); +document.getElementById('skill-install-btn').addEventListener('click', () => installSkillFromForm()); + // Auto-authenticate from URL param or saved session (function autoAuth() { const params = new URLSearchParams(window.location.search); @@ -2703,16 +2791,18 @@ function removeExtension(name) { function showConfigureModal(name) { apiFetch('/api/extensions/' + encodeURIComponent(name) + '/setup') .then((setup) => { - if (!setup.secrets || setup.secrets.length === 0) { + const secrets = Array.isArray(setup.secrets) ? setup.secrets : []; + const setupFields = Array.isArray(setup.fields) ? setup.fields : []; + if (secrets.length === 0 && setupFields.length === 0) { showToast('No configuration needed for ' + name, 'info'); return; } - renderConfigureModal(name, setup.secrets); + renderConfigureModal(name, secrets, setupFields); }) .catch((err) => showToast('Failed to load setup: ' + err.message, 'error')); } -function renderConfigureModal(name, secrets) { +function renderConfigureModal(name, secrets, setupFields) { closeConfigureModal(); const overlay = document.createElement('div'); overlay.className = 'configure-overlay'; @@ -2785,7 +2875,46 @@ function renderConfigureModal(name, secrets) { field.appendChild(inputRow); form.appendChild(field); - fields.push({ name: secret.name, input: input }); + fields.push({ kind: 'secret', name: secret.name, input: input }); + } + + for (const setupField of setupFields) { + const field = document.createElement('div'); + field.className = 'configure-field'; + + const label = document.createElement('label'); + label.textContent = setupField.prompt; + if (setupField.optional) { + const opt = document.createElement('span'); + opt.className = 'field-optional'; + opt.textContent = I18n.t('config.optional'); + label.appendChild(opt); + } + field.appendChild(label); + + const inputRow = document.createElement('div'); + inputRow.className = 'configure-input-row'; + + const input = document.createElement('input'); + input.type = setupField.input_type === 'password' ? 'password' : 'text'; + input.name = setupField.name; + input.placeholder = setupField.provided ? I18n.t('config.alreadySet') : ''; + input.addEventListener('keydown', (e) => { + if (e.key === 'Enter') submitConfigureModal(name, fields); + }); + inputRow.appendChild(input); + + if (setupField.provided) { + const badge = document.createElement('span'); + badge.className = 'field-provided'; + badge.textContent = '\u2713'; + badge.title = I18n.t('config.alreadyConfigured'); + inputRow.appendChild(badge); + } + + field.appendChild(inputRow); + form.appendChild(field); + fields.push({ kind: 'field', name: setupField.name, input: input }); } modal.appendChild(form); @@ -2927,9 +3056,16 @@ function startTelegramAutoVerify(name, fields) { function submitConfigureModal(name, fields, options) { options = options || {}; const secrets = {}; + const setupFields = {}; for (const f of fields) { - if (f.input.value.trim()) { - secrets[f.name] = f.input.value.trim(); + const value = f.input.value.trim(); + if (!value) { + continue; + } + if (f.kind === 'secret') { + secrets[f.name] = value; + } else { + setupFields[f.name] = value; } } @@ -2946,7 +3082,7 @@ function submitConfigureModal(name, fields, options) { apiFetch('/api/extensions/' + encodeURIComponent(name) + '/setup', { method: 'POST', - body: { secrets }, + body: { secrets, fields: setupFields }, }) .then((res) => { if (res.success) { @@ -2976,6 +3112,8 @@ function submitConfigureModal(name, fields, options) { showToast('Opening OAuth authorization for ' + name, 'info'); openOAuthUrl(res.auth_url); refreshCurrentSettingsTab(); + } else if (res.needs_restart) { + showToast('Configured ' + name + '. Restart IronClaw to apply all changes.', 'info'); } // For non-OAuth success: the server always broadcasts auth_completed SSE, // which will show the toast and refresh extensions โ€” no need to do it here too. @@ -3854,7 +3992,6 @@ function renderRoutineDetail(routine) { + '
' + escapeHtml(JSON.stringify(routine.trigger, null, 2)) + '
'; } - // Action config html += '

Action

' + '
' + escapeHtml(JSON.stringify(routine.action, null, 2)) + '
'; @@ -3925,7 +4062,7 @@ function formatRelativeTime(isoString) { const absDiff = Math.abs(diffMs); const future = diffMs < 0; - if (absDiff < 60000) + if (absDiff < 60000) return future ? I18n.t('time.lessThan1MinuteFromNow') : I18n.t('time.lessThan1MinuteAgo'); if (absDiff < 3600000) { const m = Math.floor(absDiff / 60000); @@ -4873,7 +5010,14 @@ function renderStructuredSettingsRow(def, value, activeValue) { inputWrap.style.gap = '8px'; var ariaLabel = I18n.t(def.label) + (def.description ? '. ' + I18n.t(def.description) : ''); - var placeholderText = activeValue ? I18n.t('settings.envValue', { value: activeValue }) : (def.placeholder || I18n.t('settings.envDefault')); + function formatSettingValue(raw) { + if (Array.isArray(raw)) return raw.join(', '); + if (raw === null || raw === undefined) return ''; + return String(raw); + } + + var activeValueText = formatSettingValue(activeValue); + var placeholderText = activeValueText ? I18n.t('settings.envValue', { value: activeValueText }) : (def.placeholder || I18n.t('settings.envDefault')); if (def.type === 'boolean') { var boolSel = document.createElement('select'); @@ -4945,6 +5089,26 @@ function renderStructuredSettingsRow(def, value, activeValue) { }; })(def.key, numInp)); inputWrap.appendChild(numInp); + } else if (def.type === 'list') { + var listInp = document.createElement('input'); + listInp.type = 'text'; + listInp.className = 'settings-input'; + listInp.setAttribute('aria-label', ariaLabel); + var listValue = ''; + if (Array.isArray(value)) listValue = value.join(', '); + else if (typeof value === 'string') listValue = value; + listInp.value = listValue; + if (!listValue) listInp.placeholder = placeholderText; + listInp.addEventListener('change', (function(k, el) { + return function() { + if (el.value.trim() === '') return saveSetting(k, null); + var items = el.value.split(/[\n,]/).map(function(item) { + return item.trim(); + }).filter(Boolean); + saveSetting(k, items); + }; + })(def.key, listInp)); + inputWrap.appendChild(listInp); } else { var textInp = document.createElement('input'); textInp.type = 'text'; diff --git a/src/channels/web/static/i18n/en.js b/src/channels/web/static/i18n/en.js index 1369b485..6029075d 100644 --- a/src/channels/web/static/i18n/en.js +++ b/src/channels/web/static/i18n/en.js @@ -24,6 +24,12 @@ I18n.register('en', { 'restart.progressSubtitle': 'Please wait for the process to restart...', 'restart.checkLogs': 'Check the Logs tab for details after restart completes.', + // Theme + 'theme.tooltipDark': 'Theme: Dark (click for Light)', + 'theme.tooltipLight': 'Theme: Light (click for System)', + 'theme.tooltipSystem': 'Theme: System (click for Dark)', + 'theme.announce': 'Theme: {mode}', + // Tabs 'tab.chat': 'Chat', 'tab.memory': 'Memory', diff --git a/src/channels/web/static/i18n/zh-CN.js b/src/channels/web/static/i18n/zh-CN.js index 6262b562..480724c9 100644 --- a/src/channels/web/static/i18n/zh-CN.js +++ b/src/channels/web/static/i18n/zh-CN.js @@ -24,6 +24,12 @@ I18n.register('zh-CN', { 'restart.progressSubtitle': '่ฏท็ญ‰ๅพ…่ฟ›็จ‹้‡ๅฏ...', 'restart.checkLogs': '้‡ๅฏๅฎŒๆˆๅŽ๏ผŒ่ฏทๆŸฅ็œ‹ๆ—ฅๅฟ—ๆ ‡็ญพ้กตไบ†่งฃ่ฏฆๆƒ…ใ€‚', + // ไธป้ข˜ + 'theme.tooltipDark': 'ไธป้ข˜๏ผšๆทฑ่‰ฒ๏ผˆ็‚นๅ‡ปๅˆ‡ๆขๆต…่‰ฒ๏ผ‰', + 'theme.tooltipLight': 'ไธป้ข˜๏ผšๆต…่‰ฒ๏ผˆ็‚นๅ‡ปๅˆ‡ๆข่ทŸ้š็ณป็ปŸ๏ผ‰', + 'theme.tooltipSystem': 'ไธป้ข˜๏ผš่ทŸ้š็ณป็ปŸ๏ผˆ็‚นๅ‡ปๅˆ‡ๆขๆทฑ่‰ฒ๏ผ‰', + 'theme.announce': 'ไธป้ข˜๏ผš{mode}', + // ๆ ‡็ญพ้กต 'tab.chat': '่Šๅคฉ', 'tab.memory': '่ฎฐๅฟ†', diff --git a/src/channels/web/static/index.html b/src/channels/web/static/index.html index b342cb53..113d144e 100644 --- a/src/channels/web/static/index.html +++ b/src/channels/web/static/index.html @@ -25,6 +25,7 @@ integrity="sha384-pN9zSKOnTZwXRtYZAu0PBPEgR2B7DOC1aeLxQ33oJ0oy5iN1we6gm57xldM2irDG" crossorigin="anonymous" > + @@ -109,6 +110,18 @@ + +