From f5995f102802a3261d559989ae27f2aeef4925f0 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Thu, 1 Oct 2026 02:46:36 +0800 Subject: [PATCH] fix(runtime): validate batched source fingerprint snapshot Concurrent prefetch could read a later source before an earlier read removed it, so the batch returned a fingerprint for a topology that no longer existed. Revalidate the snapshot after each uncached batch and retry once, while deterministic event-based tests cover one-shot and persistent churn. Signed-off-by: duanjialing.777 --- loopx/control_plane/effect_runtime.py | 35 +++++++---- .../test_turn_journal_runtime_readiness.py | 59 ++++++++++++++----- 2 files changed, 68 insertions(+), 26 deletions(-) diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index 7a7d55d73a..a815707358 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -58,6 +58,12 @@ _RuntimeSourceSnapshot = tuple[tuple[str, int, int, int], ...] +class _RuntimeSourceChanged(RuntimeError): + def __init__(self, snapshot: _RuntimeSourceSnapshot) -> None: + super().__init__("runtime source changed while hashing") + self.snapshot = snapshot + + @dataclass(frozen=True) class _RuntimeRevision: fingerprint: str @@ -290,28 +296,35 @@ def _runtime_fingerprint_for_snapshot( assert read.data is not None digest.update(relative.encode("utf-8")) digest.update(read.data) + current_snapshot = _runtime_source_snapshot(source_root) + if current_snapshot != snapshot: + raise _RuntimeSourceChanged(current_snapshot) return digest.hexdigest() def _runtime_fingerprint() -> str: root = _control_plane_root() resolved_root = os.fspath(root.resolve()) - try: - return _runtime_fingerprint_for_snapshot( - resolved_root, - _runtime_source_snapshot(root), - ) - except FileNotFoundError: + snapshot: _RuntimeSourceSnapshot | None = None + last_error: Exception | None = None + for _attempt in range(2): try: + if snapshot is None: + snapshot = _runtime_source_snapshot(root) return _runtime_fingerprint_for_snapshot( resolved_root, - _runtime_source_snapshot(root), + snapshot, ) + except _RuntimeSourceChanged as exc: + last_error = exc + snapshot = exc.snapshot except FileNotFoundError as exc: - raise EffectRuntimeStartupError( - "TypeScript Effect runtime source topology did not stabilize", - diagnostic_code="packaged_runtime_source_unstable", - ) from exc + last_error = exc + snapshot = None + raise EffectRuntimeStartupError( + "TypeScript Effect runtime source topology did not stabilize", + diagnostic_code="packaged_runtime_source_unstable", + ) from last_error @contextmanager diff --git a/tests/control_plane/test_turn_journal_runtime_readiness.py b/tests/control_plane/test_turn_journal_runtime_readiness.py index cb0a2b0165..9877af979e 100644 --- a/tests/control_plane/test_turn_journal_runtime_readiness.py +++ b/tests/control_plane/test_turn_journal_runtime_readiness.py @@ -1,6 +1,7 @@ from __future__ import annotations from pathlib import Path +from threading import Event from typing import Any import pytest @@ -104,7 +105,11 @@ def scan_then_remove(root: Path) -> tuple[str, ...]: monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path) assert len(effect_runtime._runtime_fingerprint()) == 64 - assert scans == [("kept.ts", "removed.ts"), ("kept.ts",)] + assert scans == [ + ("kept.ts", "removed.ts"), + ("kept.ts",), + ("kept.ts",), + ] def test_runtime_fingerprint_rescans_when_a_snapshotted_file_disappears_while_reading( @@ -117,19 +122,25 @@ def test_runtime_fingerprint_rescans_when_a_snapshotted_file_disappears_while_re later.write_text("export const later = true;\n", encoding="utf-8") original_read_bytes = Path.read_bytes reads: list[str] = [] + later_prefetched = Event() - def remove_later_after_first_read(path: Path) -> bytes: - reads.append(path.name) + def remove_later_after_prefetch(path: Path) -> bytes: content = original_read_bytes(path) - if path == first and later.exists(): - later.unlink() + if path == later: + reads.append(path.name) + later_prefetched.set() + else: + if later.exists(): + assert later_prefetched.wait(timeout=5) + later.unlink() + reads.append(path.name) return content monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path) - monkeypatch.setattr(Path, "read_bytes", remove_later_after_first_read) + monkeypatch.setattr(Path, "read_bytes", remove_later_after_prefetch) assert len(effect_runtime._runtime_fingerprint()) == 64 - assert reads == ["first.ts", "later.ts", "first.ts"] + assert reads == ["later.ts", "first.ts", "first.ts"] def _install_persistent_stat_read_churn( @@ -139,25 +150,35 @@ def _install_persistent_stat_read_churn( first = tmp_path / "first.ts" later = tmp_path / "later.ts" first.write_text("export const first = true;\n", encoding="utf-8") + later.write_text("export const later = true;\n", encoding="utf-8") original_scan = effect_runtime._scan_runtime_source_files original_read_bytes = Path.read_bytes scans: list[tuple[str, ...]] = [] + later_prefetched = Event() + first_reads = 0 - def restore_then_scan(root: Path) -> tuple[str, ...]: - later.write_text("export const later = true;\n", encoding="utf-8") + def record_scan(root: Path) -> tuple[str, ...]: files = original_scan(root) scans.append(files) return files - def remove_later_after_first_read(path: Path) -> bytes: + def churn_after_each_first_read(path: Path) -> bytes: + nonlocal first_reads content = original_read_bytes(path) - if path == first: + if path == later: + later_prefetched.set() + else: + first_reads += 1 + if first_reads == 1 and path == first: + assert later_prefetched.wait(timeout=5) later.unlink() + elif first_reads == 2 and path == first: + later.write_text("export const later = true;\n", encoding="utf-8") return content monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path) - monkeypatch.setattr(effect_runtime, "_scan_runtime_source_files", restore_then_scan) - monkeypatch.setattr(Path, "read_bytes", remove_later_after_first_read) + monkeypatch.setattr(effect_runtime, "_scan_runtime_source_files", record_scan) + monkeypatch.setattr(Path, "read_bytes", churn_after_each_first_read) return scans @@ -180,7 +201,11 @@ def test_runtime_source_churn_has_a_stable_readiness_diagnostic( result["runtime_lifecycle"]["diagnostic_code"] == "packaged_runtime_source_unstable" ) - assert scans == [("first.ts", "later.ts"), ("first.ts", "later.ts")] + assert scans == [ + ("first.ts", "later.ts"), + ("first.ts",), + ("first.ts", "later.ts"), + ] def test_runtime_request_source_churn_raises_a_stable_startup_diagnostic( @@ -193,7 +218,11 @@ def test_runtime_request_source_churn_raises_a_stable_startup_diagnostic( effect_runtime.effect_runtime_request("runtime.ping", {}) assert error.value.diagnostic_code == "packaged_runtime_source_unstable" - assert scans == [("first.ts", "later.ts"), ("first.ts", "later.ts")] + assert scans == [ + ("first.ts", "later.ts"), + ("first.ts",), + ("first.ts", "later.ts"), + ] def test_missing_node_blocks_the_typescript_control_plane_and_is_actionable(