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>
1350 lines
50 KiB
Python
1350 lines
50 KiB
Python
"""Unit tests for TaskService gateway-backfill methods.
|
|
|
|
These cover the methods the Choreographer calls into; full end-to-end
|
|
behavior is exercised by the gateway tests. Each test mocks the DB
|
|
session boundary and checks the method's contract.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime
|
|
from types import SimpleNamespace
|
|
from typing import Any
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
from uuid import uuid4
|
|
|
|
import pytest
|
|
from roboco.db.tables import AuditLogTable
|
|
from roboco.models.base import (
|
|
AgentRole,
|
|
AgentStatus,
|
|
BlockerResolverType,
|
|
Complexity,
|
|
TaskNature,
|
|
TaskStatus,
|
|
TaskType,
|
|
Team,
|
|
)
|
|
from roboco.models.task import TaskCreateRequest
|
|
from roboco.services.task import GatewayAgentView, TaskService
|
|
|
|
|
|
def _build_task(**overrides: object) -> MagicMock:
|
|
base: dict[str, object] = {
|
|
"id": uuid4(),
|
|
"status": TaskStatus.PENDING,
|
|
"branch_name": "feature/backend/abc12345",
|
|
"assigned_to": None,
|
|
"claimed_by": None,
|
|
"claimed_at": None,
|
|
"plan": None,
|
|
"qa_evidence_inspected": False,
|
|
"pre_block_state": None,
|
|
"pre_block_assignee": None,
|
|
"pre_block_metadata": None,
|
|
"blocker_resolver_type": None,
|
|
"blocker_raised_by": None,
|
|
"commits": [],
|
|
"dev_notes": None,
|
|
}
|
|
base.update(overrides)
|
|
return MagicMock(**base)
|
|
|
|
|
|
def _service_with(execute_returns: object) -> TaskService:
|
|
"""Build a TaskService whose session.execute returns `execute_returns`."""
|
|
session = MagicMock()
|
|
session.execute = AsyncMock(return_value=execute_returns)
|
|
session.flush = AsyncMock()
|
|
return TaskService(session)
|
|
|
|
|
|
def _bind(svc: TaskService, name: str, value: object) -> None:
|
|
"""Stub `name` on `svc` without tripping mypy's method-assign check."""
|
|
object.__setattr__(svc, name, value)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Aliases / thin wrappers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_blocked_for_team_filters_by_team() -> None:
|
|
svc = TaskService(MagicMock())
|
|
list_blocked_mock = AsyncMock(return_value=[MagicMock(id="t1")])
|
|
_bind(svc, "list_blocked", list_blocked_mock)
|
|
out = await svc.list_blocked_for_team(Team.BACKEND)
|
|
list_blocked_mock.assert_awaited_once_with(team=Team.BACKEND)
|
|
assert len(out) == 1
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_blocked_all_teams_passes_no_team() -> None:
|
|
svc = TaskService(MagicMock())
|
|
list_blocked_mock = AsyncMock(return_value=[])
|
|
_bind(svc, "list_blocked", list_blocked_mock)
|
|
await svc.list_blocked_all_teams()
|
|
list_blocked_mock.assert_awaited_once_with()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_pm_review_for_team_passes_team() -> None:
|
|
svc = TaskService(MagicMock())
|
|
list_pm_mock = AsyncMock(return_value=[])
|
|
_bind(svc, "list_awaiting_pm_review", list_pm_mock)
|
|
await svc.list_awaiting_pm_review_for_team(Team.FRONTEND)
|
|
list_pm_mock.assert_awaited_once_with(team=Team.FRONTEND)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_verification_records_progress_when_notes_given() -> None:
|
|
svc = TaskService(MagicMock())
|
|
add_progress_mock = AsyncMock()
|
|
submit_for_verification_mock = AsyncMock(return_value=MagicMock())
|
|
_bind(svc, "add_progress", add_progress_mock)
|
|
_bind(svc, "submit_for_verification", submit_for_verification_mock)
|
|
agent_id = uuid4()
|
|
task_id = uuid4()
|
|
await svc.submit_verification(agent_id, task_id, "implemented login")
|
|
add_progress_mock.assert_awaited_once_with(task_id, agent_id, "implemented login")
|
|
submit_for_verification_mock.assert_awaited_once_with(
|
|
task_id, agent_role="developer"
|
|
)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_verification_skips_progress_when_notes_empty() -> None:
|
|
svc = TaskService(MagicMock())
|
|
add_progress_mock = AsyncMock()
|
|
_bind(svc, "add_progress", add_progress_mock)
|
|
_bind(svc, "submit_for_verification", AsyncMock(return_value=MagicMock()))
|
|
await svc.submit_verification(uuid4(), uuid4(), "")
|
|
add_progress_mock.assert_not_called()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_submit_qa_records_progress_when_notes_given() -> None:
|
|
svc = TaskService(MagicMock())
|
|
add_progress_mock = AsyncMock()
|
|
submit_for_qa_mock = AsyncMock(return_value=MagicMock())
|
|
_bind(svc, "add_progress", add_progress_mock)
|
|
_bind(svc, "submit_for_qa", submit_for_qa_mock)
|
|
agent_id = uuid4()
|
|
task_id = uuid4()
|
|
await svc.submit_qa(agent_id, task_id, "ready for review")
|
|
add_progress_mock.assert_awaited_once()
|
|
submit_for_qa_mock.assert_awaited_once_with(task_id, agent_role="developer")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# list_assigned_for_agent
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_assigned_for_agent_returns_active_tasks() -> None:
|
|
expected_tasks = [MagicMock(id="t1"), MagicMock(id="t2")]
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = expected_tasks
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
out = await svc.list_assigned_for_agent(uuid4())
|
|
assert out == expected_tasks
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# agent_for
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_agent_for_returns_view_with_role_team_skills(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
fake_agent = MagicMock(
|
|
id=uuid4(), slug="be-pm", role=AgentRole.CELL_PM, team=Team.BACKEND
|
|
)
|
|
result = MagicMock()
|
|
result.scalar_one_or_none.return_value = fake_agent
|
|
svc = _service_with(result)
|
|
|
|
monkeypatch.setattr(
|
|
"roboco.agents_config.get_escalation_target", lambda _slug: "main_pm"
|
|
)
|
|
monkeypatch.setattr(
|
|
"roboco.agents_config.get_agent_skills",
|
|
lambda _slug: [{"id": "task_management"}],
|
|
)
|
|
view = await svc.agent_for(uuid4())
|
|
assert isinstance(view, GatewayAgentView)
|
|
assert view.role == "cell_pm"
|
|
assert view.team == "backend"
|
|
assert view.escalation_target == "main_pm"
|
|
assert view.skills == [{"id": "task_management"}]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_agent_for_returns_none_when_missing() -> None:
|
|
result = MagicMock()
|
|
result.scalar_one_or_none.return_value = None
|
|
svc = _service_with(result)
|
|
assert await svc.agent_for(uuid4()) is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# qa/documenter/cell_pm for_team
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_agent_for_team_finds_qa() -> None:
|
|
qa = MagicMock(id=uuid4(), role=AgentRole.QA, team=Team.BACKEND)
|
|
scalars = MagicMock()
|
|
scalars.first.return_value = qa
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
out = await svc.qa_agent_for_team(Team.BACKEND)
|
|
assert out is qa
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_documenter_for_team_returns_none_when_missing() -> None:
|
|
scalars = MagicMock()
|
|
scalars.first.return_value = None
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
assert await svc.documenter_for_team(Team.BACKEND) is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cell_pm_for_team_finds_pm() -> None:
|
|
pm = MagicMock(id=uuid4(), role=AgentRole.CELL_PM, team=Team.UX_UI)
|
|
scalars = MagicMock()
|
|
scalars.first.return_value = pm
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
assert await svc.cell_pm_for_team(Team.UX_UI) is pm
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# get_active_task_for_agent + list_paused_for_agent
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_get_active_task_for_agent_returns_top_task() -> None:
|
|
task = MagicMock(id=uuid4(), status=TaskStatus.IN_PROGRESS)
|
|
result = MagicMock()
|
|
result.scalar_one_or_none.return_value = task
|
|
svc = _service_with(result)
|
|
assert await svc.get_active_task_for_agent(uuid4()) is task
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_paused_for_agent_returns_paused_tasks() -> None:
|
|
paused = [MagicMock(id=uuid4(), status=TaskStatus.PAUSED)]
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = paused
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
out = await svc.list_paused_for_agent(uuid4())
|
|
assert out == paused
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# list_awaiting_main_pm_all
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_awaiting_main_pm_all_returns_root_tasks() -> None:
|
|
roots = [MagicMock(id=uuid4(), parent_task_id=None)]
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = roots
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
out = await svc.list_awaiting_main_pm_all()
|
|
assert out == roots
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# all_subtasks_terminal
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_true_when_all_completed() -> None:
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = [TaskStatus.COMPLETED, TaskStatus.CANCELLED]
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
assert await svc.all_subtasks_terminal(uuid4()) is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_false_when_one_active() -> None:
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = [TaskStatus.COMPLETED, TaskStatus.IN_PROGRESS]
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
assert await svc.all_subtasks_terminal(uuid4()) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_all_subtasks_terminal_true_when_no_subtasks() -> None:
|
|
scalars = MagicMock()
|
|
scalars.all.return_value = []
|
|
result = MagicMock()
|
|
result.scalars.return_value = scalars
|
|
svc = _service_with(result)
|
|
assert await svc.all_subtasks_terminal(uuid4()) is True
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# set_plan
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan_wraps_string_into_text_dict() -> None:
|
|
task = _build_task()
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.set_plan(task.id, "do the thing")
|
|
assert task.plan == {"text": "do the thing"}
|
|
assert out is task
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan_passes_dict_through() -> None:
|
|
task = _build_task()
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.set_plan(task.id, {"steps": ["a", "b"]})
|
|
assert task.plan == {"steps": ["a", "b"]}
|
|
assert out is task
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_set_plan_returns_none_when_task_missing() -> None:
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=None))
|
|
assert await svc.set_plan(uuid4(), "plan") is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# mark_evidence_inspected + mark_agent_idle
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_evidence_inspected_sets_flag() -> None:
|
|
task = _build_task(qa_evidence_inspected=False)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
await svc.mark_evidence_inspected(task.id)
|
|
assert task.qa_evidence_inspected is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_evidence_inspected_no_op_on_missing_task() -> None:
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=None))
|
|
await svc.mark_evidence_inspected(uuid4()) # must not raise
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reassign_sets_assigned_to_and_claimed_by() -> None:
|
|
task = _build_task(assigned_to=None, claimed_by=None)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
new_assignee = uuid4()
|
|
out = await svc.reassign(task.id, new_assignee)
|
|
assert out is task
|
|
assert task.assigned_to == new_assignee
|
|
assert task.claimed_by == new_assignee
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reassign_clears_assignment_when_none() -> None:
|
|
prev = uuid4()
|
|
task = _build_task(assigned_to=prev, claimed_by=prev)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.reassign(task.id, None)
|
|
assert out is task
|
|
assert task.assigned_to is None
|
|
assert task.claimed_by is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_reassign_returns_none_when_task_missing() -> None:
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=None))
|
|
out = await svc.reassign(uuid4(), uuid4())
|
|
assert out is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_mark_agent_idle_sets_status_idle() -> None:
|
|
agent = MagicMock(id=uuid4(), status=AgentStatus.ACTIVE)
|
|
result = MagicMock()
|
|
result.scalar_one_or_none.return_value = agent
|
|
svc = _service_with(result)
|
|
await svc.mark_agent_idle(agent.id)
|
|
assert agent.status == AgentStatus.IDLE
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# qa_claim / doc_claim
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_claim_sets_assignment_on_awaiting_qa() -> None:
|
|
task = _build_task(status=TaskStatus.AWAITING_QA)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
qa_id = uuid4()
|
|
out = await svc.qa_claim(qa_id, task.id)
|
|
assert out is task
|
|
assert task.assigned_to == qa_id
|
|
assert task.claimed_by == qa_id
|
|
assert isinstance(task.claimed_at, datetime)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_claim_rejects_wrong_status() -> None:
|
|
task = _build_task(status=TaskStatus.IN_PROGRESS)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.qa_claim(uuid4(), task.id)
|
|
assert out is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_doc_claim_sets_assignment_on_awaiting_documentation() -> None:
|
|
task = _build_task(status=TaskStatus.AWAITING_DOCUMENTATION)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
doc_id = uuid4()
|
|
out = await svc.doc_claim(doc_id, task.id)
|
|
assert out is task
|
|
assert task.assigned_to == doc_id
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# qa_pass / qa_fail
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_pass_delegates_to_pass_qa() -> None:
|
|
qa_id = uuid4()
|
|
task_id = uuid4()
|
|
task = _build_task(id=task_id, claimed_by=qa_id)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
pass_qa_mock = AsyncMock(return_value=MagicMock())
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "pass_qa", pass_qa_mock)
|
|
await svc.qa_pass(qa_id, task_id, "looks good")
|
|
pass_qa_mock.assert_awaited_once_with(task_id, notes="looks good", agent_role="qa")
|
|
# active_claimant_id cleared so the documenter can claim cleanly.
|
|
assert task.active_claimant_id is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_qa_fail_appends_issues_to_dev_notes() -> None:
|
|
qa_id = uuid4()
|
|
task = _build_task(dev_notes=None, claimed_by=qa_id)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
fail_qa_mock = AsyncMock(return_value=task)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "fail_qa", fail_qa_mock)
|
|
issues = ["missing test", "no docstring"]
|
|
await svc.qa_fail(qa_id, task.id, "blocking", issues)
|
|
assert task.dev_notes is not None
|
|
assert "missing test" in task.dev_notes
|
|
assert "no docstring" in task.dev_notes
|
|
fail_qa_mock.assert_awaited_once_with(task.id, notes="blocking", agent_role="qa")
|
|
assert task.active_claimant_id is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# unblock_with_restore
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_with_restore_returns_to_pre_block_state() -> None:
|
|
pre_assignee = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED,
|
|
pre_block_state="in_progress",
|
|
pre_block_assignee=pre_assignee,
|
|
pre_block_metadata={"foo": "bar"},
|
|
blocker_resolver_type=BlockerResolverType.AGENT,
|
|
blocker_raised_by=pre_assignee,
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.unblock_with_restore(uuid4(), task.id, restore=True)
|
|
assert out is task
|
|
assert task.status == TaskStatus.IN_PROGRESS
|
|
assert task.assigned_to == pre_assignee
|
|
assert task.pre_block_state is None
|
|
assert task.pre_block_assignee is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_with_restore_falls_through_when_no_snapshot() -> None:
|
|
task = _build_task(status=TaskStatus.BLOCKED, pre_block_state=None)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
unblock_mock = AsyncMock(return_value=task)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "unblock", unblock_mock)
|
|
out = await svc.unblock_with_restore(uuid4(), task.id, restore=True)
|
|
unblock_mock.assert_awaited_once()
|
|
assert out is task
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_with_restore_calls_legacy_unblock_when_restore_false() -> None:
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED, pre_block_state=TaskStatus.IN_PROGRESS.value
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
unblock_mock = AsyncMock(return_value=task)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "unblock", unblock_mock)
|
|
await svc.unblock_with_restore(uuid4(), task.id, restore=False)
|
|
unblock_mock.assert_awaited_once()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_no_branch_returns_to_pending() -> None:
|
|
# A task blocked before it was ever claimed (a dependency-gated claim that
|
|
# got escalated) has no branch — unblock must return it to pending, not a
|
|
# branchless in_progress the dispatcher refuses to spawn (spawn-loop).
|
|
raiser = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED, branch_name=None, blocker_raised_by=raiser
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "_index_lifecycle_event_background", AsyncMock())
|
|
out = await svc.unblock(task.id)
|
|
assert out is task
|
|
assert task.status == TaskStatus.PENDING
|
|
assert task.assigned_to == raiser
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_admin_set_status_out_of_blocked_restores_pre_block_owner() -> None:
|
|
# A code task a dev escalated to its cell PM (assigned_to=PM, BLOCKED,
|
|
# snapshot=dev). Taking it out of blocked via the admin override (operator
|
|
# PATCH, or the orchestrator's auto-recover/auto-resume) must hand ownership
|
|
# back to the dev — otherwise it re-enters pending/in_progress still owned by
|
|
# the PM and the dispatcher execute-spawns the PM on a dev code task (loop).
|
|
dev = uuid4()
|
|
pm = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED,
|
|
assigned_to=pm,
|
|
claimed_by=pm,
|
|
branch_name="feature/frontend/abc--def--ghi",
|
|
pre_block_state="in_progress",
|
|
pre_block_assignee=dev,
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.admin_set_status(task.id, TaskStatus.IN_PROGRESS)
|
|
assert out is task
|
|
assert task.status == TaskStatus.IN_PROGRESS
|
|
assert task.assigned_to == dev
|
|
assert task.claimed_by == dev
|
|
assert task.pre_block_assignee is None
|
|
assert task.pre_block_state is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_admin_set_status_non_blocked_is_bare_status_set() -> None:
|
|
# The restore branch fires ONLY on blocked -> pending/in_progress with a
|
|
# snapshot. Every other override stays a plain status set; the owner is
|
|
# untouched (no spurious restore/divert).
|
|
owner = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.AWAITING_PM_REVIEW,
|
|
assigned_to=owner,
|
|
claimed_by=owner,
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.admin_set_status(task.id, TaskStatus.COMPLETED)
|
|
assert out is task
|
|
assert task.status == TaskStatus.COMPLETED
|
|
assert task.assigned_to == owner
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_admin_set_status_force_emits_distinct_override_audit_row() -> None:
|
|
"""#13: a forced override (force=True) emits a distinct ``task.admin_override``
|
|
audit row marking the bypass past the lifecycle gate — distinguishable from
|
|
the in-band transition row. Without force, no override row is emitted."""
|
|
owner = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.AWAITING_PM_REVIEW,
|
|
assigned_to=owner,
|
|
claimed_by=owner,
|
|
)
|
|
added: list[object] = []
|
|
session = MagicMock()
|
|
session.flush = AsyncMock()
|
|
session.add.side_effect = added.append
|
|
svc = TaskService(session)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
|
|
await svc.admin_set_status(
|
|
task.id, TaskStatus.COMPLETED, actor_id=owner, actor_role="ceo", force=True
|
|
)
|
|
|
|
rows = [r for r in added if isinstance(r, AuditLogTable)]
|
|
override_rows = [r for r in rows if r.event_type == "task.admin_override"]
|
|
assert len(override_rows) == 1
|
|
row = override_rows[0]
|
|
assert row.target_id == task.id
|
|
assert row.severity == "warning"
|
|
assert row.details["forced"] is True
|
|
assert row.details["from_status"] == "awaiting_pm_review"
|
|
assert row.details["to_status"] == "completed"
|
|
assert row.agent_id == owner
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_admin_set_status_no_force_emits_no_override_audit_row() -> None:
|
|
"""#13: without force, only the transition row is emitted — no
|
|
``task.admin_override`` row (the bypass was not acknowledged)."""
|
|
owner = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED,
|
|
assigned_to=owner,
|
|
claimed_by=owner,
|
|
)
|
|
added: list[object] = []
|
|
session = MagicMock()
|
|
session.flush = AsyncMock()
|
|
session.add.side_effect = added.append
|
|
svc = TaskService(session)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
|
|
await svc.admin_set_status(task.id, TaskStatus.PENDING, actor_id=owner)
|
|
|
|
rows = [r for r in added if isinstance(r, AuditLogTable)]
|
|
assert not any(r.event_type == "task.admin_override" for r in rows)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pre_block_restore_skips_revision_count_bump() -> None:
|
|
"""#101 Gap B: restoring a blocked task to its snapshotted needs_revision
|
|
state is a RESTORE, not a rework bounce — ``revision_count`` must not
|
|
increment. The rework counter counts rejections INTO needs_revision, and
|
|
this task was already rejected before it was blocked; unblocking it back to
|
|
needs_revision is the same rework cycle resuming, not a new one."""
|
|
REWORK_BOUNCES_BEFORE_BLOCK = 2
|
|
dev = uuid4()
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED,
|
|
assigned_to=dev,
|
|
claimed_by=dev,
|
|
branch_name="feature/backend/abc--def--ghi",
|
|
pre_block_state="needs_revision",
|
|
pre_block_assignee=dev,
|
|
revision_count=REWORK_BOUNCES_BEFORE_BLOCK,
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.unblock_with_restore(task.id, uuid4(), restore=True)
|
|
assert out is task
|
|
assert task.status == TaskStatus.NEEDS_REVISION
|
|
# A restore must NOT bump the rework counter — only a fresh rejection does.
|
|
assert task.revision_count == REWORK_BOUNCES_BEFORE_BLOCK
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_activate_batch_root_subtasks_emits_audit_for_activated_child() -> None:
|
|
"""#101 Gap A: ``_activate_batch_root_subtasks`` sets a held root-subtask
|
|
BACKLOG→PENDING directly. No status change may bypass the audit log — the
|
|
transition journey (the metric source of truth) must record the activation,
|
|
or the child's lifecycle reconstruction silently drops its start point."""
|
|
batch = uuid4()
|
|
child = _build_task(
|
|
status=TaskStatus.BACKLOG,
|
|
batch_id=batch,
|
|
team=Team.BOARD,
|
|
task_type=TaskType.CODE,
|
|
)
|
|
umbrella = _build_task(
|
|
status=TaskStatus.PENDING,
|
|
batch_id=batch,
|
|
parent_task_id=None,
|
|
team=Team.BOARD,
|
|
task_type=TaskType.PLANNING,
|
|
)
|
|
added: list[object] = []
|
|
session = MagicMock()
|
|
session.flush = AsyncMock()
|
|
session.add.side_effect = added.append
|
|
svc = TaskService(session)
|
|
_bind(svc, "get_subtasks", AsyncMock(return_value=[child]))
|
|
|
|
await svc._activate_batch_root_subtasks(umbrella)
|
|
|
|
assert child.status == TaskStatus.PENDING
|
|
rows = [r for r in added if isinstance(r, AuditLogTable)]
|
|
assert any(
|
|
r.event_type == "task.pending" and r.target_id == child.id for r in rows
|
|
), "batch root-subtask activation must emit a task.pending audit row"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_create_generates_ac_ids_and_carries_parent_ac_refs() -> None:
|
|
# Every task gets one stable id per acceptance criterion (1:1), and a
|
|
# decomposition child carries the parent AC ids it covers — the linkage the
|
|
# coverage + roll-up gates rely on.
|
|
svc = TaskService(
|
|
MagicMock(add=MagicMock(), flush=AsyncMock(), execute=AsyncMock())
|
|
)
|
|
req = TaskCreateRequest(
|
|
title="t",
|
|
description="d",
|
|
acceptance_criteria=["crit a", "crit b", "crit c"],
|
|
team=Team.BACKEND,
|
|
created_by=uuid4(),
|
|
task_type=TaskType.CODE,
|
|
nature=TaskNature.TECHNICAL,
|
|
estimated_complexity=Complexity.MEDIUM,
|
|
project_id=uuid4(),
|
|
parent_ac_refs=["parent-ac-1", "parent-ac-2"],
|
|
)
|
|
task = await svc.create(req)
|
|
n = len(req.acceptance_criteria)
|
|
assert len(task.acceptance_criteria_ids) == n
|
|
assert len(set(task.acceptance_criteria_ids)) == n
|
|
assert list(task.parent_ac_refs) == ["parent-ac-1", "parent-ac-2"]
|
|
|
|
|
|
def _svc_with_children(parent: object, child_rows: list[tuple]) -> TaskService:
|
|
"""TaskService whose get() returns `parent` and whose execute() yields the
|
|
(status, parent_ac_refs) child rows the coverage primitive selects."""
|
|
rows = MagicMock()
|
|
rows.all.return_value = child_rows
|
|
svc = TaskService(MagicMock(execute=AsyncMock(return_value=rows)))
|
|
_bind(svc, "get", AsyncMock(return_value=parent))
|
|
return svc
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_uncovered_parent_acs_inert_without_declared_coverage() -> None:
|
|
# No child declares parent_ac_refs -> coverage tracking inactive -> the gate
|
|
# is inert (legacy/in-flight tasks are never blocked).
|
|
parent = _build_task(
|
|
acceptance_criteria=["a", "b"], acceptance_criteria_ids=["id-a", "id-b"]
|
|
)
|
|
svc = _svc_with_children(
|
|
parent, [(TaskStatus.COMPLETED, []), (TaskStatus.COMPLETED, [])]
|
|
)
|
|
assert await svc.uncovered_parent_acceptance_criteria(parent.id) == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_uncovered_parent_acs_flags_unsatisfied_and_ignores_cancelled() -> None:
|
|
parent = _build_task(
|
|
acceptance_criteria=["crit a", "crit b", "crit c"],
|
|
acceptance_criteria_ids=["id-a", "id-b", "id-c"],
|
|
)
|
|
svc = _svc_with_children(
|
|
parent,
|
|
[
|
|
(TaskStatus.COMPLETED, ["id-a"]), # covers crit a
|
|
(TaskStatus.CANCELLED, ["id-b"]), # cancelled -> does NOT cover crit b
|
|
],
|
|
)
|
|
assert await svc.uncovered_parent_acceptance_criteria(parent.id) == [
|
|
"crit b",
|
|
"crit c",
|
|
]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_uncovered_parent_acs_empty_when_all_covered() -> None:
|
|
parent = _build_task(
|
|
acceptance_criteria=["a", "b"], acceptance_criteria_ids=["id-a", "id-b"]
|
|
)
|
|
svc = _svc_with_children(parent, [(TaskStatus.COMPLETED, ["id-a", "id-b"])])
|
|
assert await svc.uncovered_parent_acceptance_criteria(parent.id) == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_uncovered_parent_acs_recognizes_text_declared_coverage() -> None:
|
|
# Regression (phantom re-delegation): a PM may declare covers_parent_criteria
|
|
# by the criterion's full TEXT instead of its id. Matching is by id, so a
|
|
# COMPLETED child that declared coverage by text used to read "uncovered" —
|
|
# and the PM re-delegated the already-finished work as an empty phantom
|
|
# subtask (0 commits, no PR) that could never close. _parent_ac_ref_sets now
|
|
# normalizes text -> id so coverage counts regardless of how it was declared.
|
|
parent = _build_task(
|
|
acceptance_criteria=["crit a", "crit b"],
|
|
acceptance_criteria_ids=["id-a", "id-b"],
|
|
)
|
|
svc = _svc_with_children(
|
|
parent,
|
|
[
|
|
(TaskStatus.COMPLETED, ["crit a"]), # declared by TEXT, not "id-a"
|
|
(TaskStatus.COMPLETED, ["id-b"]), # declared by id
|
|
],
|
|
)
|
|
assert await svc.uncovered_parent_acceptance_criteria(parent.id) == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_parent_ac_coverage_normalizes_text_refs() -> None:
|
|
# A text-declared coverage ref from a COMPLETED child surfaces as
|
|
# claimed+verified, same as an id-declared one.
|
|
parent = _build_task(
|
|
acceptance_criteria=["crit a", "crit b"],
|
|
acceptance_criteria_ids=["id-a", "id-b"],
|
|
)
|
|
svc = _svc_with_children(parent, [(TaskStatus.COMPLETED, ["crit a"])])
|
|
cov = await svc.parent_ac_coverage(parent.id)
|
|
assert cov[0] == {
|
|
"id": "id-a",
|
|
"text": "crit a",
|
|
"claimed": True,
|
|
"verified": True,
|
|
}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_parent_ac_coverage_maps_claimed_and_verified() -> None:
|
|
# Per-criterion visibility: a COMPLETED child both claims and verifies its
|
|
# criterion; an in-flight child only claims; an untouched criterion is
|
|
# neither. This is the digest a decomposing PM reads from the briefing.
|
|
parent = _build_task(
|
|
acceptance_criteria=["crit a", "crit b", "crit c"],
|
|
acceptance_criteria_ids=["id-a", "id-b", "id-c"],
|
|
)
|
|
svc = _svc_with_children(
|
|
parent,
|
|
[
|
|
(TaskStatus.COMPLETED, ["id-a"]),
|
|
(TaskStatus.IN_PROGRESS, ["id-b"]),
|
|
],
|
|
)
|
|
assert await svc.parent_ac_coverage(parent.id) == [
|
|
{"id": "id-a", "text": "crit a", "claimed": True, "verified": True},
|
|
{"id": "id-b", "text": "crit b", "claimed": True, "verified": False},
|
|
{"id": "id-c", "text": "crit c", "claimed": False, "verified": False},
|
|
]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_parent_ac_coverage_empty_without_ac_ids() -> None:
|
|
# No stable ids on the parent (e.g. created before the linkage) -> nothing to
|
|
# report; the digest stays absent rather than emitting bogus rows.
|
|
parent = _build_task(acceptance_criteria=["a"], acceptance_criteria_ids=[])
|
|
svc = _svc_with_children(parent, [(TaskStatus.IN_PROGRESS, ["id-a"])])
|
|
assert await svc.parent_ac_coverage(parent.id) == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaimed_parent_acs_inert_without_declared_coverage() -> None:
|
|
# The decomposition floor is opt-in: with no child declaring parent_ac_refs
|
|
# it returns [] so a PM who never adopts coverage is never blocked at idle.
|
|
parent = _build_task(
|
|
acceptance_criteria=["a", "b"], acceptance_criteria_ids=["id-a", "id-b"]
|
|
)
|
|
svc = _svc_with_children(
|
|
parent, [(TaskStatus.IN_PROGRESS, []), (TaskStatus.IN_PROGRESS, [])]
|
|
)
|
|
assert await svc.unclaimed_parent_acceptance_criteria(parent.id) == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unclaimed_parent_acs_counts_live_children_not_just_completed() -> None:
|
|
# The distinction from the roll-up gate: an in-flight child *claims* its
|
|
# criterion (so the decomposition floor is satisfied) even though it has not
|
|
# yet *verified* it (so the roll-up gate still flags it). A cancelled child's
|
|
# claim does not count -- its work died with it.
|
|
parent = _build_task(
|
|
acceptance_criteria=["crit a", "crit b", "crit c"],
|
|
acceptance_criteria_ids=["id-a", "id-b", "id-c"],
|
|
)
|
|
rows = [
|
|
(TaskStatus.IN_PROGRESS, ["id-a"]), # live -> claims crit a
|
|
(TaskStatus.CANCELLED, ["id-b"]), # cancelled -> claim void
|
|
]
|
|
# unclaimed: crit a is claimed by the live child; crit b (only the cancelled
|
|
# child) and crit c (nobody) remain.
|
|
assert await _svc_with_children(parent, rows).unclaimed_parent_acceptance_criteria(
|
|
parent.id
|
|
) == ["crit b", "crit c"]
|
|
# roll-up still flags crit a too: the live child has not COMPLETED it.
|
|
assert await _svc_with_children(parent, rows).uncovered_parent_acceptance_criteria(
|
|
parent.id
|
|
) == [
|
|
"crit a",
|
|
"crit b",
|
|
"crit c",
|
|
]
|
|
|
|
|
|
def _svc_with_sibling_status_seq(rows: list[tuple]) -> TaskService:
|
|
"""TaskService whose execute() yields (status, sequence) sibling rows."""
|
|
res = MagicMock()
|
|
res.all.return_value = rows
|
|
return TaskService(MagicMock(execute=AsyncMock(return_value=res)))
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_earlier_incomplete_code_sibling_true_for_live_lower_seq() -> None:
|
|
# A dev's queued code leaf (seq 2) is lane-held while its own seq-0 sibling
|
|
# is still in flight — so it must not pin the dev to idle.
|
|
task = _build_task(
|
|
task_type=TaskType.CODE.value,
|
|
parent_task_id=uuid4(),
|
|
assigned_to=uuid4(),
|
|
sequence=2,
|
|
)
|
|
svc = _svc_with_sibling_status_seq([(TaskStatus.IN_PROGRESS, 0)])
|
|
assert await svc.has_earlier_incomplete_code_sibling(task) is True
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_earlier_incomplete_code_sibling_false_when_earlier_terminal() -> None:
|
|
task = _build_task(
|
|
task_type=TaskType.CODE.value,
|
|
parent_task_id=uuid4(),
|
|
assigned_to=uuid4(),
|
|
sequence=2,
|
|
)
|
|
svc = _svc_with_sibling_status_seq(
|
|
[(TaskStatus.COMPLETED, 0), (TaskStatus.CANCELLED, 1)]
|
|
)
|
|
assert await svc.has_earlier_incomplete_code_sibling(task) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_earlier_incomplete_code_sibling_false_for_higher_seq_only() -> None:
|
|
# A LATER sibling (seq 3) does not hold an earlier leaf (seq 2).
|
|
task = _build_task(
|
|
task_type=TaskType.CODE.value,
|
|
parent_task_id=uuid4(),
|
|
assigned_to=uuid4(),
|
|
sequence=2,
|
|
)
|
|
svc = _svc_with_sibling_status_seq([(TaskStatus.IN_PROGRESS, 3)])
|
|
assert await svc.has_earlier_incomplete_code_sibling(task) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_earlier_incomplete_code_sibling_non_code_short_circuits() -> None:
|
|
# Only code queues sequence this way; a planning/doc leaf never queries.
|
|
session = MagicMock(execute=AsyncMock())
|
|
svc = TaskService(session)
|
|
task = _build_task(
|
|
task_type="planning",
|
|
parent_task_id=uuid4(),
|
|
assigned_to=uuid4(),
|
|
sequence=2,
|
|
)
|
|
assert await svc.has_earlier_incomplete_code_sibling(task) is False
|
|
session.execute.assert_not_awaited()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_earlier_incomplete_code_sibling_false_when_fields_missing() -> None:
|
|
task = _build_task(
|
|
task_type=TaskType.CODE.value,
|
|
parent_task_id=None,
|
|
assigned_to=uuid4(),
|
|
sequence=2,
|
|
)
|
|
svc = _svc_with_sibling_status_seq([(TaskStatus.IN_PROGRESS, 0)])
|
|
assert await svc.has_earlier_incomplete_code_sibling(task) is False
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_unblock_with_branch_resumes_in_progress() -> None:
|
|
# A task claimed (has a branch) before it blocked resumes in_progress.
|
|
task = _build_task(
|
|
status=TaskStatus.BLOCKED,
|
|
branch_name="feature/backend/abc12345",
|
|
blocker_raised_by=uuid4(),
|
|
)
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "_index_lifecycle_event_background", AsyncMock())
|
|
out = await svc.unblock(task.id)
|
|
assert out is task
|
|
assert task.status == TaskStatus.IN_PROGRESS
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# cell_pm_complete
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cell_pm_complete_appends_merge_commit() -> None:
|
|
task = _build_task(commits=[{"hash": "old", "message": "earlier"}])
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
complete_mock = AsyncMock(return_value=task)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "complete", complete_mock)
|
|
pm_id = uuid4()
|
|
await svc.cell_pm_complete(pm_id, task.id, "all good", merge_commit="deadbeef")
|
|
assert task.commits[-1]["hash"] == "deadbeef"
|
|
assert task.commits[-1]["kind"] == "merge"
|
|
complete_mock.assert_awaited_once_with(task.id, agent_id=pm_id)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_cell_pm_complete_skips_merge_when_none() -> None:
|
|
task = _build_task(commits=[])
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
complete_mock = AsyncMock(return_value=task)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
_bind(svc, "complete", complete_mock)
|
|
await svc.cell_pm_complete(uuid4(), task.id, "all good", merge_commit=None)
|
|
assert task.commits == []
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# escalate / escalate_up_to_role
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_escalate_returns_none_when_no_target_configured(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
task = _build_task()
|
|
agent = MagicMock(id=uuid4(), slug="lone-agent")
|
|
agent_result = MagicMock()
|
|
agent_result.scalar_one_or_none.return_value = agent
|
|
|
|
session = MagicMock()
|
|
session.execute = AsyncMock(return_value=agent_result)
|
|
session.flush = AsyncMock()
|
|
svc = TaskService(session)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
monkeypatch.setattr(
|
|
"roboco.agents_config.get_escalation_target", lambda _slug: None
|
|
)
|
|
out = await svc.escalate(uuid4(), task.id, "stuck")
|
|
assert out is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_escalate_up_to_role_returns_none_for_unknown_role() -> None:
|
|
task = _build_task()
|
|
agent = MagicMock(id=uuid4(), slug="some-agent")
|
|
agent_result = MagicMock()
|
|
agent_result.scalar_one_or_none.return_value = agent
|
|
|
|
session = MagicMock()
|
|
session.execute = AsyncMock(return_value=agent_result)
|
|
svc = TaskService(session)
|
|
_bind(svc, "get", AsyncMock(return_value=task))
|
|
out = await svc.escalate_up_to_role(uuid4(), task.id, "bogus_role", "reason")
|
|
assert out is None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# _ensure_branch_for_task — coordination/fan-out tasks do no git
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ensure_branch_returns_existing_branch() -> None:
|
|
"""An already-branched task short-circuits before any project check."""
|
|
svc = TaskService(MagicMock())
|
|
task = MagicMock(branch_name="feature/backend/abc12345", project_id=None)
|
|
assert (
|
|
await svc._ensure_branch_for_task(task, uuid4()) == "feature/backend/abc12345"
|
|
)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ensure_branch_coordination_root_cuts_integration_branch() -> None:
|
|
"""A product-backed root cuts feature/main_pm/{root} in each product repo."""
|
|
svc = TaskService(MagicMock())
|
|
task = MagicMock(branch_name=None, project_id=None, product_id=uuid4())
|
|
create_in_project = AsyncMock(return_value="feature/main_pm/root1234")
|
|
_bind(svc, "_create_branch_in_project", create_in_project)
|
|
product_svc = MagicMock(distinct_project_ids=AsyncMock(return_value=[uuid4()]))
|
|
project_svc = MagicMock(get=AsyncMock(return_value=MagicMock()))
|
|
with (
|
|
patch("roboco.services.product.get_product_service", return_value=product_svc),
|
|
patch("roboco.services.project.get_project_service", return_value=project_svc),
|
|
):
|
|
result = await svc._ensure_branch_for_task(task, uuid4())
|
|
assert result == "feature/main_pm/root1234"
|
|
create_in_project.assert_awaited_once()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ensure_branch_coordination_root_no_cell_map_stays_branchless() -> None:
|
|
"""A product with no cell->repo map yet stays branchless (graceful fallback)."""
|
|
svc = TaskService(MagicMock())
|
|
task = MagicMock(branch_name=None, project_id=None, product_id=uuid4())
|
|
product_svc = MagicMock(distinct_project_ids=AsyncMock(return_value=[]))
|
|
with patch("roboco.services.product.get_product_service", return_value=product_svc):
|
|
result = await svc._ensure_branch_for_task(task, uuid4())
|
|
assert result == ""
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ensure_branch_cell_map_root_cuts_integration_branch_per_project() -> (
|
|
None
|
|
):
|
|
"""An ad-hoc cell_projects root cuts feature/main_pm/{root} in each distinct
|
|
project the map spans — the product-root path with the map sourced from the
|
|
task instead of a Product."""
|
|
svc = TaskService(MagicMock())
|
|
be_proj, fe_proj = uuid4(), uuid4()
|
|
cell_map = [
|
|
SimpleNamespace(team=Team.BACKEND, project_id=be_proj),
|
|
SimpleNamespace(team=Team.FRONTEND, project_id=fe_proj),
|
|
]
|
|
task = MagicMock(
|
|
branch_name=None,
|
|
project_id=None,
|
|
product_id=None,
|
|
batch_id=uuid4(),
|
|
parent_task_id=uuid4(),
|
|
cell_projects=cell_map,
|
|
)
|
|
create_in_project = AsyncMock(return_value="feature/main_pm/root1234")
|
|
_bind(svc, "_create_branch_in_project", create_in_project)
|
|
project_svc = MagicMock(get=AsyncMock(return_value=MagicMock()))
|
|
with patch("roboco.services.project.get_project_service", return_value=project_svc):
|
|
result = await svc._ensure_branch_for_task(task, uuid4())
|
|
assert result == "feature/main_pm/root1234"
|
|
# one integration branch per distinct project in the map (here 2 cells, 2 projects)
|
|
assert create_in_project.await_count == len(cell_map)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_distinct_projects_for_task_dedupes_cell_map_by_project_id() -> None:
|
|
"""Two cells mapping at the same project (the monorepo case) yield ONE
|
|
integration branch, not two — mirroring product distinct_project_ids."""
|
|
svc = TaskService(MagicMock())
|
|
shared = uuid4()
|
|
task = MagicMock(
|
|
project_id=None,
|
|
product_id=None,
|
|
cell_projects=[
|
|
SimpleNamespace(team=Team.FRONTEND, project_id=shared),
|
|
SimpleNamespace(team=Team.BACKEND, project_id=shared),
|
|
],
|
|
)
|
|
ids = await svc._distinct_projects_for_task(task)
|
|
assert ids == [shared]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_ensure_branch_raises_when_neither_project_nor_product() -> None:
|
|
"""A task with neither a project, a product, nor a cell map is misconfigured."""
|
|
svc = TaskService(MagicMock())
|
|
task = MagicMock(
|
|
branch_name=None,
|
|
project_id=None,
|
|
product_id=None,
|
|
cell_projects=[],
|
|
batch_id=None,
|
|
parent_task_id=None,
|
|
)
|
|
with pytest.raises(ValueError, match="project_id"):
|
|
await svc._ensure_branch_for_task(task, uuid4())
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# _finalize_claim — branch-creation failure rollback (F060)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_finalize_claim_rollback_emits_reversal_audit() -> None:
|
|
"""When branch creation fails mid-claim, the rollback must emit a REVERSAL
|
|
audit row (CLAIMED -> original) so the audit journey matches the real
|
|
(rolled-back) task state. The audit service writes on its own connection, so
|
|
the forward `task.claimed` row is NOT undone by the rollback's flush.
|
|
"""
|
|
session = MagicMock()
|
|
session.flush = AsyncMock()
|
|
svc = TaskService(session)
|
|
|
|
task = _build_task(
|
|
status=TaskStatus.PENDING,
|
|
branch_name=None, # forces the branch-creation path
|
|
project_id=uuid4(),
|
|
product_id=None,
|
|
batch_id=None,
|
|
parent_task_id=None,
|
|
cell_projects=[], # a plain code task, not a branchless coordination root
|
|
pr_created=False,
|
|
pr_number=None,
|
|
)
|
|
agent = MagicMock(id=uuid4(), role=AgentRole.DEVELOPER)
|
|
|
|
audit_calls: list[dict[str, Any]] = []
|
|
|
|
def _capture(
|
|
_task: object, *, from_status: str, to_status: str, **_kw: object
|
|
) -> None:
|
|
audit_calls.append({"from": from_status, "to": to_status})
|
|
|
|
_bind(svc, "_emit_status_transition_audit", _capture)
|
|
_bind(
|
|
svc,
|
|
"_ensure_branch_for_task",
|
|
AsyncMock(side_effect=RuntimeError("branch boom")),
|
|
)
|
|
|
|
with pytest.raises(RuntimeError, match="branch boom"):
|
|
await svc._finalize_claim(task, agent, agent.id)
|
|
|
|
# The task reverted to its pre-claim status (the existing behavior).
|
|
assert task.status == TaskStatus.PENDING
|
|
# The forward claim row was emitted...
|
|
assert {"from": "pending", "to": "claimed"} in audit_calls
|
|
# ...AND the reversal row is emitted so the journey's last event matches
|
|
# the rolled-back state (the F060 fix). Before the fix this was missing.
|
|
assert {"from": "claimed", "to": "pending"} in audit_calls
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_emit_status_transition_audit_writes_in_session_atomically() -> None:
|
|
"""The status-transition audit row is written into the CALLER's session (same
|
|
transaction as the transition), not fire-and-forget on a separate connection,
|
|
so it commits/rolls back atomically with the transition and cannot diverge
|
|
from real state. Asserted at the unit level: the row is ``session.add``-ed
|
|
(same txn) and NO fire-and-forget background task is spawned.
|
|
"""
|
|
session = MagicMock()
|
|
added: list[object] = []
|
|
session.add.side_effect = added.append
|
|
svc = TaskService(session)
|
|
prior_bg = set(svc._background_tasks)
|
|
|
|
task = MagicMock(id=uuid4(), claimed_by=uuid4(), team=Team.BACKEND)
|
|
|
|
svc._emit_status_transition_audit(
|
|
task,
|
|
from_status="pending",
|
|
to_status="claimed",
|
|
agent_role="developer",
|
|
audit_agent_id=None,
|
|
)
|
|
|
|
# The audit row is added to the CALLER's session (same transaction) — not
|
|
# dispatched to a separate fire-and-forget connection.
|
|
rows = [r for r in added if isinstance(r, AuditLogTable)]
|
|
assert len(rows) == 1
|
|
row = rows[0]
|
|
assert row.event_type == "task.claimed"
|
|
assert row.target_type == "task"
|
|
assert row.target_id == task.id
|
|
assert row.details["from_status"] == "pending"
|
|
assert row.details["to_status"] == "claimed"
|
|
assert row.details["agent_role"] == "developer"
|
|
assert row.details["team"] == "backend"
|
|
# The claiming agent is attributed (resolved from claimed_by).
|
|
assert row.agent_id == task.claimed_by
|
|
# No fire-and-forget audit task was spawned (the old decoupled path).
|
|
assert svc._background_tasks == prior_bg
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# _resolve_doc_abspath — normalize documenter-supplied paths under /app/docs
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def test_resolve_doc_abspath_strips_redundant_docs_prefix() -> None:
|
|
"""A `docs/`-rooted relative path must not double the base segment.
|
|
|
|
DOCS_BASE_PATH is /app/docs; joining it with `docs/design/x.md` produced
|
|
/app/docs/docs/design/x.md, so the file was never found and never indexed.
|
|
"""
|
|
assert (
|
|
TaskService._resolve_doc_abspath("docs/design/spec.md")
|
|
== "/app/docs/design/spec.md"
|
|
)
|
|
|
|
|
|
def test_resolve_doc_abspath_keeps_plain_relative_path() -> None:
|
|
"""A relative path with no `docs/` prefix joins under the base unchanged."""
|
|
assert (
|
|
TaskService._resolve_doc_abspath("design/spec.md") == "/app/docs/design/spec.md"
|
|
)
|
|
|
|
|
|
def test_resolve_doc_abspath_passes_absolute_path_through() -> None:
|
|
"""An absolute path already correctly rooted under the base is unchanged."""
|
|
assert (
|
|
TaskService._resolve_doc_abspath("/app/docs/design/spec.md")
|
|
== "/app/docs/design/spec.md"
|
|
)
|
|
|
|
|
|
def test_resolve_doc_abspath_collapses_doubled_absolute_docs() -> None:
|
|
"""An absolute path that doubled the base segment is collapsed.
|
|
|
|
The documenter sometimes records `/app/docs/docs/...`; previously it was
|
|
returned verbatim and never resolved on disk (the recurring "Source not
|
|
found" warning).
|
|
"""
|
|
assert (
|
|
TaskService._resolve_doc_abspath("/app/docs/docs/backend/api/prompter.md")
|
|
== "/app/docs/backend/api/prompter.md"
|
|
)
|
|
|
|
|
|
def test_resolve_doc_abspath_leaves_external_absolute_path() -> None:
|
|
"""An absolute path outside the docs root is left as-is for the indexer to skip."""
|
|
external = "/data/workspaces/panel/frontend/fe-dev-1/src/page.tsx"
|
|
assert TaskService._resolve_doc_abspath(external) == external
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_update_skips_none_to_protect_partial_callers() -> None:
|
|
"""update() must skip None values, not write them.
|
|
|
|
Callers pass field=dict.get('x'), which is None when the key is absent —
|
|
e.g. the board-redraft update_live_draft path passes title/acceptance_criteria
|
|
that way. Without the None-skip guard those None values would null-wipe
|
|
existing data. Explicit clearing is the update ROUTE's job (a field
|
|
whitelist), never this shared service method. Locks that contract so the
|
|
guard can't be silently removed again.
|
|
"""
|
|
task = SimpleNamespace(title="original", acceptance_criteria=["keep me"])
|
|
svc = TaskService(MagicMock(flush=AsyncMock()))
|
|
with patch.object(svc, "get", AsyncMock(return_value=task)):
|
|
result: Any = await svc.update(
|
|
uuid4(), title="updated", acceptance_criteria=None
|
|
)
|
|
|
|
assert result is task
|
|
assert task.title == "updated" # explicit, non-None value is applied
|
|
assert task.acceptance_criteria == ["keep me"] # None skipped, not wiped
|