use crate::Result;
use crate::error::CoreError;
use apache_avro::Reader as AvroReader;
use apache_avro::from_value;
use serde::Deserialize;
use std::collections::HashMap;
use std::io::Cursor;
#[derive(Debug, Clone, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct RollbackMetadata {
pub start_rollback_time: String,
pub commits_rollback: Vec<String>,
}
impl RollbackMetadata {
pub fn from_avro_bytes(bytes: &[u8]) -> Result<Self> {
let reader = AvroReader::new(Cursor::new(bytes)).map_err(|e| {
CoreError::CommitMetadata(format!("Failed to create Avro reader for rollback: {e}"))
})?;
let mut records = reader;
let value = records
.next()
.ok_or_else(|| {
CoreError::CommitMetadata("Rollback metadata contains no records".to_string())
})?
.map_err(|e| {
CoreError::CommitMetadata(format!("Failed to read rollback record: {e}"))
})?;
from_value::<Self>(&value).map_err(|e| {
CoreError::CommitMetadata(format!("Failed to deserialize rollback metadata: {e}"))
})
}
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct RestoreMetadata {
pub start_restore_time: String,
pub hoodie_restore_metadata: HashMap<String, Vec<RollbackMetadata>>,
}
impl RestoreMetadata {
pub fn from_avro_bytes(bytes: &[u8]) -> Result<Self> {
let reader = AvroReader::new(Cursor::new(bytes)).map_err(|e| {
CoreError::CommitMetadata(format!("Failed to create Avro reader for restore: {e}"))
})?;
let mut records = reader;
let value = records
.next()
.ok_or_else(|| {
CoreError::CommitMetadata("Restore metadata contains no records".to_string())
})?
.map_err(|e| {
CoreError::CommitMetadata(format!("Failed to read restore record: {e}"))
})?;
from_value::<Self>(&value).map_err(|e| {
CoreError::CommitMetadata(format!("Failed to deserialize restore metadata: {e}"))
})
}
pub fn commits_rolled_back(&self) -> Vec<String> {
self.hoodie_restore_metadata
.values()
.flatten()
.flat_map(|rollback| rollback.commits_rollback.iter().cloned())
.collect()
}
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct InstantInfo {
pub commit_time: String,
pub action: String,
}
#[derive(Debug, Clone, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "camelCase")]
pub struct RollbackPlan {
pub instant_to_rollback: Option<InstantInfo>,
}
impl RollbackPlan {
pub fn from_avro_bytes(bytes: &[u8]) -> Result<Self> {
let reader = AvroReader::new(Cursor::new(bytes)).map_err(|e| {
CoreError::CommitMetadata(format!(
"Failed to create Avro reader for rollback plan: {e}"
))
})?;
let mut records = reader;
let value = records
.next()
.ok_or_else(|| {
CoreError::CommitMetadata("Rollback plan contains no records".to_string())
})?
.map_err(|e| {
CoreError::CommitMetadata(format!("Failed to read rollback plan record: {e}"))
})?;
from_value::<Self>(&value).map_err(|e| {
CoreError::CommitMetadata(format!("Failed to deserialize rollback plan: {e}"))
})
}
pub fn commit_rolled_back(&self) -> Option<&str> {
self.instant_to_rollback
.as_ref()
.map(|i| i.commit_time.as_str())
}
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
use apache_avro::types::Value;
use apache_avro::{Schema, Writer};
pub(crate) fn inlined_schema_text() -> String {
inline_named_types("HoodieRollbackMetadata.avsc", &["HoodieInstantInfo.avsc"])
}
pub(crate) fn inline_named_types(root: &str, deps: &[&str]) -> String {
let dir =
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../test/data/avro_schemas");
let read = |name: &str| {
std::fs::read_to_string(dir.join(name))
.unwrap_or_else(|e| panic!("{name} must be present: {e}"))
};
let mut text = read(root);
for dep in deps {
let type_name = dep.trim_end_matches(".avsc");
let reference = format!("\"{type_name}\"");
assert!(
text.contains(&reference),
"{root} must reference {type_name}, or this splice is stale"
);
text = text.replacen(&reference, read(dep).trim(), 1);
}
text
}
pub(crate) fn container_bytes(start: &str, rolled_back: &[&str]) -> Vec<u8> {
let schema = Schema::parse_str(&inlined_schema_text())
.expect("Hudi's schema, with its dependency inlined, must parse");
let mut writer = Writer::new(&schema, Vec::new());
let record = Value::Record(vec![
("startRollbackTime".into(), Value::String(start.into())),
("timeTakenInMillis".into(), Value::Long(42)),
("totalFilesDeleted".into(), Value::Int(3)),
(
"commitsRollback".into(),
Value::Array(
rolled_back
.iter()
.map(|c| Value::String((*c).into()))
.collect(),
),
),
("partitionMetadata".into(), Value::Map(Default::default())),
("version".into(), Value::Union(0, Box::new(Value::Int(1)))),
("instantsRollback".into(), Value::Array(vec![])),
]);
writer.append(record).expect("append");
writer.into_inner().expect("container bytes")
}
#[test]
fn a_rollback_yields_the_commits_it_rolled_back() -> Result<()> {
let bytes = container_bytes(
"20250103000000000",
&["20250101000000000", "20250102000000000"],
);
let parsed = RollbackMetadata::from_avro_bytes(&bytes)?;
assert_eq!(parsed.start_rollback_time, "20250103000000000");
assert_eq!(
parsed.commits_rollback,
vec!["20250101000000000", "20250102000000000"],
"every rolled-back commit must come back, in order"
);
Ok(())
}
#[test]
fn a_rollback_of_nothing_is_not_an_error() -> Result<()> {
let parsed = RollbackMetadata::from_avro_bytes(&container_bytes("20250103000000000", &[]))?;
assert!(parsed.commits_rollback.is_empty());
Ok(())
}
#[test]
fn garbage_is_an_error_not_an_empty_result() {
let err = RollbackMetadata::from_avro_bytes(b"not avro at all").unwrap_err();
assert!(
err.to_string().contains("rollback"),
"the error must name what failed to read, got: {err}"
);
}
#[test]
fn a_restore_yields_every_commit_its_rollbacks_rolled_back() -> Result<()> {
let text = inline_named_types(
"HoodieRestoreMetadata.avsc",
&["HoodieRollbackMetadata.avsc", "HoodieInstantInfo.avsc"],
);
let schema = Schema::parse_str(&text).expect("the restore schema must parse once inlined");
let rollback = |start: &str, commit: &str| {
Value::Record(vec![
("startRollbackTime".into(), Value::String(start.into())),
("timeTakenInMillis".into(), Value::Long(1)),
("totalFilesDeleted".into(), Value::Int(0)),
(
"commitsRollback".into(),
Value::Array(vec![Value::String(commit.into())]),
),
("partitionMetadata".into(), Value::Map(Default::default())),
("version".into(), Value::Union(0, Box::new(Value::Int(1)))),
("instantsRollback".into(), Value::Array(vec![])),
])
};
let mut nested = std::collections::HashMap::new();
nested.insert(
"20250101000000000".to_string(),
Value::Array(vec![Value::Union(
1,
Box::new(rollback("20250104000000000", "20250101000000000")),
)]),
);
nested.insert(
"20250102000000000".to_string(),
Value::Array(vec![Value::Union(
1,
Box::new(rollback("20250104000000000", "20250102000000000")),
)]),
);
let mut writer = Writer::new(&schema, Vec::new());
writer
.append(Value::Record(vec![
(
"startRestoreTime".into(),
Value::String("20250104000000000".into()),
),
("timeTakenInMillis".into(), Value::Long(9)),
("instantsToRollback".into(), Value::Array(vec![])),
("hoodieRestoreMetadata".into(), Value::Map(nested)),
("version".into(), Value::Union(0, Box::new(Value::Int(1)))),
("restoreInstantInfo".into(), Value::Array(vec![])),
]))
.expect("append restore");
let bytes = writer.into_inner().expect("container bytes");
let parsed = RestoreMetadata::from_avro_bytes(&bytes)?;
let mut got = parsed.commits_rolled_back();
got.sort();
assert_eq!(
got,
vec!["20250101000000000", "20250102000000000"],
"a restore must yield every commit across all of its rollbacks"
);
Ok(())
}
#[test]
fn a_rollback_plan_names_the_instant_it_rolls_back() -> Result<()> {
let text = inline_named_types("HoodieRollbackPlan.avsc", &["HoodieInstantInfo.avsc"]);
let schema = Schema::parse_str(&text).expect("the plan schema must parse once inlined");
let mut writer = Writer::new(&schema, Vec::new());
writer
.append(Value::Record(vec![
(
"instantToRollback".into(),
Value::Union(
1,
Box::new(Value::Record(vec![
(
"commitTime".into(),
Value::String("20250101000000000".into()),
),
("action".into(), Value::String("commit".into())),
])),
),
),
(
"RollbackRequests".into(),
Value::Union(0, Box::new(Value::Null)),
),
("version".into(), Value::Union(0, Box::new(Value::Int(1)))),
]))
.expect("append plan");
let bytes = writer.into_inner().expect("container bytes");
let parsed = RollbackPlan::from_avro_bytes(&bytes)?;
assert_eq!(parsed.commit_rolled_back(), Some("20250101000000000"));
Ok(())
}
#[test]
fn a_plan_without_an_instant_yields_nothing() -> Result<()> {
let text = inline_named_types("HoodieRollbackPlan.avsc", &["HoodieInstantInfo.avsc"]);
let schema = Schema::parse_str(&text).expect("the plan schema must parse once inlined");
let mut writer = Writer::new(&schema, Vec::new());
writer
.append(Value::Record(vec![
(
"instantToRollback".into(),
Value::Union(0, Box::new(Value::Null)),
),
(
"RollbackRequests".into(),
Value::Union(0, Box::new(Value::Null)),
),
("version".into(), Value::Union(0, Box::new(Value::Int(1)))),
]))
.expect("append plan");
let parsed = RollbackPlan::from_avro_bytes(&writer.into_inner().expect("bytes"))?;
assert_eq!(parsed.commit_rolled_back(), None);
Ok(())
}
}