mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Allow explicit conditional reads of child results
Keep the join-ledger semantic identity and current handoff facts while omitting unchanged result and trace text only when the caller supplies the known hash. Unconditional and source reads retain full access. Local checkpoint for issue #414. NOT_REVIEWED: the root coordinator owns the combined phase review. Not published.
This commit is contained in:
parent
0a9c5e9236
commit
86bdaacdc1
3 changed files with 261 additions and 4 deletions
|
|
@ -378,6 +378,7 @@ def get_tools() -> List[ToolEntry]:
|
|||
"description": "Read the effective result or exact authority of a task, including one bounded canonical work-order source range when requested.",
|
||||
"parameters": {"type": "object", "required": ["task_id"], "properties": {
|
||||
"task_id": {"type": "string", "description": "Task ID returned by scheduling or exposed by the host routing manifest."},
|
||||
"known_result_sha256": {"type": "string", "description": "Optional child_result_sha256 from a previous read. An exact match omits only unchanged result/trace text, retaining current facts and a full-read reference. Omit for full text; explicit authority/source requests always return their requested view."},
|
||||
"include_authority": {"type": "boolean", "default": False,
|
||||
"description": "Return the exact selected result, task contract, origin, artifact references, and current plan-review authority."},
|
||||
"include_work_order_source": {"type": "boolean", "default": False,
|
||||
|
|
@ -393,6 +394,7 @@ def get_tools() -> List[ToolEntry]:
|
|||
"description": "Wait for ONE subtask to reach a terminal status and return its effective result. May return EARLY (before terminal) if the child raises a tree_note blocker/question/interface_contract/review_requested/delegation_constraint beacon — the result then carries a [CHILD_BEACONS] block so you can steer, review, or override it. An unread message in your own mailbox also returns early so the ordinary loop can deliver and acknowledge it; the child keeps running. With SEVERAL children in flight, prefer wait_tasks(any_terminal) to absorb whichever finishes first rather than blocking serially on one id at a time.",
|
||||
"parameters": {"type": "object", "required": ["task_id"], "properties": {
|
||||
"task_id": {"type": "string", "description": "Task ID to check"},
|
||||
"known_result_sha256": {"type": "string", "description": "Optional child_result_sha256 already obtained for this task. An exact match returns unchanged without repeating result/trace; current facts remain. Omit to return full text. This does not change when the wait ends."},
|
||||
"timeout_sec": {"type": "integer", "default": 180, "description":
|
||||
"Maximum seconds to wait (default 180); a larger value is clamped to "
|
||||
f"{_WAIT_TASK_CLAMP_SEC}. Size the window to the child's expected life."},
|
||||
|
|
@ -403,6 +405,7 @@ def get_tools() -> List[ToolEntry]:
|
|||
"description": "Wait for MULTIPLE subtasks at once and return a compact structural projection per child (task_id, status, accounted_upper_bound_usd, cost_final, child_result_sha256, outcome_axes, result, trace_summary, capability_delta when the child has something to disclose, duplicate_of) — the right tool to ABSORB a batch of independent children you scheduled in one burst. The full per-child envelope stays on disk in task_results/<task_id>.json (child_result_sha256 pins the exact result you saw; get_task_result returns the full result text plus trace/outcome summaries). With mode=any_terminal it returns as soon as the FIRST child finishes (handle it, then call again for the rest) instead of blocking serially. The JSON also includes live_child_status (running/scheduled/terminal per child) and may early_return (before all terminal) on a child tree_note blocker/question/interface_contract/review_requested/delegation_constraint beacon so you can steer, review, or override mid-flight, or on an unread message in your own mailbox (reason=owner_mailbox_pending); the ordinary loop then handles delivery and acknowledgement. An id no surface of this tree ever minted (no task result, no queue row, no tree-ledger row) is flagged unknown_task_id — 'not yet registered or never scheduled' — and unknown_task_ids + a compact children_roster of your ACTUAL direct children are attached so you can repair the wait set instead of re-polling phantoms.",
|
||||
"parameters": {"type": "object", "required": ["task_ids"], "properties": {
|
||||
"task_ids": {"type": "array", "items": {"type": "string"}, "description": "Task IDs returned by schedule_subagent."},
|
||||
"known_result_sha256_by_task": {"type": "object", "additionalProperties": {"type": "string"}, "description": "Optional task_id to previously obtained child_result_sha256 map. Matching children omit only result/trace and return result_unchanged plus a full-read reference. Missing or different hashes return the usual complete body/trace. Current status/cost/outcome/capability facts remain; wait timing is unchanged."},
|
||||
"timeout_sec": {"type": "integer", "default": 600, "description":
|
||||
"Maximum seconds to wait (default 600); a larger value is clamped to "
|
||||
f"{_WAIT_TASKS_CLAMP_SEC}. Size the window to the children's expected life; "
|
||||
|
|
|
|||
|
|
@ -165,10 +165,25 @@ def _subtask_outcome_summary(data: Dict[str, Any], receipts: list | None = None)
|
|||
return json.dumps(summary, ensure_ascii=False, indent=2, default=str)
|
||||
|
||||
|
||||
def _unchanged_result_reference(task_id: str, current_hash: str, known_hash: Any) -> Dict[str, Any]:
|
||||
"""Omit only an explicitly matched semantic body, never its current facts.
|
||||
|
||||
This is a conditional read, not evidence that the caller still remembers or
|
||||
has accepted the result. The source request deliberately carries no condition.
|
||||
"""
|
||||
if not isinstance(known_hash, str) or known_hash != current_hash:
|
||||
return {}
|
||||
return {
|
||||
"result_unchanged": True,
|
||||
"result_source": {"tool": "get_task_result", "arguments": {"task_id": task_id}},
|
||||
}
|
||||
|
||||
|
||||
def _get_task_result(
|
||||
ctx: ToolContext, task_id: str, include_authority: bool = False,
|
||||
include_work_order_source: bool = False, source_start_char: Any = None,
|
||||
source_end_char: Any = None, include_completion_source: bool = False,
|
||||
known_result_sha256: str = "",
|
||||
) -> str:
|
||||
"""Read a task result, or a bounded canonical work-order/completion source range."""
|
||||
metadata = getattr(ctx, "task_metadata", {}) if isinstance(getattr(ctx, "task_metadata", {}), dict) else {}
|
||||
|
|
@ -246,7 +261,20 @@ def _get_task_result(
|
|||
# the stored result no longer crashes the f-string with a TypeError).
|
||||
from ouroboros.cost_projection import cost_display
|
||||
|
||||
if status == STATUS_COMPLETED:
|
||||
unchanged = _unchanged_result_reference(str(task_id), child_result_sha256, known_result_sha256)
|
||||
if unchanged:
|
||||
# Accounting, receipts, authority and capability facts are deliberately
|
||||
# outside the join-ledger result identity; keep their current projection.
|
||||
if data.get("duplicate_of"):
|
||||
unchanged["duplicate_of"] = str(data["duplicate_of"])
|
||||
output = (
|
||||
f"Task {task_id} [{status}]: cost={cost_display(data)}\n"
|
||||
f"child_result_sha256={child_result_sha256}\n\n"
|
||||
f"[SUBTASK_OUTCOME]\n{outcome_summary}\n[/SUBTASK_OUTCOME]\n\n"
|
||||
f"{json.dumps(unchanged, ensure_ascii=False)}\n"
|
||||
"Result and trace are unchanged; omit known_result_sha256 to read them in full."
|
||||
)
|
||||
elif status == STATUS_COMPLETED:
|
||||
output = (
|
||||
f"Task {task_id} [{status}]: cost={cost_display(data)}\n"
|
||||
f"child_result_sha256={child_result_sha256}\n\n"
|
||||
|
|
@ -268,7 +296,7 @@ def _get_task_result(
|
|||
f"[SUBTASK_OUTCOME]\n{outcome_summary}\n[/SUBTASK_OUTCOME]\n\n"
|
||||
f"{result or 'No details available.'}"
|
||||
)
|
||||
if trace:
|
||||
if trace and not unchanged:
|
||||
output += f"\n\n[SUBTASK_TRACE]\n{trace}\n[/SUBTASK_TRACE]"
|
||||
from ouroboros.task_finalization import provider_terminal_body, terminal_host_notice_text
|
||||
|
||||
|
|
@ -424,7 +452,9 @@ def cache_horizon_note(ctx: Any, elapsed_sec: Any) -> str:
|
|||
)
|
||||
|
||||
|
||||
def _wait_for_task(ctx: ToolContext, task_id: str, timeout_sec: int = 180) -> str:
|
||||
def _wait_for_task(
|
||||
ctx: ToolContext, task_id: str, timeout_sec: int = 180, known_result_sha256: str = "",
|
||||
) -> str:
|
||||
"""Wait for a subtask to reach a terminal status."""
|
||||
try:
|
||||
tid = validate_task_id(task_id)
|
||||
|
|
@ -465,7 +495,9 @@ def _wait_for_task(ctx: ToolContext, task_id: str, timeout_sec: int = 180) -> st
|
|||
horizon_note = cache_horizon_note(ctx, waited.get("elapsed_sec"))
|
||||
if horizon_note:
|
||||
extra += f"\n\n{horizon_note}"
|
||||
return f"{header} after {waited.get('elapsed_sec', 0):.1f}s.{extra}\n\n{_get_task_result(ctx, tid)}"
|
||||
result = (_get_task_result(ctx, tid, known_result_sha256=known_result_sha256)
|
||||
if known_result_sha256 else _get_task_result(ctx, tid))
|
||||
return f"{header} after {waited.get('elapsed_sec', 0):.1f}s.{extra}\n\n{result}"
|
||||
|
||||
|
||||
def _count_live_sibling_children(ctx: ToolContext, status_drive_root: Path, *, exclude_task_id: str) -> int:
|
||||
|
|
@ -598,6 +630,7 @@ def _wait_for_tasks(
|
|||
task_ids: List[str],
|
||||
timeout_sec: int = 600,
|
||||
mode: str = "all_terminal",
|
||||
known_result_sha256_by_task: Dict[str, str] | None = None,
|
||||
) -> str:
|
||||
"""Wait for multiple subtasks and return a compact structural projection per child.
|
||||
|
||||
|
|
@ -799,6 +832,14 @@ def _wait_for_tasks(
|
|||
# (metered) contribution beside them is unknown.
|
||||
_ee["native_contribution"] = "unknown"
|
||||
projected["execution_evidence"] = _ee
|
||||
known = (known_result_sha256_by_task.get(str(tid))
|
||||
if isinstance(known_result_sha256_by_task, dict) else None)
|
||||
unchanged = (_unchanged_result_reference(str(tid), projected["child_result_sha256"], known)
|
||||
if data else {})
|
||||
if unchanged:
|
||||
projected.pop("result", None)
|
||||
projected.pop("trace_summary", None)
|
||||
projected.update(unchanged)
|
||||
public_tasks[str(tid)] = projected
|
||||
waited["tasks"] = public_tasks
|
||||
waited["tasks_note"] = (
|
||||
|
|
|
|||
213
tests/test_conditional_child_results.py
Normal file
213
tests/test_conditional_child_results.py
Normal file
|
|
@ -0,0 +1,213 @@
|
|||
"""Conditional result reads retain current facts without recording a seen state."""
|
||||
from __future__ import annotations
|
||||
|
||||
import copy
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros.task_results import write_task_result
|
||||
from ouroboros.task_status import load_effective_task_result
|
||||
from ouroboros.tools import control_task_results as results
|
||||
from ouroboros.tools.join_ledger import _child_result_sha256
|
||||
|
||||
BODY = "Complete child analysis.\n" * 1200
|
||||
TRACE = "Exact trace summary.\n" * 400
|
||||
|
||||
|
||||
def _ctx(drive):
|
||||
return SimpleNamespace(
|
||||
drive_root=drive, budget_drive_root=drive, task_id="parent1",
|
||||
task_attempt=1, task_metadata={"root_task_id": "parent1", "budget_drive_root": str(drive)},
|
||||
_loop_mailbox_seen_ids=set(),
|
||||
)
|
||||
|
||||
|
||||
def _child(drive, task_id="child1", **changes):
|
||||
values = dict(
|
||||
result=BODY, trace_summary=TRACE, parent_task_id="parent1",
|
||||
root_task_id="parent1", delegation_role="subagent",
|
||||
)
|
||||
values.update(changes)
|
||||
write_task_result(drive, task_id, "completed", **values)
|
||||
return load_effective_task_result(drive, task_id)
|
||||
|
||||
|
||||
def _digest(drive, task_id="child1"):
|
||||
return _child_result_sha256(load_effective_task_result(drive, task_id))
|
||||
|
||||
|
||||
@pytest.mark.parametrize("surface", ["get", "wait", "batch"])
|
||||
def test_condition_is_explicit_and_full_read_survives_a_new_context(tmp_path, surface):
|
||||
_child(tmp_path)
|
||||
known = _digest(tmp_path)
|
||||
|
||||
def read(ctx, condition=False):
|
||||
if surface == "batch":
|
||||
kwargs = {"known_result_sha256_by_task": {"child1": known}} if condition else {}
|
||||
return results._wait_for_tasks(ctx, ["child1"], timeout_sec=0, **kwargs)
|
||||
handler = results._get_task_result if surface == "get" else results._wait_for_task
|
||||
kwargs = {"known_result_sha256": known} if condition else {}
|
||||
if surface == "wait":
|
||||
kwargs["timeout_sec"] = 0
|
||||
return handler(ctx, "child1", **kwargs)
|
||||
|
||||
ctx = _ctx(tmp_path)
|
||||
before = (tmp_path / "task_results" / "child1.json").read_bytes()
|
||||
first = read(ctx)
|
||||
unchanged = read(ctx, True)
|
||||
assert BODY in (json.loads(first)["tasks"]["child1"]["result"] if surface == "batch" else first)
|
||||
assert known in first and known in unchanged
|
||||
assert "result_unchanged" in unchanged
|
||||
assert "Complete child analysis." not in unchanged
|
||||
assert "Exact trace summary." not in unchanged
|
||||
assert len(unchanged) < len(first) / 5
|
||||
# A result hash is neither "seen" nor proof of a current in-context copy.
|
||||
def result_view(text):
|
||||
return json.loads(text)["tasks"] if surface == "batch" else text
|
||||
assert result_view(read(ctx)) == result_view(first)
|
||||
assert result_view(read(_ctx(tmp_path))) == result_view(first)
|
||||
assert (tmp_path / "task_results" / "child1.json").read_bytes() == before
|
||||
assert not (tmp_path / "state" / "task_trees").exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("field,new_value", [
|
||||
("result", "A revised complete answer."),
|
||||
("trace_summary", "A revised trace."),
|
||||
("status", "failed"),
|
||||
("artifact_status", "failed"),
|
||||
("artifacts", [{"name": "report.txt", "sha256": "b" * 64}]),
|
||||
("terminal_host_notice", "Unresolved delegated execution remains."),
|
||||
])
|
||||
def test_semantic_change_returns_full_single_and_batch(tmp_path, monkeypatch, field, new_value):
|
||||
data = dict(task_id="child1", status="completed", result=BODY, trace_summary=TRACE)
|
||||
original = _child_result_sha256(data)
|
||||
data[field] = new_value
|
||||
monkeypatch.setattr(results, "load_effective_task_result", lambda *_args: copy.deepcopy(data))
|
||||
monkeypatch.setattr(results, "_unminted_wait_ids", lambda *_args: [])
|
||||
monkeypatch.setattr(results, "wait_for_effective_tasks", lambda *_args, **_kwargs: {
|
||||
"tasks": {"child1": copy.deepcopy(data)}, "all_terminal": True, "elapsed_sec": 0,
|
||||
})
|
||||
single = results._get_task_result(_ctx(tmp_path), "child1", known_result_sha256=original)
|
||||
batch = json.loads(results._wait_for_tasks(
|
||||
_ctx(tmp_path), ["child1"], timeout_sec=0,
|
||||
known_result_sha256_by_task={"child1": original},
|
||||
))["tasks"]["child1"]
|
||||
assert "result_unchanged" not in single
|
||||
assert "result_unchanged" not in batch
|
||||
assert str(data["result"]) in single
|
||||
assert batch["result"] == data["result"]
|
||||
|
||||
|
||||
def test_batch_omits_only_matched_children_and_keeps_unknown_truth(tmp_path):
|
||||
_child(tmp_path)
|
||||
_child(tmp_path, "child2", result="Other complete result")
|
||||
old = _digest(tmp_path)
|
||||
result = json.loads(results._wait_for_tasks(
|
||||
_ctx(tmp_path), ["child1", "child2", "unknown"], timeout_sec=0,
|
||||
known_result_sha256_by_task={"child1": old, "unknown": _child_result_sha256({})},
|
||||
))
|
||||
assert result["tasks"]["child1"]["result_unchanged"] is True
|
||||
assert "result" not in result["tasks"]["child1"]
|
||||
assert result["tasks"]["child2"]["result"] == "Other complete result"
|
||||
assert result["tasks"]["unknown"].get("result_unchanged") is not True
|
||||
source = result["tasks"]["child1"]["result_source"]
|
||||
full = results._get_task_result(_ctx(tmp_path), **source["arguments"])
|
||||
assert BODY in full
|
||||
|
||||
|
||||
def test_accounting_and_current_handoff_facts_do_not_duplicate_body(tmp_path, monkeypatch):
|
||||
data = dict(task_id="child1", status="completed", result=BODY, trace_summary=TRACE,
|
||||
delegate_terminal_reconciliation={"audit_status": "pending", "open_run_ids": ["run-a"]})
|
||||
original = _child_result_sha256(data)
|
||||
data.update(
|
||||
cost_usd=None, cost_final=False,
|
||||
capability_delta={"reduced": True, "legacy_note": "Source unavailable now"},
|
||||
verification_ledger={"schema_version": 1, "summary": {"entry_count": 1, "failed": 1}},
|
||||
)
|
||||
assert _child_result_sha256(data) == original
|
||||
monkeypatch.setattr(results, "load_effective_task_result", lambda *_args: copy.deepcopy(data))
|
||||
monkeypatch.setattr(results, "_unminted_wait_ids", lambda *_args: [])
|
||||
monkeypatch.setattr(results, "wait_for_effective_tasks", lambda *_args, **_kwargs: {
|
||||
"tasks": {"child1": copy.deepcopy(data)}, "all_terminal": True, "elapsed_sec": 0,
|
||||
})
|
||||
shown = results._get_task_result(_ctx(tmp_path), "child1", known_result_sha256=original)
|
||||
assert "cost=unknown" in shown
|
||||
assert "Source unavailable now" in shown and "run-a" in shown
|
||||
assert '"entry_count": 1' in shown
|
||||
assert BODY not in shown
|
||||
batch = json.loads(results._wait_for_tasks(
|
||||
_ctx(tmp_path), ["child1"], timeout_sec=0,
|
||||
known_result_sha256_by_task={"child1": original},
|
||||
))["tasks"]["child1"]
|
||||
assert batch["accounted_upper_bound_usd"] is None and batch["cost_final"] is False
|
||||
assert batch["capability_delta"] == data["capability_delta"]
|
||||
|
||||
|
||||
def test_source_or_authority_request_is_not_suppressed_by_a_known_result(tmp_path, monkeypatch):
|
||||
_child(tmp_path)
|
||||
import ouroboros.agent_startup_checks as checks
|
||||
import ouroboros.task_finalization as finalization
|
||||
|
||||
monkeypatch.setattr(checks, "task_result_authority_projection", lambda *_args, **_kwargs: {"current": True})
|
||||
monkeypatch.setattr(finalization, "completion_source_projection", lambda *_args: {"text": "exact source"})
|
||||
result = json.loads(results._get_task_result(
|
||||
_ctx(tmp_path), "child1", include_authority=True, include_completion_source=True,
|
||||
known_result_sha256=_digest(tmp_path),
|
||||
))
|
||||
assert result["authority"] == {"current": True}
|
||||
assert result["completion_source"]["text"] == "exact source"
|
||||
assert "result_unchanged" not in result
|
||||
|
||||
|
||||
@pytest.mark.parametrize("known", ["", "stale", None, 17, {"child1": "bad"}])
|
||||
def test_nonmatching_condition_never_hides_result(tmp_path, known):
|
||||
_child(tmp_path)
|
||||
assert BODY in results._get_task_result(_ctx(tmp_path), "child1", known_result_sha256=known)
|
||||
|
||||
|
||||
def test_unchanged_wait_still_delivers_parent_mailbox_without_ack(tmp_path):
|
||||
from ouroboros.owner_mailbox import acknowledged_task_message_ids, write_owner_message
|
||||
|
||||
_child(tmp_path)
|
||||
known = _digest(tmp_path)
|
||||
assert write_owner_message(tmp_path, "Keep the complete original.", "parent1", msg_id="owner-followup")
|
||||
ctx = _ctx(tmp_path)
|
||||
single = results._wait_for_task(ctx, "child1", timeout_sec=0, known_result_sha256=known)
|
||||
assert "unread message" in single and "result_unchanged" in single
|
||||
assert acknowledged_task_message_ids(tmp_path, "parent1", attempt_key=1) == set()
|
||||
batch = json.loads(results._wait_for_tasks(
|
||||
_ctx(tmp_path), ["child1"], timeout_sec=0,
|
||||
known_result_sha256_by_task={"child1": known},
|
||||
))
|
||||
assert batch["early_return"]["reason"] == "owner_mailbox_pending"
|
||||
assert batch["tasks"]["child1"]["result_unchanged"] is True
|
||||
|
||||
|
||||
def test_public_schemas_offer_the_condition_without_changing_required_args():
|
||||
from ouroboros.tools.control import get_tools
|
||||
|
||||
tools = {entry.name: entry for entry in get_tools()}
|
||||
for name in ("get_task_result", "wait_task", "wait_tasks"):
|
||||
schema = tools[name].schema["parameters"]
|
||||
field = "known_result_sha256_by_task" if name == "wait_tasks" else "known_result_sha256"
|
||||
assert field in schema["properties"]
|
||||
assert field not in schema["required"]
|
||||
|
||||
|
||||
def test_duplicate_reference_and_new_receipts_remain_visible_when_body_matches(tmp_path, monkeypatch):
|
||||
import ouroboros.outcomes as outcomes
|
||||
|
||||
data = dict(task_id="child1", status="rejected_duplicate", result="Duplicate answer",
|
||||
trace_summary="trace", duplicate_of="original")
|
||||
known = _child_result_sha256(data)
|
||||
monkeypatch.setattr(results, "load_effective_task_result", lambda *_args: data)
|
||||
monkeypatch.setattr(outcomes, "read_verification_receipts_from_roots", lambda *_args: [
|
||||
{"status": "failed", "check": "a newly observed verification", "matched": False},
|
||||
])
|
||||
result = results._get_task_result(_ctx(tmp_path), "child1", known_result_sha256=known)
|
||||
assert '"duplicate_of": "original"' in result
|
||||
assert "a newly observed verification" in result
|
||||
assert '"matched": false' in result
|
||||
assert "Duplicate answer" not in result
|
||||
Loading…
Add table
Add a link
Reference in a new issue