rpo 0.1.0-beta.4

Git contribution analysis: commits, file changes, and per-line authorship over time as polars DataFrames
use std::fs::File;
use std::path::{Path, PathBuf};

use polars::prelude::*;

use super::FrameSink;
use crate::RpoError;

/// Writes a single parquet file containing all snapshots concatenated.
pub struct ParquetSink {
    path: PathBuf,
    accumulated: Option<DataFrame>,
}

impl ParquetSink {
    /// Write every snapshot into a single parquet file at `path`.
    pub fn new(path: impl AsRef<Path>) -> Result<Self, RpoError> {
        Ok(Self {
            path: path.as_ref().to_path_buf(),
            accumulated: None,
        })
    }
}

impl FrameSink for ParquetSink {
    fn write_snapshot(&mut self, frame: DataFrame) -> Result<(), RpoError> {
        match &mut self.accumulated {
            None => {
                self.accumulated = Some(frame);
            }
            Some(acc) => {
                acc.vstack_mut(&frame)?;
            }
        }
        Ok(())
    }

    fn finish(&mut self) -> Result<(), RpoError> {
        let Some(mut df) = self.accumulated.take() else {
            // No snapshots written; create an empty file so the caller sees
            // a coherent artifact.
            let mut empty = DataFrame::empty();
            let file = File::create(&self.path)?;
            ParquetWriter::new(file).finish(&mut empty)?;
            return Ok(());
        };
        df.align_chunks_par();
        let file = File::create(&self.path)?;
        ParquetWriter::new(file).finish(&mut df)?;
        Ok(())
    }
}

/// Writes one parquet file per snapshot, named `<snapshot_sha>.parquet`
/// inside the given directory. Easier for callers who want to process
/// snapshots independently (e.g., Spark/DuckDB per-file).
pub struct ParquetDirSink {
    dir: PathBuf,
}

impl ParquetDirSink {
    /// Write one parquet file per snapshot into `dir`, creating it if
    /// necessary.
    pub fn new(dir: impl AsRef<Path>) -> Result<Self, RpoError> {
        let p = dir.as_ref();
        std::fs::create_dir_all(p)?;
        Ok(Self {
            dir: p.to_path_buf(),
        })
    }
}

impl FrameSink for ParquetDirSink {
    fn write_snapshot(&mut self, mut frame: DataFrame) -> Result<(), RpoError> {
        let sha = frame
            .column("snapshot_sha")?
            .str()?
            .get(0)
            .ok_or_else(|| RpoError::Sink("empty snapshot frame".into()))?
            .to_string();
        let path = self.dir.join(format!("{sha}.parquet"));
        let file = File::create(&path)?;
        ParquetWriter::new(file).finish(&mut frame)?;
        Ok(())
    }

    fn finish(&mut self) -> Result<(), RpoError> {
        Ok(())
    }
}