diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts
index 691de82ec7..c3aec55bb7 100644
--- a/apps/presentation/dashboard/src/data/chat.ts
+++ b/apps/presentation/dashboard/src/data/chat.ts
@@ -2220,6 +2220,11 @@ const usageStatisticsSchema = z.object({
notice: z.object({ version: z.number(), endpoint: z.string(), policy: z.string() }),
automatic_notice_required: z.boolean(),
next_payload: z.unknown(), aggregate_preview: z.unknown(), goal_preview: z.unknown(),
+ diagnostic_preview: z.unknown().optional(), diagnostic_dropped: z.number().optional(),
+ identity_scope: z.string().optional(), delivery_history: z.array(z.object({
+ day: z.string(), channel: z.enum(["heartbeat", "cli", "goal"]), rows: z.number(),
+ status: z.enum(["accepted", "rejected", "unavailable"]),
+ })).optional(),
});
export type UsageStatistics = z.infer;
export async function usageStatistics(enabled?: boolean): Promise {
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 c47e0093ec..c407fce0da 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
@@ -73,11 +73,11 @@ export function UsageStatisticsNotice({ onDetails }: { onDetails: () => void })
: state.automatic_notice_required ? (zh ? "基础使用统计 · 告知后自动开启" : "Basic usage statistics · enabled after this notice")
: (zh ? "基础使用统计当前不发送,请查看详情" : "Basic usage statistics are not sending; see details")}
{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."}
+ ? "用于改进平台支持与使用体验。发送随机安装标识和环境信息,另行汇总 CLI 子操作、版本/活动日期、结果/耗时与回执信号,以及 Goal 时长。环境类型自愿声明,默认未知;不采集对话、代码、路径或参数值。可随时关闭。"
+ : "Helps improve platform support and usage. Sends a random installation ID and environment information, plus separate CLI sub-operation, release/activity day, result/timing, receipt signals and Goal duration summaries. Deployment context is voluntary, unknown by default. No conversations, code, paths or argument values. 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."}
+ ? "首个已测量的 CLI 结果立即上报,后续由使用活动触发,至少间隔 15 分钟发送一批。CLI 汇总不含安装标识。"
+ : "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."}
{zh ? "接收方:" : "Recipient: "}{state.endpoint}
{error ? {zh ? "设置未能保存,请打开详情重试。" : "Could not save this setting. Open details to retry."}
: null}
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 ed170ecae4..4703405646 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
@@ -25,7 +25,7 @@ export function UsageStatisticsSettings() {
{zh
? "用于决定平台支持和改进命令体验。每天向 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 ? "新增固定子操作、版本、UTC 活动日期、阻塞/失败分类和已回读的生命周期信号。运行环境类型仅由 LOOPX_USAGE_CONTEXT 自愿声明,默认 unknown,不推断个人或企业。不会上传提示词、代码、路径、参数值、Goal 内容或原始错误。命令成功不等于 Goal 完成。" : "Adds fixed sub-operations, release version, UTC activity date, blocked/failure classes and receipt-backed lifecycle signals. Deployment context is voluntary via LOOPX_USAGE_CONTEXT, unknown by default, never inferred. No prompts, code, paths, argument values, Goal contents or raw errors. 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 ? <>
)[state.blocked_by ?? ""] ?? "请检查配置"}` : `Not sending: ${state.blocked_by}`)}
{state.notice_required && state.consent !== "disabled" ? void update(true)}>{zh ? "已了解,启用统计" : "Understood, enable statistics"} : null}
{zh ? "接收地址:" : "Recipient: "}{state.endpoint ?? (zh ? "未配置" : "Not configured")}
- {zh ? "查看待发送数据" : "Preview outgoing data"} {JSON.stringify({ heartbeat: state.next_payload, aggregate: state.aggregate_preview, goals: state.goal_preview }, null, 2)}
+ {zh ? "查看待发送数据与本地发送摘要" : "Preview outgoing data and local delivery summaries"} {JSON.stringify({ heartbeat: state.next_payload, aggregate: state.aggregate_preview, diagnostics: state.diagnostic_preview, goals: state.goal_preview, identity_scope: state.identity_scope, dropped: state.diagnostic_dropped, delivery_history: state.delivery_history }, null, 2)}
> : null}
{error ? {zh ? "无法读取或保存;请用终端检查:" : "Could not read or save; inspect in terminal: "}loopx usage-ping status
: null}
loopx usage-ping disable · LOOPX_USAGE_PING=0
diff --git a/apps/usage-collector/README.md b/apps/usage-collector/README.md
index 80157ed628..45d1a9b030 100644
--- a/apps/usage-collector/README.md
+++ b/apps/usage-collector/README.md
@@ -2,7 +2,8 @@
Cloudflare Worker + D1 for [basic usage statistics](../../docs/reference/usage-ping.md).
The TypeScript client/collector allowlist lives in
-`loopx/control_plane/runtime/usage_statistics_contract.ts`.
+`loopx/control_plane/runtime/usage_statistics_contract.ts` and
+`usage_statistics_diagnostics.ts`.
| Endpoint | Contract |
|---|---|
@@ -12,6 +13,8 @@ The TypeScript client/collector allowlist lives in
| `GET /v1/goal-stats` | Independent 30-day duration histograms; cells below 5 omitted |
| `GET /v0/stats` | Deduplicated active/new installations, including retained v0 clients; version/OS/CPU/channel breakdown |
| `GET /v1/aggregate-stats` | Independent 30-day feature/result/duration/error totals; cells below 5 omitted |
+| `GET /v1/diagnostic-stats` | Independent versioned CLI/lifecycle marginals; cells below 5 omitted |
+| `GET /v1/adoption-stats` | Mature 1/7/30-day installation return cohorts and 30-day activity-day buckets |
| `POST /v0/ping` | Retained six-field opt-in client contract; no new default-on clients use this route |
Heartbeats are deduplicated by installation/day and retained 400 days.
@@ -24,6 +27,16 @@ 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.
+The aggregate endpoint also accepts `loopx_usage_diagnostics_v1`.
+`diagnostic_counts` stores fixed feature/sub-operation/result/reason/duration,
+numeric version, UTC activity date, receipt date, voluntary deployment context
+and receipt-backed lifecycle signal. It has no installation join key or raw
+request rows. Activity dates older than seven days or in the future are rejected.
+Keep 30 receipt days; legacy counters are never backfilled or reattributed.
+Public endpoints expose independent marginals, not multi-dimensional histories.
+Return cohorts count installation state, not people; `within_Nd` means any later
+heartbeat within a mature N-day window, not exact day-N retention.
+
Neither handler reads/stores IP, user agent or Cloudflare request metadata.
The template disables Worker observability; Cloudflare still handles network
metadata. Do not describe the identified heartbeat as fully anonymous, or the
@@ -61,6 +74,12 @@ must never be summed. The unreleased Goal payload requires measurement/Host
labels and uses `duration` instead of `execution`.
A Worker rollback can leave that additive table intact.
+Before releasing notice-v5 clients, apply `0004-diagnostics.sql` with D1
+migrations and deploy the updated Worker. It preserves existing installs,
+pings and legacy counters. An older Worker rejects new diagnostics; loss is
+not retried. Roll back the Worker/client without dropping the additive table.
+Merging this code does not deploy the collector.
+
Qualify `/v1/ping`, `/v1/aggregate`, `/v1/goals`, all stats endpoints, and invalid-field/size
rejections on a separate database first. Deploy the collector before releasing
the new client default: the v0-only Worker does not accept v1 requests. Server
@@ -83,4 +102,4 @@ node --no-warnings --experimental-strip-types --test tests/control_plane_ts/usag
```
Run from repository root. The collector suite executes actual SQL in SQLite,
-including both additive migrations; no production telemetry is needed for these tests.
+including additive migrations and mature/suppressed cohorts; no production telemetry is needed for these tests.
diff --git a/apps/usage-collector/migrations/0004-diagnostics.sql b/apps/usage-collector/migrations/0004-diagnostics.sql
new file mode 100644
index 0000000000..00299f95db
--- /dev/null
+++ b/apps/usage-collector/migrations/0004-diagnostics.sql
@@ -0,0 +1,8 @@
+-- Additive: retain v0/v1 history and rollback tables, never relabel old counts.
+CREATE TABLE IF NOT EXISTS diagnostic_counts (
+ receipt_day TEXT NOT NULL, activity_day TEXT NOT NULL, version TEXT NOT NULL,
+ context TEXT NOT NULL, feature TEXT NOT NULL, operation TEXT NOT NULL,
+ outcome TEXT NOT NULL, error TEXT NOT NULL, duration TEXT NOT NULL,
+ signal TEXT NOT NULL, count INTEGER NOT NULL,
+ PRIMARY KEY (receipt_day, activity_day, version, context, feature, operation, outcome, error, duration, signal)
+);
diff --git a/apps/usage-collector/schema.sql b/apps/usage-collector/schema.sql
index e5bb349ee1..6ae33695f8 100644
--- a/apps/usage-collector/schema.sql
+++ b/apps/usage-collector/schema.sql
@@ -39,3 +39,11 @@ CREATE TABLE IF NOT EXISTS goal_duration_counts (
span TEXT NOT NULL, duration TEXT NOT NULL, count INTEGER NOT NULL,
PRIMARY KEY (day, measurement, host, span, duration)
);
+
+CREATE TABLE IF NOT EXISTS diagnostic_counts (
+ receipt_day TEXT NOT NULL, activity_day TEXT NOT NULL, version TEXT NOT NULL,
+ context TEXT NOT NULL, feature TEXT NOT NULL, operation TEXT NOT NULL,
+ outcome TEXT NOT NULL, error TEXT NOT NULL, duration TEXT NOT NULL,
+ signal TEXT NOT NULL, count INTEGER NOT NULL,
+ PRIMARY KEY (receipt_day, activity_day, version, context, feature, operation, outcome, error, duration, signal)
+);
diff --git a/apps/usage-collector/src/basic-usage.ts b/apps/usage-collector/src/basic-usage.ts
index 6e65069b59..4a47043fe7 100644
--- a/apps/usage-collector/src/basic-usage.ts
+++ b/apps/usage-collector/src/basic-usage.ts
@@ -1,10 +1,29 @@
/** Aggregate storage has no foreign key or identifier linking it to installations. */
import { validAggregate, validPing } from "../../../loopx/control_plane/runtime/usage_statistics_contract.ts";
import type { Aggregate } from "../../../loopx/control_plane/runtime/usage_statistics_contract.ts";
+import { validDiagnostics } from "../../../loopx/control_plane/runtime/usage_statistics_diagnostics.ts";
+import type { DiagnosticAggregate } from "../../../loopx/control_plane/runtime/usage_statistics_diagnostics.ts";
type Statement = { bind(...args: unknown[]): Statement; all(): Promise<{ results: Record[] }> };
type Database = { prepare(sql: string): Statement; batch(statements: Statement[]): Promise };
export { validAggregate, validPing };
+export { validDiagnostics };
+export async function recordDiagnostics(db: Database, value: DiagnosticAggregate, receiptDay: string) {
+ await db.batch(value.counters.map(row => db.prepare(
+ "INSERT INTO diagnostic_counts (receipt_day, activity_day, version, context, feature, operation, outcome, error, duration, signal, count) " +
+ "VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) " +
+ "ON CONFLICT (receipt_day, activity_day, version, context, feature, operation, outcome, error, duration, signal) DO UPDATE SET count = count + excluded.count",
+ ).bind(receiptDay, row.activity_day, row.version, row.context, row.feature, row.operation, row.outcome, row.error, row.duration, row.signal, row.count)));
+}
+export async function diagnosticStats(db: Database, since: string) {
+ const totals: Record> = {};
+ for (const column of ["feature", "operation", "outcome", "error", "duration", "version", "context", "signal"]) {
+ const rows = await db.prepare(`SELECT ${column} AS key, SUM(count) AS n FROM diagnostic_counts WHERE receipt_day >= ?1 GROUP BY ${column}`).bind(since).all();
+ totals[column] = {};
+ for (const row of rows.results) if (Number(row.n) >= 5) totals[column][String(row.key)] = Number(row.n);
+ }
+ return { schema: "loopx_usage_diagnostic_stats_v1", definition: "Independent marginal totals from versioned, lossy observations in the last 30 receipt days. Signals describe observed transitions, not unique Goals, people, independent outcome quality or billing. Context is voluntary self-report, never inferred. Cells below 5 omitted.", totals };
+}
export async function recordAggregate(db: Database, value: Aggregate, day: string) {
await db.batch(value.counters.map(row => db.prepare(
"INSERT INTO usage_counts (day, feature, outcome, duration, error, count) VALUES (?1, ?2, ?3, ?4, ?5, ?6) " +
diff --git a/apps/usage-collector/src/collector.js b/apps/usage-collector/src/collector.js
index 80fa4778e5..8581232057 100644
--- a/apps/usage-collector/src/collector.js
+++ b/apps/usage-collector/src/collector.js
@@ -1,4 +1,5 @@
import { validAggregate, validPing, recordAggregate, aggregateStats, validGoalAggregate, recordGoals, goalStats } from "./basic-usage.ts";
+import { validDiagnostics, recordDiagnostics, diagnosticStats } from "./basic-usage.ts";
// Pure request handling for the LoopX usage collector. worker.js binds it to
// Cloudflare; tests bind it to an in-memory database.
@@ -138,13 +139,45 @@ export async function purge(db, day) {
db.prepare("DELETE FROM goal_duration_counts WHERE day < ?1").bind(shiftDays(day, -30)),
db.prepare("DELETE FROM goal_usage_counts WHERE day < ?1").bind(shiftDays(day, -30)),
db.prepare("DELETE FROM usage_counts WHERE day < ?1").bind(shiftDays(day, -30)),
+ db.prepare("DELETE FROM diagnostic_counts WHERE receipt_day < ?1").bind(shiftDays(day, -30)),
db.prepare("DELETE FROM installs WHERE install_id NOT IN (SELECT DISTINCT install_id FROM pings)"),
]);
}
+export async function adoptionStats(db, day) {
+ const cohorts = {};
+ for (const horizon of [1, 7, 30]) {
+ const row = await db.prepare(
+ "SELECT COUNT(*) AS eligible, COALESCE(SUM(EXISTS(SELECT 1 FROM pings p WHERE p.install_id = installs.install_id " +
+ "AND p.day > installs.first_day AND p.day <= date(installs.first_day, ?1))), 0) AS returned " +
+ "FROM installs WHERE first_day BETWEEN ?2 AND ?3",
+ ).bind(`+${horizon} days`, shiftDays(day, -horizon - 29), shiftDays(day, -horizon)).first();
+ // Do not publish a rate whose numerator or denominator is a small cell.
+ cohorts[`within_${horizon}d`] = Number(row?.eligible) >= MIN_BUCKET && Number(row?.returned) >= MIN_BUCKET
+ && (Number(row.eligible) === Number(row.returned) || Number(row.eligible) - Number(row.returned) >= MIN_BUCKET)
+ ? { eligible: Number(row.eligible), returned: Number(row.returned) } : null;
+ }
+ const rows = await db.prepare(
+ "SELECT CASE WHEN days = 1 THEN '1' WHEN days <= 3 THEN '2_3' WHEN days <= 7 THEN '4_7' " +
+ "WHEN days <= 14 THEN '8_14' ELSE '15_30' END AS key, COUNT(*) AS installs FROM " +
+ "(SELECT install_id, COUNT(*) AS days FROM pings WHERE day BETWEEN ?1 AND ?2 GROUP BY install_id) GROUP BY key",
+ ).bind(shiftDays(day, -29), day).all();
+ return { schema: "loopx_installation_return_stats_v1", generated_on: day,
+ definition: "Each cohort contains 30 first-seen UTC dates ending at least N days ago. Return means a heartbeat on a later day within N days, not exact day-N retention. Persistent installation state is not a person or organization; resets and ephemeral hosts affect counts. Small numerator/denominator or nonzero complement cohorts omitted. Horizon populations overlap.",
+ cohorts, active_day_distribution: suppressSmall(rows.results) };
+}
+
export async function handle(request, db, now = new Date()) {
const url = new URL(request.url);
const day = utcDay(now);
+ if (url.pathname === "/v1/adoption-stats") {
+ if (request.method !== "GET") return json({ error: "method not allowed" }, 405);
+ return json(await adoptionStats(db, day), 200, { "cache-control": "public, max-age=3600" });
+ }
+ if (url.pathname === "/v1/diagnostic-stats") {
+ if (request.method !== "GET") return json({ error: "method not allowed" }, 405);
+ return json(await diagnosticStats(db, shiftDays(day, -29)), 200, { "cache-control": "public, max-age=3600" });
+ }
if (url.pathname === "/v1/goal-stats") {
if (request.method !== "GET") return json({ error: "method not allowed" }, 405);
return json(await goalStats(db, shiftDays(day, -29)), 200, { "cache-control": "public, max-age=3600" });
@@ -187,6 +220,11 @@ export async function handle(request, db, now = new Date()) {
return new Response(null, { status: 204 });
}
if (url.pathname === "/v1/aggregate") {
+ if (validDiagnostics(parsed)) {
+ if (parsed.counters.some(row => row.activity_day > day || row.activity_day < shiftDays(day, -7))) return json({ error: "activity day outside retention window" }, 400);
+ await recordDiagnostics(db, parsed, day);
+ return new Response(null, { status: 204 });
+ }
if (!validAggregate(parsed)) return json({ error: "invalid aggregate" }, 400);
await recordAggregate(db, parsed, day);
return new Response(null, { status: 204 });
diff --git a/apps/usage-collector/test/collector.test.mjs b/apps/usage-collector/test/collector.test.mjs
index 615173d586..d787ff2fe7 100644
--- a/apps/usage-collector/test/collector.test.mjs
+++ b/apps/usage-collector/test/collector.test.mjs
@@ -204,3 +204,57 @@ test("measurement migration preserves historical counts and is safe to repeat",
assert.equal(db.prepare("SELECT count FROM goal_usage_counts").get().count,8);
db.close();
});
+
+test("diagnostics store separate activity/receipt days, preserve legacy data and reject private/late fields", async () => {
+ const db = d1();
+ const send = value => new Request("https://collector.example/v1/aggregate", {
+ method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(value),
+ });
+ const row = { feature: "pr-review", operation: "merge-readiness", outcome: "blocked", error: "not_ready", duration: "lt_1s", count: 6,
+ version: "1.2.3", activity_day: "2026-09-29", context: "maintainer", signal: "none" };
+ const value = { schema: "loopx_usage_diagnostics_v1", counters: [row] };
+ for (const field of ["install_id", "goal_id", "arguments", "ip", "raw_error"]) {
+ assert.equal((await handle(send({ ...value, [field]: "private" }), db, at("2026-09-30"))).status, 400);
+ }
+ assert.equal((await handle(send(value), db, at("2026-09-30"))).status, 204);
+ assert.equal((await handle(send({ ...value, counters: [{ ...row, activity_day: "2026-10-01" }] }), db, at("2026-09-30"))).status, 400);
+ assert.equal((await handle(send({ ...value, counters: [{ ...row, activity_day: "2026-09-01" }] }), db, at("2026-09-30"))).status, 400);
+ const stored = db.raw.get("SELECT * FROM diagnostic_counts");
+ assert.equal(stored.receipt_day, "2026-09-30"); assert.equal(stored.activity_day, "2026-09-29");
+ assert.equal("install_id" in stored, false);
+ const stats = await (await handle(new Request("https://collector.example/v1/diagnostic-stats"), db, at("2026-09-30"))).json();
+ assert.deepEqual(stats.totals.outcome, { blocked: 6 });
+ assert.deepEqual(stats.totals.context, { maintainer: 6 });
+ await purge(db, "2026-11-01");
+ assert.equal(db.raw.get("SELECT COUNT(*) n FROM diagnostic_counts").n, 0);
+});
+
+test("return cohorts require matured observation windows, deduplicate days and omit small cells", async () => {
+ const db = d1();
+ for (let n = 1; n <= 10; n++) {
+ await handle(post(ping(n)), db, at("2026-09-01"));
+ if (n <= 5) {
+ await handle(post(ping(n)), db, at("2026-09-03"));
+ await handle(post(ping(n)), db, at("2026-09-03"));
+ }
+ }
+ await handle(post(ping(11)), db, at("2026-09-29")); // too recent for 7-day eligibility
+ const stats = await (await handle(new Request("https://collector.example/v1/adoption-stats"), db, at("2026-09-30"))).json();
+ assert.deepEqual(stats.cohorts.within_7d, { eligible: 10, returned: 5 });
+ assert.equal(stats.cohorts.within_1d, null); assert.equal(stats.cohorts.within_30d, null);
+ assert.deepEqual(stats.active_day_distribution, { "1": 6, "2_3": 5 });
+ assert.equal(JSON.stringify(stats).includes(id(1)), false);
+ await handle(post(ping(6)), db, at("2026-09-04"));
+ const smallComplement = await (await handle(new Request("https://collector.example/v1/adoption-stats"), db, at("2026-09-30"))).json();
+ assert.equal(smallComplement.cohorts.within_7d, null);
+});
+
+test("diagnostics migration is repeatable and never reattributes legacy observations", () => {
+ const db = new DatabaseSync(":memory:");
+ db.exec("CREATE TABLE usage_counts (count INTEGER); INSERT INTO usage_counts VALUES (7)");
+ const migration = readFileSync(new URL("../migrations/0004-diagnostics.sql", import.meta.url), "utf8");
+ db.exec(migration); db.exec(migration);
+ assert.equal(db.prepare("SELECT count FROM usage_counts").get().count, 7);
+ assert.equal(db.prepare("SELECT COUNT(*) n FROM diagnostic_counts").get().n, 0);
+ db.close();
+});
diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md
index dd93ea9c5e..0f9a89580a 100644
--- a/docs/reference/usage-ping.md
+++ b/docs/reference/usage-ping.md
@@ -24,9 +24,10 @@ Lark has no separate switch and cannot override the machine owner's choice.
| Question | Evidence | Limit |
|---|---|---|
| Which versions/platforms need support? | Daily version, OS, CPU architecture, Python minor and install channel | Only reporting installations |
-| Do installations keep using LoopX? | Random installation ID, deduplicated by UTC day | Installations, not people; reinstall/re-enable may count again |
-| Which CLI entry points are used? | Fixed command-family counters | Polling and automation count too; not a measure of user value |
-| Which commands fail or take time? | Result, typed error category, coarse elapsed-time bucket | CLI return status is not Goal acceptance; handled domain failures may return 0 |
+| Do installations keep using LoopX? | Random installation ID, daily deduplication, mature 1/7/30-day return windows | Persistent state directories, not people or organizations; deleting state/re-enabling creates a new ID |
+| Which CLI entry points are used? | Fixed command-family/sub-operation counters, release version and UTC activity date | Polling and automation count too; not a measure of user value |
+| Which commands fail or take time? | Typed result/reason and coarse elapsed-time bucket | Merge-readiness holds are `blocked`, not failures; unclassified failures remain `command_failed` |
+| Is meaningful work observed? | Receipt-backed registration, Turn commitment, Todo completion/validation and verified return | Partial transition counts, not unique Goals or independently judged quality |
The first version measures CLI invocations, including those made by agents.
Top-level `--help`/`--version` fast paths, native exec-replaced scheduler
@@ -34,7 +35,7 @@ followups, API-only interactions and individual App/Lark actions are not
instrumented. Heartbeats start in the background before dispatch; long-running server command
results are counted only when the CLI returns. There is no claim to complete product activity or task success rates.
-## Two separate payloads
+## Separate payload contracts
**Daily heartbeat** (`POST /v1/ping`), with exactly these fields:
@@ -42,19 +43,62 @@ results are counted only when the CLI returns. There is no claim to complete pro
{"schema":"loopx_usage_ping_v1","install_id":"00000000-0000-4000-8000-000000000001","version":"1.2.0","os":"linux","arch":"x64","python":"3.13","channel":"pip"}
```
-The ID is random, local to the installation and not derived from hardware or an
-account. It enables cross-day association, so this is **not fully anonymous**.
+The ID is random, local to the persistent machine-state directory and not derived
+from hardware or an account. Sessions and disclosure upgrades retain it; explicit
+disable deletes it. Ephemeral homes, deleted/copied state and re-enabling distort
+installation counts. It enables cross-day association, so this is **not fully anonymous**.
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.
-**CLI aggregate batch** (`POST /v1/aggregate`):
+**Current CLI diagnostics** (`POST /v1/aggregate`, notice revision 5):
+
+```json
+{"schema":"loopx_usage_diagnostics_v1","counters":[{"feature":"pr-review","operation":"merge-readiness","outcome":"blocked","error":"not_ready","duration":"lt_1s","count":4,"version":"1.2.3","activity_day":"2026-09-30","context":"unknown","signal":"none"}]}
+```
+
+The new default adds numeric release version, UTC activity **date** (not event
+time), fixed sub-operation, result/reason and receipt-backed lifecycle signal.
+Deployment context is optional self-report via `LOOPX_USAGE_CONTEXT`:
+`unknown` (default), `personal`, `shared_service`, `ephemeral`,
+`organization_managed`, `maintainer`. Invalid values become `unknown`; no
+company name, person or hardware topology is inferred or sent. Set `maintainer`
+on maintainer processes to distinguish their **future diagnostics** in operator
+analysis. This does not label heartbeats or identify historical ID-free counts.
+Use the shared disable switch to exclude a machine from all channels.
+
+The collector adds a separate receipt date. It accepts activity dates from the
+previous seven days through today; old/future packets are rejected. No ID, Goal
+or request timestamp enters this table. Both ends share the typed schema.
+Operations come from parsed command structure, never argument values:
+
+- `turn`: `plan|run-once|status`; `quota`: `status|plan|should-run|spend-slot|monitor-poll`.
+- `todo`: `list|add|claim|update|complete`.
+- `project`: `register|resolve|bind-session|unbind-session`.
+- `pr-review`: `merge-readiness|check-result`; other operations are `default`.
+- Verified return observation uses `other/result-return`.
+
+Results are `ok|blocked|failed|cancelled`. Reasons are
+`none|not_ready|invalid_input|permission|not_found|timeout|connection|interrupted|command_failed`.
+An otherwise valid merge-readiness response with `ready=false` records
+`blocked/not_ready` without changing its original nonzero exit code. Errors
+come from exception classes and typed booleans, never parsed error prose.
+
+Signals are `none|project_registered|managed_turn_committed|todo_completed|todo_validated|result_returned`.
+They require existing changed/committed receipts, not mere exit success. Dry
+runs and unchanged Todo completion do not count; Turn replays do not pass the
+existing committed-current-effects predicate. Passed completion validation is
+required for `todo_validated`. Return currently covers exact-source manager-context
+deliveries verified by the provider and newly settled as delivered, not legacy
+or unverified paths. None of these observations changes work authority.
+
+**Retained legacy CLI aggregate** (same endpoint):
```json
{"schema":"loopx_usage_aggregate_v1","counters":[{"feature":"todo","outcome":"ok","duration":"lt_1s","error":"none","count":4}]}
```
-No installation ID, version, timestamp, Goal or other join key is included.
+No installation ID, version, timestamp, Goal or other join key is included in this legacy contract.
The collector adds its reception day, merges counters, and stores no individual
request rows. Both sender and collector use the same strict TS allowlist:
@@ -146,10 +190,12 @@ 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
-buffer shapes remain readable. Notice revision 4 renews disclosure before the
-faster cadence takes effect: old notice state cannot send or consume buffers.
+Legacy counts are capped at 128 distinct rows; new diagnostics at 32, with
+10,000 per row. Diagnostics expire by activity date after seven UTC days;
+legacy expiry uses the oldest buffered day. Overflow records a bounded local
+`diagnostic_dropped` count, not an unbounded queue. Existing buffers remain
+readable. Notice revision 5 renews disclosure before the expanded default 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
@@ -162,10 +208,10 @@ 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
+Current diagnostics contain version/activity date, but no installation ID or
+event time. Historical legacy counts have neither version nor activity date
+and cannot be retrospectively reattributed. Do not divide either channel's
+totals by reporting installations to infer per-install usage. 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
@@ -175,6 +221,10 @@ never enter telemetry payloads. Disable deletes the local ID and buffered
counts; an old worker cannot restore them or send the next channel. A request
already handed to the network cannot be recalled. Re-enable uses a new ID.
Corrupt or unsupported state fails closed; explicit disable is the repair path.
+Status/Settings show at most 20 local delivery summaries (UTC date, channel,
+row count, accepted/rejected/unavailable). No request bodies, errors or URLs
+enter this journal. Disable erases it. HTTP acceptance does not prove durable
+collector storage or work acceptance; an unsent preview is not a receipt.
## Collector and interpretation
@@ -184,6 +234,20 @@ The existing `/v0/stats` endpoint continues to report deduplicated installations
from old and new heartbeat clients. `/v1/aggregate-stats` publishes independent
feature/result/duration/error totals, omitting cells below five. It does not
publish cross-dimensional combinations or per-install behavior histories.
+`/v1/diagnostic-stats` publishes independent 30-receipt-day marginals for feature,
+operation, result, reason, duration, version, context and signal, suppressing
+cells below five. Operator-only queries may exclude `context='maintainer'` and
+compare activity/receipt dates; those counts still cannot join installations.
+Apply additive migration `0004-diagnostics.sql` and deploy the collector before
+releasing notice-v5 clients. An older Worker rejects new packets lossily;
+retain the table on rollback. No historical diagnostic rows are invented.
+
+`/v1/adoption-stats` uses existing heartbeat data: each 1/7/30-day cohort includes
+30 first-seen UTC dates whose full horizon has elapsed. Return means any later
+heartbeat **within** that horizon, not exact day-N retention. It also publishes
+30-day active-day buckets. Small cohorts are suppressed; the windows overlap
+and are not additive. Many installations or high activity do not prove
+enterprise adoption, distinct people or a paid customer.
Application code stores no client IP, user agent or Cloudflare request metadata;
Worker observability is disabled in the deployment template. Network providers
@@ -263,7 +327,7 @@ measurement histograms over 30 receipt days, omitting cells below five.
The existing settings switch, environment opt-outs and consent policy control
all channels and local timing reads. Settings and `loopx usage-ping status`
show `goal_preview`, a local snapshot rather than a delivery receipt. Expanded
-scope requires notice version 3; an existing explicit disable persists.
+scope requires the current notice version 5; an existing explicit disable persists.
Before shipping the client, back up D1, apply `0002-goal-usage.sql` and
`0003-goal-duration-sources.sql`, then deploy the Worker. The latter migrates
diff --git a/docs/reference/usage-ping.zh-CN.md b/docs/reference/usage-ping.zh-CN.md
index adb9acbd9e..5cd97f83d8 100644
--- a/docs/reference/usage-ping.zh-CN.md
+++ b/docs/reference/usage-ping.zh-CN.md
@@ -19,18 +19,19 @@ loopx usage-ping enable # 阅读告知后明确开启
## 能回答什么
- 每日版本、系统、CPU 架构、Python 小版本和安装渠道:哪些环境需要优先维护。
-- 随机安装 ID 的跨日心跳:有多少安装持续使用。它统计安装而非用户;重新安装或
- 关闭后再开启可能计为新安装。现有统计接口提供活跃与新增;更细留存报表未实现。
-- 固定 CLI 功能分类与次数:哪些入口常用。自动化轮询也会计数,不能当成用户价值。
-- 命令结果、错误类别和耗时区间:哪些入口失败或慢。退出码为 0 不等于 Goal 完成;
- 某些已处理的业务阻塞也可能返回 0。
+- 随机安装 ID 跨日心跳:持久机器状态目录的活跃与成熟的 1/7/30 天回访,不是
+ 用户或组织数;删除状态或关闭后重开可能计为新安装。
+- 固定 CLI 功能/子操作、版本和 UTC 活动日期:哪些入口常用,自动化轮询也计数。
+- 类型化结果、原因和耗时区间:哪些入口失败或慢;合并条件未满足单列 `blocked`。
+- 回执支持的注册、Turn 提交、Todo 完成/验证、结果回传:观察推进,不靠退出成功
+ 推断 Goal 完成,也不宣称独立验证了结果质量。
第一版只计 CLI 调用,包括 Agent 发起的命令。顶层 `--help`/`--version` 快速路径、
原生 exec 替换的 scheduler followup、纯 API 操作以及 App/Lark 内的每一次交互不计数。
心跳在命令调度前由后台发送;长驻服务的命令结果只在 CLI 返回时计数,
不宣称覆盖全部产品使用或任务成功率。
-## 两种数据包
+## 分离的数据契约
每日心跳 `POST /v1/ping`:
@@ -38,18 +39,54 @@ loopx usage-ping enable # 阅读告知后明确开启
{"schema":"loopx_usage_ping_v1","install_id":"00000000-0000-4000-8000-000000000001","version":"1.2.0","os":"linux","arch":"x64","python":"3.13","channel":"pip"}
```
-ID 随机生成,不绑定账号、不从硬件派生,但能跨天关联,因此不能称为完全匿名。
+ID 随机生成,属于持久机器状态目录,不绑定账号、不从硬件派生。普通 session 和
+告知升级保留 ID,明确关闭删除 ID。临时 home、删/复制状态或重新开启会影响统计。
+它能跨天关联,因此不能称为完全匿名。
系统只允许 `darwin|linux|windows|other`,架构只允许 `x64|arm64|x86|other`,
安装渠道只允许 `pip|local_release|source|unknown`。版本只接受数字三段式,包含
自定义后缀的版本不会上传。
-独立的 CLI 汇总批次 `POST /v1/aggregate`:
+当前 CLI 诊断 `POST /v1/aggregate`(告知版本 5):
+
+```json
+{"schema":"loopx_usage_diagnostics_v1","counters":[{"feature":"pr-review","operation":"merge-readiness","outcome":"blocked","error":"not_ready","duration":"lt_1s","count":4,"version":"1.2.3","activity_day":"2026-09-30","context":"unknown","signal":"none"}]}
+```
+
+默认新增数字三段版本、UTC 活动**日期**(非事件时间)、固定子操作、结果/原因及
+回执支持的生命周期信号。环境类型由 `LOOPX_USAGE_CONTEXT` 自愿声明:
+`unknown`(默认)、`personal`、`shared_service`、`ephemeral`、`organization_managed`、
+`maintainer`;无效值成为 `unknown`,不猜企业、人数或机器拓扑,不接受公司名称。
+维护者可声明 `maintainer`,分开**未来诊断计数**;不会标注心跳,也无法追溯识别旧
+无 ID 汇总。需要排除整机所有采集时仍用统一关闭开关。
+
+收集器另加接收日期,仅接收此前七天至当天的活动日期,拒绝过期或未来数据。
+表中没有安装 ID、Goal 或请求时间。两端共用 TS 契约;子操作来自解析器结构,
+不读取参数值:
+
+- `turn`:`plan|run-once|status`;`quota`:`status|plan|should-run|spend-slot|monitor-poll`。
+- `todo`:`list|add|claim|update|complete`。
+- `project`:`register|resolve|bind-session|unbind-session`。
+- `pr-review`:`merge-readiness|check-result`;其他操作归为 `default`。
+- 已验证回传观测使用 `other/result-return`。
+
+结果为 `ok|blocked|failed|cancelled`;原因固定为
+`none|not_ready|invalid_input|permission|not_found|timeout|connection|interrupted|command_failed`。
+合并检查正常返回 `ready=false` 时记 `blocked/not_ready`,原非零退出码不改。
+错误分类依据异常类型与明确布尔字段,不解析错误文字。
+
+信号为 `none|project_registered|managed_turn_committed|todo_completed|todo_validated|result_returned`。
+需要已有 changed/committed 回执,不靠成功退出推断;dry-run、未变化的 Todo 完成
+不计新转移;Turn 重放不满足既有“本次 effects 已提交”判定。完成回执确有 passed
+验证才记 `todo_validated`。回传目前仅覆盖 manager-context 精确来源、provider 已
+验证且本次落为 delivered 的路径,旧路径或未验证回传不覆盖,不改变工作授权。
+
+保留的旧 CLI 汇总契约(同一接收地址):
```json
{"schema":"loopx_usage_aggregate_v1","counters":[{"feature":"todo","outcome":"ok","duration":"lt_1s","error":"none","count":4}]}
```
-不带安装 ID、版本、时间戳、Goal 或其他关联键。接收端只添加接收日期并累加计数,
+这一旧契约不带安装 ID、版本、时间戳、Goal 或其他关联键。接收端只添加接收日期并累加计数,
不保存逐次请求行。发送端和收集器共用严格的 TS 白名单:
- 功能:`status|quota|todo|turn|project|connect|pr-review|version|chat|other`。
@@ -116,8 +153,9 @@ authority provider 备份、公共投影。普通命令只读取很小的本地
计数,设置和 status 查询不触发发送。此行为替代已开启统计的交互式、无人值守 CLI
原有的次日发送机制;Goal 时长快照仍按天发送。
-本地汇总最多 128 种计数组合,每项封顶 10,000;以最早积压日期计算,超过七个 UTC
-日的计数丢弃。旧按日缓冲格式仍可读取。告知版本 4 要求在新频率生效前重新告知:
+旧汇总最多 128 种组合,新诊断最多 32 种,每项封顶 10,000。诊断按活动日期过期,
+旧汇总以最早积压日期计算,超过七个 UTC 日丢弃。溢出仅记有界本地
+`diagnostic_dropped`,不建无限队列;旧缓冲可读。告知版本 5 要求扩大默认范围前重新告知:
确认前不发送、不消费缓存;确认后丢弃旧范围计数并更新 generation,阻止旧排队
观测补发,后续新测量才可发送。明确关闭继续生效,`consent_required` 仍须明确启用。
每批在发送前持久化认领时间、移除对应计数,并在同一短锁内发起请求,网络等待不占锁。失败或时钟回拨也不能绕过间隔。
@@ -125,8 +163,8 @@ authority provider 备份、公共投影。普通命令只读取很小的本地
尾部计数;优化减少对次日回访的依赖,但不承诺完整采集。锁竞争、进程退出和网络
故障也可能丢数。这是有损诊断,不能当账单或审计日志。
-CLI 汇总仍不携带安装 ID、版本或活动日期;服务端按 UTC 接收日期汇总,可能与实际
-使用日期不同。不能拿调用量除以上报安装数推算每安装用量,也不能归因到某个版本。
+新诊断带版本和活动日期,不带安装 ID 或事件时间;旧计数没有版本/活动日期,不能
+追溯补归因。两种计数都不能除以上报安装数推算每安装用量。
更频繁的请求可能增加网络时序关联的机会;载荷不带标识不代表无法关联。
每个网络请求限时 3 秒,不阻塞命令完成、不改变输出和退出码。使用受支持 Node
@@ -134,6 +172,9 @@ CLI 汇总仍不携带安装 ID、版本或活动日期;服务端按 UTC 接
关闭删除本机 ID 和
待发送计数;旧后台进程不能恢复它们,也不能继续发送下一条通道。已经交给网络的
请求无法撤回;重新开启会生成新 ID。损坏或未知格式状态拒绝发送,可明确 disable 修复。
+status/设置另显示最多 20 条本地发送摘要:UTC 日期、通道、行数和接收/拒绝/不可达。
+不记请求正文、错误或 URL,关闭清空摘要。HTTP 接收不等于收集器持久化或工作验收;
+待发送预览也不是发送回执。
## 服务端与解释边界
@@ -141,6 +182,16 @@ CLI 汇总仍不携带安装 ID、版本或活动日期;服务端按 UTC 接
汇总计数保存 30 天。原 `/v0/stats` 继续提供新旧客户端的去重活跃和新增安装数。
`/v1/aggregate-stats` 分别给出功能、结果、耗时和错误总量,低于 5 的格子不公开,
不公开多维组合或逐安装行为历史。
+`/v1/diagnostic-stats` 分别提供最近 30 个接收日的功能、子操作、结果、原因、耗时、
+版本、环境和信号总量,低于 5 的单元不公开。仅运营侧查询可排除
+`context='maintainer'`、比较活动/接收日期,仍不能关联安装。
+先应用增量迁移 `0004-diagnostics.sql` 并部署 Worker,再发布告知 v5 客户端;旧 Worker
+会有损地拒绝新包,回滚可保留新增表,不给历史数据虚构诊断行。
+
+`/v1/adoption-stats` 基于已有心跳,每个 1/7/30 天 cohort 取完整观察窗口已结束的
+30 个首次出现日期。回访指 N 天内任意后续日期有心跳,不是严格第 N 天留存。
+另给最近 30 天活跃天数分档,抑制小样本;各窗口重叠,不能相加。安装多、调用多
+不足以判断企业采用、去重人数或付费客户。
应用代码不保存 IP、User-Agent 或 Cloudflare 请求元数据;部署模板关闭 Worker
observability。但网络服务商仍处理连接信息,分开数据包不保证绝对不可关联。
@@ -199,7 +250,7 @@ Codex 发现使用既有 Goal/agent/task 绑定和所选 `CODEX_HOME` 的只读
分布,少于 5 的单元不公开。设置与 `loopx usage-ping status` 的 `goal_preview`
是本机当前快照,不是发送回执。
-沿用统一开关、环境变量和同意策略,关闭也停止本地时间读取。扩大范围需要第 3 版
+沿用统一开关、环境变量和同意策略,关闭也停止本地时间读取。扩大范围需要当前第 5 版
告知,保留原有关闭选择。发布客户端前先备份 D1,依次应用 `0002-goal-usage.sql`、
`0003-goal-duration-sources.sql` 并部署 Worker。后者将旧计数转入 `host_call` /
`unknown`,保留旧表以便回滚,不影响心跳和 CLI 计数。尚未发布的 Goal v1 协议
diff --git a/examples/loopx-update-smoke.py b/examples/loopx-update-smoke.py
index 508b92cba3..a68428fea5 100644
--- a/examples/loopx-update-smoke.py
+++ b/examples/loopx-update-smoke.py
@@ -568,7 +568,8 @@ def test_cli_rollback_previous_with_temp_home() -> None:
"previous",
],
cwd=REPO_ROOT,
- env={"HOME": str(home), "PATH": f"{home / '.local' / 'bin'}:{os.environ.get('PATH', '')}"},
+ env={"HOME": str(home), "PATH": f"{home / '.local' / 'bin'}:{os.environ.get('PATH', '')}",
+ "LOOPX_USAGE_PING": "0"},
text=True,
capture_output=True,
)
diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py
index b7bda57ea1..8cd7aac20a 100644
--- a/loopx/capabilities/manager_context/roundtrip.py
+++ b/loopx/capabilities/manager_context/roundtrip.py
@@ -568,6 +568,10 @@ def _write_exact_return_state(
else:
result.pop("admission", None)
_write(context["state_path"], result)
+ if (result.get("status") == "delivered" and result.get("reply_verified") is True
+ and current.get("status") != "delivered"):
+ from ...usage_ping import observe_verified_return
+ observe_verified_return()
def _retry_state(state, now, *, error):
diff --git a/loopx/cli.py b/loopx/cli.py
index e3c3877266..c1b28c3094 100644
--- a/loopx/cli.py
+++ b/loopx/cli.py
@@ -397,6 +397,8 @@ def main(argv: list[str] | None = None) -> int:
pass # demo package absent in installed builds; no question rewrite
parser = build_parser()
args = parser.parse_args(raw_argv)
+ from .usage_ping import select_operation
+ select_operation(args)
args.format = resolve_global_output_format(args)
guard_result = enforce_native_controller_guard(args)
if guard_result is not None:
diff --git a/loopx/cli_commands/pr_review.py b/loopx/cli_commands/pr_review.py
index 0901dd2dca..c58fe40efb 100644
--- a/loopx/cli_commands/pr_review.py
+++ b/loopx/cli_commands/pr_review.py
@@ -593,6 +593,8 @@ def handle_pr_review_command(
)
payload["request"]["local_checkpoint_write_performed"] = True
except Exception as exc:
+ from ..usage_ping import capture_failure
+ capture_failure(exc)
error = str(exc)
if checkpoint_path is not None:
path_candidates = {str(checkpoint_path)}
diff --git a/loopx/cli_commands/todo.py b/loopx/cli_commands/todo.py
index f08c84cc1d..af217661f2 100644
--- a/loopx/cli_commands/todo.py
+++ b/loopx/cli_commands/todo.py
@@ -677,6 +677,8 @@ def handle_todo_command(
else:
raise ValueError("unsupported todo command")
except Exception as exc:
+ from ..usage_ping import capture_failure
+ capture_failure(exc)
payload = todo_error_payload(args, exc)
append_todo_rollout_event(
payload,
diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py
index 0cd7c3e24d..8e9b908103 100644
--- a/loopx/cli_commands/turn.py
+++ b/loopx/cli_commands/turn.py
@@ -1103,6 +1103,8 @@ def on_managed_start_admitted() -> None:
else:
raise ValueError("turn requires the `plan` or `run-once` subcommand")
except Exception as exc: # noqa: BLE001 - CLI boundary renders typed JSON failure
+ from ..usage_ping import capture_failure
+ capture_failure(exc)
journal_readback = None
if execution_started:
transaction = payload.get("transaction") or {}
diff --git a/loopx/cli_commands/usage_ping.py b/loopx/cli_commands/usage_ping.py
index 6296a2fc3b..f31d4aecaf 100644
--- a/loopx/cli_commands/usage_ping.py
+++ b/loopx/cli_commands/usage_ping.py
@@ -19,10 +19,15 @@ def render_usage_ping_markdown(payload: dict[str, object]) -> str:
f"- Sending eligible: {payload['sending']}; blocked by: {payload['blocked_by'] or 'none'}",
f"- Endpoint: {payload['endpoint'] or 'not configured'}",
f"- Last heartbeat: {payload['last_sent_day'] or 'never'}",
- str(payload['disclosure']), "", "Payload previews (aggregate sends after the UTC day closes):"]
+ str(payload['disclosure']), "", "Payload previews (first CLI result immediately; later activity at most every 15 minutes):"]
import json
lines.append(json.dumps({"heartbeat": payload.get("next_payload"),
- "aggregate": payload.get("aggregate_preview")}, indent=2))
+ "aggregate": payload.get("aggregate_preview"),
+ "diagnostics": payload.get("diagnostic_preview"),
+ "goals": payload.get("goal_preview"),
+ "diagnostic_dropped": payload.get("diagnostic_dropped", 0),
+ "identity_scope": payload.get("identity_scope"),
+ "delivery_history": payload.get("delivery_history", [])}, indent=2))
lines.append("Details: docs/reference/usage-ping.md")
return "\n".join(lines)
diff --git a/loopx/cli_runtime.py b/loopx/cli_runtime.py
index 6666cf37ba..fa0daf7ba8 100644
--- a/loopx/cli_runtime.py
+++ b/loopx/cli_runtime.py
@@ -87,6 +87,8 @@ def print_payload(
fmt: str,
markdown_renderer: Callable[[dict[str, object]], str],
) -> None:
+ from .usage_ping import capture_result
+ capture_result(payload)
if fmt == "json":
print(json.dumps(payload, ensure_ascii=False, indent=2))
else:
@@ -427,6 +429,8 @@ def dispatch_common_command(
def _dispatch_selected(args: argparse.Namespace, raw_argv: list[str]) -> int:
+ from .usage_ping import select_operation
+ select_operation(args)
args.format = resolve_global_output_format(args)
guard_result = enforce_native_controller_guard(args)
if guard_result is not None:
diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts
index f2c3ba1fbb..9126be5712 100644
--- a/loopx/control_plane/runtime/usage_statistics.ts
+++ b/loopx/control_plane/runtime/usage_statistics.ts
@@ -6,6 +6,8 @@ import type { JsonObject } from "../effect_program.ts";
import { withFileMutationLock, atomicWriteJson } from "../effect_runtime_io.ts";
import { AGGREGATE_SCHEMA, PING_SCHEMA, MAX_COUNT, MAX_ROWS, counterKey, object, validAggregate, validCounter, validId, validPing } from "./usage_statistics_contract.ts";
import type { Aggregate, Counter, Ping } from "./usage_statistics_contract.ts";
+import { DIAGNOSTIC_SCHEMA, diagnosticKey, validDiagnostic, validDiagnostics } from "./usage_statistics_diagnostics.ts";
+import type { Diagnostic, DiagnosticAggregate } from "./usage_statistics_diagnostics.ts";
import { recordGoalUsage, goalPreview } from "./usage_statistics_goals.ts";
import { validGoalAggregate } from "./usage_statistics_goal_contract.ts";
@@ -16,15 +18,19 @@ 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 = 4;
+export const NOTICE_VERSION = 5;
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 Delivery = { day: string; channel: "heartbeat" | "cli" | "goal"; rows: number; status: "accepted" | "rejected" | "unavailable" };
+const MAX_DELIVERIES = 20;
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[]; aggregate_last_attempt_ms?: number;
+ deliveries?: Delivery[];
+ diagnostics?: Diagnostic[]; diagnostic_dropped?: number;
};
export function endpoint(env: Env): string {
try {
@@ -79,13 +85,22 @@ async function load(path: string): Promise {
|| !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");
+ || (raw.counters !== undefined && (!Array.isArray(raw.counters) || raw.counters.length > MAX_ROWS || !raw.counters.every(validCounter)))
+ || (raw.diagnostic_dropped !== undefined && (!Number.isSafeInteger(raw.diagnostic_dropped) || Number(raw.diagnostic_dropped) < 0 || Number(raw.diagnostic_dropped) > MAX_COUNT))
+ || (raw.diagnostics !== undefined && (!Array.isArray(raw.diagnostics) || raw.diagnostics.length > 32 || !raw.diagnostics.every(validDiagnostic)))) throw new Error("usage_state_invalid");
return raw as State;
}
async function save(path: string, state: State) {
await atomicWriteJson(path, state as unknown as JsonObject);
await chmod(path, 0o600);
}
+function validDelivery(value: unknown): value is Delivery {
+ return object(value) && Object.keys(value).sort().join() === "channel,day,rows,status"
+ && typeof value.day === "string" && /^\d{4}-\d{2}-\d{2}$/.test(value.day)
+ && ["heartbeat", "cli", "goal"].includes(String(value.channel))
+ && ["accepted", "rejected", "unavailable"].includes(String(value.status))
+ && Number.isSafeInteger(value.rows) && Number(value.rows) >= 1 && Number(value.rows) <= MAX_ROWS;
+}
function day(ctx: Context): string { return (ctx.now ?? new Date()).toISOString().slice(0, 10); }
function ping(state: State, ctx: Context): Ping | null {
const value = { schema: PING_SCHEMA, install_id: state.install_id, version: ctx.version,
@@ -103,9 +118,13 @@ export async function inspect(path: string, ctx: Context) {
automatic_notice_required: automaticNoticeRequired(state, ctx), last_sent_day: state.last_sent_day ?? null,
next_payload: state.consent === "disabled" ? null : ping(state, ctx),
aggregate_preview: state.consent === "disabled" || !state.counters?.length ? null : { schema: AGGREGATE_SCHEMA, counters: state.counters },
+ diagnostic_preview: state.consent === "disabled" || !state.diagnostics?.length ? null : { schema: DIAGNOSTIC_SCHEMA, counters: state.diagnostics },
+ diagnostic_dropped: state.diagnostic_dropped ?? 0,
goal_preview: state.consent === "disabled" ? null : await goalPreview(path + ".goals", state.generation).catch(() => null),
+ identity_scope: "persistent_machine_state_directory_not_person_or_session",
+ delivery_history: state.consent === "disabled" ? [] : (state.deliveries ?? []).filter(validDelivery).slice(-MAX_DELIVERIES),
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. 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") };
+ 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/sub-operation/result/duration/error counts, release version, UTC activity day, voluntary deployment context and receipt-backed lifecycle signals are sent separately without an ID. Deployment context defaults to unknown and is never inferred. 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, argument values, Goal contents or raw errors. Local status keeps at most 20 content-free delivery summaries, cleared on disable. 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 () => {
@@ -122,6 +141,7 @@ export async function configure(path: string, ctx: Context, action: "enable" | "
if (action === "enable") state.consent = "enabled";
if (state.notice && !sameNotice(state, ctx)) {
state.counters = [];
+ state.diagnostics = [];
await rm(path + ".goals", { force: true });
await rm(path + ".cycles", { force: true });
state.generation = randomUUID();
@@ -134,19 +154,22 @@ export async function configure(path: string, ctx: Context, action: "enable" | "
}, 1000);
return inspect(path, ctx);
}
-export type Post = (url: string, payload: Ping | Aggregate | GoalAggregate) => Promise;
+export type Post = (url: string, payload: Ping | Aggregate | GoalAggregate | DiagnosticAggregate) => Promise;
const post: Post = async (url, payload) => (await fetch(url, {
method: "POST", headers: { "Content-Type": "application/json", "User-Agent": "loopx-usage-ping" },
body: JSON.stringify(payload), signal: AbortSignal.timeout(3000), redirect: "error",
})).status;
/** Called in a detached process with one allowlisted observation, never raw argv/output. */
-export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation, cycle?: CycleObservation) {
+export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation, cycle?: CycleObservation, diagnostic?: Diagnostic) {
if (counter !== null && (!validCounter(counter) || counter.count !== 1)) return { sent: false, reason: "invalid_observation" };
+ if (diagnostic && (!validDiagnostic(diagnostic) || diagnostic.count !== 1)) return { sent: false, reason: "invalid_observation" };
let heartbeat: Ping | null = null;
let heartbeatRequest: Promise | undefined;
let aggregate: Aggregate | null = null;
let aggregateRequest: Promise | undefined;
+ let diagnostics: DiagnosticAggregate | null = null;
+ let diagnosticRequest: Promise | undefined;
let goals: GoalAggregate | null = null;
const today = day(ctx);
const now = (ctx.now ?? new Date()).getTime();
@@ -171,16 +194,25 @@ export async function observe(path: string, ctx: Context, generation: string, co
if (row) row.count = Math.min(MAX_COUNT, row.count + 1);
else if (state.counters.length < MAX_ROWS) state.counters.push({ ...counter });
}
+ state.diagnostics = (state.diagnostics ?? []).filter(row => Date.parse(today) - Date.parse(row.activity_day) <= 7 * 86400000);
+ if (diagnostic) {
+ const row = state.diagnostics.find(entry => diagnosticKey(entry) === diagnosticKey(diagnostic));
+ if (row) row.count = Math.min(MAX_COUNT, row.count + 1);
+ else if (state.diagnostics.length < 32) state.diagnostics.push({ ...diagnostic });
+ else state.diagnostic_dropped = Math.min(MAX_COUNT, (state.diagnostic_dropped ?? 0) + 1);
+ }
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
}
// 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
+ if ((state.counters.length || state.diagnostics.length) && (state.aggregate_last_attempt_ms === undefined
|| now - state.aggregate_last_attempt_ms >= AGGREGATE_INTERVAL_MS)) {
- aggregate = { schema: AGGREGATE_SCHEMA, counters: state.counters };
+ if (state.counters.length) aggregate = { schema: AGGREGATE_SCHEMA, counters: state.counters };
+ if (state.diagnostics.length) diagnostics = { schema: DIAGNOSTIC_SCHEMA, counters: state.diagnostics };
state.counters = [];
+ state.diagnostics = [];
state.aggregate_last_attempt_ms = now; // consume before I/O; never retry a lossy batch
}
await save(path, state);
@@ -195,15 +227,23 @@ export async function observe(path: string, ctx: Context, generation: string, co
try { aggregateRequest = send(endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), aggregate).catch(() => 0); }
catch { aggregateRequest = Promise.resolve(0); }
}
+ if (diagnostics && validDiagnostics(diagnostics)) {
+ try { diagnosticRequest = send(endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), diagnostics).catch(() => 0); }
+ catch { diagnosticRequest = Promise.resolve(0); }
+ }
return true;
}, 0); // Never queue behind business or telemetry work.
if (!allowed) return { sent: false, reason: "blocked" };
let sent = false;
- for (const [url, payload] of [[endpoint(ctx.env), heartbeat], [endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), aggregate], [endpoint(ctx.env).replace(/\/ping$/, "/goals"), goals]] as const) {
+ const outgoing: Array[1] | null]> = [
+ [endpoint(ctx.env), heartbeat], [endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), aggregate],
+ [endpoint(ctx.env).replace(/\/ping$/, "/aggregate"), diagnostics], [endpoint(ctx.env).replace(/\/ping$/, "/goals"), goals],
+ ];
+ for (const [url, payload] of outgoing) {
if (!payload) continue;
- if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload))) continue;
+ if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload) || validDiagnostics(payload))) continue;
try {
- let request = url.endsWith("/ping") ? heartbeatRequest : url.endsWith("/aggregate") ? aggregateRequest : undefined;
+ let request = url.endsWith("/ping") ? heartbeatRequest : payload.schema === DIAGNOSTIC_SCHEMA ? diagnosticRequest : 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("/goals")) {
@@ -214,16 +254,19 @@ export async function observe(path: string, ctx: Context, generation: string, co
}
if (!request) break;
const code = await request;
- if (url.endsWith("/ping") && code >= 200 && code < 300) {
- sent = true;
- await withFileMutationLock(path, async () => {
- const latest = await load(path);
- if (latest.generation === generation && !blockedBy(latest, ctx)) {
- latest.last_sent_day = today;
- await save(path, latest);
- }
- }, 0);
- }
+ const accepted = code >= 200 && code < 300;
+ if (url.endsWith("/ping") && accepted) sent = true;
+ await withFileMutationLock(path, async () => {
+ const latest = await load(path);
+ if (latest.generation !== generation || blockedBy(latest, ctx)) return;
+ if (url.endsWith("/ping") && accepted) latest.last_sent_day = today;
+ const delivery: Delivery = { day: today,
+ channel: url.endsWith("/ping") ? "heartbeat" : url.endsWith("/aggregate") ? "cli" : "goal",
+ rows: "counters" in payload ? payload.counters.length : 1,
+ status: accepted ? "accepted" : code ? "rejected" : "unavailable" };
+ latest.deliveries = [...(latest.deliveries ?? []).filter(validDelivery), delivery].slice(-MAX_DELIVERIES);
+ await save(path, latest);
+ }, 0);
} catch { /* Lossy diagnostics must never affect work. No raw errors persisted. */ }
}
return { sent };
diff --git a/loopx/control_plane/runtime/usage_statistics_cli.ts b/loopx/control_plane/runtime/usage_statistics_cli.ts
index feff1703be..ef9356efdc 100644
--- a/loopx/control_plane/runtime/usage_statistics_cli.ts
+++ b/loopx/control_plane/runtime/usage_statistics_cli.ts
@@ -2,6 +2,7 @@ import { configure, inspect, observe } from "./usage_statistics.ts";
import { durationBucket, FEATURES, object } from "./usage_statistics_contract.ts";
import type { Context } from "./usage_statistics.ts";
import type { Counter } from "./usage_statistics_contract.ts";
+import { resultDiagnostic, usageContext } from "./usage_statistics_diagnostics.ts";
import type { GoalObservation } from "./usage_statistics_goal_contract.ts";
import { validGoalObservation, hostCategory } from "./usage_statistics_goal_contract.ts";
@@ -29,10 +30,18 @@ try {
result = await observe(request.path, ctx, String(request.generation), null, undefined, { ...request.observation, host: hostCategory(request.observation.host) } as GoalObservation);
} else if (request.action === "observe" && typeof request.feature === "string"
&& typeof request.elapsed_ms === "number" && Number.isFinite(request.elapsed_ms) && request.elapsed_ms >= 0) {
- result = await observe(request.path, ctx, String(request.generation), {
- feature: (FEATURES as readonly unknown[]).includes(request.feature) ? request.feature : "other", outcome: request.outcome, error: request.error,
- duration: durationBucket(request.elapsed_ms), count: 1,
- } as Counter);
+ const feature = (FEATURES as readonly unknown[]).includes(request.feature) ? request.feature : "other";
+ if (typeof request.exit_code === "number" && Number.isInteger(request.exit_code)) {
+ result = await observe(request.path, ctx, String(request.generation), null, undefined, undefined, undefined, {
+ feature: feature as Counter["feature"], ...resultDiagnostic(feature, request.operation, request.result_facts, request.exit_code, request.failure),
+ duration: durationBucket(request.elapsed_ms), count: 1, version: ctx.version,
+ activity_day: typeof request.activity_day === "string" ? request.activity_day : new Date().toISOString().slice(0, 10), context: usageContext(ctx.env.LOOPX_USAGE_CONTEXT),
+ });
+ } else {
+ result = await observe(request.path, ctx, String(request.generation), {
+ feature, outcome: request.outcome, error: request.error, duration: durationBucket(request.elapsed_ms), count: 1,
+ } as Counter);
+ }
} else throw new Error("usage_request_invalid");
process.stdout.write(JSON.stringify(result) + "\n");
} catch {
diff --git a/loopx/control_plane/runtime/usage_statistics_diagnostics.ts b/loopx/control_plane/runtime/usage_statistics_diagnostics.ts
new file mode 100644
index 0000000000..8ad0d841ee
--- /dev/null
+++ b/loopx/control_plane/runtime/usage_statistics_diagnostics.ts
@@ -0,0 +1,74 @@
+/** Versioned, content-free diagnostics. This is observation, never work authority. */
+import { counterKey, DURATIONS, FEATURES, MAX_COUNT, MAX_ROWS, object } from "./usage_statistics_contract.ts";
+import type { Counter } from "./usage_statistics_contract.ts";
+
+export const DIAGNOSTIC_SCHEMA = "loopx_usage_diagnostics_v1";
+export const CONTEXTS = ["unknown", "personal", "shared_service", "ephemeral", "organization_managed", "maintainer"] as const;
+export const DIAGNOSTIC_OPERATIONS = ["default", "plan", "run-once", "status", "should-run", "spend-slot", "monitor-poll", "list", "add", "claim", "update", "complete", "register", "resolve", "bind-session", "unbind-session", "merge-readiness", "check-result", "result-return"] as const;
+export const REASONS = ["none", "not_ready", "invalid_input", "permission", "not_found", "timeout", "connection", "interrupted", "command_failed"] as const;
+export const SIGNALS = ["none", "project_registered", "managed_turn_committed", "todo_completed", "todo_validated", "result_returned"] as const;
+const FEATURE_OPERATIONS: Record = {
+ turn: ["plan", "run-once", "status"], quota: ["status", "plan", "should-run", "spend-slot", "monitor-poll"],
+ todo: ["list", "add", "claim", "update", "complete"], project: ["register", "resolve", "bind-session", "unbind-session"],
+ "pr-review": ["merge-readiness", "check-result"], other: ["result-return"],
+};
+const SIGNAL_SOURCES: Record = {
+ project_registered: ["project", "register"], managed_turn_committed: ["turn", "run-once"],
+ todo_completed: ["todo", "complete"], todo_validated: ["todo", "complete"], result_returned: ["other", "result-return"],
+};
+export type Diagnostic = Omit & {
+ outcome: "ok" | "blocked" | "failed" | "cancelled"; error: typeof REASONS[number];
+ operation: typeof DIAGNOSTIC_OPERATIONS[number]; signal: typeof SIGNALS[number];
+ version: string; activity_day: string; context: typeof CONTEXTS[number];
+};
+export type DiagnosticAggregate = { schema: typeof DIAGNOSTIC_SCHEMA; counters: Diagnostic[] };
+export function usageContext(value: unknown): Diagnostic["context"] {
+ return (CONTEXTS as readonly unknown[]).includes(value) ? value as Diagnostic["context"] : "unknown";
+}
+export function diagnosticKey(value: Diagnostic): string {
+ return [counterKey(value as Counter), value.operation, value.signal, value.version, value.activity_day, value.context].join(":");
+}
+export function validDiagnostic(value: unknown): value is Diagnostic {
+ if (!object(value) || Object.keys(value).sort().join() !== "activity_day,context,count,duration,error,feature,operation,outcome,signal,version") return false;
+ if (!(FEATURES as readonly unknown[]).includes(value.feature) || !(DURATIONS as readonly unknown[]).includes(value.duration)
+ || !(DIAGNOSTIC_OPERATIONS as readonly unknown[]).includes(value.operation) || !(SIGNALS as readonly unknown[]).includes(value.signal)
+ || !(CONTEXTS as readonly unknown[]).includes(value.context) || !(REASONS as readonly unknown[]).includes(value.error)
+ || typeof value.version !== "string" || !/^\d{1,3}\.\d{1,3}\.\d{1,4}$/.test(value.version)
+ || typeof value.activity_day !== "string" || !/^\d{4}-\d{2}-\d{2}$/.test(value.activity_day)
+ || !Number.isFinite(Date.parse(value.activity_day)) || new Date(value.activity_day).toISOString().slice(0, 10) !== value.activity_day
+ || !Number.isInteger(value.count) || Number(value.count) < 1 || Number(value.count) > MAX_COUNT) return false;
+ if (value.signal !== "none" && value.outcome !== "ok") return false;
+ if (value.operation !== "default" && !FEATURE_OPERATIONS[String(value.feature)]?.includes(String(value.operation))) return false;
+ if (value.signal !== "none" && SIGNAL_SOURCES[String(value.signal)]?.join(":") !== `${value.feature}:${value.operation}`) return false;
+ if (value.outcome === "blocked" && (value.feature !== "pr-review" || value.operation !== "merge-readiness")) return false;
+ return value.outcome === "ok" ? value.error === "none"
+ : value.outcome === "blocked" ? value.error === "not_ready"
+ : value.outcome === "cancelled" ? value.error === "interrupted"
+ : value.outcome === "failed" && !["none", "not_ready", "interrupted"].includes(String(value.error));
+}
+export function validDiagnostics(value: unknown): value is DiagnosticAggregate {
+ return object(value) && Object.keys(value).sort().join() === "counters,schema" && value.schema === DIAGNOSTIC_SCHEMA
+ && Array.isArray(value.counters) && value.counters.length > 0 && value.counters.length <= MAX_ROWS
+ && value.counters.every(validDiagnostic) && new Set(value.counters.map(diagnosticKey)).size === value.counters.length;
+}
+/** Only fixed metadata crosses the Python/TS bridge. Never parse error prose. */
+export function resultDiagnostic(feature: string, operation: unknown, facts: unknown, code: number, failure: unknown): Pick {
+ const op = typeof operation === "string" && FEATURE_OPERATIONS[feature]?.includes(operation) ? operation as Diagnostic["operation"] : "default";
+ const f = object(facts) ? facts : {};
+ let outcome: Diagnostic["outcome"] = code === 0 ? "ok" : "failed";
+ let error: Diagnostic["error"] = code === 0 ? "none" : "command_failed";
+ if (failure === "interrupted") { outcome = "cancelled"; error = "interrupted"; }
+ else if (["invalid_input", "permission", "not_found", "timeout", "connection"].includes(String(failure))) {
+ outcome = "failed"; error = failure as Diagnostic["error"];
+ } else if (op === "merge-readiness" && f.ok === true && f.ready === false) {
+ outcome = "blocked"; error = "not_ready";
+ } else if (f.ok === false) { outcome = "failed"; error = "command_failed"; }
+ let signal: Diagnostic["signal"] = "none";
+ if (outcome === "ok" && f.ok === true && f.dry_run !== true) {
+ if (feature === "project" && op === "register" && f.changed === true) signal = "project_registered";
+ if (feature === "todo" && op === "complete" && f.completed === true && f.changed === true) signal = f.validation_passed === true ? "todo_validated" : "todo_completed";
+ if (feature === "turn" && op === "run-once" && f.turn_committed === true) signal = "managed_turn_committed";
+ if (op === "result-return" && f.reply_verified === true && f.changed === true) signal = "result_returned";
+ }
+ return { operation: op, outcome, error, signal };
+}
diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json
index c3aad45b53..e4dcc4cf8b 100644
--- a/loopx/semantics/project_registry_io_manifest_v1.json
+++ b/loopx/semantics/project_registry_io_manifest_v1.json
@@ -255,7 +255,7 @@
},
{
"site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1",
- "line": 872,
+ "line": 876,
"column": 22,
"kind": "codec_read",
"api": "load_project_registry",
@@ -495,7 +495,7 @@
},
{
"site": "loopx/cli.py::.main::codec_read:load_project_registry#1",
- "line": 807,
+ "line": 809,
"column": 17,
"kind": "codec_read",
"api": "load_project_registry",
@@ -903,7 +903,7 @@
},
{
"site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#6",
- "line": 696,
+ "line": 698,
"column": 13,
"kind": "codec_read",
"api": "load_registry",
@@ -911,7 +911,7 @@
},
{
"site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#7",
- "line": 738,
+ "line": 740,
"column": 38,
"kind": "codec_read",
"api": "load_registry",
diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py
index 30bd782e57..b4a31ab2d6 100644
--- a/loopx/usage_ping.py
+++ b/loopx/usage_ping.py
@@ -7,6 +7,7 @@
import subprocess
import sys
import time
+from contextvars import ContextVar
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
@@ -16,8 +17,56 @@
STATE_FILENAME = "usage-ping.json"
# Scheduling hint only; keep aligned with the TypeScript notice revision.
-_NOTICE_VERSION = 4
+_NOTICE_VERSION = 5
_ENTRY = Path(__file__).parent / "control_plane/runtime/usage_statistics_cli.ts"
+_observation: ContextVar[dict[str, Any] | None] = ContextVar("usage_observation", default=None)
+
+
+def select_operation(args: Any) -> None:
+ """Parser-owned operation names only; never scan argument values."""
+ state = _observation.get()
+ if state is None:
+ return
+ command = getattr(args, "command", "")
+ operation = getattr(args, f"{command.replace('-', '_')}_command", "default")
+ if command == "pr-review":
+ operation = ("merge-readiness" if getattr(args, "check_merge_readiness", None)
+ else "check-result" if getattr(args, "check_result", None) else "default")
+ state["operation"] = operation
+
+
+def capture_result(payload: dict[str, Any]) -> None:
+ """Project booleans from existing receipts, not output text or business content."""
+ state = _observation.get()
+ if state is None:
+ return
+ try:
+ facts = {key: payload[key] for key in ('ok', 'ready', 'completed', 'changed', 'dry_run')
+ if isinstance(payload.get(key), bool)}
+ if payload.get('schema_version') == 'loopx_turn_execution_v0':
+ from .control_plane.turn_driver import loopx_turn_execution_committed
+ facts['turn_committed'] = loopx_turn_execution_committed(payload)
+ validation = payload.get('validation_receipt')
+ if isinstance(validation, dict) and validation.get('passed') is True:
+ facts['validation_passed'] = True
+ state['result_facts'] = facts
+ except Exception:
+ state.pop('result_facts', None) # Optional diagnostics cannot break output.
+
+
+def capture_failure(error: BaseException) -> None:
+ state = _observation.get()
+ if state is None:
+ return
+ category = ('interrupted' if isinstance(error, KeyboardInterrupt)
+ else 'invalid_input' if isinstance(error, SystemExit) and error.code == 2
+ else 'timeout' if isinstance(error, (TimeoutError, subprocess.TimeoutExpired))
+ else 'connection' if isinstance(error, ConnectionError)
+ else 'permission' if isinstance(error, PermissionError)
+ else 'not_found' if isinstance(error, FileNotFoundError)
+ else 'invalid_input' if isinstance(error, ValueError)
+ else 'command_failed')
+ state['failure'] = category
def state_path(runtime_root: Path | None = None) -> Path:
@@ -62,6 +111,7 @@ def control(action: str, path: Path | None = None, **fields: Any) -> dict[str, A
def begin(command: str) -> tuple[str, float] | None:
"""Read a small host hint; do not start a synchronous Node process on warm commands."""
+ _observation.set({})
# Negative-only scheduling hints for common unattended environments. These
# cannot authorize collection; TS still checks every supported switch value.
if os.environ.get("LOOPX_USAGE_PING") == "0" or os.environ.get("DO_NOT_TRACK") == "1" or os.environ.get("CI") == "true":
@@ -109,23 +159,37 @@ def begin(command: str) -> tuple[str, float] | None:
def finish(ticket: tuple[str, float] | None, command: str, code: int, error: BaseException | None = None) -> None:
"""Detach bounded local observation; never read args, output or error text."""
+ if error is not None:
+ capture_failure(error)
+ observation = _observation.get() or {}
+ _observation.set(None)
if ticket is None:
return
try:
- outcome, category = ("ok", "none") if code == 0 else ("failed", "command_failed")
- if isinstance(error, KeyboardInterrupt):
- outcome, category = "cancelled", "interrupted"
- elif isinstance(error, TimeoutError):
- outcome, category = "failed", "timeout"
- elif isinstance(error, ConnectionError):
- outcome, category = "failed", "connection"
request = _request("observe", state_path(), generation=ticket[0], feature=command if len(command) <= 64 else "other",
- outcome=outcome, error=category, elapsed_ms=max(0, (time.monotonic() - ticket[1]) * 1000))
+ exit_code=code, activity_day=datetime.now(timezone.utc).date().isoformat(),
+ **observation, elapsed_ms=max(0, (time.monotonic() - ticket[1]) * 1000))
_detach(request)
except Exception:
pass # Telemetry cannot replace the command's result.
+def observe_verified_return() -> None:
+ """An existing provider verified and durably settled a new result return."""
+ try:
+ state = json.loads(state_path().read_text(encoding='utf-8'))
+ if state.get('consent') == 'disabled' or (state.get('notice') or {}).get('version') != _NOTICE_VERSION:
+ return
+ generation = state.get('generation')
+ if not isinstance(generation, str) or not generation:
+ return
+ _detach(_request('observe', state_path(), generation=generation, feature='other', operation='result-return',
+ exit_code=0, activity_day=datetime.now(timezone.utc).date().isoformat(),
+ result_facts={'ok': True, 'changed': True, 'reply_verified': True}, elapsed_ms=0))
+ except Exception:
+ pass
+
+
def _detach(request: dict[str, Any], *, command: list[str] | None = None) -> None:
kwargs: dict[str, Any] = {"stdin": subprocess.PIPE, "stdout": subprocess.DEVNULL,
"stderr": subprocess.DEVNULL, "close_fds": True}
diff --git a/tests/canary/test_telemetry_isolation.py b/tests/canary/test_telemetry_isolation.py
index 9bfe26090a..56c368eb69 100644
--- a/tests/canary/test_telemetry_isolation.py
+++ b/tests/canary/test_telemetry_isolation.py
@@ -1,5 +1,8 @@
"""The smoke runner suppresses telemetry in actual child processes."""
import json
+import runpy
+import subprocess
+from pathlib import Path
from loopx.canary import runner
@@ -20,3 +23,19 @@ def test_smoke_subprocess_overrides_parent_telemetry_enable(tmp_path, monkeypatc
assert json.loads(result["stdout_tail"]) == {
"LOOPX_USAGE_PING": "0", "CI": None, "SYNTHETIC_VALUE": "preserved",
}
+
+
+def test_update_smoke_minimal_environment_keeps_opt_out(monkeypatch):
+ smoke = runpy.run_path(str(Path(__file__).parents[2] / 'examples/loopx-update-smoke.py'))
+ original_run = subprocess.run
+ environments = []
+
+ def actual_run(*args, **kwargs):
+ if 'env' in kwargs:
+ environments.append(kwargs['env'])
+ return original_run(*args, **kwargs)
+
+ monkeypatch.setattr(subprocess, 'run', actual_run)
+ monkeypatch.setenv('LOOPX_USAGE_PING', '1')
+ smoke['test_cli_rollback_previous_with_temp_home']()
+ assert environments and all(env.get('LOOPX_USAGE_PING') == '0' for env in environments)
diff --git a/tests/control_plane_ts/usage_statistics.test.ts b/tests/control_plane_ts/usage_statistics.test.ts
index 13121a55d4..3361f4847b 100644
--- a/tests/control_plane_ts/usage_statistics.test.ts
+++ b/tests/control_plane_ts/usage_statistics.test.ts
@@ -161,6 +161,36 @@ test("network failure is lossy and no-retry; no exception text enters local stat
assert.equal((await observe(path, ctx, generation, row, async () => { throw new Error("SECRET:/private/path"); })).sent, false);
await observe(path, ctx, generation, row, noPost);
assert.ok(!(await readFile(path, "utf8")).includes("SECRET"));
+ assert.deepEqual((await inspect(path, ctx)).delivery_history.map(row => row.status), ["unavailable", "unavailable"]);
+});
+
+test("local delivery history is bounded, content-free and cleared by disable", async t => {
+ const { path, state } = await fixture(t); const ctx = context();
+ await configure(path, ctx, "enable"); const generation = (await state()).generation;
+ for (let i = 0; i < 24; i++) {
+ await observe(path, { ...ctx, now: new Date(ctx.now!.getTime() + i * 15 * 60000) }, generation, row, async () => 400);
+ }
+ const status = await inspect(path, ctx);
+ assert.equal(status.delivery_history.length, 20);
+ assert.equal(status.identity_scope, "persistent_machine_state_directory_not_person_or_session");
+ for (const entry of status.delivery_history) {
+ assert.deepEqual(Object.keys(entry).sort(), ["channel", "day", "rows", "status"]);
+ assert.equal(entry.status, "rejected");
+ }
+ await configure(path, ctx, "disable");
+ assert.deepEqual((await inspect(path, ctx)).delivery_history, []);
+});
+
+test("scope disclosure upgrades preserve installation ID; only explicit disable resets it", async t => {
+ const { path, state } = await fixture(t); const ctx = context();
+ await configure(path, ctx, "enable"); const original = await state();
+ await writeFile(path, JSON.stringify({ ...original, notice: { ...original.notice, version: 3 } }));
+ const pending = await inspect(path, ctx);
+ assert.equal(pending.sending, false);
+ assert.equal(pending.automatic_notice_required, true);
+ await configure(path, ctx, "acknowledge", pending.notice);
+ assert.equal((await state()).install_id, original.install_id);
+ assert.notEqual((await state()).generation, original.generation);
});
test("malformed state fails closed; disable is the explicit repair", async t => {
diff --git a/tests/control_plane_ts/usage_statistics_delivery.test.ts b/tests/control_plane_ts/usage_statistics_delivery.test.ts
index a21c0f6dc4..85bc86e9d1 100644
--- a/tests/control_plane_ts/usage_statistics_delivery.test.ts
+++ b/tests/control_plane_ts/usage_statistics_delivery.test.ts
@@ -140,7 +140,7 @@ test("v3 upgrade waits for renewed notice, preserves pre-ack state and fences ol
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.notice.version, 5);
assert.equal(status.blocked_by, "notice_required");
assert.equal(status.automatic_notice_required, true);
let attempts = 0;
diff --git a/tests/control_plane_ts/usage_statistics_diagnostics.test.ts b/tests/control_plane_ts/usage_statistics_diagnostics.test.ts
new file mode 100644
index 0000000000..1f65fbd4e8
--- /dev/null
+++ b/tests/control_plane_ts/usage_statistics_diagnostics.test.ts
@@ -0,0 +1,83 @@
+import assert from "node:assert/strict";
+import test from "node:test";
+import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import { join } from "node:path";
+import { configure, inspect, observe } from "../../loopx/control_plane/runtime/usage_statistics.ts";
+import { DIAGNOSTIC_SCHEMA, resultDiagnostic, usageContext, validDiagnostics } from "../../loopx/control_plane/runtime/usage_statistics_diagnostics.ts";
+import type { Diagnostic } from "../../loopx/control_plane/runtime/usage_statistics_diagnostics.ts";
+
+const row: Diagnostic = { feature: "pr-review", operation: "merge-readiness", outcome: "blocked", error: "not_ready", duration: "lt_1s", count: 1,
+ version: "1.2.3", activity_day: "2026-09-30", context: "unknown", signal: "none" };
+test("readiness holds are not failures; errors and unknown operations remain bounded", () => {
+ assert.deepEqual(resultDiagnostic("pr-review", "merge-readiness", { ok: true, ready: false }, 1, undefined),
+ { operation: "merge-readiness", outcome: "blocked", error: "not_ready", signal: "none" });
+ assert.equal(resultDiagnostic("turn", "run-once", { ok: false }, 1, "invalid_input").error, "invalid_input");
+ assert.equal(resultDiagnostic("turn", "secret/project/path", { error: "secret" }, 1, "secret").operation, "default");
+ assert.equal(usageContext("organization-name"), "unknown");
+});
+test("signals require committed changed evidence, never mere exit success, prose or replay", () => {
+ assert.equal(resultDiagnostic("todo", "complete", { ok: true }, 0, undefined).signal, "none");
+ assert.equal(resultDiagnostic("todo", "complete", { ok: true, completed: true, changed: false, validation_passed: true }, 0, undefined).signal, "none");
+ assert.equal(resultDiagnostic("todo", "complete", { ok: true, completed: true, changed: true, validation_passed: true }, 0, undefined).signal, "todo_validated");
+ assert.equal(resultDiagnostic("turn", "run-once", { ok: true, turn_committed: true }, 0, undefined).signal, "managed_turn_committed");
+ assert.equal(resultDiagnostic("other", "result-return", { ok: true, changed: true, reply_verified: false }, 0, undefined).signal, "none");
+ assert.equal(resultDiagnostic("other", "result-return", { ok: true, changed: true, reply_verified: true }, 0, undefined).signal, "result_returned");
+});
+test("expanded contract rejects identity, raw errors, arbitrary fields and illegal result combinations", () => {
+ const payload = { schema: DIAGNOSTIC_SCHEMA, counters: [row] };
+ assert.ok(validDiagnostics(payload));
+ for (const field of ["install_id", "goal_id", "host", "ip", "timestamp", "arguments"]) {
+ assert.equal(validDiagnostics({ ...payload, [field]: "private" }), false);
+ assert.equal(validDiagnostics({ ...payload, counters: [{ ...row, [field]: "private" }] }), false);
+ }
+ for (const changes of [{ error: "raw text" }, { context: "company" }, { activity_day: "2026-02-31" }, { outcome: "ok" }, { signal: "todo_validated" }, { operation: "complete" }, { feature: "turn" }]) {
+ assert.equal(validDiagnostics({ ...payload, counters: [{ ...row, ...changes }] }), false);
+ }
+});
+test("new diagnostics flush once, preserve activity date, renew v4 notice and obey opt-out", async t => {
+ const root = await mkdtemp(join(tmpdir(), "loopx-diagnostics-")); t.after(() => rm(root, { recursive: true, force: true }));
+ const path = join(root, "state.json");
+ const ctx = { env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:1/v1/ping" }, version: "1.2.3", python: "3.13", channel: "source", now: new Date("2026-10-01T00:02:00Z") };
+ await configure(path, ctx, "enable");
+ const old = JSON.parse(await readFile(path, "utf8")); old.notice.version = 4;
+ old.diagnostics = [row]; await writeFile(path, JSON.stringify(old));
+ await observe(path, ctx, old.generation, null, async () => { assert.fail("old disclosure"); }, undefined, undefined, row);
+ assert.equal((await inspect(path, ctx)).blocked_by, "notice_required");
+ await configure(path, ctx, "acknowledge", (await inspect(path, ctx)).notice);
+ const state = JSON.parse(await readFile(path, "utf8")); const sent: unknown[] = [];
+ assert.equal(state.install_id, old.install_id); assert.notEqual(state.generation, old.generation);
+ assert.deepEqual(state.diagnostics, []);
+ await observe(path, ctx, state.generation, null, async (_url, value) => { sent.push(value); return 204; }, undefined, undefined, row);
+ assert.deepEqual(sent[1], { schema: DIAGNOSTIC_SCHEMA, counters: [row] });
+ await observe(path, ctx, state.generation, null, async () => { assert.fail("interval must hold"); }, undefined, undefined, row);
+ assert.equal((await inspect(path, ctx)).diagnostic_preview?.counters[0].count, 1);
+ await configure(path, ctx, "disable");
+ await observe(path, ctx, state.generation, null, async () => { assert.fail("disabled"); }, undefined, undefined, row);
+ assert.equal((await inspect(path, ctx)).diagnostic_preview, null);
+});
+
+test("diagnostic diversity is bounded and stale activity cannot reach a later batch", async t => {
+ const root = await mkdtemp(join(tmpdir(), "loopx-diagnostic-capacity-"));
+ t.after(() => rm(root, { recursive: true, force: true }));
+ const path = join(root, "state.json");
+ const ctx = { env: {}, version: "1.2.3", python: "3.13", channel: "source", now: new Date("2026-09-30T12:00:00Z") };
+ await configure(path, ctx, "enable");
+ const state = JSON.parse(await readFile(path, "utf8"));
+ const sent: unknown[] = [];
+ const post = async (_url: string, value: unknown) => { sent.push(value); return 204; };
+ await observe(path, ctx, state.generation, null, post, undefined, undefined, row);
+ sent.length = 0;
+ for (let n = 0; n < 34; n++) {
+ await observe(path, ctx, state.generation, null, post, undefined, undefined, { ...row, version: `1.2.${100 + n}` });
+ }
+ const full = await inspect(path, ctx);
+ assert.equal(sent.length, 0, "diversity must not bypass the delivery interval");
+ assert.equal(full.diagnostic_preview?.counters.length, 32);
+ assert.equal(full.diagnostic_dropped, 2);
+ const later = { ...ctx, now: new Date("2026-10-08T12:00:00Z") };
+ const fresh = { ...row, activity_day: "2026-10-08" };
+ await observe(path, later, state.generation, null, post, undefined, undefined, fresh);
+ assert.deepEqual(sent[1], { schema: DIAGNOSTIC_SCHEMA, counters: [fresh] });
+ assert.equal((await inspect(path, later)).diagnostic_preview, null);
+});
diff --git a/tests/control_plane_ts/usage_statistics_sources.test.ts b/tests/control_plane_ts/usage_statistics_sources.test.ts
index 3c12ed0110..87fa45b618 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,4);
+ const enabled=await configure(path,ctx,"enable"); assert.equal(enabled.notice.version,5);
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_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py
index 4f38b75557..979ea0c61c 100644
--- a/tests/test_collaboration_goal_instance.py
+++ b/tests/test_collaboration_goal_instance.py
@@ -524,8 +524,10 @@ def __call__(self, *_args):
("initial_delivery_receipt_unavailable", True),
])
def test_exact_external_return_verifies_after_recreation_without_resend(
- tmp_path: Path, verification_failure: str | None, terminal: bool,
+ tmp_path: Path, verification_failure: str | None, terminal: bool, monkeypatch,
) -> None:
+ observations = []
+ monkeypatch.setattr('loopx.usage_ping.observe_verified_return', lambda: observations.append('verified'))
registry = _create_source_registry(tmp_path)
store, _, receipt = _external_manager_request(tmp_path, registry)
admitted_at = datetime(2026, 9, 28, tzinfo=timezone.utc)
@@ -589,6 +591,7 @@ def verify(self, *_args):
assert first["status"] == "admitted"
assert first["attempt"]["message_ref"] == "om_exact_reply"
assert first["goal_ref"] == receipt["goal_ref"]
+ assert observations == [] # An unverified provider attempt is not a return.
_recreate(registry)
assert (
@@ -625,6 +628,7 @@ def verify(self, *_args):
assert transport.verify_calls == (1 if terminal else 2)
if terminal:
assert json.loads(state_path.read_text(encoding="utf-8")) == failed
+ assert observations == []
return
state = json.loads(
(
@@ -637,6 +641,8 @@ def verify(self, *_args):
assert state["status"] == "delivered"
assert state["goal_ref"] == receipt["goal_ref"]
assert state["verification"] == "reconciled_after_restart"
+ drain(tmp_path, registry, ChatSessionStore(tmp_path), transport, now=admitted_at + timedelta(days=2))
+ assert observations == ['verified']
def test_long_lived_mcp_keeps_its_captured_instance_after_recreation(
diff --git a/tests/test_manager_context_roundtrip.py b/tests/test_manager_context_roundtrip.py
index 7882587f7c..ee65dccfaa 100644
--- a/tests/test_manager_context_roundtrip.py
+++ b/tests/test_manager_context_roundtrip.py
@@ -760,7 +760,9 @@ def test_public_delivery_projection_normalizes_unknown_private_state(flow):
assert "private" not in str(state)
-def test_background_service_delivers_without_another_agent_or_query(flow):
+def test_background_service_delivers_without_another_agent_or_query(flow, monkeypatch):
+ observations = []
+ monkeypatch.setattr('loopx.usage_ping.observe_verified_return', lambda: observations.append('verified'))
root, registry, store, create = flow
_, _, receipt = create(True)
rid = receipt["request_id"]
@@ -791,6 +793,8 @@ def transport(*_):
assert state["status"] == "delivered"
assert state["provider_receipt"] == "sha256:provider-proof"
assert not service.thread.is_alive()
+ drain(root, registry, store, transport)
+ assert observations == [] # Legacy delivery is outside the exact-source telemetry contract.
@pytest.mark.parametrize("project", [False, True], ids=["steward", "project"])
diff --git a/tests/test_usage_ping.py b/tests/test_usage_ping.py
index 63c0633f6b..849a04b537 100644
--- a/tests/test_usage_ping.py
+++ b/tests/test_usage_ping.py
@@ -143,7 +143,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'] == 4
+ assert json.loads(usage_ping.state_path().read_text())['notice']['version'] == 5
@pytest.mark.parametrize('setting,value', [
@@ -261,14 +261,15 @@ def test_real_cli_first_result_reaches_http_without_next_day_return(isolated, co
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):
+ while time.monotonic() < deadline and not any(p['schema'] == 'loopx_usage_diagnostics_v1' for p in received):
time.sleep(0.02)
- aggregates = [p for p in received if p['schema'] == 'loopx_usage_aggregate_v1']
+ aggregates = [p for p in received if p['schema'] == 'loopx_usage_diagnostics_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 counters[0]['version'] and counters[0]['context'] == 'unknown'
assert usage_ping.control('status')['aggregate_preview'] is None
usage_ping.control('disable')
@@ -281,6 +282,26 @@ def test_business_failure_and_usage_failure_do_not_replace_original_result(isola
assert main(['version']) == 7
+@pytest.mark.parametrize('switch,value', [
+ ('CI', 'true'), ('CI', '1'), ('LOOPX_USAGE_PING', '0'),
+ ('LOOPX_USAGE_PING', 'off'), ('DO_NOT_TRACK', '1'),
+])
+def test_disabled_synthetic_real_cli_never_contacts_collector(isolated, collector, switch, value):
+ endpoint, received, accepted, release = collector
+ release.set()
+ env = {**os.environ, 'LOOPX_USAGE_PING_ENDPOINT': endpoint, switch: value}
+ 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)]
+ # A matching acknowledged state cannot override an environment suppressor.
+ usage_ping.control('enable')
+ for _ in range(2):
+ result = subprocess.run(command, env=env, capture_output=True, text=True, timeout=20)
+ assert result.returncode == 0 and json.loads(result.stdout)['ok'] is True
+ assert not accepted.wait(0.5)
+ assert received == []
+
+
@pytest.mark.parametrize('bypass_proxy', [False, True])
def test_detached_sender_honors_proxy_and_no_proxy(isolated, collector, monkeypatch, bypass_proxy):
endpoint, received, accepted, release = collector
@@ -327,7 +348,7 @@ def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated, u
assert usage_ping.state_path().read_bytes() == before
else:
assert not usage_ping.state_path().exists()
- assert initial['notice']['version'] == 4
+ assert initial['notice']['version'] == 5
connection.request('POST', path, json.dumps({'notice': initial['notice']}), {'Content-Type': 'application/json'})
response = connection.getresponse()
acknowledged = json.loads(response.read())
@@ -403,7 +424,7 @@ def test_v3_cli_upgrade_requires_visible_renewal_before_real_http(isolated, coll
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['notice']['version'] == 5
assert current['generation'] != old['generation']
assert current['counters'] == []
assert not accepted.wait(0.3) and received == []
@@ -411,7 +432,53 @@ def test_v3_cli_upgrade_requires_visible_renewal_before_real_http(isolated, coll
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:
+ while not any(item.get('schema') == 'loopx_usage_diagnostics_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']
+ batches = [item for item in received if item.get('schema') == 'loopx_usage_diagnostics_v1']
assert len(batches) == 1 and sum(row['count'] for row in batches[0]['counters']) == 1
+
+
+def test_actual_pr_readiness_cli_keeps_exit_code_and_projects_only_typed_facts(isolated, monkeypatch, capsys):
+ import runpy
+ from pathlib import Path
+ fixtures = runpy.run_path(str(Path(__file__).with_name('test_pr_review_github_scan.py')))
+ pr = fixtures['_merge_ready_pr']()
+ pr['statusCheckRollup'] = [{'name': 'test', 'status': 'IN_PROGRESS'}]
+ pr['review_thread_summary'] = {'complete': True, 'total_count': 0, 'unresolved_count': 0}
+ fixture = isolated / 'prs.json'
+ fixture.write_text(json.dumps({'repository': 'owner/repo', 'pull_requests': [pr]}))
+ registry = isolated / 'registry.json'
+ registry.write_text(json.dumps({'goals': [{'id': 'review-goal', 'repo': str(isolated), 'status': 'active'}]}))
+ usage_ping.control('enable')
+ requests = []
+ monkeypatch.setattr(usage_ping, '_detach', requests.append)
+ code = main(['--registry', str(registry), 'pr-review', '--goal-id', 'review-goal', '--fixture', str(fixture), '--check-merge-readiness', '4110@' + 'a' * 40, '--format', 'json'])
+ payload = json.loads(capsys.readouterr().out)
+ assert code == 1 and payload['ok'] and not payload['ready'], payload
+ observation = next(row for row in requests if row['action'] == 'observe')
+ assert observation['operation'] == 'merge-readiness'
+ assert observation['result_facts'] == {'ok': True, 'ready': False}
+ assert 'owner/repo' not in json.dumps(observation) and str(fixture) not in json.dumps(observation)
+
+
+def test_turn_capture_reuses_committed_current_effects_not_historical_receipt(isolated):
+ usage_ping.begin('turn')
+ committed = {'schema_version': 'loopx_turn_execution_v0', 'ok': True,
+ 'status': 'committed', 'receipt': {'status': 'committed'},
+ 'validation': {'status': 'passed'}, 'effects': {'state_written': True, 'quota_spent': True}}
+ usage_ping.capture_result(committed)
+ assert usage_ping._observation.get()['result_facts']['turn_committed'] is True
+ usage_ping.capture_result({**committed, 'replayed': True, 'effects': {'state_written': False, 'quota_spent': False}})
+ assert usage_ping._observation.get()['result_facts']['turn_committed'] is False
+ usage_ping.finish(None, 'turn', 0)
+
+
+def test_projection_failure_cannot_replace_business_output(isolated, monkeypatch, capsys):
+ from loopx.cli_runtime import print_payload
+ usage_ping.begin('turn')
+ monkeypatch.setattr('loopx.control_plane.turn_driver.loopx_turn_execution_committed', lambda _: (_ for _ in ()).throw(ValueError('private')))
+ payload = {'schema_version': 'loopx_turn_execution_v0', 'ok': True}
+ print_payload(payload, 'json', lambda _: 'unused')
+ assert json.loads(capsys.readouterr().out) == payload
+ assert 'result_facts' not in usage_ping._observation.get()
+ usage_ping.finish(None, 'turn', 0)