From 682c73ffae0913f9e9a4b1f3407502d3faa20ca5 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 23 Sep 2026 15:39:41 +0800 Subject: [PATCH 1/2] fix: order run history by UTC instant Signed-off-by: duanjialing.777 --- loopx/cli_commands/todo.py | 3 +- .../goals/goal_amendment_proposal.py | 3 +- .../runtime/agent_scoped_evidence_log.py | 25 +--- loopx/control_plane/runtime/time.py | 17 +++ loopx/feedback.py | 5 +- loopx/history.py | 20 +-- loopx/state_refresh.py | 7 +- .../test_canonical_planning_consumers.py | 137 ++++++++++++++++++ .../test_goal_amendment_proposal.py | 28 +++- 9 files changed, 200 insertions(+), 45 deletions(-) diff --git a/loopx/cli_commands/todo.py b/loopx/cli_commands/todo.py index d3112f0315..8acac619da 100644 --- a/loopx/cli_commands/todo.py +++ b/loopx/cli_commands/todo.py @@ -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, @@ -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, diff --git a/loopx/control_plane/goals/goal_amendment_proposal.py b/loopx/control_plane/goals/goal_amendment_proposal.py index b0021b61d0..6b15d05afe 100644 --- a/loopx/control_plane/goals/goal_amendment_proposal.py +++ b/loopx/control_plane/goals/goal_amendment_proposal.py @@ -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, ) @@ -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, diff --git a/loopx/control_plane/runtime/agent_scoped_evidence_log.py b/loopx/control_plane/runtime/agent_scoped_evidence_log.py index 7b3c8969fe..bdd15b5ae3 100644 --- a/loopx/control_plane/runtime/agent_scoped_evidence_log.py +++ b/loopx/control_plane/runtime/agent_scoped_evidence_log.py @@ -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: @@ -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 "")) @@ -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 diff --git a/loopx/control_plane/runtime/time.py b/loopx/control_plane/runtime/time.py index 78ea9b8e2c..66282db3a3 100644 --- a/loopx/control_plane/runtime/time.py +++ b/loopx/control_plane/runtime/time.py @@ -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 @@ -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) diff --git a/loopx/feedback.py b/loopx/feedback.py index 1efa077a87..1574984319 100644 --- a/loopx/feedback.py +++ b/loopx/feedback.py @@ -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, @@ -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] diff --git a/loopx/history.py b/loopx/history.py index d8638bb576..3f3a5e722c 100644 --- a/loopx/history.py +++ b/loopx/history.py @@ -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 @@ -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 @@ -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]: diff --git a/loopx/state_refresh.py b/loopx/state_refresh.py index 273b0cfb2d..268e3c4ac9 100644 --- a/loopx/state_refresh.py +++ b/loopx/state_refresh.py @@ -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, @@ -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, ) ] diff --git a/tests/control_plane/test_canonical_planning_consumers.py b/tests/control_plane/test_canonical_planning_consumers.py index 1b70cfb857..1f1ab8e7ba 100644 --- a/tests/control_plane/test_canonical_planning_consumers.py +++ b/tests/control_plane/test_canonical_planning_consumers.py @@ -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 = [ @@ -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 diff --git a/tests/control_plane/test_goal_amendment_proposal.py b/tests/control_plane/test_goal_amendment_proposal.py index 85997cc2f3..4e088c4c14 100644 --- a/tests/control_plane/test_goal_amendment_proposal.py +++ b/tests/control_plane/test_goal_amendment_proposal.py @@ -37,6 +37,7 @@ from loopx.control_plane.goals.shared_goal_alignment import ( project_shared_goal_alignment, ) +from loopx.control_plane.runtime.time import chronology_key from loopx.control_plane.status.autonomous_replan_projection import ( autonomous_replan_obligation_from_runs, ) @@ -250,7 +251,7 @@ def _newest_first_runs( 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, @@ -652,6 +653,31 @@ def test_settlement_ack_run_closes_the_derived_obligation( assert _journal_rows(paths) == [] +def test_later_utc_ack_closes_obligation_across_offsets( + tmp_path: Path, +) -> None: + stalled_runs = _stall_runs() + stalled_runs[0]["generated_at"] = "2026-09-01T08:00:00+08:00" + stalled_runs[1]["generated_at"] = "2026-09-01T08:01:00+08:00" + paths = _write_fixture( + tmp_path, + events=_default_events(), + runs=stalled_runs, + ) + proposal = _proposal(paths) + obligation_id = _derived_obligation(paths)["obligation_id"] + + _append_runs( + paths, + [_ack_run(obligation_id, generated_at="2026-09-01T01:00:00Z")], + ) + + with pytest.raises(ValueError, match="does not match an open replan obligation"): + _admit(paths, proposal) + + assert _journal_rows(paths) == [] + + def test_cross_goal_replan_obligation_fails_closed(tmp_path: Path) -> None: # The sibling Goal's run ledger never contributes to this Goal's # obligation inventory: an id derived on the peer Goal's run history From 5ce1c31dc098f27325c3dc043a4f8de50eeb7468 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 23 Sep 2026 16:09:02 +0800 Subject: [PATCH 2/2] fix: order doctor readiness by UTC instant Signed-off-by: duanjialing.777 --- loopx/doctor.py | 6 +++- tests/test_doctor_promotion_readiness.py | 39 ++++++++++++++++++++++++ 2 files changed, 44 insertions(+), 1 deletion(-) create mode 100644 tests/test_doctor_promotion_readiness.py diff --git a/loopx/doctor.py b/loopx/doctor.py index 1f2e486df7..df4fc19e2d 100644 --- a/loopx/doctor.py +++ b/loopx/doctor.py @@ -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 @@ -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 diff --git a/tests/test_doctor_promotion_readiness.py b/tests/test_doctor_promotion_readiness.py new file mode 100644 index 0000000000..b641e8f90e --- /dev/null +++ b/tests/test_doctor_promotion_readiness.py @@ -0,0 +1,39 @@ +from __future__ import annotations + +import json +from pathlib import Path + +from loopx.control_plane.runtime.promotion_readiness import ( + PROMOTION_READINESS_CLASSIFICATION, + PROMOTION_READINESS_RUNTIME_INDEX, +) +from loopx.doctor import latest_promotion_readiness_event + + +def test_latest_promotion_readiness_uses_utc_instant_across_offsets( + tmp_path: Path, +) -> None: + index_path = tmp_path / PROMOTION_READINESS_RUNTIME_INDEX + index_path.parent.mkdir(parents=True) + rows = [ + { + "classification": PROMOTION_READINESS_CLASSIFICATION, + "generated_at": "2026-01-01T08:30:00+08:00", + "recommended_action": "Older offset readiness.", + }, + { + "classification": PROMOTION_READINESS_CLASSIFICATION, + "generated_at": "2026-01-01T01:00:00Z", + "recommended_action": "Newer UTC readiness.", + }, + ] + index_path.write_text( + "".join(json.dumps(row) + "\n" for row in rows), + encoding="utf-8", + ) + + result = latest_promotion_readiness_event(tmp_path) + + assert result["source"] == "runtime_release_ledger" + assert result["generated_at"] == "2026-01-01T01:00:00Z" + assert result["recommended_action"] == "Newer UTC readiness."