From efd30749988ed0e979ae168646ded342cce986fe Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Fri, 2 Oct 2026 22:43:31 +0800 Subject: [PATCH 1/4] feat(telemetry): measure installation usage with verified independent clocks Signed-off-by: huangruiteng --- apps/presentation/dashboard/src/data/chat.ts | 8 +- .../usage-statistics-notice.tsx | 8 +- .../usage-statistics-settings.tsx | 25 ++++- .../migrations/0005-installation-usage.sql | 9 ++ apps/usage-collector/schema.sql | 9 ++ apps/usage-collector/src/collector.js | 14 ++- .../usage-collector/src/installation-usage.ts | 16 +++ loopx/chat_usage_statistics_api.py | 5 + loopx/cli_commands/usage_ping.py | 10 +- .../control_plane/runtime/usage_statistics.ts | 72 +++++++++--- .../runtime/usage_statistics_cli.ts | 11 +- .../runtime/usage_statistics_codex.ts | 23 ++-- .../runtime/usage_statistics_diagnostics.ts | 20 +++- .../runtime/usage_statistics_installation.ts | 104 ++++++++++++++++++ .../usage_statistics_installation_contract.ts | 41 +++++++ loopx/usage_ping.py | 2 +- 16 files changed, 329 insertions(+), 48 deletions(-) create mode 100644 apps/usage-collector/migrations/0005-installation-usage.sql create mode 100644 apps/usage-collector/src/installation-usage.ts create mode 100644 loopx/control_plane/runtime/usage_statistics_installation.ts create mode 100644 loopx/control_plane/runtime/usage_statistics_installation_contract.ts diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 8e6ef9e604..de6a162815 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -2259,12 +2259,18 @@ const usageStatisticsSchema = z.object({ 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(), + stored_context: z.string().optional(), effective_context: z.string().optional(), context_source: z.string().optional(), + installation_preview: z.unknown().optional(), identity_scope: z.string().optional(), delivery_history: z.array(z.object({ - day: z.string(), channel: z.enum(["heartbeat", "cli", "goal"]), rows: z.number(), + day: z.string(), channel: z.enum(["heartbeat", "cli", "goal", "installation"]), rows: z.number(), status: z.enum(["accepted", "rejected", "unavailable"]), })).optional(), }); export type UsageStatistics = z.infer; +export async function setUsageContext(context: string): Promise { + return usageStatisticsSchema.parse(await requestJson("/api/chat/usage-statistics", + { method: "POST", body: JSON.stringify({ context }) })); +} export async function usageStatistics(enabled?: boolean): Promise { return usageStatisticsSchema.parse(await requestJson("/api/chat/usage-statistics", enabled === undefined ? undefined : { method: "POST", body: JSON.stringify({ enabled }) })); 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 c407fce0da..660c2540fa 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 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."}

+ ? "用于改进平台支持与使用体验。随机安装标识会关联每日 CLI 功能计数、版本、日期、自愿环境标签与已观测运行分钟。区间按安装去重;不同计时口径不能相加,不代表机器在线或任务完成。不采集对话、代码、路径或参数值。可随时关闭。" + : "Helps improve platform support and usage. A random installation ID links to daily CLI counts, version, date, voluntary context and observed runtime minutes. Overlapping intervals are deduplicated per installation; different clocks cannot be added, and are not uptime or task completion. No conversations, code, paths or argument values. You can turn it off at any time."}

{zh - ? "首个已测量的 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."}

+ ? "独立无 ID 的 CLI 汇总保留;新增安装级概要由活动触发,至少间隔 15 分钟发送。关闭会清除本机标识和测量记录,已发送的记录不能撤回。" + : "ID-free CLI summaries remain supported. New installation profiles are activity-triggered, at least 15 minutes apart. Disabling clears local ID and measurement history; it cannot recall already-sent records."}

