diff --git a/ouroboros/context.py b/ouroboros/context.py index 9aff78897..e06e127ea 100644 --- a/ouroboros/context.py +++ b/ouroboros/context.py @@ -519,23 +519,18 @@ def build_runtime_section(env: Any, task: Dict[str, Any], *, ctx: Any = None, sc "your judgment picks the target (or none -> answer inline / promote_chat_to_task). A " "message in a project room defaults to that project unless it clearly says otherwise." ) - _main_manifest = ( - _meta.get("main_routing_manifest") - if isinstance(_meta.get("main_routing_manifest"), dict) - else None - ) - if _main_manifest: - runtime_data["main_routing_manifest"] = _main_manifest - _last_result = ( - _meta.get("project_last_task_result") - if isinstance(_meta.get("project_last_task_result"), dict) - else None - ) - if _last_result: - # Host-built ground truth about the thread's most recent task result - # (id/status/workspace facts/artifact refs) — read THIS before framing a - # "continue" promotion; never reconstruct prior work from chat memory. - runtime_data["project_last_task_result"] = _last_result + # Host-built ground truth about what THIS lane may continue: Main's bounded + # manifest, a project room's own hint (its recent ROOT results and the roots + # still live in it) and the thread's most recent result (id/status/workspace + # facts/artifact refs). Read these before framing a "continue" promotion; + # never reconstruct prior work from chat memory. What promote ACCEPTS is the + # predicate in ARCHITECTURE ch. 10 - these rows are the hint, not the door. + for _routing_key in ( + "main_routing_manifest", "project_routing_manifest", "project_last_task_result", + ): + _routing_fact = _meta.get(_routing_key) + if isinstance(_routing_fact, dict) and _routing_fact: + runtime_data[_routing_key] = _routing_fact _routing_contract = ( _meta.get("routing_contract") if isinstance(_meta.get("routing_contract"), dict) diff --git a/ouroboros/gateway/task_list_scan.py b/ouroboros/gateway/task_list_scan.py index c366ad98a..fc85c9b5e 100644 --- a/ouroboros/gateway/task_list_scan.py +++ b/ouroboros/gateway/task_list_scan.py @@ -24,6 +24,13 @@ from ouroboros.utils import read_json_dict _RAW_TS_MEMO: Dict[tuple, tuple] = {} _RESULT_FACT_KEYS = ( "task_id", "id", "ts", "updated_at", "delegation_role", "parent_task_id", + # A project room offers its OWN recent roots, so the selection needs the + # project of each row - one small scalar, no extra read. Status and cancel + # facts stay out: every row the selection keeps is then loaded WHOLE and + # carries them from there, while `cancel_state` lives in the durable + # cancel-intent projection, so a memo copy would be a second source nobody + # reads. + "project_id", "root_task_id", "child_drive_root", "headless_child_drive_root", ) diff --git a/ouroboros/runtime_limits.py b/ouroboros/runtime_limits.py index e5e4d2ced..6940eb45f 100644 --- a/ouroboros/runtime_limits.py +++ b/ouroboros/runtime_limits.py @@ -201,6 +201,17 @@ def get_acceptance_fence_ack_wait_sec() -> float: return ACCEPTANCE_FENCE_ACK_WAIT_SEC +# How many of a lane's newest ROOT results one routing manifest offers as continuation +# candidates. A HINT window: what promote ACCEPTS is a predicate (same project, a root, a +# readable result, not live), so a root older than this window stays addressable in its own +# room. Structural, not a settings key. +ROUTING_MANIFEST_RESULT_ROWS = 16 + + +def get_routing_manifest_result_rows() -> int: + return ROUTING_MANIFEST_RESULT_ROWS + + def get_vision_caption_timeout_sec() -> int: return _clamped_number_setting("OUROBOROS_VISION_CAPTION_TIMEOUT_SEC", low=1, cast=int) diff --git a/ouroboros/server_routing_context.py b/ouroboros/server_routing_context.py index fc5203f56..42dacf019 100644 --- a/ouroboros/server_routing_context.py +++ b/ouroboros/server_routing_context.py @@ -157,9 +157,119 @@ def _task_result_ground_truth(row: Dict[str, Any]) -> Dict[str, Any]: "branch": str(git.get("branch") or ""), "dirty": bool(git.get("dirty")), } + origin = row.get("cancel_origin") if isinstance(row.get("cancel_origin"), dict) else {} + if origin: + # WHY this result stopped, in the three scalars a continuation decision + # needs; the full origin (actor, request id, observation) stays in the row. + out["cancel_origin"] = { + key: str(origin.get(key) or "") + for key in ("reason", "source", "requested_at") if origin.get(key) + } return out +def _is_child_result(facts: Dict[str, Any]) -> bool: + """A result that is NOT an owner root: it has a parent, or the subagent role. + + ONE predicate for both readers of that fact - the manifest window that skips + children (owner decision batch 3, answer 6b=A) and the promote door, which + refuses them by the same rule. Reads a memoized fact row or a full result row. + """ + return bool(str(facts.get("parent_task_id") or "").strip()) or str( + facts.get("delegation_role") or "") == "subagent" + + +def _recent_root_results(ctx: Any, project_id: str = "") -> tuple: + """``(rows, omissions)``: the newest ROOT results a routing turn may continue, + optionally narrowed to ONE project. + + The one producer behind both routing manifests - Main offers every lane's + roots, a project room offers its own. Only the owner's ROOT results are + offered (owner decision batch 3, answer 6b=A): a swarm wave's children are the + newest results of ANY kind, so they evicted the owner's own roots from this + window - which is how a root the same actor had just read stopped being + offerable. The facts are already memoized, so both the filter and the count + cost no extra read. The count runs over the WHOLE candidate list, not inside + the capped loop: children older than the last shown root are skipped just the + same, and counting them only until the cap reported zero while folding them + into the cap's own number. + """ + from ouroboros.gateway.task_list_scan import raw_result_facts + from ouroboros.runtime_limits import get_routing_manifest_result_rows + from ouroboros.task_results import load_task_result, task_results_dir + + results_error = "" + try: + facts, unreadable = raw_result_facts(task_results_dir(ctx.DRIVE_ROOT, create=False)) + except OSError as exc: + facts, unreadable = {}, ["result_directory_unreadable"] + results_error = f"result_directory_unreadable: {exc}" + ordered = sorted( + facts, key=lambda name: facts[name]["ts"] or facts[name]["updated_at"], reverse=True, + ) + pool = [name for name in ordered if not project_id or facts[name]["project_id"] == project_id] + children = sum( + 1 for name in pool if not facts[name]["schema_refusal"] and _is_child_result(facts[name]) + ) + cap = get_routing_manifest_result_rows() + finals: list = [] + for name in pool: + if facts[name]["schema_refusal"] or _is_child_result(facts[name]): + continue + row = load_task_result(ctx.DRIVE_ROOT, pathlib.Path(name).stem) + if row is not None: + finals.append(_task_result_ground_truth(row)) + if len(finals) == cap: + break + return finals, { + # Kept meaning: results cut by the row cap. The children skipped above are + # a DIFFERENT omission and are counted as such, never folded in here. + "final_results": None if unreadable else max(0, len(pool) - children - len(finals)), + "final_results_error": results_error, + "children": None if unreadable else children, + } + + +def _cancel_state_facts(ctx: Any, task_id: str) -> Dict[str, Any]: + """The durable cancel intent standing over one LIVE root, or ``{}``: the typed + public projection (``cancel_state``/``cancel_reason``/``stop_policy``), never a + body. A settled result carries its own ``cancel_origin`` instead.""" + if not task_id: + return {} + try: + from ouroboros.cancel_intents import cancel_state_fields + + return cancel_state_fields(ctx.DRIVE_ROOT, task_id) + except Exception: + log.debug("cancel state unreadable for %s", task_id, exc_info=True) + return {} + + +def _project_routing_manifest(ctx: Any, project_id: str) -> Dict[str, Any]: + """The room's bounded HINT for a "continue this work" decision: the project's + recent ROOT results and the roots still live in it, each with the small typed + facts that separate the two choices - a finished root is promote's predecessor, + a live one is ``steer_task``. + + A hint, never the door: promote's predicate admits an older root of the same + project too (ch. 10), so this window may be bounded without deciding what the + room can continue. Until it existed a room saw exactly ONE candidate, the + registry pointer, so a room whose pointer had moved could not name its own + interrupted root at all. + """ + finals, omissions = _recent_root_results(ctx, project_id) + active = [ + {**row, **_cancel_state_facts(ctx, str(row.get("task_id") or ""))} + for row in _addressable_root_tasks(ctx, None) + if str(row.get("project_id") or "") == project_id + ] + return { + "final_results": finals, + "active_roots": active[:40], + "omissions": {**omissions, "active_roots": max(0, len(active) - 40)}, + } + + def _latest_project_task_result(ctx: Any, project_id: str) -> Optional[Dict[str, Any]]: """Newest task result bound to ``project_id`` WITHOUT replaying the whole store (DEVELOPMENT "Projection over replay"). The registry row's durable @@ -262,9 +372,7 @@ def _latest_project_task_result(ctx: Any, project_id: str) -> Optional[Dict[str, def _main_routing_manifest(ctx: Any) -> Dict[str, Any]: """Bounded canonical facts for one Main-chat LLM routing decision.""" from ouroboros.gateway._helpers import read_rotated_jsonl_entries - from ouroboros.gateway.task_list_scan import raw_result_facts from ouroboros.projects_registry import list_projects - from ouroboros.task_results import load_task_result, task_results_dir projects = [{ "project_id": str(row.get("id") or ""), @@ -276,35 +384,7 @@ def _main_routing_manifest(ctx: Any) -> Dict[str, Any]: "working_dir": str(row.get("working_dir") or ""), } for row in list_projects(ctx.DRIVE_ROOT)] roots = _addressable_root_tasks(ctx, None) - - results_error = "" - try: - facts, unreadable = raw_result_facts(task_results_dir(ctx.DRIVE_ROOT, create=False)) - except OSError as exc: - facts, unreadable = {}, ["result_directory_unreadable"] - results_error = f"result_directory_unreadable: {exc}" - ordered = sorted(facts, key=lambda name: facts[name]["ts"] or facts[name]["updated_at"], reverse=True) - # Only the owner's ROOT results are addressable predecessors (owner decision batch - # 3, answer 6b=A): a swarm wave's children are the newest results of ANY kind, so - # they evicted the owner's own roots from this window - which is how a root the - # same actor had just read stopped being offerable. The facts are already - # memoized, so both the filter and this count cost no extra read. The count runs - # over the WHOLE candidate list, not inside the capped loop: children older than - # the 16th root are skipped just the same, and counting them only until the cap - # reported zero while folding them into the cap's own number. - def _is_child(name: str) -> bool: - return bool(facts[name]["parent_task_id"]) or facts[name]["delegation_role"] == "subagent" - - children = sum(1 for name in ordered if not facts[name]["schema_refusal"] and _is_child(name)) - finals = [] - for name in ordered: - if facts[name]["schema_refusal"] or _is_child(name): - continue - row = load_task_result(ctx.DRIVE_ROOT, pathlib.Path(name).stem) - if row is not None: - finals.append(_task_result_ground_truth(row)) - if len(finals) == 16: - break + finals, result_omissions = _recent_root_results(ctx) dialogue_rows: list = [] root = pathlib.Path(ctx.DRIVE_ROOT) @@ -333,11 +413,7 @@ def _main_routing_manifest(ctx: Any) -> Dict[str, Any]: "omissions": { "projects": max(0, len(projects) - 40), "root_tasks": max(0, len(roots) - 40), - # Kept meaning: results cut by the 16 cap. The children skipped above are - # a DIFFERENT omission and are counted as such, never folded in here. - "final_results": None if unreadable else max(0, len(facts) - children - len(finals)), - "final_results_error": results_error, - "children": None if unreadable else children, + **result_omissions, # A bounded read cannot count bytes/rows it deliberately did not # visit. The exact historical messages remain available by id. "dialogue_rows": None, @@ -393,6 +469,10 @@ def _decision_turn_metadata(ctx: Any, chat_id: int, client_message_id: str, task md["project_last_task_result"] = _task_result_ground_truth(row) except Exception: log.debug("project last-task-result projection failed", exc_info=True) + try: + md["project_routing_manifest"] = _project_routing_manifest(ctx, project_id) + except Exception: + log.warning("Unable to build the project routing manifest", exc_info=True) if client_message_id: md["client_message_id"] = client_message_id option_roots = ( diff --git a/tests/test_context_runtime_section.py b/tests/test_context_runtime_section.py index 843c4e4e4..e5751d31f 100644 --- a/tests/test_context_runtime_section.py +++ b/tests/test_context_runtime_section.py @@ -223,6 +223,46 @@ def test_runtime_section_exposes_host_routing_manifest_and_manual_contract(tmp_p assert payload["routing_contract"]["on_uncertain_or_invalid_target"] == "needs_manual_target" +def test_runtime_section_offers_the_rooms_own_continuation_hint(tmp_path, monkeypatch): + """A project room's routing manifest has to REACH the decision turn. + + While only the pointer row travelled, a room saw exactly ONE continuation + candidate; when a child had stamped that pointer the room could not name its + own interrupted root at all and promoted again, minting a duplicate root. + """ + env = _make_health_env(tmp_path) + monkeypatch.setattr("ouroboros.config.get_runtime_mode", lambda: "advanced") + room = { + "final_results": [{"task_id": "racer-old", "status": "completed"}], + "active_roots": [{"task_id": "racer-live", "status": "running", + "cancel_state": "pending"}], + "omissions": {"final_results": 3, "children": 2, "active_roots": 0}, + } + task = { + "id": "decision-room", + "type": "task", + "metadata": { + "current_chat": { + "chat_id": 7, + "running_tasks": [], + "addressable_root_tasks": [{"task_id": "racer-live", "status": "running"}], + }, + "project_routing_manifest": room, + "project_last_task_result": {"task_id": "racer-old", "status": "completed"}, + }, + } + + payload = json.loads(build_runtime_section(env, task).split("\n\n", 1)[1]) + assert payload["project_routing_manifest"] == room + assert payload["project_last_task_result"]["task_id"] == "racer-old" + + # The quiet direction: a lane with nothing to offer states no empty hint. + task["metadata"]["project_routing_manifest"] = {} + quiet = json.loads(build_runtime_section(env, task).split("\n\n", 1)[1]) + assert "project_routing_manifest" not in quiet + assert quiet["project_last_task_result"]["task_id"] == "racer-old" + + def test_improvement_backlog_digest_is_actor_scoped(tmp_path): from ouroboros.context import build_llm_messages from ouroboros.memory import Memory