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 @@ -23,8 +23,8 @@ export function UsageStatisticsSettings() {
return <details className="personal-capability-scope-note personal-usage-statistics" data-testid="usage-statistics-settings">
<summary>{zh ? "基础使用统计 · 告知后默认开启,可关闭" : "Basic usage statistics · on after notice, optional"}</summary>
<p>{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."}</p>
? "用于决定平台支持和改进命令体验。每天向 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."}</p>
<p>{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."}</p>
<p>{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."}</p>
{state ? <>
Expand Down
5 changes: 4 additions & 1 deletion apps/usage-collector/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
37 changes: 26 additions & 11 deletions docs/reference/usage-ping.md
Original file line number Diff line number Diff line change
Expand Up @@ -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}]}
Expand Down Expand Up @@ -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
Expand Down
23 changes: 17 additions & 6 deletions docs/reference/usage-ping.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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}]}
Expand Down Expand Up @@ -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` 配置,代理地址与凭证不会进入遥测数据。
Expand Down
40 changes: 27 additions & 13 deletions loopx/control_plane/runtime/usage_statistics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, string | undefined>;
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 {
Expand Down Expand Up @@ -76,6 +77,8 @@ async function load(path: string): Promise<State> {
}
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;
}
Expand Down Expand Up @@ -143,27 +146,26 @@ export async function observe(path: string, ctx: Context, generation: string, co
let heartbeat: Ping | null = null;
let heartbeatRequest: Promise<number> | undefined;
let aggregate: Aggregate | null = null;
let aggregateRequest: Promise<number> | 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);
Expand All @@ -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" };
Expand All @@ -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);
Expand Down
22 changes: 12 additions & 10 deletions tests/control_plane_ts/usage_statistics.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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 => {
Expand Down
Loading
Loading