fix(results): retain review and completion sources through shared handles

This commit is contained in:
Ouroboros 2026-09-06 18:15:05 +00:00
parent fee1967354
commit bff9f321a4
21 changed files with 305 additions and 105 deletions

View file

@ -1282,7 +1282,7 @@ The root post-task checkpoint decides whether an error-bearing or non-trivial ru
A reflection lands where it durably belongs: a non-project root appends the full entry to the canonical `logs/task_reflections.jsonl`; a project-scoped root appends the full entry to its project drive and the canonical log receives only a bounded pointer row — full project text never enters the canonical log, which feeds future global context. A project-bound task's context includes a bounded labeled tail of its own project's reflections; the headless mirror drive of a split root is never the reflection home, and the Pattern Register update stays canonical in both cases. Every entry carries task identity, evidence, lessons, backlog candidates, and validated memory actions. `MEMORY_ACTIONS_JSON` permits only `scratchpad_append`, `knowledge_write`, and `identity_update_candidate`, at bounded count and size. `apply_memory_actions` routes accepted actions through provenance-preserving memory and knowledge APIs. An `identity_update_candidate` is recorded in the scratchpad for review and is never auto-written to `identity.md`. For a project-scoped task, only project knowledge is written; scratchpad and identity actions are skipped so local facts cannot contaminate the canonical self. Reflection may propose a future campaign or backlog item, but it cannot enqueue, review, commit, or enable one.
Only the root runs full post-task synthesis once (`root_phase_checkpoint` makes the paid phase at-most-once across restart); children contribute evidence, never a second global synthesis. The owner's final answer does not wait for blocking synthesis: after the durable result is stored, the final `send_message` is delivered immediately while the buffered-return copy is RETAINED (queue.put is not a delivery receipt); both copies carry one `delivery_id` and the supervisor suppresses the duplicate through the durable registry in `supervisor/terminal_delivery.py`. The same file holds a bounded PENDING outbox — ONE seam for the normal, cancel, and reap terminal paths: a terminal answer is recorded as owed BEFORE it is enqueued and cleared in the same write that marks it delivered, so a crash between settle and send replays it instead of losing it. Every non-ephemeral root's answer enters this outbox at durable-result persistence time under the canonical `final:<tid>:<digest>` id. Replays back off and are bounded; an exhausted or capacity-evicted row is dropped LOUDLY (full text preserved on disk, typed `terminal_delivery_exhausted` event, chat notice) — external transports stay at-least-once and that residual is disclosed. `task_done` still goes last through the buffered return — an early `task_done` would release the queue slot and start child-drive cleanup while post-task still runs — so a worker reaped during a hung synthesis has already delivered the answer. Synthesis receives a sealed final package from the durable result: submitted final text, its artifact manifest, and completion_observations. The terminal writer preserves full redacted action observations in the canonical artifact store (`task.budget_drive_root or drive_root`) before compact publication. Registered completion sources remain downloadable bookkeeping, excluded from deliverables and inferred readiness, so child cleanup never owns their only copy. Their native reader selector uses `get_task_result(include_completion_source=true)` with the source task id and its canonical drive. It first returns complete character length/hash, then explicit `source_start_char`/`source_end_char` ranges; `artifacts.text_source_range_projection` shares the unchanged work-order range contract. Source bytes, kind, path containment and SHA are checked before any excerpt. Packet-only summary/reflection receive per-send-tool counts, each family's latest recorded return, and task-related skill readiness with coverage; full-source references are for later readers, not evidence the synthesizer has read. Positive observed facts correct error-trace impressions, while tool success does not prove owner receipt, empty material does not prove absence, and skill readiness does not attribute an owner's action to the task. Recovery uses these stored observations; task-summary requests/responses use chat_observed.
Only the root runs full post-task synthesis once (`root_phase_checkpoint` makes the paid phase at-most-once across restart); children contribute evidence, never a second global synthesis. The owner's final answer does not wait for blocking synthesis: after the durable result is stored, the final `send_message` is delivered immediately while the buffered-return copy is RETAINED (queue.put is not a delivery receipt); both copies carry one `delivery_id` and the supervisor suppresses the duplicate through the durable registry in `supervisor/terminal_delivery.py`. The same file holds a bounded PENDING outbox — ONE seam for the normal, cancel, and reap terminal paths: a terminal answer is recorded as owed BEFORE it is enqueued and cleared in the same write that marks it delivered, so a crash between settle and send replays it instead of losing it. Every non-ephemeral root's answer enters this outbox at durable-result persistence time under the canonical `final:<tid>:<digest>` id. Replays back off and are bounded; an exhausted or capacity-evicted row is dropped LOUDLY (full text preserved on disk, typed `terminal_delivery_exhausted` event, chat notice) — external transports stay at-least-once and that residual is disclosed. `task_done` still goes last through the buffered return — an early `task_done` would release the queue slot and start child-drive cleanup while post-task still runs — so a worker reaped during a hung synthesis has already delivered the answer. Synthesis receives a sealed final package from the durable result: submitted final text, its artifact manifest, and completion_observations. The terminal writer preserves full redacted action observations in the canonical artifact store (`task.budget_drive_root or drive_root`) before compact publication. Completion sources use the existing write-once `source_handles/context_checkpoints` store with verified `task_source` refs; the published-ref closure includes completion observations before child cleanup. Sources stay outside deliverables and inferred readiness; legacy flat references remain readable. Their native reader selector uses `get_task_result(include_completion_source=true)` with the source task id and its canonical drive. It first returns complete character length/hash, then explicit `source_start_char`/`source_end_char` ranges; `artifacts.text_source_range_projection` shares the unchanged work-order range contract. Source bytes, kind, path containment and SHA are checked before any excerpt. Packet-only summary/reflection receive per-send-tool counts, each family's latest recorded return, and task-related skill readiness with coverage; full-source references are for later readers, not evidence the synthesizer has read. Positive observed facts correct error-trace impressions, while tool success does not prove owner receipt, empty material does not prove absence, and skill readiness does not attribute an owner's action to the task. Recovery uses these stored observations; task-summary requests/responses use chat_observed.
#### Durable memory and project focus
@ -1632,18 +1632,21 @@ otherwise the view is partial and the consumer remains non-final or abstains.
| Terminal task/project memory | Root terminal result plus existing task/project summary producers | Cognitive Main terminal summaries and the two Project-root UI lifecycle rows (started + terminal completion) | Task-result ID, project binding, and summary/source refs | Summary is a biography projection, not raw evidence. Terminal outcomes, including failed/cancelled/degraded, remain retained through their canonical result owner. |
| Background Consciousness observations | `data/state/consciousness_observations.jsonl`, append-only enqueue/ACK rows owned by `BackgroundConsciousness` | Pending count/oldest metadata and a bounded recent observation rendering | `read_file(root='runtime_data', path='state/consciousness_observations.jsonl')` | Unacknowledged rows survive restart/overflow/error. Gaps block ACK and the existing direct identity rewrite; only a settled successful cycle appends ACK. |
| Plan/review authority | Exact task-artifact/observability wave bodies, evidence selectors, reviewer route/thread receipts, and the bounded review hot index | Review status, latest wave, obligations, and compact findings; a predecessor's inherited `plan_review_state` is first projected to a compact authority core ordered around the newest wave's identity, acceptance claims, findings, and dispositions, with reviewer transport removed and `need_evidence_seen` last-priority. Every bounded collection names its total and omitted count; the projection discloses `full_chars` plus `source_ref`, and the named `include_authority` source stays complete | Exact artifact/source handle plus SHA/range/thread selectors | Missing or partial evidence is `DEGRADED`/`NOT_RUN`, never PASS. Exact artifacts remain bound to the reviewed candidate SHA; hot indexes may rotate only after the source is retained. |
| Task acceptance (three deliveries) | The FULL host packet (`review_evidence.build_task_acceptance_evidence` under the host ladder, with its `__provenance__` table), the applied host run retained in the canonical task artifact store, and the paid-identity wallet ledger | The per-delivery work order: the api pack for a packet row; the FULL packet plus absolute pointers and the access disclosure for an agent-session row; the packet without its freely degradable tail plus the real data root for a native inspection row (`loop_acceptance_review.acceptance_retrieving_work_order`) | Exact `evidence_refs` from the packet's enumerable exhibit vocabulary; absolute pointers to the task's active workspace, task result record, artifact directory, verification receipts and tool-trajectory log; `review_projection.panels[].applied_source_ref` for the complete redacted applied review | Refs resolve against the FULL packet only, never the rendered projection; a session's reads are unobserved by the host (disclosed), a native episode's are `host_observed`; the immutable-core overflow refuses every delivery, a partial tool-result projection only packet rows; one strict wallet claim per panel whatever the rows' deliveries (owner R11). |
| Task acceptance (three deliveries) | The FULL host packet (`review_evidence.build_task_acceptance_evidence` under the host ladder, with its `__provenance__` table), the applied host run retained through canonical task source handles, and the paid-identity wallet ledger | The per-delivery work order: the api pack for a packet row; the FULL packet plus absolute pointers and the access disclosure for an agent-session row; the packet without its freely degradable tail plus the real data root for a native inspection row (`loop_acceptance_review.acceptance_retrieving_work_order`) | Exact `evidence_refs` from the packet's enumerable exhibit vocabulary; absolute pointers to the task's active workspace, task result record, artifact directory, verification receipts and tool-trajectory log; `review_projection.panels[].applied_source_ref` for the complete redacted applied review | Refs resolve against the FULL packet only, never the rendered projection; a session's reads are unobserved by the host (disclosed), a native episode's are `host_observed`; the immutable-core overflow refuses every delivery, a partial tool-result projection only packet rows; one strict wallet claim per panel whatever the rows' deliveries (owner R11). |
| Canonical versus execution roots | Canonical budget/data root owns identity, authority, biography, results, and promoted observability; execution drives own tools, workspace, and transient trajectory | Project/fork/task lenses and status projections | Existing canonical-root resolver, task-result pointers, and source handles | A fork is an execution lens, not a second mind. Copy-back/promotion precedes GC for anything referenced by a canonical result; missing legacy bytes become an explicit gap. |
---
Acceptance source identity is computed before history-dependent packet budgeting.
The complete receipt/tool sources and work artifacts still invalidate the binding
when their facts change; `artifacts.is_task_bookkeeping_artifact` identifies registered
host acceptance-review and completion-observation sources. Collection and `project_deliverable_artifacts` exclude
them from user artifacts and inferred readiness, including effective reads and copy-back.
Their canonical bytes remain available through the existing flat artifact download;
that endpoint explicitly includes registered bookkeeping for its exact-name fallback.
when their facts change. Applied reviews and completion observations use the existing
write-once `source_handles/context_checkpoints` store and verified `task_source` refs;
source handles stay outside both deliverables and the acceptance artifact manifest.
`load_task_result` projects historical flat bookkeeping out without rewriting stored
bytes; raw-result/outcome boundaries retain the same compatibility projection.
The existing artifact route selects a published ref through its `source` query and
returns digest-verified bytes; legacy flat URLs keep their exact-name fallback.
`api_client.taskSourceDownloadUrl` owns the shared browser URL contract.
The final packet size includes source references and omission notes.
`review_projection.publish_acceptance_checkpoint` saves full applied host records
before updating the compact task-result field through `write_task_result` and

