Skip to content
Merged
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
5 changes: 5 additions & 0 deletions loopx/capabilities/reliability_diagnostics/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,11 @@ uncertainty are `degraded`; otherwise the receipt is `valid`.

### Diagnostic projection (`reliability_diagnostic_projection_v0`)

Events are ordered by the instant represented by `observed_at`, then by session
and sequence for equal instants. Different UTC offsets or fractional-second
formats do not change chronology. Receipt bounds retain the original timestamp
text; readback does not rewrite the ledger.

| Field | Meaning |
| --- | --- |
| `mode`, `authority`, `write_scope`, `worker_influence` | `read_only`, `none`, `diagnostic_ledger_only`, `none` |
Expand Down
4 changes: 4 additions & 0 deletions loopx/capabilities/reliability_diagnostics/README.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,10 @@ provider/observer stats 精确关联。stats 按 observer 实例累计;receipt

### Diagnostic projection(`reliability_diagnostic_projection_v0`)

事件按 `observed_at` 表示的实际时刻排序,同一时刻再按 session 和 sequence 排序。
不同 UTC 偏移或小数秒格式不会改变时间顺序。Receipt 的起止时间保留原始时间戳文本,
读回不会重写账本。

| 字段 | 含义 |
| --- | --- |
| `mode`、`authority`、`write_scope`、`worker_influence` | `read_only`、`none`、`diagnostic_ledger_only`、`none` |
Expand Down
26 changes: 25 additions & 1 deletion loopx/capabilities/reliability_diagnostics/envelope.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,10 @@
import re
from collections.abc import Mapping
from dataclasses import dataclass, field
from datetime import datetime
from datetime import datetime, timezone
from decimal import Decimal
from enum import StrEnum
from fractions import Fraction
from typing import Any

