starting phase 5; clap cli
This commit is contained in:
+32
@@ -0,0 +1,32 @@
|
||||
use clap::{Parser, Subcommand};
|
||||
use std::path::PathBuf;
|
||||
|
||||
/// Faroese fisheries data pipeline and dashboard
|
||||
#[derive(Parser)]
|
||||
#[command(name = "hagfish", version, about)]
|
||||
pub struct Cli {
|
||||
/// Path to configuration file
|
||||
#[arg(short, long, value_name = "FILE", default_value = "config.json")]
|
||||
pub config: PathBuf,
|
||||
|
||||
#[command(subcommand)]
|
||||
pub command: Command,
|
||||
}
|
||||
|
||||
#[derive(Subcommand)]
|
||||
pub enum Command {
|
||||
/// Fetch data from Hagstova API and store in DuckDB
|
||||
Ingest {
|
||||
/// Force full backfill instead of incremental
|
||||
#[arg(long)]
|
||||
full: bool,
|
||||
},
|
||||
/// Start the HTTP server
|
||||
Serve,
|
||||
/// Export landings data to Parquet
|
||||
Export {
|
||||
/// Output file path
|
||||
#[arg(short, long, value_name = "FILE")]
|
||||
out: PathBuf,
|
||||
},
|
||||
}
|
||||
+2
-2
@@ -22,7 +22,7 @@ fn is_sentinel(v: f64) -> bool {
|
||||
.any(|s| v.total_cmp(s) == std::cmp::Ordering::Equal)
|
||||
}
|
||||
|
||||
pub async fn fetch_metadata(client: &Client, url: &str) -> Result<LookupMap> {
|
||||
pub async fn fetch_metadata(client: &Client, url: &str) -> Result<(LookupMap, MetadataResponse)> {
|
||||
info!("Fetching metadata from {}", url);
|
||||
|
||||
let resp = client
|
||||
@@ -63,7 +63,7 @@ pub async fn fetch_metadata(client: &Client, url: &str) -> Result<LookupMap> {
|
||||
lookup_map.insert(var.code.clone(), lookup);
|
||||
}
|
||||
|
||||
Ok(lookup_map)
|
||||
Ok((lookup_map, meta))
|
||||
}
|
||||
|
||||
pub fn extract_available_months(meta: &MetadataResponse) -> Vec<String> {
|
||||
|
||||
+74
-3
@@ -1,10 +1,14 @@
|
||||
mod api;
|
||||
mod cli;
|
||||
mod db;
|
||||
mod ingest;
|
||||
mod types;
|
||||
|
||||
use axum::serve;
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::serve;
|
||||
use clap::Parser;
|
||||
use tokio::net::TcpListener;
|
||||
use tracing_subscriber::EnvFilter;
|
||||
|
||||
@@ -36,15 +40,24 @@ async fn main() -> anyhow::Result<()> {
|
||||
.with_env_filter(EnvFilter::from_default_env())
|
||||
.init();
|
||||
|
||||
let config = Arc::new(types::Config::default());
|
||||
let cli = cli::Cli::parse();
|
||||
let config = types::Config::load(&cli.config)?;
|
||||
|
||||
match cli.command {
|
||||
cli::Command::Serve => run_serve(&config).await,
|
||||
cli::Command::Ingest { full } => run_ingest(&config, full).await,
|
||||
cli::Command::Export { out } => run_export(&config, &out),
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_serve(config: &types::Config) -> anyhow::Result<()> {
|
||||
cleanup_stale_parquet_files();
|
||||
|
||||
let conn = db::init(&config.duckdb_path)?;
|
||||
|
||||
let state = api::AppState {
|
||||
conn: Arc::new(tokio::sync::Mutex::new(conn)),
|
||||
config: config.clone(),
|
||||
config: Arc::new(config.clone()),
|
||||
};
|
||||
|
||||
let app = api::build_router(state);
|
||||
@@ -56,3 +69,61 @@ async fn main() -> anyhow::Result<()> {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> {
|
||||
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?;
|
||||
|
||||
db::update_lookups(&conn, &lookup_maps)?;
|
||||
|
||||
let all_months = ingest::extract_available_months(&metadata);
|
||||
|
||||
let months_to_fetch: Vec<String> = if full {
|
||||
all_months
|
||||
} else {
|
||||
let last = db::get_last_month(&conn)?;
|
||||
match last {
|
||||
Some(last_month) => all_months
|
||||
.into_iter()
|
||||
.filter(|m| m.as_str() > last_month.as_str())
|
||||
.collect(),
|
||||
None => all_months,
|
||||
}
|
||||
};
|
||||
|
||||
if months_to_fetch.is_empty() {
|
||||
tracing::info!("No new months to ingest");
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
tracing::info!("Ingesting {} month(s)", months_to_fetch.len());
|
||||
|
||||
// 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?;
|
||||
|
||||
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 count = db::upsert_landings(&conn, &landings)?;
|
||||
tracing::info!("Ingested {} rows across {} month(s)", count, months_to_fetch.len());
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn run_export(config: &types::Config, out: &Path) -> anyhow::Result<()> {
|
||||
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);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
pub struct Config {
|
||||
@@ -28,6 +29,20 @@ impl Default for Config {
|
||||
}
|
||||
}
|
||||
|
||||
impl Config {
|
||||
pub fn load(path: &Path) -> anyhow::Result<Self> {
|
||||
if path.exists() {
|
||||
let contents = std::fs::read_to_string(path)?;
|
||||
let config: Config = serde_json::from_str(&contents)?;
|
||||
tracing::info!("Loaded config from {}", path.display());
|
||||
Ok(config)
|
||||
} else {
|
||||
tracing::warn!("Config file {} not found, using defaults", path.display());
|
||||
Ok(Self::default())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
pub struct MetadataResponse {
|
||||
#[serde(default)]
|
||||
|
||||
Reference in New Issue
Block a user