From 5043e1b02ebabf096be4b4c24403c5a80f2f56ee Mon Sep 17 00:00:00 2001 From: NIU-123370 <191000457+NIU-123370@users.noreply.github.com> Date: Wed, 9 Sep 2026 10:40:57 +0800 Subject: [PATCH] fix(control-plane): validate managed runtime idle timeout Signed-off-by: NIU-123370 <191000457+NIU-123370@users.noreply.github.com> --- loopx/control_plane/effect_runtime.py | 92 +++++++++++++++++-- loopx/control_plane/effect_runtime_config.ts | 22 +++++ loopx/control_plane/effect_runtime_server.ts | 5 +- .../test_effect_runtime_integration.py | 89 +++++++++++++++++- .../effect_runtime_config.test.ts | 38 ++++++++ tests/fixtures/effect_runtime_idle_ms.json | 10 ++ 6 files changed, 246 insertions(+), 10 deletions(-) create mode 100644 loopx/control_plane/effect_runtime_config.ts create mode 100644 tests/control_plane_ts/effect_runtime_config.test.ts create mode 100644 tests/fixtures/effect_runtime_idle_ms.json diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index b6a04f75d7..af116b983d 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -26,6 +26,8 @@ MINIMUM_NODE_VERSION_TEXT = ".".join(str(part) for part in MINIMUM_NODE_VERSION) MAX_RESPONSE_BYTES = 2 * 1024 * 1024 MAX_REQUEST_BYTES = 2 * 1024 * 1024 +DEFAULT_EFFECT_RUNTIME_IDLE_MS = 5 * 60 * 1_000 +MAX_EFFECT_RUNTIME_IDLE_MS = 2_147_483_647 STARTUP_LOCK_TIMEOUT_SECONDS = 15.0 STARTUP_READY_TIMEOUT_SECONDS = 15.0 STARTUP_POLL_SECONDS = 0.025 @@ -34,6 +36,38 @@ _RuntimeSourceSnapshot = tuple[tuple[str, int, int, int], ...] +def _validate_effect_runtime_idle_ms(environment: Mapping[str, str]) -> int: + raw = environment.get("LOOPX_EFFECT_RUNTIME_IDLE_MS") + if raw is None: + return DEFAULT_EFFECT_RUNTIME_IDLE_MS + normalized = raw.strip() + if not re.fullmatch(r"[0-9]+", normalized): + raise EffectRuntimeStartupError( + "LOOPX_EFFECT_RUNTIME_IDLE_MS must be a positive base-10 integer", + diagnostic_code="invalid_runtime_idle_ms", + ) + parsed = int(normalized) + if not 1 <= parsed <= MAX_EFFECT_RUNTIME_IDLE_MS: + raise EffectRuntimeStartupError( + "LOOPX_EFFECT_RUNTIME_IDLE_MS must be between " + f"1 and {MAX_EFFECT_RUNTIME_IDLE_MS}", + diagnostic_code="invalid_runtime_idle_ms", + ) + return parsed + + +def _require_matching_runtime_idle_ms( + info: Mapping[str, Any], + *, + requested_idle_ms: int, +) -> None: + if info.get("idle_ms") != requested_idle_ms: + raise EffectRuntimeStartupError( + "LOOPX_EFFECT_RUNTIME_IDLE_MS conflicts with the running managed runtime", + diagnostic_code="runtime_idle_ms_conflict", + ) + + class EffectRuntimeRemoteError(RuntimeError): """A typed exception returned by the managed TypeScript runtime.""" @@ -305,6 +339,9 @@ def _read_info(path: Path, *, fingerprint: str) -> dict[str, Any] | None: or payload.get("host") != "127.0.0.1" or not isinstance(payload.get("port"), int) or not isinstance(payload.get("token"), str) + or isinstance(payload.get("idle_ms"), bool) + or not isinstance(payload.get("idle_ms"), int) + or not 1 <= payload["idle_ms"] <= MAX_EFFECT_RUNTIME_IDLE_MS or not _pid_is_alive(payload.get("pid")) ): return None @@ -407,7 +444,14 @@ def _remote_runtime_error(value: object) -> EffectRuntimeRemoteError: ) -def _start_runtime(*, fingerprint: str, info_path: Path) -> dict[str, Any]: +def _start_runtime( + *, + fingerprint: str, + info_path: Path, + requested_idle_ms: int | None = None, +) -> dict[str, Any]: + if requested_idle_ms is None: + requested_idle_ms = _validate_effect_runtime_idle_ms(os.environ) runtime_dir = info_path.parent runtime_dir.mkdir(parents=True, exist_ok=True, mode=0o700) try: @@ -429,6 +473,10 @@ def _start_runtime(*, fingerprint: str, info_path: Path) -> dict[str, Any]: except FileExistsError: existing = _read_info(info_path, fingerprint=fingerprint) if existing is not None: + _require_matching_runtime_idle_ms( + existing, + requested_idle_ms=requested_idle_ms, + ) return existing holder_pid = _start_lock_holder_pid(lock) if holder_pid is not None and not _pid_is_alive(holder_pid): @@ -446,6 +494,10 @@ def _start_runtime(*, fingerprint: str, info_path: Path) -> dict[str, Any]: if not acquired: existing = _read_info(info_path, fingerprint=fingerprint) if existing is not None: + _require_matching_runtime_idle_ms( + existing, + requested_idle_ms=requested_idle_ms, + ) return existing raise EffectRuntimeStartupError( "TypeScript Effect runtime startup lock timed out", @@ -454,6 +506,10 @@ def _start_runtime(*, fingerprint: str, info_path: Path) -> dict[str, Any]: try: existing = _read_info(info_path, fingerprint=fingerprint) if existing is not None: + _require_matching_runtime_idle_ms( + existing, + requested_idle_ms=requested_idle_ms, + ) return existing token = secrets.token_urlsafe(32) environment = os.environ.copy() @@ -486,6 +542,10 @@ def _start_runtime(*, fingerprint: str, info_path: Path) -> dict[str, Any]: while time.monotonic() < ready_deadline: info = _read_info(info_path, fingerprint=fingerprint) if info is not None: + _require_matching_runtime_idle_ms( + info, + requested_idle_ms=requested_idle_ms, + ) return info exit_code = process.poll() if exit_code is not None: @@ -514,6 +574,7 @@ def effect_runtime_request( ) -> dict[str, Any]: """Call the managed TS runtime, retrying only idempotent typed effects.""" + requested_idle_ms = _validate_effect_runtime_idle_ms(os.environ) fingerprint = _runtime_fingerprint() info_path = _runtime_info_path(fingerprint) request_id = str(uuid.uuid4()) @@ -522,7 +583,16 @@ def effect_runtime_request( try: info = _read_info(info_path, fingerprint=fingerprint) if info is None: - info = _start_runtime(fingerprint=fingerprint, info_path=info_path) + info = _start_runtime( + fingerprint=fingerprint, + info_path=info_path, + requested_idle_ms=requested_idle_ms, + ) + else: + _require_matching_runtime_idle_ms( + info, + requested_idle_ms=requested_idle_ms, + ) return _request_with_info( info, request_id=request_id, @@ -534,6 +604,8 @@ def effect_runtime_request( raise except EffectRuntimeStartupError as exc: last_error = exc + if exc.diagnostic_code == "runtime_idle_ms_conflict": + raise if attempt == 0 and retry_safe: info_path.unlink(missing_ok=True) continue @@ -574,14 +646,20 @@ def collect_effect_runtime_readiness(*, deep: bool = False) -> dict[str, object] runtime_diagnostic_code: str | None = None if ready: try: + requested_idle_ms = _validate_effect_runtime_idle_ms(os.environ) fingerprint = _runtime_fingerprint() + info = _read_info( + _runtime_info_path(fingerprint), + fingerprint=fingerprint, + ) + if info is not None: + _require_matching_runtime_idle_ms( + info, + requested_idle_ms=requested_idle_ms, + ) runtime_state = ( "running" - if _read_info( - _runtime_info_path(fingerprint), - fingerprint=fingerprint, - ) - is not None + if info is not None else "stopped" ) except (OSError, EffectRuntimeStartupError) as exc: diff --git a/loopx/control_plane/effect_runtime_config.ts b/loopx/control_plane/effect_runtime_config.ts new file mode 100644 index 0000000000..d6e57efa38 --- /dev/null +++ b/loopx/control_plane/effect_runtime_config.ts @@ -0,0 +1,22 @@ +export const DEFAULT_EFFECT_RUNTIME_IDLE_MS = 5 * 60 * 1_000; +export const MAX_EFFECT_RUNTIME_IDLE_MS = 2_147_483_647; + +export function effectRuntimeIdleMs(value: string | undefined): number { + if (value === undefined) return DEFAULT_EFFECT_RUNTIME_IDLE_MS; + const normalized = value.trim(); + if (!/^\d+$/u.test(normalized)) { + throw new Error( + "LOOPX_EFFECT_RUNTIME_IDLE_MS must be a positive base-10 integer", + ); + } + const parsed = Number(normalized); + if ( + !Number.isSafeInteger(parsed) || parsed < 1 || + parsed > MAX_EFFECT_RUNTIME_IDLE_MS + ) { + throw new Error( + `LOOPX_EFFECT_RUNTIME_IDLE_MS must be between 1 and ${MAX_EFFECT_RUNTIME_IDLE_MS}`, + ); + } + return parsed; +} diff --git a/loopx/control_plane/effect_runtime_server.ts b/loopx/control_plane/effect_runtime_server.ts index 69dedb11e6..2d81c23537 100644 --- a/loopx/control_plane/effect_runtime_server.ts +++ b/loopx/control_plane/effect_runtime_server.ts @@ -10,6 +10,7 @@ import { EffectRuntimeRequestError, effectRuntimeErrorPayload, } from "./effect_runtime_errors.ts"; +import { effectRuntimeIdleMs } from "./effect_runtime_config.ts"; import { atomicWriteJson } from "./effect_runtime_io.ts"; import { requireJsonObject as requiredObject, @@ -20,7 +21,6 @@ const REQUEST_SCHEMA = "loopx_effect_runtime_request_v0"; const RESPONSE_SCHEMA = "loopx_effect_runtime_response_v1"; const INFO_SCHEMA = "loopx_effect_runtime_info_v0"; const MAX_REQUEST_BYTES = 2 * 1024 * 1024; -const DEFAULT_IDLE_MS = 5 * 60 * 1_000; let shutdownRequested = false; function asObject(value: unknown): JsonObject { @@ -40,7 +40,7 @@ function parseArg(name: string): string { const infoPath = parseArg("--info"); const fingerprint = parseArg("--fingerprint"); const token = requiredString(process.env.LOOPX_EFFECT_RUNTIME_TOKEN, "runtime token"); -const idleMs = Number(process.env.LOOPX_EFFECT_RUNTIME_IDLE_MS ?? DEFAULT_IDLE_MS); +const idleMs = effectRuntimeIdleMs(process.env.LOOPX_EFFECT_RUNTIME_IDLE_MS); let idleTimer: NodeJS.Timeout; const handlers = createEffectRuntimeHandlers({ fingerprint, @@ -131,6 +131,7 @@ server.listen(0, "127.0.0.1", async () => { host: "127.0.0.1", port: address.port, token, + idle_ms: idleMs, }); await chmod(infoPath, 0o600); resetIdleTimer(server); diff --git a/tests/control_plane/test_effect_runtime_integration.py b/tests/control_plane/test_effect_runtime_integration.py index 47fbc5230f..352535a551 100644 --- a/tests/control_plane/test_effect_runtime_integration.py +++ b/tests/control_plane/test_effect_runtime_integration.py @@ -26,6 +26,32 @@ _TURN_KEY = "sha256:" + "a" * 64 _TODO_ID = "todo_fixture0001" +_IDLE_TIMEOUT_FIXTURE = json.loads( + (Path(__file__).parents[1] / "fixtures" / "effect_runtime_idle_ms.json").read_text( + encoding="utf-8" + ) +) + + +def test_idle_timeout_parser_matches_shared_fixture() -> None: + assert ( + effect_runtime.DEFAULT_EFFECT_RUNTIME_IDLE_MS + == _IDLE_TIMEOUT_FIXTURE["default_ms"] + ) + assert ( + effect_runtime.MAX_EFFECT_RUNTIME_IDLE_MS + == _IDLE_TIMEOUT_FIXTURE["maximum_ms"] + ) + for case in _IDLE_TIMEOUT_FIXTURE["valid"]: + environment = ( + {} + if case["raw"] is None + else {"LOOPX_EFFECT_RUNTIME_IDLE_MS": case["raw"]} + ) + assert ( + effect_runtime._validate_effect_runtime_idle_ms(environment) + == case["value"] + ) def _effect_id(label: str) -> str: @@ -593,6 +619,7 @@ def test_runtime_ready_budget_starts_after_start_lock_acquisition( "host": "127.0.0.1", "port": 1, "token": "fixture-token", + "idle_ms": effect_runtime.DEFAULT_EFFECT_RUNTIME_IDLE_MS, } def sleep(seconds: float) -> None: @@ -674,6 +701,59 @@ def popen(*_args, **kwargs): assert launch["close_fds"] is True +@pytest.mark.parametrize( + "idle_ms", + _IDLE_TIMEOUT_FIXTURE["invalid"], +) +def test_invalid_idle_timeout_fails_before_runtime_launch( + tmp_path: Path, + monkeypatch, + idle_ms: str, +) -> None: + info_path = tmp_path / "runtime" / "effect-runtime-test.json" + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", idle_ms) + + def unexpected_popen(*_args, **_kwargs): + raise AssertionError("invalid idle timeout must fail before process launch") + + monkeypatch.setattr(effect_runtime.subprocess, "Popen", unexpected_popen) + + with pytest.raises(effect_runtime.EffectRuntimeStartupError) as exc_info: + effect_runtime._start_runtime(fingerprint="test-fingerprint", info_path=info_path) + + assert exc_info.value.diagnostic_code == "invalid_runtime_idle_ms" + assert "LOOPX_EFFECT_RUNTIME_IDLE_MS" in str(exc_info.value) + assert not info_path.exists() + + +def test_warm_runtime_rejects_invalid_and_conflicting_idle_configuration( + tmp_path: Path, + monkeypatch, +) -> None: + runtime_dir = tmp_path / "runtime" + monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: runtime_dir) + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "1000") + + original = effect_runtime.effect_runtime_result("runtime.ping", {}) + + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "not-a-number") + with pytest.raises(effect_runtime.EffectRuntimeStartupError) as invalid: + effect_runtime.effect_runtime_result("runtime.ping", {}) + assert invalid.value.diagnostic_code == "invalid_runtime_idle_ms" + + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "150") + with pytest.raises(effect_runtime.EffectRuntimeStartupError) as conflict: + effect_runtime.effect_runtime_result("runtime.ping", {}) + assert conflict.value.diagnostic_code == "runtime_idle_ms_conflict" + + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "1000") + assert ( + effect_runtime.effect_runtime_result("runtime.ping", {})["pid"] + == original["pid"] + ) + effect_runtime.effect_runtime_result("runtime.shutdown", {}, retry_safe=False) + + @pytest.mark.parametrize( ("retry_safe", "successful_attempt", "expected_attempts"), [ @@ -699,12 +779,19 @@ def test_request_startup_retry_boundary( "host": "127.0.0.1", "port": 1, "token": "fixture-token", + "idle_ms": effect_runtime.DEFAULT_EFFECT_RUNTIME_IDLE_MS, } - def start_runtime(*, fingerprint: str, info_path: Path): + def start_runtime( + *, + fingerprint: str, + info_path: Path, + requested_idle_ms: int, + ): attempts["count"] += 1 assert fingerprint == ready_info["fingerprint"] assert info_path.parent == runtime_dir + assert requested_idle_ms == effect_runtime.DEFAULT_EFFECT_RUNTIME_IDLE_MS if attempts["count"] == successful_attempt: return ready_info raise effect_runtime.EffectRuntimeStartupError( diff --git a/tests/control_plane_ts/effect_runtime_config.test.ts b/tests/control_plane_ts/effect_runtime_config.test.ts new file mode 100644 index 0000000000..01223936cd --- /dev/null +++ b/tests/control_plane_ts/effect_runtime_config.test.ts @@ -0,0 +1,38 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; + +import { + DEFAULT_EFFECT_RUNTIME_IDLE_MS, + MAX_EFFECT_RUNTIME_IDLE_MS, + effectRuntimeIdleMs, +} from "../../loopx/control_plane/effect_runtime_config.ts"; + +interface IdleFixture { + default_ms: number; + maximum_ms: number; + valid: Array<{ raw: string | null; value: number }>; + invalid: string[]; +} + +const fixture = JSON.parse(readFileSync( + new URL("../fixtures/effect_runtime_idle_ms.json", import.meta.url), + "utf8", +)) as IdleFixture; + +test("Effect runtime idle timeout matches the shared valid fixture", () => { + assert.equal(DEFAULT_EFFECT_RUNTIME_IDLE_MS, fixture.default_ms); + assert.equal(MAX_EFFECT_RUNTIME_IDLE_MS, fixture.maximum_ms); + for (const { raw, value } of fixture.valid) { + assert.equal(effectRuntimeIdleMs(raw ?? undefined), value); + } +}); + +for (const value of fixture.invalid) { + test(`Effect runtime idle timeout rejects '${value}'`, () => { + assert.throws( + () => effectRuntimeIdleMs(value), + /LOOPX_EFFECT_RUNTIME_IDLE_MS/u, + ); + }); +} diff --git a/tests/fixtures/effect_runtime_idle_ms.json b/tests/fixtures/effect_runtime_idle_ms.json new file mode 100644 index 0000000000..1ed0dd30e5 --- /dev/null +++ b/tests/fixtures/effect_runtime_idle_ms.json @@ -0,0 +1,10 @@ +{ + "default_ms": 300000, + "maximum_ms": 2147483647, + "valid": [ + { "raw": null, "value": 300000 }, + { "raw": " 250 ", "value": 250 }, + { "raw": "2147483647", "value": 2147483647 } + ], + "invalid": ["", "not-a-number", "0", "-1", "1.5", "1e3", "2147483648"] +}