fix(operator-queue): the flood guard no longer feeds on itself (#3130) - #3220
Conversation
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>
/review — automated pre-landing review (at
|
| 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 |
|
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_iteminstead ofcreate_bounded_alert: 2 red (test_it_is_created_through_the_budget,test_the_budget_is_the_hard_backstop) - full revert of
src/backendto 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:
returnIssue: 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 usesfunc.substr(...) != pwith 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 intest_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 SQLsubstrcannot 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_pendingandsince; an agent-written marker is overwritten (test_a_marker_the_agent_wrote_is_not_trusted). - Concurrency: the marker rides the existing re-read +
if_matchwrite. 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.pycomment updated, no DDL change. - Enterprise disclosure: 0 guard-pattern hits in the added doc lines.
- Design system: the pill uses the existing
state-autonomoussemantic tokens; the raw-colour ratchet passes. - Docs:
operating-room.md,api-endpoints.md,requirements/security.mdand both prompt copies are updated.
Low confidence (appendix)
- (5/10)
_flood_episodesis 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_RANKputsqueue_fullaboveinvalid_id(:782), so a malformed id is reported to the agent only after the queue drains. Arguablyinvalid_idis 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>
/review follow-ups —
|
| 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 othercreate_bounded_alertconsumers (test_2392,test_2529,test_ent499): 291 passed. - vitest: the new spec plus
operatorQueueUnknownType,operatorQueueResponseand the raw-colour, loading-gate and source-text ratchets: 99 passed. Full vitest suite: 257 files, 4580 passed.
🤖 Generated with Claude Code
…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>
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.db/operator_queue.py::_own_pending_conds: the queue-file cap and the nativeask_operatorqueue_fullcheck. 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.queue_floodtype, created throughcreate_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.platform.ingestionblock (reason:queue_full/rate_limited/invalid_id, plusmax_pendingandsince) rides the existing re-read +if_matchwrite-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).supersede_queue_flood_backlog(SQLite runner + Alembic0089, sharing one SQL statement) keeps the newest pending flood row per agent. The rest becomecancelled, withdisposed_by = 'platform'and reasonsuperseded. Platform-minted rows never wake the agent or touch its file.Acceptance criteria
platform.ingestion), matchingask_operator'squeue_fullrefusal.supersede_queue_flood_backlogmigration on upgrade, and the budget caps future flood rows at 5 pending per agent, so no operator bulk action is needed for the flood class. The operator one-action path is still blocked by Operator queue: acknowledging a platform-minted alert leaves itrespondedforever — the write-back re-injects it into the agent's queue file and the #1631 guard rejects it on every boot #2372 (acknowledging a reserved-prefix alert strands it inresponded), which remains open. The one surviving pre-PR flood row per agent (typealert) still hits Operator queue: acknowledging a platform-minted alert leaves itrespondedforever — the write-back re-injects it into the agent's queue file and the #1631 guard rejects it on every boot #2372 when acknowledged. Operator queue: acknowledging a platform-minted alert leaves itrespondedforever — the write-back re-injects it into the agent's queue file and the #1631 guard rejects it on every boot #2372 is the follow-up; AC4 should not be read as fully closed.Known residual
At its budget (5 pending
queue_floodrows),create_bounded_alertemits its ownalert-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 anexclude_request_id_prefixeskwarg oncount_pending_for_agentandcreate_native_item, threaded through both facades.src/backend/services/operator_queue_service.py: thequeue_floodbudget type, the edge-triggered emitter, the hold reason, theplatform.ingestionmarker 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: thequeue_floodlabel and pill/badge, using semantic tokens. The detail header renders the sharedqueueTypeLabelinstead of the raw type string.create_bounded_alert_outcomereturns the four-outcome result;create_bounded_alertis 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: theplatform.ingestioncontract.operating-room.md,api-endpoints.md,requirements/security.md,tables.py(thedisposed_bycomment), and a learnings fragment.Test Plan
cd tests && pytest unit/test_3130_flood_guard.py: 20 pass. 18 of them failed before the fix, and the marker-default test fails with the old default.test_ent751_gate_callers.py::TestPaidA2A(3), which fail identically on cleandev(local env: the x402/payments pins from feat(a2a): x402 payment gate on the inbound door + price on the agent card (Abilityai/trinity-enterprise#679) #3213).test_ent679_payments_pin_parity.py(7, identical on cleandev).check_alembic_heads: 1 head (0089_…).check_alembic_parity: pass. A fresh SQLiteinit_database()records the migration.operatorQueueUnknownType,operatorQueueResponseand the raw-colour, loading-gate and source-text ratchets: 97 pass.-p randomly --randomly-seed=12345was still running when this was opened.queue-flood-row per agent remains pending and the rest show cancelled/superseded in Resolved.Fixes #3130
🤖 Generated with Claude Code