{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 4703405646..3079e1e22b 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 @@ -1,5 +1,5 @@ import { useEffect, useState } from "react"; -import { usageStatistics, type UsageStatistics } from "../../data/chat"; +import { setUsageContext, usageStatistics, type UsageStatistics } from "../../data/chat"; import { useWorkspaceI18n } from "./i18n"; export function UsageStatisticsSettings() { @@ -20,14 +20,29 @@ export function UsageStatisticsSettings() { catch { setError(true); } finally { setBusy(false); } } + async function updateContext(context: string) { + setBusy(true); setError(false); + try { setState(await setUsageContext(context)); } + catch { setError(true); } + finally { setBusy(false); } + } return
{zh ? "基础使用统计 · 告知后默认开启,可关闭" : "Basic usage statistics · on after notice, optional"}

{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 ? "新增固定子操作、版本、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."}

+ ? "用于决定平台支持和改进命令体验。每天向 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. The existing separate CLI summaries remain ID-free, with fixed feature, result, duration and error counts. The first measured result attempts a send immediately; later activity sends batches at least 15 minutes apart."}

+

{zh ? "安装级每日概要将同一个随机安装标识与固定 CLI 功能计数、版本、UTC 日期、自愿环境标签和已观测运行分钟关联。设备设置可持久保存,LOOPX_USAGE_CONTEXT 优先;默认未知,不推断个人或企业。不上传对话、代码、路径、参数值或原始错误。" : "Daily profiles link the same random installation ID with fixed CLI counts, version, UTC date, voluntary context and observed runtime minutes. Device labels persist; LOOPX_USAGE_CONTEXT overrides them. Unknown by default; never inferred. No conversations, code, paths, argument values or raw errors."}

{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 ? <> + +

{zh ? "实际生效:" : "Effective: "}{state.effective_context ?? "unknown"} · {state.context_source ?? "default"}. + {zh ? " 环境变量优先;设置用途不会开启统计或改变安装标识,已记录的每日标签不改写。" : " Environment overrides the device label. This setting does not enable statistics or change the ID; recorded daily labels are not rewritten."}

{state.sending ? (zh ? "已允许发送" : "Sending allowed") @@ -35,7 +50,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" ? : null}

{zh ? "接收地址:" : "Recipient: "}{state.endpoint ?? (zh ? "未配置" : "Not configured")}

-
{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)}
+
{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, installation: state.installation_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/migrations/0005-installation-usage.sql b/apps/usage-collector/migrations/0005-installation-usage.sql new file mode 100644 index 0000000000..1990ea454f --- /dev/null +++ b/apps/usage-collector/migrations/0005-installation-usage.sql @@ -0,0 +1,9 @@ +-- Additive. No historical linkage, context or runtime backfill. +CREATE TABLE IF NOT EXISTS installation_usage ( + activity_day TEXT NOT NULL, install_id TEXT NOT NULL, + version TEXT NOT NULL, context TEXT NOT NULL, revision INTEGER NOT NULL, + cli TEXT NOT NULL, runtime TEXT NOT NULL, truncated INTEGER NOT NULL, + receipt_day TEXT NOT NULL, + PRIMARY KEY (activity_day, install_id) +); +CREATE INDEX IF NOT EXISTS installation_usage_id ON installation_usage (install_id); diff --git a/apps/usage-collector/schema.sql b/apps/usage-collector/schema.sql index 6ae33695f8..2d0b671ce4 100644 --- a/apps/usage-collector/schema.sql +++ b/apps/usage-collector/schema.sql @@ -1,4 +1,13 @@ -- LoopX usage collector (Cloudflare D1). One row per installation per UTC day. +CREATE TABLE IF NOT EXISTS installation_usage ( + activity_day TEXT NOT NULL, install_id TEXT NOT NULL, + version TEXT NOT NULL, context TEXT NOT NULL, revision INTEGER NOT NULL, + cli TEXT NOT NULL, runtime TEXT NOT NULL, truncated INTEGER NOT NULL, + receipt_day TEXT NOT NULL, + PRIMARY KEY (activity_day, install_id) +); +CREATE INDEX IF NOT EXISTS installation_usage_id ON installation_usage (install_id); + CREATE TABLE IF NOT EXISTS installs ( install_id TEXT PRIMARY KEY, first_day TEXT NOT NULL diff --git a/apps/usage-collector/src/collector.js b/apps/usage-collector/src/collector.js index 8581232057..6c6e2d87e7 100644 --- a/apps/usage-collector/src/collector.js +++ b/apps/usage-collector/src/collector.js @@ -1,5 +1,6 @@ import { validAggregate, validPing, recordAggregate, aggregateStats, validGoalAggregate, recordGoals, goalStats } from "./basic-usage.ts"; import { validDiagnostics, recordDiagnostics, diagnosticStats } from "./basic-usage.ts"; +import { validInstallationUsage, recordInstallation } from "./installation-usage.ts"; // Pure request handling for the LoopX usage collector. worker.js binds it to // Cloudflare; tests bind it to an in-memory database. @@ -140,6 +141,7 @@ export async function purge(db, day) { 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 installation_usage WHERE activity_day < ?1").bind(shiftDays(day, -29)), db.prepare("DELETE FROM installs WHERE install_id NOT IN (SELECT DISTINCT install_id FROM pings)"), ]); } @@ -186,12 +188,12 @@ export async function handle(request, db, now = new Date()) { 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", "/v1/goals"].includes(url.pathname)) { + if (["/v0/ping", "/v1/ping", "/v1/aggregate", "/v1/goals", "/v1/installation"].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 = ["/v1/aggregate", "/v1/goals"].includes(url.pathname) ? 16384 : MAX_BODY_BYTES; + const limit = ["/v1/aggregate", "/v1/goals", "/v1/installation"].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); @@ -214,6 +216,14 @@ export async function handle(request, db, now = new Date()) { } catch { return json({ error: "invalid JSON" }, 400); } + if (url.pathname === "/v1/installation") { + if (!validInstallationUsage(parsed)) return json({ error: "invalid installation profile" }, 400); + if (parsed.profiles.some(row => row.activity_day > day || row.activity_day < shiftDays(day, -7))) { + return json({ error: "activity day outside retention window" }, 400); + } + await recordInstallation(db, parsed, day); + return new Response(null, { status: 204 }); + } if (url.pathname === "/v1/goals") { if (!validGoalAggregate(parsed)) return json({ error: "invalid goal aggregate" }, 400); await recordGoals(db, parsed, day); diff --git a/apps/usage-collector/src/installation-usage.ts b/apps/usage-collector/src/installation-usage.ts new file mode 100644 index 0000000000..7ee70b128f --- /dev/null +++ b/apps/usage-collector/src/installation-usage.ts @@ -0,0 +1,16 @@ +import { validInstallationUsage } from "../../../loopx/control_plane/runtime/usage_statistics_installation_contract.ts"; +import type { InstallationUsage } from "../../../loopx/control_plane/runtime/usage_statistics_installation_contract.ts"; +export { validInstallationUsage }; + +type Statement = { bind(...values: unknown[]): Statement; run(): Promise }; +type Database = { prepare(sql: string): Statement; batch(statements: Statement[]): Promise }; +/** Full daily snapshots replace only older revisions: transport replay never adds usage. */ +export async function recordInstallation(db: Database, payload: InstallationUsage, receiptDay: string) { + await db.batch(payload.profiles.map(row => db.prepare( + "INSERT INTO installation_usage (activity_day, install_id, version, context, revision, cli, runtime, truncated, receipt_day) " + + "VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) ON CONFLICT(activity_day, install_id) DO UPDATE SET " + + "revision=excluded.revision, cli=excluded.cli, runtime=excluded.runtime, truncated=excluded.truncated, receipt_day=excluded.receipt_day " + + "WHERE excluded.revision > installation_usage.revision AND excluded.context = installation_usage.context AND excluded.version = installation_usage.version", + ).bind(row.activity_day, payload.install_id, row.version, row.context, row.revision, JSON.stringify(row.cli), + JSON.stringify(row.runtime), Number(row.truncated), receiptDay))); +} diff --git a/loopx/chat_usage_statistics_api.py b/loopx/chat_usage_statistics_api.py index 316cf7c118..49b4a101c8 100644 --- a/loopx/chat_usage_statistics_api.py +++ b/loopx/chat_usage_statistics_api.py @@ -23,6 +23,11 @@ def _usage_statistics_status(self) -> None: def _usage_statistics_update(self) -> None: try: body = self._read_json() + if set(body) == {"context"} and isinstance(body["context"], str): + if body["context"] not in {"unknown", "personal", "shared_service", "ephemeral", "organization_managed", "maintainer"}: + raise ValueError("invalid deployment context") + self._usage_statistics_request("context", context=body["context"]) + return if set(body) == {"notice"} and isinstance(body["notice"], dict): self._usage_statistics_request("acknowledge", notice=body["notice"]) return diff --git a/loopx/cli_commands/usage_ping.py b/loopx/cli_commands/usage_ping.py index f31d4aecaf..a27f7078cc 100644 --- a/loopx/cli_commands/usage_ping.py +++ b/loopx/cli_commands/usage_ping.py @@ -19,12 +19,14 @@ 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'}", + f"- Deployment context: {payload.get('effective_context', 'unknown')} ({payload.get('context_source', 'default')})", 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"), "diagnostics": payload.get("diagnostic_preview"), "goals": payload.get("goal_preview"), + "installation": payload.get("installation_preview"), "diagnostic_dropped": payload.get("diagnostic_dropped", 0), "identity_scope": payload.get("identity_scope"), "delivery_history": payload.get("delivery_history", [])}, indent=2)) @@ -43,15 +45,19 @@ def register_usage_ping_command( parser.add_argument( "action", nargs="?", - choices=("status", "enable", "disable"), + choices=("status", "enable", "disable", "context"), default="status", help="status previews payloads; enable accepts collection; disable clears the ID and pending counts.", ) + parser.add_argument("--context", choices=("unknown", "personal", "shared_service", "ephemeral", "organization_managed", "maintainer"), + help="Persistent device label; only with context. Does not enable collection; environment takes precedence.") add_subcommand_format(parser) return parser def handle_usage_ping_command(args: argparse.Namespace, print_payload: PrintPayload) -> int: - payload = usage_ping.control(args.action) + if (args.action == "context") != (args.context is not None): + raise ValueError("use usage-ping context --context ; other actions take no context") + payload = usage_ping.control(args.action, **({"context": args.context} if args.context is not None else {})) print_payload(payload, output_format(args), render_usage_ping_markdown) return 0 diff --git a/loopx/control_plane/runtime/usage_statistics.ts b/loopx/control_plane/runtime/usage_statistics.ts index 9126be5712..7fa2258042 100644 --- a/loopx/control_plane/runtime/usage_statistics.ts +++ b/loopx/control_plane/runtime/usage_statistics.ts @@ -6,7 +6,7 @@ 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 { CONTEXTS, usageContext, 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"; @@ -15,15 +15,18 @@ import type { GoalAggregate, GoalObservation } from "./usage_statistics_goal_con import { cycleObservations } from "./usage_statistics_cycles.ts"; import type { CycleObservation } from "./usage_statistics_cycles.ts"; +import { recordInstallation, installationPreview } from "./usage_statistics_installation.ts"; +import { validInstallationUsage } from "./usage_statistics_installation_contract.ts"; +import type { InstallationUsage, ProfileFeature } from "./usage_statistics_installation_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 = 5; +export const NOTICE_VERSION = 6; 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" }; +type Delivery = { day: string; channel: "heartbeat" | "cli" | "goal" | "installation"; rows: number; status: "accepted" | "rejected" | "unavailable" }; const MAX_DELIVERIES = 20; type State = { schema: typeof STATE_SCHEMA; consent: "default" | "enabled" | "disabled"; generation: string; @@ -31,6 +34,7 @@ type State = { day?: string; counters?: Counter[]; aggregate_last_attempt_ms?: number; deliveries?: Delivery[]; diagnostics?: Diagnostic[]; diagnostic_dropped?: number; + context?: Diagnostic["context"]; }; export function endpoint(env: Env): string { try { @@ -87,6 +91,7 @@ async function load(path: string): Promise { || !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))) || (raw.diagnostic_dropped !== undefined && (!Number.isSafeInteger(raw.diagnostic_dropped) || Number(raw.diagnostic_dropped) < 0 || Number(raw.diagnostic_dropped) > MAX_COUNT)) + || (raw.context !== undefined && !(CONTEXTS as readonly unknown[]).includes(raw.context)) || (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; } @@ -97,7 +102,7 @@ async function save(path: string, state: State) { 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)) + && ["heartbeat", "cli", "goal", "installation"].includes(String(value.channel)) && ["accepted", "rejected", "unavailable"].includes(String(value.status)) && Number.isSafeInteger(value.rows) && Number(value.rows) >= 1 && Number(value.rows) <= MAX_ROWS; } @@ -120,20 +125,42 @@ export async function inspect(path: string, ctx: Context) { 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, + stored_context: state.context ?? "unknown", + effective_context: effectiveContext(state, ctx), + context_source: ctx.env.LOOPX_USAGE_CONTEXT !== undefined ? "environment" : state.context ? "device" : "default", + installation_preview: state.consent === "disabled" || !state.install_id ? null + : await installationPreview(path + ".installation", state.generation, state.install_id).catch(() => null), 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/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") }; + disclosure: "LoopX basic usage statistics are on by default after this notice. A daily heartbeat sends a random installation ID, version, OS, CPU architecture, Python version and install channel to the configured LoopX Cloudflare collector. Separate ID-free CLI/Goal summaries remain supported. New daily installation profiles link that same random ID to fixed CLI family counts, UTC activity date, release version, voluntary device context and observed runtime rounded down to minutes. Runtime unions overlapping intervals across Goals within each of host_call, codex_turn and quota_cycle; these clocks overlap and cannot be added. Waiting may be included, missing instrumentation remains unobserved: this is not uptime, CPU time, completion or billing. No per-call timestamps, Goal IDs or session content are uploaded. Device context defaults to unknown; LOOPX_USAGE_CONTEXT overrides the stored setting. Profile context/version freeze at the first observation of each UTC day, never relabeling history. Profiles are lossy full snapshots, activity-triggered at most every 15 minutes, not an always-on background timer. Collector profiles expire after 30 days; local interval buffers after seven UTC days. Network services may observe connection metadata. No prompts, code, paths, argument values or raw errors. Disable all with loopx usage-ping disable or LOOPX_USAGE_PING=0; disabling clears ID, counts and local measurement history but keeps the voluntary device label. Consent-required distributions wait for explicit enable. Inspect with loopx usage-ping status. Recipient: " + (endpoint(ctx.env) || "not configured") }; +} +export function effectiveContext(state: Pick, ctx: Context): Diagnostic["context"] { + return usageContext(ctx.env.LOOPX_USAGE_CONTEXT !== undefined ? ctx.env.LOOPX_USAGE_CONTEXT : state.context); +} +/** A device label is observation metadata, never consent, identity or work authority. */ +export async function configureContext(path: string, ctx: Context, value: unknown) { + if (!(CONTEXTS as readonly unknown[]).includes(value)) throw new Error("usage_context_invalid"); + await withFileMutationLock(path, async () => { + const state = await load(path); + state.generation ||= randomUUID(); + state.context = value as Diagnostic["context"]; + await save(path, state); + }, 1000); + return inspect(path, ctx); } 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); + const previous = action === "disable" ? await load(path).catch(() => null) : null; + const state = action === "disable" ? { schema: STATE_SCHEMA, consent: "disabled", generation: randomUUID(), + ...(previous?.context ? { context: previous.context } : {}) } as State : await load(path); if (action === "disable") { await save(path, state); await rm(path + ".goals", { force: true }); await rm(path + ".cycles", { force: true }); + await rm(path + ".installation", { force: true }); return; } if (action === "acknowledge" && JSON.stringify(expectedNotice) !== JSON.stringify(notice(ctx))) throw new Error("usage_notice_changed"); @@ -144,6 +171,7 @@ export async function configure(path: string, ctx: Context, action: "enable" | " state.diagnostics = []; await rm(path + ".goals", { force: true }); await rm(path + ".cycles", { force: true }); + await rm(path + ".installation", { force: true }); state.generation = randomUUID(); if (state.notice.endpoint !== endpoint(ctx.env)) state.install_id = randomUUID(); } @@ -154,14 +182,14 @@ export async function configure(path: string, ctx: Context, action: "enable" | " }, 1000); return inspect(path, ctx); } -export type Post = (url: string, payload: Ping | Aggregate | GoalAggregate | DiagnosticAggregate) => Promise; +export type Post = (url: string, payload: Ping | Aggregate | GoalAggregate | DiagnosticAggregate | InstallationUsage) => 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, diagnostic?: Diagnostic) { +export async function observe(path: string, ctx: Context, generation: string, counter: Counter | null, send: Post = post, goal?: GoalObservation, cycle?: CycleObservation, diagnostic?: Diagnostic, feature?: ProfileFeature) { 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; @@ -171,6 +199,7 @@ export async function observe(path: string, ctx: Context, generation: string, co let diagnostics: DiagnosticAggregate | null = null; let diagnosticRequest: Promise | undefined; let goals: GoalAggregate | null = null; + let installation: InstallationUsage | null = null; const today = day(ctx); const now = (ctx.now ?? new Date()).getTime(); const allowed = await withFileMutationLock(path, async () => { @@ -179,11 +208,16 @@ export async function observe(path: string, ctx: Context, generation: string, co if (blocked || !generation || generation !== state.generation) return false; if (state.day && state.day > today) return false; if (state.aggregate_last_attempt_ms !== undefined && now < state.aggregate_last_attempt_ms) return false; + let intervals: GoalObservation[] = goal ? [goal] : []; try { - const intervals = cycle ? await cycleObservations(path + ".cycles", generation, now, cycle) : []; - goals = await recordGoalUsage(path + ".goals", generation, now, [...(goal ? [goal] : []), ...intervals]); + intervals = [...intervals, ...(cycle ? await cycleObservations(path + ".cycles", generation, now, cycle) : [])]; + goals = await recordGoalUsage(path + ".goals", generation, now, intervals); } catch { /* A damaged optional measurement cannot block other diagnostics. */ } + try { + if (state.install_id && ping(state, ctx)) installation = await recordInstallation(path + ".installation", generation, state.install_id, + now, ctx.version, effectiveContext(state, ctx), feature ?? diagnostic?.feature ?? counter?.feature, intervals); + } catch { /* Missing measurements remain unknown; never invent runtime. */ } // Keep the oldest buffered UTC day for expiry, including across midnight. // Legacy daily buffers remain readable; no event times or join keys leave. if (state.day && Date.parse(today) - Date.parse(state.day) > 7 * 86400000) state.counters = []; @@ -196,9 +230,10 @@ export async function observe(path: string, ctx: Context, generation: string, co } 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)); + const observed = { ...diagnostic, context: effectiveContext(state, ctx) }; + const row = state.diagnostics.find(entry => diagnosticKey(entry) === diagnosticKey(observed)); if (row) row.count = Math.min(MAX_COUNT, row.count + 1); - else if (state.diagnostics.length < 32) state.diagnostics.push({ ...diagnostic }); + else if (state.diagnostics.length < 32) state.diagnostics.push(observed); else state.diagnostic_dropped = Math.min(MAX_COUNT, (state.diagnostic_dropped ?? 0) + 1); } if (!state.last_attempt_day || state.last_attempt_day < today) { @@ -232,25 +267,26 @@ export async function observe(path: string, ctx: Context, generation: string, co catch { diagnosticRequest = Promise.resolve(0); } } return true; - }, 0); // Never queue behind business or telemetry work. + }, 250); // Detached workers may briefly contend with startup; never block the host command. if (!allowed) return { sent: false, reason: "blocked" }; let sent = false; 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], + [endpoint(ctx.env).replace(/\/ping$/, "/installation"), installation], ]; for (const [url, payload] of outgoing) { if (!payload) continue; - if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload) || validDiagnostics(payload))) continue; + if (!(validPing(payload) || validAggregate(payload) || validGoalAggregate(payload) || validDiagnostics(payload) || validInstallationUsage(payload))) continue; try { 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")) { + if (url.endsWith("/goals") || url.endsWith("/installation")) { await withFileMutationLock(path, async () => { const current = await load(path); if (!blockedBy(current, ctx) && current.generation === generation) request = send(url, payload).catch(() => 0); - }, 0); + }, 250); } if (!request) break; const code = await request; @@ -261,8 +297,8 @@ export async function observe(path: string, ctx: Context, generation: string, co 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, + channel: url.endsWith("/ping") ? "heartbeat" : url.endsWith("/aggregate") ? "cli" : url.endsWith("/installation") ? "installation" : "goal", + rows: "counters" in payload ? payload.counters.length : "profiles" in payload ? payload.profiles.length : 1, status: accepted ? "accepted" : code ? "rejected" : "unavailable" }; latest.deliveries = [...(latest.deliveries ?? []).filter(validDelivery), delivery].slice(-MAX_DELIVERIES); await save(path, latest); diff --git a/loopx/control_plane/runtime/usage_statistics_cli.ts b/loopx/control_plane/runtime/usage_statistics_cli.ts index ef9356efdc..19d04fa7b0 100644 --- a/loopx/control_plane/runtime/usage_statistics_cli.ts +++ b/loopx/control_plane/runtime/usage_statistics_cli.ts @@ -1,4 +1,5 @@ -import { configure, inspect, observe } from "./usage_statistics.ts"; +import { configure, configureContext, inspect, observe } from "./usage_statistics.ts"; +import { profileFeature } from "./usage_statistics_installation_contract.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"; @@ -20,6 +21,7 @@ try { const ctx: Context = { env: process.env, version: String(request.facts.version), python: String(request.facts.python), channel: String(request.facts.channel) }; let result: unknown; if (request.action === "status") result = await inspect(request.path, ctx); + else if (request.action === "context") result = await configureContext(request.path, ctx, request.context); else if (request.action === "enable" || request.action === "disable" || request.action === "acknowledge") { result = await configure(request.path, ctx, request.action, request.notice); } else if (request.action === "start") { @@ -32,15 +34,16 @@ try { && typeof request.elapsed_ms === "number" && Number.isFinite(request.elapsed_ms) && request.elapsed_ms >= 0) { const feature = (FEATURES as readonly unknown[]).includes(request.feature) ? request.feature : "other"; if (typeof request.exit_code === "number" && Number.isInteger(request.exit_code)) { + const measuredFeature = profileFeature(request.feature); 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), + feature: measuredFeature, ...resultDiagnostic(measuredFeature, 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), - }); + }, profileFeature(request.feature)); } 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); + } as Counter, undefined, undefined, undefined, undefined, profileFeature(request.feature)); } } else throw new Error("usage_request_invalid"); process.stdout.write(JSON.stringify(result) + "\n"); diff --git a/loopx/control_plane/runtime/usage_statistics_codex.ts b/loopx/control_plane/runtime/usage_statistics_codex.ts index 274a2668a7..0b0cd4aa04 100644 --- a/loopx/control_plane/runtime/usage_statistics_codex.ts +++ b/loopx/control_plane/runtime/usage_statistics_codex.ts @@ -3,7 +3,7 @@ import { open } from "node:fs/promises"; import type { FileHandle } from "node:fs/promises"; import type { GoalObservation, Host } from "./usage_statistics_goal_contract.ts"; import { object } from "./usage_statistics_contract.ts"; -export type CodexCursor = { offset: number; inode: string; since: number; seen: number; skipping?: boolean; open?: { id: string; start: number; confirmed: number } }; +export type CodexCursor = { offset: number; inode: string; since: number; seen: number; skipping?: boolean; open?: { id: string; start: number } }; const BUDGET = 1024 * 1024; // The first line is not a short id record: Codex's recorder writes // `base_instructions` and the dynamic tool list into `session_meta`, and the @@ -13,7 +13,13 @@ const BUDGET = 1024 * 1024; const HEADER_CHUNK = 65536; const HEADER_BUDGET = 2 * 1024 * 1024; type HeaderLine = { line: string } | { error: "incomplete" | "too_large" }; -function timestamp(value: unknown): number { return typeof value === "string" ? Date.parse(value) : NaN; } +function timestamp(value: unknown): number { + // Codex provider envelopes use Unix seconds; older envelopes use ISO dates. + // The recorder timestamp is not a substitute for a supplied provider time. + const milliseconds = typeof value === "number" && Number.isFinite(value) && value >= 0 + ? Math.trunc(value * 1000) : typeof value === "string" ? Date.parse(value) : NaN; + return Number.isSafeInteger(milliseconds) && milliseconds >= 0 ? milliseconds : NaN; +} /** * Read the opening record whole, framed by LF, up to a fixed budget. * @@ -78,8 +84,9 @@ export async function readCodexTiming(path: string, thread: string, previous: Co const observations: GoalObservation[] = []; const emit = (start: number, end: number) => { start = Math.max(start, cursor.since); - if (Number.isSafeInteger(start) && Number.isSafeInteger(end) && end > start && end <= now + 1000 && end - start <= 7 * 86400000) + if (Number.isSafeInteger(start) && Number.isSafeInteger(end) && end > start && end <= now + 1000 && end - start <= 7 * 86400000) { observations.push({ key, start, end, measurement: "codex_turn", host }); + } }; for (const line of body.split("\n")) { if (!line) continue; @@ -90,7 +97,7 @@ export async function readCodexTiming(path: string, thread: string, previous: Co const id = typeof payload.turn_id === "string" && payload.turn_id.length <= 128 ? payload.turn_id : ""; if (payload.type === "task_started" && id) { const start = timestamp(payload.started_at ?? event.timestamp); - if (Number.isFinite(start)) cursor.open = { id, start, confirmed: Math.max(start, cursor.since) }; + if (Number.isFinite(start) && start <= now + 1000) cursor.open = { id, start }; } else if (["task_complete", "task_completed", "turn_aborted"].includes(String(payload.type)) && id) { const start = timestamp(payload.started_at); const end = timestamp(payload.completed_at ?? event.timestamp); @@ -98,13 +105,11 @@ export async function readCodexTiming(path: string, thread: string, previous: Co // the matching open Turn. Never pair an unrelated terminal by proximity. const matched = cursor.open?.id === id ? cursor.open : undefined; if (Number.isFinite(start)) emit(start, end); - else if (matched) emit(matched.confirmed, end); + else if (matched && payload.started_at == null) emit(matched.start, end); if (matched) cursor.open = undefined; - } else if (cursor.open && payload.type === "token_count") { - const confirmed = timestamp(event.timestamp); - emit(cursor.open.confirmed, confirmed); - if (Number.isFinite(confirmed)) cursor.open.confirmed = Math.max(cursor.open.confirmed, confirmed); } + // token_count timestamps describe recording, not provider execution. + // They can arrive after completed_at, so they cannot confirm a prefix. } return { cursor, observations }; } finally { await file.close(); } diff --git a/loopx/control_plane/runtime/usage_statistics_diagnostics.ts b/loopx/control_plane/runtime/usage_statistics_diagnostics.ts index 8ad0d841ee..99c20af6c7 100644 --- a/loopx/control_plane/runtime/usage_statistics_diagnostics.ts +++ b/loopx/control_plane/runtime/usage_statistics_diagnostics.ts @@ -4,6 +4,21 @@ 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_FEATURES = [...FEATURES, "heartbeat", "state", "agent", "memory", "capability", "maintenance"] as const; +export type DiagnosticFeature = typeof DIAGNOSTIC_FEATURES[number]; +/** Fixed parser command names, not argv values or user-supplied extension names. */ +export function diagnosticFeature(command: string): DiagnosticFeature { + if ((DIAGNOSTIC_FEATURES as readonly string[]).includes(command)) return command as DiagnosticFeature; + const families: Record = { + "heartbeat-prompt": "heartbeat", "refresh-state": "state", "checkpoint-context": "state", + "agent-context": "agent", "agent-capabilities": "agent", "agent-directory": "agent", "manager-inbox": "agent", + "reward-memory": "memory", "agent-turn-recall": "memory", "semantic-preference": "memory", + extension: "capability", "workflow-skills": "capability", "project-skill": "capability", + doctor: "maintenance", update: "maintenance", "migrate-local-state": "maintenance", + "serve-status": "chat", dashboard: "chat", + }; + return Object.hasOwn(families, command) ? families[command] : "other"; +} 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; @@ -16,7 +31,8 @@ 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 & { +export type Diagnostic = Omit & { + feature: DiagnosticFeature; 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]; @@ -30,7 +46,7 @@ export function diagnosticKey(value: Diagnostic): string { } 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) + if (!(DIAGNOSTIC_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) diff --git a/loopx/control_plane/runtime/usage_statistics_installation.ts b/loopx/control_plane/runtime/usage_statistics_installation.ts new file mode 100644 index 0000000000..8ec256ce6c --- /dev/null +++ b/loopx/control_plane/runtime/usage_statistics_installation.ts @@ -0,0 +1,104 @@ +/** Bounded machine-local daily interval union; uses the common consent lock. */ +import { readFile, chmod } from "node:fs/promises"; +import { atomicWriteJson } from "../effect_runtime_io.ts"; +import type { JsonObject } from "../effect_program.ts"; +import { object, MAX_COUNT } from "./usage_statistics_contract.ts"; +import { union } from "./usage_statistics_goals.ts"; +import { validGoalObservation, MEASUREMENTS } from "./usage_statistics_goal_contract.ts"; +import type { GoalObservation, Measurement } from "./usage_statistics_goal_contract.ts"; +import { INSTALLATION_SCHEMA, validInstallationDay } from "./usage_statistics_installation_contract.ts"; +import type { InstallationDay, InstallationUsage, ProfileFeature } from "./usage_statistics_installation_contract.ts"; + +const DAY = 86400000, MAX_INTERVALS = 1024; +type LocalDay = InstallationDay & { intervals: Partial>; attempted?: number }; +type LocalState = { generation: string; days: LocalDay[]; last_attempt?: number }; +function utc(ms: number) { return new Date(ms).toISOString().slice(0, 10); } +async function load(path: string, generation: string): Promise { + try { + const text = await readFile(path, "utf8"); + if (text.length > 1024 * 1024) throw new Error("installation_usage_too_large"); + const value: unknown = JSON.parse(text); + if (!object(value)) throw new Error("installation_usage_invalid"); + if (value.generation !== generation) return { generation, days: [] }; + if (!Array.isArray(value.days) || value.days.length > 8 + || (value.last_attempt !== undefined && (!Number.isSafeInteger(value.last_attempt) || Number(value.last_attempt) < 0)) + || value.days.some(raw => { + if (!object(raw)) return true; + const { intervals, attempted, ...profile } = raw; + if (!validInstallationDay(profile) || !object(intervals) + || (attempted !== undefined && (!Number.isSafeInteger(attempted) || Number(attempted) < 0 || Number(attempted) > profile.revision))) return true; + return Object.entries(intervals).some(([measurement, rows]) => { + if (!(MEASUREMENTS as readonly string[]).includes(measurement) || !Array.isArray(rows) || rows.length > MAX_INTERVALS) return true; + const start = Date.parse(profile.activity_day); + return rows.some((row, index) => !Array.isArray(row) || row.length !== 2 || !row.every(Number.isSafeInteger) + || row[0] < start || row[1] > start + DAY || row[0] > row[1] || (index > 0 && row[0] <= rows[index - 1][1])); + }); + }) || new Set(value.days.map(raw => raw.activity_day)).size !== value.days.length) throw new Error("installation_usage_invalid"); + return value as LocalState; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return { generation, days: [] }; + throw error; + } +} +function publicDay(row: LocalDay): InstallationDay { + const { intervals: _intervals, attempted: _attempted, ...profile } = row; + return profile; +} +export async function installationPreview(path: string, generation: string, installId: string): Promise { + const profiles = (await load(path, generation)).days.filter(row => row.attempted !== row.revision).map(publicDay); + return profiles.length ? { schema: INSTALLATION_SCHEMA, install_id: installId, profiles } : null; +} +export async function recordInstallation(path: string, generation: string, installId: string, now: number, + version: string, context: InstallationDay["context"], feature: ProfileFeature | undefined, + observations: GoalObservation[]): Promise { + const state = await load(path, generation); + const today = utc(now), oldest = Date.parse(today) - 7 * DAY; + state.days = state.days.filter(row => Date.parse(row.activity_day) >= oldest); + const getDay = (activityDay: string): LocalDay => { + let row = state.days.find(item => item.activity_day === activityDay); + if (!row) { + row = { activity_day: activityDay, version, context, revision: 1, cli: [], runtime: [], truncated: false, intervals: {} }; + state.days.push(row); + } + return row; + }; + // A daily profile's context/version is frozen on first observation, not relabeled on settings changes. + const current = getDay(today); + if (feature) { + const row = current.cli.find(item => item.feature === feature); + if (row && row.count < MAX_COUNT) row.count++; + else if (row) current.truncated = true; + else current.cli.push({ feature, count: 1 }); + current.revision++; + } + for (const observation of observations.filter(item => validGoalObservation(item, now))) { + const start = Math.max(observation.start, oldest); + for (let boundary = Date.parse(utc(start)); boundary < observation.end; boundary += DAY) { + const row = getDay(utc(boundary)); + const before = row.intervals[observation.measurement] ?? []; + const after = union(before, [Math.max(start, boundary), Math.min(observation.end, boundary + DAY)]); + if (JSON.stringify(after) === JSON.stringify(before)) continue; // replay never extends measured time + if (after.length > MAX_INTERVALS) row.truncated = true; + else { + row.intervals[observation.measurement] = after; + const minutes = Math.floor(after.reduce((sum, [a, b]) => sum + b - a, 0) / 60000); + const measured = row.runtime.find(item => item.measurement === observation.measurement); + if (measured) measured.observed_minutes = minutes; + else row.runtime.push({ measurement: observation.measurement, observed_minutes: minutes }); + } + row.revision++; + } + } + let payload: InstallationUsage | null = null; + if (state.days.some(row => row.cli.length > 0 || row.runtime.length > 0 || row.activity_day < today) + && (state.last_attempt === undefined || now - state.last_attempt >= 15 * 60000)) { + const pending = state.days.filter(row => row.attempted !== row.revision); + if (pending.length) { + payload = { schema: INSTALLATION_SCHEMA, install_id: installId, profiles: pending.map(publicDay) }; + for (const row of pending) row.attempted = row.revision; + state.last_attempt = now; // lossy claim, never retry or consume on settings reads + } + } + await atomicWriteJson(path, state as unknown as JsonObject); await chmod(path, 0o600); + return payload; +} diff --git a/loopx/control_plane/runtime/usage_statistics_installation_contract.ts b/loopx/control_plane/runtime/usage_statistics_installation_contract.ts new file mode 100644 index 0000000000..7079f1f334 --- /dev/null +++ b/loopx/control_plane/runtime/usage_statistics_installation_contract.ts @@ -0,0 +1,41 @@ +/** Installation-linked daily observations. Never uptime, billing or work authority. */ +import { object, validId, MAX_COUNT } from "./usage_statistics_contract.ts"; +import { CONTEXTS, DIAGNOSTIC_FEATURES, diagnosticFeature } from "./usage_statistics_diagnostics.ts"; +import { MEASUREMENTS } from "./usage_statistics_goal_contract.ts"; + +export const INSTALLATION_SCHEMA = "loopx_installation_usage_v1"; +export const PROFILE_FEATURES = DIAGNOSTIC_FEATURES; +export type ProfileFeature = typeof PROFILE_FEATURES[number]; +/** Reuse the diagnostic owner's fixed families; no parallel classification. */ +export { diagnosticFeature as profileFeature }; +export type InstallationDay = { + activity_day: string; version: string; context: typeof CONTEXTS[number]; revision: number; + cli: { feature: ProfileFeature; count: number }[]; + runtime: { measurement: typeof MEASUREMENTS[number]; observed_minutes: number }[]; + truncated: boolean; +}; +export type InstallationUsage = { schema: typeof INSTALLATION_SCHEMA; install_id: string; profiles: InstallationDay[] }; +function exact(value: Record, fields: string): boolean { return Object.keys(value).sort().join() === fields; } +export function validInstallationDay(value: unknown): value is InstallationDay { + if (!object(value) || !exact(value, "activity_day,cli,context,revision,runtime,truncated,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 + || typeof value.version !== "string" || !/^\d{1,3}\.\d{1,3}\.\d{1,4}$/.test(value.version) + || !(CONTEXTS as readonly unknown[]).includes(value.context) || typeof value.truncated !== "boolean" + || !Number.isSafeInteger(value.revision) || Number(value.revision) < 1 + || !Array.isArray(value.cli) || value.cli.length > PROFILE_FEATURES.length + || !Array.isArray(value.runtime) || value.runtime.length > MEASUREMENTS.length) return false; + const cli = new Set(), runtime = new Set(); + return value.cli.every(row => object(row) && exact(row, "count,feature") + && (PROFILE_FEATURES as readonly unknown[]).includes(row.feature) && !cli.has(row.feature) && !!cli.add(row.feature) + && Number.isInteger(row.count) && Number(row.count) >= 1 && Number(row.count) <= MAX_COUNT) + && value.runtime.every(row => object(row) && exact(row, "measurement,observed_minutes") + && (MEASUREMENTS as readonly unknown[]).includes(row.measurement) && !runtime.has(row.measurement) && !!runtime.add(row.measurement) + && Number.isInteger(row.observed_minutes) && Number(row.observed_minutes) >= 0 && Number(row.observed_minutes) <= 1440); +} +export function validInstallationUsage(value: unknown): value is InstallationUsage { + return object(value) && exact(value, "install_id,profiles,schema") && value.schema === INSTALLATION_SCHEMA + && validId(value.install_id) && Array.isArray(value.profiles) && value.profiles.length >= 1 && value.profiles.length <= 8 + && value.profiles.every(validInstallationDay) + && new Set(value.profiles.map(row => row.activity_day)).size === value.profiles.length; +} diff --git a/loopx/usage_ping.py b/loopx/usage_ping.py index 55ae73fd23..91d003a35d 100644 --- a/loopx/usage_ping.py +++ b/loopx/usage_ping.py @@ -17,7 +17,7 @@ STATE_FILENAME = "usage-ping.json" # Scheduling hint only; keep aligned with the TypeScript notice revision. -_NOTICE_VERSION = 5 +_NOTICE_VERSION = 6 _ENTRY = Path(__file__).parent / "control_plane/runtime/usage_statistics_cli.ts" _observation: ContextVar[dict[str, Any] | None] = ContextVar("usage_observation", default=None) From 225b0e3f259789ef37fb67b04a0b76bb02c51ee0 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Fri, 2 Oct 2026 22:43:51 +0800 Subject: [PATCH 2/4] docs(telemetry): explain runtime coverage and collector rollout Signed-off-by: huangruiteng --- apps/usage-collector/README.md | 13 +++ .../queries/installation-usage.sql | 35 ++++++ docs/reference/usage-ping.md | 103 +++++++++++++++++- docs/reference/usage-ping.zh-CN.md | 65 ++++++++++- 4 files changed, 207 insertions(+), 9 deletions(-) create mode 100644 apps/usage-collector/queries/installation-usage.sql diff --git a/apps/usage-collector/README.md b/apps/usage-collector/README.md index 45d1a9b030..accfa63158 100644 --- a/apps/usage-collector/README.md +++ b/apps/usage-collector/README.md @@ -10,6 +10,7 @@ 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` | Independent Goal/measurement/Host-day span/duration buckets, no identity; ≤16 KiB | +| `POST /v1/installation` | Installation-linked fixed daily CLI counts and independent interval-union minutes; ≤16 KiB; operator-only rows | | `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 | @@ -80,6 +81,18 @@ 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. +Before releasing notice-v6 clients, back up and apply additive migration +`0005-installation-usage.sql`, then deploy the Worker. `installation_usage` +keeps one snapshot per random installation/UTC activity day for 30 activity +days. Newer revisions replace counts, never add them; context and version +freeze per day. No historical profiles or runtime are backfilled. Rollback +retains the table, and all existing endpoint contracts remain supported. +Daily runtime is partial instrumented interval union, **not uptime**. Different +clocks cannot be added. [Operator read queries](queries/installation-usage.sql) +separate span, active days, fixed-family usage and measured minutes. +Only authorized D1/Access-protected operator surfaces may render linked rows; +do not add per-ID results to unauthenticated public stats. + 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 diff --git a/apps/usage-collector/queries/installation-usage.sql b/apps/usage-collector/queries/installation-usage.sql new file mode 100644 index 0000000000..18205881dd --- /dev/null +++ b/apps/usage-collector/queries/installation-usage.sql @@ -0,0 +1,35 @@ +-- Owner-only D1 read queries. Bind :from_day and :through_day (UTC dates), +-- including the published CI-exclusion cutover and latest partial day. +-- Missing historical profiles are UNKNOWN, never zero runtime or no usage. +-- Do not publish installation IDs or cross-dimensional small cells. + +-- Reporting state continuity, not people, work time or continuous uptime. +SELECT install_id, MIN(day) AS first_observed_day, MAX(day) AS last_observed_day, + COUNT(*) AS observed_days, + CAST(julianday(MAX(day)) - julianday(MIN(day)) AS INTEGER) + 1 AS calendar_span_days +FROM pings WHERE day BETWEEN :from_day AND :through_day +GROUP BY install_id; + +-- Fixed CLI distribution. Exclusion covers explicitly declared days only. +SELECT json_extract(j.value, '$.feature') AS feature, + COUNT(DISTINCT i.install_id) AS reporting_installations, + SUM(json_extract(j.value, '$.count')) AS observed_calls +FROM installation_usage i, json_each(i.cli) j +WHERE i.activity_day BETWEEN :from_day AND :through_day AND i.context != 'maintainer' +GROUP BY feature; + +-- Windowed installation runtime: deduplicated daily union, separate clocks. +-- Presence below one minute may report zero; absence must remain unknown. +SELECT i.install_id, json_extract(j.value, '$.measurement') AS measurement, + COUNT(DISTINCT i.activity_day) AS measured_days, + SUM(json_extract(j.value, '$.observed_minutes')) AS observed_minutes, + MAX(i.truncated) AS contains_truncated_day +FROM installation_usage i, json_each(i.runtime) j +WHERE i.activity_day BETWEEN :from_day AND :through_day AND i.context != 'maintainer' +GROUP BY i.install_id, measurement; + +-- Observe coverage before interpreting runtime/CLI as adoption. +SELECT context, COUNT(DISTINCT install_id) AS installations, + COUNT(*) AS profile_days, SUM(truncated) AS truncated_days +FROM installation_usage WHERE activity_day BETWEEN :from_day AND :through_day +GROUP BY context; diff --git a/docs/reference/usage-ping.md b/docs/reference/usage-ping.md index 65477bdbd9..af7ea48655 100644 --- a/docs/reference/usage-ping.md +++ b/docs/reference/usage-ping.md @@ -21,6 +21,80 @@ Lark has no separate switch and cannot override the machine owner's choice. ## Product questions and exact scope +### Device context and installation runtime (notice revision 6) + +Configure once for the machine, not once per Goal or terminal: + +```bash +loopx usage-ping context --context maintainer +loopx usage-ping status --format json +loopx usage-ping context --context unknown # remove the voluntary label +``` + +The Workspace device settings use the same TypeScript owner. Precedence is an +explicit `LOOPX_USAGE_CONTEXT` environment value, then the stored device label, +then `unknown`. Invalid environment values remain unknown; invalid settings are +rejected. Context configuration never enables collection, changes the random +installation ID, or grants work authority. Disable retains only that optional +label, clearing identity and measurement history. No generic repository `.env` +is auto-loaded. Existing ID-free diagnostic rows remain ID-free. + +The **new** `POST /v1/installation` contract, +`loopx_installation_usage_v1`, links the existing random installation ID to up +to eight UTC daily snapshots. Each has activity date, numeric version, +voluntary context, revision, fixed CLI family counts, independent observed +runtime minutes, and a truncation flag. The context and version freeze on that +day's first observation; later settings do not relabel that day or history. +Counts are capped at 10,000 per family. New diagnostics and profiles distinguish +`heartbeat`, `state`, `agent`, `memory`, `capability` and `maintenance` +families using fixed parsed command names; custom names remain `other`. +The legacy feature contract is unchanged. + +For each of `host_call`, `codex_turn` and `quota_cycle`, union actual +observation intervals **across all Goals and Hosts in the same installation**. +Split at UTC midnight and round each daily union down to whole minutes. +Parallel/nested overlap counts once; replay adds nothing; inter-interval idle +time adds nothing. Runtime rows are absent when no interval was observed; a +present row with zero minutes means less than one observed minute, not no work. +Do not add the three clocks: quota/Turn intervals can include pauses, approval +and tool waits, while direct Host calls use the existing short-checkpoint gap +rules. This measures partial instrumented elapsed time, **not process uptime**, +CPU time, completion, billable time or model thinking time. A process that only +heartbeats has unknown runtime, not 24 hours/day. + +Intervals remain local and bounded to 1,024 disjoint intervals per clock/day +and eight UTC days. Overflow is marked `truncated` rather than extrapolated. +Only whole-minute daily totals leave; no interval timestamps, Goal/Agent/Host +identity or transcripts. Full snapshots attempt at most every 15 minutes on +supported activity, not an autonomous timer. A quiet final partial day may +remain unsent; late observations can update the previous seven days. The +collector replaces only newer revisions for the exact installation/date and +same frozen context/version; duplicate or out-of-order transport never adds +counts. Client state loss without changing the ID can cause undercounting; +copied IDs can conflate machines. These are diagnostics, not a ledger. + +Installation profiles are a **new association**, not anonymous aggregate +data. Renewed disclosure fences old workers before it begins; old scope +buffers are discarded, not backfilled. `CI`, `DO_NOT_TRACK`, explicit disable +and consent-required policy keep their precedence. Disable removes local +`.installation` interval buffers; a request already started cannot be recalled. +Collector profile retention is 30 activity days, independent of 400-day +heartbeat retention. No public endpoint exposes per-ID profiles. + +Operators may sum daily observed minutes within a specified retained window +for each installation/clock, and report active days, calendar span and runtime +**separately**. This is observed time in that window, not lifetime uptime. +See [fixed read queries](../../apps/usage-collector/queries/installation-usage.sql). +Do not divide old anonymous CLI counts by ID totals or fill missing days with +zero. Excluding `maintainer` applies only to explicit future profile labels; +`unknown` does not prove external or personal use. + +Deploy additive migration `0005-installation-usage.sql` and the Worker before +releasing notice-v6 clients. Validate against a disposable database/collector, +never by writing synthetic events to production. Worker rollback keeps the +additive table; older clients/endpoints remain supported. Merging does not +deploy the Worker or change installed clients. + | Question | Evidence | Limit | |---|---|---| | Which versions/platforms need support? | Daily version, OS, CPU architecture, Python minor and install channel | Only reporting installations | @@ -51,7 +125,7 @@ OS is `darwin|linux|windows|other`, CPU is `x64|arm64|x86|other`, and channel is `pip|local_release|source|unknown`. Version accepts only numeric major.minor.patch; a custom version containing a private suffix is not sent. -**Current CLI diagnostics** (`POST /v1/aggregate`, notice revision 5): +**ID-free CLI diagnostics** (`POST /v1/aggregate`, introduced in 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"}]} @@ -59,12 +133,14 @@ a custom version containing a private suffix is not sent. 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`: +Deployment context is optional self-report via the shared device setting or +`LOOPX_USAGE_CONTEXT` (environment takes precedence): `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. +on maintainer devices/processes to distinguish their **future diagnostics and +daily installation profiles** in operator analysis. It does not label the +heartbeat table 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 @@ -200,7 +276,7 @@ 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 +readable. The current notice revision 6 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. @@ -291,6 +367,21 @@ up on a later observation; this is not a global session watcher. Missing or ambiguous bindings and unavailable files leave the common quota cycle working. No historical backfill or extrapolation of crashed sessions occurs. +Codex provider `started_at`/`completed_at` values accept Unix seconds or legacy +ISO dates. Explicit provider times take precedence over the recorder timestamp; +without an explicit start, a terminal must match the observed open Turn. Only +completed or aborted Turns produce intervals. `token_count` record timestamps +may be delayed beyond the provider's completion, so they no longer extend an +unfinished Turn. An ongoing native session can report its completed Turns, but +its current unfinished Turn remains unknown until a terminal is read. This +corrects prior prefix timing; it does not rewrite previously sent aggregates. + +These populations have no enforced size ordering. A quota cycle commonly +encloses several Turns and waits; a managed Host call can enclose a Turn plus +startup/cleanup. Native sessions, incomplete cycles and coverage gaps can reverse +their observed totals. Compare each clock over the same window and inspect +coverage rather than assuming `host_call <= codex_turn <= quota_cycle`. + Within each Goal/measurement/Host series, **span** is first to most recent observed activity, including pauses; **duration** is the union of observed intervals. Parallel or nested overlap counts once within that series. Both stop @@ -333,7 +424,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 the current notice version 5; an existing explicit disable persists. +scope requires the current notice version 6; 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 91c915d828..08bffeefcb 100644 --- a/docs/reference/usage-ping.zh-CN.md +++ b/docs/reference/usage-ping.zh-CN.md @@ -18,6 +18,54 @@ loopx usage-ping enable # 阅读告知后明确开启 ## 能回答什么 +### 设备用途与安装级运行时长(告知版本 6) + +一次设置即可跨 Goal 和终端持久使用,不自动加载仓库里的任意 `.env`: + +```bash +loopx usage-ping context --context maintainer +loopx usage-ping status --format json +loopx usage-ping context --context unknown +``` + +Workspace 设备设置复用同一个 TypeScript owner。优先级是显式 +`LOOPX_USAGE_CONTEXT` 环境变量 → 保存的设备标签 → `unknown`。非法环境值保持 +未知,非法设置拒绝写入。设置用途不等于开启统计,不更换安装 ID、不授予工作权限。 +关闭统计清除标识与测量记录,但保留自愿用途标签。 + +新增 `POST /v1/installation`,契约 `loopx_installation_usage_v1`。 +同一个随机安装 ID 会与每日固定 CLI 功能计数、数字版本、UTC 活动日期、自愿环境 +标签、已观测运行分钟、快照 revision 和截断标志关联。每天首次观察后固定该日 +标签和版本,后续配置不能反向改写历史。这是**新增关联**,不是匿名汇总。 +CLI 诊断及安装概要增加 `heartbeat|state|agent|memory|capability|maintenance` +固定分类;来自解析器命令名,不上传参数或自定义插件名,旧汇总契约不变。 + +对 `host_call`、`codex_turn`、`quota_cycle` 分别取**同一安装所有 Goal/Host** +实际区间的并集,跨 UTC 午夜拆分,每天向下取整为分钟。并行和嵌套重叠只算一次, +重放不增加时长,区间之间的空闲不增加时长。没有区间时不填 runtime 行;已有行的 +零分钟表示观测不足一分钟,不是没有工作。三种口径重叠,不能相加;轮次与推进 +周期可包含暂停、工具和审批等待,直接 Host 调用沿用短检查点的休眠间隔排除规则。 +这不是机器在线、进程 uptime、CPU、任务完成、模型思考或计费用时。 + +本机最多保留八个 UTC 日期,每个日期/口径最多 1,024 个离散区间,溢出明确标记 +`truncated`,不外推。只上传每日整分钟,不上传区间时间戳、Goal/Agent/Host +身份或对话。完整快照由活动触发,至少间隔 15 分钟,不增加常驻计时器;安静结束的 +当天末段可能未发出,前七天的迟到区间可以修正。收集端按安装/日期只替换更新 +revision,并保留首次标签和版本,传输重试或乱序不会累加。状态损坏或丢失可能 +漏计;复制标识可能合并多台机器,不能作为账本。 + +扩大的范围必须重新告知,旧 worker 被 generation 隔离,旧缓冲丢弃,不回填历史。 +CI、请勿追踪、明确关闭与需要明确同意的策略继续生效。关闭删除本机 +`.installation` 区间缓冲,已经开始的网络请求不能撤回。收集端保留 30 个活动日, +和 400 天心跳保留期独立。公共接口不暴露每个 ID 的概要。 + +维护者可按**实际保留窗口**统计每个安装/口径的已观测分钟,同时分开看活跃日、 +日历跨度和 CLI 功能分布,不称为终身运行时长。缺失日期不可补零,旧匿名 CLI +计数不可按安装数分摊。`maintainer` 排除只适用于以后明确标注的概要;未知不等于 +外部用户或个人用户。[固定只读查询](../../apps/usage-collector/queries/installation-usage.sql)。 +先在隔离数据库验证,再应用增量迁移 `0005-installation-usage.sql` 并部署 Worker, +然后发布告知版本 6 客户端;合并代码不等于已经部署,回滚保留新增表。 + - 每日版本、系统、CPU 架构、Python 小版本和安装渠道:哪些环境需要优先维护。 - 随机安装 ID 跨日心跳:持久机器状态目录的活跃与成熟的 1/7/30 天回访,不是 用户或组织数;删除状态或关闭后重开可能计为新安装。 @@ -46,17 +94,17 @@ ID 随机生成,属于持久机器状态目录,不绑定账号、不从硬 安装渠道只允许 `pip|local_release|source|unknown`。版本只接受数字三段式,包含 自定义后缀的版本不会上传。 -当前 CLI 诊断 `POST /v1/aggregate`(告知版本 5): +无 ID 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` 自愿声明: +回执支持的生命周期信号。环境类型由设备设置或优先级更高的 `LOOPX_USAGE_CONTEXT` 自愿声明: `unknown`(默认)、`personal`、`shared_service`、`ephemeral`、`organization_managed`、 `maintainer`;无效值成为 `unknown`,不猜企业、人数或机器拓扑,不接受公司名称。 -维护者可声明 `maintainer`,分开**未来诊断计数**;不会标注心跳,也无法追溯识别旧 +维护者可声明 `maintainer`,分开**未来诊断计数和每日安装概要**;不会标注心跳表,也无法追溯识别旧 无 ID 汇总。需要排除整机所有采集时仍用统一关闭开关。 收集器另加接收日期,仅接收此前七天至当天的活动日期,拒绝过期或未来数据。 @@ -223,6 +271,17 @@ Codex 发现使用既有 Goal/agent/task 绑定和所选 `CODEX_HOME` 的只读 原始会话内容不上传。spend 后才写入的结束事件需要等下次观测,并非全局实时监听。 绑定缺失、歧义或文件不可用不会阻断通用周期统计;不回填历史,不外推崩溃后的时间。 +Codex 的 `started_at`/`completed_at` 支持 Unix 秒数和旧版 ISO 日期。 +优先使用明确的 provider 时间;缺少起点时,结束事件必须匹配已观测的同一 Turn。 +仅已结束或中止的 Turn 生成区间:`token_count` 的记录时间可能晚于实际完成, +不再用于外推未结束 Turn。本机原生会话可统计已完成的轮次,当前未结束轮次须 +等结束事件被读取后计入。这修正了原来的前缀计时,不改写已发送的历史聚合。 + +三种口径没有强制大小关系。quota 周期通常包含多轮工作和等待;受管 Host 调用 +可能比对应 Codex Turn 多出启动、收尾时间。原生会话、未结算周期和观测缺口 +会改变总量关系。应在同一窗口分别核验覆盖,不能假定 +`host_call <= codex_turn <= quota_cycle`,也不能将三者相加。 + 每个 Goal/口径/Host 分别计算 **span**(首次至最近观测活动,包含中间暂停)和 **duration**(已观测区间的并集)。同口径同 Host 内并行重叠只计一次;没有新证据 就不增长。它们不是 Goal 年龄、CPU 用时、完成证据或计费时长。Host 只允许固定枚举 From e1d4ca7582c9b20e676ab685178e59d91da04a08 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Fri, 2 Oct 2026 22:43:51 +0800 Subject: [PATCH 3/4] test(telemetry): qualify interval union and real timing sources Signed-off-by: huangruiteng --- apps/usage-collector/test/collector.test.mjs | 81 +++++++++ .../control_plane_ts/usage_statistics.test.ts | 11 +- .../usage_statistics_delivery.test.ts | 2 +- .../usage_statistics_installation.test.ts | 157 ++++++++++++++++++ .../usage_statistics_sources.test.ts | 50 +++++- tests/test_usage_goal.py | 63 ++++++- tests/test_usage_ping.py | 65 +++++++- 7 files changed, 414 insertions(+), 15 deletions(-) create mode 100644 tests/control_plane_ts/usage_statistics_installation.test.ts diff --git a/apps/usage-collector/test/collector.test.mjs b/apps/usage-collector/test/collector.test.mjs index d787ff2fe7..20e42e2ae5 100644 --- a/apps/usage-collector/test/collector.test.mjs +++ b/apps/usage-collector/test/collector.test.mjs @@ -3,6 +3,11 @@ import assert from "node:assert/strict"; import { readFileSync } from "node:fs"; import { DatabaseSync } from "node:sqlite"; import test from "node:test"; +import { createServer } from "node:http"; +import { mkdtemp, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { configure, configureContext, observe } from "../../../loopx/control_plane/runtime/usage_statistics.ts"; import { MAX_BODY_BYTES, handle, purge, suppressSmall, validatePing } from "../src/collector.js"; @@ -22,6 +27,7 @@ function d1() { const raw = { all: (sql) => db.prepare(sql).all().map(plain), get: (sql) => plain(db.prepare(sql).get()), + query: (sql, params) => db.prepare(sql).all(params).map(plain), }; return { raw, @@ -133,6 +139,81 @@ test("unknown paths are 404", async () => { assert.equal((await handle(new Request("https://collector.example/"), d1())).status, 404); }); +test("installation profiles replace exact days idempotently, reject private fields and never expose IDs publicly", async () => { + const db = d1(); + const profile = { activity_day: "2026-10-01", version: "1.2.4", context: "maintainer", + revision: 1, cli: [{ feature: "memory", count: 2 }], + runtime: [{ measurement: "codex_turn", observed_minutes: 120 }], truncated: false }; + const payload = { schema: "loopx_installation_usage_v1", install_id: id(1), profiles: [profile] }; + const request = value => new Request("https://collector.example/v1/installation", { + method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(value), + }); + for (let i = 0; i < 2; i++) assert.equal((await handle(request(payload), db, at("2026-10-02"))).status, 204); + assert.equal(db.raw.get("SELECT COUNT(*) AS n FROM installation_usage").n, 1); + assert.equal(JSON.parse(db.raw.get("SELECT cli FROM installation_usage").cli)[0].count, 2); + const newer = { ...profile, revision: 3, cli: [{ feature: "memory", count: 4 }] }; + await handle(request({ ...payload, profiles: [newer] }), db, at("2026-10-02")); + await handle(request(payload), db, at("2026-10-02")); + assert.equal(JSON.parse(db.raw.get("SELECT cli FROM installation_usage").cli)[0].count, 4); + await handle(request({ ...payload, profiles: [{ ...newer, revision: 4, context: "personal" }] }), db, at("2026-10-02")); + assert.equal(db.raw.get("SELECT context FROM installation_usage").context, "maintainer", "historical context never relabeled"); + assert.equal((await handle(request({ ...payload, prompt: "private" }), db, at("2026-10-02"))).status, 400); + assert.equal((await handle(request({ ...payload, profiles: [{ ...profile, activity_day: "2026-09-01" }] }), db, at("2026-10-02"))).status, 400); + assert.equal((await handle(new Request("https://collector.example/v1/installation"), db)).status, 405); + const stats = await (await handle(new Request("https://collector.example/v0/stats"), db, at("2026-10-02"))).json(); + assert.equal(JSON.stringify(stats).includes(id(1)), false); + await purge(db, "2026-10-30"); + assert.equal(db.raw.get("SELECT COUNT(*) AS n FROM installation_usage").n, 1, "30 activity days including today are retained"); + await purge(db, "2026-10-31"); + assert.equal(db.raw.get("SELECT COUNT(*) AS n FROM installation_usage").n, 0); +}); + +test("installation migration is additive and preserves the existing ping history", () => { + const db = new DatabaseSync(":memory:"); + db.exec("CREATE TABLE pings(day TEXT, install_id TEXT); INSERT INTO pings VALUES ('2026-09-30', 'legacy');"); + const migration = readFileSync(new URL("../migrations/0005-installation-usage.sql", import.meta.url), "utf8"); + db.exec(migration); db.exec(migration); + assert.equal(db.prepare("SELECT COUNT(*) AS n FROM pings").get().n, 1); + assert.equal(db.prepare("SELECT COUNT(*) AS n FROM installation_usage").get().n, 0); +}); + +test("real client HTTP to SQLite preserves interval union and fixed operator query semantics", async t => { + const db = d1(), now = at("2026-10-01"); + const server = createServer(async (request, response) => { + const chunks = []; + for await (const chunk of request) chunks.push(chunk); + const result = await handle(new Request("http://collector.example" + request.url, { + method: request.method, headers: request.headers, body: Buffer.concat(chunks), + }), db, now); + response.writeHead(result.status); response.end(await result.text()); + }); + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + const root = await mkdtemp(join(tmpdir(), "loopx-real-collector-")), path = join(root, "usage.json"); + t.after(async () => { server.closeAllConnections(); server.close(); await rm(root, { recursive: true, force: true }); }); + const ctx = { env: { LOOPX_USAGE_PING_ENDPOINT: "http://127.0.0.1:" + server.address().port + "/v1/ping" }, + version: "1.2.4", python: "3.13", channel: "source", now }; + await configureContext(path, ctx, "personal"); await configure(path, ctx, "enable"); + const generation = JSON.parse(await readFile(path, "utf8")).generation; + const interval = { key: "a".repeat(64), measurement: "codex_turn", host: "codex_app", + start: now.getTime() - 60 * 60000, end: now.getTime() }; + await observe(path, ctx, generation, { feature: "turn", count: 1, outcome: "ok", error: "none", duration: "lt_1s" }, undefined, interval); + const row = db.raw.get("SELECT * FROM installation_usage"); + assert.equal(row.context, "personal"); + assert.deepEqual(JSON.parse(row.runtime), [{ measurement: "codex_turn", observed_minutes: 60 }]); + // Real SQL consumer: independent mathematical oracle is one 60-minute interval, not Goal span. + const source = readFileSync(new URL("../queries/installation-usage.sql", import.meta.url), "utf8"); + const queries = source.replace(/--[^\n]*/g, "").split(";").filter(sql => sql.trim()); + const bindings = { from_day: "2026-10-01", through_day: "2026-10-02" }; + for (const sql of queries) { + // Query uses named bind parameters; qualify against real SQLite. + assert.equal(db.raw.query(sql, bindings).length, 1); + } + const runtime = db.raw.all( + "SELECT SUM(json_extract(j.value, '$.observed_minutes')) AS minutes FROM installation_usage i, json_each(i.runtime) j", + ); + assert.equal(runtime[0].minutes, 60); +}); + test("v1 heartbeat has architecture; aggregates reject identifiers and store only counters", async () => { const db = d1(); diff --git a/tests/control_plane_ts/usage_statistics.test.ts b/tests/control_plane_ts/usage_statistics.test.ts index 714159350f..a4917b9133 100644 --- a/tests/control_plane_ts/usage_statistics.test.ts +++ b/tests/control_plane_ts/usage_statistics.test.ts @@ -96,7 +96,7 @@ test("one heartbeat per UTC day; same-day aggregates stay separate and identifie const { path, state } = await fixture(t); const ctx = context(); await configure(path, ctx, "enable"); const generation = (await state()).generation; const sent: { url: string; payload: unknown }[] = []; - const post: Post = async (url, payload) => { sent.push({ url, payload }); return 204; }; + const post: Post = async (url, payload) => { if (!url.endsWith("/installation")) sent.push({ url, payload }); return 204; }; await observe(path, ctx, generation, row, post); await observe(path, ctx, generation, row, post); assert.equal(sent.length, 2); assert.ok(validPing(sent[0].payload)); assert.deepEqual(sent[1].payload, { schema: AGGREGATE_SCHEMA, counters: [row] }); @@ -161,7 +161,7 @@ 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"]); + assert.deepEqual((await inspect(path, ctx)).delivery_history.map(row => row.status), ["unavailable", "unavailable", "unavailable"]); }); test("local delivery history is bounded, content-free and cleared by disable", async t => { @@ -219,7 +219,7 @@ test("real HTTP sender does not follow redirects to another recipient", async t ctx.env.LOOPX_USAGE_PING_ENDPOINT = `http://127.0.0.1:${(server.address() as { port: number }).port}/v1/ping`; await configure(path, ctx, "enable"); assert.equal((await observe(path, ctx, (await state()).generation, row)).sent, false); - assert.equal(requests, 2, "heartbeat and aggregate must each stop at the redirect"); + assert.equal(requests, 3, "heartbeat, aggregate and installation profile each stop at the redirect"); }); test("startup heartbeat does not invent a successful command result", async t => { @@ -231,7 +231,7 @@ test("startup heartbeat does not invent a successful command result", async t => assert.equal((await inspect(path, ctx)).aggregate_preview, null); }); -test("a real unresponsive collector aborts both requests without retaining error details", { timeout: 30_000 }, async t => { +test("a real unresponsive collector aborts all three channels without retaining error details", { timeout: 30_000 }, async t => { const server = createServer(() => {}); await new Promise(r => server.listen(0, "127.0.0.1", r)); t.after(() => { server.closeAllConnections(); server.close(); }); @@ -250,7 +250,7 @@ test("a real unresponsive collector aborts both requests without retaining error return signal; }); assert.equal((await observe(path, ctx, (await state()).generation, row)).sent, false); - assert.deepEqual(deadlines, [3000, 3000]); + assert.deepEqual(deadlines, [3000, 3000, 3000]); assert.ok(signals.every(signal => signal.aborted && signal.reason.name === "TimeoutError")); const saved = await state(); assert.doesNotMatch(JSON.stringify(saved), /TimeoutError|aborted due to timeout/); @@ -258,5 +258,6 @@ test("a real unresponsive collector aborts both requests without retaining error assert.deepEqual(saved.deliveries, [ { day: "2026-09-26", channel: "heartbeat", rows: 1, status: "unavailable" }, { day: "2026-09-26", channel: "cli", rows: 1, status: "unavailable" }, + { day: "2026-09-26", channel: "installation", rows: 1, status: "unavailable" }, ]); }); diff --git a/tests/control_plane_ts/usage_statistics_delivery.test.ts b/tests/control_plane_ts/usage_statistics_delivery.test.ts index 85bc86e9d1..15d5ff7c0b 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, 5); + assert.equal(status.notice.version, 6); 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_installation.test.ts b/tests/control_plane_ts/usage_statistics_installation.test.ts new file mode 100644 index 0000000000..0d451ed0cb --- /dev/null +++ b/tests/control_plane_ts/usage_statistics_installation.test.ts @@ -0,0 +1,157 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { mkdtemp, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { configure, configureContext, inspect, observe } from "../../loopx/control_plane/runtime/usage_statistics.ts"; +import type { Context, Post } from "../../loopx/control_plane/runtime/usage_statistics.ts"; +import { recordInstallation, installationPreview } from "../../loopx/control_plane/runtime/usage_statistics_installation.ts"; +import { validInstallationUsage, profileFeature } from "../../loopx/control_plane/runtime/usage_statistics_installation_contract.ts"; +import type { GoalObservation } from "../../loopx/control_plane/runtime/usage_statistics_goal_contract.ts"; +import { acquireFileMutationLock, releaseFileMutationLock } from "../../loopx/control_plane/effect_runtime_io.ts"; +const id = "00000000-0000-4000-8000-000000000001"; +const at = (value: string) => Date.parse("2026-10-" + value + "Z"); +const observation = (start: number, end: number, key = "a", measurement: GoalObservation["measurement"] = "codex_turn"): GoalObservation => + ({ key: key.repeat(64), start, end, measurement, host: "codex_app" }); +async function fixture(t: test.TestContext) { + const root = await mkdtemp(join(tmpdir(), "loopx-installation-")); + t.after(() => rm(root, { recursive: true, force: true })); + return join(root, "usage.json"); +} +test("device context is shared, environment wins, and labeling never enables or rotates identity", async t => { + const path = await fixture(t); + const ctx: Context = { env: {}, version: "1.2.4", python: "3.13", channel: "source" }; + await configureContext(path, ctx, "maintainer"); + assert.equal((await inspect(path, ctx)).sending, false); + assert.equal((await inspect(path, ctx)).effective_context, "maintainer"); + assert.equal((await inspect(path, ctx)).context_source, "device"); + assert.equal((await inspect(path, { ...ctx, env: { LOOPX_USAGE_CONTEXT: "personal" } })).effective_context, "personal"); + assert.equal((await inspect(path, { ...ctx, env: { LOOPX_USAGE_CONTEXT: "invalid" } })).effective_context, "unknown"); + await assert.rejects(configureContext(path, ctx, "company-name")); + await configure(path, ctx, "enable"); + const original = (await inspect(path, ctx)).next_payload?.install_id; + await configureContext(path, ctx, "personal"); + assert.equal((await inspect(path, ctx)).next_payload?.install_id, original); + await configure(path, ctx, "disable"); + await configureContext(path, ctx, "maintainer"); + assert.equal((await inspect(path, ctx)).blocked_by, "disabled"); + assert.equal((await inspect(path, ctx)).next_payload, null); +}); +test("installation union crosses Goals and midnight, excludes idle gaps and keeps clocks separate", async t => { + const path = await fixture(t), now = at("02T01:00:00"); + await recordInstallation(path, id, id, now, "1.2.4", "personal", "turn", [ + observation(at("01T23:30:00"), at("02T00:30:00")), + observation(at("02T00:00:00"), at("02T00:45:00"), "b"), + observation(at("02T00:50:00"), at("02T00:51:30"), "c", "host_call"), + ]); + const saved = JSON.parse(await readFile(path, "utf8")); + assert.deepEqual(saved.days.find((row: {activity_day: string}) => row.activity_day === "2026-10-01").runtime, + [{ measurement: "codex_turn", observed_minutes: 30 }]); + assert.deepEqual(saved.days.find((row: {activity_day: string}) => row.activity_day === "2026-10-02").runtime, + [{ measurement: "codex_turn", observed_minutes: 45 }, { measurement: "host_call", observed_minutes: 1 }]); + const replay = await recordInstallation(path, id, id, now + 16 * 60000, "1.2.4", "maintainer", undefined, [ + observation(at("01T23:30:00"), at("02T00:30:00")), observation(at("02T00:00:00"), at("02T00:45:00"), "b"), + ]); + assert.equal(replay, null); + assert.equal((await installationPreview(path, id, id)), null); + assert.ok(saved.days.every((row: {context: string}) => row.context === "personal"), "history is not relabeled"); +}); + +test("each daily clock agrees with an independent endpoint sweep, including sub-minute fragments", async t => { + const observations: GoalObservation[] = []; + const midnight=at("02T00:00:00"), now=midnight+600000; + for(const measurement of ["host_call","codex_turn","quota_cycle"] as const) { + observations.push(observation(midnight-60000,midnight+60000,"a",measurement)); + observations.push(observation(midnight-30000,midnight+30000,"b",measurement)); + observations.push(observation(midnight+120000,midnight+149000,"c",measurement)); + observations.push(observation(midnight+180000,midnight+211000,"d",measurement)); + } + // Unlike the product's interval-merging algorithm, this oracle sweeps signed + // endpoints. Duration is integrated while population > 0, then rounded once. + const expected = (measurement: GoalObservation["measurement"], day: number) => { + const points = observations.filter(o=>o.measurement===measurement).flatMap(o=> { + const start=Math.max(o.start,day), end=Math.min(o.end,day+86400000); + return end>start ? [[start,1],[end,-1]] as [number,number][] : []; + }).sort(([a],[b])=>a-b); + let active=0, previous=day, elapsed=0; + for(const [time,delta] of points) { if(active>0) elapsed+=time-previous; active+=delta; previous=time; } + return Math.floor(elapsed/60000); + }; + for(const ordered of [observations,[...observations].reverse()]) { + const path=await fixture(t); + await recordInstallation(path,id,id,now,"1.2.4","maintainer",undefined,ordered); + const before=JSON.parse(await readFile(path,"utf8")); + await recordInstallation(path,id,id,now,"1.2.4","maintainer",undefined,ordered); + assert.deepEqual(JSON.parse(await readFile(path,"utf8")),before,"runtime replay changes neither totals nor revision"); + for(const row of before.days) { + for(const clock of row.runtime) { + assert.equal(clock.observed_minutes,expected(clock.measurement,Date.parse(row.activity_day))); + } + } + assert.ok(before.days.find((row:{activity_day:string})=>row.activity_day==="2026-10-02").runtime.every((clock:{observed_minutes:number})=>clock.observed_minutes===2)); + } +}); +test("restart, partial updates and expiration do not extrapolate missing runtime", async t => { + const path = await fixture(t), now = at("01T12:00:00"); + const first = await recordInstallation(path, id, id, now, "1.2.4", "maintainer", "memory", []); + assert.ok(validInstallationUsage(first)); + assert.deepEqual(first!.profiles[0].runtime, [], "no observed runtime is unknown, not zero uptime"); + await recordInstallation(path, id, id, now + 60000, "1.2.4", "personal", "state", []); + const later = await recordInstallation(path, id, id, now + 16 * 60000, "1.2.4", "personal", undefined, []); + assert.ok(validInstallationUsage(later)); + assert.deepEqual(later!.profiles[0].cli, [{ feature: "memory", count: 1 }, { feature: "state", count: 1 }]); + assert.equal(later!.profiles[0].context, "maintainer"); + const expired = await recordInstallation(path, id, id, at("10T12:00:00"), "1.2.4", "personal", "todo", []); + assert.equal(expired!.profiles.length, 1); + assert.equal(expired!.profiles[0].context, "personal"); +}); +test("expanded fixed families reject custom names and transport rejects content or impossible daily time", () => { + assert.equal(profileFeature("refresh-state"), "state"); + assert.equal(profileFeature("heartbeat-prompt"), "heartbeat"); + assert.equal(profileFeature("reward-memory"), "memory"); + assert.equal(profileFeature("my-private-plugin"), "other"); + const profile = { activity_day: "2026-10-01", version: "1.2.4", context: "unknown", revision: 1, + cli: [], runtime: [{ measurement: "quota_cycle", observed_minutes: 60 }], truncated: false }; + const payload = { schema: "loopx_installation_usage_v1", install_id: id, profiles: [profile] }; + assert.ok(validInstallationUsage(payload)); + assert.equal(validInstallationUsage({ ...payload, goal: "private" }), false); + assert.equal(validInstallationUsage({ ...payload, profiles: [{ ...profile, runtime: [{ measurement: "quota_cycle", observed_minutes: 1441 }] }] }), false); +}); +test("suppression and disable fence all new channels and erase local interval history", async t => { + const path = await fixture(t); + const ctx: Context = { env: {}, version: "1.2.4", python: "3.13", channel: "source", now: new Date(at("01T12:00:00")) }; + await configureContext(path, ctx, "maintainer"); await configure(path, ctx, "enable"); + const generation = JSON.parse(await readFile(path, "utf8")).generation; + const sent: unknown[] = [], post: Post = async (_url, payload) => { sent.push(payload); return 204; }; + await observe(path, ctx, generation, null, post, observation(at("01T11:00:00"), at("01T12:00:00"))); + assert.ok(sent.some(payload => validInstallationUsage(payload) && payload.profiles[0].context === "maintainer")); + const stateBefore = await readFile(path + ".installation", "utf8"); + await observe(path, { ...ctx, env: { CI: "true" } }, generation, null, async () => assert.fail("CI must not send")); + assert.equal(await readFile(path + ".installation", "utf8"), stateBefore); + await configure(path, ctx, "disable"); + await assert.rejects(readFile(path + ".installation"), /ENOENT/); + await observe(path, ctx, generation, null, async () => assert.fail("stale observer must not send")); + assert.equal((await inspect(path, ctx)).installation_preview, null); + assert.equal((await inspect(path, ctx)).stored_context, "maintainer"); +}); +test("short-lived CLI completion survives a competing detached startup observation", async t => { + const path = await fixture(t); + const ctx: Context = { env: {}, version: "1.2.4", python: "3.13", channel: "source" }; + await configure(path, ctx, "enable"); + const generation = JSON.parse(await readFile(path, "utf8")).generation; + const sent: Parameters[1][] = []; + const send: Post = async (_url, payload) => { sent.push(payload); return 204; }; + // The daily startup worker and command-completion worker share this local + // telemetry lock. A briefly occupied lock must not lose the first CLI count. + const lock = await acquireFileMutationLock(path); + const completed = observe(path, ctx, generation, null, send, undefined, undefined, undefined, "version") + .then(value => ({ value }), error => ({ error })); + await new Promise(resolve => setTimeout(resolve, 40)); + await releaseFileMutationLock(path, lock.token); + const result = await completed; + assert.ok("value" in result, "detached observer must wait for a short telemetry-only lock hold"); + await observe(path, ctx, generation, null, send); + const profile = sent.find(validInstallationUsage); + assert.deepEqual(profile?.profiles[0].cli, [{ feature: "version", count: 1 }]); + assert.equal(sent.filter(row => row.schema === "loopx_usage_ping_v1").length, 1); +}); diff --git a/tests/control_plane_ts/usage_statistics_sources.test.ts b/tests/control_plane_ts/usage_statistics_sources.test.ts index 87fa45b618..78db4c7704 100644 --- a/tests/control_plane_ts/usage_statistics_sources.test.ts +++ b/tests/control_plane_ts/usage_statistics_sources.test.ts @@ -17,6 +17,54 @@ async function fixture(t: test.TestContext) { } const event = (type: string, time: number, fields = {}) => JSON.stringify({ type: "event_msg", timestamp: new Date(time).toISOString(), payload: { type, ...fields } }) + "\n"; +test("numeric provider seconds measure the Turn, not a delayed recorder timestamp", async t => { + const path=join(await fixture(t),"numeric.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + await appendFile(path,event("task_started",at+2000,{turn_id:"numeric",started_at:at/1000})); + result=await readCodexTiming(path,"thread",result.cursor,at+2001,key,"codex_app"); + assert.equal(result.cursor.open?.start,at); + await appendFile(path,event("token_count",at+394000)+event("task_complete",at+394000,{turn_id:"numeric",started_at:at/1000,completed_at:at/1000+329})); + result=await readCodexTiming(path,"thread",result.cursor,at+394001,key,"codex_app"); + assert.equal(result.observations.length,1); + assert.equal(result.observations[0].end-result.observations[0].start,329000); + assert.equal(result.cursor.open,undefined); + assert.deepEqual((await readCodexTiming(path,"thread",result.cursor,at+400000,key,"codex_app")).observations,[]); +}); + +test("numeric starts support legacy terminals without treating token logs or idle gaps as work", async t => { + const path=join(await fixture(t),"numeric-prefix.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + await appendFile(path,event("task_started",at,{turn_id:"first",started_at:at/1000})+event("token_count",at+10000)); + result=await readCodexTiming(path,"thread",result.cursor,at+10001,key,"codex_app"); + assert.deepEqual(result.observations,[],"an unfinished Turn has no terminal duration yet"); + await appendFile(path,event("task_complete",at+20000,{turn_id:"first"})+event("task_started",at+120000,{turn_id:"second",started_at:at/1000+120})); + result=await readCodexTiming(path,"thread",result.cursor,at+120001,key,"codex_app"); + assert.equal(result.observations[0].end-result.observations[0].start,20000); + await appendFile(path,event("turn_aborted",at+150000,{turn_id:"second",completed_at:at/1000+150})); + result=await readCodexTiming(path,"thread",result.cursor,at+150001,key,"codex_app"); + assert.equal(result.observations[0].start,at+120000); + assert.equal(result.observations[0].end-result.observations[0].start,30000); +}); + +test("invalid numeric envelopes cannot become execution or poison later valid timing", async t => { + const path=join(await fixture(t),"invalid-time.jsonl"); + await writeFile(path,JSON.stringify({type:"session_meta",payload:{id:"thread"}})+"\n"); + let result=await readCodexTiming(path,"thread",undefined,at,key,"codex_app"); + for(const [start,end] of [[-1,at/1000+5],[at,at+5],[at/1000+10,at/1000+5],[at/1000,at/1000+10000]]) { + await appendFile(path,event("task_complete",at+5000,{turn_id:"bad",started_at:start,completed_at:end})); + } + await appendFile(path,event("task_started",at,{turn_id:"good",started_at:at/1000})+event("token_count",at+86400000)+event("token_count",at+5000)); + result=await readCodexTiming(path,"thread",result.cursor,at+5001,key,"codex_app"); + assert.deepEqual(result.observations,[]); + assert.equal(result.cursor.open?.start,at); + await appendFile(path,event("task_complete",at+6000,{turn_id:"good",completed_at:at/1000+6})); + result=await readCodexTiming(path,"thread",result.cursor,at+6001,key,"codex_app"); + assert.equal(result.observations.length,1); + assert.equal(result.observations[0].end-result.observations[0].start,6000); +}); + test("universal cycles preserve first quota and first successful spend across all Host labels and reordered transport", async t => { const root = await fixture(t); for (const host of ["codex-app", "claude-code", "dsh", "opencode", "generic-cli"]) { @@ -163,7 +211,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,5); + const enabled=await configure(path,ctx,"enable"); assert.equal(enabled.notice.version,6); const current=JSON.parse(await readFile(path,"utf8")); assert.notEqual(current.generation,prior.generation); await configure(path,ctx,"disable"); assert.equal((await inspect(path,ctx)).consent,"disabled"); diff --git a/tests/test_usage_goal.py b/tests/test_usage_goal.py index 2a233387a0..053d3fccac 100644 --- a/tests/test_usage_goal.py +++ b/tests/test_usage_goal.py @@ -67,6 +67,44 @@ def set(self): assert str(tmp_path) not in json.dumps([row["observation"] for row in observed]) +def test_real_host_call_matches_wall_duration_through_detached_ts(tmp_path, monkeypatch): + import subprocess + import sys + from pathlib import Path + monkeypatch.setattr(usage_ping, "select_default_runtime_root", lambda: tmp_path) + 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") + before = time.monotonic_ns() + with usage_goal.observe_goal_execution(tmp_path / "runtime", "fixture-goal", host="generic-cli"): + child_start = time.time_ns() // 1_000_000 + subprocess.run([sys.executable, "-c", "import time; time.sleep(0.25)"], check=True) + child_end = time.time_ns() // 1_000_000 + duration = (time.monotonic_ns() - before) // 1_000_000 + goals = Path(str(usage_ping.state_path()) + ".goals") + installation_path = Path(str(usage_ping.state_path()) + ".installation") + deadline = time.monotonic() + 8 + while time.monotonic() < deadline: + if goals.exists() and installation_path.exists(): + state = json.loads(goals.read_text()) + installation = json.loads(installation_path.read_text()) + if state["goals"] and installation["days"][0]["runtime"]: + break + time.sleep(0.03) + else: + pytest.fail("real host observation did not reach detached TS") + intervals = state["goals"][0]["intervals"] + measured = sum(end - start for start, end in intervals) + assert state["goals"][0]["measurement"] == "host_call" + assert 250 <= measured <= duration + 1 + # The wall origin and monotonic elapsed duration are each floored to + # milliseconds; their sum can be one millisecond below floor(wall end). + assert intervals[0][0] <= child_start and intervals[-1][1] + 1 >= child_end + assert installation["days"][0]["runtime"] == [{"measurement": "host_call", "observed_minutes": 0}] + usage_ping.control("disable") + + def _bound_fixture(tmp_path, monkeypatch): import sqlite3 home = tmp_path / "codex-home" @@ -118,7 +156,7 @@ def test_metadata_cannot_redirect_timing_read_outside_selected_home(tmp_path, mo assert usage_goal._bound_codex_session(registry, "goal", "agent") is None -def _bound_cycle_through_detached_ts(tmp_path, monkeypatch, *, header_characters: int = 0): +def _bound_cycle_through_detached_ts(tmp_path, monkeypatch, *, header_characters: int = 0, numeric_seconds: bool = False): """Drive the real quota observer -> detached TS chain over one bound session. `header_characters` widens the session_meta line the way Codex's recorder @@ -156,11 +194,17 @@ def publish(phase): at=time.time_ns() // 1_000_000, host="unknown") publish("start") wait_for(lambda: bool(json.loads(cycle_path.read_text())["cursors"])) + if numeric_seconds: + # Wait for binding before choosing a whole-second provider interval. + # Slow subprocess startup must not put its start before the cursor's + # consent boundary and legitimately clip the expected two seconds. + now = (time.time_ns() // 1_000_000_000 + 1) * 1000 + time.sleep(max(0, (now + 2000 - time.time_ns() // 1_000_000) / 1000)) end = time.time_ns() // 1_000_000 with rollout.open("a") as stream: stream.write(json.dumps({"type": "event_msg", "timestamp": datetime.now(timezone.utc).isoformat(), "payload": { - "type": "task_complete", "turn_id": "turn-a", "started_at": datetime.fromtimestamp(now / 1000, timezone.utc).isoformat(), - "completed_at": datetime.fromtimestamp(end / 1000, timezone.utc).isoformat(), + "type": "task_complete", "turn_id": "turn-a", "started_at": now / 1000 if numeric_seconds else datetime.fromtimestamp(now / 1000, timezone.utc).isoformat(), + "completed_at": (now + 2000) / 1000 if numeric_seconds else datetime.fromtimestamp(end / 1000, timezone.utc).isoformat(), "last_agent_message": "PRIVATE CONTENT MUST NOT LEAVE THE SESSION", }}) + "\n") publish("spend") @@ -182,6 +226,19 @@ def test_real_detached_cycle_and_bound_codex_event_reach_shared_ts_aggregator(tm assert not cycle_path.exists() and not goals.exists() +def test_provider_numeric_timing_reaches_real_detached_installation_aggregator(tmp_path, monkeypatch): + from pathlib import Path + preview, _, _, _ = _bound_cycle_through_detached_ts(tmp_path, monkeypatch, numeric_seconds=True) + assert {row["measurement"] for row in preview["counters"]} == {"quota_cycle", "codex_turn"} + installation = json.loads(Path(str(usage_ping.state_path()) + ".installation").read_text()) + timed = [day for day in installation["days"] if "codex_turn" in day["intervals"]] + assert len(timed) == 1 + start, end = timed[0]["intervals"]["codex_turn"][0] + assert end - start == 2000 + assert next(row for row in timed[0]["runtime"] if row["measurement"] == "codex_turn")["observed_minutes"] == 0 + usage_ping.control("disable") + + def test_bound_session_with_a_long_metadata_header_still_reaches_ts(tmp_path, monkeypatch): """Codex records base_instructions in session_meta; a long header must bind.""" characters = 72_000 diff --git a/tests/test_usage_ping.py b/tests/test_usage_ping.py index 88eb7a6de1..3d4bcb684a 100644 --- a/tests/test_usage_ping.py +++ b/tests/test_usage_ping.py @@ -175,7 +175,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'] == 5 + assert json.loads(usage_ping.state_path().read_text())['notice']['version'] == 6 @pytest.mark.parametrize('setting,value', [ @@ -380,7 +380,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'] == 5 + assert initial['notice']['version'] == 6 connection.request('POST', path, json.dumps({'notice': initial['notice']}), {'Content-Type': 'application/json'}) response = connection.getresponse() acknowledged = json.loads(response.read()) @@ -417,6 +417,61 @@ def test_real_chat_settings_share_cli_choice_and_reject_cross_origin(isolated, u server.server_close() +def test_context_setting_real_cli_and_http_share_state_without_enabling(isolated, capsys): + import http.client + from loopx.chat_server import ChatHTTPServer, ChatRequestHandler + assert main(['usage-ping', 'context', '--context', 'maintainer', '--format', 'json']) == 0 + initial = json.loads(capsys.readouterr().out) + assert initial['stored_context'] == 'maintainer' and not initial['sending'] + server = ChatHTTPServer(('127.0.0.1', 0), ChatRequestHandler) + server.verbose = False + threading.Thread(target=server.serve_forever, daemon=True).start() + connection = http.client.HTTPConnection(*server.server_address, timeout=8) + path = '/api/chat/usage-statistics' + try: + connection.request('GET', path) + response = connection.getresponse() + assert json.loads(response.read())['effective_context'] == 'maintainer' + connection.request('POST', path, json.dumps({'context': 'personal'}), {'Content-Type': 'application/json'}) + response = connection.getresponse() + updated = json.loads(response.read()) + assert response.status == 200 and updated['effective_context'] == 'personal' and not updated['sending'] + assert usage_ping.control('status')['stored_context'] == 'personal' + connection.request('POST', path, json.dumps({'context': 'my-company'}), {'Content-Type': 'application/json'}) + response = connection.getresponse() + assert response.status == 400 + response.read() + connection.request('POST', path, json.dumps({'context': 'maintainer'}), {'Content-Type': 'application/json', 'Origin': 'https://evil.example'}) + response = connection.getresponse() + assert response.status == 403 + response.read() + assert usage_ping.control('status')['stored_context'] == 'personal' + finally: + connection.close() + server.shutdown() + server.server_close() + + +def test_real_cli_installation_profile_tags_context_and_contains_no_work_identity(isolated, collector, monkeypatch, capsys): + endpoint, received, accepted, release = collector + release.set() + monkeypatch.setenv('LOOPX_USAGE_PING_ENDPOINT', endpoint) + assert main(['usage-ping', 'context', '--context', 'maintainer', '--format', 'json']) == 0 + capsys.readouterr() + usage_ping.control('enable') + assert main(['version', '--format', 'json']) == 0 + assert json.loads(capsys.readouterr().out)['ok'] + deadline = time.monotonic() + 8 + while time.monotonic() < deadline and not any(row['schema'] == 'loopx_installation_usage_v1' for row in received): + time.sleep(0.05) + payload = next(row for row in received if row['schema'] == 'loopx_installation_usage_v1') + assert payload['profiles'][0]['context'] == 'maintainer' + assert payload['profiles'][0]['cli'] == [{'feature': 'version', 'count': 1}] + assert payload['profiles'][0]['runtime'] == [] + assert set(payload) == {'schema', 'install_id', 'profiles'} + assert not any(word in json.dumps(payload) for word in ('goal_id', 'agent_id', 'prompt', str(isolated))) + + 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 @@ -453,10 +508,10 @@ def test_v3_cli_upgrade_requires_visible_renewal_before_real_http(isolated, coll assert not accepted.wait(0.3) and received == [] visible = subprocess.run(command, capture_output=True, text=True, timeout=30) assert visible.returncode == 0 and json.loads(visible.stdout)['ok'] - assert 'first measured CLI result' in visible.stderr and '15 minutes' in visible.stderr - assert 'network timing' in visible.stderr + assert 'daily installation profiles' in visible.stderr and '15 minutes' in visible.stderr + assert 'connection metadata' in visible.stderr current = json.loads(path.read_text()) - assert current['notice']['version'] == 5 + assert current['notice']['version'] == 6 assert current['generation'] != old['generation'] assert current['counters'] == [] assert not accepted.wait(0.3) and received == [] From c1e0abe179b8595eb9db1cae8ecf5d02f418dcf7 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Fri, 2 Oct 2026 23:48:49 +0800 Subject: [PATCH 4/4] refactor(usage): reuse typed device context validation Signed-off-by: huangruiteng --- apps/presentation/dashboard/src/data/chat.ts | 1 + .../personal-workspace/usage-statistics-settings.tsx | 4 ++-- loopx/chat_usage_statistics_api.py | 7 ++++--- loopx/cli_commands/usage_ping.py | 9 +++++++-- loopx/control_plane/runtime/usage_statistics.ts | 6 +++++- loopx/control_plane/runtime/usage_statistics_cli.ts | 7 ++++--- .../control_plane/runtime/usage_statistics_contract.ts | 1 + .../runtime/usage_statistics_diagnostics.ts | 4 ++-- loopx/usage_ping.py | 6 ++++++ tests/test_usage_ping.py | 10 ++++++++++ 10 files changed, 42 insertions(+), 13 deletions(-) diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index de6a162815..0c03438b52 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -2251,6 +2251,7 @@ export async function disconnectLarkGoalTopic(goalId: string, connectionId: stri ); } +export { CONTEXTS as usageContexts } from "../../../../../loopx/control_plane/runtime/usage_statistics_contract"; const usageStatisticsSchema = z.object({ consent: z.enum(["default", "enabled", "disabled"]), sending: z.boolean(), blocked_by: z.string().nullable(), endpoint: z.string().nullable(), 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 3079e1e22b..bfe25840d7 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 @@ -1,5 +1,5 @@ import { useEffect, useState } from "react"; -import { setUsageContext, usageStatistics, type UsageStatistics } from "../../data/chat"; +import { setUsageContext, usageContexts, usageStatistics, type UsageStatistics } from "../../data/chat"; import { useWorkspaceI18n } from "./i18n"; export function UsageStatisticsSettings() { @@ -36,7 +36,7 @@ export function UsageStatisticsSettings() { {state ? <>