1use std::collections::HashMap;
22use std::sync::Arc;
23
24use chrono::{DateTime, Utc};
25use serde::{Deserialize, Serialize};
26use typed_builder::TypedBuilder;
27
28use crate::error::{Result, invalid_data, timestamp_ms_to_utc};
29use crate::spec::{SchemaId, SchemaRef, TableMetadata};
30
31pub const MAIN_BRANCH: &str = "main";
33
34pub type SnapshotRef = Arc<Snapshot>;
36#[derive(Debug, Default, Serialize, Deserialize, PartialEq, Eq, Clone)]
37#[serde(rename_all = "lowercase")]
38pub enum Operation {
40 #[default]
42 Append,
43 Replace,
46 Overwrite,
48 Delete,
50}
51
52impl Operation {
53 pub fn as_str(&self) -> &str {
55 match self {
56 Operation::Append => "append",
57 Operation::Replace => "replace",
58 Operation::Overwrite => "overwrite",
59 Operation::Delete => "delete",
60 }
61 }
62}
63
64#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
65pub struct Summary {
67 pub operation: Operation,
69 #[serde(flatten)]
71 pub additional_properties: HashMap<String, String>,
72}
73
74#[derive(Debug, PartialEq, Eq, Clone)]
75pub struct SnapshotRowRange {
77 pub first_row_id: u64,
79 pub added_rows: u64,
81}
82
83#[derive(Debug, PartialEq, Eq, Clone, TypedBuilder)]
84#[builder(field_defaults(setter(prefix = "with_")))]
85pub struct Snapshot {
87 pub(crate) snapshot_id: i64,
89 #[builder(default = None)]
92 pub(crate) parent_snapshot_id: Option<i64>,
93 pub(crate) sequence_number: i64,
96 pub(crate) timestamp_ms: i64,
99 #[builder(setter(into))]
103 pub(crate) manifest_list: String,
104 pub(crate) summary: Summary,
106 #[builder(setter(strip_option(fallback = schema_id_opt)), default = None)]
108 pub(crate) schema_id: Option<SchemaId>,
109 #[builder(default)]
111 pub(crate) encryption_key_id: Option<String>,
112 #[builder(default, setter(!strip_option, transform = |first_row_id: u64, added_rows: u64| Some(SnapshotRowRange { first_row_id, added_rows })))]
115 pub(crate) row_range: Option<SnapshotRowRange>,
120}
121
122impl Snapshot {
123 #[inline]
125 pub fn snapshot_id(&self) -> i64 {
126 self.snapshot_id
127 }
128
129 #[inline]
131 pub fn parent_snapshot_id(&self) -> Option<i64> {
132 self.parent_snapshot_id
133 }
134
135 #[inline]
137 pub fn sequence_number(&self) -> i64 {
138 self.sequence_number
139 }
140 #[inline]
142 pub fn manifest_list(&self) -> &str {
143 &self.manifest_list
144 }
145
146 #[inline]
148 pub fn summary(&self) -> &Summary {
149 &self.summary
150 }
151 #[inline]
153 pub fn timestamp(&self) -> Result<DateTime<Utc>> {
154 timestamp_ms_to_utc(self.timestamp_ms)
155 }
156
157 #[inline]
159 pub fn timestamp_ms(&self) -> i64 {
160 self.timestamp_ms
161 }
162
163 #[inline]
165 pub fn schema_id(&self) -> Option<SchemaId> {
166 self.schema_id
167 }
168
169 pub fn schema(&self, table_metadata: &TableMetadata) -> Result<SchemaRef> {
171 Ok(match self.schema_id() {
172 Some(schema_id) => table_metadata
173 .schema_by_id(schema_id)
174 .ok_or_else(|| invalid_data!("Schema with id {schema_id} not found"))?
175 .clone(),
176 None => table_metadata.current_schema().clone(),
177 })
178 }
179
180 #[cfg(test)]
182 pub(crate) fn parent_snapshot(&self, table_metadata: &TableMetadata) -> Option<SnapshotRef> {
183 match self.parent_snapshot_id {
184 Some(id) => table_metadata.snapshot_by_id(id).cloned(),
185 None => None,
186 }
187 }
188
189 pub fn first_row_id(&self) -> Option<u64> {
196 self.row_range.as_ref().map(|r| r.first_row_id)
197 }
198
199 pub fn added_rows_count(&self) -> Option<u64> {
204 self.row_range.as_ref().map(|r| r.added_rows)
205 }
206
207 pub fn row_range(&self) -> Option<(u64, u64)> {
210 self.row_range
211 .as_ref()
212 .map(|r| (r.first_row_id, r.added_rows))
213 }
214
215 pub fn encryption_key_id(&self) -> Option<&str> {
217 self.encryption_key_id.as_deref()
218 }
219}
220
221pub(super) mod _serde {
222 use std::collections::HashMap;
227
228 use serde::{Deserialize, Serialize};
229
230 use super::{Operation, Snapshot, Summary};
231 use crate::Error;
232 use crate::error::invalid_data;
233 use crate::spec::SchemaId;
234 use crate::spec::snapshot::SnapshotRowRange;
235
236 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
237 #[serde(rename_all = "kebab-case")]
238 pub(crate) struct SnapshotV3 {
240 pub snapshot_id: i64,
241 #[serde(skip_serializing_if = "Option::is_none")]
242 pub parent_snapshot_id: Option<i64>,
243 pub sequence_number: i64,
244 pub timestamp_ms: i64,
245 pub manifest_list: String,
246 pub summary: Summary,
247 #[serde(skip_serializing_if = "Option::is_none")]
248 pub schema_id: Option<SchemaId>,
249 #[serde(skip_serializing_if = "Option::is_none")]
250 pub first_row_id: Option<u64>,
251 #[serde(skip_serializing_if = "Option::is_none")]
252 pub added_rows: Option<u64>,
253 #[serde(skip_serializing_if = "Option::is_none")]
254 pub key_id: Option<String>,
255 }
256
257 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
258 #[serde(rename_all = "kebab-case")]
259 pub(crate) struct SnapshotV2 {
261 pub snapshot_id: i64,
262 #[serde(skip_serializing_if = "Option::is_none")]
263 pub parent_snapshot_id: Option<i64>,
264 #[serde(default)]
265 pub sequence_number: i64,
266 pub timestamp_ms: i64,
267 pub manifest_list: String,
268 pub summary: Summary,
269 #[serde(skip_serializing_if = "Option::is_none")]
270 pub schema_id: Option<SchemaId>,
271 }
272
273 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
274 #[serde(rename_all = "kebab-case")]
275 pub(crate) struct SnapshotV1 {
277 pub snapshot_id: i64,
278 #[serde(skip_serializing_if = "Option::is_none")]
279 pub parent_snapshot_id: Option<i64>,
280 pub timestamp_ms: i64,
281 #[serde(skip_serializing_if = "Option::is_none")]
282 pub manifest_list: Option<String>,
283 #[serde(skip_serializing_if = "Option::is_none")]
284 pub manifests: Option<Vec<String>>,
285 #[serde(skip_serializing_if = "Option::is_none")]
286 pub summary: Option<Summary>,
287 #[serde(skip_serializing_if = "Option::is_none")]
288 pub schema_id: Option<SchemaId>,
289 }
290
291 impl From<SnapshotV3> for Snapshot {
292 fn from(s: SnapshotV3) -> Self {
293 Snapshot {
294 snapshot_id: s.snapshot_id,
295 parent_snapshot_id: s.parent_snapshot_id,
296 sequence_number: s.sequence_number,
297 timestamp_ms: s.timestamp_ms,
298 manifest_list: s.manifest_list,
299 summary: s.summary,
300 schema_id: s.schema_id,
301 encryption_key_id: s.key_id,
302 row_range: match (s.first_row_id, s.added_rows) {
303 (Some(first_row_id), Some(added_rows)) => Some(SnapshotRowRange {
304 first_row_id,
305 added_rows,
306 }),
307 _ => None,
308 },
309 }
310 }
311 }
312
313 impl TryFrom<Snapshot> for SnapshotV3 {
314 type Error = Error;
315
316 fn try_from(s: Snapshot) -> Result<Self, Self::Error> {
317 let (first_row_id, added_rows) = match s.row_range {
318 Some(row_range) => (Some(row_range.first_row_id), Some(row_range.added_rows)),
319 None => (None, None),
320 };
321
322 Ok(SnapshotV3 {
323 snapshot_id: s.snapshot_id,
324 parent_snapshot_id: s.parent_snapshot_id,
325 sequence_number: s.sequence_number,
326 timestamp_ms: s.timestamp_ms,
327 manifest_list: s.manifest_list,
328 summary: s.summary,
329 schema_id: s.schema_id,
330 first_row_id,
331 added_rows,
332 key_id: s.encryption_key_id,
333 })
334 }
335 }
336
337 impl From<SnapshotV2> for Snapshot {
338 fn from(v2: SnapshotV2) -> Self {
339 Snapshot {
340 snapshot_id: v2.snapshot_id,
341 parent_snapshot_id: v2.parent_snapshot_id,
342 sequence_number: v2.sequence_number,
343 timestamp_ms: v2.timestamp_ms,
344 manifest_list: v2.manifest_list,
345 summary: v2.summary,
346 schema_id: v2.schema_id,
347 encryption_key_id: None,
348 row_range: None,
349 }
350 }
351 }
352
353 impl From<Snapshot> for SnapshotV2 {
354 fn from(v2: Snapshot) -> Self {
355 SnapshotV2 {
356 snapshot_id: v2.snapshot_id,
357 parent_snapshot_id: v2.parent_snapshot_id,
358 sequence_number: v2.sequence_number,
359 timestamp_ms: v2.timestamp_ms,
360 manifest_list: v2.manifest_list,
361 summary: v2.summary,
362 schema_id: v2.schema_id,
363 }
364 }
365 }
366
367 impl TryFrom<SnapshotV1> for Snapshot {
368 type Error = Error;
369
370 fn try_from(v1: SnapshotV1) -> Result<Self, Self::Error> {
371 Ok(Snapshot {
372 snapshot_id: v1.snapshot_id,
373 parent_snapshot_id: v1.parent_snapshot_id,
374 sequence_number: 0,
375 timestamp_ms: v1.timestamp_ms,
376 manifest_list: match (v1.manifest_list, v1.manifests) {
377 (Some(file), None) => file,
378 (Some(_), Some(_)) => {
379 return Err(invalid_data!(
380 "Invalid v1 snapshot, when manifest list provided, manifest files should be omitted"
381 ));
382 }
383 (None, _) => {
384 return Err(invalid_data!(
385 "Unsupported v1 snapshot, only manifest list is supported"
386 ));
387 }
388 },
389 summary: v1.summary.unwrap_or(Summary {
390 operation: Operation::default(),
391 additional_properties: HashMap::new(),
392 }),
393 schema_id: v1.schema_id,
394 encryption_key_id: None,
395 row_range: None,
396 })
397 }
398 }
399
400 impl From<Snapshot> for SnapshotV1 {
401 fn from(v2: Snapshot) -> Self {
402 SnapshotV1 {
403 snapshot_id: v2.snapshot_id,
404 parent_snapshot_id: v2.parent_snapshot_id,
405 timestamp_ms: v2.timestamp_ms,
406 manifest_list: Some(v2.manifest_list),
407 summary: Some(v2.summary),
408 schema_id: v2.schema_id,
409 manifests: None,
410 }
411 }
412 }
413}
414
415#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
416#[serde(rename_all = "kebab-case")]
417pub struct SnapshotReference {
419 pub snapshot_id: i64,
421 #[serde(flatten)]
422 pub retention: SnapshotRetention,
424}
425
426impl SnapshotReference {
427 pub fn is_branch(&self) -> bool {
429 matches!(self.retention, SnapshotRetention::Branch { .. })
430 }
431}
432
433impl SnapshotReference {
434 pub fn new(snapshot_id: i64, retention: SnapshotRetention) -> Self {
436 SnapshotReference {
437 snapshot_id,
438 retention,
439 }
440 }
441}
442
443#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
444#[serde(rename_all = "lowercase", tag = "type")]
445pub enum SnapshotRetention {
447 #[serde(rename_all = "kebab-case")]
448 Branch {
451 #[serde(skip_serializing_if = "Option::is_none")]
454 min_snapshots_to_keep: Option<i32>,
455 #[serde(skip_serializing_if = "Option::is_none")]
458 max_snapshot_age_ms: Option<i64>,
459 #[serde(skip_serializing_if = "Option::is_none")]
462 max_ref_age_ms: Option<i64>,
463 },
464 #[serde(rename_all = "kebab-case")]
465 Tag {
467 #[serde(skip_serializing_if = "Option::is_none")]
470 max_ref_age_ms: Option<i64>,
471 },
472}
473
474impl SnapshotRetention {
475 pub fn branch(
477 min_snapshots_to_keep: Option<i32>,
478 max_snapshot_age_ms: Option<i64>,
479 max_ref_age_ms: Option<i64>,
480 ) -> Self {
481 SnapshotRetention::Branch {
482 min_snapshots_to_keep,
483 max_snapshot_age_ms,
484 max_ref_age_ms,
485 }
486 }
487}
488
489#[cfg(test)]
490mod tests {
491 use std::collections::HashMap;
492
493 use chrono::{TimeZone, Utc};
494
495 use crate::spec::TableMetadata;
496 use crate::spec::snapshot::_serde::SnapshotV1;
497 use crate::spec::snapshot::{Operation, Snapshot, Summary};
498
499 #[test]
500 fn schema() {
501 let record = r#"
502 {
503 "snapshot-id": 3051729675574597004,
504 "timestamp-ms": 1515100955770,
505 "summary": {
506 "operation": "append"
507 },
508 "manifest-list": "s3://b/wh/.../s1.avro",
509 "schema-id": 0
510 }
511 "#;
512
513 let result: Snapshot = serde_json::from_str::<SnapshotV1>(record)
514 .unwrap()
515 .try_into()
516 .unwrap();
517 assert_eq!(3051729675574597004, result.snapshot_id());
518 assert_eq!(
519 Utc.timestamp_millis_opt(1515100955770).unwrap(),
520 result.timestamp().unwrap()
521 );
522 assert_eq!(1515100955770, result.timestamp_ms());
523 assert_eq!(
524 Summary {
525 operation: Operation::Append,
526 additional_properties: HashMap::new()
527 },
528 *result.summary()
529 );
530 assert_eq!("s3://b/wh/.../s1.avro".to_string(), *result.manifest_list());
531 }
532
533 #[test]
534 fn test_snapshot_v1_to_v2_projection() {
535 use crate::spec::snapshot::_serde::SnapshotV1;
536
537 let v1_snapshot = SnapshotV1 {
539 snapshot_id: 1234567890,
540 parent_snapshot_id: Some(987654321),
541 timestamp_ms: 1515100955770,
542 manifest_list: Some("s3://bucket/manifest-list.avro".to_string()),
543 manifests: None, summary: Some(Summary {
545 operation: Operation::Append,
546 additional_properties: HashMap::from([
547 ("added-files".to_string(), "5".to_string()),
548 ("added-records".to_string(), "100".to_string()),
549 ]),
550 }),
551 schema_id: Some(1),
552 };
553
554 let v2_snapshot: Snapshot = v1_snapshot.try_into().unwrap();
556
557 assert_eq!(
559 v2_snapshot.sequence_number(),
560 0,
561 "V1 snapshot sequence_number should default to 0"
562 );
563
564 assert_eq!(v2_snapshot.snapshot_id(), 1234567890);
566 assert_eq!(v2_snapshot.parent_snapshot_id(), Some(987654321));
567 assert_eq!(v2_snapshot.timestamp_ms(), 1515100955770);
568 assert_eq!(
569 v2_snapshot.manifest_list(),
570 "s3://bucket/manifest-list.avro"
571 );
572 assert_eq!(v2_snapshot.schema_id(), Some(1));
573 assert_eq!(v2_snapshot.summary().operation, Operation::Append);
574 assert_eq!(
575 v2_snapshot
576 .summary()
577 .additional_properties
578 .get("added-files"),
579 Some(&"5".to_string())
580 );
581 }
582
583 #[test]
584 fn test_v1_snapshot_with_manifest_list_and_manifests() {
585 {
586 let metadata = r#"
587 {
588 "format-version": 1,
589 "table-uuid": "d20125c8-7284-442c-9aea-15fee620737c",
590 "location": "s3://bucket/test/location",
591 "last-updated-ms": 1700000000000,
592 "last-column-id": 1,
593 "schema": {
594 "type": "struct",
595 "fields": [
596 {"id": 1, "name": "x", "required": true, "type": "long"}
597 ]
598 },
599 "partition-spec": [],
600 "properties": {},
601 "current-snapshot-id": 111111111,
602 "snapshots": [
603 {
604 "snapshot-id": 111111111,
605 "timestamp-ms": 1600000000000,
606 "summary": {"operation": "append"},
607 "manifest-list": "s3://bucket/metadata/snap-123.avro",
608 "manifests": ["s3://bucket/metadata/manifest-1.avro"]
609 }
610 ]
611 }
612 "#;
613
614 let result_both_manifest_list_and_manifest_set =
615 serde_json::from_str::<TableMetadata>(metadata);
616 assert!(result_both_manifest_list_and_manifest_set.is_err());
617 assert_eq!(
618 result_both_manifest_list_and_manifest_set
619 .unwrap_err()
620 .to_string(),
621 "DataInvalid => Invalid v1 snapshot, when manifest list provided, manifest files should be omitted"
622 )
623 }
624
625 {
626 let metadata = r#"
627 {
628 "format-version": 1,
629 "table-uuid": "d20125c8-7284-442c-9aea-15fee620737c",
630 "location": "s3://bucket/test/location",
631 "last-updated-ms": 1700000000000,
632 "last-column-id": 1,
633 "schema": {
634 "type": "struct",
635 "fields": [
636 {"id": 1, "name": "x", "required": true, "type": "long"}
637 ]
638 },
639 "partition-spec": [],
640 "properties": {},
641 "current-snapshot-id": 111111111,
642 "snapshots": [
643 {
644 "snapshot-id": 111111111,
645 "timestamp-ms": 1600000000000,
646 "summary": {"operation": "append"},
647 "manifests": ["s3://bucket/metadata/manifest-1.avro"]
648 }
649 ]
650 }
651 "#;
652 let result_missing_manifest_list = serde_json::from_str::<TableMetadata>(metadata);
653 assert!(result_missing_manifest_list.is_err());
654 assert_eq!(
655 result_missing_manifest_list.unwrap_err().to_string(),
656 "DataInvalid => Unsupported v1 snapshot, only manifest list is supported"
657 )
658 }
659 }
660
661 #[test]
662 fn test_snapshot_v1_to_v2_with_missing_summary() {
663 use crate::spec::snapshot::_serde::SnapshotV1;
664
665 let v1_snapshot = SnapshotV1 {
667 snapshot_id: 1111111111,
668 parent_snapshot_id: None,
669 timestamp_ms: 1515100955770,
670 manifest_list: Some("s3://bucket/manifest-list.avro".to_string()),
671 manifests: None,
672 summary: None, schema_id: None,
674 };
675
676 let v2_snapshot: Snapshot = v1_snapshot.try_into().unwrap();
678
679 assert_eq!(
681 v2_snapshot.sequence_number(),
682 0,
683 "V1 snapshot sequence_number should default to 0"
684 );
685 assert_eq!(
686 v2_snapshot.summary().operation,
687 Operation::Append,
688 "Missing V1 summary should default to Append operation"
689 );
690 assert!(
691 v2_snapshot.summary().additional_properties.is_empty(),
692 "Default summary should have empty additional_properties"
693 );
694
695 assert_eq!(v2_snapshot.snapshot_id(), 1111111111);
697 assert_eq!(v2_snapshot.parent_snapshot_id(), None);
698 assert_eq!(v2_snapshot.schema_id(), None);
699 }
700}