phase 5 complete, sending to QA and code review
This commit is contained in:
+150
-18
@@ -6,12 +6,26 @@ mod types;
|
||||
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
use axum::serve;
|
||||
use clap::Parser;
|
||||
use tokio::net::TcpListener;
|
||||
use tracing_subscriber::EnvFilter;
|
||||
|
||||
fn init_logging() {
|
||||
let filter = EnvFilter::from_default_env();
|
||||
|
||||
if std::env::var("HAGFISH_LOG_FORMAT").as_deref() == Ok("json") {
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(filter)
|
||||
.json()
|
||||
.init();
|
||||
} else {
|
||||
tracing_subscriber::fmt().with_env_filter(filter).init();
|
||||
}
|
||||
}
|
||||
|
||||
fn cleanup_stale_parquet_files() {
|
||||
let tmp_dir = std::env::temp_dir();
|
||||
match std::fs::read_dir(&tmp_dir) {
|
||||
@@ -30,15 +44,13 @@ fn cleanup_stale_parquet_files() {
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => tracing::warn!("Failed to list temp dir: {e}"),
|
||||
Err(e) => tracing::warn!(error = %e, "Failed to list temp dir"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> anyhow::Result<()> {
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(EnvFilter::from_default_env())
|
||||
.init();
|
||||
init_logging();
|
||||
|
||||
let cli = cli::Cli::parse();
|
||||
let config = types::Config::load(&cli.config)?;
|
||||
@@ -63,7 +75,7 @@ async fn run_serve(config: &types::Config) -> anyhow::Result<()> {
|
||||
let app = api::build_router(state);
|
||||
|
||||
let listener = TcpListener::bind(&config.bind_address).await?;
|
||||
tracing::info!("hagfish listening on http://{}", config.bind_address);
|
||||
tracing::info!(address = %config.bind_address, "hagfish server started");
|
||||
|
||||
serve(listener, app).await?;
|
||||
|
||||
@@ -71,10 +83,21 @@ async fn run_serve(config: &types::Config) -> anyhow::Result<()> {
|
||||
}
|
||||
|
||||
async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> {
|
||||
let run_start = Instant::now();
|
||||
|
||||
tracing::info!(
|
||||
mode = if full { "full" } else { "incremental" },
|
||||
"starting ingestion run"
|
||||
);
|
||||
|
||||
let conn = db::init(&config.duckdb_path)?;
|
||||
let client = reqwest::Client::new();
|
||||
|
||||
let (lookup_maps, metadata) = ingest::fetch_metadata(&client, &config.data_source_url).await?;
|
||||
let (lookup_maps, metadata) = ingest::fetch_metadata(&client, &config.data_source_url)
|
||||
.await
|
||||
.inspect_err(|e| {
|
||||
tracing::error!(error = %e, stage = "fetch_metadata", "ingestion failed");
|
||||
})?;
|
||||
|
||||
db::update_lookups(&conn, &lookup_maps)?;
|
||||
|
||||
@@ -94,36 +117,145 @@ async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> {
|
||||
};
|
||||
|
||||
if months_to_fetch.is_empty() {
|
||||
tracing::info!("No new months to ingest");
|
||||
let elapsed = run_start.elapsed();
|
||||
tracing::info!(
|
||||
mode = if full { "full" } else { "incremental" },
|
||||
months_requested = 0,
|
||||
total_rows = 0,
|
||||
batches = 0,
|
||||
errors = 0,
|
||||
elapsed_ms = elapsed.as_millis() as u64,
|
||||
"ingestion complete — no new months"
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
tracing::info!("Ingesting {} month(s)", months_to_fetch.len());
|
||||
const BATCH_SIZE: usize = 12;
|
||||
let batches = ingest::chunk_months(&months_to_fetch, BATCH_SIZE);
|
||||
|
||||
// TODO: 5.2 — batch into chunks of <=12 months when count exceeds 12
|
||||
let query = ingest::build_query(&months_to_fetch);
|
||||
let data = ingest::fetch_data(&client, &config.data_source_url, &query).await?;
|
||||
tracing::info!(
|
||||
mode = if full { "full" } else { "incremental" },
|
||||
months_requested = months_to_fetch.len(),
|
||||
batches = batches.len(),
|
||||
"beginning batched ingestion"
|
||||
);
|
||||
|
||||
let mut landings = Vec::with_capacity(data.dataset.value.len());
|
||||
for i in 0..data.dataset.value.len() {
|
||||
let row = ingest::parse_row(i, &data, &lookup_maps)?;
|
||||
landings.push(ingest::data_row_to_landing(&row));
|
||||
let mut total_rows = 0;
|
||||
let mut total_skipped = 0;
|
||||
let mut errors = 0;
|
||||
|
||||
for (i, batch) in batches.iter().enumerate() {
|
||||
let batch_start = Instant::now();
|
||||
|
||||
tracing::info!(
|
||||
batch = i + 1,
|
||||
total_batches = batches.len(),
|
||||
months_in_batch = batch.len(),
|
||||
months = %batch.join(", "),
|
||||
"fetching batch"
|
||||
);
|
||||
|
||||
let query = ingest::build_query(batch);
|
||||
|
||||
let data = match ingest::fetch_data(&client, &config.data_source_url, &query).await {
|
||||
Ok(d) => d,
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
batch = i + 1,
|
||||
error = %e,
|
||||
stage = "fetch_data",
|
||||
"batch failed, continuing to next batch"
|
||||
);
|
||||
errors += 1;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let mut landings = Vec::with_capacity(data.dataset.value.len());
|
||||
let mut parse_errors = 0;
|
||||
|
||||
for j in 0..data.dataset.value.len() {
|
||||
match ingest::parse_row(j, &data.dataset, &lookup_maps) {
|
||||
Ok(row) => {
|
||||
if row.value.is_some() {
|
||||
landings.push(ingest::data_row_to_landing(&row));
|
||||
} else {
|
||||
total_skipped += 1;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
batch = i + 1,
|
||||
row_index = j,
|
||||
error = %e,
|
||||
"failed to parse row"
|
||||
);
|
||||
parse_errors += 1;
|
||||
errors += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let count = match db::upsert_landings(&conn, &landings) {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
batch = i + 1,
|
||||
error = %e,
|
||||
stage = "upsert_landings",
|
||||
"batch upsert failed"
|
||||
);
|
||||
errors += 1;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
total_rows += count;
|
||||
let batch_elapsed = batch_start.elapsed();
|
||||
|
||||
tracing::info!(
|
||||
batch = i + 1,
|
||||
total_batches = batches.len(),
|
||||
rows_inserted = count,
|
||||
null_rows_skipped = total_skipped,
|
||||
parse_errors = parse_errors,
|
||||
elapsed_ms = batch_elapsed.as_millis() as u64,
|
||||
"batch complete"
|
||||
);
|
||||
}
|
||||
|
||||
let count = db::upsert_landings(&conn, &landings)?;
|
||||
tracing::info!("Ingested {} rows across {} month(s)", count, months_to_fetch.len());
|
||||
let elapsed = run_start.elapsed();
|
||||
|
||||
tracing::info!(
|
||||
mode = if full { "full" } else { "incremental" },
|
||||
months_requested = months_to_fetch.len(),
|
||||
batches = batches.len(),
|
||||
total_rows = total_rows,
|
||||
null_rows_skipped = total_skipped,
|
||||
errors = errors,
|
||||
elapsed_ms = elapsed.as_millis() as u64,
|
||||
"ingestion run complete"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn run_export(config: &types::Config, out: &Path) -> anyhow::Result<()> {
|
||||
let run_start = Instant::now();
|
||||
|
||||
let conn = db::init(&config.duckdb_path)?;
|
||||
let path_str = out
|
||||
.to_str()
|
||||
.ok_or_else(|| anyhow::anyhow!("Output path contains invalid UTF-8"))?;
|
||||
|
||||
db::export_parquet(&conn, path_str)?;
|
||||
tracing::info!("Exported landings to {}", path_str);
|
||||
|
||||
let elapsed = run_start.elapsed();
|
||||
tracing::info!(
|
||||
path = path_str,
|
||||
elapsed_ms = elapsed.as_millis() as u64,
|
||||
"export complete"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user