From 6d73a78701defc3ed1553897765eb6e2c00b2aa0 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:10:22 +0800 Subject: [PATCH 1/3] fix(usage): send bounded CLI batches on the same day Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../control_plane/runtime/usage_statistics.ts | 40 ++++-- .../control_plane_ts/usage_statistics.test.ts | 22 +-- .../usage_statistics_delivery.test.ts | 133 ++++++++++++++++++ tests/test_usage_ping.py | 24 ++++ tsconfig.control-plane.json | 1 + 5 files changed, 197 insertions(+), 23 deletions(-) create mode 100644 tests/control_plane_ts/usage_statistics_delivery.test.ts diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts index b5a79ebb2..c4b5789cc 100644 --- a/loopx/control_plane/runtime/usage_statistics.ts +++ b/loopx/control_plane/runtime/usage_statistics.ts @@ -17,13 +17,14 @@ import type { CycleObservation } from "./usage_statistics_cycles.ts"; export const STATE_SCHEMA = "loopx_usage_ping_state_v1"; export const DEFAULT_ENDPOINT = "https://loopx-usage-collector.huangrt01.workers.dev/v1/ping"; export const NOTICE_VERSION = 3; +const AGGREGATE_INTERVAL_MS = 15 * 60 * 1000; export type Env = Record; export type Context = { env: Env; version: string; python: string; channel: string; now?: Date }; type Notice = { version: number; endpoint: string; policy: string }; type State = { schema: typeof STATE_SCHEMA; consent: "default" | "enabled" | "disabled"; generation: string; install_id?: string; notice?: Notice; last_attempt_day?: string; last_sent_day?: string; - day?: string; counters?: Counter[]; + day?: string; counters?: Counter[]; aggregate_last_attempt_ms?: number; }; export function endpoint(env: Env): string { try { @@ -76,6 +77,8 @@ async function load(path: string): Promise { } if (raw.schema !== STATE_SCHEMA || !["default", "enabled", "disabled"].includes(String(raw.consent)) || !validId(raw.generation) || (raw.install_id !== undefined && !validId(raw.install_id)) + || (raw.aggregate_last_attempt_ms !== undefined && (typeof raw.aggregate_last_attempt_ms !== "number" + || !Number.isSafeInteger(raw.aggregate_last_attempt_ms) || raw.aggregate_last_attempt_ms < 0)) || (raw.counters !== undefined && (!Array.isArray(raw.counters) || raw.counters.length > MAX_ROWS || !raw.counters.every(validCounter)))) throw new Error("usage_state_invalid"); return raw as State; } @@ -143,27 +146,26 @@ export async function observe(path: string, ctx: Context, generation: string, co let heartbeat: Ping | null = null; let heartbeatRequest: Promise | undefined; let aggregate: Aggregate | null = null; + let aggregateRequest: Promise | undefined; let goals: GoalAggregate | null = null; const today = day(ctx); + const now = (ctx.now ?? new Date()).getTime(); const allowed = await withFileMutationLock(path, async () => { const state = await load(path); const blocked = blockedBy(state, ctx); if (blocked || !generation || generation !== state.generation) return false; if (state.day && state.day > today) return false; + if (state.aggregate_last_attempt_ms !== undefined && now < state.aggregate_last_attempt_ms) return false; try { - const now = (ctx.now ?? new Date()).getTime(); const intervals = cycle ? await cycleObservations(path + ".cycles", generation, now, cycle) : []; goals = await recordGoalUsage(path + ".goals", generation, now, [...(goal ? [goal] : []), ...intervals]); } catch { /* A damaged optional measurement cannot block other diagnostics. */ } - // Flush only a completed day's local aggregate. No event times or per-install key leave this boundary. - if (state.day && state.day < today && state.counters?.length) { - const age = Date.parse(today) - Date.parse(state.day); - if (age <= 7 * 86400000) aggregate = { schema: AGGREGATE_SCHEMA, counters: state.counters }; - state.counters = []; - } - state.day = today; + // Keep the oldest buffered UTC day for expiry, including across midnight. + // Legacy daily buffers remain readable; no event times or join keys leave. + if (state.day && Date.parse(today) - Date.parse(state.day) > 7 * 86400000) state.counters = []; state.counters ??= []; + if (!state.counters.length) state.day = today; if (counter) { const row = state.counters.find((entry) => counterKey(entry) === counterKey(counter)); if (row) row.count = Math.min(MAX_COUNT, row.count + 1); @@ -172,15 +174,27 @@ export async function observe(path: string, ctx: Context, generation: string, co if (!state.last_attempt_day || state.last_attempt_day < today) { heartbeat = ping(state, ctx); state.last_attempt_day = today; // claim before I/O; failures are not retried - } else aggregate = null; + } + // First result is eligible immediately. Later activity flushes deltas at + // most once per interval, independently of heartbeat success or UTC rollover. + if (state.counters.length && (state.aggregate_last_attempt_ms === undefined + || now - state.aggregate_last_attempt_ms >= AGGREGATE_INTERVAL_MS)) { + aggregate = { schema: AGGREGATE_SCHEMA, counters: state.counters }; + state.counters = []; + state.aggregate_last_attempt_ms = now; // consume before I/O; never retry a lossy batch + } await save(path, state); - // Persist the daily claim, then initiate its request before releasing this + // Persist both claims, then initiate their requests 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); } } + if (aggregate && validAggregate(aggregate)) { + try { aggregateRequest = send(endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), aggregate).catch(() => 0); } + catch { aggregateRequest = Promise.resolve(0); } + } return true; }, 0); // Never queue behind business or telemetry work. if (!allowed) return { sent: false, reason: "blocked" }; @@ -189,10 +203,10 @@ 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 = url.endsWith("/ping") ? heartbeatRequest : undefined; + let request = url.endsWith("/ping") ? heartbeatRequest : url.endsWith("/aggregate") ? aggregateRequest : 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. - if (!url.endsWith("/ping")) { + if (url.endsWith("/goals")) { await withFileMutationLock(path, async () => { const current = await load(path); if (!blockedBy(current, ctx) && current.generation === generation) request = send(url, payload).catch(() => 0); diff --git a/tests/control_plane_ts/usage_statistics.test.ts b/tests/control_plane_ts/usage_statistics.test.ts index ed7784b4f..13121a55d 100644 --- a/tests/control_plane_ts/usage_statistics.test.ts +++ b/tests/control_plane_ts/usage_statistics.test.ts @@ -92,22 +92,24 @@ test("automatic acknowledgment cannot override a choice, changed policy or recip assert.equal(await readFile(path, "utf8"), disabled); }); -test("one heartbeat per UTC day; closed-day aggregation is separate and identifier-free", async t => { +test("one heartbeat per UTC day; same-day aggregates stay separate and identifier-free", async t => { const { path, state } = await fixture(t); const ctx = context(); await configure(path, ctx, "enable"); const generation = (await state()).generation; const sent: { url: string; payload: unknown }[] = []; const post: Post = async (url, payload) => { sent.push({ url, payload }); return 204; }; await observe(path, ctx, generation, row, post); await observe(path, ctx, generation, row, post); - assert.equal(sent.length, 1); assert.ok(validPing(sent[0].payload)); - assert.equal((await inspect(path, ctx)).aggregate_preview?.counters[0].count, 2); + assert.equal(sent.length, 2); assert.ok(validPing(sent[0].payload)); + assert.deepEqual(sent[1].payload, { schema: AGGREGATE_SCHEMA, counters: [row] }); + assert.equal((await inspect(path, ctx)).aggregate_preview?.counters[0].count, 1); await observe(path, context("2026-09-27"), generation, row, post); - assert.equal(sent.length, 3); - assert.ok(sent[2].url.endsWith("/aggregate")); - assert.deepEqual(sent[2].payload, { schema: AGGREGATE_SCHEMA, counters: [{ ...row, count: 2 }] }); - assert.equal((await state()).counters[0].count, 1); + assert.equal(sent.length, 4); + assert.ok(validPing(sent[2].payload)); + assert.ok(sent[3].url.endsWith("/aggregate")); + assert.deepEqual(sent[3].payload, { schema: AGGREGATE_SCHEMA, counters: [{ ...row, count: 2 }] }); + assert.deepEqual((await state()).counters, []); }); -test("a daily claim starts its request before a competing observer can take the released lock", async t => { +test("heartbeat and aggregate claims start 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; @@ -128,7 +130,7 @@ test("a daily claim starts its request before a competing observer can take the 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(requests, 2, "durable heartbeat and aggregate claims must not need another lock acquisition to start"); assert.equal(result.sent, true); hook.mock.restore(); syncBuiltinESMExports(); await releaseFileMutationLock(path, contender.token); contender = undefined; @@ -187,7 +189,7 @@ test("real HTTP sender does not follow redirects to another recipient", async t ctx.env.LOOPX_USAGE_PING_ENDPOINT = `http://127.0.0.1:${(server.address() as { port: number }).port}/v1/ping`; await configure(path, ctx, "enable"); assert.equal((await observe(path, ctx, (await state()).generation, row)).sent, false); - assert.equal(requests, 1); + assert.equal(requests, 2, "heartbeat and aggregate must each stop at the redirect"); }); test("startup heartbeat does not invent a successful command result", async t => { diff --git a/tests/control_plane_ts/usage_statistics_delivery.test.ts b/tests/control_plane_ts/usage_statistics_delivery.test.ts new file mode 100644 index 000000000..873708704 --- /dev/null +++ b/tests/control_plane_ts/usage_statistics_delivery.test.ts @@ -0,0 +1,133 @@ +import assert from "node:assert/strict"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { configure, inspect, observe } from "../../loopx/control_plane/runtime/usage_statistics.ts"; +import type { Context, Post } from "../../loopx/control_plane/runtime/usage_statistics.ts"; +import { AGGREGATE_SCHEMA } from "../../loopx/control_plane/runtime/usage_statistics_contract.ts"; +import type { Aggregate, Counter } from "../../loopx/control_plane/runtime/usage_statistics_contract.ts"; + +const row: Counter = { feature: "todo", outcome: "ok", duration: "lt_1s", error: "none", count: 1 }; +async function fixture(t: test.TestContext) { + const root = await mkdtemp(join(tmpdir(), "loopx-usage-delivery-")); + t.after(() => rm(root, { recursive: true, force: true })); + const path = join(root, "usage-ping.json"); + const ctx: Context = { env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:8787/v1/ping" }, + version: "1.2.2", python: "3.13", channel: "source", now: new Date("2026-09-28T23:59:00Z") }; + await configure(path, ctx, "enable"); + const state = async () => JSON.parse(await readFile(path, "utf8")); + const generation = (await state()).generation; + const batches: Aggregate[] = []; + const post: Post = async (_url, payload) => { + if (payload.schema === AGGREGATE_SCHEMA) batches.push(payload); + return 204; + }; + const at = (minutes: number): Context => ({ ...ctx, now: new Date(ctx.now!.getTime() + minutes * 60000) }); + return { path, ctx, state, generation, batches, post, at }; +} + +test("first completed command sends today; startup and settings invent no counts", async t => { + const { path, ctx, generation, batches, post } = await fixture(t); + await observe(path, ctx, generation, null, post); + assert.deepEqual(batches, []); + assert.equal((await inspect(path, ctx)).aggregate_preview, null); + const failure: Counter = { ...row, outcome: "failed", error: "command_failed" }; + await observe(path, ctx, generation, failure, post); + assert.deepEqual(batches, [{ schema: AGGREGATE_SCHEMA, counters: [failure] }]); + assert.equal((await inspect(path, ctx)).aggregate_preview, null); +}); + +test("midnight and high call volume cannot bypass 15-minute spacing; each batch is a delta", async t => { + const { path, ctx, generation, batches, post, at } = await fixture(t); + await observe(path, ctx, generation, row, post); + for (let i = 0; i < 20; i++) await observe(path, at(1), generation, row, post); + await observe(path, at(14.999), generation, row, post); + assert.equal(batches.length, 1); + assert.equal((await inspect(path, at(14.999))).aggregate_preview?.counters[0].count, 21); + await observe(path, at(15), generation, row, post); + assert.deepEqual(batches.map(b => b.counters[0].count), [1, 22]); + await observe(path, at(16), generation, row, post); + await observe(path, at(30), generation, null, post); + assert.deepEqual(batches.map(b => b.counters[0].count), [1, 22, 1]); + await observe(path, at(45), generation, null, post); + assert.equal(batches.length, 3, "no empty periodic requests"); +}); + +test("upgrade flushes retained legacy counts; expired counts are dropped before a new result", async t => { + for (const [bufferDay, count] of [["2026-09-27", 8], ["2026-09-20", 1]] as const) { + const { path, ctx, state, generation, batches, post } = await fixture(t); + await writeFile(path, JSON.stringify({ ...await state(), day: bufferDay, counters: [{ ...row, count: 7 }], last_attempt_day: "2026-09-28" })); + await observe(path, ctx, generation, row, post); + assert.deepEqual(batches, [{ schema: AGGREGATE_SCHEMA, counters: [{ ...row, count }] }]); + } +}); + +test("failed sends consume the batch without retrying or uploading errors; later batches contain new counts", async t => { + const { path, ctx, state, generation, batches, post, at } = await fixture(t); + let attempts = 0; + await observe(path, ctx, generation, row, async (_url, payload) => { + if (payload.schema === AGGREGATE_SCHEMA) { attempts++; throw new Error("PRIVATE-ERROR"); } + return 204; + }); + assert.equal(attempts, 1); + assert.deepEqual((await state()).counters, []); + await observe(path, at(1), generation, row, post); + assert.deepEqual(batches, []); + await observe(path, at(15), generation, null, post); + assert.deepEqual(batches, [{ schema: AGGREGATE_SCHEMA, counters: [row] }]); + assert.ok(!(await readFile(path, "utf8")).includes("PRIVATE-ERROR")); +}); + +test("clock rollback neither reopens a send nor mutates buffered counts", async t => { + const { path, ctx, generation, batches, post, at } = await fixture(t); + await observe(path, ctx, generation, row, post); + await observe(path, at(15), generation, row, post); + const before = await readFile(path, "utf8"); + await observe(path, at(14), generation, row, post); + assert.equal(await readFile(path, "utf8"), before); + assert.equal(batches.length, 2); +}); + +test("concurrent observers cannot claim duplicate aggregates while HTTP remains pending", async t => { + const { path, ctx, generation, batches, post, at } = await fixture(t); + let release!: () => void; + let started!: () => void; + const held = new Promise(r => { release = r; }); + const ready = new Promise(r => { started = r; }); + const sending = observe(path, ctx, generation, row, async (url, payload) => { + await post(url, payload); + if (payload.schema === AGGREGATE_SCHEMA) { started(); await held; } + return 204; + }); + await ready; + try { + const results = await Promise.allSettled(Array.from({ length: 8 }, () => observe(path, at(1), generation, row, post))); + for (const result of results) { + if (result.status === "rejected") assert.equal(result.reason.code, "mutation_lock_timeout"); + } + assert.equal(batches.length, 1); + await configure(path, at(1), "disable"); + } finally { release(); await sending; } + assert.equal((await inspect(path, at(1))).aggregate_preview, null); + await observe(path, at(15), generation, row, post); + assert.equal(batches.length, 1); +}); + +test("invalid persisted delivery timestamp fails closed and disable repairs it", async t => { + const { path, ctx, state } = await fixture(t); + const initial = await state(); + for (const value of [-1, "yesterday", 1.5, null]) { + await writeFile(path, JSON.stringify({ ...initial, aggregate_last_attempt_ms: value })); + await assert.rejects(inspect(path, ctx), /usage_state_invalid/); + } + await configure(path, ctx, "disable"); + assert.equal((await inspect(path, ctx)).blocked_by, "disabled"); +}); + +test("a malformed legacy buffer cannot bypass outgoing aggregate validation", async t => { + const { path, ctx, state, generation, batches, post } = await fixture(t); + await writeFile(path, JSON.stringify({ ...await state(), day: "2026-09-27", counters: [row, row] })); + await observe(path, ctx, generation, null, post); + assert.deepEqual(batches, [], "duplicate counter keys are not a valid wire payload"); +}); diff --git a/tests/test_usage_ping.py b/tests/test_usage_ping.py index 2b92b5275..eea9f78f3 100644 --- a/tests/test_usage_ping.py +++ b/tests/test_usage_ping.py @@ -204,6 +204,30 @@ def test_real_cli_returns_while_http_response_is_held_and_disable_survives(isola assert str(isolated) not in json.dumps(received) +def test_real_cli_first_result_reaches_http_without_next_day_return(isolated, collector, monkeypatch): + endpoint, received, accepted, release = collector + monkeypatch.setenv('LOOPX_USAGE_PING_ENDPOINT', endpoint) + release.set() + setup = 'import sys; from pathlib import Path; from loopx import usage_ping; usage_ping.DEFAULT_RUNTIME_ROOT=Path(sys.argv[1]); from loopx.cli_runtime import main; ' + command = [sys.executable, '-c', setup + 'raise SystemExit(main(["version", "--format", "json"]))', str(isolated)] + first = subprocess.run(command, capture_output=True, text=True, timeout=30) + assert first.returncode == 0 and 'random installation ID' in first.stderr + assert received == [] + second = subprocess.run(command, capture_output=True, text=True, timeout=30) + assert second.returncode == 0 and json.loads(second.stdout) == json.loads(first.stdout) + deadline = time.monotonic() + 5 + while time.monotonic() < deadline and not any(p['schema'] == 'loopx_usage_aggregate_v1' for p in received): + time.sleep(0.02) + aggregates = [p for p in received if p['schema'] == 'loopx_usage_aggregate_v1'] + assert len(aggregates) == 1, 'one completed command must not depend on a next-day invocation' + assert set(aggregates[0]) == {'schema', 'counters'} + counters = aggregates[0]['counters'] + assert len(counters) == 1 and counters[0]['feature'] == 'version' + assert counters[0]['outcome'] == 'ok' and counters[0]['count'] == 1 + assert usage_ping.control('status')['aggregate_preview'] is None + usage_ping.control('disable') + + def test_business_failure_and_usage_failure_do_not_replace_original_result(isolated, monkeypatch): usage_ping.control('enable') import loopx.cli_runtime as cli diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 4ec2310ab..ed7151069 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -17,6 +17,7 @@ "tests/control_plane_ts/explore_research.test.ts", "loopx/control_plane/runtime/usage_statistics*.ts", "tests/control_plane_ts/usage_statistics.test.ts", + "tests/control_plane_ts/usage_statistics_delivery.test.ts", "tests/control_plane_ts/usage_statistics_goals.test.ts", "tests/control_plane_ts/usage_statistics_sources.test.ts", "loopx/control_plane/effect_program.ts", From 1e58f13170e35aa1cf5d5fe2ba9cf72d99b8b07f Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:10:33 +0800 Subject: [PATCH 2/3] docs(usage): explain same-day delivery and remaining sampling limits Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../usage-statistics-settings.tsx | 4 +- apps/usage-collector/README.md | 5 ++- docs/reference/usage-ping.md | 37 +++++++++++++------ docs/reference/usage-ping.zh-CN.md | 23 +++++++++--- 4 files changed, 49 insertions(+), 20 deletions(-) diff --git a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx index 9e33b9264..ed170ecae 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-settings.tsx @@ -23,8 +23,8 @@ export function UsageStatisticsSettings() { return
{zh ? "基础使用统计 · 告知后默认开启,可关闭" : "Basic usage statistics · on after notice, optional"}

{zh - ? "用于决定平台支持和改进命令体验。每天向 LoopX 的 Cloudflare 收集服务发送随机安装标识、版本、系统、CPU 架构、Python 版本和安装渠道;固定的 CLI 功能、结果、耗时区间和错误类别在本机按天汇总后另行发送,不带安装标识。" - : "Helps prioritize platform support and CLI improvements. A daily heartbeat sends a random installation ID, version, OS, CPU architecture, Python version and install channel to the LoopX Cloudflare collector. Fixed CLI feature, result, duration and error counts are aggregated locally by day and sent separately without the ID."}

+ ? "用于决定平台支持和改进命令体验。每天向 LoopX 的 Cloudflare 收集服务发送随机安装标识、版本、系统、CPU 架构、Python 版本和安装渠道;固定的 CLI 功能、结果、耗时区间和错误类别在本机汇总,不带安装标识。首个可采集的命令结果立即尝试发送,之后有活动时每隔至少 15 分钟发送一批。" + : "Helps prioritize platform support and CLI improvements. A daily heartbeat sends a random installation ID, version, OS, CPU architecture, Python version and install channel to the LoopX Cloudflare collector. Fixed CLI feature, result, duration and error counts are aggregated locally without the ID. The first measured result attempts a send immediately; later activity sends batches at least 15 minutes apart."}

{zh ? "不采集提示词、代码、路径、命令参数、Goal 内容或原始错误。当前功能计数仅覆盖 CLI;命令成功不等于 Goal 完成。" : "No prompts, code, paths, arguments, Goal contents or raw errors. Feature counts currently cover CLI only; command success is not Goal completion."}

{zh ? "按天分别汇总所有 Host 的 quota→spend 推进周期、已绑定 Codex 任务的本地轮次时间、受管 Turn 与普通 Goal 对话的 Host 调用时间。上传固定 Host 类别及跨度/时长区间,不上传会话内容、Goal 或安装标识。三种口径重叠,不能相加;可能漏计,不代表完成、CPU 用时或计费。" : "Daily, separate span/duration buckets for all Hosts using quota→spend, local timing events from bound Codex tasks, and direct Host calls in managed Turns and regular owner Goal chat. Sends fixed Host categories, never session contents, Goal or installation IDs. The three overlapping populations cannot be added; partial observations are not completion, CPU time or billing."}

{state ? <> diff --git a/apps/usage-collector/README.md b/apps/usage-collector/README.md index dda2a0189..80157ed62 100644 --- a/apps/usage-collector/README.md +++ b/apps/usage-collector/README.md @@ -19,7 +19,10 @@ Aggregate requests merge directly into `usage_counts(day, feature, outcome, duration, error, count)` and are retained 30 days. No raw request rows, ID, version or per-request timestamps enter that table. Aggregate writes are lossy, not idempotent: clients make no retry. The server uses its UTC reception date. -Counters are estimates, not people, accepted Goal outcomes or billing records. +Clients may send multiple non-overlapping CLI batches within a day; the collector +adds each delta without requiring a schema migration. Delivery cadence does not +add a version or installation join key. Counters are estimates, not people, +accepted Goal outcomes or billing records. Neither handler reads/stores IP, user agent or Cloudflare request metadata. The template disables Worker observability; Cloudflare still handles network diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md index 776e3c272..1d460afac 100644 --- a/docs/reference/usage-ping.md +++ b/docs/reference/usage-ping.md @@ -48,7 +48,7 @@ OS is `darwin|linux|windows|other`, CPU is `x64|arm64|x86|other`, and channel is `pip|local_release|source|unknown`. Version accepts only numeric major.minor.patch; a custom version containing a private suffix is not sent. -**Closed-day CLI aggregate** (`POST /v1/aggregate`): +**CLI aggregate batch** (`POST /v1/aggregate`): ```json {"schema":"loopx_usage_aggregate_v1","counters":[{"feature":"todo","outcome":"ok","duration":"lt_1s","error":"none","count":4}]} @@ -128,16 +128,31 @@ 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. 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. -Lock contention, crashes and failed requests can lose counts. There are no -immediate retries and no durable network queue. Clock rollback does not reopen -a daily attempt. These are **lossy diagnostics**, not billing or audit records. +Each installation attempts at most one heartbeat per UTC day. The first measured +CLI result attempts an aggregate send immediately, including a failed result. +Later eligible activity sends buffered deltas at most once every 15 minutes; +UTC midnight does not reset that interval. Heartbeats and CLI batches have +independent claims, so a heartbeat attempt cannot suppress a completed result. +Startup alone never invents a result. Settings/status operations never flush. +This replaces next-day-only CLI delivery for enabled installations across +interactive and unattended CLI lanes; Goal-duration snapshots remain daily. + +Counts are capped at 128 distinct rows and 10,000 per row. Buffered counts expire +after seven UTC days measured from the oldest buffered day. Existing daily +buffers remain readable on upgrade. Each batch is removed and its attempt time +persisted before the request starts under the same short lock; network waiting +happens after release. Attempts are spaced even on failure or clock rollback. +There are no immediate retries, background timers or durable network queues. +A session that stops within the interval can still lose its unsent tail; this +reduces dependence on next-day return without promising complete coverage. +Lock contention, crashes and failed requests can also lose counts. These are +**lossy diagnostics**, not billing or audit records. + +CLI batches still contain no installation ID, version or event date. The +collector groups them by UTC reception date, which can differ from the activity +date. Do not divide their totals by reporting installations to infer per-install +usage, or attribute them to a release version. More frequent requests can make +network timing correlation easier; identity-free payloads do not prevent that. Requests have a three-second deadline, do not block command completion, and cannot change its output or exit code. They use the supported Node runtime's diff --git a/docs/reference/usage-ping.zh-CN.md b/docs/reference/usage-ping.zh-CN.md index 23d06d92d..0dad74e35 100644 --- a/docs/reference/usage-ping.zh-CN.md +++ b/docs/reference/usage-ping.zh-CN.md @@ -43,7 +43,7 @@ ID 随机生成,不绑定账号、不从硬件派生,但能跨天关联, 安装渠道只允许 `pip|local_release|source|unknown`。版本只接受数字三段式,包含 自定义后缀的版本不会上传。 -独立的 CLI 日汇总 `POST /v1/aggregate`: +独立的 CLI 汇总批次 `POST /v1/aggregate`: ```json {"schema":"loopx_usage_aggregate_v1","counters":[{"feature":"todo","outcome":"ok","duration":"lt_1s","error":"none","count":4}]} @@ -103,11 +103,22 @@ App 不再要求首次点击启用。环境变量覆盖和 `consent_required` authority provider 备份、公共投影。普通命令只读取很小的本地提示;独立 Node 后台 进程负责计数、锁和网络。首次告知和设置操作可能等待本机 Node,不等待收集服务。 -每天最多尝试一次心跳;本地汇总最多 128 种计数组合,每项封顶 10,000。 -UTC 日期结束后的下一次合格调用发送上一日汇总,超过七天的积压丢弃;最后一天 -之后不再运行的安装不会发送最后一日计数。锁竞争、进程退出和网络故障可能丢数, -不会立即重试,也没有持久网络队列。时钟回拨不重新开放当日尝试。 -这是有损诊断,不能当账单或审计日志。 +每天最多尝试一次心跳。首次可采集的 CLI 命令结束后立即尝试发送汇总,包括失败结果; +后续有合格活动时,每隔至少 15 分钟发送一批新增计数,UTC 换日不重置间隔。 +心跳和 CLI 汇总分别认领发送机会,心跳尝试不会压掉命令结果。单独启动不产生成功 +计数,设置和 status 查询不触发发送。此行为替代已开启统计的交互式、无人值守 CLI +原有的次日发送机制;Goal 时长快照仍按天发送。 + +本地汇总最多 128 种计数组合,每项封顶 10,000;以最早积压日期计算,超过七个 UTC +日的计数丢弃。升级可继续读取旧的按日缓冲。每批在发送前持久化认领时间、移除对应 +计数,并在同一短锁内发起请求,网络等待不占锁。失败或时钟回拨也不能绕过间隔。 +不会立即重试,不增加后台定时器或持久网络队列。间隔内停止使用仍可能丢失未发送的 +尾部计数;优化减少对次日回访的依赖,但不承诺完整采集。锁竞争、进程退出和网络 +故障也可能丢数。这是有损诊断,不能当账单或审计日志。 + +CLI 汇总仍不携带安装 ID、版本或活动日期;服务端按 UTC 接收日期汇总,可能与实际 +使用日期不同。不能拿调用量除以上报安装数推算每安装用量,也不能归因到某个版本。 +更频繁的请求可能增加网络时序关联的机会;载荷不带标识不代表无法关联。 每个网络请求限时 3 秒,不阻塞命令完成、不改变输出和退出码。使用受支持 Node 运行时的 `HTTP_PROXY`、`HTTPS_PROXY`、`NO_PROXY` 配置,代理地址与凭证不会进入遥测数据。 From f9735814e09020113eb73eb13354dfaaa432e84e Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:05:43 +0800 Subject: [PATCH 3/3] fix(usage): renew visible notice before faster delivery Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../usage-statistics-notice.tsx | 3 + docs/reference/usage-ping.md | 12 +++- docs/reference/usage-ping.zh-CN.md | 9 ++- .../control_plane/runtime/usage_statistics.ts | 4 +- loopx/usage_ping.py | 4 +- .../test_source_cli_entrypoint.py | 2 +- .../usage_statistics_delivery.test.ts | 59 ++++++++++++++++++- .../usage_statistics_sources.test.ts | 2 +- tests/test_usage_ping.py | 56 +++++++++++++++++- 9 files changed, 136 insertions(+), 15 deletions(-) diff --git a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-notice.tsx b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-notice.tsx index 02c65676a..c47e0093e 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-notice.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/usage-statistics-notice.tsx @@ -75,6 +75,9 @@ export function UsageStatisticsNotice({ onDetails }: { onDetails: () => void })

{zh ? "用于改进平台支持与使用体验。发送随机安装标识和环境信息,以及另行汇总的 CLI 使用次数、结果、耗时和 Goal 时长区间;不采集对话、代码、路径或命令参数。可随时关闭。" : "Helps improve platform support and usage. Sends a random installation ID and environment information, plus separate CLI usage, result, timing and Goal duration summaries. No conversations, code, paths or command arguments. You can turn it off at any time."}

+

{zh + ? "首个已测量的 CLI 结果立即上报,后续由使用活动触发,至少间隔 15 分钟发送一批。CLI 汇总不含安装标识;更频繁的请求仍可能让网络服务通过 IP 和请求时间关联活动。" + : "The first measured CLI result is sent immediately; later activity sends buffered counts at most once every 15 minutes. CLI summaries contain no installation ID; more frequent requests may still let network services correlate activity using IP addresses and request timing."}

{zh ? "接收方:" : "Recipient: "}{state.endpoint}

{error ?

{zh ? "设置未能保存,请打开详情重试。" : "Could not save this setting. Open details to retry."}

: null} diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md index 1d460afac..4c3123838 100644 --- a/docs/reference/usage-ping.md +++ b/docs/reference/usage-ping.md @@ -75,8 +75,9 @@ Additional payload fields and invalid enum combinations are rejected. Interactive CLI, unattended scripts/agents and the App use the same **first disclosure → automatic activation → subsequent measurement** policy. -The first ordinary CLI command prints the recipient, fields, purpose and -both disable mechanisms to stderr, records the disclosure, and sends nothing. +The first ordinary CLI command prints the recipient, fields, purpose, CLI delivery +cadence, network timing correlation boundary and both disable mechanisms to +stderr, records the disclosure, and sends nothing. This also applies to captured stderr in scripts and Agent tool calls; JSON stdout is unaffected. Discarded stderr (the null device) or a failed write cannot acknowledge a notice. Background `chat`/`serve-status` services defer @@ -139,7 +140,12 @@ interactive and unattended CLI lanes; Goal-duration snapshots remain daily. Counts are capped at 128 distinct rows and 10,000 per row. Buffered counts expire after seven UTC days measured from the oldest buffered day. Existing daily -buffers remain readable on upgrade. Each batch is removed and its attempt time +buffer shapes remain readable. Notice revision 4 renews disclosure before the +faster cadence takes effect: old notice state cannot send or consume buffers. +Acknowledging the renewed notice discards old-scope counters and fences queued +observations with a new generation; only subsequent measurements can send. +Explicit disable remains disabled, and acknowledgment cannot replace explicit +enable under `consent_required`. Each batch is removed and its attempt time persisted before the request starts under the same short lock; network waiting happens after release. Attempts are spaced even on failure or clock rollback. There are no immediate retries, background timers or durable network queues. diff --git a/docs/reference/usage-ping.zh-CN.md b/docs/reference/usage-ping.zh-CN.md index 0dad74e35..a251ce199 100644 --- a/docs/reference/usage-ping.zh-CN.md +++ b/docs/reference/usage-ping.zh-CN.md @@ -66,7 +66,8 @@ ID 随机生成,不绑定账号、不从硬件派生,但能跨天关联, ## 告知、设置与升级 交互式 CLI、后台脚本/Agent 和 App 统一采用 **首次告知 → 自动开启 → 后续采集**。 -首次普通 CLI 命令向 stderr 显示接收方、字段、用途和关闭方法,然后记录告知; +首次普通 CLI 命令向 stderr 显示接收方、字段、用途、CLI 发送频率、网络时序关联 +边界和关闭方法,然后记录告知; 这一轮不计数、不发送。脚本和 Agent 工具调用捕获的 stderr 也适用,JSON stdout 保持不变。stderr 指向空设备或输出失败时不记录告知。后台 `chat`/`serve-status` 服务把首次告知交给 App,不以服务日志代替 App 界面。 @@ -110,8 +111,10 @@ authority provider 备份、公共投影。普通命令只读取很小的本地 原有的次日发送机制;Goal 时长快照仍按天发送。 本地汇总最多 128 种计数组合,每项封顶 10,000;以最早积压日期计算,超过七个 UTC -日的计数丢弃。升级可继续读取旧的按日缓冲。每批在发送前持久化认领时间、移除对应 -计数,并在同一短锁内发起请求,网络等待不占锁。失败或时钟回拨也不能绕过间隔。 +日的计数丢弃。旧按日缓冲格式仍可读取。告知版本 4 要求在新频率生效前重新告知: +确认前不发送、不消费缓存;确认后丢弃旧范围计数并更新 generation,阻止旧排队 +观测补发,后续新测量才可发送。明确关闭继续生效,`consent_required` 仍须明确启用。 +每批在发送前持久化认领时间、移除对应计数,并在同一短锁内发起请求,网络等待不占锁。失败或时钟回拨也不能绕过间隔。 不会立即重试,不增加后台定时器或持久网络队列。间隔内停止使用仍可能丢失未发送的 尾部计数;优化减少对次日回访的依赖,但不承诺完整采集。锁竞争、进程退出和网络 故障也可能丢数。这是有损诊断,不能当账单或审计日志。 diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts index c4b5789cc..f2c3ba1fb 100644 --- a/loopx/control_plane/runtime/usage_statistics.ts +++ b/loopx/control_plane/runtime/usage_statistics.ts @@ -16,7 +16,7 @@ import type { CycleObservation } from "./usage_statistics_cycles.ts"; export const STATE_SCHEMA = "loopx_usage_ping_state_v1"; export const DEFAULT_ENDPOINT = "https://loopx-usage-collector.huangrt01.workers.dev/v1/ping"; -export const NOTICE_VERSION = 3; +export const NOTICE_VERSION = 4; const AGGREGATE_INTERVAL_MS = 15 * 60 * 1000; export type Env = Record; export type Context = { env: Env; version: string; python: string; channel: string; now?: Date }; @@ -105,7 +105,7 @@ export async function inspect(path: string, ctx: Context) { aggregate_preview: state.consent === "disabled" || !state.counters?.length ? null : { schema: AGGREGATE_SCHEMA, counters: state.counters }, goal_preview: state.consent === "disabled" ? null : await goalPreview(path + ".goals", state.generation).catch(() => null), aggregate_day: state.day ?? null, - disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. Goal span/duration buckets and fixed Host labels are aggregated without Goal or installation IDs. Common quota-to-spend cycles cover every Host using the quota CLI; bound Codex tasks add local timing-event reads; managed Turns and regular owner Goal chat add direct Host-call timing. These overlapping measurements are separate, partial and not completion or billing evidence. Raw session content is never uploaded. No prompts, code, paths, arguments, Goal contents or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") }; + disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. The first measured CLI result is sent immediately; later activity sends buffered counts at most once every 15 minutes. More frequent requests can make network timing correlation easier: network services may observe IP addresses and request times even though CLI summaries have no installation ID. Goal span/duration buckets and fixed Host labels are aggregated without Goal or installation IDs. Common quota-to-spend cycles cover every Host using the quota CLI; bound Codex tasks add local timing-event reads; managed Turns and regular owner Goal chat add direct Host-call timing. These overlapping measurements are separate, partial and not completion or billing evidence. Raw session content is never uploaded. No prompts, code, paths, arguments, Goal contents or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") }; } export async function configure(path: string, ctx: Context, action: "enable" | "disable" | "acknowledge", expectedNotice?: unknown) { await withFileMutationLock(path, async () => { diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py index 40bb1d20e..30bd782e5 100644 --- a/loopx/usage_ping.py +++ b/loopx/usage_ping.py @@ -15,6 +15,8 @@ from .paths import DEFAULT_RUNTIME_ROOT STATE_FILENAME = "usage-ping.json" +# Scheduling hint only; keep aligned with the TypeScript notice revision. +_NOTICE_VERSION = 4 _ENTRY = Path(__file__).parent / "control_plane/runtime/usage_statistics_cli.ts" @@ -71,7 +73,7 @@ def begin(command: str) -> tuple[str, float] | None: state = json.loads(path.read_text(encoding="utf-8")) if path.exists() else {} if state.get("consent") == "disabled": return None - if (state.get("notice") or {}).get("version") != 3: + if (state.get("notice") or {}).get("version") != _NOTICE_VERSION: # App services defer disclosure to the visible frontend. Ordinary # script/Agent calls disclose on stderr too; JSON stdout stays clean. stream = sys.stderr diff --git a/tests/control_plane/test_source_cli_entrypoint.py b/tests/control_plane/test_source_cli_entrypoint.py index 28535cd23..fdb5fcd15 100644 --- a/tests/control_plane/test_source_cli_entrypoint.py +++ b/tests/control_plane/test_source_cli_entrypoint.py @@ -311,7 +311,7 @@ def test_source_first_usage_disclosure_keeps_json_pure_and_does_not_send(tmp_pat assert json.loads(result.stdout)["ok"] is True assert "random installation ID" in result.stderr stored = json.loads((state / "usage-ping.json").read_text()) - assert stored["notice"]["version"] == 3 + assert stored["notice"]["version"] == 4 assert "last_attempt_day" not in stored and "counters" not in stored diff --git a/tests/control_plane_ts/usage_statistics_delivery.test.ts b/tests/control_plane_ts/usage_statistics_delivery.test.ts index 873708704..a21c0f6dc 100644 --- a/tests/control_plane_ts/usage_statistics_delivery.test.ts +++ b/tests/control_plane_ts/usage_statistics_delivery.test.ts @@ -54,7 +54,7 @@ test("midnight and high call volume cannot bypass 15-minute spacing; each batch assert.equal(batches.length, 3, "no empty periodic requests"); }); -test("upgrade flushes retained legacy counts; expired counts are dropped before a new result", async t => { +test("current-notice legacy buffer shape remains readable; expired counts are dropped before a new result", async t => { for (const [bufferDay, count] of [["2026-09-27", 8], ["2026-09-20", 1]] as const) { const { path, ctx, state, generation, batches, post } = await fixture(t); await writeFile(path, JSON.stringify({ ...await state(), day: bufferDay, counters: [{ ...row, count: 7 }], last_attempt_day: "2026-09-28" })); @@ -131,3 +131,60 @@ test("a malformed legacy buffer cannot bypass outgoing aggregate validation", as await observe(path, ctx, generation, null, post); assert.deepEqual(batches, [], "duplicate counter keys are not a valid wire payload"); }); + +test("v3 upgrade waits for renewed notice, preserves pre-ack state and fences old observations", async t => { + for (const consent of ["default", "enabled"] as const) { + const { path, ctx, state, generation, batches, post } = await fixture(t); + await writeFile(path, JSON.stringify({ ...await state(), consent, + notice: { version: 3, endpoint: ctx.env.LOOPX_USAGE_PING_ENDPOINT, policy: "opt_out" }, + day: "2026-09-28", counters: [{ ...row, count: 7 }] })); + const before = await readFile(path, "utf8"); + const status = await inspect(path, ctx); + assert.equal(status.notice.version, 4); + assert.equal(status.blocked_by, "notice_required"); + assert.equal(status.automatic_notice_required, true); + let attempts = 0; + for (const counter of [null, row]) { + await observe(path, ctx, generation, counter, async () => { attempts++; return 204; }); + } + assert.equal(attempts, 0, "neither heartbeat nor aggregate can start before renewed notice"); + assert.equal(await readFile(path, "utf8"), before, "unacknowledged buffers remain untouched"); + await assert.rejects(configure(path, ctx, "acknowledge", { ...status.notice, version: 3 }), /usage_notice_changed/); + await configure(path, ctx, "acknowledge", status.notice); + const renewed = await state(); + assert.notEqual(renewed.generation, generation); + assert.deepEqual(renewed.counters, [], "existing notice migration discards old-scope counts"); + await observe(path, ctx, generation, row, post); + assert.deepEqual(batches, [], "queued old-generation work cannot send after acknowledgment"); + await observe(path, ctx, renewed.generation, row, post); + assert.deepEqual(batches, [{ schema: AGGREGATE_SCHEMA, counters: [row] }]); + } +}); + +test("v3 renewal cannot undo disable or replace explicit consent", async t => { + for (const [consent, policy, reason] of [ + ["disabled", "opt_out", "disabled"], + ["default", "consent_required", "consent_required"], + ] as const) { + const { path, ctx, state, generation } = await fixture(t); + const restricted = { ...ctx, env: { ...ctx.env, LOOPX_USAGE_POLICY: policy } }; + await writeFile(path, JSON.stringify({ ...await state(), consent, + notice: { version: 3, endpoint: ctx.env.LOOPX_USAGE_PING_ENDPOINT, policy }, + day: "2026-09-28", counters: [{ ...row, count: 7 }] })); + const before = await readFile(path, "utf8"); + const status = await inspect(path, restricted); + assert.equal(status.blocked_by, reason); + assert.equal(status.automatic_notice_required, false); + await configure(path, restricted, "acknowledge", status.notice); + let attempts = 0; + await observe(path, restricted, generation, row, async () => { attempts++; return 204; }); + assert.equal(attempts, 0); + assert.equal(await readFile(path, "utf8"), before); + assert.equal((await inspect(path, restricted)).sending, false); + if (policy === "consent_required") { + await configure(path, restricted, "enable"); + await observe(path, restricted, (await state()).generation, row, async () => { attempts++; return 204; }); + assert.ok(attempts > 0, "only explicit enable authorizes consent-required collection"); + } + } +}); diff --git a/tests/control_plane_ts/usage_statistics_sources.test.ts b/tests/control_plane_ts/usage_statistics_sources.test.ts index db847705b..3c12ed011 100644 --- a/tests/control_plane_ts/usage_statistics_sources.test.ts +++ b/tests/control_plane_ts/usage_statistics_sources.test.ts @@ -163,7 +163,7 @@ test("scope expansion renews disclosure and fences old observations without undo assert.equal((await inspect(path,ctx)).blocked_by,"notice_required"); await observe(path,ctx,prior.generation,null,async()=>{throw new Error("unexpected send");},undefined,cycle("start",at)); await assert.rejects(readFile(path+".cycles"),/ENOENT/); - const enabled=await configure(path,ctx,"enable"); assert.equal(enabled.notice.version,3); + const enabled=await configure(path,ctx,"enable"); assert.equal(enabled.notice.version,4); const current=JSON.parse(await readFile(path,"utf8")); assert.notEqual(current.generation,prior.generation); await configure(path,ctx,"disable"); assert.equal((await inspect(path,ctx)).consent,"disabled"); diff --git a/tests/test_usage_ping.py b/tests/test_usage_ping.py index eea9f78f3..7620b3728 100644 --- a/tests/test_usage_ping.py +++ b/tests/test_usage_ping.py @@ -141,7 +141,7 @@ def test_absent_stderr_keeps_real_cli_json_pure_until_a_stream_discloses(isolate assert main(['version', '--format', 'json']) == 0 assert 'random installation ID' in stderr.getvalue() assert json.loads(capsys.readouterr().out)['ok'] is True - assert json.loads(usage_ping.state_path().read_text())['notice']['version'] == 3 + assert json.loads(usage_ping.state_path().read_text())['notice']['version'] == 4 @pytest.mark.parametrize('setting,value', [ @@ -257,7 +257,15 @@ def test_detached_sender_honors_proxy_and_no_proxy(isolated, collector, monkeypa release.set() -def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated): +@pytest.mark.parametrize('upgrade', [False, True]) +def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated, upgrade): + if upgrade: + usage_ping.control('enable') + prior = json.loads(usage_ping.state_path().read_text()) + prior['consent'] = 'default' + prior['notice']['version'] = 3 + usage_ping.state_path().write_text(json.dumps(prior)) + before = usage_ping.state_path().read_bytes() import http.client from loopx.chat_server import ChatHTTPServer, ChatRequestHandler server = ChatHTTPServer(('127.0.0.1', 0), ChatRequestHandler) @@ -270,7 +278,11 @@ def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated): response = connection.getresponse() initial = json.loads(response.read()) assert response.status == 200 and initial['automatic_notice_required'] - assert not usage_ping.state_path().exists() + if upgrade: + assert usage_ping.state_path().read_bytes() == before + else: + assert not usage_ping.state_path().exists() + assert initial['notice']['version'] == 4 connection.request('POST', path, json.dumps({'notice': initial['notice']}), {'Content-Type': 'application/json'}) response = connection.getresponse() acknowledged = json.loads(response.read()) @@ -320,3 +332,41 @@ def test_goal_observer_does_not_wait_for_unresponsive_collector(isolated, collec assert not release.is_set(), 'host returned before collector released HTTP response' usage_ping.control('disable') release.set() + + +def test_v3_cli_upgrade_requires_visible_renewal_before_real_http(isolated, collector, monkeypatch): + endpoint, received, accepted, release = collector + release.set() + monkeypatch.setenv('LOOPX_USAGE_PING_ENDPOINT', endpoint) + usage_ping.control('enable') + path = usage_ping.state_path() + old = json.loads(path.read_text()) + old.update(consent='default', day='2026-09-28', counters=[{ + 'feature': 'todo', 'outcome': 'ok', 'duration': 'lt_1s', 'error': 'none', 'count': 7, + }]) + old['notice']['version'] = 3 + path.write_text(json.dumps(old)) + before = path.read_bytes() + setup = 'import sys; from pathlib import Path; from loopx import usage_ping; usage_ping.DEFAULT_RUNTIME_ROOT=Path(sys.argv[1]); from loopx.cli_runtime import main; ' + command = [sys.executable, '-c', setup + 'raise SystemExit(main(["version", "--format", "json"]))', str(isolated)] + hidden = subprocess.run(command, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, text=True, timeout=30) + assert hidden.returncode == 0 and json.loads(hidden.stdout)['ok'] + assert path.read_bytes() == before + assert not accepted.wait(0.3) and received == [] + visible = subprocess.run(command, capture_output=True, text=True, timeout=30) + assert visible.returncode == 0 and json.loads(visible.stdout)['ok'] + assert 'first measured CLI result' in visible.stderr and '15 minutes' in visible.stderr + assert 'network timing' in visible.stderr + current = json.loads(path.read_text()) + assert current['notice']['version'] == 4 + assert current['generation'] != old['generation'] + assert current['counters'] == [] + assert not accepted.wait(0.3) and received == [] + subsequent = subprocess.run(command, capture_output=True, text=True, timeout=30) + assert subsequent.returncode == 0 and subsequent.stderr == '' + assert accepted.wait(5) + deadline = time.monotonic() + 5 + while not any(item.get('schema') == 'loopx_usage_aggregate_v1' for item in received) and time.monotonic() < deadline: + time.sleep(0.02) + batches = [item for item in received if item.get('schema') == 'loopx_usage_aggregate_v1'] + assert len(batches) == 1 and sum(row['count'] for row in batches[0]['counters']) == 1