From 506e362fc7872314ed14e73e8956b8fb3822f71c Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:35:31 +0800
Subject: [PATCH] feat(telemetry): measure observed Goal execution duration
without identities
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
apps/presentation/dashboard/src/data/chat.ts | 2 +-
.../usage-statistics-settings.tsx | 3 +-
apps/usage-collector/README.md | 10 +-
.../migrations/0002-goal-usage.sql | 4 +
apps/usage-collector/schema.sql | 5 +
apps/usage-collector/src/basic-usage.ts | 18 +++
apps/usage-collector/src/collector.js | 16 ++-
apps/usage-collector/test/collector.test.mjs | 29 ++++
docs/reference/usage-ping.md | 62 +++++++-
docs/reference/usage-ping.zh-CN.md | 47 +++++-
loopx/chat_runtime.py | 16 ++-
.../control_plane/runtime/usage_statistics.ts | 31 ++--
.../runtime/usage_statistics_cli.ts | 4 +
.../runtime/usage_statistics_goal_contract.ts | 34 +++++
.../runtime/usage_statistics_goals.ts | 92 ++++++++++++
loopx/control_plane/turn_driver/executor.py | 24 ++--
loopx/usage_goal.py | 76 ++++++++++
loopx/usage_ping.py | 2 +-
.../usage_statistics_goals.test.ts | 134 ++++++++++++++++++
tests/test_loopx_turn_executor.py | 36 +++++
tests/test_usage_goal.py | 59 ++++++++
tests/test_usage_ping.py | 15 ++
tsconfig.control-plane.json | 1 +
23 files changed, 683 insertions(+), 37 deletions(-)
create mode 100644 apps/usage-collector/migrations/0002-goal-usage.sql
create mode 100644 loopx/control_plane/runtime/usage_statistics_goal_contract.ts
create mode 100644 loopx/control_plane/runtime/usage_statistics_goals.ts
create mode 100644 loopx/usage_goal.py
create mode 100644 tests/control_plane_ts/usage_statistics_goals.test.ts
create mode 100644 tests/test_usage_goal.py
diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts
index 470383e611..e8937682dd 100644
--- a/apps/presentation/dashboard/src/data/chat.ts
+++ b/apps/presentation/dashboard/src/data/chat.ts
@@ -2070,7 +2070,7 @@ const usageStatisticsSchema = z.object({
consent: z.enum(["default", "enabled", "disabled"]),
sending: z.boolean(), blocked_by: z.string().nullable(), endpoint: z.string().nullable(),
policy: z.string(), notice_required: z.boolean(),
- next_payload: z.unknown(), aggregate_preview: z.unknown(),
+ next_payload: z.unknown(), aggregate_preview: z.unknown(), goal_preview: z.unknown(),
});
export type UsageStatistics = z.infer;
export async function usageStatistics(enabled?: boolean): Promise {
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 ff78b9517a..dde07c2004 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
@@ -26,6 +26,7 @@ export function UsageStatisticsSettings() {
? "用于决定平台支持和改进命令体验。每天向 LoopX 的 Cloudflare 收集服务发送随机安装标识、版本、系统、CPU 架构、Python 版本和安装渠道;固定的 CLI 功能、结果、耗时区间和错误类别在本机按天汇总后另行发送,不带安装标识。"
: "Helps prioritize platform support and CLI improvements. A daily heartbeat sends a random installation ID, version, OS, CPU architecture, Python version and install channel to the LoopX Cloudflare collector. Fixed CLI feature, result, duration and error counts are aggregated locally by day and sent separately without the ID."}
{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 ? "另按天汇总受管 Turn 与普通 Goal 对话的执行跨度和 Host 调用累计时长区间(并行重叠只计一次),包括未结束的 Goal,不带 Goal 或安装标识。从本机开始观测后计量,可能漏计;不代表 Goal 完成、CPU 用时或计费时长。原生 /goal 和外部附加任务暂不覆盖。" : "Also aggregates observed Goal execution-span and Host-call duration buckets by day, including unfinished Goals; overlaps count once. No Goal or installation ID. Measurement begins locally and may undercount; it is not completion evidence, CPU time or billing. Covers managed Turns and regular owner Goal chat; native /goal and externally attached tasks are excluded."}
{state ? <>
void update(event.target.checked)} /> {zh ? "允许基础使用统计(整台机器)" : "Allow basic usage statistics (this machine)"}
@@ -34,7 +35,7 @@ export function UsageStatisticsSettings() {
: (zh ? `当前不发送:${({disabled:"已关闭",CI:"CI 环境",DO_NOT_TRACK:"请勿追踪开关",LOOPX_USAGE_PING:"环境变量已关闭",consent_required:"需要明确同意",invalid_policy:"策略配置无效",invalid_endpoint:"接收地址无效",notice_required:"需要重新告知"} as Record)[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 }, null, 2)}
+ {zh ? "查看待发送数据" : "Preview outgoing data"} {JSON.stringify({ heartbeat: state.next_payload, aggregate: state.aggregate_preview, goals: state.goal_preview }, 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 b9f33645ba..fb26213807 100644
--- a/apps/usage-collector/README.md
+++ b/apps/usage-collector/README.md
@@ -8,6 +8,8 @@ The TypeScript client/collector allowlist lives in
|---|---|
| `POST /v1/ping` | Daily random-ID heartbeat with version/OS/CPU/Python/channel; ≤1 KiB |
| `POST /v1/aggregate` | Fixed CLI counts, no installation ID or join key; ≤16 KiB |
+| `POST /v1/goals` | Observed Goal-day span/execution buckets, no identity; ≤16 KiB |
+| `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 |
| `POST /v0/ping` | Retained six-field opt-in client contract; no new default-on clients use this route |
@@ -47,7 +49,11 @@ npx wrangler d1 migrations apply loopx-usage --remote
npx wrangler deploy
```
-Qualify `/v1/ping`, `/v1/aggregate`, both stats endpoints, and invalid-field/size
+Existing v1 installations also apply `0002-goal-usage.sql` before deploying
+the Goal-duration Worker. Back up first; this only adds `goal_usage_counts`.
+A Worker rollback can leave that additive table intact.
+
+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
rollback can restore the prior Worker without dropping the additive columns
@@ -69,4 +75,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 the v0 migration; no production telemetry is needed for these tests.
+including both additive migrations; no production telemetry is needed for these tests.
diff --git a/apps/usage-collector/migrations/0002-goal-usage.sql b/apps/usage-collector/migrations/0002-goal-usage.sql
new file mode 100644
index 0000000000..aeaaa25d82
--- /dev/null
+++ b/apps/usage-collector/migrations/0002-goal-usage.sql
@@ -0,0 +1,4 @@
+CREATE TABLE IF NOT EXISTS goal_usage_counts (
+ day TEXT NOT NULL, span TEXT NOT NULL, execution TEXT NOT NULL, count INTEGER NOT NULL,
+ PRIMARY KEY (day, span, execution)
+);
diff --git a/apps/usage-collector/schema.sql b/apps/usage-collector/schema.sql
index 9c2ec6d78e..5999f44bcb 100644
--- a/apps/usage-collector/schema.sql
+++ b/apps/usage-collector/schema.sql
@@ -28,3 +28,8 @@ CREATE TABLE IF NOT EXISTS usage_counts (
count INTEGER NOT NULL,
PRIMARY KEY (day, feature, outcome, duration, error)
);
+
+CREATE TABLE IF NOT EXISTS goal_usage_counts (
+ day TEXT NOT NULL, span TEXT NOT NULL, execution TEXT NOT NULL, count INTEGER NOT NULL,
+ PRIMARY KEY (day, span, execution)
+);
diff --git a/apps/usage-collector/src/basic-usage.ts b/apps/usage-collector/src/basic-usage.ts
index 8d9453fffb..9fc079c7f1 100644
--- a/apps/usage-collector/src/basic-usage.ts
+++ b/apps/usage-collector/src/basic-usage.ts
@@ -23,3 +23,21 @@ export async function aggregateStats(db: Database, since: string) {
}
return { schema: "loopx_usage_aggregate_stats_v1", definition: "Lossy CLI invocation counts received in the last 30 UTC days; not people, installations or accepted Goal outcomes. Cells below 5 omitted.", totals };
}
+
+export { validGoalAggregate } from "../../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts";
+import type { GoalAggregate } from "../../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts";
+export async function recordGoals(db: Database, value: GoalAggregate, day: string) {
+ await db.batch(value.counters.map(row => db.prepare(
+ "INSERT INTO goal_usage_counts (day, span, execution, count) VALUES (?1, ?2, ?3, ?4) " +
+ "ON CONFLICT (day, span, execution) DO UPDATE SET count = count + excluded.count",
+ ).bind(day, row.span, row.execution, row.count)));
+}
+export async function goalStats(db: Database, since: string) {
+ const totals: Record> = {};
+ for (const column of ["span", "execution"]) {
+ const rows = await db.prepare(`SELECT ${column} AS key, SUM(count) AS n FROM goal_usage_counts WHERE 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_goal_usage_stats_v1", definition: "Lossy observed Goal-day snapshots received in the last 30 UTC days, including unfinished Goals. Span is first-to-last observed execution; execution is union of Host-call intervals, not CPU time. Lower bounds since measurement began; managed Turns and regular owner Goal chat only. Not unique Goals, completion evidence or billing. Cells below 5 omitted.", totals };
+}
diff --git a/apps/usage-collector/src/collector.js b/apps/usage-collector/src/collector.js
index fe64d10a70..a880d51e03 100644
--- a/apps/usage-collector/src/collector.js
+++ b/apps/usage-collector/src/collector.js
@@ -1,4 +1,4 @@
-import { validAggregate, validPing, recordAggregate, aggregateStats } from "./basic-usage.ts";
+import { validAggregate, validPing, recordAggregate, aggregateStats, validGoalAggregate, recordGoals, goalStats } 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.
@@ -135,6 +135,7 @@ export async function purge(db, day) {
const cutoff = shiftDays(day, -RETENTION_DAYS);
await db.batch([
db.prepare("DELETE FROM pings WHERE day < ?1").bind(cutoff),
+ 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 installs WHERE install_id NOT IN (SELECT DISTINCT install_id FROM pings)"),
]);
@@ -143,16 +144,20 @@ export async function purge(db, day) {
export async function handle(request, db, now = new Date()) {
const url = new URL(request.url);
const day = utcDay(now);
+ 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" });
+ }
if (url.pathname === "/v1/aggregate-stats") {
if (request.method !== "GET") return json({ error: "method not allowed" }, 405);
return json(await aggregateStats(db, shiftDays(day, -29)), 200, { "cache-control": "public, max-age=3600" });
}
- if (["/v0/ping", "/v1/ping", "/v1/aggregate"].includes(url.pathname)) {
+ if (["/v0/ping", "/v1/ping", "/v1/aggregate", "/v1/goals"].includes(url.pathname)) {
if (request.method !== "POST") return json({ error: "method not allowed" }, 405, { allow: "POST" });
if (!(request.headers.get("content-type") ?? "").startsWith("application/json")) {
return json({ error: "content-type must be application/json" }, 415);
}
- const limit = url.pathname === "/v1/aggregate" ? 16384 : MAX_BODY_BYTES;
+ const limit = ["/v1/aggregate", "/v1/goals"].includes(url.pathname) ? 16384 : MAX_BODY_BYTES;
// Bound streaming reads too: Content-Length can be absent or untrusted.
const reader = request.body?.getReader();
if (!reader) return json({ error: "missing body" }, 400);
@@ -175,6 +180,11 @@ export async function handle(request, db, now = new Date()) {
} catch {
return json({ error: "invalid JSON" }, 400);
}
+ if (url.pathname === "/v1/goals") {
+ if (!validGoalAggregate(parsed)) return json({ error: "invalid goal aggregate" }, 400);
+ await recordGoals(db, parsed, day);
+ return new Response(null, { status: 204 });
+ }
if (url.pathname === "/v1/aggregate") {
if (!validAggregate(parsed)) return json({ error: "invalid aggregate" }, 400);
await recordAggregate(db, parsed, day);
diff --git a/apps/usage-collector/test/collector.test.mjs b/apps/usage-collector/test/collector.test.mjs
index f2be6971bd..b07debf541 100644
--- a/apps/usage-collector/test/collector.test.mjs
+++ b/apps/usage-collector/test/collector.test.mjs
@@ -162,3 +162,32 @@ test("existing v0 database upgrades without deleting heartbeat history", () => {
assert.equal(db.prepare("SELECT count(*) n FROM pings").get().n, 1);
db.close();
});
+
+test("Goal duration ingestion uses real SQL, rejects identity, suppresses small cells and expires counts", async () => {
+ const db = d1();
+ const send = value => new Request("https://collector.example/v1/goals", {
+ method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(value),
+ });
+ const value = { schema: "loopx_goal_usage_aggregate_v1", counters: [{ span: "lt_30d", execution: "lt_1d", count: 3 }] };
+ for (const field of ["goal_id", "install_id", "timestamp", "prompt"]) {
+ assert.equal((await handle(send({ ...value, [field]: "private" }), db)).status, 400);
+ }
+ assert.equal((await handle(send(value), db, at("2026-09-01"))).status, 204);
+ const stats = () => handle(new Request("https://collector.example/v1/goal-stats"), db, at("2026-09-02")).then(r => r.json());
+ assert.deepEqual((await stats()).totals, { span: {}, execution: {} });
+ await handle(send(value), db, at("2026-09-01"));
+ assert.deepEqual((await stats()).totals, { span: { lt_30d: 6 }, execution: { lt_1d: 6 } });
+ assert.deepEqual(Object.keys(db.raw.get("SELECT * FROM goal_usage_counts")).sort(), ["count", "day", "execution", "span"]);
+ await purge(db, "2026-10-02");
+ assert.equal(db.raw.get("SELECT COUNT(*) n FROM goal_usage_counts").n, 0);
+});
+
+test("Goal migration is additive and preserves existing aggregate counters", () => {
+ 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/0002-goal-usage.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 goal_usage_counts").get().n, 0);
+ db.close();
+});
diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md
index 974194e2b7..158e62b98d 100644
--- a/docs/reference/usage-ping.md
+++ b/docs/reference/usage-ping.md
@@ -9,8 +9,8 @@ outcomes. No content collection is implemented.
```bash
loopx usage-ping status # current policy, recipient and outgoing payload previews
-loopx usage-ping disable # stop both channels; delete local ID and pending counts
-loopx usage-ping enable # explicitly allow both channels after reading the disclosure
+loopx usage-ping disable # stop all channels; delete local ID and pending counts
+loopx usage-ping enable # explicitly allow all channels after reading the disclosure
```
Settings → Capability Center exposes the same machine-wide switch and previews.
@@ -80,8 +80,8 @@ silently opt itself in: use the visible App setting or explicit CLI enable.
JSON stdout is unaffected. Previously enabled v0 clients keep their random ID
but must see the expanded-scope disclosure; previously disabled clients stay off.
-An explicit stored disable blocks both channels. The following environment
-settings also block both channels, even after explicit enable:
+An explicit stored disable blocks all channels. The following environment
+settings also block all channels, even after explicit enable:
- `LOOPX_USAGE_PING=0|false|no|off`
- `DO_NOT_TRACK` set to a nonempty value other than `0`
@@ -142,3 +142,57 @@ requests can never be correlated. Public unauthenticated counters can be
inflated, and suppression/loss makes these estimates unsuitable for billing.
The service may be unreachable on some networks; the owner can supply a reachable
collector, and LoopX continues to work without telemetry.
+
+## Observed Goal duration
+
+The Goal channel adds two duration histograms to basic statistics. It helps
+answer whether observed work continues for hours or days, and how much Host
+execution those spans contain. It does not identify a person or a Goal.
+
+- **Span:** first to most recent observed Host execution, including intervening
+ pauses. It stops growing while no execution is observed.
+- **Execution:** union of observed Host-call intervals for one Goal on one
+ machine. Concurrent or nested calls overlap only once; retry execution counts,
+ settlement-only replay does not. Network/tool/approval waits inside a Host call
+ are included; this is neither CPU time nor billing time.
+- **Sampling:** one cumulative snapshot per locally observed Goal-day, flushed
+ after that UTC day closes when another observation or normal usage occurs.
+ Unfinished Goals are included. A continuously executing Host checkpoints every
+ minute and can flush without another CLI command. Quiet Goals are not counted
+ again every day. These counts are **Goal-day observations, not unique Goals**;
+ the collector cannot join a Goal across days or machines.
+
+Coverage is managed `turn run-once` Host execution and regular owner Goal chat.
+Native `/goal`, externally attached agent sessions, manager and external-audience
+conversations are excluded because they do not share these timing boundaries.
+File, SQLite and PostgreSQL use the same observer; no provider state is queried
+or changed. Measurement starts when first observed after notice acknowledgment,
+not at historical Goal creation. Disabling, changing recipient, or clearing local
+state restarts measurement. There is no historical backfill or completion claim.
+
+All durations are lower bounds on observed work: confirmed prefixes survive a
+crash; missing final checkpoints, lock contention, network failure, suspended
+hosts and collection limits can lose observations. Never extrapolate a crashed
+Host as still executing. Local storage holds at most 64 Goals and 512 disjoint
+recent intervals per Goal; old intervals compact into totals, and 90-day inactive
+Goals expire. Delayed observations older than one day are discarded; pending
+snapshots older than seven days are discarded. Do not use this channel for
+liveness detection, quotas, acceptance, or accounting.
+
+The only outgoing Goal payload is:
+
+```json
+{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"span":"lt_7d","execution":"lt_6h","count":1}]}
+```
+
+Both durations use `lt_1m`, `lt_10m`, `lt_1h`, `lt_6h`, `lt_1d`, `lt_7d`,
+`lt_30d`, `gte_30d`. No Goal ID, installation ID, source path, name, event time
+or free text is sent. `/v1/goals` accepts the strict payload; `/v1/goal-stats`
+returns independent marginal histograms for the last 30 receipt days and omits
+cells below five. The existing usage settings switch, environment opt-outs and
+consent policy control all three channels. Settings and `loopx usage-ping status`
+show `goal_preview`; this is a current local snapshot, not a delivery receipt.
+The expanded scope requires notice version 2; previous explicit disable persists.
+
+Deploy collector migration `0002-goal-usage.sql` and its Worker before shipping
+the client. This additive table leaves existing heartbeats and CLI counts intact.
diff --git a/docs/reference/usage-ping.zh-CN.md b/docs/reference/usage-ping.zh-CN.md
index fe81152b14..bcd04d32ba 100644
--- a/docs/reference/usage-ping.zh-CN.md
+++ b/docs/reference/usage-ping.zh-CN.md
@@ -7,7 +7,7 @@
```bash
loopx usage-ping status # 查看策略、接收方、待发送数据
-loopx usage-ping disable # 关闭两条通道,删除本机 ID 和待发送计数
+loopx usage-ping disable # 关闭所有通道,删除本机 ID 和待发送计数
loopx usage-ping enable # 阅读告知后明确开启
```
@@ -70,7 +70,7 @@ ID 随机生成,不绑定账号、不从硬件派生,但能跨天关联,
需要所有者在 App 设置或 CLI 中明确开启。JSON stdout 保持不变。
旧版明确关闭的选择继续生效;旧版已经开启的保留随机 ID,但扩大范围前重新告知。
-以下任一设置都会压过“已开启”,同时关闭两条通道:
+以下任一设置都会压过“已开启”,同时关闭所有通道:
- `LOOPX_USAGE_PING=0|false|no|off`
- `DO_NOT_TRACK` 非空且不是 `0`
@@ -115,3 +115,46 @@ UTC 日期结束后的下一次合格调用发送上一日汇总,超过七天
observability。但网络服务商仍处理连接信息,分开数据包不保证绝对不可关联。
接口未认证,统计可能被灌水;缺失、抑制和采样偏差也使其不适用于计费。
部分网络无法访问该服务时可配置可达收集器,LoopX 的正常使用不受影响。
+
+## Goal 执行时长观测
+
+Goal 统计用于判断实际观测到的工作能否跨小时、跨天持续,以及其中有多少
+Host 执行时间。它不会识别某个用户或 Goal。
+
+- **执行跨度(span)**:首次到最近一次观测执行的时间,包含中间的暂停;
+ 没有新执行时不会继续增长,不是从历史创建日期推算的 Goal 年龄。
+- **累计执行(execution)**:同机同 Goal 的 Host 调用时间区间求并集,
+ 并行和嵌套重叠只算一次,真实重试计入,仅补结算的重放不计入。
+ 调用中的网络、工具和审批等待也包含在内,不是 CPU 用时或计费时长。
+- **采样**:每个本机有执行观测的 Goal 每 UTC 日最多产生一份累计快照,
+ 次日发生观测或普通使用时发送。持续运行的 Host 每分钟提交检查点,
+ 无须等待下一条 CLI 命令;未完成的 Goal 也会计入。长期静默的 Goal
+ 不会每天重复计数。统计单位是 **Goal-day 观测,不是去重 Goal 数**。
+
+覆盖受管 `turn run-once` 与普通用户 Goal 对话。原生 `/goal`、外部附加的
+Agent、管理者和外部受众对话暂不覆盖,避免将不同执行边界混作同一种统计。
+File、SQLite、PostgreSQL 共用这层旁观机制,不读取或改写 provider 状态。
+从告知确认后第一次实际观测开始计量;关闭、变更收集地址、清空本机状态会
+重新开始,不回填历史,也不把停止或命令成功解释为 Goal 完成。
+
+这些时长是已观测工作的下界:崩溃只保留已确认的检查点,缺少尾部检查点、
+锁竞争、网络失败、主机挂起和容量限制都可能漏计;不会把崩溃后的时间继续
+算作执行。本机最多保留 64 个 Goal、每个 Goal 512 个近期不相交区间,旧区间
+压缩为累计数,90 天无执行的 Goal 到期删除。超过一天的迟到观测、超过七天的
+待发快照丢弃。不能用此统计判定存活、验收、配额或计费。
+
+出站样例:
+
+```json
+{"schema":"loopx_goal_usage_aggregate_v1","counters":[{"span":"lt_7d","execution":"lt_6h","count":1}]}
+```
+
+两个时长均分为小于 1 分钟、10 分钟、1 小时、6 小时、1 天、7 天、30 天及
+不少于 30 天八档。不发送 Goal ID、安装 ID、名称、路径、事件时间或自由文本。
+`/v1/goals` 接收严格白名单数据;`/v1/goal-stats` 提供最近 30 个接收日的独立
+分布,少于 5 的单元不公开。设置页与 `loopx usage-ping status` 的 `goal_preview`
+可查看本机当前快照;它不是发送回执。
+
+沿用同一统计开关、环境变量和同意策略。扩大统计范围需要重新展示第 2 版告知,
+不会重置原有的关闭选择。发布客户端前先备份 D1,应用 `0002-goal-usage.sql` 并
+部署 Worker;新增表不会删除既有心跳或 CLI 计数。
diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py
index dfb1896891..0188bdf4ad 100644
--- a/loopx/chat_runtime.py
+++ b/loopx/chat_runtime.py
@@ -1274,12 +1274,18 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None:
if execution_lock is not None:
execution_lock.__exit__(None, None, None)
adapter.goal_driver = None
- elif attachments:
- if not isinstance(adapter, CodexAppServerAdapter):
- raise ValueError("image attachments currently require the Codex Agent endpoint")
- response = adapter.start_turn_with_attachments(message, event_sink, attachments)
else:
- response = adapter.start_turn(message, event_sink)
+ from .usage_goal import observe_goal_execution
+ # Manager/native/external conversations have different execution
+ # and waiting boundaries. Do not invent Goal timing for them.
+ observed_goal = str(session.get("goal_id") or "") if scope["kind"] == "owner_goal" else ""
+ with observe_goal_execution(self.store.root.parent, observed_goal):
+ if attachments:
+ if not isinstance(adapter, CodexAppServerAdapter):
+ raise ValueError("image attachments currently require the Codex Agent endpoint")
+ response = adapter.start_turn_with_attachments(message, event_sink, attachments)
+ else:
+ response = adapter.start_turn(message, event_sink)
if consume_interrupted():
event_buffer.close()
return
diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts
index 7846dc4bed..79cc679fef 100644
--- a/loopx/control_plane/runtime/usage_statistics.ts
+++ b/loopx/control_plane/runtime/usage_statistics.ts
@@ -1,5 +1,5 @@
/** Machine-local telemetry owner. Never opens Goal/provider state. */
-import { readFile, chmod } from "node:fs/promises";
+import { readFile, chmod, rm } from "node:fs/promises";
import { randomUUID } from "node:crypto";
import { arch, platform } from "node:os";
import type { JsonObject } from "../effect_program.ts";
@@ -7,9 +7,13 @@ 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 { recordGoalUsage, goalPreview } from "./usage_statistics_goals.ts";
+import { validGoalAggregate } from "./usage_statistics_goal_contract.ts";
+import type { GoalAggregate, GoalObservation } from "./usage_statistics_goal_contract.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 = 1;
+export const NOTICE_VERSION = 2;
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 };
@@ -85,19 +89,25 @@ export async function inspect(path: string, ctx: Context) {
notice: notice(ctx), notice_required: !sameNotice(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 },
+ goal_preview: state.consent === "disabled" ? null : await goalPreview(path + ".goals", state.generation).catch(() => null),
aggregate_day: state.day ?? null,
- disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. No prompts, code, paths, arguments, Goal data or raw errors. Disable both with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") };
+ disclosure: "LoopX basic usage statistics are on by default after this notice. Daily heartbeats send a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX collector (Cloudflare). Fixed CLI feature/result/duration/error counts are sent separately without an ID. Observed Goal execution-span and Host-call duration buckets are separately aggregated without Goal or installation IDs. Measurements cover managed Turns and regular owner Goal chat, from observation onward; they are lower bounds, not completion or billing evidence. No prompts, code, paths, arguments, Goal contents or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; inspect with loopx usage-ping status. Consent-required distributions wait for explicit enable. Recipient: " + (endpoint(ctx.env) || "not configured") };
}
export async function configure(path: string, ctx: Context, action: "enable" | "disable" | "acknowledge", expectedNotice?: unknown) {
await withFileMutationLock(path, async () => {
// Explicit disable can repair malformed state without permitting a send.
const state = action === "disable" ? { schema: STATE_SCHEMA, consent: "disabled", generation: randomUUID() } as State : await load(path);
- if (action === "disable") return save(path, state);
+ if (action === "disable") {
+ await save(path, state);
+ await rm(path + ".goals", { force: true });
+ return;
+ }
if (action === "acknowledge" && JSON.stringify(expectedNotice) !== JSON.stringify(notice(ctx))) throw new Error("usage_notice_changed");
if (action === "acknowledge" && state.consent === "disabled") return;
if (action === "enable") state.consent = "enabled";
if (state.notice && !sameNotice(state, ctx)) {
state.counters = [];
+ await rm(path + ".goals", { force: true });
state.generation = randomUUID();
if (state.notice.endpoint !== endpoint(ctx.env)) state.install_id = randomUUID();
}
@@ -108,29 +118,32 @@ export async function configure(path: string, ctx: Context, action: "enable" | "
}, 1000);
return inspect(path, ctx);
}
-export type Post = (url: string, payload: Ping | Aggregate) => Promise;
+export type Post = (url: string, payload: Ping | Aggregate | GoalAggregate) => 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) {
+export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation) {
if (counter !== null && (!validCounter(counter) || counter.count !== 1)) return { sent: false, reason: "invalid_observation" };
let heartbeat: Ping | null = null;
let aggregate: Aggregate | null = null;
+ let goals: GoalAggregate | null = null;
const today = day(ctx);
const allowed = await withFileMutationLock(path, async () => {
const state = await load(path);
const blocked = blockedBy(state, ctx);
if (blocked || !generation || generation !== state.generation) return false;
+ if (state.day && state.day > today) return false;
+ try { goals = await recordGoalUsage(path + ".goals", generation, (ctx.now ?? new Date()).getTime(), goal); }
+ catch { /* A damaged optional measurement cannot block other diagnostics. */ }
// Flush only a completed day's local aggregate. No event times or per-install key leave this boundary.
if (state.day && state.day < today && state.counters?.length) {
const age = Date.parse(today) - Date.parse(state.day);
if (age <= 7 * 86400000) aggregate = { schema: AGGREGATE_SCHEMA, counters: state.counters };
state.counters = [];
}
- if (state.day && state.day > today) return false; // backward clock: do not resend or mislabel counts
state.day = today;
state.counters ??= [];
if (counter) {
@@ -147,9 +160,9 @@ export async function observe(path: string, ctx: Context, generation: string, co
}, 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]] as const) {
+ 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) {
if (!payload) continue;
- if (!(validPing(payload) || validAggregate(payload))) continue;
+ if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload))) continue;
try {
let request: Promise | undefined;
// Start under the same short lock as disable, but never hold it while
diff --git a/loopx/control_plane/runtime/usage_statistics_cli.ts b/loopx/control_plane/runtime/usage_statistics_cli.ts
index 54988b2850..58c56b0c8f 100644
--- a/loopx/control_plane/runtime/usage_statistics_cli.ts
+++ b/loopx/control_plane/runtime/usage_statistics_cli.ts
@@ -3,6 +3,8 @@ 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 { validGoalObservation } from "./usage_statistics_goal_contract.ts";
+
try {
let input = "";
for await (const chunk of process.stdin) {
@@ -18,6 +20,8 @@ try {
result = await configure(request.path, ctx, request.action, request.notice);
} else if (request.action === "start") {
result = await observe(request.path, ctx, String(request.generation), null);
+ } else if (request.action === "goal" && validGoalObservation(request.observation, Date.now())) {
+ result = await observe(request.path, ctx, String(request.generation), null, undefined, request.observation);
} 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), {
diff --git a/loopx/control_plane/runtime/usage_statistics_goal_contract.ts b/loopx/control_plane/runtime/usage_statistics_goal_contract.ts
new file mode 100644
index 0000000000..8b05a65811
--- /dev/null
+++ b/loopx/control_plane/runtime/usage_statistics_goal_contract.ts
@@ -0,0 +1,34 @@
+import { object } from "./usage_statistics_contract.ts";
+
+export const GOAL_SCHEMA = "loopx_goal_usage_aggregate_v1";
+export const GOAL_DURATIONS = ["lt_1m", "lt_10m", "lt_1h", "lt_6h", "lt_1d", "lt_7d", "lt_30d", "gte_30d"] as const;
+export type GoalDuration = typeof GOAL_DURATIONS[number];
+export type GoalCount = { span: GoalDuration; execution: GoalDuration; count: number };
+export type GoalAggregate = { schema: typeof GOAL_SCHEMA; counters: GoalCount[] };
+export type GoalObservation = { key: string; start: number; end: number };
+const DAY = 86400000;
+export function goalDuration(ms: number): GoalDuration {
+ const limits = [60000, 600000, 3600000, 21600000, DAY, 7 * DAY, 30 * DAY];
+ return GOAL_DURATIONS[limits.findIndex(limit => ms < limit)] ?? "gte_30d";
+}
+export function validGoalAggregate(value: unknown): value is GoalAggregate {
+ if (!object(value) || Object.keys(value).sort().join() !== "counters,schema" || value.schema !== GOAL_SCHEMA
+ || !Array.isArray(value.counters) || !value.counters.length || value.counters.length > 64) return false;
+ const keys = new Set();
+ return value.counters.every(row => {
+ if (!object(row) || Object.keys(row).sort().join() !== "count,execution,span"
+ || !(GOAL_DURATIONS as readonly unknown[]).includes(row.span) || !(GOAL_DURATIONS as readonly unknown[]).includes(row.execution)
+ || !Number.isInteger(row.count) || Number(row.count) < 1 || Number(row.count) > 128) return false;
+ const key = `${row.span}:${row.execution}`;
+ if (keys.has(key)) return false;
+ keys.add(key); return true;
+ });
+}
+export function validGoalObservation(value: unknown, now: number): value is GoalObservation {
+ return object(value) && Object.keys(value).sort().join() === "end,key,start"
+ && typeof value.key === "string" && /^[a-f0-9]{64}$/.test(value.key)
+ && Number.isSafeInteger(value.start) && Number.isSafeInteger(value.end)
+ && Number(value.start) > 0 && Number(value.start) <= Number(value.end)
+ && Number(value.end) <= now + 1000 && Number(value.end) >= now - DAY
+ && Number(value.end) - Number(value.start) <= 120000;
+}
diff --git a/loopx/control_plane/runtime/usage_statistics_goals.ts b/loopx/control_plane/runtime/usage_statistics_goals.ts
new file mode 100644
index 0000000000..d70f82c02f
--- /dev/null
+++ b/loopx/control_plane/runtime/usage_statistics_goals.ts
@@ -0,0 +1,92 @@
+/** Local, lossy execution observations; never a Goal lifecycle or billing authority. */
+import { readFile, chmod } from "node:fs/promises";
+import { atomicWriteJson } from "../effect_runtime_io.ts";
+import type { JsonObject } from "../effect_program.ts";
+
+import { object } from "./usage_statistics_contract.ts";
+import { GOAL_SCHEMA, goalDuration, validGoalObservation } from "./usage_statistics_goal_contract.ts";
+import type { GoalAggregate, GoalObservation, GoalCount } from "./usage_statistics_goal_contract.ts";
+type Interval = [number, number];
+type MeasuredGoal = { key: string; first: number; last: number; intervals: Interval[]; total: number; watermark: number; day: string; reported?: string };
+type GoalState = { generation: string; goals: MeasuredGoal[] };
+const DAY = 86400000;
+const MAX_GOALS = 64;
+const MAX_INTERVALS = 512;
+/** Interval union is idempotent, order-independent and excludes inter-Turn idle time. */
+export function union(intervals: Interval[], next: Interval): Interval[] {
+ const result: Interval[] = [];
+ for (const [start, end] of [...intervals, next].sort((a, b) => a[0] - b[0])) {
+ const tail = result.at(-1);
+ if (tail && start <= tail[1]) tail[1] = Math.max(tail[1], end);
+ else result.push([start, end]);
+ }
+ return result;
+}
+async function load(path: string, generation: string): Promise {
+ try {
+ const text = await readFile(path, "utf8");
+ if (text.length > 4 * 1024 * 1024) throw new Error("goal_usage_too_large");
+ const value = JSON.parse(text) as GoalState;
+ if (value.generation !== generation) return { generation, goals: [] };
+ if (!Array.isArray(value.goals) || value.goals.length > MAX_GOALS || value.goals.some(g => !object(g)
+ || typeof g.key !== "string" || !/^[a-f0-9]{64}$/.test(g.key)
+ || ![g.first, g.last, g.total, g.watermark].every(n => Number.isSafeInteger(n) && n >= 0)
+ || g.first > g.last || typeof g.day !== "string" || !/^\d{4}-\d{2}-\d{2}$/.test(g.day)
+ || !Array.isArray(g.intervals) || g.intervals.length > MAX_INTERVALS
+ || g.intervals.some((v, i) => !Array.isArray(v) || v.length !== 2 || !v.every(Number.isSafeInteger)
+ || v[0] < g.first || v[1] > g.last || v[0] > v[1] || (i > 0 && v[0] <= g.intervals[i - 1][1])))) throw new Error("goal_usage_invalid");
+ return value;
+ } catch (error) {
+ if ((error as NodeJS.ErrnoException).code === "ENOENT") return { generation, goals: [] };
+ throw error;
+ }
+}
+function snapshot(goals: MeasuredGoal[]): GoalAggregate | null {
+ const rows = new Map();
+ for (const goal of goals) {
+ const span = goalDuration(goal.last - goal.first);
+ const execution = goalDuration(goal.total + goal.intervals.reduce((sum, [a, b]) => sum + b - a, 0));
+ const key = `${span}:${execution}`;
+ const row = rows.get(key) ?? { span, execution, count: 0 };
+ row.count++; rows.set(key, row);
+ }
+ return rows.size ? { schema: GOAL_SCHEMA, counters: [...rows.values()] } : null;
+}
+export async function goalPreview(path: string, generation: string): Promise {
+ return snapshot((await load(path, generation)).goals.filter(g => g.reported !== g.day));
+}
+/** Caller holds the common consent lock. Claims before sending: no retry identifiers. */
+export async function recordGoalUsage(path: string, generation: string, now: number, observation?: GoalObservation): Promise {
+ const state = await load(path, generation);
+ const today = new Date(now).toISOString().slice(0, 10);
+ const ready = state.goals.filter(g => g.day < today && g.reported !== g.day && now - Date.parse(g.day) < 8 * DAY);
+ const payload = snapshot(ready);
+ for (const goal of ready) goal.reported = goal.day;
+ // A quiet Goal is not reported again each day. Keep observed cumulative time
+ // for up to 90 days of inactivity, then restart measurement if it returns.
+ state.goals = state.goals.filter(g => now - g.last < 90 * DAY);
+ if (observation && validGoalObservation(observation, now)) {
+ const observedDay = new Date(observation.end).toISOString().slice(0, 10);
+ let goal = state.goals.find(g => g.key === observation.key);
+ if (!goal && state.goals.length < MAX_GOALS) {
+ goal = { key: observation.key, first: observation.start, last: observation.end, intervals: [], total: 0, watermark: 0, day: observedDay };
+ state.goals.push(goal);
+ }
+ if (goal && observation.start >= goal.watermark && goal.day <= observedDay && goal.reported !== observedDay) {
+ // Compact only intervals older than the accepted late-observation window.
+ const boundary = now - DAY - 120000;
+ const retired = goal.intervals.filter(([, b]) => b < boundary);
+ goal.total += retired.reduce((sum, [a, b]) => sum + b - a, 0);
+ goal.intervals = goal.intervals.filter(([, b]) => b >= boundary);
+ goal.watermark = Math.max(goal.watermark, boundary);
+ const intervals = union(goal.intervals, [observation.start, observation.end]);
+ if (intervals.length <= MAX_INTERVALS) {
+ goal.intervals = intervals;
+ goal.first = Math.min(goal.first, observation.start); goal.last = Math.max(goal.last, observation.end);
+ goal.day = observedDay;
+ }
+ }
+ }
+ await atomicWriteJson(path, state as unknown as JsonObject); await chmod(path, 0o600);
+ return payload;
+}
diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py
index 61ada9ba7a..bb0fbd1c64 100644
--- a/loopx/control_plane/turn_driver/executor.py
+++ b/loopx/control_plane/turn_driver/executor.py
@@ -755,6 +755,8 @@ def _host_result_stage(
journal_path: Path,
effects: dict[str, bool],
confirm_start: Callable[[], None] | None = None,
+ usage_runtime_root: Path | None = None,
+ usage_goal_id: str = "",
) -> tuple[dict[str, Any] | None, list[str], dict[str, Any] | None]:
completed_phases = list(journal.get("completed_phases") or [])
result = (
@@ -769,16 +771,18 @@ def _host_result_stage(
# reservation. Confirmation failure stops before the host starts.
if confirm_start is not None:
confirm_start()
- host_observation = (
- _run_host_runner(request, runner=host_runner)
- if host_runner is not None
- else _run_host(
- request,
- argv=argv or [],
- project=project,
- timeout_seconds=timeout_seconds,
+ from ...usage_goal import observe_goal_execution
+ with observe_goal_execution(usage_runtime_root or project, usage_goal_id):
+ host_observation = (
+ _run_host_runner(request, runner=host_runner)
+ if host_runner is not None
+ else _run_host(
+ request,
+ argv=argv or [],
+ project=project,
+ timeout_seconds=timeout_seconds,
+ )
)
- )
effects["host_invoked"] = True
if not host_observation.get("ok"):
failure = _host_failure(
@@ -1418,6 +1422,8 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]:
plan,
request,
host_runner=host_runner,
+ usage_runtime_root=runtime_root,
+ usage_goal_id=goal_id,
argv=argv,
completion_lifecycle_configured=all(
callback is not None
diff --git a/loopx/usage_goal.py b/loopx/usage_goal.py
new file mode 100644
index 0000000000..ec67210e54
--- /dev/null
+++ b/loopx/usage_goal.py
@@ -0,0 +1,76 @@
+"""Best-effort Host observation transport; TS owns union, limits and consent.
+
+Not an execution controller. Only confirmed 60-second prefixes survive a crash;
+there is no open interval that a later process can extrapolate indefinitely.
+"""
+from __future__ import annotations
+
+from contextlib import contextmanager
+import hashlib
+import json
+import os
+from pathlib import Path
+import threading
+import time
+from collections.abc import Iterator
+
+from . import usage_ping
+
+
+@contextmanager
+def observe_goal_execution(runtime_root: Path, goal_id: str) -> Iterator[None]:
+ stop = threading.Event()
+ publish = None
+ try:
+ # This is only a scheduling hint. TS rechecks consent, environment,
+ # notice and generation under the same lock used by disable.
+ path = usage_ping.state_path()
+ state = json.loads(path.read_text())
+ generation = state.get("generation")
+ if (goal_id and generation and state.get("consent") != "disabled"
+ and os.environ.get("LOOPX_USAGE_PING") != "0"
+ and os.environ.get("DO_NOT_TRACK") != "1"
+ and os.environ.get("CI") != "true"):
+ key = hashlib.sha256(json.dumps([generation, str(runtime_root.resolve()), goal_id]).encode()).hexdigest()
+ wall = time.time_ns() // 1_000_000
+ origin = time.monotonic()
+ previous = 0
+ lock = threading.Lock()
+
+ def checkpoint() -> None:
+ nonlocal previous
+ try:
+ if not lock.acquire(blocking=False):
+ return
+ try:
+ elapsed = int((time.monotonic() - origin) * 1000)
+ # Scheduling suspension is not proven execution. Drop a
+ # delayed prefix rather than calling hours asleep work.
+ start = wall + previous
+ previous = elapsed
+ usage_ping._detach(usage_ping._request(
+ "goal", path, generation=generation,
+ observation={"key": key, "start": start, "end": wall + elapsed},
+ ))
+ finally:
+ lock.release()
+ except Exception:
+ pass
+
+ def periodically() -> None:
+ while not stop.wait(60):
+ checkpoint()
+
+ publish = checkpoint
+ worker = threading.Thread(target=periodically, daemon=True, name="loopx-usage-goal")
+ worker.start()
+ except Exception:
+ pass
+ try:
+ yield
+ finally:
+ stop.set()
+ # No join or HTTP await on the business path. A racing checkpoint is
+ # harmless: the TS interval union deduplicates it.
+ if publish is not None:
+ publish()
diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py
index aea710393f..d21f6c2a48 100644
--- a/loopx/usage_ping.py
+++ b/loopx/usage_ping.py
@@ -71,7 +71,7 @@ def begin(command: str) -> tuple[str, float] | None:
state = json.loads(path.read_text()) if path.exists() else {}
if state.get("consent") == "disabled":
return None
- if not state.get("notice"):
+ if (state.get("notice") or {}).get("version") != 2:
# Unattended machines remain silent until the owner sees the notice
# or explicitly enables from CLI/settings. JSON stdout stays clean.
if not sys.stderr.isatty():
diff --git a/tests/control_plane_ts/usage_statistics_goals.test.ts b/tests/control_plane_ts/usage_statistics_goals.test.ts
new file mode 100644
index 0000000000..410ab6275f
--- /dev/null
+++ b/tests/control_plane_ts/usage_statistics_goals.test.ts
@@ -0,0 +1,134 @@
+import assert from "node:assert/strict";
+import test from "node:test";
+import { mkdtemp, readFile, rm } 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 type { Context, Post } from "../../loopx/control_plane/runtime/usage_statistics.ts";
+import { union, recordGoalUsage, goalPreview } from "../../loopx/control_plane/runtime/usage_statistics_goals.ts";
+import { GOAL_SCHEMA, validGoalAggregate, validGoalObservation, goalDuration } from "../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts";
+const base = Date.parse("2026-09-01T12:00:00Z");
+const day = 86400000;
+const key = "a".repeat(64);
+async function fixture(t: test.TestContext) {
+ const root = await mkdtemp(join(tmpdir(), "loopx-goal-usage-"));
+ t.after(() => rm(root, { recursive: true, force: true }));
+ return join(root, "usage.json");
+}
+const ctx = (now: number): Context => ({ env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:1/v1/ping" }, version: "1.0.0", python: "3.13", channel: "source", now: new Date(now) });
+const checkpoint = (start: number, end: number) => ({ key, start, end });
+
+test("union counts concurrent/nested intervals once and preserves idle gaps under replay and reordering", () => {
+ const observations: [number, number][] = [[0, 10], [4, 6], [8, 12], [20, 25], [0, 10]];
+ for (const order of [observations, [...observations].reverse()]) {
+ assert.deepEqual(order.reduce((all, next) => union(all, next), [] as [number, number][]), [[0, 12], [20, 25]]);
+ }
+});
+
+test("unfinished Goals report once per observed day; concurrent retries add only real execution", async t => {
+ const path = await fixture(t);
+ await recordGoalUsage(path, "generation", base + 60000, checkpoint(base, base + 60000));
+ await recordGoalUsage(path, "generation", base + 90000, checkpoint(base + 30000, base + 90000));
+ await recordGoalUsage(path, "generation", base + 90000, checkpoint(base, base + 60000));
+ // Retry after a ten-hour pause: span includes the pause, execution does not.
+ await recordGoalUsage(path, "generation", base + 10 * 3600000 + 30000, checkpoint(base + 10 * 3600000, base + 10 * 3600000 + 30000));
+ const expected = { schema: GOAL_SCHEMA, counters: [{ span: "lt_1d", execution: "lt_10m", count: 1 }] };
+ assert.deepEqual(await goalPreview(path, "generation"), expected);
+ assert.deepEqual(await recordGoalUsage(path, "generation", base + day), expected);
+ assert.equal(await recordGoalUsage(path, "generation", base + day), null);
+ assert.equal(await recordGoalUsage(path, "generation", base + 2 * day), null);
+ assert.equal(await goalPreview(path, "generation"), null);
+});
+
+test("restart retains measurements, compaction retains totals and rejects old replay", async t => {
+ const path = await fixture(t);
+ await recordGoalUsage(path, "g", base + 60000, checkpoint(base, base + 60000));
+ await recordGoalUsage(path, "g", base + 3 * day, checkpoint(base + 3 * day - 60000, base + 3 * day));
+ const state = JSON.parse(await readFile(path, "utf8"));
+ assert.equal(state.goals[0].total, 60000);
+ assert.equal(state.goals[0].intervals.length, 1);
+ await recordGoalUsage(path, "g", base + 3 * day, checkpoint(base, base + 60000));
+ assert.equal(JSON.parse(await readFile(path, "utf8")).goals[0].total, 60000);
+ assert.deepEqual((await goalPreview(path, "g"))?.counters, [{ span: "lt_7d", execution: "lt_10m", count: 1 }]);
+ // Consent generation reset has no continuity with the previous measurement.
+ assert.equal(await goalPreview(path, "new-generation"), null);
+});
+
+test("no execution after last confirmed prefix is inferred when host disappears", async t => {
+ const path = await fixture(t);
+ await recordGoalUsage(path, "g", base + 10000, checkpoint(base, base + 10000));
+ const result = await recordGoalUsage(path, "g", base + 5 * day);
+ assert.deepEqual(result?.counters, [{ span: "lt_1m", execution: "lt_1m", count: 1 }]);
+ assert.equal(await recordGoalUsage(path, "g", base + 6 * day), null);
+ assert.equal(validGoalObservation(checkpoint(base, base + 10 * day), base + 10 * day), false);
+ assert.equal(validGoalObservation(checkpoint(base + 10, base), base), false);
+});
+
+test("all channels obey consent; disable removes local measurements and rejects stale work", async t => {
+ const path = await fixture(t); const now = ctx(base + 60000);
+ await configure(path, now, "enable");
+ const generation = JSON.parse(await readFile(path, "utf8")).generation;
+ const sent: unknown[] = [];
+ const post: Post = async (_, payload) => { sent.push(payload); return 204; };
+ await observe(path, now, generation, null, post, checkpoint(base, base + 60000));
+ assert.ok((await inspect(path, now)).goal_preview);
+ await observe(path, ctx(base + day), generation, null, post);
+ assert.equal(sent.filter(validGoalAggregate).length, 1);
+ const wire = JSON.stringify(sent.filter(validGoalAggregate));
+ assert.ok(!wire.includes(key) && !wire.includes(generation) && !wire.includes("install_id"));
+ await configure(path, now, "disable");
+ await assert.rejects(readFile(path + ".goals"), /ENOENT/);
+ await observe(path, ctx(base + 2 * day), generation, null, post, checkpoint(base + 2 * day - 10000, base + 2 * day));
+ await assert.rejects(readFile(path + ".goals"), /ENOENT/);
+ assert.equal((await inspect(path, now)).goal_preview, null);
+});
+
+test("new scope needs notice again without undoing explicit disable", async t => {
+ const path = await fixture(t); const now = ctx(base);
+ await configure(path, now, "enable");
+ const state = JSON.parse(await readFile(path, "utf8"));
+ state.notice.version = 1;
+ const { writeFile } = await import("node:fs/promises");
+ await writeFile(path, JSON.stringify(state));
+ assert.equal((await inspect(path, now)).blocked_by, "notice_required");
+ await configure(path, now, "disable");
+ await configure(path, now, "acknowledge", (await inspect(path, now)).notice);
+ assert.equal((await inspect(path, now)).blocked_by, "disabled");
+});
+
+test("closed contract rejects identifiers, timestamps and extra fields; long durations are not clipped at a minute", () => {
+ const valid = { schema: GOAL_SCHEMA, counters: [{ span: "gte_30d", execution: "lt_7d", count: 5 }] };
+ assert.ok(validGoalAggregate(valid));
+ for (const extra of [{ goal_id: "private" }, { install_id: key }, { timestamp: base }]) {
+ assert.equal(validGoalAggregate({ ...valid, ...extra }), false);
+ assert.equal(validGoalAggregate({ ...valid, counters: [{ ...valid.counters[0], ...extra }] }), false);
+ }
+ assert.equal(validGoalAggregate({ ...valid, counters: [valid.counters[0], valid.counters[0]] }), false);
+ assert.equal(goalDuration(30 * day), "gte_30d");
+});
+
+test("late replay after daily claim cannot invent a second observed day", async t => {
+ const path = await fixture(t);
+ const observation = checkpoint(base, base + 10000);
+ await recordGoalUsage(path, "g", base + 10000, observation);
+ assert.ok(await recordGoalUsage(path, "g", base + day - 1000, observation));
+ assert.equal(await recordGoalUsage(path, "g", base + day), null);
+ assert.equal(await goalPreview(path, "g"), null);
+});
+
+test("disable during heartbeat prevents a queued Goal aggregate from starting", async t => {
+ const path = await fixture(t);
+ await configure(path, ctx(base), "enable");
+ const generation = JSON.parse(await readFile(path, "utf8")).generation;
+ await observe(path, ctx(base + 10000), generation, null, async () => 204, checkpoint(base, base + 10000));
+ const sent: string[] = [];
+ await observe(path, ctx(base + day), generation, null, async (url) => {
+ sent.push(url);
+ // Send starts under the lock; asynchronous completion runs outside it.
+ await new Promise(resolve => setTimeout(resolve, 10));
+ await configure(path, ctx(base + day), "disable");
+ return 204;
+ });
+ assert.deepEqual(sent, ["http://127.0.0.1:1/v1/ping"]);
+ await assert.rejects(readFile(path + ".goals"), /ENOENT/);
+});
diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py
index 0e8c755579..bdefa2b32e 100644
--- a/tests/test_loopx_turn_executor.py
+++ b/tests/test_loopx_turn_executor.py
@@ -2964,3 +2964,39 @@ def test_run_once_does_not_refuse_a_launchable_managed_executor(tmp_path):
timeout_seconds=5,
execute=True,
)
+
+
+def test_real_host_duration_is_observed_but_settlement_replay_is_not(tmp_path, monkeypatch):
+ """Production Turn entrypoint, actual Host subprocess and detached TS state."""
+ import time
+ from loopx import usage_ping
+
+ machine = tmp_path / "machine"
+ monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", machine)
+ for name in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING", "LOOPX_USAGE_POLICY"):
+ monkeypatch.delenv(name, raising=False)
+ monkeypatch.setenv("LOOPX_USAGE_PING_ENDPOINT", "http://127.0.0.1:1/v1/ping")
+ usage_ping.control("enable")
+ plan = _plan()
+ host_file = tmp_path / "host.py"
+ host_file.write_text("import json, sys, time\njson.load(sys.stdin)\ntime.sleep(0.08)\nprint(" + repr(json.dumps(_host_result(plan))) + ")\n")
+ calls = {"writeback": 0, "spend": 0, "scheduler": 0}
+ writeback, spend, scheduler = _callbacks(calls)
+ options = dict(host_argv=[sys.executable, str(host_file)], project=tmp_path,
+ runtime_root=tmp_path / "runtime", goal_id="fixture-goal",
+ timeout_seconds=5, execute=True, task_validator=_passing_validator,
+ writeback=writeback, spend=spend, scheduler=scheduler)
+ result = run_loopx_turn_once(plan, **options)
+ assert result["ok"], result
+ local = machine / "usage-ping.json.goals"
+ deadline = time.monotonic() + 5
+ while not local.exists() and time.monotonic() < deadline:
+ time.sleep(0.02)
+ before = json.loads(local.read_text())
+ intervals = before["goals"][0]["intervals"]
+ assert sum(end - start for start, end in intervals) >= 80
+ assert "fixture-goal" not in local.read_text()
+ assert run_loopx_turn_once(plan, **options)["ok"]
+ time.sleep(0.15)
+ assert json.loads(local.read_text()) == before
+ assert calls == {"writeback": 1, "spend": 1, "scheduler": 1}
diff --git a/tests/test_usage_goal.py b/tests/test_usage_goal.py
new file mode 100644
index 0000000000..bd33e638f1
--- /dev/null
+++ b/tests/test_usage_goal.py
@@ -0,0 +1,59 @@
+"""Host observer remains outside execution authority and never waits for HTTP."""
+import json
+import threading
+import time
+
+import pytest
+
+from loopx import usage_goal, usage_ping
+
+
+def test_telemetry_failure_cannot_replace_host_exception(tmp_path, monkeypatch):
+ monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path)
+ usage_ping.state_path().write_text(json.dumps({"generation": "fixture", "consent": "enabled"}))
+ monkeypatch.delenv("CI", raising=False)
+ monkeypatch.setattr(usage_ping, "_detach", lambda _: (_ for _ in ()).throw(OSError("fixture failure")))
+ with pytest.raises(ValueError, match="host failure"):
+ with usage_goal.observe_goal_execution(tmp_path, "fixture-goal"):
+ raise ValueError("host failure")
+
+
+def test_disabled_observer_starts_no_worker_or_process(tmp_path, monkeypatch):
+ monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path)
+ usage_ping.state_path().write_text(json.dumps({"generation": "fixture", "consent": "disabled"}))
+ monkeypatch.setattr(threading.Thread, "start", lambda _: pytest.fail("disabled worker"))
+ monkeypatch.setattr(usage_ping, "_detach", lambda _: pytest.fail("disabled process"))
+ with usage_goal.observe_goal_execution(tmp_path, "fixture-goal"):
+ pass
+
+
+def test_periodic_observation_does_not_need_turn_completion(tmp_path, monkeypatch):
+ monkeypatch.setattr(usage_ping, "DEFAULT_RUNTIME_ROOT", tmp_path)
+ usage_ping.state_path().write_text(json.dumps({"generation": "fixture", "consent": "enabled"}))
+ for name in ("CI", "DO_NOT_TRACK", "LOOPX_USAGE_PING"):
+ monkeypatch.delenv(name, raising=False)
+ observed = []
+ monkeypatch.setattr(usage_ping, "_detach", observed.append)
+ # Accelerate only the observer's checkpoint interval, not clocks or threads.
+ real_event = threading.Event
+ from types import SimpleNamespace
+
+ class CheckpointEvent:
+ def __init__(self):
+ self.event = real_event()
+
+ def wait(self, seconds):
+ return self.event.wait(0.01)
+
+ def set(self):
+ self.event.set()
+
+ monkeypatch.setattr(usage_goal, "threading", SimpleNamespace(Event=CheckpointEvent, Lock=threading.Lock, Thread=threading.Thread))
+ with usage_goal.observe_goal_execution(tmp_path, "private-goal"):
+ deadline = time.monotonic() + 1
+ while not observed and time.monotonic() < deadline:
+ time.sleep(0.01)
+ assert observed, "unfinished Host should checkpoint"
+ assert all(row["observation"]["start"] <= row["observation"]["end"] for row in observed)
+ assert "private-goal" not in json.dumps(observed)
+ assert str(tmp_path) not in json.dumps([row["observation"] for row in observed])
diff --git a/tests/test_usage_ping.py b/tests/test_usage_ping.py
index d8422ab570..c971a9c7a3 100644
--- a/tests/test_usage_ping.py
+++ b/tests/test_usage_ping.py
@@ -187,3 +187,18 @@ def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated):
connection.close()
server.shutdown()
server.server_close()
+
+
+def test_goal_observer_does_not_wait_for_unresponsive_collector(isolated, collector, monkeypatch):
+ from loopx.usage_goal import observe_goal_execution
+ endpoint, received, accepted, release = collector
+ monkeypatch.setenv('LOOPX_USAGE_PING_ENDPOINT', endpoint)
+ usage_ping.control('enable')
+ started = time.monotonic()
+ with observe_goal_execution(isolated / 'runtime', 'synthetic-goal'):
+ time.sleep(0.01)
+ assert time.monotonic() - started < 0.8
+ assert accepted.wait(4)
+ assert not release.is_set(), 'host returned before collector released HTTP response'
+ usage_ping.control('disable')
+ release.set()
diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json
index 4fcbaedad6..ad12babe1f 100644
--- a/tsconfig.control-plane.json
+++ b/tsconfig.control-plane.json
@@ -15,6 +15,7 @@
"include": [
"loopx/control_plane/runtime/usage_statistics*.ts",
"tests/control_plane_ts/usage_statistics.test.ts",
+ "tests/control_plane_ts/usage_statistics_goals.test.ts",
"loopx/control_plane/effect_program.ts",
"loopx/control_plane/effect_runtime_errors.ts",
"loopx/control_plane/effect_runtime_handlers.ts",