use std::path::PathBuf;
use std::str::FromStr;
use anyhow::Context;
use chrono::{DateTime, Utc};
use clap::{ArgGroup, Args, Subcommand};
use nominal::core::{
CsvIngest, DatasetCreate, DatasetTarget, IngestJob, IngestJobStatus, NominalClient,
ParquetIngest, TimeUnit, Timestamp,
};
#[derive(Subcommand)]
pub enum IngestCommands {
Csv(CsvArgs),
Parquet(ParquetArgs),
}
#[derive(Args)]
pub struct CsvArgs {
#[command(flatten)]
common: UploadArgs,
}
#[derive(Args)]
pub struct ParquetArgs {
#[command(flatten)]
common: UploadArgs,
#[arg(long)]
archive: bool,
}
#[derive(Args)]
#[command(group(
ArgGroup::new("target").required(true).args(["dataset", "name"])
))]
struct UploadArgs {
path: PathBuf,
#[arg(long, value_name = "RID")]
dataset: Option<String>,
#[arg(long)]
name: Option<String>,
#[arg(long, requires = "name")]
description: Option<String>,
#[arg(long = "label", value_name = "LABEL", requires = "name")]
labels: Vec<String>,
#[arg(
long = "property",
value_names = ["KEY", "VALUE"],
num_args = 2,
action = clap::ArgAction::Append,
requires = "name",
)]
properties: Vec<String>,
#[arg(long, value_name = "COLUMN")]
timestamp_column: String,
#[arg(
long,
value_name = "SPEC",
long_help = "Timestamp encoding. One of:\n\
\x20\x20iso8601\n\
\x20\x20<unit> epoch timestamps (e.g. milliseconds, ns)\n\
Combine a unit with --relative-to to treat values as offsets from a start time."
)]
timestamp_type: TimestampSpec,
#[arg(long, value_name = "RFC3339")]
relative_to: Option<DateTime<Utc>>,
#[arg(long, value_name = "PREFIX")]
channel_prefix: Option<String>,
#[arg(
long = "tag-column",
value_names = ["TAG", "COLUMN"],
num_args = 2,
action = clap::ArgAction::Append,
)]
tag_columns: Vec<String>,
#[arg(
long = "file-tag",
value_names = ["KEY", "VALUE"],
num_args = 2,
action = clap::ArgAction::Append,
)]
file_tags: Vec<String>,
#[arg(long = "exclude-column", value_name = "COLUMN")]
exclude_columns: Vec<String>,
#[arg(long)]
wait: bool,
}
#[derive(Clone, Debug)]
enum TimestampSpec {
Iso8601,
Epoch(TimeUnit),
}
impl FromStr for TimestampSpec {
type Err = String;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let lowered = s.to_ascii_lowercase();
if lowered == "iso8601" {
return Ok(Self::Iso8601);
}
parse_time_unit(&lowered).map(Self::Epoch)
}
}
fn parse_time_unit(s: &str) -> Result<TimeUnit, String> {
match s.trim().to_ascii_lowercase().as_str() {
"ns" | "nanos" | "nanoseconds" => Ok(TimeUnit::Nanoseconds),
"us" | "micros" | "microseconds" => Ok(TimeUnit::Microseconds),
"ms" | "millis" | "milliseconds" => Ok(TimeUnit::Milliseconds),
"s" | "secs" | "seconds" => Ok(TimeUnit::Seconds),
other => Err(format!(
"unknown timestamp type '{other}': expected iso8601 or a time unit (nanoseconds, microseconds, milliseconds, seconds)"
)),
}
}
pub async fn handle(cmd: IngestCommands, client: NominalClient) -> anyhow::Result<()> {
match cmd {
IngestCommands::Csv(args) => handle_csv(args, client).await,
IngestCommands::Parquet(args) => handle_parquet(args, client).await,
}
}
async fn handle_csv(args: CsvArgs, client: NominalClient) -> anyhow::Result<()> {
let CsvArgs { common } = args;
let target = build_target(&common);
let timestamp = build_timestamp(&common)?;
let mut ingest = CsvIngest::new(timestamp);
if let Some(prefix) = &common.channel_prefix {
ingest = ingest.channel_prefix(prefix);
}
for pair in common.tag_columns.chunks(2) {
ingest = ingest.tag_column(&pair[0], &pair[1]);
}
for pair in common.file_tags.chunks(2) {
ingest = ingest.additional_file_tag(&pair[0], &pair[1]);
}
for col in &common.exclude_columns {
ingest = ingest.exclude_column(col);
}
let path = common.path.clone();
let job = client
.ingest()
.upload_csv(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload CSV '{}'", path.display()))?;
print_result(&job, common.wait, client).await
}
async fn handle_parquet(args: ParquetArgs, client: NominalClient) -> anyhow::Result<()> {
let ParquetArgs { common, archive } = args;
let target = build_target(&common);
let timestamp = build_timestamp(&common)?;
let mut ingest = ParquetIngest::new(timestamp);
if let Some(prefix) = &common.channel_prefix {
ingest = ingest.channel_prefix(prefix);
}
for pair in common.tag_columns.chunks(2) {
ingest = ingest.tag_column(&pair[0], &pair[1]);
}
for pair in common.file_tags.chunks(2) {
ingest = ingest.additional_file_tag(&pair[0], &pair[1]);
}
for col in &common.exclude_columns {
ingest = ingest.exclude_column(col);
}
if archive {
ingest = ingest.is_archive(true);
}
let path = common.path.clone();
let job = client
.ingest()
.upload_parquet(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload Parquet '{}'", path.display()))?;
print_result(&job, common.wait, client).await
}
fn build_target(common: &UploadArgs) -> DatasetTarget {
if let Some(rid) = &common.dataset {
return DatasetTarget::Existing(rid.clone());
}
let name = common
.name
.as_ref()
.expect("ArgGroup requires --dataset or --name");
let mut create = DatasetCreate::new(name);
if let Some(d) = &common.description {
create = create.description(d);
}
if !common.labels.is_empty() {
create = create.labels(common.labels.clone());
}
if !common.properties.is_empty() {
let pairs: Vec<(String, String)> = common
.properties
.chunks(2)
.map(|p| (p[0].clone(), p[1].clone()))
.collect();
create = create.properties(pairs);
}
DatasetTarget::New(create)
}
fn build_timestamp(common: &UploadArgs) -> anyhow::Result<Timestamp> {
let col = common.timestamp_column.clone();
match (common.timestamp_type.clone(), common.relative_to) {
(TimestampSpec::Iso8601, None) => Ok(Timestamp::iso8601(col)),
(TimestampSpec::Iso8601, Some(_)) => {
anyhow::bail!("--relative-to cannot be combined with iso8601 timestamps")
}
(TimestampSpec::Epoch(unit), None) => Ok(Timestamp::epoch(col, unit)),
(TimestampSpec::Epoch(unit), Some(offset)) => {
Ok(Timestamp::relative(col, unit).with_offset(offset))
}
}
}
async fn print_result(job: &IngestJob, wait: bool, client: NominalClient) -> anyhow::Result<()> {
println!("Ingest job RID: {}", job.rid());
println!("Status: {}", status_str(job.status()));
if wait {
let terminal = client
.ingest()
.wait_for_ingest_job(job.rid())
.await
.with_context(|| format!("ingest job '{}' did not complete successfully", job.rid()))?;
println!("Final status: {}", status_str(terminal.status()));
}
Ok(())
}
fn status_str(status: &IngestJobStatus) -> String {
match status {
IngestJobStatus::Submitted => "Submitted".into(),
IngestJobStatus::Queued => "Queued".into(),
IngestJobStatus::InProgress => "InProgress".into(),
IngestJobStatus::Completed => "Completed".into(),
IngestJobStatus::Failed => "Failed".into(),
IngestJobStatus::Cancelled => "Cancelled".into(),
IngestJobStatus::Unknown(s) => format!("Unknown({s})"),
}
}