From 6d9b34cf8de44acbd722670a75b68d6a2db4d586 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Wed, 24 Jun 2026 15:01:19 +0300 Subject: [PATCH 1/2] test: add crash recovery to real-life validation lab --- docs/real-life-validation-plan.md | 1 + scripts/real-life-validation-lab.ts | 168 ++++++++++++++++++++++++++++ 2 files changed, 169 insertions(+) diff --git a/docs/real-life-validation-plan.md b/docs/real-life-validation-plan.md index 759bd4b..3387bc4 100644 --- a/docs/real-life-validation-plan.md +++ b/docs/real-life-validation-plan.md @@ -23,6 +23,7 @@ For each scenario: These runs are useful evidence, but they do not replace the larger targets below: - `bun run validate:real-life -- --keep` retained `/tmp/open-logs-real-life-lab-NJZ95y/real-life-validation-report.json`, proving one isolated token-secured server, CLI project creation, remote CLI live watch, cross-process CLI log/span writes, HTTP exception/metric writes, one captured `logs run` test command with Linux procfs resource and process-tree summaries, one raw-backed indexed `process.resource.peak_rss` metric event mapped to that run, one raw-backed indexed process-tree event mapped to that run, one isolated artifact-producing build command with generated JavaScript and source-map outputs, two raw-backed indexed `artifact` events, two SQLite `artifacts` rows with relative paths and SHA-256 hashes, one SQLite `source_maps` row, one SQLite `source_map_sources` row, source-map JavaScript linkage, source content represented only by a hash in collector-owned evidence, source-map content canaries absent from retained collector-owned report/export/raw/SQLite evidence, forbidden source-map raw/projection container keys absent from retained scans, one real Bun JUnit reporter command captured through `logs run`, one raw-backed indexed `category=test_report` event with relative path, parsed JUnit counts, SHA-256 hash, one live SQLite `test_reports` projection row, zero SQLite `test_cases` rows for the passing report, authenticated API test-report list/get access, local CLI test-report list/get access, real MCP stdio test-report search/get access, event export with raw envelopes, `Last-Event-ID` catch-up for a pre-existing post-anchor event, a live follow-up event after catch-up, and `doctor segments` with 60 checked raw events and 0 unindexed raw events. +- `bun run validate:real-life -- --keep` retained `/tmp/open-logs-real-life-lab-6uRbnp/real-life-validation-report.json` after adding a process-level crash-recovery drill: a child producer appended raw log evidence and exited with code 42 before SQLite indexing, `doctor segments` detected exactly 1 unindexed raw event, `doctor rebuild-index` reindexed 61 raw-backed events, verification returned to `ok` with 0 unindexed raw events, the crash log projection was reconstructed from raw evidence, final export included the crash event raw envelope, and final `doctor segments` reported 61 checked raw events with 0 unindexed raw events. - `bun run validate:stress -- --keep` retained `/tmp/open-logs-high-volume-stress-AtCHHl/high-volume-ingest-stress-report.json`, proving 10 producer processes, 5,000 mixed event records, 81 segment rows/files, no missing expected IDs, no duplicate event IDs, full raw pointer reconstruction before and after rebuild, SQLite integrity and foreign-key checks, `doctor segments`, and `doctor rebuild-index`. - `bun run validate:streams -- --keep` retained `/tmp/open-logs-stream-load-q2aFBk/stream-load-validation-report.json`, proving direct generic SSE delivery, `Last-Event-ID` buffer-miss SQLite catch-up with explicit `buffer_miss_sqlite_catchup`, forced slow-subscriber overflow with explicit `subscriber_queue_overflow`, bounded multi-consumer API SSE fanout with 8 consumers receiving 80 burst events each, remote `logs watch --server`, local `logs watch --events`, real MCP `event_watch` over stdio including raw envelopes and missing-cursor overflow, and `doctor segments` with 141 checked raw events and 0 unindexed raw events. - `bun run validate:dashboard-stream -- --keep` retained `/tmp/open-logs-dashboard-stream-lab-f1Ji0B/dashboard-stream-validation-report.json`, proving real Chromium sessions against the built dashboard served by both the token-secured source API server and `dist/server/index.js` from an extracted npm package, unauthenticated stream blocking, authorized fetch-backed SSE after dashboard token entry, rendered live event records, dashboard paused state with the paused event not rendered before resume, pause/resume reconnect with query `last_event_id`, catch-up of an event written while paused, post-resume live delivery, and `doctor segments` with 6 checked raw events and 0 unindexed raw events. diff --git a/scripts/real-life-validation-lab.ts b/scripts/real-life-validation-lab.ts index 6b1d246..9296941 100644 --- a/scripts/real-life-validation-lab.ts +++ b/scripts/real-life-validation-lab.ts @@ -58,6 +58,12 @@ interface LabReport { event_type: string; message: string | null; }>; + crash_recovery: { + event_id: string; + crash_exit_code: number; + before_rebuild: Record; + rebuild: Record; + }; run_summary: Record; artifact_run_summary: Record; test_report_run_summary: Record; @@ -958,6 +964,86 @@ try { ); const remoteWatchEvents = [...remoteCatchupEvents, ...remoteLiveEvents]; + const crashEventId = `${labId}-crash-raw-before-index`; + expectedEventIds.push(crashEventId); + const crash = await runRawAppendCrashWorker(crashEventId); + commands.push(crash); + assert( + crash.exit_code === 42, + "crash drill producer exited after raw append before indexing", + ); + assert(crash.stderr === "", "crash drill producer wrote no stderr"); + + const preRebuildDoctor = await runCliWithAllowedExitCodes( + "logs doctor segments before crash rebuild", + ["doctor", "segments", "--json"], + [1], + ); + commands.push(preRebuildDoctor); + const preRebuildDoctorResult = JSON.parse(preRebuildDoctor.stdout) as Record< + string, + unknown + >; + assert( + preRebuildDoctorResult.ok === false, + "doctor segments detected crash-appended raw event before rebuild", + ); + assert( + readRequiredNumber(preRebuildDoctorResult, "unindexed_raw_events") === 1, + "doctor segments reported exactly one unindexed crash-drill raw event", + ); + assert( + JSON.stringify(preRebuildDoctorResult.errors ?? []).includes( + "Raw event is not indexed in SQLite", + ), + "doctor segments explained the unindexed crash-drill raw event", + ); + + const crashRebuild = await runCli("logs doctor rebuild-index after crash", [ + "doctor", + "rebuild-index", + "--json", + ]); + commands.push(crashRebuild); + const crashRebuildResult = JSON.parse(crashRebuild.stdout) as Record< + string, + unknown + >; + const crashRebuildStats = objectValue(crashRebuildResult.rebuild); + const crashRebuildVerification = objectValue(crashRebuildResult.verification); + assert( + crashRebuildStats?.errors && + Array.isArray(crashRebuildStats.errors) && + crashRebuildStats.errors.length === 0, + "crash rebuild completed without rebuild errors", + ); + assert( + crashRebuildVerification?.ok === true, + "crash rebuild verification returned to ok", + ); + assert( + readRequiredNumber( + crashRebuildVerification ?? {}, + "unindexed_raw_events", + ) === 0, + "crash rebuild cleared all unindexed raw events", + ); + const crashRecord = readEventRecord(dbPath, crashEventId); + assert( + crashRecord?.event_type === "log", + "crash rebuild indexed the raw-appended log event", + ); + assert( + readLogMessage(dbPath, crashEventId) === "producer died after raw append", + "crash rebuild reconstructed the log projection from raw evidence", + ); + const crashRecovery = { + event_id: crashEventId, + crash_exit_code: crash.exit_code, + before_rebuild: preRebuildDoctorResult, + rebuild: crashRebuildResult, + }; + const exportFile = join(dataDir, "events-export.json"); commands.push( await runCli("logs events export", [ @@ -1039,6 +1125,7 @@ try { commands, streamed_events: streamedEvents, remote_watch_events: remoteWatchEvents, + crash_recovery: crashRecovery, run_summary: runSummary, artifact_run_summary: artifactRunSummary, test_report_run_summary: testReportRunSummary, @@ -1155,6 +1242,21 @@ async function runCli(label: string, args: string[]): Promise { ); } +async function runCliWithAllowedExitCodes( + label: string, + args: string[], + allowedExitCodes: number[], +): Promise { + const child = spawnCli(args); + return await waitForProcess( + label, + [process.execPath, "src/cli/index.ts", ...args], + child, + 20_000, + allowedExitCodes, + ); +} + function spawnCli( args: string[], ): ReturnType> { @@ -1167,6 +1269,60 @@ function spawnCli( }); } +async function runRawAppendCrashWorker( + eventId: string, +): Promise { + const script = ` + import { getDb } from "./src/db/index.ts"; + import { appendRawEvent } from "./src/lib/event-store.ts"; + + const db = getDb(); + const now = new Date().toISOString(); + appendRawEvent(db, { + schema_version: 1, + event_id: ${JSON.stringify(eventId)}, + source_event_id: "producer-crash-raw-before-index", + event_time: now, + ingest_time: now, + type: "log", + source: "sdk", + severity: "error", + privacy: "internal", + message: "producer died after raw append", + body: { + log: { + id: ${JSON.stringify(eventId)}, + timestamp: now, + level: "error", + source: "sdk", + service: "crash-drill", + message: "producer died after raw append", + metadata: { crash_drill: true }, + }, + }, + attributes: { + service: "crash-drill", + privacy_tier: "internal", + }, + }); + process.exit(42); + `; + const child = Bun.spawn([process.execPath, "-e", script], { + cwd: repoRoot, + env, + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + return await waitForProcess( + "raw append crash worker", + [process.execPath, "-e", "/* raw append crash worker */"], + child, + 20_000, + [42], + ); +} + async function waitForProcess( label: string, command: string[], @@ -1508,6 +1664,18 @@ function readEventRecord( } } +function readLogMessage(dbFile: string, logId: string): string | null { + const db = new Database(dbFile, { readonly: true }); + try { + const row = db + .prepare("SELECT message FROM logs WHERE id = ?") + .get(logId) as { message: string | null } | null; + return row?.message ?? null; + } finally { + db.close(); + } +} + function readArtifactRow( dbFile: string, artifactId: string, From e03aae5a9a61b687415e546ae858b0c8d4d0567f Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Wed, 24 Jun 2026 17:16:08 +0300 Subject: [PATCH 2/2] test: validate collector restart during ingest --- docs/real-life-validation-plan.md | 3 +- scripts/real-life-validation-lab.ts | 243 ++++++++++++++++++++++++++++ 2 files changed, 245 insertions(+), 1 deletion(-) diff --git a/docs/real-life-validation-plan.md b/docs/real-life-validation-plan.md index 3387bc4..ccbd600 100644 --- a/docs/real-life-validation-plan.md +++ b/docs/real-life-validation-plan.md @@ -1,6 +1,6 @@ # Real-Life Telemetry Validation Plan -Last updated: 2026-06-19 +Last updated: 2026-06-24 This plan defines the evidence required before open-logs can be considered a robust universal telemetry data substrate. @@ -24,6 +24,7 @@ These runs are useful evidence, but they do not replace the larger targets below - `bun run validate:real-life -- --keep` retained `/tmp/open-logs-real-life-lab-NJZ95y/real-life-validation-report.json`, proving one isolated token-secured server, CLI project creation, remote CLI live watch, cross-process CLI log/span writes, HTTP exception/metric writes, one captured `logs run` test command with Linux procfs resource and process-tree summaries, one raw-backed indexed `process.resource.peak_rss` metric event mapped to that run, one raw-backed indexed process-tree event mapped to that run, one isolated artifact-producing build command with generated JavaScript and source-map outputs, two raw-backed indexed `artifact` events, two SQLite `artifacts` rows with relative paths and SHA-256 hashes, one SQLite `source_maps` row, one SQLite `source_map_sources` row, source-map JavaScript linkage, source content represented only by a hash in collector-owned evidence, source-map content canaries absent from retained collector-owned report/export/raw/SQLite evidence, forbidden source-map raw/projection container keys absent from retained scans, one real Bun JUnit reporter command captured through `logs run`, one raw-backed indexed `category=test_report` event with relative path, parsed JUnit counts, SHA-256 hash, one live SQLite `test_reports` projection row, zero SQLite `test_cases` rows for the passing report, authenticated API test-report list/get access, local CLI test-report list/get access, real MCP stdio test-report search/get access, event export with raw envelopes, `Last-Event-ID` catch-up for a pre-existing post-anchor event, a live follow-up event after catch-up, and `doctor segments` with 60 checked raw events and 0 unindexed raw events. - `bun run validate:real-life -- --keep` retained `/tmp/open-logs-real-life-lab-6uRbnp/real-life-validation-report.json` after adding a process-level crash-recovery drill: a child producer appended raw log evidence and exited with code 42 before SQLite indexing, `doctor segments` detected exactly 1 unindexed raw event, `doctor rebuild-index` reindexed 61 raw-backed events, verification returned to `ok` with 0 unindexed raw events, the crash log projection was reconstructed from raw evidence, final export included the crash event raw envelope, and final `doctor segments` reported 61 checked raw events with 0 unindexed raw events. +- `bun run validate:real-life -- --keep` retained `/tmp/open-logs-real-life-lab-ku8uZn/real-life-validation-report.json` after adding a true collector-restart/high-volume drill: the lab stopped the live collector process, queued 120 Pino structured records through the SDK transport against the down collector, observed a failed flush, persisted 120 records to the SDK file spool, restarted the collector on the same port and data dir, loaded and replayed all 120 spooled records, cleared the spool file, verified 120 matching SQLite log records with raw segment pointers and SHA-256 record hashes, sampled replayed raw envelopes in export verification, and final `doctor segments` reported 181 checked raw events with 0 unindexed raw events. - `bun run validate:stress -- --keep` retained `/tmp/open-logs-high-volume-stress-AtCHHl/high-volume-ingest-stress-report.json`, proving 10 producer processes, 5,000 mixed event records, 81 segment rows/files, no missing expected IDs, no duplicate event IDs, full raw pointer reconstruction before and after rebuild, SQLite integrity and foreign-key checks, `doctor segments`, and `doctor rebuild-index`. - `bun run validate:streams -- --keep` retained `/tmp/open-logs-stream-load-q2aFBk/stream-load-validation-report.json`, proving direct generic SSE delivery, `Last-Event-ID` buffer-miss SQLite catch-up with explicit `buffer_miss_sqlite_catchup`, forced slow-subscriber overflow with explicit `subscriber_queue_overflow`, bounded multi-consumer API SSE fanout with 8 consumers receiving 80 burst events each, remote `logs watch --server`, local `logs watch --events`, real MCP `event_watch` over stdio including raw envelopes and missing-cursor overflow, and `doctor segments` with 141 checked raw events and 0 unindexed raw events. - `bun run validate:dashboard-stream -- --keep` retained `/tmp/open-logs-dashboard-stream-lab-f1Ji0B/dashboard-stream-validation-report.json`, proving real Chromium sessions against the built dashboard served by both the token-secured source API server and `dist/server/index.js` from an extracted npm package, unauthenticated stream blocking, authorized fetch-backed SSE after dashboard token entry, rendered live event records, dashboard paused state with the paused event not rendered before resume, pause/resume reconnect with query `last_event_id`, catch-up of an event written while paused, post-resume live delivery, and `doctor segments` with 6 checked raw events and 0 unindexed raw events. diff --git a/scripts/real-life-validation-lab.ts b/scripts/real-life-validation-lab.ts index 9296941..0b4cdb7 100644 --- a/scripts/real-life-validation-lab.ts +++ b/scripts/real-life-validation-lab.ts @@ -4,6 +4,7 @@ import { existsSync, mkdirSync, mkdtempSync, + readdirSync, readFileSync, rmSync, writeFileSync, @@ -14,6 +15,7 @@ import { dirname, join, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import { Client } from "@modelcontextprotocol/sdk/client/index.js"; import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; +import { createPinoOpenLogsTransport } from "../sdk/src/index.ts"; const repoRoot = resolve(dirname(fileURLToPath(import.meta.url)), ".."); @@ -64,6 +66,17 @@ interface LabReport { before_rebuild: Record; rebuild: Record; }; + collector_restart_high_volume: { + event_count: number; + spool_file: string; + spooled_records: number; + spool_loaded: number; + replay_sent: number; + matching_logs: number; + raw_backed_events: number; + sample_event_ids: string[]; + spool_file_cleared: boolean; + }; run_summary: Record; artifact_run_summary: Record; test_report_run_summary: Record; @@ -964,6 +977,19 @@ try { ); const remoteWatchEvents = [...remoteCatchupEvents, ...remoteLiveEvents]; + const collectorRestartHighVolume = + await runCollectorRestartHighVolumeScenario({ + baseUrl, + port, + token, + projectId, + labId, + dataDir, + dbPath, + eventCount: 120, + }); + expectedEventIds.push(...collectorRestartHighVolume.sample_event_ids); + const crashEventId = `${labId}-crash-raw-before-index`; expectedEventIds.push(crashEventId); const crash = await runRawAppendCrashWorker(crashEventId); @@ -1126,6 +1152,7 @@ try { streamed_events: streamedEvents, remote_watch_events: remoteWatchEvents, crash_recovery: crashRecovery, + collector_restart_high_volume: collectorRestartHighVolume, run_summary: runSummary, artifact_run_summary: artifactRunSummary, test_report_run_summary: testReportRunSummary, @@ -1472,6 +1499,222 @@ function mcpToolText(result: unknown): string { return typeof text === "string" ? text : ""; } +async function runCollectorRestartHighVolumeScenario(input: { + baseUrl: string; + port: number; + token: string; + projectId: string; + labId: string; + dataDir: string; + dbPath: string; + eventCount: number; +}): Promise { + assert(server, "collector restart drill starts with a live server process"); + const spoolDirectory = join(input.dataDir, "collector-restart-sdk-spool"); + const messagePrefix = `${input.labId} collector restart high-volume`; + + await stopServer(server); + server = undefined; + await sleep(100); + + const stoppedTransport = createPinoOpenLogsTransport({ + url: input.baseUrl, + apiKey: input.token, + projectId: input.projectId, + service: "real-life-restart-producer", + environment: "test", + maxBatchSize: input.eventCount + 1, + maxQueueSize: input.eventCount + 10, + maxRetries: 0, + retryBaseDelayMs: 0, + flushIntervalMs: 60_000, + sourceEventPrefix: `${input.labId}:collector-restart`, + spoolDirectory, + metadata: { + validation_id: input.labId, + scenario: "collector_restart_high_volume", + }, + }); + for (let index = 0; index < input.eventCount; index += 1) { + stoppedTransport.write( + `${JSON.stringify({ + level: 30, + time: Date.now() + index, + msg: `${messagePrefix} ${String(index).padStart(3, "0")}`, + traceId: `${input.labId}-collector-restart-trace`, + sequence: index, + validation_id: input.labId, + restart_drill: true, + })}\n`, + ); + } + assert( + stoppedTransport.stats().pending === input.eventCount, + "collector restart drill queued high-volume records while collector was stopped", + ); + let failedFlush = false; + try { + await stoppedTransport.flush(); + } catch { + failedFlush = true; + } + assert( + failedFlush, + "collector restart drill observed a failed flush while collector was stopped", + ); + assert( + stoppedTransport.stats().spool_pending === input.eventCount, + "collector restart drill retained all stopped-collector records in the SDK spool", + ); + stoppedTransport.stop(); + await sleep(250); + + const spoolFile = findStructuredSpoolFile(spoolDirectory); + const spooledRecords = countSpoolRecords(spoolFile); + assert( + spooledRecords === input.eventCount, + "collector restart drill wrote every high-volume record to the SDK spool file", + ); + + server = Bun.spawn([process.execPath, "src/server/index.ts"], { + cwd: repoRoot, + env: { ...env, LOGS_PORT: String(input.port) }, + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + await waitForHealth(input.baseUrl, server); + + const replayTransport = createPinoOpenLogsTransport({ + url: input.baseUrl, + apiKey: input.token, + projectId: input.projectId, + service: "wrong-restart-service", + environment: "wrong-restart-env", + maxBatchSize: input.eventCount + 1, + maxQueueSize: input.eventCount + 10, + maxRetries: 0, + retryBaseDelayMs: 0, + flushIntervalMs: 60_000, + sourceEventPrefix: "wrong-restart-prefix", + spoolDirectory, + }); + const loadedStats = replayTransport.stats(); + assert( + loadedStats.spool_loaded === input.eventCount, + "collector restart drill loaded every spooled record after process restart", + ); + await replayTransport.flush(); + const replayStats = replayTransport.stats(); + replayTransport.stop(); + assert( + replayStats.sent === input.eventCount, + "collector restart drill replayed every high-volume record after restart", + ); + assert( + replayStats.pending === 0, + "collector restart drill left no pending records after replay", + ); + const spoolFileCleared = !structuredSpoolFileExists(spoolDirectory); + assert( + spoolFileCleared, + "collector restart drill cleared the SDK spool file after replay", + ); + + const replayRows = readCollectorRestartRows(input.dbPath, messagePrefix); + assert( + replayRows.length === input.eventCount, + "collector restart drill indexed every replayed high-volume log", + ); + const rawBackedEvents = replayRows.filter( + (row) => + row.segment_path.length > 0 && + row.byte_length > 0 && + /^[a-f0-9]{64}$/.test(row.record_hash), + ).length; + assert( + rawBackedEvents === input.eventCount, + "collector restart drill preserved raw segment pointers for every replayed high-volume log", + ); + const sampleEventIds = Array.from( + new Set( + [ + ...replayRows.slice(0, 3).map((row) => row.event_id), + replayRows.at(-1)?.event_id, + ].filter((eventId): eventId is string => Boolean(eventId)), + ), + ); + + return { + event_count: input.eventCount, + spool_file: spoolFile, + spooled_records: spooledRecords, + spool_loaded: loadedStats.spool_loaded, + replay_sent: replayStats.sent, + matching_logs: replayRows.length, + raw_backed_events: rawBackedEvents, + sample_event_ids: sampleEventIds, + spool_file_cleared: spoolFileCleared, + }; +} + +function findStructuredSpoolFile(spoolDirectory: string): string { + assert(existsSync(spoolDirectory), "SDK spool directory exists"); + const spoolFileName = readdirSync(spoolDirectory).find((file) => + file.endsWith("structured-spool.jsonl"), + ); + assert(spoolFileName, "SDK structured spool file exists"); + return join(spoolDirectory, spoolFileName); +} + +function structuredSpoolFileExists(spoolDirectory: string): boolean { + return ( + existsSync(spoolDirectory) && + readdirSync(spoolDirectory).some((file) => + file.endsWith("structured-spool.jsonl"), + ) + ); +} + +function countSpoolRecords(spoolFile: string): number { + return readFileSync(spoolFile, "utf8") + .split(/\n/) + .filter((line) => line.trim().length > 0).length; +} + +function readCollectorRestartRows( + dbFile: string, + messagePrefix: string, +): Array<{ + event_id: string; + segment_path: string; + byte_length: number; + record_hash: string; +}> { + const db = new Database(dbFile, { readonly: true }); + try { + return db + .prepare( + ` + SELECT + event_id, segment_path, byte_length, record_hash + FROM event_records + WHERE event_type = 'log' + AND message LIKE ? + ORDER BY message + `, + ) + .all(`${messagePrefix}%`) as Array<{ + event_id: string; + segment_path: string; + byte_length: number; + record_hash: string; + }>; + } finally { + db.close(); + } +} + function assertNoRawTestReportPayload(value: unknown, label: string): void { const text = JSON.stringify(value).toLowerCase(); assert(