Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -1354,6 +1354,15 @@ If an authority daemon owns a registry/workspace, a CLI process must connect
to it instead of opening a second direct writer. Runtime discovery and startup
are automatic; users do not configure ports or supervise processes.

For the managed loopback runtime, a visible locator is discovery evidence;
request dispatch and successful replies wait until locator publication and its
awaited lock cleanup finish. Otherwise an immediate exit after the first reply
can leave an incomplete cleanup claim and prevent the next retry-safe write
from restarting within the existing lock budget. The real-Node publication
regression covers first ping and typed write, then abrupt exit and receipt
replay. This repair preserves lock reclaim ages, startup deadlines and retry
classification; broader process/storage recovery remains separately qualified.

### 2.3 TypeScript owns migrated effects

The target is not “TypeScript decides, Python always executes”. TypeScript may
Expand Down
11 changes: 10 additions & 1 deletion loopx/control_plane/effect_runtime_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,9 @@ const idleMs = parseIdleMs(process.env.LOOPX_EFFECT_RUNTIME_IDLE_MS);
let idleTimer: NodeJS.Timeout;
let pendingRequests = 0;
let resolveDrain: (() => void) | undefined;
let publicationComplete = false;
let resolvePublication: () => void;
const publicationReady = new Promise<void>((resolve) => { resolvePublication = resolve; });
const handlers = createEffectRuntimeHandlers({
fingerprint,
requestShutdown: () => {
Expand All @@ -119,7 +122,7 @@ const handlers = createEffectRuntimeHandlers({
function resetIdleTimer(server: ReturnType<typeof createServer>): void {
clearTimeout(idleTimer);
// Idleness starts after effects finish, not when their sockets connect.
if (pendingRequests > 0 || shutdownRequested || !server.listening) return;
if (!publicationComplete || pendingRequests > 0 || shutdownRequested || !server.listening) return;
idleTimer = setTimeout(() => server.close(), idleMs);
idleTimer.unref();
}
Expand Down Expand Up @@ -178,6 +181,10 @@ const server = createServer((socket) => {
let sink: FileHandle | null = null;
let dispatched = false;
try {
// The locator becomes visible inside the publication lock. A first
// response must wait for its release: otherwise an immediate process
// death can strand an incomplete cleanup claim before the next write.
await publicationReady;
let parsed: unknown;
try {
parsed = JSON.parse(
Expand Down Expand Up @@ -299,5 +306,7 @@ server.listen(0, "127.0.0.1", async () => {
});
await chmod(infoPath, 0o600);
});
publicationComplete = true;
resolvePublication();
resetIdleTimer(server);
});
120 changes: 120 additions & 0 deletions tests/control_plane/test_effect_runtime_publication_ready.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
"""A visible locator must not admit effects before its publication settles."""
from __future__ import annotations

from concurrent.futures import ThreadPoolExecutor
import json
import os
from pathlib import Path
import signal
import time

import pytest

from loopx.control_plane import effect_runtime


@pytest.mark.parametrize("method", ["runtime.ping", "turn_journal.write"])
def test_first_request_waits_for_real_publication_cleanup_and_recovers(
tmp_path: Path, monkeypatch, request: pytest.FixtureRequest, method: str,
) -> None:
runtime_dir = tmp_path / "runtime"
opened = tmp_path / "claim-opened"
release = tmp_path / "release"
preload = tmp_path / "pause-publication.mjs"
# Pause the real Node file handle after exclusive claim creation, before
# owner bytes or lock cleanup. No mocked transport, dispatch or lock owner.
preload.write_text("""import fs from 'node:fs';
import {syncBuiltinESMExports} from 'node:module';
const original = fs.promises.open;
fs.promises.open = async function(path, ...args) {
const handle = await original.call(this, path, ...args);
if (String(path).includes('.ts-effect.lock.claim.') &&
String(path).includes('runtime-') &&
!fs.existsSync(process.env.LOOPX_TEST_PUBLICATION_OPENED)) {
fs.writeFileSync(process.env.LOOPX_TEST_PUBLICATION_OPENED, 'opened');
while (!fs.existsSync(process.env.LOOPX_TEST_PUBLICATION_RELEASE)) {
await new Promise(resolve => setTimeout(resolve, 5));
}
}
return handle;
};
syncBuiltinESMExports();
""", encoding="utf-8")
monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: runtime_dir)
request.addfinalizer(effect_runtime.restart_effect_runtime)
monkeypatch.setenv("NODE_OPTIONS", f"--import={preload.as_uri()}")
monkeypatch.setenv("LOOPX_TEST_PUBLICATION_OPENED", str(opened))
monkeypatch.setenv("LOOPX_TEST_PUBLICATION_RELEASE", str(release))
monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "60000")
effect_id = "fixture-goal:fixture-agent:todo_fixture0001:publication"
journal_path = tmp_path / "turn.json"
journal = {
"schema_version": "loopx_turn_journal_v0",
"goal_id": "fixture-goal",
"turn_key": "sha256:" + "a" * 64,
"status": "in_progress",
"completed_phases": [],
"plan": {
"turn_envelope": {
"goal_id": "fixture-goal", "agent_id": "fixture-agent",
"action": {"selected_todo": {"todo_id": "todo_fixture0001"}},
},
"transaction": {
"turn_key": "sha256:" + "a" * 64,
"turn_instance_id": "publication",
"settlement_plan": {
"schema_version": "quota_settlement_plan_v1",
"identity": {
"schema_version": "quota_settlement_identity_v0",
"effect_id": effect_id, "goal_id": "fixture-goal",
"agent_id": "fixture-agent", "todo_id": "todo_fixture0001",
"turn_instance_id": "publication",
},
},
},
},
}
params = {"path": str(journal_path), "journal": journal,
"expected_effect_id": effect_id}
info_path = effect_runtime._runtime_info_path(effect_runtime._runtime_fingerprint())
with ThreadPoolExecutor(max_workers=1) as executor:
response = executor.submit(
effect_runtime.effect_runtime_result, method,
params if method == "turn_journal.write" else {}, retry_safe=False,
)
try:
deadline = time.monotonic() + 5
while not opened.exists() and time.monotonic() < deadline:
time.sleep(0.01)
assert opened.exists(), "real publication did not reach its cleanup claim"
assert info_path.exists()
lock_path = Path(str(info_path) + ".ts-effect.lock")
claims = list(runtime_dir.glob("*.ts-effect.lock.claim.*"))
assert lock_path.exists() and len(claims) == 1
assert claims[0].stat().st_size == 0
# The old implementation replies/commits while this claim is blank.
time.sleep(0.2)
assert not response.done(), "first request acknowledged unfinished publication"
assert not journal_path.exists(), "typed write ran before publication settled"
finally:
release.touch()
result = response.result(timeout=5)
assert not lock_path.exists()
assert not list(runtime_dir.glob("*.ts-effect.lock.claim.*"))
if method == "turn_journal.write":
assert result["appended"] is True and result["replayed"] is False
assert json.loads(journal_path.read_text(encoding="utf-8")) == journal
original = effect_runtime.effect_runtime_result("runtime.ping", {})
try:
os.kill(int(original["pid"]), signal.SIGTERM)
time.sleep(0.1)
recovered = effect_runtime.effect_runtime_result(
"turn_journal.write", params, retry_safe=True,
)
assert recovered["replayed"] is (method == "turn_journal.write")
assert recovered["appended"] is (method == "runtime.ping")
assert json.loads(journal_path.read_text(encoding="utf-8")) == journal
replacement = effect_runtime.effect_runtime_result("runtime.ping", {})
assert replacement["pid"] != original["pid"]
finally:
effect_runtime.effect_runtime_result("runtime.shutdown", {}, retry_safe=False)
Loading