Repository navigation
Conversation
Bulk-write the full ASN info dataset into PostgreSQL via streaming COPY with an atomic staging-table swap, so API readers never see an empty table. Reuses the same bgpkit-commons source build as generate, so field parity is by construction. - asninfo.current: one row per ASN; typed search columns (asn, name, country, country_name, org_id, org_name) plus per-source JSONB columns (population, hegemony, peeringdb; SQL NULL when absent) and data_as_of/source_revision provenance - streaming CSV COPY into an index-free staging table; indexes built on staging (PK, country, trigram GIN on name/org_name when pg_trgm is installed), then DROP current + RENAME in a single transaction - trigram indexes degrade to a warning when pg_trgm is unavailable - every run recorded in asninfo.ingest_run (success and failure) - connection via DATABASE_URL env var or --database-url flag - dataset load runs on the blocking pool: the blocking HTTP clients used by bgpkit-commons own nested tokio runtimes that must not be dropped inside an async context (tokio >= 1.48 panics) - tokio-postgres 0.7.18 (fixes RUSTSEC-2026-0178) - also fix three pre-existing clippy warnings in main.rs so the crate passes clippy -D warnings
There was a problem hiding this comment.
🟡 Changes recommended
Critical concurrency, TLS, and empty-dataset safeguards remain unresolved.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
Adds a pg-write command for bulk ASN ingestion into PostgreSQL using streaming COPY and staging-table replacement.
Changes:
- Adds PostgreSQL ingestion, indexing, and provenance tracking.
- Registers and documents the new command.
- Adds serialization tests and PostgreSQL dependencies.
Review findings:
- Critical (2 votes),
src/pg_write.rs:151: Concurrent runs can interfere through the shared staging table. Hold a session-level advisory lock for the complete load. - Moderate (2 votes),
src/pg_write.rs:249: Replacingcurrentloses grants and fails with dependent objects. Preserve the live relation or define an explicit grant/dependency strategy. - Critical (2 votes),
src/pg_write.rs:111:NoTlsbreaks TLS-required URLs and can expose traffic. Use TLS that honorssslmode. - Moderate (1 vote),
src/pg_write.rs:359: Failure records can contain fabricated provenance and row counts. Store NULL or actual progress. - Nit (2 votes),
src/pg_write.rs:245: Add lifecycle integration tests covering refreshes, rollback, indexes, and provenance. - Moderate (2 votes),
src/pg_write.rs:230: Schema-qualifygin_trgm_opsusing the extension namespace. - Critical (1 vote),
src/pg_write.rs:198: Reject unexpectedly empty datasets before replacing production data.
File summaries
| File | Description |
|---|---|
src/pg_write.rs |
Implements PostgreSQL ingestion, staging-table publication, indexing, provenance, and serialization tests. |
src/main.rs |
Registers and dispatches pg-write. |
README.md |
Documents PostgreSQL usage and schema. |
CHANGELOG.md |
Records the new feature. |
Cargo.toml |
Adds PostgreSQL and streaming dependencies. |
Cargo.lock |
Locks the updated dependency graph. |
Review details
- Files reviewed: 5/6 changed files
- Comments generated: 7
- Review effort level: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| let data_as_of = Utc::now(); | ||
|
|
||
| info!("connecting to PostgreSQL ..."); | ||
| let (mut client, connection) = match tokio_postgres::connect(database_url, NoTls).await { |
| data_as_of: DateTime<Utc>, | ||
| started_at: DateTime<Utc>, | ||
| ) -> Result<u64, (i32, String)> { | ||
| ensure_schema(client).await?; |
| let copied_rows = sink | ||
| .finish() | ||
| .await | ||
| .map_err(|e| (15, format!("failed to finish COPY: {e}")))?; | ||
| info!("COPY complete: {copied_rows} rows in staging table"); |
| "CREATE INDEX current_staging_name_trgm_idx ON asninfo.current_staging USING gin (name gin_trgm_ops)", | ||
| "CREATE INDEX current_staging_org_name_trgm_idx ON asninfo.current_staging USING gin (org_name gin_trgm_ops)", |
| .transaction() | ||
| .await | ||
| .map_err(|e| (15, format!("failed to begin swap transaction: {e}")))?; | ||
| tx.batch_execute("DROP TABLE IF EXISTS asninfo.current") |
| &0_i64, | ||
| &started_at, |
| let tx = client | ||
| .transaction() | ||
| .await | ||
| .map_err(|e| (15, format!("failed to begin swap transaction: {e}")))?; | ||
| tx.batch_execute("DROP TABLE IF EXISTS asninfo.current") |
Accepted and fixed: - serialize concurrent runs with a session-level advisory lock (0x41534E494E464F21) held for the whole load - honest failure provenance: error ingest_run rows record NULL row_count/data_as_of instead of fabricated values - schema-qualify gin_trgm_ops from the extension's actual namespace - refuse to swap when the loaded dataset is empty or has fewer than half the rows currently loaded (pure fn, unit-tested) - document that per-table grants do not survive the swap and must be granted via default privileges Declined: - TLS support (NoTls): the loader is designed for a local/trusted connection; a TLS-required server fails the connection safely. Adding sslmode support would pull in a TLS stack, out of scope. - lifecycle integration tests in CI: the repo has no PostgreSQL service; lifecycle behavior is exercised against a local PostgreSQL (two consecutive refreshes, advisory-lock blocking with a competing session, pg_trgm present/absent, provenance rows).
|
Thanks for the review. Responses per finding:
|
|
Superseded by direction change (2026-09-08): instead of extending the asninfo crate, the ASN/PeeringDB/IRR data load into PostgreSQL moves to a new standalone inserter with per-source schemas. The validated PG mechanics (advisory lock, staging swap, empty-dataset guard, CSV streaming, ingest_run provenance) will be carried over. Reference implementation kept on branch |
Adds an
asninfo pg-writesubcommand that bulk-writes the full ASN-info dataset into PostgreSQL via streaming COPY with an atomic staging-table swap, so API readers never see an empty table.What it does
bgpkit-commonssource build asgenerate(ASN names, as2org, population, hegemony, PeeringDB, countries), so field parity is by construction.asninfo.current_stagingvia CSV-format COPY, builds indexes on staging (PK onasn,countrybtree, trigram GIN onname/org_namewhenpg_trgmis installed), then atomically swaps it overasninfo.currentin one transaction (DROP current+RENAME), carrying index/constraint names to canonical names.pg_trgmis not installed, the trigram indexes are skipped with a warning instead of failing.asninfo.ingest_run(success and failure rows): task, status, row count,data_as_of,source_revision, timing, and error message.Row shape (
asninfo.current)asn bigint PK,name,country,country_name,org_id,org_namepopulation,hegemony,peeringdb(SQL NULL when a source has no data for the ASN)data_as_of,source_revisionUsage
--database-urlis accepted too; the env var is preferred (keeps credentials out of process listings).Notes
bgpkit-commonsown nested tokio runtimes, and dropping such a runtime inside an async context panics on tokio >= 1.48. Note this same nested-runtime panic currently affects the pre-existinggenerate/servedata loads on main; this PR applies the blocking-pool fix to the newpg-writepath only, fixinggenerate/servewould be a separate change.main.rsso the crate passescargo clippy --all-targets -- -D warnings.pg_trgm, verifying the swap, index renames, provenance rows, and the no-pg_trgmdegrade path.