From 2a97ccdd2985d058b9a19c9ef20c93a9c4bfc291 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 11 Sep 2026 02:53:23 +0800 Subject: [PATCH] feat(telemetry): add aggregate capture pilot --- CHANGELOG.md | 1 + README.md | 95 +++ package.json | 1 + scripts/lib/telemetry-aggregate-capture.mjs | 860 ++++++++++++++++++++ scripts/telemetry-aggregate-capture.mjs | 34 + test/telemetry-aggregate-capture.test.mjs | 727 +++++++++++++++++ 6 files changed, 1718 insertions(+) create mode 100644 scripts/lib/telemetry-aggregate-capture.mjs create mode 100644 scripts/telemetry-aggregate-capture.mjs create mode 100644 test/telemetry-aggregate-capture.test.mjs diff --git a/CHANGELOG.md b/CHANGELOG.md index cd3c051..68ecb9b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ **Highlights:** Public update checks and seven-day usage aggregates, with optional feature statistics and no stored install identifiers. +- Add an operator-invoked, dry-run-default aggregate capture pilot for one closed UTC day: hourly Analytics Engine reports and separate HTTP country estimates in a private, exclusively created bundle with offline verification. No scheduled backup, upload, raw-event restoration, storage provisioning, or retention change. - Export one hash-pinned archived hourly Analytics Engine query to private daily JSON, CSV, and provenance/coverage artifacts. Preserve exact sampled report counts, flag partial days and missing hours, and restrict comparison totals to complete closed UTC days. No raw-event restoration, live queries, or backup job. - Record bounded Cloudflare-derived request-origin country, region code, city, and timezone in baseline update-request metadata, without raw IPs or direct identifiers. Replace the unshipped client UTC-offset proposal, preserve existing columns and public stats, and disclose the unchanged three-month retention and feature opt-out boundary. - Add offline, manifest-bound npm comparison tooling with matched UTC windows, captured feed cutoffs, anomaly checks, exclusion sensitivity, and separate undated version totals. Download events are never converted to installations or users. diff --git a/README.md b/README.md index 558d4f6..9414afa 100644 --- a/README.md +++ b/README.md @@ -365,6 +365,101 @@ Only `complete_closed` days are `comparisonEligible` and contribute to `summary.completeClosedTotals`. This certifies UTC calendar coverage, not complete events. Geography is unknown in this export, including rows collected before geography was recorded. +## Aggregate capture pilot + +`npm run telemetry:capture` plans or captures **one explicit closed UTC day** of hourly +Analytics Engine report aggregates and a **separate** HTTP country query. It is an +operator-invoked pilot, not a scheduled or full backup. No version/plugin distributions, +finer geography, raw events, replay, uploads, storage bindings, or retention policy are added. +The existing historical exporter is unchanged. + +The default is a dry run. It reads no credentials, makes no network requests, and creates +no files: + +```bash +npm run telemetry:capture -- --day 2025-02-03 +``` + +Choose and approve a private storage location and its access/retention policy before +performing an actual capture. There is no default output location. The output parent must +already exist; explicitly resolve symlinked parents to their real paths. Execution needs +separate operator-provided `TELEMETRY_AE_READ_TOKEN` and `TELEMETRY_HTTP_READ_TOKEN` +environment variables, plus explicit noncredential account and zone IDs: + +```bash +npm run telemetry:capture -- --execute --day "$UTC_DAY" \ + --account-id "$ACCOUNT_ID" --zone-id "$ZONE_ID" \ + --output "$(realpath "$APPROVED_PARENT")/day-$UTC_DAY" +``` + +Provide an account-scoped Analytics Engine SQL read token and a separate zone-scoped +HTTP analytics read token. The tool never discovers credentials, opens an auth store, +refreshes OAuth, or provisions permissions. It posts only to the fixed Cloudflare +API endpoints, rejects redirects, and never retries automatically. + +Current HTTP dataset settings are requested first. The dataset must be enabled and +advertise all required metric, dimension and filter fields, a full-day duration, a +10,000-row page, and sufficient field capacity. The requested day must fit the returned +`notOlderThan` lookback. Lookback is checked again immediately before the HTTP country +request, after AE finishes; expiry stops the request and leaves an incomplete bundle. +Each request has a 45-second deadline. Responses are bounded to 256 KiB for settings, +4 MiB for AE and 2 MiB for HTTP countries. Truncated/encoded bodies, GraphQL errors, +unsafe counts, duplicate groups, and reached row limits fail closed. +Unsupported/error payloads are rejected before their bodies or complete receipts enter +the bundle; failure diagnostics omit upstream text. + +The AE statement is the historical exporter's audited hourly query, with a structural +maximum of 25 rows and SQL limit 26. It preserves the distinct submitted SQL and saved +SQL-plus-LF hashes, raw response bytes, plan, attempt, receipt, and their digest pins. +`count()` remains a query row count, not a stored-row census. + +The HTTP query uses `httpRequestsAdaptiveGroups` with the exact host +`telemetry.openclaw.ai`, path `/api/latest-version`, `requestSource: eyeball`, and +half-open UTC day. It requests only country, `count`, and average `sampleInterval`. +All methods, statuses and bots within that scope remain included. HTTP `count` is +already estimated; the sampling interval is diagnostic, **not another multiplier**. +Country counts must be nonnegative safe JSON integers and are exported and summed as +exact decimal strings. Special, empty and null country labels remain distinct. +Absent countries and empty results are unknown, not zero. + +### Bundle and offline verification + +The private bundle contains `bundle-plan.json`, the `ae/` capture, unchanged exporter +outputs in `ae-daily/`, HTTP settings/query responses and receipts in `http/`, and +normalized `http/country.json`. Directories are `0700` and files `0600`. Both sources' +local prerequisites are validated before exclusively claiming the output directory, +which happens before the first request. A concurrent loser makes no request. + +`manifest.json` is written last, only after both captures validate and offline +regeneration matches all derived outputs byte-for-byte. Completion means **query/wire +completeness**, not complete events or census coverage. Missing AE hours remain unknown +through the unchanged exporter; HTTP requests are never joined to AE reports or used +to infer geography for pre-geography AE rows. + +```bash +npm run telemetry:capture -- --verify --day "$UTC_DAY" \ + --account-id "$ACCOUNT_ID" --zone-id "$ZONE_ID" \ + --output "$(realpath "$APPROVED_PARENT")/day-$UTC_DAY" +``` + +Verification needs no credentials or network. It binds the requested day, source +identities, exact queries, raw bytes, receipts and hashes, then reruns the unchanged +historical exporter in a private owned temporary directory outside the bundle. +All three regenerated AE files must match; only the verifier's scratch is removed. +Hashes establish internally consistent operator-selected evidence, not independent +server authentication. + +An execute rerun checks an existing bundle **before** reading credentials or fetching +metadata, returning `unchanged` only after the same offline verification. It still works +after upstream retention expires. Partial, conflicting, extra, symlinked, hardlinked or +non-private files fail without overwrite, repair, resume or network access. Failed +captures retain their incomplete evidence; absence of a valid completion manifest +means the bundle is not complete. Review failures and choose a new explicit destination +only after resolving the cause. + +This pilot changes neither Analytics Engine retention nor the receiver's collection +or runtime behavior. + ## License MIT © OpenClaw Foundation diff --git a/package.json b/package.json index cce2d4f..88d0690 100644 --- a/package.json +++ b/package.json @@ -13,6 +13,7 @@ "check": "npm run vocabulary:check && npm run typecheck && npm run test", "npm:quality": "node scripts/npm-quality.mjs", "telemetry:history": "node scripts/telemetry-history.mjs", + "telemetry:capture": "node scripts/telemetry-aggregate-capture.mjs", "vocabulary:update": "node scripts/public-vocabulary.mjs", "vocabulary:check": "node scripts/public-vocabulary.mjs --check" }, diff --git a/scripts/lib/telemetry-aggregate-capture.mjs b/scripts/lib/telemetry-aggregate-capture.mjs new file mode 100644 index 0000000..df08bb8 --- /dev/null +++ b/scripts/lib/telemetry-aggregate-capture.mjs @@ -0,0 +1,860 @@ +import { createHash } from "node:crypto"; +import { + closeSync, + constants, + fstatSync, + lstatSync, + mkdirSync, + mkdtempSync, + openSync, + opendirSync, + readSync, + realpathSync, + rmSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { dirname, isAbsolute, join, parse, relative, resolve, sep } from "node:path"; +import { TextDecoder } from "node:util"; +import { runTelemetryHistory } from "./telemetry-history.mjs"; + +const DAY = 86_400_000; +const HTTP_LIMIT = 10_000; +const TIMEOUT = 45_000; +const GRAPHQL = "https://api.cloudflare.com/client/v4/graphql"; +const AE_ENDPOINT = "/client/v4/accounts//analytics_engine/sql"; +const CONTRACT = "telemetry-aggregate-capture-pilot-v1"; +const AE_COLUMNS = { + bucket: "DateTime", + weightedReports: "UInt64", + queryRows: "UInt64", + featureReports: "UInt64", + featureQueryRows: "UInt64", + minSampleInterval: "UInt32", + maxSampleInterval: "UInt32", + latestEventAt: "DateTime", + latestFeatureAt: "DateTime", +}; +const REQUIRED_FIELDS = [ + "count", + "avg_sampleInterval", + "dimensions_clientCountryName", + "dimensions_clientRequestHTTPHost", + "dimensions_clientRequestPath", + "dimensions_requestSource", + "dimensions_datetime", +]; +const SETTINGS_QUERY = `query TelemetryMetadata($zoneTag: string) { + viewer { zones(filter: {zoneTag: $zoneTag}) { + settings { httpRequestsAdaptiveGroups { + enabled notOlderThan maxDuration maxPageSize maxNumberOfFields availableFields + } } + } } +}`; +const COUNTRY_QUERY = `query CountryDay($zoneTag: string, $start: Time, $end: Time) { + viewer { + zones(filter: {zoneTag: $zoneTag}) { + httpRequestsAdaptiveGroups( + limit: 10000 + filter: { + clientRequestHTTPHost: "telemetry.openclaw.ai" + clientRequestPath: "/api/latest-version" + requestSource: "eyeball" + datetime_geq: $start + datetime_lt: $end + } + ) { + dimensions { clientCountryName } + count + avg { sampleInterval } + } + } + } +}`; +const FILES = { + "bundle-plan.json": 16_384, + "ae/capture-plan.json": 16_384, + "ae/q2/query.sql": 9500, + "ae/q2/attempt.json": 16_384, + "ae/q2/response.json": 4 * 1024 * 1024, + "ae/q2/receipt.json": 65_536, + "http/settings-request.json": 16_384, + "http/settings-response.json": 256 * 1024, + "http/settings-receipt.json": 16_384, + "http/country-request.json": 16_384, + "http/country-response.json": 2 * 1024 * 1024, + "http/country-receipt.json": 16_384, + "http/country.json": 2 * 1024 * 1024, + "ae-daily/daily.json": 32_768, + "ae-daily/daily.csv": 32_768, + "ae-daily/manifest.json": 65_536, +}; +const json = (value) => `${JSON.stringify(value, null, 2)}\n`; +const hash = (bytes) => createHash("sha256").update(bytes).digest("hex"); +const iso = (ms) => new Date(ms).toISOString(); +const object = (value) => value !== null && typeof value === "object" && !Array.isArray(value); +class CaptureError extends Error {} +function requireValue(condition, message) { + if (!condition) throw new CaptureError(message); +} +function keys(value, names) { + requireValue( + object(value) && + Object.keys(value).length === names.length && + names.every((name) => Object.hasOwn(value, name)), + "unexpected capture schema", + ); +} +function decode(bytes) { + try { + return JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(bytes)); + } catch { + throw new CaptureError("response is not valid UTF-8 JSON"); + } +} +function instant(value) { + const ms = typeof value === "string" ? Date.parse(value) : NaN; + requireValue(Number.isFinite(ms) && iso(ms) === value, "invalid capture clock"); + return ms; +} +function dayWindow(day) { + requireValue( + typeof day === "string" && /^\d{4}-\d{2}-\d{2}$/u.test(day), + "day must be YYYY-MM-DD", + ); + const start = Date.parse(`${day}T00:00:00.000Z`); + requireValue( + Number.isFinite(start) && start >= 0 && iso(start).slice(0, 10) === day, + "invalid UTC calendar day", + ); + const end = start + DAY; + requireValue(end <= Date.now() && !iso(end).startsWith("+"), "day must be closed in UTC"); + return { startInclusive: iso(start), endExclusive: iso(end) }; +} +function hourlySql(window) { + const at = (value) => value.slice(0, 19).replace("T", " "); + return ( + "SELECT toStartOfInterval(timestamp,INTERVAL '1' HOUR) AS bucket," + + "sum(_sample_interval) AS weightedReports,count() AS queryRows," + + "sumIf(_sample_interval,double1=1) AS featureReports," + + "sum(if(double1=1,1,0)) AS featureQueryRows," + + "min(_sample_interval) AS minSampleInterval,max(_sample_interval) AS maxSampleInterval," + + "max(timestamp) AS latestEventAt," + + "max(if(double1=1,timestamp,toDateTime('1970-01-01 00:00:00'))) AS latestFeatureAt " + + `FROM openclaw_telemetry WHERE timestamp>=toDateTime('${at(window.startInclusive)}') ` + + `AND timestamp 0 && + path.length <= 4096 && + !path.includes("\0") && + !path.includes("\\") && + !path.split("/").includes(".."), + "an output path without traversal is required", + ); +} +function realDirectory(path) { + const absolute = resolve(path); + let current = parse(absolute).root; + for (const part of relative(current, absolute).split(sep).filter(Boolean)) { + current = join(current, part); + const info = lstatSync(current); + requireValue(info.isDirectory() && !info.isSymbolicLink(), "symlink or non-directory rejected"); + } + return absolute; +} +function privateMode(info, mode) { + requireValue( + (info.mode & 0o777) === mode && + (typeof process.getuid !== "function" || info.uid === process.getuid()), + "bundle must be private and owned by the current operator", + ); +} +function readBytes(path, maximum) { + const before = lstatSync(path); + requireValue( + before.isFile() && !before.isSymbolicLink() && before.nlink === 1, + "unsafe bundle file", + ); + privateMode(before, 0o600); + requireValue(before.size <= maximum, "bundle byte limit exceeded"); + const fd = openSync(path, constants.O_RDONLY | constants.O_NOFOLLOW); + try { + const opened = fstatSync(fd); + requireValue( + opened.isFile() && + opened.nlink === 1 && + opened.dev === before.dev && + opened.ino === before.ino && + opened.size === before.size, + "bundle changed before reading", + ); + privateMode(opened, 0o600); + const buffer = Buffer.alloc(opened.size + 1); + let used = 0; + while (used < buffer.length) { + const n = readSync(fd, buffer, used, buffer.length - used, used); + if (!n) break; + used += n; + } + const after = fstatSync(fd); + requireValue( + used === opened.size && + after.size === opened.size && + after.mtimeMs === opened.mtimeMs && + after.ctimeMs === opened.ctimeMs, + "bundle changed while reading", + ); + return buffer.subarray(0, used); + } finally { + closeSync(fd); + } +} +function bundleIO(root) { + const claimed = lstatSync(root); + const guard = (directory = root) => { + realDirectory(directory); + const now = lstatSync(root); + requireValue(now.ino === claimed.ino && now.dev === claimed.dev, "bundle directory changed"); + privateMode(now, 0o700); + privateMode(lstatSync(directory), 0o700); + }; + guard(); + return { + root, + guard, + read(name, maximum = FILES[name]) { + guard(dirname(join(root, name))); + return readBytes(join(root, name), maximum); + }, + write(name, bytes) { + requireValue( + Buffer.byteLength(bytes) <= (FILES[name] ?? 65_536), + "output byte limit exceeded", + ); + guard(dirname(join(root, name))); + const fd = openSync( + join(root, name), + constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL | constants.O_NOFOLLOW, + 0o600, + ); + try { + privateMode(fstatSync(fd), 0o600); + writeFileSync(fd, bytes); + } finally { + closeSync(fd); + } + guard(); + }, + }; +} +function validateLayout(bundle, completed) { + const expected = new Map(); + for (const file of [...Object.keys(FILES), ...(completed ? ["manifest.json"] : [])]) { + const parts = file.split("/"); + for (let i = 0; i < parts.length; i++) { + const parent = parts.slice(0, i).join("/"); + if (!expected.has(parent)) expected.set(parent, new Set()); + expected.get(parent).add(parts[i]); + } + } + for (const [path, names] of expected) { + bundle.guard(join(bundle.root, path)); + const dir = opendirSync(join(bundle.root, path)); + let count = 0; + try { + for (let entry; (entry = dir.readSync());) { + requireValue(names.has(entry.name) && ++count <= names.size, "bundle conflicts"); + } + } finally { + dir.closeSync(); + } + requireValue(count === names.size, "bundle is incomplete; no resume or overwrite"); + } +} +function zoneResult(value, field) { + keys(value, ["data", "errors"]); + requireValue( + value.errors === null || (Array.isArray(value.errors) && value.errors.length === 0), + "GraphQL returned errors", + ); + keys(value.data, ["viewer"]); + keys(value.data.viewer, ["zones"]); + const zones = value.data.viewer.zones; + requireValue(Array.isArray(zones) && zones.length === 1, "expected exactly one HTTP zone"); + keys(zones[0], [field]); + return zones[0][field]; +} +function settingsFrom(value) { + const wrapper = zoneResult(value, "settings"); + keys(wrapper, ["httpRequestsAdaptiveGroups"]); + const settings = wrapper.httpRequestsAdaptiveGroups; + keys(settings, [ + "enabled", + "notOlderThan", + "maxDuration", + "maxPageSize", + "maxNumberOfFields", + "availableFields", + ]); + requireValue(settings.enabled === true, "HTTP dataset is not enabled"); + for (const key of ["notOlderThan", "maxDuration", "maxPageSize", "maxNumberOfFields"]) + requireValue( + Number.isSafeInteger(settings[key]) && settings[key] > 0, + "invalid HTTP dataset limits", + ); + const fields = settings.availableFields; + requireValue( + Array.isArray(fields) && + fields.length <= 4096 && + fields.every((field) => typeof field === "string" && field.length <= 160) && + REQUIRED_FIELDS.every((field) => fields.includes(field)), + "HTTP dataset lacks required fields or filters", + ); + requireValue( + settings.maxDuration >= DAY / 1000 && + settings.maxPageSize >= HTTP_LIMIT && + settings.maxNumberOfFields >= REQUIRED_FIELDS.length, + "HTTP dataset duration, page, or field cap is insufficient", + ); + return settings; +} +function checkLookback(settings, window, now) { + requireValue( + now >= instant(window.endExclusive) && + (now - instant(window.startInclusive)) / 1000 <= settings.notOlderThan, + "requested day is outside current HTTP lookback; no country request was made", + ); +} +function countryOutput(value, window) { + const rows = zoneResult(value, "httpRequestsAdaptiveGroups"); + requireValue( + Array.isArray(rows) && rows.length < HTTP_LIMIT, + "HTTP country limit reached or invalid rows", + ); + const seen = new Set(); + const countries = rows.map((row) => { + keys(row, ["dimensions", "count", "avg"]); + keys(row.dimensions, ["clientCountryName"]); + keys(row.avg, ["sampleInterval"]); + const country = row.dimensions.clientCountryName; + requireValue( + country === null || + (typeof country === "string" && + Buffer.byteLength(country) <= 128 && + !/[\u0000-\u001f\u007f]/u.test(country)), + "invalid country label", + ); + requireValue(!seen.has(country), "duplicate country group"); + seen.add(country); + requireValue( + Number.isSafeInteger(row.count) && row.count >= 0, + "HTTP count must be a nonnegative safe integer", + ); + const interval = row.avg.sampleInterval; + requireValue( + interval === null || + (typeof interval === "number" && Number.isFinite(interval) && interval >= 1), + "invalid HTTP sample interval", + ); + return { country, estimatedRequests: String(row.count), sampleInterval: interval }; + }); + countries.sort((a, b) => + a.country === b.country + ? 0 + : a.country === null + ? -1 + : b.country === null + ? 1 + : a.country < b.country + ? -1 + : 1, + ); + return { + schemaVersion: 1, + queryContract: "http-country-eyeball-v1", + window, + coverage: countries.length ? "returned_groups_only" : "unknown", + estimatedRequests: countries.length + ? countries.reduce((sum, row) => sum + BigInt(row.estimatedRequests), 0n).toString() + : null, + countries, + limitations: [ + "HTTP count is already estimated; sampleInterval is diagnostic, never another multiplier.", + "Missing countries and empty results are unknown, not measured zero.", + "All methods, statuses and bots within the exact host, path and eyeball scope are included.", + "HTTP requests are not accepted AE reports, unique installations, users, or an opt-in rate.", + "Country labels, including null and special values, are preserved without inference or joins.", + ], + }; +} +async function post(url, token, body, maximum, contentType, stage) { + const controller = new AbortController(); + let timer; + let reader; + try { + return await Promise.race([ + (async () => { + const response = await fetch(url, { + method: "POST", + headers: { + Authorization: `Bearer ${token}`, + "Content-Type": contentType, + Accept: "application/json", + "Accept-Encoding": "identity", + }, + body, + redirect: "error", + signal: controller.signal, + }); + requireValue( + !response.redirected && response.status === 200, + `${stage} request failed (HTTP ${Number(response.status)}); redirects are rejected`, + ); + const encoding = response.headers.get("content-encoding"); + requireValue(!encoding || encoding === "identity", "encoded wire response is unsupported"); + const length = response.headers.get("content-length"); + requireValue( + length === null || (/^\d{1,10}$/u.test(length) && Number(length) <= maximum), + "response byte limit exceeded", + ); + requireValue(response.body, "response body is missing"); + reader = response.body.getReader(); + const bytes = Buffer.alloc(maximum); + let used = 0; + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + requireValue(used + value.byteLength <= maximum, "response byte limit exceeded"); + bytes.set(value, used); + used += value.byteLength; + } + requireValue(length === null || Number(length) === used, "response body is truncated"); + return { bytes: bytes.subarray(0, used), receivedAt: iso(Date.now()) }; + })(), + new Promise((_, reject) => { + timer = setTimeout(() => { + controller.abort(); + reject(new CaptureError(`${stage} request timed out; no retry`)); + }, TIMEOUT); + }), + ]); + } catch (error) { + if (error instanceof CaptureError) throw error; + throw new CaptureError( + `${stage} request failed; check explicit read credentials and connectivity; no retry`, + ); + } finally { + clearTimeout(timer); + controller.abort(); + if (reader) reader.cancel().catch(() => {}); + } +} +function httpReceipt(request, wire, requestedAt) { + return { + method: "POST", + endpoint: GRAPHQL, + requestSha256: hash(request), + requestedAt, + receivedAt: wire.receivedAt, + status: 200, + complete: true, + wireBodyTruncated: false, + responseRedacted: false, + wireBytesRead: wire.bytes.length, + wireSha256: hash(wire.bytes), + }; +} +function aePlan(spec, captureStartedAt) { + const sql = hourlySql(spec.window); + return { + captureStartedAt, + windowStartInclusive: spec.window.startInclusive, + windowEndExclusive: spec.window.endExclusive, + queries: [{ id: "q2", sql, sqlSha256: hash(sql), structuralMaxRows: 25, sqlLimit: 26 }], + }; +} +function aeAttempt(spec, requestedAt) { + return { + queryId: "q2", + ordinal: 1, + method: "POST", + endpointTemplate: AE_ENDPOINT, + accountSha256: spec.sourceIdentity.aeAccountSha256, + sqlSha256: spec.ae.submittedSqlSha256, + requestedAt, + }; +} +function validateAeContent(response) { + keys(response, ["meta", "data", "rows", "rows_before_limit_at_least"]); + const columns = Object.entries(AE_COLUMNS); + requireValue( + Array.isArray(response.meta) && + response.meta.length === columns.length && + response.meta.every((column, index) => { + keys(column, ["name", "type"]); + return column.name === columns[index][0] && column.type === columns[index][1]; + }), + "unsupported AE response metadata", + ); + requireValue( + Array.isArray(response.data) && + response.data.length <= 25 && + response.rows === response.data.length && + response.rows_before_limit_at_least === response.data.length, + "incomplete AE response rows", + ); + // Reject unrequested fields and freeform error text before archiving any upstream bytes. + // The unchanged exporter still owns hourly, sampling, watermark and coverage semantics. + for (const row of response.data) { + keys(row, Object.keys(AE_COLUMNS)); + for (const [name, type] of columns) { + const value = row[name]; + const valid = + type === "DateTime" + ? typeof value === "string" && /^\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}$/u.test(value) + : type === "UInt64" + ? typeof value === "string" && + /^(?:0|[1-9]\d{0,19})$/u.test(value) && + BigInt(value) <= 18_446_744_073_709_551_615n + : Number.isInteger(value) && value >= 1 && value <= 4_294_967_295; + requireValue(valid, "unsupported AE response field representation"); + } + } +} +function aeReceipt(spec, wire) { + const response = decode(wire.bytes); + validateAeContent(response); + return { + queryId: "q2", + sqlSha256: spec.ae.submittedSqlSha256, + complete: true, + status: 200, + receivedAt: wire.receivedAt, + wireBytesRead: wire.bytes.length, + wireSha256: hash(wire.bytes), + wireBodyTruncated: false, + responseRedacted: false, + hardLimitReached: false, + rowsReturned: response.data.length, + declaredRows: response.rows, + rowsBeforeLimitAtLeast: response.rows_before_limit_at_least, + meta: response.meta, + }; +} +function history(bundle, output, pins) { + try { + bundle.guard(join(bundle.root, "ae/q2")); + return runTelemetryHistory({ + archive: join(bundle.root, "ae"), + query: "q2", + ...pins, + output, + }); + } catch { + throw new CaptureError("AE archive failed unchanged history validation"); + } +} +function verifyBundle(bundle, spec, requestBodies, completed) { + validateLayout(bundle, completed); + const bytes = Object.fromEntries(Object.keys(FILES).map((name) => [name, bundle.read(name)])); + const expectBytes = (name, expected) => + requireValue( + bytes[name].equals(Buffer.from(expected)), + "bundle content conflicts with its contract", + ); + expectBytes("bundle-plan.json", json(spec)); + const files = Object.entries(bytes).map(([file, data]) => ({ + file, + bytes: data.length, + sha256: hash(data), + })); + let savedManifest; + if (completed) { + savedManifest = bundle.read("manifest.json", 65_536); + const saved = decode(savedManifest); + requireValue( + saved?.complete === true && json(saved.files) === json(files), + "bundle hashes conflict or completion is missing", + ); + } + const parsed = (name) => decode(bytes[name]); + const capture = parsed("ae/capture-plan.json"); + const attempt = parsed("ae/q2/attempt.json"); + const receipt = parsed("ae/q2/receipt.json"); + expectBytes("ae/capture-plan.json", json(aePlan(spec, capture.captureStartedAt))); + expectBytes("ae/q2/attempt.json", json(aeAttempt(spec, attempt.requestedAt))); + expectBytes("ae/q2/query.sql", `${hourlySql(spec.window)}\n`); + expectBytes( + "ae/q2/receipt.json", + json( + aeReceipt(spec, { + bytes: bytes["ae/q2/response.json"], + receivedAt: receipt.receivedAt, + }), + ), + ); + const clocks = [instant(spec.window.endExclusive), instant(capture.captureStartedAt)]; + for (const name of ["settings", "country"]) { + const path = `http/${name}`; + const record = parsed(`${path}-receipt.json`); + expectBytes(`${path}-request.json`, requestBodies[name]); + expectBytes( + `${path}-receipt.json`, + json( + httpReceipt( + requestBodies[name], + { + bytes: bytes[`${path}-response.json`], + receivedAt: record.receivedAt, + }, + record.requestedAt, + ), + ), + ); + if (name === "country") clocks.push(instant(attempt.requestedAt), instant(receipt.receivedAt)); + clocks.push(instant(record.requestedAt), instant(record.receivedAt)); + } + requireValue( + clocks.every((time, index) => !index || time >= clocks[index - 1]), + "capture clock order mismatch", + ); + const settings = settingsFrom(parsed("http/settings-response.json")); + for (const name of ["settings", "country"]) + checkLookback(settings, spec.window, instant(parsed(`http/${name}-receipt.json`).requestedAt)); + const country = countryOutput(parsed("http/country-response.json"), spec.window); + expectBytes("http/country.json", json(country)); + const pins = { + planSha256: hash(bytes["ae/capture-plan.json"]), + receiptSha256: hash(bytes["ae/q2/receipt.json"]), + }; + // The unchanged exporter writes only to our scratch, never into the saved archive. + const scratchParent = realDirectory(realpathSync(tmpdir())); + const rel = relative(bundle.root, scratchParent); + requireValue( + rel === ".." || rel.startsWith(`..${sep}`) || isAbsolute(rel), + "verification scratch must be outside the bundle", + ); + const scratch = mkdtempSync(join(scratchParent, "telemetry-capture-verify-")); + const claimed = lstatSync(scratch); + let summary; + try { + privateMode(claimed, 0o700); + const output = join(scratch, "daily"); + summary = history(bundle, output, pins).summary; + for (const name of ["daily.json", "daily.csv", "manifest.json"]) + expectBytes(`ae-daily/${name}`, readBytes(join(output, name), FILES[`ae-daily/${name}`])); + } finally { + const now = lstatSync(scratch); + requireValue( + !now.isSymbolicLink() && now.ino === claimed.ino && now.dev === claimed.dev, + "verification scratch changed", + ); + rmSync(scratch, { recursive: true }); + } + const manifest = { + schemaVersion: 1, + contract: CONTRACT, + day: spec.day, + sourceIdentity: spec.sourceIdentity, + complete: true, + completeness: "query_and_wire_only", + pins, + files, + summary: { + ae: summary, + http: { + countryGroups: country.countries.length, + coverage: country.coverage, + estimatedRequests: country.estimatedRequests, + }, + }, + limitations: [ + "Pilot captures only hourly AE report aggregates and separate HTTP country aggregates.", + "Completion certifies query and wire validation, not complete events, unique installations, or a census.", + "Missing AE hours and empty HTTP results remain unknown, not measured zero.", + "Pre-geography AE rows have unknown geography; HTTP countries do not enrich or join AE reports.", + "No raw events, version/plugin distributions, finer geography, or full-dimensional/lossless backup.", + "No replay, restore writes, upload, scheduling, storage provisioning, or retention policy is implemented.", + "Hashes bind operator-selected local evidence, not independent server authentication.", + ], + }; + const encoded = json(manifest); + if (completed) + requireValue(savedManifest.equals(Buffer.from(encoded)), "completion manifest conflicts"); + for (const [name, original] of Object.entries(bytes)) + requireValue(bundle.read(name).equals(original), "bundle changed during verification"); + validateLayout(bundle, completed); + return { manifest: encoded, summary: manifest.summary }; +} + +/** Plan without side effects, capture one explicit day, or verify an existing bundle offline. */ +export async function runTelemetryAggregateCapture({ + day, + output, + accountId, + zoneId, + execute = false, + verify = false, +}) { + try { + requireValue( + typeof execute === "boolean" && typeof verify === "boolean" && !(execute && verify), + "select either execute or verify", + ); + const window = dayWindow(day); + if (!execute && !verify) { + return { + status: "planned", + contract: CONTRACT, + day, + window, + ae: { + endpointTemplate: AE_ENDPOINT, + sql: hourlySql(window), + structuralMaxRows: 25, + sqlLimit: 26, + }, + http: { endpoint: GRAPHQL, ...requests(window, "") }, + requires: [ + "explicit operator output", + "separate AE and HTTP read tokens", + "current HTTP capability and lookback checks", + ], + }; + } + for (const value of [accountId, zoneId]) + requireValue( + typeof value === "string" && /^[a-f0-9]{32}$/u.test(value), + "explicit account and zone identities must be 32 lowercase hex characters", + ); + noTraversal(output); + const root = resolve(output); + realDirectory(dirname(root)); + const spec = plan(day, window, accountId, zoneId); + const requestBodies = requests(window, zoneId); + let exists = false; + try { + lstatSync(root); + exists = true; + } catch (error) { + if (error.code !== "ENOENT") throw error; + } + // Reruns must work offline even after retention or credentials have expired. + if (exists) { + const result = verifyBundle(bundleIO(root), spec, requestBodies, true); + return { status: "unchanged", day, summary: result.summary }; + } + requireValue(!verify, "bundle is missing; offline verification cannot create it"); + const aeToken = process.env.TELEMETRY_AE_READ_TOKEN; + const httpToken = process.env.TELEMETRY_HTTP_READ_TOKEN; + for (const token of [aeToken, httpToken]) + requireValue( + typeof token === "string" && /^[\x21-\x7e]{1,4096}$/u.test(token), + "separate explicit AE and HTTP read credentials are required", + ); + requireValue(aeToken !== httpToken, "AE and HTTP read credentials must be separate"); + requireValue(typeof globalThis.fetch === "function", "fetch is unavailable"); + mkdirSync(root, { mode: 0o700 }); + const bundle = bundleIO(root); + for (const path of ["ae", "ae/q2", "http"]) { + bundle.guard(); + mkdirSync(join(root, path), { mode: 0o700 }); + } + bundle.write("bundle-plan.json", json(spec)); + bundle.write("ae/capture-plan.json", json(aePlan(spec, iso(Date.now())))); + bundle.write("ae/q2/query.sql", `${hourlySql(window)}\n`); + bundle.write("http/settings-request.json", requestBodies.settings); + const metadataAt = iso(Date.now()); + const metadataWire = await post( + GRAPHQL, + httpToken, + requestBodies.settings, + FILES["http/settings-response.json"], + "application/json", + "HTTP settings", + ); + const settings = settingsFrom(decode(metadataWire.bytes)); + bundle.write("http/settings-response.json", metadataWire.bytes); + bundle.write( + "http/settings-receipt.json", + json(httpReceipt(requestBodies.settings, metadataWire, metadataAt)), + ); + checkLookback(settings, window, Date.now()); + const requestedAt = iso(Date.now()); + bundle.write("ae/q2/attempt.json", json(aeAttempt(spec, requestedAt))); + const aeWire = await post( + `https://api.cloudflare.com/client/v4/accounts/${accountId}/analytics_engine/sql`, + aeToken, + hourlySql(window), + FILES["ae/q2/response.json"], + "text/plain", + "AE hourly", + ); + const receipt = aeReceipt(spec, aeWire); + bundle.write("ae/q2/response.json", aeWire.bytes); + bundle.write("ae/q2/receipt.json", json(receipt)); + history(bundle, join(root, "ae-daily"), { + planSha256: hash(bundle.read("ae/capture-plan.json")), + receiptSha256: hash(bundle.read("ae/q2/receipt.json")), + }); + bundle.write("http/country-request.json", requestBodies.country); + const countryAt = iso(Date.now()); + checkLookback(settings, window, instant(countryAt)); + const countryWire = await post( + GRAPHQL, + httpToken, + requestBodies.country, + FILES["http/country-response.json"], + "application/json", + "HTTP country", + ); + const country = countryOutput(decode(countryWire.bytes), window); + bundle.write("http/country-response.json", countryWire.bytes); + bundle.write( + "http/country-receipt.json", + json(httpReceipt(requestBodies.country, countryWire, countryAt)), + ); + bundle.write("http/country.json", json(country)); + const result = verifyBundle(bundle, spec, requestBodies, false); + bundle.write("manifest.json", result.manifest); + return { status: "written", day, summary: result.summary }; + } catch (error) { + if (error instanceof CaptureError) throw error; + throw new CaptureError( + error?.code === "EEXIST" + ? "output already exists; no resume or overwrite" + : error?.code === "ENOENT" + ? "bundle is incomplete or output parent is missing" + : "unable to read or write private bundle; incomplete evidence is retained", + ); + } +} diff --git a/scripts/telemetry-aggregate-capture.mjs b/scripts/telemetry-aggregate-capture.mjs new file mode 100644 index 0000000..598e14f --- /dev/null +++ b/scripts/telemetry-aggregate-capture.mjs @@ -0,0 +1,34 @@ +import { parseArgs } from "node:util"; +import { runTelemetryAggregateCapture } from "./lib/telemetry-aggregate-capture.mjs"; + +const usage = + "Usage: npm run telemetry:capture -- --day YYYY-MM-DD [--execute | --verify] [--output --account-id --zone-id ]"; + +try { + const { values } = parseArgs({ + options: { + day: { type: "string" }, + output: { type: "string" }, + "account-id": { type: "string" }, + "zone-id": { type: "string" }, + execute: { type: "boolean", default: false }, + verify: { type: "boolean", default: false }, + }, + }); + if (!values.day) throw new Error(usage); + console.log( + JSON.stringify( + await runTelemetryAggregateCapture({ + day: values.day, + output: values.output, + accountId: values["account-id"], + zoneId: values["zone-id"], + execute: values.execute, + verify: values.verify, + }), + ), + ); +} catch (error) { + console.error(error?.code?.startsWith("ERR_PARSE_ARGS") ? usage : error.message); + process.exitCode = 1; +} diff --git a/test/telemetry-aggregate-capture.test.mjs b/test/telemetry-aggregate-capture.test.mjs new file mode 100644 index 0000000..27e21fc --- /dev/null +++ b/test/telemetry-aggregate-capture.test.mjs @@ -0,0 +1,727 @@ +import { spawnSync } from "node:child_process"; +import { createHash } from "node:crypto"; +import { + chmodSync, + existsSync, + linkSync, + lstatSync, + mkdirSync, + mkdtempSync, + readFileSync, + readdirSync, + realpathSync, + rmSync, + symlinkSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { dirname, join } from "node:path"; +import { fileURLToPath } from "node:url"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { runTelemetryAggregateCapture } from "../scripts/lib/telemetry-aggregate-capture.mjs"; +import { runTelemetryHistory } from "../scripts/lib/telemetry-history.mjs"; + +const CLI = fileURLToPath(new URL("../scripts/telemetry-aggregate-capture.mjs", import.meta.url)); +const DAY = "2025-02-03"; +const NOW = Date.parse("2025-02-05T12:00:00Z"); +const ACCOUNT = "a".repeat(32); +const ZONE = "b".repeat(32); +const AE_TOKEN = "synthetic-ae-read-token"; +const HTTP_TOKEN = "synthetic-http-read-token"; +const owned = []; +const json = (value) => `${JSON.stringify(value, null, 2)}\n`; +const digest = (value) => createHash("sha256").update(value).digest("hex"); +let fetchMock; + +function options() { + const parent = realpathSync(mkdtempSync(join(tmpdir(), "telemetry-capture-test-"))); + owned.push(parent); + return { + day: DAY, + output: join(parent, "bundle"), + accountId: ACCOUNT, + zoneId: ZONE, + execute: true, + }; +} +function settings() { + return { + data: { + viewer: { + zones: [ + { + settings: { + httpRequestsAdaptiveGroups: { + enabled: true, + notOlderThan: 10 * 86_400, + maxDuration: 86_400, + maxPageSize: 10_000, + maxNumberOfFields: 20, + availableFields: [ + "count", + "avg_sampleInterval", + "dimensions_clientCountryName", + "dimensions_clientRequestHTTPHost", + "dimensions_clientRequestPath", + "dimensions_requestSource", + "dimensions_datetime", + ], + }, + }, + }, + ], + }, + }, + errors: null, + }; +} +function hourly() { + const meta = [ + { name: "bucket", type: "DateTime" }, + { name: "weightedReports", type: "UInt64" }, + { name: "queryRows", type: "UInt64" }, + { name: "featureReports", type: "UInt64" }, + { name: "featureQueryRows", type: "UInt64" }, + { name: "minSampleInterval", type: "UInt32" }, + { name: "maxSampleInterval", type: "UInt32" }, + { name: "latestEventAt", type: "DateTime" }, + { name: "latestFeatureAt", type: "DateTime" }, + ]; + const data = Array.from({ length: 24 }, (_, hour) => { + const prefix = `${DAY} ${String(hour).padStart(2, "0")}`; + return { + bucket: `${prefix}:00:00`, + weightedReports: "6", + queryRows: "3", + featureReports: "2", + featureQueryRows: "1", + minSampleInterval: 2, + maxSampleInterval: 2, + latestEventAt: `${prefix}:59:00`, + latestFeatureAt: `${prefix}:58:00`, + }; + }); + return { meta, data, rows: data.length, rows_before_limit_at_least: data.length }; +} +const group = (country, count, sampleInterval) => ({ + dimensions: { clientCountryName: country }, + count, + avg: { sampleInterval }, +}); +function country( + rows = [group("XX", 7, 9), group(null, 3, 1), group("T1", 0, null), group("", 1, 2)], +) { + return { data: { viewer: { zones: [{ httpRequestsAdaptiveGroups: rows }] } }, errors: null }; +} +const response = (body, init = {}) => + new Response(typeof body === "string" || body instanceof Uint8Array ? body : json(body), init); +function transport({ metadata = settings(), ae = hourly(), http = country(), afterAe } = {}) { + let index = 0; + fetchMock.mockImplementation(async () => { + const next = [metadata, ae, http][index++]; + if (index === 2) afterAe?.(); + if (next === undefined) throw new Error("unexpected extra request"); + return next instanceof Response ? next : response(next); + }); + return fetchMock; +} +const read = (opts, name) => JSON.parse(readFileSync(join(opts.output, name))); +function snapshot(root) { + return Object.fromEntries( + readdirSync(root, { recursive: true }) + .sort() + .map((name) => { + const path = join(root, name); + const stat = lstatSync(path); + return [ + name, + { + mode: stat.mode & 0o777, + mtime: stat.mtimeMs, + hash: stat.isFile() ? digest(readFileSync(path)) : null, + }, + ]; + }), + ); +} +function retainedText(opts) { + return readdirSync(opts.output, { recursive: true }) + .map((name) => join(opts.output, name)) + .filter((path) => lstatSync(path).isFile()) + .map((path) => readFileSync(path, "utf8")) + .join("\n"); +} +function forbidCredentials() { + const original = process; + vi.stubGlobal("process", { + ...original, + env: new Proxy(original.env, { + get(target, key) { + if (key === "TELEMETRY_AE_READ_TOKEN" || key === "TELEMETRY_HTTP_READ_TOKEN") + throw new Error("credentials accessed"); + return Reflect.get(target, key); + }, + }), + }); +} +function repin(opts, name, contents) { + writeFileSync(join(opts.output, name), contents); + const manifest = read(opts, "manifest.json"); + const entry = manifest.files.find((file) => file.file === name); + entry.bytes = Buffer.byteLength(contents); + entry.sha256 = digest(contents); + writeFileSync(join(opts.output, "manifest.json"), json(manifest)); +} + +beforeEach(() => { + vi.useFakeTimers({ toFake: ["Date"] }); + vi.setSystemTime(NOW); + vi.stubEnv("TELEMETRY_AE_READ_TOKEN", AE_TOKEN); + vi.stubEnv("TELEMETRY_HTTP_READ_TOKEN", HTTP_TOKEN); + fetchMock = vi.spyOn(globalThis, "fetch").mockRejectedValue(new Error("unexpected network")); +}); +afterEach(() => { + vi.useRealTimers(); + vi.unstubAllGlobals(); + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + for (const path of owned.splice(0)) rmSync(path, { recursive: true, force: true }); +}); + +describe("planning and exclusive capture", () => { + it("defaults to a closed-day plan without reading credentials, touching output, or requesting data", async () => { + const opts = options(); + opts.output = join(opts.output, "missing-parent", "result"); + forbidCredentials(); + const result = await runTelemetryAggregateCapture({ ...opts, execute: false }); + expect(result).toMatchObject({ + status: "planned", + window: { startInclusive: `${DAY}T00:00:00.000Z`, endExclusive: "2025-02-04T00:00:00.000Z" }, + ae: { structuralMaxRows: 25, sqlLimit: 26 }, + }); + expect(result.ae.sql).toContain("LIMIT 26 FORMAT JSON"); + expect(JSON.stringify(result)).not.toContain(ACCOUNT); + expect(JSON.stringify(result)).not.toContain(ZONE); + expect(existsSync(dirname(dirname(opts.output)))).toBe(false); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it.each(["2025-02-30", "2025-02-05", "2025-02-06", "2025-02-03T00:00:00Z"])( + "rejects an invalid or unclosed day: %s", + async (day) => { + await expect(runTelemetryAggregateCapture({ day })).rejects.toThrow(/day/); + expect(fetchMock).not.toHaveBeenCalled(); + }, + ); + + it.each(["missing HTTP", "shared token", "header injection", "bad identity", "missing parent"])( + "validates local prerequisites before claiming output: %s", + async (kind) => { + const opts = options(); + if (kind === "missing HTTP") vi.stubEnv("TELEMETRY_HTTP_READ_TOKEN", ""); + if (kind === "shared token") vi.stubEnv("TELEMETRY_HTTP_READ_TOKEN", AE_TOKEN); + if (kind === "header injection") vi.stubEnv("TELEMETRY_AE_READ_TOKEN", "bad\r\nheader"); + if (kind === "bad identity") opts.accountId = "../invalid"; + if (kind === "missing parent") opts.output = join(opts.output, "child"); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(); + expect(existsSync(opts.output)).toBe(false); + expect(fetchMock).not.toHaveBeenCalled(); + }, + ); + + it("makes a concurrent output loser fail without a request or repair", async () => { + const opts = options(); + let release; + fetchMock + .mockImplementationOnce( + () => + new Promise((resolve) => { + release = resolve; + }), + ) + .mockResolvedValueOnce(response(hourly())) + .mockResolvedValueOnce(response(country())); + const winner = runTelemetryAggregateCapture(opts); + expect(fetchMock).toHaveBeenCalledTimes(1); + const pending = snapshot(opts.output); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/incomplete/); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(snapshot(opts.output)).toEqual(pending); + release(response(settings())); + await expect(winner).resolves.toMatchObject({ status: "written" }); + expect(fetchMock).toHaveBeenCalledTimes(3); + }); + + it("captures separate estimates, exact SQL representations and a PR17-compatible archive", async () => { + const opts = options(); + const ae = hourly(); + Object.assign(ae.data[0], { + weightedReports: "18446744073709551615", + queryRows: "18446744073709551615", + featureReports: "0", + featureQueryRows: "0", + minSampleInterval: 1, + maxSampleInterval: 1, + latestFeatureAt: "1970-01-01 00:00:00", + }); + transport({ ae }); + const result = await runTelemetryAggregateCapture(opts); + expect(result).toMatchObject({ + status: "written", + summary: { + ae: { + completeClosedDays: 1, + completeClosedTotals: { + weightedReports: "18446744073709551753", + queryRows: "18446744073709551684", + featureReports: "46", + featureQueryRows: "23", + }, + }, + http: { estimatedRequests: "11", countryGroups: 4, coverage: "returned_groups_only" }, + }, + }); + expect(read(opts, "http/country.json").countries).toEqual([ + { country: null, estimatedRequests: "3", sampleInterval: 1 }, + { country: "", estimatedRequests: "1", sampleInterval: 2 }, + { country: "T1", estimatedRequests: "0", sampleInterval: null }, + { country: "XX", estimatedRequests: "7", sampleInterval: 9 }, + ]); + const calls = fetchMock.mock.calls; + expect(calls.map(([url]) => url)).toEqual([ + "https://api.cloudflare.com/client/v4/graphql", + `https://api.cloudflare.com/client/v4/accounts/${ACCOUNT}/analytics_engine/sql`, + "https://api.cloudflare.com/client/v4/graphql", + ]); + expect(calls.map(([, init]) => init.headers.Authorization)).toEqual([ + `Bearer ${HTTP_TOKEN}`, + `Bearer ${AE_TOKEN}`, + `Bearer ${HTTP_TOKEN}`, + ]); + expect(calls.every(([, init]) => init.method === "POST" && init.redirect === "error")).toBe( + true, + ); + const request = JSON.parse(calls[2][1].body); + expect(request.variables).toEqual({ + zoneTag: ZONE, + start: `${DAY}T00:00:00.000Z`, + end: "2025-02-04T00:00:00.000Z", + }); + for (const filter of [ + 'clientRequestHTTPHost: "telemetry.openclaw.ai"', + 'clientRequestPath: "/api/latest-version"', + 'requestSource: "eyeball"', + "datetime_geq: $start", + "datetime_lt: $end", + ]) + expect(request.query).toContain(filter); + const manifest = read(opts, "manifest.json"); + expect(manifest).toMatchObject({ complete: true, completeness: "query_and_wire_only" }); + const sql = readFileSync(join(opts.output, "ae/q2/query.sql"), "utf8"); + expect(sql).toBe(`${calls[1][1].body}\n`); + expect(read(opts, "ae/q2/receipt.json").sqlSha256).toBe(digest(calls[1][1].body)); + expect(digest(sql)).not.toBe(digest(calls[1][1].body)); + const standalone = join(dirname(opts.output), "standalone"); + runTelemetryHistory({ + archive: join(opts.output, "ae"), + query: "q2", + ...manifest.pins, + output: standalone, + }); + for (const name of ["daily.json", "daily.csv", "manifest.json"]) + expect(readFileSync(join(standalone, name))).toEqual( + readFileSync(join(opts.output, "ae-daily", name)), + ); + expect(lstatSync(opts.output).mode & 0o777).toBe(0o700); + for (const [name, stat] of Object.entries(snapshot(opts.output))) { + expect(stat.mode).toBe(stat.hash === null ? 0o700 : 0o600); + if (stat.hash !== null) { + const text = readFileSync(join(opts.output, name), "utf8"); + expect(text).not.toContain(AE_TOKEN); + expect(text).not.toContain(HTTP_TOKEN); + expect(text).not.toContain(opts.output); + } + } + }); + + it.each([0, 23])( + "keeps %s observed hours distinct from complete events or measured zero", + async (hours) => { + const opts = options(); + const ae = hourly(); + ae.data = ae.data.slice(0, hours); + ae.rows = ae.rows_before_limit_at_least = hours; + transport({ ae, http: country([]) }); + await runTelemetryAggregateCapture(opts); + expect(read(opts, "manifest.json").complete).toBe(true); + expect(read(opts, "http/country.json")).toMatchObject({ + coverage: "unknown", + estimatedRequests: null, + countries: [], + }); + const day = read(opts, "ae-daily/daily.json").days[0]; + expect(day).toMatchObject({ + coverage: "missing_hours", + comparisonEligible: false, + observedHours: hours, + weightedReports: hours === 0 ? null : "138", + }); + expect(day.missingHours).toHaveLength(24 - hours); + }, + ); + + it("sums safe HTTP integers exactly without multiplying the sampling diagnostic", async () => { + const opts = options(); + transport({ + http: country([ + group("AA", Number.MAX_SAFE_INTEGER, 5), + group("BB", Number.MAX_SAFE_INTEGER, 7), + ]), + }); + await runTelemetryAggregateCapture(opts); + expect(read(opts, "http/country.json").estimatedRequests).toBe("18014398509481982"); + }); +}); + +describe("offline immutable bundle verification", () => { + it("verifies byte-identical reruns after retention expires without credentials or network", async () => { + const opts = options(); + transport(); + await runTelemetryAggregateCapture(opts); + const original = snapshot(opts.output); + fetchMock.mockClear(); + forbidCredentials(); + vi.setSystemTime("2026-02-05T12:00:00Z"); + await expect(runTelemetryAggregateCapture(opts)).resolves.toMatchObject({ + status: "unchanged", + }); + await expect( + runTelemetryAggregateCapture({ ...opts, execute: false, verify: true }), + ).resolves.toMatchObject({ status: "unchanged" }); + expect(snapshot(opts.output)).toEqual(original); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it.each([ + "partial", + "different day", + "different account", + "different zone", + "changed bytes", + "extra file", + "nonprivate", + "symlink", + "hardlink", + ])("rejects %s output offline without overwriting evidence", async (kind) => { + const opts = options(); + transport(); + await runTelemetryAggregateCapture(opts); + if (kind === "partial") rmSync(join(opts.output, "manifest.json")); + if (kind === "different day") opts.day = "2025-02-02"; + if (kind === "different account") opts.accountId = "c".repeat(32); + if (kind === "different zone") opts.zoneId = "d".repeat(32); + const path = join(opts.output, "http/country.json"); + if (kind === "changed bytes") writeFileSync(path, "changed"); + if (kind === "extra file") writeFileSync(join(opts.output, "extra"), "keep"); + if (kind === "nonprivate") chmodSync(path, 0o644); + if (kind === "symlink" || kind === "hardlink") { + rmSync(path); + (kind === "symlink" ? symlinkSync : linkSync)(join(opts.output, "bundle-plan.json"), path); + } + fetchMock.mockClear(); + forbidCredentials(); + const original = snapshot(opts.output); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow( + /incomplete|conflict|private|unsafe|symlink/, + ); + expect(fetchMock).not.toHaveBeenCalled(); + expect(snapshot(opts.output)).toEqual(original); + }); + + it.each(["daily output", "SQL", "HTTP query", "AE count"])( + "rejects repinned %s tampering by contract or offline regeneration", + async (kind) => { + const opts = options(); + transport(); + await runTelemetryAggregateCapture(opts); + if (kind === "daily output") { + const data = read(opts, "ae-daily/daily.json"); + data.days[0].weightedReports = "999"; + repin(opts, "ae-daily/daily.json", json(data)); + } + if (kind === "SQL") + repin( + opts, + "ae/q2/query.sql", + `${readFileSync(join(opts.output, "ae/q2/query.sql"), "utf8")}\n`, + ); + if (kind === "HTTP query") { + const body = read(opts, "http/country-request.json"); + body.query = body.query.replace('requestSource: "eyeball"', 'requestSource: "all"'); + repin(opts, "http/country-request.json", json(body)); + } + if (kind === "AE count") { + const body = read(opts, "ae/q2/response.json"); + body.data[0].weightedReports = 6; + const encoded = json(body); + repin(opts, "ae/q2/response.json", encoded); + const receipt = read(opts, "ae/q2/receipt.json"); + receipt.wireSha256 = digest(encoded); + receipt.wireBytesRead = Buffer.byteLength(encoded); + repin(opts, "ae/q2/receipt.json", json(receipt)); + } + fetchMock.mockClear(); + forbidCredentials(); + const original = snapshot(opts.output); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow( + /conflict|history validation|AE response/, + ); + expect(snapshot(opts.output)).toEqual(original); + expect(fetchMock).not.toHaveBeenCalled(); + }, + ); + + it("rejects symlink parents and scratch within the bundle without writing into them", async () => { + const opts = options(); + const alias = join(dirname(opts.output), "alias"); + symlinkSync(dirname(opts.output), alias); + await expect( + runTelemetryAggregateCapture({ ...opts, output: join(alias, "bundle") }), + ).rejects.toThrow(/symlink/); + expect(fetchMock).not.toHaveBeenCalled(); + transport(); + await runTelemetryAggregateCapture(opts); + const original = snapshot(opts.output); + fetchMock.mockClear(); + vi.stubEnv("TMPDIR", opts.output); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/scratch/); + expect(snapshot(opts.output)).toEqual(original); + expect(fetchMock).not.toHaveBeenCalled(); + }); +}); + +describe("fail-closed HTTP and AE boundaries", () => { + it.each([ + [ + "disabled", + (s) => { + s.enabled = false; + }, + ], + [ + "unsupported filter", + (s) => { + s.availableFields = s.availableFields.filter((v) => v !== "dimensions_clientRequestPath"); + }, + ], + [ + "expired lookback", + (s) => { + s.notOlderThan = 86_400; + }, + ], + [ + "duration", + (s) => { + s.maxDuration = 86_399; + }, + ], + [ + "page cap", + (s) => { + s.maxPageSize = 9999; + }, + ], + [ + "field cap", + (s) => { + s.maxNumberOfFields = 2; + }, + ], + ])("stops after metadata for %s", async (_name, mutate) => { + const opts = options(); + const metadata = settings(); + mutate(metadata.data.viewer.zones[0].settings.httpRequestsAdaptiveGroups); + transport({ metadata }); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/HTTP/); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + expect(existsSync(join(opts.output, "ae/q2/response.json"))).toBe(false); + }); + + it("rechecks lookback immediately after a delayed AE request and retains incomplete evidence", async () => { + const opts = options(); + const metadata = settings(); + metadata.data.viewer.zones[0].settings.httpRequestsAdaptiveGroups.notOlderThan = + (NOW - Date.parse(`${DAY}T00:00:00Z`)) / 1000 + 2; + transport({ metadata, afterAe: () => vi.setSystemTime(NOW + 3000) }); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/lookback/); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(existsSync(join(opts.output, "ae-daily/manifest.json"))).toBe(true); + expect(existsSync(join(opts.output, "http/country-response.json"))).toBe(false); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + fetchMock.mockClear(); + forbidCredentials(); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/incomplete/); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it.each([ + "redirect", + "expired token", + "GraphQL errors", + "unrequested metadata field", + "oversize header", + "oversize stream", + "truncated", + "invalid UTF-8", + "read failure", + ])("retains an incomplete claim and never retries %s", async (kind) => { + const opts = options(); + let metadata; + if (kind === "redirect") + metadata = response("redirect", { + status: 302, + headers: { Location: "https://example.invalid/" }, + }); + if (kind === "expired token") metadata = response(HTTP_TOKEN, { status: 401 }); + if (kind === "GraphQL errors") + metadata = { data: null, errors: [{ message: `upstream echo: ${HTTP_TOKEN}` }] }; + if (kind === "unrequested metadata field") + metadata = { ...settings(), debug: { echo: HTTP_TOKEN } }; + if (kind === "oversize header") + metadata = response("{}", { headers: { "content-length": String(256 * 1024 + 1) } }); + if (kind === "oversize stream") metadata = response(Buffer.alloc(256 * 1024 + 1, 32)); + if (kind === "truncated") metadata = response("{}", { headers: { "content-length": "3" } }); + if (kind === "invalid UTF-8") metadata = response(Buffer.from([0xff])); + if (kind === "read failure") + metadata = new Response( + new ReadableStream({ + start(controller) { + controller.error(new Error(`transport ${HTTP_TOKEN}`)); + }, + }), + ); + transport({ metadata }); + const error = await runTelemetryAggregateCapture(opts).catch((error) => error); + expect(error).toBeInstanceOf(Error); + if (kind === "expired token") + expect(error.message).toMatch(/^HTTP settings request failed \(HTTP 401\)/); + for (const secret of [HTTP_TOKEN, AE_TOKEN, ZONE, ACCOUNT, opts.output]) + expect(error.message).not.toContain(secret); + expect(retainedText(opts).includes(HTTP_TOKEN)).toBe(false); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + }); + + it("bounds a stalled request and aborts it without retry", async () => { + vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] }); + const opts = options(); + fetchMock.mockImplementation(() => new Promise(() => {})); + const result = expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/timed out/); + await vi.advanceTimersByTimeAsync(45_000); + await result; + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(fetchMock.mock.calls[0][1].signal.aborted).toBe(true); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + }); + + it("does not relax the unchanged AE exporter to complete a bundle", async () => { + const opts = options(); + const ae = hourly(); + ae.data[0].featureReports = "4"; + transport({ ae }); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(/history validation/); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + }); + + it.each(["error envelope", "unrequested row field", "invalid typed field"])( + "rejects upstream AE %s before retaining an echoed token", + async (kind) => { + const opts = options(); + let ae = hourly(); + if (kind === "error envelope") ae = { errors: [{ message: AE_TOKEN }] }; + if (kind === "unrequested row field") ae.data[0].debug = AE_TOKEN; + if (kind === "invalid typed field") ae.data[0].weightedReports = AE_TOKEN; + transport({ ae }); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow(); + expect(retainedText(opts).includes(AE_TOKEN)).toBe(false); + expect(fetchMock).toHaveBeenCalledTimes(2); + expect(existsSync(join(opts.output, "ae/q2/response.json"))).toBe(false); + expect(existsSync(join(opts.output, "ae/q2/receipt.json"))).toBe(false); + expect(existsSync(join(opts.output, "ae/q2/attempt.json"))).toBe(true); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + }, + ); + + it.each([ + "duplicate country", + "unsafe integer", + "fractional count", + "page limit", + "GraphQL errors", + "unrequested country field", + ])("rejects %s without a completion manifest", async (kind) => { + const opts = options(); + let http = country(); + if (kind === "duplicate country") http = country([group(null, 1, 1), group(null, 2, 1)]); + if (kind === "unsafe integer") http = country([group("XX", Number.MAX_SAFE_INTEGER + 1, 1)]); + if (kind === "fractional count") http = country([group("XX", 1.5, 1)]); + if (kind === "page limit") + http = response( + JSON.stringify(country(Array.from({ length: 10_000 }, () => group("XX", 1, 1)))), + ); + if (kind === "GraphQL errors") + http = { data: null, errors: [{ message: `upstream echo: ${HTTP_TOKEN}` }] }; + if (kind === "unrequested country field") + http.data.viewer.zones[0].httpRequestsAdaptiveGroups[0].debug = HTTP_TOKEN; + transport({ http }); + await expect(runTelemetryAggregateCapture(opts)).rejects.toThrow( + /country|integer|GraphQL|schema/, + ); + expect(retainedText(opts).includes(HTTP_TOKEN)).toBe(false); + expect(fetchMock).toHaveBeenCalledTimes(3); + expect(existsSync(join(opts.output, "http/country-response.json"))).toBe(false); + expect(existsSync(join(opts.output, "http/country-receipt.json"))).toBe(false); + expect(existsSync(join(opts.output, "manifest.json"))).toBe(false); + expect(existsSync(join(opts.output, "ae/q2/response.json"))).toBe(true); + }); +}); + +it("runs the real CLI in default-plan and offline-verify modes without leaking arguments", async () => { + const opts = options(); + transport(); + await runTelemetryAggregateCapture(opts); + const original = snapshot(opts.output); + const run = (args) => + spawnSync(process.execPath, [CLI, ...args], { + encoding: "utf8", + timeout: 10_000, + maxBuffer: 65_536, + env: { PATH: dirname(process.execPath), TZ: "Pacific/Honolulu" }, + }); + const planned = run(["--day", DAY]); + expect(planned.status, planned.stderr).toBe(0); + expect(JSON.parse(planned.stdout).status).toBe("planned"); + const verified = run([ + "--verify", + "--day", + DAY, + "--output", + opts.output, + "--account-id", + ACCOUNT, + "--zone-id", + ZONE, + ]); + expect(verified.status, verified.stderr).toBe(0); + expect(JSON.parse(verified.stdout).status).toBe("unchanged"); + expect(snapshot(opts.output)).toEqual(original); + for (const privateValue of [opts.output, ACCOUNT, ZONE, AE_TOKEN, HTTP_TOKEN]) + expect(verified.stdout + verified.stderr).not.toContain(privateValue); + const invalid = run(["--unknown", "private-marker"]); + expect(invalid.status).toBe(1); + expect(invalid.stderr).toMatch(/Usage:/); + expect(invalid.stderr).not.toContain("private-marker"); +});