Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 7 additions & 29 deletions src/catalog.rs
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
//! JSON-backed local dataset catalog.
//! JSON-backed local dataset catalog.
//!
//! Tracks metadata about fetched datasets including ETags, quality status,
//! output paths, and modification timestamps. Supports atomic saves and
//! search operations.

use std::collections::BTreeMap;
use std::fs::{self, File};
use std::fs::File;
use std::io::{BufReader, BufWriter};
use std::path::{Path, PathBuf};

Expand Down Expand Up @@ -103,24 +103,11 @@ impl LocalCatalog {
}

pub fn save_atomic(&self, path: impl AsRef<Path>) -> Result<()> {
let path = path.as_ref();
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}

let tmp_path = tmp_path_for(path);
if tmp_path.exists() {
fs::remove_file(&tmp_path)?;
}

let writer = BufWriter::new(File::create(&tmp_path)?);
serde_json::to_writer_pretty(writer, self)?;

if path.exists() {
fs::remove_file(path)?;
}
fs::rename(tmp_path, path)?;
Ok(())
crate::utils::atomic_write(path, "catalog.json", |tmp_path| {
let writer = BufWriter::new(File::create(tmp_path)?);
serde_json::to_writer_pretty(writer, self)?;
Ok(())
})
}

pub fn upsert(&mut self, dataset: CachedDataset) {
Expand Down Expand Up @@ -254,15 +241,6 @@ pub fn dataset_key(provider: &str, dataset_id: &str) -> String {
format!("{}:{}", provider.trim(), dataset_id.trim())
}

fn tmp_path_for(path: &Path) -> PathBuf {
let mut name = path
.file_name()
.map(|file_name| file_name.to_os_string())
.unwrap_or_else(|| "catalog.json".into());
name.push(".tmp");
path.with_file_name(name)
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
2 changes: 2 additions & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ pub mod quality;
pub mod registry;
pub mod sqlite_catalog;
pub mod traits;
pub mod utils;

// Re-export core types for convenience
pub use catalog::*;
Expand All @@ -48,3 +49,4 @@ pub use quality::*;
pub use registry::*;
pub use sqlite_catalog::*;
pub use traits::*;
pub use utils::*;
43 changes: 10 additions & 33 deletions src/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@
//! [`write_parquet_atomic`] for crash-safe Parquet file writes.

use std::collections::HashMap;
use std::fs::{self, File};
use std::path::{Path, PathBuf};
use std::fs::File;
use std::path::Path;

use polars::prelude::*;

Expand Down Expand Up @@ -130,28 +130,14 @@ pub fn validate_schema(frame: &DataFrame, expected: &[ExpectedColumn]) -> Result
}

pub fn write_parquet_atomic(frame: &DataFrame, output_path: impl AsRef<Path>) -> Result<()> {
let output_path = output_path.as_ref();
if let Some(parent) = output_path.parent() {
fs::create_dir_all(parent)?;
}

let tmp_path = tmp_path_for(output_path);
if tmp_path.exists() {
fs::remove_file(&tmp_path)?;
}

let mut file = File::create(&tmp_path)?;
let mut frame = frame.clone();
ParquetWriter::new(&mut file)
.finish(&mut frame)
.map_err(|error| CoreError::TransformationError(error.to_string()))?;
drop(file);

if output_path.exists() {
fs::remove_file(output_path)?;
}
fs::rename(&tmp_path, output_path)?;
Ok(())
crate::utils::atomic_write(output_path, "output.parquet", |tmp_path| {
let mut file = File::create(tmp_path)?;
let mut frame = frame.clone();
ParquetWriter::new(&mut file)
.finish(&mut frame)
.map_err(|error| CoreError::TransformationError(error.to_string()))?;
Ok(())
})
}

/// Reads a DataFrame from a Parquet file path.
Expand All @@ -162,15 +148,6 @@ pub fn read_parquet(path: impl AsRef<Path>) -> Result<DataFrame> {
.map_err(|e| CoreError::TransformationError(e.to_string()))
}

fn tmp_path_for(output_path: &Path) -> PathBuf {
let mut name = output_path
.file_name()
.map(|file_name| file_name.to_os_string())
.unwrap_or_else(|| "output.parquet".into());
name.push(".tmp");
output_path.with_file_name(name)
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
36 changes: 7 additions & 29 deletions src/quality.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,9 @@
//! incremental delta updates to existing Parquet datasets.

use std::collections::HashSet;
use std::fs::{self, File};
use std::fs::File;
use std::io::BufWriter;
use std::path::{Path, PathBuf};
use std::path::Path;

use polars::prelude::*;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -111,24 +111,11 @@ impl QualityReport {
}

pub fn save_atomic(&self, path: impl AsRef<Path>) -> Result<()> {
let path = path.as_ref();
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}

let tmp_path = tmp_path_for(path);
if tmp_path.exists() {
fs::remove_file(&tmp_path)?;
}

let writer = BufWriter::new(File::create(&tmp_path)?);
serde_json::to_writer_pretty(writer, self)?;

if path.exists() {
fs::remove_file(path)?;
}
fs::rename(tmp_path, path)?;
Ok(())
crate::utils::atomic_write(path, "quality-report.json", |tmp_path| {
let writer = BufWriter::new(File::create(tmp_path)?);
serde_json::to_writer_pretty(writer, self)?;
Ok(())
})
}
}

Expand Down Expand Up @@ -260,15 +247,6 @@ fn unique_stringified_count(series: &Column) -> Result<usize> {
Ok(values.len())
}

fn tmp_path_for(path: &Path) -> PathBuf {
let mut name = path
.file_name()
.map(|file_name| file_name.to_os_string())
.unwrap_or_else(|| "quality-report.json".into());
name.push(".tmp");
path.with_file_name(name)
}

pub fn provider_payload_assertions() -> Vec<QualityAssertion> {
vec![
QualityAssertion::non_null("provider"),
Expand Down
36 changes: 36 additions & 0 deletions src/utils.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
use std::fs;
use std::path::{Path, PathBuf};

use crate::error::Result;

pub fn tmp_path_for(path: &Path, fallback_name: &str) -> PathBuf {
let mut name = path
.file_name()
.map(|file_name| file_name.to_os_string())
.unwrap_or_else(|| fallback_name.into());
name.push(".tmp");
path.with_file_name(name)
}

pub fn atomic_write<F>(path: impl AsRef<Path>, fallback_name: &str, write_fn: F) -> Result<()>
where
F: FnOnce(&Path) -> Result<()>,
{
let path = path.as_ref();
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}

let tmp_path = tmp_path_for(path, fallback_name);
if tmp_path.exists() {
fs::remove_file(&tmp_path)?;
}

write_fn(&tmp_path)?;

if path.exists() {
fs::remove_file(path)?;
}
fs::rename(tmp_path, path)?;
Ok(())
}
Loading