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
35 changes: 24 additions & 11 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
59 changes: 44 additions & 15 deletions tests/control_plane/test_turn_journal_runtime_readiness.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from __future__ import annotations

from pathlib import Path
from threading import Event
from typing import Any

import pytest
Expand Down Expand Up @@ -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(
Expand All @@ -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(
Expand All @@ -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


Expand All @@ -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(
Expand All @@ -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(
Expand Down
Loading