diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 9e52835b9..5adb0baa3 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -56,7 +56,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de │ ├── evolution_lifecycle.py ← Evolution campaign state + transaction lifecycle: campaign file IO, start/pause, begin/update transaction, cycle-outcome recording, deterministic worktree cleanup, owner cycle reports, idle dispatch over queue-owned state, supervisor auto-restart request │ ├── events.py ← Worker→supervisor event dispatcher with exact attempt/execution/round/call correlation for the process-local active main-LLM row; composes the frozen subagent task text (`_compose_subagent_text`; the acting `[WRITE SURFACE]` block states only the write-root authority boundary — actor identity comes from the immutable configured snapshot, startup/wake facts from bootstrap); an event type absent from `EVENT_HANDLERS` is DROPPED into a truncated `unknown_worker_event` row, and `tests/test_worker_event_registry.py` pins the registry by AST scan — an unexplained allowlist entry is exactly the silent blessing the scan exists to end; the scan is shape-bounded, and outside its reach the discipline is code review; holds shrink-only byte debt above the module byte ceiling │ ├── subagent_task_truth.py ← Delegation-truth enrichment of the subagent `task_done` transport frame (`enrich_task_done_event`) - │ ├── event_taxonomy.py, events_budget.py, events_chat_delivery.py, events_coop_checkpoint.py, events_evolution_done.py, events_project_routing.py, events_runtime_controls.py, events_schedule_task.py, events_subagent_admission.py, events_task_done.py, events_worker_reports.py ← The handler leaves `events.py` merges into `EVENT_HANDLERS`, one owner per family: the declared disposition of every event kind the runtime puts on `EVENT_Q` (`event_taxonomy.py`, the registry the AST scan reads); usage-accounting and budget-pause reports; owner-facing chat delivery (text, media, typing); cooperative repository checkpoints when a task tree goes quiescent; terminal handling of an evolution task and its campaign; where a chat turn becomes a task and where a project scope is bound; posture changes that are not one task's state; the `schedule_task` admission and duplicate gates and their refusals; admission facts for a requested subagent (census, caps, constraint); resolution of a terminal event into durable truth and delivery; and what a running worker reports about itself + │ ├── event_taxonomy.py, events_budget.py, events_chat_delivery.py, events_coop_checkpoint.py, events_evolution_done.py, events_project_routing.py, events_runtime_controls.py, events_schedule_task.py, events_subagent_admission.py, events_task_done.py, events_worker_reports.py ← The handler leaves `events.py` merges into `EVENT_HANDLERS`, one owner per family: the declared disposition of every event kind the runtime puts on `EVENT_Q` (`event_taxonomy.py`, the registry the AST scan reads); usage-accounting and budget-pause reports; owner-facing chat delivery (text, media, typing); cooperative repository checkpoints when a task tree goes quiescent; terminal handling of an evolution task and its campaign; where a chat turn becomes a task and where a project scope is bound; posture changes that are not one task's state; the `schedule_subagent` admission gates (chat target, depth, worker pool, active-child cap) and their refusals, with deliberately no semantic duplicate gate — exact task-id fencing, caps and cost ceilings are the floor and the parent decides what to spawn; admission facts for a requested subagent (census, caps, constraint); resolution of a terminal event into durable truth and delivery; and what a running worker reports about itself │ ├── task_dispatch.py ← Pure admitted-event→worker-payload construction, including the identical top-level/metadata depth and configured-route projections consumed by workers │ ├── log_addressing.py ← Explicit audience for task-scoped live log events: `address_task_event` (lineage from the RUNNING row, project binding wins, explicit chat_id preserved — 0 is `HIDDEN_CHAT_ID`, the hidden partition of the Skill Review panel and headless runs without a registered project; A2A frames are suppressed at the `push_log` choke, not by dishonest addressing), `make_server_log_sink`, `address_handler_push` │ ├── steering.py ← Owner steering-message delivery to running tasks: mailbox routing to the drive the worker drains, plus a typed refusal while a cancel intent is pending (steering is fenced during a stop — what makes the owner-stop single-turn rail safe) diff --git a/docs/CHECKLISTS_ARCHIVE.md b/docs/CHECKLISTS_ARCHIVE.md index 32fb08569..316e15930 100644 --- a/docs/CHECKLISTS_ARCHIVE.md +++ b/docs/CHECKLISTS_ARCHIVE.md @@ -23,3 +23,5 @@ superseding entries, not edits to old ones. **Standing note for item 21 (2026-08-15), do not re-raise as a missing successor class:** the delegated-coding target class lost in the D10 migration — editing one exact non-Git installed skill payload — is restored through `delegate_start(subagent_id=..., prompt=..., root="skill_payload", bucket=..., skill_name=...)` (private standalone snapshot, explicit parent CAS apply, review goes stale), with a golden registry-level test pinning it. **SUPERSEDING entry (campaign owner decision Q10=A, v7.0 ABI window; batch #9 item 6=A), read this instead of the v6.80.0 floor clause above:** By the campaign owner decision Q10=A (v7.0) the `OUROBOROS_SCOPE_REVIEW_FLOOR` gateway surface (key, endpoint, contract field, route, merge-skip, web client, self-lowering guards) IS removed in 7.0 — this supersedes the v6.80.0 standing disclosure above; the key is retired via `RETIRED_SETTING_KEYS` (a stored value is stripped on load). This is an INTENTIONAL, disclosed capability removal, not an accident to re-raise. Everything else the v6.80.0 entry records is unchanged and still binding: scope-review applicability is decided solely by the owner-only `OUROBOROS_CONTEXT_MODE`, `low` and the owner-declared RETRIEVING scope reviewer remain the two admitted replacement paths, and the accepted costs stand. One clause of that entry narrows with the removal: the INVERTED-polarity read-carve it describes for "the key or the owner endpoint" survives family-wide as `_owner_control_mention_blocks` (a provably read-only inspection reaches any owner-control key or endpoint; an interpreter or HTTP client naming one is refused whatever verb spelling it carries), while the floor-SPECIFIC detector was retired together with its setting. The `grep OUROBOROS_SCOPE_REVIEW_FLOOR data/settings.json` example in that clause therefore names a setting that no longer exists; the carve it illustrates does. + +**Standing disclosure (owner decision 2026-09-14, #884), do not re-raise as an undisclosed removal:** the semantic duplicate-task gate on `schedule_subagent` admission (`supervisor/events_schedule_task.py::_find_duplicate_task`, a light-model judge that rejected a child as `rejected_duplicate` for resembling an active task) is REMOVED. Admission keeps exact task-id fencing, the per-root active-child cap, the hard cap, depth caps and cost ceilings; the parent decides what to spawn and receives `task_id`, `objective` and `subagent_id` for every child. Byte-identical siblings (multi-model cross-checks, majority votes) are admitted by design. The `rejected_duplicate` status constant and its projections remain for old task records only. Residual risk, disclosed: a parent re-queued after a mid-wave crash may pay a wave twice, bounded by the cap and the per-tree ceiling. diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index 8e9fd32a0..9c6a5bdcd 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -195,6 +195,7 @@ BAND_PATHS = { "tests/test_loop_transport_wait.py": "Contract suite for the transport-wait episode: classification, custody, round-level wait, terminals, and the final-review regression pins live together as one coherent surface.", "tests/test_managed_review_subject.py": "Lane L-review contract suite: the managed resolution-delta subject (gate/advisory surfaces, M0 fallback), Q25-A admission, Q28-A yield outcomes and enforcement-honest advisory texts grew past 1000 across the adversarial fix round; one subject, one suite.", "tests/test_native_tool_round_executor.py": "One suite per executor contract: the native tool-round episode's bounds, floors, custody facts and delivery shapes are one behaviour pinned together; split at the next natural seam (custody vs bounds), not by size.", + "tests/test_nested_rights_depth.py": "Nested-rights depth suite shrank into the band when the semantic duplicate-gate stubs left (#884); one contract, one file; shrink next touch.", "tests/test_observability_outcomes_v2.py": None, "tests/test_onboarding_wizard.py": None, "tests/test_owner_stop_s3.py": "Entered the band from 821 lines: the S3 contract suite now covers retry-root aliasing, graceful-to-immediate hardening, stale-control drain races, hard deadline preservation, descendant settlement failure, and late resweep exactly-once root finalization.", diff --git a/supervisor/events.py b/supervisor/events.py index 823e3d092..206f6f629 100644 --- a/supervisor/events.py +++ b/supervisor/events.py @@ -205,13 +205,8 @@ from supervisor.events_runtime_controls import ( # noqa: E402, F401 -- intentio ) from supervisor.events_schedule_task import ( # noqa: E402, F401 -- intentional public re-exports VALID_SUBAGENT_MEMORY_MODES, - _PARENT_CONTEXT_END, - _PARENT_CONTEXT_MARKER, _cleanup_rejected_worktree, - _extract_task_description_and_context, - _find_duplicate_task, _handle_schedule_task, - _format_task_for_dedup, _reject_schedule_task, ) from supervisor.events_subagent_admission import ( # noqa: E402, F401 -- intentional public re-exports diff --git a/supervisor/events_schedule_task.py b/supervisor/events_schedule_task.py index 2a7cd2e53..3606c3cc5 100644 --- a/supervisor/events_schedule_task.py +++ b/supervisor/events_schedule_task.py @@ -1,14 +1,14 @@ -"""The schedule_task admission gates, the duplicate gate, and its refusals. +"""The schedule_task admission gates and their refusals. One owner for the facts the dispatch parent's schedule handler needs: the -chat-target gate, the semantic duplicate gate, the composed queue payload, -and every refusal path including worktree cleanup for a rejected subagent. +chat-target gate, the depth and worker-pool gates, the active-child cap, the +composed queue payload, and every refusal path including worktree cleanup for +a rejected subagent. """ from __future__ import annotations import logging -import os from typing import Any, Dict, Optional from ouroboros.tool_capabilities import ACTING_SUBAGENT_MODE from ouroboros.task_results import ( @@ -24,7 +24,6 @@ from ouroboros.config import get_max_active_subagents_per_root from ouroboros.contracts.task_contract import build_task_contract from ouroboros.contracts.task_contract import normalize_allowed_resources from ouroboros.subagents import intended_lane as intended_subagent_lane -from ouroboros.task_results import STATUS_REJECTED_DUPLICATE from ouroboros.task_results import STATUS_SCHEDULED from ouroboros.tools.control_delegation import admitted_depth_cap from ouroboros.tools.control_delegation import check_delegation_admission @@ -58,257 +57,9 @@ def _events(): return events -_PARENT_CONTEXT_MARKER = "[BEGIN_PARENT_CONTEXT" - - -_PARENT_CONTEXT_END = "[END_PARENT_CONTEXT]" - - VALID_SUBAGENT_MEMORY_MODES = frozenset({"forked", "empty"}) -def _extract_task_description_and_context(task: Dict[str, Any]) -> tuple[str, str]: - description = str(task.get("description") or "").strip() - context = str(task.get("context") or "").strip() - if description or context: - return description, context - - text = str(task.get("text") or task.get("description") or "").strip() - if not text: - return "", "" - if _PARENT_CONTEXT_MARKER not in text or _PARENT_CONTEXT_END not in text: - return text, "" - - before_marker, after_marker = text.split(_PARENT_CONTEXT_MARKER, 1) - description = before_marker.split("\n\n---\n", 1)[0].strip() - if "]\n" in after_marker: - after_marker = after_marker.split("]\n", 1)[1] - context = after_marker.rsplit(_PARENT_CONTEXT_END, 1)[0].strip() - return description, context - - -def _format_task_for_dedup( - task_id: str, - description: str, - context: str, - *, - expected_output: str = "", - constraints: str = "", - role: str = "", -) -> str: - sections = [ - f"Task ID: {task_id}\n" - f"Description:\n{description or '(empty)'}\n\n" - f"Context:\n{context or '(none)'}" - ] - if expected_output: - sections.append(f"Expected output:\n{expected_output}") - if constraints: - sections.append(f"Constraints:\n{constraints}") - if role: - sections.append(f"Role:\n{role}") - return "\n\n".join(sections) - - -def _find_duplicate_task( - desc: str, - task_context: str, - pending: list, - running: dict, - *, - expected_output: str = "", - constraints: str = "", - role: str = "", - dedupe_identity: Optional[Dict[str, str]] = None, -) -> Optional[str]: - """Use a scoped light-model attempt to reject only true duplicate active tasks. - - Provider/parse failures remain fail-soft, but monetary-accounting rails propagate - so an unavailable budget can never be mistaken for a semantic non-duplicate. - """ - identity = dedupe_identity if isinstance(dedupe_identity, dict) else {} - - def _task_identifier(existing_task: Dict[str, Any]) -> str: - return str(existing_task.get("id") or existing_task.get("task_id") or "").strip() - - def _is_subagent_ancestor_task(existing_task: Dict[str, Any]) -> bool: - delegation_role = str(identity.get("delegation_role") or "") - if delegation_role != "subagent": - return False - existing_id = _task_identifier(existing_task) - parent = str(identity.get("parent_task_id") or "").strip() - root = str(identity.get("root_task_id") or "").strip() - if existing_id and existing_id in {parent, root}: - return True - existing_role = str(existing_task.get("delegation_role") or "") - existing_root = str(existing_task.get("root_task_id") or "").strip() - return bool(existing_role == "root" and root and existing_root == root) - - def _is_distinct_parallel_subagent(existing_task: Dict[str, Any]) -> bool: - # Lineage/role are scheduler identity facts for parallel swarm slots; - # semantic duplicate judgment still belongs to the LLM for remaining cases. - delegation_role = str(identity.get("delegation_role") or "") - if str(delegation_role or "") != "subagent": - return False - if str(existing_task.get("delegation_role") or "") != "subagent": - return False - root = str(identity.get("root_task_id") or "") - if not root or str(existing_task.get("root_task_id") or "") != root: - return False - parent = str(identity.get("parent_task_id") or "") - existing_parent = str(existing_task.get("parent_task_id") or "") - if parent != existing_parent: - return True - new_role = str(role or "").strip() - existing_role = str(existing_task.get("role") or "").strip() - return bool(new_role and existing_role and new_role != existing_role) - - existing = [] - for task in pending: - description, context = _extract_task_description_and_context(task) - if ( - description.strip() - and not _is_subagent_ancestor_task(task) - and not _is_distinct_parallel_subagent(task) - ): - existing.append({ - "id": str(task.get("id", "?")), - "description": description, - "context": context, - "expected_output": str(task.get("expected_output") or ""), - "constraints": str(task.get("constraints") or ""), - "role": str(task.get("role") or ""), - "delegation_role": str(task.get("delegation_role") or ""), - "parent_task_id": str(task.get("parent_task_id") or ""), - "root_task_id": str(task.get("root_task_id") or ""), - }) - for task_id, meta in running.items(): - task_data = meta.get("task") if isinstance(meta, dict) else None - if not isinstance(task_data, dict): - continue - description, context = _extract_task_description_and_context(task_data) - if ( - description.strip() - and not _is_subagent_ancestor_task({"id": task_id, **task_data}) - and not _is_distinct_parallel_subagent(task_data) - ): - existing.append({ - "id": str(task_id), - "description": description, - "context": context, - "expected_output": str(task_data.get("expected_output") or ""), - "constraints": str(task_data.get("constraints") or ""), - "role": str(task_data.get("role") or ""), - "delegation_role": str(task_data.get("delegation_role") or ""), - "parent_task_id": str(task_data.get("parent_task_id") or ""), - "root_task_id": str(task_data.get("root_task_id") or ""), - }) - - if not existing: - return None - - existing_lines = "\n\n".join( - _format_task_for_dedup( - e["id"], - e["description"], - e["context"], - expected_output=e.get("expected_output", ""), - constraints=e.get("constraints", ""), - role=e.get("role", ""), - ) - for e in existing - ) - prompt = ( - "Determine whether the NEW task is a true duplicate of any EXISTING active task.\n" - "Only return a task ID if the requested work is materially the same.\n" - "Tasks that share a broad goal but differ in target model, creative focus, " - "scope, parent context, or intended output are NOT duplicates.\n\n" - "NEW TASK\n" - f"{_format_task_for_dedup('NEW', desc, task_context, expected_output=expected_output, constraints=constraints, role=role)}\n\n" - f"EXISTING ACTIVE TASKS\n{existing_lines}\n\n" - "Reply ONLY with the task ID if duplicate, or NONE if not." - ) - - from dataclasses import replace - - from ouroboros.usage_accounting import ( - BudgetExceeded, - UsageAccountingError, - UsageScope, - current_usage_scope, - usage_scope, - ) - - base_scope = current_usage_scope() - prospective_task_id = str(identity.get("task_id") or (base_scope.task_id if base_scope else "")) - prospective_root_id = str( - identity.get("root_task_id") - or (base_scope.root_task_id if base_scope else "") - or prospective_task_id - ) - prospective_parent_id = str( - identity.get("parent_task_id") - or (base_scope.parent_task_id if base_scope else "") - ) - prospective_budget_root: Any = identity.get("budget_drive_root") or ( - base_scope.drive_root if base_scope else None - ) - if base_scope is not None: - duplicate_scope = replace( - base_scope, - drive_root=prospective_budget_root, - task_id=prospective_task_id, - root_task_id=prospective_root_id, - parent_task_id=prospective_parent_id, - category="planning", - source="task_duplicate_check", - ) - else: - from ouroboros.settings_setup_contract import resolve_total_budget_usd - global_limit = resolve_total_budget_usd() - try: - root_limit = float(os.environ.get("OUROBOROS_PER_TASK_COST_USD", "0") or 0) - except (TypeError, ValueError): - root_limit = 0.0 - duplicate_scope = UsageScope( - drive_root=prospective_budget_root, - task_id=prospective_task_id, - root_task_id=prospective_root_id, - parent_task_id=prospective_parent_id, - category="planning", - source="task_duplicate_check", - global_limit_usd=global_limit, - root_limit_usd=root_limit if root_limit > 0 else None, - ) - - try: - from ouroboros.config import get_light_model - from ouroboros.llm import LLMClient - light_model = get_light_model() - client = LLMClient() - with usage_scope(duplicate_scope): - resp_msg, _usage = client.chat( - messages=[{"role": "user", "content": prompt}], - model=light_model, - model_role="light", - reasoning_effort="low", - max_tokens=50, - ) - answer = (resp_msg.get("content") or "NONE").strip() - if answer.upper() == "NONE" or not answer: - return None - answer_lower = answer.lower() - for e in existing: - if e["id"].lower() in answer_lower: - return e["id"] - return None - except (BudgetExceeded, UsageAccountingError): - raise - except Exception as exc: - log.warning("LLM dedup unavailable, accepting task: %s", exc) - return None - - def _cleanup_rejected_worktree(tid: str, result_fields: Dict[str, Any]) -> None: """Tear down a write surface provisioned for an acting subagent that is then rejected by a later gate, so rejected schedules never leak a worktree or an @@ -688,7 +439,6 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None: return if desc: - # Bible P5: duplicate judgment stays LLM-first, not hardcoded. from supervisor.queue import PENDING as QUEUE_PENDING, RUNNING as QUEUE_RUNNING pending_ref = getattr(ctx, "PENDING", QUEUE_PENDING) running_ref = getattr(ctx, "RUNNING", QUEUE_RUNNING) @@ -718,35 +468,6 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None: rationale=f"Queued behind active subagent cap {max_active}; wait for a slot before additional fan-out.", advisory=True, ) - dup_id = _find_duplicate_task( - desc, - task_context, - pending_ref, - running_ref, - expected_output=expected_output, - constraints=constraints, - role=role, - dedupe_identity={ - "delegation_role": delegation_role, - "task_id": tid, - "parent_task_id": str(parent_id or ""), - "root_task_id": root_task_id, - "budget_drive_root": budget_drive_root or str(ctx.DRIVE_ROOT), - }, - ) - if dup_id: - log.info("Rejected duplicate task: new='%s' duplicates='%s'", desc[:100], dup_id) - detail = f"Task was rejected as semantically similar to already active task {dup_id}." - _reject_schedule_task( - ctx, tid=tid, chat_id=chat_id, delegation_role=delegation_role, - parent_id=parent_id, root_task_id=root_task_id, role=role, - result_fields=result_fields, - detail=detail, - status=STATUS_REJECTED_DUPLICATE, - extra_fields={"duplicate_of": dup_id}, - fallback_message=f"⚠️ Task rejected: semantically similar to already active task {dup_id}", - ) - return # Assignment, not admission, proves achieved depth. admitted_task_contract, admitted_depth_provenance = stamp_depth_provenance( diff --git a/tests/test_budget_tracking.py b/tests/test_budget_tracking.py index 8f02abbd1..d1a8d3a9b 100644 --- a/tests/test_budget_tracking.py +++ b/tests/test_budget_tracking.py @@ -286,111 +286,6 @@ class TestUpdatePatternsCostTracking: assert usage_arg.get("prompt_tokens") == 300 -class TestSupervisorDedupCostTracking: - """Supervisor duplicate checks bind the physical-attempt ledger exactly once.""" - - def test_dedup_check_binds_prospective_scope_without_legacy_increment(self, tmp_path): - import supervisor.events as ev_mod - from ouroboros.usage_accounting import current_usage_scope - - usage = {"prompt_tokens": 50, "completion_tokens": 10, "cost": 0.0001} - # Need at least one existing task so the early-return guard doesn't skip the LLM call. - pending = [{"id": "existing-1", "type": "task", "text": "some other task"}] - captured = [] - - with patch("ouroboros.llm.LLMClient") as mock_cls, \ - patch("supervisor.state.update_budget_from_usage") as mock_budget: - inst = MagicMock() - inst.chat.side_effect = lambda **_kwargs: ( - captured.append(current_usage_scope()) or {"content": "NONE"}, - usage, - ) - mock_cls.return_value = inst - result = ev_mod._find_duplicate_task( - "Deploy new feature", - "", - pending, - {}, - dedupe_identity={ - "task_id": "prospective", - "root_task_id": "root-1", - "parent_task_id": "parent-1", - "budget_drive_root": str(tmp_path), - }, - ) - mock_budget.assert_not_called() - assert result is None # "NONE" response = no duplicate found - scope = captured[0] - assert scope.drive_root == str(tmp_path) - assert scope.task_id == "prospective" - assert scope.root_task_id == "root-1" - assert scope.parent_task_id == "parent-1" - assert scope.category == "planning" - assert scope.source == "task_duplicate_check" - - @pytest.mark.parametrize("raw, expected_default", [(None, True), ("0", False)]) - def test_dedup_scope_uses_total_budget_resolver(self, tmp_path, monkeypatch, raw, expected_default): - import supervisor.events as ev_mod - from ouroboros.config import SETTINGS_DEFAULTS - from ouroboros.usage_accounting import current_usage_scope - - if raw is None: - monkeypatch.delenv("TOTAL_BUDGET", raising=False) - else: - monkeypatch.setenv("TOTAL_BUDGET", raw) - captured = [] - with patch("ouroboros.llm.LLMClient") as mock_cls: - mock_cls.return_value.chat.side_effect = lambda **_k: ( - captured.append(current_usage_scope()) or {"content": "NONE"}, {} - ) - ev_mod._find_duplicate_task( - "new", "", [{"id": "old", "text": "old"}], {}, - dedupe_identity={"task_id": "new", "budget_drive_root": str(tmp_path)}, - ) - - expected = float(SETTINGS_DEFAULTS["TOTAL_BUDGET"]) if expected_default else None - assert captured[0].global_limit_usd == expected - - def test_no_budget_call_when_no_usage(self): - import supervisor.events as ev_mod - - pending = [{"id": "existing-1", "type": "task", "text": "some task"}] - - with patch("ouroboros.llm.LLMClient") as mock_cls, \ - patch("supervisor.state.update_budget_from_usage") as mock_budget: - inst = MagicMock() - inst.chat.return_value = ({"content": "NONE"}, None) - mock_cls.return_value = inst - ev_mod._find_duplicate_task("test", "", pending, {}) - mock_budget.assert_not_called() - - def test_no_budget_call_when_no_existing_tasks(self): - """Empty pending+running — LLM not called at all, no budget update.""" - import supervisor.events as ev_mod - - with patch("ouroboros.llm.LLMClient") as mock_cls, \ - patch("supervisor.state.update_budget_from_usage") as mock_budget: - result = ev_mod._find_duplicate_task("test", "", [], {}) - mock_cls.assert_not_called() - mock_budget.assert_not_called() - assert result is None - - @pytest.mark.parametrize("error_type", [ - pytest.param("budget", id="budget_exceeded"), - pytest.param("accounting", id="accounting_error"), - ]) - def test_accounting_rails_are_not_downgraded_to_accept(self, error_type): - import supervisor.events as ev_mod - from ouroboros.usage_accounting import BudgetExceeded, UsageAccountingError - - error = BudgetExceeded("rail") if error_type == "budget" else UsageAccountingError("ledger") - pending = [{"id": "existing-1", "type": "task", "text": "some task"}] - with patch("ouroboros.llm.LLMClient") as mock_cls: - mock_cls.return_value.chat.side_effect = error - with pytest.raises(type(error), match=str(error)): - ev_mod._find_duplicate_task("new task", "", pending, {}) - - class TestAdvisoryCostAccounting: """Advisory spend is accounted inside the review substrate: the native episode executor stamps every paid send into the usage ledger under diff --git a/tests/test_events_extraction.py b/tests/test_events_extraction.py index 3cbd0dd59..bad0838bf 100644 --- a/tests/test_events_extraction.py +++ b/tests/test_events_extraction.py @@ -72,12 +72,7 @@ _MOVED_OWNERS = { "_task_own_id": events_subagent_admission, "_validate_external_workspace": events_subagent_admission, "VALID_SUBAGENT_MEMORY_MODES": events_schedule_task, - "_PARENT_CONTEXT_END": events_schedule_task, - "_PARENT_CONTEXT_MARKER": events_schedule_task, "_cleanup_rejected_worktree": events_schedule_task, - "_extract_task_description_and_context": events_schedule_task, - "_find_duplicate_task": events_schedule_task, - "_format_task_for_dedup": events_schedule_task, "_reject_schedule_task": events_schedule_task, "_emit_routing_receipt": events_project_routing, "_handle_ensure_project_scope": events_project_routing, diff --git a/tests/test_model_slot_role_model.py b/tests/test_model_slot_role_model.py index 0fc9221d8..fae9aa308 100644 --- a/tests/test_model_slot_role_model.py +++ b/tests/test_model_slot_role_model.py @@ -511,7 +511,6 @@ def _enqueue_through_supervisor(tmp_path, monkeypatch, *, parent_lane: str = "", from types import SimpleNamespace from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.tools.control import _schedule_task from tests._shared import configure_test_subagent @@ -551,7 +550,6 @@ def _enqueue_through_supervisor(tmp_path, monkeypatch, *, parent_lane: str = "", event["depth"] = 0 event["delegation_role"] = "" - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *a, **k: None) enqueued = [] class FakeCtx: diff --git a/tests/test_nested_coordination_acceptance.py b/tests/test_nested_coordination_acceptance.py index dd15df06b..3570b9d0c 100644 --- a/tests/test_nested_coordination_acceptance.py +++ b/tests/test_nested_coordination_acceptance.py @@ -378,11 +378,6 @@ def test_depth3_control_plane_reaches_root_acceptance(tmp_path, monkeypatch): monkeypatch.setenv("OUROBOROS_MAX_SUBAGENT_DEPTH", "3") monkeypatch.setenv("OUROBOROS_MAX_ACTIVE_SUBAGENTS_PER_ROOT", "6") monkeypatch.setattr(control, "load_settings", lambda: settings) - monkeypatch.setattr( - events, - "_find_duplicate_task", - lambda *_args, **_kwargs: None, - ) monkeypatch.setattr(task_tree_ledger, "DATA_DIR", tmp_path) monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path) diff --git a/tests/test_nested_rights_depth.py b/tests/test_nested_rights_depth.py index e7628a047..4c48c63f5 100644 --- a/tests/test_nested_rights_depth.py +++ b/tests/test_nested_rights_depth.py @@ -404,9 +404,7 @@ def test_supervisor_schedule_path_preserves_admitted_cap_after_live_depth_decrea tmp_path, monkeypatch, ): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) monkeypatch.setattr(events, "get_max_subagent_depth", lambda: 1) contract = build_task_contract({ "delegation_budget": { @@ -508,9 +506,7 @@ def test_supervisor_ingress_bounds_legacy_permission_by_admitted_remaining_envel def test_supervisor_ingress_records_explicit_root_depth_request(tmp_path, monkeypatch): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) contract = build_task_contract({ "delegation_budget": {"depth_remaining": 3}, }) @@ -621,9 +617,7 @@ def _fake_ctx(tmp_path, enqueued): def test_supervisor_admission_enforces_parent_rights_and_allows_one_non_fanout_child(tmp_path, monkeypatch): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) parent_contract = build_task_contract({"delegation_budget": {"may_fan_out": False}}) write_task_result(tmp_path, "parent", STATUS_RUNNING, parent_task_id="", root_task_id="parent", delegation_role="root", task_contract=parent_contract) @@ -641,9 +635,7 @@ def test_supervisor_admission_enforces_parent_rights_and_allows_one_non_fanout_c def test_supervisor_rejects_invalid_depth_before_provisioning_or_enqueue(tmp_path, monkeypatch): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) for index, raw_depth in enumerate((-1, -0.5, "-1", "not-a-depth")): task_id = f"invalid-depth-{index}" enqueued = [] @@ -661,9 +653,7 @@ def test_supervisor_rolls_back_subagent_when_scheduled_result_write_fails( tmp_path, monkeypatch, ): from supervisor import events, task_admission - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) parent_contract = build_task_contract({"delegation_budget": {"may_fan_out": False}}) write_task_result( tmp_path, "parent", STATUS_RUNNING, @@ -724,9 +714,7 @@ def test_supervisor_receipt_rollback_removes_only_its_enqueue_identity( tmp_path, monkeypatch, ): from supervisor import events, task_admission - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) parent_contract = build_task_contract({"delegation_budget": {"may_fan_out": True}}) write_task_result( tmp_path, "parent", STATUS_RUNNING, @@ -794,7 +782,6 @@ def test_replayed_schedule_event_keeps_one_physical_task_and_transition( from supervisor import events, queue, state, workers from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) write_task_result( tmp_path, "parent", STATUS_RUNNING, root_task_id="parent", delegation_role="root", @@ -1125,9 +1112,7 @@ def test_supervisor_keeps_admission_when_scheduled_write_raises_after_commit( tmp_path, monkeypatch, ): from supervisor import events, task_admission - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) parent_contract = build_task_contract({"delegation_budget": {"may_fan_out": False}}) write_task_result( tmp_path, "parent", STATUS_RUNNING, @@ -1171,9 +1156,7 @@ def test_supervisor_rolls_back_when_monotonic_writer_returns_old_terminal( tmp_path, monkeypatch, ): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) parent_contract = build_task_contract({"delegation_budget": {"may_fan_out": True}}) write_task_result( tmp_path, "parent", STATUS_RUNNING, @@ -1247,9 +1230,7 @@ def test_schedule_exact_id_preserves_unreadable_result_and_does_not_enqueue( tmp_path, monkeypatch, ): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) result_path = tmp_path / "task_results" / "child-malformed.json" result_path.parent.mkdir() malformed = b'{"status":"scheduled"' @@ -1273,9 +1254,7 @@ def test_schedule_exact_id_preserves_unreadable_result_and_does_not_enqueue( def test_late_schedule_lookup_failure_preserves_exact_result(tmp_path, monkeypatch): from supervisor import events, task_admission - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) result_path = tmp_path / "task_results" / "late-corrupt.json" result_path.parent.mkdir() original = b"{late-corrupt" @@ -1296,9 +1275,7 @@ def test_late_schedule_lookup_failure_preserves_exact_result(tmp_path, monkeypat def test_generic_late_schedule_lookup_failure_preserves_exact_result(tmp_path, monkeypatch): from supervisor import events, queue - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) result_path = tmp_path / "task_results" / "generic-late-corrupt.json" result_path.parent.mkdir() original = b"{generic-late-corrupt" @@ -1324,9 +1301,7 @@ def test_generic_late_schedule_lookup_failure_preserves_exact_result(tmp_path, m def test_generic_schedule_replay_preserves_valid_exact_result(tmp_path, monkeypatch): from supervisor import events, queue - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) write_task_result( tmp_path, "generic-existing", "completed", root_task_id="generic-existing", delegation_role="root", result="keep me", @@ -1361,9 +1336,7 @@ def test_generic_malformed_replay_preserves_live_exact_result( tmp_path, monkeypatch, status, location, ): from supervisor import events, queue - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) tid = f"generic-live-{location}" write_task_result( tmp_path, @@ -1487,9 +1460,7 @@ def test_supervisor_rejects_count_bounded_child_when_count_scan_fails( tmp_path, monkeypatch, ): from supervisor import events - from supervisor import events_schedule_task as schedule_module - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) monkeypatch.setattr( "ouroboros.task_results.list_task_results", lambda *_args, **_kwargs: (_ for _ in ()).throw(OSError("unreadable")), diff --git a/tests/test_subscription_model_roles.py b/tests/test_subscription_model_roles.py index b7921eebb..40737535e 100644 --- a/tests/test_subscription_model_roles.py +++ b/tests/test_subscription_model_roles.py @@ -33,18 +33,6 @@ def test_equal_model_names_keep_distinct_main_and_light_accounts(): assert resolve_model_target(settings["OUROBOROS_MODEL"]).credential_ref == "" -def test_supervisor_duplicate_check_uses_the_saved_light_account(subscription_transport, monkeypatch): - from supervisor.events_schedule_task import _find_duplicate_task - - _, gateway, _ = subscription_transport - monkeypatch.setenv("OUROBOROS_MODEL_LIGHT", MODEL) - monkeypatch.setenv(MODEL_ACCOUNTS_KEY, json.dumps({"main": "main-pin", "light": "light-pin"})) - gateway.results[0]["message"] = {"role": "assistant", "content": "NONE"} - assert _find_duplicate_task("New task", "", [{"id": "existing", "description": "Another task"}], {}) is None - assert len(gateway.creates) == 1 - assert gateway.uploads[0][0]["account"] == {"mode": "pin", "profileId": "light-pin"} - - def test_fallback_options_preserve_order_and_explicit_auto(): raw = {"main": "", "fallback": ["profile-B", "", "profile-A"]} parsed, encoded = normalize_model_role_options(MODEL_ACCOUNTS_KEY, raw) diff --git a/tests/test_task_creation_facts.py b/tests/test_task_creation_facts.py index ce83d157b..7a41971f0 100644 --- a/tests/test_task_creation_facts.py +++ b/tests/test_task_creation_facts.py @@ -83,7 +83,6 @@ def test_supervisor_stamps_only_a_locally_allocated_id(tmp_path, monkeypatch, su from tests.test_nested_rights_depth import _fake_ctx, _schedule_event monkeypatch.setattr(events_schedule_task, "utc_now_iso", lambda: CREATED) - monkeypatch.setattr(events_schedule_task, "_find_duplicate_task", lambda *args, **kwargs: None) event = _schedule_event("supplied-task" if supplied_id else "", "", depth=0, drive_root=tmp_path) event["delegation_role"] = "root" enqueued = [] diff --git a/tests/test_task_status_flow.py b/tests/test_task_status_flow.py index 374449a81..79ac9ddb7 100644 --- a/tests/test_task_status_flow.py +++ b/tests/test_task_status_flow.py @@ -1929,329 +1929,100 @@ def test_wait_for_task_reports_rejected_duplicate(tmp_path): assert "duplicate_of=orig999" in output -def test_handle_schedule_task_duplicate_writes_rejected_status(tmp_path, monkeypatch): +def test_handle_schedule_task_admits_identical_siblings_without_semantic_veto(tmp_path, monkeypatch): + """No semantic duplicate judge stands between a parent and its siblings. + + Exact task-id fencing, the active-child cap and cost ceilings are the floor; + identical objectives under one parent are the parent's call, so every sibling + is admitted as ``scheduled`` and none is stamped ``duplicate_of``. + """ from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module - from ouroboros.task_results import STATUS_REJECTED_DUPLICATE + import ouroboros.llm as llm_module + from ouroboros.task_results import STATUS_SCHEDULED - captured_identity = {} + # Admission must not consult any model: a judge that merely failed open + # (provider unreachable in a keyless run) would make this test pass on a + # tree that still carries the veto, so constructing a client is the failure. + constructed = [] - def _duplicate(*args, **kwargs): - captured_identity.update(kwargs.get("dedupe_identity") or {}) - return "orig111" + def _no_client(*args, **kwargs): + constructed.append((args, kwargs)) + raise AssertionError("admission constructed an LLM client") - monkeypatch.setattr(schedule_module, "_find_duplicate_task", _duplicate) + monkeypatch.setattr(llm_module, "LLMClient", _no_client) sent = [] + enqueued = [] class FakeCtx: DRIVE_ROOT = tmp_path - PENDING = [] - RUNNING = {} WORKERS = {0: SimpleNamespace(busy_task_id=None)} + def __init__(self): + self.PENDING = [] + self.RUNNING = {} + def load_state(self): return {"owner_chat_id": 1} def send_with_budget(self, chat_id, text, **kwargs): sent.append((chat_id, text, kwargs)) - ev_module._handle_schedule_task( - { - "type": "schedule_subagent", - "task_id": "dup222", - "objective": "Do the thing", - "expected_output": "Duplicate verdict", - "context": "Model focus B", - "depth": 1, - "memory_mode": "forked", - "parent_task_id": "parent111", - "root_task_id": "root111", - "drive_root": str(tmp_path / "state" / "headless_tasks" / "dup222" / "data"), - "child_drive_root": str(tmp_path / "state" / "headless_tasks" / "dup222" / "data"), - "budget_drive_root": str(tmp_path), - }, - FakeCtx(), - ) + def enqueue_task(self, task): + enqueued.append(task) + self.PENDING.append(task) + return task - path = tmp_path / "task_results" / "dup222.json" - data = json.loads(path.read_text(encoding="utf-8")) - assert data["status"] == STATUS_REJECTED_DUPLICATE - assert data["duplicate_of"] == "orig111" - assert sent and "semantically similar" in sent[0][1] - assert sent[0][2]["is_progress"] is True - assert sent[0][2]["progress_meta"]["delegation_role"] == "subagent" - assert sent[0][2]["progress_meta"]["parent_task_id"] == "parent111" - assert sent[0][2]["progress_meta"]["status"] == STATUS_REJECTED_DUPLICATE - assert captured_identity == { - "delegation_role": "subagent", - "task_id": "dup222", - "parent_task_id": "parent111", - "root_task_id": "root111", - "budget_drive_root": str(tmp_path), - } + def persist_queue_snapshot(self, reason=""): + self.snapshot_reason = reason + ctx = FakeCtx() -def test_find_duplicate_task_includes_subagent_handoff_fields(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - captured = {} - - class FakeClient: - def chat(self, messages, **kwargs): - captured["prompt"] = messages[0]["content"] - return {"content": "NONE"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "Review shared surface", - "same context", - [ + def _dispatch(tid, configured_subagent): + ev_module._handle_schedule_task( { - "id": "pending1", - "description": "Review shared surface", - "context": "same context", - "expected_output": "Docs table", - "constraints": "docs only", - "role": "docs reviewer", - } - ], - {}, - expected_output="Security table", - constraints="security only", - role="security reviewer", - ) + "type": "schedule_subagent", + "task_id": tid, + "objective": "Do the thing", + "expected_output": "Duplicate verdict", + "context": "Model focus B", + "depth": 1, + "memory_mode": "forked", + "parent_task_id": "parent111", + "root_task_id": "root111", + "drive_root": str(tmp_path / "state" / "headless_tasks" / tid / "data"), + "child_drive_root": str(tmp_path / "state" / "headless_tasks" / tid / "data"), + "budget_drive_root": str(tmp_path), + "configured_subagent": configured_subagent, + }, + ctx, + ) - assert result is None - prompt = captured["prompt"] - assert "Expected output:\nSecurity table" in prompt - assert "Expected output:\nDocs table" in prompt - assert "Constraints:\nsecurity only" in prompt - assert "Constraints:\ndocs only" in prompt - assert "Role:\nsecurity reviewer" in prompt - assert "Role:\ndocs reviewer" in prompt + # Same objective, same lineage, same (default) role; only the selected + # configured subagent differs. + _dispatch("sib1", {"selected_subagent_id": "primary-builder"}) + _dispatch("sib2", {"selected_subagent_id": "fast-scout"}) + # ... and a pair whose configured subagent is identical too. + _dispatch("sib3", {"selected_subagent_id": "primary-builder"}) + _dispatch("sib4", {"selected_subagent_id": "primary-builder"}) - -def test_find_duplicate_task_allows_distinct_subagent_roles(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - calls = [] - - class FakeClient: - def chat(self, messages, **kwargs): - calls.append(messages[0]["content"]) - return {"content": "pending1"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "Run nested smoke slot", - "", - [ - { - "id": "pending1", - "description": "Run nested smoke slot", - "expected_output": "Smoke handoff", - "role": "l1-alpha-coordinator", - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - } - ], - {}, - expected_output="Smoke handoff", - role="l1-beta-coordinator", - dedupe_identity={ - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - }, - ) - - assert result is None - assert calls == [] - - -def test_find_duplicate_task_keeps_same_role_subagent_dedupe(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - class FakeClient: - def chat(self, messages, **kwargs): - return {"content": "pending1"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "Run nested smoke slot", - "", - [ - { - "id": "pending1", - "description": "Run nested smoke slot", - "expected_output": "Smoke handoff", - "role": "l1-alpha-coordinator", - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - } - ], - {}, - expected_output="Smoke handoff", - role="l1-alpha-coordinator", - dedupe_identity={ - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - }, - ) - - assert result == "pending1" - - -def test_find_duplicate_task_allows_distinct_subagent_parent_branches(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - calls = [] - - class FakeClient: - def chat(self, messages, **kwargs): - calls.append(messages[0]["content"]) - return {"content": "pending1"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "Run nested branch smoke slot", - "", - [ - { - "id": "pending1", - "description": "Run nested branch smoke slot", - "expected_output": "Smoke handoff", - "role": "shared-l2-role", - "delegation_role": "subagent", - "parent_task_id": "l1-alpha", - "root_task_id": "root1", - } - ], - {}, - expected_output="Smoke handoff", - role="shared-l2-role", - dedupe_identity={ - "delegation_role": "subagent", - "parent_task_id": "l1-beta", - "root_task_id": "root1", - }, - ) - - assert result is None - assert calls == [] - - -def test_find_duplicate_task_allows_subagent_against_running_root_ancestor(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - calls = [] - - class FakeClient: - def chat(self, messages, **kwargs): - calls.append(messages[0]["content"]) - return {"content": "root1"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "You are l1-alpha-coordinator; schedule L2 smoke agents", - "", - [], - { - "root1": { - "task": { - "id": "root1", - "description": "Root coordinator: schedule l1-alpha, l1-beta, l1-gamma subagents", - "delegation_role": "root", - "parent_task_id": "", - "root_task_id": "root1", - } - } - }, - expected_output="L1 handoff", - role="l1-alpha-coordinator", - dedupe_identity={ - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - }, - ) - - assert result is None - assert calls == [] - - -def test_find_duplicate_task_allows_subagent_against_pending_parent_ancestor(monkeypatch): - from supervisor import events as ev_module - import ouroboros.config as config_module - import ouroboros.llm as llm_module - - calls = [] - - class FakeClient: - def chat(self, messages, **kwargs): - calls.append(messages[0]["content"]) - return {"content": "parent1"}, {} - - monkeypatch.setattr(config_module, "get_light_model", lambda: "test-light") - monkeypatch.setattr(llm_module, "LLMClient", lambda: FakeClient()) - - result = ev_module._find_duplicate_task( - "You are l1-alpha-coordinator-l2-1; return a smoke handoff", - "", - [ - { - "id": "parent1", - "description": "You are l1-alpha-coordinator; schedule three L2 smoke subagents", - "role": "l1-alpha-coordinator", - "delegation_role": "subagent", - "parent_task_id": "root1", - "root_task_id": "root1", - } - ], - {}, - expected_output="L2 handoff", - role="l1-alpha-coordinator-l2-1", - dedupe_identity={ - "delegation_role": "subagent", - "parent_task_id": "parent1", - "root_task_id": "root1", - }, - ) - - assert result is None - assert calls == [] + assert [task["id"] for task in enqueued] == ["sib1", "sib2", "sib3", "sib4"] + for tid in ("sib1", "sib2", "sib3", "sib4"): + path = tmp_path / "task_results" / f"{tid}.json" + assert path.exists(), f"{tid} was not admitted" + data = json.loads(path.read_text(encoding="utf-8")) + assert data["status"] == STATUS_SCHEDULED + assert "duplicate_of" not in data + assert not [text for _chat, text, _kw in sent if "semantically similar" in text] + # The judge caught every exception and failed open, so a raise alone would + # not distinguish a tree that still consults a model: the record does. + assert constructed == [] def test_handle_schedule_task_accepts_unique_subagent_with_lineage_and_constraint(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_SCHEDULED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) enqueued = [] sent = [] @@ -2323,10 +2094,8 @@ def test_handle_schedule_task_accepts_unique_subagent_with_lineage_and_constrain def test_handle_schedule_task_rejects_internal_subagent_without_child_drive_contract(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_FAILED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) sent = [] class FakeCtx: @@ -2368,10 +2137,8 @@ def test_handle_schedule_task_rejects_internal_subagent_without_child_drive_cont def test_handle_schedule_task_uses_event_chat_id_without_owner(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_SCHEDULED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) enqueued = [] sent = [] @@ -2447,11 +2214,9 @@ def test_handle_schedule_task_uses_event_chat_id_without_owner(tmp_path, monkeyp def test_handle_schedule_task_depth_rejection_writes_failed_status(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.config import get_max_subagent_depth from ouroboros.task_results import STATUS_FAILED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) sent = [] class FakeCtx: @@ -2496,7 +2261,6 @@ def test_handle_schedule_task_depth_rejection_writes_failed_status(tmp_path, mon def test_configured_zero_subagent_depth_truly_disables_delegation(tmp_path, monkeypatch): """A configured depth of zero disables child delegation, not the root task.""" from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.config import get_max_subagent_depth from ouroboros.task_results import STATUS_FAILED @@ -2527,7 +2291,6 @@ def test_configured_zero_subagent_depth_truly_disables_delegation(tmp_path, monk "achieved_depth": None, } - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) enqueued = [] class FakeCtx: @@ -2607,10 +2370,8 @@ def test_settings_ui_carries_a_configured_zero_subagent_depth(): def test_handle_schedule_task_rejects_legacy_subagent_event_schema(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_FAILED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) enqueued = [] sent = [] @@ -2657,10 +2418,8 @@ def test_handle_schedule_task_rejects_legacy_subagent_event_schema(tmp_path, mon def test_handle_schedule_task_queues_when_active_subagent_cap_is_full(tmp_path, monkeypatch): from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_COMPLETED, STATUS_FAILED, STATUS_SCHEDULED, load_task_result, write_task_result - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) monkeypatch.setenv("OUROBOROS_MAX_ACTIVE_SUBAGENTS_PER_ROOT", "3") # pin cap (v6.20.0 raised default to 6) sent = [] enqueued = [] @@ -2832,10 +2591,8 @@ def test_handle_schedule_task_fails_fast_when_worker_pool_unavailable(tmp_path, schedule must NOT be left as a 'scheduled' ghost — it gets a terminal workers_unavailable result so the parent can act.""" from supervisor import events as ev_module - from supervisor import events_schedule_task as schedule_module from ouroboros.task_results import STATUS_FAILED - monkeypatch.setattr(schedule_module, "_find_duplicate_task", lambda *args, **kwargs: None) sent = [] class FakeCtx: