diff --git a/roboco/agent_sdk/intake_driver.py b/roboco/agent_sdk/intake_driver.py index 7d21785b..974853b8 100644 --- a/roboco/agent_sdk/intake_driver.py +++ b/roboco/agent_sdk/intake_driver.py @@ -103,22 +103,30 @@ def _coerce_draft(data: Any) -> dict[str, Any] | None: if not (isinstance(data, dict) and isinstance(data.get("title"), str)): return None coerced = dict(data) + _coerce_spec_fields(coerced, coerce_str_list) + return coerced + + +def _coerce_spec_fields( + coerced: dict[str, Any], coerce: Callable[[Any], list[str]] +) -> None: + """Coerce the list-shaped spec fields in place: wrap the_work, flatten the + string-list fields, and flatten each the_work unit's items to ``list[str]``.""" if "the_work" in coerced: coerced["the_work"] = _coerce_to_list(coerced["the_work"]) for key in _STR_LIST_FIELDS: if key in coerced: - coerced[key] = coerce_str_list(coerced[key]) + coerced[key] = coerce(coerced[key]) work = coerced.get("the_work") if isinstance(work, list): coerced["the_work"] = [ { **unit, - "items": coerce_str_list(unit.get("items")), + "items": coerce(unit.get("items")), } for unit in work if isinstance(unit, dict) ] - return coerced def _extract_draft(text: str) -> dict[str, Any] | None: diff --git a/roboco/foundation/policy/content/validators.py b/roboco/foundation/policy/content/validators.py index f47b2177..a41fe8d5 100644 --- a/roboco/foundation/policy/content/validators.py +++ b/roboco/foundation/policy/content/validators.py @@ -116,13 +116,7 @@ def _extract_strs(item: Any) -> list[str]: stripped = item.strip() return [stripped] if stripped else [] if isinstance(item, dict): - for key in _TEXT_KEYS: - if key in item: - return _extract_strs(item[key]) - # No recognized key: keep any bare string values the dict carries. - return [ - str(v).strip() for v in item.values() if isinstance(v, str) and v.strip() - ] + return _extract_strs_from_dict(item) if isinstance(item, list): out: list[str] = [] for sub in item: @@ -131,6 +125,15 @@ def _extract_strs(item: Any) -> list[str]: return [] +def _extract_strs_from_dict(item: dict) -> list[str]: + """Resolve a dict element to strings via the first recognized text key, + else any bare string values it carries (recurses through SDK ``$text`` wrappers).""" + for key in _TEXT_KEYS: + if key in item: + return _extract_strs(item[key]) + return [str(v).strip() for v in item.values() if isinstance(v, str) and v.strip()] + + def coerce_str_list(value: Any) -> list[str]: """Coerce a list-of-strings field to a flat ``list[str]``. diff --git a/roboco/runtime/orchestrator.py b/roboco/runtime/orchestrator.py index 5221b427..6de59376 100644 --- a/roboco/runtime/orchestrator.py +++ b/roboco/runtime/orchestrator.py @@ -5961,6 +5961,54 @@ Start by: return None return rc == 0 + async def _maybe_park_for_exit_error( + self, agent_id: str, instance: Any, graceful: bool + ) -> bool: + """Park the provider on a session/usage limit or a server overload detected + in the dead container's output, instead of crash-retrying into it. Returns + True when parked (caller returns); False to proceed with normal handling. + The probe-resume loop revives the task when the limit lifts / overload clears. + """ + if graceful: + return False + rate_limited_provider = await self._provider_rate_limit_park_target( + agent_id, instance + ) + if rate_limited_provider is not None: + logger.warning( + "Session/usage limit detected in agent output; parking provider", + agent_id=agent_id, + provider=rate_limited_provider, + task_id=instance.current_task_id, + ) + await self._park_provider_unavailable( + agent_id, + instance, + provider=rate_limited_provider, + retry_after=_RATE_LIMIT_RETRY_AFTER_S, + kind="rate_limited", + ) + return True + overloaded_provider = await self._provider_overload_park_target( + agent_id, instance + ) + if overloaded_provider is not None: + logger.warning( + "Provider overload detected in agent output; parking provider", + agent_id=agent_id, + provider=overloaded_provider, + task_id=instance.current_task_id, + ) + await self._park_provider_unavailable( + agent_id, + instance, + provider=overloaded_provider, + retry_after=_OVERLOAD_RETRY_AFTER_S, + kind="overloaded", + ) + return True + return False + async def _handle_stopped_container( self, agent_id: str, instance: Any, exit_code: int | None ) -> None: @@ -5990,53 +6038,11 @@ Start by: await self._park_grok_auth_unavailable(agent_id, instance) return graceful = exit_code == 0 - # Session/usage-limit parking: the Claude session ("5-hour") limit is a - # 429 the SDK does not retry — the container exits non-zero with a - # 0-token rejection. Detect it in the dead container's output and park - # the provider (instead of crash-respawning straight back into the - # limit); the probe-resume loop revives the task when the quota resets. - if not graceful: - rate_limited_provider = await self._provider_rate_limit_park_target( - agent_id, instance - ) - if rate_limited_provider is not None: - logger.warning( - "Session/usage limit detected in agent output; parking provider", - agent_id=agent_id, - provider=rate_limited_provider, - task_id=instance.current_task_id, - ) - await self._park_provider_unavailable( - agent_id, - instance, - provider=rate_limited_provider, - retry_after=_RATE_LIMIT_RETRY_AFTER_S, - kind="rate_limited", - ) - return - # Server-overload parking: a persistent 529/500/503 from the model API - # kills the run (the SDK already retries transient ones). Detect the - # overload marker in the dead container's output and park the provider — - # the same break as a 429 — instead of crash-retrying into the overload. - if not graceful: - overloaded_provider = await self._provider_overload_park_target( - agent_id, instance - ) - if overloaded_provider is not None: - logger.warning( - "Provider overload detected in agent output; parking provider", - agent_id=agent_id, - provider=overloaded_provider, - task_id=instance.current_task_id, - ) - await self._park_provider_unavailable( - agent_id, - instance, - provider=overloaded_provider, - retry_after=_OVERLOAD_RETRY_AFTER_S, - kind="overloaded", - ) - return + # Park the provider on a session/usage limit or a server overload detected + # in the dead container's output instead of crash-retrying into it. The + # probe-resume loop revives the task when the limit lifts / overload clears. + if await self._maybe_park_for_exit_error(agent_id, instance, graceful): + return if graceful: logger.info( "Agent container exited gracefully", @@ -8590,6 +8596,26 @@ Start now: evidence(task_id="{task_id}") grace = settings.gateway_health_grace_seconds return (now - first_seen).total_seconds() >= grace + async def _should_skip_live_reap(self, t: Any, ts: Any) -> bool: + """True when a live container should be spared from reaping. + + A live container normally protects its task; on a registry MISS (e.g. the + orchestrator restarted and forgot a still-running container) fall back to + asking Docker. A live container is spared UNLESS it is wedged (grok) or + its gateway is broken-but-alive past the grace window — both checks kill + + evict it (returning False here) so the caller falls through to release + + respawn. Short-circuits like the original ``and``: when not live, neither + kill nor recovery check is awaited. + """ + live = self._assignee_has_active_instance( + t + ) or await self._assignee_container_running(t) + return ( + live + and not await self._maybe_kill_wedged_grok(t, ts) + and not await self._maybe_recover_broken_gateway(t) + ) + async def _reap_with_service(self, svc: "TaskService") -> None: """Inner reap loop, parameterized by the TaskService to use. @@ -8606,26 +8632,11 @@ Start now: evidence(task_id="{task_id}") for t in candidates: ts = t.last_heartbeat_at if ts is None or ts < cutoff: - # A live container normally protects its task. Prefer the - # in-memory registry; on a registry MISS (e.g. the orchestrator - # restarted and forgot a still-running container) fall back to - # asking Docker, so we don't reap a task out from under a live - # agent. The sole exception is a wedged GROK container — ACTIVE - # yet firing no verb — which the live skip would shield forever: - # kill + evict it past the grok-idle TTL (then fall through to - # release); a live non-grok agent, or a grok within the TTL, is - # skipped. - live = self._assignee_has_active_instance( - t - ) or await self._assignee_container_running(t) - # A live container is spared UNLESS it is wedged (grok) or its - # gateway is broken-but-alive past the grace window — both get - # killed + evicted here so we fall through to release + respawn. - if ( - live - and not await self._maybe_kill_wedged_grok(t, ts) - and not await self._maybe_recover_broken_gateway(t) - ): + # A live container is spared unless it is wedged (grok) or its + # gateway is broken-but-alive past the grace window — see + # _should_skip_live_reap, which kills + evicts those so we fall + # through to release + respawn. + if await self._should_skip_live_reap(t, ts): continue # A provider-parked agent (session-limit / overload / grok-429) # is OFFLINE with a dead container and a ``rate_limit_lifted`` diff --git a/roboco/services/gateway/choreographer/_impl.py b/roboco/services/gateway/choreographer/_impl.py index 26d1816e..5480d4c2 100644 --- a/roboco/services/gateway/choreographer/_impl.py +++ b/roboco/services/gateway/choreographer/_impl.py @@ -61,6 +61,7 @@ if TYPE_CHECKING: # etc.) via MRO, but ``_LegacyChoreographer`` itself does not inherit # ``ChoreographerHelpers`` — so mypy can't see those names on ``self`` here. # The cast below reaches the typed view the mixins use (``_Base`` pattern). + from roboco.models.base import TaskNature from roboco.services.gateway.choreographer._protocol import ChoreographerHelpers # Minimum character length enforced on rich_plan["approach"] by the PM @@ -5136,34 +5137,19 @@ class Choreographer: message=f"cannot resolve a project for the {inputs.team} subtask", ) - async def _create_subtask_from_inputs( - self, - pm_agent_id: UUID, - parent_task_id: UUID, - parent: Any, + @staticmethod + def _require_subtask_completeness( inputs: DelegateInputs, - ) -> Any: - """Resolve enums + AGENT_UUIDS slug and call TaskService.create_subtask. - - By contract, callers (the `delegate` verb body) MUST run - `_delegate_completeness_check` first, so `inputs.acceptance_criteria` - and `inputs.nature` are guaranteed non-None / non-empty here. The - defensive `TaskCompletenessError` raises preserve correctness if - a future caller bypasses the gateway path — defense-in-depth in - line with the service-layer raise. + ) -> tuple[TaskNature, list[str]]: + """Defensive completeness check (callers run ``_delegate_completeness_check`` + first); returns the validated ``TaskNature`` and the non-empty + ``acceptance_criteria``. Raises ``TaskCompletenessError`` with field hints + if a non-gateway caller bypassed the gateway check — never silently + substitutes an empty acceptance list. """ from roboco.foundation.policy.task_completeness import TaskCompletenessError from roboco.models.base import TaskNature - from roboco.models.task import TaskCreateRequest - from roboco.seeds.initial_data import AGENT_UUIDS - team_enum, type_enum, complexity_enum = self._resolve_delegate_enums(inputs) - assignee_id = UUID(AGENT_UUIDS[inputs.assigned_to]) - # The `or []` collapse was removed. The gateway runs - # `_delegate_completeness_check` BEFORE this helper, so empty/None - # acceptance_criteria here means a non-gateway caller bypassed the - # check. Raise so the service-layer raise can attach the - # field hints — never silently substitute. if not inputs.acceptance_criteria: raise TaskCompletenessError( missing=["acceptance_criteria"], @@ -5180,9 +5166,7 @@ class Choreographer: if inputs.nature is None: raise TaskCompletenessError( missing=["nature"], - field_hints={ - "nature": "one of: technical | non_technical", - }, + field_hints={"nature": "one of: technical | non_technical"}, message=( "_create_subtask_from_inputs called with no nature — " "completeness check must run first" @@ -5196,11 +5180,39 @@ class Choreographer: field_hints={"nature": "one of: technical | non_technical"}, message=f"invalid nature {inputs.nature!r}: {exc}", ) from exc + return nature_enum, inputs.acceptance_criteria + + async def _create_subtask_from_inputs( + self, + pm_agent_id: UUID, + parent_task_id: UUID, + parent: Any, + inputs: DelegateInputs, + ) -> Any: + """Resolve enums + AGENT_UUIDS slug and call TaskService.create_subtask. + + By contract, callers (the `delegate` verb body) MUST run + `_delegate_completeness_check` first, so `inputs.acceptance_criteria` + and `inputs.nature` are guaranteed non-None / non-empty here. The + defensive `TaskCompletenessError` raises preserve correctness if + a future caller bypasses the gateway path — defense-in-depth in + line with the service-layer raise. + """ + from roboco.models.task import TaskCreateRequest + from roboco.seeds.initial_data import AGENT_UUIDS + + team_enum, type_enum, complexity_enum = self._resolve_delegate_enums(inputs) + assignee_id = UUID(AGENT_UUIDS[inputs.assigned_to]) + # The gateway runs `_delegate_completeness_check` BEFORE this helper, so + # empty/None acceptance_criteria or nature here means a non-gateway caller + # bypassed the check — _require_subtask_completeness raises with field + # hints rather than silently substituting. + nature_enum, acceptance_criteria = self._require_subtask_completeness(inputs) resolved_project_id = await self._resolve_subtask_project(parent, inputs) req = TaskCreateRequest( title=inputs.title, description=inputs.description, - acceptance_criteria=inputs.acceptance_criteria, + acceptance_criteria=acceptance_criteria, parent_ac_refs=inputs.covers_parent_criteria or [], team=team_enum, created_by=pm_agent_id, @@ -6542,6 +6554,18 @@ class Choreographer: context_briefing=await self._briefing_for(main_pm_agent_id, root_task_id), ).with_introspection(task=t, role="main_pm") + @staticmethod + def _is_umbrella_in_progress(t: Any, role_str: str) -> bool: + """A MegaTask umbrella sits in_progress branchless (no submit_root/pr_pass), + so the complete spec gate (AWAITING_PM_REVIEW only) would reject it before + main_pm_complete's branchless-aware guard runs — this lets it fall through + to main_pm_complete (CEO merges the root PR; no agent touches master).""" + return ( + role_str == "main_pm" + and str(t.status) == "in_progress" + and is_batch_umbrella(batch_id=t.batch_id, parent_task_id=t.parent_task_id) + ) + async def complete(self, agent_id: UUID, task_id: UUID, notes: str) -> Envelope: """Dispatch to cell_pm_complete or main_pm_complete based on agent role. @@ -6596,19 +6620,11 @@ class Choreographer: verb="complete", ): return soup - # A MegaTask umbrella is branchless by design and never goes through - # submit_root / pr_pass, so it sits in in_progress with no branch/PR. - # The ``complete`` action's source_statuses={AWAITING_PM_REVIEW} spec - # gate would reject it before main_pm_complete's branchless-aware guard - # can run; skip the spec gate for an in_progress batch umbrella and fall - # through to main_pm_complete (CEO merges the root PR; no agent touches - # master). Role membership is preserved; main_pm_complete re-checks - # assignment, subtasks-terminal, and the journal:decision gate. - umbrella_in_progress = ( - role_str == "main_pm" - and str(t.status) == "in_progress" - and is_batch_umbrella(batch_id=t.batch_id, parent_task_id=t.parent_task_id) - ) + # An in_progress batch umbrella is branchless and skips the complete spec + # gate (AWAITING_PM_REVIEW only) to fall through to main_pm_complete — + # see _is_umbrella_in_progress. Role membership is preserved; main_pm_complete + # re-checks assignment, subtasks-terminal, and the journal:decision gate. + umbrella_in_progress = self._is_umbrella_in_progress(t, role_str) if not umbrella_in_progress: decision = spec_module.can_invoke_intent(role, "complete", t, spec_ctx) if not decision.allowed: diff --git a/roboco/services/gateway/choreographer/pr_gate.py b/roboco/services/gateway/choreographer/pr_gate.py index acdf4b12..c1846ab4 100644 --- a/roboco/services/gateway/choreographer/pr_gate.py +++ b/roboco/services/gateway/choreographer/pr_gate.py @@ -226,6 +226,69 @@ class PRGateMixin(_Base): ) return (t, agent, role_str, briefing, spec_ctx) + async def _record_gate_verdict_for( + self, verb: str, t: Any, notes: str, *, issues: tuple[str, ...] + ) -> None: + """Author the canonical pr_review verdict note before the transition. + + On pr_fail also stamp the assembled PR's head SHA so the next submit_root + can structurally refuse to re-submit the unchanged root (the 2026-06-27 + infinite pr_fail re-submit loop). Best-effort: a capture failure leaves + head_sha absent and submit_root fails open rather than wedging the PM. + """ + if verb == "pr_fail": + head_sha = await self._capture_pr_head_sha(t) + self._record_gate_verdict(t, verb, notes, issues=issues, head_sha=head_sha) + else: + self._record_gate_verdict(t, verb, notes, issues=issues) + + async def _post_gate_review( + self, t: Any, agent: Any, role_str: str, verb: str, notes: str + ) -> None: + """Post the gate verdict to the PR itself (best-effort, after the DB + transition — a GitHub failure must not roll back the gate decision).""" + reviewer_slug = getattr(agent, "slug", None) or role_str + await self._post_gate_review_to_pr(t, verb, reviewer_slug, notes) + + async def _deliver_pr_fail_to_owner( + self, t: Any, reviewer_agent_id: UUID, task_id: UUID, notes: str + ) -> None: + """a2a the pr_fail change-requests to the owning PM (best-effort). + + The verdict is posted on the PR but never reaches a PM-readable channel + (no a2a, and _briefing_for / build_task_handoff read neither + pr_reviewer_notes nor notes_structured.pr_review), so without this the + owning PM respawned into needs_revision re-submits the same PR blind — + an infinite pr_fail loop (live on 9980d0a0 / PR #138). Mirrors QA's + fail_review a2a. A Main-PM branch-bearing root is assembled cell work the + Main PM can't fix directly, so steer it to re-delegate + wait for + re-assembly rather than re-submit the unchanged root. + """ + if t.assigned_to is None: + return + team = getattr(t, "team", None) + team_value = str(getattr(team, "value", team)) + is_main_pm_root = team_value == spec_module.Team.MAIN_PM.value and bool( + getattr(t, "branch_name", None) + ) + steer = ( + " Assembled cell work failed — re-delegate the fixes to the" + " owning cell PM(s) and wait for re-assembly; do NOT re-submit" + " the root." + if is_main_pm_root + else "" + ) + try: + await self.a2a.send( + from_agent=reviewer_agent_id, + to_agent=t.assigned_to, + skill="code_review", + task_id=task_id, + body=f"PR review needs changes. {notes}{steer}", + ) + except Exception: + logger.exception("pr_fail a2a to owning PM failed", task_id=str(task_id)) + async def _gate_decision( self, reviewer_agent_id: UUID, @@ -253,21 +316,10 @@ class PRGateMixin(_Base): ) if blocked is not None: return blocked - # Author the canonical pr_review verdict note BEFORE the transition so - # it is persisted by the same commit (mirrors post_pr_review). This is - # what keeps notes_structured.pr_review in lock-step with the decision — - # a later pr_fail overwrites an earlier pr_pass verdict instead of - # leaving a stale "passed" on a task that was just sent back. On pr_fail - # also stamp the assembled PR's head SHA so the next submit_root can - # structurally refuse to re-submit the unchanged root (the 2026-06-27 - # infinite pr_fail re-submit loop). Best-effort: a capture failure - # (no token, no PR, GitHub error) leaves head_sha absent and the - # submit_root gate fails open rather than wedging the PM. - if verb == "pr_fail": - head_sha = await self._capture_pr_head_sha(t) - self._record_gate_verdict(t, verb, notes, issues=issues, head_sha=head_sha) - else: - self._record_gate_verdict(t, verb, notes, issues=issues) + # Author the canonical pr_review verdict note BEFORE the transition so it + # is persisted by the same commit (mirrors post_pr_review) and stays in + # lock-step with the decision (pr_fail overwrites an earlier pr_pass). + await self._record_gate_verdict_for(verb, t, notes, issues=issues) runner = self._verb_runner() try: t = await runner.run_intent(verb, t, agent, spec_ctx) @@ -304,54 +356,13 @@ class PRGateMixin(_Base): task_id=task_id, verb=verb, ) - # Leave the gate verdict on the PR itself so there's a visible trail on - # the very PR the PM (or CEO) merges. Best-effort and AFTER the DB - # transition — a GitHub failure must not roll back the gate decision. - reviewer_slug = getattr(agent, "slug", None) or role_str - await self._post_gate_review_to_pr(t, verb, reviewer_slug, notes) - # Deliver the change-requests to the owner that now has to act on them - # — the cell PM the runner just re-assigned via _revision_pm_for_task. - # The reviewer posts the verdict on the PR itself but that never reaches - # any PM-readable channel (no a2a, and _briefing_for / build_task_handoff - # read neither pr_reviewer_notes nor notes_structured.pr_review). Without - # this the owning PM respawned into needs_revision saw a generic "needs - # revision" with zero concrete issues, concluded nothing to rework, and - # re-submitted the same PR — an infinite pr_fail loop (live on - # 9980d0a0 / PR #138). Mirrors QA's fail_review a2a to the dev (qa.py:671). - # Best-effort: the transition already committed, so a delivery failure - # must not roll the verdict back or 500 the reviewer. - if verb == "pr_fail" and t.assigned_to is not None: - # A Main-PM branch-bearing root is an assembled cell→root / root→master - # PR — coordination, not the Main PM's own code. The rejection is - # about the cells' merged code, which the Main PM cannot fix directly - # (no code verb). Steer the a2a body to re-delegate + wait for - # re-assembly so the PM doesn't re-submit the unchanged root (the - # 2026-06-27 infinite pr_fail loop). The Envelope ``next`` hint makes - # the same steer via _next_hint_pr_fail. - team = getattr(t, "team", None) - team_value = str(getattr(team, "value", team)) - is_main_pm_root = team_value == spec_module.Team.MAIN_PM.value and bool( - getattr(t, "branch_name", None) - ) - steer = ( - " Assembled cell work failed — re-delegate the fixes to the" - " owning cell PM(s) and wait for re-assembly; do NOT re-submit" - " the root." - if is_main_pm_root - else "" - ) - try: - await self.a2a.send( - from_agent=reviewer_agent_id, - to_agent=t.assigned_to, - skill="code_review", - task_id=task_id, - body=f"PR review needs changes. {notes}{steer}", - ) - except Exception: - logger.exception( - "pr_fail a2a to owning PM failed", task_id=str(task_id) - ) + # Post the gate verdict on the PR itself (best-effort, after the DB + # transition — a GitHub failure must not roll back the gate decision). + await self._post_gate_review(t, agent, role_str, verb, notes) + # a2a the pr_fail change-requests to the owning PM (best-effort) — see + # _deliver_pr_fail_to_owner for the rationale and the Main-PM-root steer. + if verb == "pr_fail": + await self._deliver_pr_fail_to_owner(t, reviewer_agent_id, task_id, notes) return Envelope.ok( status=str(t.status), task_id=str(task_id), diff --git a/roboco/services/gateway/evidence_builder.py b/roboco/services/gateway/evidence_builder.py index a6693997..59116213 100644 --- a/roboco/services/gateway/evidence_builder.py +++ b/roboco/services/gateway/evidence_builder.py @@ -86,6 +86,27 @@ def _typed(value: Any, expected: type | tuple[type, ...], default: Any) -> Any: return value if isinstance(value, expected) else default +def _has_prior_work( + commits: list, + acceptance: list, + highlights: list, + pr_number: int | None, + dev_summary: str | None, + completed_deps: list, + pr_review: dict[str, Any] | None, +) -> bool: + """True when any resumable prior-work signal is present on the task.""" + return bool( + commits + or acceptance + or highlights + or pr_number is not None + or dev_summary + or completed_deps + or pr_review is not None + ) + + def build_task_handoff( task: Any, journal_highlights: list[dict[str, Any]] ) -> dict[str, Any] | None: @@ -111,16 +132,15 @@ def build_task_handoff( # the concrete issues in every PM briefing so a respawned PM doesn't # re-submit the same PR blind. pr_review = _extract_pr_review(getattr(task, "notes_structured", None)) - has_prior = bool( - commits - or acceptance - or highlights - or pr_number is not None - or dev_summary - or completed_deps - or pr_review is not None - ) - if not has_prior: + if not _has_prior_work( + commits, + acceptance, + highlights, + pr_number, + dev_summary, + completed_deps, + pr_review, + ): return None handoff: dict[str, Any] = { "pr_number": pr_number, diff --git a/roboco/services/prompter.py b/roboco/services/prompter.py index d234b3e9..96d587a4 100644 --- a/roboco/services/prompter.py +++ b/roboco/services/prompter.py @@ -193,6 +193,55 @@ class PrompterService: return Team.BOARD return Team.MAIN_PM + def _validate_and_coerce_draft(self, draft_data: dict[str, Any]) -> None: + """Validate title + acceptance criteria, then flatten the list-shaped + fields (acceptance_criteria / what_this_builds / notes / each the_work + unit's items) to ``list[str]`` in place. + + Raises ``ValidationError`` (clean 400) for a missing title or empty / + missing acceptance criteria — a malformed draft (e.g. an incomplete + ``propose_batch`` item) would otherwise hit a bare ``KeyError`` and + surface as an opaque 500. Coercion runs here too because a draft can + arrive via re-draft / localStorage, not only the intake choke point. + """ + if not draft_data.get("title"): + raise ValidationError( + message="This task draft is missing a title.", field="title" + ) + if not draft_data.get("acceptance_criteria"): + raise ValidationError( + message="This task draft is missing acceptance criteria.", + field="acceptance_criteria", + ) + draft_data["acceptance_criteria"] = coerce_str_list( + draft_data.get("acceptance_criteria") + ) + draft_data["what_this_builds"] = coerce_str_list( + draft_data.get("what_this_builds") + ) + draft_data["notes"] = coerce_str_list(draft_data.get("notes")) + for unit in draft_data.get("the_work") or []: + if isinstance(unit, dict): + unit["items"] = coerce_str_list(unit.get("items")) + if not draft_data["acceptance_criteria"]: + raise ValidationError( + message="This task draft is missing acceptance criteria.", + field="acceptance_criteria", + ) + + def _resolve_draft_assignee( + self, assigned_to: UUID | None, draft_data: dict[str, Any] + ) -> UUID | None: + """Explicit confirm-button assignment wins; else fall back to any assignee + carried on the draft.""" + if assigned_to is not None: + return assigned_to + if not draft_data.get("assigned_to"): + return None + with contextlib.suppress(ValueError): + return UUID(str(draft_data["assigned_to"])) + return None + async def create_task_from_draft( self, draft_data: dict[str, Any], @@ -223,40 +272,7 @@ class PrompterService: draft through to the task so the analyzer's surface is persisted. """ place = placement or BatchPlacement() - # A draft must carry a title + acceptance criteria — without this a - # malformed draft (e.g. an agent's incomplete propose_batch item) hits a - # bare KeyError below and surfaces as an opaque 500 instead of a clean, - # actionable 400. - if not draft_data.get("title"): - raise ValidationError( - message="This task draft is missing a title.", field="title" - ) - if not draft_data.get("acceptance_criteria"): - raise ValidationError( - message="This task draft is missing acceptance criteria.", - field="acceptance_criteria", - ) - # Flatten the string-list fields to list[str]. The agent sometimes emits - # these as XML-ish elements the SDK parses into dict - # wrappers ({"item": {"$text": "…"}}); left as-is they crash the - # VARCHAR[] insert (asyncpg: "expected str, got dict") and dump str(dict) - # into the rendered description. Coerce here too — a draft can arrive - # via redraft/localStorage, not only the intake choke point. - draft_data["acceptance_criteria"] = coerce_str_list( - draft_data.get("acceptance_criteria") - ) - draft_data["what_this_builds"] = coerce_str_list( - draft_data.get("what_this_builds") - ) - draft_data["notes"] = coerce_str_list(draft_data.get("notes")) - for unit in draft_data.get("the_work") or []: - if isinstance(unit, dict): - unit["items"] = coerce_str_list(unit.get("items")) - if not draft_data["acceptance_criteria"]: - raise ValidationError( - message="This task draft is missing acceptance criteria.", - field="acceptance_criteria", - ) + self._validate_and_coerce_draft(draft_data) # Recompose the description from the (possibly edited) structured fields — # the task always carries a freshly-composed, consistent description. draft_data["description"] = compose_description(draft_data) @@ -291,10 +307,7 @@ class PrompterService: # Explicit assignment (from the confirm button) wins; else fall back to # any assignee carried on the draft. Resolved before team routing — the # owner decides the team for a product. - resolved_assigned_to: UUID | None = assigned_to - if resolved_assigned_to is None and draft_data.get("assigned_to"): - with contextlib.suppress(ValueError): - resolved_assigned_to = UUID(str(draft_data["assigned_to"])) + resolved_assigned_to = self._resolve_draft_assignee(assigned_to, draft_data) team = await self._resolve_owning_team( draft_data, diff --git a/roboco/services/sequencing.py b/roboco/services/sequencing.py index 2cc90acf..2fa0aee9 100644 --- a/roboco/services/sequencing.py +++ b/roboco/services/sequencing.py @@ -215,6 +215,21 @@ class SequencingService: _MIN_COLLISION_PAIR = 2 +def _surfaced_siblings(siblings: list) -> list: + """Siblings carrying a collision surface: a project to collide within and at + least one of intends_to_touch / adds_migration / touches_shared.""" + return [ + s + for s in siblings + if getattr(s, "project_id", None) + and ( + getattr(s, "intends_to_touch", None) + or getattr(s, "adds_migration", False) + or getattr(s, "touches_shared", False) + ) + ] + + def dev_task_collision_edges(siblings: list) -> list[tuple[object, object]]: """Wire the dev-task collision DAG for a parent's surfaced siblings. @@ -237,16 +252,7 @@ def dev_task_collision_edges(siblings: list) -> list[tuple[object, object]]: reverse edge (which would cycle). ``add_dependency`` dedupes, so repeated wiring is a no-op on already-wired pairs. """ - surfaced = [ - s - for s in siblings - if getattr(s, "project_id", None) - and ( - getattr(s, "intends_to_touch", None) - or getattr(s, "adds_migration", False) - or getattr(s, "touches_shared", False) - ) - ] + surfaced = _surfaced_siblings(siblings) if len(surfaced) < _MIN_COLLISION_PAIR: return [] # Stable order across incremental re-runs: priority is set at creation, diff --git a/roboco/services/task.py b/roboco/services/task.py index d61bbcab..daaa4d49 100644 --- a/roboco/services/task.py +++ b/roboco/services/task.py @@ -6481,6 +6481,38 @@ class TaskService(BaseService): def _task_status_value(task: TaskTable) -> str: return task.status.value if hasattr(task.status, "value") else str(task.status) + @staticmethod + def _raise_if_self_review(agent: AgentContext, task: TaskTable) -> None: + """Refuse a QA/Documenter claim of a task it itself developed.""" + if agent.role in (AgentRole.QA, AgentRole.DOCUMENTER): + original_dev = extract_original_developer(task) + if original_dev and str(agent.agent_id) == original_dev: + raise UnauthorizedError( + action="claim", + reason=( + "SELF_REVIEW: Cannot claim a task that you developed. " + f"Leave it for another {agent.role.value}." + ), + ) + + @staticmethod + def _raise_if_main_pm_code_claim( + claimant_is_main_pm: bool, task: TaskTable + ) -> None: + """Refuse a Main-PM claim of a code task it would have to execute.""" + if ( + claimant_is_main_pm + and _task_type_is_code(task.task_type) + and task.status != TaskStatus.NEEDS_REVISION + ): + raise UnauthorizedError( + action="claim", + reason=( + "MAIN_PM_NO_CODE: A Main PM coordinates — it does not own a" + " code task. Leave it for a developer; delegate instead." + ), + ) + async def claim_task_for_agent( self, task_id: UUID, @@ -6497,16 +6529,7 @@ class TaskService(BaseService): ) # QA / Documenter cannot claim what they themselves developed. - if agent.role in (AgentRole.QA, AgentRole.DOCUMENTER): - original_dev = extract_original_developer(task) - if original_dev and str(agent.agent_id) == original_dev: - raise UnauthorizedError( - action="claim", - reason=( - "SELF_REVIEW: Cannot claim a task that you developed. " - f"Leave it for another {agent.role.value}." - ), - ) + self._raise_if_self_review(agent, task) can_assign = permissions.can_perform_task_action( agent, TaskAction.ASSIGN, task.team @@ -6549,18 +6572,7 @@ class TaskService(BaseService): claimant_is_main_pm = await self._is_main_pm_agent(claim_agent_id) else: claimant_is_main_pm = agent.role == AgentRole.MAIN_PM - if ( - claimant_is_main_pm - and _task_type_is_code(task.task_type) - and task.status != TaskStatus.NEEDS_REVISION - ): - raise UnauthorizedError( - action="claim", - reason=( - "MAIN_PM_NO_CODE: A Main PM coordinates — it does not own a" - " code task. Leave it for a developer; delegate instead." - ), - ) + self._raise_if_main_pm_code_claim(claimant_is_main_pm, task) claimed = await self.claim( task_id, claim_agent_id, allow_reassign=allow_reassign diff --git a/tests/unit/services/test_conventions_cache_put.py b/tests/unit/services/test_conventions_cache_put.py index 8ae1dbb1..f46a55a5 100644 --- a/tests/unit/services/test_conventions_cache_put.py +++ b/tests/unit/services/test_conventions_cache_put.py @@ -41,7 +41,7 @@ class _FakeNested: async def __aenter__(self) -> None: self._session.savepoint_started += 1 - async def __aexit__(self, exc_type: Any, exc: Any, tb: Any) -> None: + async def __aexit__(self, exc_type: Any, exc: Any, _tb: Any) -> None: if self._session._duplicate: self._session.savepoint_rolled_back = True raise IntegrityError(