import fs from 'fs-extra'; import path from 'node:path'; import type { AxiosInstance } from 'axios'; import type { LogEntry, LogsQueryParams, LogsResponse } from '../types.js'; import { getLogs } from '../api.js'; import { EVIDENCE_SCHEMA_VERSION, createRunId, sha256Bytes, writeJsonFile, type EvidenceEnvelope } from './contract.js'; type GetLogsFn = ( client: AxiosInstance, params: LogsQueryParams, signal?: AbortSignal ) => Promise; export interface LogCaptureOptions { maxItems: number; intervalMs: number; timeoutMs: number; maxAttempts: number; getLogsFn?: GetLogsFn; sleepFn?: (ms: number) => Promise; nowFn?: () => number; } export interface CaptureAttempt { attempt: number; started_at: string; completed_at: string; pages_fetched: number; rows_fetched: number; new_rows: number; exhausted: boolean; cap_saturated: boolean; failure?: { page: number; message: string; retryable: boolean; }; } export interface SettledLogCapture { rows: LogEntry[]; attempts: CaptureAttempt[]; settled: boolean; complete: boolean; capSaturated: boolean; timedOut: boolean; attemptsExhausted: boolean; } interface LogCaptureDetails extends Record { query_scope: LogsQueryParams; client_filter: { name: string | null }; row_count: number; empty_capture: boolean; complete: boolean; settled: boolean; total_cap: number; total_cap_saturated: boolean; settle: { interval_ms: number; timeout_ms: number; max_attempts: number; attempts_used: number; timed_out: boolean; attempts_exhausted: boolean; }; attempts: CaptureAttempt[]; } function delay(ms: number): Promise { return new Promise(resolve => setTimeout(resolve, ms)); } function errorMessage(error: unknown): string { if (error instanceof Error) return error.message; return String(error); } function isRetryable(error: unknown): boolean { const status = (error as { response?: { status?: number } }).response?.status; return status === undefined || status === 408 || status === 429 || status >= 500; } function logIdentity(log: LogEntry): string { if (log.log_id) return `id:${log.log_id}`; return `sha256:${sha256Bytes(JSON.stringify(log))}`; } async function fetchSnapshot( client: AxiosInstance, params: LogsQueryParams, maxItems: number, getLogsFn: GetLogsFn, deadlineMs: number, nowFn: () => number ): Promise<{ rows: LogEntry[]; pagesFetched: number; exhausted: boolean; capSaturated: boolean; failure?: CaptureAttempt['failure']; }> { const pageSize = Number.isFinite(params.per) && params.per && params.per > 0 ? params.per : 50; let page = Number.isFinite(params.page) && params.page && params.page > 0 ? params.page : 1; const rows: LogEntry[] = []; let pagesFetched = 0; while (rows.length < maxItems) { const remaining = maxItems - rows.length; const per = Math.min(pageSize, remaining); let response: LogsResponse; const remainingMs = deadlineMs - nowFn(); if (remainingMs <= 0) { return { rows, pagesFetched, exhausted: false, capSaturated: false, failure: { page, message: 'settle timeout before page completed', retryable: true } }; } const controller = new AbortController(); const timeoutId = setTimeout(() => controller.abort(), remainingMs); try { response = await getLogsFn(client, { ...params, page, per }, controller.signal); } catch (error: unknown) { return { rows, pagesFetched, exhausted: false, capSaturated: false, failure: { page, message: errorMessage(error), retryable: isRetryable(error) } }; } finally { clearTimeout(timeoutId); } pagesFetched++; rows.push(...response.items); if (response.items.length < per) { return { rows, pagesFetched, exhausted: true, capSaturated: false }; } page++; } return { rows, pagesFetched, exhausted: false, capSaturated: true }; } export async function captureSettledLogs( client: AxiosInstance, params: LogsQueryParams, options: LogCaptureOptions ): Promise { const getLogsFn = options.getLogsFn ?? getLogs; const sleepFn = options.sleepFn ?? delay; const nowFn = options.nowFn ?? Date.now; const startedMs = nowFn(); const deadlineMs = startedMs + options.timeoutMs; const observed = new Map(); const attempts: CaptureAttempt[] = []; let settled = false; let finalSuccessfulExhausted = false; let capSaturated = false; let timedOut = false; for (let attemptNumber = 1; attemptNumber <= options.maxAttempts; attemptNumber++) { if (attemptNumber > 1 && nowFn() >= deadlineMs) { timedOut = true; break; } const attemptStartedMs = nowFn(); const snapshot = await fetchSnapshot( client, params, options.maxItems, getLogsFn, deadlineMs, nowFn ); let newRows = 0; for (const row of snapshot.rows) { const key = logIdentity(row); if (!observed.has(key)) newRows++; observed.set(key, row); } const attemptCompletedMs = nowFn(); const attempt: CaptureAttempt = { attempt: attemptNumber, started_at: new Date(attemptStartedMs).toISOString(), completed_at: new Date(attemptCompletedMs).toISOString(), pages_fetched: snapshot.pagesFetched, rows_fetched: snapshot.rows.length, new_rows: newRows, exhausted: snapshot.exhausted, cap_saturated: snapshot.capSaturated, ...(snapshot.failure ? { failure: snapshot.failure } : {}) }; attempts.push(attempt); capSaturated ||= snapshot.capSaturated; finalSuccessfulExhausted = !snapshot.failure && snapshot.exhausted; if (!snapshot.failure && attemptNumber > 1 && newRows === 0) { settled = true; break; } if (attemptNumber < options.maxAttempts) { const remainingMs = deadlineMs - nowFn(); if (remainingMs <= 0) { timedOut = true; break; } await sleepFn(Math.min(options.intervalMs, remainingMs)); } } if (!settled && nowFn() >= deadlineMs) timedOut = true; const attemptsExhausted = !settled && attempts.length >= options.maxAttempts; const rows = [...observed.values()].sort((a, b) => new Date(a.datetime).getTime() - new Date(b.datetime).getTime() ); return { rows, attempts, settled, complete: settled && finalSuccessfulExhausted && !capSaturated, capSaturated, timedOut, attemptsExhausted }; } export async function writeLogCapture( capturePath: string, manifestPath: string, capture: SettledLogCapture, context: { account: Record; accountLimitations: string[]; baseUrl: string; params: LogsQueryParams; nameFilter: string | null; options: LogCaptureOptions; startedAt: string; } ): Promise> { const absoluteCapturePath = path.resolve(capturePath); await fs.ensureDir(path.dirname(absoluteCapturePath)); const jsonl = capture.rows.map(row => JSON.stringify(row)).join('\n'); const captureBytes = jsonl ? `${jsonl}\n` : ''; await fs.writeFile(absoluteCapturePath, captureBytes, 'utf8'); const completedAt = new Date().toISOString(); const failures = capture.attempts.filter(attempt => attempt.failure); const hasActiveFault = !capture.complete && failures.length > 0; const limitations: string[] = [...context.accountLimitations]; if (capture.capSaturated) limitations.push('The analytics endpoint exposes no exact total; reaching the caller cap makes the capture incomplete.'); if (!capture.settled) limitations.push('The capture did not reach a no-new-rows settle observation within its explicit bounds.'); if (failures.length > 0) limitations.push('One or more capture attempts had a page or network failure; later attempts may have recovered the rows.'); if (capture.rows.length === 0) limitations.push('An empty capture cannot prove that an event did not fire.'); if (context.nameFilter) limitations.push('The capture file retains unfiltered rows; the --name filter applies only to display output so anchors are preserved.'); const envelope: EvidenceEnvelope = { schema_version: EVIDENCE_SCHEMA_VERSION, result_kind: 'log_capture', run_id: createRunId(), account: context.account, environment: { base_url: context.baseUrl }, phase: 'observe', state: capture.complete && capture.rows.length > 0 ? 'SUCCEEDED' : 'PARTIAL', fault_owner: hasActiveFault ? failures.some(attempt => attempt.failure?.retryable) ? 'platform' : 'client' : !capture.complete ? 'harness' : null, retryable: !capture.complete ? hasActiveFault ? failures.some(attempt => attempt.failure?.retryable) : true : null, ...(hasActiveFault ? { fault_message: failures.map(attempt => `attempt ${attempt.attempt}: ${attempt.failure?.message}`).join('; ') } : !capture.complete ? { fault_message: 'The capture is empty, unsettled, capped, or attempts-exhausted.' } : {}), evidence: [{ kind: 'analytics_log_jsonl', path: absoluteCapturePath, sha256: sha256Bytes(captureBytes), media_type: 'application/x-ndjson', row_count: capture.rows.length }], side_effects: [], cleanup_owner: 'none', timestamps: { started_at: context.startedAt, completed_at: completedAt }, limitations, details: { query_scope: context.params, client_filter: { name: context.nameFilter }, row_count: capture.rows.length, empty_capture: capture.rows.length === 0, complete: capture.complete, settled: capture.settled, total_cap: context.options.maxItems, total_cap_saturated: capture.capSaturated, settle: { interval_ms: context.options.intervalMs, timeout_ms: context.options.timeoutMs, max_attempts: context.options.maxAttempts, attempts_used: capture.attempts.length, timed_out: capture.timedOut, attempts_exhausted: capture.attemptsExhausted }, attempts: capture.attempts } }; await writeJsonFile(manifestPath, envelope); return envelope; }