Skip to content
188 changes: 169 additions & 19 deletions loopx/cli_commands/post_writeback.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,125 @@
PostWritebackHookRegistration,
dispatch_post_writeback_hooks,
)
from ..control_plane.post_writeback_composition_retry import (
append_composition_retry_receipt,
build_composition_retry_receipt,
composition_retry_receipt_id,
composition_retry_receipt_log_path,
composition_retry_receipt_ref,
settle_composition_retry_receipt,
)
from ..history import load_registry
from ..paths import resolve_runtime_root


PostWritebackProjectionBuilder = Callable[..., Mapping[str, object]]

def _composition_hook_identities(
hooks: Sequence[PostWritebackHookRegistration],
) -> list[dict[str, str]]:
return [
{
"hook_id": str(registration.hook_id or ""),
"capability_id": str(registration.capability_id or ""),
"policy_version": str(
getattr(registration, "policy_version", "") or ""
),
}
for registration in hooks
]


def _composition_failure_dispatch(
hooks: Sequence[PostWritebackHookRegistration],
*,
error_code: str,
receipt_ref: str | None,
) -> dict[str, Any]:
"""Report the concrete public-safe hook identities behind one composition failure."""

failures: list[dict[str, str]] = []
for registration in hooks:
failure = {
"hook_id": str(registration.hook_id or ""),
"capability_id": str(registration.capability_id or ""),
"error_code": error_code,
}
if receipt_ref is not None:
failure["durable_receipt_ref"] = receipt_ref
failures.append(failure)
return {
"schema_version": POST_WRITEBACK_HOOK_DISPATCH_SCHEMA_VERSION,
"phase": "post_writeback",
"registered_count": len(hooks),
"invoked_count": 0,
"replayed_hooks": [],
"retried_hooks": [],
"intent_count": 0,
"intents": [],
"failures": failures,
"primary_writeback_preserved": True,
"external_writes_performed": False,
}


def _recorded_composition_failure(
journal_path: Path,
*,
hooks: Sequence[PostWritebackHookRegistration],
goal_id: str,
event_kind: str,
identity: Mapping[str, Any],
state_version: str,
committed_at: str,
error_code: str,
) -> dict[str, Any]:
"""Persist one retryable receipt, degrading to identity-only on journal errors."""

receipt = build_composition_retry_receipt(
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
committed_at=committed_at,
hook_identities=_composition_hook_identities(hooks),
error_code=error_code,
)
receipt_ref: str | None = composition_retry_receipt_ref(
str(receipt["receipt_id"])
)
try:
append_composition_retry_receipt(journal_path, receipt)
except (OSError, ValueError):
receipt_ref = None
return _composition_failure_dispatch(
hooks, error_code=error_code, receipt_ref=receipt_ref
)


def _settle_composition_retry_quietly(
journal_path: Path,
*,
goal_id: str,
event_kind: str,
identity: Mapping[str, Any],
state_version: str,
committed_at: str,
hooks: Sequence[PostWritebackHookRegistration],
) -> None:
try:
settle_composition_retry_receipt(
journal_path,
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
committed_at=committed_at,
hook_identities=_composition_hook_identities(hooks),
)
except (OSError, ValueError):
return


def dispatch_committed_cli_post_writeback_hooks(
*,
Expand All @@ -31,14 +144,31 @@ def dispatch_committed_cli_post_writeback_hooks(
) -> dict[str, Any]:
"""Bridge one committed CLI mutation into the TS-owned hook lifecycle.

Projection and provider failures are isolated from the primary mutation.
The helper intentionally owns no capability policy and grants no effects.
Projection and provider failures are isolated from the primary mutation:
the primary writeback is never rolled back or repeated, a failed
composition records one retryable receipt bound to that writeback, and a
later replay of the same writeback settles the receipt once the
projection composes cleanly. The helper intentionally owns no capability
policy and grants no effects.
"""

try:
runtime_root = resolve_runtime_root(
load_registry(registry_path), runtime_root_arg
)
journal_path = composition_retry_receipt_log_path(runtime_root, goal_id)
composition_retry_receipt_id(
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
hook_identities=_composition_hook_identities(hooks),
)
except Exception: # Optional hooks never alter primary truth.
return _composition_failure_dispatch(
hooks, error_code="source_projection_failed", receipt_ref=None
)
try:
projection = (
dict(
projection_builder(
Expand All @@ -52,7 +182,19 @@ def dispatch_committed_cli_post_writeback_hooks(
if projection_builder is not None
else {}
)
return dispatch_post_writeback_hooks(
except Exception: # The source projection stays retryable, never primary.
return _recorded_composition_failure(
journal_path,
hooks=hooks,
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
committed_at=committed_at,
error_code="source_projection_failed",
)
try:
result = dispatch_post_writeback_hooks(
hooks,
source={
"schema_version": "loopx_post_writeback_hook_source_v0",
Expand All @@ -74,22 +216,30 @@ def dispatch_committed_cli_post_writeback_hooks(
},
runtime_root=runtime_root,
)
except Exception: # Optional hooks never alter primary truth.
return {
"schema_version": POST_WRITEBACK_HOOK_DISPATCH_SCHEMA_VERSION,
"phase": "post_writeback",
"registered_count": len(hooks),
"intent_count": 0,
"failures": [
{
"hook_id": "composition",
"capability_id": "unknown",
"error_code": "source_projection_failed",
}
],
"primary_writeback_preserved": True,
"external_writes_performed": False,
}
except Exception: # An unexpected transport collapse stays retryable.
return _recorded_composition_failure(
journal_path,
hooks=hooks,
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
committed_at=committed_at,
error_code="dispatch_failed",
)
# The composition receipt tracks projection composition only: it settles
# once the projection composed and the hook lifecycle returned, while any
# hook-level failure keeps its own per-hook failure trail in `result`.
_settle_composition_retry_quietly(
journal_path,
goal_id=goal_id,
event_kind=event_kind,
identity=identity,
state_version=state_version,
committed_at=committed_at,
hooks=hooks,
)
return result


__all__ = [
Expand Down
10 changes: 10 additions & 0 deletions loopx/cli_commands/status.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,9 @@
resolve_status_projection_cache_runtime_root,
write_status_projection_cache,
)
from ..control_plane.post_writeback_composition_retry import (
collect_pending_composition_retry_projection,
)
from ..control_plane.status.agent_lane_projection import (
compact_agent_lane_status_payload_for_display,
)
Expand Down Expand Up @@ -232,6 +235,13 @@ def handle_status_command(
agent_id=args.agent_id,
)
compact_agent_lane_todo_index_for_status_display(payload)
pending_composition_retries = collect_pending_composition_retry_projection(
runtime_root, args.goal_id, agent_id=args.agent_id
)
if pending_composition_retries is not None:
payload["pending_composition_retry_receipts"] = (
pending_composition_retries
)
except Exception as exc:
payload = {
"ok": False,
Expand Down
Loading