timeseries-table-format 0.4.0

Append-only time-series table format with gap/overlap tracking
Documentation
//! Append profiling binary: create a fresh benchmark table and append a segment.

use std::path::{Path, PathBuf};

use clap::Parser;

use timeseries_table_format::{
    metadata::table_metadata::{IndexKind, IndexSpec, TableMeta, TimeBucket},
    storage::TableLocation,
    table::{
        TimeSeriesTable,
        append_report::{AppendReport, AppendStep},
    },
};

#[derive(Debug, Parser)]
struct Args {
    #[arg(long)]
    table: PathBuf,

    #[arg(long)]
    parquet: PathBuf,

    /// Timestamp column name for create + append
    #[arg(long = "time-column")]
    time_column: String,

    /// Time bucket (e.g. 1s, 1m, 1h, 1d)
    #[arg(long)]
    bucket: String,

    /// Optional IANA timezone string
    #[arg(long)]
    timezone: Option<String>,

    /// Repeatable entity column names
    #[arg(long = "entity")]
    entity: Vec<String>,

    /// Optional CSV output path
    #[arg(long = "csv-out")]
    csv_out: Option<PathBuf>,
}

fn format_fields(step: &AppendStep) -> String {
    if step.fields.is_empty() {
        return String::new();
    }
    step.fields
        .iter()
        .map(|(k, v)| format!("{k}={v}"))
        .collect::<Vec<_>>()
        .join(";")
}

fn get_context_value<'a>(report: &'a AppendReport, key: &str) -> &'a str {
    report
        .context
        .iter()
        .find(|(k, _)| k == key)
        .map(|(_, v)| v.as_str())
        .unwrap_or("")
}

fn csv_escape(raw: &str) -> String {
    if raw.contains(',') || raw.contains('"') || raw.contains('\n') {
        let escaped = raw.replace('"', "\"\"");
        format!("\"{escaped}\"")
    } else {
        raw.to_string()
    }
}

fn write_csv(report: &AppendReport, out_path: &Path) -> std::io::Result<()> {
    let mut out = String::new();
    out.push_str(
        "relative_path,time_column,file_size_bytes,step,elapsed_ms,percent,total_ms,fields\n",
    );

    let rel = get_context_value(report, "relative_path");
    let time_column = get_context_value(report, "time_column");
    let file_size_bytes = get_context_value(report, "file_size_bytes");

    for step in &report.steps {
        let percent = if report.total_ms == 0 {
            0.0
        } else {
            (step.elapsed_ms as f64) * 100.0 / (report.total_ms as f64)
        };

        let fields = format_fields(step);
        let line = format!(
            "{},{},{},{},{},{:.2},{},{}\n",
            csv_escape(rel),
            csv_escape(time_column),
            csv_escape(file_size_bytes),
            csv_escape(&step.name),
            step.elapsed_ms,
            percent,
            report.total_ms,
            csv_escape(&fields),
        );
        out.push_str(&line);
    }

    std::fs::write(out_path, out)
}

fn print_report(report: &AppendReport) {
    let rel = get_context_value(report, "relative_path");
    let file_size_bytes = get_context_value(report, "file_size_bytes");
    println!(
        "Append profile: {rel} (file size: {file_size_bytes} bytes, total_ms: {})",
        report.total_ms
    );

    for step in &report.steps {
        let percent = if report.total_ms == 0 {
            0.0
        } else {
            (step.elapsed_ms as f64) * 100.0 / (report.total_ms as f64)
        };

        let fields = format_fields(step);
        if fields.is_empty() {
            println!(
                "{:<24} {:>8} ms  {:>6.2}%",
                step.name, step.elapsed_ms, percent
            );
        } else {
            println!(
                "{:<24} {:>8} ms  {:>6.2}%  {}",
                step.name, step.elapsed_ms, percent, fields
            );
        }
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let args = Args::parse();

    let benchmark_root = args.table.join("benchmark");
    if tokio::fs::metadata(&benchmark_root).await.is_ok() {
        tokio::fs::remove_dir_all(&benchmark_root).await?;
    }
    tokio::fs::create_dir_all(&benchmark_root).await?;

    let location = TableLocation::parse(benchmark_root.to_string_lossy().as_ref())?;

    let bucket = args.bucket.parse::<TimeBucket>()?;
    let index = IndexSpec {
        column: args.time_column.clone(),
        entity_columns: args.entity.clone(),
        kind: IndexKind::Timestamp {
            bucket,
            timezone: args.timezone.clone(),
        },
    };
    let meta = TableMeta::new_time_series(index);
    TimeSeriesTable::create(location.clone(), meta).await?;

    let mut table = TimeSeriesTable::open(location).await?;

    let (_version, _relative_path, report) = table
        .append_parquet_from_path_with_report(&args.parquet)
        .await?;

    print_report(&report);

    if let Some(csv_out) = args.csv_out.as_deref() {
        write_csv(&report, csv_out)?;
        println!("Wrote CSV report to {}", csv_out.display());
    }

    Ok(())
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn csv_uses_relative_path_as_segment_identity() {
        let tmp = tempfile::tempdir().expect("create temp dir");
        let output = tmp.path().join("report.csv");
        let report = AppendReport {
            context: vec![
                ("relative_path".to_string(), "data/seg.parquet".to_string()),
                ("time_column".to_string(), "ts".to_string()),
                ("file_size_bytes".to_string(), "42".to_string()),
            ],
            steps: vec![AppendStep {
                name: "segment_meta".to_string(),
                elapsed_ms: 1,
                fields: Vec::new(),
            }],
            total_ms: 2,
        };

        write_csv(&report, &output).expect("write CSV");
        let csv = std::fs::read_to_string(output).expect("read CSV");

        assert_eq!(
            csv,
            "relative_path,time_column,file_size_bytes,step,elapsed_ms,percent,total_ms,fields\n\
             data/seg.parquet,ts,42,segment_meta,1,50.00,2,\n"
        );
    }
}