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
3 changes: 2 additions & 1 deletion loopx/cli_commands/todo.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
read_heartbeat_settlement,
settlement_result_payload,
)
from ..control_plane.runtime.time import chronology_key
from ..control_plane.todos.markdown import render_todo_markdown
from ..control_plane.todos.provider_projection import (
project_current_canonical_todos,
Expand Down Expand Up @@ -137,7 +138,7 @@ def _validated_replan_successor_obligation(
for _, run in sorted(
enumerate(existing_runs),
key=lambda item: (
str(item[1].get("generated_at") or ""),
*chronology_key(item[1].get("generated_at")),
item[0],
),
reverse=True,
Expand Down
3 changes: 2 additions & 1 deletion loopx/control_plane/goals/goal_amendment_proposal.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@
from ...history import load_index, load_registry
from ...runtime import validate_goal_id_path_segment
from ..effect_runtime import EffectRuntimeRejected, effect_runtime_result
from ..runtime.time import chronology_key
from ..status.autonomous_replan_projection import (
autonomous_replan_obligation_from_runs,
)
Expand Down Expand Up @@ -315,7 +316,7 @@ def _derive_open_replan_obligation_inventory(
for _, run in sorted(
enumerate(runs),
key=lambda item: (
str(item[1].get("generated_at") or ""),
*chronology_key(item[1].get("generated_at")),
item[0],
),
reverse=True,
Expand Down
25 changes: 5 additions & 20 deletions loopx/control_plane/runtime/agent_scoped_evidence_log.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,19 +2,18 @@

import shlex
from collections.abc import Iterable, Mapping
from datetime import datetime, timezone
from datetime import datetime
from typing import Any

from ..todos.contract import normalize_todo_id_list
from .public_safety import public_safe_compact_text
from .time import parse_timestamp
from .time import chronology_key, parse_timestamp


SCHEMA_VERSION = "agent_scoped_evidence_log_v0"
REQUIRED_READ_SCHEMA_VERSION = "loopx_agent_required_read_v0"
READ_RECEIPT_SCHEMA_VERSION = "evidence_log_read_receipt_v0"
MAX_PROJECTED_READ_RECEIPTS = 12
_MIN_TIMESTAMP = datetime.min.replace(tzinfo=timezone.utc)


def _compact_text(value: Any, *, limit: int = 220) -> str | None:
Expand Down Expand Up @@ -211,22 +210,8 @@ def _run_matches(
return True


def _chronology_key(value: Any) -> tuple[int, datetime, str]:
raw = str(value or "")
try:
parsed = parse_timestamp(value)
except OverflowError:
# UTC conversion can overflow at datetime's representable boundaries.
parsed = None
if parsed is None:
# Keep malformed legacy rows deterministic, but never let them outrank
# a row with a valid timestamp.
return (0, _MIN_TIMESTAMP, raw)
return (1, parsed, raw)


def _sort_key(row: Mapping[str, Any]) -> tuple[int, datetime, str, str]:
rank, recorded_at, raw = _chronology_key(row.get("recorded_at"))
rank, recorded_at, raw = chronology_key(row.get("recorded_at"))
return (rank, recorded_at, raw, str(row.get("source") or ""))


Expand Down Expand Up @@ -347,9 +332,9 @@ def _other_agent_frontier(
if not other_agent or other_agent == agent_id:
continue
current = latest_by_agent.get(other_agent)
is_newer = current is None or _chronology_key(
is_newer = current is None or chronology_key(
run.get("generated_at")
) > _chronology_key(current.get("recorded_at"))
) > chronology_key(current.get("recorded_at"))
if is_newer:
row = _safe_run_history_row(run)
row["agent_id"] = other_agent
Expand Down
17 changes: 17 additions & 0 deletions loopx/control_plane/runtime/time.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@
from typing import Any


_MIN_TIMESTAMP = datetime.min.replace(tzinfo=timezone.utc)


def parse_timestamp(value: Any) -> datetime | None:
if not value:
return None
Expand All @@ -19,6 +22,20 @@ def parse_timestamp(value: Any) -> datetime | None:
return parsed.astimezone(timezone.utc)


def chronology_key(value: Any) -> tuple[int, datetime, str]:
"""Order timestamps by UTC instant with deterministic legacy fallbacks."""

raw = str(value or "")
try:
parsed = parse_timestamp(value)
except OverflowError:
# UTC conversion can overflow at datetime's representable boundaries.
parsed = None
if parsed is None:
return (0, _MIN_TIMESTAMP, raw)
return (1, parsed, raw)


def now_utc() -> datetime:
return datetime.now(timezone.utc).replace(microsecond=0)

Expand Down
6 changes: 5 additions & 1 deletion loopx/doctor.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
PROMOTION_READINESS_CLASSIFICATION,
PROMOTION_READINESS_RUNTIME_INDEX,
)
from .control_plane.runtime.time import chronology_key
from .install_contract import NO_CLONE_INSTALL_URL
from .paths import DEFAULT_RUNTIME_ROOT, global_registry_path
from .python_install_owner import PythonInstallOwner, python_distribution_upgrade_command, resolve_python_install_owner
Expand Down Expand Up @@ -678,7 +679,10 @@ def latest_promotion_readiness_event(runtime_root: Path, goal_id: str | None = N
else "no canary promotion readiness run found"
),
}
matches.sort(key=lambda item: str(item.get("generated_at") or ""), reverse=True)
matches.sort(
key=lambda item: chronology_key(item.get("generated_at")),
reverse=True,
)
latest = matches[0]
latest["runtime_root"] = str(runtime_root)
return latest
Expand Down
5 changes: 3 additions & 2 deletions loopx/feedback.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@
from pathlib import Path
from typing import Any

from .history import _chronology_key, load_index, load_registry
from .control_plane.runtime.time import chronology_key
from .history import load_index, load_registry
from .paths import resolve_runtime_root
from .public_safe_text import (
PRIVATE_TEXT_PATTERNS as SHARED_PRIVATE_TEXT_PATTERNS,
Expand Down Expand Up @@ -215,7 +216,7 @@ def select_run(runs: list[dict[str, Any]], run_generated_at: str | None) -> dict
return max(
enumerate(runs),
key=lambda item: (
*_chronology_key(item[1].get("generated_at")),
*chronology_key(item[1].get("generated_at")),
item[0],
),
)[1]
Expand Down
20 changes: 2 additions & 18 deletions loopx/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@
from contextlib import nullcontext
from collections.abc import Callable
from dataclasses import dataclass
from datetime import datetime, timezone
from heapq import merge
from itertools import islice
from pathlib import Path
Expand Down Expand Up @@ -47,7 +46,7 @@
collision_review_groups,
validate_reviewed_collision_plan,
)
from .control_plane.runtime.time import now_local_iso, parse_timestamp
from .control_plane.runtime.time import chronology_key, now_local_iso
from .doctor import PROMOTION_READINESS_CLASSIFICATIONS
from .execution_profile import compact_execution_profile
from .explore_graph import compact_explore_graph_policy
Expand Down Expand Up @@ -129,22 +128,7 @@ def now_local() -> str:
return now_local_iso()


_MIN_TIMESTAMP = datetime.min.replace(tzinfo=timezone.utc)


def _chronology_key(value: Any) -> tuple[int, datetime, str]:
"""Return a UTC-aware ordering key while keeping legacy rows deterministic."""

raw = str(value or "")
try:
parsed = parse_timestamp(value)
except OverflowError:
# UTC conversion can overflow at datetime's representable boundaries.
parsed = None
if parsed is None:
# Malformed or missing legacy rows must never outrank valid timestamps.
return (0, _MIN_TIMESTAMP, raw)
return (1, parsed, raw)
_chronology_key = chronology_key


def unique_run_paths(runs_dir: Path, generated_at: str) -> tuple[Path, Path]:
Expand Down
7 changes: 5 additions & 2 deletions loopx/state_refresh.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from pathlib import Path
from typing import Any

from .control_plane.runtime.time import now_local_iso
from .control_plane.runtime.time import chronology_key, now_local_iso
from .control_plane.work_items.delivery_history import require_consistent_delivery_claim
from .control_plane.work_items.delivery_batch_scale import (
DELIVERY_BATCH_SCALE_CHOICES as DELIVERY_BATCH_SCALE_CHOICES,
Expand Down Expand Up @@ -1026,7 +1026,10 @@ def refresh_state_run(
run
for _, run in sorted(
enumerate(existing_runs),
key=lambda item: (str(item[1].get("generated_at") or ""), item[0]),
key=lambda item: (
*chronology_key(item[1].get("generated_at")),
item[0],
),
reverse=True,
)
]
Expand Down
137 changes: 137 additions & 0 deletions tests/control_plane/test_canonical_planning_consumers.py
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,15 @@ def _fixture(root: Path, *, state: str | None = None) -> tuple[Path, Path, dict]
return registry, state_path, goal


def _write_run_index(registry: Path, runs: list[dict[str, object]]) -> None:
index_path = registry.parent / "runtime/goals/goal-a/runs/index.jsonl"
index_path.parent.mkdir(parents=True, exist_ok=True)
index_path.write_text(
"".join(json.dumps(run, sort_keys=True) + "\n" for run in runs),
encoding="utf-8",
)


def _promote(registry: Path, path: Path, goal: dict) -> dict:
fields = parse_active_state_todos(path.read_text(), goal=goal, item_limit=None)
todos = [
Expand Down Expand Up @@ -299,6 +308,134 @@ def test_todo_add_replan_guard_binds_canonical_obligation_without_display(
assert not path.exists()


def test_refresh_uses_newest_vision_by_utc_instant_across_offsets(
tmp_path: Path,
) -> None:
registry, _, _ = _fixture(tmp_path)
_write_run_index(
registry,
[
{
"generated_at": "2026-01-01T08:30:00+08:00",
"goal_id": "goal-a",
"agent_id": "agent-a",
"agent_vision": {
"agent_id": "agent-a",
"state": "vision_active",
"vision_patch": {"vision_summary": "Older offset vision."},
},
},
{
"generated_at": "2026-01-01T01:00:00Z",
"goal_id": "goal-a",
"agent_id": "agent-a",
"agent_vision": {
"agent_id": "agent-a",
"state": "vision_active",
"vision_patch": {"vision_summary": "Newer UTC vision."},
},
},
],
)

result = _refresh(
registry,
vision_unchanged_reason="Validated evidence keeps the newer vision unchanged.",
)

assert result["vision_checkpoint"]["continuity_basis"] == {
"kind": "existing_vision_unchanged",
"vision_generated_at": "2026-01-01T01:00:00Z",
}


def test_todo_replan_guard_closes_later_utc_ack_across_offsets(
tmp_path: Path,
) -> None:
from argparse import Namespace
from loopx.cli_commands.todo import _validated_replan_successor_obligation

registry, _, goal = _fixture(tmp_path)
stalled_runs = [
{
"generated_at": f"2026-01-01T08:0{index}:00+08:00",
"goal_id": "goal-a",
"agent_id": "agent-a",
"classification": "bounded_replan_progress",
"turn_instance_id": f"turn-stalled-{index}",
"progress_observation": {
"schema_version": "typed_progress_observation_v0",
"result_class": "blocked",
"surface_id": "surface-a",
"hypothesis_id": "hypothesis-a",
"probe_kind": "probe-a",
"evidence_ids": ["evidence-a"],
},
}
for index in range(2)
]
obligation, _ = qualify_replan_writeback(
todo_fields={},
newest_first_runs=list(reversed(stalled_runs)),
state_text="",
agent_id="agent-a",
goal_id="goal-a",
registry_goal=goal,
)
assert obligation is not None
_write_run_index(
registry,
[
*stalled_runs,
{
"generated_at": "2026-01-01T01:00:00Z",
"goal_id": "goal-a",
"agent_id": "agent-a",
"classification": "autonomous_replan_recorded",
"turn_instance_id": "turn-replan-ack",
"autonomous_replan_ack": {
"schema_version": "autonomous_replan_ack_v0",
"recorded": True,
"source": "refresh_state",
"semantic_delta": {
"schema_version": "replan_semantic_delta_v0",
"accepted": True,
"obligation_id": obligation["obligation_id"],
"outcomes": ["new_runnable_successor"],
},
},
"progress_observation": {
"schema_version": "typed_progress_observation_v0",
"result_class": "advanced",
"surface_id": "surface-a",
"hypothesis_id": "hypothesis-b",
"probe_kind": "probe-a",
"evidence_ids": ["evidence-b"],
},
},
],
)
args = Namespace(
replan_obligation_id=obligation["obligation_id"],
role="agent",
task_class="advancement_task",
claimed_by="agent-a",
action_kind="implementation",
monitor_target_key="feature-a",
explore_result_node_refs=[],
goal_id="goal-a",
project=None,
state_file=None,
)

with pytest.raises(ValueError, match="no open replan obligation"):
_validated_replan_successor_obligation(
args,
registry_path=registry,
runtime_root_arg=None,
)


@pytest.mark.parametrize("promoted", [False, True])
def test_public_refresh_retains_legacy_parity_and_uses_one_provider_read(
tmp_path: Path, monkeypatch, promoted: bool
Expand Down
Loading
Loading