Skip to content

Commit 647756d

Browse files
authored
fix(runtime): validate batched source fingerprint snapshot (#5367)
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 <duanjialing.777@bytedance.com>
1 parent 322b48e commit 647756d

2 files changed

Lines changed: 68 additions & 26 deletions

File tree

‎loopx/control_plane/effect_runtime.py‎

Lines changed: 24 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,12 @@
5858
_RuntimeSourceSnapshot = tuple[tuple[str, int, int, int], ...]
5959

6060

61+
class _RuntimeSourceChanged(RuntimeError):
62+
def __init__(self, snapshot: _RuntimeSourceSnapshot) -> None:
63+
super().__init__("runtime source changed while hashing")
64+
self.snapshot = snapshot
65+
66+
6167
@dataclass(frozen=True)
6268
class _RuntimeRevision:
6369
fingerprint: str
@@ -290,28 +296,35 @@ def _runtime_fingerprint_for_snapshot(
290296
assert read.data is not None
291297
digest.update(relative.encode("utf-8"))
292298
digest.update(read.data)
299+
current_snapshot = _runtime_source_snapshot(source_root)
300+
if current_snapshot != snapshot:
301+
raise _RuntimeSourceChanged(current_snapshot)
293302
return digest.hexdigest()
294303

295304

296305
def _runtime_fingerprint() -> str:
297306
root = _control_plane_root()
298307
resolved_root = os.fspath(root.resolve())
299-
try:
300-
return _runtime_fingerprint_for_snapshot(
301-
resolved_root,
302-
_runtime_source_snapshot(root),
303-
)
304-
except FileNotFoundError:
308+
snapshot: _RuntimeSourceSnapshot | None = None
309+
last_error: Exception | None = None
310+
for _attempt in range(2):
305311
try:
312+
if snapshot is None:
313+
snapshot = _runtime_source_snapshot(root)
306314
return _runtime_fingerprint_for_snapshot(
307315
resolved_root,
308-
_runtime_source_snapshot(root),
316+
snapshot,
309317
)
318+
except _RuntimeSourceChanged as exc:
319+
last_error = exc
320+
snapshot = exc.snapshot
310321
except FileNotFoundError as exc:
311-
raise EffectRuntimeStartupError(
312-
"TypeScript Effect runtime source topology did not stabilize",
313-
diagnostic_code="packaged_runtime_source_unstable",
314-
) from exc
322+
last_error = exc
323+
snapshot = None
324+
raise EffectRuntimeStartupError(
325+
"TypeScript Effect runtime source topology did not stabilize",
326+
diagnostic_code="packaged_runtime_source_unstable",
327+
) from last_error
315328

316329

317330
@contextmanager

‎tests/control_plane/test_turn_journal_runtime_readiness.py‎

Lines changed: 44 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
from __future__ import annotations
22

33
from pathlib import Path
4+
from threading import Event
45
from typing import Any
56

67
import pytest
@@ -104,7 +105,11 @@ def scan_then_remove(root: Path) -> tuple[str, ...]:
104105
monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path)
105106

106107
assert len(effect_runtime._runtime_fingerprint()) == 64
107-
assert scans == [("kept.ts", "removed.ts"), ("kept.ts",)]
108+
assert scans == [
109+
("kept.ts", "removed.ts"),
110+
("kept.ts",),
111+
("kept.ts",),
112+
]
108113

109114

110115
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
117122
later.write_text("export const later = true;\n", encoding="utf-8")
118123
original_read_bytes = Path.read_bytes
119124
reads: list[str] = []
125+
later_prefetched = Event()
120126

121-
def remove_later_after_first_read(path: Path) -> bytes:
122-
reads.append(path.name)
127+
def remove_later_after_prefetch(path: Path) -> bytes:
123128
content = original_read_bytes(path)
124-
if path == first and later.exists():
125-
later.unlink()
129+
if path == later:
130+
reads.append(path.name)
131+
later_prefetched.set()
132+
else:
133+
if later.exists():
134+
assert later_prefetched.wait(timeout=5)
135+
later.unlink()
136+
reads.append(path.name)
126137
return content
127138

128139
monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path)
129-
monkeypatch.setattr(Path, "read_bytes", remove_later_after_first_read)
140+
monkeypatch.setattr(Path, "read_bytes", remove_later_after_prefetch)
130141

131142
assert len(effect_runtime._runtime_fingerprint()) == 64
132-
assert reads == ["first.ts", "later.ts", "first.ts"]
143+
assert reads == ["later.ts", "first.ts", "first.ts"]
133144

134145

135146
def _install_persistent_stat_read_churn(
@@ -139,25 +150,35 @@ def _install_persistent_stat_read_churn(
139150
first = tmp_path / "first.ts"
140151
later = tmp_path / "later.ts"
141152
first.write_text("export const first = true;\n", encoding="utf-8")
153+
later.write_text("export const later = true;\n", encoding="utf-8")
142154
original_scan = effect_runtime._scan_runtime_source_files
143155
original_read_bytes = Path.read_bytes
144156
scans: list[tuple[str, ...]] = []
157+
later_prefetched = Event()
158+
first_reads = 0
145159

146-
def restore_then_scan(root: Path) -> tuple[str, ...]:
147-
later.write_text("export const later = true;\n", encoding="utf-8")
160+
def record_scan(root: Path) -> tuple[str, ...]:
148161
files = original_scan(root)
149162
scans.append(files)
150163
return files
151164

152-
def remove_later_after_first_read(path: Path) -> bytes:
165+
def churn_after_each_first_read(path: Path) -> bytes:
166+
nonlocal first_reads
153167
content = original_read_bytes(path)
154-
if path == first:
168+
if path == later:
169+
later_prefetched.set()
170+
else:
171+
first_reads += 1
172+
if first_reads == 1 and path == first:
173+
assert later_prefetched.wait(timeout=5)
155174
later.unlink()
175+
elif first_reads == 2 and path == first:
176+
later.write_text("export const later = true;\n", encoding="utf-8")
156177
return content
157178

158179
monkeypatch.setattr(effect_runtime, "_control_plane_root", lambda: tmp_path)
159-
monkeypatch.setattr(effect_runtime, "_scan_runtime_source_files", restore_then_scan)
160-
monkeypatch.setattr(Path, "read_bytes", remove_later_after_first_read)
180+
monkeypatch.setattr(effect_runtime, "_scan_runtime_source_files", record_scan)
181+
monkeypatch.setattr(Path, "read_bytes", churn_after_each_first_read)
161182
return scans
162183

163184

@@ -180,7 +201,11 @@ def test_runtime_source_churn_has_a_stable_readiness_diagnostic(
180201
result["runtime_lifecycle"]["diagnostic_code"]
181202
== "packaged_runtime_source_unstable"
182203
)
183-
assert scans == [("first.ts", "later.ts"), ("first.ts", "later.ts")]
204+
assert scans == [
205+
("first.ts", "later.ts"),
206+
("first.ts",),
207+
("first.ts", "later.ts"),
208+
]
184209

185210

186211
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(
193218
effect_runtime.effect_runtime_request("runtime.ping", {})
194219

195220
assert error.value.diagnostic_code == "packaged_runtime_source_unstable"
196-
assert scans == [("first.ts", "later.ts"), ("first.ts", "later.ts")]
221+
assert scans == [
222+
("first.ts", "later.ts"),
223+
("first.ts",),
224+
("first.ts", "later.ts"),
225+
]
197226

198227

199228
def test_missing_node_blocks_the_typescript_control_plane_and_is_actionable(

0 commit comments

Comments
 (0)