diff --git a/projects/onchain-analytics/README.md b/projects/onchain-analytics/README.md index 74c56c7..11db1d1 100644 --- a/projects/onchain-analytics/README.md +++ b/projects/onchain-analytics/README.md @@ -130,7 +130,7 @@ The pipeline and dbt are independent — the pipeline writes raw tables, dbt rea | [`gd_dbt/`](gd_dbt/) | dbt project — all warehouse models, tests, docs, macros | | [`pipeline-v5/`](pipeline-v5/) | The ingestion pipeline (TypeScript). The only one. Runbook in its own README | | [`warehouse/L1/`](warehouse/L1/) | Raw table DDL (pipeline-written tables, dbt *sources*) | -| [`scripts/`](scripts/) | L1 bootstrap script (`deploy-warehouse.ps1`) | +| [`scripts/`](scripts/) | Explicitly allowlisted L1 migration helper and sandbox validator | | [`contracts/`](contracts/) | ABI files, deployment block numbers, contract reference | | [`docs/`](docs/) | System documentation, data model, operations guide, governance | @@ -142,7 +142,8 @@ The pipeline and dbt are independent — the pipeline writes raw tables, dbt rea - Node.js LTS (v20+) - Google Cloud SDK (`gcloud`, `bq`) -- `gcloud auth application-default login` with BigQuery Job User + Data Editor on `gooddollar` +- `gcloud auth application-default login` for metadata reads and labelled sandbox validation +- Production L1 DDL and ingestion use separately approved impersonated identities; see [`03_OPERATIONS.md`](docs/03_OPERATIONS.md) - Python 3.9+ with dbt-bigquery (`pip install dbt-bigquery`) ### Run the warehouse @@ -164,12 +165,16 @@ npx tsx src/index.ts daily # Ingest from chain into the BigQuery raw tables npx tsx src/index.ts verify # Reconcile what was ingested against the contracts ``` -### Bootstrap L1 raw tables (first time only) +### Inspect the prepared L1 migration (plan-only) ```powershell -.\scripts\deploy-warehouse.ps1 +.\scripts\deploy-warehouse.ps1 -Migration 09_CreateRawLogs_v1.sql ``` +This only prints the selected target. Validate migrations in a labelled sandbox first. Production +execution requires separate approval, an explicit allowlisted migration, and service-account +impersonation; see [`03_OPERATIONS.md`](docs/03_OPERATIONS.md). + --- ## Documentation diff --git a/projects/onchain-analytics/docs/03_OPERATIONS.md b/projects/onchain-analytics/docs/03_OPERATIONS.md index b2f08a8..14d4242 100644 --- a/projects/onchain-analytics/docs/03_OPERATIONS.md +++ b/projects/onchain-analytics/docs/03_OPERATIONS.md @@ -97,16 +97,48 @@ nonzero. Two layers, two tools. -### L1 raw tables — one-time bootstrap (PowerShell) +### L1 raw tables -- allowlisted additive migrations -The raw event tables (`BlockchainEvents.*`) are what the pipeline streams into. They are dbt -*sources* (pipeline-written, dbt-read), not dbt models, so their DDL still lives in `warehouse/L1/`. -Create them once: +Raw tables are pipeline-written dbt sources, not dbt models. `scripts/deploy-warehouse.ps1` accepts +one named migration from a fixed allowlist; it never scans `warehouse/L1/`. Its default mode only +prints the target and migration name: +```powershell +.\scripts\deploy-warehouse.ps1 -Migration 09_CreateRawLogs_v1.sql ``` -.\scripts\deploy-warehouse.ps1 # creates the L1 raw tables + +Before any production change, validate the same migration files against a fresh labelled sandbox. +From `pipeline-v5/`: + +```powershell +node --import tsx ..\scripts\ops\validate-l0-migrations.mjs ..\..\_scratch\unit-07a-commissioning\sandbox-validation.json ``` +This sandbox check exercises the old `PipelineRuns` and `OracleReconciliation` shapes, repeats the +migrations, checks their statement types and byte caps, verifies historical fixture rows remain, +and proves the sandbox is absent after cleanup. It does not write production tables or ingest chain +data. + +Production DDL is a separate operation and is not authorized by running the validator or plan mode. +Only after separate approval of the exact migration and access list may the named administrator run +one migration at a time: + +```powershell +.\scripts\deploy-warehouse.ps1 -Migration 09_CreateRawLogs_v1.sql -Execute -AllowProduction ` + -ImpersonateServiceAccount schema-commissioner@gooddollar.iam.gserviceaccount.com +``` + +The service-account address above is illustrative; replace it only with the approved identity. The +helper refuses production execution without an explicit account. It sets +`CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT` only in the query child process's environment; +persistent gcloud configuration and the caller's environment are never modified. Separate +deployments therefore cannot overwrite each other's selected identity. The +helper also requires a typed confirmation for `gooddollar.BlockchainEvents` and applies a 10 GiB +per-job bytes cap. Stop if a live object differs from the measured schema baseline, an object that +should be absent already exists, a migration returns `SCRIPT`, a legacy row count changes, or an +effective permission is broader than the approved list. Never run `04_L0Contract_v3.sql`, +`06_L0Contract_v4.sql`, or `07_RetireV3EventTables.sql` through this path. + ### Staging, Semantic, Marts — dbt Everything above raw is managed by dbt. There are no numbered files to run by hand — dbt resolves @@ -154,8 +186,8 @@ model/column docs with `dbt docs serve` (opens ). |---|---|---| | `ENVIO_API_TOKEN is missing` | `pipeline-v5/.env` not created or empty | `cp pipeline-v5/.env.example pipeline-v5/.env` and fill in the token | | `Could not authenticate to Google` | gcloud ADC expired | `gcloud auth application-default login` again | -| `Table not found: gooddollar.BlockchainEvents.…` | L1 DDL not run yet | Apply `warehouse/L1/04_L0Contract_v3.sql` | -| `SCHEMA_MISMATCH: has no column(s) …` | The pipeline writes a column the live table lacks | Apply the L0 contract. The pipeline refuses to write rather than corrupting a MERGE | +| `Table not found: gooddollar.BlockchainEvents.…` | The prepared L1 schema migration has not been commissioned | Stop and check the approved migration list; do not run a historical contract file | +| `SCHEMA_MISMATCH:
has no column(s) …` | The live schema differs from the runtime contract | Stop ingestion. Re-measure the schema and approve a new additive migration; do not recreate the table | | Run exits 1 with skipped chunks | HyperSync rate limiting or a timeout | Read `IngestionCoverage` for the exact ranges, then `backfill --from --to` over them | | `UNCONFIRMED EMPTY RANGE` | A range came back empty and no independent endpoint could confirm it | Not an error to clear by retrying. The watermark deliberately did not advance. Re-run when the endpoints recover | | `REORG SUSPECTED` | An existing key now sits under a different block hash | Delete and re-ingest that block range | diff --git a/projects/onchain-analytics/pipeline-v5/README.md b/projects/onchain-analytics/pipeline-v5/README.md index 9c39bbf..9dfbb85 100644 --- a/projects/onchain-analytics/pipeline-v5/README.md +++ b/projects/onchain-analytics/pipeline-v5/README.md @@ -86,10 +86,12 @@ cp .env.example .env # then fill in ENVIO_API_TOKEN gcloud auth application-default login ``` -The BigQuery tables are created by the DDL in `../warehouse/L1/`, not by the pipeline. Apply -`06_L0Contract_v4.sql` before ingesting anything. The pipeline refuses to write a column the -live table does not have, and checks at startup that the bookkeeping tables carry the columns it -is about to write, rather than failing part way into a backfill. +The BigQuery tables are created by allowlisted migrations in `../warehouse/L1/`, not by the +pipeline. The multi-statement `06_L0Contract_v4.sql` is a reference, not a deployment command. +Before production commissioning, run the labelled-sandbox validator described in +`../docs/03_OPERATIONS.md`. The pipeline refuses to write a column the live table does not have, +and checks at startup that bookkeeping tables carry the columns it is about to write, rather than +failing part way into a backfill. ### Two datasets, and why staging is not one of them diff --git a/projects/onchain-analytics/scripts/README.md b/projects/onchain-analytics/scripts/README.md index c17ed4f..9759b64 100644 --- a/projects/onchain-analytics/scripts/README.md +++ b/projects/onchain-analytics/scripts/README.md @@ -5,14 +5,30 @@ One PowerShell helper remains. The warehouse (Semantic + Marts) is managed by ** ## Prerequisites - Google Cloud SDK installed: -- `gcloud auth application-default login` already run -- Authenticated user has BQ Data Editor + Job User on the `gooddollar` project +- `gcloud auth application-default login` already run for metadata reads and sandbox validation +- Production schema changes require a separately authorized administrator; do not grant the ordinary + analytics identity raw-dataset write permissions ## What's here | Script | Purpose | When to run | |---|---|---| -| [`deploy-warehouse.ps1`](deploy-warehouse.ps1) | Creates the **L1 raw tables** (`BlockchainEvents.*`) from the DDL in [`warehouse/L1/`](../warehouse/L1/). These are dbt *sources* (pipeline-written, dbt-read), not dbt models, so their bootstrap DDL still lives here. | Once after a clone, or if a raw table schema changes. | +| [`deploy-warehouse.ps1`](deploy-warehouse.ps1) | Applies one named migration from a fixed allowlist. Default is plan-only; production execution requires `-Execute`, `-AllowProduction`, explicit service-account impersonation, and a typed confirmation. | Only after the exact migration and production access are separately approved. | + +The L1 SQL folder is not an execution queue. `04_L0Contract_v3.sql`, `06_L0Contract_v4.sql`, +`07_RetireV3EventTables.sql`, and any unlisted file are refused. Run the labelled-sandbox migration +validator from `pipeline-v5/` with `node --import tsx ..\scripts\ops\validate-l0-migrations.mjs ..\..\_scratch\unit-07a-commissioning\sandbox-validation.json`. + +Local regression checks from the project root, with no BigQuery jobs or credential acquisition: + +```powershell +.\scripts\tests\deploy-warehouse.Tests.ps1 +node --test scripts/tests/validate-l0-migrations.test.mjs +``` + +The deployment checks verify overlapping child processes retain separate identities without +changing persistent gcloud settings. The parser checks distinguish required and nullable columns +without importing or executing the live sandbox validator. ## Everything else is dbt diff --git a/projects/onchain-analytics/scripts/deploy-warehouse.ps1 b/projects/onchain-analytics/scripts/deploy-warehouse.ps1 index 83458aa..a819900 100644 --- a/projects/onchain-analytics/scripts/deploy-warehouse.ps1 +++ b/projects/onchain-analytics/scripts/deploy-warehouse.ps1 @@ -1,34 +1,100 @@ # deploy-warehouse.ps1 -# Creates the L1 raw event tables (BlockchainEvents.*) from the DDL in warehouse/L1/. -# These are the tables pipeline-v5 writes into and dbt reads as sources; they are NOT managed -# by dbt, so this bootstrap DDL still lives here. -# -# The Semantic (L2) and Marts (L3) layers are managed by dbt. Use `dbt run`, not this script. -# See gd_dbt/ and docs/03_OPERATIONS.md. -# -# SAFETY: files whose header carries a DO NOT RUN or NOT THE LIVE SHAPE banner are skipped, and -# -Force deliberately does not override that. Two files in warehouse/L1 are CREATE OR REPLACE -# against tables holding 2.6 million rows of production data. +# Applies one explicitly allowlisted L1 migration to gooddollar.BlockchainEvents. +# It never scans the SQL directory. Unknown and historical files are refused by filename. +# Default execution is plan-only. Production execution requires both switches and a typed prompt. +# See docs/03_OPERATIONS.md for the migration and approval requirements. # # Usage: -# .\scripts\deploy-warehouse.ps1 # applies the L1 DDL that is safe to re-apply +# .\scripts\deploy-warehouse.ps1 -Migration 09_CreateRawLogs_v1.sql +# .\scripts\deploy-warehouse.ps1 -Migration 09_CreateRawLogs_v1.sql -Execute -AllowProduction # # Requires: # - Google Cloud SDK installed (provides the `bq` CLI) # - `gcloud auth application-default login` already run +# - Production execution is only for an administrator after separate approval param( - [Parameter(Position = 0)] - [ValidateSet("L1")] - [string]$Layer = "L1", + [Parameter(Mandatory = $true)] + [string]$Migration, + + [switch]$Execute, - [switch]$Force + [switch]$AllowProduction, + + [string]$ImpersonateServiceAccount ) +function New-BqProcessStartInfo { + param( + [string]$BqExe, + [string[]]$Arguments, + [string]$ImpersonateServiceAccount + ) + + $startInfo = New-Object System.Diagnostics.ProcessStartInfo + $startInfo.FileName = $env:ComSpec + $startInfo.Arguments = '/d /s /c ""' + $BqExe + '" ' + ($Arguments -join ' ') + '"' + $startInfo.UseShellExecute = $false + $startInfo.CreateNoWindow = $true + $startInfo.RedirectStandardInput = $true + $startInfo.RedirectStandardOutput = $true + $startInfo.RedirectStandardError = $true + $startInfo.EnvironmentVariables['CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT'] = $ImpersonateServiceAccount + return $startInfo +} + $ErrorActionPreference = "Stop" $ScriptDir = Split-Path -Parent $MyInvocation.MyCommand.Path $RepoRoot = Split-Path -Parent $ScriptDir $WarehouseDir = Join-Path $RepoRoot "warehouse" +$MigrationDir = Join-Path $WarehouseDir "L1" +$AllowedMigrations = @( + "08_PipelineRunsOutcome_v1.sql", + "09_CreateRawLogs_v1.sql", + "10_AddOracleReconciliationCompatibility_v1.sql", + "11_CreateRawLogsAllHistory_v1.sql", + "12_CreateTransactionsAllHistory_v1.sql" +) + +if ($Migration -notin $AllowedMigrations) { + throw "REFUSED_UNLISTED_MIGRATION: '$Migration' is not in the deployment allowlist. Historical and unknown SQL is never executed by this helper." +} + +$MigrationPath = Join-Path $MigrationDir $Migration +if (-not (Test-Path -LiteralPath $MigrationPath -PathType Leaf)) { + throw "Allowlisted migration is missing: $MigrationPath" +} + +$Sql = Get-Content -LiteralPath $MigrationPath -Raw +if (-not $Sql.Contains('${PROJECT}') -or -not $Sql.Contains('${DATASET}')) { + throw "Migration must use the literal project and dataset placeholders: $Migration" +} +$Sql = $Sql.Replace('${PROJECT}', 'gooddollar').Replace('${DATASET}', 'BlockchainEvents') +if ($Sql -match '\$\{(PROJECT|DATASET)\}') { + throw "Unresolved identifier placeholder in $Migration" +} + +Write-Host "Migration: $Migration" +Write-Host "Target: gooddollar.BlockchainEvents" + +if (-not $Execute) { + Write-Host "PLAN ONLY. No BigQuery client was resolved and no query was submitted." + return +} + +if (-not $AllowProduction) { + throw "REFUSED_PRODUCTION_EXECUTION: production execution requires -AllowProduction after separate production authorization." +} + +if ($ImpersonateServiceAccount -notmatch '^[A-Za-z0-9._+-]+@[A-Za-z0-9.-]+\.iam\.gserviceaccount\.com$') { + throw "REFUSED_IMPERSONATION_REQUIRED: provide the separately approved commissioner service account with -ImpersonateServiceAccount." +} + +$ExpectedConfirmation = "APPLY APPROVED DDL TO gooddollar.BlockchainEvents" +$Confirmation = Read-Host "Type '$ExpectedConfirmation' to continue" +if ($Confirmation -cne $ExpectedConfirmation) { + throw "Production confirmation did not match. Nothing was submitted." +} # Resolve the bq.cmd location (gcloud SDK ships it as bq.cmd on Windows). # Try common paths; fall back to PATH lookup. @@ -49,71 +115,27 @@ if (-not $BqExe) { exit 1 } Write-Host "Using bq: $BqExe" - -function Invoke-SqlFile { - param([string]$Path) - Write-Host "" - Write-Host "==== $Path ====" -ForegroundColor Cyan - $sql = Get-Content -Raw -Path $Path - # bq query reads SQL from stdin - $sql | & $BqExe query --use_legacy_sql=false --format=none --project_id=gooddollar - if ($LASTEXITCODE -ne 0) { - Write-Error "bq query failed for $Path" - exit $LASTEXITCODE - } -} - -function Deploy-Layer { - param([string]$LayerName) - $layerDir = Join-Path $WarehouseDir $LayerName - if (-not (Test-Path $layerDir)) { - Write-Error "Layer folder not found: $layerDir" - exit 1 - } - $files = Get-ChildItem -Path $layerDir -Filter "*.sql" | Sort-Object Name - if ($files.Count -eq 0) { - Write-Warning "No .sql files in $layerDir" - return - } - - # Refuse anything that would drop a table holding production data. warehouse/L1 now contains - # historical DDL that is NOT the live shape: 01 and 02 are CREATE OR REPLACE against the two - # tables holding 2.6 million rows, and 03 was superseded. Running this folder end to end used - # to be safe and no longer is. Each such file carries a banner and is skipped by name. - $skipped = @() - $toRun = @() - foreach ($f in $files) { - $head = Get-Content -Path $f.FullName -TotalCount 20 -Raw - if ($head -match 'DO NOT RUN|NOT THE LIVE SHAPE') { - $skipped += $f.Name - } else { - $toRun += $f - } - } - - if ($skipped.Count -gt 0) { - Write-Host "" - Write-Host "SKIPPED (superseded or destructive, banner in file header):" -ForegroundColor Yellow - foreach ($s in $skipped) { Write-Host " $s" -ForegroundColor Yellow } - if ($Force) { - Write-Error "-Force does not override this. These files would drop tables holding production data. Run the individual statements you actually want, by hand." - exit 1 - } - } - - if ($toRun.Count -eq 0) { - Write-Warning "Nothing to run in $LayerName after skips." - return - } - - Write-Host "Deploying $($toRun.Count) file(s) in $LayerName..." -ForegroundColor Green - foreach ($f in $toRun) { - Invoke-SqlFile -Path $f.FullName +$bqProcess = New-Object System.Diagnostics.Process +$bqProcess.StartInfo = New-BqProcessStartInfo -BqExe $BqExe -ImpersonateServiceAccount $ImpersonateServiceAccount -Arguments @( + 'query', '--use_legacy_sql=false', '--format=none', '--project_id=gooddollar', '--maximum_bytes_billed=10737418240' +) +try { + Write-Host "Executing $Migration as $ImpersonateServiceAccount" -ForegroundColor Cyan + [void]$bqProcess.Start() + $stdoutTask = $bqProcess.StandardOutput.ReadToEndAsync() + $stderrTask = $bqProcess.StandardError.ReadToEndAsync() + $bqProcess.StandardInput.WriteLine($Sql) + $bqProcess.StandardInput.Close() + $bqProcess.WaitForExit() + $stdout = $stdoutTask.Result + $stderr = $stderrTask.Result + if ($stdout) { Write-Host $stdout } + if ($bqProcess.ExitCode -ne 0) { + throw "bq query failed for $Migration with exit code $($bqProcess.ExitCode): $stderr" } - Write-Host "$LayerName complete." -ForegroundColor Green + if ($stderr) { Write-Host $stderr } +} finally { + $bqProcess.Dispose() } -Deploy-Layer $Layer - -Write-Host "" -Write-Host "Done." -ForegroundColor Green +Write-Host "Migration completed." -ForegroundColor Green diff --git a/projects/onchain-analytics/scripts/ops/validate-l0-migrations.mjs b/projects/onchain-analytics/scripts/ops/validate-l0-migrations.mjs new file mode 100644 index 0000000..e266e02 --- /dev/null +++ b/projects/onchain-analytics/scripts/ops/validate-l0-migrations.mjs @@ -0,0 +1,279 @@ +#!/usr/bin/env node +// Applies the exact allowlisted migration files to a fresh labelled sandbox, then proves schema, +// idempotency, legacy-row preservation and cleanup. No production table is written or copied. + +import { readFileSync, mkdirSync, writeFileSync } from 'node:fs'; +import { dirname, resolve } from 'node:path'; +import { fileURLToPath } from 'node:url'; + +process.env.ENVIO_API_TOKEN ??= 'stage-a-no-chain-reader'; +process.env.DATASET_ID ??= 'BlockchainEvents'; + +const projectRoot = resolve(dirname(fileURLToPath(import.meta.url)), '../..'); +const migrationDir = resolve(projectRoot, 'warehouse/L1'); +const runtimeSource = readFileSync(resolve(projectRoot, 'pipeline-v5/src/bq.ts'), 'utf8'); +const outputPath = resolve(process.argv[2] ?? resolve(projectRoot, '../../_scratch/schema-migration-validation.json')); +const allowed = [ + '08_PipelineRunsOutcome_v1.sql', + '09_CreateRawLogs_v1.sql', + '10_AddOracleReconciliationCompatibility_v1.sql', + '11_CreateRawLogsAllHistory_v1.sql', + '12_CreateTransactionsAllHistory_v1.sql', +]; +const maxBytesPerJob = 10 * 1024 ** 3; +const maxBytesTotal = 50 * 1024 ** 3; + +const [{ CONFIG, RAW_LOGS_SCHEMA, TRANSACTIONS_SCHEMA }, { getBigQueryClient }, sandbox] = await Promise.all([ + import('../../pipeline-v5/src/config.ts'), + import('../../pipeline-v5/src/adapters.ts'), + import('../../pipeline-v5/src/sandbox.ts'), +]); +await import('../../pipeline-v5/src/bq.ts'); + +const report = { + project: CONFIG.GCP_PROJECT_ID, + access_mode: 'labelled sandbox only; no production DDL, DML, or chain ingestion', + migrations: [], + jobs: [], + errors: [], +}; +let handle; +let cleanup; + +function assert(condition, message) { + if (!condition) throw new Error(message); +} + +function assertSchemaMatches(tableId, fields, expected) { + assert(expected.length > 0, `Runtime schema for ${tableId} could not be extracted`); + const actualShape = fields.map((field) => `${field.name}:${field.type}:${field.mode ?? 'NULLABLE'}`); + const expectedShape = expected.map((field) => `${field.name}:${field.type}:${field.mode ?? 'NULLABLE'}`); + assert(JSON.stringify(actualShape) === JSON.stringify(expectedShape), + `${tableId} schema differs from the runtime contract; expected=${JSON.stringify(expectedShape)} actual=${JSON.stringify(actualShape)}`); +} + +function runtimeTableSchema(tableId) { + const marker = 'CREATE TABLE IF NOT EXISTS ${fullTableName("' + tableId + '")} ('; + const start = runtimeSource.indexOf(marker); + assert(start >= 0, `Runtime CREATE TABLE definition for ${tableId} was not found`); + const bodyStart = start + marker.length; + const bodyEnd = runtimeSource.indexOf('\n )', bodyStart); + assert(bodyEnd >= 0, `Runtime CREATE TABLE definition for ${tableId} is unterminated`); + const typeMap = { INT64: 'INTEGER', BOOL: 'BOOLEAN', FLOAT64: 'FLOAT' }; + return runtimeSource.slice(bodyStart, bodyEnd).split(/\r?\n/).flatMap((sourceLine) => { + const line = sourceLine.trim(); + const match = line.match(/^([A-Za-z_][A-Za-z0-9_]*)\s+(STRING|INT64|TIMESTAMP|JSON|FLOAT64|BOOL)\b/i); + if (!match) return []; + const sourceType = match[2].toUpperCase(); + return [{ + name: match[1], + type: typeMap[sourceType] ?? sourceType, + mode: /\bNOT NULL\b/i.test(line) ? 'REQUIRED' : 'NULLABLE', + }]; + }); +} + +function renderMigration(name, datasetId) { + assert(allowed.includes(name), `REFUSED_UNLISTED_MIGRATION: ${name}`); + const source = readFileSync(resolve(migrationDir, name), 'utf8'); + assert(source.includes('${PROJECT}') && source.includes('${DATASET}'), `${name} lacks identifier placeholders`); + const rendered = source.replaceAll('${PROJECT}', CONFIG.GCP_PROJECT_ID).replaceAll('${DATASET}', datasetId); + assert(!rendered.includes('${PROJECT}') && !rendered.includes('${DATASET}'), `${name} has unresolved placeholders`); + return rendered; +} + +async function submit(sql, label, dryRun = false) { + const [job, apiResponse] = await getBigQueryClient().createQueryJob({ + query: sql, + projectId: CONFIG.GCP_PROJECT_ID, + useLegacySql: false, + dryRun, + maximumBytesBilled: String(maxBytesPerJob), + }); + if (dryRun) { + return { rows: [], metadata: apiResponse ?? job.metadata ?? {}, job }; + } + const [rows] = await job.getQueryResults(); + const [metadata] = await job.getMetadata(); + const queryStats = metadata.statistics?.query ?? {}; + const billed = Number(queryStats.totalBytesBilled ?? 0); + const entry = { + label, + job_id: job.id ?? metadata.jobReference?.jobId ?? null, + statement_type: queryStats.statementType ?? null, + bytes_processed: Number(queryStats.totalBytesProcessed ?? 0), + bytes_billed: billed, + maximum_bytes_billed: maxBytesPerJob, + }; + report.jobs.push(entry); + assert(billed <= maxBytesPerJob, `${label} billed ${billed} bytes above the per-job cap`); + assert(report.jobs.reduce((sum, item) => sum + item.bytes_billed, 0) <= maxBytesTotal, 'Sandbox work exceeded the 50 GiB total cap'); + return { rows, metadata, job }; +} + +async function applyMigration(name, datasetId) { + const sql = renderMigration(name, datasetId); + let dryRun; + try { + const checked = await submit(sql, `${name}:dry-run`, true); + dryRun = { supported: true, statement_type: checked.metadata.statistics?.query?.statementType ?? null }; + if (dryRun.statement_type === 'SCRIPT') throw new Error(`${name} dry-run returned SCRIPT, expected one statement`); + } catch (error) { + const message = String(error?.message ?? error); + if (!/dry.?run.{0,50}not supported|not supported.{0,50}dry.?run/i.test(message)) throw error; + dryRun = { supported: false, error: message.slice(0, 500) }; + } + + const result = await submit(sql, name); + const statementType = result.metadata.statistics?.query?.statementType ?? null; + assert(statementType !== 'SCRIPT', `${name} submitted as SCRIPT, expected a single statement`); + report.migrations.push({ name, dry_run: dryRun, executed_statement_type: statementType }); +} + +async function tableMetadata(datasetId, tableId) { + const [metadata] = await getBigQueryClient().dataset(datasetId).table(tableId).getMetadata(); + return metadata; +} + +async function scalar(sql) { + const { rows } = await submit(sql, 'sandbox assertion'); + return rows[0]; +} + +try { + handle = await sandbox.createSandbox({ purpose: 'schema-migration-validation', tableExpirationHours: 2 }); + report.sandbox_dataset = handle.datasetId; + report.sandbox_label = handle.purposeLabel; + assert(!sandbox.PROTECTED_DATASETS.includes(handle.datasetId), 'Sandbox guard returned a protected dataset'); + + const dataset = `\`${CONFIG.GCP_PROJECT_ID}.${handle.datasetId}\``; + const production = `\`${CONFIG.GCP_PROJECT_ID}.BlockchainEvents\``; + + const productionTransactions = await tableMetadata('BlockchainEvents', 'Transactions'); + assertSchemaMatches('Production Transactions', productionTransactions.schema?.fields ?? [], TRANSACTIONS_SCHEMA); + const productionCoverage = await tableMetadata('BlockchainEvents', 'IngestionCoverage'); + const runtimeCoverageSchema = runtimeTableSchema('IngestionCoverage'); + assertSchemaMatches('Production IngestionCoverage', productionCoverage.schema?.fields ?? [], runtimeCoverageSchema); + const runtimePipelineSchema = runtimeTableSchema('PipelineRuns'); + const runtimeReconciliationSchema = runtimeTableSchema('OracleReconciliation'); + + await submit(`CREATE TABLE ${dataset}.PipelineRuns LIKE ${production}.PipelineRuns`, 'clone legacy PipelineRuns schema'); + await submit(` + INSERT INTO ${dataset}.PipelineRuns + (run_id, mode, started_at, completed_at, exit_code, total_rows_merged, + contracts_processed, contracts_failed, host, error_message, chains_processed, + captures_planned, captures_ok, captures_failed, pipeline_version) + SELECT CONCAT('stage-a-', CAST(n AS STRING)), 'backfill', + TIMESTAMP('2026-01-01 00:00:00+00'), TIMESTAMP('2026-01-01 00:01:00+00'), + 0, 0, 1, 0, 'sandbox-fixture', '', 'XDC', 1, 1, 0, 'fixture' + FROM UNNEST(GENERATE_ARRAY(1, 22)) AS n`, 'seed PipelineRuns fixture rows'); + + await submit(`CREATE TABLE ${dataset}.OracleReconciliation LIKE ${production}.OracleReconciliation`, 'clone legacy OracleReconciliation schema'); + await submit(` + INSERT INTO ${dataset}.OracleReconciliation + (run_id, network, table_id, protocol_day, oracle_block, oracle_count, + oracle_amount_raw, warehouse_stored, warehouse_distinct, warehouse_amount_raw, + count_gap, amount_gap_raw, verdict, checked_at) + SELECT CONCAT('legacy-', CAST(n AS STRING)), 'XDC', 'legacy_fixture', n, + 1000 + n, 1, '1', 1, 1, '1', 0, '0', 'exact', + TIMESTAMP('2026-01-01 00:00:00+00') + FROM UNNEST(GENERATE_ARRAY(1, 265)) AS n`, 'seed OracleReconciliation legacy rows'); + + await submit(`CREATE TABLE ${dataset}.Transactions LIKE ${production}.Transactions`, 'clone Transactions schema only'); + report.prepared_migrations = [ + { name: '09_CreateRawLogs_v1.sql', change: 'create RawLogs if absent' }, + { name: '08_PipelineRunsOutcome_v1.sql', change: 'add 26 nullable PipelineRuns columns' }, + { name: '10_AddOracleReconciliationCompatibility_v1.sql', change: 'add nullable chain_id and contract_address' }, + { name: '11_CreateRawLogsAllHistory_v1.sql', change: 'create RawLogsAllHistory if absent' }, + { name: '12_CreateTransactionsAllHistory_v1.sql', change: 'create TransactionsAllHistory if absent' }, + ]; + + await applyMigration('09_CreateRawLogs_v1.sql', handle.datasetId); + await applyMigration('09_CreateRawLogs_v1.sql', handle.datasetId); + const rawLogs = await tableMetadata(handle.datasetId, 'RawLogs'); + assertSchemaMatches('RawLogs', rawLogs.schema?.fields ?? [], RAW_LOGS_SCHEMA); + assert(rawLogs.timePartitioning?.type === 'MONTH', 'RawLogs is not monthly partitioned'); + assert(rawLogs.requirePartitionFilter === true, 'RawLogs is missing require_partition_filter'); + + await applyMigration('08_PipelineRunsOutcome_v1.sql', handle.datasetId); + await applyMigration('08_PipelineRunsOutcome_v1.sql', handle.datasetId); + const pipelineRuns = await tableMetadata(handle.datasetId, 'PipelineRuns'); + assertSchemaMatches('PipelineRuns', pipelineRuns.schema?.fields ?? [], runtimePipelineSchema); + const pipelinePreservation = await scalar(` + SELECT COUNT(*) AS rows_preserved, + COUNTIF(run_id IS NOT NULL) AS fixture_rows, + COUNTIF(release_sha IS NULL) AS old_rows_without_release_sha + FROM ${dataset}.PipelineRuns`); + assert(Number(pipelinePreservation.rows_preserved) === 22, 'PipelineRuns row count changed'); + assert(Number(pipelinePreservation.fixture_rows) === 22, 'PipelineRuns fixture rows changed'); + assert(Number(pipelinePreservation.old_rows_without_release_sha) === 22, 'PipelineRuns historical release_sha values changed'); + + await applyMigration('10_AddOracleReconciliationCompatibility_v1.sql', handle.datasetId); + await applyMigration('10_AddOracleReconciliationCompatibility_v1.sql', handle.datasetId); + const reconciliation = await tableMetadata(handle.datasetId, 'OracleReconciliation'); + const reconciliationFields = new Set((reconciliation.schema?.fields ?? []).map((field) => field.name)); + const reconciliationByName = new Map((reconciliation.schema?.fields ?? []).map((field) => [field.name, field])); + for (const field of runtimeReconciliationSchema) { + const actual = reconciliationByName.get(field.name); + assert(actual && actual.type === field.type && (actual.mode ?? 'NULLABLE') === field.mode, + `OracleReconciliation runtime field ${field.name} differs from its schema contract`); + } + assert(reconciliationFields.has('table_id'), 'OracleReconciliation legacy table_id was removed'); + const reconciliationPreservation = await scalar(` + SELECT COUNT(*) AS rows_preserved, + COUNTIF(table_id = 'legacy_fixture') AS legacy_rows, + COUNTIF(chain_id IS NULL AND contract_address IS NULL) AS uninferred_rows + FROM ${dataset}.OracleReconciliation`); + assert(Number(reconciliationPreservation.rows_preserved) === 265, 'OracleReconciliation row count changed'); + assert(Number(reconciliationPreservation.legacy_rows) === 265, 'OracleReconciliation table_id history changed'); + assert(Number(reconciliationPreservation.uninferred_rows) === 265, 'Legacy reconciliation dimensions were inferred'); + + const transactions = await tableMetadata(handle.datasetId, 'Transactions'); + assertSchemaMatches('Transactions', transactions.schema?.fields ?? [], TRANSACTIONS_SCHEMA); + assert(Number(transactions.numRows ?? 0) === 0, 'Transactions schema-only clone unexpectedly contains rows'); + + await applyMigration('11_CreateRawLogsAllHistory_v1.sql', handle.datasetId); + await applyMigration('11_CreateRawLogsAllHistory_v1.sql', handle.datasetId); + await applyMigration('12_CreateTransactionsAllHistory_v1.sql', handle.datasetId); + await applyMigration('12_CreateTransactionsAllHistory_v1.sql', handle.datasetId); + for (const viewId of ['RawLogsAllHistory', 'TransactionsAllHistory']) { + const view = await tableMetadata(handle.datasetId, viewId); + assert(view.type === 'VIEW', `${viewId} is not a view`); + assert(view.view?.query?.includes('block_timestamp'), `${viewId} is missing its partition-bounded definition`); + } + + report.row_preservation = { PipelineRuns: pipelinePreservation, OracleReconciliation: reconciliationPreservation }; + report.schema_checks = { + RawLogs: { fields: rawLogs.schema?.fields?.length, matches_runtime_contract: true, partition: rawLogs.timePartitioning, require_partition_filter: rawLogs.requirePartitionFilter }, + Transactions: { fields: transactions.schema?.fields?.length, rows: Number(transactions.numRows ?? 0), matches_runtime_contract: true }, + PipelineRuns: { fields: pipelineRuns.schema?.fields?.length, matches_runtime_contract: true }, + IngestionCoverage: { fields: productionCoverage.schema?.fields?.length, matches_runtime_contract: true }, + OracleReconciliation: { fields: reconciliation.schema?.fields?.length, matches_runtime_fields: true, legacy_table_id_retained: reconciliationFields.has('table_id') }, + views: ['RawLogsAllHistory', 'TransactionsAllHistory'], + }; +} catch (error) { + report.errors.push(String(error?.message ?? error)); +} finally { + if (handle) { + try { + cleanup = await sandbox.dropSandbox(handle); + report.cleanup = cleanup; + } catch (error) { + report.cleanup = { provenAbsent: false, error: String(error?.message ?? error) }; + report.errors.push(`Sandbox cleanup failed: ${report.cleanup.error}`); + } + } +} + +report.bytes_billed_total = report.jobs.reduce((sum, item) => sum + item.bytes_billed, 0); +report.finished_at_utc = new Date().toISOString(); +mkdirSync(dirname(outputPath), { recursive: true }); +writeFileSync(outputPath, `${JSON.stringify(report, null, 2)}\n`); +console.log(`sandbox dataset : ${report.sandbox_dataset ?? 'not created'}`); +console.log(`jobs : ${report.jobs.length}`); +console.log(`bytes billed : ${report.bytes_billed_total}`); +console.log(`cleanup proven : ${report.cleanup?.provenAbsent === true}`); +console.log(`report : ${outputPath}`); +for (const error of report.errors) console.error(`validation error: ${error}`); + +if (report.errors.length || report.cleanup?.provenAbsent !== true) process.exit(1); \ No newline at end of file diff --git a/projects/onchain-analytics/scripts/tests/deploy-warehouse.Tests.ps1 b/projects/onchain-analytics/scripts/tests/deploy-warehouse.Tests.ps1 new file mode 100644 index 0000000..4e17f8c --- /dev/null +++ b/projects/onchain-analytics/scripts/tests/deploy-warehouse.Tests.ps1 @@ -0,0 +1,108 @@ +$ErrorActionPreference = "Continue" +$testDir = Split-Path -Parent $MyInvocation.MyCommand.Path +$repoRoot = Split-Path -Parent (Split-Path -Parent $testDir) +$deployScript = Join-Path $repoRoot "scripts\deploy-warehouse.ps1" +$refused = @( + "99_UnlistedUnbannered.sql", + "04_L0Contract_v3.sql", + "06_L0Contract_v4.sql", + "07_RetireV3EventTables.sql" +) + +foreach ($migration in $refused) { + $output = & powershell.exe -NoProfile -NonInteractive -File $deployScript -Migration $migration 2>&1 + $exitCode = $LASTEXITCODE + $text = $output -join "`n" + if ($exitCode -eq 0 -or $text -notmatch "REFUSED_UNLISTED_MIGRATION") { + throw "Expected refusal for $migration; exit=$exitCode output=$text" + } + Write-Host "PASS refused $migration before query submission" +} + +foreach ($migration in @( + "08_PipelineRunsOutcome_v1.sql", + "09_CreateRawLogs_v1.sql", + "10_AddOracleReconciliationCompatibility_v1.sql", + "11_CreateRawLogsAllHistory_v1.sql", + "12_CreateTransactionsAllHistory_v1.sql" +)) { + $plan = & powershell.exe -NoProfile -NonInteractive -File $deployScript -Migration $migration 2>&1 + if ($LASTEXITCODE -ne 0 -or ($plan -join "`n") -notmatch "PLAN ONLY") { + throw "An allowlisted migration must default to plan-only: $migration. output=$($plan -join "`n")" + } +} +Write-Host "PASS all allowlisted migrations default to plan-only" + +$impersonation = & powershell.exe -NoProfile -NonInteractive -File $deployScript -Migration "09_CreateRawLogs_v1.sql" -Execute -AllowProduction 2>&1 +$impersonationExit = $LASTEXITCODE +if ($impersonationExit -eq 0 -or ($impersonation -join "`n") -notmatch "REFUSED_IMPERSONATION_REQUIRED") { + throw "Production execution without an approved service account must refuse. output=$($impersonation -join "`n")" +} +Write-Host "PASS production execution requires an explicit service account" + +$ErrorActionPreference = "Stop" +$tokens = $null +$parseErrors = $null +$deployAst = [System.Management.Automation.Language.Parser]::ParseFile($deployScript, [ref]$tokens, [ref]$parseErrors) +if ($parseErrors.Count -gt 0) { + throw "Deployment script must parse before identity checks" +} +$scopedBuilders = $deployAst.FindAll({ + param($node) + $node -is [System.Management.Automation.Language.FunctionDefinitionAst] -and + $node.Name -eq 'New-BqProcessStartInfo' +}, $true) +if ($scopedBuilders.Count -ne 1) { + throw "Deployment must build a child process with its own approved identity" +} +if ($deployAst.Extent.Text -match 'config\s+(set|unset)\s+auth/impersonate_service_account') { + throw "Deployment must not mutate shared gcloud impersonation settings" +} +Invoke-Expression $scopedBuilders[0].Extent.Text +$parentIdentity = $env:CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT +$first = New-Object System.Diagnostics.Process +$second = New-Object System.Diagnostics.Process +try { + $first.StartInfo = New-BqProcessStartInfo -BqExe $env:ComSpec -Arguments @('/d', '/c', 'echo', '%CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT%') -ImpersonateServiceAccount 'review-a@example.iam.gserviceaccount.com' + $second.StartInfo = New-BqProcessStartInfo -BqExe $env:ComSpec -Arguments @('/d', '/c', 'echo', '%CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT%') -ImpersonateServiceAccount 'review-b@example.iam.gserviceaccount.com' + [void]$first.Start() + [void]$second.Start() + $first.StandardInput.Close() + $second.StandardInput.Close() + $firstOutput = $first.StandardOutput.ReadToEnd() + $secondOutput = $second.StandardOutput.ReadToEnd() + $first.WaitForExit() + $second.WaitForExit() + if ($first.ExitCode -ne 0 -or $firstOutput.Trim() -cne 'review-a@example.iam.gserviceaccount.com') { + throw "The first deployment child lost its approved identity: $firstOutput" + } + if ($second.ExitCode -ne 0 -or $secondOutput.Trim() -cne 'review-b@example.iam.gserviceaccount.com') { + throw "The second deployment child lost its approved identity: $secondOutput" + } + if ($env:CLOUDSDK_AUTH_IMPERSONATE_SERVICE_ACCOUNT -cne $parentIdentity) { + throw "Deployment changed the caller's impersonation environment" + } +} finally { + $first.Dispose() + $second.Dispose() +} +Write-Host "PASS overlapping children keep separate identities without modifying the caller" + +$gcloudCommand = Get-Command gcloud.cmd -ErrorAction SilentlyContinue +if ($gcloudCommand) { + $sdkProbe = New-Object System.Diagnostics.Process + try { + $sdkProbe.StartInfo = New-BqProcessStartInfo -BqExe $gcloudCommand.Source -Arguments @('config', 'get-value', 'auth/impersonate_service_account') -ImpersonateServiceAccount 'review-a@example.iam.gserviceaccount.com' + [void]$sdkProbe.Start() + $sdkProbe.StandardInput.Close() + $sdkOutputTask = $sdkProbe.StandardOutput.ReadToEndAsync() + $sdkErrorTask = $sdkProbe.StandardError.ReadToEndAsync() + $sdkProbe.WaitForExit() + if ($sdkProbe.ExitCode -ne 0 -or $sdkOutputTask.Result.Trim() -cne 'review-a@example.iam.gserviceaccount.com') { + throw "The installed SDK did not recognize the child-only identity: $($sdkErrorTask.Result)" + } + } finally { + $sdkProbe.Dispose() + } + Write-Host "PASS installed SDK recognizes child-only impersonation without token acquisition" +} \ No newline at end of file diff --git a/projects/onchain-analytics/scripts/tests/validate-l0-migrations.test.mjs b/projects/onchain-analytics/scripts/tests/validate-l0-migrations.test.mjs new file mode 100644 index 0000000..6d26acb --- /dev/null +++ b/projects/onchain-analytics/scripts/tests/validate-l0-migrations.test.mjs @@ -0,0 +1,66 @@ +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import { runInNewContext } from 'node:vm'; +import { test } from 'node:test'; + +const validator = readFileSync(new URL('../ops/validate-l0-migrations.mjs', import.meta.url), 'utf8'); +const parserSource = validator.match(/(function runtimeTableSchema\(tableId\) \{[\s\S]*?\r?\n\})\r?\n\r?\nfunction renderMigration/); +assert.ok(parserSource, 'The runtime schema parser must be found without executing the live validator'); + +function parseColumns(columns) { + const runtimeSource = [ + 'CREATE TABLE IF NOT EXISTS ${fullTableName("Fixture")} (', + ...columns.map((column) => ` ${column}`), + ' )', + ].join('\n'); + const fields = runInNewContext(`${parserSource[1]}; runtimeTableSchema("Fixture")`, { + runtimeSource, + assert: (condition, message) => assert.ok(condition, message), + }); + return JSON.parse(JSON.stringify(fields)); +} + +test('retains required columns from the runtime schema', () => { + assert.deepEqual(parseColumns(['run_id STRING NOT NULL,']), [ + { name: 'run_id', type: 'STRING', mode: 'REQUIRED' }, + ]); +}); + +test('keeps optional columns nullable', () => { + assert.deepEqual(parseColumns(['chain_id INT64,']), [ + { name: 'chain_id', type: 'INTEGER', mode: 'NULLABLE' }, + ]); +}); + +test('handles required and optional fields together, including case-insensitive declarations', () => { + assert.deepEqual(parseColumns(['started_at TIMESTAMP not null,', 'finished_at TIMESTAMP,']), [ + { name: 'started_at', type: 'TIMESTAMP', mode: 'REQUIRED' }, + { name: 'finished_at', type: 'TIMESTAMP', mode: 'NULLABLE' }, + ]); +}); + +const preservationChecks = validator.match(/ assert\(Number\(pipelinePreservation\.rows_preserved\)[\s\S]*?(?=\r?\n\r?\n await applyMigration)/); +assert.ok(preservationChecks, 'The live historical-row assertions must be found without executing migrations'); + +function checkHistoricalRows(unknownHashes) { + runInNewContext(preservationChecks[0], { + pipelinePreservation: { + rows_preserved: '22', + fixture_rows: '22', + old_rows_without_release_sha: String(unknownHashes), + }, + assert: (condition, message) => assert.ok(condition, message), + }); +} + +test('accepts preserved historical rows with all release hashes unknown', () => { + assert.doesNotThrow(() => checkHistoricalRows(22)); +}); + +test('rejects a release hash populated on even one historical row', () => { + assert.throws(() => checkHistoricalRows(21), /historical release_sha values changed/); +}); + +test('rejects populated release hashes even when all historical rows remain', () => { + assert.throws(() => checkHistoricalRows(0), /historical release_sha values changed/); +}); \ No newline at end of file diff --git a/projects/onchain-analytics/warehouse/L1/04_L0Contract_v3.sql b/projects/onchain-analytics/warehouse/L1/04_L0Contract_v3.sql index 98a5883..41d54a3 100644 --- a/projects/onchain-analytics/warehouse/L1/04_L0Contract_v3.sql +++ b/projects/onchain-analytics/warehouse/L1/04_L0Contract_v3.sql @@ -1,4 +1,6 @@ --- ================================================================================================= +-- DO NOT RUN. Historical v3 schema; retained as a reference only. +-- Apply only the explicitly allowlisted additive migrations through scripts/deploy-warehouse.ps1. +-- ================================================================================================= -- L0 INGESTION CONTRACT v3.0 -- ================================================================================================= -- The raw-event and state-snapshot tables for the Celo Phase 1 expansion. Numbered 04 because it diff --git a/projects/onchain-analytics/warehouse/L1/06_L0Contract_v4.sql b/projects/onchain-analytics/warehouse/L1/06_L0Contract_v4.sql index 499c2ef..cfd96b1 100644 --- a/projects/onchain-analytics/warehouse/L1/06_L0Contract_v4.sql +++ b/projects/onchain-analytics/warehouse/L1/06_L0Contract_v4.sql @@ -1,3 +1,5 @@ +-- DO NOT RUN DIRECTLY. Multi-statement historical schema reference, not a deployment command. +-- Apply only the explicitly allowlisted single-statement migrations through scripts/deploy-warehouse.ps1. -- ================================================================================================= -- L0 INGESTION CONTRACT v4.0 -- ================================================================================================= diff --git a/projects/onchain-analytics/warehouse/L1/08_PipelineRunsOutcome_v1.sql b/projects/onchain-analytics/warehouse/L1/08_PipelineRunsOutcome_v1.sql index 6d136b3..7a472e0 100644 --- a/projects/onchain-analytics/warehouse/L1/08_PipelineRunsOutcome_v1.sql +++ b/projects/onchain-analytics/warehouse/L1/08_PipelineRunsOutcome_v1.sql @@ -1,112 +1,62 @@ --- 08_PipelineRunsOutcome_v1.sql +-- PipelineRuns outcome and job-lineage columns. -- --- DO NOT RUN. This file is a migration DEFINITION, not a deployment step. +-- One additive statement. Every column is nullable so historical runs remain valid and read NULL +-- for values they did not record. ADD COLUMN IF NOT EXISTS preserves rows and makes reapplication +-- a no-op. Existing host and pipeline_version columns are reused. -- --- That banner is load bearing and it is the FIRST thing in this file on purpose. --- `scripts/deploy-warehouse.ps1` executes every .sql file in this folder in filename order unless --- its first twenty lines carry `DO NOT RUN` or `NOT THE LIVE SHAPE`. That bootstrap script is the --- readiness audit's finding C3, it is owned by Phase 8, and it is not fixed yet. Adding a file --- here without this banner would hand a live migration to a script that is already known to run --- things nobody asked it to run. The banner is the only control that exists today, so this file --- carries it, and it is not a substitute for fixing C3. --- --- ONE STATEMENT. ADDITIVE ONLY. NOT APPLIED BY THIS PHASE. --- --- WHAT THIS IS. The remediation plan's Phase 3 task 11 migration: the columns that let a --- PipelineRuns row say what a run actually did, rather than carrying two counters that were --- incremented separately from the outcomes they claim to describe. --- --- WHY IT EXISTS. The readiness audit's finding C2 has two runtime receipts and both end the same --- way. A bare `daily` planned 142 targets against a limit of 12, logged REFUSED, read nothing, --- and persisted ZERO planned and ZERO failed captures. An XDC backfill requested 12,394,045 --- blocks against a limit of 1,296,000, wrote one refusal row, and persisted ONE planned and ONE --- successful capture. In both cases the run record agreed with the exit code and both were wrong, --- because the only shapes available were "succeeded" and "failed" and a refusal is neither. These --- columns are the vocabulary that was missing. --- --- SAFETY PROPERTIES, each one deliberate. --- --- * Every clause is ADD COLUMN IF NOT EXISTS. There is no DROP, no ALTER COLUMN, no rename and --- no data movement anywhere in this file. Running it twice is a no-op. This matters here --- specifically: `04_L0Contract_v3.sql` carries three unconditional DROP TABLE statements and --- one of them was executed against production by accident in this project's own history. --- * One statement in one file, per plan section 1.2. Nothing extracts a fragment of this file --- by delimiter or string slicing: the incident above was caused by exactly that, on a CRLF --- file where the delimiter search returned -1 and took the rest of the file with it. --- * `host` and `pipeline_version` are REUSED for runner identity and release version. Plan task --- 11 says so explicitly and forbids adding aliases for them, so neither appears below. --- * Adding a column to a BigQuery table rewrites no rows and scans no bytes. Existing rows read --- NULL for every column below, which is correct: those runs genuinely did not record this. --- --- WHO APPLIES IT. Not this change. The migration is rehearsed against a sandbox dataset and --- applied to production only inside an authorised additive-DDL window. What lands here is the --- definition plus the code that reads and writes it, and `ensureInfraTables` asserts every --- column at startup so a dataset that has not had this applied fails immediately, naming this --- file, instead of failing at the INSERT after a whole run has already happened. --- --- TARGET. Substitute the dataset deliberately. This file names no project so it cannot be run --- against production by copy-paste alone. - +-- Apply through the explicit migration helper after sandbox validation and separate production +-- approval. The pipeline checks these columns at startup and refuses an older table shape. +-- Identifiers are literal placeholders rendered by the migration helper. ALTER TABLE `${PROJECT}.${DATASET}.PipelineRuns` - -- The outcome model. `execution_status` is the word form of the exit code, because exit 2 does - -- not distinguish a refusal from a crash and those need different responses. ADD COLUMN IF NOT EXISTS execution_status STRING - OPTIONS(description="completed, partial, refused, unsupported, empty or failed. Never plain success: a refused run and a clean run must not share a word."), + OPTIONS(description="completed, partial, refused, unsupported, empty or failed. A refusal is not a completed run."), ADD COLUMN IF NOT EXISTS units_planned INT64 - OPTIONS(description="Work units this run intended, recorded BEFORE any guard ran. A globally refused run has a real planned count here; that number does not exist after the refusal."), + OPTIONS(description="Work units declared before guards run."), ADD COLUMN IF NOT EXISTS units_attempted INT64 - OPTIONS(description="Units the pipeline actually tried to read. planned minus attempted is exactly the work this run declined."), + OPTIONS(description="Work units the pipeline attempted to read."), ADD COLUMN IF NOT EXISTS units_completed INT64 - OPTIONS(description="Units that finished their whole range. The only value that is success."), + OPTIONS(description="Work units whose complete range was captured."), ADD COLUMN IF NOT EXISTS units_noop INT64 - OPTIONS(description="Units whose range was empty by arithmetic. Not work, and not a completion: a run of pure no-ops exits nonzero."), + OPTIONS(description="Work units whose requested range was empty by arithmetic."), ADD COLUMN IF NOT EXISTS units_refused INT64 - OPTIONS(description="Units a budget guard declined before reading. The range is unread and recorded as unread in IngestionCoverage with status refused_budget."), + OPTIONS(description="Work units refused before any chain read."), ADD COLUMN IF NOT EXISTS units_unsupported INT64 - OPTIONS(description="Units this pipeline cannot do: no reader for the chain, chain outside the frozen release scope, or an address the control plane does not know."), + OPTIONS(description="Work units the configured readers or release scope do not support."), ADD COLUMN IF NOT EXISTS units_failed INT64 - OPTIONS(description="Units that ran and did not finish, whether they threw or returned an incomplete range."), + OPTIONS(description="Work units that ran but did not complete."), ADD COLUMN IF NOT EXISTS outcome_counts_by_grain JSON - OPTIONS(description="The same counters per target grain, RawLogs and Transactions separately, plus RawLogs+Transactions for run-level parent units. Reconciles exactly to the units_ columns."), + OPTIONS(description="Outcome counters for RawLogs and Transactions separately."), ADD COLUMN IF NOT EXISTS release_sha STRING - OPTIONS(description="40-character lowercase Git SHA of the release that produced this run. NULL until a release process sets it; a hash invented by the program it is supposed to bind would bind nothing."), + OPTIONS(description="40-character lowercase Git SHA of the release that produced this run."), ADD COLUMN IF NOT EXISTS plan_hash STRING - OPTIONS(description="64-character lowercase hex SHA-256 of the remediation plan this writer is bound to, per plan section 1.1."), - - -- The child-plan lineage. A run may be a child of another run's plan, and the readiness audit's - -- reconciliation requirement is that a parent's outcome be derivable from its children rather - -- than asserted alongside them. + OPTIONS(description="64-character lowercase SHA-256 of the immutable execution plan."), ADD COLUMN IF NOT EXISTS parent_run_id STRING - OPTIONS(description="The run that planned this one. NULL for an ordinary top-level run."), + OPTIONS(description="Parent run identifier, NULL for a top-level run."), ADD COLUMN IF NOT EXISTS child_plan_stage STRING - OPTIONS(description="Fixed lowercase stage enum this child belongs to."), + OPTIONS(description="Execution stage assigned to this child run."), ADD COLUMN IF NOT EXISTS child_plan_root_hash STRING - OPTIONS(description="Hash of the whole plan tree this child was derived from."), + OPTIONS(description="Hash of the complete child plan tree."), ADD COLUMN IF NOT EXISTS child_plan_hash STRING - OPTIONS(description="Hash of this child's own plan. Part of the canonical input to a deterministic child run id."), + OPTIONS(description="Hash of this child run's plan."), ADD COLUMN IF NOT EXISTS child_ordinal INT64 - OPTIONS(description="Position of this child in the global plan ordering."), + OPTIONS(description="Position in the complete child plan."), ADD COLUMN IF NOT EXISTS stage_child_ordinal INT64 - OPTIONS(description="Position of this child within its own stage."), + OPTIONS(description="Position within this stage's child plan."), ADD COLUMN IF NOT EXISTS stage_child_count INT64 - OPTIONS(description="How many children the stage declared, so a missing child is detectable without a scan."), - - -- The job ledger and the closure chain. Written by the mutation broker and the terminalizer in - -- later phases; defined here because plan task 11 fixes the column set and forbids Phase 8 from - -- redefining it. + OPTIONS(description="Number of children declared for this stage."), ADD COLUMN IF NOT EXISTS bigquery_job_ledger JSON - OPTIONS(description="Every WORK job this run submitted, with its submit intent and terminal state. The terminalizer is never in its own ledger."), + OPTIONS(description="Submitted work jobs and their terminal states."), ADD COLUMN IF NOT EXISTS job_ledger_hash STRING - OPTIONS(description="Hash over the complete ledger, computed before terminalization."), + OPTIONS(description="Hash of the complete job ledger."), ADD COLUMN IF NOT EXISTS work_jobs_terminal BOOL - OPTIONS(description="True only when every parent and generated child work job reached a terminal state. Parent status alone is insufficient."), + OPTIONS(description="True only when every recorded work job is terminal."), ADD COLUMN IF NOT EXISTS terminalizer_job_id STRING - OPTIONS(description="The single committed terminalizer attempt. Once a terminalized row exists, no second terminalizer patches it."), + OPTIONS(description="Identifier of the job that finalized this run record."), ADD COLUMN IF NOT EXISTS closure_status STRING - OPTIONS(description="terminalized, closed or closure_invalidated. terminalized is provisional and never means closed."), + OPTIONS(description="terminalized, closed or closure_invalidated."), ADD COLUMN IF NOT EXISTS terminalized_at TIMESTAMP - OPTIONS(description="Predetermined at intent creation and passed in as a typed parameter. The terminalizer may not call CURRENT_TIMESTAMP()."), + OPTIONS(description="Timestamp assigned when terminalization was requested."), ADD COLUMN IF NOT EXISTS terminalized_row_hash STRING - OPTIONS(description="Hash over the canonical final fields excluding this field itself, so a recovered row can be proven byte-identical to the intended one."), + OPTIONS(description="Hash of the canonical terminal row, excluding this field."), ADD COLUMN IF NOT EXISTS closure_receipt_uri STRING - OPTIONS(description="Predetermined URI of the closure object. That object alone makes a run closed."); + OPTIONS(description="URI of the record that proves closure."); diff --git a/projects/onchain-analytics/warehouse/L1/09_CreateRawLogs_v1.sql b/projects/onchain-analytics/warehouse/L1/09_CreateRawLogs_v1.sql new file mode 100644 index 0000000..a870dad --- /dev/null +++ b/projects/onchain-analytics/warehouse/L1/09_CreateRawLogs_v1.sql @@ -0,0 +1,36 @@ +-- One additive statement. Does not alter or replace an existing table. +-- Identifiers are rendered from literal placeholders by the deployment helper. +CREATE TABLE IF NOT EXISTS `${PROJECT}.${DATASET}.RawLogs` +( + chain_id INT64 NOT NULL, + block_number INT64 NOT NULL, + block_timestamp TIMESTAMP NOT NULL, + block_hash STRING NOT NULL, + tx_hash STRING NOT NULL, + tx_index INT64 NOT NULL, + log_index INT64 NOT NULL, + contract_address STRING NOT NULL, + implementation_address STRING, + era_index INT64, + era_resolution STRING NOT NULL, + topic0 STRING, + topic1 STRING, + topic2 STRING, + topic3 STRING, + topic_count INT64 NOT NULL, + log_data STRING NOT NULL, + removed BOOL, + source_kind STRING NOT NULL, + source_id STRING NOT NULL, + assurance STRING NOT NULL, + confirmations_at_capture INT64, + capture_id STRING NOT NULL, + ingestion_run_id STRING NOT NULL, + ingested_at TIMESTAMP NOT NULL +) +PARTITION BY TIMESTAMP_TRUNC(block_timestamp, MONTH) +CLUSTER BY chain_id, contract_address, topic0, block_number +OPTIONS( + require_partition_filter = TRUE, + description = "L0 raw log store. One row per log, keyed by (chain_id, tx_hash, log_index)." +); \ No newline at end of file diff --git a/projects/onchain-analytics/warehouse/L1/10_AddOracleReconciliationCompatibility_v1.sql b/projects/onchain-analytics/warehouse/L1/10_AddOracleReconciliationCompatibility_v1.sql new file mode 100644 index 0000000..679fcec --- /dev/null +++ b/projects/onchain-analytics/warehouse/L1/10_AddOracleReconciliationCompatibility_v1.sql @@ -0,0 +1,7 @@ +-- One additive statement. Existing table_id values and historical rows are retained. +-- Legacy rows remain NULL in the new dimensions; no chain or contract is inferred. +ALTER TABLE `${PROJECT}.${DATASET}.OracleReconciliation` + ADD COLUMN IF NOT EXISTS chain_id INT64 + OPTIONS(description="EVM chain id. NULL on historical rows written before this dimension existed."), + ADD COLUMN IF NOT EXISTS contract_address STRING + OPTIONS(description="Oracle contract address. NULL on historical rows written before this dimension existed."); \ No newline at end of file diff --git a/projects/onchain-analytics/warehouse/L1/11_CreateRawLogsAllHistory_v1.sql b/projects/onchain-analytics/warehouse/L1/11_CreateRawLogsAllHistory_v1.sql new file mode 100644 index 0000000..2182e7f --- /dev/null +++ b/projects/onchain-analytics/warehouse/L1/11_CreateRawLogsAllHistory_v1.sql @@ -0,0 +1,7 @@ +-- One additive statement. IF NOT EXISTS leaves an existing view unchanged for preflight review. +CREATE VIEW IF NOT EXISTS `${PROJECT}.${DATASET}.RawLogsAllHistory` +OPTIONS(description="All RawLogs rows through the supported timestamp range. Use for whole-history reads; filter the base table for a bounded window.") +AS SELECT * +FROM `${PROJECT}.${DATASET}.RawLogs` +WHERE block_timestamp >= TIMESTAMP('2000-01-01 00:00:00') + AND block_timestamp < TIMESTAMP('2100-01-01 00:00:00'); \ No newline at end of file diff --git a/projects/onchain-analytics/warehouse/L1/12_CreateTransactionsAllHistory_v1.sql b/projects/onchain-analytics/warehouse/L1/12_CreateTransactionsAllHistory_v1.sql new file mode 100644 index 0000000..4d527f6 --- /dev/null +++ b/projects/onchain-analytics/warehouse/L1/12_CreateTransactionsAllHistory_v1.sql @@ -0,0 +1,7 @@ +-- One additive statement. IF NOT EXISTS leaves an existing view unchanged for preflight review. +CREATE VIEW IF NOT EXISTS `${PROJECT}.${DATASET}.TransactionsAllHistory` +OPTIONS(description="All Transactions rows through the supported timestamp range. Use for whole-history reads; filter the base table for a bounded window.") +AS SELECT * +FROM `${PROJECT}.${DATASET}.Transactions` +WHERE block_timestamp >= TIMESTAMP('2000-01-01 00:00:00') + AND block_timestamp < TIMESTAMP('2100-01-01 00:00:00'); \ No newline at end of file