from ...control_plane.runtime.public_safety import (
Expand Down Expand Up @@ -277,6 +279,28 @@ def parse_observed_at(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00"))


def observed_at_microseconds(value: str) -> int | Fraction:
"""Exact UTC microseconds, including precision datetime would truncate."""

parsed = parse_observed_at(value)
elapsed = parsed - datetime.min.replace(tzinfo=timezone.utc)
key: int | Fraction = (elapsed.days * 86_400 + elapsed.seconds) * 1_000_000 + elapsed.microseconds
zone_start = max(value.rfind("+"), value.rfind("-"), value.rfind("Z"))
for match in re.finditer(r"[.,](\d+)", value):
digits = match.group(1)
remainder = digits[6:]
offset_fraction = match.start() > zone_start
# fromisoformat also discards the entire fraction of a zero-second offset.
if offset_fraction and not parsed.utcoffset():
adjustment = Fraction(Decimal("0." + digits)) * 1_000_000
elif remainder:
adjustment = Fraction(Decimal("0." + remainder))
else:
continue
key += -adjustment if offset_fraction and value[zone_start] == "+" else adjustment
return key


def _clock(value: Any) -> ObserverClock:
if not isinstance(value, Mapping):
raise ObserverEnvelopeError(
Expand Down
4 changes: 3 additions & 1 deletion loopx/capabilities/reliability_diagnostics/projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,14 @@

from dataclasses import dataclass, field
from enum import StrEnum
from fractions import Fraction
from typing import Any

from .envelope import (
CAPABILITY_ID,
ObserverEnvelope,
ObserverEventKind,
observed_at_microseconds,
parse_observed_at,
)
from .receipt import LedgerReading, build_integrity_receipt
Expand Down Expand Up @@ -75,7 +77,7 @@ def _stage_after(envelope: ObserverEnvelope) -> DiagnosticStage:

def _ms_between(earlier: str, later: str) -> int:
return int(
(parse_observed_at(later) - parse_observed_at(earlier)).total_seconds() * 1000
Fraction(observed_at_microseconds(later) - observed_at_microseconds(earlier), 1000)
)


Expand Down
3 changes: 2 additions & 1 deletion loopx/capabilities/reliability_diagnostics/receipt.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
ObserverEnvelope,
ObserverEnvelopeError,
normalize_observer_envelope,
observed_at_microseconds,
)
from .intake import ObserverStats, normalize_observer_stats

Expand Down Expand Up @@ -90,7 +91,7 @@ class LedgerReading:
def ordered_envelopes(self) -> list[ObserverEnvelope]:
return sorted(
self.envelopes,
key=lambda item: (item.observed_at, item.session_id, item.sequence),
key=lambda item: (observed_at_microseconds(item.observed_at), item.session_id, item.sequence),
)


Expand Down
133 changes: 133 additions & 0 deletions tests/capabilities/test_reliability_diagnostics.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@

import json
import re
import subprocess
import sys
from pathlib import Path
from typing import Any

Expand Down Expand Up @@ -577,6 +579,137 @@ def projection_for(*records: dict[str, Any], **kwargs: Any) -> dict[str, Any]:
return build_diagnostic_projection(read_ledger(records, goal_id=GOAL), **kwargs)


@pytest.mark.parametrize("first,last,gap_ms", [
("2026-09-01T12:00:00+02:00", "2026-09-01T10:01:00Z", 60_000),
("2026-09-02T00:00:00+14:00", "2026-09-01T10:01:00Z", 60_000),
("2026-09-01T10:00:00Z", "2026-09-01T10:00:00.500Z", 500),
("2026-09-01T10:00:00Z", "2026-09-01T10:00:00.000+00:00", 0),
("2026-09-01T10:00:00+00:00", "2026-09-01T10:01:00+00:00", 60_000),
])
def test_diagnostics_order_instants_without_rewriting_timestamps(
first: str, last: str, gap_ms: int,
) -> None:
# UTC instants, then session/sequence, determine order, not text or append order.
records = [
envelope(1, ObserverEventKind.TURN_ENDED, observed_at=last),
envelope(0, ObserverEventKind.AGENT_ERROR, observed_at=first),
stats(accepted_event_count=2),
]
reading = read_ledger(records, goal_id=GOAL)
assert [item.sequence for item in reading.ordered_envelopes] == [0, 1]
receipt = build_integrity_receipt(reading)
assert receipt["status"] == "valid"
assert receipt["observed_from"] == first
assert receipt["observed_until"] == last
projection = build_diagnostic_projection(reading, as_of="2026-09-01T10:06:00Z")
assert projection["stage"] == "idle"
assert projection["recovery"] == {
"error_count": 1, "recovered_error_count": 1, "unrecovered_error_count": 0,
}
assert projection["stall"]["max_inter_event_gap_ms"] == gap_ms
assert projection["stall"]["last_event_age_ms"] == 360_000 - gap_ms
assert projection["signals"] == []


def test_diagnostics_equal_instants_keep_session_then_sequence_order() -> None:
reading = read_ledger([
envelope(0, observed_at="2026-09-01T10:00:00+00:00", session_id="session-b"),
envelope(1, observed_at="2026-09-01T10:00:00.000Z", session_id="session-a"),
envelope(0, observed_at="2026-09-01T12:00:00+02:00", session_id="session-a"),
], goal_id=GOAL)
assert [(item.session_id, item.sequence) for item in reading.ordered_envelopes] == [
("session-a", 0), ("session-a", 1), ("session-b", 0),
]


@pytest.mark.parametrize("earlier,later", [
("2026-09-01T10:00:00.0000001Z", "2026-09-01T10:00:00.0000009Z"),
("20260901T100000,0000001Z", "2026-09-01T12:00:00.0000009+02:00"),
("2026-09-01T10:00:00+00:00:01.0000009", "2026-09-01T10:00:00+00:00:01.0000001"),
("2026-09-01T10:00:00-00:00:01.0000001", "2026-09-01T10:00:00-00:00:01.0000009"),
("2026-09-01T10:00:00+00:00:00.5", "2026-09-01T09:59:59.6Z"),
("2026-09-01T09:59:59.4Z", "2026-09-01T10:00:00+00:00:00.5"),
])
def test_diagnostics_preserves_accepted_fractional_precision(earlier: str, later: str) -> None:
reading = read_ledger([
envelope(0, ObserverEventKind.AGENT_ERROR, observed_at=later),
envelope(1, ObserverEventKind.TURN_ENDED, observed_at=earlier),
stats(accepted_event_count=2),
], goal_id=GOAL)
assert reading.invalid_record_count == 0
assert [item.sequence for item in reading.ordered_envelopes] == [1, 0]
receipt = build_integrity_receipt(reading)
assert (receipt["observed_from"], receipt["observed_until"]) == (earlier, later)
projection = build_diagnostic_projection(reading)
assert projection["stage"] == "errored"
assert projection["recovery"]["unrecovered_error_count"] == 1
assert projection["recovery"]["recovered_error_count"] == 0


@pytest.mark.parametrize("first,last,elapsed_ms", [
("2026-09-01T10:00:00.0000009Z", "2026-09-01T10:00:00.0010001Z", 0),
("2026-09-01T10:00:00+00:00:00.5", "2026-09-01T10:00:00Z", 500),
("2026-09-01T10:00:00Z", "2026-09-01T10:00:00-00:00:00.5", 500),
("0001-01-01T00:00:00Z", "9999-12-31T23:59:59.999999Z", 315_537_897_599_999),
])
def test_diagnostic_age_and_gap_use_the_same_instant_precision(
first: str, last: str, elapsed_ms: int,
) -> None:
projection = projection_for(
envelope(0, ObserverEventKind.STEP_STARTED, observed_at=first),
stats(), as_of=last, stall_threshold_ms=1,
)
assert projection["stall"]["last_event_age_ms"] == elapsed_ms
assert projection["stall"]["detected"] is (elapsed_ms >= 1)
interval = projection_for(
envelope(0, observed_at=first), envelope(1, observed_at=last),
stats(accepted_event_count=2),
)
assert interval["stall"]["max_inter_event_gap_ms"] == elapsed_ms


def test_cli_diagnostics_replays_mixed_offsets_without_mutating_ledger(tmp_path: Path) -> None:
first, recovered, last = (
"2026-09-01T12:00:00+02:00", "2026-09-01T10:01:00Z", "2026-09-01T10:02:00Z",
)
records = [
envelope(0, ObserverEventKind.AGENT_ERROR, observed_at=first),
envelope(1, ObserverEventKind.STEP_ENDED, observed_at=recovered),
envelope(2, ObserverEventKind.TURN_ENDED, observed_at=last),
stats(accepted_event_count=3),
]
command = [
sys.executable, "-m", "loopx.cli", "--registry", str(tmp_path / "registry.json"),
"--runtime-root", str(tmp_path), "--format", "json", "reliability-diagnostics",
]

def run(*args: str, source: str | None = None) -> dict[str, Any]:
result = subprocess.run(
[*command, *args, "--goal-id", GOAL], input=source,
capture_output=True, text=True, encoding="utf-8", check=True, timeout=30,
)
return json.loads(result.stdout)

ingest = run("ingest", "--input", "-", source="\n".join(map(json.dumps, records)))
assert ingest["accepted_envelope_count"] == 3
assert ingest["rejected_event_count"] == 0
path = tmp_path / ingest["ledger_ref"]
before = path.read_bytes()
combined = run("status", "--with-receipt", "--as-of", "2026-09-01T10:10:00Z")
receipt = run("receipt")["receipt"]
assert receipt == combined["receipt"]
assert receipt["status"] == "valid"
assert (receipt["observed_from"], receipt["observed_until"]) == (first, last)
projection = combined["projection"]
assert projection["stage"] == "idle"
assert projection["recovery"]["recovered_error_count"] == 1
assert projection["recovery"]["unrecovered_error_count"] == 0
assert projection["stall"]["last_event_age_ms"] == 480_000
assert projection["signals"] == []
assert projection["authority"] == "none"
assert path.read_bytes() == before


def test_projection_declares_read_only_boundary() -> None:
projection = projection_for(envelope(0), stats())
assert projection["mode"] == "read_only"
Expand Down
Loading