View file

@ -2185,16 +2185,17 @@ by "Provider Independence" above. Call-site imperatives:
budgeting; recording a review must not change the facts it reviewed. Complete
applied host records (including resolved criteria, decisions and supersession)
are saved by `review_projection.publish_acceptance_checkpoint` through the
existing artifact store before compact publication. Artifact registration uses
a short locked manifest merge shared by all three writers; copying/hashing
finishes before the lock, which never acquires a task-result lock. The live
existing write-once source handles before compact publication. Ordinary artifact
registration keeps its short locked manifest merge; copying/hashing finishes
before that lock, which never acquires a task-result lock. The live
publication changes only `review_projection`, preserving lifecycle and other
writers' fields. Test delayed snapshots and child replicas through the same
central merge, and verify that the full source downloads while the task is
still running. Registered review bookkeeping remains downloadable without entering
user deliverables or making an otherwise artifact-free task ready. Completion
observations use the same classifier only after canonical-first full-source
persistence (`task.budget_drive_root or drive_root`). Give native readers an
still running. Review/completion sources use `source_handles/context_checkpoints`,
outside deliverables and the acceptance artifact manifest. Canonical-first
persistence (`task.budget_drive_root or drive_root`) and the existing published-ref
closure preserve them before child cleanup. Legacy flat refs and their read-only
bookkeeping projection remain compatible; do not migrate live task data. Give native readers an
executable get_task_result selector for that source task; reuse explicit character
ranges and complete-source hashes so a split or later task never resolves the
basename against its own artifact directory. Preserve that

