ソースを参照

编排:建立任务闭环领域模型与文件存储

新增 CompletionPolicy、AgentRole、Task/Attempt/Validation/Decision 等通用模型和显式状态机。定义 TaskStore、ArtifactStore、AgentExecutor、ToolPolicy、EventSink 端口,并实现带原子替换、revision 乐观锁、单进程锁、不可变快照哈希和事件日志的文件系统适配器。
SamLee 4 日 前
コミット
99cd3e04dd

+ 17 - 0
agent/agent/orchestration/config.py

@@ -0,0 +1,17 @@
+"""Orchestration configuration."""
+
+from dataclasses import dataclass
+
+
+@dataclass(frozen=True)
+class OrchestrationConfig:
+    max_parallel_tasks: int = 4
+    max_repair_continuations: int = 1
+    default_worker_preset: str = "worker"
+    default_validator_preset: str = "validator"
+
+    def __post_init__(self) -> None:
+        if self.max_parallel_tasks < 1:
+            raise ValueError("max_parallel_tasks must be at least 1")
+        if self.max_repair_continuations < 0:
+            raise ValueError("max_repair_continuations cannot be negative")

+ 443 - 0
agent/agent/orchestration/models.py

@@ -0,0 +1,443 @@
+"""Domain models for explicit task execution and independent validation.
+
+The orchestration ledger is intentionally independent from GoalTree and the
+runtime implementation.  It is the source of truth for explicit-validation
+runs; GoalTree is only a compatibility projection.
+"""
+
+from __future__ import annotations
+
+from dataclasses import asdict, dataclass, field
+from datetime import datetime, timezone
+from enum import Enum
+from typing import Any, Dict, List, Optional
+from uuid import uuid4
+
+
+def utc_now() -> str:
+    return datetime.now(timezone.utc).isoformat()
+
+
+def new_id() -> str:
+    return str(uuid4())
+
+
+class StrEnum(str, Enum):
+    def __str__(self) -> str:
+        return self.value
+
+
+class CompletionPolicy(StrEnum):
+    LEGACY_AUTO = "legacy_auto"
+    EXPLICIT_VALIDATION = "explicit_validation"
+
+
+class AgentRole(StrEnum):
+    LEGACY = "legacy"
+    PLANNER = "planner"
+    WORKER = "worker"
+    VALIDATOR = "validator"
+
+
+class TaskStatus(StrEnum):
+    PENDING = "pending"
+    RUNNING = "running"
+    AWAITING_VALIDATION = "awaiting_validation"
+    VALIDATING = "validating"
+    AWAITING_DECISION = "awaiting_decision"
+    NEEDS_REPLAN = "needs_replan"
+    WAITING_CHILDREN = "waiting_children"
+    COMPLETED = "completed"
+    SUPERSEDED = "superseded"
+    BLOCKED = "blocked"
+    CANCELLED = "cancelled"
+
+
+class AttemptStatus(StrEnum):
+    RUNNING = "running"
+    SUBMITTED = "submitted"
+    FAILED = "failed"
+    STOPPED = "stopped"
+    EXPIRED = "expired"
+
+
+class ValidationRunStatus(StrEnum):
+    PENDING = "pending"
+    RUNNING = "running"
+    COMPLETED = "completed"
+    ERROR = "error"
+    STOPPED = "stopped"
+    EXPIRED = "expired"
+
+
+class ValidationVerdict(StrEnum):
+    PASSED = "passed"
+    FAILED = "failed"
+    INCONCLUSIVE = "inconclusive"
+
+
+class DecisionAction(StrEnum):
+    ACCEPT = "accept"
+    REPAIR = "repair"
+    RETRY = "retry"
+    REVALIDATE = "revalidate"
+    REVISE = "revise"
+    SPLIT = "split"
+    BLOCK = "block"
+    UNBLOCK = "unblock"
+    CANCEL = "cancel"
+    SUPERSEDE = "supersede"
+
+
+@dataclass(frozen=True)
+class AcceptanceCriterion:
+    criterion_id: str
+    description: str
+    hard: bool = True
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "AcceptanceCriterion":
+        return cls(
+            criterion_id=str(data.get("criterion_id") or data.get("id") or new_id()),
+            description=str(data.get("description", "")),
+            hard=bool(data.get("hard", data.get("required", True))),
+        )
+
+
+@dataclass(frozen=True)
+class TaskSpec:
+    version: int
+    objective: str
+    acceptance_criteria: List[AcceptanceCriterion] = field(default_factory=list)
+    context_refs: List[str] = field(default_factory=list)
+    created_at: str = field(default_factory=utc_now)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "TaskSpec":
+        return cls(
+            version=int(data.get("version", 1)),
+            objective=str(data.get("objective", "")),
+            acceptance_criteria=[AcceptanceCriterion.from_dict(x) for x in data.get("acceptance_criteria", [])],
+            context_refs=list(data.get("context_refs", [])),
+            created_at=data.get("created_at", utc_now()),
+        )
+
+
+@dataclass
+class TaskRecord:
+    task_id: str
+    goal_id: Optional[str]
+    parent_task_id: Optional[str]
+    display_path: str
+    specs: List[TaskSpec]
+    current_spec_version: int = 1
+    status: TaskStatus = TaskStatus.PENDING
+    child_task_ids: List[str] = field(default_factory=list)
+    attempt_ids: List[str] = field(default_factory=list)
+    validation_ids: List[str] = field(default_factory=list)
+    decision_ids: List[str] = field(default_factory=list)
+    repair_count_by_version: Dict[str, int] = field(default_factory=dict)
+    blocked_reason: Optional[str] = None
+    superseded_by: Optional[str] = None
+    created_at: str = field(default_factory=utc_now)
+    updated_at: str = field(default_factory=utc_now)
+
+    @property
+    def current_spec(self) -> TaskSpec:
+        for spec in reversed(self.specs):
+            if spec.version == self.current_spec_version:
+                return spec
+        raise ValueError(f"Task {self.task_id} is missing spec version {self.current_spec_version}")
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "TaskRecord":
+        return cls(
+            task_id=data["task_id"],
+            goal_id=data.get("goal_id"),
+            parent_task_id=data.get("parent_task_id"),
+            display_path=data.get("display_path", ""),
+            specs=[TaskSpec.from_dict(x) for x in data.get("specs", [])],
+            current_spec_version=int(data.get("current_spec_version", 1)),
+            status=TaskStatus(data.get("status", TaskStatus.PENDING.value)),
+            child_task_ids=list(data.get("child_task_ids", [])),
+            attempt_ids=list(data.get("attempt_ids", [])),
+            validation_ids=list(data.get("validation_ids", [])),
+            decision_ids=list(data.get("decision_ids", [])),
+            repair_count_by_version={str(k): int(v) for k, v in data.get("repair_count_by_version", {}).items()},
+            blocked_reason=data.get("blocked_reason"),
+            superseded_by=data.get("superseded_by"),
+            created_at=data.get("created_at", utc_now()),
+            updated_at=data.get("updated_at", utc_now()),
+        )
+
+
+@dataclass(frozen=True)
+class ArtifactRef:
+    uri: str
+    kind: str = "generic"
+    version: Optional[str] = None
+    digest: Optional[str] = None
+    summary: Optional[str] = None
+    metadata: Dict[str, Any] = field(default_factory=dict)
+
+    def __post_init__(self) -> None:
+        if not self.uri:
+            raise ValueError("ArtifactRef.uri is required")
+        if not self.version and not self.digest:
+            raise ValueError("ArtifactRef requires an immutable version or digest")
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "ArtifactRef":
+        return cls(
+            uri=str(data.get("uri", "")),
+            kind=str(data.get("kind", "generic")),
+            version=data.get("version"),
+            digest=data.get("digest") or data.get("hash"),
+            summary=data.get("summary"),
+            metadata=dict(data.get("metadata", {})),
+        )
+
+
+@dataclass(frozen=True)
+class AttemptSubmission:
+    summary: str
+    artifact_refs: List[ArtifactRef] = field(default_factory=list)
+    evidence_refs: List[ArtifactRef] = field(default_factory=list)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "AttemptSubmission":
+        return cls(
+            summary=str(data.get("summary", "")),
+            artifact_refs=[ArtifactRef.from_dict(x) for x in data.get("artifact_refs", [])],
+            evidence_refs=[ArtifactRef.from_dict(x) for x in data.get("evidence_refs", [])],
+        )
+
+
+@dataclass(frozen=True)
+class ArtifactSnapshot:
+    snapshot_id: str
+    attempt_id: str
+    normalized_content: Dict[str, Any]
+    sha256: str
+    artifact_refs: List[ArtifactRef]
+    evidence_refs: List[ArtifactRef]
+    created_at: str = field(default_factory=utc_now)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "ArtifactSnapshot":
+        return cls(
+            snapshot_id=data["snapshot_id"],
+            attempt_id=data["attempt_id"],
+            normalized_content=dict(data.get("normalized_content", {})),
+            sha256=data["sha256"],
+            artifact_refs=[ArtifactRef.from_dict(x) for x in data.get("artifact_refs", [])],
+            evidence_refs=[ArtifactRef.from_dict(x) for x in data.get("evidence_refs", [])],
+            created_at=data.get("created_at", utc_now()),
+        )
+
+
+@dataclass
+class TaskAttempt:
+    attempt_id: str
+    task_id: str
+    spec_version: int
+    worker_trace_id: str
+    worker_preset: str
+    execution_mode: str
+    status: AttemptStatus = AttemptStatus.RUNNING
+    continue_from_trace_id: Optional[str] = None
+    snapshot_id: Optional[str] = None
+    submission: Optional[AttemptSubmission] = None
+    error: Optional[str] = None
+    created_at: str = field(default_factory=utc_now)
+    updated_at: str = field(default_factory=utc_now)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "TaskAttempt":
+        submission = data.get("submission")
+        return cls(
+            attempt_id=data["attempt_id"],
+            task_id=data["task_id"],
+            spec_version=int(data["spec_version"]),
+            worker_trace_id=data["worker_trace_id"],
+            worker_preset=data.get("worker_preset", "worker"),
+            execution_mode=data.get("execution_mode", "new"),
+            status=AttemptStatus(data.get("status", AttemptStatus.RUNNING.value)),
+            continue_from_trace_id=data.get("continue_from_trace_id"),
+            snapshot_id=data.get("snapshot_id"),
+            submission=AttemptSubmission.from_dict(submission) if submission else None,
+            error=data.get("error"),
+            created_at=data.get("created_at", utc_now()),
+            updated_at=data.get("updated_at", utc_now()),
+        )
+
+
+@dataclass(frozen=True)
+class CriterionResult:
+    criterion_id: str
+    verdict: ValidationVerdict
+    reason: str
+    evidence_refs: List[ArtifactRef] = field(default_factory=list)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "CriterionResult":
+        return cls(
+            criterion_id=str(data.get("criterion_id", "")),
+            verdict=ValidationVerdict(data.get("verdict", ValidationVerdict.INCONCLUSIVE.value)),
+            reason=str(data.get("reason", "")),
+            evidence_refs=[ArtifactRef.from_dict(x) for x in data.get("evidence_refs", [])],
+        )
+
+
+@dataclass
+class ValidationReport:
+    validation_id: str
+    task_id: str
+    attempt_id: str
+    spec_version: int
+    snapshot_id: str
+    validator_trace_id: str
+    validator_preset: str = "validator"
+    status: ValidationRunStatus = ValidationRunStatus.PENDING
+    verdict: Optional[ValidationVerdict] = None
+    criterion_results: List[CriterionResult] = field(default_factory=list)
+    summary: str = ""
+    evidence_refs: List[ArtifactRef] = field(default_factory=list)
+    unverified_claims: List[str] = field(default_factory=list)
+    risks: List[str] = field(default_factory=list)
+    recommendation: str = ""
+    error: Optional[str] = None
+    created_at: str = field(default_factory=utc_now)
+    updated_at: str = field(default_factory=utc_now)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "ValidationReport":
+        verdict = data.get("verdict")
+        return cls(
+            validation_id=data["validation_id"],
+            task_id=data["task_id"],
+            attempt_id=data["attempt_id"],
+            spec_version=int(data["spec_version"]),
+            snapshot_id=data["snapshot_id"],
+            validator_trace_id=data["validator_trace_id"],
+            validator_preset=data.get("validator_preset", "validator"),
+            status=ValidationRunStatus(data.get("status", ValidationRunStatus.PENDING.value)),
+            verdict=ValidationVerdict(verdict) if verdict else None,
+            criterion_results=[CriterionResult.from_dict(x) for x in data.get("criterion_results", [])],
+            summary=data.get("summary", ""),
+            evidence_refs=[ArtifactRef.from_dict(x) for x in data.get("evidence_refs", [])],
+            unverified_claims=list(data.get("unverified_claims", [])),
+            risks=list(data.get("risks", [])),
+            recommendation=data.get("recommendation", ""),
+            error=data.get("error"),
+            created_at=data.get("created_at", utc_now()),
+            updated_at=data.get("updated_at", utc_now()),
+        )
+
+
+@dataclass(frozen=True)
+class PlannerDecision:
+    decision_id: str
+    task_id: str
+    action: DecisionAction
+    reason: str
+    from_status: TaskStatus
+    to_status: TaskStatus
+    attempt_id: Optional[str] = None
+    validation_id: Optional[str] = None
+    payload: Dict[str, Any] = field(default_factory=dict)
+    created_at: str = field(default_factory=utc_now)
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "PlannerDecision":
+        return cls(
+            decision_id=data["decision_id"],
+            task_id=data["task_id"],
+            action=DecisionAction(data["action"]),
+            reason=data.get("reason", ""),
+            from_status=TaskStatus(data["from_status"]),
+            to_status=TaskStatus(data["to_status"]),
+            attempt_id=data.get("attempt_id"),
+            validation_id=data.get("validation_id"),
+            payload=dict(data.get("payload", {})),
+            created_at=data.get("created_at", utc_now()),
+        )
+
+
+@dataclass
+class TaskLedger:
+    root_trace_id: str
+    mission: str
+    revision: int = -1
+    tasks: Dict[str, TaskRecord] = field(default_factory=dict)
+    attempts: Dict[str, TaskAttempt] = field(default_factory=dict)
+    validations: Dict[str, ValidationReport] = field(default_factory=dict)
+    decisions: Dict[str, PlannerDecision] = field(default_factory=dict)
+    focused_task_id: Optional[str] = None
+    idempotency_results: Dict[str, Dict[str, Any]] = field(default_factory=dict)
+    created_at: str = field(default_factory=utc_now)
+    updated_at: str = field(default_factory=utc_now)
+
+    def to_dict(self) -> Dict[str, Any]:
+        return _enum_values(asdict(self))
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "TaskLedger":
+        return cls(
+            root_trace_id=data["root_trace_id"],
+            mission=data.get("mission", ""),
+            revision=int(data.get("revision", -1)),
+            tasks={k: TaskRecord.from_dict(v) for k, v in data.get("tasks", {}).items()},
+            attempts={k: TaskAttempt.from_dict(v) for k, v in data.get("attempts", {}).items()},
+            validations={k: ValidationReport.from_dict(v) for k, v in data.get("validations", {}).items()},
+            decisions={k: PlannerDecision.from_dict(v) for k, v in data.get("decisions", {}).items()},
+            focused_task_id=data.get("focused_task_id"),
+            idempotency_results=dict(data.get("idempotency_results", {})),
+            created_at=data.get("created_at", utc_now()),
+            updated_at=data.get("updated_at", utc_now()),
+        )
+
+
+@dataclass(frozen=True)
+class TaskCycleResult:
+    task_id: str
+    task_status: TaskStatus
+    attempt_id: Optional[str] = None
+    validation_id: Optional[str] = None
+    validation: Optional[ValidationReport] = None
+    error: Optional[str] = None
+
+    def to_dict(self) -> Dict[str, Any]:
+        return _enum_values(asdict(self))
+
+    @classmethod
+    def from_dict(cls, data: Dict[str, Any]) -> "TaskCycleResult":
+        validation = data.get("validation")
+        return cls(
+            task_id=data["task_id"],
+            task_status=TaskStatus(data["task_status"]),
+            attempt_id=data.get("attempt_id"),
+            validation_id=data.get("validation_id"),
+            validation=ValidationReport.from_dict(validation) if validation else None,
+            error=data.get("error"),
+        )
+
+
+def _enum_values(value: Any) -> Any:
+    if isinstance(value, Enum):
+        return value.value
+    if isinstance(value, dict):
+        return {str(k): _enum_values(v) for k, v in value.items()}
+    if isinstance(value, (list, tuple)):
+        return [_enum_values(v) for v in value]
+    return value
+
+
+__all__ = [
+    "CompletionPolicy", "AgentRole", "TaskStatus", "AttemptStatus",
+    "ValidationRunStatus", "ValidationVerdict", "DecisionAction",
+    "AcceptanceCriterion", "TaskSpec", "TaskRecord", "TaskAttempt",
+    "ArtifactRef", "ArtifactSnapshot", "AttemptSubmission", "CriterionResult",
+    "ValidationReport", "PlannerDecision", "TaskLedger", "TaskCycleResult",
+    "new_id", "utc_now",
+]

