Skip to content
Closed
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
92 changes: 85 additions & 7 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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."""

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand All @@ -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):
Expand All @@ -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",
Expand All @@ -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()
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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())
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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:
Expand Down
22 changes: 22 additions & 0 deletions loopx/control_plane/effect_runtime_config.ts
Original file line number Diff line number Diff line change
@@ -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;
}
5 changes: 3 additions & 2 deletions loopx/control_plane/effect_runtime_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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 {
Expand All @@ -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,
Expand Down Expand Up @@ -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);
Expand Down
89 changes: 88 additions & 1 deletion tests/control_plane/test_effect_runtime_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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"),
[
Expand All @@ -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(
Expand Down
38 changes: 38 additions & 0 deletions tests/control_plane_ts/effect_runtime_config.test.ts
Original file line number Diff line number Diff line change
@@ -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,
);
});
}
Loading
Loading