Skip to content

fix(operator-queue): the flood guard no longer feeds on itself (#3130) - #3220

Merged
vybe merged 3 commits into
devfrom
feature/3130-flood-guard-self-count
Oct 5, 2026
Merged

vybe merged 3 commits into
devfrom
feature/3130-flood-guard-self-count

Conversation

@dolho

@dolho dolho commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Summary

The #1632 flood guard fed on itself. Its alert was created as a pending row against the same agent whose pending count tripped it, and that count includes the alert. Once an agent reached the cap (25), the condition could never clear. The 300 s cooldown then created a new timestamped alert every window. On one install, 386 of an agent's 435 queue rows were flood alerts. The same count gates ask_operator (queue_full), so the backlog also refused the agent's real asks. Held queue-file requests were never reported back to the agent.

  • The cap counts only the agent's own rows. Both depth checks now share one predicate, db/operator_queue.py::_own_pending_conds: the queue-file cap and the native ask_operator queue_full check. Rows carrying a platform-reserved id prefix are excluded. Agents can't create those prefixes on either path, so the exclusion can't hide an agent's own row.
  • One alert per over-cap episode. The next alert needs a cycle that holds nothing first. The cooldown is now only the minimum gap between episodes.
  • Budgeted, as vybe's issue comment proposed. The flood alert is the registered queue_flood type, created through create_bounded_alert (bug: platform operator-queue alert emitters bypass #1632 ingestion caps — skill-not-found is agent-triggerable (flood residual, gates pull default-ON) #1677). The budget is a backstop in case of leader failover. The UI still labels it "Heads up" with the alert pill, on both the card and the detail panel.
  • A held queue file is told why. A file-level platform.ingestion block (reason: queue_full / rate_limited / invalid_id, plus max_pending and since) rides the existing re-read + if_match write-back. It is written only when it changes and removed when nothing is held. A refused write is retried only after the file changes. The block is documented in the agent prompt (both copies, prompt parity green).
  • Backlog cleared once. Migration supersede_queue_flood_backlog (SQLite runner + Alembic 0089, sharing one SQL statement) keeps the newest pending flood row per agent. The rest become cancelled, with disposed_by = 'platform' and reason superseded. Platform-minted rows never wake the agent or touch its file.

Acceptance criteria

Known residual

At its budget (5 pending queue_flood rows), create_bounded_alert emits its own alert-budget- alert, at most once per cooldown window. With the edge trigger in place this path is only reachable after repeated leader failovers. The emitter is unchanged (#1677 behaviour), and it is the same window-repeating shape, so it may deserve its own follow-up.

Changes

  • src/backend/db/operator_queue.py, src/backend/database.py: _own_pending_conds, plus an exclude_request_id_prefixes kwarg on count_pending_for_agent and create_native_item, threaded through both facades.
  • src/backend/services/operator_queue_service.py: the queue_flood budget type, the edge-triggered emitter, the hold reason, the platform.ingestion marker on the write-back, and a write-back return value. The marker argument defaults to "unchanged".
  • src/backend/services/ask_service.py: the native create passes the reserved prefixes.
  • src/backend/db/migrations.py, src/backend/migrations/versions/0089_supersede_queue_flood_backlog.py: the backlog migration (data only).
  • src/frontend/src/utils/operatorQueue.js, QueueCard.vue, QueueItemDetail.vue: the queue_flood label and pill/badge, using semantic tokens. The detail header renders the shared queueTypeLabel instead of the raw type string.
  • Review follow-up (I1): create_bounded_alert_outcome returns the four-outcome result; create_bounded_alert is its bool view (callers unchanged). The flood emitter releases its episode on a failed count/create so the alarm retries once per cooldown window; a budget refusal still consumes the episode.
  • config/trinity-meta-prompt/prompt.md, platform_prompt_service.py: the platform.ingestion contract.
  • Docs: operating-room.md, api-endpoints.md, requirements/security.md, tables.py (the disposed_by comment), and a learnings fragment.

Test Plan

Fixes #3130

🤖 Generated with Claude Code

dolho and others added 2 commits October 5, 2026 10:38
The #1632 flood alert was a raw platform create filed against the agent
whose depth it reported, the depth read counted it, and a sustained
over-cap condition minted a fresh timestamped alert every cooldown
window. Measured on one install: 386 of 435 rows were flood alerts, and
four file-seam requests were held with nothing telling the agent.

- The agent's own budget excludes platform-minted rows. Both depth caps
  (file seam, native ask_operator queue_full) read one predicate,
  _own_pending_conds, with _RESERVED_ID_PREFIXES excluded.
- The flood alert is edge-triggered: one per over-cap episode, re-armed
  only when a cycle holds nothing; the cooldown is now only the minimum
  spacing between episodes.
- It is a budgeted platform alert (#1677): type queue_flood through
  create_bounded_alert, as the backstop for failover. The UI labels it
  "Heads up" with the alert pill, as before.
- A held file is told why: a file-level platform.ingestion block
  (queue_full / rate_limited / invalid_id, max_pending, since) rides the
  guarded write-back, written only on change, removed when nothing is
  held; a refused write is retried only after the file changes.
- Migration supersede_queue_flood_backlog (SQLite + Alembic 0089) keeps
  the newest pending flood row per agent and cancels the rest
  (disposed_by 'platform', reason 'superseded').

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ow-ups (#3130)

- _write_responses_to_agent's ingestion argument defaults to "unchanged",
  never "remove": a caller that omits it must not strip a hold or force a
  write. Pinned by a test that fails with the old default.
- tables.py: disposed_by now also records 'platform' (the backlog
  migration's superseded flood alarms).
- Learnings fragment: an alarm must not count toward the limit it reports.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@dolho dolho added the ui PR touches the frontend UI — triggers Playwright e2e tests label Oct 5, 2026
@dolho

dolho commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

/review — automated pre-landing review (at d9dcc36325caf73dde275f9a70a3a9b8eaa85300)

Branch: feature/3130-flood-guard-self-count → dev (merge-base 482ad0742, no drift on dev in the touched files)
Files Changed: 20 (+742/-65)
Scope: CLEAN. Every change sits on the flood-guard / cap / write-back / backlog path, or is the matching doc, prompt or UI label.
Plan Completion (issue ACs + vybe's 10-02 "do both halves" ruling on AC1): 4 done, 1 changed, 0 partial, 0 not done

AC Status Evidence
AC1a: emitter routed through create_bounded_alert with a registered type DONE operator_queue_service.py:152 ("queue_flood" in _BUDGETED_ALERT_TYPES), :2253 created = await create_bounded_alert(agent_name, alert); removed from the platform-only allowlist in test_1677_operator_alert_emitters.py
AC1b: cap excludes platform rows (both caps) DONE db/operator_queue.py:49 _own_pending_conds, used at :356 (native queue_full) and :1645; operator_queue_service.py:1593 and ask_service.py:470 pass _RESERVED_ID_PREFIXES
AC2: edge-triggered re-alert DONE :2209 episode check, :1880 episode ends on held == 0, :1522 oversize episode ends on a sane-size read
AC3: held file reported to the agent DONE platform.ingestion marker (:775-835, :1920-1939), documented in both prompt copies
AC4: operator clears the backlog in one action CHANGED Cleared automatically, once, by migration supersede_queue_flood_backlog (SQLite migrations.py:5133-5178 + Alembic 0089). It is not an operator action. #2372 (acknowledging a reserved-prefix alert strands it in responded) is still open. That is acceptable now that new flood rows are capped at 5 per agent, but the deviation should be stated in the PR.
AC5: tests (bounded alerts, held vs delivered) DONE tests/unit/test_3130_flood_guard.py

Execution coverage (Step 2.5)

changed symbol / test file executed by live consumer verdict
_own_pending_conds / count_pending_for_agent(exclude_request_id_prefixes=) TestOwnBudgetCountRealDb (real SQLite) operator_queue_service.py:1593, create_native_item :356 ✅ executed
create_native_item(exclude_request_id_prefixes=) test_platform_rows_never_make_ask_operator_queue_full (real DB) ask_service.py:459-470 ✅ executed
_maybe_emit_flood_alert episode + budget TestFloodAlertIsBoundedAndEdgeTriggered (drives _sync_agent) _sync_agent :1509, :1877 ✅ executed
_ingestion_marker / _apply_ingestion_marker / write-back return TestIngestionMarker (drives _sync_agent, asserts on the written file) the agent, via the prompt contract ✅ executed
_migrate_supersede_queue_flood_backlog (SQLite) TestBacklogMigration (hand-rolled table) plus the real init_database() in the unit island MIGRATIONS ✅ executed
Alembic 0089 (PG) — (only the pg-migrations upgrade on an empty DB) alembic_runner.upgrade_to_head ⚠️ syntax only (see I3)
QUEUE_TYPE_LABELS.queue_flood operatorQueueResponse.spec.js, operatorQueueUnknownType.spec.js (behavioural queueTypeLabel) QueueCard.vue ✅ executed
QueueCard.vue typePill.queue_flood — rendered card n/a (style map)

Source-text grep: no new source-text assertion in the changed tests. The readFileSync in operatorQueueUnknownType.spec.js already existed, and the new block calls queueTypeLabel directly.

Fix mutation: run from a scratch copy, against test_3130_flood_guard.py (20 tests):

  • drop the depth-read exclusion at :1593: 1 red (test_the_depth_read_excludes_platform_minted_rows)
  • drop the episode check at :2209: 1 red (test_a_sustained_over_cap_condition_alerts_once)
  • raw db.create_operator_queue_item instead of create_bounded_alert: 2 red (test_it_is_created_through_the_budget, test_the_budget_is_the_hard_backstop)
  • full revert of src/backend to the merge-base: 18 red, 2 green

Local runs: test_3130 + test_1632 + test_1677: 80 passed. Vitest (operatorQueueUnknownType, operatorQueueResponse, rawColorRatchet): 83 passed. check_alembic_heads: 1 head, PASS. Enterprise-docs grep on the added doc lines: 0 hits.

Critical Findings (block merge)

None.

Informational Findings (review required)

[I1] Error handling: one transient create failure silences the flood alarm for the whole episode (Confidence: 9/10)
File: src/backend/services/operator_queue_service.py:2219-2220, :2253-2258
Evidence:

        self._flood_alert_cooldown[agent_name] = now
        self._flood_episodes.add(key)
...
            created = await create_bounded_alert(agent_name, alert)
...
        if not created:
            return

Issue: The episode key is stamped before the create, and nothing removes it when the create fails. create_bounded_alert turns a create exception into False, and every later cycle then exits at if key in self._flood_episodes: return. On dev, the cooldown was the only gate, so a failed create was retried once per window. Now a single failure, for example a SQLite database is locked at the onset, suppresses the alarm until the over-cap condition clears, which for a runaway agent may be never. The docstring's "backs off instead of retrying every 5s" no longer holds: it never retries. Verified with a probe against the PR's harness: the first create raised, the clock advanced over 20 cooldown windows, and there was 1 create attempt in total.
Suggestion: Separate "refused by the budget" from "create failed". For example, have create_bounded_alert return an outcome, or discard key from _flood_episodes on the create-failure path while keeping the cooldown stamp, so a failure retries once per window. Discarding on every not created would make the budget-full path re-enter _maybe_emit_alert_budget_episode every window, which is the residual the PR already names. Add the probe as a test.

[I2] Enum completeness: QueueItemDetail.vue was not given the new type (Confidence: 8/10)
File: src/frontend/src/components/operator/QueueItemDetail.vue:22-24, :247-252
Evidence: :class="typeBadge(item.type)" / {{ item.type }} with badges = { approval: …, question: …, alert: … } and return badges[type] || ''
Issue: QueueCard.vue got the queue_flood pill and the shared label, but the detail panel renders the raw string queue_flood with no badge colour. The PR body says "The UI still labels it 'Heads up' with the alert pill", which is true for the card only. Also, any consumer that filters the list by type=alert (routers/operator_queue.py:142, MCP or CLI callers) no longer returns new flood alarms. That is probably fine, but it is a silent change in meaning.
Suggestion: Use queueTypeLabel(item.type) in the detail header and add queue_flood to typeBadge. Better, share one pill map between the card and the detail so the next budgeted type cannot drift between them.

[I3] Migration test depth: the PostgreSQL half is never run against data (Confidence: 6/10)
File: src/backend/migrations/versions/0089_supersede_queue_flood_backlog.py:26-32, tests/unit/test_3130_flood_guard.py TestBacklogMigration
Evidence: from db.migrations import SUPERSEDE_QUEUE_FLOOD_BACKLOG_SQL … op.get_bind().execute(text(SUPERSEDE_QUEUE_FLOOD_BACKLOG_SQL), {...})
Issue: The shared SQL is a reasonable way to keep both tracks identical, and 0041 imports from the live tree the same way. But the survivor rule (correlated EXISTS on the table being updated, newer.id > operator_queue.id tie-break) is only exercised on SQLite, against a hand-written 12-column table rather than the real DDL. The pg-migrations job runs upgrade head on an empty database, so it proves the syntax, not the rule. A second consequence: because the revision imports a live constant, a later edit to SUPERSEDE_QUEUE_FLOOD_BACKLOG_SQL silently changes what historical revision 0089 does.
Suggestion: Run the migration test against the real init_database() schema, as TestOwnBudgetCountRealDb already does. Add a comment on the constant: "frozen — referenced by Alembic 0089; never edit, add a new revision".

[I4] Product quality: the scope of the AC4 deviation is not stated (Confidence: 6/10)
File: PR body, "Backlog cleared once"
Issue: AC4 asked that an operator can clear an accumulated backlog in one action, and named #2372 as the blocker. The migration clears today's backlog automatically and the budget caps future rows at 5, which makes an operator action unnecessary. #2372 is still open, though, and the surviving row per agent (type alert, pre-PR) still hits #2372's strand-in-responded path when acknowledged. Say so explicitly in the PR, or link #2372 as the follow-up, so AC4 is not read as fully closed.

Clean Categories

  • SQL safety: the migration SQL is static with bound :now / :batch_id; the prefix exclusion uses func.substr(...) != p with constants (db/operator_queue.py:1499-1506).
  • Data loss: only pending + queue-flood-% rows that have a newer pending sibling are touched. The agent's own rows, non-pending rows and single-row agents are untouched (asserted in test_keeps_the_newest_pending_flood_alert_per_agent). The downgrade no-op is justified.
  • Exclusion soundness: an agent cannot mint a reserved prefix on either seam. The file seam folds case before checking (operator_queue_service.py:1616) and so does the native seam (ask_service.py:549-555), so the case-sensitive SQL substr cannot hide an agent's own row.
  • Auth boundaries: no new endpoints or principals.
  • Credential exposure: none. The marker carries only a closed-vocabulary reason, max_pending and since; an agent-written marker is overwritten (test_a_marker_the_agent_wrote_is_not_trusted).
  • Concurrency: the marker rides the existing re-read + if_match write. A refused write is retried only after the sha changes (:1927, :1937).
  • Dual-track migrations (Invariant Fix git pushing bug #9/Feature/vector log retention #3): both tracks are present, one Alembic head, tables.py comment updated, no DDL change.
  • Enterprise disclosure: 0 guard-pattern hits in the added doc lines.
  • Design system: the pill uses the existing state-autonomous semantic tokens; the raw-colour ratchet passes.
  • Docs: operating-room.md, api-endpoints.md, requirements/security.md and both prompt copies are updated.

Low confidence (appendix)

  • (5/10) _flood_episodes is per-instance and is not cleared on a leadership transition (:1276-1280). Worker A keeps a stale (agent, "ingestion_cap") key if the episode ended while worker B led. If a new episode is already running when A regains the lease, A suppresses its alert. Narrow window; the budget does not help here because this path under-alerts rather than over-alerts. Clearing the episode set when leadership is gained would close it.
  • (4/10) _HOLD_RANK puts queue_full above invalid_id (:782), so a malformed id is reported to the agent only after the queue drains. Arguably invalid_id is the more actionable fix, since it never self-clears.

Suggested deferred debt: the alert-budget- episode emitter is still once per cooldown window while a budget stays full (the PR's "Known residual"), the same window-repeating shape this PR removes from the flood alarm. Track it as a follow-up against #1677.

Summary

  • Critical: 0 — none found
  • Informational: 4 — review recommended. I1 should be fixed before merge: it is a regression in alarm delivery compared with dev.
  • Scope: clean

🤖 Generated with Claude Code

…he episode (#3130)

Review follow-ups on #3220:

- I1: create_bounded_alert gains an outcome-returning twin
  (create_bounded_alert_outcome; the bool API is unchanged). The flood
  emitter releases its episode on a FAILED count read or create so the
  alarm is retried once per cooldown window, instead of one transient
  failure suppressing it until the over-cap condition clears. A budget
  refusal still consumes the episode: the budget emits its own episode
  alert, and re-knocking every window adds only churn.
- I2: QueueItemDetail gives queue_flood the alert badge and renders the
  shared queueTypeLabel instead of the raw type string.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@dolho

dolho commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

/review follow-ups — c29566f76

Finding Status Evidence
I1 a failed create silenced the episode Fixed New create_bounded_alert_outcome (four-outcome result; create_bounded_alert stays the bool view, so callers are unchanged). The flood emitter now releases the episode on a failed count read or create, so it retries once per cooldown window. A budget refusal still consumes the episode, because the budget emits its own alert-budget- episode alert and retrying every window would only add churn. New TestAFailedCreateDoesNotConsumeTheEpisode in test_3130_flood_guard.py: (1) the first create fails and a later window retries → 2 attempts, 1 success, then silent; (2) continuous failure over 6 windows × 5 cycles → exactly 6 attempts; (3) budget-full → 1 budget read in 10 windows. (1) and (2) failed before the fix and pass after it. The _ALLOWED_CALLERS key in test_1677_operator_alert_emitters.py was renamed to the function that now holds the create.
I2 QueueItemDetail lacked queue_flood Fixed queue_flood gets the alert badge (semantic state-autonomous tokens), and the header renders queueTypeLabel instead of the raw type. New jsdom mount spec queueItemDetailType.mount.spec.js failed before the fix (raw queue_flood, no badge) and passes after it.
I3 PG arm syntax-only Left as is Alembic 0089 is still exercised only by pg-migrations on an empty DB. The survivor rule is tested on SQLite only.
I4 AC4 deviation Stated in the PR body AC4 is now marked changed: the migration clears the backlog automatically, and the operator one-action path is still blocked by #2372 (open).

Tests

  • pytest test_3130 + test_1632 + test_1677_budget + test_1677_emitters + test_2915: 205 passed (-p no:randomly). Adding the other create_bounded_alert consumers (test_2392, test_2529, test_ent499): 291 passed.
  • vitest: the new spec plus operatorQueueUnknownType, operatorQueueResponse and the raw-colour, loading-gate and source-text ratchets: 99 passed. Full vitest suite: 257 files, 4580 passed.

🤖 Generated with Claude Code

@vybe vybe left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

merge-train: batch validated on train/20261005-1214 (#3228, all gates green)

@vybe
vybe merged commit 1e77bf3 into dev Oct 5, 2026
29 checks passed
trinity-ability pushed a commit that referenced this pull request Oct 5, 2026
…m ending (#3246)

C3 of #3246 — the DB seam the platform-alert service (C4) will stand on.

db/operator_queue.py:
- `create_platform_item(agent, item, *, subject, max_pending_for_type)`:
  find → touch → count → insert in ONE `_lock_agent_for_create` transaction.
  A reading of a subject with a pending row updates it in place by a
  compare-and-set on `status='pending'` (title / question / priority /
  context with `seen_count`+1 / `last_seen_at` / `expires_at`); a lost CAS
  falls through to a fresh row, so a row a person ended is never
  overwritten. The find runs BEFORE the #1677 per-type count, so an update
  is never refused at budget. An `IntegrityError` from the partial unique
  index (a lock that failed open) re-finds and touches the winner. Returns
  `{"outcome": created|updated|refused_at_budget, "row", "changed"}`;
  `changed` is false on a bare repeat reading so nothing broadcasts.
- `find_pending_by_subject`, `find_person_ended_by_subject` (the seam's
  snooze read, newest person ending at or after `since`).
- `end_items_by_platform(ids, *, reason, batch_id)` — the `bulk_cancel_items`
  shape with `disposed_by='platform'`, NULL email, re-selected by batch id
  so only CAS-won rows come back.
- `mark_expired`'s per-id CAS gains `expires_at < now`: a row refreshed
  between the candidate select and the CAS keeps its new deadline.
- `_insert_values` takes keyword-only `subject` / `last_seen_at`;
  `_row_to_item` / `_SELECT_COLS` carry both.

services/ask_service.py: `clear_platform(ids, *, reason, batch_id=None)`
mirroring `expire()` — no Actor, `PLATFORM_ENDING_REASONS =
(condition_cleared, superseded)` enforced, one `platform_cleared` audit row,
one thin `operator_queue_cancelled` trigger per agent, observers get only
the rows this call won. database.py facade for the four accessors.
`test_ent329_operator_resume.py` G2 list gains the new CAS accessor.

#3130 / #3220: `_own_pending_conds`, the flood guard and
`create_bounded_alert_outcome` are untouched; their suites stay green.

Test (red first — the accessors did not exist):
tests/unit/test_3246_platform_alerts_db.py on the unit island's real
SQLite: two readings → one row, bare repeat → unchanged, person-ended row
never overwritten, budget refuses only with nothing to update, the index
refuses a second pending row, `mark_expired` keeps a refreshed row and still
expires an overdue one, the platform ending skips a person-ended row and
hands observers only CAS-won rows.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

ui PR touches the frontend UI — triggers Playwright e2e tests

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants