diff --git a/.github/workflows/review-reliability-observer.yml b/.github/workflows/review-reliability-observer.yml new file mode 100644 index 0000000000..75a7ad15d7 --- /dev/null +++ b/.github/workflows/review-reliability-observer.yml @@ -0,0 +1,33 @@ +name: Observe review reliability + +on: + workflow_run: + workflows: [ClawSweeper] + types: [completed] + +permissions: + actions: read + contents: read + +concurrency: + group: review-reliability-observer-${{ github.event.workflow_run.id }}-${{ github.event.workflow_run.run_attempt }} + cancel-in-progress: false + +jobs: + observe: + name: Observe terminal ClawSweeper review run + runs-on: ubuntu-latest + timeout-minutes: 5 + steps: + - uses: actions/checkout@v7 + # workflow_run may expose secrets, so execute observer code only from the trusted default branch. + with: + ref: ${{ github.event.repository.default_branch }} + persist-credentials: false + + - name: Record terminal review wave + env: + CLAWSWEEPER_WEBHOOK_SECRET: ${{ secrets.CLAWSWEEPER_WEBHOOK_SECRET }} + GH_TOKEN: ${{ github.token }} + QUEUE_URL: ${{ vars.CLAWSWEEPER_EXACT_REVIEW_QUEUE_URL || 'https://clawsweeper.openclaw.ai' }} + run: node scripts/review-run-observer.mjs --event-file "$GITHUB_EVENT_PATH" diff --git a/dashboard/dashboard-health.ts b/dashboard/dashboard-health.ts index 3e8d8ad0fb..42fadc97dd 100644 --- a/dashboard/dashboard-health.ts +++ b/dashboard/dashboard-health.ts @@ -45,6 +45,19 @@ export function summarizeDashboardHealth(snapshot: Record): Das else if (reviewTelemetryStatus !== "healthy") { raise("amber", "review_telemetry_unavailable"); } + + if (queue.review_execution_health != null) { + const reviewExecution = objectValue(queue.review_execution_health); + const reviewExecutionStatus = String(reviewExecution.health || ""); + // Passive mode is a deliberate pre-producer rollout state and must not page operators + // before PR 674 turns enforcement on. + if (reviewExecutionStatus === "critical") raise("red", "review_execution_critical"); + else if (reviewExecutionStatus === "degraded") { + raise("amber", "review_execution_degraded"); + } else if (!["healthy", "passive"].includes(reviewExecutionStatus)) { + raise("amber", "review_execution_degraded"); + } + } } const operationalStatus = String(objectValue(snapshot.operational_health).status || ""); diff --git a/dashboard/exact-review-queue.ts b/dashboard/exact-review-queue.ts index a8b015527e..771c41dd89 100644 --- a/dashboard/exact-review-queue.ts +++ b/dashboard/exact-review-queue.ts @@ -11,6 +11,14 @@ import { type ReviewTelemetryHealth, normalizeReviewTelemetry, } from "./review-telemetry.ts"; +import { + REVIEW_OBSERVABILITY_RANGES, + summarizeReviewObservability, +} from "./review-observability.ts"; +import { + type DurableReviewRunTelemetry, + normalizeReviewRunTelemetry, +} from "./review-run-telemetry.ts"; type GithubAppJsonOptions = { method?: string; body?: BodyInit; errorLabel?: string }; const GITHUB_TIMEOUT_MS = 4500; @@ -272,6 +280,8 @@ const EXACT_REVIEW_QUEUE_METRICS_TABLE = "exact_review_queue_metrics"; const EXACT_REVIEW_QUEUE_METRIC_BUCKET_TABLE = "exact_review_queue_metric_buckets"; const EXACT_REVIEW_QUEUE_DEAD_LETTER_TABLE = "exact_review_queue_dead_letters"; const EXACT_REVIEW_REVIEW_TELEMETRY_TABLE = "exact_review_review_telemetry"; +const EXACT_REVIEW_RUN_TELEMETRY_TABLE = "exact_review_run_telemetry"; +const REVIEW_OBSERVABILITY_SCAN_LIMIT = 10_000; const EXACT_REVIEW_QUEUE_DEAD_LETTER_LIMIT = 5_000; const EXACT_REVIEW_QUEUE_DEAD_LETTER_RESOLVED_TTL_MS = 30 * 24 * 60 * 60 * 1000; const EXACT_REVIEW_QUEUE_METRIC_BUCKET_MS = 5 * 60 * 1000; @@ -860,10 +870,18 @@ export class ExactReviewQueue { return this.recordReviewTelemetry(await request.json().catch(() => null)); } + if (request.method === "POST" && url.pathname === "/review-run-telemetry") { + return this.recordReviewRunTelemetry(await request.json().catch(() => null)); + } + if (request.method === "GET" && url.pathname === "/review-telemetry") { return this.listReviewTelemetry(url.searchParams); } + if (request.method === "GET" && url.pathname === "/review-observability") { + return this.reviewObservability(url.searchParams); + } + if (request.method === "GET" && url.pathname === "/item-status") { const targetRepo = String(url.searchParams.get("target_repo") || "").trim(); const itemNumber = Number(url.searchParams.get("item_number")); @@ -974,6 +992,7 @@ export class ExactReviewQueue { this.pruneDeliveryReceiptsSync(now); this.pruneQueueTelemetrySync(now); this.pruneReviewTelemetrySync(now); + this.reconcileStoredReviewRunsSync(now); const current = this.readStateSync(); // Dashboard reads are also the operational heartbeat. Reclaim leases and // restore the alarm here so a deploy or lost alarm cannot strand backlog. @@ -992,10 +1011,22 @@ export class ExactReviewQueue { publicationFlow: this.publicationFlowSummarySync(now), deadLetters: this.deadLetterStatsSync(), reviewTelemetryHealth: this.reviewTelemetryHealthSync(now), + reviewExecutionHealth: this.reviewObservabilitySync({ + range: "24h", + repo: null, + now, + }), }; }); - const { state, metrics, reviewFlow, publicationFlow, deadLetters, reviewTelemetryHealth } = - snapshot; + const { + state, + metrics, + reviewFlow, + publicationFlow, + deadLetters, + reviewTelemetryHealth, + reviewExecutionHealth, + } = snapshot; const publicationControl = this.refreshPublicationControlSync(state, now); await this.scheduleNext(state, now); const stats = exactReviewQueueStats( @@ -1047,6 +1078,7 @@ export class ExactReviewQueue { }, delivery_receipts: this.deliveryReceiptCountSync(), review_telemetry_health: reviewTelemetryHealth, + review_execution_health: reviewExecutionHealth, storage_schema_version: EXACT_REVIEW_QUEUE_STORAGE_SCHEMA_VERSION, legacy_rollback_available: !this.legacyMirrorDisabled && @@ -1064,6 +1096,7 @@ export class ExactReviewQueue { await this.storage.deleteAlarm(); this.storage.transactionSync(() => { this.pruneDeliveryReceiptsSync(startedAt); + this.reconcileStoredReviewRunsSync(startedAt); this.syncLegacyCompatibilitySync(this.readStateSync()); }); const snapshot = this.readStateSync(); @@ -1564,24 +1597,37 @@ export class ExactReviewQueue { }); } - private recordReviewTelemetry(value: unknown) { + private async recordReviewTelemetry(value: unknown) { const record = normalizeReviewTelemetry(value); if (!record) return json({ error: "invalid_review_telemetry" }, 400); + const now = Date.now(); const updatedAt = Date.parse(record.updated_at); this.storage.transactionSync(() => { - this.pruneReviewTelemetrySync(Date.now()); + this.pruneReviewTelemetrySync(now); // Terminal truth is first-writer immutable. Retries may replay the same // payload, but neither a heartbeat nor a conflicting terminal delivery // may make durable observations depend on arrival order. this.storage.sql.exec( `INSERT INTO ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} - (repo, item_number, run_id, run_attempt, status, updated_at, lease_expires_at, - record_json) - VALUES (?, ?, ?, ?, ?, ?, ?, ?) + (repo, item_number, run_id, run_attempt, status, outcome, trigger_lane, + trigger_origin, terminal_at, updated_at, lease_expires_at, generation, + operation_id, queue_ms, claim_ms, review_ms, publication_ms, total_ms, record_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(repo, item_number, run_id, run_attempt) DO UPDATE SET status = excluded.status, + outcome = excluded.outcome, + trigger_lane = excluded.trigger_lane, + trigger_origin = excluded.trigger_origin, + terminal_at = excluded.terminal_at, updated_at = excluded.updated_at, lease_expires_at = excluded.lease_expires_at, + generation = excluded.generation, + operation_id = excluded.operation_id, + queue_ms = excluded.queue_ms, + claim_ms = excluded.claim_ms, + review_ms = excluded.review_ms, + publication_ms = excluded.publication_ms, + total_ms = excluded.total_ms, record_json = excluded.record_json WHERE ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE}.status != 'completed' AND (excluded.status = 'completed' @@ -1591,14 +1637,168 @@ export class ExactReviewQueue { record.run_id, record.run_attempt, record.status, + record.outcome, + record.trigger_lane ?? null, + record.trigger_origin ?? null, + record.terminal_at ? Date.parse(record.terminal_at) : null, updatedAt, record.lease_expires_at === null ? null : Date.parse(record.lease_expires_at), + record.generation ?? null, + record.operation_id ?? null, + record.phase_durations_ms.queue ?? null, + record.phase_durations_ms.claim ?? null, + record.phase_durations_ms.review ?? null, + record.phase_durations_ms.publication ?? null, + record.phase_durations_ms.total ?? null, JSON.stringify(record), ); + // workflow_run can arrive before a delayed producer write. Re-check the + // durable terminal evidence here so delivery order cannot strand a row. + this.reconcileStoredReviewRunsSync(now); }); + await this.scheduleNext(this.readStateSync(), now); return json({ ok: true }); } + private async recordReviewRunTelemetry(value: unknown) { + const record = normalizeReviewRunTelemetry(value); + if (!record) return json({ error: "invalid_review_run_telemetry" }, 400); + const completedAt = Date.parse(record.completed_at); + this.storage.transactionSync(() => { + this.pruneReviewTelemetrySync(Date.now()); + // workflow_run deliveries can be replayed, but GitHub's first terminal tuple is immutable. + this.storage.sql.exec( + `INSERT OR IGNORE INTO ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} + (run_id, run_attempt, workflow_outcome, trigger_lane, trigger_origin, target_repo, + completed_at, record_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, + record.run_id, + record.run_attempt, + record.workflow_outcome, + record.trigger_lane, + record.trigger_origin, + record.target_repo, + completedAt, + JSON.stringify(record), + ); + const storedRow = Array.from( + this.storage.sql.exec( + `SELECT record_json FROM ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} + WHERE run_id = ? AND run_attempt = ?`, + record.run_id, + record.run_attempt, + ), + )[0] as { record_json?: unknown } | undefined; + const storedRecord = normalizeReviewRunTelemetry( + JSON.parse(String(storedRow?.record_json || "null")), + ); + if (storedRecord) this.reconcileReviewTelemetryFromRunSync(storedRecord, Date.now()); + }); + await this.scheduleNext(this.readStateSync(), Date.now()); + return json({ ok: true }); + } + + private reconcileReviewTelemetryFromRunSync(run: DurableReviewRunTelemetry, now: number) { + const rows = this.reviewTelemetryRowsSync({ + runId: run.run_id, + runAttempt: run.run_attempt, + status: "refreshing", + }); + for (const record of rows) { + if (record.lease_expires_at !== null && Date.parse(record.lease_expires_at) > now) { + continue; + } + const repoAttributed = run.target_repo === record.repo; + const itemJob = repoAttributed + ? run.review_jobs?.find((job) => job.item_number === record.item_number) + : undefined; + const onlyJob = run.review_jobs?.length === 1 ? run.review_jobs[0] : undefined; + // A generic matrix job is item evidence only when the entire wave has one + // item. Applying one arbitrary shard conclusion to siblings is less safe + // than falling back to the immutable workflow conclusion. + const attributableJob = + itemJob ?? + (repoAttributed && run.item_count === 1 && onlyJob?.item_number === null + ? onlyJob + : undefined); + const outcome = attributableJob + ? attributableJob.conclusion === "success" + ? "succeeded" + : attributableJob.conclusion === "cancelled" + ? "cancelled" + : "interrupted" + : null; + // A workflow terminal proves wave health, not which unattributed matrix + // item succeeded or failed. Keep that row visible to the watchdog. + if (outcome === null) continue; + const terminal: DurableReviewTelemetry = { + ...record, + status: "completed", + outcome, + updated_at: run.completed_at, + lease_expires_at: null, + terminal_at: run.completed_at, + terminal_reason: + outcome === "succeeded" + ? "workflow_job_succeeded" + : outcome === "cancelled" + ? "workflow_cancelled" + : "workflow_terminal", + }; + this.recordReviewTelemetrySync(terminal); + } + } + + private reconcileStoredReviewRunsSync(now: number) { + const rows = this.storage.sql.exec( + `SELECT DISTINCT runs.record_json + FROM ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} AS runs + JOIN ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} AS reviews + ON reviews.run_id = runs.run_id AND reviews.run_attempt = runs.run_attempt + WHERE reviews.status = 'refreshing' + AND (reviews.lease_expires_at IS NULL OR reviews.lease_expires_at <= ?)`, + now, + ) as Iterable<{ record_json?: unknown }>; + for (const row of rows) { + const run = normalizeReviewRunTelemetry(JSON.parse(String(row.record_json || "null"))); + if (run) this.reconcileReviewTelemetryFromRunSync(run, now); + } + } + + private nextReviewReconcileAtSync(now: number) { + const row = Array.from( + this.storage.sql.exec( + `SELECT MIN(reviews.lease_expires_at) AS next_at + FROM ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} AS reviews + JOIN ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} AS runs + ON runs.run_id = reviews.run_id AND runs.run_attempt = reviews.run_attempt + WHERE reviews.status = 'refreshing' + AND reviews.lease_expires_at > ?`, + now, + ), + )[0] as { next_at?: number } | undefined; + const next = Number(row?.next_at || 0); + return next > 0 ? next : null; + } + + private recordReviewTelemetrySync(record: DurableReviewTelemetry) { + this.storage.sql.exec( + `UPDATE ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} + SET status = 'completed', outcome = ?, terminal_at = ?, updated_at = ?, + lease_expires_at = NULL, record_json = ? + WHERE repo = ? AND item_number = ? AND run_id = ? AND run_attempt = ? + AND status != 'completed'`, + record.outcome, + Date.parse(record.terminal_at ?? record.updated_at), + Date.parse(record.updated_at), + JSON.stringify(record), + record.repo, + record.item_number, + record.run_id, + record.run_attempt, + ); + } + private listReviewTelemetry(search: URLSearchParams) { const repo = String(search.get("repo") || "").trim(); const itemNumber = Number(search.get("item_number")); @@ -1625,6 +1825,8 @@ export class ExactReviewQueue { private reviewTelemetryRowsSync(options: { repo?: string; itemNumber?: number; + runId?: string; + runAttempt?: number; status?: DurableReviewTelemetry["status"]; limit?: number; }) { @@ -1638,6 +1840,14 @@ export class ExactReviewQueue { predicates.push("item_number = ?"); bindings.push(options.itemNumber); } + if (options.runId !== undefined) { + predicates.push("run_id = ?"); + bindings.push(options.runId); + } + if (options.runAttempt !== undefined) { + predicates.push("run_attempt = ?"); + bindings.push(options.runAttempt); + } if (options.status !== undefined) { predicates.push("status = ?"); bindings.push(options.status); @@ -1656,6 +1866,98 @@ export class ExactReviewQueue { .filter((record): record is DurableReviewTelemetry => record !== null); } + private reviewRunTelemetryRowsSync(options: { repo?: string; from: number; limit?: number }) { + const predicates = ["completed_at >= ?"]; + const bindings: unknown[] = [options.from]; + if (options.repo !== undefined) { + predicates.push("(target_repo = ? OR target_repo IS NULL)"); + bindings.push(options.repo); + } + const rows = this.storage.sql.exec( + `SELECT record_json FROM ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} + WHERE ${predicates.join(" AND ")} + ORDER BY CASE WHEN workflow_outcome IN ('failure', 'cancelled') THEN 0 ELSE 1 END, + completed_at DESC, run_id, run_attempt LIMIT ?`, + ...bindings, + options.limit ?? 10_000, + ) as Iterable<{ record_json?: unknown }>; + return Array.from(rows) + .map((row) => normalizeReviewRunTelemetry(JSON.parse(String(row.record_json || "null")))) + .filter((record): record is DurableReviewRunTelemetry => record !== null); + } + + private reviewObservabilityTelemetryRowsSync(options: { repo?: string; from: number }) { + const predicates = ["(status = 'refreshing' OR COALESCE(terminal_at, updated_at) >= ?)"]; + const bindings: unknown[] = [options.from]; + if (options.repo !== undefined) { + predicates.push("repo = ?"); + bindings.push(options.repo); + } + const rows = this.storage.sql.exec( + `SELECT record_json FROM ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} + WHERE ${predicates.join(" AND ")} + ORDER BY CASE WHEN status = 'refreshing' THEN 0 ELSE 1 END, + CASE WHEN status = 'refreshing' THEN updated_at END ASC, + CASE WHEN status = 'completed' THEN COALESCE(terminal_at, updated_at) END DESC + LIMIT ?`, + ...bindings, + REVIEW_OBSERVABILITY_SCAN_LIMIT + 1, + ) as Iterable<{ record_json?: unknown }>; + return Array.from(rows) + .map((row) => normalizeReviewTelemetry(JSON.parse(String(row.record_json || "null")))) + .filter((record): record is DurableReviewTelemetry => record !== null); + } + + private reviewObservability(search: URLSearchParams) { + const range = String(search.get("range") || "24h") as keyof typeof REVIEW_OBSERVABILITY_RANGES; + const repoValue = String(search.get("repo") || "all").trim(); + if ( + !Object.hasOwn(REVIEW_OBSERVABILITY_RANGES, range) || + (repoValue !== "all" && !/^[A-Za-z0-9_.-]+\/[A-Za-z0-9_.-]+$/.test(repoValue)) + ) { + return json({ error: "invalid_review_observability_query" }, 400); + } + return json( + this.reviewObservabilitySync({ + range, + repo: repoValue === "all" ? null : repoValue, + now: Date.now(), + }), + ); + } + + private reviewObservabilitySync(options: { + range: keyof typeof REVIEW_OBSERVABILITY_RANGES; + repo: string | null; + now: number; + }) { + const from = options.now - REVIEW_OBSERVABILITY_RANGES[options.range]; + const records = this.reviewObservabilityTelemetryRowsSync({ + ...(options.repo ? { repo: options.repo } : {}), + from, + }); + const runs = this.reviewRunTelemetryRowsSync({ + ...(options.repo ? { repo: options.repo } : {}), + from, + limit: REVIEW_OBSERVABILITY_SCAN_LIMIT + 1, + }); + const telemetryComplete = + records.length <= REVIEW_OBSERVABILITY_SCAN_LIMIT && + runs.length <= REVIEW_OBSERVABILITY_SCAN_LIMIT; + const requiredSinceRaw = Date.parse(String(this.env.REVIEW_OBSERVABILITY_REQUIRED_SINCE || "")); + return summarizeReviewObservability({ + records: records.slice(0, REVIEW_OBSERVABILITY_SCAN_LIMIT), + runs: runs.slice(0, REVIEW_OBSERVABILITY_SCAN_LIMIT), + range: options.range, + repo: options.repo, + required: String(this.env.REVIEW_OBSERVABILITY_REQUIRED || "") === "1", + ...(Number.isFinite(requiredSinceRaw) ? { requiredSince: requiredSinceRaw } : {}), + recoveryEnabled: String(this.env.REVIEW_RECOVERY_ENABLED || "") === "1", + telemetryComplete, + now: options.now, + }); + } + private reviewTelemetryHealthSync(now: number): ReviewTelemetryHealth { const counts = Array.from( this.storage.sql.exec( @@ -1715,6 +2017,10 @@ export class ExactReviewQueue { WHERE status = 'completed' AND updated_at <= ?`, now - REVIEW_TELEMETRY_RETENTION_MS, ); + this.storage.sql.exec( + `DELETE FROM ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} WHERE completed_at <= ?`, + now - REVIEW_TELEMETRY_RETENTION_MS, + ); } private async supersedePublicationCandidates(value: unknown) { @@ -1974,28 +2280,80 @@ export class ExactReviewQueue { run_id TEXT NOT NULL, run_attempt INTEGER NOT NULL CHECK (run_attempt >= 1), status TEXT NOT NULL CHECK (status IN ('refreshing', 'completed')), + outcome TEXT, + trigger_lane TEXT, + trigger_origin TEXT, + terminal_at INTEGER, updated_at INTEGER NOT NULL, lease_expires_at INTEGER, + generation INTEGER, + operation_id TEXT, + queue_ms INTEGER, + claim_ms INTEGER, + review_ms INTEGER, + publication_ms INTEGER, + total_ms INTEGER, record_json TEXT NOT NULL, PRIMARY KEY (repo, item_number, run_id, run_attempt) ) STRICT`, ); - const hasLeaseExpiry = Array.from( - this.storage.sql.exec( - `SELECT name FROM pragma_table_info('${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE}') - WHERE name = 'lease_expires_at'`, - ), - ).length; - if (!hasLeaseExpiry) { - this.storage.sql.exec( - `ALTER TABLE ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} - ADD COLUMN lease_expires_at INTEGER`, - ); + for (const [column, definition] of [ + ["lease_expires_at", "INTEGER"], + ["outcome", "TEXT"], + ["trigger_lane", "TEXT"], + ["trigger_origin", "TEXT"], + ["terminal_at", "INTEGER"], + ["generation", "INTEGER"], + ["operation_id", "TEXT"], + ["queue_ms", "INTEGER"], + ["claim_ms", "INTEGER"], + ["review_ms", "INTEGER"], + ["publication_ms", "INTEGER"], + ["total_ms", "INTEGER"], + ]) { + const present = Array.from( + this.storage.sql.exec( + `SELECT name FROM pragma_table_info('${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE}') WHERE name = ?`, + column, + ), + ).length; + if (!present) { + this.storage.sql.exec( + `ALTER TABLE ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} ADD COLUMN ${column} ${definition}`, + ); + } } this.storage.sql.exec( `CREATE INDEX IF NOT EXISTS exact_review_review_telemetry_status ON ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} (status, updated_at)`, ); + this.storage.sql.exec( + `CREATE INDEX IF NOT EXISTS exact_review_review_telemetry_aggregate + ON ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} + (trigger_lane, repo, terminal_at, outcome)`, + ); + this.storage.sql.exec( + `CREATE INDEX IF NOT EXISTS exact_review_review_telemetry_operation + ON ${EXACT_REVIEW_REVIEW_TELEMETRY_TABLE} (operation_id, terminal_at)`, + ); + this.storage.sql.exec( + `CREATE TABLE IF NOT EXISTS ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} ( + run_id TEXT NOT NULL, + run_attempt INTEGER NOT NULL CHECK (run_attempt >= 1), + workflow_outcome TEXT NOT NULL, + trigger_lane TEXT NOT NULL, + trigger_origin TEXT NOT NULL, + target_repo TEXT, + completed_at INTEGER NOT NULL, + record_json TEXT NOT NULL, + PRIMARY KEY (run_id, run_attempt) + ) STRICT`, + ); + this.storage.sql.exec( + `CREATE INDEX IF NOT EXISTS exact_review_run_telemetry_aggregate + ON ${EXACT_REVIEW_RUN_TELEMETRY_TABLE} + (trigger_lane, target_repo, completed_at, workflow_outcome)`, + ); } private readStorageMetaSync() { @@ -2846,7 +3204,7 @@ export class ExactReviewQueue { private async scheduleNext(state: ExactReviewQueueState, now: number) { const publicationControl = this.refreshPublicationControlSync(state, now); - const next = exactReviewQueueNextWakeAt( + const queueNext = exactReviewQueueNextWakeAt( state, now, exactReviewQueueCapacity(this.env), @@ -2862,6 +3220,13 @@ export class ExactReviewQueue { exactReviewPublicationDispatchLeaseMs(this.env), exactReviewHeartbeatGraceMs(this.env), ); + const reviewNext = this.nextReviewReconcileAtSync(now); + const next = + queueNext === null + ? reviewNext + : reviewNext === null + ? queueNext + : Math.min(queueNext, reviewNext); if (next === null) { await this.storage.deleteAlarm(); return; diff --git a/dashboard/review-observability.ts b/dashboard/review-observability.ts new file mode 100644 index 0000000000..057dfb9a8b --- /dev/null +++ b/dashboard/review-observability.ts @@ -0,0 +1,369 @@ +import type { DurableReviewRunTelemetry } from "./review-run-telemetry.ts"; +import { + REVIEW_TELEMETRY_DEGRADED_MS, + REVIEW_TELEMETRY_ORPHAN_MS, + type DurableReviewTelemetry, + type ReviewTriggerLane, +} from "./review-telemetry.ts"; + +export const REVIEW_OBSERVABILITY_RANGES = { + "6h": 6 * 60 * 60 * 1000, + "24h": 24 * 60 * 60 * 1000, + "7d": 7 * 24 * 60 * 60 * 1000, +} as const; +export const REVIEW_OBSERVABILITY_WARMUP_MS = 30 * 60 * 1000; + +const LANE_POLICY: Record< + ReviewTriggerLane, + { label: string; cadenceMs: number | null; optional?: boolean } +> = { + exact_event: { label: "Exact event", cadenceMs: null }, + hot_intake: { label: "Hot intake", cadenceMs: 5 * 60 * 1000 }, + normal_backfill: { label: "Normal backfill", cadenceMs: 5 * 60 * 1000 }, + recovery: { label: "Recovery", cadenceMs: null, optional: true }, +}; + +export type ReviewObservability = ReturnType; + +export function summarizeReviewObservability(options: { + records: readonly DurableReviewTelemetry[]; + runs: readonly DurableReviewRunTelemetry[]; + range: keyof typeof REVIEW_OBSERVABILITY_RANGES; + repo: string | null; + required: boolean; + requiredSince?: number; + recoveryEnabled?: boolean; + telemetryComplete?: boolean; + now?: number; +}) { + const now = options.now ?? Date.now(); + const rangeMs = REVIEW_OBSERVABILITY_RANGES[options.range]; + const from = now - rangeMs; + const records = options.records.filter( + (record) => + (record.status === "refreshing" || reviewTerminalTime(record) >= from) && + (!options.repo || record.repo === options.repo), + ); + const runs = options.runs.filter( + (run) => + Date.parse(run.completed_at) >= from && + (!options.repo || run.target_repo === null || run.target_repo === options.repo), + ); + const refreshing = records.filter((record) => record.status === "refreshing"); + const completed = records.filter((record) => record.status === "completed"); + const slow = refreshing.filter( + (record) => now - Date.parse(record.updated_at) >= REVIEW_TELEMETRY_DEGRADED_MS, + ); + const orphans = refreshing.filter( + (record) => + now - Date.parse(record.updated_at) >= REVIEW_TELEMETRY_ORPHAN_MS && + (record.lease_expires_at === null || Date.parse(record.lease_expires_at) <= now), + ); + const recoveredFailures = recoveredFailureKeys(completed); + const outcomes = Object.fromEntries( + ["succeeded", "failed", "interrupted", "cancelled", "superseded"].map((outcome) => [ + outcome, + completed.filter((record) => record.outcome === outcome).length, + ]), + ) as Record, number>; + const unresolvedFailures = completed.filter( + (record) => record.outcome === "failed" && !recoveredFailures.has(attemptKey(record)), + ); + const itemTerminalAnomalyRuns = new Set( + completed + .filter((record) => ["failed", "interrupted", "cancelled"].includes(String(record.outcome))) + .map(runAttemptKey), + ); + const runTerminalAnomalies = runs.filter( + (run) => + ["failure", "cancelled"].includes(run.workflow_outcome) && + !itemTerminalAnomalyRuns.has(runAttemptKey(run)), + ); + const { expectedAttempts, terminalAttempts } = reviewCoverage( + records, + runs.filter((run) => !options.repo || run.target_repo === options.repo), + ); + const terminalCoverage = expectedAttempts ? terminalAttempts / expectedAttempts : null; + const abnormalCount = + unresolvedFailures.length + + outcomes.interrupted + + outcomes.cancelled + + runTerminalAnomalies.length; + const abnormalSamples = completed.length + runTerminalAnomalies.length; + const abnormalRate = abnormalSamples ? abnormalCount / abnormalSamples : 0; + const warmup = + options.required && + options.requiredSince !== undefined && + now - options.requiredSince < REVIEW_OBSERVABILITY_WARMUP_MS; + const sources = (Object.keys(LANE_POLICY) as ReviewTriggerLane[]).map((lane) => + summarizeLane({ + lane, + runs: runs.filter((run) => run.trigger_lane === lane), + now, + required: options.required, + warmup, + recoveryEnabled: options.recoveryEnabled === true, + repo: options.repo, + }), + ); + + let health: "passive" | "healthy" | "degraded" | "critical" = options.required + ? "healthy" + : "passive"; + const reasons: string[] = []; + const raise = (next: "degraded" | "critical", reason: string) => { + reasons.push(reason); + if (options.required && (next === "critical" || health === "healthy")) health = next; + }; + if (options.required && !warmup) { + if (options.telemetryComplete === false) raise("degraded", "telemetry_unavailable"); + if (expectedAttempts >= 10 && terminalCoverage !== null && terminalCoverage < 0.9) { + raise("critical", "terminal_coverage_critical"); + } else if (expectedAttempts > 0 && terminalCoverage !== null && terminalCoverage < 0.98) { + raise("degraded", "terminal_coverage_degraded"); + } + if (orphans.length) raise("critical", "orphan_review_attempt"); + else if (slow.length) raise("degraded", "slow_review_attempt"); + if (abnormalSamples >= 5 && abnormalRate >= 0.2) { + raise("critical", "review_abnormal_rate_critical"); + } else if (abnormalCount) { + raise("degraded", "review_terminal_anomaly"); + } + for (const source of sources) { + if (source.status === "critical") raise("critical", `${source.lane}_missed_cadence`); + else if (source.status === "degraded") raise("degraded", `${source.lane}_degraded`); + } + } + + return { + mode: options.required ? (warmup ? "warmup" : "required") : "passive", + health, + reasons: [...new Set(reasons)], + range: options.range, + repo: options.repo ?? "all", + generated_at: new Date(now).toISOString(), + telemetry_complete: options.telemetryComplete !== false, + terminal_coverage: terminalCoverage === null ? null : round(terminalCoverage * 100, 1), + expected_attempts: expectedAttempts, + terminal_attempts: terminalAttempts, + success_rate_percent: reviewSuccessRate(outcomes, unresolvedFailures.length), + outcomes, + recovered_failures: recoveredFailures.size, + unresolved_failures: unresolvedFailures.length, + expected_superseded: outcomes.superseded, + unexpected_cancelled: outcomes.cancelled, + refreshing: refreshing.length, + slow: slow.length, + orphan: orphans.length, + abnormal_rate_percent: round(abnormalRate * 100, 1), + phases: phasePercentiles(completed), + sources, + anomalies: anomalyRows({ records, runs: runTerminalAnomalies, recoveredFailures, now }).slice( + 0, + 20, + ), + }; +} + +function reviewTerminalTime(record: DurableReviewTelemetry) { + // A range describes when an attempt reached terminal truth. Using its start + // would make attempts crossing the boundary look like missing telemetry. + return Date.parse(record.terminal_at ?? record.updated_at); +} + +function summarizeLane(options: { + lane: ReviewTriggerLane; + runs: readonly DurableReviewRunTelemetry[]; + now: number; + required: boolean; + warmup: boolean; + recoveryEnabled: boolean; + repo: string | null; +}) { + const policy = LANE_POLICY[options.lane]; + const attributedRuns = options.repo + ? options.runs.filter((run) => run.target_repo === options.repo) + : options.runs; + const newest = [...attributedRuns].sort( + (left, right) => Date.parse(right.completed_at) - Date.parse(left.completed_at), + ); + const hasOnlyUnattributedRuns = + options.repo !== null && !newest.length && options.runs.some((run) => run.target_repo === null); + const lastRun = newest[0]; + const lastSuccess = newest.find((run) => run.workflow_outcome === "success"); + let status: "passive" | "disabled" | "idle" | "healthy" | "degraded" | "critical"; + if (!options.required) status = "passive"; + else if (options.lane === "recovery" && !options.recoveryEnabled) status = "disabled"; + else if (options.warmup) status = "idle"; + else if (hasOnlyUnattributedRuns) status = "degraded"; + else if (!lastRun) status = policy.cadenceMs === null ? "idle" : "critical"; + else if (policy.cadenceMs === null) status = "healthy"; + else { + const age = options.now - Date.parse(lastRun.completed_at); + status = + age > policy.cadenceMs * 3 + ? "critical" + : lastRun.workflow_outcome !== "success" || age > policy.cadenceMs * 2 + ? "degraded" + : "healthy"; + } + return { + lane: options.lane, + label: policy.label, + status, + last_run_at: lastRun?.completed_at ?? null, + last_success_at: lastSuccess?.completed_at ?? null, + item_count: options.runs.reduce((total, run) => total + run.item_count, 0), + run_count: options.runs.length, + attribution: hasOnlyUnattributedRuns ? "unavailable" : "available", + }; +} + +function reviewCoverage( + records: readonly DurableReviewTelemetry[], + runs: readonly DurableReviewRunTelemetry[], +) { + const groups = new Map(); + for (const record of records) { + const key = `${record.run_id}:${record.run_attempt}`; + const group = groups.get(key) ?? { records: 0, terminal: 0, observed: 0 }; + group.records += 1; + if (record.status === "completed") group.terminal += 1; + groups.set(key, group); + } + for (const run of runs) { + const key = `${run.run_id}:${run.run_attempt}`; + const group = groups.get(key) ?? { records: 0, terminal: 0, observed: 0 }; + group.observed = Math.max(group.observed, run.item_count); + groups.set(key, group); + } + let expectedAttempts = 0; + let terminalAttempts = 0; + for (const group of groups.values()) { + const expected = Math.max(group.records, group.observed); + expectedAttempts += expected; + terminalAttempts += Math.min(group.terminal, expected); + } + return { expectedAttempts, terminalAttempts }; +} + +function phasePercentiles(records: readonly DurableReviewTelemetry[]) { + return Object.fromEntries( + ["queue", "claim", "review", "publication", "total"].map((phase) => { + const values = records + .map((record) => record.phase_durations_ms[phase as keyof typeof record.phase_durations_ms]) + .filter((value): value is number => Number.isFinite(value)) + .sort((left, right) => left - right); + return [phase, { p50_ms: percentile(values, 0.5), p95_ms: percentile(values, 0.95) }]; + }), + ); +} + +function percentile(values: readonly number[], quantile: number) { + if (!values.length) return null; + return values[Math.ceil(quantile * values.length) - 1] ?? values.at(-1) ?? null; +} + +function reviewSuccessRate( + outcomes: Record, number>, + unresolvedFailures: number, +) { + const denominator = + outcomes.succeeded + unresolvedFailures + outcomes.cancelled + outcomes.interrupted; + return denominator ? round((outcomes.succeeded / denominator) * 100, 1) : null; +} + +function recoveredFailureKeys(records: readonly DurableReviewTelemetry[]) { + const successes = new Map(); + for (const record of records) { + if (record.outcome === "succeeded" && record.operation_id) { + const operationKey = `${record.repo}\u0000${record.operation_id}`; + successes.set( + operationKey, + Math.max( + successes.get(operationKey) ?? 0, + Date.parse(record.terminal_at ?? record.updated_at), + ), + ); + } + } + return new Set( + records + .filter( + (record) => + record.outcome === "failed" && + record.operation_id && + (successes.get(`${record.repo}\u0000${record.operation_id}`) ?? 0) > + Date.parse(record.terminal_at ?? record.updated_at), + ) + .map(attemptKey), + ); +} + +function anomalyRows(options: { + records: readonly DurableReviewTelemetry[]; + runs: readonly DurableReviewRunTelemetry[]; + recoveredFailures: ReadonlySet; + now: number; +}) { + const itemRows = options.records.flatMap((record) => { + const age = options.now - Date.parse(record.updated_at); + const orphan = + record.status === "refreshing" && + age >= REVIEW_TELEMETRY_ORPHAN_MS && + (record.lease_expires_at === null || Date.parse(record.lease_expires_at) <= options.now); + const slow = record.status === "refreshing" && age >= REVIEW_TELEMETRY_DEGRADED_MS; + const unresolved = + record.outcome === "failed" && !options.recoveredFailures.has(attemptKey(record)); + if ( + !orphan && + !slow && + !unresolved && + !["cancelled", "interrupted"].includes(String(record.outcome)) + ) { + return []; + } + return [ + { + kind: orphan ? "orphan" : slow ? "slow" : record.outcome, + repo: record.repo, + item_number: record.item_number, + item_url: `https://github.com/${record.repo}/issues/${record.item_number}`, + run_url: `https://github.com/openclaw/clawsweeper/actions/runs/${record.run_id}`, + run_id: record.run_id, + run_attempt: record.run_attempt, + at: record.terminal_at ?? record.updated_at, + reason: record.terminal_reason ?? null, + }, + ]; + }); + const runRows = options.runs + .filter((run) => ["failure", "cancelled"].includes(run.workflow_outcome)) + .map((run) => ({ + kind: `workflow_${run.workflow_outcome}`, + repo: run.target_repo, + item_number: null, + item_url: null, + run_url: run.run_url, + run_id: run.run_id, + run_attempt: run.run_attempt, + at: run.completed_at, + reason: null, + })); + return [...itemRows, ...runRows].sort( + (left, right) => Date.parse(right.at) - Date.parse(left.at), + ); +} + +function attemptKey(record: DurableReviewTelemetry) { + return `${record.repo}#${record.item_number}:${record.run_id}:${record.run_attempt}`; +} + +function runAttemptKey(record: { run_id: string; run_attempt: number }) { + return `${record.run_id}:${record.run_attempt}`; +} + +function round(value: number, digits: number) { + const scale = 10 ** digits; + return Math.round(value * scale) / scale; +} diff --git a/dashboard/review-run-telemetry.ts b/dashboard/review-run-telemetry.ts new file mode 100644 index 0000000000..bbcafdffb4 --- /dev/null +++ b/dashboard/review-run-telemetry.ts @@ -0,0 +1,132 @@ +import type { ReviewTriggerLane, ReviewTriggerOrigin } from "./review-telemetry.ts"; + +const RUN_OUTCOMES = new Set(["success", "failure", "cancelled", "skipped"]); +const TRIGGER_LANES = new Set(["exact_event", "hot_intake", "normal_backfill", "recovery"]); +const TRIGGER_ORIGINS = new Set(["webhook", "command", "schedule", "manual", "system"]); +const JOB_OUTCOMES = new Set(["success", "failure", "cancelled", "skipped"]); + +export type DurableReviewRunTelemetry = { + run_id: string; + run_attempt: number; + workflow_outcome: "success" | "failure" | "cancelled" | "skipped"; + trigger_lane: ReviewTriggerLane; + trigger_origin: ReviewTriggerOrigin; + target_repo: string | null; + started_at: string; + completed_at: string; + run_url: string; + plan_count: number; + item_count: number; + publication_count: number; + source_event?: string; + source_action?: string; + review_jobs?: Array<{ + name: string; + conclusion: "success" | "failure" | "cancelled" | "skipped"; + item_number: number | null; + }>; +}; + +export function normalizeReviewRunTelemetry(value: unknown): DurableReviewRunTelemetry | null { + const record = objectValue(value); + const runId = String(record.run_id || "").trim(); + const runAttempt = Number(record.run_attempt); + const workflowOutcome = String(record.workflow_outcome || ""); + const triggerLane = String(record.trigger_lane || ""); + const triggerOrigin = String(record.trigger_origin || ""); + const targetRepo = record.target_repo == null ? null : String(record.target_repo).trim(); + const startedAt = timestamp(record.started_at); + const completedAt = timestamp(record.completed_at); + const runUrl = String(record.run_url || "").trim(); + const planCount = nonNegativeInteger(record.plan_count); + const itemCount = nonNegativeInteger(record.item_count); + const publicationCount = nonNegativeInteger(record.publication_count); + const sourceEvent = optionalString(record.source_event, 100); + const sourceAction = optionalString(record.source_action, 200); + const reviewJobs = normalizeReviewJobs(record.review_jobs); + if ( + !/^\d+$/.test(runId) || + !Number.isSafeInteger(runAttempt) || + runAttempt < 1 || + !RUN_OUTCOMES.has(workflowOutcome) || + !TRIGGER_LANES.has(triggerLane) || + !TRIGGER_ORIGINS.has(triggerOrigin) || + (targetRepo !== null && !/^[A-Za-z0-9_.-]+\/[A-Za-z0-9_.-]+$/.test(targetRepo)) || + !startedAt || + !completedAt || + Date.parse(completedAt) < Date.parse(startedAt) || + !/^https:\/\/github\.com\/[^/]+\/[^/]+\/actions\/runs\/\d+$/.test(runUrl) || + planCount === null || + itemCount === null || + publicationCount === null || + sourceEvent === null || + sourceAction === null || + reviewJobs === null + ) { + return null; + } + return { + run_id: runId, + run_attempt: runAttempt, + workflow_outcome: workflowOutcome as DurableReviewRunTelemetry["workflow_outcome"], + trigger_lane: triggerLane as ReviewTriggerLane, + trigger_origin: triggerOrigin as ReviewTriggerOrigin, + target_repo: targetRepo, + started_at: startedAt, + completed_at: completedAt, + run_url: runUrl, + plan_count: planCount, + item_count: itemCount, + publication_count: publicationCount, + ...(sourceEvent === undefined ? {} : { source_event: sourceEvent }), + ...(sourceAction === undefined ? {} : { source_action: sourceAction }), + ...(reviewJobs === undefined ? {} : { review_jobs: reviewJobs }), + }; +} + +function normalizeReviewJobs(value: unknown): DurableReviewRunTelemetry["review_jobs"] | null { + if (value === undefined || value === null) return undefined; + if (!Array.isArray(value) || value.length > 1_000) return null; + const jobs: NonNullable = []; + for (const valueJob of value) { + const job = objectValue(valueJob); + const name = optionalString(job.name, 200); + const conclusion = String(job.conclusion || ""); + const itemNumber = job.item_number == null ? null : Number(job.item_number); + if ( + !name || + !JOB_OUTCOMES.has(conclusion) || + (itemNumber !== null && (!Number.isSafeInteger(itemNumber) || itemNumber < 1)) + ) { + return null; + } + jobs.push({ + name, + conclusion: conclusion as NonNullable< + DurableReviewRunTelemetry["review_jobs"] + >[number]["conclusion"], + item_number: itemNumber, + }); + } + return jobs; +} + +function nonNegativeInteger(value: unknown) { + const number = Number(value); + return Number.isSafeInteger(number) && number >= 0 ? number : null; +} + +function optionalString(value: unknown, maxLength: number) { + if (value === undefined || value === null) return undefined; + const string = String(value).trim(); + return string && string.length <= maxLength ? string : null; +} + +function timestamp(value: unknown) { + const parsed = Date.parse(String(value || "")); + return Number.isFinite(parsed) ? new Date(parsed).toISOString() : null; +} + +function objectValue(value: unknown): Record { + return value && typeof value === "object" ? (value as Record) : {}; +} diff --git a/dashboard/review-telemetry.ts b/dashboard/review-telemetry.ts index 26b80be8b4..00c8648518 100644 --- a/dashboard/review-telemetry.ts +++ b/dashboard/review-telemetry.ts @@ -6,6 +6,11 @@ export const REVIEW_TELEMETRY_MAX_LEASE_HORIZON_MS = 24 * 60 * 60 * 1000; const OUTCOMES = new Set(["succeeded", "failed", "cancelled", "interrupted", "superseded"]); const PHASES = ["queue", "claim", "review", "publication", "total"] as const; +const TRIGGER_LANES = new Set(["exact_event", "hot_intake", "normal_backfill", "recovery"]); +const TRIGGER_ORIGINS = new Set(["webhook", "command", "schedule", "manual", "system"]); + +export type ReviewTriggerLane = "exact_event" | "hot_intake" | "normal_backfill" | "recovery"; +export type ReviewTriggerOrigin = "webhook" | "command" | "schedule" | "manual" | "system"; export type DurableReviewTelemetry = { repo: string; @@ -20,6 +25,12 @@ export type DurableReviewTelemetry = { phase_durations_ms: Partial>; generation?: number; operation_id?: string; + trigger_lane?: ReviewTriggerLane; + trigger_origin?: ReviewTriggerOrigin; + source_event?: string; + source_action?: string; + terminal_reason?: string; + terminal_at?: string; }; export type ReviewTelemetryHealth = { @@ -78,7 +89,26 @@ export function normalizeReviewTelemetry( } const generation = optionalPositiveInteger(record.generation); const operationId = optionalBoundedString(record.operation_id, 200); - if (generation === null || operationId === null) return null; + const triggerLane = optionalEnum(record.trigger_lane, TRIGGER_LANES); + const triggerOrigin = optionalEnum(record.trigger_origin, TRIGGER_ORIGINS); + const sourceEvent = optionalBoundedString(record.source_event, 100); + const sourceAction = optionalBoundedString(record.source_action, 200); + const terminalReason = optionalBoundedString(record.terminal_reason, 200); + const terminalAt = record.terminal_at == null ? undefined : timestamp(record.terminal_at); + if ( + generation === null || + operationId === null || + triggerLane === null || + triggerOrigin === null || + sourceEvent === null || + sourceAction === null || + terminalReason === null || + terminalAt === null || + (status === "refreshing" && (terminalReason !== undefined || terminalAt !== undefined)) || + (terminalAt !== undefined && Date.parse(terminalAt) > now + REVIEW_TELEMETRY_CLOCK_SKEW_MS) + ) { + return null; + } const phaseDurations = normalizePhaseDurations(record.phase_durations_ms); if (!phaseDurations) return null; return { @@ -94,6 +124,14 @@ export function normalizeReviewTelemetry( phase_durations_ms: phaseDurations, ...(generation === undefined ? {} : { generation }), ...(operationId === undefined ? {} : { operation_id: operationId }), + ...(triggerLane === undefined ? {} : { trigger_lane: triggerLane as ReviewTriggerLane }), + ...(triggerOrigin === undefined + ? {} + : { trigger_origin: triggerOrigin as ReviewTriggerOrigin }), + ...(sourceEvent === undefined ? {} : { source_event: sourceEvent }), + ...(sourceAction === undefined ? {} : { source_action: sourceAction }), + ...(terminalReason === undefined ? {} : { terminal_reason: terminalReason }), + ...(terminalAt === undefined ? {} : { terminal_at: terminalAt }), }; } @@ -163,6 +201,12 @@ function optionalBoundedString(value: unknown, maxLength: number) { return string && string.length <= maxLength ? string : null; } +function optionalEnum(value: unknown, allowed: ReadonlySet) { + if (value === undefined || value === null) return undefined; + const string = String(value); + return allowed.has(string) ? string : null; +} + function objectValue(value: unknown): Record { return value && typeof value === "object" ? (value as Record) : {}; } diff --git a/dashboard/worker.ts b/dashboard/worker.ts index 5f2951694a..2f40e35b07 100644 --- a/dashboard/worker.ts +++ b/dashboard/worker.ts @@ -529,6 +529,8 @@ export default { return authenticatedExactReviewQueueRequest(request, env, "/publications/supersede"); if (url.pathname === "/internal/exact-review/review-telemetry" && request.method === "POST") return authenticatedExactReviewQueueRequest(request, env, "/review-telemetry"); + if (url.pathname === "/internal/exact-review/review-run-telemetry" && request.method === "POST") + return authenticatedExactReviewQueueRequest(request, env, "/review-run-telemetry"); if (url.pathname === "/internal/exact-review/reconcile" && request.method === "POST") return authenticatedExactReviewReconcile(request, env); if (url.pathname === "/api/exact-review-queue" && request.method === "GET") @@ -537,6 +539,8 @@ export default { return exactReviewQueueRequest(env, `/item-status?${url.searchParams.toString()}`); if (url.pathname === "/api/exact-review-queue/reviews" && request.method === "GET") return exactReviewQueueRequest(env, `/review-telemetry?${url.searchParams.toString()}`); + if (url.pathname === "/api/review-observability" && request.method === "GET") + return exactReviewQueueRequest(env, `/review-observability?${url.searchParams.toString()}`); if (url.pathname === "/api/health-history" && request.method === "GET") return healthHistoryJson(request, env); if (url.pathname === "/api/automerge-metrics" && request.method === "GET") @@ -7176,6 +7180,28 @@ h2::before { content: ""; flex: 0 0 auto; width: 14px; height: 2px; border-radiu .execution-alert-title strong { font-size: 13px; } .execution-alert-title span, .execution-alert-toggle, .execution-alert-body { color: var(--muted); font-size: 11px; } .execution-alert-body { padding: 0 15px 13px; } +.review-reliability { margin-top: 22px; padding: 16px; border: 1px solid var(--line); border-radius: 12px; background: color-mix(in srgb, var(--panel) 94%, var(--claw)); } +.review-reliability-head { display: flex; align-items: flex-start; justify-content: space-between; gap: 16px; } +.review-reliability-title { display: grid; gap: 5px; } +.review-reliability-title h3 { margin: 0; font-size: 14px; } +.review-reliability-title span { color: var(--muted); font-size: 11px; } +.review-reliability-controls { display: flex; align-items: center; gap: 8px; flex-wrap: wrap; justify-content: flex-end; } +.review-reliability-controls select { max-width: 220px; padding: 5px 8px; border: 1px solid var(--line); border-radius: 8px; background: var(--panel); color: var(--text); font: inherit; font-size: 11px; } +.review-reliability-kpis, .review-reliability-sources { display: grid; grid-template-columns: repeat(4, minmax(0, 1fr)); gap: 1px; margin-top: 14px; background: var(--line-soft); border: 1px solid var(--line-soft); } +.review-reliability-kpi, .review-reliability-source { padding: 11px 12px; background: var(--panel); min-width: 0; } +.review-reliability-kpi span, .review-reliability-source span { display: block; color: var(--muted); font-size: 10px; } +.review-reliability-kpi strong, .review-reliability-source strong { display: block; margin-top: 5px; font-size: 18px; } +.review-reliability-source strong { font-size: 12px; } +.review-reliability-source small { display: block; margin-top: 5px; color: var(--muted); font-size: 10px; } +.review-status { display: inline-flex; align-items: center; gap: 6px; font-weight: 650; } +.review-status::before { content: ""; width: 7px; height: 7px; border-radius: 50%; background: var(--muted); } +.review-status.healthy::before { background: var(--green); } +.review-status.degraded::before { background: var(--amber); } +.review-status.critical::before { background: var(--red); } +.review-anomalies { margin-top: 12px; font-size: 11px; } +.review-anomalies summary { cursor: pointer; color: var(--muted); } +.review-anomaly-list { display: grid; gap: 7px; margin-top: 9px; } +.review-anomaly { display: flex; justify-content: space-between; gap: 12px; padding-top: 7px; border-top: 1px solid var(--line-soft); } .exact-trend { margin: 14px 0 16px; } .exact-trend-status { font-size: 12px; font-weight: 650; } .exact-trend-status.growing { color: var(--amber); } @@ -7993,6 +8019,7 @@ a.pill:hover { color: var(--claw); text-decoration: none; } } @media (max-width: 760px) { .automerge-kpis { grid-template-columns: repeat(2, minmax(0, 1fr)); } + .review-reliability-kpis, .review-reliability-sources { grid-template-columns: repeat(2, minmax(0, 1fr)); } .automerge-kpi:nth-child(3) { border-left: 0; border-top: 1px solid var(--line-soft); } .automerge-kpi:nth-child(4) { border-top: 1px solid var(--line-soft); } .automerge-details { grid-template-columns: 1fr; } @@ -8011,6 +8038,9 @@ a.pill:hover { color: var(--claw); text-decoration: none; } .hero-headline { font-size: 23px; gap: 10px; } .hero-dot { width: 10px; height: 10px; } .exact-review-head { align-items: flex-start; flex-direction: column; } + .review-reliability-head { flex-direction: column; } + .review-reliability-controls { justify-content: flex-start; } + .review-reliability-kpis, .review-reliability-sources { grid-template-columns: 1fr; } .grid, .drawer-grid { grid-template-columns: 1fr; } .metric, .metric:nth-child(3n + 1) { border-left: 0; border-top: 1px solid var(--line-soft); padding-left: 0; } .metric:first-child { border-top: 0; } @@ -8054,6 +8084,23 @@ a.pill:hover { color: var(--claw); text-decoration: none; }

