Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion docs/memory/architecture/background-services.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ Services that run continuously in the backend process:

| Service | Module | Description |
|---------|--------|-------------|
| **Cleanup Service** | `cleanup_service.py` | Every 5 min: active watchdog reconciliation against agent process registries (orphan recovery, auto-terminate timeouts) + passive stale recovery (CLEANUP-001, #129). Also runs retention + soft-delete purge sweeps, the **expired-SSH sweep** (`_sweep_expired_ssh_credentials` → `SshService.cleanup_expired_credentials` — removes an expired ephemeral key's line from the container `authorized_keys` sshd reads; TTL was previously enforced only on Redis metadata, #1616), and the #740 startup orphan-loop hook, and the **agent_reminders retention sweep** (`_sweep_agent_reminders_retention` — DELETEs terminal `fired`/`cancelled`/`failed` reminders past `agent_reminders_retention_days`, #1296) — see [Soft Delete & Retention](reliability.md#soft-delete-retention--recovery-834-772). Runs the additive **lease-reaper** (`lease_reaper_service`) each cycle — re-queues (preserving `execution_id`) or poison-parks expired pull leases (#1081 Phase 3, #429/#1402; inert until an agent is piloted). #1804: every recovery path also closes its execution's dispatch activity (`_close_bulk_swept_activities` for the bulk sweeps, the shared helper elsewhere), counted in `activities_closed_on_recovery`; the 120-minute activity backstop now runs **last** in the cycle. **#2433:** orphan = the agent does not know the execution **and** no live backend dispatcher owns it — see [In-Flight Dispatch Proof-of-Life](execution.md#in-flight-dispatch-proof-of-life-2433) |
| **Cleanup Service** | `cleanup_service.py` | Every 5 min: active watchdog reconciliation against agent process registries (orphan recovery, auto-terminate timeouts) + passive stale recovery (CLEANUP-001, #129). Also runs retention + soft-delete purge sweeps, the **expired-SSH sweep** (`_sweep_expired_ssh_credentials` → `SshService.cleanup_expired_credentials` — removes an expired ephemeral key's line from the container `authorized_keys` sshd reads; TTL was previously enforced only on Redis metadata, #1616), and the #740 startup loop hook (since #2523 `LoopService.reconcile_after_restart`: re-arms or advances rather than interrupting — see the Loop Due-Run Sweep row), and the **agent_reminders retention sweep** (`_sweep_agent_reminders_retention` — DELETEs terminal `fired`/`cancelled`/`failed` reminders past `agent_reminders_retention_days`, #1296) — see [Soft Delete & Retention](reliability.md#soft-delete-retention--recovery-834-772). Runs the additive **lease-reaper** (`lease_reaper_service`) each cycle — re-queues (preserving `execution_id`) or poison-parks expired pull leases (#1081 Phase 3, #429/#1402; inert until an agent is piloted). #1804: every recovery path also closes its execution's dispatch activity (`_close_bulk_swept_activities` for the bulk sweeps, the shared helper elsewhere), counted in `activities_closed_on_recovery`; the 120-minute activity backstop now runs **last** in the cycle. **#2433:** orphan = the agent does not know the execution **and** no live backend dispatcher owns it — see [In-Flight Dispatch Proof-of-Life](execution.md#in-flight-dispatch-proof-of-life-2433) |
| **Operator Queue Sync** | `operator_queue_service.py` | Polls running agents every 5s, reads `~/.trinity/operator-queue.json`, syncs to DB, writes responses back (OPS-001). The item `id` is a platform-minted uuid; the agent's correlation string is `request_id` with `(agent_name, request_id)` uniqueness, so all sync reads/writes (exists, acknowledge, response write-back) are agent-scoped and two agents can't collide (#1631). **Leader-locked (#1632):** only the holder of `opqueue:leader` (SET NX, TTL `max(3×interval, 30s)` floor so a slow-write cycle can't flap leadership, own-lease refresh, fail-open — mirror monitoring #1464) runs a cycle, so `--workers 2` doesn't double-charge the ingestion rate limiter or double-broadcast the flood alert. Ingestion is capped per agent (depth + rate + fleet + field hygiene, #1632) — see [Operator Queue](api-endpoints.md#operator-queue-ops-001). **Honest since #2915:** the cycle reads the tri-state `agent_container_states()` (`None` ⇒ nothing swept or synced — the #2196 class), runs expiry BEFORE the no-running-agents return (since trinity-enterprise#611 through the ask sink, `ask_service.expire()`: per-id CAS with the endings ledger `disposed_by = 'timeout'`, one `expired` audit row per ask, the ending wake for opted-in running filers; `_changed_this_cycle` is reset above it so an expiry rides the cycle's one trigger, even when Docker is unreadable), sweeps open rows of stopped agents to `unconfirmed:agent_not_running` and their owed answers / terminal flips to `undelivered:agent_not_running` (both sweeps skip platform-minted rows in SQL), reconciles every file entry against the row (`sync_state`/`sync_detail`, edge-triggered by rowcount, 3-cycle hysteresis on read failures, no write-back on an unparseable or wrong-shape file, first entry wins on a duplicated id), delivers an answer only into a matching still-pending entry after a pre-write re-read with `if_match` (412 = retry, never clobber), records every delivery outcome, audits accountability transitions only, and broadcasts ONE `operator_queue_sync` per cycle — see [operating-room.md → Sync honesty](../feature-flows/operating-room.md#sync-honesty-2915). **Native asks (trinity-enterprise#611):** a row raised over the platform (`channel=mcp`) is outside the file contract — the sync index names its ids as `foreign`, and the cycle skips a file entry re-using one before any branch reads it (logged once per agent and id); the cycle creates through `create_operator_queue_item_with_outcome` and never counts, announces or audits a row it did not insert (a native ask raised between the index read and the create); the first file ingest per agent per process logs the file channel's deprecation (`ask_operator`, two releases) |
| **Sync Health Service** | `sync_health_service.py` | Polls git-enabled agents every 60s (`SYNC_HEALTH_POLL_INTERVAL_SECONDS`, default unchanged) — see [Git Sync Health](agent-lifecycle.md#git-sync-health-389390). **Leader-locked (#2742):** `synchealth:leader`, SET NX, TTL `max(3×interval, 30s)`, own-lease refresh, compare-and-delete release (a GET-then-DEL lets an expired worker delete a sibling's fresh lease), **fail-open** — Redis down ⇒ every worker polls, i.e. the pre-#2742 behaviour, and the agent-side lock-free + single-flighted status (#2742) makes a duplicated poll harmless. Failing closed was rejected: this feed is how `sync_failing` is ever raised, and darkening it exactly when infra is degraded is the worse error. **Named side effect:** `upsert_sync_state` *increments* `consecutive_failures` per failed poll, so it is not idempotent — one leader instead of two doubles time-to-`sync_failing` from ~90s to ~180s (arguably making the counter mean what its name says: 3 consecutive failed polls = 3 minutes). **Also raises the `sync_diverged` episode item (trinity-enterprise#706)** under the same lease; its id is deterministic per episode (`sync-diverged-{agent}-{diverged_since}`), so a fail-open double poll converges on one row through the `(agent_name, request_id)` conflict target |
| **Skills Library Sync** | `skills_sync_service.py` | Scheduled skills-library `git pull` + optional fleet-wide skill re-inject (ent#236). Runs in every worker but only the `skills:sync:leader` lease-holder performs a cycle (fail-open, mirrors #1464); self-gates on the default-OFF `skills_library_auto_sync_enabled` setting, re-read each cycle so an interval change needs no restart. A sweep fires only on a changed library commit, targets running non-ghost agents at `SKILLS_FLEET_INJECT_CONCURRENCY` (5), and persists an honest per-agent report + raises an operator alarm on any failure |
Expand All @@ -25,6 +25,7 @@ Services that run continuously in the backend process:
| **Heartbeat Watch Loop** | `heartbeat_service.py` | 5s loop acting on missed agent heartbeats — see [Heartbeat Liveness](reliability.md#heartbeat-liveness-reliability-004-307) |
| **Scheduler Service** | `scheduler_service.py` | APScheduler cron execution; async fire-and-forget with DB polling for status. On each cron fire, optionally invokes the agent's `~/.trinity/pre-check` (see Agent Containers), preceded — for a seat-delivery schedule — by the readiness gate (ent#689, `GET /api/internal/agents/{name}/brief-readiness`; both skips share `_record_gate_skip`, both fail open). Also owns one-shot `DateTrigger`s for RETRY-001 retries and **agent self-reminders** (#1296): `_reconcile_reminders` arms pending reminders + reclaims stale `firing` rows at boot, in the 60s sync loop (own try/except), and on full reload — see [Agent Self-Reminders](execution.md#agent-self-reminders-1296) |
| **Capacity Maintenance** | `capacity_manager.py` | `run_maintenance()` every 60s — see [Capacity & Backlog](execution.md#capacity--backlog-428) |
| **Loop Due-Run Sweep** | `loop_service.py` | #2523: brings back a loop parked on `agent_loops.next_run_at` by `delay_seconds`. Started in `main.py::_start_capacity_and_canary` (`_loop_due_sweep`, first tick after 8–10s), then every 5–6s: `LoopService.dispatch_due_loops` → `db.list_due_loops(now)`. Runs in every uvicorn worker with **no leader lease** — the per-loop claim (`claim_due_loop`) is a CAS on the exact `next_run_at` read, so one worker dispatches and the rest see nothing. **Idle cost (#3436):** with no loop parked each tick is one index seek on `idx_loops_next_run` returning no rows (≈6 µs measured on a development database, ≈15,700 queries/day per worker); deliberately not gated — see [scheduling §38.9](../requirements/scheduling.md#389-idle-cost-of-the-due-loop-sweep-3436) |
| **Audit Retention** | `audit_retention_service.py` | Daily 04:15 UTC: DELETEs `audit_log` rows past retention. `AUDIT_LOG_RETENTION_DAYS` (default 365, floored at 365 — the `audit_log_no_delete` trigger refuses younger rows). Pruning ages out hash-chain history past the cutoff by design (#552) |
| **DB Vacuum** | `db_vacuum_service.py` | Daily 04:30 UTC: `VACUUM` on `/data/trinity.db` to reclaim pages freed by retention sweeps. `DB_VACUUM_ENABLED`/`DB_VACUUM_HOUR`/`DB_VACUUM_MINUTE`. Autocommit connection (VACUUM can't run in a transaction); accepts rare BUSY rather than retrying (#772) |
| **DB Backup** | `db_backup_service.py` | Daily **03:30 UTC** (before the destructive 04:15/04:30 jobs — capture-more-data ordering): a verified recovery point under `/data/backups/` for **both** backends — SQLite via to_thread'd stdlib `Connection.backup()`, PostgreSQL via `pg_dump -Fc` (`postgresql-client-17` baked into the backend image). Day-keyed artifacts + a fail-open SETNX lease (`db_backup:running`, duplicate-I/O suppression only). Prune (window + fixed `MIN_KEEP=3` floor) + staleness check run in the tail of EVERY attempt. `DB_BACKUP_ENABLED`/`DB_BACKUP_HOUR`/`DB_BACKUP_MINUTE`/`DB_BACKUP_PG_DUMP_TIMEOUT_SECONDS` (forwarded in all three compose files — dev, prod and the #2280 hosted file, whose prod parity is CI-guarded by `tests/unit/test_2280_hosted_compose_parity.py`). Default ON — see [Automatic Database Backups](reliability.md#automatic-database-backups-2216) (#2216) |
Expand Down
2 changes: 1 addition & 1 deletion docs/memory/architecture/execution.md
Original file line number Diff line number Diff line change
Expand Up @@ -155,7 +155,7 @@ Trigger-boundary dedup — policy in Architectural Invariant #18, table DDL unde

### Sequential Agent Loops (#740, UI #1106)

Bounded sequential task execution against one agent. Runner is an in-process `asyncio.Task` spawned by `loop_service.py`; each iteration dispatches through `task_execution_service.execute_task()` with `triggered_by="loop"` and the parent `loop_id` carried on the resulting `schedule_executions` row — iterations go through the standard `capacity_manager` admit/slot path, sharing the agent's `max_parallel_tasks` budget. Message template supports `{{run}}` and `{{previous_response}}`; `max_runs` 1–100 hard cap; optional `stop_signal` (until-mode), `delay_seconds`, `timeout_per_run`, `max_duration_seconds`, `model`, `allowed_tools`. Stop is cooperative: `POST /api/loops/{id}/stop` flips an in-process `should_stop` flag; the current iteration finishes and the runner exits with `stop_reason="user_stopped"`. **Wall-clock deadline (#1156):** optional `max_duration_seconds` (≤7 days) measured from `started_at`, checked only at iteration boundaries (before the next run and before/after the inter-run delay, which is capped to the remaining budget) — an in-flight run is never killed mid-turn, so overshoot is bounded by one `timeout_per_run`; expiry stops the loop with `stop_reason="deadline_exceeded"`. Rejected at create (400) when smaller than the effective per-run timeout (`timeout_per_run`, else the agent's `execution_timeout_seconds`). **Cost budget (#1155):** optional `max_cost_usd` (`gt=0`, no upper cap) — an iteration-boundary gate enforced *after* the deadline check: the runner accumulates each completed run's cost (only finite, positive values; NULL/unknown counts as 0 fail-open; NaN/inf ignored so it can't poison the accumulator; both unusable-cost cases WARN under an active budget) and stops *before the next run* with `stop_reason="budget_exhausted"` once accumulated cost meets/exceeds the budget. **No-progress / doom-loop detection (#1157):** optional `no_progress_threshold` (`0` disables; **default 3** for new loops via the API/MCP; NULL ⇒ disabled, so in-flight loops created before this change are unaffected). The runner fingerprints each successful run's full response — SHA-256 of normalized text (`" ".join(text.split())`, so word boundaries are preserved and whitespace-only/empty all collapse to one fingerprint) — and stops the loop with `stop_reason="no_progress"` (status `stopped`) once K consecutive runs share a fingerprint. Counter + last-fingerprint are runner-local (no persistence). Detection is **exact-hash only** (no fuzzy/semantic similarity). The validator rejects `1` (422 — "repeated identical" needs ≥2). **Boundary-only precedence** (per iteration: `user_stopped` → `deadline_exceeded` → `budget_exhausted` → run → `stop_signal_matched` → `no_progress`; natural exit `max_runs_reached`): the current run always finishes, so one run — including the first — can overshoot; a run that crosses the budget but is also the final `max_runs` run or matches `stop_signal` yields those reasons instead, and a pending `user_stopped`/`deadline_exceeded` outranks `no_progress`. `GET /api/loops/{id}` returns `max_duration_seconds` + computed `elapsed_seconds`, plus `max_cost_usd` + `total_cost` (computed on read = sum of `agent_loop_runs.cost`, NULL→0; `0.0` for a zero-run loop). Restart recovery via the cleanup-service startup hook (above); no auto-resume. WS events `loop_run_completed`/`loop_completed`.
Bounded sequential task execution against one agent. Since #2523 there is no runner: the `agent_loops` row is the loop, and `loop_service.py` advances it from each iteration's execution terminal (and from the due-loop sweep after a `delay_seconds` park) — see [run-agent-loop.md](../feature-flows/run-agent-loop.md); each iteration dispatches through `task_execution_service.execute_task()` with `triggered_by="loop"` and the parent `loop_id` carried on the resulting `schedule_executions` row — iterations go through the standard `capacity_manager` admit/slot path, sharing the agent's `max_parallel_tasks` budget. Message template supports `{{run}}` and `{{previous_response}}`; `max_runs` 1–100 hard cap; optional `stop_signal` (until-mode), `delay_seconds`, `timeout_per_run`, `max_duration_seconds`, `model`, `allowed_tools`. Stop is cooperative: `POST /api/loops/{id}/stop` stamps `agent_loops.stop_requested_at`, so any worker can serve it; the current iteration finishes and the loop finalizes with `stop_reason="user_stopped"` (a loop parked on `next_run_at` is finalized at once). **Wall-clock deadline (#1156):** optional `max_duration_seconds` (≤7 days) measured from `started_at`, checked only at iteration boundaries (before the next run and before/after the inter-run delay, which is capped to the remaining budget) — an in-flight run is never killed mid-turn, so overshoot is bounded by one `timeout_per_run`; expiry stops the loop with `stop_reason="deadline_exceeded"`. Rejected at create (400) when smaller than the effective per-run timeout (`timeout_per_run`, else the agent's `execution_timeout_seconds`). **Cost budget (#1155):** optional `max_cost_usd` (`gt=0`, no upper cap) — an iteration-boundary gate enforced *after* the deadline check: the advance accumulates each completed run's cost (only finite, positive values; NULL/unknown counts as 0 fail-open; NaN/inf ignored so it can't poison the accumulator; both unusable-cost cases WARN under an active budget) and stops *before the next run* with `stop_reason="budget_exhausted"` once accumulated cost meets/exceeds the budget. **No-progress / doom-loop detection (#1157):** optional `no_progress_threshold` (`0` disables; **default 3** for new loops via the API/MCP; NULL ⇒ disabled, so in-flight loops created before this change are unaffected). The advance fingerprints each successful run's full response — SHA-256 of normalized text (`" ".join(text.split())`, so word boundaries are preserved and whitespace-only/empty all collapse to one fingerprint) — and stops the loop with `stop_reason="no_progress"` (status `stopped`) once K consecutive runs share a fingerprint. Counter + last-fingerprint are not stored; since #2523 they are rebuilt from the loop's `agent_loop_runs` rows (`response` of each completed run) on every advance. Detection is **exact-hash only** (no fuzzy/semantic similarity). The validator rejects `1` (422 — "repeated identical" needs ≥2). **Boundary-only precedence** (per iteration: `user_stopped` → `deadline_exceeded` → `budget_exhausted` → run → `stop_signal_matched` → `no_progress`; natural exit `max_runs_reached`): the current run always finishes, so one run — including the first — can overshoot; a run that crosses the budget but is also the final `max_runs` run or matches `stop_signal` yields those reasons instead, and a pending `user_stopped`/`deadline_exceeded` outranks `no_progress`. `GET /api/loops/{id}` returns `max_duration_seconds` + computed `elapsed_seconds`, plus `max_cost_usd` + `total_cost` (computed on read = sum of `agent_loop_runs.cost`, NULL→0; `0.0` for a zero-run loop). Restart recovery via the cleanup-service startup hook (`cleanup_service._cleanup_loop` → `LoopService.reconcile_after_restart`, see [background-services.md](background-services.md)), which since #2523 re-arms or advances a loop rather than interrupting it. WS events `loop_run_completed`/`loop_completed`.

**Failure policy (#1167):** per-loop `on_failure` — `abort` (default; fail-fast, first failed iteration ends the loop `failed`/`stop_reason=error`) or `continue` (tolerate a failed iteration and proceed). Both failure surfaces are gated: a raised exception from `execute_task` and a non-success `TaskExecutionResult`. Continue mode is bounded by `max_consecutive_failures` (default 3) — once that many iterations fail in a row the loop aborts `failed`/`stop_reason=max_consecutive_failures`; a success resets the streak. A continue-mode loop that reaches `max_runs` (or matches its stop-signal) with ≥1 tolerated failure finalizes as `completed_with_errors`, with the `failed_runs` count surfaced. `{{previous_response}}` always carries the last *successful* response (a failed iteration never overwrites it).

Expand Down
Loading
Loading