Skip to content
Merged
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
5 changes: 4 additions & 1 deletion docs/reference/usage-ping.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
8 changes: 8 additions & 0 deletions examples/project/global-registry-writability-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
27 changes: 27 additions & 0 deletions examples/release/release-promotion-concurrency-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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:
Expand All @@ -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:
Expand Down
20 changes: 15 additions & 5 deletions loopx/control_plane/runtime/usage_statistics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<number> | undefined;
let aggregate: Aggregate | null = null;
let goals: GoalAggregate | null = null;
const today = day(ctx);
Expand Down Expand Up @@ -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" };
Expand All @@ -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<number> | 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) {
Expand Down
25 changes: 18 additions & 7 deletions loopx/global_registry.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from __future__ import annotations

import copy
from contextlib import ExitStack
from dataclasses import dataclass
import errno
import json
Expand Down Expand Up @@ -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():
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
)


Expand Down
34 changes: 25 additions & 9 deletions loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -407,15 +407,15 @@
},
{
"site": "loopx/chat_status_api.py::<module>.ChatStatusRequestMixin._status::codec_read:load_registry#1",
"line": 110,
"line": 112,
"column": 28,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/chat_status_api.py::<module>.ChatStatusRequestMixin._status::codec_read:load_registry#2",
"line": 137,
"line": 139,
"column": 21,
"kind": "codec_read",
"api": "load_registry",
Expand All @@ -439,36 +439,44 @@
},
{
"site": "loopx/cli_commands/authority_archive.py::<module>.authority_upgrade_roots::codec_read:load_registry#1",
"line": 86,
"line": 101,
"column": 16,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/cli_commands/authority_archive.py::<module>.authority_upgrade_roots::codec_read:load_registry#2",
"line": 94,
"line": 109,
"column": 27,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/cli_commands/authority_archive.py::<module>.authority_upgrade_roots::codec_read:load_registry#3",
"line": 104,
"line": 119,
"column": 48,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/cli_commands/authority_archive.py::<module>.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::<module>.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::<module>.handle_automation_cadence_command::codec_read:load_registry#1",
"line": 63,
Expand Down Expand Up @@ -687,7 +695,7 @@
},
{
"site": "loopx/cli_commands/quota.py::<module>._dispatch_quota_turn_start_hooks::codec_read:load_registry#1",
"line": 267,
"line": 270,
"column": 37,
"kind": "codec_read",
"api": "load_registry",
Expand Down Expand Up @@ -1647,15 +1655,15 @@
},
{
"site": "loopx/global_registry.py::<module>._sync_project_registry_to_global_once::codec_read:load_registry#1",
"line": 647,
"line": 649,
"column": 24,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/global_registry.py::<module>.sync_project_registry_to_global::codec_read:load_registry#1",
"line": 855,
"line": 857,
"column": 22,
"kind": "codec_read",
"api": "load_registry",
Expand Down Expand Up @@ -1964,6 +1972,14 @@
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/usage_goal.py::<module>._bound_codex_session::codec_read:load_registry#1",
"line": 125,
"column": 31,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
}
]
}
2 changes: 1 addition & 1 deletion loopx/usage_ping.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.")
Expand Down
16 changes: 9 additions & 7 deletions scripts/install-local.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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 "$@"

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Complete concurrent-install acceptance remains unqualified at this head. The original native smoke fails with the 60-second guard timeout on both current base and head, using the same locked dependency workload and original 120-second child deadline. Because this call moves preparation into the guarded interval, identical errors alone do not establish unaffected guard residence or two-install completion. Please attribute preparation, copying/cleanup and doctor costs under the same supported workload, then show both real installs, distinct release manifests and installed standard/deep doctor succeed. This is an evidence hold, not a proven introduced defect; preserve the original budgets and failed sample.

fi
chat_bundle_args=(ensure)
if [[ -L "$bin_dir/loopx" ]]; then
previous_chat_assets="$("${LOOPX_PYTHON:-python3}" - "$bin_dir/loopx" <<'PYTHON'
Expand All @@ -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
Expand Down Expand Up @@ -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"
Expand Down
32 changes: 32 additions & 0 deletions tests/control_plane_ts/usage_statistics.test.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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") });
Expand Down Expand Up @@ -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<typeof rename>) => {
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();
Expand Down
Loading
Loading