Codex Capacity

+
+
+
+

Review reliability

+ Loading durable review telemetry… +
+
+ +
+ + + +
+
+
+
Loading review reliability…
+

Exact Review

@@ -8210,6 +8257,8 @@ let automaticIndex = new Map(); let activeHealthRange = "6h"; let healthHistoryLoadedAt = 0; let healthHistorySamples = []; +let activeReviewRange = "24h"; +let reviewObservabilityRequestGeneration = 0; function exactReviewHistory(lane) { return healthHistorySamples.flatMap(sample => { @@ -8492,6 +8541,56 @@ function renderExecutionAlert(current) { target.innerHTML = '
⚠ Work execution needs attention' + esc(parts.join(" · ")) + 'Details ▾
' + esc(details) + '
'; } +function reviewMetric(label, value, detail) { + return '
' + esc(label) + '' + esc(value) + '' + (detail ? '' + esc(detail) + '' : '') + '
'; +} + +function renderReviewReliability(payload) { + const summary = document.getElementById("review-reliability-summary"); + const target = document.getElementById("review-reliability-body"); + if (!summary || !target) return; + const passive = payload.mode === "passive"; + const status = passive ? "Awaiting v2 producers" : payload.mode === "warmup" ? "Producer warm-up" : payload.health === "healthy" ? "Healthy" : payload.health === "critical" ? "Critical" : "Needs attention"; + summary.innerHTML = '' + esc(status) + '' + (passive ? ' · PR 674 will enable enforcement' : ' · ' + esc(payload.range) + ' window'); + const coverage = payload.terminal_coverage == null ? "n/a" : payload.terminal_coverage + "%"; + const p95 = payload.phases?.total?.p95_ms; + const outcomes = payload.outcomes || {}; + const kpis = [ + reviewMetric("Terminal coverage", coverage, fmt.format(payload.terminal_attempts || 0) + " / " + fmt.format(payload.expected_attempts || 0)), + reviewMetric("Succeeded / failed / interrupted", fmt.format(outcomes.succeeded || 0) + " / " + fmt.format(outcomes.failed || 0) + " / " + fmt.format(outcomes.interrupted || 0), fmt.format(payload.unresolved_failures || 0) + " unresolved · success " + (payload.success_rate_percent == null ? "n/a" : payload.success_rate_percent + "%")), + reviewMetric("Superseded / cancelled", fmt.format(payload.expected_superseded || 0) + " / " + fmt.format(payload.unexpected_cancelled || 0), "expected / unexpected"), + reviewMetric("p95 total", p95 == null ? "n/a" : elapsed(p95), "refreshing " + fmt.format(payload.refreshing || 0) + " · slow " + fmt.format(payload.slow || 0) + " · orphan " + fmt.format(payload.orphan || 0)) + ].join(""); + const sources = (payload.sources || []).map(source => '
' + esc(source.label) + '' + esc(source.status) + 'last run ' + esc(source.last_run_at ? since(source.last_run_at) : "none") + ' · success ' + esc(source.last_success_at ? since(source.last_success_at) : "none") + ' · ' + fmt.format(source.item_count || 0) + ' items
').join(""); + const anomalies = (payload.anomalies || []).map(row => '
' + esc(row.kind) + ' ' + (row.item_url ? link(row.item_url, row.repo + "#" + row.item_number) : esc(row.repo || "workflow")) + (row.reason ? ' · ' + esc(row.reason) : '') + '' + link(row.run_url, "Actions run") + '
').join(""); + target.innerHTML = '
' + kpis + '
' + sources + '
' + (anomalies ? '
' + fmt.format(payload.anomalies.length) + ' anomalies
' + anomalies + '
' : '
No review anomalies in this window.
'); +} + +function populateReviewRepoFilter(repositories) { + const select = document.getElementById("review-reliability-repo"); + if (!select) return; + const selected = select.value || "all"; + const values = [...new Set((repositories || []).filter(Boolean))].sort(); + select.innerHTML = '' + values.map(repo => '').join(""); + select.value = values.includes(selected) ? selected : "all"; +} + +async function loadReviewReliability() { + const generation = ++reviewObservabilityRequestGeneration; + const repo = document.getElementById("review-reliability-repo")?.value || "all"; + try { + const response = await fetch("/api/review-observability?range=" + encodeURIComponent(activeReviewRange) + "&repo=" + encodeURIComponent(repo), { cache: "no-store" }); + if (!response.ok) throw new Error("review observability returned " + response.status); + const payload = await response.json(); + if (generation !== reviewObservabilityRequestGeneration) return; + renderReviewReliability(payload); + } catch { + if (generation !== reviewObservabilityRequestGeneration) return; + document.getElementById("review-reliability-summary").innerHTML = 'Telemetry unavailable'; + document.getElementById("review-reliability-body").innerHTML = '
Durable review telemetry could not be loaded.
'; + } +} + function formatAgeMinutes(value) { const minutes = Number(value); if (!Number.isFinite(minutes)) return "unknown"; @@ -8888,6 +8987,7 @@ async function load() { : "", ); loadHealthHistory(activeHealthRange, false).catch(() => undefined); + loadReviewReliability().catch(() => undefined); loadAutomergeMetrics().catch(() => undefined); } catch (error) { if (lastData) { @@ -8936,6 +9036,7 @@ function renderDashboard(data, note) { metric("Codex Capacity", fleet.budget_used_percent + "%", "Codex slot utilization", fleet.budget_used_percent, "var(--green)") ].join(""); renderExecutionAlert(data.operational_health); + populateReviewRepoFilter(data.source?.target_repositories || []); renderSystemMap(data); renderExactReviewLanes(data.exact_review_queue); renderExactReviewHandoff(data.exact_review_queue); @@ -9437,6 +9538,14 @@ document.getElementById("trend-ranges").addEventListener("click", event => { document.querySelectorAll("button[data-trend-range]").forEach(item => item.classList.toggle("active", item === button)); loadHealthHistory(button.dataset.trendRange || "6h", true).catch(() => undefined); }); +document.getElementById("review-reliability-ranges").addEventListener("click", event => { + const button = event.target.closest("button[data-review-range]"); + if (!button) return; + activeReviewRange = button.dataset.reviewRange || "24h"; + document.querySelectorAll("button[data-review-range]").forEach(item => item.classList.toggle("active", item === button)); + loadReviewReliability().catch(() => undefined); +}); +document.getElementById("review-reliability-repo").addEventListener("change", () => loadReviewReliability().catch(() => undefined)); document.getElementById("automerge-ranges").addEventListener("click", event => { const button = event.target.closest("button[data-automerge-range]"); if (!button) return; diff --git a/dashboard/wrangler.toml b/dashboard/wrangler.toml index 537896ec21..cbf84b4586 100644 --- a/dashboard/wrangler.toml +++ b/dashboard/wrangler.toml @@ -56,3 +56,5 @@ EXACT_REVIEW_WORKFLOW_PAUSED_RETRY_MS = "60000" EXACT_REVIEW_DISPATCH_DEBOUNCE_MS = "45000" EXACT_REVIEW_DISPATCH_DEBOUNCE_MAX_MS = "180000" EXACT_REVIEW_PENDING_SOFT_LIMIT = "300" +REVIEW_OBSERVABILITY_REQUIRED = "0" +REVIEW_RECOVERY_ENABLED = "0" diff --git a/docs/pr-674-observation-runbook.md b/docs/pr-674-observation-runbook.md index d7658fe8ac..e90420d67f 100644 --- a/docs/pr-674-observation-runbook.md +++ b/docs/pr-674-observation-runbook.md @@ -17,6 +17,7 @@ uses these stable rules: | Queue telemetry | current snapshot available | unavailable or incomplete | n/a | | Durable review status | refreshing under 30m | refreshing for at least 30m | refreshing for at least 150m without a provably active lease | | Queue/workflow execution | healthy or idle | degraded/unknown | stalled | +| Review execution | coverage at least 98% | coverage 90–98% or anomaly | coverage under 90% with 10 attempts or 20% anomaly rate | Open DLQ is amber even when another compatibility producer has not populated publication health. `critical` and `stalled` always dominate and render the top-level indicator red. @@ -33,30 +34,84 @@ GET /api/exact-review-queue/reviews?repo=openclaw/openclaw&item_number=123&limit Every row carries `status`, terminal `outcome`, `started_at`, `updated_at`, optional lease expiry, and bounded phase durations for `queue`, `claim`, `review`, `publication`, and `total`. -`generation` and `operation_id` are optional until PR 674 supplies authoritative values. Terminal -rows cannot be reopened by a delayed refreshing heartbeat. Completed rows are retained for 30 -days; active refreshing rows are retained until a producer records a terminal result. +`generation` and `operation_id` are optional until PR 674 supplies authoritative values. The v2 +shape also adds indexed `trigger_lane`, `trigger_origin`, terminal reason/time, outcome, operation +identity, and phase durations while continuing to accept v1 rows. Terminal rows cannot be reopened +by a delayed refreshing heartbeat. Completed rows are retained for 30 days; active refreshing rows +are retained until a producer or terminal observer records a result. + +The four review lanes are `exact_event`, `hot_intake`, `normal_backfill`, and `recovery`. Origins +are `webhook`, `command`, `schedule`, `manual`, or `system`; producers preserve the raw source +event/action as optional diagnostic context. Generation replacement is +`superseded/generation_superseded` and is displayed without degrading health. An unexplained +GitHub cancellation is `cancelled/workflow_cancelled` and does degrade health. A failed attempt +followed by success for the same `operation_id` remains in the raw failure count but is no longer +unresolved. The watchdog is deliberately read-only. An aged refreshing row without an active lease appears in `review_telemetry_health.orphans` and turns dashboard health red, but observation alone never changes the outcome to `interrupted`. A producer may record `interrupted` only after it proves the GitHub run is terminal and no current lease owns the attempt. +## Review reliability API and card + +The main dashboard loads: + +```text +GET /api/review-observability?range=24h&repo=all +``` + +Accepted ranges are `6h`, `24h`, and `7d`; repository values are `all` or one `owner/repo` slug. +The bounded response includes terminal coverage, outcome totals, expected supersession, +unexpected cancellation, recovered/unresolved failures, phase p50/p95, four lane freshness rows, +and at most 20 anomalies with complete item and Actions URLs. The card sits after Work execution +and before Exact Review. `normal_backfill` has its own row so exact-event or hot-intake success +cannot conceal global review failure. + +Coverage uses the observer's paginated GitHub job count as an independent denominator, so an +entirely missing item-producer record lowers coverage instead of disappearing from both sides of +the ratio. Observer reconciliation is bound to the exact run attempt, and operation recovery is +scoped by repository plus operation ID. + +The prerequisite deployment sets `REVIEW_OBSERVABILITY_REQUIRED=0`, and the card displays +`Awaiting v2 producers` instead of green. PR 674 sets the flag to `1` and records +`REVIEW_OBSERVABILITY_REQUIRED_SINCE` at rollout; the first 30 minutes are warm-up. After warm-up: + +- Green requires terminal coverage at least 98%, no slow/orphan/unexpected cancellation or + unresolved failure, and periodic lanes within two cadences. +- Amber covers 90–98% coverage, any missing terminal record in a small sample, a periodic lane + beyond two cadences, a refreshing item at 30 minutes, an unexpected terminal anomaly, or + unavailable telemetry. +- Red covers under 90% coverage with at least 10 expected attempts, a periodic lane beyond three + cadences, an orphan at 150 minutes without active lease, at least five terminal samples with a + 20% anomaly rate, or stalled workflow execution. + +Recovery is neutral `disabled` until explicitly enabled. Exact event is on-demand and remains +`idle` when it has no traffic. Expected superseded records and their duration are always visible +but never affect health. + ## Rollout and rollback 1. Deploy this change and verify `/api/status` includes `dashboard_health` and `exact_review_queue.review_telemetry_health` while existing review traffic remains unchanged. 2. Confirm healthy traffic remains green, then exercise fixture or staging records at the 30m and 150m boundaries. Do not create synthetic records in the live queue. -3. Land PR 674 and have each item shard write `refreshing` after its authoritative claim, update - phase durations at transitions, then write exactly one terminal outcome. Populate its generation - and operation identity without changing this endpoint shape. +3. Rebase PR 674 onto this prerequisite. Have both its exact-event and per-item matrix jobs write + `refreshing` only after authoritative claim, update queue/claim/review/publication durations, + and write exactly one terminal outcome. Populate generation, operation identity, lane/origin, + and terminal reason without changing these endpoints. Telemetry write failures warn but do not + fail review. 4. Compare each sampled row with its full GitHub Actions run URL and verify repo, item, run attempt, generation, terminal outcome, and total duration. Confirm a live lease keeps an old heartbeat out of the orphan list. 5. Verify publication backlog, open DLQ, queue read failure, one degraded sample, and one critical/stalled sample all reach the dashboard hero with the expected amber/red severity. +6. In proposal-only mode, trigger two generations for one item. Verify the older generation is + expected superseded, a sibling succeeds, terminal coverage is 100%, and item/run/generation/ + operation identities map to GitHub. Wait for at least one hot-intake and normal-backfill cadence + before removing PR 674's draft status. -If PR 674 must roll back, stop its telemetry writes only; this additive table and query remain safe -for older producers. Do not pause the live sweep. Investigate an orphan from its repo/item/run tuple, -then record `interrupted` only with terminal-run and lease evidence. +If PR 674 must roll back, first set `REVIEW_OBSERVABILITY_REQUIRED=0`, then stop its producer +writes. Keep observer data and 30-day item history; the card returns to passive/legacy behavior. +Do not pause the live sweep. Investigate an orphan from its repo/item/run tuple, then record +`interrupted` only with terminal-run and lease evidence. diff --git a/scripts/review-run-observer.mjs b/scripts/review-run-observer.mjs new file mode 100644 index 0000000000..421f99504e --- /dev/null +++ b/scripts/review-run-observer.mjs @@ -0,0 +1,229 @@ +#!/usr/bin/env node + +/** + * 定义:只读 GitHub workflow_run 终态并写入 ClawSweeper review 观测接口。 + * 参数:--event-file 必填;--dry-run 可选。认证和 API 地址来自环境变量。 + * 输出:stdout 打印 skipped 或写入摘要;错误写 stderr,退出码 1。 + * 决策:无法证明是 review 的 sweep 运行直接跳过,避免把 apply/audit/router 支持任务计入成功率。 + */ + +import { createHmac } from "node:crypto"; +import { readFile } from "node:fs/promises"; +import process from "node:process"; +import { pathToFileURL } from "node:url"; + +export function usage() { + return `Usage: + node scripts/review-run-observer.mjs --event-file [--dry-run] + +Description: + Observe one completed ClawSweeper workflow_run and publish bounded run-level telemetry. + +Options: + --event-file GitHub event JSON containing workflow_run (required) + --dry-run Print the normalized record without posting it + -h, --help Show this help + +Outputs: + Prints a JSON record in dry-run mode or a one-line publish/skip summary. Exit 1 means + the event, GitHub lookup, configuration, or telemetry write was invalid. + +Examples: + node scripts/review-run-observer.mjs --event-file "$GITHUB_EVENT_PATH" --dry-run + node scripts/review-run-observer.mjs --event-file event.json +`; +} + +export function classifyReviewRun(run) { + const title = String(run.display_title || run.name || "").trim(); + const event = String(run.event || ""); + if (/^(Apply |Sync |Audit |Fan out )/.test(title)) return null; + let triggerLane; + if (title.startsWith("Review event item")) triggerLane = "exact_event"; + else if (/^Review hot (?:ClawSweeper items|target repo)/.test(title)) triggerLane = "hot_intake"; + else if (title.startsWith("Retry failed Codex reviews")) triggerLane = "recovery"; + else if (/^Review (?:target repo|ClawSweeper items)/.test(title)) triggerLane = "normal_backfill"; + else return null; + + const command = /\[(?:router-|command:)/i.test(title); + const triggerOrigin = command + ? "command" + : event === "schedule" + ? "schedule" + : event === "workflow_dispatch" + ? "manual" + : event === "repository_dispatch" + ? "webhook" + : "system"; + const targetMatch = title.match( + /\b(?:repo |item |items |for )([A-Za-z0-9_.-]+\/[A-Za-z0-9_.-]+)(?:#|\b)/, + ); + return { + trigger_lane: triggerLane, + trigger_origin: triggerOrigin, + target_repo: targetMatch?.[1] ?? null, + }; +} + +export function buildReviewRunTelemetry(run, jobs) { + const classification = classifyReviewRun(run); + if (!classification) return null; + const names = jobs.map((job) => String(job.name || "")); + const titleItem = String(run.display_title || "").match(/#(\d+)\b/)?.[1]; + const rawReviewJobs = jobs.flatMap((job) => { + const name = String(job.name || ""); + if (!/^Review (?:exact event item|item|shard)\b/i.test(name)) { + return []; + } + const conclusion = ["success", "failure", "cancelled", "skipped"].includes(job.conclusion) + ? job.conclusion + : "failure"; + const itemMatch = name.match(/(?:#|\bitem\s+)(\d+)\b/i)?.[1]; + return [{ name, conclusion, item_number: itemMatch ? Number(itemMatch) : null }]; + }); + // Exact-event rolling deploys can expose both old and new compute jobs for + // one title-bound item. Collapse them to the durable item identity instead + // of manufacturing multiple attempts for one run tuple. + const reviewJobs = + titleItem && rawReviewJobs.length + ? [mergeTitleItemJobs(rawReviewJobs, Number(titleItem))] + : rawReviewJobs; + const count = (pattern) => names.filter((name) => pattern.test(name)).length; + const activePlanCount = jobs.filter( + (job) => + /\b(?:plan|select).*review/i.test(String(job.name || "")) && job.conclusion !== "skipped", + ).length; + const activeReviewCount = reviewJobs.filter((job) => job.conclusion !== "skipped").length; + // Queue-intake and publication-only sweep runs share the review title. Only an active plan or + // review-compute job proves that the run belongs in the review reliability denominator. + if (!activePlanCount && !activeReviewCount) return null; + const startedAt = run.run_started_at || run.created_at; + return { + run_id: String(run.id || ""), + run_attempt: Number(run.run_attempt), + workflow_outcome: + run.conclusion === "success" + ? "success" + : run.conclusion === "cancelled" + ? "cancelled" + : run.conclusion === "skipped" + ? "skipped" + : "failure", + ...classification, + started_at: startedAt, + completed_at: run.updated_at, + run_url: run.html_url, + plan_count: activePlanCount, + item_count: reviewJobs.length, + publication_count: count(/\bPublish .*review/i), + source_event: String(run.event || "system"), + review_jobs: reviewJobs.slice(0, 1_000), + }; +} + +function mergeTitleItemJobs(jobs, itemNumber) { + const rank = { skipped: 0, success: 1, cancelled: 2, failure: 3 }; + const selected = jobs.reduce((current, candidate) => + rank[candidate.conclusion] > rank[current.conclusion] ? candidate : current, + ); + return { ...selected, item_number: itemNumber }; +} + +async function main(argv) { + const args = parseArgs(argv); + if (args.help) { + process.stdout.write(usage()); + return; + } + if (!args.eventFile) throw new Error("--event-file is required; use --help for examples"); + const event = JSON.parse(await readFile(args.eventFile, "utf8")); + const run = event.workflow_run; + if (!run || run.status !== "completed") + throw new Error("event does not contain a completed workflow_run"); + const classification = classifyReviewRun(run); + if (!classification) { + process.stdout.write(`skipped non-review run ${String(run.id || "unknown")}\n`); + return; + } + const jobs = await fetchJobs(run); + const record = buildReviewRunTelemetry(run, jobs); + if (!record) { + process.stdout.write(`skipped support-only review run ${String(run.id || "unknown")}\n`); + return; + } + if (args.dryRun) { + process.stdout.write(`${JSON.stringify(record, null, 2)}\n`); + return; + } + await publish(record); + process.stdout.write( + `observed review run ${record.run_id}/${record.run_attempt} lane=${record.trigger_lane} items=${record.item_count}\n`, + ); +} + +function parseArgs(argv) { + const result = { eventFile: "", dryRun: false, help: false }; + for (let index = 0; index < argv.length; index += 1) { + const arg = argv[index]; + if (arg === "-h" || arg === "--help") result.help = true; + else if (arg === "--dry-run") result.dryRun = true; + else if (arg === "--event-file") result.eventFile = String(argv[++index] || ""); + else throw new Error(`unknown option ${arg}; use --help`); + } + return result; +} + +export async function fetchJobs(run, options = {}) { + const token = options.token ?? process.env.GH_TOKEN ?? ""; + const repository = options.repository ?? process.env.GITHUB_REPOSITORY ?? ""; + const apiUrl = options.apiUrl ?? process.env.GITHUB_API_URL ?? "https://api.github.com"; + if (!token || !repository) throw new Error("GH_TOKEN and GITHUB_REPOSITORY are required"); + const maxJobs = 1_000; + const jobs = []; + for (let page = 1; page <= 10; page += 1) { + const response = await fetch( + `${apiUrl}/repos/${repository}/actions/runs/${run.id}/attempts/${run.run_attempt}/jobs?per_page=100&page=${page}`, + { headers: { authorization: `Bearer ${token}`, accept: "application/vnd.github+json" } }, + ); + if (!response.ok) throw new Error(`GitHub jobs lookup returned ${response.status}`); + const body = await response.json(); + const pageJobs = Array.isArray(body.jobs) ? body.jobs : []; + const totalCount = Number.isSafeInteger(body.total_count) ? body.total_count : null; + if (totalCount !== null && totalCount > maxJobs) + throw new Error(`workflow job list exceeds observer bound of ${maxJobs}`); + jobs.push(...pageJobs); + // GitHub's total_count distinguishes an exact full final page from a truncated run. + // Keep the short-page fallback for test doubles and older compatible API proxies. + if (pageJobs.length < 100 || (totalCount !== null && jobs.length >= totalCount)) return jobs; + } + throw new Error(`workflow job list exceeds observer bound of ${maxJobs}`); +} + +async function publish(record) { + const secret = process.env.CLAWSWEEPER_WEBHOOK_SECRET || ""; + const queueUrl = String(process.env.QUEUE_URL || "").replace(/\/$/, ""); + if (!secret || !queueUrl) + throw new Error("CLAWSWEEPER_WEBHOOK_SECRET and QUEUE_URL are required"); + const body = JSON.stringify(record); + const signature = `sha256=${createHmac("sha256", secret).update(body).digest("hex")}`; + const response = await fetch(`${queueUrl}/internal/exact-review/review-run-telemetry`, { + method: "POST", + headers: { + "content-type": "application/json", + "x-clawsweeper-exact-review-signature": signature, + }, + body, + signal: AbortSignal.timeout(20_000), + }); + if (!response.ok) + throw new Error(`review telemetry write returned ${response.status}: ${await response.text()}`); +} + +if (import.meta.url === pathToFileURL(process.argv[1] || "").href) { + main(process.argv.slice(2)).catch((error) => { + process.stderr.write( + `review-run-observer: ${error instanceof Error ? error.message : String(error)}\n`, + ); + process.exitCode = 1; + }); +} diff --git a/test/dashboard-health.test.ts b/test/dashboard-health.test.ts index 6af339862e..ba10865065 100644 --- a/test/dashboard-health.test.ts +++ b/test/dashboard-health.test.ts @@ -68,6 +68,34 @@ test("dashboard health maps critical and stalled signals to red", () => { }); }); +test("dashboard health rolls required review execution up while passive rollout stays neutral", () => { + const passive = healthySnapshot(); + (passive.exact_review_queue as Record).review_execution_health = { + health: "passive", + }; + assert.equal(summarizeDashboardHealth(passive).severity, "green"); + + const degraded = healthySnapshot(); + (degraded.exact_review_queue as Record).review_execution_health = { + health: "degraded", + }; + assert.deepEqual(summarizeDashboardHealth(degraded), { + conclusion: "needs_attention", + severity: "amber", + reasons: ["review_execution_degraded"], + }); + + const critical = healthySnapshot(); + (critical.exact_review_queue as Record).review_execution_health = { + health: "critical", + }; + assert.deepEqual(summarizeDashboardHealth(critical), { + conclusion: "needs_attention", + severity: "red", + reasons: ["review_execution_critical"], + }); +}); + test("dashboard health fails amber when a required signal is absent", () => { const snapshot = healthySnapshot(); const queue = snapshot.exact_review_queue as Record; diff --git a/test/dashboard-worker.test.ts b/test/dashboard-worker.test.ts index dc8281a352..2940b0be70 100644 --- a/test/dashboard-worker.test.ts +++ b/test/dashboard-worker.test.ts @@ -6610,6 +6610,9 @@ test("dashboard HTML preserves UTF-8 emoji labels", async () => { assert.doesNotMatch(html, /id="health-trend-grid"/); assert.match(html, /\/api\/health-history\?range=/); assert.match(html, /Work execution needs attention/); + assert.match(html, /Review reliability/); + assert.match(html, /Awaiting v2 producers/); + assert.match(html, /\/api\/review-observability\?range=/); assert.match(html, /data-trend-range="6h"/); assert.match(html, /
/); assert.match(html, /Error Rate/); @@ -6625,6 +6628,7 @@ test("dashboard HTML preserves UTF-8 emoji labels", async () => { assert.doesNotMatch(html, /id="control-plane"/); assert.match(html, /\.exact-lanes \{ grid-template-columns: 1fr; \}/); assert.ok(html.indexOf("Codex Capacity") < html.indexOf('id="exact-review-lanes"')); + assert.ok(html.indexOf("Review reliability") < html.indexOf('id="exact-review-lanes"')); assert.ok(html.indexOf('id="exact-review-lanes"') < html.indexOf("Handoff Health")); assert.match(html, /Live terminals/); assert.match(html, /href="https:\/\/fleet\.example\.test\/terminal\?view=live&mode=all"/); @@ -7042,6 +7046,58 @@ test("dashboard hero treats apply and exact-review handoff health as attention", context.renderExecutionAlert({ ...healthyOperational, telemetry_complete: false }); assert.match(elementFor("execution-alert").innerHTML, /telemetry is incomplete/); + const reviewReliability = { + mode: "passive", + health: "passive", + range: "24h", + terminal_coverage: null, + terminal_attempts: 0, + expected_attempts: 0, + success_rate_percent: null, + outcomes: {}, + unresolved_failures: 0, + expected_superseded: 0, + unexpected_cancelled: 0, + refreshing: 0, + slow: 0, + orphan: 0, + phases: { total: { p95_ms: null } }, + sources: [], + anomalies: [], + }; + context.renderReviewReliability(reviewReliability); + assert.match(elementFor("review-reliability-summary").innerHTML, /Awaiting v2 producers/); + context.renderReviewReliability({ + ...reviewReliability, + mode: "required", + health: "healthy", + terminal_coverage: 100, + }); + assert.match(elementFor("review-reliability-summary").innerHTML, /review-status healthy/); + context.renderReviewReliability({ + ...reviewReliability, + mode: "required", + health: "degraded", + anomalies: [ + { + kind: "cancelled", + repo: "openclaw/openclaw", + item_number: 674, + item_url: "https://github.com/openclaw/openclaw/pull/674", + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/123", + reason: "workflow_cancelled", + }, + ], + }); + assert.match(elementFor("review-reliability-summary").innerHTML, /review-status degraded/); + assert.match(elementFor("review-reliability-body").innerHTML, /actions\/runs\/123/); + context.renderReviewReliability({ + ...reviewReliability, + mode: "required", + health: "critical", + }); + assert.match(elementFor("review-reliability-summary").innerHTML, /review-status critical/); + status.diagnostics.exact_review_queue_error = null; status.exact_review_queue = { handoff_health: { status: "healthy", phases: {} } }; status.operational_health = { @@ -11285,9 +11341,462 @@ test("exact-review queue durably stores and queries per-item review telemetry", }); }); -test("exact-review health includes the oldest refreshing row beyond operator page bounds", async () => { +test("exact-review queue incrementally indexes the v1 review telemetry table", async () => { + const storage = new MemoryDurableStorage(); + storage.sql.exec( + `CREATE TABLE exact_review_review_telemetry ( + repo TEXT NOT NULL, + item_number INTEGER NOT NULL, + run_id TEXT NOT NULL, + run_attempt INTEGER NOT NULL, + status TEXT NOT NULL, + updated_at INTEGER NOT NULL, + lease_expires_at INTEGER, + record_json TEXT NOT NULL, + PRIMARY KEY (repo, item_number, run_id, run_attempt) + ) STRICT`, + ); + const queue = new ExactReviewQueue({ storage }, {}); + const now = new Date().toISOString(); + const response = await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify({ + repo: "openclaw/openclaw", + item_number: 674, + run_id: "44000", + run_attempt: 1, + status: "completed", + outcome: "superseded", + started_at: now, + updated_at: now, + lease_expires_at: null, + phase_durations_ms: { total: 12_000 }, + generation: 2, + operation_id: "review:674:2", + trigger_lane: "exact_event", + trigger_origin: "webhook", + terminal_at: now, + terminal_reason: "generation_superseded", + }), + }), + ); + assert.equal(response.status, 200); + const columns = Array.from( + storage.sql.exec( + "SELECT name FROM pragma_table_info('exact_review_review_telemetry') ORDER BY name", + ), + (row) => row.name, + ); + for (const column of [ + "generation", + "operation_id", + "outcome", + "terminal_at", + "total_ms", + "trigger_lane", + "trigger_origin", + ]) { + assert.ok(columns.includes(column)); + } +}); + +test("terminal workflow observer conservatively completes only unleased refreshing telemetry", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue({ storage }, { REVIEW_OBSERVABILITY_REQUIRED: "1" }); + const now = Date.now(); + const refreshing = (itemNumber: number, leaseExpiresAt: string | null) => ({ + repo: "openclaw/openclaw", + item_number: itemNumber, + run_id: "55555", + run_attempt: 1, + status: "refreshing", + outcome: null, + started_at: new Date(now - 10 * 60_000).toISOString(), + updated_at: new Date(now - 5 * 60_000).toISOString(), + lease_expires_at: leaseExpiresAt, + phase_durations_ms: { queue: 1_000 }, + trigger_lane: "exact_event", + trigger_origin: "webhook", + }); + for (const record of [ + refreshing(674, null), + refreshing(675, new Date(now + 60 * 60_000).toISOString()), + { ...refreshing(676, null), run_attempt: 2 }, + ]) { + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify(record), + }), + ) + ).status, + 200, + ); + } + const run = { + run_id: "55555", + run_attempt: 1, + workflow_outcome: "cancelled", + trigger_lane: "exact_event", + trigger_origin: "webhook", + target_repo: "openclaw/openclaw", + started_at: new Date(now - 10 * 60_000).toISOString(), + completed_at: new Date(now).toISOString(), + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/55555", + plan_count: 1, + item_count: 2, + publication_count: 0, + source_event: "repository_dispatch", + review_jobs: [ + { name: "Review item 674", conclusion: "cancelled", item_number: 674 }, + { name: "Review item 675", conclusion: "cancelled", item_number: 675 }, + ], + }; + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-run-telemetry", { + method: "POST", + body: JSON.stringify(run), + }), + ) + ).status, + 200, + ); + const first = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=674"), + ) + ).json()) as { reviews: Array> }; + assert.equal(first.reviews[0].outcome, "cancelled"); + assert.equal(first.reviews[0].terminal_reason, "workflow_cancelled"); + const leased = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=675"), + ) + ).json()) as { reviews: Array> }; + assert.equal(leased.reviews[0].status, "refreshing"); + const scheduledReconcile = await storage.getAlarm(); + assert.ok(scheduledReconcile !== null && scheduledReconcile <= now + 60 * 60_000); + const otherAttempt = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=676"), + ) + ).json()) as { reviews: Array> }; + assert.equal(otherAttempt.reviews[0].status, "refreshing"); + + const aggregate = (await ( + await queue.fetch( + new Request("https://queue/review-observability?range=24h&repo=openclaw%2Fopenclaw"), + ) + ).json()) as Record; + assert.equal(aggregate.outcomes.cancelled, 1); + assert.equal(aggregate.unexpected_cancelled, 1); + assert.match(aggregate.anomalies[0].run_url, /\/actions\/runs\/55555$/); + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-observability?range=30d&repo=openclaw%2Fopenclaw"), + ) + ).status, + 400, + ); + + const originalNow = Date.now; + Date.now = () => now + 61 * 60_000; + try { + await queue.alarm(); + } finally { + Date.now = originalNow; + } + const reconciled = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=675"), + ) + ).json()) as { reviews: Array> }; + assert.equal(reconciled.reviews[0].outcome, "cancelled"); + assert.equal(reconciled.reviews[0].terminal_reason, "workflow_cancelled"); +}); + +test("terminal workflow evidence reconciles telemetry that arrives later", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue({ storage }, { REVIEW_OBSERVABILITY_REQUIRED: "1" }); + const now = Date.now(); + const run = { + run_id: "56565", + run_attempt: 1, + workflow_outcome: "failure", + trigger_lane: "normal_backfill", + trigger_origin: "schedule", + target_repo: "openclaw/openclaw", + started_at: new Date(now - 10 * 60_000).toISOString(), + completed_at: new Date(now - 60_000).toISOString(), + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/56565", + plan_count: 1, + item_count: 2, + publication_count: 0, + review_jobs: [ + { name: "Review item 701", conclusion: "failure", item_number: 701 }, + { name: "Review shard 2", conclusion: "success", item_number: null }, + ], + }; + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-run-telemetry", { + method: "POST", + body: JSON.stringify(run), + }), + ) + ).status, + 200, + ); + for (const itemNumber of [701, 702]) { + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify({ + repo: "openclaw/openclaw", + item_number: itemNumber, + run_id: run.run_id, + run_attempt: 1, + status: "refreshing", + outcome: null, + started_at: new Date(now - 9 * 60_000).toISOString(), + updated_at: new Date(now - 2 * 60_000).toISOString(), + lease_expires_at: null, + phase_durations_ms: {}, + trigger_lane: "normal_backfill", + trigger_origin: "schedule", + }), + }), + ) + ).status, + 200, + ); + const telemetry = (await ( + await queue.fetch( + new Request( + `https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=${itemNumber}`, + ), + ) + ).json()) as { reviews: Array> }; + if (itemNumber === 701) { + assert.equal(telemetry.reviews[0].outcome, "interrupted"); + assert.equal(telemetry.reviews[0].terminal_reason, "workflow_terminal"); + } else { + assert.equal(telemetry.reviews[0].status, "refreshing"); + } + } + assert.equal(await storage.getAlarm(), null); +}); + +test("terminal workflow evidence records only attributable success", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue({ storage }, { REVIEW_OBSERVABILITY_REQUIRED: "1" }); + const now = Date.now(); + const refreshing = (itemNumber: number, repo = "openclaw/openclaw") => ({ + repo, + item_number: itemNumber, + run_id: "57575", + run_attempt: 1, + status: "refreshing", + outcome: null, + started_at: new Date(now - 5 * 60_000).toISOString(), + updated_at: new Date(now - 2 * 60_000).toISOString(), + lease_expires_at: null, + phase_durations_ms: {}, + trigger_lane: "normal_backfill", + trigger_origin: "schedule", + }); + for (const itemNumber of [801, 802]) { + await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify(refreshing(itemNumber)), + }), + ); + } + const run = { + run_id: "57575", + run_attempt: 1, + workflow_outcome: "success", + trigger_lane: "normal_backfill", + trigger_origin: "schedule", + target_repo: "openclaw/openclaw", + started_at: new Date(now - 6 * 60_000).toISOString(), + completed_at: new Date(now - 60_000).toISOString(), + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/57575", + plan_count: 1, + item_count: 2, + publication_count: 1, + review_jobs: [ + { name: "Review item 801", conclusion: "success", item_number: 801 }, + { name: "Review shard 2", conclusion: "success", item_number: null }, + ], + }; + assert.equal( + ( + await queue.fetch( + new Request("https://queue/review-run-telemetry", { + method: "POST", + body: JSON.stringify(run), + }), + ) + ).status, + 200, + ); + for (const [itemNumber, expectedStatus] of [ + [801, "completed"], + [802, "refreshing"], + ] as const) { + const telemetry = (await ( + await queue.fetch( + new Request( + `https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=${itemNumber}`, + ), + ) + ).json()) as { reviews: Array> }; + assert.equal(telemetry.reviews[0].status, expectedStatus); + if (itemNumber === 801) { + assert.equal(telemetry.reviews[0].outcome, "succeeded"); + assert.equal(telemetry.reviews[0].terminal_reason, "workflow_job_succeeded"); + } + } + await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify(refreshing(801, "openclaw/clawhub")), + }), + ); + const otherRepo = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fclawhub&item_number=801"), + ) + ).json()) as { reviews: Array> }; + assert.equal(otherRepo.reviews[0].status, "refreshing"); + assert.equal(await storage.getAlarm(), null); +}); + +test("workflow observer replay reconciles only from immutable first-writer evidence", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue({ storage }, { REVIEW_OBSERVABILITY_REQUIRED: "1" }); + const now = Date.now(); + const baseRun = { + run_id: "58585", + run_attempt: 1, + workflow_outcome: "success", + trigger_lane: "exact_event", + trigger_origin: "webhook", + target_repo: "openclaw/clawhub", + started_at: new Date(now - 5 * 60_000).toISOString(), + completed_at: new Date(now - 60_000).toISOString(), + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/58585", + plan_count: 1, + item_count: 1, + publication_count: 0, + review_jobs: [{ name: "Review item 901", conclusion: "success", item_number: 901 }], + }; + await queue.fetch( + new Request("https://queue/review-run-telemetry", { + method: "POST", + body: JSON.stringify(baseRun), + }), + ); + await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify({ + repo: "openclaw/openclaw", + item_number: 901, + run_id: "58585", + run_attempt: 1, + status: "refreshing", + outcome: null, + started_at: new Date(now - 4 * 60_000).toISOString(), + updated_at: new Date(now - 2 * 60_000).toISOString(), + lease_expires_at: null, + phase_durations_ms: {}, + trigger_lane: "exact_event", + trigger_origin: "webhook", + }), + }), + ); + await queue.fetch( + new Request("https://queue/review-run-telemetry", { + method: "POST", + body: JSON.stringify({ + ...baseRun, + workflow_outcome: "cancelled", + target_repo: "openclaw/openclaw", + review_jobs: [{ name: "Review item 901", conclusion: "cancelled", item_number: 901 }], + }), + }), + ); + const telemetry = (await ( + await queue.fetch( + new Request("https://queue/review-telemetry?repo=openclaw%2Fopenclaw&item_number=901"), + ) + ).json()) as { reviews: Array> }; + assert.equal(telemetry.reviews[0].status, "refreshing"); +}); + +test("review observer write is signed while aggregate telemetry remains read-only", async () => { const storage = new MemoryDurableStorage(); const queue = new ExactReviewQueue({ storage }, {}); + const record = JSON.stringify({ + run_id: "60000", + run_attempt: 1, + workflow_outcome: "success", + trigger_lane: "normal_backfill", + trigger_origin: "schedule", + target_repo: "openclaw/openclaw", + started_at: new Date(Date.now() - 60_000).toISOString(), + completed_at: new Date().toISOString(), + run_url: "https://github.com/openclaw/clawsweeper/actions/runs/60000", + plan_count: 1, + item_count: 4, + publication_count: 1, + }); + const env = { + CLAWSWEEPER_WEBHOOK_SECRET: "test-token-placeholder", + EXACT_REVIEW_QUEUE: new MemoryDurableNamespace(queue), + }; + const denied = await worker.fetch( + new Request("https://clawsweeper.openclaw.ai/internal/exact-review/review-run-telemetry", { + method: "POST", + body: record, + }), + env, + ); + assert.equal(denied.status, 401); + + const signature = `sha256=${createHmac("sha256", "test-token-placeholder").update(record).digest("hex")}`; + const accepted = await worker.fetch( + new Request("https://clawsweeper.openclaw.ai/internal/exact-review/review-run-telemetry", { + method: "POST", + headers: { "x-clawsweeper-exact-review-signature": signature }, + body: record, + }), + env, + ); + assert.equal(accepted.status, 200); + const aggregate = await worker.fetch( + new Request("https://clawsweeper.openclaw.ai/api/review-observability?range=24h&repo=all"), + env, + ); + assert.equal(aggregate.status, 200); + assert.equal(((await aggregate.json()) as { mode: string }).mode, "passive"); +}); + +test("exact-review health includes the oldest refreshing row beyond operator page bounds", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue({ storage }, { REVIEW_OBSERVABILITY_REQUIRED: "1" }); await queue.fetch(new Request("https://queue/stats")); const now = Date.now(); for (let index = 0; index <= 10_000; index += 1) { @@ -11342,4 +11851,47 @@ test("exact-review health includes the oldest refreshing row beyond operator pag }, ], }); + const aggregate = (await ( + await queue.fetch(new Request("https://queue/review-observability?range=24h&repo=all")) + ).json()) as { health: string; orphan: number; telemetry_complete: boolean }; + assert.equal(aggregate.telemetry_complete, false); + assert.equal(aggregate.orphan, 1); + assert.equal(aggregate.health, "critical"); +}); + +test("review observability retains refreshing attempts older than the selected range", async () => { + const storage = new MemoryDurableStorage(); + const queue = new ExactReviewQueue( + { storage }, + { REVIEW_OBSERVABILITY_REQUIRED: "1", REVIEW_RECOVERY_ENABLED: "0" }, + ); + const now = Date.now(); + const response = await queue.fetch( + new Request("https://queue/review-telemetry", { + method: "POST", + body: JSON.stringify({ + repo: "openclaw/openclaw", + item_number: 674, + run_id: "70000", + run_attempt: 1, + status: "refreshing", + outcome: null, + started_at: new Date(now - 26 * 60 * 60_000).toISOString(), + updated_at: new Date(now - 25 * 60 * 60_000).toISOString(), + lease_expires_at: null, + phase_durations_ms: {}, + trigger_lane: "exact_event", + trigger_origin: "webhook", + }), + }), + ); + assert.equal(response.status, 200); + const aggregate = (await ( + await queue.fetch( + new Request("https://queue/review-observability?range=24h&repo=openclaw%2Fopenclaw"), + ) + ).json()) as { refreshing: number; orphan: number; health: string }; + assert.equal(aggregate.refreshing, 1); + assert.equal(aggregate.orphan, 1); + assert.equal(aggregate.health, "critical"); }); diff --git a/test/review-observability.test.ts b/test/review-observability.test.ts new file mode 100644 index 0000000000..03d081b41d --- /dev/null +++ b/test/review-observability.test.ts @@ -0,0 +1,285 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import type { DurableReviewRunTelemetry } from "../dashboard/review-run-telemetry.ts"; +import { + summarizeReviewObservability, + type ReviewObservability, +} from "../dashboard/review-observability.ts"; +import type { DurableReviewTelemetry } from "../dashboard/review-telemetry.ts"; + +const NOW = Date.parse("2026-07-19T12:00:00Z"); + +function item( + number: number, + outcome: DurableReviewTelemetry["outcome"] = "succeeded", + overrides: Partial = {}, +): DurableReviewTelemetry { + return { + repo: "openclaw/openclaw", + item_number: number, + run_id: String(9000 + number), + run_attempt: 1, + status: outcome === null ? "refreshing" : "completed", + outcome, + started_at: "2026-07-19T11:00:00.000Z", + updated_at: "2026-07-19T11:10:00.000Z", + lease_expires_at: outcome === null ? "2026-07-19T13:00:00.000Z" : null, + phase_durations_ms: { total: number * 1_000 }, + trigger_lane: "exact_event", + trigger_origin: "webhook", + ...(outcome === null + ? {} + : { terminal_at: "2026-07-19T11:10:00.000Z", terminal_reason: "completed" }), + ...overrides, + }; +} + +function wave( + lane: DurableReviewRunTelemetry["trigger_lane"], + minutesAgo: number, + overrides: Partial = {}, +): DurableReviewRunTelemetry { + const completedAt = new Date(NOW - minutesAgo * 60_000).toISOString(); + return { + run_id: String(10000 + minutesAgo), + run_attempt: 1, + workflow_outcome: "success", + trigger_lane: lane, + trigger_origin: "schedule", + target_repo: "openclaw/openclaw", + started_at: new Date(NOW - (minutesAgo + 5) * 60_000).toISOString(), + completed_at: completedAt, + run_url: `https://github.com/openclaw/clawsweeper/actions/runs/${10000 + minutesAgo}`, + plan_count: 1, + item_count: 0, + publication_count: 1, + ...overrides, + }; +} + +function summary( + records: DurableReviewTelemetry[], + runs = [wave("hot_intake", 2), wave("normal_backfill", 2)], +): ReviewObservability { + return summarizeReviewObservability({ + records, + runs, + range: "24h", + repo: null, + required: true, + recoveryEnabled: false, + now: NOW, + }); +} + +test("review observability is passive before v2 producers are required", () => { + const result = summarizeReviewObservability({ + records: [], + runs: [], + range: "24h", + repo: null, + required: false, + now: NOW, + }); + assert.equal(result.mode, "passive"); + assert.equal(result.health, "passive"); + assert.ok(result.sources.every((source) => source.status === "passive")); +}); + +test("review observability reports green coverage, percentiles, and lane freshness", () => { + const result = summary([item(1), item(2), item(3), item(4)]); + assert.equal(result.health, "healthy"); + assert.equal(result.terminal_coverage, 100); + assert.equal(result.success_rate_percent, 100); + assert.deepEqual(result.phases.total, { p50_ms: 2_000, p95_ms: 4_000 }); + assert.equal(result.sources.find((source) => source.lane === "exact_event")?.status, "idle"); + assert.equal(result.sources.find((source) => source.lane === "recovery")?.status, "disabled"); +}); + +test("workflow jobs provide an independent denominator for missing item telemetry", () => { + const result = summary( + [item(1, "succeeded", { run_id: "10002" })], + [wave("hot_intake", 2, { item_count: 10 }), wave("normal_backfill", 2)], + ); + assert.equal(result.expected_attempts, 10); + assert.equal(result.terminal_attempts, 1); + assert.equal(result.terminal_coverage, 10); + assert.equal(result.health, "critical"); +}); + +test("coverage sums disjoint run populations instead of masking missing telemetry", () => { + const result = summary( + Array.from({ length: 5 }, (_, index) => item(index + 1, "succeeded", { run_id: "20000" })), + [wave("hot_intake", 2, { run_id: "20001", item_count: 5 }), wave("normal_backfill", 2)], + ); + assert.equal(result.expected_attempts, 10); + assert.equal(result.terminal_attempts, 5); + assert.equal(result.terminal_coverage, 50); +}); + +test("completed attempts are ranged by terminal time instead of start time", () => { + const result = summary([ + item(1, "succeeded", { + run_id: "10002", + started_at: "2026-07-18T11:59:00.000Z", + updated_at: "2026-07-18T12:01:00.000Z", + terminal_at: "2026-07-18T12:01:00.000Z", + }), + ]); + assert.equal(result.expected_attempts, 1); + assert.equal(result.terminal_attempts, 1); + assert.equal(result.terminal_coverage, 100); +}); + +test("expected supersession is health-neutral while cancellation and missed cadence are amber", () => { + assert.equal(summary([item(1, "superseded")]).health, "healthy"); + const cancelled = summary([item(1, "cancelled", { terminal_reason: "workflow_cancelled" })]); + assert.equal(cancelled.health, "degraded"); + assert.equal(cancelled.unexpected_cancelled, 1); + assert.equal(cancelled.success_rate_percent, 0); + const staleLane = summary([item(1)], [wave("hot_intake", 11), wave("normal_backfill", 11)]); + assert.equal(staleLane.health, "degraded"); +}); + +test("workflow failure evidence degrades health when item terminals do not explain it", () => { + const result = summary( + [item(1, "succeeded")], + [ + wave("hot_intake", 4, { run_id: "30000", workflow_outcome: "failure" }), + wave("hot_intake", 2), + wave("normal_backfill", 2), + ], + ); + assert.equal(result.health, "degraded"); + assert.ok(result.reasons.includes("review_terminal_anomaly")); + assert.ok(result.anomalies.some((row) => row.kind === "workflow_failure")); +}); + +test("item terminal anomalies suppress duplicate workflow anomaly rows", () => { + const result = summary( + [item(1, "failed", { run_id: "31000" })], + [ + wave("hot_intake", 2, { run_id: "31000", workflow_outcome: "failure", item_count: 1 }), + wave("normal_backfill", 2), + ], + ); + assert.ok(result.anomalies.some((row) => row.kind === "failed")); + assert.ok(!result.anomalies.some((row) => row.kind === "workflow_failure")); +}); + +test("periodic lanes become critical when no run or only a stale failed run remains", () => { + assert.equal(summary([item(1)], []).health, "critical"); + const staleFailure = summary( + [item(1)], + [ + wave("hot_intake", 16, { workflow_outcome: "failure" }), + wave("normal_backfill", 16, { workflow_outcome: "failure" }), + ], + ); + assert.equal(staleFailure.health, "critical"); +}); + +test("review observability makes low coverage, orphan attempts, and high anomaly rates red", () => { + const lowCoverage = summary([ + ...Array.from({ length: 8 }, (_, index) => item(index + 1)), + item(9, null), + item(10, null), + ]); + assert.equal(lowCoverage.health, "critical"); + assert.ok(lowCoverage.reasons.includes("terminal_coverage_critical")); + + const orphan = item(20, null, { + started_at: "2026-07-19T08:00:00.000Z", + updated_at: "2026-07-19T09:00:00.000Z", + lease_expires_at: null, + }); + assert.equal(summary([orphan]).health, "critical"); + + const highAnomaly = summary([item(30, "cancelled"), item(31), item(32), item(33), item(34)]); + assert.equal(highAnomaly.health, "critical"); + assert.ok(highAnomaly.reasons.includes("review_abnormal_rate_critical")); + + const missedCadence = summary([item(40)], [wave("hot_intake", 16), wave("normal_backfill", 16)]); + assert.equal(missedCadence.health, "critical"); +}); + +test("a later success recovers the same operation failure without erasing its count", () => { + const result = summary([ + item(1, "failed", { + operation_id: "review:674", + terminal_at: "2026-07-19T11:05:00.000Z", + updated_at: "2026-07-19T11:05:00.000Z", + }), + item(1, "succeeded", { + run_id: "9999", + operation_id: "review:674", + terminal_at: "2026-07-19T11:15:00.000Z", + updated_at: "2026-07-19T11:15:00.000Z", + }), + ]); + assert.equal(result.outcomes.failed, 1); + assert.equal(result.recovered_failures, 1); + assert.equal(result.unresolved_failures, 0); + assert.equal(result.health, "healthy"); +}); + +test("operation recovery never crosses repository identity", () => { + const result = summary([ + item(1, "failed", { + operation_id: "shared-operation", + terminal_at: "2026-07-19T11:05:00.000Z", + updated_at: "2026-07-19T11:05:00.000Z", + }), + item(2, "succeeded", { + repo: "openclaw/clawhub", + operation_id: "shared-operation", + terminal_at: "2026-07-19T11:15:00.000Z", + updated_at: "2026-07-19T11:15:00.000Z", + }), + ]); + assert.equal(result.recovered_failures, 0); + assert.equal(result.unresolved_failures, 1); + assert.equal(result.health, "degraded"); +}); + +test("range and repo filters exclude unrelated telemetry and anomalies stay bounded", () => { + const records = Array.from({ length: 25 }, (_, index) => item(index + 1, "cancelled")); + records.push( + item(99, "failed", { + repo: "openclaw/clawhub", + started_at: "2026-07-01T00:00:00.000Z", + updated_at: "2026-07-01T01:00:00.000Z", + terminal_at: "2026-07-01T01:00:00.000Z", + }), + ); + const result = summarizeReviewObservability({ + records, + runs: [wave("hot_intake", 2), wave("normal_backfill", 2)], + range: "6h", + repo: "openclaw/openclaw", + required: true, + now: NOW, + }); + assert.equal(result.expected_attempts, 25); + assert.equal(result.anomalies.length, 20); +}); + +test("repo filters preserve unattributed run evidence without assigning its attempts", () => { + const result = summarizeReviewObservability({ + records: [item(1)], + runs: [ + wave("hot_intake", 2, { target_repo: null, item_count: 9 }), + wave("normal_backfill", 2, { target_repo: "openclaw/openclaw" }), + ], + range: "24h", + repo: "openclaw/openclaw", + required: true, + now: NOW, + }); + const hotIntake = result.sources.find((source) => source.lane === "hot_intake"); + assert.equal(result.expected_attempts, 1); + assert.equal(hotIntake?.status, "degraded"); + assert.equal(hotIntake?.attribution, "unavailable"); + assert.equal(hotIntake?.item_count, 9); +}); diff --git a/test/review-reliability-workflow.test.ts b/test/review-reliability-workflow.test.ts new file mode 100644 index 0000000000..ea75573ad3 --- /dev/null +++ b/test/review-reliability-workflow.test.ts @@ -0,0 +1,28 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; +import { parse } from "yaml"; + +test("review reliability observer listens only to terminal ClawSweeper runs with read permissions", () => { + const source = readFileSync(".github/workflows/review-reliability-observer.yml", "utf8"); + const workflow = parse(source) as Record; + assert.deepEqual(workflow.on.workflow_run, { + workflows: ["ClawSweeper"], + types: ["completed"], + }); + assert.deepEqual(workflow.permissions, { actions: "read", contents: "read" }); + const checkout = workflow.jobs.observe.steps.find((candidate: Record) => + String(candidate.uses || "").startsWith("actions/checkout@"), + ); + assert.equal(checkout.with.ref, "${{ github.event.repository.default_branch }}"); + assert.equal(checkout.with["persist-credentials"], false); + const step = workflow.jobs.observe.steps.find((candidate: Record) => + String(candidate.run || "").includes("review-run-observer.mjs"), + ); + assert.ok(step); + assert.match(step.run, /--event-file/); + assert.ok(step.env.CLAWSWEEPER_WEBHOOK_SECRET); + assert.ok(step.env.GH_TOKEN); + assert.ok(step.env.QUEUE_URL); + assert.doesNotMatch(source, /workflow_dispatch|schedule|apply-existing|apply-decisions/); +}); diff --git a/test/review-run-telemetry.test.ts b/test/review-run-telemetry.test.ts new file mode 100644 index 0000000000..f5b677ea6f --- /dev/null +++ b/test/review-run-telemetry.test.ts @@ -0,0 +1,243 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { normalizeReviewRunTelemetry } from "../dashboard/review-run-telemetry.ts"; +import { + buildReviewRunTelemetry, + classifyReviewRun, + fetchJobs, +} from "../scripts/review-run-observer.mjs"; + +function run(overrides: Record = {}) { + return { + id: 1234, + run_attempt: 2, + status: "completed", + conclusion: "success", + event: "repository_dispatch", + display_title: "Review event item openclaw/openclaw#674 [router-command]", + run_started_at: "2026-07-19T10:00:00Z", + updated_at: "2026-07-19T10:05:00Z", + html_url: "https://github.com/openclaw/clawsweeper/actions/runs/1234", + ...overrides, + }; +} + +test("review observer attributes each review entry path without counting support runs", () => { + assert.deepEqual(classifyReviewRun(run()), { + trigger_lane: "exact_event", + trigger_origin: "command", + target_repo: "openclaw/openclaw", + }); + assert.equal( + classifyReviewRun(run({ display_title: "Review event item openclaw/openclaw#674" })) + ?.trigger_origin, + "webhook", + ); + assert.equal( + classifyReviewRun( + run({ + display_title: "Review event item openclaw/openclaw#674", + event: "workflow_dispatch", + }), + )?.trigger_origin, + "manual", + ); + assert.equal( + classifyReviewRun(run({ display_title: "Review hot ClawSweeper items", event: "schedule" })) + ?.trigger_lane, + "hot_intake", + ); + assert.equal( + classifyReviewRun(run({ display_title: "Review ClawSweeper items", event: "schedule" })) + ?.trigger_lane, + "normal_backfill", + ); + assert.equal( + classifyReviewRun(run({ display_title: "Retry failed Codex reviews", event: "schedule" })) + ?.trigger_lane, + "recovery", + ); + assert.equal( + classifyReviewRun(run({ display_title: "Retry failed Codex reviews", event: "workflow_call" })) + ?.trigger_origin, + "system", + ); + assert.equal( + classifyReviewRun(run({ display_title: "Apply default ClawSweeper closures" })), + null, + ); + assert.equal(classifyReviewRun(run({ display_title: "Reconcile exact-review leases" })), null); +}); + +test("review observer records bounded plan, item, and publication counts", () => { + const record = buildReviewRunTelemetry(run(), [ + { name: "Plan review items" }, + { name: "Review item 674", conclusion: "success" }, + { name: "Review shard 2", conclusion: "cancelled" }, + { name: "Publish exact review" }, + ]); + assert.equal(record?.plan_count, 1); + assert.equal(record?.item_count, 1); + assert.equal(record?.publication_count, 1); + assert.deepEqual(record?.review_jobs, [ + { name: "Review shard 2", conclusion: "cancelled", item_number: 674 }, + ]); + assert.deepEqual(normalizeReviewRunTelemetry(record), { + ...record, + started_at: "2026-07-19T10:00:00.000Z", + completed_at: "2026-07-19T10:05:00.000Z", + }); +}); + +test("review observer excludes queue and publication support runs from review attempts", () => { + assert.equal( + buildReviewRunTelemetry(run(), [ + { name: "Queue legacy exact-review event", conclusion: "success" }, + { name: "Review exact event item", conclusion: "skipped" }, + ]), + null, + ); + assert.equal( + buildReviewRunTelemetry( + run({ display_title: "Review event item openclaw/openclaw#674@publish:1:1" }), + [{ name: "Publish exact review artifact", conclusion: "success" }], + ), + null, + ); + assert.deepEqual( + buildReviewRunTelemetry(run(), [{ name: "Review exact event item", conclusion: "success" }]) + ?.review_jobs, + [{ name: "Review exact event item", conclusion: "success", item_number: 674 }], + ); +}); + +test("review observer counts skipped matrix jobs after an active plan", () => { + const record = buildReviewRunTelemetry( + run({ conclusion: "failure", display_title: "Review ClawSweeper items" }), + [ + { name: "Plan review items", conclusion: "failure" }, + { name: "Review item 674", conclusion: "skipped" }, + { name: "Review item 675", conclusion: "skipped" }, + ], + ); + assert.equal(record?.item_count, 2); + assert.deepEqual(record?.review_jobs, [ + { name: "Review item 674", conclusion: "skipped", item_number: 674 }, + { name: "Review item 675", conclusion: "skipped", item_number: 675 }, + ]); +}); + +test("review run telemetry rejects unsafe identities and nonterminal time order", () => { + const record = buildReviewRunTelemetry(run(), [ + { name: "Review exact event item", conclusion: "success" }, + ]); + assert.ok(record); + assert.equal(normalizeReviewRunTelemetry({ ...record, run_id: "not-a-run" }), null); + assert.equal( + normalizeReviewRunTelemetry({ + ...record, + completed_at: "2026-07-19T09:59:59Z", + }), + null, + ); +}); + +test("review observer paginates jobs instead of silently undercounting large matrices", async () => { + const originalFetch = globalThis.fetch; + const pages: number[] = []; + globalThis.fetch = async (input) => { + const page = Number(new URL(String(input)).searchParams.get("page")); + pages.push(page); + const jobs = + page === 1 + ? Array.from({ length: 100 }, (_, index) => ({ + name: `Review item ${index + 1}`, + conclusion: "success", + })) + : [{ name: "Publish review artifacts", conclusion: "success" }]; + return new Response(JSON.stringify({ jobs }), { + status: 200, + headers: { "content-type": "application/json" }, + }); + }; + try { + const jobs = await fetchJobs(run(), { + token: "test-token-placeholder", + repository: "openclaw/clawsweeper", + apiUrl: "https://api.github.test", + }); + assert.equal(jobs.length, 101); + assert.deepEqual(pages, [1, 2]); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("review observer accepts exactly its 1000-job fetch bound", async () => { + const originalFetch = globalThis.fetch; + const pages: number[] = []; + globalThis.fetch = async (input) => { + const page = Number(new URL(String(input)).searchParams.get("page")); + pages.push(page); + return new Response( + JSON.stringify({ + total_count: 1_000, + jobs: Array.from({ length: 100 }, (_, index) => ({ + name: `Review item ${(page - 1) * 100 + index + 1}`, + conclusion: "success", + })), + }), + { status: 200, headers: { "content-type": "application/json" } }, + ); + }; + try { + const jobs = await fetchJobs(run(), { + token: "test-token-placeholder", + repository: "openclaw/clawsweeper", + apiUrl: "https://api.github.test", + }); + assert.equal(jobs.length, 1_000); + assert.deepEqual(pages, [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("review observer rejects runs beyond its 1000-job fetch bound", async () => { + const originalFetch = globalThis.fetch; + globalThis.fetch = async () => + new Response(JSON.stringify({ total_count: 1_001, jobs: [] }), { + status: 200, + headers: { "content-type": "application/json" }, + }); + try { + await assert.rejects( + fetchJobs(run(), { + token: "test-token-placeholder", + repository: "openclaw/clawsweeper", + apiUrl: "https://api.github.test", + }), + /exceeds observer bound of 1000/, + ); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("review observer preserves every job identity within its fetch bound", () => { + const jobs = Array.from({ length: 101 }, (_, index) => ({ + name: `Review item ${index + 1}`, + conclusion: "success", + })); + const record = buildReviewRunTelemetry( + run({ + display_title: "Review ClawSweeper items", + event: "schedule", + }), + jobs, + ); + assert.equal(record?.item_count, 101); + assert.equal(record?.review_jobs?.length, 101); + assert.ok(normalizeReviewRunTelemetry(record)); +}); diff --git a/test/review-telemetry.test.ts b/test/review-telemetry.test.ts index a63439f6a8..c77cc6e069 100644 --- a/test/review-telemetry.test.ts +++ b/test/review-telemetry.test.ts @@ -39,6 +39,29 @@ test("review telemetry contract accepts optional generation and operation identi ); }); +test("review telemetry v2 accepts trigger attribution and terminal reasons without requiring them from v1", () => { + const completed = telemetry({ + status: "completed", + outcome: "superseded", + lease_expires_at: null, + trigger_lane: "exact_event", + trigger_origin: "command", + source_event: "issue_comment", + source_action: "rereview", + terminal_reason: "generation_superseded", + terminal_at: "2026-07-19T11:59:00.000Z", + }); + assert.deepEqual(normalizeReviewTelemetry(completed, NOW), completed); + assert.equal( + normalizeReviewTelemetry({ ...completed, trigger_lane: "router_workflow" }, NOW), + null, + ); + assert.equal( + normalizeReviewTelemetry({ ...completed, status: "refreshing", outcome: null }, NOW), + null, + ); +}); + test("review telemetry rejects identifiers outside SQLite's reliable integer range", () => { assert.equal(normalizeReviewTelemetry(telemetry({ item_number: 1e20 }), NOW), null); assert.equal(normalizeReviewTelemetry(telemetry({ run_attempt: 1e20 }), NOW), null);