mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-14 17:02:19 +00:00
Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3dc1a95592 | ||
|
|
312c449439 | ||
|
|
3a409ce0d2 | ||
|
|
dfc5284005 | ||
|
|
b0421574d2 | ||
|
|
a9af1fae7d | ||
|
|
ca02a5f51b | ||
|
|
f63d444453 |
@@ -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 <id>` 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.**
|
||||
|
||||
@@ -62,12 +62,13 @@ strict behavior when unset.
|
||||
- `src/commands/graph-query.ts` — `gbrain graph-query <slug> [--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 <version>` 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 <id>`.
|
||||
- `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=<ms>` 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.
|
||||
|
||||
@@ -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 <id>
|
||||
```
|
||||
|
||||
## Rescue actions (in order of escalation)
|
||||
|
||||
```bash
|
||||
# Force-kill a single stuck job:
|
||||
gbrain jobs cancel <id>
|
||||
|
||||
# Clear a specific job entirely (last resort):
|
||||
gbrain jobs delete <id>
|
||||
|
||||
# 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.
|
||||
+3
-2
@@ -141,12 +141,13 @@ strict behavior when unset.
|
||||
- `src/commands/graph-query.ts` — `gbrain graph-query <slug> [--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 <version>` 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 <id>`.
|
||||
- `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=<ms>` 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.
|
||||
|
||||
+1
-1
@@ -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",
|
||||
|
||||
@@ -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 <path>] [--interval N] [--json]\n' +
|
||||
'Usage: gbrain autopilot [--repo <path>] [--interval N] [--json] [--no-worker]\n' +
|
||||
' gbrain autopilot --install [--repo <path>]\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
|
||||
|
||||
@@ -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 <id>` to inspect
|
||||
// or `gbrain jobs cancel <id>` 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 <id> to inspect; gbrain jobs cancel <id> 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
|
||||
|
||||
+118
-10
@@ -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<void
|
||||
USAGE
|
||||
gbrain jobs submit <name> [--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
|
||||
|
||||
@@ -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<BackpressureAuditEvent, 'ts' | 'decision'>): 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`);
|
||||
}
|
||||
}
|
||||
@@ -207,6 +207,16 @@ class TailBuffer {
|
||||
|
||||
/** The shell handler itself. */
|
||||
export async function shellHandler(ctx: MinionJobContext): Promise<ShellJobResult> {
|
||||
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();
|
||||
|
||||
@@ -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<Record<string, unknown>>(
|
||||
`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<MinionJob[]> {
|
||||
return this.engine.transaction(async (tx) => {
|
||||
const rows = await tx.executeRaw<Record<string, unknown>>(
|
||||
`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<number>();
|
||||
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
|
||||
|
||||
@@ -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
|
||||
/**
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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<max and both insert — the TOCTOU
|
||||
// bug codex caught in D2/H2.
|
||||
const results = await Promise.all([
|
||||
queue.add('race', {}, { maxWaiting: 2 }),
|
||||
queue.add('race', {}, { maxWaiting: 2 }),
|
||||
queue.add('race', {}, { maxWaiting: 2 }),
|
||||
]);
|
||||
expect(results.length).toBe(3);
|
||||
const rows = await engine.executeRaw<{ count: string }>(
|
||||
`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');
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user