From 579213f4f407edf9834c5c5957596b83ca960293 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 05:16:50 +0800 Subject: [PATCH 1/6] fix(registry): report write denial during global sync lock acquisition Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/global_registry.py | 25 ++++++--- ...est_global_registry_write_serialization.py | 53 +++++++++++++++++++ 2 files changed, 71 insertions(+), 7 deletions(-) diff --git a/loopx/global_registry.py b/loopx/global_registry.py index bc4602dbc..071a37e68 100644 --- a/loopx/global_registry.py +++ b/loopx/global_registry.py @@ -1,6 +1,7 @@ from __future__ import annotations import copy +from contextlib import ExitStack from dataclasses import dataclass import errno import json @@ -640,6 +641,7 @@ def _sync_project_registry_to_global_once( allow_route_replacement: bool = False, _global_registry_lock_held: bool = False, _expected_global_registry: Path | None = None, + _lock_write_denied: OSError | None = None, ) -> dict[str, Any]: registry_path = registry_path.expanduser() if not registry_path.exists(): @@ -722,8 +724,8 @@ def _sync_project_registry_to_global_once( writability = None if not dry_run: writability = probe_registry_write_path(global_path, create_parent=True) - if not writability.get("ok"): - exc = PermissionError( + if _lock_write_denied is not None or not writability.get("ok"): + exc = _lock_write_denied or PermissionError( writability.get("errno") if isinstance(writability.get("errno"), int) else errno.EPERM, @@ -870,18 +872,27 @@ def sync_project_registry_to_global( allow_route_replacement=allow_route_replacement, _global_registry_lock_held=_global_registry_lock_held, ) - with exclusive_cross_runtime_file_lock( - target_registry, - operation="sync_global_registry", - ): + with ExitStack() as lock: + lock_write_denied: OSError | None = None + try: + lock.enter_context(exclusive_cross_runtime_file_lock( + target_registry, + operation="sync_global_registry", + )) + except OSError as exc: + if not is_write_denied_error(exc): + raise + # A later writable probe cannot authorize an unlocked retry. + lock_write_denied = exc return _sync_project_registry_to_global_once( registry_path=source_registry, runtime_root_override=runtime_root_override, goal_id=goal_id, dry_run=False, allow_route_replacement=allow_route_replacement, - _global_registry_lock_held=True, + _global_registry_lock_held=lock_write_denied is None, _expected_global_registry=target_registry, + _lock_write_denied=lock_write_denied, ) diff --git a/tests/test_global_registry_write_serialization.py b/tests/test_global_registry_write_serialization.py index df0efe830..732cb8718 100644 --- a/tests/test_global_registry_write_serialization.py +++ b/tests/test_global_registry_write_serialization.py @@ -1,5 +1,6 @@ from __future__ import annotations +import errno import json import subprocess import sys @@ -128,6 +129,58 @@ def recording_write(path: Path, payload: dict[str, Any]) -> None: assert events.index("read:locked") < events.index("write:locked"), events +@pytest.mark.parametrize("denied_errno", [errno.EACCES, errno.EPERM, errno.EROFS]) +def test_sync_lock_write_denial_is_terminal_even_when_probe_is_writable( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, denied_errno: int +) -> None: + runtime_root = tmp_path / "runtime" + registry_path = _project_registry(tmp_path, "alpha", runtime_root) + original = registry_path.read_bytes() + failure = OSError(denied_errno, "lock write denied") + + @contextmanager + def denied_lock(path: Path, **kwargs: Any) -> Iterator[Path]: + raise failure + yield path + + monkeypatch.setattr(global_registry, "exclusive_cross_runtime_file_lock", denied_lock) + result = sync_project_registry_to_global( + registry_path=registry_path, runtime_root_override=str(runtime_root) + ) + + assert result["ok"] is False + assert result["error_kind"] == "global_registry_write_denied" + assert result["errno"] == denied_errno + assert result["wrote"] is False + assert result["synced_goal_ids"] == [] + assert result["requires_global_registry_repair"] is True + assert result["requires_host_permission"] is True + assert result["global_registry_writability"]["ok"] is True + assert registry_path.read_bytes() == original + assert not global_registry_path(runtime_root).exists() + + +def test_sync_lock_other_io_error_is_not_permission_failure( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + runtime_root = tmp_path / "runtime" + registry_path = _project_registry(tmp_path, "alpha", runtime_root) + failure = OSError(errno.ENOSPC, "device full") + + @contextmanager + def failed_lock(path: Path, **kwargs: Any) -> Iterator[Path]: + raise failure + yield path + + monkeypatch.setattr(global_registry, "exclusive_cross_runtime_file_lock", failed_lock) + with pytest.raises(OSError) as caught: + sync_project_registry_to_global( + registry_path=registry_path, runtime_root_override=str(runtime_root) + ) + assert caught.value is failure + assert not global_registry_path(runtime_root).exists() + + def test_dry_run_sync_does_not_take_the_write_lock( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: From ef9839bd2884cf2e3ea8e29649a2d521fb1c79e1 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 05:38:28 +0800 Subject: [PATCH 2/6] fix(install): guard shared Chat preparation before promotion Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../release-promotion-concurrency-smoke.py | 27 +++++++++++++++++++ scripts/install-local.sh | 16 ++++++----- 2 files changed, 36 insertions(+), 7 deletions(-) diff --git a/examples/release/release-promotion-concurrency-smoke.py b/examples/release/release-promotion-concurrency-smoke.py index 598f97f74..96e1e1be5 100644 --- a/examples/release/release-promotion-concurrency-smoke.py +++ b/examples/release/release-promotion-concurrency-smoke.py @@ -71,6 +71,31 @@ def write_lock_test_python(path: Path) -> None: def assert_install_waits_for_promotion_guard(root: Path) -> None: env = {**install_env(root), "LOOPX_RELEASE_ID": "guarded"} + python_wrapper = root / "guard-python" + # Record preparation without skipping the real guard, manifest, or install. + python_wrapper.write_text( + "#!/bin/sh\n" + "if [ \"${1:-}\" = \"-\" ]; then\n" + " source=$(cat)\n" + " case \"$source\" in\n" + " *\"print(sys.executable)\"*) printf '%s\\n' \"$0\"; exit 0 ;;\n" + " esac\n" + " printf '%s' \"$source\" | \"$LOOPX_TEST_REAL_PYTHON\" \"$@\"\n" + " exit $?\n" + "fi\n" + "case \"${1:-}\" in\n" + " */scripts/chat_bundle.py) printf '%s\\n' prepared >\"$LOOPX_TEST_CHAT_PREP_MARKER\" ;;\n" + "esac\n" + "exec \"$LOOPX_TEST_REAL_PYTHON\" \"$@\"\n", + encoding="utf-8", + ) + python_wrapper.chmod(0o755) + chat_marker = root / "chat-prepared" + env.update( + LOOPX_PYTHON=str(python_wrapper), + LOOPX_TEST_REAL_PYTHON=sys.executable, + LOOPX_TEST_CHAT_PREP_MARKER=str(chat_marker), + ) releases_dir = Path(env["LOOPX_RELEASES_DIR"]) releases_dir.mkdir(parents=True) guard_path = releases_dir / ".install-guard" @@ -89,6 +114,7 @@ def assert_install_waits_for_promotion_guard(root: Path) -> None: deadline = time.monotonic() + 2 while time.monotonic() < deadline and process.poll() is None: assert not legacy_lock.exists(), legacy_lock + assert not chat_marker.exists(), "Chat assets prepared before promotion guard" time.sleep(0.05) assert process.poll() is None, process.communicate() except Exception: @@ -105,6 +131,7 @@ def assert_install_waits_for_promotion_guard(root: Path) -> None: stdout, stderr = process.communicate(timeout=120) assert process.returncode == 0, (stdout, stderr) assert release_path(stdout).is_dir(), stdout + assert chat_marker.is_file(), "Guarded install did not prepare Chat assets" def assert_legacy_lock_timeout_preserves_live_owner(root: Path) -> None: diff --git a/scripts/install-local.sh b/scripts/install-local.sh index eaa7a669d..90290b532 100755 --- a/scripts/install-local.sh +++ b/scripts/install-local.sh @@ -679,6 +679,15 @@ if [[ -z "$shell_profile" ]]; then fi configure_python_runtime +promote_default=0 +if resolve_default_promotion; then + promote_default=1 +fi +if [[ "$promote_default" == "1" ]]; then + # Preparing shared Chat assets is part of the guarded installation. + mkdir -p "$releases_dir" + run_under_install_guard "$@" +fi chat_bundle_args=(ensure) if [[ -L "$bin_dir/loopx" ]]; then previous_chat_assets="$("${LOOPX_PYTHON:-python3}" - "$bin_dir/loopx" <<'PYTHON' @@ -693,11 +702,6 @@ PYTHON fi "${LOOPX_PYTHON:-python3}" "$repo_root/scripts/chat_bundle.py" "${chat_bundle_args[@]}" -promote_default=0 -if resolve_default_promotion; then - promote_default=1 -fi - if [[ "$promote_default" == "0" ]]; then if [[ "$install_canary" == "0" ]]; then echo "loopx installer error: default promotion is guarded and LOOPX_INSTALL_CANARY=0 leaves no install target" >&2 @@ -730,8 +734,6 @@ fi export LOOPX_PROMOTION_MODE="$promotion_mode" -mkdir -p "$releases_dir" -run_under_install_guard "$@" warn_stale_promotion_readiness acquire_install_lock mkdir -p "$bin_dir" From 46d6e8d652aab5a03b06d6eadc96e88c0771dca9 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 05:39:43 +0800 Subject: [PATCH 3/6] test(registry): verify repaired write access resumes global sync Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- examples/project/global-registry-writability-smoke.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/examples/project/global-registry-writability-smoke.py b/examples/project/global-registry-writability-smoke.py index d0243d7c7..fc7f870e3 100644 --- a/examples/project/global-registry-writability-smoke.py +++ b/examples/project/global-registry-writability-smoke.py @@ -102,6 +102,14 @@ def assert_sync_reports_write_denied(root: Path) -> None: assert result["wrote"] is False, result assert result["global_registry_writability"]["ok"] is False, result assert result["requires_global_registry_repair"] is True, result + recovered = sync_project_registry_to_global( + registry_path=registry, + runtime_root_override=None, + goal_id=GOAL_ID, + dry_run=False, + ) + assert recovered["ok"] is True and recovered["wrote"] is True, recovered + assert only_goal(runtime / "registry.global.json")["id"] == GOAL_ID def assert_connect_fails_without_partial_local_state(root: Path) -> None: From 76fd74330cdd436bcf200829a79c457f896d17c9 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 09:43:30 +0800 Subject: [PATCH 4/6] fix(usage): pin UTF-8 for Node control replies Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/usage_ping.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py index 402dae5d3..438c90f79 100644 --- a/loopx/usage_ping.py +++ b/loopx/usage_ping.py @@ -51,7 +51,7 @@ def _command() -> list[str]: def control(action: str, path: Path | None = None, **fields: Any) -> dict[str, Any]: result = subprocess.run(_command(), input=json.dumps(_request(action, path or state_path(), **fields)), - capture_output=True, text=True, timeout=4, check=False) + capture_output=True, text=True, encoding="utf-8", timeout=4, check=False) payload = json.loads(result.stdout) if result.returncode or not isinstance(payload, dict) or "error" in payload: raise RuntimeError("Usage settings unavailable. Inspect the local usage-ping.json; disable can repair invalid state.") From b903eee0e2a72bb5f7b57988eacad60c2e18c3aa Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 09:43:30 +0800 Subject: [PATCH 5/6] chore(registry): refresh the current IO census Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../project_registry_io_manifest_v1.json | 34 ++++++++++++++----- 1 file changed, 25 insertions(+), 9 deletions(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 99b16fb2c..b0290626d 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -407,7 +407,7 @@ }, { "site": "loopx/chat_status_api.py::.ChatStatusRequestMixin._status::codec_read:load_registry#1", - "line": 110, + "line": 112, "column": 28, "kind": "codec_read", "api": "load_registry", @@ -415,7 +415,7 @@ }, { "site": "loopx/chat_status_api.py::.ChatStatusRequestMixin._status::codec_read:load_registry#2", - "line": 137, + "line": 139, "column": 21, "kind": "codec_read", "api": "load_registry", @@ -439,7 +439,7 @@ }, { "site": "loopx/cli_commands/authority_archive.py::.authority_upgrade_roots::codec_read:load_registry#1", - "line": 86, + "line": 101, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -447,7 +447,7 @@ }, { "site": "loopx/cli_commands/authority_archive.py::.authority_upgrade_roots::codec_read:load_registry#2", - "line": 94, + "line": 109, "column": 27, "kind": "codec_read", "api": "load_registry", @@ -455,7 +455,7 @@ }, { "site": "loopx/cli_commands/authority_archive.py::.authority_upgrade_roots::codec_read:load_registry#3", - "line": 104, + "line": 119, "column": 48, "kind": "codec_read", "api": "load_registry", @@ -463,12 +463,20 @@ }, { "site": "loopx/cli_commands/authority_archive.py::.handle_authority_archive_command::codec_read:load_registry#1", - "line": 64, + "line": 70, "column": 17, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/cli_commands/authority_archive.py::.handle_authority_archive_command::codec_read:load_registry#2", + "line": 81, + "column": 21, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" + }, { "site": "loopx/cli_commands/automation_cadence.py::.handle_automation_cadence_command::codec_read:load_registry#1", "line": 63, @@ -687,7 +695,7 @@ }, { "site": "loopx/cli_commands/quota.py::._dispatch_quota_turn_start_hooks::codec_read:load_registry#1", - "line": 267, + "line": 270, "column": 37, "kind": "codec_read", "api": "load_registry", @@ -1647,7 +1655,7 @@ }, { "site": "loopx/global_registry.py::._sync_project_registry_to_global_once::codec_read:load_registry#1", - "line": 647, + "line": 649, "column": 24, "kind": "codec_read", "api": "load_registry", @@ -1655,7 +1663,7 @@ }, { "site": "loopx/global_registry.py::.sync_project_registry_to_global::codec_read:load_registry#1", - "line": 855, + "line": 857, "column": 22, "kind": "codec_read", "api": "load_registry", @@ -1964,6 +1972,14 @@ "kind": "codec_read", "api": "load_registry", "classification": "codec_api" + }, + { + "site": "loopx/usage_goal.py::._bound_codex_session::codec_read:load_registry#1", + "line": 125, + "column": 31, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" } ] } From 3c1320b29987f16a9ab3f86d50a63212034e851c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 10:41:11 +0800 Subject: [PATCH 6/6] fix(usage): initiate daily heartbeat before releasing its claim Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/reference/usage-ping.md | 5 ++- .../control_plane/runtime/usage_statistics.ts | 20 +++++++++--- .../control_plane_ts/usage_statistics.test.ts | 32 +++++++++++++++++++ 3 files changed, 51 insertions(+), 6 deletions(-) diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md index e688a202e..fad6d6091 100644 --- a/docs/reference/usage-ping.md +++ b/docs/reference/usage-ping.md @@ -110,7 +110,10 @@ CLI invocation reads only a small local hint; a detached Node process owns measurement, locks and network I/O. A first-use/settings operation may wait for local Node execution, never for a collector connection. -Each installation attempts at most one heartbeat per UTC day. Counts are capped +Each installation attempts at most one heartbeat per UTC day. The detached +sender persists that daily claim and starts the request under one short lock; +network waiting happens after release, so another observer cannot consume the +claim between persistence and request initiation. Counts are capped at 128 distinct rows and 10,000 per row, then flushed on the first eligible command after the UTC day closes. Unsent counts older than seven days are discarded. An installation that never runs again will not flush its final day. diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts index 1c66d9d00..321d3b623 100644 --- a/loopx/control_plane/runtime/usage_statistics.ts +++ b/loopx/control_plane/runtime/usage_statistics.ts @@ -133,6 +133,7 @@ const post: Post = async (url, payload) => (await fetch(url, { export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation, cycle?: CycleObservation) { if (counter !== null && (!validCounter(counter) || counter.count !== 1)) return { sent: false, reason: "invalid_observation" }; let heartbeat: Ping | null = null; + let heartbeatRequest: Promise | undefined; let aggregate: Aggregate | null = null; let goals: GoalAggregate | null = null; const today = day(ctx); @@ -165,6 +166,13 @@ export async function observe(path: string, ctx: Context, generation: string, co state.last_attempt_day = today; // claim before I/O; failures are not retried } else aggregate = null; await save(path, state); + // Persist the daily claim, then initiate its request before releasing this + // lock. A competing observer must not consume a second acquisition between + // the claim and the request. Network waiting stays outside the lock. + if (heartbeat) { + try { heartbeatRequest = send(endpoint(ctx.env), heartbeat).catch(() => 0); } + catch { heartbeatRequest = Promise.resolve(0); } + } return true; }, 0); // Never queue behind business or telemetry work. if (!allowed) return { sent: false, reason: "blocked" }; @@ -173,13 +181,15 @@ export async function observe(path: string, ctx: Context, generation: string, co if (!payload) continue; if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload))) continue; try { - let request: Promise | undefined; + let request = url.endsWith("/ping") ? heartbeatRequest : undefined; // Start under the same short lock as disable, but never hold it while // awaiting network I/O. Once disable returns, no new channel can start. - await withFileMutationLock(path, async () => { - const current = await load(path); - if (!blockedBy(current, ctx) && current.generation === generation) request = send(url, payload).catch(() => 0); - }, 0); + if (!url.endsWith("/ping")) { + await withFileMutationLock(path, async () => { + const current = await load(path); + if (!blockedBy(current, ctx) && current.generation === generation) request = send(url, payload).catch(() => 0); + }, 0); + } if (!request) break; const code = await request; if (url.endsWith("/ping") && code >= 200 && code < 300) { diff --git a/tests/control_plane_ts/usage_statistics.test.ts b/tests/control_plane_ts/usage_statistics.test.ts index e16d9852a..b3a5316c8 100644 --- a/tests/control_plane_ts/usage_statistics.test.ts +++ b/tests/control_plane_ts/usage_statistics.test.ts @@ -1,5 +1,7 @@ import assert from "node:assert/strict"; +import fs from "node:fs/promises"; import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { syncBuiltinESMExports } from "node:module"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { createServer } from "node:http"; @@ -8,6 +10,8 @@ import { configure, inspect, observe, endpoint } from "../../loopx/control_plane import type { Context, Post } from "../../loopx/control_plane/runtime/usage_statistics.ts"; import { validAggregate, validPing, AGGREGATE_SCHEMA, durationBucket } from "../../loopx/control_plane/runtime/usage_statistics_contract.ts"; import type { Counter } from "../../loopx/control_plane/runtime/usage_statistics_contract.ts"; +import { acquireFileMutationLock, releaseFileMutationLock } from "../../loopx/control_plane/effect_runtime_io.ts"; +import type { FileMutationLock } from "../../loopx/control_plane/effect_runtime_io.ts"; const row: Counter = { feature: "todo", outcome: "ok", duration: "lt_1s", error: "none", count: 1 }; const context = (day = "2026-09-26"): Context => ({ env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:8787/v1/ping" }, version: "1.2.0", python: "3.13", channel: "source", now: new Date(day + "T12:00:00Z") }); @@ -79,6 +83,34 @@ test("one heartbeat per UTC day; closed-day aggregation is separate and identifi assert.equal((await state()).counters[0].count, 1); }); +test("a daily claim starts its request before a competing observer can take the released lock", async t => { + const { path, state } = await fixture(t); const ctx = context(); + await configure(path, ctx, "enable"); const generation = (await state()).generation; + let contender: FileMutationLock | undefined; + const rename = fs.rename; + // Control only the scheduling boundary; state, acquisition and retirement + // still use the real filesystem lock. A second worker wins immediately after + // the daily claim is persisted and its owner retires the lock. + const hook = t.mock.method(fs, "rename", async (...args: Parameters) => { + await rename(...args); + if (args[0] === path + ".ts-effect.lock" && !contender) contender = await acquireFileMutationLock(path, process.pid, 0); + }); + syncBuiltinESMExports(); + t.after(async () => { + hook.mock.restore(); syncBuiltinESMExports(); + if (contender) await releaseFileMutationLock(path, contender.token); + }); + let requests = 0; + const result = await observe(path, ctx, generation, row, async () => { requests++; return 204; }); + assert.ok(contender); + assert.equal((await state()).last_attempt_day, "2026-09-26"); + assert.equal(requests, 1, "a durable daily claim must not need another lock acquisition to initiate its request"); + assert.equal(result.sent, true); + hook.mock.restore(); syncBuiltinESMExports(); + await releaseFileMutationLock(path, contender.token); contender = undefined; + await observe(path, ctx, generation, row, noPost); +}); + test("disable clears ID and pending counts; queued observers and in-flight completion cannot resurrect either", async t => { const { path, state } = await fixture(t); const ctx = context(); await configure(path, ctx, "enable"); const old = await state();