mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
fix: converge child reference promotion recovery
Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
parent
fefbdac4d4
commit
72b502228f
4 changed files with 710 additions and 101 deletions
|
|
@ -233,9 +233,19 @@ def prune_headless_task_drives(
|
|||
base = parent / HEADLESS_TASKS_DIR
|
||||
days = _resolve_retention_days(retention_days)
|
||||
cutoff = age_cutoff(days, now)
|
||||
report: Dict[str, Any] = {"retention_days": days, "scanned": 0, "pruned": [], "skipped": [], "errors": []}
|
||||
report: Dict[str, Any] = {
|
||||
"retention_days": days,
|
||||
"scanned": 0,
|
||||
"pruned": [],
|
||||
"skipped": [],
|
||||
"errors": [],
|
||||
"promotion_retry": {},
|
||||
}
|
||||
if not base.is_dir():
|
||||
return report
|
||||
from ouroboros.observability import retry_pending_child_ref_promotions
|
||||
|
||||
report["promotion_retry"] = retry_pending_child_ref_promotions(parent)
|
||||
for task_dir in sorted(base.iterdir()):
|
||||
if not task_dir.is_dir():
|
||||
continue
|
||||
|
|
|
|||
|
|
@ -19,7 +19,7 @@ import uuid
|
|||
from dataclasses import dataclass, field
|
||||
from typing import Any, Callable, Dict, List, Optional, Tuple
|
||||
|
||||
from ouroboros.utils import atomic_write_json, replace_atomic, utc_now_iso, write_bytes_atomic
|
||||
from ouroboros.utils import atomic_write_json, replace_atomic, utc_now_iso
|
||||
|
||||
|
||||
OBSERVABILITY_DIR = "observability"
|
||||
|
|
@ -481,6 +481,7 @@ def promote_call_manifest_ref(
|
|||
_PUBLISHED_CHILD_REF_FIELDS = frozenset(
|
||||
{
|
||||
"trace_refs",
|
||||
"loop_outcome",
|
||||
"review_evidence",
|
||||
"review_projection",
|
||||
"verification_ledger",
|
||||
|
|
@ -491,7 +492,7 @@ _PUBLISHED_CHILD_REF_FIELDS = frozenset(
|
|||
}
|
||||
)
|
||||
_SOURCE_HANDLES_SUBDIR = "source_handles"
|
||||
_SOURCE_HANDLE_DIGEST_RE = re.compile(r"-([0-9a-f]{64})\.[A-Za-z0-9]+$")
|
||||
_TASK_SOURCE_MARKER = "FULL_RESULT_SOURCE_JSON="
|
||||
_SERVICE_REF_TOOLS = frozenset({"service_logs", "stop_service"})
|
||||
|
||||
|
||||
|
|
@ -548,13 +549,109 @@ def _is_manifest_ref(value: Any) -> bool:
|
|||
|
||||
|
||||
def _is_task_source_ref(value: Any) -> bool:
|
||||
return bool(
|
||||
isinstance(value, dict)
|
||||
and value.get("kind") == "task_source"
|
||||
and value.get("root") == "artifact_store"
|
||||
and pathlib.PurePosixPath(str(value.get("path") or "")).parts[:1]
|
||||
== (_SOURCE_HANDLES_SUBDIR,)
|
||||
return bool(isinstance(value, dict) and value.get("kind") == "task_source")
|
||||
|
||||
|
||||
def _task_source_contract_valid(ref: Dict[str, Any]) -> bool:
|
||||
read = ref.get("read") if isinstance(ref.get("read"), dict) else {}
|
||||
arguments = (
|
||||
read.get("arguments") if isinstance(read.get("arguments"), dict) else {}
|
||||
)
|
||||
path = str(ref.get("path") or "")
|
||||
return bool(
|
||||
ref.get("root") == "artifact_store"
|
||||
and read.get("tool") == "read_file"
|
||||
and arguments.get("root") == ref.get("root")
|
||||
and str(arguments.get("path") or "") == path
|
||||
)
|
||||
|
||||
|
||||
def _task_source_failure_reason(exc: Exception) -> str:
|
||||
message = str(exc)
|
||||
if isinstance(exc, FileNotFoundError):
|
||||
return "source_missing"
|
||||
if "size verification" in message or "sha256 verification" in message:
|
||||
return "digest_mismatch"
|
||||
if "escapes" in message or "symlink" in message:
|
||||
return "invalid_scope"
|
||||
if isinstance(exc, (TypeError, ValueError)):
|
||||
return "invalid_ref"
|
||||
return "source_unreadable"
|
||||
|
||||
|
||||
def _promote_task_source_ref(
|
||||
parent_root: pathlib.Path,
|
||||
child_root: pathlib.Path,
|
||||
task_id: str,
|
||||
ref: Dict[str, Any],
|
||||
state: Dict[str, Any],
|
||||
) -> Dict[str, Any]:
|
||||
"""Promote one exact Phase3B actor ref through its own read/write seams."""
|
||||
|
||||
if not _task_source_contract_valid(ref):
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(ref, "invalid_ref")
|
||||
)
|
||||
return _typed_unavailable_ref(ref, "invalid_ref")
|
||||
try:
|
||||
from ouroboros.artifacts import (
|
||||
read_actor_source_bytes,
|
||||
store_actor_source_bytes,
|
||||
)
|
||||
|
||||
read_actor_source_bytes(parent_root, task_id, ref)
|
||||
return dict(ref)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
raw = read_actor_source_bytes(child_root, task_id, ref)
|
||||
except Exception as exc:
|
||||
reason = _task_source_failure_reason(exc)
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(ref, reason)
|
||||
)
|
||||
return _typed_unavailable_ref(ref, reason)
|
||||
|
||||
rel = pathlib.PurePosixPath(str(ref.get("path") or ""))
|
||||
expected_sha = str(ref.get("sha256") or "")
|
||||
name_match = re.fullmatch(
|
||||
rf"(.+)-{re.escape(expected_sha)}\.([A-Za-z0-9]+)",
|
||||
rel.name,
|
||||
)
|
||||
if len(rel.parts) != 3 or rel.parts[0] != _SOURCE_HANDLES_SUBDIR or not name_match:
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(ref, "invalid_ref")
|
||||
)
|
||||
return _typed_unavailable_ref(ref, "invalid_ref")
|
||||
|
||||
source = _task_artifact_dir(child_root, task_id, create=False).joinpath(
|
||||
*rel.parts
|
||||
)
|
||||
try:
|
||||
promoted = store_actor_source_bytes(
|
||||
parent_root,
|
||||
task_id,
|
||||
category=rel.parts[1],
|
||||
source_id=name_match.group(1),
|
||||
data=raw,
|
||||
extension=name_match.group(2),
|
||||
)
|
||||
if str(promoted.get("path") or "") != rel.as_posix():
|
||||
raise OSError("canonical task source path changed during promotion")
|
||||
read_actor_source_bytes(parent_root, task_id, ref)
|
||||
state["promoted_source_handle_count"] += 1
|
||||
return dict(ref)
|
||||
except Exception as exc:
|
||||
_append_promotion_fact(
|
||||
state["pending_refs"],
|
||||
_promotion_fact(
|
||||
{**ref, "path": str(source)},
|
||||
f"{type(exc).__name__}: {exc}",
|
||||
),
|
||||
)
|
||||
state["status"] = "incomplete"
|
||||
return dict(ref)
|
||||
|
||||
|
||||
def _promote_known_observability_ref(
|
||||
|
|
@ -571,12 +668,25 @@ def _promote_known_observability_ref(
|
|||
raise OSError("embedded child observability ref promotion is pending")
|
||||
return rewritten
|
||||
|
||||
source_root = child_root
|
||||
try:
|
||||
ref_path = pathlib.Path(str(ref.get("path") or "")).resolve(strict=False)
|
||||
ref_path.relative_to(_observability_root(parent_root).resolve(strict=False))
|
||||
source_root = parent_root
|
||||
except (OSError, ValueError):
|
||||
pass
|
||||
|
||||
try:
|
||||
promoted = (
|
||||
promote_blob_ref(child_root, parent_root, ref, transform_json=transform_json)
|
||||
promote_blob_ref(
|
||||
source_root,
|
||||
parent_root,
|
||||
ref,
|
||||
transform_json=transform_json,
|
||||
)
|
||||
if _is_blob_ref(ref)
|
||||
else promote_call_manifest_ref(
|
||||
child_root, parent_root, ref, task_id=task_id,
|
||||
source_root, parent_root, ref, task_id=task_id,
|
||||
transform_json=transform_json,
|
||||
)
|
||||
)
|
||||
|
|
@ -613,6 +723,46 @@ def _rewrite_service_result(
|
|||
return prefix + json.dumps(rewritten, ensure_ascii=False, indent=2)
|
||||
|
||||
|
||||
def _rewrite_task_source_markers(
|
||||
text: str,
|
||||
parent_root: pathlib.Path,
|
||||
child_root: pathlib.Path,
|
||||
task_id: str,
|
||||
state: Dict[str, Any],
|
||||
) -> str:
|
||||
"""Rewrite only Phase3B's explicit actor-source envelope inside tool text."""
|
||||
|
||||
rewritten_lines: List[str] = []
|
||||
for line in str(text).splitlines(keepends=True):
|
||||
body = line.rstrip("\r\n")
|
||||
newline = line[len(body):]
|
||||
if not body.startswith(_TASK_SOURCE_MARKER):
|
||||
rewritten_lines.append(line)
|
||||
continue
|
||||
try:
|
||||
ref = json.loads(body[len(_TASK_SOURCE_MARKER):])
|
||||
except (TypeError, ValueError):
|
||||
rewritten_lines.append(line)
|
||||
continue
|
||||
if not _is_task_source_ref(ref):
|
||||
rewritten_lines.append(line)
|
||||
continue
|
||||
promoted = _promote_task_source_ref(
|
||||
parent_root, child_root, task_id, ref, state
|
||||
)
|
||||
rewritten_lines.append(
|
||||
_TASK_SOURCE_MARKER
|
||||
+ json.dumps(
|
||||
promoted,
|
||||
ensure_ascii=False,
|
||||
sort_keys=True,
|
||||
separators=(",", ":"),
|
||||
)
|
||||
+ newline
|
||||
)
|
||||
return "".join(rewritten_lines)
|
||||
|
||||
|
||||
def _rewrite_service_payload(
|
||||
payload: Any,
|
||||
parent_root: pathlib.Path,
|
||||
|
|
@ -648,71 +798,9 @@ def _rewrite_child_ref_tree(
|
|||
if _is_blob_ref(value) or _is_manifest_ref(value):
|
||||
return _promote_known_observability_ref(parent_root, child_root, task_id, value, state)
|
||||
if _is_task_source_ref(value):
|
||||
rel = pathlib.PurePosixPath(str(value.get("path") or ""))
|
||||
try:
|
||||
expected_size = int(value["size"])
|
||||
except (KeyError, TypeError, ValueError):
|
||||
expected_size = -1
|
||||
expected_sha = str(value.get("sha256") or "")
|
||||
if expected_size < 0 or not re.fullmatch(r"[0-9a-f]{64}", expected_sha):
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(value, "invalid_ref")
|
||||
)
|
||||
return _typed_unavailable_ref(value, "invalid_ref")
|
||||
source = _task_artifact_dir(child_root, task_id, create=False).joinpath(*rel.parts)
|
||||
target = _task_artifact_dir(parent_root, task_id, create=False).joinpath(*rel.parts)
|
||||
if target.is_file():
|
||||
raw = target.read_bytes()
|
||||
if (
|
||||
len(raw) == expected_size
|
||||
and hashlib.sha256(raw).hexdigest() == expected_sha
|
||||
):
|
||||
return dict(value)
|
||||
if source.is_file():
|
||||
try:
|
||||
raw = source.read_bytes()
|
||||
except OSError as exc:
|
||||
reason = f"source_unreadable:{type(exc).__name__}"
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(value, reason)
|
||||
)
|
||||
return _typed_unavailable_ref(value, reason)
|
||||
match = _SOURCE_HANDLE_DIGEST_RE.search(source.name)
|
||||
if (
|
||||
len(raw) != expected_size
|
||||
or hashlib.sha256(raw).hexdigest() != expected_sha
|
||||
or match is None
|
||||
or match.group(1) != expected_sha
|
||||
):
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"],
|
||||
_promotion_fact(value, "digest_mismatch"),
|
||||
)
|
||||
return _typed_unavailable_ref(value, "digest_mismatch")
|
||||
try:
|
||||
write_bytes_atomic(target, raw)
|
||||
copied = target.read_bytes()
|
||||
if (
|
||||
len(copied) != expected_size
|
||||
or hashlib.sha256(copied).hexdigest() != expected_sha
|
||||
):
|
||||
raise OSError("canonical task source verification failed")
|
||||
state["promoted_source_handle_count"] += 1
|
||||
return dict(value)
|
||||
except Exception as exc:
|
||||
_append_promotion_fact(
|
||||
state["pending_refs"],
|
||||
_promotion_fact(
|
||||
{**value, "path": str(source)},
|
||||
f"{type(exc).__name__}: {exc}",
|
||||
),
|
||||
)
|
||||
state["status"] = "incomplete"
|
||||
return dict(value)
|
||||
_append_promotion_fact(
|
||||
state["unavailable_refs"], _promotion_fact(value, "source_missing")
|
||||
return _promote_task_source_ref(
|
||||
parent_root, child_root, task_id, value, state
|
||||
)
|
||||
return _typed_unavailable_ref(value, "source_missing")
|
||||
if isinstance(value, dict):
|
||||
return {
|
||||
key: _rewrite_child_ref_tree(
|
||||
|
|
@ -725,6 +813,10 @@ def _rewrite_child_ref_tree(
|
|||
_rewrite_child_ref_tree(item, parent_root, child_root, task_id, state)
|
||||
for item in value
|
||||
]
|
||||
if isinstance(value, str) and _TASK_SOURCE_MARKER in value:
|
||||
return _rewrite_task_source_markers(
|
||||
value, parent_root, child_root, task_id, state
|
||||
)
|
||||
return value
|
||||
|
||||
|
||||
|
|
@ -756,6 +848,121 @@ def promote_child_task_refs(
|
|||
return rewritten, state
|
||||
|
||||
|
||||
def promote_child_task_ref_patch(
|
||||
parent_drive_root: pathlib.Path,
|
||||
child_drive_root: pathlib.Path,
|
||||
task_id: str,
|
||||
canonical_result: Dict[str, Any],
|
||||
) -> Dict[str, Any]:
|
||||
"""Return only ref-bearing canonical fields for a pending retry write."""
|
||||
|
||||
rewritten, state = promote_child_task_refs(
|
||||
parent_drive_root,
|
||||
child_drive_root,
|
||||
task_id,
|
||||
canonical_result,
|
||||
)
|
||||
patch = {
|
||||
key: rewritten[key]
|
||||
for key in _PUBLISHED_CHILD_REF_FIELDS
|
||||
if key in rewritten
|
||||
}
|
||||
patch["child_ref_promotion"] = state
|
||||
return patch
|
||||
|
||||
|
||||
def _has_pending_ref_promotion(promotion: Any) -> bool:
|
||||
if not isinstance(promotion, dict):
|
||||
return False
|
||||
try:
|
||||
version = int(promotion.get("schema_version") or 0)
|
||||
except (TypeError, ValueError):
|
||||
return False
|
||||
return bool(
|
||||
version == 1
|
||||
and str(promotion.get("status") or "") != "complete"
|
||||
and isinstance(promotion.get("pending_refs"), list)
|
||||
and promotion.get("pending_refs")
|
||||
)
|
||||
|
||||
|
||||
def _retry_pending_child_ref_promotion(
|
||||
parent: pathlib.Path,
|
||||
child: pathlib.Path,
|
||||
task_id: str,
|
||||
loaded_result: Dict[str, Any],
|
||||
) -> Dict[str, Any]:
|
||||
"""Retry ref fields from CURRENT canonical authority under its result lock."""
|
||||
|
||||
from ouroboros.task_results import write_task_result
|
||||
|
||||
def _project(current: Dict[str, Any], _incoming: Dict[str, Any]) -> Dict[str, Any]:
|
||||
current_status = str(current.get("status") or loaded_result.get("status") or "")
|
||||
if not _has_pending_ref_promotion(current.get("child_ref_promotion")):
|
||||
return {"status": current_status}
|
||||
patch = promote_child_task_ref_patch(parent, child, task_id, current)
|
||||
patch["status"] = current_status
|
||||
return patch
|
||||
|
||||
return write_task_result(
|
||||
parent,
|
||||
task_id,
|
||||
str(loaded_result.get("status") or ""),
|
||||
_field_projector=_project,
|
||||
)
|
||||
|
||||
|
||||
def retry_pending_child_ref_promotions(
|
||||
parent_drive_root: pathlib.Path,
|
||||
) -> Dict[str, Any]:
|
||||
"""Retry only newly ledgered pending refs, never the stale child result."""
|
||||
|
||||
from ouroboros.headless import HEADLESS_TASKS_DIR
|
||||
from ouroboros.task_status import SETTLED_STATUSES
|
||||
from ouroboros.task_results import load_task_result, validate_task_id
|
||||
|
||||
parent = pathlib.Path(parent_drive_root)
|
||||
base = parent / HEADLESS_TASKS_DIR
|
||||
report: Dict[str, Any] = {
|
||||
"scanned": 0,
|
||||
"retried": [],
|
||||
"completed": [],
|
||||
"pending": [],
|
||||
"errors": [],
|
||||
}
|
||||
if not base.is_dir():
|
||||
return report
|
||||
for task_dir in sorted(base.iterdir()):
|
||||
if not task_dir.is_dir():
|
||||
continue
|
||||
task_id = task_dir.name
|
||||
report["scanned"] += 1
|
||||
try:
|
||||
validate_task_id(task_id)
|
||||
result = load_task_result(parent, task_id) or {}
|
||||
if str(result.get("status") or "").lower() not in SETTLED_STATUSES:
|
||||
continue
|
||||
if not _has_pending_ref_promotion(result.get("child_ref_promotion")):
|
||||
continue
|
||||
settled = _retry_pending_child_ref_promotion(
|
||||
parent, task_dir / "data", task_id, result
|
||||
)
|
||||
report["retried"].append(task_id)
|
||||
promotion = settled.get("child_ref_promotion") or {}
|
||||
destination = (
|
||||
"completed"
|
||||
if str(promotion.get("status") or "") == "complete"
|
||||
else "pending"
|
||||
)
|
||||
report[destination].append(task_id)
|
||||
except Exception as exc:
|
||||
report["errors"].append({
|
||||
"task_id": task_id,
|
||||
"error": f"{type(exc).__name__}: {exc}",
|
||||
})
|
||||
return report
|
||||
|
||||
|
||||
def write_call_manifest(
|
||||
drive_root: pathlib.Path,
|
||||
*,
|
||||
|
|
|
|||
12
server.py
12
server.py
|
|
@ -1031,9 +1031,9 @@ _LAST_CANCEL_INTENT_SWEEP = [0.0]
|
|||
|
||||
def _periodic_supervisor_maintenance(last_custody_reap: list, last_review_reconcile: list) -> None:
|
||||
"""Throttled periodic upkeep extracted from the supervisor loop: cancel-intent
|
||||
watchdog (every 20s), custody reap of orphaned task-scoped processes (every
|
||||
600s) + review-job zombie reconcile (every 300s). Each cadence gates itself
|
||||
via its own last-run marker."""
|
||||
watchdog and pending child-ref promotion replay (every 20s), custody reap of
|
||||
orphaned task-scoped processes (every 600s) + review-job zombie reconcile
|
||||
(every 300s). Each cadence gates itself via its own last-run marker."""
|
||||
if time.time() - _LAST_CANCEL_INTENT_SWEEP[0] > 20:
|
||||
_LAST_CANCEL_INTENT_SWEEP[0] = time.time()
|
||||
try:
|
||||
|
|
@ -1056,6 +1056,12 @@ def _periodic_supervisor_maintenance(last_custody_reap: list, last_review_reconc
|
|||
replay_pending_deliveries(DATA_DIR)
|
||||
except Exception:
|
||||
log.debug("Pending terminal-delivery replay failed", exc_info=True)
|
||||
try:
|
||||
from ouroboros.observability import retry_pending_child_ref_promotions
|
||||
|
||||
retry_pending_child_ref_promotions(DATA_DIR)
|
||||
except Exception:
|
||||
log.debug("Pending child-ref promotion retry failed", exc_info=True)
|
||||
if time.time() - last_custody_reap[0] > 600:
|
||||
last_custody_reap[0] = time.time()
|
||||
try:
|
||||
|
|
|
|||
|
|
@ -3,10 +3,12 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import gzip
|
||||
import hashlib
|
||||
import json
|
||||
import pathlib
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros.headless import (
|
||||
copy_child_task_result,
|
||||
|
|
@ -15,7 +17,7 @@ from ouroboros.headless import (
|
|||
remove_subagent_task_drive,
|
||||
)
|
||||
from ouroboros.observability import persist_call, read_blob_ref, write_blob
|
||||
from ouroboros.task_results import STATUS_COMPLETED, write_task_result
|
||||
from ouroboros.task_results import STATUS_COMPLETED, load_task_result, write_task_result
|
||||
|
||||
|
||||
def _child(tmp_path: pathlib.Path, task_id: str) -> tuple[pathlib.Path, pathlib.Path]:
|
||||
|
|
@ -34,6 +36,12 @@ def _manifest(ref: dict) -> dict:
|
|||
return json.loads(pathlib.Path(ref["path"]).read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def _source_ref_from_visible_result(text: str) -> dict:
|
||||
prefix = "FULL_RESULT_SOURCE_JSON="
|
||||
line = next(line for line in text.splitlines() if line.startswith(prefix))
|
||||
return json.loads(line[len(prefix):])
|
||||
|
||||
|
||||
def test_copyback_promotes_trace_manifest_and_blobs_before_headless_gc(tmp_path):
|
||||
task_id = "phase3c-trace"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
|
|
@ -108,6 +116,204 @@ def test_copyback_promotes_trace_manifest_and_blobs_before_headless_gc(tmp_path)
|
|||
] == "read_file"
|
||||
|
||||
|
||||
def test_pipeline_loop_outcome_trace_refs_are_rebased_and_readable_after_gc(tmp_path):
|
||||
from ouroboros.agent_task_pipeline import _store_task_result
|
||||
from ouroboros.outcomes import derive_loop_outcome
|
||||
|
||||
task_id = "phase3c-loop-outcome"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
repo = tmp_path / "repo"
|
||||
repo.mkdir()
|
||||
request = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="pipeline-request",
|
||||
call_type="llm_request",
|
||||
payload={"messages": [{"role": "user", "content": "exact pipeline prompt"}]},
|
||||
)
|
||||
response = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="pipeline-response",
|
||||
call_type="llm_response",
|
||||
payload={"message": {"role": "assistant", "content": "exact answer"}},
|
||||
)
|
||||
tool = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="pipeline-tool",
|
||||
call_type="tool_call",
|
||||
payload={"tool": "read_file", "result": "pipeline tool result"},
|
||||
)
|
||||
usage = {
|
||||
"execution_id": "phase3c-execution",
|
||||
"rounds": 1,
|
||||
"llm_call_refs": [{
|
||||
"llm_call_id": "phase3c-llm",
|
||||
"request_ref": request["manifest_ref"],
|
||||
"response_ref": response["manifest_ref"],
|
||||
}],
|
||||
}
|
||||
trace = {
|
||||
"tool_calls": [{
|
||||
"tool": "read_file",
|
||||
"tool_call_id": "pipeline-tool-call",
|
||||
"result": "pipeline tool result",
|
||||
"is_error": False,
|
||||
"trace_ref": tool,
|
||||
}],
|
||||
"reasoning_notes": [],
|
||||
}
|
||||
outcome = derive_loop_outcome("FINAL ANSWER: exact answer", usage, trace)
|
||||
_store_task_result(
|
||||
SimpleNamespace(drive_root=child, repo_dir=repo),
|
||||
{"id": task_id, "type": "task", "text": "pipeline task"},
|
||||
"FINAL ANSWER: exact answer",
|
||||
usage,
|
||||
trace,
|
||||
review_evidence={},
|
||||
loop_outcome=outcome,
|
||||
)
|
||||
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
|
||||
assert copied is not None
|
||||
nested_refs = copied["loop_outcome"]["trace_refs"]
|
||||
nested_request = nested_refs["llm_call_refs"][0]["request_ref"]
|
||||
nested_tool = nested_refs["tool_call_refs"][0]["manifest_ref"]
|
||||
assert pathlib.Path(nested_request["path"]).is_relative_to(parent / "observability")
|
||||
assert pathlib.Path(nested_tool["path"]).is_relative_to(parent / "observability")
|
||||
prune_headless_task_drives(parent, retention_days=0, now=_future_now())
|
||||
assert not child.exists()
|
||||
assert read_blob_ref(parent, _manifest(nested_request)["full_payload_ref"])[
|
||||
"messages"
|
||||
][0]["content"] == "exact pipeline prompt"
|
||||
assert read_blob_ref(parent, _manifest(nested_tool)["full_payload_ref"])[
|
||||
"result"
|
||||
] == "pipeline tool result"
|
||||
|
||||
|
||||
def test_real_truncated_tool_source_envelope_remains_actor_readable_after_gc(tmp_path):
|
||||
from ouroboros.agent_task_pipeline import _store_task_result
|
||||
from ouroboros.loop_tool_execution import process_tool_results
|
||||
from ouroboros.outcomes import derive_loop_outcome
|
||||
from ouroboros.tools.core import _read_file
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
||||
task_id = "phase3c-real-source"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
repo = tmp_path / "repo-real-source"
|
||||
repo.mkdir()
|
||||
ctx = ToolContext(repo_dir=repo, drive_root=child, task_id=task_id)
|
||||
messages: list[dict] = []
|
||||
trace = {"tool_calls": [], "reasoning_notes": []}
|
||||
exact_tail = "\nDECISIVE_SUFFIX=FAILED_AFTER_ONE_SHOT"
|
||||
full_result = "one-shot output\n" + ("x" * 100_000) + exact_tail
|
||||
process_tool_results(
|
||||
[{
|
||||
"fn_name": "run_command",
|
||||
"tool_call_id": "one-shot-call",
|
||||
"result": full_result,
|
||||
"is_error": False,
|
||||
"tool_args": {"cmd": "non-idempotent-operation"},
|
||||
"args_for_log": {"cmd": "non-idempotent-operation"},
|
||||
"trace_ref": {},
|
||||
"result_meta": {"status": "ok"},
|
||||
}],
|
||||
messages,
|
||||
trace,
|
||||
emit_progress=lambda _message: None,
|
||||
tools=SimpleNamespace(_ctx=ctx),
|
||||
)
|
||||
produced_ref = _source_ref_from_visible_result(messages[0]["content"])
|
||||
request = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="source-envelope-request",
|
||||
call_type="llm_request",
|
||||
payload={"messages": messages},
|
||||
)
|
||||
usage = {
|
||||
"execution_id": "source-envelope-execution",
|
||||
"rounds": 1,
|
||||
"llm_call_refs": [{
|
||||
"llm_call_id": "source-envelope-llm",
|
||||
"request_ref": request["manifest_ref"],
|
||||
}],
|
||||
}
|
||||
outcome = derive_loop_outcome("FINAL ANSWER: inspected", usage, trace)
|
||||
_store_task_result(
|
||||
SimpleNamespace(drive_root=child, repo_dir=repo),
|
||||
{"id": task_id, "type": "task", "text": "one-shot"},
|
||||
"FINAL ANSWER: inspected",
|
||||
usage,
|
||||
trace,
|
||||
review_evidence={},
|
||||
loop_outcome=outcome,
|
||||
)
|
||||
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
|
||||
assert copied is not None
|
||||
request_ref = copied["loop_outcome"]["trace_refs"]["llm_call_refs"][0][
|
||||
"request_ref"
|
||||
]
|
||||
prune_headless_task_drives(parent, retention_days=0, now=_future_now())
|
||||
payload = read_blob_ref(parent, _manifest(request_ref)["full_payload_ref"])
|
||||
promoted_ref = _source_ref_from_visible_result(payload["messages"][0]["content"])
|
||||
assert promoted_ref == produced_ref
|
||||
read_args = dict(promoted_ref["read"]["arguments"])
|
||||
read_args["start_char"] = 95_000
|
||||
canonical_ctx = ToolContext(repo_dir=repo, drive_root=parent, task_id=task_id)
|
||||
assert exact_tail in _read_file(canonical_ctx, **read_args)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mismatch", ["tool", "root", "path"])
|
||||
def test_task_source_read_contract_mismatch_is_typed_unavailable(tmp_path, mismatch):
|
||||
from ouroboros.artifacts import store_actor_source_bytes
|
||||
|
||||
task_id = f"phase3c-source-contract-{mismatch}"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
ref = store_actor_source_bytes(
|
||||
child,
|
||||
task_id,
|
||||
category="tool_results",
|
||||
source_id="contract",
|
||||
data=b"exact source",
|
||||
extension="txt",
|
||||
)
|
||||
malformed = json.loads(json.dumps(ref))
|
||||
if mismatch == "tool":
|
||||
malformed["read"]["tool"] = "run_command"
|
||||
elif mismatch == "root":
|
||||
malformed["read"]["arguments"]["root"] = "runtime_data"
|
||||
else:
|
||||
malformed["read"]["arguments"]["path"] = (
|
||||
"source_handles/tool_results/other.txt"
|
||||
)
|
||||
write_task_result(
|
||||
child,
|
||||
task_id,
|
||||
STATUS_COMPLETED,
|
||||
result="done",
|
||||
artifact_status="ready",
|
||||
review_evidence={"exact_source_ref": malformed},
|
||||
)
|
||||
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
|
||||
assert copied is not None
|
||||
gap = copied["review_evidence"]["exact_source_ref"]
|
||||
assert gap["availability"] == "unavailable"
|
||||
assert gap["reason"] == "invalid_ref"
|
||||
assert not (
|
||||
parent / "task_results" / "artifacts" / task_id / pathlib.Path(ref["path"])
|
||||
).exists()
|
||||
assert prune_headless_task_drives(
|
||||
parent, retention_days=0, now=_future_now()
|
||||
)["pruned"]
|
||||
|
||||
|
||||
def test_copyback_promotes_service_full_log_refs_in_durable_evidence_and_tool_payload(tmp_path):
|
||||
task_id = "phase3c-service"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
|
|
@ -233,19 +439,36 @@ def test_digest_mismatch_becomes_typed_unavailable_and_does_not_pin_drive(tmp_pa
|
|||
def test_concurrent_copyback_is_idempotent_and_copies_only_referenced_source_handle(
|
||||
tmp_path,
|
||||
):
|
||||
from ouroboros.artifacts import store_actor_source_bytes
|
||||
|
||||
task_id = "phase3c-source-handles"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
source_dir = child / "task_results" / "artifacts" / task_id / "source_handles" / "tool_results"
|
||||
source_dir.mkdir(parents=True)
|
||||
source_bytes = b"actor promised source"
|
||||
source_digest = hashlib.sha256(source_bytes).hexdigest()
|
||||
source = source_dir / ("tool-" + source_digest + ".txt")
|
||||
source.write_bytes(source_bytes)
|
||||
unreferenced_source_bytes = b"unreferenced source handle"
|
||||
unreferenced_source = source_dir / (
|
||||
"unused-" + hashlib.sha256(unreferenced_source_bytes).hexdigest() + ".txt"
|
||||
source_ref = store_actor_source_bytes(
|
||||
child,
|
||||
task_id,
|
||||
category="tool_results",
|
||||
source_id="tool",
|
||||
data=source_bytes,
|
||||
extension="txt",
|
||||
)
|
||||
source = child / "task_results" / "artifacts" / task_id / source_ref["path"]
|
||||
unreferenced_source_bytes = b"unreferenced source handle"
|
||||
unreferenced_source_ref = store_actor_source_bytes(
|
||||
child,
|
||||
task_id,
|
||||
category="tool_results",
|
||||
source_id="unused",
|
||||
data=unreferenced_source_bytes,
|
||||
extension="txt",
|
||||
)
|
||||
unreferenced_source = (
|
||||
child
|
||||
/ "task_results"
|
||||
/ "artifacts"
|
||||
/ task_id
|
||||
/ unreferenced_source_ref["path"]
|
||||
)
|
||||
unreferenced_source.write_bytes(unreferenced_source_bytes)
|
||||
unrelated = child / "task_results" / "artifacts" / task_id / "unrelated.txt"
|
||||
unrelated.write_text("must not copy", encoding="utf-8")
|
||||
unreferenced_blob = write_blob(child, {"unreferenced": True})
|
||||
|
|
@ -256,18 +479,6 @@ def test_concurrent_copyback_is_idempotent_and_copies_only_referenced_source_han
|
|||
call_type="tool_call",
|
||||
payload={"result": "copy once by content identity"},
|
||||
)
|
||||
source_ref = {
|
||||
"kind": "task_source",
|
||||
"root": "artifact_store",
|
||||
"path": f"source_handles/tool_results/{source.name}",
|
||||
"size": len(source_bytes),
|
||||
"sha256": source_digest,
|
||||
"read": {
|
||||
"tool": "read_file",
|
||||
"root": "artifact_store",
|
||||
"path": f"source_handles/tool_results/{source.name}",
|
||||
},
|
||||
}
|
||||
write_task_result(
|
||||
child,
|
||||
task_id,
|
||||
|
|
@ -337,3 +548,178 @@ def test_legacy_missing_child_ref_is_typed_gap_without_permanent_retention(tmp_p
|
|||
assert gap["reason"] == "source_missing"
|
||||
assert copied["child_ref_promotion"]["status"] == "complete"
|
||||
assert prune_headless_task_drives(parent, retention_days=0, now=_future_now())["pruned"]
|
||||
|
||||
|
||||
def test_startup_sweep_retries_only_pending_refs_then_prunes_without_manual_copyback(
|
||||
tmp_path, monkeypatch,
|
||||
):
|
||||
import ouroboros.observability as observability
|
||||
import server
|
||||
|
||||
task_id = "phase3c-startup-retry"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
first_trace = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="startup-retry-first",
|
||||
call_type="tool_call",
|
||||
payload={"result": "already promoted before interruption"},
|
||||
)
|
||||
pending_trace = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="startup-retry-pending",
|
||||
call_type="tool_call",
|
||||
payload={"result": "survive startup retry"},
|
||||
)
|
||||
write_task_result(
|
||||
child,
|
||||
task_id,
|
||||
STATUS_COMPLETED,
|
||||
result="stale child result",
|
||||
artifact_status="ready",
|
||||
artifacts=[{"kind": "stale_child_artifact", "path": "child-only"}],
|
||||
root_phase_checkpoint={"post_task_synthesis": "pending_once"},
|
||||
trace_refs={
|
||||
"tool_call_refs": [
|
||||
{"manifest_ref": first_trace["manifest_ref"]},
|
||||
{"manifest_ref": pending_trace["manifest_ref"]},
|
||||
]
|
||||
},
|
||||
)
|
||||
real = observability.promote_call_manifest_ref
|
||||
|
||||
def _interrupt_pending(*args, **kwargs):
|
||||
ref = args[2] if len(args) > 2 else kwargs.get("ref") or {}
|
||||
if ref.get("call_id") == "startup-retry-pending":
|
||||
raise OSError("first copy interrupted")
|
||||
return real(*args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(
|
||||
observability,
|
||||
"promote_call_manifest_ref",
|
||||
_interrupt_pending,
|
||||
)
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
assert copied is not None
|
||||
assert copied["child_ref_promotion"]["status"] == "incomplete"
|
||||
|
||||
canonical_artifact = {
|
||||
"kind": "canonical_newer_artifact",
|
||||
"path": str(parent / "canonical-newer.txt"),
|
||||
}
|
||||
write_task_result(
|
||||
parent,
|
||||
task_id,
|
||||
STATUS_COMPLETED,
|
||||
result="canonical newer result",
|
||||
artifact_status="ready_with_changes",
|
||||
artifact_bundle={"status": "ready_with_changes", "artifacts": [canonical_artifact]},
|
||||
artifacts=[canonical_artifact],
|
||||
artifact_finalized_at="2000-01-01T00:00:00+00:00",
|
||||
root_phase_checkpoint={
|
||||
"post_task_synthesis": "completed",
|
||||
"canonical_newer": True,
|
||||
},
|
||||
)
|
||||
monkeypatch.setattr(observability, "promote_call_manifest_ref", real)
|
||||
monkeypatch.setattr(server, "DATA_DIR", parent)
|
||||
monkeypatch.setenv("OUROBOROS_GC_RETENTION_DAYS", "1")
|
||||
|
||||
server._startup_prune_sweeps()
|
||||
|
||||
settled = load_task_result(parent, task_id) or {}
|
||||
assert settled["child_ref_promotion"]["status"] == "complete"
|
||||
assert settled["result"] == "canonical newer result"
|
||||
assert settled["artifact_status"] == "ready_with_changes"
|
||||
assert settled["artifacts"] == [canonical_artifact]
|
||||
assert settled["artifact_bundle"] == {
|
||||
"status": "ready_with_changes",
|
||||
"artifacts": [canonical_artifact],
|
||||
}
|
||||
assert settled["root_phase_checkpoint"] == {
|
||||
"post_task_synthesis": "completed",
|
||||
"canonical_newer": True,
|
||||
}
|
||||
assert not child.exists()
|
||||
promoted_first = settled["trace_refs"]["tool_call_refs"][0]["manifest_ref"]
|
||||
promoted_pending = settled["trace_refs"]["tool_call_refs"][1]["manifest_ref"]
|
||||
assert read_blob_ref(parent, _manifest(promoted_first)["full_payload_ref"])[
|
||||
"result"
|
||||
] == "already promoted before interruption"
|
||||
assert read_blob_ref(parent, _manifest(promoted_pending)["full_payload_ref"])[
|
||||
"result"
|
||||
] == "survive startup retry"
|
||||
|
||||
|
||||
def test_startup_prune_retries_missing_pending_source_into_typed_gap(
|
||||
tmp_path, monkeypatch,
|
||||
):
|
||||
import ouroboros.observability as observability
|
||||
|
||||
task_id = "phase3c-startup-missing"
|
||||
parent, child = _child(tmp_path, task_id)
|
||||
trace = persist_call(
|
||||
child,
|
||||
task_id=task_id,
|
||||
call_id="startup-missing",
|
||||
call_type="tool_call",
|
||||
payload={"result": "lost before retry"},
|
||||
)
|
||||
write_task_result(
|
||||
child,
|
||||
task_id,
|
||||
STATUS_COMPLETED,
|
||||
result="done",
|
||||
artifact_status="ready",
|
||||
trace_refs={"tool_call_refs": [{"manifest_ref": trace["manifest_ref"]}]},
|
||||
)
|
||||
real = observability.promote_call_manifest_ref
|
||||
monkeypatch.setattr(
|
||||
observability,
|
||||
"promote_call_manifest_ref",
|
||||
lambda *_args, **_kwargs: (_ for _ in ()).throw(OSError("interrupted")),
|
||||
)
|
||||
copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)})
|
||||
assert copied is not None
|
||||
assert copied["child_ref_promotion"]["status"] == "incomplete"
|
||||
pathlib.Path(trace["manifest_ref"]["path"]).unlink()
|
||||
monkeypatch.setattr(observability, "promote_call_manifest_ref", real)
|
||||
|
||||
report = prune_headless_task_drives(
|
||||
parent, retention_days=0, now=_future_now()
|
||||
)
|
||||
|
||||
assert report["promotion_retry"]["completed"] == [task_id]
|
||||
assert report["pruned"][0]["task_id"] == task_id
|
||||
settled = load_task_result(parent, task_id) or {}
|
||||
gap = settled["trace_refs"]["tool_call_refs"][0]["manifest_ref"]
|
||||
assert gap["availability"] == "unavailable"
|
||||
assert gap["reason"] == "source_missing"
|
||||
assert settled["child_ref_promotion"]["status"] == "complete"
|
||||
|
||||
|
||||
def test_periodic_maintenance_invokes_pending_ref_promotion_sweep(
|
||||
tmp_path, monkeypatch,
|
||||
):
|
||||
import ouroboros.observability as observability
|
||||
import server
|
||||
import supervisor.task_lifecycle as task_lifecycle
|
||||
import supervisor.terminal_delivery as terminal_delivery
|
||||
|
||||
calls: list[pathlib.Path] = []
|
||||
monkeypatch.setattr(server, "DATA_DIR", tmp_path)
|
||||
monkeypatch.setattr(server.time, "time", lambda: 10_000.0)
|
||||
monkeypatch.setattr(server, "_LAST_CANCEL_INTENT_SWEEP", [0.0])
|
||||
monkeypatch.setattr(task_lifecycle, "sweep_cancel_intents", lambda: {})
|
||||
monkeypatch.setattr(terminal_delivery, "replay_pending_deliveries", lambda _root: None)
|
||||
monkeypatch.setattr(
|
||||
observability,
|
||||
"retry_pending_child_ref_promotions",
|
||||
lambda root: calls.append(pathlib.Path(root)) or {},
|
||||
raising=False,
|
||||
)
|
||||
|
||||
server._periodic_supervisor_maintenance([10_000.0], [10_000.0])
|
||||
|
||||
assert calls == [tmp_path]
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue