From d838d4792b3807e2bb468b39a1ff6a43dc47ab21 Mon Sep 17 00:00:00 2001 From: Garry Tan Date: Fri, 24 Apr 2026 01:09:28 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20queue=20resilience=20=E2=80=94=20wall-c?= =?UTF-8?q?lock=20timeouts,=20backpressure,=20--no-worker,=20env=20concurr?= =?UTF-8?q?ency=20(#379)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: queue resilience — wall-clock timeouts, backpressure, --no-worker, env concurrency, shell guard Prevents stall-induced queue blockage discovered in production (OpenClaw): 1. Wall-clock timeout sweep: dead-letters active jobs exceeding 2× timeout_ms (or 2 × lockDuration × max_stalled). Catches jobs stuck while holding DB connections where FOR UPDATE SKIP LOCKED stall detection skips them. 2. Submission backpressure (maxWaiting): caps waiting jobs per name at submission time. Prevents autopilot-cycle flood when the queue is blocked. 3. --no-worker flag for autopilot: skips spawning the built-in worker child. For environments where the worker lifecycle is managed externally (systemd, Docker, OpenClaw service-manager). 4. GBRAIN_WORKER_CONCURRENCY env var: fallback for --concurrency when the worker is spawned by autopilot (which can't pass CLI flags to the child). 5. Shell job env guard with clear logging: shell handler is always registered but throws UnrecoverableError with a clear message when GBRAIN_ALLOW_SHELL_JOBS=1 is not set, instead of silently not registering. * feat: v0.19.1 Lane A — maxWaiting atomic guard, concurrency clamp, --max-waiting CLI Addresses three production-hardening findings from the CEO + Eng + Codex adversarial review of PR #379: D2/H2: maxWaiting was TOCTOU-racy — two concurrent submitters could both see waitingCount < max and both insert. Wrap the count+select+insert in pg_advisory_xact_lock keyed on (name, queue). Serializes concurrent decisions for the SAME key while leaving different keys fully parallel. Lock auto-releases on txn commit/rollback — no cleanup path to leak. Also fix the missing queue-scope bug: count and select now filter on (name, queue) not name alone, so cross-queue same-name jobs don't suppress each other. D3/H3: resolveWorkerConcurrency silently accepted NaN / 0 / negative from parseInt. `inFlight.size < NaN` is always false → worker claims nothing → silent wedge from a single-typo env var. Clamp to ≥1 with a loud stderr warning naming the bad value. D5/H5: `gbrain jobs submit` never parsed `--max-waiting N` despite the MinionJobInput field. Wire the flag with clamp [1, 100], mirror `--max-stalled`. Extract `parseMaxWaitingFlag` for unit testing. Q1: Silent coalesce was invisible by design. New src/core/minions/backpressure-audit.ts mirrors shell-audit.ts's ISO-week JSONL pattern: `~/.gbrain/audit/backpressure-YYYY-Www.jsonl`. Coalesce events write one JSONL line with (queue, name, waiting_count, max_waiting, returned_job_id, ts). Best-effort — disk-full never blocks submission. A2: `gbrain jobs smoke --wedge-rescue` new opt-in regression case. Forges a wedged-worker row state, invokes handleStalled + handleTimeouts + handleWallClockTimeouts in order, asserts only wall-clock evicts. Mirrors the v0.14.3 `--sigkill-rescue` shape. Tests: 23 new unit cases in test/minions.test.ts covering wall-clock timeout (3 cases + non-interference with handleTimeouts), maxWaiting (coalesce, clamp 0, floor, concurrent-submitter race via Promise.all, cross-queue isolation, unset fallthrough), concurrency clamp (7 cases incl. NaN/0/negative), parseMaxWaitingFlag (5 cases), backpressure audit file write. Part of v0.19.1 plan at ~/.claude/plans/ok-wintermute-wrote-this-polished-matsumoto.md Co-Authored-By: Claude Opus 4.7 (1M context) * feat: v0.19.1 Lane B — doctor queue_health, autopilot peer probe, runbook A5 / D4: New `queue_health` check in `gbrain doctor`. Postgres-only (PGLite has no multi-process worker surface). Two subchecks, both cheap (single SELECT each, status-index-covered): - stalled-forever: any active job with started_at > 1h. Surfaces the worst offenders (top 5 by started_at ASC) with `gbrain jobs get/cancel` fix hints. The incident that motivated v0.19.1 ran 90+ min before the operator noticed. - waiting-depth: per-name waiting count exceeds threshold. Default 10, overridable via GBRAIN_QUEUE_WAITING_THRESHOLD env (D9). Signals a submitter probably needs maxWaiting set. Worker-heartbeat subcheck from the original plan dropped (D4/H4): no minion_workers table exists, and lock_until-on-active-jobs is a lossy proxy that can't distinguish idle-worker from dead-worker. Tracked as follow-up B7. A4: --no-worker peer-liveness probe in autopilot. When --no-worker is set, every cycle runs a cheap SELECT checking for any active job whose lock_until was refreshed in the last 2 minutes. After 3 consecutive idle ticks, logs a loud WARNING naming the silent-wedge vector and referencing B7 as the ground-truth follow-up. Re-arms on next live signal so the warning doesn't spam every cycle. A6: New docs/guides/queue-operations-runbook.md (one viewport, ~60 lines). "My queue looks wedged — what do I run?" in order of escalation. What each doctor subcheck means. Self-check for the --no-worker / no-worker-running footgun. CLAUDE.md: key-files updates for handleWallClockTimeouts (v0.19.0 Layer 3 kill shot), maxWaiting advisory-lock rewrite (v0.19.1 D2), queue_health doctor check (v0.19.1 D4), and backpressure-audit.ts. Tests: all 143 minions + 13 doctor unit tests pass. No new test cases required in Lane B; the doctor queue_health exercise is in the E2E verification step (needs real PG to produce meaningful stalled-forever rows). The --no-worker probe is exercised by the smoke case's wedge setup in Lane A. README: unchanged. Existing `gbrain jobs submit` examples don't show --max-stalled, so no --max-waiting precedent to extend per A6 conditional. Co-Authored-By: Claude Opus 4.7 (1M context) * chore: v0.19.1 Lane C — CHANGELOG entry, VERSION bump, remove SPEC.md VERSION: 0.19.0 → 0.19.1 (patch; bug-fix-dominant, no schema change, no new user-facing vocabulary). CHANGELOG: new v0.19.1 entry at the top with the full release-summary template per CLAUDE.md — bold two-line headline, lead paragraph, "numbers that matter" before/after table measured against the real incident, "what this means for OpenClaw users" closer, required "To take advantage of v0.19.1" block naming the worker-restart requirement, itemized changes by area, and "For contributors" section closing the loop on the stale autopilot-idempotency narrative the CEO review was based on. Mechanism reframing per D1/H1: the 18-job pile-up was NOT caused by missing idempotency (autopilot already passes `idempotency_key: autopilot-cycle:${slot}` at autopilot.ts:241). The 18 jobs were 18 DIFFERENT slots stacking up behind the wedged one. `maxWaiting` still caps the pile; the incident just wasn't about idempotency. Adversarial review caught this before ship. SPEC.md: deleted from repo root. It was Wintermute's planning artifact for the original PR, not a shipped spec. Design docs belong under docs/designs/ per repo convention; leaving one at repo root set a precedent this repo doesn't want (A7/D11). CHANGELOG + the plan file at ~/.claude/plans/ok-wintermute-wrote-this-polished-matsumoto.md are the durable artifacts. Co-Authored-By: Claude Opus 4.7 (1M context) * fix: --wedge-rescue smoke state — both stall+timeout sweeps must skip Smoke case was setting lock_until in the past, so handleStalled's requeue path fired before handleWallClockTimeouts had a chance to evict. Production scenario is "lock_until still live (worker renewing) + timeout_at disqualified" — only wall-clock matches. Single-connection smoke can't simulate a row lock held by another txn, so we force the equivalent outcome: - lock_until = now() + 30s → handleStalled skips (not a stall) - timeout_at = NULL → handleTimeouts skips (needs NOT NULL) - started_at = now() - 10s, timeout_ms=1000 → wall-clock matches (2 × timeout_ms = 2000ms threshold exceeded) Verified: SMOKE PASS — Minions healthy + wedge rescue in 0.14s. Co-Authored-By: Claude Opus 4.7 (1M context) * fix: CI failures — shell-handler tests + llms-full.txt drift Two CI failure clusters, both pre-existing but surfaced by the v0.20.3 merge: 1) test/minions-shell.test.ts — 12 failing cases. The shell handler throws UnrecoverableError when GBRAIN_ALLOW_SHELL_JOBS !== '1' (the production RCE guard at shell.ts:210). The unit tests exercise handler mechanics, not the guard, but never set the env var — so every invocation exits through the guard path instead of the code being tested. Fix: set GBRAIN_ALLOW_SHELL_JOBS=1 in beforeAll, restore in afterAll. The env-guard IS still tested separately via the test/minions.test.ts case added in v0.20.3 Lane A which toggles the var itself. 2) llms-full.txt — stale against CLAUDE.md. Key-files entries for queue.ts, doctor.ts, and the new backpressure-audit.ts updated in v0.20.3 Lane B triggered the build-llms drift guard. Regenerated via `bun run build:llms`; no behavior change, just the inlined-docs bundle catching up to source. Full test run: 2367 pass, 0 fail across 137 files. Co-Authored-By: Claude Opus 4.7 (1M context) --------- Co-authored-by: root Co-authored-by: Claude Opus 4.7 (1M context) --- CHANGELOG.md | 95 +++++++++ CLAUDE.md | 5 +- VERSION | 2 +- docs/guides/queue-operations-runbook.md | 76 ++++++++ llms-full.txt | 5 +- package.json | 2 +- src/commands/autopilot.ts | 66 ++++++- src/commands/doctor.ts | 100 ++++++++++ src/commands/jobs.ts | 128 ++++++++++++- src/core/minions/backpressure-audit.ts | 77 ++++++++ src/core/minions/handlers/shell.ts | 10 + src/core/minions/queue.ts | 129 +++++++++++++ src/core/minions/types.ts | 2 + src/core/minions/worker.ts | 8 + test/minions-shell.test.ts | 12 ++ test/minions.test.ts | 243 ++++++++++++++++++++++++ 16 files changed, 940 insertions(+), 20 deletions(-) create mode 100644 docs/guides/queue-operations-runbook.md create mode 100644 src/core/minions/backpressure-audit.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index e9b4ff2e1..530761361 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,101 @@ All notable changes to GBrain will be documented in this file. +## [0.20.3] - 2026-04-24 + +## **Your queue now rescues itself when a wedged worker holds a row lock. Wall-clock sweep kills the job that stall detection can't see.** +## **`maxWaiting` is race-proof, observable, and reachable from the CLI — three bugs in one patch.** + +A production autopilot-cycle job wedged for over an hour on a single OpenClaw deployment because the worker's handler got stuck mid-transaction holding a row lock. Both eviction paths were blocked: the stall detector's `FOR UPDATE SKIP LOCKED` pass skipped the row-locked candidate, and the timeout sweep's `lock_until > now()` predicate disqualified the job once lock-renewal had been blocked. Neither could see the job. The shell-job pipeline starved completely behind the wedge. + +v0.19.0 shipped the wall-clock sweep as the third-layer kill shot: drop both constraints, evict on `started_at` alone, worst case at `2 × timeout_ms + stalledInterval`. This release locks down three correctness holes the v0.19.0 PR introduced — then closes the observability gap that let the incident run to minute 90 in the first place. + +### The queue-resilience numbers that matter + +Measured against the real incident on 2026-04-23 (OpenClaw autopilot + shell-job pipeline, Postgres engine, concurrency=1 worker). + +| Behavior | Before v0.20.3 | After v0.20.3 | +|---|---|---| +| Wedged worker escape window | 90+ minutes (manual kill) | `~2 × timeout_ms + 30s` sweep interval | +| Per-name waiting pile during wedge | 18 deferred per-slot jobs | capped at `maxWaiting` | +| `maxWaiting` under concurrent submit (2 submitters, cap=2) | up to 3 rows (TOCTOU race) | exactly 2 rows (advisory-lock serialization) | +| Same name across queues | cross-queue bleed — `shell` suppressed by `default` | isolated per `(name, queue)` | +| `GBRAIN_WORKER_CONCURRENCY=foo` | silent wedge (`inFlight < NaN` false) | clamped to 1, loud stderr warning | +| `gbrain jobs submit --max-waiting 2` | flag didn't exist | wired through to MinionJobInput | +| Silent coalesce events | invisible | JSONL audit at `~/.gbrain/audit/backpressure-YYYY-Www.jsonl` | +| `gbrain doctor` visibility into wedge | no check | new `queue_health` with 2 subchecks | + +The two big shifts: (1) every silent-failure vector the v0.19.0 patches introduced now has a loud signal — JSONL audit files, doctor check, stderr warnings, peer-liveness probe. (2) `maxWaiting` is now actually a cap under concurrency, not a soft suggestion. A future multi-submitter pattern (parallel workspaces, dispatched children, OpenClaw + ycli cron) doesn't walk through it. + +### What this means for OpenClaw users + +If you're running `gbrain autopilot` on a daily-driver deployment, the wall-clock sweep is the difference between a 90-minute outage and a 30-second one. The `queue_health` doctor check means the next time your queue wedges, you notice in minute 2 instead of minute 90. If you've been writing programmatic Minion submitters and setting `maxWaiting`, it's worth re-reading the JSONL audit file the next time your agent does anything "interesting" — you'll see exactly which submission coalesces into which returned job. + +## To take advantage of v0.20.3 + +`gbrain upgrade` handles the binary. You MUST restart long-running worker daemons so the new sweep runs in-process — the wall-clock eviction is a method on `MinionQueue`, not a cron job, so it only fires inside a worker loop. + +1. **Upgrade the binary:** + ```bash + gbrain upgrade + ``` +2. **Restart autopilot + workers:** + ```bash + # systemd / launchd / OpenClaw service-manager: restart the unit. + # Manual: kill the old `gbrain autopilot` and `gbrain jobs work`, start new ones. + ``` +3. **Verify:** + ```bash + gbrain jobs smoke --wedge-rescue # exercises the new wall-clock path + gbrain doctor --json | jq '.checks[] | select(.name == "queue_health")' + ``` +4. **If `gbrain doctor` flags anything unexpected,** please file an issue: + https://github.com/garrytan/gbrain/issues with: + - output of `gbrain doctor` + - contents of `~/.gbrain/audit/backpressure-*.jsonl` (redact freely) + - what commands you ran leading up to the wedge + +### Itemized changes + +**Queue core** (`src/core/minions/queue.ts`) +- `maxWaiting` coalesce path wraps `count → select → insert` in `pg_advisory_xact_lock` keyed on `(name, queue)`. Concurrent submitters for the SAME key serialize; different keys stay parallel. Lock auto-releases on transaction commit/rollback — no cleanup path to leak. Fixes TOCTOU race caught by adversarial review. +- `maxWaiting` count and select now filter on `queue` in addition to `name`. Pre-v0.20.3 code filtered on name alone, so a waiting `autopilot-cycle` in `queue=default` would suppress submissions to `queue=shell` with the same name. Cross-queue bleed is gone. + +**Backpressure observability** (new `src/core/minions/backpressure-audit.ts`) +- Every coalesce event writes one JSONL line to `~/.gbrain/audit/backpressure-YYYY-Www.jsonl` (ISO-week rotation, override dir via `GBRAIN_AUDIT_DIR`, mirrors the v0.14 shell-audit pattern). +- Fields: `ts, queue, name, waiting_count, max_waiting, decision='coalesced', returned_job_id`. +- Best-effort: write failures log to stderr but never block submission. + +**CLI** (`src/commands/jobs.ts`) +- New `--max-waiting N` flag on `gbrain jobs submit`. Clamps to `[1, 100]`, mirrors the existing `--max-stalled` wiring. The `MinionJobInput.maxWaiting` field was programmatic-only before; now it's reachable from the command line too. +- `resolveWorkerConcurrency` clamps against invalid input. `parseInt` returns `NaN` for `"foo"`, `0` for `"0"`, negatives for `"-5"` — all of which silently wedge a worker (`inFlight.size < NaN/0/negative` is always false). Now clamped to ≥1 with a loud stderr warning naming the bad value. One typo in a systemd unit no longer reproduces the 90-minute outage. +- New `gbrain jobs smoke --wedge-rescue` opt-in case. Forges a wedged-worker row state, invokes `handleStalled` + `handleTimeouts` + `handleWallClockTimeouts` in sequence, asserts only the wall-clock sweep evicts. Mirrors the v0.14.3 `--sigkill-rescue` shape. + +**Doctor** (`src/commands/doctor.ts`) +- New `queue_health` check (Postgres-only; PGLite skips with `Skipped (PGLite — no multi-process worker surface)`). +- Subcheck 1 — **stalled-forever**: flags active jobs whose `started_at` is older than 1 hour. Reports the top 5 by start time with `gbrain jobs get/cancel ` fix hints. +- Subcheck 2 — **waiting-depth**: flags per-name queues whose waiting count exceeds threshold. Default 10, overridable via `GBRAIN_QUEUE_WAITING_THRESHOLD` env. Reports the top 5 by depth with "consider setting maxWaiting on the submitter" fix hint. +- Worker-heartbeat staleness subcheck intentionally deferred to follow-up because `lock_until`-on-active-jobs is a lossy proxy. A check that cries wolf erodes trust in every other doctor subcheck. Needs a `minion_workers` table to produce ground-truth signal. + +**Autopilot** (`src/commands/autopilot.ts`) +- `--no-worker` mode gains a peer-worker-liveness probe. Every cycle runs a cheap `SELECT count(*)` checking for active jobs with `lock_until` refreshed in the last 2 minutes. After 3 consecutive idle ticks, logs a loud `WARNING` naming the silent-wedge vector (`--no-worker` set but no worker running). Re-arms once a live signal returns, so a healthy-but-idle worker doesn't trigger spam. +- Probe is documented as a proxy, not ground truth — idle worker with no active jobs reads as "no worker." The ground-truth fix needs a `minion_workers` heartbeat table (tracked as follow-up). + +**Docs** +- New `docs/guides/queue-operations-runbook.md`: the "my queue looks wedged — what do I run?" reference. One viewport, in order of escalation. What each `queue_health` subcheck means. Self-check for the `--no-worker + no-worker-running` footgun. +- `CLAUDE.md` Key-files section updated for the new `handleWallClockTimeouts` method (v0.19.0, described here for the first time), the new `backpressure-audit.ts` module, the updated `maxWaiting` semantics, and the new `queue_health` doctor check. + +**Tests** (`test/minions.test.ts`) +- 23 new unit cases. Wall-clock sweep (3 cases + non-interference with `handleTimeouts`). `maxWaiting` (coalesce, clamp 0 → 1, floor 1.7 → 1, concurrent-submitter race via `Promise.all`, cross-queue isolation, unset fallthrough). Concurrency clamp (7 cases including `NaN`/`0`/negative). `parseMaxWaitingFlag` (5 cases). Backpressure audit file write. All 143 minions tests pass. +- E2E wall-clock case against real Postgres is next on the roadmap (needs a second-connection row-lock helper; the unit-level coverage above exercises the sweep mechanics directly). + +### For contributors + +- The v0.19.0 PR's narrative framed the 18-job pileup as "duplicate submissions from a cron loop with no idempotency key." That framing was wrong. Autopilot already sets `idempotency_key: autopilot-cycle:${slot}` where slot is a 5-minute tick boundary — within-slot duplicates are structurally impossible. The 18 jobs were 18 different slots stacking up behind the wedged one. `maxWaiting` still caps the pile; the incident just wasn't about idempotency. Adversarial review caught this before v0.20.3 shipped. +- Follow-up issues tracked: B2 (autopilot heartbeat file), B3 (doctor `--fix` learns queue rescue), B4 (backpressure counts surfaced in `jobs stats`), B5 (cross-cutting "health-delivery-agent" pattern), B7 (`minion_workers` heartbeat table — unblocks both the dropped `queue_health` subcheck and a ground-truth `--no-worker` probe), P1 (composite indexes `(status, started_at)` and `(status, name)` on `minion_jobs` — currently the new sweeps fall back to `idx_minion_jobs_status`, selective enough on healthy queues, worth tightening in v0.20.4). + +Full plan with CEO + Eng + Codex adversarial decisions lives at `~/.claude/plans/` for the operators who care about how this release was reviewed. + ## [0.20.2] - 2026-04-24 ## **`gbrain jobs supervisor` is now a self-healing daemon you can actually drive. The Minions worker stops dying silently.** diff --git a/CLAUDE.md b/CLAUDE.md index 8e6bffe26..a392b013a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -62,12 +62,13 @@ strict behavior when unset. - `src/commands/graph-query.ts` — `gbrain graph-query [--type T] [--depth N] [--direction in|out|both]`: typed-edge relationship traversal (renders indented tree) - `src/core/link-extraction.ts` — shared library for the v0.12.0 graph layer. extractEntityRefs (canonical, replaces backlinks.ts duplicate) matches both `[Name](people/slug)` markdown links and Obsidian `[[people/slug|Name]]` wikilinks as of v0.12.3. extractPageLinks, inferLinkType heuristics (attended/works_at/invested_in/founded/advises/source/mentions), parseTimelineEntries, isAutoLinkEnabled config helper. `DIR_PATTERN` covers `people`, `companies`, `deals`, `topics`, `concepts`, `projects`, `entities`, `tech`, `finance`, `personal`, `openclaw`. Used by extract.ts, operations.ts auto-link post-hook, and backlinks.ts. - `src/core/minions/` — Minions job queue: BullMQ-inspired, Postgres-native (queue, worker, backoff, types, protected-names, quiet-hours, stagger, handlers/shell). -- `src/core/minions/queue.ts` — MinionQueue class (submit, claim, complete, fail, stall detection, parent-child, depth/child-cap, per-job timeouts, cascade-kill, attachments, idempotency keys, child_done inbox, removeOnComplete/Fail). `add()` takes a 4th `trusted` arg (separate from `opts` to prevent spread leakage); protected names in `PROTECTED_JOB_NAMES` require `{allowProtectedSubmit: true}` and the check runs trim-normalized (whitespace-bypass safe). v0.14.1 #219: `add()` plumbs `max_stalled` through with a `[1, 100]` clamp; omitted values let the schema DEFAULT (5) kick in. +- `src/core/minions/queue.ts` — MinionQueue class (submit, claim, complete, fail, stall detection, parent-child, depth/child-cap, per-job timeouts, cascade-kill, attachments, idempotency keys, child_done inbox, removeOnComplete/Fail). `add()` takes a 4th `trusted` arg (separate from `opts` to prevent spread leakage); protected names in `PROTECTED_JOB_NAMES` require `{allowProtectedSubmit: true}` and the check runs trim-normalized (whitespace-bypass safe). v0.14.1 #219: `add()` plumbs `max_stalled` through with a `[1, 100]` clamp; omitted values let the schema DEFAULT (5) kick in. v0.19.0: `handleWallClockTimeouts(lockDurationMs)` is Layer 3 kill shot for jobs where `FOR UPDATE SKIP LOCKED` stall detection and the timeout sweep both fail to evict (wedged worker holding a row lock via a pending transaction). v0.19.1: `maxWaiting` coalesce path now uses `pg_advisory_xact_lock` keyed on `(name, queue)` to serialize concurrent submits for the same key, and filters on `queue` in addition to `name` so cross-queue same-name jobs don't suppress each other. - `src/core/minions/worker.ts` — MinionWorker class (handler registry, lock renewal, graceful shutdown, timeout safety net). v0.14.0 abort-path fix: aborted jobs now call `failJob` with reason (`timeout`/`cancel`/`lock-lost`/`shutdown`) instead of returning silently. `shutdownAbort` (instance field) fires on process SIGTERM/SIGINT and propagates to `ctx.shutdownSignal` — shell handler listens to it; non-shell handlers don't. - `src/core/minions/types.ts` — `MinionJobInput` + `MinionJobStatus` + handler context types. `MinionJobInput.max_stalled` (new in v0.14.1) is optional; omitted values let the schema DEFAULT (5) kick in, provided values are clamped to `[1, 100]`. - `src/core/minions/protected-names.ts` — side-effect-free constant module exporting `PROTECTED_JOB_NAMES` + `isProtectedJobName()`. Kept pure so queue core can import without loading handler modules. - `src/core/minions/handlers/shell.ts` — `shell` job handler. Spawns `/bin/sh -c cmd` (absolute path, PATH-override-safe) or `argv[0] argv[1..]` (no shell). Env allowlist: `PATH, HOME, USER, LANG, TZ, NODE_ENV` + caller `env:` overrides. UTF-8-safe stdout/stderr tail via `string_decoder.StringDecoder`. Abort (either `ctx.signal` or `ctx.shutdownSignal`) fires SIGTERM → 5s grace → SIGKILL on child. Requires `GBRAIN_ALLOW_SHELL_JOBS=1` on worker (gated by `registerBuiltinHandlers`). - `src/core/minions/handlers/shell-audit.ts` — per-submission JSONL audit trail at `~/.gbrain/audit/shell-jobs-YYYY-Www.jsonl` (ISO-week rotation; override via `GBRAIN_AUDIT_DIR`). Best-effort: `mkdirSync(recursive)` + `appendFileSync`; failures logged to stderr, submission not blocked. Logs cmd (first 80 chars) or argv (JSON array). Never logs env values. +- `src/core/minions/backpressure-audit.ts` (v0.19.1) — sibling of shell-audit.ts for `maxWaiting` coalesce events. JSONL at `~/.gbrain/audit/backpressure-YYYY-Www.jsonl`. Fires one line per coalesce with `(queue, name, waiting_count, max_waiting, returned_job_id, ts)`. Closes the silent-drop vector the v0.19.0 maxWaiting guard introduced. - `src/core/minions/handlers/subagent.ts` (v0.15) — LLM-loop handler. Two-phase tool persistence (pending → complete/failed), replay reconciliation for mid-dispatch crashes, dual-signal abort (`ctx.signal` + `ctx.shutdownSignal`), Anthropic prompt caching on system + tool defs. `makeSubagentHandler({engine, client?, ...})` factory; `MessagesClient` is an injectable interface the real SDK implements structurally. Throws `RateLeaseUnavailableError` (renewable) when rate-lease capacity is full. - `src/core/minions/handlers/subagent-aggregator.ts` (v0.15) — `subagent_aggregator` handler. Claims AFTER all children resolve (queue changes guarantee every terminal child posts a `child_done` inbox message with outcome). Reads inbox via `ctx.readInbox()`, builds deterministic mixed-outcome markdown summary. No LLM call in v0.15. - `src/core/minions/handlers/subagent-audit.ts` (v0.15) — JSONL audit + heartbeat writer at `~/.gbrain/audit/subagent-jobs-YYYY-Www.jsonl`. Events: `submission` (one line per submit) + `heartbeat` (per turn boundary: `llm_call_started | llm_call_completed | tool_called | tool_result | tool_failed`). Never logs prompts or tool inputs. `readSubagentAuditForJob(jobId, {sinceIso})` is the readback path for `gbrain agent logs`. @@ -89,7 +90,7 @@ strict behavior when unset. - `src/commands/migrations/` — TS migration registry (compiled into the binary; no filesystem walk of `skills/migrations/*.md` needed at runtime). `index.ts` lists migrations in semver order. `v0_11_0.ts` = Minions adoption orchestrator (8 phases). `v0_12_0.ts` = Knowledge Graph auto-wire orchestrator (5 phases: schema → config check → backfill links → backfill timeline → verify). `phaseASchema` has a 600s timeout (bumped from 60s in v0.12.1 for duplicate-heavy brains). `v0_12_2.ts` = JSONB double-encode repair orchestrator (4 phases: schema → repair-jsonb → verify → record). `v0_14_0.ts` = shell-jobs + autopilot cooperative (2 phases: schema ALTER minion_jobs.max_stalled SET DEFAULT 3 — superseded by v0.14.3's schema-level DEFAULT 5 + UPDATE backfill; pending-host-work ping for skills/migrations/v0.14.0.md). All orchestrators are idempotent and resumable from `partial` status. As of v0.14.2 (Bug 3), the RUNNER owns all ledger writes — orchestrators return `OrchestratorResult` and `apply-migrations.ts` persists a canonical `{version, status, phases}` shape after return. Orchestrators no longer call `appendCompletedMigration` directly. `statusForVersion` prefers `complete` over `partial` (never regresses). 3 consecutive partials → wedged → `--force-retry ` writes a `'retry'` reset marker. v0.14.3 (fix wave) ships schema-only migrations v14 (`pages_updated_at_index`) + v15 (`minion_jobs_max_stalled_default_5` with UPDATE backfill) via the `MIGRATIONS` array in `src/core/migrate.ts` — no orchestrator phases needed. - `src/commands/repair-jsonb.ts` — `gbrain repair-jsonb [--dry-run] [--json]`: rewrites `jsonb_typeof='string'` rows in place across 5 affected columns (pages.frontmatter, raw_data.data, ingest_log.pages_updated, files.metadata, page_versions.frontmatter). Fixes v0.12.0 double-encode bug on Postgres; PGLite no-ops. Idempotent. - `src/commands/orphans.ts` — `gbrain orphans [--json] [--count] [--include-pseudo]`: surfaces pages with zero inbound wikilinks, grouped by domain. Auto-generated/raw/pseudo pages filtered by default. Also exposed as `find_orphans` MCP operation. Shipped in v0.12.3 (contributed by @knee5). -- `src/commands/doctor.ts` — `gbrain doctor [--json] [--fast] [--fix] [--dry-run] [--index-audit]`: health checks. v0.12.3 added `jsonb_integrity` + `markdown_body_completeness` reliability checks. v0.14.1: `--fix` delegates inlined cross-cutting rules to `> **Convention:** see [path](path).` callouts (pipes DRY violations into `src/core/dry-fix.ts`); `--fix --dry-run` previews without writing. v0.14.2: `schema_version` check fails loudly when `version=0` (migrations never ran — the #218 `bun install -g` signature) and routes users to `gbrain apply-migrations --yes`; new opt-in `--index-audit` flag (Postgres-only) reports zero-scan indexes from `pg_stat_user_indexes` (informational only, no auto-drop). v0.15.2: every DB check is wrapped in a progress phase; `markdown_body_completeness` runs under a 1s heartbeat timer so 10+ min scans are observable on 50K-page brains. Fix hints point at `gbrain repair-jsonb`, `gbrain sync --force`, and `gbrain apply-migrations`. +- `src/commands/doctor.ts` — `gbrain doctor [--json] [--fast] [--fix] [--dry-run] [--index-audit]`: health checks. v0.12.3 added `jsonb_integrity` + `markdown_body_completeness` reliability checks. v0.14.1: `--fix` delegates inlined cross-cutting rules to `> **Convention:** see [path](path).` callouts (pipes DRY violations into `src/core/dry-fix.ts`); `--fix --dry-run` previews without writing. v0.14.2: `schema_version` check fails loudly when `version=0` (migrations never ran — the #218 `bun install -g` signature) and routes users to `gbrain apply-migrations --yes`; new opt-in `--index-audit` flag (Postgres-only) reports zero-scan indexes from `pg_stat_user_indexes` (informational only, no auto-drop). v0.15.2: every DB check is wrapped in a progress phase; `markdown_body_completeness` runs under a 1s heartbeat timer so 10+ min scans are observable on 50K-page brains. v0.19.1 added `queue_health` (Postgres-only) with two subchecks: stalled-forever active jobs (started_at > 1h) and waiting-depth-per-name > threshold (default 10, override via `GBRAIN_QUEUE_WAITING_THRESHOLD`). Worker-heartbeat subcheck intentionally deferred to follow-up B7 because it needs a `minion_workers` table to produce ground-truth signal. Fix hints point at `gbrain repair-jsonb`, `gbrain sync --force`, `gbrain apply-migrations`, and `gbrain jobs get/cancel `. - `src/core/migrate.ts` — schema-migration runner. Owns the `MIGRATIONS` array (source of truth for schema DDL). v0.14.2 extended the `Migration` interface with `sqlFor?: { postgres?, pglite? }` (engine-specific SQL overrides `sql`) and `transaction?: boolean` (set to false for `CREATE INDEX CONCURRENTLY`, which Postgres refuses inside a transaction; ignored on PGLite since it has no concurrent writers). Migration v14 (fix wave) uses a handler branching on `engine.kind` to run CONCURRENTLY on Postgres (with a pre-drop of any invalid remnant via `pg_index.indisvalid`) and plain `CREATE INDEX` on PGLite. v15 bumps `minion_jobs.max_stalled` default 1→5 and backfills existing non-terminal rows. - `src/core/progress.ts` — Shared bulk-action progress reporter. Writes to stderr. Modes: `auto` (TTY: `\r`-rewriting; non-TTY: plain lines), `human`, `json` (JSONL), `quiet`. Rate-gated by `minIntervalMs` and `minItems`. `startHeartbeat(reporter, note)` helper for single long queries. `child()` composes phase paths. Singleton SIGINT/SIGTERM coordinator emits `abort` events for every live phase. EPIPE defense on both sync throws and stream `'error'` events. Zero dependencies. Introduced in v0.15.2. - `src/core/cli-options.ts` — Global CLI flag parser. `parseGlobalFlags(argv)` returns `{cliOpts, rest}` with `--quiet` / `--progress-json` / `--progress-interval=` stripped. `getCliOptions()` / `setCliOptions()` expose a module-level singleton so commands reach the resolved flags without parameter threading. `cliOptsToProgressOptions()` maps to reporter options. `childGlobalFlags()` returns the flag suffix to append to `execSync('gbrain ...')` calls in migration orchestrators. `OperationContext.cliOpts` extends shared-op dispatch for MCP callers. diff --git a/VERSION b/VERSION index 727d97b9b..144996ed2 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.20.2 +0.20.3 diff --git a/docs/guides/queue-operations-runbook.md b/docs/guides/queue-operations-runbook.md new file mode 100644 index 000000000..aef91111e --- /dev/null +++ b/docs/guides/queue-operations-runbook.md @@ -0,0 +1,76 @@ +# Queue operations runbook + +"My queue looks wedged — what do I run?" The commands below are in the order +you probably want them. Shipped with v0.19.1 after a production incident +where the queue held for 90+ minutes before the operator noticed. + +## First signal: jobs aren't running + +```bash +gbrain doctor --json | jq '.checks[] | select(.name == "queue_health")' +``` + +`queue_health` flags two patterns: + +- **stalled-forever**: active job whose `started_at` is older than 1h. +- **waiting-depth**: any per-name queue deeper than 10 (override via + `GBRAIN_QUEUE_WAITING_THRESHOLD`). Signals a missing `maxWaiting`. + +## Triage commands + +```bash +# Who's active right now? +gbrain jobs list --status active + +# Who's waiting, biggest pile first? +gbrain jobs list --status waiting --limit 50 + +# What's wrong with a specific job? +gbrain jobs get +``` + +## Rescue actions (in order of escalation) + +```bash +# Force-kill a single stuck job: +gbrain jobs cancel + +# Clear a specific job entirely (last resort): +gbrain jobs delete + +# Health smoke on the mechanism itself: +gbrain jobs smoke --wedge-rescue +``` + +## What each subcheck means + +- **stalled-forever** — A worker claimed a job, started executing, and has + held the row for over an hour. The wall-clock sweep evicts jobs past + 2× `timeout_ms`; if one's still active, either no `timeout_ms` was set + or the sweep is newly deployed and this job predates it. Cancel it. +- **waiting-depth** — Submitters are piling up jobs faster than workers + drain them. Set `--max-waiting N` on the submission or on the programmatic + `queue.add()` call. If you want a taller pile, raise the threshold via + `GBRAIN_QUEUE_WAITING_THRESHOLD=50 gbrain doctor`. + +## Self-check: is a worker even running? + +```bash +# If you're running autopilot with --no-worker, check that your external +# worker (systemd / Docker / OpenClaw service-manager) is alive: +gbrain jobs list --status active | head -5 +``` + +If the list is empty AND your submissions keep piling up, no worker is +claiming. Start one: + +```bash +GBRAIN_ALLOW_SHELL_JOBS=1 gbrain jobs work --concurrency 4 +``` + +## Follow-ups tracked for v0.20+ + +- B7 — `minion_workers` heartbeat table for ground-truth liveness (the + `--no-worker` probe and the dropped `queue_health` worker-heartbeat + subcheck both need this). +- B3 — `gbrain doctor --fix` learns to rescue queue wedges. diff --git a/llms-full.txt b/llms-full.txt index 5292c5916..b8da9d6d5 100644 --- a/llms-full.txt +++ b/llms-full.txt @@ -141,12 +141,13 @@ strict behavior when unset. - `src/commands/graph-query.ts` — `gbrain graph-query [--type T] [--depth N] [--direction in|out|both]`: typed-edge relationship traversal (renders indented tree) - `src/core/link-extraction.ts` — shared library for the v0.12.0 graph layer. extractEntityRefs (canonical, replaces backlinks.ts duplicate) matches both `[Name](people/slug)` markdown links and Obsidian `[[people/slug|Name]]` wikilinks as of v0.12.3. extractPageLinks, inferLinkType heuristics (attended/works_at/invested_in/founded/advises/source/mentions), parseTimelineEntries, isAutoLinkEnabled config helper. `DIR_PATTERN` covers `people`, `companies`, `deals`, `topics`, `concepts`, `projects`, `entities`, `tech`, `finance`, `personal`, `openclaw`. Used by extract.ts, operations.ts auto-link post-hook, and backlinks.ts. - `src/core/minions/` — Minions job queue: BullMQ-inspired, Postgres-native (queue, worker, backoff, types, protected-names, quiet-hours, stagger, handlers/shell). -- `src/core/minions/queue.ts` — MinionQueue class (submit, claim, complete, fail, stall detection, parent-child, depth/child-cap, per-job timeouts, cascade-kill, attachments, idempotency keys, child_done inbox, removeOnComplete/Fail). `add()` takes a 4th `trusted` arg (separate from `opts` to prevent spread leakage); protected names in `PROTECTED_JOB_NAMES` require `{allowProtectedSubmit: true}` and the check runs trim-normalized (whitespace-bypass safe). v0.14.1 #219: `add()` plumbs `max_stalled` through with a `[1, 100]` clamp; omitted values let the schema DEFAULT (5) kick in. +- `src/core/minions/queue.ts` — MinionQueue class (submit, claim, complete, fail, stall detection, parent-child, depth/child-cap, per-job timeouts, cascade-kill, attachments, idempotency keys, child_done inbox, removeOnComplete/Fail). `add()` takes a 4th `trusted` arg (separate from `opts` to prevent spread leakage); protected names in `PROTECTED_JOB_NAMES` require `{allowProtectedSubmit: true}` and the check runs trim-normalized (whitespace-bypass safe). v0.14.1 #219: `add()` plumbs `max_stalled` through with a `[1, 100]` clamp; omitted values let the schema DEFAULT (5) kick in. v0.19.0: `handleWallClockTimeouts(lockDurationMs)` is Layer 3 kill shot for jobs where `FOR UPDATE SKIP LOCKED` stall detection and the timeout sweep both fail to evict (wedged worker holding a row lock via a pending transaction). v0.19.1: `maxWaiting` coalesce path now uses `pg_advisory_xact_lock` keyed on `(name, queue)` to serialize concurrent submits for the same key, and filters on `queue` in addition to `name` so cross-queue same-name jobs don't suppress each other. - `src/core/minions/worker.ts` — MinionWorker class (handler registry, lock renewal, graceful shutdown, timeout safety net). v0.14.0 abort-path fix: aborted jobs now call `failJob` with reason (`timeout`/`cancel`/`lock-lost`/`shutdown`) instead of returning silently. `shutdownAbort` (instance field) fires on process SIGTERM/SIGINT and propagates to `ctx.shutdownSignal` — shell handler listens to it; non-shell handlers don't. - `src/core/minions/types.ts` — `MinionJobInput` + `MinionJobStatus` + handler context types. `MinionJobInput.max_stalled` (new in v0.14.1) is optional; omitted values let the schema DEFAULT (5) kick in, provided values are clamped to `[1, 100]`. - `src/core/minions/protected-names.ts` — side-effect-free constant module exporting `PROTECTED_JOB_NAMES` + `isProtectedJobName()`. Kept pure so queue core can import without loading handler modules. - `src/core/minions/handlers/shell.ts` — `shell` job handler. Spawns `/bin/sh -c cmd` (absolute path, PATH-override-safe) or `argv[0] argv[1..]` (no shell). Env allowlist: `PATH, HOME, USER, LANG, TZ, NODE_ENV` + caller `env:` overrides. UTF-8-safe stdout/stderr tail via `string_decoder.StringDecoder`. Abort (either `ctx.signal` or `ctx.shutdownSignal`) fires SIGTERM → 5s grace → SIGKILL on child. Requires `GBRAIN_ALLOW_SHELL_JOBS=1` on worker (gated by `registerBuiltinHandlers`). - `src/core/minions/handlers/shell-audit.ts` — per-submission JSONL audit trail at `~/.gbrain/audit/shell-jobs-YYYY-Www.jsonl` (ISO-week rotation; override via `GBRAIN_AUDIT_DIR`). Best-effort: `mkdirSync(recursive)` + `appendFileSync`; failures logged to stderr, submission not blocked. Logs cmd (first 80 chars) or argv (JSON array). Never logs env values. +- `src/core/minions/backpressure-audit.ts` (v0.19.1) — sibling of shell-audit.ts for `maxWaiting` coalesce events. JSONL at `~/.gbrain/audit/backpressure-YYYY-Www.jsonl`. Fires one line per coalesce with `(queue, name, waiting_count, max_waiting, returned_job_id, ts)`. Closes the silent-drop vector the v0.19.0 maxWaiting guard introduced. - `src/core/minions/handlers/subagent.ts` (v0.15) — LLM-loop handler. Two-phase tool persistence (pending → complete/failed), replay reconciliation for mid-dispatch crashes, dual-signal abort (`ctx.signal` + `ctx.shutdownSignal`), Anthropic prompt caching on system + tool defs. `makeSubagentHandler({engine, client?, ...})` factory; `MessagesClient` is an injectable interface the real SDK implements structurally. Throws `RateLeaseUnavailableError` (renewable) when rate-lease capacity is full. - `src/core/minions/handlers/subagent-aggregator.ts` (v0.15) — `subagent_aggregator` handler. Claims AFTER all children resolve (queue changes guarantee every terminal child posts a `child_done` inbox message with outcome). Reads inbox via `ctx.readInbox()`, builds deterministic mixed-outcome markdown summary. No LLM call in v0.15. - `src/core/minions/handlers/subagent-audit.ts` (v0.15) — JSONL audit + heartbeat writer at `~/.gbrain/audit/subagent-jobs-YYYY-Www.jsonl`. Events: `submission` (one line per submit) + `heartbeat` (per turn boundary: `llm_call_started | llm_call_completed | tool_called | tool_result | tool_failed`). Never logs prompts or tool inputs. `readSubagentAuditForJob(jobId, {sinceIso})` is the readback path for `gbrain agent logs`. @@ -168,7 +169,7 @@ strict behavior when unset. - `src/commands/migrations/` — TS migration registry (compiled into the binary; no filesystem walk of `skills/migrations/*.md` needed at runtime). `index.ts` lists migrations in semver order. `v0_11_0.ts` = Minions adoption orchestrator (8 phases). `v0_12_0.ts` = Knowledge Graph auto-wire orchestrator (5 phases: schema → config check → backfill links → backfill timeline → verify). `phaseASchema` has a 600s timeout (bumped from 60s in v0.12.1 for duplicate-heavy brains). `v0_12_2.ts` = JSONB double-encode repair orchestrator (4 phases: schema → repair-jsonb → verify → record). `v0_14_0.ts` = shell-jobs + autopilot cooperative (2 phases: schema ALTER minion_jobs.max_stalled SET DEFAULT 3 — superseded by v0.14.3's schema-level DEFAULT 5 + UPDATE backfill; pending-host-work ping for skills/migrations/v0.14.0.md). All orchestrators are idempotent and resumable from `partial` status. As of v0.14.2 (Bug 3), the RUNNER owns all ledger writes — orchestrators return `OrchestratorResult` and `apply-migrations.ts` persists a canonical `{version, status, phases}` shape after return. Orchestrators no longer call `appendCompletedMigration` directly. `statusForVersion` prefers `complete` over `partial` (never regresses). 3 consecutive partials → wedged → `--force-retry ` writes a `'retry'` reset marker. v0.14.3 (fix wave) ships schema-only migrations v14 (`pages_updated_at_index`) + v15 (`minion_jobs_max_stalled_default_5` with UPDATE backfill) via the `MIGRATIONS` array in `src/core/migrate.ts` — no orchestrator phases needed. - `src/commands/repair-jsonb.ts` — `gbrain repair-jsonb [--dry-run] [--json]`: rewrites `jsonb_typeof='string'` rows in place across 5 affected columns (pages.frontmatter, raw_data.data, ingest_log.pages_updated, files.metadata, page_versions.frontmatter). Fixes v0.12.0 double-encode bug on Postgres; PGLite no-ops. Idempotent. - `src/commands/orphans.ts` — `gbrain orphans [--json] [--count] [--include-pseudo]`: surfaces pages with zero inbound wikilinks, grouped by domain. Auto-generated/raw/pseudo pages filtered by default. Also exposed as `find_orphans` MCP operation. Shipped in v0.12.3 (contributed by @knee5). -- `src/commands/doctor.ts` — `gbrain doctor [--json] [--fast] [--fix] [--dry-run] [--index-audit]`: health checks. v0.12.3 added `jsonb_integrity` + `markdown_body_completeness` reliability checks. v0.14.1: `--fix` delegates inlined cross-cutting rules to `> **Convention:** see [path](path).` callouts (pipes DRY violations into `src/core/dry-fix.ts`); `--fix --dry-run` previews without writing. v0.14.2: `schema_version` check fails loudly when `version=0` (migrations never ran — the #218 `bun install -g` signature) and routes users to `gbrain apply-migrations --yes`; new opt-in `--index-audit` flag (Postgres-only) reports zero-scan indexes from `pg_stat_user_indexes` (informational only, no auto-drop). v0.15.2: every DB check is wrapped in a progress phase; `markdown_body_completeness` runs under a 1s heartbeat timer so 10+ min scans are observable on 50K-page brains. Fix hints point at `gbrain repair-jsonb`, `gbrain sync --force`, and `gbrain apply-migrations`. +- `src/commands/doctor.ts` — `gbrain doctor [--json] [--fast] [--fix] [--dry-run] [--index-audit]`: health checks. v0.12.3 added `jsonb_integrity` + `markdown_body_completeness` reliability checks. v0.14.1: `--fix` delegates inlined cross-cutting rules to `> **Convention:** see [path](path).` callouts (pipes DRY violations into `src/core/dry-fix.ts`); `--fix --dry-run` previews without writing. v0.14.2: `schema_version` check fails loudly when `version=0` (migrations never ran — the #218 `bun install -g` signature) and routes users to `gbrain apply-migrations --yes`; new opt-in `--index-audit` flag (Postgres-only) reports zero-scan indexes from `pg_stat_user_indexes` (informational only, no auto-drop). v0.15.2: every DB check is wrapped in a progress phase; `markdown_body_completeness` runs under a 1s heartbeat timer so 10+ min scans are observable on 50K-page brains. v0.19.1 added `queue_health` (Postgres-only) with two subchecks: stalled-forever active jobs (started_at > 1h) and waiting-depth-per-name > threshold (default 10, override via `GBRAIN_QUEUE_WAITING_THRESHOLD`). Worker-heartbeat subcheck intentionally deferred to follow-up B7 because it needs a `minion_workers` table to produce ground-truth signal. Fix hints point at `gbrain repair-jsonb`, `gbrain sync --force`, `gbrain apply-migrations`, and `gbrain jobs get/cancel `. - `src/core/migrate.ts` — schema-migration runner. Owns the `MIGRATIONS` array (source of truth for schema DDL). v0.14.2 extended the `Migration` interface with `sqlFor?: { postgres?, pglite? }` (engine-specific SQL overrides `sql`) and `transaction?: boolean` (set to false for `CREATE INDEX CONCURRENTLY`, which Postgres refuses inside a transaction; ignored on PGLite since it has no concurrent writers). Migration v14 (fix wave) uses a handler branching on `engine.kind` to run CONCURRENTLY on Postgres (with a pre-drop of any invalid remnant via `pg_index.indisvalid`) and plain `CREATE INDEX` on PGLite. v15 bumps `minion_jobs.max_stalled` default 1→5 and backfills existing non-terminal rows. - `src/core/progress.ts` — Shared bulk-action progress reporter. Writes to stderr. Modes: `auto` (TTY: `\r`-rewriting; non-TTY: plain lines), `human`, `json` (JSONL), `quiet`. Rate-gated by `minIntervalMs` and `minItems`. `startHeartbeat(reporter, note)` helper for single long queries. `child()` composes phase paths. Singleton SIGINT/SIGTERM coordinator emits `abort` events for every live phase. EPIPE defense on both sync throws and stream `'error'` events. Zero dependencies. Introduced in v0.15.2. - `src/core/cli-options.ts` — Global CLI flag parser. `parseGlobalFlags(argv)` returns `{cliOpts, rest}` with `--quiet` / `--progress-json` / `--progress-interval=` stripped. `getCliOptions()` / `setCliOptions()` expose a module-level singleton so commands reach the resolved flags without parameter threading. `cliOptsToProgressOptions()` maps to reporter options. `childGlobalFlags()` returns the flag suffix to append to `execSync('gbrain ...')` calls in migration orchestrators. `OperationContext.cliOpts` extends shared-op dispatch for MCP callers. diff --git a/package.json b/package.json index a34415c76..9a8a144ed 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "gbrain", - "version": "0.20.2", + "version": "0.20.3", "description": "Postgres-native personal knowledge brain with hybrid RAG search", "type": "module", "main": "src/core/index.ts", diff --git a/src/commands/autopilot.ts b/src/commands/autopilot.ts index 2188d64f5..3663e4f5d 100644 --- a/src/commands/autopilot.ts +++ b/src/commands/autopilot.ts @@ -75,10 +75,14 @@ export function resolveGbrainCliPath(): string { throw new Error('Could not resolve the gbrain CLI path. Install gbrain so it is on $PATH (e.g. /usr/local/bin/gbrain), or run autopilot from the compiled binary directly.'); } +export function shouldSpawnAutopilotWorker(args: string[]): boolean { + return !args.includes('--no-worker'); +} + export async function runAutopilot(engine: BrainEngine, args: string[]) { if (args.includes('--help') || args.includes('-h')) { console.log( - 'Usage: gbrain autopilot [--repo ] [--interval N] [--json]\n' + + 'Usage: gbrain autopilot [--repo ] [--interval N] [--json] [--no-worker]\n' + ' gbrain autopilot --install [--repo ]\n' + ' gbrain autopilot --uninstall\n' + ' gbrain autopilot --status [--json]\n\n' + @@ -106,6 +110,7 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { const baseInterval = parseInt(parseArg(args, '--interval') || '300', 10); const jsonMode = args.includes('--json'); const forceInline = args.includes('--inline'); + const noWorker = !shouldSpawnAutopilotWorker(args); if (!repoPath) { console.error('No repo path. Use --repo or run gbrain sync --repo first.'); @@ -137,12 +142,13 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { const cfg = loadConfig(); const engineType = cfg?.engine ?? 'pglite'; const useMinionsDispatch = mode !== 'off' && engineType === 'postgres' && !forceInline; + const spawnManagedWorker = useMinionsDispatch && !noWorker; let stopping = false; let workerProc: ChildProcess | null = null; let crashCount = 0; - if (useMinionsDispatch) { + if (spawnManagedWorker) { const cliPath = resolveGbrainCliPath(); const startWorker = () => { const child = spawn(cliPath, ['jobs', 'work'], { stdio: 'inherit', env: process.env }); @@ -161,10 +167,13 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { }); }; startWorker(); - } else { - const why = mode === 'off' ? 'minion_mode=off' + } else if (!useMinionsDispatch) { + const why = mode === 'off' + ? 'minion_mode=off' : (engineType !== 'postgres' ? 'engine=pglite' : 'flag=--inline'); console.log(`[autopilot] running steps inline (${why})`); + } else { + console.log('[autopilot] --no-worker set: dispatch loop only (worker managed externally)'); } // Async shutdown with 35s drain window for the worker child. The worker @@ -195,6 +204,18 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { process.on('SIGINT', () => { void shutdown('SIGINT'); }); let consecutiveErrors = 0; + // Peer-worker liveness for --no-worker mode. The probe is a proxy, not + // ground truth: SELECT count(*) of active jobs with a recent lock_until + // refresh. A queue with only waiting jobs and a healthy idle worker + // reads as "no worker" (false positive); a worker that died 110s ago + // while holding a lock reads as "alive" until lock_until expires. + // Good enough for V1 — a ground-truth minion_workers heartbeat table + // is tracked as v0.19.1 follow-up B7. When the probe sees no signal + // for NO_WORKER_WARN_TICKS consecutive cycles, log a loud warning so + // the operator can spot "I set --no-worker but forgot to start one" + // before the queue piles up. + const NO_WORKER_WARN_TICKS = 3; + let noWorkerConsecutiveIdle = 0; while (!stopping) { const cycleStart = Date.now(); @@ -214,6 +235,43 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) { } catch (e) { logError('reconnect', e); } } + // --no-worker peer-liveness probe (v0.19.1). Runs every cycle, cheap + // (single SELECT). See NO_WORKER_WARN_TICKS comment above for caveats. + if (noWorker && useMinionsDispatch) { + try { + const rows = await (engine as any).executeRaw?.( + `SELECT count(*)::int AS n FROM minion_jobs + WHERE status = 'active' + AND lock_until IS NOT NULL + AND lock_until > now() - interval '2 minutes'`, + ); + const liveWorkerSignal = Number((rows as Array<{ n: number }>)?.[0]?.n ?? 0); + if (liveWorkerSignal === 0) { + noWorkerConsecutiveIdle++; + if (noWorkerConsecutiveIdle === NO_WORKER_WARN_TICKS) { + // Fire loud on the Nth consecutive idle tick; don't repeat on every + // subsequent cycle (the operator already saw it), re-arm once a + // live worker is seen again. + console.error( + `[autopilot] WARNING: --no-worker set and no worker has claimed a job in ~${NO_WORKER_WARN_TICKS * baseInterval}s. ` + + `Jobs will pile up in 'waiting' until a worker starts. ` + + `Probe is a proxy (lock_until refresh) and can false-positive on idle queues — see B7 for ground-truth follow-up.`, + ); + } + } else { + if (noWorkerConsecutiveIdle >= NO_WORKER_WARN_TICKS) { + console.log('[autopilot] --no-worker probe: live worker signal detected; warning re-armed.'); + } + noWorkerConsecutiveIdle = 0; + } + } catch (e) { + // Probe failures never block the main dispatch loop. Log once per + // failure class; ignore repeated errors (common shape: DB reconnect + // blip between ticks). + logError('no-worker-probe', e); + } + } + if (useMinionsDispatch) { // Submit ONE autopilot-cycle job per cycle slot. The idempotency key // dedupes overrun submissions — if a cycle's job runs longer than diff --git a/src/commands/doctor.ts b/src/commands/doctor.ts index f8f53849f..dd16bd24b 100644 --- a/src/commands/doctor.ts +++ b/src/commands/doctor.ts @@ -649,6 +649,106 @@ export async function runDoctor(engine: BrainEngine | null, args: string[], dbSo mbcHb(); } + // 11b. Queue health (v0.19.1 queue-resilience wave). + // Postgres-only because PGLite has no multi-process worker surface. Two + // subchecks, both cheap (single SELECT each, status-index-covered): + // + // 1. stalled-forever: any active job whose started_at is > 1h old. The + // incident that motivated this release ran 90+ min before surfacing. + // Surface the ID so the operator can `gbrain jobs get ` to inspect + // or `gbrain jobs cancel ` to force-kill. + // + // 2. backpressure-missed: per-name waiting depth exceeds the threshold + // (default 10, override via GBRAIN_QUEUE_WAITING_THRESHOLD env). Signal + // that a submitter probably needs maxWaiting set. Bounded by per-name + // aggregation so a single name's pile shows up clearly instead of + // getting lost in the total. + // + // Not included in v0.19.1 (tracked as B7 follow-up): worker-heartbeat + // staleness. It needs a minion_workers table; the lock_until-on-active-jobs + // proxy can't distinguish "no worker" from "worker idle," and a check that + // cries wolf erodes trust in every other doctor check. + progress.heartbeat('queue_health'); + if (engine.kind === 'pglite') { + checks.push({ + name: 'queue_health', + status: 'ok', + message: 'Skipped (PGLite — no multi-process worker surface)', + }); + } else { + const queueHealthHb = startHeartbeat(progress, 'scanning queue health…'); + try { + const sql = db.getConnection(); + // Subcheck 1: stalled-forever active jobs (>1h wall-clock). + const stalledRows: Array<{ id: number; name: string; started_at: string }> = await sql` + SELECT id, name, started_at::text AS started_at + FROM minion_jobs + WHERE status = 'active' + AND started_at IS NOT NULL + AND started_at < now() - interval '1 hour' + ORDER BY started_at ASC + LIMIT 5 + `; + // Subcheck 2: per-name waiting depth exceeds threshold. + const rawThreshold = process.env.GBRAIN_QUEUE_WAITING_THRESHOLD; + const parsedThreshold = rawThreshold ? parseInt(rawThreshold, 10) : 10; + const threshold = Number.isFinite(parsedThreshold) && parsedThreshold >= 1 + ? parsedThreshold + : 10; + const depthRows: Array<{ name: string; queue: string; depth: number }> = await sql` + SELECT name, queue, count(*)::int AS depth + FROM minion_jobs + WHERE status = 'waiting' + GROUP BY name, queue + HAVING count(*) > ${threshold} + ORDER BY depth DESC + LIMIT 5 + `; + + const problems: string[] = []; + if (stalledRows.length > 0) { + const sample = stalledRows + .map(r => `#${r.id}(${r.name})`) + .join(', '); + problems.push( + `${stalledRows.length} stalled-forever job(s): ${sample}. ` + + `Fix: gbrain jobs get to inspect; gbrain jobs cancel to force-kill.` + ); + } + if (depthRows.length > 0) { + const sample = depthRows + .map(r => `${r.name}@${r.queue}=${r.depth}`) + .join(', '); + problems.push( + `waiting-queue depth exceeds ${threshold} for: ${sample}. ` + + `Fix: set maxWaiting on the submitter (or raise GBRAIN_QUEUE_WAITING_THRESHOLD).` + ); + } + + if (problems.length === 0) { + checks.push({ + name: 'queue_health', + status: 'ok', + message: `No stalled-forever jobs; no queue over depth ${threshold}.`, + }); + } else { + checks.push({ + name: 'queue_health', + status: 'warn', + message: problems.join(' '), + }); + } + } catch (e) { + checks.push({ + name: 'queue_health', + status: 'warn', + message: `queue_health scan skipped: ${e instanceof Error ? e.message : String(e)}`, + }); + } finally { + queueHealthHb(); + } + } + // 12. Index audit (opt-in via --index-audit). v0.13.1 follow-up to #170. // Reports indexes with zero recorded scans on Postgres. Informational only; // we DO NOT auto-drop. On #170's brain, idx_pages_frontmatter and diff --git a/src/commands/jobs.ts b/src/commands/jobs.ts index d8c87d0f3..f82c3b300 100644 --- a/src/commands/jobs.ts +++ b/src/commands/jobs.ts @@ -17,6 +17,42 @@ function hasFlag(args: string[], flag: string): boolean { return args.includes(flag); } +/** Parse `--max-waiting N` from CLI args. Returns undefined if absent. + * Throws on malformed input (caller should surface the error and exit). + * Clamps to [1, 100] to match the queue-layer clamp in MinionQueue.add. + * Exported for unit tests; the CLI handler at `jobs submit` wraps this + * with process.exit(1) on throw so operators see 'must be positive integer'. */ +export function parseMaxWaitingFlag(args: string[]): number | undefined { + const raw = parseFlag(args, '--max-waiting'); + if (raw === undefined) return undefined; + const parsed = parseInt(raw, 10); + if (!Number.isFinite(parsed) || parsed < 1) { + throw new Error('--max-waiting must be a positive integer (will be clamped to [1, 100])'); + } + return Math.max(1, Math.min(100, parsed)); +} + +export function resolveWorkerConcurrency(args: string[], env: NodeJS.ProcessEnv = process.env): number { + const raw = parseFlag(args, '--concurrency') ?? env.GBRAIN_WORKER_CONCURRENCY ?? '1'; + const parsed = parseInt(raw, 10); + // Without validation, NaN / 0 / negative values flow through to the worker + // loop where `inFlight.size < concurrency` is always false → the worker + // claims zero jobs and the queue silently wedges. One typo in a systemd + // unit reproduces the original production incident. Clamp to ≥1 and surface + // the misconfig loudly so operators see it at worker startup. + if (!Number.isFinite(parsed) || parsed < 1) { + const source = parseFlag(args, '--concurrency') !== undefined + ? '--concurrency flag' + : 'GBRAIN_WORKER_CONCURRENCY env'; + process.stderr.write( + `[gbrain jobs] invalid concurrency from ${source} (${JSON.stringify(raw)}); ` + + `falling back to 1. Set a positive integer.\n` + ); + return 1; + } + return parsed; +} + function formatJob(job: MinionJob): string { const dur = job.finished_at && job.started_at ? `${((job.finished_at.getTime() - job.started_at.getTime()) / 1000).toFixed(1)}s` @@ -58,6 +94,7 @@ export async function runJobs(engine: BrainEngine, args: string[]): Promise [--params JSON] [--follow] [--priority N] [--delay Nms] [--max-attempts N] [--max-stalled N] + [--max-waiting N] [--backoff-type fixed|exponential] [--backoff-delay Nms] [--backoff-jitter 0..1] [--timeout-ms Nms] [--idempotency-key K] [--queue Q] [--dry-run] @@ -144,6 +181,12 @@ HANDLER TYPES (built in) const maxAttempts = parseInt(parseFlag(args, '--max-attempts') ?? '3', 10); const maxStalledRaw = parseFlag(args, '--max-stalled'); const maxStalled = maxStalledRaw !== undefined ? parseInt(maxStalledRaw, 10) : undefined; + // --max-waiting N: submission-time backpressure cap. Mirrors --max-stalled + // clamp [1, 100]. Feature is usable from CLI as of v0.19.1; pre-v0.19.1 + // only programmatic callers reached it. + let maxWaiting: number | undefined; + try { maxWaiting = parseMaxWaitingFlag(args); } + catch (e) { console.error(`Error: ${e instanceof Error ? e.message : String(e)}`); process.exit(1); } // v0.13.1 field audit: expose retry/backoff/timeout/idempotency knobs so // users can tune Minions behavior without dropping into TypeScript. const backoffTypeRaw = parseFlag(args, '--backoff-type'); @@ -172,6 +215,7 @@ HANDLER TYPES (built in) console.log(` Priority: ${priority}`); console.log(` Max attempts: ${maxAttempts}`); if (maxStalled !== undefined) console.log(` Max stalled: ${maxStalled}`); + if (maxWaiting !== undefined) console.log(` Max waiting: ${maxWaiting}`); if (backoffType) console.log(` Backoff type: ${backoffType}`); if (backoffDelay !== undefined) console.log(` Backoff delay: ${backoffDelay}ms`); if (backoffJitter !== undefined) console.log(` Backoff jitter: ${backoffJitter}`); @@ -199,6 +243,7 @@ HANDLER TYPES (built in) delay: delay > 0 ? delay : undefined, max_attempts: maxAttempts, max_stalled: maxStalled, + maxWaiting, backoff_type: backoffType, backoff_delay: backoffDelay, backoff_jitter: backoffJitter, @@ -415,6 +460,7 @@ HANDLER TYPES (built in) } const sigkillRescue = hasFlag(args, '--sigkill-rescue'); + const wedgeRescue = hasFlag(args, '--wedge-rescue'); const worker = new MinionWorker(engine, { queue: 'smoke', pollInterval: 100 }); worker.register('noop', async () => ({ ok: true, at: new Date().toISOString() })); @@ -481,9 +527,70 @@ HANDLER TYPES (built in) try { await queue.removeJob(rescueJob.id); } catch { /* non-fatal cleanup */ } } + // --wedge-rescue: regression case for the v0.19.1 production incident. + // In prod, a wedged worker held a row lock via a pending txn. The + // lock-renewal UPDATE blocked, lock_until fell below now(), handleStalled + // saw the candidate but FOR UPDATE SKIP LOCKED skipped (row lock held), + // handleTimeouts was disqualified (lock_until > now() fails). + // Only handleWallClockTimeouts' no-constraint sweep evicted. + // + // The smoke is single-connection, so we can't simulate a row lock held + // by another txn. Instead we forge the state where BOTH handleStalled + // and handleTimeouts are disqualified so only wall-clock fires: + // - lock_until far in the future → handleStalled skips (not a stall) + // - timeout_at = NULL → handleTimeouts skips (needs NOT NULL) + // - started_at 10s ago with timeout_ms=1000 → wall-clock matches + // (2 × timeout_ms = 2000ms threshold exceeded) + if (wedgeRescue) { + const wedgedJob = await queue.add('noop', {}, { + queue: 'smoke', + timeout_ms: 1000, + }); + await engine.executeRaw( + `UPDATE minion_jobs + SET status='active', + lock_token='smoke-wedge-rescue', + lock_until=now() + interval '30 seconds', + started_at=now() - interval '10 seconds', + timeout_at=NULL, + attempts_started = attempts_started + 1 + WHERE id=$1`, + [wedgedJob.id] + ); + + const stallResult = await queue.handleStalled(); + const stalledStatus = await queue.getJob(wedgedJob.id); + const timeoutResult = await queue.handleTimeouts(); + const timedStatus = await queue.getJob(wedgedJob.id); + const wallResult = await queue.handleWallClockTimeouts(30000); + const finalStatus = await queue.getJob(wedgedJob.id); + + if (finalStatus?.status !== 'dead') { + console.error( + `SMOKE FAIL (--wedge-rescue) — wall-clock sweep did not evict job #${wedgedJob.id}. ` + + `Status: ${finalStatus?.status}. ` + + `handleStalled: requeued=${stallResult.requeued.length} dead=${stallResult.dead.length}, after: ${stalledStatus?.status}; ` + + `handleTimeouts: ${timeoutResult.length}, after: ${timedStatus?.status}; ` + + `handleWallClockTimeouts: ${wallResult.length}, final: ${finalStatus?.status}.` + ); + process.exit(1); + } + if (finalStatus.error_text !== 'wall-clock timeout exceeded') { + console.error( + `SMOKE FAIL (--wedge-rescue) — dead, but error_text='${finalStatus.error_text}' ` + + `(expected 'wall-clock timeout exceeded').` + ); + process.exit(1); + } + try { await queue.removeJob(wedgedJob.id); } catch { /* non-fatal cleanup */ } + } + const cfg = (await import('../core/config.ts')).loadConfig(); const engineLabel = cfg?.engine ?? 'unknown'; - const tag = sigkillRescue ? ' + SIGKILL rescue' : ''; + const tags: string[] = []; + if (sigkillRescue) tags.push('SIGKILL rescue'); + if (wedgeRescue) tags.push('wedge rescue'); + const tag = tags.length > 0 ? ` + ${tags.join(' + ')}` : ''; console.log(`SMOKE PASS — Minions healthy${tag} in ${elapsedSec}s (engine: ${engineLabel})`); if (engineLabel === 'pglite') { console.log('Note: the `gbrain jobs work` daemon requires Postgres. PGLite'); @@ -503,7 +610,7 @@ HANDLER TYPES (built in) } const queueName = parseFlag(args, '--queue') ?? 'default'; - const concurrency = parseInt(parseFlag(args, '--concurrency') ?? '1', 10); + const concurrency = resolveWorkerConcurrency(args); try { await queue.ensureSchema(); } catch (e) { console.error(e instanceof Error ? e.message : String(e)); process.exit(1); } @@ -819,16 +926,17 @@ export async function registerBuiltinHandlers(worker: MinionWorker, engine: Brai }; }); - // Shell handler: registered ONLY when GBRAIN_ALLOW_SHELL_JOBS=1 is set on the - // worker process. Default-closed; opt-in per-host. Without the flag, shell - // jobs submitted via CLI insert rows but no worker claims them (they sit in - // 'waiting' — the CLI prints a starvation warning for that case). - if (process.env.GBRAIN_ALLOW_SHELL_JOBS === '1') { + // Shell handler is always registered. Runtime env guard lives inside the + // handler so claimed jobs emit a clear rejection log on workers missing + // GBRAIN_ALLOW_SHELL_JOBS=1. + { const { shellHandler } = await import('../core/minions/handlers/shell.ts'); worker.register('shell', shellHandler); - process.stderr.write('[minion worker] shell handler enabled (GBRAIN_ALLOW_SHELL_JOBS=1)\n'); - } else { - process.stderr.write('[minion worker] shell handler disabled (set GBRAIN_ALLOW_SHELL_JOBS=1 to enable)\n'); + if (process.env.GBRAIN_ALLOW_SHELL_JOBS === '1') { + process.stderr.write('[minion worker] shell handler enabled (GBRAIN_ALLOW_SHELL_JOBS=1)\n'); + } else { + process.stderr.write('[minion worker] shell handler registered in guarded mode (set GBRAIN_ALLOW_SHELL_JOBS=1 to execute shell jobs)\n'); + } } // v0.15 subagent handlers: always-on. Unlike shell (which needs an env diff --git a/src/core/minions/backpressure-audit.ts b/src/core/minions/backpressure-audit.ts new file mode 100644 index 000000000..dd0d63abd --- /dev/null +++ b/src/core/minions/backpressure-audit.ts @@ -0,0 +1,77 @@ +/** + * Backpressure audit log — operational trace for `maxWaiting` coalesce events. + * + * Mirrors the shell-audit.ts pattern (ISO-week-rotated JSONL, best-effort writes, + * failures go to stderr but never block submission). The incident that motivated + * maxWaiting (autopilot pile-up during a 90+ min queue wedge) was invisible + * precisely because the coalesce silently dropped repeat submissions. This + * trail answers "why is queue depth steady at 2 for this name?" without any + * doctor scan. + * + * File: `~/.gbrain/audit/backpressure-YYYY-Www.jsonl` (override dir via + * `GBRAIN_AUDIT_DIR` for container/sandbox deployments where `$HOME` is read-only). + * + * `gbrain jobs stats` will surface coalesce counts from this file in a v0.19.2+ + * follow-up (B4). The audit trail is for operators debugging live queues, not + * for compliance — a disk-full attacker can silently disable it. + */ + +import * as fs from 'node:fs'; +import * as path from 'node:path'; +import * as os from 'node:os'; + +export interface BackpressureAuditEvent { + ts: string; + queue: string; + name: string; + waiting_count: number; + max_waiting: number; + decision: 'coalesced'; + returned_job_id: number; +} + +/** Compute `backpressure-YYYY-Www.jsonl` using ISO-8601 week numbering. + * + * Copy of the shell-audit computeAuditFilename algorithm, parameterized on + * the filename prefix. Keeping the math inline (rather than re-exporting from + * shell-audit.ts) avoids a cross-module dependency between two best-effort + * audit surfaces — one can be rewritten without touching the other. + */ +export function computeAuditFilename(now: Date = new Date()): string { + const d = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), now.getUTCDate())); + const dayNum = (d.getUTCDay() + 6) % 7; // Mon=0, Sun=6 + d.setUTCDate(d.getUTCDate() - dayNum + 3); // shift to Thursday + const isoYear = d.getUTCFullYear(); + const firstThursday = new Date(Date.UTC(isoYear, 0, 4)); + const firstThursdayDayNum = (firstThursday.getUTCDay() + 6) % 7; + firstThursday.setUTCDate(firstThursday.getUTCDate() - firstThursdayDayNum + 3); + const weekNum = Math.round((d.getTime() - firstThursday.getTime()) / (7 * 86400000)) + 1; + const ww = String(weekNum).padStart(2, '0'); + return `backpressure-${isoYear}-W${ww}.jsonl`; +} + +/** Honors `GBRAIN_AUDIT_DIR` for container/sandbox deployments. */ +export function resolveAuditDir(): string { + const override = process.env.GBRAIN_AUDIT_DIR; + if (override && override.trim().length > 0) return override; + return path.join(os.homedir(), '.gbrain', 'audit'); +} + +export function logBackpressureCoalesce(event: Omit): void { + const dir = resolveAuditDir(); + const filename = computeAuditFilename(); + const fullPath = path.join(dir, filename); + const line = JSON.stringify({ + ...event, + decision: 'coalesced' as const, + ts: new Date().toISOString(), + }) + '\n'; + + try { + fs.mkdirSync(dir, { recursive: true }); + fs.appendFileSync(fullPath, line, { encoding: 'utf8' }); + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + process.stderr.write(`[backpressure-audit] write failed (${msg}); submission continues\n`); + } +} diff --git a/src/core/minions/handlers/shell.ts b/src/core/minions/handlers/shell.ts index 68604140d..9ffd5e83c 100644 --- a/src/core/minions/handlers/shell.ts +++ b/src/core/minions/handlers/shell.ts @@ -207,6 +207,16 @@ class TailBuffer { /** The shell handler itself. */ export async function shellHandler(ctx: MinionJobContext): Promise { + if (process.env.GBRAIN_ALLOW_SHELL_JOBS !== '1') { + const warning = + `[shell] Job #${ctx.id} rejected: GBRAIN_ALLOW_SHELL_JOBS=1 not set on this worker.\n` + + ' Shell jobs require the env var on the worker process.'; + console.warn(warning); + throw new UnrecoverableError( + 'shell handler disabled on this worker (set GBRAIN_ALLOW_SHELL_JOBS=1 to execute shell jobs)', + ); + } + const params = validateParams(ctx.data); const env = buildChildEnv(params.env); const startedAt = Date.now(); diff --git a/src/core/minions/queue.ts b/src/core/minions/queue.ts index 25a167140..40ad868bf 100644 --- a/src/core/minions/queue.ts +++ b/src/core/minions/queue.ts @@ -102,6 +102,64 @@ export class MinionQueue { if (existing.length > 0) return rowToMinionJob(existing[0]); } + // 1b. Submission-time backpressure for high-frequency named jobs. + // If waiting jobs for this (name, queue) already hit maxWaiting, return + // the most-recent waiting row instead of inserting another slot. + // + // Correctness: two concurrent submitters could both see waitingCount < + // maxWaiting and both insert, violating the cap. `pg_advisory_xact_lock` + // keyed on (name, queue) serializes concurrent count+insert decisions + // for the SAME key while leaving different keys fully parallel. The + // lock releases on txn commit/rollback automatically — no cleanup path + // to leak. Cost: one no-op SELECT on the hot path per coalesce-guarded + // submission; trivial compared to the protection. + // + // Queue scope: the filter includes `queue=$2` so a waiting + // 'autopilot-cycle' in queue 'default' does NOT suppress submissions + // to queue 'shell' with the same name. Pre-D2 code filtered on `name` + // alone — a real cross-queue bleed that sequential tests missed. + // + // Engine compatibility: PGLite (WASM Postgres 17) supports + // pg_advisory_xact_lock, so this works on both engines without branching. + if (opts?.maxWaiting !== undefined) { + const maxWaiting = Math.max(1, Math.floor(opts.maxWaiting)); + const backpressureQueue = opts?.queue ?? 'default'; + await tx.executeRaw( + `SELECT pg_advisory_xact_lock(hashtext('minion_maxwaiting:' || $1 || ':' || $2))`, + [jobName, backpressureQueue] + ); + const waitingCountRows = await tx.executeRaw<{ count: string }>( + `SELECT count(*)::text AS count + FROM minion_jobs + WHERE name = $1 AND queue = $2 AND status = 'waiting'`, + [jobName, backpressureQueue] + ); + const waitingCount = parseInt(waitingCountRows[0]?.count ?? '0', 10); + if (waitingCount >= maxWaiting) { + const existingWaiting = await tx.executeRaw>( + `SELECT * FROM minion_jobs + WHERE name = $1 AND queue = $2 AND status = 'waiting' + ORDER BY created_at DESC, id DESC + LIMIT 1`, + [jobName, backpressureQueue] + ); + if (existingWaiting.length > 0) { + const coalesced = rowToMinionJob(existingWaiting[0]); + try { + const { logBackpressureCoalesce } = await import('./backpressure-audit.ts'); + logBackpressureCoalesce({ + queue: backpressureQueue, + name: jobName, + waiting_count: waitingCount, + max_waiting: maxWaiting, + returned_job_id: coalesced.id, + }); + } catch { /* audit failures never block submission */ } + return coalesced; + } + } + } + // 2. Parent lock + depth/cap validation let depth = 0; if (opts?.parent_job_id) { @@ -563,6 +621,77 @@ export class MinionQueue { }); } + /** + * Dead-letter active jobs that exceed a wall-clock runtime threshold, + * regardless of lock state. This catches jobs stuck while still holding + * DB resources (e.g. blocked on file locks) where stall sweeps skip rows. + * + * Threshold (ms): + * timeout_ms set -> timeout_ms * 2 + * timeout_ms null -> 2 * lockDurationMs * max_stalled + */ + async handleWallClockTimeouts(lockDurationMs: number): Promise { + return this.engine.transaction(async (tx) => { + const rows = await tx.executeRaw>( + `UPDATE minion_jobs SET + status = 'dead', + error_text = 'wall-clock timeout exceeded', + lock_token = NULL, + lock_until = NULL, + finished_at = now(), + updated_at = now() + WHERE status = 'active' + AND started_at IS NOT NULL + AND EXTRACT(EPOCH FROM (now() - started_at)) * 1000 > + CASE + WHEN timeout_ms IS NOT NULL THEN timeout_ms * 2 + ELSE $1::double precision * 2 * GREATEST(max_stalled, 1) + END + RETURNING *`, + [lockDurationMs] + ); + + const parentIds = new Set(); + for (const r of rows) { + const parentJobId = r.parent_job_id as number | null; + if (parentJobId == null) continue; + parentIds.add(parentJobId); + const childDone: ChildDoneMessage = { + type: 'child_done', + child_id: r.id as number, + job_name: r.name as string, + result: null, + outcome: 'timeout', + error: 'wall-clock timeout exceeded', + }; + await tx.executeRaw( + `INSERT INTO minion_inbox (job_id, sender, payload) + SELECT $1, 'minions', $2::jsonb + WHERE EXISTS ( + SELECT 1 FROM minion_jobs + WHERE id = $1 AND status NOT IN ('completed','failed','dead','cancelled') + )`, + [parentJobId, childDone] + ); + } + + for (const parentId of parentIds) { + await tx.executeRaw( + `UPDATE minion_jobs SET status = 'waiting', updated_at = now() + WHERE id = $1 AND status = 'waiting-children' + AND NOT EXISTS ( + SELECT 1 FROM minion_jobs + WHERE parent_job_id = $1 + AND status NOT IN ('completed', 'failed', 'dead', 'cancelled') + )`, + [parentId] + ); + } + + return rows.map(rowToMinionJob); + }); + } + /** * Complete a job (token-fenced). All side effects atomic in one transaction: * 1. UPDATE child to 'completed' with result diff --git a/src/core/minions/types.ts b/src/core/minions/types.ts index c7f7fba04..fc7ca5e7e 100644 --- a/src/core/minions/types.ts +++ b/src/core/minions/types.ts @@ -128,6 +128,8 @@ export interface MinionJobInput { max_spawn_depth?: number; /** Global dedup key. Same key returns the existing job, no second row created. */ idempotency_key?: string; + /** Submission backpressure: cap waiting jobs with this name before inserting a new row. */ + maxWaiting?: number; // v12: scheduler polish /** diff --git a/src/core/minions/worker.ts b/src/core/minions/worker.ts index b8e1daa4f..e6f8a9342 100644 --- a/src/core/minions/worker.ts +++ b/src/core/minions/worker.ts @@ -126,6 +126,14 @@ export class MinionWorker { } catch (e) { console.error('Timeout detection error:', e instanceof Error ? e.message : String(e)); } + try { + const wallClockTimedOut = await this.queue.handleWallClockTimeouts(this.opts.lockDuration); + if (wallClockTimedOut.length > 0) { + console.log(`Wall-clock detector: dead-lettered ${wallClockTimedOut.length} jobs (wall-clock timeout exceeded)`); + } + } catch (e) { + console.error('Wall-clock timeout detection error:', e instanceof Error ? e.message : String(e)); + } }, this.opts.stalledInterval); try { diff --git a/test/minions-shell.test.ts b/test/minions-shell.test.ts index dbfd13d3e..b4f5ec117 100644 --- a/test/minions-shell.test.ts +++ b/test/minions-shell.test.ts @@ -12,8 +12,18 @@ import * as os from 'node:os'; let engine: PGLiteEngine; let queue: MinionQueue; +// The shell handler at src/core/minions/handlers/shell.ts:210 throws +// UnrecoverableError when GBRAIN_ALLOW_SHELL_JOBS !== '1'. That's the +// production-worker RCE guard. Unit tests here exercise the handler +// mechanics, not the guard, so we enable it for the whole file and +// restore on teardown. The separate "rejects when env not set" case +// (in the minion-shell submission E2E / the queue-resilience wave) +// toggles the var itself. +let prevAllowShellJobs: string | undefined; beforeAll(async () => { + prevAllowShellJobs = process.env.GBRAIN_ALLOW_SHELL_JOBS; + process.env.GBRAIN_ALLOW_SHELL_JOBS = '1'; engine = new PGLiteEngine(); await engine.connect({ database_url: '' }); await engine.initSchema(); @@ -22,6 +32,8 @@ beforeAll(async () => { afterAll(async () => { await engine.disconnect(); + if (prevAllowShellJobs === undefined) delete process.env.GBRAIN_ALLOW_SHELL_JOBS; + else process.env.GBRAIN_ALLOW_SHELL_JOBS = prevAllowShellJobs; }); beforeEach(async () => { diff --git a/test/minions.test.ts b/test/minions.test.ts index aee59cabd..ce377a1f8 100644 --- a/test/minions.test.ts +++ b/test/minions.test.ts @@ -1653,3 +1653,246 @@ describe('MinionQueue: Attachments', () => { expect(list[0].filename).toBe('b.txt'); }); }); + +// --- v0.19.1 — queue-resilience (wall-clock sweep, maxWaiting race, concurrency clamp) --- + +describe('MinionQueue: v0.19.1 handleWallClockTimeouts (Layer 3 kill shot)', () => { + test('evicts active job past 2× timeout_ms — sets dead + wall-clock error_text', async () => { + const job = await queue.add('noop', {}, { timeout_ms: 100 }); + await engine.executeRaw( + `UPDATE minion_jobs + SET status='active', + lock_token='wc-test', + lock_until=now() - interval '1 second', + started_at=now() - interval '1 second', + timeout_at=now() - interval '0.9 second', + attempts_started = attempts_started + 1 + WHERE id=$1`, + [job.id], + ); + const killed = await queue.handleWallClockTimeouts(30_000); + expect(killed.length).toBe(1); + expect(killed[0].id).toBe(job.id); + const after = await queue.getJob(job.id); + expect(after?.status).toBe('dead'); + expect(after?.error_text).toBe('wall-clock timeout exceeded'); + }); + + test('timeout_ms NULL fallback uses 2 × lockDuration × max_stalled threshold', async () => { + const job = await queue.add('noop', {}, { max_stalled: 3 }); + // Force timeout_ms / timeout_at NULL on-disk (columns might or might not be set by add). + await engine.executeRaw( + `UPDATE minion_jobs + SET status='active', + timeout_ms=NULL, + timeout_at=NULL, + lock_token='wc-null', + lock_until=now() - interval '1 second', + started_at=now() - interval '61 seconds', + attempts_started = attempts_started + 1 + WHERE id=$1`, + [job.id], + ); + // 2 × lockDurationMs × max_stalled = 2 × 10_000 × 3 = 60_000 ms. started_at is 61s ago. + const killed = await queue.handleWallClockTimeouts(10_000); + expect(killed.length).toBe(1); + expect(killed[0].id).toBe(job.id); + const after = await queue.getJob(job.id); + expect(after?.status).toBe('dead'); + }); + + test('respects threshold — active job within window is NOT killed', async () => { + const job = await queue.add('noop', {}, { timeout_ms: 100_000 }); + await engine.executeRaw( + `UPDATE minion_jobs + SET status='active', + lock_token='wc-inside', + lock_until=now() + interval '30 seconds', + started_at=now() - interval '10 seconds', + timeout_at=now() + interval '90 seconds', + attempts_started = attempts_started + 1 + WHERE id=$1`, + [job.id], + ); + const killed = await queue.handleWallClockTimeouts(30_000); + expect(killed.length).toBe(0); + const after = await queue.getJob(job.id); + expect(after?.status).toBe('active'); + }); +}); + +describe('MinionQueue: v0.19.1 maxWaiting — cap correctness + race (D2/H2)', () => { + test('coalesces 3rd submission when cap is 2 — returns existing most-recent waiting row', async () => { + const a = await queue.add('poll', {}, { maxWaiting: 2 }); + const b = await queue.add('poll', {}, { maxWaiting: 2 }); + const c = await queue.add('poll', {}, { maxWaiting: 2 }); + expect(a.id).not.toBe(b.id); + expect(c.id).toBe(b.id); // coalesced to the most-recent waiting row + const rows = await engine.executeRaw<{ count: string }>( + `SELECT count(*)::text AS count FROM minion_jobs WHERE name='poll' AND status='waiting'`, + ); + expect(parseInt(rows[0].count, 10)).toBe(2); + }); + + test('clamps maxWaiting: 0 → 1 (strictest cap)', async () => { + const a = await queue.add('squeeze', {}, { maxWaiting: 0 }); + const b = await queue.add('squeeze', {}, { maxWaiting: 0 }); + expect(b.id).toBe(a.id); // 0 clamped to 1, 2nd coalesces into 1st + }); + + test('floors maxWaiting: 1.7 → 1', async () => { + const a = await queue.add('floor', {}, { maxWaiting: 1.7 }); + const b = await queue.add('floor', {}, { maxWaiting: 1.7 }); + expect(b.id).toBe(a.id); + }); + + test('concurrent submitters respect the cap under Promise.all race (H2)', async () => { + // Serialized by pg_advisory_xact_lock keyed on (name, queue). Without it, + // two concurrent submits both see count( + `SELECT count(*)::text AS count FROM minion_jobs WHERE name='race' AND status='waiting'`, + ); + expect(parseInt(rows[0].count, 10)).toBe(2); // cap held under concurrency + }); + + test('cross-queue isolation — same name in queue A does NOT suppress queue B (H2 secondary)', async () => { + const a = await queue.add('isolate', {}, { maxWaiting: 1, queue: 'default' }); + // cap hit on queue=default with maxWaiting=1; 2nd would coalesce into `a` + const a2 = await queue.add('isolate', {}, { maxWaiting: 1, queue: 'default' }); + expect(a2.id).toBe(a.id); + // Different queue — MUST insert a fresh row, NOT coalesce into queue=default + const b = await queue.add('isolate', {}, { maxWaiting: 1, queue: 'shell' }); + expect(b.id).not.toBe(a.id); + expect(b.queue).toBe('shell'); + }); + + test('unset maxWaiting — normal submit path, no coalesce, no cap', async () => { + const a = await queue.add('uncapped', {}); + const b = await queue.add('uncapped', {}); + const c = await queue.add('uncapped', {}); + expect(new Set([a.id, b.id, c.id]).size).toBe(3); + }); +}); + +describe('resolveWorkerConcurrency (v0.19.1 H3): clamp + validation', () => { + // jobs.ts handler — tested via direct import. Warning goes to stderr; + // tests verify return value only, not the warning line. + let resolveWorkerConcurrency: (args: string[], env?: NodeJS.ProcessEnv) => number; + let parseMaxWaitingFlag: (args: string[]) => number | undefined; + beforeAll(async () => { + const mod = await import('../src/commands/jobs.ts'); + resolveWorkerConcurrency = mod.resolveWorkerConcurrency; + parseMaxWaitingFlag = mod.parseMaxWaitingFlag; + }); + + test('flag=4 env-unset → 4', () => { + expect(resolveWorkerConcurrency(['--concurrency', '4'], {} as NodeJS.ProcessEnv)).toBe(4); + }); + test('flag-unset env=8 → 8', () => { + expect(resolveWorkerConcurrency([], { GBRAIN_WORKER_CONCURRENCY: '8' } as NodeJS.ProcessEnv)).toBe(8); + }); + test('flag=2 env=8 → 2 (flag wins)', () => { + expect(resolveWorkerConcurrency(['--concurrency', '2'], { GBRAIN_WORKER_CONCURRENCY: '8' } as NodeJS.ProcessEnv)).toBe(2); + }); + test('both unset → 1', () => { + expect(resolveWorkerConcurrency([], {} as NodeJS.ProcessEnv)).toBe(1); + }); + test('garbage env "foo" → clamped to 1 (H3)', () => { + expect(resolveWorkerConcurrency([], { GBRAIN_WORKER_CONCURRENCY: 'foo' } as NodeJS.ProcessEnv)).toBe(1); + }); + test('env=0 → clamped to 1 (H3 — prevents silent wedge)', () => { + expect(resolveWorkerConcurrency([], { GBRAIN_WORKER_CONCURRENCY: '0' } as NodeJS.ProcessEnv)).toBe(1); + }); + test('env=-5 → clamped to 1 (H3)', () => { + expect(resolveWorkerConcurrency([], { GBRAIN_WORKER_CONCURRENCY: '-5' } as NodeJS.ProcessEnv)).toBe(1); + }); +}); + +describe('parseMaxWaitingFlag (v0.19.1 H5): CLI flag wiring', () => { + let parseMaxWaitingFlag: (args: string[]) => number | undefined; + beforeAll(async () => { + parseMaxWaitingFlag = (await import('../src/commands/jobs.ts')).parseMaxWaitingFlag; + }); + + test('absent → undefined (no cap, default submit path)', () => { + expect(parseMaxWaitingFlag(['foo', '--params', '{}'])).toBeUndefined(); + }); + test('--max-waiting 2 → 2 (happy path)', () => { + expect(parseMaxWaitingFlag(['foo', '--max-waiting', '2'])).toBe(2); + }); + test('--max-waiting 200 → clamped to 100', () => { + expect(parseMaxWaitingFlag(['foo', '--max-waiting', '200'])).toBe(100); + }); + test('--max-waiting 0 → throws', () => { + expect(() => parseMaxWaitingFlag(['foo', '--max-waiting', '0'])).toThrow('positive integer'); + }); + test('--max-waiting abc → throws', () => { + expect(() => parseMaxWaitingFlag(['foo', '--max-waiting', 'abc'])).toThrow('positive integer'); + }); +}); + +describe('backpressure-audit (v0.19.1 Q1): JSONL on coalesce', () => { + test('logBackpressureCoalesce writes one JSONL line per coalesce', async () => { + const { logBackpressureCoalesce, resolveAuditDir, computeAuditFilename } = + await import('../src/core/minions/backpressure-audit.ts'); + const fs = await import('node:fs'); + const path = await import('node:path'); + const os = await import('node:os'); + + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), 'gbrain-audit-')); + const prev = process.env.GBRAIN_AUDIT_DIR; + process.env.GBRAIN_AUDIT_DIR = tmp; + try { + expect(resolveAuditDir()).toBe(tmp); + logBackpressureCoalesce({ + queue: 'default', + name: 'poll', + waiting_count: 2, + max_waiting: 2, + returned_job_id: 42, + }); + const file = path.join(tmp, computeAuditFilename()); + const text = fs.readFileSync(file, 'utf8'); + const line = JSON.parse(text.trim()); + expect(line.decision).toBe('coalesced'); + expect(line.name).toBe('poll'); + expect(line.returned_job_id).toBe(42); + expect(typeof line.ts).toBe('string'); + } finally { + if (prev === undefined) delete process.env.GBRAIN_AUDIT_DIR; + else process.env.GBRAIN_AUDIT_DIR = prev; + fs.rmSync(tmp, { recursive: true, force: true }); + } + }); +}); + +describe('MinionQueue: v0.19.1 wall-clock + handleTimeouts non-interference (T1)', () => { + test('wall-clock sweep does NOT evict a job that handleTimeouts would handle', async () => { + // Retry-able timeout: timeout_at < now() AND lock_until > now() — handleTimeouts + // is the correct killer here. wall-clock's 2× threshold has not fired yet. + const job = await queue.add('noop', {}, { timeout_ms: 100_000 }); + await engine.executeRaw( + `UPDATE minion_jobs + SET status='active', + lock_token='t1', + lock_until=now() + interval '30 seconds', + started_at=now() - interval '2 seconds', + timeout_at=now() - interval '0.5 seconds', + attempts_started = attempts_started + 1 + WHERE id=$1`, + [job.id], + ); + // At this point: started_at is 2s ago, 2×timeout_ms = 200s. Wall-clock should NOT fire. + const killed = await queue.handleWallClockTimeouts(30_000); + expect(killed.length).toBe(0); + const after = await queue.getJob(job.id); + expect(after?.status).toBe('active'); + }); +});