use std::path::PathBuf;
use std::str::FromStr;
use anyhow::Context;
use chrono::{DateTime, Utc};
use clap::{ArgGroup, Args, Subcommand};
use nominal::core::{
AvroStreamIngest, CsvIngest, DataflashIngest, DatasetCreate, DatasetTarget, IngestJob,
JournalJsonIngest, McapIngest, NominalClient, ParquetIngest, TimeUnit, Timestamp, VideoCreate,
VideoIngest, VideoTarget,
};
#[derive(Subcommand)]
pub enum IngestCommands {
Csv(CsvArgs),
Parquet(ParquetArgs),
Mcap(McapArgs),
JournalJson(JournalJsonArgs),
AvroStream(AvroStreamArgs),
ArdupilotDataflash(DataflashArgs),
Video(VideoArgs),
McapVideo(McapVideoArgs),
}
#[derive(Args)]
pub struct CsvArgs {
#[command(flatten)]
common: UploadArgs,
}
#[derive(Args)]
pub struct ParquetArgs {
#[command(flatten)]
common: UploadArgs,
#[arg(long)]
archive: bool,
}
#[derive(Args)]
pub struct McapArgs {
#[command(flatten)]
target: TargetArgs,
#[arg(long = "include-topic", value_name = "TOPIC")]
include_topics: Vec<String>,
#[arg(long = "exclude-topic", value_name = "TOPIC")]
exclude_topics: Vec<String>,
#[arg(
long = "file-tag",
value_names = ["KEY", "VALUE"],
num_args = 2,
action = clap::ArgAction::Append,
)]
file_tags: Vec<String>,
#[arg(long)]
ignore_invalid_topics: bool,
}
#[derive(Args)]
pub struct JournalJsonArgs {
#[command(flatten)]
target: TargetArgs,
#[arg(long, value_name = "NAME")]
channel: Option<String>,
}
#[derive(Args)]
pub struct AvroStreamArgs {
#[command(flatten)]
target: TargetArgs,
}
#[derive(Args)]
pub struct DataflashArgs {
#[command(flatten)]
target: TargetArgs,
#[arg(
long = "file-tag",
value_names = ["KEY", "VALUE"],
num_args = 2,
action = clap::ArgAction::Append,
)]
file_tags: Vec<String>,
}
#[derive(Args)]
pub struct VideoArgs {
#[command(flatten)]
target: VideoTargetArgs,
#[arg(long, value_name = "RFC3339")]
start: DateTime<Utc>,
}
#[derive(Args)]
pub struct McapVideoArgs {
#[command(flatten)]
target: VideoTargetArgs,
#[arg(long)]
topic: String,
}
#[derive(Args)]
#[command(group(
ArgGroup::new("video_target").required(true).args(["video", "name"])
))]
struct VideoTargetArgs {
path: PathBuf,
#[arg(long, value_name = "RID")]
video: 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)]
no_wait: bool,
}
#[derive(Args)]
#[command(group(
ArgGroup::new("target").required(true).args(["dataset", "name"])
))]
struct TargetArgs {
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)]
no_wait: bool,
}
#[derive(Args)]
struct UploadArgs {
#[command(flatten)]
target: TargetArgs,
#[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>,
}
#[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,
IngestCommands::Mcap(args) => handle_mcap(args, client).await,
IngestCommands::JournalJson(args) => handle_journal_json(args, client).await,
IngestCommands::AvroStream(args) => handle_avro_stream(args, client).await,
IngestCommands::ArdupilotDataflash(args) => handle_dataflash(args, client).await,
IngestCommands::Video(args) => handle_video(args, client).await,
IngestCommands::McapVideo(args) => handle_mcap_video(args, client).await,
}
}
async fn handle_csv(args: CsvArgs, client: NominalClient) -> anyhow::Result<()> {
let CsvArgs { common } = args;
let target = build_target(&common.target);
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.target.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_csv(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload CSV '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", common.target.no_wait, client).await
}
async fn handle_parquet(args: ParquetArgs, client: NominalClient) -> anyhow::Result<()> {
let ParquetArgs { common, archive } = args;
let target = build_target(&common.target);
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.target.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_parquet(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload Parquet '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", common.target.no_wait, client).await
}
async fn handle_mcap(args: McapArgs, client: NominalClient) -> anyhow::Result<()> {
let McapArgs {
target: target_args,
include_topics,
exclude_topics,
file_tags,
ignore_invalid_topics,
} = args;
let target = build_target(&target_args);
let mut ingest = McapIngest::new();
for topic in include_topics {
ingest = ingest.include_topic(topic);
}
for topic in exclude_topics {
ingest = ingest.exclude_topic(topic);
}
for pair in file_tags.chunks(2) {
ingest = ingest.additional_file_tag(&pair[0], &pair[1]);
}
if ignore_invalid_topics {
ingest = ingest.ignore_invalid_topics(true);
}
let path = target_args.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_mcap(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload MCAP '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", target_args.no_wait, client).await
}
async fn handle_journal_json(args: JournalJsonArgs, client: NominalClient) -> anyhow::Result<()> {
let JournalJsonArgs {
target: target_args,
channel,
} = args;
let target = build_target(&target_args);
let mut ingest = JournalJsonIngest::new();
if let Some(name) = channel {
ingest = ingest.channel(name);
}
let path = target_args.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_journal_json(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload journal JSON '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", target_args.no_wait, client).await
}
async fn handle_avro_stream(args: AvroStreamArgs, client: NominalClient) -> anyhow::Result<()> {
let AvroStreamArgs {
target: target_args,
} = args;
let target = build_target(&target_args);
let path = target_args.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_avro_stream(&path, target, AvroStreamIngest::new())
.await
.with_context(|| format!("Failed to upload Avro stream '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", target_args.no_wait, client).await
}
async fn handle_dataflash(args: DataflashArgs, client: NominalClient) -> anyhow::Result<()> {
let DataflashArgs {
target: target_args,
file_tags,
} = args;
let target = build_target(&target_args);
let mut ingest = DataflashIngest::new();
for pair in file_tags.chunks(2) {
ingest = ingest.additional_file_tag(&pair[0], &pair[1]);
}
let path = target_args.path.clone();
let (job, dataset_rid) = client
.ingest()
.upload_ardupilot_dataflash(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload DataFlash '{}'", path.display()))?;
print_result(&job, &dataset_rid, "Dataset", target_args.no_wait, client).await
}
async fn handle_video(args: VideoArgs, client: NominalClient) -> anyhow::Result<()> {
let VideoArgs {
target: target_args,
start,
} = args;
let target = build_video_target(&target_args);
let ingest = VideoIngest::starting_at(start);
let path = target_args.path.clone();
let (job, video_rid) = client
.ingest()
.upload_video(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload video '{}'", path.display()))?;
print_result(&job, &video_rid, "Video", target_args.no_wait, client).await
}
async fn handle_mcap_video(args: McapVideoArgs, client: NominalClient) -> anyhow::Result<()> {
let McapVideoArgs {
target: target_args,
topic,
} = args;
let target = build_video_target(&target_args);
let ingest = VideoIngest::mcap_topic(topic);
let path = target_args.path.clone();
let (job, video_rid) = client
.ingest()
.upload_video(&path, target, ingest)
.await
.with_context(|| format!("Failed to upload MCAP video '{}'", path.display()))?;
print_result(&job, &video_rid, "Video", target_args.no_wait, client).await
}
fn build_target(target: &TargetArgs) -> DatasetTarget {
if let Some(rid) = &target.dataset {
return DatasetTarget::Existing(rid.clone());
}
let name = target
.name
.as_ref()
.expect("ArgGroup requires --dataset or --name");
let mut create = DatasetCreate::new(name);
if let Some(d) = &target.description {
create = create.description(d);
}
if !target.labels.is_empty() {
create = create.labels(target.labels.clone());
}
if !target.properties.is_empty() {
let pairs: Vec<(String, String)> = target
.properties
.chunks(2)
.map(|p| (p[0].clone(), p[1].clone()))
.collect();
create = create.properties(pairs);
}
DatasetTarget::New(create)
}
fn build_video_target(target: &VideoTargetArgs) -> VideoTarget {
if let Some(rid) = &target.video {
return VideoTarget::Existing(rid.clone());
}
let name = target
.name
.as_ref()
.expect("ArgGroup requires --video or --name");
let mut create = VideoCreate::new(name);
if let Some(d) = &target.description {
create = create.description(d);
}
if !target.labels.is_empty() {
create = create.labels(target.labels.clone());
}
if !target.properties.is_empty() {
let pairs: Vec<(String, String)> = target
.properties
.chunks(2)
.map(|p| (p[0].clone(), p[1].clone()))
.collect();
create = create.properties(pairs);
}
VideoTarget::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,
resource_rid: &str,
resource_label: &str,
no_wait: bool,
client: NominalClient,
) -> anyhow::Result<()> {
if no_wait {
println!("Ingest job RID: {}", job.rid());
println!("{resource_label} RID: {resource_rid}");
return Ok(());
}
client
.ingest()
.wait_for_ingest_job(job.rid())
.await
.with_context(|| format!("ingest job '{}' did not complete successfully", job.rid()))?;
println!("{resource_label} RID: {resource_rid}");
Ok(())
}