+ 65 - 0
agent/agent/orchestration/protocols.py

@@ -0,0 +1,65 @@
+"""Ports used by the orchestration domain."""
+
+from __future__ import annotations
+
+from dataclasses import dataclass
+from typing import Any, Dict, Optional, Protocol
+
+from .models import ArtifactSnapshot, AttemptSubmission, TaskLedger
+
+
+@dataclass(frozen=True)
+class CommitResult:
+    revision: int
+    ledger: TaskLedger
+
+
+@dataclass(frozen=True)
+class WorkerRunResult:
+    trace_id: str
+    status: str
+    summary: str = ""
+    error: Optional[str] = None
+
+
+@dataclass(frozen=True)
+class ValidatorRunResult:
+    trace_id: str
+    status: str
+    summary: str = ""
+    error: Optional[str] = None
+
+
+class TaskStore(Protocol):
+    async def load(self, root_trace_id: str) -> TaskLedger: ...
+    async def commit(
+        self,
+        ledger: TaskLedger,
+        expected_revision: int,
+        idempotency_key: Optional[str] = None,
+    ) -> CommitResult: ...
+
+
+class ArtifactStore(Protocol):
+    async def freeze(self, attempt_id: str, submission: AttemptSubmission) -> ArtifactSnapshot: ...
+    async def get(self, snapshot_id: str) -> ArtifactSnapshot: ...
+
+
+class AgentExecutor(Protocol):
+    async def run_worker(self, context: Dict[str, Any]) -> WorkerRunResult: ...
+    async def run_validator(self, context: Dict[str, Any]) -> ValidatorRunResult: ...
+
+
+class ToolPolicy(Protocol):
+    def resolve(self, config: Any, preset: Any, registry: Any) -> Any: ...
+    def authorize(self, role: Any, tool_name: str, resolved_policy: Any) -> Any: ...
+
+
+class EventSink(Protocol):
+    async def emit(self, root_trace_id: str, event_type: str, payload: Dict[str, Any]) -> None: ...
+
+
+__all__ = [
+    "CommitResult", "WorkerRunResult", "ValidatorRunResult", "TaskStore",
+    "ArtifactStore", "AgentExecutor", "ToolPolicy", "EventSink",
+]

