Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1453,12 +1453,12 @@ facade-exit condition in §4), TS gathers the facts directly and the adapter is
deleted.

Rollout. `status` (including `--goal-id` and projection-cache hits),
`global-summary` and `global-gates` now carry the envelope. Every other
`collect_status` caller receives the status envelope in its payload but does
not yet emit its own. Next, in order: `global-todos` and `global-risks` (the
same composition, one call each), `quota should-run`, `review-packet`, and
Decision Context packets. A read model added to or migrated into TypeScript
emits the envelope in the same PR; §6 makes this a promotion gate.
`global-summary`, `global-gates`, `global-todos` and `global-risks` now carry
the envelope. Every other `collect_status` caller receives the status envelope
in its payload but does not yet emit its own. Next, in order: `quota
should-run`, `review-packet`, and Decision Context packets. A read model added
to or migrated into TypeScript emits the envelope in the same PR; §6 makes
this a promotion gate.

Consumers treat a missing envelope as unknown freshness, and disclose an
alerting one before stating any conclusion that depends on it. Field
Expand Down Expand Up @@ -1530,7 +1530,7 @@ choice is now implemented rather than hypothetical.
| Quota monitor-poll commit transaction | TypeScript owns monitor admission revalidation, target/event/result construction, effect replay/index CAS, provider intent, and repairable JSON/Markdown/index persistence | Python projects compact `should-run` facts, invokes the real Todo provider between at most two reductions, reloads legacy status, and holds the cross-writer index lock |
| Runtime decoders ([#3443](https://github.com/loopx-project/loopx/pull/3443)) | Stable primitive decoding has one small shared module; domain decoders remain local | No larger schema framework is justified |
| Transaction payoff ([#3464](https://github.com/loopx-project/loopx/pull/3464), [#3481](https://github.com/loopx-project/loopx/pull/3481), and Todo completion) | Turn settlement, quota delivery routing, and Todo completion each cross one coarse TS boundary; the Todo transaction owns identity, replay fencing, validation planning/result reduction, continuation/recovery, and completion metadata | Python still executes explicitly external providers and materializes legacy Markdown results; other domains still need their own bounded cutovers |
| Projection envelope | TypeScript owns decoding of `loopx_projection_envelope_v0` and every freshness, alert, completeness and replay decision | Python gathers read facts for `status`, `global-summary` and `global-gates` until those projections migrate |
| Projection envelope | TypeScript owns decoding of `loopx_projection_envelope_v0` and every freshness, alert, completeness and replay decision | Python gathers read facts for `status`, `global-summary`, `global-gates`, `global-todos` and `global-risks` until those projections migrate |
| Promoted-authority Todo claim | TypeScript owns the provider-head read, lifecycle validation, complete-record update, hard-lease check, CAS, receipt, and readback-safe result for claims after authority promotion | Default local Markdown mode remains on the legacy writer; other Todo mutations and Markdown regeneration remain bounded follow-ups |

The scheduler facade exit now includes its first bounded Stage 3 route. A
Expand Down
2 changes: 2 additions & 0 deletions docs/reference/contracts/projection-envelope-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ Python-owned projections do so through `loopx/control_plane/projection_envelope_
| `loopx status --use-projection-cache` hit | `status`, `served_from_cache: true` | as stored |
| `loopx global-summary` | `global_summary` | `global` |
| `loopx global-gates` | `global_gates` | `global` |
| `loopx global-todos` | `global_todos` | `global` |
| `loopx global-risks` | `global_risks` | `global` |

Other read models adopt the envelope in the order listed in the RFC; until
then they carry no freshness guarantee and consumer rule 1 below applies.
Expand Down
146 changes: 139 additions & 7 deletions loopx/global_risks.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,12 @@
from pathlib import Path
from typing import Any, Protocol, cast

from .control_plane.projection_envelope_facts import (
render_projection_envelope_markdown,
seal_projection_envelope,
source_fact,
)

class _CollectStatus(Protocol):
def __call__(
self,
Expand Down Expand Up @@ -157,6 +163,8 @@ def build_global_risks_error(
*,
time_range: str = "24h",
error_code: str = "global_risks_unavailable",
status_payload: dict[str, Any] | None = None,
status_read_at: str | None = None,
) -> dict[str, Any]:
return {
"ok": False,
Expand All @@ -168,6 +176,12 @@ def build_global_risks_error(
"omissions": [
"Raw/private failure details and local paths were intentionally omitted."
],
"projection_envelope": _risks_envelope(
status_payload or {}, status_read_at=status_read_at,
host_poll_status="not_read", host_poll_read_at=None,
host_poll_goal_count=None, source_rows_omitted=0,
shown_count=0, available_count=0, agent_id=None,
),
"boundary": public_safe_boundary(),
}

Expand Down Expand Up @@ -390,14 +404,14 @@ def _collect_stale_host_poll_risks(
registry_path: Path,
scan_limit: int,
warnings: list[dict[str, Any]],
) -> list[dict[str, Any]]:
) -> tuple[list[dict[str, Any]], str, int | None, int]:
"""Scan project-local host poll receipts for loops that died mid-wait."""

risks: list[dict[str, Any]] = []
try:
registry_payload = load_registry(registry_path)
except FileNotFoundError:
return risks
return risks, "missing", None, 0
except (OSError, ValueError) as exc:
warnings.append(
_warning(
Expand All @@ -406,13 +420,15 @@ def _collect_stale_host_poll_risks(
detail=str(exc)[:200],
)
)
return risks
return risks, "unreadable", None, 0
goals = registry_payload.get("goals") if isinstance(registry_payload, dict) else None
if not isinstance(goals, list):
return risks
return risks, "unreadable", None, 0
inspected_count = 0
for source_index, goal in enumerate(goals):
if len(risks) >= scan_limit:
break
inspected_count += 1
if not isinstance(goal, dict) or not goal.get("id"):
continue
state_path = resolve_goal_local_path(
Expand Down Expand Up @@ -460,7 +476,7 @@ def _collect_stale_host_poll_risks(
"requires_user_approval": True,
}
)
return risks
return risks, "read", len(goals), len(goals) - inspected_count

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] This claims the required receipt aggregate was read even when the individual reads were discarded. read_host_poll_receipt returns None for broken JSON, non-object/unsupported schema and I/O errors; invalid timestamps are also skipped. With valid registry/status sources, real CLI probes for broken JSON, [] and wrong schema/time return zero risks with fresh=true, complete=true, alert=false and no warnings. Preserve invalid/unreadable outcomes at the existing host-poll read boundary and pass aggregate unreadable/omission facts to the TS envelope owner. Keep normal absent receipts compatible, and prove mixed-Goal corruption and repair through the real CLI.



def _normalize_stale_warning(
Expand Down Expand Up @@ -688,15 +704,108 @@ def _groups(risks: list[dict[str, Any]]) -> dict[str, list[dict[str, Any]]]:
}


def _risks_envelope(
status_payload: dict[str, Any],
*,
status_read_at: str | None,
host_poll_status: str,
host_poll_read_at: str | None,
host_poll_goal_count: int | None,
source_rows_omitted: int,
host_poll_goals_omitted: int = 0,
shown_count: int,
available_count: int,
agent_id: str | None,
) -> dict[str, Any]:
global_registry = as_dict(status_payload.get("global_registry"))
status_envelope = as_dict(status_payload.get("projection_envelope"))
global_source = next(
(row for row in as_list(status_envelope.get("sources"))
if isinstance(row, dict) and row.get("source_id") == "global_registry"),
{},
)
available = global_registry.get("available") is True
count = global_registry.get("global_goal_count")
expected = count if available and isinstance(count, int) and not isinstance(count, bool) else None
excluded = int(global_registry.get("current_registry_excluded_goal_count") or 0) if available else 0
omitted = []
if excluded:
omitted.append({
"reason": "outside_current_registry", "count": excluded,
"refs": [str(ref) for ref in as_list(global_registry.get("current_registry_excluded_goal_ids"))[:8]],
})
if source_rows_omitted:
omitted.append({"reason": "source_rows_not_scanned", "count": source_rows_omitted})
if host_poll_goals_omitted:
omitted.append({"reason": "host_poll_goals_not_scanned", "count": host_poll_goals_omitted})
if status_payload and status_payload.get("ok") is not True:
omitted.append({"reason": "status_unavailable", "count": 1})
contract = as_dict(status_payload.get("contract"))
queue = as_dict(status_payload.get("attention_queue"))
history = as_dict(status_payload.get("run_history"))
contract_read = isinstance(status_payload.get("contract"), dict) and isinstance(contract.get("error_diagnostics", []), list)
queue_read = isinstance(status_payload.get("attention_queue"), dict) and isinstance(queue.get("items", []), list)
sources = [
source_fact(
"global_registry",
read_status=(
str(global_source.get("read_status") or "read") if available
else str(global_registry.get("read_status") or ("missing" if status_read_at else "not_read"))
),
last_read_at=global_source.get("last_read_at") or (status_read_at if available else None),
item_count=expected,
),
source_fact(
"status_contract",
read_status="read" if contract_read else "not_read",
last_read_at=status_read_at if contract_read else None,
item_count=len(as_list(contract.get("error_diagnostics"))) if contract_read else None,
),
source_fact(
"attention_queue",
read_status="read" if queue_read else "not_read",
last_read_at=status_read_at if queue_read else None,
item_count=len(as_list(queue.get("items"))) if queue_read else None,
),
source_fact(
"host_poll_receipts", read_status=host_poll_status,
last_read_at=host_poll_read_at, item_count=host_poll_goal_count,
),
*([] if status_envelope else [source_fact("status", read_status="not_read")]),
]
if agent_id:
sources.append(source_fact(
"agent_scope",
read_status="read" if isinstance(history.get("goals"), list) else "not_read",
last_read_at=status_read_at if isinstance(history.get("goals"), list) else None,
item_count=len(as_list(history.get("goals"))) if history else None,
))
return seal_projection_envelope(
projection="global_risks", observed_at=now_utc_iso(),
sources=sources,
coverage={
"scope": "global", "expected_count": expected,
"included_count": max(0, expected - excluded) if expected is not None else 0,
"omitted": omitted,
"shown_count": shown_count, "available_count": available_count,
},
upstream=[status_envelope] if status_envelope else [],
)


def _malformed_projection_error(
error: object,
*,
time_range: str,
status_payload: dict[str, Any] | None = None,
status_read_at: str | None = None,
) -> dict[str, Any]:
return build_global_risks_error(
error,
time_range=time_range,
error_code="malformed_status_projection",
status_payload=status_payload,
status_read_at=status_read_at,
)


Expand Down Expand Up @@ -731,6 +840,7 @@ def build_global_risks(
time_range=normalized_time_range,
)
status_payload: dict[str, Any] = status_result
status_read_at = now_utc_iso()
try:
contract = _required_container(status_payload, "contract")
global_registry = _required_container(status_payload, "global_registry")
Expand All @@ -739,7 +849,10 @@ def build_global_risks(
findings = _optional_source_list(global_registry, "findings")
items = _optional_source_list(attention_queue, "items")
except ValueError as exc:
return _malformed_projection_error(exc, time_range=normalized_time_range)
return _malformed_projection_error(
exc, time_range=normalized_time_range,
status_payload=status_payload, status_read_at=status_read_at,
)

warnings: list[dict[str, Any]] = []
risks: list[dict[str, Any]] = []
Expand Down Expand Up @@ -800,12 +913,14 @@ def build_global_risks(
)
)

host_poll_risks = _collect_stale_host_poll_risks(
host_poll_risks, host_poll_status, host_poll_goal_count, host_poll_goals_omitted = _collect_stale_host_poll_risks(
registry_path=registry_path,
scan_limit=scan_limit,
warnings=warnings,
)
host_poll_read_at = now_utc_iso() if host_poll_status == "read" else None
risks.extend(host_poll_risks)
source_rows_truncated = source_rows_truncated or host_poll_goals_omitted > 0

try:
risks, history_truncated = _filter_for_agent(
Expand All @@ -820,6 +935,8 @@ def build_global_risks(
exc,
time_range=normalized_time_range,
error_code="agent_scope_unavailable",
status_payload=status_payload,
status_read_at=status_read_at,
)
source_rows_truncated = source_rows_truncated or history_truncated

Expand Down Expand Up @@ -859,6 +976,19 @@ def build_global_risks(
"risks": retained,
"source_warnings": warnings[:SOURCE_WARNING_LIMIT],
"source_warnings_truncated": warning_count > SOURCE_WARNING_LIMIT,
"projection_envelope": _risks_envelope(
status_payload, status_read_at=status_read_at,
host_poll_status=host_poll_status,
host_poll_read_at=host_poll_read_at,
host_poll_goal_count=host_poll_goal_count,
source_rows_omitted=sum(
max(0, int(warning.get("available_count") or 0) - int(warning.get("inspected_count") or 0))
for warning in warnings if warning.get("reason_code") == "source_rows_truncated"
),
host_poll_goals_omitted=host_poll_goals_omitted,
shown_count=returned_risk_count, available_count=matched_risk_count,
agent_id=agent_id,
),
"omissions": [dict(_ROLLBACK_OMISSION)],
"boundary": public_safe_boundary(),
}
Expand Down Expand Up @@ -911,6 +1041,7 @@ def render_global_risks_markdown(payload: dict[str, Any]) -> str:
"- ok: `False`",
f"- error_code: `{_redact_text(payload.get('error_code'), limit=120)}`",
f"- error: {_redact_text(payload.get('error'))}",
*render_projection_envelope_markdown(payload.get("projection_envelope")),
]
omissions = [
_redact_text(item)
Expand All @@ -932,6 +1063,7 @@ def render_global_risks_markdown(payload: dict[str, Any]) -> str:
f"- matched: `{summary.get('matched_risk_count')}`",
f"- returned: `{summary.get('returned_risk_count')}`",
f"- truncated: `{bool(summary.get('truncated'))}`",
*render_projection_envelope_markdown(payload.get("projection_envelope")),
"",
(
"No current accepted source proves a rollback candidate; "
Expand Down
Loading
Loading