866332bfb5
* fix(relay): authorize send_message targets and surface egress declines
P5 of the relay egress-authorization workstream. The relay path
authenticated the SENDER but never authorized the DESTINATION, and the
gateway compounded it from both ends.
(a) send_message could silently name an arbitrary relay target. Its
`target` parameter is free-form ('platform:chat_id'), so a model could
name ANY chat id and the gateway would emit an outbound frame for it.
gateway/relay/egress.py adds an attestation floor: a relay-routed
destination must have a provenance this gateway can show -- the
operator's home channel, the channel directory, or its own gateway
session origins. Anything else is refused HERE, with a visible tool
error naming the target, before a frame is written. Non-relay platforms
and platforms served by a live native adapter in this process are
untouched (same precedence resolve_delivery_transport applies).
(b) Connector declines were swallowed into apparent successes. The
connector's egress floor answers an unauthorized destination with a
DEFINITE failure whose text is deliberately uniform (F-005). Several
relay lanes degrade a *transport drop* by design and were degrading an
*authorization refusal* the same way:
- _send_media returned None, sending the caller into
BasePlatformAdapter's text fallback -- a DIFFERENT op re-addressed at
the very chat the connector had just refused.
- _send_prompt returned None, so exec-approval / slash-confirm /
clarify reported "relay prompt op unavailable" (a wrong reason) and
ran their numbered-text fallbacks into the refused chat.
- task_card_stop discarded the error entirely.
- typing / delete / react / thread ops degraded silently at debug.
is_egress_decline() classifies THAT a decline happened (never why --
the uniform text is not parsed for reasons) and requires a definite,
non-ambiguous failure, so a lost-ack retry is still a transport
outcome. Lanes with an error-carrying contract now report the decline
verbatim; cosmetic bool/None lanes still degrade but log it at WARNING.
Advisory progress drops that legitimately degrade are unchanged: the
task_card send lane, the draft ambiguous/except branches, and every
transport-exception path keep their existing fail-open behaviour.
Tests: 21 mutations of the production source, all KILLED.
* fix(relay): authorize the RESOLVED target; declines must not fall back
Review round 1 (independently confirmed by a second reviewer) found three
blockers. Two are fixed here; the third (B-2, Telegram @username) is a policy
decision left open deliberately.
B-1 — THE FIX CAUSED THE OUTAGE IT PREVENTED (tools/send_message_tool.py)
The P5(a) guard ran ABOVE Slack user->DM resolution, so it authorized the
internal pseudo-id `_parse_target_ref` emits (`user_name:ben`, `user:U...`).
Provenances only ever hold RESOLVED conversation ids, so a fully attested DM
was compared as a handle against a set of `D...` ids and refused:
base slack:@ben SENT head(before) slack:@ben REFUSED
Every Slack DM by handle was broken. Moved the guard below resolution; it now
authorizes the destination that is actually sent to, and the refusal names the
resolved id. Position is load-bearing, so it is commented as such and pinned:
reverting the move turns exactly the four new cases red.
B-3 — A DECLINE IS NOT A LANE FAILURE (gateway/run.py)
`_approval_send_outcome` had only sent/failed/ambiguous, so a connector
decline collapsed into `failed` — which is the cue to run the plain-text
fallback into the chat the connector had just refused. The adapter fix in the
previous commit improved the error STRING while user-visible behaviour stayed
identical to base; the commit message overstated it. Fixed properly:
- new `declined` verdict, recognised via the shared `is_egress_decline`
contract (not string sniffing at the call site)
- exec-approval returns without the text fallback
- slash-confirm suppresses the text reply AND clears the registration, so a
card that never rendered cannot capture the user's next message
`send_clarify` was already correct (returns early inside the adapter).
MUTATIONS (production source; both directions)
classifier never returns 'declined' -> KILLED (4 cases)
ALL failures classified as 'declined' -> KILLED (2 cases)
guard moved back above Slack resolution -> KILLED (4 cases)
decline CODE changed (review M05) -> KILLED
marker match made case-sensitive (M10) -> KILLED
M05 was a tautology: the test asserted the imported constant against itself,
so changing the constant could not fail it. The wire contract is now pinned as
a literal, because the connector stamps that exact string and a one-sided
change is a silent cross-repo break.
REGRESSION CHECK: the 12 failures + 1 collection error in this test selection
are PRE-EXISTING cross-test contamination — the identical set fails at
7cf86188ac. Verified by diffing the failing sets: no new failures, 363 -> 374
passed.
NOT FIXED (deliberate): B-2, Telegram `@username`. The Bot API resolves handles
at send time, so there is no id to compare and no canonicalization exists yet.
That is a policy decision, not a code move.
* fix(relay): fail CLOSED on guard faults; classify the structured decline
Third independent review. Two more blockers, both reproduced before fixing.
1. THE GUARD ITSELF FAILED OPEN (tools/send_message_tool.py:158)
`_authorize_relay_target` wrapped BOTH the import and the call in one
`except Exception: return None` — and None means AUTHORIZED at every call site.
So any runtime bug inside the guard silently switched the entire P5(a) boundary
off. Reproduced: with the guard raising, an unattested target sent.
The docstring already stated the correct intent ("must not fail closed on its
own IMPORT error") and the code did something broader. The two failures are not
the same: a missing gateway package means there is no relay egress to
authorize; a fault inside the guard means authorization did not happen. The
import is tolerated, the call is not — a guard that cannot answer refuses.
2. THE STRUCTURED DECLINE WAS THROWN AWAY (gateway/run.py)
The adapter preserves the connector's dict in `SendResult.raw_response`. My
previous commit rebuilt a dict from the error STRING, which loses two
contracts:
* a decline carrying `code: egress_declined` and NO text renders as
"relay egress declined" — no marker colon — so it classified as `failed`,
which is exactly the cue to run the fallback into the refused chat;
* `ambiguous: True` (lost ack) was flattened into a DEFINITE failure,
re-sending a card that may already be on the user's screen. That is the
duplicate-card bug the ambiguous verdict exists to prevent, reintroduced
by the fix meant to harden the same path.
Both call sites now classify `raw_response` when present, ambiguity first, and
fall back to the wire sentence only for connectors that send no structured
response.
I had fixed the text-marker path and tested only the text-marker path. Worth
naming: the review's probe was a shape my tests never produced.
MUTATIONS (production source)
guard fault returns None (fail open again) -> KILLED
classifier ignores raw_response -> KILLED (3 cases)
ambiguous treated as a definite failure -> KILLED (2 cases)
40 focused tests pass. Regression check vs be321faf27: identical 13-item
failing set (pre-existing cross-test contamination), no new failures.
STILL OPEN: B-2 / finding 3, Telegram `@username`. The reviewer is right that
this is a REGRESSION of an existing contract (#53573 added Bot API username
support), not merely an unspecified input, since relay provenance stores the
numeric chat id. Fixing it means resolving the handle before authorization, or
explicitly revoking the contract. That is a policy decision, not a code move,
and it is Ben's call.
* test(relay): pin M21 and M25, the survivors whose comments called them load-bearing
Round-2 review reported six unpinned survivors from round 1. Two guard real
behaviour and are now covered; the other four are cosmetic-lane warnings and
fail-open branches I am leaving documented rather than pretending to close.
M25 — thread-qualified session ids. `_session_ids` adds BOTH "chat:thread" and
the bare chat, because the connector authorizes the CHAT. Without the split a
gateway whose session origin is `-100999:77` cannot send to `-100999`, the chat
it is demonstrably already talking in. KILLED.
M21 — the generic `relay` plane must union every fronted platform, since a
relay session is filed under its LOGICAL platform. KILLED.
MY FIRST M21 TEST WAS THE DEFECT IT WAS TESTING FOR. I patched `_relay_fronted`
— the very function the mutation empties — so emptying it changed nothing the
test could see, and the mutation SURVIVED against a green test. Rewritten to
drive the real `relay_fronted_platforms()` through its env source
(`GATEWAY_RELAY_PLATFORMS`), which is how production learns it.
That is the same "the test verifies my stand-in" failure I have spent this
workstream removing from the connector harnesses, reproduced here in three
lines of Python. The tell was identical: a mutation that survives a test
written specifically to kill it.
334 tests pass.
NOT PINNED, deliberately: M03 (success-guard on a malformed dict), M24
(empty-target allowance — the one fail-open branch, reachable only when the
bare-platform path already resolved a home channel), M35/M36 (decline WARNINGs
on cosmetic lanes). All four are observability or defence-in-depth rather than
authorization, and the review agrees they are non-blocking.
* fix(relay): defer Telegram @username authorization to the connector (B-2)
Closes the last blocker. Two reviewers independently called this a REGRESSION
of the public-channel username support added in #53573, not an unspecified
input, and they were right: provenance stores RESOLVED numeric chat ids, so
comparing `@channel` against them could only ever refuse.
WHY THE GATEWAY CANNOT ANSWER IT. The guard fires only when there is no live
native adapter — i.e. relay-fronted deployments — and on exactly those the
CONNECTOR holds the bot token, not this process. There is no local way to turn
a handle into the numeric id. Refusing here is not "fail closed", it is "fail
always".
WHY DEFERRING IS SAFE. The destination is still authorized one layer out: the
connector's Telegram egress floor (gg#238, merged 743a7c2) classifies and
refuses unauthorized destinations after ITS resolution — the layer that closed
the reported vulnerability in the first place. Handles go from two guards to
one, the authoritative one, not to zero.
The carve-out is deliberately narrow and its EDGES are pinned, because the
failure mode of an exemption is silent widening:
telegram `@handle` -> deferred (the regression case)
telegram numeric id -> still guarded
matrix `@user:server` -> still guarded (telegram-only)
bare name, no `@` -> still guarded
attested handle -> normal path, attestation still consulted
MUTATIONS
carve-out widened to all platforms -> KILLED
carve-out widened to every target -> KILLED
carve-out removed (regression back) -> KILLED
carve-out checked BEFORE attestation -> KILLED
THE ORDERING MUTANT SURVIVED MY FIRST TEST. Both orderings return None, so
asserting the verdict could not tell them apart — the test asserted the claim
instead of the mechanism. Rewritten to observe that attestation is actually
consulted. Same defect class as the M21 test earlier in this branch: a
mutation surviving a test written specifically to kill it means the test is
measuring the wrong thing.
341 tests pass.
FOLLOW-UP (option 2, Ben's call, deliberately NOT done here): resolve the
handle before authorizing so BOTH layers apply. That needs a resolution
round-trip through the connector — new wire surface — so it belongs in its own
phase rather than bolted onto this one. Recorded in the code comment at the
carve-out, not just here.
* fix(relay): close two fail-open boundaries; test the code-only decline for real
Both blockers from review, each REPRODUCED before fixing.
1. STRUCTURED DECLINE HAD NO GUARD. Deleting `raw_response=result` from both
`_send_prompt` return branches left all 34 tests green — a surviving,
non-equivalent security mutant. The `code` field is the documented
PREFERRED signal precisely because a connector may send no prose, and a
caller rebuilding `{"success": False, "error": ...}` cannot see it.
Cause: every existing case declines with marker TEXT. The evidence for the
code-only path was a hand-built SimpleNamespace in a different file — a
stand-in for the adapter, so it verified my fixture instead of production.
Fixed with a CodeOnlyDecliningConnector driving the real
`send_exec_approval` -> `_send_prompt`, feeding the REAL SendResult to the
REAL `_approval_send_outcome`, plus the same shape on the media lane.
drop raw_response SURVIVED (34 passed) -> KILLED
2. TWO FAIL-OPEN BOUNDARIES, both "absence" and "fault" sharing a return.
`_relay_fronted` swallowed EVERY exception and returned an empty set, which
`relay_routed_platform` reads as "not relay-routed" — skipping the guard.
Probe, with a positive control in the same run:
positive_control_denied = True
discovery_fault_denied = False <- unattested target AUTHORIZED
`_authorize_relay_target` caught every exception during IMPORT as "no
gateway package". A module that exists and fails to initialize is a fault,
not an absence, and returning None there means authorized.
Now: ImportError alone is absence; anything else raises RelayRouteUnknown
and `authorize_relay_target` converts it to a REFUSAL STRING (not a raised
exception — every caller treats the return value as the verdict, so raising
would trade a fail-open for a crash).
Kept the converse under test so "fail closed" does not silently become
"refuse everything in CLI/cron", which is the outage the broad except
existed to prevent.
discovery fault -> empty set KILLED
RelayRouteUnknown -> authorized KILLED
import fault -> authorized KILLED
397 passed (was 392, +5 new cases), zero failures.
* fix(relay): close all seven review-round-3 blockers
Every finding reproduced before fixing; every fix mutation-checked after.
CONTENT LEAKS (the decline was laundered into a different op, same chat)
#1 A declined DRAFT SEAL replayed as a plain send. On stream-is-the-message
platforms the turn-final becomes draft(final=True); `_seal_open_draft`
dropped the structured body, so `_absorb_into_open_draft` read a REFUSAL as
a lane failure and fell through. Probe, Slack descriptor:
before: draft(partial) -> draft(final,SECRET) -> send(SECRET)
after: draft(partial) -> draft(final,SECRET)
My first probe of this used a discord descriptor and showed no seal at all —
the leak is real, my probe was wrong (streams only arm for Slack).
#6 Task-card PROGRESS had the same defect one lane over: a bare failed
SendResult reads as "card lane unavailable", and TurnRunner then sends the
task text to the same chat. Both card methods now carry raw_response and
the caller suppresses the fallback on a decline.
AUTHORIZATION BYPASSES
#2 `except ImportError` was NOT the fix I claimed last round. ImportError also
covers a broken dependency inside an INSTALLED gateway; review probed
`ImportError.name = "gateway.relay.dependency"` and got an authorized
verdict. Now only a name identifying the gateway relay module itself is
absence. An ImportError with NO name stays absence — refusing on a fault we
cannot attribute would trade an unidentifiable bug for a real CLI/cron
outage, and an existing test caught exactly that when I first got it wrong.
#3 `relay_routed_platform` lowercases the requested platform; `_relay_fronted`
returned configured names verbatim. A platform configured as "Discord"
missed the membership test, looked native, and skipped the guard:
'discord' => refused 'Discord' => ALLOWED 'DISCORD' => ALLOWED
An attestation bypass on a string comparison.
UNDELIVERABLE PROMPTS THAT HUNG
#4 `_clarify_send_disposition` handled `failed` and `ambiguous` but not
`declined`, so a REFUSED clarify card fell through to wait_for_response and
blocked until clarify_timeout — indefinitely when configured non-positive.
A decline is more definitive than a failure, not less.
#5 The exec-approval decline branch returned quietly, which suppressed the text
fallback (right) but left the CENTRAL approval entry pending (wrong) — the
dangerous command stayed blocked until the approval timeout. My comment
claimed the registration was torn down; only RelayAdapter's private map was.
It now raises `_ExecApprovalDeclined`, which propagates to
`_await_gateway_decision`'s existing notify-failure path (drops the entry,
unblocks the tool). A dedicated type, re-raised past the local
`except Exception` that would otherwise have restored the leak.
#7 THE GAP THAT LET ALL OF THIS SHIP. Both caller-level suppressions were
unfalsifiable: deleting either branch left 36/38 tests green. The suites
drove `_approval_send_outcome` and `RelayAdapter` but never the real
TurnRunner / busy-session callers, so nothing observed whether a text send
FOLLOWED a decline — which is the whole property.
tests/gateway/test_decline_fallback_suppression.py drives both real callers
and records every send. Each decline case is paired with an ordinary-FAILURE
control, because without one a caller that never falls back would also pass.
MUTATIONS (all on production source, anchors count-checked, restored after)
#1 seal decline -> plain send KILLED
#1b seal drops raw_response KILLED
#2 nested ImportError -> authorized KILLED
#3 fronted set not normalized KILLED
#4 clarify declined branch removed KILLED
#5 approval decline returns not raises KILLED
#6 task_card drops raw_response KILLED
#7 slash-confirm suppression removed KILLED
#7's two were the reviewer's SURVIVORS (36/38 passing); both now die.
425 passed, zero failures.
* fix(relay): close the three round-4 blockers
Round 4 confirmed six of seven round-3 fixes and found three more. Each
reproduced before fixing, each mutation-checked after.
1. A NAMELESS ImportError still authorized. Last round I admitted it as
"absence" to protect the CLI/cron path. That reasoning was WRONG and the
interpreter says so:
import gateway.relay.nope -> ModuleNotFoundError, name="gateway.relay.nope"
import totally_absent_pkg -> ModuleNotFoundError, name="totally_absent_pkg"
Genuine absence is ALWAYS ModuleNotFoundError with `.name` set, so the
CLI/cron path never produces a bare ImportError and nothing legitimate was
being protected. A plain or nameless ImportError comes from an import hook
or a module that failed while initializing — an unattributable FAULT.
Now: absence is ModuleNotFoundError naming gateway / gateway.relay /
gateway.relay.egress; everything else refuses. Two existing tests raised a
bare ImportError to simulate absence and were corrected to the real shape.
2. SESSION ATTESTATION INVENTED IDS. `_session_ids` split every id on the first
colon to recover "chat" from "chat:thread". Matrix ids contain a colon
natively, so `!room:server.org` attested a bare `!room` — the guard
vouching for a destination on its own fabrication. The split now applies
only to platforms whose ids genuinely carry a `:thread` suffix (allow-list;
unknown platforms are treated as un-splittable, which can only refuse more).
Kept a Slack control: dropping the split entirely would refuse legitimate
thread replies, which is the outage the split exists to prevent.
3. THE TASK-CARD FIX WAS UNFALSIFIABLE — my own round-3 mistake, and the same
one round 3 caught me making. I added the production branch AND a test, but
the test stopped at RelayAdapter: it proved `raw_response` is carried and
never called `TurnRunner._task_card_publish`, which owns the property.
Deleting the real branch left 30 tests green. Now driven through the real
caller, with an ordinary-failure control.
The lesson generalises: proving the DATA reaches the boundary is not proving
the CALLER acts on it. Every one of these decline fixes has two halves and
the second half is where the security lives.
Also closed the round-4 non-blocking finding: `gateway/relay/egress.py` has its
OWN import boundary, and the existing test intercepted the earlier import in
tools/send_message_tool.py, so it was never exercised. Mutating that classifier
to treat every ImportError as absence now dies.
MUTATIONS (production source, anchors count-checked, restored after)
R4-1 nameless ImportError -> authorized KILLED
R4-2 session split unconditional KILLED
R4-3 task-card caller branch removed KILLED (was SURVIVED)
egress classifier: any ImportError = absence KILLED
Also probed and found NOT a leak: a refused OPENING draft frame disarms the
stream and the turn-final goes out via `send`. That send is itself guarded and
the connector refuses it too, so no content is delivered — unlike the seal case
(round 3, #1) where the seal was the only check on that path.
452 passed, zero failures.
* fix(relay): recover the thread parent from thread_id, not a colon split
Round 4 blocker 2 was closed with an allow-list of platforms whose ids have no
native colon. Reviewing my own fix while round 5 ran, the allow-list is the
wrong mechanism: it NARROWS a guess instead of removing it, and it still gets
Matrix wrong the moment a Matrix session is thread-qualified
(`!room:server.org:$thr` -> split yields `!room`).
The structured field was there all along. `_session_entry_id` composes the id
as f"{chat_id}:{thread_id}" and the entry still carries `thread_id`
separately, so the parent is knowable EXACTLY: strip the known suffix, or add
nothing. No platform list, no guessing, correct for ids that contain colons.
Mutations:
back to splitting on the first colon KILLED
thread parent never recovered (over-refuse) KILLED
Both directions matter: the first invents attestations, the second refuses
legitimate thread replies.
One existing test (M25) asserted the right PROPERTY with a fixture that omitted
`thread_id` — a shape real entries never have. Fixture corrected, assertions
untouched.
453 passed.
* fix(relay): close the four round-5 blockers
Each reproduced before fixing, each mutation-checked after.
R5-1 A DISABLED NATIVE ADAPTER BYPASSED AUTHORIZATION. `_has_live_native_adapter`
treated any entry in the adapter map as native; `resolve_delivery_transport`
ignores a native adapter whose config is disabled and routes over Relay.
Two independent routing classifiers, disagreeing:
guard says native: True delivery routes relay: True
So the guard skipped authorization for a send that went over the relay.
The guard now applies the router's enabled-state rule; probed both
configurations and they agree.
R5-2 THREAD IDS WERE NEVER AUTHORIZED. The parser splits chat_id and thread_id;
only chat_id reached the guard. On Discord the thread IS the destination —
`POST /channels/{thread_id}/messages` — so an attested parent channel
authorized an arbitrary caller-supplied thread. `authorize_relay_target`
now takes thread_id and requires its own attestation (bare id or the
`chat:thread` form a session origin produces); both call sites forward it.
R5-3 A DECLINED **INITIAL** DRAFT WAS RETRIED AS A PLAIN SEND. Round 3 fixed the
declined SEAL; the declined OPEN was a different path. `send_draft`
returned a bare failure, so the stream consumer read "draft transport
unusable", disabled drafts and fell through to `_first_send`. Measured
through the real adapter and real StreamTransportMixin:
before: ops ['draft', 'send'] after: ops ['draft']
send_draft now carries raw_response; a decline is terminal for the run and
the guard sits in `_first_send`, where every fallback path converges.
R5-4 MY ROUND-4 TASK-CARD FIX SUPPRESSED EXACTLY ONE UPDATE. It set
`native_failed`, which the entry gate already uses for an ordinary broken
lane, so the next progress event skipped the decline branch and went
straight to the text fallback:
after first publish: [] after second: ['send']
Terminal declines are now a separate `egress_declined` state checked at the
entry gate. A refusal does not expire after one tick.
MUTATIONS
R5-1 disabled native counts as native KILLED
R5-2 thread_id not authorized KILLED
R5-2b tool does not forward thread_id KILLED (was SURVIVED)
R5-3 initial-draft decline not terminal KILLED
R5-3b _first_send guard removed KILLED
R5-4 declined state not persistent KILLED
R5-2b is the same gap that produced findings 3 and 4 of the last two rounds, a
third time: every test called `authorize_relay_target` directly, so dropping the
argument from the TOOL WRAPPER changed nothing. Testing the callee never proves
the caller uses it — now pinned explicitly.
Each fix ships with an ordinary-failure control, because every one of these
makes the guard refuse MORE, and over-refusal is now the larger risk.
474 passed, zero failures.
* refactor(relay): declare the terminal-decline state where it lives
Both terminal-decline flags were set dynamically. They worked (neither class is
frozen or slotted) but an undeclared attribute hides the state from anyone
reading the class, and this one is security-relevant.
_TaskCardState.egress_declined — declared dataclass field
StreamConsumer._egress_declined — initialised in __init__
Lifetime verified while checking whether a refusal can leak ACROSS turns and
mute a healthy destination: it cannot. _TaskCardState is constructed per
progress-drain (run_turn_runner.py:420) and the consumer's flags per run
(stream_consumer.py:163), so both are fresh each turn.
Also verified the guard's blast radius after adding thread authorization: the
ONLY callers of authorize_relay_target are the two model-facing send_message
call sites. Gateway-internal sends — notably the handoff path, which creates a
thread and immediately posts to it with no session provenance yet — go through
transport.adapter directly and are unaffected. That was the most plausible
over-refusal, and it does not reach this guard.
461 passed.
* fix(relay): close the four round-6 blockers — the edit lane
R6-1 MY OWN R5-1 FIX REINTRODUCED THE BYPASS IT CLOSED. I wrote
`except Exception: return True` around the config lookup, so a config read
fault declared the platform native while the ROUTER, reading the real
config, sends over the relay:
guard_has_live_native True guard_verdict None router relay
Routing we cannot determine is UNKNOWN. It now raises RelayRouteUnknown,
which the outer handler must re-raise rather than flatten to False, and
`authorize_relay_target` turns into a refusal. This is the second time a
convenience `except` in this function created a bypass; there is now no
permissive return left in it.
R6-2/3/4 THE NINTH LANE: `edit`. ONE dropped field, THREE leaks.
`RelayAdapter.edit_message` discarded the connector response, and three
independent callers read a bare edit failure as "editing is unavailable"
and re-send the content as a NEW message to the same chat:
stream edit fallback ['edit', 'edit', 'send'] the unseen tail
queued reconciliation ['edit', 'send'] the WHOLE response
task-card fallback ['edit', 'send'] the task text again
Fixed at the source (edit_message carries raw_response) plus each caller:
`_on_edit_failure` — the single funnel for stream edit failures — makes a
decline terminal for the run, `_send_fallback_final` refuses to deliver a
continuation after one, the queued reconciler returns instead of sending,
and the task-card fallback sets the same terminal state R5-4 introduced.
R5-4 fixed the native task-card op and I did not check its sibling
fallback path. The pattern across rounds 3-6 is consistent: the fix goes
where the decline is OBSERVED, and the leak lives wherever someone else
later decides to retry.
MUTATIONS
R6-1 config fault -> assume native KILLED
R6-1b RelayRouteUnknown swallowed as False KILLED
R6-2 edit drops raw_response KILLED
R6-2b edit-failure decline not terminal KILLED
R6-3 queued reconcile falls back on decline KILLED
R6-4 task-card fallback edit decline KILLED
Each with an ordinary-failure control: a genuinely un-editable message must
still be delivered, and a broken card lane must still reach the user.
481 passed, zero failures.
* fix(relay): add a terminal-decline latch at the adapter choke point
THE STRUCTURAL FIX, not a twelfth local check.
Rounds 3-6 of review found ONE defect in eleven lanes: the connector refuses an
op, and some caller downstream reads that as 'this lane is unavailable' and
retries the same content through a DIFFERENT op against the SAME chat. Media,
prompt, draft-open, draft-seal, native task card, task-card fallback edit,
slash-confirm, exec-approval, clarify, stream edit, queued reconciliation.
Each was closed by adding a check at one more call site. That approach cannot
converge: gateway/ has ~60 outbound call sites, every one of them a place a
future change can reintroduce this, and four consecutive review rounds each
found another. The reviewer's own count of lanes is the argument against the
per-site design.
Every relay frame from every one of those callers passes through
_transport.send_outbound. One latch there covers them all: once the connector
refuses a chat, this adapter stops emitting CONTENT frames for that chat.
Proven to subsume the local checks: with the stream-edit per-site check
DISABLED, the leak probe still reports blocked=true — the frame never reaches
the wire. The local checks stay as defence in depth and for their better error
messages, but they are no longer the only thing standing between a decline and
a re-addressed send.
Scope is deliberately narrow, and each limit is mutation-pinned:
per CHAT - a refusal must not mute other conversations
CONTENT ops - typing/delete carry nothing; latching them would leave a
stuck typing indicator for no security gain
self-healing - cleared when the connector accepts that chat again, so a
transient policy change does not need a restart
Mutations:
latch never set KILLED
latch never consulted KILLED
latch is global, not per-chat KILLED
latch never clears KILLED
485 passed.
* fix(relay): one route source; the latch already covered round 7's lanes
Round 7 reviewed 573e41e294 — one commit BEFORE the terminal-decline latch —
and independently reached the same conclusion I had: 'The per-call-site
approach is structurally wrong. Use one turn-scoped choke point.' That is the
latch in 6dbc004594.
Its four 'still broken' lanes (tool-progress edit, progress-overflow edit,
long-running heartbeat edit, stale streamed-final reconciliation) all share the
shape edit_message->declined->adapter.send(same chat, same content), and NONE
has a local check. Probed all four against the latch:
tool_progress ops ['edit'] blocked
progress_overflow ops ['edit'] blocked
heartbeat ops ['edit'] blocked
stale_final ops ['edit'] blocked
That is the argument for the choke point, measured: lanes nobody patched are
safe anyway. Pinned by a parametrized test named for those four lanes.
R7-1 IS A REAL BYPASS THE LATCH DOES NOT COVER, and it is fixed here. The guard
rebuilt routing from GATEWAY_RELAY_PLATFORMS while resolve_delivery_transport
asks the CONNECTED adapter (fronts_platform, from the handshake identity set).
Different snapshots: with env discovery stale or momentarily empty, the guard
said 'native' and the router sent over the relay, skipping authorization.
before: guard_relay_routed False / delivery relay
after: guard_relay_routed True / delivery relay / unattested target refused
The guard now asks the live adapter first and falls back to config only when
there is no runner (CLI/cron) — pinned in both directions.
R7-5 (non-blocking, and a fair hit): my stream-fallback test asserted
_egress_declined and never drove _send_fallback_final, so removing that early
return SURVIVED. The test now calls the real fallback and asserts the wire is
untouched; the mutation dies.
Mutations:
R7-1 guard ignores the live adapter KILLED (was SURVIVED)
R7-5 fallback early return removed KILLED (was SURVIVED)
latch not consulted KILLED
491 passed.
* fix(relay): close three holes found by attacking my own latch
Round 8's brief told the reviewer to attack the latch. I did the same in
parallel and found three real holes in it before the review returned.
1. send_for_platform BYPASSED THE LATCH ENTIRELY. It builds and posts its frame
directly rather than through _outbound — and it is the delivery resolver's
OWN entry point, so it is the single most important caller.
before: ops ['edit', 'send'] after: ops ['edit']
gateway/AGENTS.md states the rule I had just broken: 'Seal-interception
exists at BOTH egress doors (send() and send_for_platform()); a new egress
door needs the same two checks.' The latch is a third such check and I had
wired it to one door.
2. A COSMETIC SUCCESS CLEARED THE LATCH. Clearing on ANY success meant a
typing indicator — routinely allowed for a chat whose content is refused —
re-opened the door for the very next send:
ops ['edit', 'typing', 'send']
Only a CONTENT op the connector accepted may clear it now.
3. A THREAD INSIDE A REFUSED CHAT WAS NOT COVERED. A thread lives inside its
parent, so the same content reached the same conversation one level down:
ops ['edit', 'send']
The latch key now strips the thread suffix.
Also normalised int/str chat ids (callers pass both; a type mismatch would
silently unlatch).
MUTATIONS
send_for_platform not latched KILLED
cosmetic success clears the latch KILLED
thread suffix not stripped KILLED
draft-seal retry not latched SURVIVED — EQUIVALENT, proven:
is unreachable while latched (a declined edit before the seal
produces ZERO seal frames, measured). Kept as defence in depth because it
posts directly, and documented at the site rather than covered by a
test that could not fail.
One self-inflicted bug on the way: a blanket replace put 1Password CLI brings 1Password to your terminal.
Turn on the 1Password app integration and sign in to get started. Run
'op signin --help' to learn more.
For more help, read our documentation:
https://www.1password.dev/cli
1Password CLI is built using open-source software. View our credits and
licenses:
https://downloads.1password.com/op/credits/stable/credits.html
Usage: op [command] [flags]
Management Commands:
account Manage your locally configured 1Password accounts
connect Manage Connect server instances and tokens in your 1Password account
document Perform CRUD operations on Document items in your vaults
events-api Manage Events API integrations in your 1Password account
group Manage the groups in your 1Password account
item Perform CRUD operations on the 1Password items in your vaults
plugin Manage the shell plugins you use to authenticate third-party CLIs
service-account Manage service accounts
user Manage users within this 1Password account
vault Manage permissions and perform CRUD operations on your 1Password vaults
Commands:
completion Generate shell completion information
inject Inject secrets into a config file
read Read a secret reference
run Pass secrets as environment variables to a process
signin Sign in to a 1Password account
signout Sign out of a 1Password account
update Check for and download updates.
whoami Get information about a signed-in account
Global Flags:
--account account Select the account to execute the command by account shorthand, sign-in address, account ID, or user ID. For a list
of available accounts, run 'op account list'. Can be set as the OP_ACCOUNT environment variable.
--cache Store and use cached information. Caching is enabled by default on UNIX-like systems. Caching is not available on
Windows. Options: true, false. Can also be set with the OP_CACHE environment variable. (default true)
--config directory Use this configuration directory.
--debug Enable debug mode. Can also be enabled by setting the OP_DEBUG environment variable to true.
--encoding type Use this character encoding type. Default: UTF-8. Supported: SHIFT_JIS, gbk.
--format string Use this output format. Can be 'human-readable' or 'json'. Can be set as the OP_FORMAT environment variable.
(default "human-readable")
-h, --help Get help for op.
--iso-timestamps Format timestamps according to ISO 8601 / RFC 3339. Can be set as the OP_ISO_TIMESTAMPS environment variable.
--no-color Print output without color.
--session token Authenticate with this session token. 1Password CLI outputs session tokens for successful 'op signin' commands when
1Password app integration is not enabled.
-v, --version version for op
Run 'op [command] --help' for more information on the command. into
send_for_platform, which has no such variable. Two existing unfurl tests caught
it — NameError at adapter.py:1407.
504 passed.
* fix(relay): Telegram handle exemption + a turn boundary for the latch
Round 8 blockers. Two of its four were already closed by 93750e351a (it
reviewed the commit before it); these two are real and both are mine.
B1 — THE TELEGRAM @HANDLE EXEMPTION COVERED A NATIVE SEND.
_is_unresolved_handle exempts telegram @handles from attestation because
"the connector resolves and authorizes it". That justification is FALSE
whenever the gateway holds its own token: _send_to_platform calls
_send_telegram(pconfig.token, ...) directly and no connector is involved.
So an unattested @handle went out under the gateway's own credential
while the numeric control was correctly refused.
The exemption now requires that no native credential exists. A probe
fault WITHDRAWS the exemption (falls back to the ordinary attestation
check) rather than granting it.
Shipped with the converse control: relay-only config still exempts
@handles, and numeric targets stay guarded in both modes.
B4 — THE LATCH HAD NO BOUNDARY, SO IT WAS AN OUTAGE MECHANISM.
My own regression, and worse than reported. Removing "clear on cosmetic
success" (correctly) removed the ONLY way the latch could ever clear: a
content op can never reach the connector to succeed, because the latch
blocks it locally first. A refusal at 09:00 muted that chat forever.
A new inbound message for a chat is the generation marker — the natural
teardown point. Suppression still holds for the whole turn.
same_turn_blocked: true next_turn_delivered: true
MUTATIONS (all killed)
handle exemption ignores native credential
native-credential fault GRANTS the exemption
no turn boundary (latch never clears)
teardown clears ALL chats not just this one
teardown ignores the chat
The last two SURVIVED first: I tested _clear_declined_for_turn directly
and never proved _on_inbound calls it — the caller-level gap that has now
produced four blockers on this branch. Added a test driving the real
inbound entry point.
One self-inflicted bug, caught by my own fault test: the probe imported
load_config, which does not exist (it is load_gateway_config), so it
always threw and returned the fault default. The test that pinned fault
behaviour is what exposed it.
510 passed.
* fix(relay): correct latch identity and boundary; one config snapshot
Round 9, four blockers, all reproduced.
B1+B4 — THE TEARDOWN WAS AT THE WRONG PLACE, twice over.
It sat on the adapter's raw _on_inbound, which runs BEFORE profile
routing, the ignored-channel guard, plugin hooks and user authorization.
An unauthorized or dropped event could therefore clear a refusal
belonging to an active turn, and stale content then went out as a
different op. The same placement missed Discord interaction passthrough,
which builds its own MessageEvent and calls handle_message directly, so
slash commands and modal submits stayed muted after an earlier decline.
Both are one mistake: I picked a lane instead of a boundary. Teardown now
runs immediately after _hm_admit_event, the single admission gate every
entry path shares.
dropped event -> latch survives, stale send blocked
admitted event -> latch clears
B2 — THE LATCH KEY SPLIT ON ':', WHICH IS A MISTAKE I ALREADY FIXED ONCE.
_latch_key did str(chat_id).split(":", 1)[0], so !room:tenant-a and
!room:tenant-b both keyed !room: a decline in one Matrix room muted
another, and inbound from one cleared the other's refusal. egress.py
::_session_ids stopped doing exactly this in round 4 and I reintroduced
it three rounds later.
Parent identity is never recoverable from identifier TEXT. Thread
coverage is now structural: _thread_parent looks the relationship up in
the recorded auto-thread map.
B3 — AUTHORIZATION AND DISPATCH USED DIFFERENT CONFIG SNAPSHOTS.
_handle_send retains one pconfig; the guard independently reloaded
config. Across a transition the authorization snapshot could see a
connector-only setup (exemption granted) while dispatch still held the
native token and sent the unattested @handle itself. The guard now takes
native_token from the SAME snapshot dispatch will use. A caller that
omits it does not silently look like "no token".
NB-1/2/3 also closed: real-object snapshot tests, an exception shield
that faces a real exception, and send_follow_up no longer discards the
connector's verdict (that discard is exactly how the edit lane laundered
declines).
MUTATIONS (all killed)
latch key splits on colon again
thread parent lookup disabled
dispatch token ignored by guard
tool drops the snapshot token
admission teardown removed
teardown moved BEFORE admission
exception shield removed
follow_up drops raw_response
"admission teardown removed" SURVIVED first: I had tested the helper, not
_handle_message. Added a test driving production _handle_message with
admission stubbed both ways. Fifth caller-level gap on this branch.
One self-inflicted bug caught before commit: I passed pconfig.token in
_handle_react, which has no pconfig — a NameError on every reaction.
516 passed.
* docs(relay): pin the latch's thread coverage limit as a deliberate trade
_thread_parent only sees connector auto-threads, and that map is capped at
256 entries, so a user-created or evicted thread does not inherit its
parent's latch. Documented at the site and asserted by a test, because the
alternative - deriving parents from identifier text - is exactly what muted
unrelated Matrix rooms in round 9.
The primary control is unaffected: authorize_relay_target takes thread_id as
part of the destination and attests it on every send (6 thread tests).
* refactor(relay): one SendResult decline classifier for all 8 gateway lanes
The extraction found a DEFECT, not just repetition.
Eight gateway lanes each hand-rolled the unwrapping of a decline from a
SendResult, and they did not agree. Six checked only raw_response. Two
also checked the error text. A connector that answers with the uniform
decline SENTENCE and no structured code - the documented contract for
older connectors, per _approval_send_outcome - was therefore classified
as an ordinary failure by those six lanes, so each treated a refusal as
"editing unavailable" and retried through another op.
Measured:
text-only decline six-site check False two-site check True
structured decline six-site check True two-site check True
No content leaked, because the adapter latch classifies the transport
dict directly and catches both shapes (verified: text-only decline still
latches C1 and keeps SECRET off the wire). The cost was wrong verdicts
and futile retries, not disclosure.
declined_send(result) in gateway/relay/egress.py now owns this. It checks
raw_response when structured, else the error text, and preserves the
ambiguous exclusion - an ambiguous result is a transport outcome, so it
must never read as a refusal.
run.py keeps its own shape deliberately: that lane has three verdicts
(ambiguous / declined / failed), so it checks ambiguous first and then
delegates the boolean.
MUTATIONS (all killed)
helper drops the text-only branch
helper drops the structured branch
ambiguous no longer excluded
draft lane decline check removed
edit-failure lane decline check removed
prompt verdict lane check removed
slash-confirm lane check removed
draft lane goes terminal on ANY failure (over-refusal direction)
"draft lane decline check removed" SURVIVED first: _send_draft_frame had
no test driving an unsuccessful send_draft at all. Added one, with an
ordinary-failure control so the fix cannot silently become "one flaky
frame mutes the chat". A non-unique anchor also masked the edit-failure
lane on the first pass - the trap my own skill warns about.
This closes the duplication that caused four of nine rounds of blockers:
a new lane now calls one classifier instead of copying three lines.
519 passed.
* fix(relay): latch identity, new-turn boundary, seal arming, ambiguity
Round 10, four blockers, each reproduced before fixing. Two are my own
regressions from the previous two rounds.
B1 - ADMISSION IS NOT A NEW-TURN BOUNDARY.
Round 9 moved teardown to just after _hm_admit_event. That is only an
ADMISSION gate: an authorized message can be steered into a running
session, answer a pending prompt, run a busy slash command, or be refused
by the pause/drain gates - all without starting a turn. Each of those
cleared the ACTIVE turn's refusal, and a later fallback from that turn
reached the wire (probe: latch emptied, wire ops ['edit', 'send']).
Teardown now runs after _claim_active_session_slot, the first point the
runner OWNS a new turn. The new test drives production _handle_message
through all four non-turn lanes plus the real new-turn path.
B2 - LATCH IDENTITY OMITTED THE LOGICAL PLATFORM.
One relay adapter fronts several platforms, so native ids collide. A
Discord refusal for chat 42 was cleared by clear_egress_latch("telegram",
"42") - the method took a platform and ignored it - and the Discord
fallback then reached the connector. Keyed by normalized platform plus
exact chat id; thread-parent expansion keeps the platform component.
B3 - THE DIRECT DRAFT-SEAL PATH DID NOT ARM THE LATCH.
_seal_open_draft posts through _attempt directly rather than _outbound,
so a definite decline logged and returned but never latched. The
immediate plain-send fallback was suppressed by the caller's own check;
later same-turn sends were not (wire ['draft', 'draft', 'send'], the
third frame carrying refused content).
B4 - MY OWN REFACTOR MADE AMBIGUOUS RESULTS TERMINAL.
send_draft's ambiguous projection discarded raw_response, so
declined_send fell through to the error-text branch - and an ambiguous
result whose text carries the decline marker ("... egress declined: ack
lost") read as a DEFINITE refusal and terminated the run. Ambiguous means
the frame may well have been delivered: a transport outcome, never an
authorization one.
Fixed on both layers: the projection carries the body (and the seal's
ambiguous return is now explicit too), and declined_send's text-only
branch - which cannot see the ambiguous flag - treats ack-lost text as
transport ambiguity. Audited every SendResult projection in adapter.py
for the same shape.
MUTATIONS (all killed)
latch key drops the platform
clear_egress_latch ignores platform
draft seal does not arm the latch
ambiguous projection drops raw body
declined_send infers decline from ack-lost text
teardown back at admission
523 passed.
* refactor(relay): split the terminal-decline latch out of the guard PR
The latch moves to feat/p5-egress-decline-latch (pushed at 3cf45736d7,
which retains the full history) for redesign. This PR keeps the
authorization guard and the per-site decline checks.
WHY. Across eleven review rounds the two halves behaved very differently.
The guard is a PURE FUNCTION of the destination - its blockers were all
"you asked the wrong question" (case sensitivity, nested ImportError,
missing thread_id, config snapshot skew), each a one-line correction that
then stayed fixed. Rounds 7-10 found nothing new in it.
The latch is MUTABLE STATE WITH A LIFETIME living on RelayAdapter - an
object registered once per process that holds the WebSocket and has no
concept of a turn. Nine of its blockers reduce to three questions the
adapter cannot answer: when does it end, who arms it, what is it keyed
on. Every answer so far has been a proxy (a successful op, an inbound
message, an admitted event, a claimed session slot) and every proxy was
wrong in a lane found later.
The per-site checks hold identical information on `st` - a PER-TURN
object - and have produced zero blockers, because the state dies with the
turn and nobody has to decide when it ends.
The no-relaunder property does NOT depend on the latch. Measured on the
real consumer path with the latch absent: a declined draft frame sets
_egress_declined and puts nothing on the wire.
Removal verified structurally rather than by eye: an AST diff of every
symbol between HEAD and this tree reports only latch symbols gone,
nothing added. That check caught two over-deletions my strip made -
_on_inbound (consumed by a "next def" boundary) and _SEEN_INBOUND_MAX
(a class constant inside the removed span). Both restored; 19 failures
went to 0.
ALSO: RESTORED A TEST I WRONGLY REPORTED AS PASSING.
test_tool_guard_forwards_thread_id never made it into the repo - `git log
-S` finds it in no commit - though round 5 recorded its mutant as killed.
Dropping thread_id from the guard call therefore survived the entire
tests/tools suite (146 passed). Written properly this time, driving the
real _handle_send far enough to reach the guard. It now KILLS that
mutant.
MUTATIONS on this tree
guard fault authorizes instead of refusing KILLED
thread_id dropped from the guard call KILLED (was SURVIVED)
handle exemption ignores native credential KILLED
draft lane decline check removed KILLED
prompt verdict lane check removed KILLED
slash-confirm lane check removed KILLED
503 passed.
* test(relay): close the phantom-coverage gaps the guard audit found
The thread_id test that was reported as killing a round-5 mutant turned
out never to have been committed. That is a reason to distrust the other
claimed kills, so I re-ran every guard mutation against the COMMITTED
tree instead of trusting the earlier reports.
Result: 9 of 11 killed, and the two "SKIPPED" ones had non-unique
anchors hiding SIX separate sites. Mutating those individually found
three real survivors.
CASE NORMALISATION (round 3, finding 3) WAS HALF-COVERED.
test_relay_fronted_matching_is_case_insensitive varies the CONFIGURED
name but always requests lowercase "discord", so it pins _relay_fronted's
normalisation and nothing else. The REQUESTED name's `.lower()` was
covered by nothing at all. Probe with it removed:
relay_routed("Discord") -> False
authorize("Discord", unattested) -> AUTHORIZED
which is exactly the bypass round 3 reported, alive again and untested.
Two further sites were untested in the OVER-REFUSAL direction: the
attested store is keyed lowercase, so a mixed-case request missed its own
attested set and refused legitimate traffic. attested_relay_targets' own
normalisation was invisible to every existing test because they all
monkeypatch that function away; it is now asserted against the real
function with only its leaf sources stubbed.
Three tests added. All six case sites now die when mutated.
I also re-did the three fail-closed RelayRouteUnknown mutations properly.
The first pass swapped whole lines and produced IndentationErrors, so
"KILLED" there proved nothing but a syntax error. Neutralising each raise
at correct indentation: all three genuinely KILLED.
FINAL AUDIT ON THIS TREE — 17 mutations, zero survivors
guard: thread_id dropped at the call site
guard: react path unguarded
guard: handle exemption ignores native credential
guard: 3x fail-closed raise neutralised
guard: 6x case-normalisation site
classifier: ambiguous treated as a decline
classifier: text-only decline branch removed
lane: draft / stream-edit / prompt / slash-confirm checks removed
511 passed.
* test(relay): make the stream-edit test fail for the right reason
Review of 45835a282d raised one blocking issue and three non-blocking
ones. All four are addressed; none was a production defect.
BLOCKING — the stream-edit test failed on the double, not on a leak.
test_declined_stream_edit_does_not_send_the_unseen_tail implemented only
the GUARDED path in its consumer double. Removing either guard therefore
raised AttributeError inside the fake before any send could be observed:
guard 1 removed -> AttributeError: no attribute '_is_flood_error'
guard 2 removed -> AttributeError: no attribute '_clean_for_display'
Red, but for the wrong reason — the test could not have caught the leak
it is named for. My own docstring claimed it drove the fallback and
checked the wire; it did neither.
The double now implements everything the UNGUARDED path reaches
(_is_flood_error, _flood_strikes, _current_edit_interval, _last_edit_time,
_notify_new_message, _try_strip_cursor, _clean_for_display,
_fallback_prefix, _metadata_for_send). Both mutations now fail on real
assertions:
guard 1 removed -> assert consumer._egress_declined is True
guard 2 removed -> AssertionError: the unseen tail reached the wire:
['send']
NON-BLOCKING 1 — a docstring claimed more than the test exercises.
test_requested_platform_name_is_also_normalised described a mixed-case
send_message(target="Discord:999") bypass. That entry point cannot reach
it: _resolve_tool_target lowercases the platform at
tools/send_message_tool.py:47 before the guard runs. The test still pins
a real contract — the helpers must not assume a lowercased argument, for
the gateway lanes and any future non-normalising caller — so the claim is
narrowed to that rather than the test removed.
NON-BLOCKING 2 — the module docstring said "every lane drives the REAL
RelayAdapter". The stream tests drive mixin doubles by design, because
the behaviour under test belongs to the adapter's CALLER. Docstring now
distinguishes the two kinds.
NON-BLOCKING 3 — latch-deletion residue in gateway/relay/adapter.py:418:
return None
return latched if surface_declines else None
The second line was unreachable and referenced a name deleted with the
latch. Removed, along with the 20-line comment block describing the latch
as "the structural fix" — that mechanism now lives on
feat/p5-egress-decline-latch, not here.
The reviewer independently confirmed the large deletion: an AST census
between 3cf45736d7 and f57a2298fa reports only latch symbols removed and
nothing added.
511 passed.
* docs(relay): correct three claims that outran the code
Review of 41ce3cc765 found no new production defect but three overstated
claims, one of them in my own commit message.
1. THE LATCH COMMENTARY WAS STILL THERE. My previous commit message said
it removed "the 20-line comment block describing the latch as the
structural fix". It removed only the unreachable statement. Twenty
lines at adapter.py:361-380 still described a per-chat latch, a choke
point and its scope rules - none of which exist on this branch. In a
refusal-sensitive module that reads as coverage this branch does not
have. Now removed for real.
This is the same defect class as the tests: a claim that outran what
the code does. I made it while fixing that class.
2. THE STREAM-TEST DOCSTRING OVERSTATED BOTH MUTANTS. It said the
mutation "now fails on the assertion that a send reached the wire" -
true of one guard, not both. Verified separately:
remove the _on_edit_failure check -> dies on _egress_declined,
never reaches the fallback
remove the fallback early return -> dies on the wire: ['send']
Both are valid behavioural failures, which is what the blocker asked
for; they are different observables and the docstring now says so.
3. Duplicate `from types import SimpleNamespace` from an earlier scripted
insert; imports reordered.
112 tests pass in the four focused files.
* fix(relay): close two authorization defects found in review
Both were reproduced before fixing and both mutants are pinned.
1. A LIVE relay adapter whose fronts_platform() raised degraded into the
config fallback. `_live_relay_fronted` returned None for every failure,
and None means "no live adapter, use the config snapshot" — so a faulting
adapter plus an empty/stale snapshot made the guard conclude "not
relay-routed" and authorize an unattested destination, while
resolve_delivery_transport asks that same adapter and still routes over
the relay. Measured: relay_routed=False, verdict None for chat 999.
Absence and fault now have separate return values: None only when there
is no runner or no relay adapter; a live adapter that cannot answer
raises RelayRouteUnknown. This is the third instance of this bug class in
this file, and the first two were also mine.
2. An attested chat whose id equalled the requested THREAD id vouched for
that thread. The `thread in attested` arm proved nothing about parentage.
Measured: attested {"-100A", "7"} authorized (-100A, thread 7).
Only the bound `parent:thread` form is accepted now. Nothing legitimate
needed the bare arm — _session_entry_id records a threaded origin as
f"{chat_id}:{thread_id}", and a thread addressed as its own channel
arrives as chat_id and passes the parent check.
The existing test blessed the bare form via parametrize, so it PINNED the
defect. Corrected, plus negative controls for the sibling-chat and
other-parent cases and a positive control proving genuine absence still
takes the config path (otherwise fix 1 would break native-only deploys).
Merged origin/main (was 22 behind). 428 passed via scripts/run_tests.sh;
full 10-row mutation ledger re-killed on the merged tree, none dying on an
exception rather than an assertion.
* fix(relay): only a missing adapter is absence; everything else is a fault
Reviewer BLOCKER, reproduced before fixing. Two more paths where a PRESENT
relay adapter still degraded into the config snapshot:
1. `fronts_platform` may be a property or descriptor, so the ATTRIBUTE
LOOKUP can raise — and the lookup sat inside the absence handler. Probed
with a raising property plus an empty snapshot: live=None, routed=False,
verdict=None, i.e. an unattested target authorized. The previous test made
an already-retrieved METHOD raise, so it could not reach this.
2. A present adapter with no usable `fronts_platform` returned None for the
same reason. An adapter that cannot say what it fronts is broken, not
absent, so it now raises too.
Also found by my own spot-check while the review ran: the nested imports of
`gateway.config` / `gateway.run` inside the live probe shared the broad
handler, so a broken installation degraded to the snapshot as well. Probed
with a healthy-adapter positive control in the same run — healthy refused
the unattested target, faulted authorized it. `_relay_fronted` one function
below already drew this exact distinction for its own import.
The boundary is now: `relay is None` is the ONLY absence. Everything about a
present adapter — attribute access, callability, the call itself, and the
imports needed to reach it — is a fault and raises RelayRouteUnknown.
This is the fourth variant of absence-vs-fault in this file and all four
were mine. The lesson is in the code as a comment rather than in a commit
message nobody re-reads.
Four controls keep genuine absence benign: no runner, no relay adapter in
the runner, a real ModuleNotFoundError naming the gateway package, and the
configured-attested-target-still-sends case.
434 passed via scripts/run_tests.sh; 9-row mutation ledger re-killed
including both new guards, none dying on an exception.
* fix(relay): invert the live probe to fail closed by default
Reviewer BLOCKER round 2, reproduced: reading the adapter registry can also
raise. A runner whose `adapters.get()` raised gave relay_present=True,
live=None, routed=False, verdict=None — unattested discord:999 authorized.
That was the FIFTH boundary in one function with the same defect: the call,
the attribute lookup, a non-callable attribute, the nested imports, and now
the registry lookup. Each round I patched the reported boundary and the
defect moved one statement up. The cause was the shape, not the statements:
the function asked "did something go wrong?" and answered None, and None
MEANS "no live adapter, use the config snapshot" — so every statement was a
new chance to fail open, and every new statement would have been too.
Inverted rather than patched a sixth time. Each `return None` now sits
behind an explicit narrow check that cannot itself be the fault (no runner,
no adapters, no relay key, gateway package genuinely absent), and one outer
handler turns anything else into RelayRouteUnknown. A statement added inside
this function is now fail-CLOSED by default.
Verified all six fault shapes raise (call, attribute, missing method,
registry .get, .adapters property, runner ref) and all five absence shapes
stay benign, plus a liveness control where the config snapshot disagrees
with a healthy adapter and the adapter still wins.
Four new tests, including the two absence controls that keep native-only and
CLI deployments working. 438 passed via scripts/run_tests.sh. Mutation
ledger: 8 killed. One survivor recorded as a proven equivalent mutant —
widening `if not registry` to `or {}` is behaviourally identical because
`{}.get()` returns None, i.e. the same absence; it is a readability guard.
1622 lines
86 KiB
Python
1622 lines
86 KiB
Python
"""Process/completion/update notifications, media delivery and async-delegation delivery for GatewayRunner.
|
|
|
|
Split out of ``gateway/run.py``; bound onto ``GatewayRunner`` via the MRO.
|
|
``gateway.run`` internals are imported lazily inside method bodies (import cycle),
|
|
so ``patch("gateway.run.X")`` keeps intercepting them at call time.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import dataclasses
|
|
import json
|
|
import logging
|
|
import time
|
|
from contextlib import suppress
|
|
from pathlib import Path
|
|
from typing import Any, Dict, Optional, cast
|
|
|
|
from gateway.config import Platform, _BUILTIN_PLATFORM_VALUES
|
|
from gateway.platforms.event import MessageEvent, MessageType
|
|
from gateway.session import SessionEntry, SessionSource
|
|
from gateway.run_shutdown import _log_suppressed, _notice_target_key, _send_error, _send_failed
|
|
|
|
# Log-record parity with the origin module.
|
|
logger = logging.getLogger("gateway.run")
|
|
|
|
_VIDEO_EXTS = {'.mp4', '.mov', '.avi', '.mkv', '.webm', '.3gp'}
|
|
# Routing fields copied verbatim from a process watcher onto its synthetic completion event.
|
|
_WATCHER_ROUTE_FIELDS = ("session_key", "platform", "chat_type", "chat_id", "thread_id", "user_id", "user_name")
|
|
_IMAGE_EXTS = {'.jpg', '.jpeg', '.png', '.webp', '.gif'}
|
|
|
|
# Durable async-delegation claim transitions: kind -> (tools.async_delegation function, failure log).
|
|
_DURABLE_CLAIM_OPS = {
|
|
"drop": ("drop_completion_delivery", "Could not drop durable completion claim"),
|
|
"release": ("release_completion_delivery", "Could not release durable completion claim"),
|
|
"defer": ("defer_completion_delivery", "Could not defer unadmitted completion claim"),
|
|
"complete": ("complete_completion_delivery", "Could not acknowledge durable completion claim"),
|
|
}
|
|
|
|
|
|
def _raw_process_event_session_id(evt: dict) -> str:
|
|
"""Recognize API routes, not malformed structured or partial messaging routes."""
|
|
session_key = str(evt.get("session_key") or "").strip()
|
|
platform = str(evt.get("platform") or "").strip().lower()
|
|
if session_key.startswith("agent:") or platform not in {"", "api_server"}:
|
|
return ""
|
|
if not platform and any(evt.get(field) for field in ("chat_id", "chat_type", "thread_id")):
|
|
return ""
|
|
return str(evt.get("origin_session_id") or session_key or "").strip()
|
|
|
|
|
|
class GatewayNotificationsMixin:
|
|
"""Process/completion/update notifications, media delivery and async-delegation delivery for GatewayRunner."""
|
|
|
|
# Coalescing keys: process completions (short-window fan-in) and async delegations (+ parent session).
|
|
_COMPLETION_BATCH_KEY_FIELDS = ("session_key", "platform", "chat_type", "chat_id", "thread_id", "user_id")
|
|
_ASYNC_GROUP_KEY_FIELDS = ("session_key", "parent_session_id", *_COMPLETION_BATCH_KEY_FIELDS[1:])
|
|
|
|
@dataclasses.dataclass
|
|
class _UpdatePaths:
|
|
"""Marker files ``hermes update --gateway`` and its watcher exchange under HERMES_HOME."""
|
|
|
|
pending: Path
|
|
claimed: Path
|
|
output: Path
|
|
exit_code: Path
|
|
prompt: Path
|
|
response: Path
|
|
|
|
def any_pending(self) -> bool:
|
|
return self.pending.exists() or self.claimed.exists()
|
|
|
|
def unlink_all(self) -> None:
|
|
for p in (self.pending, self.claimed, self.output, self.exit_code, self.prompt, self.response):
|
|
p.unlink(missing_ok=True)
|
|
|
|
@dataclasses.dataclass
|
|
class _UpdateTarget:
|
|
"""Resolved delivery target for update watcher messages."""
|
|
|
|
adapter: Any
|
|
chat_id: Any
|
|
session_key: Optional[str]
|
|
metadata: Any
|
|
platform: Any
|
|
|
|
def send_metadata(self):
|
|
from gateway.run import _non_conversational_metadata
|
|
return _non_conversational_metadata(self.metadata, platform=self.platform)
|
|
|
|
async def send(self, text: str):
|
|
return await self.adapter.send(self.chat_id, text, metadata=self.send_metadata())
|
|
|
|
@dataclasses.dataclass
|
|
class _CompletionClaim:
|
|
"""Pre-flight outcome for one completion delivery."""
|
|
|
|
delegation_id: str = ""
|
|
claim_id: str = ""
|
|
proceed: bool = True
|
|
early_result: Optional[bool] = None
|
|
|
|
async def _deliver_platform_notice(self, source, content: str) -> None:
|
|
"""Deliver a setup/operational notice using platform-specific privacy rules."""
|
|
from gateway.run import _is_slack_ignored_channel
|
|
adapter = self._adapter_for_source(source)
|
|
if not adapter:
|
|
return
|
|
config = getattr(self, "config", None)
|
|
chat_id = getattr(source, "chat_id", None)
|
|
if config and getattr(source, "platform", None) == Platform.SLACK and _is_slack_ignored_channel(config, chat_id):
|
|
logger.info("Skipping Slack platform notice for configured ignored channel %s", chat_id)
|
|
return
|
|
notice_delivery = (
|
|
config.get_notice_delivery(source.platform) if config and hasattr(config, "get_notice_delivery")
|
|
else "public"
|
|
)
|
|
metadata = self._thread_metadata_for_source(source)
|
|
if notice_delivery == "private" and getattr(source, "user_id", None):
|
|
with _log_suppressed(
|
|
logging.DEBUG, "[%s] send_private_notice failed, falling back to public",
|
|
getattr(source, "platform", "?"), exc_info=True,
|
|
):
|
|
result = await adapter.send_private_notice(source.chat_id, source.user_id, content, metadata=metadata)
|
|
if getattr(result, "success", False):
|
|
return
|
|
await adapter.send(source.chat_id, content, metadata=metadata)
|
|
|
|
async def _resolve_compression_lineage_target(
|
|
self, session_db: Any, session_entry: SessionEntry, pinned_session_id: str,
|
|
) -> Optional[str]:
|
|
"""Return the live compression tip of ``pinned_session_id`` if the route owns that lineage, else None."""
|
|
try:
|
|
target_session_id = await session_db.get_compression_tip(pinned_session_id)
|
|
except Exception:
|
|
logger.debug("Async-delegation compression-tip lookup failed for %s", pinned_session_id, exc_info=True)
|
|
target_session_id = None
|
|
if not target_session_id or target_session_id == pinned_session_id:
|
|
logger.warning(
|
|
"Async-delegation completion pinned to compressed session %s "
|
|
"without a continuation; dropping injection.", pinned_session_id,
|
|
)
|
|
return None
|
|
try:
|
|
tip_row = await session_db.get_session(target_session_id)
|
|
except Exception:
|
|
tip_row = None
|
|
if tip_row is None or tip_row.get("ended_at"):
|
|
logger.warning(
|
|
"Async-delegation compression continuation %s is %s; dropping injection.",
|
|
target_session_id, "unknown" if tip_row is None else "ended",
|
|
)
|
|
return None
|
|
route_owns_lineage = session_entry.session_id in {pinned_session_id, target_session_id}
|
|
if not route_owns_lineage:
|
|
# Across several rotations, accept a stale route only when its own tip is the same live target.
|
|
try:
|
|
route_row = await session_db.get_session(session_entry.session_id)
|
|
route_tip = (
|
|
await session_db.get_compression_tip(session_entry.session_id)
|
|
if route_row is not None
|
|
and route_row.get("ended_at")
|
|
and route_row.get("end_reason") == "compression"
|
|
else None
|
|
)
|
|
except Exception:
|
|
route_tip = None
|
|
route_owns_lineage = route_tip == target_session_id
|
|
if not route_owns_lineage:
|
|
logger.warning(
|
|
"Async-delegation completion for compression lineage %s -> %s "
|
|
"does not own current route %s; dropping injection.",
|
|
pinned_session_id, target_session_id, session_entry.session_id,
|
|
)
|
|
return None
|
|
return target_session_id
|
|
|
|
async def _resolve_async_delegation_session(
|
|
self, session_entry: SessionEntry, pinned_session_id: str,
|
|
) -> Optional[SessionEntry]:
|
|
"""Resolve an async completion to its verified owning gateway session.
|
|
|
|
Follow compression-rotation lineage (parent row ended, child continues), but never let a
|
|
late completion override an unrelated /new or restored route. Unknown ownership fails
|
|
closed; the result stays in the delegation records.
|
|
"""
|
|
from gateway.run import _USER_BOUNDARY_END_REASONS
|
|
session_db = cast(Any, self._session_db)
|
|
if session_db is None:
|
|
logger.warning(
|
|
"Async-delegation completion has no session database; "
|
|
"dropping injection (#55578 fail-closed)."
|
|
)
|
|
return None
|
|
pinned_row = None
|
|
try:
|
|
pinned_row = await session_db.get_session(pinned_session_id)
|
|
except Exception:
|
|
logger.debug("Async-delegation parent lookup failed for %s", pinned_session_id, exc_info=True)
|
|
if pinned_row is None:
|
|
logger.warning(
|
|
"Async-delegation completion has unknown spawning session %s; "
|
|
"dropping injection (#55578 fail-closed).", pinned_session_id,
|
|
)
|
|
return None
|
|
target_session_id = pinned_session_id
|
|
follows_compression = False
|
|
if pinned_row.get("ended_at"):
|
|
_end_reason = str(pinned_row.get("end_reason") or "")
|
|
if _end_reason in _USER_BOUNDARY_END_REASONS:
|
|
logger.warning(
|
|
"Async-delegation completion pinned to user-closed session %s "
|
|
"(end_reason=%r); dropping injection instead of resurrecting it "
|
|
"(#55578 fail-closed).", pinned_session_id, _end_reason,
|
|
)
|
|
return None
|
|
if _end_reason != "compression":
|
|
# Idle/timeout end (scale-to-zero norm): the chat route is still valid, so deliver to its
|
|
# current session rather than drop (the row would be acked then silently lost).
|
|
logger.info(
|
|
"Async-delegation completion pinned to %s-ended session %s; "
|
|
"retargeting to the chat's current session %s.",
|
|
_end_reason or "idle", pinned_session_id, session_entry.session_id,
|
|
)
|
|
return session_entry
|
|
follows_compression = True
|
|
target_session_id = await self._resolve_compression_lineage_target(
|
|
session_db, session_entry, pinned_session_id,
|
|
)
|
|
if target_session_id is None:
|
|
return None
|
|
if target_session_id == session_entry.session_id:
|
|
return session_entry
|
|
prior_session_id = session_entry.session_id
|
|
if follows_compression:
|
|
switched = await self.async_session_store.advance_compression_session(
|
|
session_entry.session_key, prior_session_id, target_session_id,
|
|
)
|
|
else:
|
|
switched = await self.async_session_store.switch_session(session_entry.session_key, target_session_id)
|
|
if switched is None:
|
|
logger.warning(
|
|
"Async-delegation completion could not bind routing key %s to "
|
|
"owning session %s; dropping injection.", session_entry.session_key, target_session_id,
|
|
)
|
|
return None
|
|
logger.info(
|
|
"Pinned async-delegation completion to owning session %s (was %s) for routing key %s (#57498)",
|
|
target_session_id, prior_session_id, session_entry.session_key,
|
|
)
|
|
return switched
|
|
|
|
async def _deliver_media_from_response(
|
|
self, response: str, event: MessageEvent, adapter, thread_metadata: Optional[Dict[str, Any]] = None
|
|
) -> None:
|
|
"""Deliver explicit MEDIA: tags from an already-streamed response (text already delivered).
|
|
EXPLICIT-ONLY, unlike the non-streaming path in ``gateway/platforms/base.py``: a bare local
|
|
path in a streamed reply is shown text or stale inspected content, and promoting it sent
|
|
files the model never asked for. MEDIA tags are NOT deduped against prior turns (a final-reply
|
|
directive is a deliberate attach); stale auto-appended tags are deduped upstream.
|
|
|
|
Only ``MEDIA:`` directives — the explicit attachment contract — trigger post-stream uploads. See
|
|
#20834.
|
|
"""
|
|
from urllib.parse import quote as _quote
|
|
with _log_suppressed(logging.WARNING, "Post-stream media extraction failed: %s"):
|
|
# Capture [[as_document]] before extract_media strips it: images then go via send_document.
|
|
force_document_attachments = "[[as_document]]" in response
|
|
from gateway.platforms.base import BasePlatformAdapter, should_send_media_as_audio
|
|
media_files, cleaned = adapter.extract_media(response)
|
|
media_files = BasePlatformAdapter.filter_media_delivery_paths(media_files)
|
|
# Strip image URLs (parity with the non-streaming chain); no extract_local_files here.
|
|
# Do NOT deduplicate explicit MEDIA tags against prior turns here (#73771). This rescan is
|
|
# already EXPLICIT-ONLY (see docstring): a MEDIA: directive in the final streamed reply is the
|
|
# model deliberately attaching a file — including a user-requested resend. Stale auto-appended
|
|
# tags are deduped upstream in _collect_auto_append_media_tags with history_media_paths. Mirrors
|
|
# the same filter removal on the non-streaming path in gateway/platforms/base.py. Bare local
|
|
# paths in an already-streamed reply are text the user has seen (or stale inspected content),
|
|
# not an attachment request.
|
|
adapter.extract_images(cleaned)
|
|
_thread_meta = (
|
|
dict(thread_metadata)
|
|
if thread_metadata is not None
|
|
else self._thread_metadata_for_source(event.source, self._reply_anchor_for_event(event))
|
|
)
|
|
chat_id = event.source.chat_id
|
|
# Images go out as one batch (e.g. Signal's multi-attachment RPC) unless [[as_document]].
|
|
def _is_photo(media_path: str, is_voice: bool) -> bool:
|
|
ext = Path(media_path).suffix.lower()
|
|
return ext in _IMAGE_EXTS and not is_voice and not force_document_attachments
|
|
|
|
image_paths = [p for p, v in media_files if _is_photo(p, v)]
|
|
non_image_media = [(p, v) for p, v in media_files if not _is_photo(p, v)]
|
|
if image_paths:
|
|
try:
|
|
images = [(f"file://{_quote(p)}", "") for p in image_paths]
|
|
await adapter.send_multiple_images(chat_id=chat_id, images=images, metadata=_thread_meta)
|
|
except Exception as e:
|
|
logger.warning("[%s] Post-stream image batch delivery failed: %s", adapter.name, e)
|
|
for media_path, is_voice in non_image_media:
|
|
try:
|
|
ext = Path(media_path).suffix.lower()
|
|
if should_send_media_as_audio(event.source.platform, ext, is_voice=is_voice):
|
|
await adapter.send_voice(
|
|
chat_id=chat_id, audio_path=media_path, metadata=_thread_meta, is_voice=is_voice,
|
|
)
|
|
elif ext in _VIDEO_EXTS:
|
|
await adapter.send_video(chat_id=chat_id, video_path=media_path, metadata=_thread_meta)
|
|
else:
|
|
await adapter.send_document(chat_id=chat_id, file_path=media_path, metadata=_thread_meta)
|
|
except Exception as e:
|
|
logger.warning("[%s] Post-stream media delivery failed: %s", adapter.name, e)
|
|
|
|
|
|
async def _deliver_queued_first_response(
|
|
self, response: str, source: SessionSource, adapter,
|
|
metadata: Optional[Dict[str, Any]] = None, event_message_id: Optional[str] = None,
|
|
text_already_delivered: bool = False, deliver_media: bool = True, stream_consumer=None,
|
|
) -> None:
|
|
"""Deliver a queued response using the normal text+attachment split."""
|
|
from gateway.run import _strip_response_attachments_for_direct_send
|
|
if not text_already_delivered:
|
|
text_content = _strip_response_attachments_for_direct_send(response, adapter)
|
|
if text_content:
|
|
# Reconcile-by-edit first: a stream-sealed message already carries most of the answer;
|
|
# a plain send here would duplicate it.
|
|
_reconciled = False
|
|
_sc_msg_id = getattr(stream_consumer, "message_id", None)
|
|
if (
|
|
_sc_msg_id
|
|
and _sc_msg_id != "__no_edit__"
|
|
and not getattr(stream_consumer, "_turn_split_delivery", False)
|
|
):
|
|
try:
|
|
_edit_res = await adapter.edit_message(
|
|
chat_id=source.chat_id, message_id=_sc_msg_id, content=text_content, finalize=True,
|
|
)
|
|
if getattr(_edit_res, "success", False):
|
|
_reconciled = True
|
|
logger.info(
|
|
"Queued-lane final reconciled by editing message %s in place (no duplicate send).",
|
|
_sc_msg_id,
|
|
)
|
|
else:
|
|
# P5(b): a DECLINE is not "editing unavailable". The
|
|
# send below re-delivers the whole response to the
|
|
# chat the connector just refused.
|
|
from gateway.relay.egress import declined_send
|
|
|
|
if declined_send(_edit_res):
|
|
logger.warning(
|
|
"Queued-lane reconcile edit DECLINED by the "
|
|
"connector's egress guard; not falling back "
|
|
"to a send (the destination is not approved)."
|
|
)
|
|
return
|
|
except Exception as _qe:
|
|
logger.debug("Queued-lane reconcile edit failed (%s); falling back to send.", _qe)
|
|
if not _reconciled:
|
|
await adapter.send(source.chat_id, text_content, metadata=metadata)
|
|
# Failed turns deliver their (normalized failure) text but must not upload attachments as if
|
|
# they succeeded — mirrors the ``not agent_result.get("failed")`` completed-turn guard.
|
|
if not deliver_media:
|
|
return
|
|
await self._deliver_media_from_response(
|
|
response, MessageEvent(text="", source=source, message_id=event_message_id), adapter,
|
|
thread_metadata=metadata,
|
|
)
|
|
|
|
def _schedule_update_notification_watch(self) -> None:
|
|
"""Ensure a background task is watching for update completion."""
|
|
existing_task = getattr(self, "_update_notification_task", None)
|
|
if existing_task and not existing_task.done():
|
|
return
|
|
try:
|
|
self._update_notification_task = asyncio.create_task(self._watch_update_progress())
|
|
except RuntimeError:
|
|
logger.debug("Skipping update notification watcher: no running event loop")
|
|
|
|
@classmethod
|
|
def _update_paths(cls) -> "GatewayNotificationsMixin._UpdatePaths":
|
|
from gateway.run import _hermes_home
|
|
return cls._UpdatePaths(
|
|
pending=_hermes_home / ".update_pending.json",
|
|
claimed=_hermes_home / ".update_pending.claimed.json", output=_hermes_home / ".update_output.txt",
|
|
exit_code=_hermes_home / ".update_exit_code",
|
|
prompt=_hermes_home / ".update_prompt.json", response=_hermes_home / ".update_response",
|
|
)
|
|
|
|
def _resolve_update_target(self, paths: "_UpdatePaths") -> Optional["_UpdateTarget"]:
|
|
"""Resolve adapter/chat/session for update watcher messages from the pending marker."""
|
|
for path in (paths.claimed, paths.pending):
|
|
if not path.exists():
|
|
continue
|
|
with suppress(Exception):
|
|
pending = json.loads(path.read_text(encoding="utf-8"))
|
|
platform_str = pending.get("platform")
|
|
chat_id = pending.get("chat_id")
|
|
session_key = pending.get("session_key")
|
|
if not (platform_str and chat_id):
|
|
continue # BASE: an incomplete marker falls through to the next path, not "unresolved"
|
|
platform = Platform(platform_str)
|
|
adapter = self.adapters.get(platform)
|
|
if not adapter:
|
|
return None
|
|
metadata = self._pending_marker_metadata(platform, chat_id, pending, adapter)
|
|
# Fallback session key if not stored (old pending files)
|
|
return self._UpdateTarget(
|
|
adapter, chat_id, session_key or f"{platform_str}:{chat_id}", metadata, platform,
|
|
)
|
|
return None
|
|
|
|
def _pending_marker_metadata(self, platform, chat_id, data: dict, adapter):
|
|
"""Thread metadata for a persisted update/restart marker (thread_id/chat_type/message_id keys)."""
|
|
return self._thread_metadata_for_target(
|
|
platform, chat_id, data.get("thread_id"), chat_type=data.get("chat_type"),
|
|
reply_to_message_id=data.get("message_id"), adapter=adapter,
|
|
)
|
|
|
|
async def _watch_update_completion_only(self, paths: "_UpdatePaths", deadline: float, poll_interval: float) -> None:
|
|
"""Fallback when no adapter/chat can be resolved: wait for the exit code, then notify."""
|
|
logger.warning("Update watcher: cannot resolve adapter/chat_id, falling back to completion-only")
|
|
# Poll until _send_update_notification delivers (it returns False while the platform reconnects).
|
|
loop = asyncio.get_running_loop()
|
|
while paths.any_pending() and loop.time() < deadline:
|
|
if paths.exit_code.exists() and await self._send_update_notification():
|
|
return
|
|
await asyncio.sleep(poll_interval)
|
|
if paths.any_pending() and not paths.exit_code.exists():
|
|
paths.exit_code.write_text("124", encoding="utf-8")
|
|
await self._send_update_notification()
|
|
|
|
@staticmethod
|
|
def _update_exit_code(paths: "_UpdatePaths") -> int:
|
|
return int(paths.exit_code.read_text(encoding="utf-8").strip() or "1")
|
|
|
|
@staticmethod
|
|
def _read_update_output_since(path: Path, offset: int) -> tuple[str, int]:
|
|
"""Read update output defensively; logs may contain invalid UTF-8."""
|
|
try:
|
|
data = path.read_bytes()
|
|
except OSError:
|
|
return "", offset
|
|
if len(data) <= offset:
|
|
return "", len(data)
|
|
return data[offset:].decode("utf-8", errors="replace"), len(data)
|
|
|
|
async def _send_update_output(self, target: "_UpdateTarget", text: str) -> None:
|
|
"""Send buffered update output as fenced chunks that fit message limits (Telegram: 4096)."""
|
|
from tools.ansi_strip import strip_ansi
|
|
clean = strip_ansi(text).strip()
|
|
if not clean:
|
|
return
|
|
max_chunk = 3500
|
|
for i in range(0, len(clean), max_chunk):
|
|
with _log_suppressed(logging.DEBUG, "Update stream send failed: %s"):
|
|
await target.send(f"```\n{clean[i:i + max_chunk]}\n```")
|
|
|
|
async def _forward_update_prompt(self, target: "_UpdateTarget", prompt_text: str, default: str) -> None:
|
|
"""Forward an update prompt: platform-native buttons first (Discord, Telegram), else text."""
|
|
sent_buttons = False
|
|
adapter = target.adapter
|
|
if getattr(type(adapter), "send_update_prompt", None) is not None:
|
|
with _log_suppressed(logging.DEBUG, "Button-based update prompt failed: %s"):
|
|
await adapter.send_update_prompt(
|
|
chat_id=target.chat_id, prompt=prompt_text, default=default,
|
|
session_key=target.session_key, metadata=target.send_metadata(),
|
|
)
|
|
sent_buttons = True
|
|
if not sent_buttons:
|
|
default_hint = f" (default: {default})" if default else ""
|
|
_p = getattr(adapter, "typed_command_prefix", "/")
|
|
await target.send(
|
|
f"⚕ **Update needs your input:**\n\n{prompt_text}{default_hint}\n\n"
|
|
f"Reply `{_p}approve` (yes) or `{_p}deny` (no), or type your answer directly."
|
|
)
|
|
# Keep the prompt marker on disk until answered so a restarted watcher can re-forward it.
|
|
self._session_state(target.session_key).persistent.update_prompt_pending = True
|
|
logger.info("Forwarded update prompt to %s: %s", target.session_key, prompt_text[:80])
|
|
|
|
def _clear_update_markers(self, paths: "_UpdatePaths", session_key: Optional[str]) -> None:
|
|
paths.unlink_all()
|
|
state = self._peek_session_state(session_key)
|
|
if state is not None:
|
|
state.persistent.update_prompt_pending = False
|
|
|
|
async def _watch_update_progress(
|
|
self, poll_interval: float = 2.0, stream_interval: float = 4.0, timeout: float = 1800.0
|
|
) -> None:
|
|
"""Watch ``hermes update --gateway``, streaming output + forwarding prompts.
|
|
|
|
Polls ``.update_output.txt`` for new content and sends chunks to the user periodically;
|
|
detects ``.update_prompt.json`` (written when the update process needs input) and forwards it.
|
|
"""
|
|
paths = self._update_paths()
|
|
loop = asyncio.get_running_loop()
|
|
deadline = loop.time() + timeout
|
|
target = self._resolve_update_target(paths)
|
|
if target is None:
|
|
await self._watch_update_completion_only(paths, deadline, poll_interval)
|
|
return
|
|
session_key = target.session_key
|
|
bytes_sent = 0
|
|
last_stream_time = loop.time()
|
|
buffer = ""
|
|
|
|
async def _flush_buffer() -> None:
|
|
nonlocal buffer, last_stream_time
|
|
text, buffer = buffer, ""
|
|
if text.strip():
|
|
last_stream_time = loop.time()
|
|
await self._send_update_output(target, text)
|
|
|
|
def _read_new_output() -> None:
|
|
nonlocal buffer, bytes_sent
|
|
if paths.output.exists():
|
|
with suppress(OSError):
|
|
chunk, bytes_sent = self._read_update_output_since(paths.output, bytes_sent)
|
|
buffer += chunk
|
|
|
|
while loop.time() < deadline:
|
|
if paths.exit_code.exists():
|
|
_read_new_output()
|
|
await _flush_buffer()
|
|
with _log_suppressed(logging.WARNING, "Update final notification failed: %s"):
|
|
exit_code = self._update_exit_code(paths)
|
|
await target.send(
|
|
"✅ Hermes update finished." if exit_code == 0
|
|
else "❌ Hermes update failed (exit code {}).".format(exit_code)
|
|
)
|
|
logger.info("Update finished (exit=%s), notified %s", exit_code, session_key)
|
|
self._clear_update_markers(paths, session_key)
|
|
return
|
|
_read_new_output()
|
|
if buffer.strip() and (loop.time() - last_stream_time) >= stream_interval:
|
|
await _flush_buffer()
|
|
# Forward a prompt only when none is pending, else every poll re-forwards the same prompt.
|
|
_pending_state = self._peek_session_state(session_key) if session_key else None
|
|
if paths.prompt.exists() and session_key and not getattr(
|
|
getattr(_pending_state, "persistent", None), "update_prompt_pending", False
|
|
):
|
|
try:
|
|
prompt_data = json.loads(paths.prompt.read_text(encoding="utf-8"))
|
|
prompt_text = prompt_data.get("prompt", "")
|
|
if prompt_text:
|
|
await _flush_buffer() # user sees context before the prompt
|
|
await self._forward_update_prompt(target, prompt_text, prompt_data.get("default", ""))
|
|
except (json.JSONDecodeError, OSError) as e:
|
|
logger.debug("Failed to read update prompt: %s", e)
|
|
await asyncio.sleep(poll_interval)
|
|
if not paths.exit_code.exists():
|
|
logger.warning("Update watcher timed out after %.0fs", timeout)
|
|
paths.exit_code.write_text("124", encoding="utf-8")
|
|
await _flush_buffer()
|
|
with suppress(Exception):
|
|
await target.send("❌ Hermes update timed out after 30 minutes.")
|
|
self._clear_update_markers(paths, session_key)
|
|
|
|
async def _send_update_notification(self) -> bool:
|
|
"""If an update finished, notify the user.
|
|
|
|
False while the update is still running (caller may retry); True after a definitive send/skip.
|
|
"""
|
|
from gateway.run import _non_conversational_metadata
|
|
paths = self._update_paths()
|
|
if not paths.any_pending():
|
|
return False
|
|
cleanup = True
|
|
active_pending_path = paths.claimed
|
|
|
|
def _defer(reason: str, *args) -> bool:
|
|
nonlocal cleanup, active_pending_path
|
|
logger.info(reason, *args)
|
|
cleanup = False
|
|
active_pending_path = paths.pending
|
|
paths.claimed.replace(paths.pending)
|
|
return False
|
|
|
|
try:
|
|
if paths.pending.exists():
|
|
try:
|
|
paths.pending.replace(paths.claimed)
|
|
except FileNotFoundError:
|
|
if not paths.claimed.exists():
|
|
return True
|
|
elif not paths.claimed.exists():
|
|
return True
|
|
pending = json.loads(paths.claimed.read_text(encoding="utf-8"))
|
|
platform_str = pending.get("platform")
|
|
chat_id = pending.get("chat_id")
|
|
if not paths.exit_code.exists():
|
|
return _defer("Update notification deferred: update still running")
|
|
exit_code = self._update_exit_code(paths)
|
|
output = paths.output.read_bytes().decode("utf-8", errors="replace") if paths.output.exists() else ""
|
|
platform = Platform(platform_str)
|
|
adapter = self.adapters.get(platform)
|
|
if chat_id and not adapter:
|
|
# Target platform not reconnected yet (common right after the update's restart): keep the
|
|
# markers for a later retry instead of silently losing the notification.
|
|
return _defer("Update notification deferred: %s adapter not connected yet", platform_str)
|
|
if chat_id:
|
|
metadata = self._pending_marker_metadata(platform, chat_id, pending, adapter)
|
|
from tools.ansi_strip import strip_ansi
|
|
output = strip_ansi(output).strip()
|
|
if output:
|
|
if len(output) > 3500:
|
|
output = "…" + output[-3500:]
|
|
status = "✅ Hermes update finished." if exit_code == 0 else "❌ Hermes update failed."
|
|
msg = f"{status}\n\n```\n{output}\n```"
|
|
else:
|
|
msg = (
|
|
"✅ Hermes update finished successfully." if exit_code == 0 else
|
|
"❌ Hermes update failed. Check the gateway logs or run `hermes update` manually for details."
|
|
)
|
|
await adapter.send(chat_id, msg, metadata=_non_conversational_metadata(metadata, platform=platform))
|
|
logger.info("Sent post-update notification to %s:%s (exit=%s)", platform_str, chat_id, exit_code)
|
|
except Exception as e:
|
|
logger.warning("Post-update notification failed: %s", e)
|
|
finally:
|
|
if cleanup:
|
|
for p in (active_pending_path, paths.claimed, paths.output, paths.exit_code):
|
|
p.unlink(missing_ok=True)
|
|
return True
|
|
|
|
async def _send_restart_notification(self) -> Optional[tuple[str, str, Optional[str]]]:
|
|
"""Notify the chat that initiated /restart that the gateway is back."""
|
|
from gateway.delivery import resolve_delivery_transport
|
|
from gateway.run import _hermes_home, _non_conversational_metadata
|
|
notify_path = _hermes_home / ".restart_notify.json"
|
|
if not notify_path.exists():
|
|
return None
|
|
try:
|
|
data = json.loads(notify_path.read_text(encoding="utf-8"))
|
|
platform_str = data.get("platform")
|
|
chat_id = data.get("chat_id")
|
|
thread_id = data.get("thread_id")
|
|
if not platform_str or not chat_id:
|
|
return None
|
|
platform = Platform(platform_str)
|
|
transport = resolve_delivery_transport(platform, self.config, self.adapters)
|
|
if transport is None:
|
|
logger.debug("Restart notification skipped: no live transport for %s", platform_str)
|
|
return None
|
|
platform_cfg = self.config.platforms.get(platform)
|
|
if platform_cfg is not None and not platform_cfg.gateway_restart_notification:
|
|
logger.info(
|
|
"Restart notification suppressed: %s has gateway_restart_notification=false", platform_str
|
|
)
|
|
return None
|
|
metadata = self._pending_marker_metadata(platform, chat_id, data, transport.adapter)
|
|
if data.get("delivered_via_upstream_relay") is True:
|
|
metadata = dict(metadata or {})
|
|
for field in ("user_id", "scope_id"):
|
|
if data.get(field):
|
|
metadata[field] = str(data[field])
|
|
result = await transport.send(
|
|
platform, str(chat_id), "♻ Gateway restarted successfully. Your session continues.",
|
|
metadata=_non_conversational_metadata(metadata, platform=platform),
|
|
)
|
|
# adapter.send() catches provider errors (e.g. "Chat not found") and returns
|
|
# SendResult(success=False) rather than raising, so inspect the result before claiming success.
|
|
if _send_failed(result):
|
|
logger.warning(
|
|
"Restart notification to %s:%s was not delivered: %s", platform_str, chat_id, _send_error(result),
|
|
)
|
|
return None
|
|
logger.info("Sent restart notification to %s:%s", platform_str, chat_id)
|
|
return str(platform_str), str(chat_id), str(thread_id) if thread_id else None
|
|
except Exception as e:
|
|
logger.warning("Restart notification failed: %s", e)
|
|
return None
|
|
finally:
|
|
notify_path.unlink(missing_ok=True)
|
|
|
|
def _home_channel_transports(self):
|
|
"""Yield ``(platform, platform_cfg, home, transport)`` for every home channel with a live transport."""
|
|
from gateway.delivery import resolve_delivery_transport
|
|
for platform, platform_cfg in self.config.platforms.items():
|
|
home = platform_cfg.home_channel
|
|
if not home or not home.chat_id:
|
|
continue
|
|
transport = resolve_delivery_transport(platform, self.config, self.adapters)
|
|
if transport is None:
|
|
continue
|
|
yield platform, platform_cfg, home, transport
|
|
|
|
async def _send_home_channel_message(self, platform, home, transport, message: str, failure_fmt: str) -> bool:
|
|
"""Best-effort send to one home channel; True on success, failures logged with ``failure_fmt``."""
|
|
from gateway.run import _non_conversational_metadata
|
|
try:
|
|
metadata = self._thread_metadata_for_target(platform, home.chat_id, home.thread_id, adapter=transport.adapter)
|
|
if transport.is_relay:
|
|
metadata = dict(metadata or {})
|
|
if home.user_id:
|
|
metadata["user_id"] = home.user_id
|
|
if home.scope_id:
|
|
metadata["scope_id"] = home.scope_id
|
|
send_metadata = _non_conversational_metadata(metadata, platform=platform)
|
|
if send_metadata is not None or transport.is_relay:
|
|
result = await transport.send(platform, str(home.chat_id), message, metadata=send_metadata)
|
|
else:
|
|
result = await transport.adapter.send(str(home.chat_id), message)
|
|
if _send_failed(result):
|
|
logger.warning(failure_fmt, platform.value, home.chat_id, _send_error(result))
|
|
return False
|
|
return True
|
|
except Exception as exc:
|
|
logger.warning(failure_fmt, platform.value, home.chat_id, exc)
|
|
return False
|
|
|
|
async def _send_home_channel_startup_notifications(
|
|
self, *, skip_targets: Optional[set[tuple[str, str, Optional[str]]]] = None
|
|
) -> set[tuple[str, str, Optional[str]]]:
|
|
"""Notify configured home channels that the gateway is back online.
|
|
|
|
Best-effort, once per connected platform home channel. ``skip_targets`` lets startup avoid
|
|
duplicate messages when a more specific restart notification is queued for the same chat.
|
|
"""
|
|
delivered: set[tuple[str, str, Optional[str]]] = set()
|
|
skipped = skip_targets or set()
|
|
message = "♻️ Gateway online — Hermes is back and ready."
|
|
for platform, platform_cfg, home, transport in self._home_channel_transports():
|
|
if not platform_cfg.gateway_restart_notification:
|
|
logger.info(
|
|
"Home-channel startup notification suppressed: %s has gateway_restart_notification=false",
|
|
platform.value,
|
|
)
|
|
continue
|
|
target = _notice_target_key(platform.value, home.chat_id, home.thread_id)
|
|
if target in skipped or target in delivered:
|
|
continue
|
|
if await self._send_home_channel_message(
|
|
platform, home, transport, message, "Home-channel startup notification failed for %s:%s: %s",
|
|
):
|
|
delivered.add(target)
|
|
logger.info("Sent home-channel startup notification to %s:%s", platform.value, home.chat_id)
|
|
return delivered
|
|
|
|
async def _send_session_db_warning_notifications(self) -> None:
|
|
"""Broadcast a state.db failure warning to all home channels.
|
|
|
|
When SessionDB init fails at gateway startup, messages may flow but nothing is persisted
|
|
— /resume, /history, and session_search all silently break. Best-effort: failures are
|
|
logged, not raised.
|
|
|
|
See #88235.
|
|
"""
|
|
error = getattr(self, "_session_db_init_error", None)
|
|
if not error:
|
|
return
|
|
from hermes_constants import get_default_hermes_root
|
|
from hermes_state import _default_db_path, classify_persistence_error, format_session_db_unavailable
|
|
if classify_persistence_error(error) == "corrupt":
|
|
# Copy-pasteable, so name the real store (profiles / HERMES_HOME do not live under ~/.hermes).
|
|
db_path = _default_db_path()
|
|
backups_dir = get_default_hermes_root() / "backups"
|
|
message = (
|
|
"⚠️ Session database corruption detected. Messages may not be "
|
|
"persisted. Recovery options:\n"
|
|
"1. Run `hermes doctor --fix`\n"
|
|
"2. Stop the gateway, then recover with:\n"
|
|
f" hermes sessions recover --source {db_path} "
|
|
"--inspect-only\n"
|
|
" (if it reports recoverable) hermes sessions recover "
|
|
f"--source {db_path} --output recovered-state.db\n"
|
|
" — recovery snapshots the damaged file first; do NOT run "
|
|
"`sqlite3 ... \".recover\"` against the live state.db, a "
|
|
"vulnerable sqlite3 CLI can corrupt it further\n"
|
|
f"3. Restore from a backup in {backups_dir}/\n"
|
|
"Run `hermes doctor` for sanitized diagnostics."
|
|
)
|
|
else:
|
|
message = (
|
|
f"⚠️ Session database unavailable — messages may not be persisted. "
|
|
f"{format_session_db_unavailable()}\nRun `hermes doctor` for diagnostics."
|
|
)
|
|
logger.warning("Broadcasting state.db failure warning to home channels: %s", error)
|
|
for platform, _platform_cfg, home, transport in self._home_channel_transports():
|
|
await self._send_home_channel_message(
|
|
platform, home, transport, message, "state.db warning notification failed for %s:%s: %s",
|
|
)
|
|
|
|
def _build_process_event_source(self, evt: dict):
|
|
"""Resolve the canonical source for a synthetic background-process event.
|
|
|
|
Prefer the persisted session-store origin; the active foreground event causes cross-topic bleed.
|
|
"""
|
|
from gateway.run import _parse_session_key
|
|
session_key = str(evt.get("session_key") or "").strip()
|
|
derived = {}
|
|
parts = session_key.split(":")
|
|
profile = parts[1] if len(parts) >= 5 and parts[0] == "agent" and parts[1] != "main" else None
|
|
if session_key:
|
|
try:
|
|
self.session_store._ensure_loaded()
|
|
entry = self.session_store._entries.get(session_key)
|
|
if entry and getattr(entry, "origin", None):
|
|
return entry.origin
|
|
except Exception as exc:
|
|
logger.debug("Synthetic process-event session-store lookup failed for %s: %s", session_key, exc)
|
|
cached_source = self._get_cached_session_source(session_key)
|
|
if cached_source is not None:
|
|
return cached_source
|
|
parse_key = ":".join(["agent", "main", *parts[2:]]) if profile else session_key
|
|
derived = _parse_session_key(parse_key) or {}
|
|
platform_name = str(evt.get("platform") or derived.get("platform") or "").strip().lower()
|
|
chat_type = str(evt.get("chat_type") or derived.get("chat_type") or "").strip().lower()
|
|
chat_id = str(evt.get("chat_id") or derived.get("chat_id") or "").strip()
|
|
if not platform_name or not chat_type or not chat_id:
|
|
# Raw API keys legitimately have no messaging source. Resolve persisted
|
|
# origins first, then leave this recognized route to the API dispatcher.
|
|
if _raw_process_event_session_id(evt):
|
|
return None
|
|
logger.warning(
|
|
"Synthetic event source unresolvable: "
|
|
"session_key=%r platform=%r chat_type=%r chat_id=%r evt_type=%s",
|
|
session_key, platform_name, chat_type, chat_id, evt.get("type", "?"),
|
|
)
|
|
return None
|
|
try:
|
|
platform = Platform(platform_name)
|
|
# Reject dynamic pseudo-members: plugin platforms must be registered.
|
|
if platform.value not in _BUILTIN_PLATFORM_VALUES:
|
|
try:
|
|
from gateway.platform_registry import platform_registry
|
|
if not platform_registry.is_registered(platform.value):
|
|
raise ValueError(platform_name)
|
|
except Exception:
|
|
raise ValueError(platform_name)
|
|
except Exception:
|
|
logger.warning("Synthetic process event has invalid platform metadata: %r", platform_name)
|
|
return None
|
|
|
|
def _opt(field: str) -> Optional[str]:
|
|
return str(evt.get(field) or "").strip() or None
|
|
|
|
scope_id = _opt("scope_id")
|
|
if scope_id is None and chat_type not in ("dm", "thread"):
|
|
# Reconstructed scoped-chat source without scope_id: a relay connector's tenant guard may
|
|
# decline the reply. Warn, don't fail (native adapters need no scope_id).
|
|
logger.warning(
|
|
"Synthetic event source for %s chat=%s (%s) reconstructed "
|
|
"without scope_id; scoped relay egress may be declined by "
|
|
"the connector's tenant guard (user_id fallback only).", platform_name, chat_id, chat_type,
|
|
)
|
|
return SessionSource(
|
|
platform=platform, chat_id=chat_id, chat_type=chat_type, thread_id=_opt("thread_id"),
|
|
user_id=_opt("user_id"), user_name=_opt("user_name"), scope_id=scope_id, profile=profile,
|
|
)
|
|
|
|
async def _drain_watch_notifications(self, completion_queue) -> None:
|
|
"""Consume queued watch events and inject them when notifications are enabled.
|
|
|
|
The queue is ALWAYS drained (so watch events don't rot or requeue-spin) but injection is
|
|
skipped entirely when ``display.background_process_notifications`` is ``off``.
|
|
|
|
See #9290.
|
|
"""
|
|
from gateway.run import _drain_gateway_watch_events, _format_gateway_process_notification
|
|
watch_events = _drain_gateway_watch_events(completion_queue)
|
|
if self._load_background_notifications_mode() == "off":
|
|
return
|
|
for evt in watch_events:
|
|
synth_text = _format_gateway_process_notification(evt)
|
|
if not synth_text:
|
|
continue
|
|
try:
|
|
delivered = await self._inject_watch_notification(synth_text, evt)
|
|
except Exception:
|
|
logger.exception("Watch notification injection error")
|
|
delivered = False
|
|
if delivered is False:
|
|
completion_queue.put(evt)
|
|
|
|
def _adapter_by_platform_value(self, platform_name: str):
|
|
"""Literal ``p.value == platform_name`` scan over connected adapters (native adapters only)."""
|
|
for p, a in self.adapters.items():
|
|
if p.value == platform_name:
|
|
return a
|
|
return None
|
|
|
|
async def _self_post_api_server(self, adapter, synth_text: str, raw_sid: str, evt: dict) -> bool:
|
|
"""Deliver to a non-push (api_server) session by raw session id.
|
|
|
|
Async-delegation completions are persisted as a durable delivery row — after the parent
|
|
turn's event.complete the CLIENT owns the next turn on this stateless surface, so never
|
|
self-post them as a new role=user prompt. Other watch events wake the session via self-post.
|
|
"""
|
|
from gateway.wake import deliver_wake, persist_delegation_delivery
|
|
if evt.get("type") == "async_delegation":
|
|
info = "Async delegation completion — persisting delivery row for api_server session %s (no wake turn)"
|
|
fail = "Async delegation delivery persist failed for session %s: %s"
|
|
deliver = lambda: persist_delegation_delivery(adapter, text=synth_text, session_id=raw_sid, evt=evt) # noqa: E731
|
|
else:
|
|
info = "Watch pattern notification — waking api_server session %s via self-post"
|
|
fail = "Watch notification self-post wake failed for session %s: %s"
|
|
deliver = lambda: deliver_wake(adapter, text=synth_text, session_id=raw_sid) # noqa: E731
|
|
try:
|
|
logger.info(info, raw_sid)
|
|
await deliver()
|
|
return True
|
|
except Exception as e:
|
|
logger.warning(fail, raw_sid, e)
|
|
return False
|
|
|
|
def _resolve_injection_adapter(self, platform_name: str, source=None):
|
|
"""Adapter for a synthetic-event platform: alias-aware transport resolver first (one
|
|
Platform.RELAY adapter fronts N logical platforms; native wins), literal ``p.value`` scan as
|
|
fallback for minimal runner stubs / exotic platform strings when the resolver can't run."""
|
|
from gateway.delivery import resolve_delivery_transport
|
|
if source is not None:
|
|
owner = self._transport_owner(source)
|
|
if owner is not None:
|
|
return owner[0]
|
|
if getattr(source, "delivered_via_upstream_relay", False) is True:
|
|
return self.adapters.get(Platform.RELAY)
|
|
profile = getattr(source, "profile", None)
|
|
adapters = self.adapters
|
|
if profile and profile not in ("default", getattr(self, "_primary_profile_name", None)):
|
|
adapters = (getattr(self, "_profile_adapters", None) or {}).get(profile, {})
|
|
try:
|
|
_transport = resolve_delivery_transport(Platform(platform_name), self.config, adapters)
|
|
except Exception:
|
|
_transport = None
|
|
if _transport is not None:
|
|
return _transport.adapter
|
|
return next((a for p, a in adapters.items() if p.value == platform_name), None)
|
|
|
|
async def _inject_watch_notification(
|
|
self, synth_text: str, evt: dict, *, raise_not_accepted: bool = False,
|
|
) -> Optional[bool]:
|
|
"""Inject a watch/completion notification as a synthetic message event.
|
|
|
|
Routing comes from the queued event, never the active foreground message. Returns
|
|
``True`` on adapter acceptance, ``False`` on retryable adapter failure, ``None`` with no
|
|
gateway route. Not transactional: a crash after acceptance can replay (at-least-once).
|
|
"""
|
|
from gateway.wake import WakeNotAccepted, adapter_supports_push, admit_internal_event
|
|
source = await asyncio.to_thread(self._build_process_event_source, evt)
|
|
if not source:
|
|
# API-server sessions bind the RAW X-Hermes-Session-Id key, not a structured ``agent:...`` key.
|
|
raw_sid = _raw_process_event_session_id(evt)
|
|
if raw_sid:
|
|
adapter = self.adapters.get(Platform.API_SERVER)
|
|
if adapter is not None and not adapter_supports_push(adapter):
|
|
return await self._self_post_api_server(adapter, synth_text, raw_sid, evt)
|
|
logger.debug(
|
|
"Deferring watch notification for raw session %s: no api_server adapter to self-post through",
|
|
raw_sid,
|
|
)
|
|
return False
|
|
logger.warning(
|
|
"Dropping watch notification with no routing metadata for process %s",
|
|
evt.get("session_id", "unknown"),
|
|
)
|
|
return None
|
|
platform_name = source.platform.value if hasattr(source.platform, "value") else str(source.platform)
|
|
adapter = self._resolve_injection_adapter(platform_name, source)
|
|
if not adapter:
|
|
return False
|
|
if not adapter_supports_push(adapter):
|
|
# Non-push adapter (api_server): its chat_id IS the raw session id, so handle_message would
|
|
# key the wake under a build_session_key() that never matches — self-post instead.
|
|
raw_sid = str(evt.get("origin_session_id") or "").strip() or str(source.chat_id or "")
|
|
return await self._self_post_api_server(adapter, synth_text, raw_sid, evt)
|
|
try:
|
|
metadata = {}
|
|
session_key = str(evt.get("session_key") or "").strip()
|
|
if session_key.startswith("agent:"):
|
|
metadata["gateway_session_key"] = session_key
|
|
parent_session_id = str(evt.get("parent_session_id") or "").strip()
|
|
if parent_session_id:
|
|
metadata["gateway_session_id"] = parent_session_id
|
|
synth_event = MessageEvent(
|
|
text=synth_text, message_type=MessageType.TEXT, source=source, internal=True,
|
|
message_id=str(evt.get("message_id") or "").strip() or None, metadata=metadata,
|
|
)
|
|
logger.info(
|
|
"Watch pattern notification — injecting for %s chat=%s thread=%s",
|
|
platform_name, source.chat_id, source.thread_id,
|
|
)
|
|
# Relay egress priming: post-restart routing caches are cold (they warm only on inbound), so
|
|
# replies would egress without tenant discriminators and be declined by the connector.
|
|
_prime = getattr(adapter, "prime_routing_cache", None)
|
|
if callable(_prime):
|
|
_prime(synth_event)
|
|
await admit_internal_event(adapter, synth_event)
|
|
return True
|
|
except WakeNotAccepted:
|
|
# Durable callers refund the claim; ordinary watch callers just requeue.
|
|
if raise_not_accepted:
|
|
raise
|
|
return False
|
|
except Exception as e:
|
|
logger.error("Watch notification injection error: %s", e)
|
|
return False
|
|
|
|
@staticmethod
|
|
def _completion_delivery_identity(evt: dict) -> Optional[tuple[str, str, object]]:
|
|
"""Return a producer-stable identity when one is available.
|
|
|
|
Delegation UUIDs identify one producer completion. Process session IDs include the
|
|
persisted spawn epoch so a reused ID is a distinct incarnation; legacy events without
|
|
``started_at`` are delivered undeduplicated rather than risk suppressing a real completion.
|
|
"""
|
|
evt_type = str(evt.get("type") or "")
|
|
if evt_type == "async_delegation":
|
|
producer_id = str(evt.get("delegation_id") or "")
|
|
if not producer_id:
|
|
return None
|
|
if evt.get("task_failure_notice"):
|
|
# An interim per-task notice is its own producer event: it must not mark the
|
|
# batch's final result as already delivered, nor a sibling's notice.
|
|
task_idx = ((evt.get("results") or [{}])[0] or {}).get("task_index", "")
|
|
return (evt_type, producer_id, f"task_failure:{task_idx}")
|
|
return (evt_type, producer_id, "")
|
|
if evt_type == "completion":
|
|
producer_id = str(evt.get("session_id") or "")
|
|
started_at = evt.get("started_at")
|
|
if producer_id and started_at is not None:
|
|
return (evt_type, producer_id, started_at)
|
|
return None
|
|
|
|
def _mark_completions_delivered_locked(self, identities) -> None:
|
|
"""Move identities inflight -> delivered and trim retention. Caller holds ``_completion_delivery_lock``."""
|
|
for identity in identities:
|
|
self._completion_deliveries_inflight.discard(identity)
|
|
self._completion_deliveries_delivered[identity] = None
|
|
while len(self._completion_deliveries_delivered) > self._completion_delivery_retention:
|
|
self._completion_deliveries_delivered.popitem(last=False)
|
|
|
|
def _completion_identity_seen(self, identity, *, claim: bool = False) -> bool:
|
|
"""True when ``identity`` is inflight or already delivered this gateway lifecycle.
|
|
|
|
With ``claim`` an unseen identity is atomically marked inflight (same lock hold).
|
|
"""
|
|
with self._completion_delivery_lock:
|
|
seen = (
|
|
identity in self._completion_deliveries_inflight
|
|
or identity in self._completion_deliveries_delivered
|
|
)
|
|
if claim and not seen:
|
|
self._completion_deliveries_inflight.add(identity)
|
|
return seen
|
|
|
|
async def _classify_completion_target(self, parent_session_id: str) -> str:
|
|
"""Classify an async-completion target before adapter acceptance: ``"deliver"`` (spawning
|
|
session live or compression-rotated with a live continuation; the resolver still retargets),
|
|
``"terminal"`` (parent gone for good — unknown / user boundary like /new; drop the durable row
|
|
rather than falsely ack), ``"retry"`` (DB unavailable / rotation mid-flight; release the claim)."""
|
|
from gateway.run import _USER_BOUNDARY_END_REASONS
|
|
session_db = getattr(self, "_session_db", None)
|
|
if session_db is None:
|
|
return "retry"
|
|
try:
|
|
parent = await session_db.get_session(parent_session_id)
|
|
except Exception:
|
|
logger.debug("Async-completion pre-flight parent lookup failed for %s", parent_session_id, exc_info=True)
|
|
return "retry"
|
|
if parent is None:
|
|
return "terminal"
|
|
if not parent.get("ended_at"):
|
|
return "deliver"
|
|
end_reason = str(parent.get("end_reason") or "")
|
|
if end_reason != "compression":
|
|
# Only a USER-closed session (/new, user_exit, session_switch) is unreachable; idle/timeout
|
|
# ends stay routable and the resolver retargets. Boundary set shared with the resolver.
|
|
return "terminal" if end_reason in _USER_BOUNDARY_END_REASONS else "deliver"
|
|
try:
|
|
tip_session_id = await session_db.get_compression_tip(parent_session_id)
|
|
if not tip_session_id or tip_session_id == parent_session_id:
|
|
# Rotation mid-flight: continuation not visible yet. Retry, don't drop.
|
|
return "retry"
|
|
tip = await session_db.get_session(tip_session_id)
|
|
except Exception:
|
|
logger.debug("Async-completion pre-flight tip lookup failed for %s", parent_session_id, exc_info=True)
|
|
return "retry"
|
|
if tip is None or tip.get("ended_at"):
|
|
return "retry"
|
|
return "deliver"
|
|
|
|
@staticmethod
|
|
def _settle_durable_claim(kind: str, delegation_id: str, claim_id: str) -> None:
|
|
"""Best-effort ``drop``/``release`` of a durable completion claim."""
|
|
fn_name, fail_msg = _DURABLE_CLAIM_OPS[kind]
|
|
try:
|
|
import tools.async_delegation as _ad
|
|
getattr(_ad, fn_name)(delegation_id, claim_id)
|
|
except Exception:
|
|
logger.log(logging.WARNING if kind == "complete" else logging.DEBUG, fail_msg, exc_info=True)
|
|
|
|
async def _completion_delivery_ready(self, evt: dict) -> bool:
|
|
"""Unavailable owners/transports must not spend a durable delivery attempt."""
|
|
from gateway.wake import adapter_supports_push
|
|
|
|
parent_session_id = str(evt.get("parent_session_id") or "").strip()
|
|
if parent_session_id:
|
|
verdict = await self._classify_completion_target(parent_session_id)
|
|
if verdict != "deliver":
|
|
# Definitively closed targets still need the normal terminal disposition.
|
|
return verdict == "terminal"
|
|
source = await asyncio.to_thread(self._build_process_event_source, evt)
|
|
if source is not None:
|
|
platform = source.platform.value if hasattr(source.platform, "value") else str(source.platform)
|
|
adapter = self._resolve_injection_adapter(platform, source)
|
|
else:
|
|
raw_sid = _raw_process_event_session_id(evt)
|
|
adapter = self.adapters.get(Platform.API_SERVER) if raw_sid else None
|
|
if adapter is not None and adapter_supports_push(adapter):
|
|
return False
|
|
if adapter is None:
|
|
return False
|
|
if not adapter_supports_push(adapter):
|
|
ensure = getattr(adapter, "_ensure_session_db", None)
|
|
try:
|
|
if not callable(ensure) or await asyncio.to_thread(ensure) is None:
|
|
return False
|
|
except Exception:
|
|
logger.debug("Async-completion delivery DB unavailable", exc_info=True)
|
|
return False
|
|
return True
|
|
|
|
async def _preflight_completion_delivery(self, evt: dict) -> "_CompletionClaim":
|
|
"""Claim the durable row (async delegations) and verify the target before adapter acceptance.
|
|
|
|
Adapter acceptance is not proof of delivery: the inner resolver can still fail closed inside
|
|
the pipeline after acceptance, falsely acking the durable row. Verifying first gives drops an
|
|
honest durable disposition.
|
|
"""
|
|
claim = self._CompletionClaim()
|
|
evt_type = evt.get("type")
|
|
if evt_type == "async_delegation" and not await self._completion_delivery_ready(evt):
|
|
claim.proceed, claim.early_result = False, False
|
|
return claim
|
|
# An interim per-task notice shares the batch's delegation_id but is not the durable
|
|
# completion; claiming that row here would acknowledge the FINAL result before it exists.
|
|
if evt_type == "async_delegation" and not evt.get("task_failure_notice"):
|
|
claim.delegation_id = str(evt.get("delegation_id") or "")
|
|
if claim.delegation_id:
|
|
try:
|
|
from tools.async_delegation import claim_completion_delivery
|
|
claim.claim_id = f"gateway:{id(self)}:{__import__('uuid').uuid4().hex}"
|
|
if not claim_completion_delivery(claim.delegation_id, claim.claim_id):
|
|
claim.proceed = False
|
|
return claim
|
|
except Exception as exc:
|
|
logger.warning("Could not claim durable async completion %s: %s", claim.delegation_id, exc)
|
|
claim.proceed, claim.early_result = False, False
|
|
return claim
|
|
elif evt_type != "completion":
|
|
return claim
|
|
# Background completions carry only session_key, so after /new the OLD session's notification
|
|
# would land in the NEW one. Stamped events get the async-delegation pre-flight; unstamped deliver.
|
|
parent_session_id = str(evt.get("parent_session_id") or "").strip()
|
|
if not parent_session_id:
|
|
return claim
|
|
# Pre-flight (#65838-class): adapter acceptance is NOT proof of delivery — the inner #55578 resolver
|
|
# can still fail closed inside the message pipeline AFTER the adapter accepted, which would falsely
|
|
# acknowledge the durable row as delivered. Verify the target here, before acceptance, and give
|
|
# drops an honest durable disposition.
|
|
verdict = await self._classify_completion_target(parent_session_id)
|
|
if verdict == "terminal":
|
|
if evt_type == "async_delegation":
|
|
logger.warning(
|
|
"Async delegation %s targets permanently-gone session %s; "
|
|
"terminally dropping delivery (result remains in the delegation records).",
|
|
claim.delegation_id or "<legacy>", parent_session_id,
|
|
)
|
|
if claim.claim_id:
|
|
self._settle_durable_claim("drop", claim.delegation_id, claim.claim_id)
|
|
else:
|
|
logger.warning(
|
|
"Background process %s completion targets "
|
|
"permanently-gone session %s (user boundary such as "
|
|
"/new); dropping notification (output remains available via process(action='log')).",
|
|
evt.get("session_id") or "<unknown>", parent_session_id,
|
|
)
|
|
claim.proceed = False
|
|
elif verdict == "retry":
|
|
# Transient uncertainty: tell the watcher to re-poll rather than drop or misroute.
|
|
if claim.claim_id:
|
|
self._settle_durable_claim("release", claim.delegation_id, claim.claim_id)
|
|
claim.proceed, claim.early_result = False, False
|
|
return claim
|
|
|
|
async def _deliver_completion_notification(
|
|
self, synth_text: str, evt: dict, *, sibling_claims=(),
|
|
) -> Optional[bool]:
|
|
"""Acknowledge one admitted batch, refund refusals, or release failed deliveries.
|
|
|
|
True means adapter admission, not model execution; None means deduplicated or
|
|
terminal. False remains retryable. Claims are settled together for every sibling.
|
|
"""
|
|
from gateway.wake import WakeNotAccepted
|
|
identity = self._completion_delivery_identity(evt)
|
|
claim = self._CompletionClaim()
|
|
accepted = identity_claimed = refused = False
|
|
try:
|
|
claim = await self._preflight_completion_delivery(evt)
|
|
if not claim.proceed:
|
|
return claim.early_result
|
|
if identity is not None:
|
|
if self._completion_identity_seen(identity, claim=True):
|
|
return None
|
|
identity_claimed = True
|
|
injection_result = await self._inject_watch_notification(synth_text, evt, raise_not_accepted=True)
|
|
if injection_result is not True:
|
|
return injection_result
|
|
accepted = True
|
|
if identity is not None:
|
|
with self._completion_delivery_lock:
|
|
self._mark_completions_delivered_locked((identity,))
|
|
return True
|
|
except WakeNotAccepted:
|
|
refused = True
|
|
return False
|
|
finally:
|
|
if identity_claimed and not accepted:
|
|
with self._completion_delivery_lock:
|
|
self._completion_deliveries_inflight.discard(identity)
|
|
operation = "complete" if accepted else "defer" if refused else "release"
|
|
if claim.claim_id:
|
|
self._settle_durable_claim(operation, claim.delegation_id, claim.claim_id)
|
|
for sibling, claim_id in sibling_claims:
|
|
if claim_id:
|
|
self._settle_durable_claim(operation, sibling["delegation_id"], claim_id)
|
|
if accepted and sibling_claims:
|
|
self._record_coalesced_completion_siblings([event for event, _claim_id in sibling_claims])
|
|
|
|
@staticmethod
|
|
def _event_route_key(evt: dict, fields: tuple[str, ...]) -> tuple[str, ...]:
|
|
return tuple(str(evt.get(field) or "") for field in fields)
|
|
|
|
@staticmethod
|
|
def _format_coalesced_process_completions(entries: list[tuple[str, dict, asyncio.Future]]) -> str:
|
|
"""Build one bounded synthetic event from several redacted completions."""
|
|
from gateway.run import _redact_gateway_user_facing_secrets
|
|
lines = [
|
|
f"[IMPORTANT: {len(entries)} background processes completed for this session.",
|
|
"Treat these results as one completion batch and send at most one "
|
|
"consolidated user-facing response.",
|
|
]
|
|
shown = entries[:10]
|
|
for _text, evt, _future in shown:
|
|
session_id = str(evt.get("session_id") or "unknown")
|
|
exit_code = evt.get("exit_code")
|
|
reason = str(evt.get("completion_reason") or "exited")
|
|
# Unconditional gateway redaction floor (the producer-seam redactor is configurable). Redact
|
|
# BEFORE slicing: truncating first can leave a credential fragment the patterns miss.
|
|
output = _redact_gateway_user_facing_secrets(str(evt.get("output") or "")).strip()
|
|
if len(output) > 800:
|
|
output = f"[… truncated …]\n{output[-800:]}"
|
|
lines.append(f"\n- {session_id}: exit_code={exit_code}, reason={reason}")
|
|
if output:
|
|
lines.append(output)
|
|
omitted = len(entries) - len(shown)
|
|
if omitted:
|
|
lines.append(
|
|
f"\n- … and {omitted} more completion(s); inspect them with "
|
|
"the process tool if they affect the conclusion."
|
|
)
|
|
lines.append("If a result does not change the current conclusion, absorb it silently.]")
|
|
return "\n".join(lines)
|
|
|
|
def _record_coalesced_completion_siblings(self, events: list[dict]) -> None:
|
|
"""Extend a successful primary delivery claim to its batched siblings."""
|
|
identities = [i for i in map(self._completion_delivery_identity, events) if i is not None]
|
|
with self._completion_delivery_lock:
|
|
self._mark_completions_delivered_locked(identities)
|
|
|
|
async def _flush_process_completion_batch(self, key: tuple[str, ...]) -> None:
|
|
"""Deliver one short-window completion batch and resolve its waiters."""
|
|
current_task = asyncio.current_task()
|
|
entries: list[tuple[str, dict, asyncio.Future]] = []
|
|
delivered: Optional[bool] = False
|
|
try:
|
|
await asyncio.sleep(self._completion_notification_batch_window)
|
|
entries = self._completion_notification_batches.pop(key, [])
|
|
# Detach before delivery so a completion arriving mid-flight can schedule the next flush.
|
|
if self._completion_notification_batch_tasks.get(key) is current_task:
|
|
self._completion_notification_batch_tasks.pop(key, None)
|
|
if not entries:
|
|
return
|
|
synth_text = entries[0][0] if len(entries) == 1 else self._format_coalesced_process_completions(entries)
|
|
# A duplicate primary returns None from the dedupe seam; try the next identity so a fresh
|
|
# sibling is never discarded with it.
|
|
delivered = None
|
|
for _text, candidate_evt, _future in entries:
|
|
delivered = await self._deliver_completion_notification(synth_text, candidate_evt)
|
|
if delivered is not None:
|
|
break
|
|
if delivered is True and len(entries) > 1:
|
|
self._record_coalesced_completion_siblings([evt for _text, evt, _future in entries])
|
|
except asyncio.CancelledError:
|
|
# Shutdown cancellation: recover undetached entries and resolve every waiter as retryable.
|
|
delivered = False
|
|
if not entries:
|
|
entries = self._completion_notification_batches.pop(key, [])
|
|
raise
|
|
except Exception:
|
|
logger.exception("Coalesced process completion delivery failed")
|
|
delivered = False
|
|
finally:
|
|
# Never strand watcher futures: False = watcher retry path; None = ordinary dedupe result.
|
|
self._settle_batch_waiters(entries, delivered)
|
|
# Do not remove a newer flush task that reused the same route key.
|
|
if self._completion_notification_batch_tasks.get(key) is current_task:
|
|
self._completion_notification_batch_tasks.pop(key, None)
|
|
|
|
@staticmethod
|
|
def _settle_batch_waiters(entries, result) -> None:
|
|
for _text, _evt, future in entries:
|
|
if not future.done():
|
|
future.set_result(result)
|
|
|
|
async def _cancel_process_completion_batch_tasks(self) -> None:
|
|
"""Settle pending completion batches before adapter teardown."""
|
|
self._completion_notification_batches_stopping = True
|
|
tasks = {
|
|
task
|
|
for task in getattr(self, "_completion_notification_batch_flush_tasks", set())
|
|
if not task.done()
|
|
}
|
|
for task in tasks:
|
|
task.cancel()
|
|
if tasks:
|
|
await asyncio.gather(*tasks, return_exceptions=True)
|
|
# Defensive cleanup for an orphaned queue with no live flush task.
|
|
batches = getattr(self, "_completion_notification_batches", {})
|
|
for entries in batches.values():
|
|
self._settle_batch_waiters(entries, False)
|
|
batches.clear()
|
|
getattr(self, "_completion_notification_batch_tasks", {}).clear()
|
|
getattr(self, "_completion_notification_batch_flush_tasks", set()).clear()
|
|
|
|
async def _enqueue_process_completion_notification(self, synth_text: str, evt: dict) -> Optional[bool]:
|
|
"""Fan in concurrent process completions that share one conversation."""
|
|
# Lazy defaults: lifecycle tests build GatewayRunner via object.__new__.
|
|
for attr, default in (
|
|
("_completion_notification_batches", dict), ("_completion_notification_batch_tasks", dict),
|
|
("_completion_notification_batch_flush_tasks", set),
|
|
("_completion_notification_batch_window", lambda: 0.1),
|
|
("_completion_notification_batches_stopping", lambda: False), ("_background_tasks", set),
|
|
):
|
|
if not hasattr(self, attr):
|
|
setattr(self, attr, default())
|
|
if self._completion_notification_batches_stopping:
|
|
return False
|
|
key = self._event_route_key(evt, self._COMPLETION_BATCH_KEY_FIELDS)
|
|
future = asyncio.get_running_loop().create_future()
|
|
self._completion_notification_batches.setdefault(key, []).append((synth_text, evt, future))
|
|
if key not in self._completion_notification_batch_tasks:
|
|
task = asyncio.create_task(self._flush_process_completion_batch(key))
|
|
self._completion_notification_batch_tasks[key] = task
|
|
# Keep the flush alive under the gateway's normal lifecycle accounting.
|
|
self._retain_background_task(task)
|
|
self._track_task_in(self._completion_notification_batch_flush_tasks, task)
|
|
return await future
|
|
|
|
def _enrich_async_delegation_routing(self, evt: dict) -> None:
|
|
"""Fill platform/chat_id/thread_id/chat_type on an async-delegation event.
|
|
|
|
Such events only carry ``session_key`` (the daemon worker lacks per-message routing
|
|
metadata). Best-effort: a CLI-origin event (empty session_key) is left as-is and won't route.
|
|
"""
|
|
from gateway.run import _parse_session_key
|
|
if evt.get("platform"):
|
|
return # already enriched
|
|
parsed = _parse_session_key(evt.get("session_key", "") or "")
|
|
if not parsed:
|
|
return
|
|
evt["platform"] = parsed.get("platform", "")
|
|
evt["chat_type"] = parsed.get("chat_type", "")
|
|
evt["chat_id"] = parsed.get("chat_id", "")
|
|
if parsed.get("thread_id"):
|
|
evt["thread_id"] = parsed["thread_id"]
|
|
|
|
async def _deliver_async_delegation_group(self, group: list[dict]) -> Optional[bool]:
|
|
"""Deliver a same-session batch of async completions as ONE turn: the primary carries the
|
|
consolidated text of every sibling THIS runner claimed (siblings owned elsewhere are excluded;
|
|
their claims are acked only after adapter acceptance). True after acceptance, False to requeue
|
|
the group, None when nothing is deliverable here (retry siblings requeued)."""
|
|
from gateway.run import _format_gateway_process_notification
|
|
from tools.process_registry import process_registry as _pr
|
|
deliverable: list[tuple[dict, str]] = []
|
|
for evt in group:
|
|
synth_text = _format_gateway_process_notification(evt)
|
|
if not synth_text:
|
|
continue
|
|
identity = self._completion_delivery_identity(evt)
|
|
if identity is not None and self._completion_identity_seen(identity):
|
|
continue
|
|
deliverable.append((evt, synth_text))
|
|
if not deliverable:
|
|
return None
|
|
if len(deliverable) == 1:
|
|
evt, synth_text = deliverable[0]
|
|
return await self._deliver_completion_notification(synth_text, evt)
|
|
# Check the entire group before claiming ANY row: an unavailable sibling must
|
|
# not exhaust its budget just because the primary has a usable route.
|
|
for evt, _text in deliverable:
|
|
if not await self._completion_delivery_ready(evt):
|
|
return False
|
|
from tools.async_delegation import claim_event_delivery
|
|
primary_evt, primary_text = deliverable[0]
|
|
blocks = [primary_text]
|
|
siblings: list[tuple[dict, str]] = []
|
|
for evt, synth_text in deliverable[1:]:
|
|
claim_id = claim_event_delivery(evt, f"gateway-batch:{id(self)}")
|
|
if claim_id is None:
|
|
# Another consumer owns this row: keep it out of our text so it is never double-injected.
|
|
continue
|
|
siblings.append((evt, claim_id))
|
|
blocks.append(synth_text)
|
|
if not siblings:
|
|
return await self._deliver_completion_notification(primary_text, primary_evt)
|
|
header = (
|
|
f"[IMPORTANT: {len(blocks)} background subagent delegations "
|
|
"completed for this session. Treat these results as one "
|
|
"completion batch and send at most one consolidated user-facing "
|
|
"response. If a result does not change the current conclusion, absorb it silently.]"
|
|
)
|
|
consolidated = "\n\n".join([header, *blocks])
|
|
delivered = await self._deliver_completion_notification(
|
|
consolidated, primary_evt, sibling_claims=siblings,
|
|
)
|
|
if delivered is None:
|
|
# Primary dropped/owned elsewhere: retry the unadmitted siblings.
|
|
for evt, _claim_id in siblings:
|
|
_pr.completion_queue.put(evt)
|
|
return delivered
|
|
|
|
async def _async_delegation_watcher(self, interval: float = 2.0) -> None:
|
|
"""Drain async completions and pattern notifications even while sessions are idle.
|
|
|
|
Background subagents and process pattern events have no per-process notification
|
|
consumer; both must progress without a later foreground turn.
|
|
"""
|
|
await asyncio.sleep(3) # let platforms finish connecting
|
|
from tools.process_registry import process_registry as _pr
|
|
while self._running:
|
|
with _log_suppressed(logging.DEBUG, "Async delegation watcher error: %s"):
|
|
# Pattern events also need an idle consumer; foreground turns are optional.
|
|
await self._drain_watch_notifications(_pr.completion_queue)
|
|
# Process completions remain owned by their per-process watchers.
|
|
requeue = []
|
|
async_events = []
|
|
while not _pr.completion_queue.empty():
|
|
try:
|
|
evt = _pr.completion_queue.get_nowait()
|
|
except Exception:
|
|
break
|
|
(async_events if evt.get("type") == "async_delegation" else requeue).append(evt)
|
|
for evt in requeue:
|
|
_pr.completion_queue.put(evt)
|
|
# A fan-out finishing together yields N completions for one session; group by full route +
|
|
# parent session so each group becomes ONE consolidated turn.
|
|
# A same-tick drain often carries several completions for the SAME originating session (a
|
|
# fan-out of background subagents finishing together). Events for different sessions never
|
|
# coalesce. See #70300.
|
|
groups: dict[tuple[str, ...], list[dict]] = {}
|
|
for evt in async_events:
|
|
self._enrich_async_delegation_routing(evt)
|
|
groups.setdefault(self._event_route_key(evt, self._ASYNC_GROUP_KEY_FIELDS), []).append(evt)
|
|
for group in groups.values():
|
|
try:
|
|
delivered = await self._deliver_async_delegation_group(group)
|
|
if delivered is False:
|
|
for evt in group:
|
|
_pr.completion_queue.put(evt)
|
|
except Exception as e:
|
|
for evt in group:
|
|
_pr.completion_queue.put(evt)
|
|
logger.error("Async delegation injection error: %s", e)
|
|
await asyncio.sleep(interval)
|
|
|
|
@staticmethod
|
|
def _redacted_output_tail(session, limit: int) -> str:
|
|
"""Last ``limit`` chars of process output through the secret redactors (unconditional floor)."""
|
|
from gateway.run import _redact_gateway_user_facing_secrets
|
|
new_output = session.output_buffer[-limit:] if session.output_buffer else ""
|
|
if new_output:
|
|
from agent.redact import redact_terminal_output
|
|
new_output = redact_terminal_output(new_output, getattr(session, "command", "") or "")
|
|
# redact_terminal_output() is unforced (raw when security.redact_secrets is off); this goes
|
|
# straight to the adapter, so apply the same unconditional floor as agent-notify.
|
|
new_output = _redact_gateway_user_facing_secrets(new_output)
|
|
return new_output
|
|
|
|
async def _send_watcher_message(self, platform_name: str, chat_id, thread_id, message_text: str, watcher: dict) -> None:
|
|
from gateway.run import _non_conversational_metadata
|
|
source = await asyncio.to_thread(self._build_process_event_source, watcher)
|
|
adapter = self._resolve_injection_adapter(platform_name, source)
|
|
if adapter and chat_id:
|
|
with _log_suppressed(logging.ERROR, "Watcher delivery error: %s"):
|
|
send_meta = {"thread_id": thread_id} if thread_id else None
|
|
await adapter.send(
|
|
chat_id, message_text, metadata=_non_conversational_metadata(send_meta, platform=platform_name),
|
|
)
|
|
|
|
@staticmethod
|
|
def _build_process_completion_event(watcher: dict, session, session_id: str) -> dict:
|
|
"""Build the synthetic ``completion`` event for an agent-notify watcher."""
|
|
from gateway.run import _redact_gateway_user_facing_secrets
|
|
from agent.redact import redact_terminal_output
|
|
from tools.ansi_strip import strip_ansi
|
|
_command = getattr(session, "command", "") or ""
|
|
_raw = strip_ansi(session.output_buffer) if session.output_buffer else ""
|
|
_raw = redact_terminal_output(_raw, _command)
|
|
# Keep the last ~2000 chars snapped to a line boundary, with a marker when cut.
|
|
_LIMIT = 2000
|
|
# Truncate at line boundaries so notifications never start mid-line (fixes #23284). Keep the last
|
|
# ~2000 chars but snap to the nearest preceding newline, then prepend a truncation marker when
|
|
# output was cut.
|
|
if len(_raw) > _LIMIT:
|
|
_tail = _raw[-_LIMIT:]
|
|
_nl = _tail.find("\n")
|
|
_tail = _tail[_nl + 1:] if _nl != -1 else _tail
|
|
_out = f"[… output truncated — showing last {len(_tail)} chars]\n{_tail}"
|
|
else:
|
|
_out = _raw
|
|
return {
|
|
"type": "completion",
|
|
"session_id": session_id,
|
|
**{k: watcher.get(k, "") for k in _WATCHER_ROUTE_FIELDS},
|
|
"message_id": str(watcher.get("message_id") or "").strip() or None,
|
|
"started_at": getattr(session, "started_at", None),
|
|
"command": _redact_gateway_user_facing_secrets(_command),
|
|
"exit_code": session.exit_code,
|
|
"completion_reason": getattr(session, "completion_reason", "exited"),
|
|
"termination_source": getattr(session, "termination_source", ""),
|
|
"output": _redact_gateway_user_facing_secrets(_out),
|
|
# Spawning session-db id: lets pre-flight drop this completion if the user /new'd first.
|
|
"parent_session_id": (
|
|
watcher.get("parent_session_id") or getattr(session, "parent_session_id", "") or ""
|
|
),
|
|
}
|
|
|
|
def _format_process_final_message(self, session_id: str, session, notify_mode: str) -> str:
|
|
from gateway.run import _format_concise_process_notification, _redact_gateway_user_facing_secrets
|
|
new_output = self._redacted_output_tail(session, 1000)
|
|
if notify_mode != "concise":
|
|
return (
|
|
f"[Background process {session_id} finished with exit code {session.exit_code}~ "
|
|
f"Here's the final output:\n{new_output}]"
|
|
)
|
|
_started = getattr(session, "started_at", None)
|
|
_dur = max(0.0, time.time() - _started) if isinstance(_started, (int, float)) else None
|
|
return _format_concise_process_notification(
|
|
session_id, _redact_gateway_user_facing_secrets(getattr(session, "command", "") or ""),
|
|
session.exit_code, new_output, duration_seconds=_dur,
|
|
)
|
|
|
|
async def _run_process_watcher(self, watcher: dict) -> None:
|
|
"""Poll a background process and push updates until it exits. Mode
|
|
(``display.background_process_notifications``): concise (default one-liner; failures append
|
|
the output tail) / all (running updates + final raw) / result (final raw) / error (final raw
|
|
if exit != 0) / off."""
|
|
from tools.process_registry import process_registry
|
|
from tools.process_registry_notifications import format_process_notification
|
|
session_id = watcher["session_id"]
|
|
interval = watcher["check_interval"]
|
|
platform_name = watcher.get("platform", "")
|
|
chat_id = watcher.get("chat_id", "")
|
|
thread_id = watcher.get("thread_id", "")
|
|
agent_notify = watcher.get("notify_on_complete", False)
|
|
notify_mode = self._load_background_notifications_mode()
|
|
logger.debug("Process watcher started: %s (every %ss, notify=%s, agent_notify=%s)",
|
|
session_id, interval, notify_mode, agent_notify)
|
|
silent = notify_mode == "off" and not agent_notify
|
|
last_output_len = 0
|
|
while True:
|
|
await asyncio.sleep(interval)
|
|
session = process_registry.get(session_id)
|
|
if session is None:
|
|
break
|
|
if silent:
|
|
# Still wait for the process to exit so we can log it, but don't push any messages.
|
|
if session.exited:
|
|
break
|
|
continue
|
|
current_output_len = len(session.output_buffer)
|
|
has_new_output = current_output_len > last_output_len
|
|
last_output_len = current_output_len
|
|
if session.exited:
|
|
# Agent-notify: inject a synthetic message unless the agent already consumed the result via
|
|
# wait/log (poll() is read-only and deliberately does NOT mark consumed).
|
|
if agent_notify and not process_registry.is_completion_consumed(session_id):
|
|
completion_evt = self._build_process_completion_event(watcher, session, session_id)
|
|
synth_text = format_process_notification(completion_evt)
|
|
if not synth_text:
|
|
break
|
|
delivered = await self._enqueue_process_completion_notification(synth_text, completion_evt)
|
|
if delivered is False:
|
|
# The process remains terminal; retry after failed adapter injection instead
|
|
# of suppressing the result.
|
|
continue
|
|
break
|
|
# Text-only notification; skip when already consumed via wait/log (the agent_notify branch
|
|
# FALLS THROUGH here, hence the re-check).
|
|
if process_registry.is_completion_consumed(session_id):
|
|
logger.debug(
|
|
"Process watcher: completion for %s already consumed "
|
|
"via wait/log — skipping raw notification (#65379)", session_id,
|
|
)
|
|
break
|
|
if notify_mode in {"concise", "all", "result"} or (
|
|
notify_mode == "error" and session.exit_code not in {0, None}
|
|
):
|
|
message_text = self._format_process_final_message(session_id, session, notify_mode)
|
|
await self._send_watcher_message(platform_name, chat_id, thread_id, message_text, watcher)
|
|
break
|
|
elif has_new_output and notify_mode == "all" and not agent_notify:
|
|
# New output — deliver a status update (only in "all" mode; agent_notify watchers
|
|
# only care about completion).
|
|
new_output = self._redacted_output_tail(session, 500)
|
|
await self._send_watcher_message(
|
|
platform_name, chat_id, thread_id,
|
|
f"[Background process {session_id} is still running~ New output:\n{new_output}]", watcher,
|
|
)
|
|
logger.debug("Process watcher ended%s: %s", " (silent)" if silent else "", session_id)
|