View file

@ -18,7 +18,7 @@ from ouroboros.task_results import (
load_task_result,
write_task_result,
)
from ouroboros.artifacts import collect_task_artifact_records, merge_artifact_records, project_deliverable_artifacts
from ouroboros.artifacts import collect_task_artifact_records, merge_artifact_records
from ouroboros.outcomes import (
EXECUTION_BEST_EFFORT,
EXECUTION_FAILED,
@ -837,7 +837,7 @@ def _store_task_result(env: Any, task: Dict[str, Any], text: str,
"reserved_usd": None, "unresolved_upper_bound_usd": None,
"unknown_unmetered": None,
})
existing = project_deliverable_artifacts(load_task_result(env.drive_root, str(task.get("id") or "")) or {})
existing = load_task_result(env.drive_root, str(task.get("id") or "")) or {}
if loop_outcome is None:
loop_outcome = _derive_host_bound_loop_outcome(env, task, text, usage, llm_trace)
# Apply FR3 before normalization so the persisted axes and ledger agree.

View file

@ -627,12 +627,17 @@ def read_actor_source_bytes(
) -> bytes:
"""Resolve and verify one task-local actor source ref or raise explicitly."""
if not isinstance(ref, dict) or ref.get("kind") != "task_source":
legacy = isinstance(ref, dict) and ref.get("kind") in {"task_acceptance_review", "task_completion_observations"}
if not isinstance(ref, dict) or (ref.get("kind") != "task_source" and not legacy):
raise ValueError("actor source ref has an unexpected kind")
if ref.get("root") != "artifact_store":
raise ValueError("actor source ref has an unexpected root")
rel = pathlib.PurePosixPath(str(ref.get("path") or ""))
if not rel.parts or rel.parts[0] != _SOURCE_HANDLES_SUBDIR or rel.is_absolute():
valid_path = (
len(rel.parts) == 1 and rel.name not in {".", ".."} and "\\" not in str(ref.get("path") or "")
if legacy else bool(rel.parts and rel.parts[0] == _SOURCE_HANDLES_SUBDIR)
)
if not valid_path or rel.is_absolute():
raise ValueError("actor source ref has an invalid path")
base = task_artifact_dir_path(drive_root, task_id, create=False).resolve(strict=False)
target = base.joinpath(*rel.parts)
@ -647,7 +652,7 @@ def read_actor_source_bytes(
raise ValueError("actor source ref escapes its task artifact root") from exc
raw = target.read_bytes()
try:
expected_size = int(ref["size"])
expected_size = int(ref["bytes" if legacy else "size"])
except (KeyError, TypeError, ValueError) as exc:
raise ValueError("actor source ref has no valid size") from exc
if len(raw) != expected_size:
@ -657,6 +662,22 @@ def read_actor_source_bytes(
return raw
def read_task_result_source_bytes(
drive_root: Any, result: Dict[str, Any], name: str, source_path: str,
) -> bytes:
"""Read an exact published source ref, never a caller-selected filesystem path."""
review = result.get("review_projection")
panels = review.get("panels") if isinstance(review, dict) else []
refs = [row.get("applied_source_ref") for row in (panels if isinstance(panels, list) else []) if isinstance(row, dict)]
observations = result.get("completion_observations")
if isinstance(observations, dict):
refs.append(observations.get("source_ref"))
for ref in refs:
if isinstance(ref, dict) and ref.get("path") == source_path and pathlib.PurePosixPath(source_path).name == name:
return read_actor_source_bytes(drive_root, validate_task_id(result.get("task_id")), ref)
raise ValueError("the requested source is not published by this task result")
def persist_exact_text_source(
drive_root: Union[pathlib.Path, str], task_id: str, *,
source_id: str, text: str,

View file

@ -13,7 +13,7 @@ from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
from starlette.requests import Request
from starlette.responses import FileResponse, JSONResponse
from starlette.responses import FileResponse, JSONResponse, Response
from ouroboros.gateway._helpers import coerce_int, json_error, json_exception, request_drive_root, request_json_or, request_repo_dir, stage_initial_task_attachments
from ouroboros.gateway.contracts import TaskCreateRequest
@ -54,7 +54,7 @@ from ouroboros.contracts.task_contract import (
normalize_resource_policy,
)
from ouroboros.outcomes import public_task_result
from ouroboros.artifacts import resolve_chat_media_path
from ouroboros import artifacts as artifact_store
from ouroboros.task_result_schema import (
emit_quarantine_event,
quarantine_task_result,
@ -990,30 +990,32 @@ async def api_task_artifact(request: Request):
if not name or "/" in name or "\\" in name or name in {".", ".."} or ".." in pathlib.PurePosixPath(name).parts:
return json_error("artifact name must be a simple filename", 400)
drive_root = request_drive_root(request)
chat_media = resolve_chat_media_path(drive_root, task_id, name)
if chat_media is not None:
return FileResponse(chat_media)
result = load_effective_task_result(drive_root, task_id)
if not result:
return json_error("task not found", 404)
artifact = _artifact_by_name(result, name)
if artifact is None:
from ouroboros.artifacts import collect_task_artifact_records, is_task_bookkeeping_artifact
records = collect_task_artifact_records(drive_root, task_id, include_bookkeeping=True)
artifact = _artifact_by_name({"artifacts": [row for row in records if is_task_bookkeeping_artifact(row)]}, name)
if artifact is None:
return json_error("artifact not found", 404, task_id=task_id, artifact=name)
base = task_artifacts_dir(drive_root, task_id).resolve(strict=False)
path = pathlib.Path(str(artifact.get("path") or "")).resolve(strict=False)
if path.name != name:
return json_error("artifact metadata path does not match requested name", 500)
try:
path.relative_to(base)
except ValueError:
return json_error("artifact path is outside task artifact directory", 500)
if not path.is_file():
return json_error("artifact file is missing", 404, task_id=task_id, artifact=name)
path = artifact_store.resolve_chat_media_path(drive_root, task_id, name)
if path is None:
result = load_effective_task_result(drive_root, task_id)
if not result:
return json_error("task not found", 404)
if source := request.query_params.get("source"):
try:
return Response(artifact_store.read_task_result_source_bytes(drive_root, result, name, source), media_type="application/json")
except (OSError, ValueError, RuntimeError):
return json_error("task source is unavailable or does not match its recorded identity", 404)
artifact = _artifact_by_name(result, name)
if artifact is None:
records = artifact_store.collect_task_artifact_records(drive_root, task_id, include_bookkeeping=True)
artifact = _artifact_by_name({"artifacts": [row for row in records if artifact_store.is_task_bookkeeping_artifact(row)]}, name)
if artifact is None:
return json_error("artifact not found", 404, task_id=task_id, artifact=name)
base = task_artifacts_dir(drive_root, task_id).resolve(strict=False)
path = pathlib.Path(str(artifact.get("path") or "")).resolve(strict=False)
if path.name != name:
return json_error("artifact metadata path does not match requested name", 500)
try:
path.relative_to(base)
except ValueError:
return json_error("artifact path is outside task artifact directory", 500)
if not path.is_file():
return json_error("artifact file is missing", 404, task_id=task_id, artifact=name)
return FileResponse(path)

View file

@ -427,12 +427,10 @@ def remove_subagent_task_drive(parent_drive_root: pathlib.Path, task_id: str) ->
def copy_child_task_result(parent_drive_root: pathlib.Path, task: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""Copy a child-drive task result back to the parent data root."""
from ouroboros.artifacts import project_deliverable_artifacts
task_id = str(task.get("id") or "")
if not task_id:
return None
canonical_existing = project_deliverable_artifacts(load_task_result(parent_drive_root, task_id) or {})
canonical_existing = load_task_result(parent_drive_root, task_id) or {}
# Cancellation is authoritative before any child-root read or artifact copy.
if cancellation_blocks_child_result(canonical_existing):
return canonical_existing
@ -442,7 +440,6 @@ def copy_child_task_result(parent_drive_root: pathlib.Path, task: Dict[str, Any]
child_result = load_task_result(child_drive, task_id)
if not isinstance(child_result, dict):
return None
child_result = project_deliverable_artifacts(child_result)
# Canonical terminal results must not retain child-local forensic/source
# pointers that the startup GC is about to delete. Promotion happens before
# the one atomic canonical result write; a live ref whose destination write
@ -535,14 +532,11 @@ def _copy_child_artifacts_to_parent(
) -> List[Dict[str, Any]]:
"""Rebase child-drive artifact files into the parent task artifact store."""
from ouroboros.artifacts import is_task_bookkeeping_artifact
from ouroboros.outcome_receipt_store import is_verification_receipts_path
parent_dir = task_artifacts_dir(parent_drive_root, task_id)
rebased: List[Dict[str, Any]] = []
for artifact in artifacts:
if is_task_bookkeeping_artifact(artifact):
continue
item = dict(artifact)
raw_path = str(item.get("path") or "").strip()
if not raw_path:
@ -921,9 +915,7 @@ def _merge_artifacts(
*,
drop_kinds: Optional[set[str]] = None,
) -> List[Dict[str, Any]]:
from ouroboros.artifacts import is_task_bookkeeping_artifact
new_items = [item for item in new_items if not is_task_bookkeeping_artifact(item)]
merged: List[Dict[str, Any]] = []
drop = drop_kinds or set()
key_for = lambda item: (
@ -932,7 +924,7 @@ def _merge_artifacts(
)
keys = {key_for(item) for item in new_items if isinstance(item, dict)}
for item in existing:
if not isinstance(item, dict) or is_task_bookkeeping_artifact(item):
if not isinstance(item, dict):
continue
key = key_for(item)
if key[0] not in drop and key not in keys:

View file

@ -476,6 +476,7 @@ _PUBLISHED_CHILD_REF_FIELDS = frozenset(
"loop_outcome",
"review_evidence",
"review_projection",
"completion_observations",
"verification_ledger",
"root_phase_checkpoint",
"plan_review_state",

View file

@ -661,7 +661,7 @@ def _accept_trajectory(tool_calls: list, drive_root: Any = None, task_id: str =
def _accept_artifact_manifest(drive_root: Any, task_id: str, protected: set) -> list:
"""Return a leak-safe manifest; protected, large and binary artifacts stay manifest-only."""
from ouroboros.task_results import validate_task_id
from ouroboros.artifacts import _ARTIFACT_MANIFEST, is_task_bookkeeping_artifact
from ouroboros.artifacts import _ARTIFACT_MANIFEST, _SOURCE_HANDLES_SUBDIR, is_task_bookkeeping_artifact
from ouroboros.utils import read_json_dict
out: list = []
@ -686,7 +686,7 @@ def _accept_artifact_manifest(drive_root: Any, task_id: str, protected: set) ->
except OSError:
continue
rel = str(p.relative_to(base))
if rel in {_ARTIFACT_MANIFEST, _ARTIFACT_MANIFEST + ".lock"} or rel in review_sources:
if p.relative_to(base).parts[0] == _SOURCE_HANDLES_SUBDIR or rel in {_ARTIFACT_MANIFEST, _ARTIFACT_MANIFEST + ".lock"} or rel in review_sources:
continue # host bookkeeping is available through its own review refs
entry: Dict[str, Any] = {"name": rel, "size": size, "provenance": "artifact"}
# Match protected paths by artifact path, prefix, or basename.

View file

@ -331,7 +331,7 @@ def publish_acceptance_checkpoint(
"""
from pathlib import Path
from ouroboros.artifacts import store_task_artifact_bytes
from ouroboros.artifacts import store_actor_source_bytes
from ouroboros.task_results import write_task_result
from ouroboros.tools.plan_review_references import _emit_review_reference
@ -359,9 +359,9 @@ def publish_acceptance_checkpoint(
run.pop("applied_source_ref", None)
run["applied_source_status"] = "unavailable"
try:
run["applied_source_ref"] = store_task_artifact_bytes(
Path(root), task_id, f"acceptance-{hashlib.sha256(raw).hexdigest()}.json",
raw, kind="task_acceptance_review",
run["applied_source_ref"] = store_actor_source_bytes(
Path(root), task_id, category="context_checkpoints",
source_id="acceptance", data=raw, extension="json",
)
run["applied_source_status"] = "available"
except (OSError, ValueError, TimeoutError):

View file

@ -267,7 +267,7 @@ def build_completion_observations(drive_root: Any, task: Dict[str, Any], trace:
so an earlier photo is not hidden by a tail of later text sends. Skill state
is the existing task-related readiness projection, not owner-click authorship.
"""
from ouroboros.artifacts import store_task_artifact_bytes
from ouroboros.artifacts import store_actor_source_bytes
from ouroboros.observability import redact_projection
from ouroboros.skill_readiness import acceptance_skill_lifecycle
from ouroboros.tool_capabilities import OWNER_DELIVERY_TOOL_NAMES
@ -312,9 +312,9 @@ def build_completion_observations(drive_root: Any, task: Dict[str, Any], trace:
return projection # no full action record to store for an ordinary empty turn
raw = json.dumps(snapshot, ensure_ascii=False, sort_keys=True, default=str).encode("utf-8")
try:
projection["source_ref"] = store_task_artifact_bytes(
task.get("budget_drive_root") or drive_root, str(task.get("id") or ""), f"completion-{hashlib.sha256(raw).hexdigest()}.json",
raw, kind="task_completion_observations",
projection["source_ref"] = store_actor_source_bytes(
task.get("budget_drive_root") or drive_root, str(task.get("id") or ""),
category="context_checkpoints", source_id="completion", data=raw, extension="json",
)
projection["source_ref"]["reader"] = {
"tool": "get_task_result",
@ -331,26 +331,20 @@ def completion_source_projection(
drive_root: Any, task_id: str, result: Dict[str, Any], start_char: Any = None, end_char: Any = None,
) -> Dict[str, Any]:
"""Read the selected task's complete stored observations through its canonical root."""
from ouroboros.artifacts import task_artifact_dir_path, text_source_range_projection
from ouroboros.artifacts import read_actor_source_bytes, text_source_range_projection
unavailable = {"schema": 1, "kind": "task_completion_observations", "status": "unavailable"}
observations = result.get("completion_observations")
ref = observations.get("source_ref") if isinstance(observations, dict) else None
if not isinstance(ref, dict) or ref.get("root") != "artifact_store" or ref.get("kind") != unavailable["kind"]:
if not isinstance(ref, dict) or ref.get("kind") not in {"task_source", unavailable["kind"]}:
return {**unavailable, "reason": "source_unavailable"}
name = ref.get("path")
if not isinstance(name, str) or not name or pathlib.Path(name).name != name or "\\" in name or name in {".", ".."}:
return {**unavailable, "reason": "source_ref_invalid"}
try:
root = task_artifact_dir_path(drive_root, task_id).resolve()
path = (root / name).resolve()
if not path.is_relative_to(root):
return {**unavailable, "reason": "source_ref_invalid"}
raw = path.read_bytes()
if len(raw) != ref.get("bytes") or hashlib.sha256(raw).hexdigest() != ref.get("sha256"):
return {**unavailable, "reason": "source_identity_mismatch"}
raw = read_actor_source_bytes(drive_root, str(result.get("task_id") or task_id), ref)
projection, reason = text_source_range_projection(raw.decode("utf-8"), unavailable["kind"], start_char, end_char)
except (OSError, ValueError, RuntimeError):
except ValueError as exc:
reason = "source_identity_mismatch" if "verification" in str(exc) else "source_ref_invalid"
return {**unavailable, "reason": reason}
except (OSError, RuntimeError):
return {**unavailable, "reason": "source_unavailable"}
payload = projection or unavailable
return {**payload, **({"reason": reason} if reason else {})}

View file

@ -749,7 +749,9 @@ def load_task_result(
``task_result_schema_refusal`` — is QUARANTINED by the fail-soft path
(moved under ``task_results/quarantine/``, one batched durable event) and
the read reports no result; the strict path raises WITHOUT moving, so an
authority probe never mutates storage.
authority probe never mutates storage. Historical review-source artifacts
are projected out of deliverables here without rewriting their saved bytes;
lifecycle, cost, review references and independent artifact states remain.
"""
try:
tid = validate_task_id(task_id)
@ -773,17 +775,21 @@ def load_task_result(
outcome = _quarantine_task_result(path, refusal)
if outcome == "kept_admissible":
data = read_json_dict(path)
return data if not task_result_schema_refusal(data) else None
if outcome == "moved":
_emit_quarantine_event(drive_root, [{"task_id": tid, "reason": refusal}])
return None
if task_result_schema_refusal(data):
return None
else:
if outcome == "moved":
_emit_quarantine_event(drive_root, [{"task_id": tid, "reason": refusal}])
return None
if strict and (
str(data.get("task_id") or "") != tid
or not isinstance(data.get("status"), str)
or not str(data.get("status") or "").strip()
):
raise ValueError(f"task result authority is unreadable or invalid: {path}")
return data
from ouroboros.artifacts import project_deliverable_artifacts
return project_deliverable_artifacts(data)
def list_task_results(

View file

@ -736,7 +736,7 @@ def effective_task_result(
if merged_is_workspace and child_status in FINAL_STATUSES and (parent_status not in {STATUS_FAILED, STATUS_CANCELLED, STATUS_REJECTED_DUPLICATE} or copied_child_terminal):
merged = _normalize_workspace_artifact_status(merged)
merged = _normalize_workspace_artifact_status(project_deliverable_artifacts(merged))
merged = _normalize_workspace_artifact_status(merged)
parent_status = str(merged.get("status") or "").lower()
if parent_status not in FINAL_STATUSES:

View file

@ -4,6 +4,7 @@ import copy
import hashlib
import shutil
from types import SimpleNamespace
from pathlib import Path
import pytest
from starlette.applications import Starlette
@ -53,10 +54,11 @@ def _download(root, stored):
app = Starlette(routes=[Route("/api/tasks/{task_id}/artifacts/{name}", api_task_artifact)])
app.state.drive_root = root
with TestClient(app) as client:
response = client.get(f"/api/tasks/applied/artifacts/{ref['path']}")
response = client.get(f"/api/tasks/applied/artifacts/{Path(ref['path']).name}",
params={"source": ref["path"]} if ref["kind"] == "task_source" else {})
assert response.status_code == 200
assert hashlib.sha256(response.content).hexdigest() == ref["sha256"]
assert len(response.content) == ref["bytes"]
assert len(response.content) == ref.get("size", ref.get("bytes"))
assert len(response.json()["actors"][0]["parsed"]["findings"]) == 80
for name in (artifacts._ARTIFACT_MANIFEST, artifacts._ARTIFACT_MANIFEST + ".lock", "verification_receipts.jsonl"):
assert client.get(f"/api/tasks/applied/artifacts/{name}").status_code == 404

View file

@ -7,6 +7,7 @@ import queue
import threading
from concurrent.futures import ThreadPoolExecutor
from types import SimpleNamespace
from pathlib import Path
import pytest
@ -42,7 +43,7 @@ def _source(root, panel):
path = artifacts.task_artifact_dir_path(root, "applied") / ref["path"]
raw = path.read_bytes()
assert hashlib.sha256(raw).hexdigest() == ref["sha256"]
assert len(raw) == ref["bytes"]
assert len(raw) == ref["size"]
return json.loads(raw)
@ -68,7 +69,9 @@ def test_full_applied_source_downloads_while_task_is_running(tmp_path):
app = Starlette(routes=[Route("/api/tasks/{task_id}/artifacts/{name}", api_task_artifact)])
app.state.drive_root = tmp_path
with TestClient(app) as client:
response = client.get(f"/api/tasks/applied/artifacts/{panel['applied_source_ref']['path']}")
ref = panel["applied_source_ref"]
response = client.get(f"/api/tasks/applied/artifacts/{Path(ref['path']).name}",
params={"source": ref["path"]})
assert response.status_code == 200
assert response.json() == full
assert client.get("/api/tasks/applied/artifacts/missing.json").status_code == 404
@ -84,7 +87,7 @@ def test_delayed_publication_cannot_replace_supersession_or_terminal_fields(tmp_
ctx = _context(tmp_path)
trace = {"review_runs": [_run()]}
reached, release = threading.Event(), threading.Event()
actual_store = artifacts.store_task_artifact_bytes
actual_store = artifacts.store_actor_source_bytes
first_thread = []
def delayed_store(*args, **kwargs):
@ -94,7 +97,7 @@ def test_delayed_publication_cannot_replace_supersession_or_terminal_fields(tmp_
assert release.wait(10)
return actual_store(*args, **kwargs)
monkeypatch.setattr(artifacts, "store_task_artifact_bytes", delayed_store)
monkeypatch.setattr(artifacts, "store_actor_source_bytes", delayed_store)
with ThreadPoolExecutor(max_workers=1) as executor:
older = executor.submit(review_projection.publish_acceptance_checkpoint, ctx, trace)
try:
@ -162,7 +165,7 @@ def test_publishing_legacy_run_does_not_invent_its_task_attempt(tmp_path):
def test_source_failure_discloses_unavailable_without_changing_verdict(tmp_path, monkeypatch):
ctx, trace = _context(tmp_path), {"review_runs": [_run()]}
monkeypatch.setattr(artifacts, "store_task_artifact_bytes", lambda *a, **k: (_ for _ in ()).throw(OSError("disk unavailable")))
monkeypatch.setattr(artifacts, "store_actor_source_bytes", lambda *a, **k: (_ for _ in ()).throw(OSError("disk unavailable")))
review_projection.publish_acceptance_checkpoint(ctx, trace)
panel = load_task_result(tmp_path, "applied")["review_projection"]["panels"][0]
assert panel["aggregate_signal"] == "PASS"

View file

@ -89,3 +89,13 @@ def test_unavoidable_overflow_reports_the_complete_final_packet_size():
assert overflow["budget_chars"] == 1000
assert overflow["packet_chars"] == len(json.dumps(packet, ensure_ascii=False))
assert overflow["packet_chars"] > 1000
def test_source_handles_do_not_change_task_work_identity(review_context):
from ouroboros.artifacts import store_actor_source_bytes
before = review_evidence.task_acceptance_evidence_revision(loop._build_host_acceptance_evidence(review_context))
for category in ["context_checkpoints", "tool_results"]:
store_actor_source_bytes(review_context.drive_root, "identity", category=category,
source_id="recorded-review", data=b"complete review source", extension="txt")
after = review_evidence.task_acceptance_evidence_revision(loop._build_host_acceptance_evidence(review_context))
assert after == before

View file

@ -5,6 +5,7 @@ import gzip
import json
import shutil
from types import SimpleNamespace
from pathlib import Path
import pytest
from starlette.applications import Starlette
@ -74,7 +75,8 @@ def test_sealed_observations_survive_copyback_and_recovery_without_global_attrib
app = Starlette(routes=[Route("/api/tasks/{task_id}/artifacts/{name}", api_task_artifact)])
app.state.drive_root = canonical
with TestClient(app) as client:
response = client.get(f"/api/tasks/observation/artifacts/{ref['path']}")
response = client.get(f"/api/tasks/observation/artifacts/{Path(ref['path']).name}",
params={"source": ref["path"]})
assert response.status_code == 200 and response.content == raw
download() # canonical custody exists before any child artifact copy-back
if source_root != canonical:
@ -108,7 +110,7 @@ def test_full_source_failure_is_disclosed_without_inventing_a_ref(tmp_path, monk
monkeypatch.setattr("ouroboros.skill_readiness.acceptance_skill_lifecycle", lambda *_a, **_k: [])
def fail(*args, **kwargs):
raise failure("source unavailable")
monkeypatch.setattr("ouroboros.artifacts.store_task_artifact_bytes", fail)
monkeypatch.setattr("ouroboros.artifacts.store_actor_source_bytes", fail)
trace = {"tool_calls": [{"tool": "send_photo", "status": "error", "is_error": True,
"result_partial": True, "result": "submission incomplete"}]}
result = build_completion_observations(tmp_path, {"id": "partial"}, trace)

View file

@ -85,7 +85,7 @@ def test_completion_reader_does_not_claim_a_different_or_missing_source(tmp_path
if fault == 'digest':
ref['sha256'] = '0' * 64
elif fault == 'bytes':
ref['bytes'] += 1
ref['size'] += 1
elif fault == 'traversal':
ref['path'] = '../outside.json'
else:
@ -97,3 +97,15 @@ def test_completion_reader_does_not_claim_a_different_or_missing_source(tmp_path
}))['completion_source']
assert result['status'] == 'unavailable' and 'text' not in result
assert result['reason'] in {'source_unavailable', 'source_identity_mismatch', 'source_ref_invalid'}
def test_completion_source_follows_the_effective_retry_identity(tmp_path, monkeypatch):
reg, ref, path, canonical, _execution, _row = source_reader(tmp_path, monkeypatch, 'split_readonly')
write_task_result(canonical, 'original-task', 'interrupted', superseded_by='source-task')
text = path.read_text(encoding='utf-8')
result = json.loads(reg.execute('get_task_result', {
'task_id': 'original-task', 'include_completion_source': True,
'source_start_char': 0, 'source_end_char': len(text),
}))['completion_source']
assert result['complete_sha256'] == ref['sha256']
assert result['text'] == text

View file

@ -0,0 +1,125 @@
"""Review/completion handles retain exact bytes through publication and child cleanup."""
import hashlib
import json
from pathlib import Path
import pytest
from starlette.applications import Starlette
from starlette.routing import Route
from starlette.testclient import TestClient
from ouroboros import artifacts, review_projection
from ouroboros.gateway.tasks import api_task_artifact
from ouroboros.headless import copy_child_task_result, prepare_task_drive, remove_subagent_task_drive
from ouroboros.task_finalization import completion_source_projection
from ouroboros.task_results import load_task_result, task_result_path, write_task_result
from tests.test_acceptance_publication import _context, _run
def _field(ref, source):
if source == "review":
return {"review_projection": {"panels": [{"surface": "task_acceptance", "applied_source_ref": ref}]}}
return {"completion_observations": {"source_ref": ref, "source_status": "available"}}
@pytest.mark.parametrize("source", ["review", "completion"])
def test_child_source_closure_survives_real_cleanup(tmp_path, source):
parent = tmp_path / "canonical"
child = prepare_task_drive(parent, "source", "empty")
raw = b'{"full":"retained evidence"}'
ref = artifacts.store_actor_source_bytes(child, "source", category="context_checkpoints",
source_id=source, data=raw, extension="json")
write_task_result(child, "source", "completed", **_field(ref, source))
copied = copy_child_task_result(parent, {"id": "source", "drive_root": str(child)})
assert copied["child_ref_promotion"]["promoted_source_handle_count"] == 1
assert remove_subagent_task_drive(parent, "source") is True
assert not child.exists()
assert artifacts.read_actor_source_bytes(parent, "source", ref) == raw
assert artifacts.collect_task_artifact_records(parent, "source") == []
def test_failed_completion_promotion_retains_child_until_retry(tmp_path, monkeypatch):
parent = tmp_path / "canonical"
child = prepare_task_drive(parent, "source", "empty")
raw = b'{"full":"survives failed copy"}'
ref = artifacts.store_actor_source_bytes(child, "source", category="context_checkpoints",
source_id="completion", data=raw, extension="json")
write_task_result(child, "source", "completed", **_field(ref, "completion"))
with monkeypatch.context() as patch:
patch.setattr(artifacts, "store_actor_source_bytes", lambda *_a, **_k: (_ for _ in ()).throw(OSError("copy failed")))
copied = copy_child_task_result(parent, {"id": "source", "drive_root": str(child)})
assert copied["child_ref_promotion"]["status"] == "incomplete"
assert remove_subagent_task_drive(parent, "source") is False
assert child.exists()
copied = copy_child_task_result(parent, {"id": "source", "drive_root": str(child)})
assert copied["child_ref_promotion"]["status"] == "complete"
assert remove_subagent_task_drive(parent, "source") is True
assert artifacts.read_actor_source_bytes(parent, "source", ref) == raw
def test_unchanged_review_source_is_write_once(tmp_path, monkeypatch):
ctx = _context(tmp_path)
trace = {"review_runs": [_run()]}
review_projection.publish_acceptance_checkpoint(ctx, trace)
ref = trace["review_runs"][0]["applied_source_ref"]
path = artifacts.task_artifact_dir_path(tmp_path, "applied") / ref["path"]
stamp = path.stat().st_mtime_ns
monkeypatch.setattr(artifacts, "write_bytes_atomic", lambda *_a, **_k: pytest.fail("same source rewritten"))
review_projection.publish_acceptance_checkpoint(ctx, trace)
assert trace["review_runs"][0]["applied_source_ref"] == ref
assert path.stat().st_mtime_ns == stamp
assert not (path.parents[2] / artifacts._ARTIFACT_MANIFEST).exists()
@pytest.mark.serial
@pytest.mark.parametrize("kind", ["task_acceptance_review", "task_completion_observations"])
def test_legacy_flat_sources_remain_readable_without_data_migration(tmp_path, kind):
raw = b'{"delivery_results": [], "legacy": true}'
ref = artifacts.store_task_artifact_bytes(tmp_path, "source", "old.json", raw, kind=kind)
record = artifacts.artifact_record(artifacts.task_artifact_dir_path(tmp_path, "source") / "old.json", kind=kind)
write_task_result(tmp_path, "source", "completed", artifacts=[record], artifact_status="ready",
**_field(ref, "completion" if kind == "task_completion_observations" else "review"))
path = task_result_path(tmp_path, "source", create=False)
original = path.read_bytes()
loaded = load_task_result(tmp_path, "source")
assert loaded["artifacts"] == [] and loaded["artifact_status"] == "not_applicable"
assert artifacts.read_actor_source_bytes(tmp_path, "source", ref) == raw
if kind == "task_completion_observations":
assert completion_source_projection(tmp_path, "source", loaded, 0, len(raw))["text"] == raw.decode()
app = Starlette(routes=[Route("/api/tasks/{task_id}/artifacts/{name}", api_task_artifact)])
app.state.drive_root = tmp_path
with TestClient(app) as client:
assert client.get("/api/tasks/source/artifacts/old.json").content == raw
assert path.read_bytes() == original
@pytest.mark.serial
def test_source_download_is_bound_and_distinct_from_same_named_user_file(tmp_path):
ctx = _context(tmp_path)
trace = {"review_runs": [_run()]}
review_projection.publish_acceptance_checkpoint(ctx, trace)
ref = trace["review_runs"][0]["applied_source_ref"]
name = Path(ref["path"]).name
artifacts.store_task_artifact_bytes(tmp_path, "applied", name, b"user result", kind="user_file")
app = Starlette(routes=[Route("/api/tasks/{task_id}/artifacts/{name}", api_task_artifact)])
app.state.drive_root = tmp_path
with TestClient(app) as client:
url = f"/api/tasks/applied/artifacts/{name}"
assert client.get(url).content == b"user result"
source = client.get(url, params={"source": ref["path"]})
assert source.status_code == 200
assert hashlib.sha256(source.content).hexdigest() == ref["sha256"]
assert len(json.loads(source.content)["actors"][0]["parsed"]["findings"]) == 80
assert client.get(url, params={"source": "source_handles/context_checkpoints/not-published.json"}).status_code == 404
(artifacts.task_artifact_dir_path(tmp_path, "applied") / ref["path"]).write_bytes(b"corrupt")
assert client.get(url, params={"source": ref["path"]}).status_code == 404
@pytest.mark.parametrize("panels", [True, 7, {"legacy": "unknown"}])
def test_unavailable_review_projection_does_not_hide_valid_completion_source(tmp_path, panels):
raw = b'{"delivery_results": []}'
ref = artifacts.store_actor_source_bytes(tmp_path, "source", category="context_checkpoints",
source_id="completion", data=raw, extension="json")
result = {"task_id": "source", "review_projection": {"panels": panels},
"completion_observations": {"source_ref": ref}}
assert artifacts.read_task_result_source_bytes(tmp_path, result, Path(ref["path"]).name, ref["path"]) == raw

View file

@ -97,6 +97,22 @@ export function cancelTask(taskId, { cascade = false, stopPolicy = '' } = {}) {
return Object.keys(body).length ? jsonPost(url, body) : fetchJson(url, { method: 'POST' });
}
/** URL for one published immutable source, retaining the legacy flat-ref contract. */
export function taskSourceDownloadUrl(taskId, ref, legacyKind = '') {
const source = ref?.kind === 'task_source';
const size = source ? ref?.size : ref?.bytes;
const path = typeof ref?.path === 'string' ? ref.path : '';
const allowedPath = source
? /^source_handles\/(tool_results|context_checkpoints)\/[A-Za-z0-9][A-Za-z0-9._-]*$/
: /^[A-Za-z0-9][A-Za-z0-9._-]*$/;
if (!taskId || ref?.root !== 'artifact_store' || (!source && ref?.kind !== legacyKind)
|| !/^[0-9a-f]{64}$/.test(ref?.sha256 || '') || !allowedPath.test(path)
|| !Number.isSafeInteger(size) || size < 0) return '';
const name = path.split('/').at(-1);
const url = `/api/tasks/${encodeURIComponent(taskId)}/artifacts/${encodeURIComponent(name)}`;
return source ? `${url}?source=${encodeURIComponent(path)}` : url;
}
export async function resumeTask(taskId) {
return fetchJson(`/api/tasks/${encodeURIComponent(taskId)}/resume`, { method: 'POST' });
}

View file

@ -1,4 +1,5 @@
import { escapeHtmlAttr } from './utils.js';
import { taskSourceDownloadUrl } from './api_client.js';
import { harnessIdentityMarkup } from './harness_presentation.js';
import { reconcileReviewMarkup } from './review_dom_patch.js';
@ -767,15 +768,6 @@ export function formatReviewProjection(projection) {
return lines.join('\n');
}
function appliedReviewUrl(panel, owner) {
const ref = panel.applied_source_ref;
if (panel.applied_source_status !== 'available' || ref?.root !== 'artifact_store'
|| ref.kind !== 'task_acceptance_review' || !/^[0-9a-f]{64}$/.test(ref.sha256 || '')
|| !Number.isSafeInteger(ref.bytes) || ref.bytes < 0
|| !/^[A-Za-z0-9][A-Za-z0-9._-]*$/.test(ref.path || '')) return '';
return `/api/tasks/${encodeURIComponent(owner)}/artifacts/${encodeURIComponent(ref.path)}`;
}
export function taskAcceptanceGroupFromTaskDetail(detail, ownerTaskId = '') {
const owner = text(ownerTaskId || detail?.task_id);
const projection = detail?.review_projection;
@ -802,7 +794,8 @@ export function taskAcceptanceGroupFromTaskDetail(detail, ownerTaskId = '') {
initiatorTaskId: owner,
executions: executionsFromReviewRecord(panel),
execution: null,
detailRef: { surface: 'task_acceptance', url: appliedReviewUrl(panel, owner) },
detailRef: { surface: 'task_acceptance', url: panel.applied_source_status === 'available'
? taskSourceDownloadUrl(owner, panel.applied_source_ref, 'task_acceptance_review') : '' },
detailText: `${formatReviewProjection({ panels: [panel] })}\nCost unavailable`.trim(),
};
});

View file

@ -100,3 +100,20 @@ test('acceptance invalidation joins a current read and refreshes to the newly ap
assert.equal(release.length, 2);
assert.equal(reviewReferenceFromRow({ type: 'review_reference', surface: 'commit', task_id: 'root' }), null);
});
test('source handles use the bound-source selector on the existing download route', () => {
const path = `source_handles/context_checkpoints/acceptance-${hash}.json`;
const ref = { kind: 'task_source', root: 'artifact_store', path, size: 75000, sha256: hash };
const projected = group(panel({ applied_source_ref: ref }));
assert.equal(projected.attempts[0].detailRef.url,
`/api/tasks/root/artifacts/acceptance-${hash}.json?source=${encodeURIComponent(path)}`);
assert.match(renderReviewsSection([projected]), /Download full applied review/);
for (const bad of [
{ ...ref, path: 'source_handles/context_checkpoints/../private.json' },
{ ...ref, path: 'source_handles/other/private.json' },
{ ...ref, size: '75000' }, { ...ref, sha256: '' },
{ ...ref, root: 'runtime_data' },
]) {
assert.equal(group(panel({ applied_source_ref: bad })).attempts[0].detailRef.url, '');
}
});