mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
* feat: complete multi-tenant isolation — per-user budgets, model selection, heartbeat cycling Finishes the remaining isolation work from phases 2–4 of #59: Phase 2 (DB scoping): Fix /status and /list commands to use _for_user DB variants instead of global queries that leaked cross-user job data. Phase 3 (Runtime isolation): Per-user workspace in routine engine's spawn_fire so lightweight routines run in the correct user context. Per-user daily cost tracking in CostGuard with configurable budget via MAX_COST_PER_USER_PER_DAY_CENTS. Multi-user heartbeat that cycles through all users with routines, auto-detected from GATEWAY_USER_TOKENS. Phase 4 (Provider/tools): Per-user model selection via preferred_model setting — looked up from SettingsStore on first iteration, threaded through ReasoningContext.model_override to CompletionRequest. Works with providers that support per-request model overrides (NearAI). Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: use selected_model setting key to match /model command persistence The dispatcher was reading "preferred_model" but the /model command (merged from staging) persists to "selected_model". Since set_setting is already per-user scoped, using the same key makes /model work as the per-user model override in multi-tenant mode. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: heartbeat hygiene, /model multi-tenant guard, RigAdapter model override Three follow-up fixes for multi-tenant isolation: 1. Multi-user heartbeat now runs memory hygiene per user before each heartbeat check, matching single-user heartbeat behavior. 2. /model command in multi-tenant mode only persists to per-user settings (selected_model) without calling set_model() on the shared LlmProvider. The per-request model_override in the dispatcher reads from the same setting. Added multi_tenant flag to AgentConfig (auto-detected from GATEWAY_USER_TOKENS). 3. RigAdapter now supports per-request model overrides by injecting the model name into rig-core's additional_params. OpenAI/Anthropic/Ollama API servers use last-key-wins for duplicate JSON keys, so the override takes effect via serde's flatten serialization order. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address PR review — cost model attribution, heartbeat concurrency, pruning Fixes from review comments on #1614: - Cost tracking now uses the override model name (not active_model_name) when a per-user model override is active, for accurate attribution. - Multi-user heartbeat runs per-user checks concurrently via JoinSet instead of sequentially, preventing one slow user from blocking others. - Per-user failure counts tracked independently; users exceeding max_failures are skipped (matching single-user semantics). - per_user_daily_cost HashMap pruned on day rollover to prevent unbounded growth in long-lived deployments. - Doc comment fixed: says "routines" not "active routines". Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: /status ownership, model persistence scoping, heartbeat robustness Addresses second round of PR review on #1614: - /status <job_id> DB path now validates job.user_id == requesting user before returning data (was missing ownership check, security fix). - persist_selected_model takes user_id param instead of owner_id, and skips .env/TOML writes in multi-tenant mode (these are shared global files). handle_system_command now receives user_id from caller. - JoinSet collection handles Err(JoinError) explicitly instead of silently dropping panicked tasks. - Notification forwarder extracts owner_id from response metadata in multi-tenant mode for per-user routing instead of broadcasting to the agent owner. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: cost pricing, fire_manual workspace, heartbeat concurrency cap Round 3 review fixes: - Cost tracking passes None for cost_per_token when model override is active, letting CostGuard look up pricing by model name instead of using the default provider's rates (serrrfirat). - fire_manual() now uses per-user workspace, matching spawn_fire() pattern (serrrfirat). - Removed MULTI_TENANT env var — multi-tenant mode is auto-detected solely from GATEWAY_USER_TOKENS presence (serrrfirat + Copilot). - Multi-user heartbeat capped at 8 concurrent tasks to avoid flooding the LLM provider (serrrfirat + Copilot). - Fixed inject_model_override doc comment accuracy (Copilot). - Added comment explaining multi-tenant notification routing priority (Copilot). Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: user-scoped webhook endpoint for multi-tenant isolation Adds POST /api/webhooks/u/{user_id}/{path} — a user-scoped webhook endpoint that filters the routine lookup by user_id, preventing cross-user webhook triggering when paths collide. The existing /api/webhooks/{path} endpoint remains unchanged for backward compatibility in single-user deployments. Changes: - get_webhook_routine_by_path gains user_id: Option<&str> param - Both postgres and libsql implementations add AND user_id = ? filter when user_id is provided - New webhook_trigger_user_scoped_handler extracts (user_id, path) from URL and passes to shared fire_webhook_inner logic - Route registered on public router (webhooks are called by external services that can't send bearer tokens) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(db): add UserStore trait with users, api_tokens, invitations tables Foundation for DB-backed user management (#1605): - UserRecord, ApiTokenRecord, InvitationRecord types in db/mod.rs - UserStore sub-trait (17 methods) added to Database supertrait - PostgreSQL migration V14__users.sql (users, api_tokens, invitations) - libSQL schema + incremental migration V14 - Full implementations for both PgBackend (via Store delegation) and LibSqlBackend (direct SQL in libsql/users.rs) - authenticate_token JOINs api_tokens+users with active/non-revoked checks; has_any_users for bootstrap detection Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(web): DB-backed auth, user/token/invitation API handlers Adds the web gateway layer for DB-backed user management (#1605): Auth refactor: - CombinedAuthState wraps env-var tokens (MultiAuthState) + optional DbAuthenticator for DB-backed token lookup with LRU cache (60s TTL, 1024 max entries) - auth_middleware tries env-var tokens first, then DB fallback - From<MultiAuthState> impl for backward compatibility - main.rs wires with_db_auth when database is available API handlers (12 new endpoints): - /api/admin/users — CRUD: create, list, detail, update, suspend, activate - /api/tokens — create (returns plaintext once), list, revoke - /api/invitations — create, list, accept (creates user + first token) Token creation: 32 random bytes → hex plaintext, SHA-256 hash stored. Invitation accept: validates hash + pending + not expired, creates user record and first API token atomically. All test files updated for CombinedAuthState type change. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: startup env-var user migration + UserStore integration tests Completes the DB-backed user management feature (#1605): - Startup migration: when GATEWAY_USER_TOKENS is set and the users table is empty, inserts env-var users + hashed tokens into DB. Logs deprecation notice when DB already has users. - hash_token made pub for reuse in migration code. - 10 integration tests for UserStore (libsql file-backed): - has_any_users bootstrap detection - create/get/get_by_email/list/update user lifecycle - token create → authenticate → revoke → reject cycle - suspended user tokens rejected - wrong-user token revoke returns false - invitation create → accept → user created - record_login and record_token_usage timestamps - libSQL migration: removed FK constraints from V14 (incompatible with execute_batch inside transactions). Tables in both base SCHEMA and incremental migration for fresh and existing databases. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * refactor: remove GATEWAY_USER_TOKENS, fix review feedback GATEWAY_USER_TOKENS never went to production — replaced entirely by DB-backed user management via /api/admin/users and /api/tokens. Removed: - UserTokenConfig struct and GATEWAY_USER_TOKENS env var parsing - user_tokens field from GatewayConfig - GatewayChannel::new_multi_auth() constructor - Env-var user migration block in main.rs (~90 lines) - multi_tenant auto-detection from GATEWAY_USER_TOKENS (now runtime via db.has_any_users() in app.rs) Review fixes (zmanian): - User ID generation: UUID instead of display-name derivation (#1) - Invitation accept moved to public router (no auth needed) (#3) - libSQL get_invitation_by_hash aligned with postgres: filters status='pending' AND expires_at > now (#4) - UUID parse: returns DatabaseError::Serialization instead of unwrap_or_default (#7) - PostgreSQL SELECT * replaced with explicit column lists (#8) - Sort order aligned (both backends use DESC) (#6) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: add role-based access control (admin/member) Adds a `role` field (admin|member) to user management: Schema: - `role TEXT NOT NULL DEFAULT 'member'` added to users table in both PostgreSQL V14 migration and libSQL schema/incremental migration - UserRecord gains `role: String` field - UserIdentity gains `role: String` field, populated from DB in DbAuthenticator and defaulting to "admin" for single-user mode Access control: - AdminUser extractor: returns 403 Forbidden if role != "admin" - /api/admin/users/* handlers: require AdminUser (create, list, detail, update, suspend, activate) - POST /api/invitations: requires AdminUser (only admins can invite) - User creation accepts optional "role" param (defaults to "member") - Invitation acceptance creates users with "member" role Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(web): add Users admin tab to web UI Adds a Users tab to the web gateway UI for managing users, tokens, and roles without needing direct API calls. Features: - User list table with ID, name, email, role, status, created date - Create user form with display name, email, role selector - Suspend/activate actions per user - Create API token for any user (shows plaintext once with copy button) - Role badges (admin highlighted, member muted) - Non-admin users see "Admin access required" message - Keyboard shortcut: Cmd/Ctrl+5 switches to Users tab CSS: - Reuses routines-table styles for the user list - Badge, token-display, btn-small, btn-danger, btn-primary components Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: move Users to Settings subtab, bootstrap admin user on first run - Moved Users from top-level tab to Settings sidebar subtab (under Skills, before Theme toggle) - On first startup with empty users table, automatically creates an admin user from GATEWAY_USER_ID config with a corresponding API token from GATEWAY_AUTH_TOKEN. This ensures the owner appears in the Users panel immediately. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: user creation shows token, + Token works, no password save popup Three UI/UX fixes: 1. Create user now generates an initial API token and shows it in a copy-able banner instead of triggering the browser's password save dialog. Uses autocomplete="off" and type="text" for email field. 2. "+ Token" button works: exposed createTokenForUser/suspendUser/ activateUser on window for inline onclick handlers in dynamically generated table rows. Token creation uses showTokenBanner helper. 3. Admin token creation: POST /api/tokens now accepts optional "user_id" field when the requesting user is admin, allowing token creation for other users from the Users panel. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: use event delegation for user action buttons (CSP compliance) Inline onclick handlers are blocked by the Content-Security-Policy (script-src 'self' without 'unsafe-inline'). Switched to data-action attributes with a delegated click listener on the users table. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: add i18n for Users subtab, show login link on user creation - Added 'settings.users' i18n key for English and Chinese - Token banner now shows a full login link (domain/?token=xxx) with a Copy Link button, plus the raw token below - Login link works automatically via existing ?token= auto-auth Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: token hash mismatch — hash hex string, not raw bytes Critical auth bug: token creation hashed the raw 32 bytes (hasher.update(token_bytes)) but authentication hashed the hex-encoded string (hash_token(candidate) where candidate is the hex string the user sends). This meant newly created tokens could never authenticate. Fixed all 4 token creation sites (users, tokens, invitations create, invitations accept) to use hash_token(&plaintext_token) which hashes the hex string consistently with the auth lookup path. Removed now-unused sha2::Digest imports from handlers. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * refactor: remove invitation system The invitation flow is redundant — admin create user already generates a token and shows a login link. Invitations add complexity without value until email integration exists. Removed: - InvitationRecord struct and 4 UserStore trait methods - invitations table from V14 migration (postgres + both libsql schemas) - PostgreSQL Store methods (create/get/accept/list invitations) - libSQL UserStore invitation methods + row_to_invitation helper - invitations.rs handler file (212 lines) - /api/invitations routes (create, list, accept) - test_invitation_lifecycle test Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: user deletion, self-service profile, per-user job limits, usage API Four multi-tenancy improvements: 1. User deletion cascade (DELETE /api/admin/users/{id}): Deletes user and all data across 11 user-scoped tables (settings, secrets, routines, memory, jobs, conversations, etc.). Admin only. 2. Self-service profile (GET/PATCH /api/profile): Users can read and update their own display_name and metadata without admin privileges. 3. Per-user job concurrency (MAX_JOBS_PER_USER env var): Scheduler checks active_jobs_for(user_id) before dispatch. Prevents one user from exhausting all job slots. 4. Usage reporting (GET /api/admin/usage?user_id=X&period=day|week|month): Aggregates LLM costs from llm_calls via agent_jobs.user_id. Returns per-user, per-model breakdown of calls, tokens, and cost. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: add TenantCtx for compile-time tenant isolation Implements zmanian's architectural proposal from #1614 review: two-tier scoped database access (TenantScope/AdminScope) so handler code cannot accidentally bypass tenant scoping. TenantScope (default): wraps user_id + Arc<dyn Database>, auto-binds user_id on every operation. ID-based lookups return None for cross- tenant resources. No escape hatch — forgetting to scope is a compile error. AdminScope (explicit opt-in): cross-tenant access for system-level components (heartbeat, routine engine, self-repair, scheduler, worker). TenantCtx bundles TenantScope + workspace + cost guard + per-user rate limiting. Constructed once per request in handle_message, threaded through all command handlers and ChatDelegate. Key changes: - New src/tenant.rs (~920 lines): TenantScope, AdminScope, TenantCtx, TenantRateState, TenantRateRegistry - All command handlers: user_id: &str → ctx: &TenantCtx - ChatDelegate: cost check/record/settings via self.tenant - System components: store field changed to AdminScope - Config: TENANT_MAX_LLM_CONCURRENT, TENANT_MAX_JOBS_CONCURRENT env vars - Fixes bug: /status <job_id> cross-tenant leak (now auto-filtered) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address PR #1626 review feedback — bounded LRU cache, admin auth, FK cleanup - Replace HashMap with lru::LruCache in DbAuthenticator so the token cache is hard-bounded at 1024 entries (evicts LRU, not just expired) - Gate admin user endpoints (list/detail/update/suspend/activate) with AdminUser extractor so members get 403 instead of full access - Add api_tokens to libSQL delete_user cleanup list to prevent orphaned tokens (libSQL has no FK cascade) - Add regression tests for all three fixes Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: update CA certificates in runtime Docker image Ensures the root certificate bundle is current so TLS handshakes to services like Supabase succeed on Railway. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: resolve CI failures — formatting, no-panics check - Run cargo fmt on test code - Replace .expect() with const NonZeroUsize in DbAuthenticator - Add // safety: comments for test-only code in multi_tenant.rs Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: switch PostgreSQL TLS from rustls to native-tls rustls with rustls-native-certs fails TLS handshake on Railway's slim container (empty or stale root cert store). native-tls delegates to OpenSSL on Linux which handles system certs more reliably. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * Adding user management api * feat: admin secrets provisioning API + API documentation - Add PUT/GET/DELETE /api/admin/users/{id}/secrets/{name} endpoints for application backends to provision per-user secrets (AES-256-GCM encrypted) - Add secrets_store field to GatewayState with builder wiring - Create docs/USER_MANAGEMENT_API.md with full API spec covering users, secrets, tokens, profile, and usage endpoints - Update web gateway CLAUDE.md route table Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: add CatchPanicLayer to capture handler panics Without this, panics in async handlers silently drop the connection and the edge proxy returns a generic 503. Now panics are caught, logged, and returned as 500 with the panic message. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address second-round review — transactional delete, overflow, error logging - C1: Wrap PostgreSQL delete_user() in a transaction so partial cleanup can't leave users in a half-deleted state - M2: Add job_events to delete cleanup (both backends) — FK to agent_jobs without CASCADE would cause FK violation - H1/M4: Cap expires_in_days to 36500 before i64 cast (tokens + secrets) - H2: Validate target user exists before creating admin token to prevent orphan tokens on libSQL - H3: Log DB errors in DbAuthenticator::authenticate() instead of silently swallowing them as 401 Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: revert to rustls with webpki-roots fallback for PostgreSQL TLS native-tls/OpenSSL caused silent crashes (segfaults in C code) during DB writes on Railway containers. Switch back to rustls but add webpki-roots as a fallback when system certs are missing, which was the original TLS handshake failure on slim container images. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * chore: update Cargo.lock for rustls + webpki-roots Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * debug: add /api/debug/db-write endpoint to diagnose user insert failure Temporary diagnostic endpoint that tests DB INSERT to users table with full error logging. No auth required. Will be removed after debugging. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * perf: use cargo-chef in Dockerfile for dependency caching Splits the build into planner/deps/builder stages. Dependencies are only recompiled when Cargo.toml or Cargo.lock change. Source-only changes skip straight to the final build stage. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * debug: add tracing to users_create_handler Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: guard created_by FK in user creation handler The auth identity user_id (from owner_id scope) may not match any user row in the DB, causing a FK violation on the created_by column. Check that the referenced user exists before setting created_by. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * refactor: collapse GATEWAY_USER_ID into IRONCLAW_OWNER_ID Remove the separate GATEWAY_USER_ID config. The gateway now uses IRONCLAW_OWNER_ID (config.owner_id) directly for auth identity, bootstrap user creation, and workspace scoping. Previously, with_owner_scope() rebinds the auth identity to owner_id while keeping default_sender_id as the gateway user_id. This caused a FK constraint violation when creating users because the auth identity ("default") didn't match any user in the DB ("nearai"). Changes: - Remove GATEWAY_USER_ID env var and gateway_user_id from settings - Remove user_id field from GatewayConfig - Add owner_id parameter to GatewayChannel::new() - Remove with_owner_scope() method - Remove default_sender_id from GatewayState - Remove sender override logic in chat/approval handlers - Remove debug endpoint and tracing from prior debugging - Update all tests and E2E fixtures Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: hide Users tab for non-admins, remove auth hint text - Fetch /api/profile after login and hide the Users settings tab when the user's role is not admin - Remove the "Enter the GATEWAY_AUTH_TOKEN" hint from the login page since tokens are now managed via the admin panel, not .env files Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address review feedback (auth 503, token expiry, CORS PATCH) - DB auth errors now return 503 instead of 401 so outages are distinguishable from invalid tokens (serrrfirat H3) - Cap expires_in_days to 36500 before i64 cast to prevent negative duration from u64 overflow (serrrfirat H1) - Add PATCH to CORS allowed methods for profile/user update endpoints (Copilot) - Stop leaking panic details in CatchPanicLayer response body Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: harden multi-tenant isolation — review fixes from #1614 - Add conversation ownership checks in TenantScope: add_conversation_message, touch_conversation, list_conversation_messages (+ paginated), update_conversation_metadata_field, get_conversation_metadata now return NotFound for conversations not owned by the tenant (cross-tenant data leak) - Fix multi-user heartbeat: clear notify_user_id per runner so notifications persist to the correct user, not the shared config target - Move hygiene tasks into bounded JoinSet instead of unbounded tokio::spawn - Revert send_notification to private visibility (only used within module) - Use effective_model_name() for cost attribution in dispatcher so providers that ignore per-request model overrides report the actual model used - Fix inject_model_override doc comment; add 3 unit tests - Fix heartbeat doc comment ("routines" not "active routines") Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: add Jobs, Cost, Last Active columns to admin Users table Add UserSummaryStats struct and user_summary_stats() batch query to the UserStore trait (both PostgreSQL and libSQL backends). The admin users list endpoint now fetches per-user aggregates (job count, total LLM spend, most recent activity) in a single query and includes them inline in the response. The frontend Users table displays three new columns. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address review comments and CI formatting failures CI fixes: - cargo fmt fixes in cli/mod.rs and db/tls.rs Security/correctness (from Copilot + serrrfirat + pranavraja99 reviews): - Token create: reject expires_in_days > 36500 with 400 instead of silent clamp - Token create: return 404 when admin targets non-existent user - User create: map duplicate email constraint violations to 409 Conflict - User create: remove unnecessary DB roundtrip for created_by (use AdminUser directly) - DB auth: log warn on DB lookup failures instead of silently swallowing errors - libSQL: add FK constraints on users.created_by and api_tokens.user_id Config fixes: - agent.multi_tenant: resolve from AGENT_MULTI_TENANT env var instead of hardcoding false - heartbeat.multi_tenant: fix doc comment to match actual env-var-based behavior UI fix: - showTokenBanner: pass correct title ("Token created!" vs "User created!") Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address remaining review comments (round 2) - Secrets handlers: normalize name to lowercase before store operations, validate target user_id exists (returns 404 if not found) - libSQL: propagate cost parsing errors instead of unwrap_or_default() in both user_usage_stats and user_summary_stats - users_list_handler: propagate user_summary_stats DB errors (was silently swallowed with unwrap_or_default) - loadUsers: distinguish 401/403 (admin required) from other errors - Docs: fix users.id type (TEXT not UUID), remove "invitation flow" from V14 migration comment Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: i18n for Users tab, atomic user+token creation, transactional delete_user i18n: - Add 31 translation keys for all Users tab strings (en + zh-CN) - Wire data-i18n attributes on HTML elements (headings, buttons, inputs, table headers, empty state) - Replace all hard-coded strings in app.js with I18n.t() calls Atomic user+token creation: - Add create_user_with_token() to UserStore trait - PostgreSQL: wraps both INSERTs in conn.transaction() with auto-rollback - libSQL: wraps in explicit BEGIN/COMMIT with ROLLBACK on error - Handler uses single atomic call instead of two separate operations Transactional delete_user for libSQL: - Wrap multi-table DELETE cascade in BEGIN/COMMIT transaction - ROLLBACK on any error to prevent partial cleanup / inconsistent state - Matches the PostgreSQL implementation which already used transactions Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: revert V14 migration to match deployed checksum [skip-regression-check] Refinery checksums applied migrations — editing V14__users.sql after it was already applied causes deployment failures. Revert the cosmetic comment changes (added in df40b22f) to restore the original checksum. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: bootstrap onboarding flow for multi-tenant users The bootstrap greeting and workspace seeding only ran for the owner workspace at startup, so new users created via the admin API never received the welcome message or identity files (BOOTSTRAP.md, SOUL.md, AGENTS.md, USER.md, etc.). Three fixes: - tenant_ctx(): seed per-user workspace on first creation via seed_if_empty(), which writes identity files and sets bootstrap_pending when the workspace is truly fresh - handle_message(): check take_bootstrap_pending() on the tenant workspace (not the owner workspace) and persist the greeting to the user's own assistant conversation + broadcast via SSE - WorkspacePool: seed new per-user workspaces in the web gateway so memory tools also see identity files immediately The existing single-user bootstrap in Agent::run() is preserved for non-multi-tenant deployments. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address remaining PR review comments (round 3) - Docs: fix metadata description from "merge patch" to "full replacement" - Secrets: reject expires_in_days > 36500 with 400 (was silently clamped) - libSQL: CAST(SUM(cost) AS TEXT) in user_usage_stats and user_summary_stats to prevent SQLite numeric coercion from crashing get_text() — this was the root cause of the Copilot "SUM returns numeric type" comments - Add 3 regression tests: user_summary_stats (empty + with data) and user_usage_stats (multi-model aggregation) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: add role change support for users (admin/member toggle) - Add update_user_role() to UserStore trait + both backends (PostgreSQL and libSQL) - Extend PATCH /api/admin/users/{id} to accept optional "role" field with validation (must be "admin" or "member") - Add "Make Admin" / "Make Member" toggle button in Users table actions - Add i18n keys for role change (en + zh-CN) - Update API docs to document the role field on PATCH - Fix test helpers to use fmt_ts() for timestamps (was using SQLite datetime('now') which produces incompatible format for string comparison) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: show live LLM spend in Users table instead of only DB-recorded costs [skip-regression-check] Chat turns record LLM cost in CostGuard (in-memory) but don't create agent_jobs/llm_calls DB rows — those are only written for background jobs. The Users table was querying only from DB, so it showed $0.00 for users who only chatted. Now supplements DB stats with CostGuard.daily_spend_for_user() — the same source displayed in the status bar token counter. Shows whichever is larger (DB historical total vs live daily spend). Also falls back to last_login_at for "Last Active" when no DB job activity exists. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: persist chat LLM calls to DB and fix usage stats query Two root causes for zero usage stats: 1. ChatDelegate only recorded LLM costs to CostGuard (in-memory) — never to the llm_calls DB table. Added DB persistence via TenantScope.record_llm_call() after each chat LLM call, with job_id=NULL and conversation_id=thread_id. 2. user_summary_stats query only joined agent_jobs→llm_calls, missing chat calls (which have job_id=NULL). Redesigned query to start from llm_calls and resolve user_id via COALESCE(agent_jobs.user_id, conversations.user_id) — covers both job and chat LLM calls. Both PostgreSQL and libSQL queries updated. TenantScope gets record_llm_call() method. Tests updated for new query semantics. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address review comments — input validation, cost semantics, panic safety [skip-regression-check] - Validate display_name: trim whitespace, reject empty strings (create + update) - Validate metadata: must be a JSON object, return 400 if not (admin + profile) - secrets_list_handler: verify target user_id exists before listing - Cost display: use DB total directly (chat calls now persist to DB), remove confusing max(db,live) CostGuard fallback - CatchPanicLayer: truncate panic payload to 200 chars in log to limit potential sensitive data exposure Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address Copilot round 5 — docs, secrets consistency, token name, provider field [skip-regression-check] - Docs: users.id note updated to "typically UUID v4 strings (bootstrap admin may use a custom ID)" - secrets_list_handler: return 503 when DB store is None (was falling through to list secrets without user validation) - tokens_create: trim + reject empty token name (matching display_name pattern) - LlmCallRecord.provider: use llm_backend ("nearai","openai") instead of model_name() which returns the model identifier - user_summary_stats zero-LLM users: acceptable — handler already falls back to 0 cost and last_login_at for missing entries Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: DB auth returns 503 on outage, scheduler counts only blocking jobs From serrrfirat review: - DB auth: return Err(()) on database errors so middleware returns 503 instead of silently returning Ok(None) → 401 (auth miss) - Scheduler: add parallel_blocking_count_for() that uses is_parallel_blocking() (Pending/InProgress/Stuck) instead of is_active() for per-user concurrency — Completed/Submitted jobs no longer count against MAX_JOBS_PER_USER From Copilot: - CLAUDE.md: fix secrets route paths from {id} to {user_id} - token_hash: use .as_slice() instead of .to_vec() to avoid heap allocation on every token auth/creation call Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: immediate auth cache invalidation on security-critical actions (zmanian review #6) Add DbAuthenticator::invalidate_user() that evicts all cached entries for a user. Called after: - Suspend user (immediate lockout, was 60s delay) - Activate user (immediate access restoration) - Role change (admin↔member takes effect immediately) - Token revocation (revoked token can't be reused from cache) The DbAuthenticator is shared (via Clone, which Arc-clones the cache) between the auth middleware and GatewayState, so handlers can evict entries from the same cache the middleware reads. Also from zmanian's review: - Items 1-5, 7-11 were already resolved in prior commits - Item 12 (String→enum for status/role) is deferred as a broader refactor Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: last-admin protection, usage stats for chat calls, UTF-8 safe panic truncation Last-admin protection: - Suspend, delete, and role-demotion of the last active admin now return 409 Conflict instead of succeeding and locking out the admin API - Helper is_last_admin() checks active admin count before destructive ops Usage stats: - user_usage_stats() now includes chat LLM calls (job_id=NULL) by joining via conversations.user_id, matching user_summary_stats() - Both PostgreSQL and libSQL queries updated Panic handler: - Use floor_char_boundary(200) instead of byte-index [..200] to prevent panic on multi-byte UTF-8 characters in panic messages Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: workspace seed race, bootstrap atomicity, email trim, secrets upsert response [skip-regression-check] - WorkspacePool: await seed_if_empty() synchronously after inserting into cache (drop lock first to avoid blocking), so callers see identity files immediately instead of racing a background task - Bootstrap admin: use create_user_with_token() for atomic user+token creation, matching the admin create endpoint - Email: trim whitespace, treat empty as None to prevent " " being stored and breaking uniqueness - Secrets PUT: report "updated" vs "created" based on prior existence - Last token_hash.to_vec() → .as_slice() in authenticate_token Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: disable unscoped webhook endpoint in multi-tenant mode [skip-regression-check] The original /api/webhooks/{path} endpoint looks up routines across all users. In multi-tenant mode, anyone who knows the webhook path + secret could trigger another user's routine. Now returns 410 Gone with a message pointing to the scoped endpoint /api/webhooks/u/{user_id}/{path}. Detection uses state.db_auth.is_some() — present only when DB-backed auth is enabled (multi-tenant). Single-user deployments are unaffected. From: standardtoaster review comment Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: webhook multi-tenant check, secrets error propagation, stale doc comment [skip-regression-check] - Webhook: use workspace_pool.is_some() instead of db_auth.is_some() for multi-tenant detection — db_auth is set for any DB deployment, workspace_pool is only set when has_any_users() was true at startup - Secrets: propagate exists() errors instead of unwrap_or(false) so backend outages surface as 500 rather than incorrect "created" status - Config: fix stale workspace_read_scopes comment referencing user_id Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> --------- Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
1705 lines
70 KiB
Rust
1705 lines
70 KiB
Rust
//! Main agent loop.
|
|
//!
|
|
//! Contains the `Agent` struct, `AgentDeps`, and the core event loop (`run`).
|
|
//! The heavy lifting is delegated to sibling modules:
|
|
//!
|
|
//! - `dispatcher` - Tool dispatch (agentic loop, tool execution)
|
|
//! - `commands` - System commands and job handlers
|
|
//! - `thread_ops` - Thread/session operations (user input, undo, approval, persistence)
|
|
|
|
use std::sync::Arc;
|
|
|
|
use futures::StreamExt;
|
|
use uuid::Uuid;
|
|
|
|
use crate::agent::context_monitor::ContextMonitor;
|
|
use crate::agent::heartbeat::{spawn_heartbeat, spawn_multi_user_heartbeat};
|
|
use crate::agent::routine_engine::{RoutineEngine, spawn_cron_ticker};
|
|
use crate::agent::self_repair::{DefaultSelfRepair, RepairResult, SelfRepair};
|
|
use crate::agent::session::ThreadState;
|
|
use crate::agent::session_manager::SessionManager;
|
|
use crate::agent::submission::{Submission, SubmissionParser, SubmissionResult};
|
|
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;
|
|
use crate::db::Database;
|
|
use crate::error::{ChannelError, Error};
|
|
use crate::extensions::ExtensionManager;
|
|
use crate::hooks::HookRegistry;
|
|
use crate::llm::LlmProvider;
|
|
use crate::safety::SafetyLayer;
|
|
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
|
|
.chars()
|
|
.take(max_chars + 50)
|
|
.map(|c| if c == '\n' { ' ' } else { c })
|
|
.collect::<String>()
|
|
.split_whitespace()
|
|
.collect::<Vec<_>>()
|
|
.join(" ");
|
|
// char_indices gives us byte offsets at char boundaries, so the slice is always valid UTF-8.
|
|
if collapsed.chars().count() > max_chars {
|
|
let byte_offset = collapsed
|
|
.char_indices()
|
|
.nth(max_chars)
|
|
.map(|(i, _)| i)
|
|
.unwrap_or(collapsed.len());
|
|
format!("{}...", &collapsed[..byte_offset])
|
|
} else {
|
|
collapsed
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
fn resolve_routine_notification_user(metadata: &serde_json::Value) -> Option<String> {
|
|
resolve_owner_scope_notification_user(
|
|
metadata.get("notify_user").and_then(|value| value.as_str()),
|
|
metadata.get("owner_id").and_then(|value| value.as_str()),
|
|
)
|
|
}
|
|
|
|
fn trimmed_option(value: Option<&str>) -> Option<String> {
|
|
value
|
|
.map(str::trim)
|
|
.filter(|value| !value.is_empty())
|
|
.map(ToOwned::to_owned)
|
|
}
|
|
|
|
fn resolve_owner_scope_notification_user(
|
|
explicit_user: Option<&str>,
|
|
owner_fallback: Option<&str>,
|
|
) -> Option<String> {
|
|
trimmed_option(explicit_user).or_else(|| trimmed_option(owner_fallback))
|
|
}
|
|
|
|
fn is_single_message_repl(message: &IncomingMessage) -> bool {
|
|
message.channel == "repl"
|
|
&& message
|
|
.metadata
|
|
.get("single_message_mode")
|
|
.and_then(|value| value.as_bool())
|
|
.unwrap_or(false)
|
|
}
|
|
|
|
async fn resolve_channel_notification_user(
|
|
extension_manager: Option<&Arc<ExtensionManager>>,
|
|
channel: Option<&str>,
|
|
explicit_user: Option<&str>,
|
|
owner_fallback: Option<&str>,
|
|
) -> Option<String> {
|
|
if let Some(user) = trimmed_option(explicit_user) {
|
|
return Some(user);
|
|
}
|
|
|
|
if let Some(channel_name) = trimmed_option(channel)
|
|
&& let Some(extension_manager) = extension_manager
|
|
&& let Some(target) = extension_manager
|
|
.notification_target_for_channel(&channel_name)
|
|
.await
|
|
{
|
|
return Some(target);
|
|
}
|
|
|
|
resolve_owner_scope_notification_user(explicit_user, owner_fallback)
|
|
}
|
|
|
|
async fn resolve_routine_notification_target(
|
|
extension_manager: Option<&Arc<ExtensionManager>>,
|
|
metadata: &serde_json::Value,
|
|
) -> Option<String> {
|
|
resolve_channel_notification_user(
|
|
extension_manager,
|
|
metadata
|
|
.get("notify_channel")
|
|
.and_then(|value| value.as_str()),
|
|
metadata.get("notify_user").and_then(|value| value.as_str()),
|
|
metadata.get("owner_id").and_then(|value| value.as_str()),
|
|
)
|
|
.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 { .. })
|
|
}
|
|
|
|
/// Core dependencies for the agent.
|
|
///
|
|
/// Bundles the shared components to reduce argument count.
|
|
pub struct AgentDeps {
|
|
/// Resolved durable owner scope for the instance.
|
|
pub owner_id: String,
|
|
pub store: Option<Arc<dyn Database>>,
|
|
pub llm: Arc<dyn LlmProvider>,
|
|
/// Cheap/fast LLM for lightweight tasks (heartbeat, routing, evaluation).
|
|
/// Falls back to the main `llm` if None.
|
|
pub cheap_llm: Option<Arc<dyn LlmProvider>>,
|
|
pub safety: Arc<SafetyLayer>,
|
|
pub tools: Arc<ToolRegistry>,
|
|
pub workspace: Option<Arc<Workspace>>,
|
|
pub extension_manager: Option<Arc<ExtensionManager>>,
|
|
pub skill_registry: Option<Arc<std::sync::RwLock<SkillRegistry>>>,
|
|
pub skill_catalog: Option<Arc<crate::skills::catalog::SkillCatalog>>,
|
|
pub skills_config: SkillsConfig,
|
|
pub hooks: Arc<HookRegistry>,
|
|
/// Cost enforcement guardrails (daily budget, hourly rate limits).
|
|
pub cost_guard: Arc<crate::agent::cost_guard::CostGuard>,
|
|
/// SSE manager for live job event streaming to the web gateway.
|
|
pub sse_tx: Option<Arc<crate::channels::web::sse::SseManager>>,
|
|
/// HTTP interceptor for trace recording/replay.
|
|
pub http_interceptor: Option<Arc<dyn crate::llm::recording::HttpInterceptor>>,
|
|
/// Audio transcription middleware for voice messages.
|
|
pub transcription: Option<Arc<crate::llm::transcription::TranscriptionMiddleware>>,
|
|
/// Document text extraction middleware for PDF, DOCX, PPTX, etc.
|
|
pub document_extraction: Option<Arc<crate::document_extraction::DocumentExtractionMiddleware>>,
|
|
/// 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<Arc<dyn crate::tools::SoftwareBuilder>>,
|
|
/// Resolved LLM backend identifier (e.g., "nearai", "openai", "groq").
|
|
/// Used by `/model` persistence to determine which env var to update.
|
|
pub llm_backend: String,
|
|
/// Per-tenant rate limiting registry (lazily creates rate state per user).
|
|
pub tenant_rates: Arc<crate::tenant::TenantRateRegistry>,
|
|
}
|
|
|
|
/// The main agent that coordinates all components.
|
|
pub struct Agent {
|
|
pub(super) config: AgentConfig,
|
|
pub(super) deps: AgentDeps,
|
|
pub(super) channels: Arc<ChannelManager>,
|
|
pub(super) context_manager: Arc<ContextManager>,
|
|
pub(super) scheduler: Arc<Scheduler>,
|
|
pub(super) router: Router,
|
|
pub(super) session_manager: Arc<SessionManager>,
|
|
pub(super) context_monitor: ContextMonitor,
|
|
pub(super) heartbeat_config: Option<HeartbeatConfig>,
|
|
pub(super) hygiene_config: Option<crate::config::HygieneConfig>,
|
|
pub(super) routine_config: Option<RoutineConfig>,
|
|
/// Shared routine-engine slot used for internal event matching and for exposing
|
|
/// the engine to gateway/manual trigger entry points.
|
|
pub(super) routine_engine_slot:
|
|
Arc<tokio::sync::RwLock<Option<Arc<crate::agent::routine_engine::RoutineEngine>>>>,
|
|
}
|
|
|
|
impl Agent {
|
|
pub(super) fn owner_id(&self) -> &str {
|
|
if let Some(workspace) = self.deps.workspace.as_ref() {
|
|
debug_assert_eq!(
|
|
workspace.user_id(),
|
|
self.deps.owner_id,
|
|
"workspace.user_id() must stay aligned with deps.owner_id"
|
|
);
|
|
}
|
|
|
|
&self.deps.owner_id
|
|
}
|
|
|
|
/// Create a new agent.
|
|
///
|
|
/// Optionally accepts pre-created `ContextManager` and `SessionManager` for sharing
|
|
/// with external components (job tools, web gateway). Creates new ones if not provided.
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub fn new(
|
|
config: AgentConfig,
|
|
deps: AgentDeps,
|
|
channels: Arc<ChannelManager>,
|
|
heartbeat_config: Option<HeartbeatConfig>,
|
|
hygiene_config: Option<crate::config::HygieneConfig>,
|
|
routine_config: Option<RoutineConfig>,
|
|
context_manager: Option<Arc<ContextManager>>,
|
|
session_manager: Option<Arc<SessionManager>>,
|
|
) -> Self {
|
|
let context_manager = context_manager
|
|
.unwrap_or_else(|| Arc::new(ContextManager::new(config.max_parallel_jobs)));
|
|
|
|
let session_manager = session_manager.unwrap_or_else(|| Arc::new(SessionManager::new()));
|
|
|
|
let mut scheduler = Scheduler::new(
|
|
config.clone(),
|
|
context_manager.clone(),
|
|
deps.llm.clone(),
|
|
deps.safety.clone(),
|
|
SchedulerDeps {
|
|
tools: deps.tools.clone(),
|
|
extension_manager: deps.extension_manager.clone(),
|
|
store: deps
|
|
.store
|
|
.as_ref()
|
|
.map(|db| crate::tenant::AdminScope::new(Arc::clone(db))),
|
|
hooks: deps.hooks.clone(),
|
|
},
|
|
);
|
|
if let Some(ref sse) = deps.sse_tx {
|
|
scheduler.set_sse_sender(Arc::clone(sse));
|
|
}
|
|
if let Some(ref interceptor) = deps.http_interceptor {
|
|
scheduler.set_http_interceptor(Arc::clone(interceptor));
|
|
}
|
|
let scheduler = Arc::new(scheduler);
|
|
|
|
Self {
|
|
config,
|
|
deps,
|
|
channels,
|
|
context_manager,
|
|
scheduler,
|
|
router: Router::new(),
|
|
session_manager,
|
|
context_monitor: ContextMonitor::new(),
|
|
heartbeat_config,
|
|
hygiene_config,
|
|
routine_config,
|
|
routine_engine_slot: Arc::new(tokio::sync::RwLock::new(None)),
|
|
}
|
|
}
|
|
|
|
/// Replace the routine-engine slot with a shared one so the gateway and
|
|
/// agent reference the same engine.
|
|
pub fn set_routine_engine_slot(
|
|
&mut self,
|
|
slot: Arc<tokio::sync::RwLock<Option<Arc<crate::agent::routine_engine::RoutineEngine>>>>,
|
|
) {
|
|
self.routine_engine_slot = slot;
|
|
}
|
|
|
|
async fn routine_engine(&self) -> Option<Arc<crate::agent::routine_engine::RoutineEngine>> {
|
|
self.routine_engine_slot.read().await.clone()
|
|
}
|
|
|
|
// Convenience accessors
|
|
|
|
/// Get the scheduler (for external wiring, e.g. CreateJobTool).
|
|
pub fn scheduler(&self) -> Arc<Scheduler> {
|
|
Arc::clone(&self.scheduler)
|
|
}
|
|
|
|
pub(super) fn store(&self) -> Option<&Arc<dyn Database>> {
|
|
self.deps.store.as_ref()
|
|
}
|
|
|
|
pub(super) fn llm(&self) -> &Arc<dyn LlmProvider> {
|
|
&self.deps.llm
|
|
}
|
|
|
|
/// Get the cheap/fast LLM provider, falling back to the main one.
|
|
pub(super) fn cheap_llm(&self) -> &Arc<dyn LlmProvider> {
|
|
self.deps.cheap_llm.as_ref().unwrap_or(&self.deps.llm)
|
|
}
|
|
|
|
pub(super) fn safety(&self) -> &Arc<SafetyLayer> {
|
|
&self.deps.safety
|
|
}
|
|
|
|
pub(super) fn tools(&self) -> &Arc<ToolRegistry> {
|
|
&self.deps.tools
|
|
}
|
|
|
|
pub(super) fn workspace(&self) -> Option<&Arc<Workspace>> {
|
|
self.deps.workspace.as_ref()
|
|
}
|
|
|
|
pub(super) fn hooks(&self) -> &Arc<HookRegistry> {
|
|
&self.deps.hooks
|
|
}
|
|
|
|
pub(super) fn cost_guard(&self) -> &Arc<crate::agent::cost_guard::CostGuard> {
|
|
&self.deps.cost_guard
|
|
}
|
|
|
|
/// Build a tenant-scoped execution context for the given user.
|
|
///
|
|
/// This is the standard entry point for per-user operations. The returned
|
|
/// [`TenantCtx`] provides a [`TenantScope`] that auto-binds `user_id` on
|
|
/// every database operation and a per-user rate limiter.
|
|
pub(super) async fn tenant_ctx(&self, user_id: &str) -> crate::tenant::TenantCtx {
|
|
let rate = self.deps.tenant_rates.get_or_create(user_id).await;
|
|
|
|
let store = self
|
|
.deps
|
|
.store
|
|
.as_ref()
|
|
.map(|db| crate::tenant::TenantScope::new(user_id, Arc::clone(db)));
|
|
|
|
// Reuse the owner workspace if user matches, otherwise create per-user.
|
|
// Per-user workspaces are seeded on first creation so they get identity
|
|
// files and BOOTSTRAP.md (which triggers the onboarding greeting).
|
|
let workspace = match &self.deps.workspace {
|
|
Some(ws) if ws.user_id() == user_id => Some(Arc::clone(ws)),
|
|
_ => {
|
|
if let Some(db) = self.deps.store.as_ref() {
|
|
let ws = Arc::new(Workspace::new_with_db(user_id, Arc::clone(db)));
|
|
if let Err(e) = ws.seed_if_empty().await {
|
|
tracing::warn!(
|
|
user_id = user_id,
|
|
"Failed to seed per-user workspace: {}",
|
|
e
|
|
);
|
|
}
|
|
Some(ws)
|
|
} else {
|
|
None
|
|
}
|
|
}
|
|
};
|
|
|
|
crate::tenant::TenantCtx::new(
|
|
user_id,
|
|
store,
|
|
workspace,
|
|
Arc::clone(&self.deps.cost_guard),
|
|
rate,
|
|
)
|
|
}
|
|
|
|
/// Get an admin-scoped database accessor for cross-tenant operations.
|
|
///
|
|
/// Only for system-level components (heartbeat, routine engine, self-repair,
|
|
/// scheduler). Handler code should use [`tenant_ctx()`](Self::tenant_ctx) instead.
|
|
pub(super) fn admin_store(&self) -> Option<crate::tenant::AdminScope> {
|
|
self.deps
|
|
.store
|
|
.as_ref()
|
|
.map(|db| crate::tenant::AdminScope::new(Arc::clone(db)))
|
|
}
|
|
|
|
pub(super) fn skill_registry(&self) -> Option<&Arc<std::sync::RwLock<SkillRegistry>>> {
|
|
self.deps.skill_registry.as_ref()
|
|
}
|
|
|
|
pub(super) fn skill_catalog(&self) -> Option<&Arc<crate::skills::catalog::SkillCatalog>> {
|
|
self.deps.skill_catalog.as_ref()
|
|
}
|
|
|
|
/// Select active skills for a message using deterministic prefiltering.
|
|
pub(super) fn select_active_skills(
|
|
&self,
|
|
message_content: &str,
|
|
) -> Vec<crate::skills::LoadedSkill> {
|
|
if let Some(registry) = self.skill_registry() {
|
|
let guard = match registry.read() {
|
|
Ok(g) => g,
|
|
Err(e) => {
|
|
tracing::error!("Skill registry lock poisoned: {}", e);
|
|
return vec![];
|
|
}
|
|
};
|
|
let available = guard.skills();
|
|
let skills_cfg = &self.deps.skills_config;
|
|
let selected = crate::skills::prefilter_skills(
|
|
message_content,
|
|
available,
|
|
skills_cfg.max_active_skills,
|
|
skills_cfg.max_context_tokens,
|
|
);
|
|
|
|
if !selected.is_empty() {
|
|
tracing::debug!(
|
|
"Selected {} skill(s) for message: {}",
|
|
selected.len(),
|
|
selected
|
|
.iter()
|
|
.map(|s| s.name())
|
|
.collect::<Vec<_>>()
|
|
.join(", ")
|
|
);
|
|
}
|
|
|
|
selected.into_iter().cloned().collect()
|
|
} else {
|
|
vec![]
|
|
}
|
|
}
|
|
|
|
/// 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?;
|
|
|
|
// Start self-repair task with notification forwarding
|
|
let mut self_repair = DefaultSelfRepair::new(
|
|
self.context_manager.clone(),
|
|
self.config.stuck_threshold,
|
|
self.config.max_repair_attempts,
|
|
);
|
|
if let Some(admin) = self.admin_store() {
|
|
self_repair = self_repair.with_store(admin);
|
|
}
|
|
if let Some(ref builder) = self.deps.builder {
|
|
self_repair = self_repair.with_builder(Arc::clone(builder), Arc::clone(self.tools()));
|
|
}
|
|
let repair = Arc::new(self_repair);
|
|
let repair_interval = self.config.repair_check_interval;
|
|
let repair_channels = self.channels.clone();
|
|
let repair_owner_id = self.owner_id().to_string();
|
|
let repair_handle = tokio::spawn(async move {
|
|
loop {
|
|
tokio::time::sleep(repair_interval).await;
|
|
|
|
// Check stuck jobs
|
|
let stuck_jobs = repair.detect_stuck_jobs().await;
|
|
for job in stuck_jobs {
|
|
tracing::info!("Attempting to repair stuck job {}", job.job_id);
|
|
let result = repair.repair_stuck_job(&job).await;
|
|
let notification = match &result {
|
|
Ok(RepairResult::Success { message }) => {
|
|
tracing::info!("Repair succeeded: {}", message);
|
|
Some(format!(
|
|
"Job {} was stuck for {}s, recovery succeeded: {}",
|
|
job.job_id,
|
|
job.stuck_duration.as_secs(),
|
|
message
|
|
))
|
|
}
|
|
Ok(RepairResult::Failed { message }) => {
|
|
tracing::error!("Repair failed: {}", message);
|
|
Some(format!(
|
|
"Job {} was stuck for {}s, recovery failed permanently: {}",
|
|
job.job_id,
|
|
job.stuck_duration.as_secs(),
|
|
message
|
|
))
|
|
}
|
|
Ok(RepairResult::ManualRequired { message }) => {
|
|
tracing::warn!("Manual intervention needed: {}", message);
|
|
Some(format!(
|
|
"Job {} needs manual intervention: {}",
|
|
job.job_id, message
|
|
))
|
|
}
|
|
Ok(RepairResult::Retry { message }) => {
|
|
tracing::warn!("Repair needs retry: {}", message);
|
|
None // Don't spam the user on retries
|
|
}
|
|
Err(e) => {
|
|
tracing::error!("Repair error: {}", e);
|
|
None
|
|
}
|
|
};
|
|
|
|
if let Some(msg) = notification {
|
|
let response = OutgoingResponse::text(format!("Self-Repair: {}", msg));
|
|
let _ = repair_channels
|
|
.broadcast_all(&repair_owner_id, response)
|
|
.await;
|
|
}
|
|
}
|
|
|
|
// Check broken tools
|
|
let broken_tools = repair.detect_broken_tools().await;
|
|
for tool in broken_tools {
|
|
tracing::info!("Attempting to repair broken tool: {}", tool.name);
|
|
match repair.repair_broken_tool(&tool).await {
|
|
Ok(RepairResult::Success { message }) => {
|
|
let response = OutgoingResponse::text(format!(
|
|
"Self-Repair: Tool '{}' repaired: {}",
|
|
tool.name, message
|
|
));
|
|
let _ = repair_channels
|
|
.broadcast_all(&repair_owner_id, response)
|
|
.await;
|
|
}
|
|
Ok(result) => {
|
|
tracing::info!("Tool repair result: {:?}", result);
|
|
}
|
|
Err(e) => {
|
|
tracing::error!("Tool repair error: {}", e);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
// Spawn session pruning task
|
|
let session_mgr = self.session_manager.clone();
|
|
let session_idle_timeout = self.config.session_idle_timeout;
|
|
let pruning_handle = tokio::spawn(async move {
|
|
let mut interval = tokio::time::interval(std::time::Duration::from_secs(600)); // Every 10 min
|
|
interval.tick().await; // Skip immediate first tick
|
|
loop {
|
|
interval.tick().await;
|
|
session_mgr.prune_stale_sessions(session_idle_timeout).await;
|
|
}
|
|
});
|
|
|
|
// Spawn heartbeat if enabled
|
|
let heartbeat_handle = if let Some(ref hb_config) = self.heartbeat_config {
|
|
if hb_config.enabled {
|
|
if let Some(workspace) = self.workspace() {
|
|
let mut config = AgentHeartbeatConfig::default()
|
|
.with_interval(std::time::Duration::from_secs(hb_config.interval_secs));
|
|
config.quiet_hours_start = hb_config.quiet_hours_start;
|
|
config.quiet_hours_end = hb_config.quiet_hours_end;
|
|
config.multi_tenant = hb_config.multi_tenant;
|
|
config.timezone = hb_config
|
|
.timezone
|
|
.clone()
|
|
.or_else(|| Some(self.config.default_timezone.clone()));
|
|
let heartbeat_notify_user = resolve_owner_scope_notification_user(
|
|
hb_config.notify_user.as_deref(),
|
|
Some(self.owner_id()),
|
|
);
|
|
if let Some(channel) = &hb_config.notify_channel
|
|
&& let Some(user) = heartbeat_notify_user.as_deref()
|
|
{
|
|
config = config.with_notify(user, channel);
|
|
}
|
|
|
|
// Set up notification channel
|
|
let (notify_tx, mut notify_rx) =
|
|
tokio::sync::mpsc::channel::<OutgoingResponse>(16);
|
|
|
|
// Spawn notification forwarder that routes through channel manager
|
|
let notify_channel = hb_config.notify_channel.clone();
|
|
let notify_target = resolve_channel_notification_user(
|
|
self.deps.extension_manager.as_ref(),
|
|
hb_config.notify_channel.as_deref(),
|
|
hb_config.notify_user.as_deref(),
|
|
Some(self.owner_id()),
|
|
)
|
|
.await;
|
|
let notify_user = heartbeat_notify_user;
|
|
let channels = self.channels.clone();
|
|
let is_multi_tenant = hb_config.multi_tenant;
|
|
tokio::spawn(async move {
|
|
while let Some(response) = notify_rx.recv().await {
|
|
// In multi-tenant mode, extract the owning user_id from
|
|
// the response metadata so notifications reach the
|
|
// correct user rather than the agent's owner.
|
|
// This intentionally overrides the configured notify_target
|
|
// because each user's heartbeat should notify that user.
|
|
let effective_user = if is_multi_tenant {
|
|
response
|
|
.metadata
|
|
.get("owner_id")
|
|
.and_then(|v| v.as_str())
|
|
.map(String::from)
|
|
} else {
|
|
None
|
|
};
|
|
|
|
// Try the configured channel first, fall back to
|
|
// broadcasting on all channels.
|
|
let targeted_ok = if let Some(ref channel) = notify_channel {
|
|
let target = effective_user.as_deref().or(notify_target.as_deref());
|
|
if let Some(user) = target {
|
|
channels
|
|
.broadcast(channel, user, response.clone())
|
|
.await
|
|
.is_ok()
|
|
} else {
|
|
false
|
|
}
|
|
} else {
|
|
false
|
|
};
|
|
|
|
if !targeted_ok {
|
|
let fallback = effective_user.as_deref().or(notify_user.as_deref());
|
|
if let Some(user) = fallback {
|
|
let results = channels.broadcast_all(user, response).await;
|
|
for (ch, result) in results {
|
|
if let Err(e) = result {
|
|
tracing::warn!(
|
|
"Failed to broadcast heartbeat to {}: {}",
|
|
ch,
|
|
e
|
|
);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
let hygiene = self
|
|
.hygiene_config
|
|
.as_ref()
|
|
.map(|h| h.to_workspace_config())
|
|
.unwrap_or_default();
|
|
|
|
if config.multi_tenant {
|
|
if let Some(admin) = self.admin_store() {
|
|
Some(spawn_multi_user_heartbeat(
|
|
config,
|
|
hygiene,
|
|
self.cheap_llm().clone(),
|
|
Some(notify_tx),
|
|
admin,
|
|
))
|
|
} else {
|
|
tracing::warn!("Multi-tenant heartbeat requires a database store");
|
|
None
|
|
}
|
|
} else {
|
|
Some(spawn_heartbeat(
|
|
config,
|
|
hygiene,
|
|
workspace.clone(),
|
|
self.cheap_llm().clone(),
|
|
Some(notify_tx),
|
|
self.admin_store(),
|
|
))
|
|
}
|
|
} else {
|
|
tracing::warn!("Heartbeat enabled but no workspace available");
|
|
None
|
|
}
|
|
} else {
|
|
None
|
|
}
|
|
} else {
|
|
None
|
|
};
|
|
|
|
// Spawn routine engine if enabled
|
|
let routine_handle = if let Some(ref rt_config) = self.routine_config {
|
|
if rt_config.enabled {
|
|
if let (Some(store), Some(workspace)) = (self.store(), self.workspace()) {
|
|
// Set up notification channel (same pattern as heartbeat)
|
|
let (notify_tx, mut notify_rx) =
|
|
tokio::sync::mpsc::channel::<OutgoingResponse>(32);
|
|
|
|
let engine = Arc::new(RoutineEngine::new(
|
|
rt_config.clone(),
|
|
crate::tenant::AdminScope::new(Arc::clone(store)),
|
|
self.llm().clone(),
|
|
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
|
|
self.deps
|
|
.tools
|
|
.register_routine_tools(Arc::clone(store), Arc::clone(&engine));
|
|
|
|
// Load initial event cache
|
|
engine.refresh_event_cache().await;
|
|
|
|
// Spawn notification forwarder (mirrors heartbeat pattern)
|
|
let channels = self.channels.clone();
|
|
let extension_manager = self.deps.extension_manager.clone();
|
|
tokio::spawn(async move {
|
|
while let Some(response) = notify_rx.recv().await {
|
|
let notify_channel = response
|
|
.metadata
|
|
.get("notify_channel")
|
|
.and_then(|v| v.as_str())
|
|
.map(|s| s.to_string());
|
|
let fallback_user = resolve_owner_scope_notification_user(
|
|
response
|
|
.metadata
|
|
.get("notify_user")
|
|
.and_then(|v| v.as_str()),
|
|
response.metadata.get("owner_id").and_then(|v| v.as_str()),
|
|
);
|
|
let Some(user) = resolve_routine_notification_target(
|
|
extension_manager.as_ref(),
|
|
&response.metadata,
|
|
)
|
|
.await
|
|
else {
|
|
tracing::warn!(
|
|
notify_channel = ?notify_channel,
|
|
"Skipping routine notification with no explicit target or owner scope"
|
|
);
|
|
continue;
|
|
};
|
|
|
|
// Try the configured channel first, fall back to
|
|
// broadcasting on all channels.
|
|
let targeted_ok = if let Some(ref channel) = notify_channel {
|
|
match channels.broadcast(channel, &user, response.clone()).await {
|
|
Ok(()) => true,
|
|
Err(e) => {
|
|
let should_fallback =
|
|
should_fallback_routine_notification(&e);
|
|
tracing::warn!(
|
|
channel = %channel,
|
|
user = %user,
|
|
error = %e,
|
|
should_fallback,
|
|
"Failed to send routine notification to configured channel"
|
|
);
|
|
if !should_fallback {
|
|
continue;
|
|
}
|
|
false
|
|
}
|
|
}
|
|
} else {
|
|
false
|
|
};
|
|
|
|
if !targeted_ok && let Some(user) = fallback_user {
|
|
let results = channels.broadcast_all(&user, response).await;
|
|
for (ch, result) in results {
|
|
if let Err(e) = result {
|
|
tracing::warn!(
|
|
"Failed to broadcast routine notification to {}: {}",
|
|
ch,
|
|
e
|
|
);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
// Spawn cron ticker
|
|
let cron_interval =
|
|
std::time::Duration::from_secs(rt_config.cron_check_interval_secs);
|
|
let cron_handle = spawn_cron_ticker(Arc::clone(&engine), cron_interval);
|
|
|
|
// Store engine reference for event trigger checking
|
|
// Safety: we're in run() which takes self, no other reference exists
|
|
let engine_ref = Arc::clone(&engine);
|
|
// SAFETY: self is consumed by run(), we can smuggle the engine in
|
|
// via a local to use in the message loop below.
|
|
|
|
// Expose engine to gateway for manual triggering
|
|
*self.routine_engine_slot.write().await = Some(Arc::clone(&engine));
|
|
|
|
tracing::debug!(
|
|
"Routines enabled: cron ticker every {}s, max {} concurrent",
|
|
rt_config.cron_check_interval_secs,
|
|
rt_config.max_concurrent_routines
|
|
);
|
|
|
|
Some((cron_handle, engine_ref))
|
|
} else {
|
|
tracing::warn!("Routines enabled but store/workspace not available");
|
|
None
|
|
}
|
|
} else {
|
|
None
|
|
}
|
|
} else {
|
|
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);
|
|
|
|
loop {
|
|
let message = tokio::select! {
|
|
biased;
|
|
_ = tokio::signal::ctrl_c() => {
|
|
tracing::debug!("Ctrl+C received, shutting down...");
|
|
break;
|
|
}
|
|
msg = message_stream.next() => {
|
|
match msg {
|
|
Some(m) => m,
|
|
None => {
|
|
tracing::debug!("All channel streams ended, shutting down...");
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
};
|
|
|
|
// Apply transcription middleware to audio attachments
|
|
let mut message = message;
|
|
if let Some(ref transcription) = self.deps.transcription {
|
|
transcription.process(&mut message).await;
|
|
}
|
|
|
|
// Apply document extraction middleware to document attachments
|
|
if let Some(ref doc_extraction) = self.deps.document_extraction {
|
|
doc_extraction.process(&mut message).await;
|
|
}
|
|
|
|
// Store successfully extracted document text in workspace for indexing
|
|
self.store_extracted_documents(&message).await;
|
|
|
|
match self.handle_message(&message).await {
|
|
Ok(Some(response)) if !response.is_empty() => {
|
|
// Hook: BeforeOutbound — allow hooks to modify or suppress outbound
|
|
let event = crate::hooks::HookEvent::Outbound {
|
|
user_id: message.user_id.clone(),
|
|
channel: message.channel.clone(),
|
|
content: response.clone(),
|
|
thread_id: message.thread_id.clone(),
|
|
};
|
|
match self.hooks().run(&event).await {
|
|
Err(err) => {
|
|
tracing::warn!("BeforeOutbound hook blocked response: {}", err);
|
|
}
|
|
Ok(crate::hooks::HookOutcome::Continue {
|
|
modified: Some(new_content),
|
|
}) => {
|
|
if let Err(e) = self
|
|
.channels
|
|
.respond(&message, OutgoingResponse::text(new_content))
|
|
.await
|
|
{
|
|
tracing::error!(
|
|
channel = %message.channel,
|
|
error = %e,
|
|
"Failed to send response to channel"
|
|
);
|
|
}
|
|
}
|
|
_ => {
|
|
if let Err(e) = self
|
|
.channels
|
|
.respond(&message, OutgoingResponse::text(response))
|
|
.await
|
|
{
|
|
tracing::error!(
|
|
channel = %message.channel,
|
|
error = %e,
|
|
"Failed to send response to channel"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
Ok(Some(empty)) => {
|
|
// Empty response, nothing to send (e.g. approval handled via send_status)
|
|
tracing::debug!(
|
|
channel = %message.channel,
|
|
user = %message.user_id,
|
|
empty_len = empty.len(),
|
|
"Suppressed empty response (not sent to channel)"
|
|
);
|
|
}
|
|
Ok(None) => {
|
|
// Shutdown signal received (/quit, /exit, /shutdown)
|
|
tracing::debug!("Shutdown command received, exiting...");
|
|
break;
|
|
}
|
|
Err(e) => {
|
|
tracing::error!("Error handling message: {}", e);
|
|
if let Err(send_err) = self
|
|
.channels
|
|
.respond(&message, OutgoingResponse::text(format!("Error: {}", e)))
|
|
.await
|
|
{
|
|
tracing::error!(
|
|
channel = %message.channel,
|
|
error = %send_err,
|
|
"Failed to send error response to channel"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Cleanup
|
|
tracing::debug!("Agent shutting down...");
|
|
repair_handle.abort();
|
|
pruning_handle.abort();
|
|
if let Some(handle) = heartbeat_handle {
|
|
handle.abort();
|
|
}
|
|
if let Some((cron_handle, _)) = routine_handle {
|
|
cron_handle.abort();
|
|
}
|
|
self.scheduler.stop_all().await;
|
|
self.channels.shutdown_all().await?;
|
|
|
|
Ok(())
|
|
}
|
|
|
|
/// Store extracted document text in workspace memory for future search/recall.
|
|
async fn store_extracted_documents(&self, message: &IncomingMessage) {
|
|
let workspace = match self.workspace() {
|
|
Some(ws) => ws,
|
|
None => return,
|
|
};
|
|
|
|
for attachment in &message.attachments {
|
|
if attachment.kind != crate::channels::AttachmentKind::Document {
|
|
continue;
|
|
}
|
|
let text = match &attachment.extracted_text {
|
|
Some(t) if !t.starts_with('[') => t, // skip error messages like "[Failed to..."
|
|
_ => continue,
|
|
};
|
|
|
|
// Sanitize filename: strip path separators to prevent directory traversal
|
|
let raw_name = attachment.filename.as_deref().unwrap_or("unnamed_document");
|
|
let filename: String = raw_name
|
|
.chars()
|
|
.map(|c| {
|
|
if c == '/' || c == '\\' || c == '\0' {
|
|
'_'
|
|
} else {
|
|
c
|
|
}
|
|
})
|
|
.collect();
|
|
let filename = filename.trim_start_matches('.');
|
|
let filename = if filename.is_empty() {
|
|
"unnamed_document"
|
|
} else {
|
|
filename
|
|
};
|
|
let date = chrono::Utc::now().format("%Y-%m-%d");
|
|
let path = format!("documents/{date}/{filename}");
|
|
|
|
let header = format!(
|
|
"# {filename}\n\n\
|
|
> Uploaded by **{}** via **{}** on {date}\n\
|
|
> MIME: {} | Size: {} bytes\n\n---\n\n",
|
|
message.user_id,
|
|
message.channel,
|
|
attachment.mime_type,
|
|
attachment.size_bytes.unwrap_or(0),
|
|
);
|
|
let content = format!("{header}{text}");
|
|
|
|
match workspace.write(&path, &content).await {
|
|
Ok(_) => {
|
|
tracing::info!(
|
|
path = %path,
|
|
text_len = text.len(),
|
|
"Stored extracted document in workspace memory"
|
|
);
|
|
}
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
path = %path,
|
|
error = %e,
|
|
"Failed to store extracted document in workspace"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
async fn handle_message(&self, message: &IncomingMessage) -> Result<Option<String>, Error> {
|
|
// Log sensitive details at debug level for troubleshooting
|
|
tracing::debug!(
|
|
message_id = %message.id,
|
|
user_id = %message.user_id,
|
|
channel = %message.channel,
|
|
thread_id = ?message.thread_id,
|
|
"Message details"
|
|
);
|
|
|
|
// Internal messages (e.g. job-monitor notifications) are already
|
|
// rendered text and should be forwarded directly to the user without
|
|
// entering the normal user-input pipeline (LLM/tool loop).
|
|
// The `is_internal` field and `into_internal()` setter are pub(crate),
|
|
// so external channels cannot spoof this flag.
|
|
if message.is_internal {
|
|
tracing::debug!(
|
|
message_id = %message.id,
|
|
channel = %message.channel,
|
|
"Forwarding internal message"
|
|
);
|
|
return Ok(Some(message.content.clone()));
|
|
}
|
|
|
|
// Set message tool context for this turn (current channel and target)
|
|
// For Signal, use signal_target from metadata (group:ID or phone number),
|
|
// otherwise fall back to user_id
|
|
let target = message
|
|
.routing_target()
|
|
.unwrap_or_else(|| message.user_id.clone());
|
|
self.tools()
|
|
.set_message_tool_context(Some(message.channel.clone()), Some(target))
|
|
.await;
|
|
|
|
// Parse submission type first
|
|
let mut submission = SubmissionParser::parse(&message.content);
|
|
tracing::trace!(
|
|
"[agent_loop] Parsed submission: {:?}",
|
|
std::any::type_name_of_val(&submission)
|
|
);
|
|
|
|
// Hook: BeforeInbound — allow hooks to modify or reject user input
|
|
if let Submission::UserInput { ref content } = submission {
|
|
let event = crate::hooks::HookEvent::Inbound {
|
|
user_id: message.user_id.clone(),
|
|
channel: message.channel.clone(),
|
|
content: content.clone(),
|
|
thread_id: message.thread_id.clone(),
|
|
};
|
|
match self.hooks().run(&event).await {
|
|
Err(crate::hooks::HookError::Rejected { reason }) => {
|
|
return Ok(Some(format!("[Message rejected: {}]", reason)));
|
|
}
|
|
Err(err) => {
|
|
return Ok(Some(format!("[Message blocked by hook policy: {}]", err)));
|
|
}
|
|
Ok(crate::hooks::HookOutcome::Continue {
|
|
modified: Some(new_content),
|
|
}) => {
|
|
submission = Submission::UserInput {
|
|
content: new_content,
|
|
};
|
|
}
|
|
_ => {} // Continue, fail-open errors already logged in registry
|
|
}
|
|
}
|
|
|
|
// Hydrate thread from DB if it's a historical thread not in memory
|
|
if let Some(external_thread_id) = message.conversation_scope() {
|
|
tracing::trace!(
|
|
message_id = %message.id,
|
|
thread_id = %external_thread_id,
|
|
"Hydrating thread from DB"
|
|
);
|
|
if let Some(rejection) = self.maybe_hydrate_thread(message, external_thread_id).await {
|
|
return Ok(Some(format!("Error: {}", rejection)));
|
|
}
|
|
}
|
|
|
|
// 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_with_parsed_uuid(
|
|
&message.user_id,
|
|
&message.channel,
|
|
message.conversation_scope(),
|
|
approval_thread_uuid,
|
|
)
|
|
.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,
|
|
"Resolved session and thread"
|
|
);
|
|
|
|
// Auth mode interception: if the thread is awaiting a token, route
|
|
// the message directly to the credential store. Nothing touches
|
|
// logs, turns, history, or compaction.
|
|
let pending_auth = {
|
|
let sess = session.lock().await;
|
|
sess.threads
|
|
.get(&thread_id)
|
|
.and_then(|t| t.pending_auth.clone())
|
|
};
|
|
|
|
if let Some(pending) = pending_auth {
|
|
if pending.is_expired() {
|
|
// TTL exceeded — clear stale auth mode
|
|
tracing::warn!(
|
|
extension = %pending.extension_name,
|
|
"Auth mode expired after TTL, clearing"
|
|
);
|
|
{
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.pending_auth = None;
|
|
}
|
|
}
|
|
// If this was a user message (possibly a pasted token), return an
|
|
// explicit error instead of forwarding it to the LLM/history.
|
|
if matches!(submission, Submission::UserInput { .. }) {
|
|
return Ok(Some(format!(
|
|
"Authentication for **{}** expired. Please try again.",
|
|
pending.extension_name
|
|
)));
|
|
}
|
|
// Control submissions (interrupt, undo, etc.) fall through to normal handling
|
|
} else {
|
|
match &submission {
|
|
Submission::UserInput { content } => {
|
|
return self
|
|
.process_auth_token(message, &pending, content, session, thread_id)
|
|
.await;
|
|
}
|
|
_ => {
|
|
// Any control submission (interrupt, undo, etc.) cancels auth mode
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.pending_auth = None;
|
|
}
|
|
// Fall through to normal handling
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
tracing::trace!(
|
|
"Received message from {} on {} ({} chars)",
|
|
message.user_id,
|
|
message.channel,
|
|
message.content.len()
|
|
);
|
|
|
|
if !message.is_internal
|
|
&& let Submission::UserInput { ref content } = submission
|
|
&& let Some(engine) = self.routine_engine().await
|
|
{
|
|
let single_message_repl = is_single_message_repl(message);
|
|
// Use post-hook content so that BeforeInbound hooks that rewrite
|
|
// input are respected by event trigger matching.
|
|
let fired = if single_message_repl {
|
|
engine.check_event_triggers_and_wait(message, content).await
|
|
} else {
|
|
engine.check_event_triggers(message, content).await
|
|
};
|
|
if fired > 0 {
|
|
tracing::debug!(
|
|
channel = %message.channel,
|
|
user = %message.user_id,
|
|
fired,
|
|
"Consumed inbound user message with matching event-triggered routine(s)"
|
|
);
|
|
return if single_message_repl {
|
|
Ok(None)
|
|
} else {
|
|
Ok(Some(String::new()))
|
|
};
|
|
}
|
|
}
|
|
|
|
// Build per-tenant execution context once; threaded through all handlers.
|
|
let tenant = self.tenant_ctx(&message.user_id).await;
|
|
|
|
// Per-user bootstrap: if this user's workspace was just seeded (fresh),
|
|
// persist the static greeting to their assistant conversation and
|
|
// broadcast it so the web client shows it immediately.
|
|
if tenant
|
|
.workspace()
|
|
.is_some_and(|ws| ws.take_bootstrap_pending())
|
|
{
|
|
tracing::info!(
|
|
user_id = message.user_id,
|
|
"Fresh user workspace — persisting bootstrap greeting"
|
|
);
|
|
if let Some(store) = tenant.store()
|
|
&& let Ok(conv_id) = store
|
|
.get_or_create_assistant_conversation(&message.channel)
|
|
.await
|
|
{
|
|
let _ = store
|
|
.add_conversation_message(conv_id, "assistant", BOOTSTRAP_GREETING)
|
|
.await;
|
|
let mut out = OutgoingResponse::text(BOOTSTRAP_GREETING.to_string());
|
|
out.thread_id = Some(conv_id.to_string());
|
|
let _ = self
|
|
.channels
|
|
.broadcast(&message.channel, &message.user_id, out)
|
|
.await;
|
|
}
|
|
}
|
|
|
|
let session_for_empty_exit = Arc::clone(&session);
|
|
|
|
// Process based on submission type
|
|
let result = match submission {
|
|
Submission::UserInput { content } => {
|
|
let mut result = self
|
|
.process_user_input(
|
|
message,
|
|
tenant.clone(),
|
|
session.clone(),
|
|
thread_id,
|
|
&content,
|
|
)
|
|
.await;
|
|
|
|
// Drain any messages queued during processing.
|
|
// Messages are merged (newline-separated) so the LLM receives
|
|
// full context from rapid consecutive inputs instead of
|
|
// processing each as a separate turn with partial context (#259).
|
|
//
|
|
// Only `Response` continues the drain — the user got a normal
|
|
// reply and there may be more queued messages to process.
|
|
//
|
|
// Everything else stops the loop:
|
|
// - `NeedApproval`: thread is blocked on user approval
|
|
// - `Interrupted`: turn was cancelled
|
|
// - `Ok`: control-command acknowledgment (including the "queued"
|
|
// ack returned when a message arrives during Processing)
|
|
// - `Error`: soft error — draining more messages after an error
|
|
// would produce confusing interleaved output
|
|
// - `Err(_)`: hard error
|
|
while let Ok(SubmissionResult::Response { content: outgoing }) = &result {
|
|
let merged = {
|
|
let mut sess = session.lock().await;
|
|
sess.threads
|
|
.get_mut(&thread_id)
|
|
.and_then(|t| t.drain_pending_messages())
|
|
};
|
|
let Some(next_content) = merged else {
|
|
break;
|
|
};
|
|
|
|
tracing::debug!(
|
|
thread_id = %thread_id,
|
|
merged_len = next_content.len(),
|
|
"Drain loop: processing merged queued messages"
|
|
);
|
|
|
|
// Send the completed turn's response before starting the next.
|
|
//
|
|
// Known limitations:
|
|
// - One-shot channels (HttpChannel) consume the response
|
|
// sender on the first respond() call keyed by msg.id.
|
|
// Subsequent calls (including the outer handler's final
|
|
// respond) are silently dropped. For one-shot channels
|
|
// only this intermediate response is delivered.
|
|
// - All drain-loop responses are routed via the original
|
|
// `message`, so channels that key routing on message
|
|
// identity will attribute every response to the first
|
|
// message. This is acceptable for the current
|
|
// single-user-per-thread model.
|
|
if let Err(e) = self
|
|
.channels
|
|
.respond(message, OutgoingResponse::text(outgoing.clone()))
|
|
.await
|
|
{
|
|
tracing::warn!(
|
|
thread_id = %thread_id,
|
|
"Failed to send intermediate drain-loop response: {e}"
|
|
);
|
|
}
|
|
|
|
// Process merged queued messages as a single turn.
|
|
// Use a message clone with cleared attachments so
|
|
// augment_with_attachments doesn't re-apply the original
|
|
// message's attachments to unrelated queued text.
|
|
let mut queued_msg = message.clone();
|
|
queued_msg.attachments.clear();
|
|
result = self
|
|
.process_user_input(
|
|
&queued_msg,
|
|
tenant.clone(),
|
|
session.clone(),
|
|
thread_id,
|
|
&next_content,
|
|
)
|
|
.await;
|
|
|
|
// If processing failed, re-queue the drained content so it
|
|
// isn't lost. It will be picked up on the next successful turn.
|
|
if !matches!(&result, Ok(SubmissionResult::Response { .. })) {
|
|
let mut sess = session.lock().await;
|
|
if let Some(thread) = sess.threads.get_mut(&thread_id) {
|
|
thread.requeue_drained(next_content);
|
|
tracing::debug!(
|
|
thread_id = %thread_id,
|
|
"Re-queued drained content after non-Response result"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
|
|
result
|
|
}
|
|
Submission::SystemCommand { command, args } => {
|
|
tracing::debug!(
|
|
"[agent_loop] SystemCommand: command={}, channel={}",
|
|
command,
|
|
message.channel
|
|
);
|
|
// /reasoning is special-cased here (not in handle_system_command)
|
|
// because it needs the session + thread_id to read turn reasoning
|
|
// data, which handle_system_command's signature doesn't provide.
|
|
if command == "reasoning" {
|
|
let result = self
|
|
.handle_reasoning_command(&args, &session, thread_id)
|
|
.await;
|
|
return match result {
|
|
SubmissionResult::Response { content } => Ok(Some(content)),
|
|
SubmissionResult::Ok { message } => Ok(message),
|
|
SubmissionResult::Error { message } => {
|
|
Ok(Some(format!("Error: {}", message)))
|
|
}
|
|
_ => {
|
|
if is_single_message_repl(message) {
|
|
Ok(None)
|
|
} else {
|
|
Ok(Some(String::new()))
|
|
}
|
|
}
|
|
};
|
|
}
|
|
// Authorization checks (including restart channel check) are enforced in handle_system_command
|
|
self.handle_system_command(&command, &args, &message.channel, &tenant)
|
|
.await
|
|
}
|
|
Submission::Undo => self.process_undo(session, thread_id).await,
|
|
Submission::Redo => self.process_redo(session, thread_id).await,
|
|
Submission::Interrupt => self.process_interrupt(session, thread_id).await,
|
|
Submission::Compact => self.process_compact(session, thread_id).await,
|
|
Submission::Clear => self.process_clear(session, thread_id).await,
|
|
Submission::NewThread => self.process_new_thread(message).await,
|
|
Submission::Heartbeat => self.process_heartbeat().await,
|
|
Submission::Summarize => self.process_summarize(session, thread_id).await,
|
|
Submission::Suggest => self.process_suggest(session, thread_id).await,
|
|
Submission::JobStatus { job_id } => {
|
|
self.process_job_status(&tenant, job_id.as_deref()).await
|
|
}
|
|
Submission::JobCancel { job_id } => self.process_job_cancel(&tenant, &job_id).await,
|
|
Submission::Quit => return Ok(None),
|
|
Submission::SwitchThread { thread_id: target } => {
|
|
self.process_switch_thread(message, target).await
|
|
}
|
|
Submission::Resume { checkpoint_id } => {
|
|
self.process_resume(session, thread_id, checkpoint_id).await
|
|
}
|
|
Submission::ExecApproval {
|
|
request_id,
|
|
approved,
|
|
always,
|
|
} => {
|
|
self.process_approval(
|
|
message,
|
|
session,
|
|
thread_id,
|
|
Some(request_id),
|
|
approved,
|
|
always,
|
|
)
|
|
.await
|
|
}
|
|
Submission::ApprovalResponse { approved, always } => {
|
|
self.process_approval(message, session, thread_id, None, approved, always)
|
|
.await
|
|
}
|
|
};
|
|
|
|
// Convert SubmissionResult to response string
|
|
match result? {
|
|
SubmissionResult::Response { content } => {
|
|
// Suppress silent replies (e.g. from group chat "nothing to say" responses)
|
|
if crate::llm::is_silent_reply(&content) {
|
|
tracing::debug!("Suppressing silent reply token");
|
|
Ok(None)
|
|
} else {
|
|
Ok(Some(content))
|
|
}
|
|
}
|
|
SubmissionResult::Ok {
|
|
message: output_message,
|
|
} => {
|
|
let should_exit =
|
|
if output_message.as_deref() == Some("") && is_single_message_repl(message) {
|
|
let sess = session_for_empty_exit.lock().await;
|
|
sess.threads
|
|
.get(&thread_id)
|
|
.map(|thread| thread.state != ThreadState::AwaitingApproval)
|
|
.unwrap_or(true)
|
|
} else {
|
|
false
|
|
};
|
|
|
|
if should_exit {
|
|
Ok(None)
|
|
} else {
|
|
Ok(output_message)
|
|
}
|
|
}
|
|
SubmissionResult::Error { message } => Ok(Some(format!("Error: {}", message))),
|
|
SubmissionResult::Interrupted => Ok(Some("Interrupted.".into())),
|
|
SubmissionResult::NeedApproval { .. } => {
|
|
// ApprovalNeeded status was already sent by thread_ops.rs before
|
|
// returning this result. Empty string signals the caller to skip
|
|
// respond() (no duplicate text).
|
|
Ok(Some(String::new()))
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::{
|
|
chat_tool_execution_metadata, is_single_message_repl, resolve_routine_notification_user,
|
|
should_fallback_routine_notification, truncate_for_preview,
|
|
};
|
|
use crate::channels::IncomingMessage;
|
|
use crate::error::ChannelError;
|
|
|
|
#[test]
|
|
fn test_truncate_short_input() {
|
|
assert_eq!(truncate_for_preview("hello", 10), "hello");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_empty_input() {
|
|
assert_eq!(truncate_for_preview("", 10), "");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_exact_length() {
|
|
assert_eq!(truncate_for_preview("hello", 5), "hello");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_over_limit() {
|
|
let result = truncate_for_preview("hello world, this is long", 10);
|
|
assert!(result.ends_with("..."));
|
|
// "hello worl" = 10 chars + "..."
|
|
assert_eq!(result, "hello worl...");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_collapses_newlines() {
|
|
let result = truncate_for_preview("line1\nline2\nline3", 100);
|
|
assert!(!result.contains('\n'));
|
|
assert_eq!(result, "line1 line2 line3");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_collapses_whitespace() {
|
|
let result = truncate_for_preview("hello world", 100);
|
|
assert_eq!(result, "hello world");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_multibyte_utf8() {
|
|
// Each emoji is 4 bytes. Truncating at char boundary must not panic.
|
|
let input = "😀😁😂🤣😃😄😅😆😉😊";
|
|
let result = truncate_for_preview(input, 5);
|
|
assert!(result.ends_with("..."));
|
|
// First 5 chars = 5 emoji
|
|
assert_eq!(result, "😀😁😂🤣😃...");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_cjk_characters() {
|
|
// CJK chars are 3 bytes each in UTF-8.
|
|
let input = "你好世界测试数据很长的字符串";
|
|
let result = truncate_for_preview(input, 4);
|
|
assert_eq!(result, "你好世界...");
|
|
}
|
|
|
|
#[test]
|
|
fn test_truncate_mixed_multibyte_and_ascii() {
|
|
let input = "hello 世界 foo";
|
|
let result = truncate_for_preview(input, 8);
|
|
// 'h','e','l','l','o',' ','世','界' = 8 chars
|
|
assert_eq!(result, "hello 世界...");
|
|
}
|
|
|
|
#[test]
|
|
fn resolve_routine_notification_user_prefers_explicit_target() {
|
|
let metadata = serde_json::json!({
|
|
"notify_user": "12345",
|
|
"owner_id": "owner-scope",
|
|
});
|
|
|
|
let resolved = resolve_routine_notification_user(&metadata);
|
|
assert_eq!(resolved.as_deref(), Some("12345")); // safety: test-only assertion
|
|
}
|
|
|
|
#[test]
|
|
fn resolve_routine_notification_user_falls_back_to_owner_scope() {
|
|
let metadata = serde_json::json!({
|
|
"notify_user": null,
|
|
"owner_id": "owner-scope",
|
|
});
|
|
|
|
let resolved = resolve_routine_notification_user(&metadata);
|
|
assert_eq!(resolved.as_deref(), Some("owner-scope")); // safety: test-only assertion
|
|
}
|
|
|
|
#[test]
|
|
fn resolve_routine_notification_user_rejects_missing_values() {
|
|
let metadata = serde_json::json!({
|
|
"notify_user": " ",
|
|
});
|
|
|
|
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 {
|
|
name: "telegram".to_string(),
|
|
reason: "No stored owner routing target for channel 'telegram'.".to_string(),
|
|
};
|
|
|
|
assert!(!should_fallback_routine_notification(&error)); // safety: test-only assertion
|
|
}
|
|
|
|
#[test]
|
|
fn targeted_routine_notifications_may_fallback_for_other_errors() {
|
|
let error = ChannelError::SendFailed {
|
|
name: "telegram".to_string(),
|
|
reason: "timeout talking to channel".to_string(),
|
|
};
|
|
|
|
assert!(should_fallback_routine_notification(&error)); // safety: test-only assertion
|
|
}
|
|
|
|
#[test]
|
|
fn single_message_repl_detection_requires_repl_channel_and_metadata_flag() {
|
|
let repl = IncomingMessage::new("repl", "owner-scope", "hello")
|
|
.with_metadata(serde_json::json!({ "single_message_mode": true }));
|
|
let gateway = IncomingMessage::new("gateway", "owner-scope", "hello")
|
|
.with_metadata(serde_json::json!({ "single_message_mode": true }));
|
|
let plain_repl = IncomingMessage::new("repl", "owner-scope", "hello");
|
|
|
|
assert!(is_single_message_repl(&repl)); // safety: test-only assertion
|
|
assert!(!is_single_message_repl(&gateway)); // safety: test-only assertion
|
|
assert!(!is_single_message_repl(&plain_repl)); // safety: test-only assertion
|
|
}
|
|
}
|