+ 35 - 0
agent/agent/orchestration/state_machine.py

@@ -0,0 +1,35 @@
+"""Legal task transitions for explicit validation."""
+
+from __future__ import annotations
+
+from typing import Dict, FrozenSet
+
+from .models import TaskStatus
+
+
+class InvalidTaskTransition(ValueError):
+    pass
+
+
+LEGAL_TRANSITIONS: Dict[TaskStatus, FrozenSet[TaskStatus]] = {
+    TaskStatus.PENDING: frozenset({TaskStatus.RUNNING, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.RUNNING: frozenset({TaskStatus.AWAITING_VALIDATION, TaskStatus.NEEDS_REPLAN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.AWAITING_VALIDATION: frozenset({TaskStatus.VALIDATING, TaskStatus.NEEDS_REPLAN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.VALIDATING: frozenset({TaskStatus.AWAITING_DECISION, TaskStatus.NEEDS_REPLAN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.AWAITING_DECISION: frozenset({TaskStatus.COMPLETED, TaskStatus.PENDING, TaskStatus.WAITING_CHILDREN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.NEEDS_REPLAN: frozenset({TaskStatus.PENDING, TaskStatus.VALIDATING, TaskStatus.WAITING_CHILDREN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.WAITING_CHILDREN: frozenset({TaskStatus.NEEDS_REPLAN, TaskStatus.BLOCKED, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.BLOCKED: frozenset({TaskStatus.NEEDS_REPLAN, TaskStatus.CANCELLED, TaskStatus.SUPERSEDED}),
+    TaskStatus.COMPLETED: frozenset(),
+    TaskStatus.CANCELLED: frozenset(),
+    TaskStatus.SUPERSEDED: frozenset(),
+}
+
+
+def transition(current: TaskStatus, target: TaskStatus) -> TaskStatus:
+    if target not in LEGAL_TRANSITIONS[current]:
+        raise InvalidTaskTransition(f"Illegal task transition: {current.value} -> {target.value}")
+    return target
+
+
+__all__ = ["InvalidTaskTransition", "LEGAL_TRANSITIONS", "transition"]

+ 199 - 0
agent/agent/orchestration/store.py

@@ -0,0 +1,199 @@
+"""Single-process filesystem adapters for orchestration ports."""
+
+from __future__ import annotations
+
+import asyncio
+import hashlib
+import json
+import os
+import tempfile
+from dataclasses import asdict
+from pathlib import Path
+from typing import Any, Dict, Optional
+
+from .models import ArtifactSnapshot, AttemptSubmission, TaskLedger, new_id, utc_now
+from .protocols import CommitResult
+
+
+class TaskStoreError(RuntimeError):
+    pass
+
+
+class TaskStoreNotFound(TaskStoreError):
+    pass
+
+
+class RevisionConflict(TaskStoreError):
+    pass
+
+
+class FileSystemTaskStore:
+    """Atomic JSON ledger store with optimistic revisions and per-root locks."""
+
+    _locks: Dict[str, asyncio.Lock] = {}
+
+    def __init__(self, base_path: str = ".trace") -> None:
+        self.base_path = Path(base_path)
+
+    def orchestration_dir(self, root_trace_id: str) -> Path:
+        return self.base_path / root_trace_id / "orchestration"
+
+    def ledger_path(self, root_trace_id: str) -> Path:
+        return self.orchestration_dir(root_trace_id) / "ledger.json"
+
+    def _lock(self, root_trace_id: str) -> asyncio.Lock:
+        key = str(self.ledger_path(root_trace_id).resolve())
+        if key not in self._locks:
+            self._locks[key] = asyncio.Lock()
+        return self._locks[key]
+
+    async def load(self, root_trace_id: str) -> TaskLedger:
+        async with self._lock(root_trace_id):
+            return self._load_unlocked(root_trace_id)
+
+    def _load_unlocked(self, root_trace_id: str) -> TaskLedger:
+        path = self.ledger_path(root_trace_id)
+        if not path.exists():
+            raise TaskStoreNotFound(f"No task ledger for root trace {root_trace_id}")
+        try:
+            return TaskLedger.from_dict(json.loads(path.read_text(encoding="utf-8")))
+        except (OSError, ValueError, TypeError, json.JSONDecodeError) as exc:
+            raise TaskStoreError(f"Cannot load task ledger {path}: {exc}") from exc
+
+    async def commit(
+        self,
+        ledger: TaskLedger,
+        expected_revision: int,
+        idempotency_key: Optional[str] = None,
+    ) -> CommitResult:
+        root_trace_id = ledger.root_trace_id
+        async with self._lock(root_trace_id):
+            path = self.ledger_path(root_trace_id)
+            current_revision = -1
+            if path.exists():
+                current_revision = self._load_unlocked(root_trace_id).revision
+            if current_revision != expected_revision:
+                raise RevisionConflict(
+                    f"Task ledger {root_trace_id} revision conflict: "
+                    f"expected {expected_revision}, actual {current_revision}"
+                )
+
+            ledger.revision = expected_revision + 1
+            ledger.updated_at = utc_now()
+            path.parent.mkdir(parents=True, exist_ok=True)
+            self._atomic_json_write(path, ledger.to_dict())
+            return CommitResult(revision=ledger.revision, ledger=ledger)
+
+    @staticmethod
+    def _atomic_json_write(path: Path, data: Dict[str, Any]) -> None:
+        fd, tmp_name = tempfile.mkstemp(prefix=f".{path.name}.", suffix=".tmp", dir=str(path.parent))
+        try:
+            with os.fdopen(fd, "w", encoding="utf-8") as handle:
+                json.dump(data, handle, ensure_ascii=False, indent=2, sort_keys=True)
+                handle.flush()
+                os.fsync(handle.fileno())
+            os.replace(tmp_name, path)
+        except Exception:
+            try:
+                os.unlink(tmp_name)
+            except OSError:
+                pass
+            raise
+
+
+class FileSystemArtifactStore:
+    """Immutable snapshot store bound to a root trace."""
+
+    _locks: Dict[str, asyncio.Lock] = {}
+
+    def __init__(self, base_path: str = ".trace", root_trace_id: Optional[str] = None) -> None:
+        self.base_path = Path(base_path)
+        self.root_trace_id = root_trace_id
+
+    def for_root(self, root_trace_id: str) -> "FileSystemArtifactStore":
+        return FileSystemArtifactStore(str(self.base_path), root_trace_id=root_trace_id)
+
+    def _artifacts_dir(self) -> Path:
+        if not self.root_trace_id:
+            raise ValueError("FileSystemArtifactStore must be bound with for_root(root_trace_id)")
+        return self.base_path / self.root_trace_id / "orchestration" / "artifacts"
+
+    def _lock(self) -> asyncio.Lock:
+        key = str(self._artifacts_dir().resolve())
+        if key not in self._locks:
+            self._locks[key] = asyncio.Lock()
+        return self._locks[key]
+
+    async def freeze(self, attempt_id: str, submission: AttemptSubmission) -> ArtifactSnapshot:
+        normalized = {
+            "summary": submission.summary.strip(),
+            "artifact_refs": [asdict(x) for x in submission.artifact_refs],
+            "evidence_refs": [asdict(x) for x in submission.evidence_refs],
+        }
+        canonical = json.dumps(normalized, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
+        digest = hashlib.sha256(canonical.encode("utf-8")).hexdigest()
+        snapshot = ArtifactSnapshot(
+            snapshot_id=new_id(),
+            attempt_id=attempt_id,
+            normalized_content=normalized,
+            sha256=digest,
+            artifact_refs=list(submission.artifact_refs),
+            evidence_refs=list(submission.evidence_refs),
+        )
+        path = self._artifacts_dir() / f"{snapshot.snapshot_id}.json"
+        async with self._lock():
+            path.parent.mkdir(parents=True, exist_ok=True)
+            if path.exists():
+                raise TaskStoreError(f"Snapshot already exists: {snapshot.snapshot_id}")
+            FileSystemTaskStore._atomic_json_write(path, _json_values(asdict(snapshot)))
+        return snapshot
+
+    async def get(self, snapshot_id: str) -> ArtifactSnapshot:
+        path = self._artifacts_dir() / f"{snapshot_id}.json"
+        async with self._lock():
+            if not path.exists():
+                raise FileNotFoundError(f"Artifact snapshot not found: {snapshot_id}")
+            return ArtifactSnapshot.from_dict(json.loads(path.read_text(encoding="utf-8")))
+
+
+class TraceEventSink:
+    """Append-only orchestration event stream."""
+
+    _locks: Dict[str, asyncio.Lock] = {}
+
+    def __init__(self, base_path: str = ".trace") -> None:
+        self.base_path = Path(base_path)
+
+    async def emit(self, root_trace_id: str, event_type: str, payload: Dict[str, Any]) -> None:
+        path = self.base_path / root_trace_id / "orchestration" / "events.jsonl"
+        key = str(path.resolve())
+        lock = self._locks.setdefault(key, asyncio.Lock())
+        event = {
+            "event_id": new_id(),
+            "event_type": event_type,
+            "root_trace_id": root_trace_id,
+            "created_at": utc_now(),
+            "payload": payload,
+        }
+        async with lock:
+            path.parent.mkdir(parents=True, exist_ok=True)
+            with path.open("a", encoding="utf-8") as handle:
+                handle.write(json.dumps(event, ensure_ascii=False, sort_keys=True) + "\n")
+                handle.flush()
+
+
+def _json_values(value: Any) -> Any:
+    from enum import Enum
+    if isinstance(value, Enum):
+        return value.value
+    if isinstance(value, dict):
+        return {str(k): _json_values(v) for k, v in value.items()}
+    if isinstance(value, (list, tuple)):
+        return [_json_values(v) for v in value]
+    return value
+
+
+__all__ = [
+    "TaskStoreError", "TaskStoreNotFound", "RevisionConflict",
+    "FileSystemTaskStore", "FileSystemArtifactStore", "TraceEventSink",
+]