mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
* release-manager: fencing-token mutex + executor/readiness hardening
Closes the release-mutex TTL race (#17, HIGH) and the remaining
release-manager gaps (#88, #89, #201, #202):
- #17: the release mutex is now acquired with a uuid4 fencing token and
released via Lua compare-and-del; a background asyncio heartbeat
compare-and-expires the TTL ~every 60s while the execute owns the lock,
so a live execute no longer expires and a crashed one auto-releases
<=3000s. A second approve after TTL expiry cannot usurp and rm -rf the
in-flight clone — the fenced first-finally keeps its lock.
- #89: a Redis outage during acquire stays fail-closed (the execute never
runs) but now returns a distinct redis_unavailable result + log so the
CEO sees the cause instead of a false already_in_progress.
- #88: commit_and_push RuntimeError is wrapped into a structured
ReleaseResult(commit_failed) instead of a 500.
- #201: first-release fallback still emits untracked version-ref files as
gaps (no longer silenced by the first-release branch).
- #202: _await_proc awaits proc.wait() after kill() so a timeout cannot
leak a zombie.
TDD: tests/unit/services/test_release_proposal_concurrency.py extends
_FakeRedis with eval/get/expire and pins the fencing/heartbeat/usurper
invariants + the redis_unavailable result.
* PM/code-task creation guard + main_pm coverage + issue carve-out
Closes the creation-time role x task_type gap (the user's explicit example)
and the main_pm delegate hole:
- New pure helper `pm_cannot_own_code(role, task_type, is_issue_resolution)`
in foundation/policy/batch.py — single source of truth. Both PM roles
(cell_pm + main_pm) coordinate; a `code` task assigned/claimed by a PM is
a structural mismatch, EXCEPT a PM taking a code task in needs_revision to
resolve review/QA issues directly (the carve-out).
- Creation-time guard: TaskService.create calls the helper (closes the
create-with-cell-PM-assignee hole the team-based check misses).
- Delegate path + spec claim gate consult the same helper.
`_validate_assignee_task_type` / `_task_type_hint_for` now key on the
Role (CELL_PM OR MAIN_PM), not the cell-PM slug set — closes the
delegate-to-main-pm-as-code hole.
- identity.role_for_uuid_or_none is None-tolerant (treats None as "not a
PM" and proceeds) so a malformed/missing assignee cannot crash the guard.
- prompter.create_task_from_draft reuses the guard at draft-create.
TDD: test_batch.py (helper matrix + carve-out), test_main_pm_code_guard.py
(main_pm coverage), test_delegate_assignee_task_type.py (delegate parity),
test_lifecycle_spec.py (claim gate: rejects PM claiming code from pending,
allows from needs_revision + PM claiming planning + dev claiming code).
* task-service: completion hooks + escalation/cancel/audit hardening
Closes the task-service cluster (#21/#98, #99, #100, #101, #103, #216; #102
verified already-covered, #217 verified already-guarded):
- #21/#98: ceo_approve now closes the work session + triggers completion
hooks before worktree removal (no-op when work_session_id is None), so a
CEO-approved task lands the same close-path as PM-completed.
- #99: apply_escalation routes through the transition validator with an
enumerated escalation exemption (_ESCALATABLE_TO_BLOCKED) instead of an
arbitrary source->BLOCKED write; BACKLOG is refused.
- #100: branchless ceo_reject awaiting_ceo_approval->pending gets a real
spec edge (ceo_reject_to_pool ActionSpec + _STATUS_TRANSITIONS entry) so
future admin-override tightening can't wedge the path.
- #101: revision_count bump is documented as the single chokepoint, with
the pre-block RESTORE path undoing it when restoring a snapshotted
needs_revision (same cycle resuming, not a new rejection).
- #103: cancel cascade surfaces non-terminal orphans instead of swallowing
the role violation.
- #216: _remove_task_worktree_on_terminal escalates recurring FS/permission
failure (audit/notify after N) instead of silent-failing forever.
- #102: pinned in test_verb_runner_midverb_invalid_state.py (committed with
the choreographer cluster) — verb-runner savepoints already surface a
concurrent mid-verb state change as INVALID_STATE.
- #217: submit_for_qa claimed_by guard verified intact.
TDD: test_task.py, test_worktree_cleanup_on_complete.py,
test_escalation_board_guard.py (#99), test_task_service_* integration,
test_lifecycle_spec.py.
* choreographer: gate-claim guards + pr-gate hardening + fail-open logging
Closes the choreographer cluster (#5/#222, #29, #30, #82, #188, #189,
#192; #157/#187 verified already-fixed/pinned; #102 pin lives here):
- #5/#222: the unchanged-PR guard's fail-open head_sha lookup now logs
(warning) on a slug-resolver/git-helper error so a regression cannot
silently turn the pr_fail re-submit loop-stopper into a no-op. Stays
fail-open (never wedges the PM).
- #29: pinned (REFUTED-with-pin) — _lane_claim_guard already returns the
error envelope without releasing the claim on a transient lookup error.
- #30: pinned (REFUTED-with-negative-pin) — a non-batch branchless main_pm
root cannot bypass the complete spec gate (is_batch_umbrella requires
batch_id set).
- #82: _post_gate_review_to_pr wraps the slug-resolution call in try/except
(mirrors _capture_pr_head_sha) so a malformed cell_map AttributeError no
longer 500s the reviewer after a committed gate transition.
- #188: _is_hand_formatted_verdict anchors the header regex to line-start,
so a quoted (> ## Summary) or inline (mid-prose) header mention no longer
false-refuses a hand-formatted verdict.
- #189: pr_fail re-captures the PR head SHA after the transition commits and
re-stamps the verdict note only when it advanced (closes the stale-SHA
false-allow loop-hole); no-advance stays a single note write.
- #192: claim_gate_review skips the dev claim guards (already_active/paused/
lane) via a new skip_dev_guards param — a pr_reviewer inspecting an
assembled PR does not start work, so the single-active-task / code-lane
invariants do not apply; the dependency guard is kept, and QA's
claim_review parity is preserved.
- #157/#187: verified in tree — pr_review-only handoff is intentionally
prior-work-worth-resuming; self_review_block wiring (reviewer != dev)
holds on assembled tasks with 4 existing pin tests.
TDD: test_choreographer_*, test_pr_gate_posts_review (#82),
test_pr_review_hand_format_guard (#188), test_submit_root_unchanged_pr_guard
(#189), test_claim_gate_review_guards (#192), test_verb_runner_midverb
_invalid_state (#102 pin).
* playbook curate: guard the gating commit against a poisoned session (#55)
The explicit `session.commit()` that gates the RAG index (commit-before-index
so an uncommitted playbook cannot land in the corpus) raised PendingRollbackError
when a prior mid-verb failure had rolled the caller's session back — 500ing the
whole curation verb instead of returning a clean envelope, and (worse) risking a
fall-through to index an uncommitted playbook. Wrap the commit: on
PendingRollbackError, log + return invalid_state with a re-fetch/retry remediate
and skip the index. The happy path still commits exactly once then indexes.
TDD: test_playbook_verbs.py — poisoned-session returns a clean invalid_state and
does NOT index; clean-session still commits once + indexes (pins no fail-closed
inversion / no double-commit).
* [chore] gateway: atomic activate merge — preserve probe_failures across re-park (#156)
activate() was a blind SET that reset probe_failures to 0, so a probe-failure
increment that just landed (or was in flight) could be wiped by a concurrent
re-park — resetting the give-up / CEO-notify count mid-episode. Route activate
through a server-side Lua merge (roboco:activate_rate_limit) that refreshes the
episode metadata (kind / activated_at / retry_after / affected_agents) while
carrying over the previous probe_failures count. Indivisible w.r.t. the
increment/reset scripts (Redis single-threads an EVAL).
#56 (notify to prompter/secretary refused) verified SAFE — the pin tests
(test_notify_rejects_prompter_recipient / _secretary_recipient /
_allows_ceo_recipient) already cover the only human notify target invariant;
no legitimate send is dropped, no code change.
* [chore] foundation/policy: spec gates + QA retry-key pin (Cluster F)
#50 sync_branch composes=() so the spec gate accepted a terminal/paused/
blocked task and the handler rebased a dead/parked branch — add a
PRECONDITION_SYNC_BRANCH_STATE (claimed/in_progress/verifying/needs_revision
only), rejection_kind=invalid_state. TDD: 28 spec tests.
#148 submit_root's prose asserts 'a Main-PM root is planning-typed, never
code' but only the creation path (main_pm_cannot_own_code) backed it — add
PRECONDITION_ROOT_NOT_CODE on the submit_root IntentSpec (defense in depth),
scoped to submit_root only so the shared submit_for_review action keeps
cell_pm+code submit_up parity. Graceful on Mock/None task_type so the
choreographer Mock-task tests don't crash. TDD: 2 spec tests.
#150 VERB_RETRY_LIMITS is keyed by the MCP-exposed names (pass/fail), not
the IntentSpec-internal pass_review/fail_review — already correct; add a
pin test so a one-sided rename can't silently drop the QA-handoff cap.
#142 main_pm_cannot_own_code/pm_cannot_own_code already normalize casing
(.lower()) — no-op, pin test test_main_pm_cannot_own_code_is_case_insensitive
already in tree.
* [chore] worksession-git: 405 merge-method fallback + non-destructive close (Cluster W)
#108 _merge_with_retry hardcoded 'squash' and raised MergeConflictError on a
405 with no method fallback — wedging the PM on an open, mergeable PR whose
repo merely had the squash button off. Add a 405 fallback to a permitted
method (via _first_allowed_merge_method, exclude='squash'), mirroring the CEO
merge_pull_request path. A 405 with no permitted fallback (or a second 405)
still falls through to the already-merged disambiguation / MergeConflictError.
TDD: 2 new tests (fallback-success, no-permitted-method-raises).
#109 close_pull_request defaulted delete_branch=True, so the choreographer
supersede path deleted a superseded PR's branch while the orchestrator
supersede path explicitly preserved it — the two disagreed, and the
destructive default ran on the 'close the dead PR' path where the branch may
still be referenced / useful for audit. Flip the default to False (opt-in
deletion) and make the choreographer caller explicit (parity with the
orchestrator). TDD: 1 new test (default preserves branch); existing
deletion-when-requested test now passes delete_branch=True explicitly.
Dispositions verified against current code (no silent drops):
- #27 REFUTED/FIXED-UNDEPLOYED: work_session.merge_pr resolves by session_id
(no global pr_number lookup); the real cross-repo collision fix
(project_id scoping on pr_merge/close_pull_request/rebase_pr_for_task/
pr_target) is already in tree + tested (test_pr_merge_scopes_task_lookup_
by_project_id, test_close_pull_request_scopes_task_lookup_by_project_id,
test_git_pr_target_scoping). Verify-only.
- #106 REFUTED: a guard exists (rev-list --count {base_ref}..{branch} == 0)
before reset --hard + base_ref falls back to default_branch; tests lock
the safety (test_create_branch_never_repoints_branch_with_real_work,
test_create_branch_does_not_reset_or_checkout_shared_clone).
- #218 BY-DESIGN: the merge_pr idempotent guard intentionally preserves the
audit trail (docstring + test_merge_pr_idempotent_on_already_completed_
preserves_audit_trail); a COMPLETED session always carries attribution
(COMPLETED only via merge_pr), so the NULL-COMPLETED case is unreachable.
- #104 BY-DESIGN: agents never merge to the repo default branch in RoboCo's
model (root→master is CEO-only); the guard is a correct CEO-only rail,
locked by test_pr_merge_into_default_branch_is_ceo_only.
* [chore] llm: surface disabled-provider downgrade + scrub probe log (#20/#3/#211)
#20/#3 resolve_for_agent silently fell through to the legacy Anthropic path
when a configured provider was disabled — indistinguishable from 'no
assignment', so the operator got no signal that spawns bypassed the
provider. Surface the bypass with a warning (graceful degradation stays the
default — a stalled spawn is worse than a routing miss) and add an opt-in
ROBOCO_ROUTING_STRICT (default-off) that fail-closes instead. Wired into the
panel Feature Flags card. TDD: 3 unit tests (warn-on-disabled, strict-raises,
no-assignment-stays-silent).
#211 probe_ollama_tags logged str(exc) raw on the generic-exception branch —
structured log could carry connection internals / stack traces. Log the
exception class name only. TDD: existing generic-branch test strengthened to
assert the log kwargs don't leak the raw text.
* [chore] support/stream/optimal/playbook/comms hardening (Cluster S)
Logical-gaps sweep, Cluster S (TDD, red→green per item):
#64 notification_delivery.acknowledge published the NOTIFICATION_ACKED bus
event directly (bypassing the outbox) — a rollback left a phantom ACK. Route
it through defer_bus_publish (after_commit), mirroring deliver.
#76 playbook.archive()/reject() stamped the archiver into approved_by/
approved_at, overwriting approval provenance (and fabricating approval for a
rejected draft). Add archived_by/archived_at (migration 053 + table + model)
and write those on archive/reject, leaving approval attribution intact.
#181 vector_store.replace_chunks wiped existing index rows even when every
chunk lacked an embedding (embedder failure). Skip the wipe when chunks is
non-empty but records is empty — preserve good rows for nothing.
#182/#183 optimal.record_learning recomputed a learn-{md5(full_content)}
tracking source that never matched the URI the plugin embedded chunks under
(roboco://learnings/{doc_id}, doc_id=lrn-{hash100}). Use the plugin's
returned doc_id so de-index/lookup-by-source finds the chunk rows.
#96/#97 transcription periodic flush only peeked ready buffers (unbounded
map growth) and ran sync callbacks on the event loop (a slow callback
blocked the flush task). Flush (remove) each ready buffer after notifying,
and offload each callback to a thread.
#212 _TEAM_SCOPED_ROLES was duplicated across communications/agents_config/
seeds. Single-source it in foundation.policy.communications; consumers
reference that object (identity-tested).
#19 stream_bus._dispatch_event re-ran already-succeeded handlers on a
recover_pending replay (duplicate side effects). Add a per-(event.id,
handler) SET-NX idempotency guard: skip on a hit, clear the key on handler
failure so a replay re-runs it, fail-open when redis is unavailable.
Dispositions (no code change): #77 approve() index-write pair asserted
BY-DESIGN; #62/#63 notification DB-dedup verified pinned; #184/#185 REFUTED;
#214 REFUTED; #215 BY-DESIGN.
* [chore] db/migrations: graph-integrity guard + conftest unreachable-DB warning (Cluster D)
Logical-gaps sweep, Cluster D (TDD + real alembic upgrade head verification):
#16/#37 add tests/unit/test_migration_graph_integrity.py — a static guard that
the alembic migration graph has exactly one head, every down_revision resolves,
every revision is reachable from a root, and no revision id is duplicated. The
suite builds its DB via Base.metadata.create_all (not alembic upgrade head), so
a forked head / dangling down_revision / duplicate id would otherwise ship
silently and break a real deploy mid-stream.
Caught a real bug in the process: migration 053's revision id
"053_playbook_archived_attribution" (33 chars) exceeded alembic's
alembic_version.version_num VARCHAR(32) — a fresh `alembic upgrade head` raised
"value too long for type character varying(32)" at the 053 stamp. Renamed to
"053_playbook_archived_attr" (26 chars). Verified end-to-end on a scratch PG:
upgrade head stamps 053, downgrade -1 returns to 052. (The pre-existing
test_every_migration_revision_id_fits_the_alembic_version_column guard is now
green too; it had been red on the 33-char id.)
#90 conftest silently pytest.skip'd every DB test when Postgres was unreachable
— a non-Docker box reported a green run of all-skips. Extract the warning into
_warn_if_pg_unavailable and fire it at import so the operator sees the DB is
down (the per-test skip path is unchanged). Test: warns when unavailable, silent
when reachable (verified under -W error::UserWarning).
Dispositions (verified against real code + a fresh alembic upgrade head, no code
change): #6 REFUTED — sa.Enum(create_type=False) at 001:119/304 does NOT break a
fresh upgrade head (001→052 applied cleanly on a scratch DB); #8 REFUTED — the
upgrade passed 030/031 (RAG chunk tables) without pgvector installed; pgvector is
a runtime concern handled by roboco/db/base.py, not a migration prerequisite;
#40 REFUTED — the `|| echo` mask was already removed and partial-schema drift
reports exit 1 (only by-design unreachable/unmigrated skips remain); #204 REFUTED
— the property walk seed IS pinned (random.Random(20260504), line 97); #205/#206
REFUTED — the smoke-trace fixture IS wired via
test_lifecycle_smoke_replay.py (8 passed); no shell smoke scripts exist in the
tree to wire; #137 BY-DESIGN — pyproject version 0.14.0 is an operational note,
no code gate.
* [chore] panel: admin-override force flag + kanban subtask_count + ws cleanup + ui-store dedupe (Cluster P)
#13: kanban admin-override into a hatch state (completed / awaiting_qa /
awaiting_pm_review) now requires an explicit force=true from the panel and
emits a dedicated task.admin_override audit row server-side; non-hatch
overrides need no force. Backend gate in tasks route + admin_set_status;
panel kanban-board sends force for hatch targets; TaskUpdate carries force.
#198: kanban service threads the real subtask_count (one grouped query) into
dev + priority-swimlane + main-pm-flat boards instead of a hardcoded 0.
#79: useWebSocket cleanup clears messages/lastMessage/state on unmount or
endpoint change so a dep-change (navigating to another stream) can't leak the
prior subscription's stale snapshot as live.
#186: disambiguate the duplicate ui-store modules -- the session/scroll store
in lib/stores renamed to useScrollRestorationStore / scroll-restoration-store
(barrel + 2 consumers updated); the sidebar/theme useUIStore in @/store is now
the sole useUIStore.
#12: verified already in-tree (release-proposal-card surfaces non-404 errors
with retry; getProposal maps only 404->null). #80 by-design (handleTransportError
already resets isSending on a no-payload SSE drop). #81 docs (streamUrl docstring
records that live-intake SSE auth is session-id-based bearer-style).
Backend: ruff+mypy clean, 278 tests green. Panel: lint+typecheck clean, 159 tests.
* orchestrator: park/reaper/readopt/a2a hardening + self-heal/ci-watch dedupe (Cluster O)
Closes the orchestrator-side logical gaps from the sweep:
- #75 a2a human-only drop surfaced: _dispatch_a2a_work logs the skip
("a2a request targets a human-only role; left as a notification for the
human (not spawned)") instead of silently dropping the target — the
CEO/secretary/prompter still see the notification; only the spawn is
suppressed. (orchestrator.py)
- #72 readopt liveness: _readopt_running_agents requires a non-stale live
claim (via _agent_holds_live_claim) and skips a zombie container so a
reaped-but-restart-readopted agent isn't double-counted as active.
- #74 shutdown drain: stop() calls _flush_respawn_tracker so the durable
respawn counter write-throughs aren't lost on a clean stop.
- #71 resolve_wait active-guard + deferred liveness: a rate_limit_lifted
WaitingRecord is only confirmed-live after a _confirm_resume_liveness
probe (deferred deletion _resume_confirm_delay=30.0), and an
already-active agent short-circuits the repark. Scoped to
rate_limit_lifted records (the only ones at risk of a false lift).
- #73 stuck-Claude kill: _maybe_kill_stuck_claude + _claude_stuck_kill_ttl
(config.claude_stuck_kill_seconds) — a live container whose heartbeat is
stale past the grace AND whose gateway probe is broken is killed+evicted,
not protected forever by the reaper's live-skip.
- #230 verified FIXED-UNDEPLOYED: _gateway_broken_past_grace already
requires N consecutive false-broken probes (not one flaky streak); no
change, test added to pin the N-consecutive invariant.
- #43 self-heal per-observation dedupe: a fingerprint collapses repeat
CEO notifications for the same CI regression.
- #44 ci_watch dedupe by (git_url, workflow): a monorepo's multiple
workflows each get their own fix task (was collapsed by git_url alone).
- #49 identity.role_for_slug_or_none None-hardening: a stale/malformed
slug resolves to None and the human-only skip falls through to the safe
"not spawnable" path instead of crashing.
- #193 strategy engine: notify the CEO on a persistent assess failure
instead of failing silently in the background loop.
TDD: test_no_spawn_human_roles (a2a skip surfaced), test_orchestrator_
shutdown_drain (#74), test_provider_overload_break (#71), test_readopt_
running_agents (#72), test_resolve_wait_repark (#71), test_stale_claim_
reaper (#73/#230), test_strategy_engine_loop (#193, new),
test_self_heal_engine (#43), test_ci_watch_engine (#44), test_identity
(#49). All red->green.
* chore: make-quality green — xenon complexity refactors + mypy test fixes + lifecycle regen
No behavior changes. Brings the tree to a fully green `make quality` (the
base branch never passed the xenon B-rank gate on several blocks; the
lifecycle artifacts had drifted from the committed ceo_reject_to_pool edge).
Xenon B-rank refactors (extract a helper; preserve semantics exactly):
- api/routes/tasks.py: _apply_forced_status_override + _StatusOverride
dataclass bundle (update_task override block).
- services/task.py: _enforce_no_pm_code_on_create (create guards) +
_escalation_diverts_to_pool (collapses the two board/advisory +
main_pm+code divert branches into one predicate).
- services/prompter.py: _coerce_pm_code_to_planning (create_task_from_draft).
- services/notification.py: _duplicate_unacked_exists (_create_notification
purpose-based dedup query + ACK_REQUIRED_BY_TYPE gate).
- services/sequencing.py: _same_assignee_lane_edges (the undeclared-surface
same-assignee lane fallback at the tail of dev_task_collision_edges).
- gateway/choreographer/_impl.py: _pm_task_type_error static helper
(_validate_assignee_task_type compound PM guard).
- gateway/choreographer/pr_gate.py: _gate_review_event_verdict +
_gate_review_body static helpers (_post_gate_review_to_pr).
mypy test fixes (no type:ignore — banned; use typing.cast with quoted
strings per TC006):
- test_task_update_completeness: TaskUpdate(acceptance_criteria=None).
- test_bus: cast("Redis", _FakeRedis()); Redis import under TYPE_CHECKING.
- test_pr_merge_concurrency: capture AsyncMocks into locals before asserting.
- test_notification_delivery_phantom: cast("UUID", to_agents[0]).
Lifecycle artifact regen (owed from Cluster T #100 — the
awaiting_ceo_approval -> pending `ceo_reject_to_pool` edge was added to the
spec in 3d633084 without regenerating the derived artifacts the
foundation-check gate diffs against): docs/rag/lifecycle/intent-verbs.md,
docs/rag/lifecycle/status-transitions.md, panel/lib/lifecycle.json.
services/kanban.py: ruff format only (collapses the _load_subtask_counts
signature that drifted unformatted from Cluster P).
* [chore] logical-gaps sweep — Cluster I (intake/product/pitch)
#57/#58 prompter: preserve a top-level product_id with a 1-cell map
(prompter.py create_task_from_draft — top-level target wins over a
redundant 1-cell map instead of dropping product_id); reject — not
silently skip — a malformed project_id in the_work cell entries
(prompter.py _draft_cell_map raises ValidationError).
#59/#159 prompter: create_task_from_draft now operates on a copy
(_copy_draft) so _validate_and_coerce_draft / _clean_list never mutate
the caller's draft dict.
#160 prompter: _resolve_owning_team consults product/board routing
before forcing MAIN_PM on a multi-cell map (product root stays Board,
product+assignee-is-board stays Board).
#83/#84 github_provisioning: create_repo is idempotent by GitHub name
— a 422 "name already exists" (orphaned repo from a rolled-back prior
approval) is fetched and reused instead of erroring; pitch re-approval
now reuses the orphaned repo end-to-end.
#196 kanban: flat main-PM board has a "coordination" column for
non-cell teams (MAIN_PM/Board) instead of dropping their cards.
#197 project update: an explicit null in the PATCH body now clears the
stored field, distinct from an absent field (leave unchanged).
ProjectService.update drops exclude_none so explicit-None applies; the
PATCH route uses ProjectUpdate.model_validate(data.model_dump(
exclude_unset=True)) to preserve the request's unset-tracking (the old
field-by-field construction marked every field set and defeated the
distinction — nulling NOT-NULL git_url).
TDD: prompter 47, github_provisioning+pitch 16, kanban+project 67,
project routes 37 — all green; ruff + mypy clean.
* [chore] logical-gaps sweep — Cluster M (mcp-servers)
#60 flow_server/do_server: the circuit-breaker substitution no longer
erases the fixable rejection — the original envelope (kind/message/
remediate) is nested as inner on a copy of the SDK's circuit_open
envelope (the SDK dict is not mutated in place). The agent still sees
WHY the verb failed, not just that the breaker tripped.
#61 flow_server/do_server: a 404 carrying a *descriptive* detail
(not FastAPI's bare default {"detail":"Not Found"}) is now a
real resource not_found, surfaced as not_found so the agent
re-fetches state — instead of a misleading "server-side wiring gap"
invalid_state. The bare default and unparseable 404s still synthesize
the wiring-gap envelope; a 404 with a real Envelope (error field)
is still surfaced as-is.
#161 flow_server/do_server: dict error.code classification now uses
an exact-code map (authoritative for the codes the handlers emit) with
a substring fallback for unknown codes. Fixes the real regression:
AUTHENTICATION_REQUIRED carries no AUTHORIZED/DENIED/PERMISSION
substring, so the old substring-only rule dropped it to invalid_state
instead of not_authorized — an auth storm attributed as a state storm.
The fallback also adds AUTH so future AUTH-prefixed codes classify.
#162 flow_server/do_server: _register_tools gains a
ROBOCO_ALLOW_FULL_TOOLSET env override (default-off) so a missing
manifest falls back to the full tool set instead of raising — a
dev/test escape hatch. Production fail-loud behaviour is unchanged.
#163 intake_server: propose_batch accepts name as well as
title (intake drafts in the wild have used both), normalizing a
name-only draft onto a copy as title (caller's dict never mutated),
and reports the dropped count + reason in the return instead of
silently vanishing malformed drafts. The empty-batch hint now names
name as an alternative.
TDD: 123 mcp_servers tests green (14 new + 2 updated); ruff + mypy clean.
* [chore] Cluster N — conventions/docs logical-gaps sweep
#33: _create_new_doc/_update_existing_doc now resolve via
_resolve_contained_path (the RAG-returned update path was not containment-
checked — an escaping source could write/overwrite outside the docs dir).
#34: _commit_doc_to_repo returns committed/skipped/failed instead of
swallowing all exceptions; surfaced on DocRef.commit_status, the write
response, and the docs MCP guidance so a failed repo commit is fail-loud.
#35: write_doc only updates the similar doc when its filename matches — a
different filename creates a new file instead of collapsing onto the
similar doc's path (the dedup-overwrite defect codified by the old tests).
#129: a custom rule scoped to a language the validator never reports (a
typo) is surfaced as a warn finding on .roboco/conventions.yml via the
runner's once-per-run validation; #32 (tsx->typescript dialect) stays
BY-DESIGN.
#130: _cache_put only swallows a UNIQUE violation (23505) as a concurrent
duplicate; a non-unique IntegrityError (FK/NOT NULL/check) is log-errored
and re-raised instead of being silently misattributed.
#132: health re-reads the live file status (a cached degraded row hid an
in-place repair at a stale head key); get_map skips cached degraded rows
and stops caching degraded so a repaired file re-derives. #134 BY-DESIGN.
#133: _DB_METHODS gains stream/stream_scalars (SQLAlchemy 2.0 streaming
constructs are data access too — a route calling them is not thin).
#199: regenerate_verb_tables._annot_str strips Annotated[...] metadata
(BeforeValidator) before rendering; regenerated verbs.md + per-role
prompts so the BeforeValidator(func=...) repr (with a memory address) no
longer leaks into agent-facing prompt text.
TDD: 206 conventions/docs tests green (incl. 5 new files / appended
cases); ruff + mypy clean.
* [chore] Cluster 16 — cross-cutting hygiene logical-gaps sweep
Disposition + fix the 10 cross-cutting-hygiene gaps, TDD. make quality green
(ruff, mypy 944 files, pytest, xenon, vulture, foundation-check, enum-parity).
FIX:
- #24 /ws/system now gated by _require_panel_token (matches every sibling
/ws/* stream); rejects a missing token in strict mode and a forged token
even in dev. (roboco/api/websocket.py)
- #25 two drifted _require_ceo implementations (orchestrator router vs release
handler) unified on a single require_ceo_role helper in deps — same 403,
same role set, accepts Role/AgentRole/"ceo". (roboco/api/deps.py,
routes/orchestrator.py, routes/release.py)
- #11 a spawn session for a delivery role (developer/qa/documenter) with no
task_id now logs an unattributed-usage warning via is_unattributed_delivery_spawn.
(roboco/runtime/orchestrator.py)
- #65 pricing returns a structured CostResult(cost_usd, unpriced, is_anthropic)
so an unpriced Anthropic model (real spend we'd undercount) is flagged
instead of silently $0; calculate_cost stays a thin float wrapper.
(roboco/billing/pricing.py, billing/__init__.py)
- #67 blocker-metrics "blocked since" reads the task.blocked audit transition
(indexed on target_id/event_type/timestamp), not updated_at — which
over-counted when a blocked task was touched for a non-blocking reason.
Falls back to updated_at/created_at only with no audit row.
(roboco/services/metrics.py)
- #94 grok refresh_if_stale uses double-checked locking (_refresh_lock +
_recheck_or_refresh) so two concurrent callers don't both POST the
single-use refresh grant and burn the credential. (grok_auth.py)
DOCS (fix the doc, behavior already correct/pinned by tests):
- #66 get_summary docstring corrected — it sums raw agent_spawn_sessions rows
(sub-day precise); daily_usage_rollups/get_today_summary can diverge for
"today" until the sweeper catches up. (roboco/services/usage.py)
- #68 DashboardStorage is a documented in-memory stub; added a test pinning
that auditor flags are lost on storage reset (persisting = a migration +
service refactor, out of scope as a half-implementation).
(tests/integration/test_dashboard_service.py)
BY-DESIGN (no code change, with file:line evidence):
- #28 dashboard reads are open to the authenticated operator (dashboard.py:36
documents this); mutating auditor routes already gate via
_require_auditor_or_ceo. Role-gating reads would break the panel (no
X-Agent-ID on dashboard reads) and CEO-token-gating the router would block
the Auditor (auditor token != CEO token). nginx is the prod boundary.
REFUTED (narrowing would reintroduce a documented hang):
- #93 the ~/.grok directory mount (vs a single auth.json file) is load-bearing
— a single-file bind mount pins the inode so the atomic tmp.replace refresh
doesn't propagate to running containers (they hang at grok's login prompt).
Already documented in grok.py:161 and locked by
test_intake_grok_mounts_subscription_auth_when_present.
Incidental gate-greening (mypy errors a stale .mypy_cache had hidden in
earlier-cluster test files; xenon refactors for the new B-threshold):
- tests/unit/test_regenerate_verb_tables.py: type the dynamic-module loader.
- tests/unit/services/test_prompter.py: annotate the draft dict as dict[str,Any].
- tests/unit/services/test_conventions_cache_put.py: _FakeOrig is a real Exception
(IntegrityError's orig arg requires BaseException).
- metrics._blocked_since_map extracted from get_blocker_metrics (complexity).
- intake_server._normalize_batch_drafts extracted from propose_batch (complexity).
* Docs update
* [bug] spawn: self-heal vanished clone + branch ref before worktree ensure (be-dev-1 fatal loop)
A vanished clone_root (disk loss / /data/workspaces wipe / manual cleanup)
fatal-looped the resume path: _ensure_worktree_before_spawn ran
`git -C <missing>` and released the claim, but the reaper-style release
preserves assigned_to + branch_name so the next dispatch is a RESUME
(create_branch never re-runs to re-clone) and the same missing clone failed
every ~30s.
- workspace.py: ensure_worktree_self_heal re-attaches a present worktree +
symlinks the shared .venv; on a missing local branch ref it fetches from
origin (create_branch pushes at claim time, so pushed work survives) and
re-creates the ref, falling back to -b origin/HEAD only when the branch
was never pushed. _fetch_branch_ref is the token-aware fetch helper.
- orchestrator.py: _ensure_worktree_before_spawn health-checks the clone
and re-clones via ensure_workspace BEFORE the worktree self-heal. Fatal
git-state (WorkspaceError) still releases the claim + aborts; transient
failures abort without releasing (a fresh claim wouldn't help and
re-cloning is destructive).
TDD: 21 new + 61 related worktree/git/cancel/cleanup tests green; ruff +
mypy clean.
---------
Co-authored-by: Renn F <rennf93@users.noreply.github.com>
699 lines
26 KiB
Python
699 lines
26 KiB
Python
"""
|
|
WebSocket Handlers
|
|
|
|
Real-time communication via WebSocket connections for:
|
|
- Channel streams (all messages in a channel)
|
|
- Agent streams (individual agent output)
|
|
- Session streams (messages in a session)
|
|
|
|
Security Note:
|
|
WebSocket connections validate agent_id via query params and verify
|
|
the agent exists in the database. In production, this should be
|
|
enhanced with proper token-based authentication (JWT, etc.).
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
from uuid import UUID
|
|
|
|
import structlog
|
|
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, status
|
|
|
|
from roboco.agents_config import CEO_AGENT_ID, verify_agent_token
|
|
from roboco.api.deps import _auth_required
|
|
from roboco.db.base import get_db
|
|
from roboco.services.repositories import resolve_agent_uuid
|
|
|
|
router = APIRouter()
|
|
log = structlog.get_logger()
|
|
|
|
# Server-side idle timeout for WS receive loops. A half-open socket (dead
|
|
# agent container, silent client) blocks ``receive_text()`` forever;
|
|
# ``asyncio.wait_for`` reaps the socket after this many seconds of silence.
|
|
IDLE_TIMEOUT_SECONDS: float = 90.0
|
|
|
|
# Per-connection send queue + send timeout. Each registered connection owns a
|
|
# bounded ``asyncio.Queue`` drained by a sender task, so a slow client can't
|
|
# back-pressure the fan-out: broadcast enqueues (non-blocking) and returns
|
|
# immediately; a full queue drops + logs (client lagging, not the fan-out).
|
|
MAX_SEND_QUEUE: int = 256
|
|
SEND_TIMEOUT_SECONDS: float = 10.0
|
|
|
|
|
|
class _ClientConnection:
|
|
"""Per-connection send queue + sender task.
|
|
|
|
Holds the bounded outbound queue drained by ``sender``; broadcast enqueues
|
|
here instead of awaiting ``send_text`` directly, so one slow client cannot
|
|
block the fan-out to every other client.
|
|
"""
|
|
|
|
__slots__ = ("queue", "sender", "websocket")
|
|
|
|
def __init__(self, websocket: WebSocket, maxsize: int) -> None:
|
|
self.websocket = websocket
|
|
self.queue: asyncio.Queue[str] = asyncio.Queue(maxsize=maxsize)
|
|
self.sender: asyncio.Task[None] | None = None
|
|
|
|
|
|
async def _require_panel_token(websocket: WebSocket) -> bool:
|
|
"""Bind a per-agent WS upgrade to the panel/CEO HMAC token.
|
|
|
|
/ws/* streams are operator-only (the panel is the sole WS client; agents
|
|
use MCP verbs). nginx injects the CEO panel token as ``X-Agent-Token``.
|
|
In strict mode (``ROBOCO_AGENT_AUTH_REQUIRED=true``) the token is required
|
|
+ verified against the CEO identity; a presented-but-forged token is
|
|
rejected even in dev mode. Returns True to proceed, False to close.
|
|
"""
|
|
token = websocket.headers.get("x-agent-token")
|
|
if _auth_required() and not token:
|
|
return False
|
|
# A missing token in dev mode proceeds; a presented token must verify.
|
|
return not (token and not verify_agent_token(token, CEO_AGENT_ID, "ceo", ""))
|
|
|
|
|
|
# =============================================================================
|
|
# Connection Manager
|
|
# =============================================================================
|
|
|
|
|
|
class ConnectionManager:
|
|
"""
|
|
Manages WebSocket connections organized by type and ID.
|
|
|
|
Supports:
|
|
- Channel subscriptions
|
|
- Agent output streams
|
|
- Session streams
|
|
"""
|
|
|
|
def __init__(self) -> None:
|
|
# channel_id -> set of websockets
|
|
self.channel_connections: dict[UUID, set[WebSocket]] = {}
|
|
|
|
# agent_id -> set of websockets
|
|
self.agent_connections: dict[UUID, set[WebSocket]] = {}
|
|
|
|
# session_id -> set of websockets
|
|
self.session_connections: dict[UUID, set[WebSocket]] = {}
|
|
|
|
# agent_id -> set of websockets (for notifications)
|
|
self.notification_connections: dict[UUID, set[WebSocket]] = {}
|
|
|
|
# Operator/system-wide stream (rate limits, etc.) — no per-agent keying.
|
|
self.system_connections: set[WebSocket] = set()
|
|
|
|
# websocket -> agent_id (for tracking who is connected)
|
|
self.connection_agents: dict[WebSocket, UUID] = {}
|
|
|
|
# websocket -> per-connection send queue + sender task. Every connect_*
|
|
# registers here; disconnect cancels + removes. Broadcast enqueues into
|
|
# these queues instead of awaiting send_text directly so one slow client
|
|
# can't block the fan-out.
|
|
self.connection_senders: dict[WebSocket, _ClientConnection] = {}
|
|
|
|
# Fire-and-forget fallback send tasks for unregistered sockets (legacy
|
|
# path). Held to satisfy ruff RUF006 + allow clean shutdown; each task
|
|
# removes itself on completion.
|
|
self._pending_sends: set[asyncio.Task[None]] = set()
|
|
|
|
def _register_sender(self, websocket: WebSocket) -> _ClientConnection:
|
|
"""Create the per-connection send queue + start its sender task."""
|
|
conn = _ClientConnection(websocket, maxsize=MAX_SEND_QUEUE)
|
|
conn.sender = asyncio.create_task(self._run_sender(conn))
|
|
self.connection_senders[websocket] = conn
|
|
return conn
|
|
|
|
async def _run_sender(self, conn: _ClientConnection) -> None:
|
|
"""Drain the per-connection send queue; each send is timeout-bounded."""
|
|
ws = conn.websocket
|
|
while True:
|
|
data = await conn.queue.get()
|
|
try:
|
|
await asyncio.wait_for(ws.send_text(data), timeout=SEND_TIMEOUT_SECONDS)
|
|
except TimeoutError:
|
|
log.warning(
|
|
"WebSocket send timeout — dropping message to slow client",
|
|
timeout=SEND_TIMEOUT_SECONDS,
|
|
)
|
|
except Exception as exc:
|
|
# Transport closed / hard send error — the socket is provably
|
|
# dead on the send side. Reap it from every subscription set
|
|
# now instead of waiting for the receive loop's idle timeout to
|
|
# notice: otherwise the dead socket lingers in the sets and
|
|
# every subsequent broadcast enqueues into this connection's
|
|
# queue whose consumer has just exited (queue-overflow log spam
|
|
# then silent drops) for up to IDLE_TIMEOUT_SECONDS. A send
|
|
# TIMEOUT alone (slow client) does NOT reach here — it is
|
|
# caught above and the live socket is kept.
|
|
log.debug(
|
|
"WebSocket sender disconnecting on send error",
|
|
error=str(exc),
|
|
)
|
|
self.disconnect(ws)
|
|
return
|
|
|
|
async def connect_channel(
|
|
self, websocket: WebSocket, channel_id: UUID, agent_id: UUID
|
|
) -> None:
|
|
"""Connect to a channel stream."""
|
|
await websocket.accept()
|
|
|
|
if channel_id not in self.channel_connections:
|
|
self.channel_connections[channel_id] = set()
|
|
|
|
self.channel_connections[channel_id].add(websocket)
|
|
self.connection_agents[websocket] = agent_id
|
|
self._register_sender(websocket)
|
|
|
|
async def connect_agent(
|
|
self, websocket: WebSocket, target_agent_id: UUID, viewer_agent_id: UUID
|
|
) -> None:
|
|
"""Connect to an agent's output stream."""
|
|
await websocket.accept()
|
|
|
|
if target_agent_id not in self.agent_connections:
|
|
self.agent_connections[target_agent_id] = set()
|
|
|
|
self.agent_connections[target_agent_id].add(websocket)
|
|
self.connection_agents[websocket] = viewer_agent_id
|
|
self._register_sender(websocket)
|
|
|
|
async def connect_session(
|
|
self, websocket: WebSocket, session_id: UUID, agent_id: UUID
|
|
) -> None:
|
|
"""Connect to a session stream."""
|
|
await websocket.accept()
|
|
|
|
if session_id not in self.session_connections:
|
|
self.session_connections[session_id] = set()
|
|
|
|
self.session_connections[session_id].add(websocket)
|
|
self.connection_agents[websocket] = agent_id
|
|
self._register_sender(websocket)
|
|
|
|
async def connect_notifications(self, websocket: WebSocket, agent_id: UUID) -> None:
|
|
"""Connect to an agent's notification stream."""
|
|
await websocket.accept()
|
|
|
|
if agent_id not in self.notification_connections:
|
|
self.notification_connections[agent_id] = set()
|
|
|
|
self.notification_connections[agent_id].add(websocket)
|
|
self.connection_agents[websocket] = agent_id
|
|
self._register_sender(websocket)
|
|
|
|
async def connect_system(self, websocket: WebSocket) -> None:
|
|
"""Connect to the operator/system-wide stream (rate limits, etc.)."""
|
|
await websocket.accept()
|
|
self.system_connections.add(websocket)
|
|
self._register_sender(websocket)
|
|
|
|
def disconnect(self, websocket: WebSocket) -> None:
|
|
"""Remove a websocket from all subscriptions."""
|
|
# Remove from channel connections
|
|
for connections in self.channel_connections.values():
|
|
connections.discard(websocket)
|
|
|
|
# Remove from agent connections
|
|
for connections in self.agent_connections.values():
|
|
connections.discard(websocket)
|
|
|
|
# Remove from session connections
|
|
for connections in self.session_connections.values():
|
|
connections.discard(websocket)
|
|
|
|
# Remove from notification connections
|
|
for connections in self.notification_connections.values():
|
|
connections.discard(websocket)
|
|
|
|
# Remove from the system-wide stream
|
|
self.system_connections.discard(websocket)
|
|
|
|
# Remove from tracking
|
|
self.connection_agents.pop(websocket, None)
|
|
|
|
# Cancel + drop the per-connection sender task so a slow/stale client's
|
|
# queue doesn't leak after the socket is removed.
|
|
conn = self.connection_senders.pop(websocket, None)
|
|
if conn is not None and conn.sender is not None:
|
|
conn.sender.cancel()
|
|
|
|
def _enqueue_or_send(self, websocket: WebSocket, data: str) -> None:
|
|
"""Fan out one message to one connection without blocking.
|
|
|
|
Registered connections get the message enqueued into their bounded send
|
|
queue (non-blocking, drop + warn on overflow). An unregistered socket
|
|
falls back to a timeout-bounded ``send_text`` scheduled on the loop, so
|
|
the broadcast never blocks on a single slow client.
|
|
"""
|
|
conn = self.connection_senders.get(websocket)
|
|
if conn is not None:
|
|
try:
|
|
conn.queue.put_nowait(data)
|
|
except asyncio.QueueFull:
|
|
log.warning(
|
|
"WebSocket send queue overflow — dropping message",
|
|
queue_size=conn.queue.maxsize,
|
|
)
|
|
return
|
|
# Legacy fallback: schedule a timeout-bounded send so a slow
|
|
# unregistered client can't wedge the fan-out either. Keep a strong
|
|
# reference so the task isn't GC'd mid-flight (ruff RUF006); it
|
|
# discards itself on completion.
|
|
task = asyncio.create_task(self._send_with_timeout(websocket, data))
|
|
self._pending_sends.add(task)
|
|
task.add_done_callback(self._pending_sends.discard)
|
|
|
|
async def _send_with_timeout(self, websocket: WebSocket, data: str) -> None:
|
|
try:
|
|
await asyncio.wait_for(
|
|
websocket.send_text(data), timeout=SEND_TIMEOUT_SECONDS
|
|
)
|
|
except TimeoutError:
|
|
log.warning(
|
|
"WebSocket send timeout — dropping message to slow client",
|
|
timeout=SEND_TIMEOUT_SECONDS,
|
|
)
|
|
except Exception as exc: # transport closed / cancelled
|
|
log.debug("WebSocket send failed", error=str(exc))
|
|
|
|
async def broadcast_to_channel(
|
|
self, channel_id: UUID, message: dict[str, Any]
|
|
) -> None:
|
|
"""Broadcast a message to all channel subscribers."""
|
|
connections = self.channel_connections.get(channel_id, set())
|
|
if not connections:
|
|
return
|
|
data = json.dumps(message, default=str)
|
|
for conn in connections:
|
|
self._enqueue_or_send(conn, data)
|
|
|
|
async def broadcast_to_agent_watchers(
|
|
self, agent_id: UUID, message: dict[str, Any]
|
|
) -> None:
|
|
"""Broadcast a message to all watching an agent's stream."""
|
|
connections = self.agent_connections.get(agent_id, set())
|
|
if not connections:
|
|
return
|
|
data = json.dumps(message, default=str)
|
|
for conn in connections:
|
|
self._enqueue_or_send(conn, data)
|
|
|
|
async def broadcast_to_session(
|
|
self, session_id: UUID, message: dict[str, Any]
|
|
) -> None:
|
|
"""Broadcast a message to all session subscribers."""
|
|
connections = self.session_connections.get(session_id, set())
|
|
if not connections:
|
|
return
|
|
data = json.dumps(message, default=str)
|
|
for conn in connections:
|
|
self._enqueue_or_send(conn, data)
|
|
|
|
async def broadcast_system(self, message: dict[str, Any]) -> None:
|
|
"""Broadcast a message to all operator/system-wide subscribers."""
|
|
if not self.system_connections:
|
|
return
|
|
data = json.dumps(message, default=str)
|
|
for conn in self.system_connections:
|
|
self._enqueue_or_send(conn, data)
|
|
|
|
def get_channel_subscriber_count(self, channel_id: UUID) -> int:
|
|
"""Get number of subscribers to a channel."""
|
|
return len(self.channel_connections.get(channel_id, set()))
|
|
|
|
def get_agent_watcher_count(self, agent_id: UUID) -> int:
|
|
"""Get number of watchers of an agent's stream."""
|
|
return len(self.agent_connections.get(agent_id, set()))
|
|
|
|
|
|
# Global connection manager
|
|
manager = ConnectionManager()
|
|
|
|
|
|
async def validate_agent_exists(agent_id: UUID | str) -> bool:
|
|
"""
|
|
Validate that an agent exists in the database.
|
|
|
|
This provides basic security by ensuring the claimed agent_id
|
|
is a valid agent, not just a valid UUID format.
|
|
|
|
TODO: Enhance with token-based authentication (JWT) for production.
|
|
"""
|
|
try:
|
|
async for db in get_db():
|
|
result = await resolve_agent_uuid(db, str(agent_id))
|
|
return result is not None
|
|
except Exception:
|
|
return False
|
|
return False
|
|
|
|
|
|
# =============================================================================
|
|
# WebSocket Routes
|
|
# =============================================================================
|
|
|
|
|
|
@router.websocket("/channels/{channel_id}")
|
|
async def channel_stream(
|
|
websocket: WebSocket,
|
|
channel_id: UUID,
|
|
) -> None:
|
|
"""
|
|
WebSocket endpoint for channel message streams.
|
|
|
|
Clients receive real-time messages for the channel.
|
|
"""
|
|
# Verify the panel/CEO token before any subject lookup.
|
|
if not await _require_panel_token(websocket):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
# Get agent ID from query params (or auth in production)
|
|
agent_id_str = websocket.query_params.get("agent_id")
|
|
if not agent_id_str:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
try:
|
|
agent_id = UUID(agent_id_str)
|
|
except ValueError:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
await manager.connect_channel(websocket, channel_id, agent_id)
|
|
|
|
try:
|
|
# Send connection confirmation
|
|
await websocket.send_json(
|
|
{
|
|
"type": "connected",
|
|
"channel_id": str(channel_id),
|
|
"subscriber_count": manager.get_channel_subscriber_count(channel_id),
|
|
}
|
|
)
|
|
|
|
# Keep connection alive and handle incoming messages
|
|
while True:
|
|
data = await asyncio.wait_for(
|
|
websocket.receive_text(), timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
|
|
# Handle ping/pong for keepalive
|
|
if data == "ping":
|
|
await websocket.send_text("pong")
|
|
continue
|
|
|
|
# Handle other client messages if needed
|
|
# For now, channels are primarily for receiving
|
|
|
|
except WebSocketDisconnect:
|
|
# Clean client-initiated disconnect — handled here for clarity; the
|
|
# finally below also disconnects (idempotent) to cover every other
|
|
# exit path (anyio closed-resource, CancelledError, transport errors).
|
|
pass
|
|
except TimeoutError:
|
|
# Idle timeout — the client has been silent for IDLE_TIMEOUT_SECONDS
|
|
# (likely a half-open socket from a dead container). Fall through to
|
|
# the finally so the socket is removed from every subscription set.
|
|
log.warning(
|
|
"WebSocket idle timeout — disconnecting", timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
finally:
|
|
manager.disconnect(websocket)
|
|
|
|
|
|
@router.websocket("/agents/{agent_id}")
|
|
async def agent_stream(
|
|
websocket: WebSocket,
|
|
agent_id: UUID,
|
|
) -> None:
|
|
"""
|
|
WebSocket endpoint for an agent's output stream.
|
|
|
|
Clients receive real-time LLM output from the agent.
|
|
"""
|
|
# Verify the panel/CEO token before any subject lookup.
|
|
if not await _require_panel_token(websocket):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
# Get viewer agent ID
|
|
viewer_id_str = websocket.query_params.get("viewer_id")
|
|
if not viewer_id_str:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
try:
|
|
viewer_id = UUID(viewer_id_str)
|
|
except ValueError:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
# Validate viewer agent exists in database
|
|
if not await validate_agent_exists(viewer_id):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
await manager.connect_agent(websocket, agent_id, viewer_id)
|
|
|
|
try:
|
|
await websocket.send_json(
|
|
{
|
|
"type": "connected",
|
|
"agent_id": str(agent_id),
|
|
"watcher_count": manager.get_agent_watcher_count(agent_id),
|
|
}
|
|
)
|
|
|
|
while True:
|
|
data = await asyncio.wait_for(
|
|
websocket.receive_text(), timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
if data == "ping":
|
|
await websocket.send_text("pong")
|
|
|
|
except WebSocketDisconnect:
|
|
# Clean client-initiated disconnect — handled here for clarity; the
|
|
# finally below also disconnects (idempotent) to cover every other
|
|
# exit path (anyio closed-resource, CancelledError, transport errors).
|
|
pass
|
|
except TimeoutError:
|
|
# Idle timeout — the client has been silent for IDLE_TIMEOUT_SECONDS
|
|
# (likely a half-open socket from a dead container). Fall through to
|
|
# the finally so the socket is removed from every subscription set.
|
|
log.warning(
|
|
"WebSocket idle timeout — disconnecting", timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
finally:
|
|
manager.disconnect(websocket)
|
|
|
|
|
|
@router.websocket("/sessions/{session_id}")
|
|
async def session_stream(
|
|
websocket: WebSocket,
|
|
session_id: UUID,
|
|
) -> None:
|
|
"""
|
|
WebSocket endpoint for session message streams.
|
|
|
|
Clients receive real-time messages for a specific session.
|
|
"""
|
|
# Verify the panel/CEO token before any subject lookup.
|
|
if not await _require_panel_token(websocket):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
agent_id_str = websocket.query_params.get("agent_id")
|
|
if not agent_id_str:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
try:
|
|
agent_id = UUID(agent_id_str)
|
|
except ValueError:
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
# Validate agent exists in database
|
|
if not await validate_agent_exists(agent_id):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
await manager.connect_session(websocket, session_id, agent_id)
|
|
|
|
try:
|
|
await websocket.send_json(
|
|
{
|
|
"type": "connected",
|
|
"session_id": str(session_id),
|
|
}
|
|
)
|
|
|
|
while True:
|
|
data = await asyncio.wait_for(
|
|
websocket.receive_text(), timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
if data == "ping":
|
|
await websocket.send_text("pong")
|
|
|
|
except WebSocketDisconnect:
|
|
# Clean client-initiated disconnect — handled here for clarity; the
|
|
# finally below also disconnects (idempotent) to cover every other
|
|
# exit path (anyio closed-resource, CancelledError, transport errors).
|
|
pass
|
|
except TimeoutError:
|
|
# Idle timeout — the client has been silent for IDLE_TIMEOUT_SECONDS
|
|
# (likely a half-open socket from a dead container). Fall through to
|
|
# the finally so the socket is removed from every subscription set.
|
|
log.warning(
|
|
"WebSocket idle timeout — disconnecting", timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
finally:
|
|
manager.disconnect(websocket)
|
|
|
|
|
|
@router.websocket("/notifications/{agent_id}")
|
|
async def notification_stream(
|
|
websocket: WebSocket,
|
|
agent_id: UUID,
|
|
) -> None:
|
|
"""
|
|
WebSocket endpoint for agent notifications.
|
|
|
|
Agents receive real-time notifications via this stream.
|
|
"""
|
|
# Verify the panel/CEO token before any subject lookup.
|
|
if not await _require_panel_token(websocket):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
# Validate agent exists in database
|
|
if not await validate_agent_exists(agent_id):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
|
|
await manager.connect_notifications(websocket, agent_id)
|
|
|
|
try:
|
|
await websocket.send_json(
|
|
{
|
|
"type": "connected",
|
|
"agent_id": str(agent_id),
|
|
}
|
|
)
|
|
|
|
while True:
|
|
data = await asyncio.wait_for(
|
|
websocket.receive_text(), timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
if data == "ping":
|
|
await websocket.send_text("pong")
|
|
|
|
except WebSocketDisconnect:
|
|
# Clean client-initiated disconnect — handled here for clarity; the
|
|
# finally below also disconnects (idempotent) to cover every other
|
|
# exit path (anyio closed-resource, CancelledError, transport errors).
|
|
pass
|
|
except TimeoutError:
|
|
# Idle timeout — the client has been silent for IDLE_TIMEOUT_SECONDS
|
|
# (likely a half-open socket from a dead container). Fall through to
|
|
# the finally so the socket is removed from every subscription set.
|
|
log.warning(
|
|
"WebSocket idle timeout — disconnecting", timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
finally:
|
|
manager.disconnect(websocket)
|
|
|
|
|
|
@router.websocket("/system")
|
|
async def system_stream(websocket: WebSocket) -> None:
|
|
"""Operator/system-wide WebSocket stream.
|
|
|
|
Carries system-level events for the control panel — currently the
|
|
rate-limit lifecycle (``RATE_LIMIT_HIT`` / ``RATE_LIMIT_LIFTED``), bridged
|
|
from the event bus by ``websocket_bridge``. No per-agent keying; the
|
|
panel/CEO token gate matches every sibling /ws/* stream (#24 — this was the
|
|
only ungated /ws endpoint): in strict mode a missing CEO token closes with
|
|
policy-violation, a presented-but-forged token is rejected even in dev.
|
|
"""
|
|
# Verify the panel/CEO token before subscribing — same gate as every other
|
|
# /ws/* handler (#24).
|
|
if not await _require_panel_token(websocket):
|
|
await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
|
|
return
|
|
await manager.connect_system(websocket)
|
|
|
|
try:
|
|
await websocket.send_json({"type": "connected"})
|
|
|
|
while True:
|
|
data = await asyncio.wait_for(
|
|
websocket.receive_text(), timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
if data == "ping":
|
|
await websocket.send_text("pong")
|
|
|
|
except WebSocketDisconnect:
|
|
# Clean client-initiated disconnect — handled here for clarity; the
|
|
# finally below also disconnects (idempotent) to cover every other
|
|
# exit path (anyio closed-resource, CancelledError, transport errors).
|
|
pass
|
|
except TimeoutError:
|
|
# Idle timeout — the client has been silent for IDLE_TIMEOUT_SECONDS
|
|
# (likely a half-open socket from a dead container). Fall through to
|
|
# the finally so the socket is removed from every subscription set.
|
|
log.warning(
|
|
"WebSocket idle timeout — disconnecting", timeout=IDLE_TIMEOUT_SECONDS
|
|
)
|
|
finally:
|
|
manager.disconnect(websocket)
|
|
|
|
|
|
# =============================================================================
|
|
# Helper Functions for Broadcasting
|
|
# =============================================================================
|
|
|
|
|
|
async def broadcast_agent_chunk(
|
|
agent_id: str, chunk: str, metadata: dict[str, Any]
|
|
) -> None:
|
|
"""Broadcast an agent stream chunk to watchers."""
|
|
event = {
|
|
"type": "agent.stream",
|
|
"agent_id": agent_id,
|
|
"chunk": chunk,
|
|
"timestamp": datetime.now(UTC).isoformat(),
|
|
**metadata,
|
|
}
|
|
|
|
await manager.broadcast_to_agent_watchers(UUID(agent_id), event)
|
|
|
|
|
|
async def broadcast_notification(
|
|
agent_ids: list[UUID],
|
|
notification_id: UUID,
|
|
notification_type: str,
|
|
subject: str,
|
|
priority: str,
|
|
) -> None:
|
|
"""
|
|
Broadcast notification to specific agents.
|
|
|
|
Sends to all agents that have notification websocket connections.
|
|
"""
|
|
event = {
|
|
"type": "notification",
|
|
"notification_id": str(notification_id),
|
|
"notification_type": notification_type,
|
|
"subject": subject,
|
|
"priority": priority,
|
|
"timestamp": datetime.now(UTC).isoformat(),
|
|
}
|
|
data = json.dumps(event)
|
|
|
|
for agent_id in agent_ids:
|
|
connections = manager.notification_connections.get(agent_id, set())
|
|
if connections:
|
|
for conn in connections:
|
|
manager._enqueue_or_send(conn, data)
|