Skip to main content

froe/writer/
commit.rs

1//! The commit path: editing the content tree and managing checkpoints.
2//!
3//! Every commit rewrites the *super-root* — the node the journal points
4//! at, whose children are `root` (the content tree) and `checkpoints` —
5//! and advances the head with a compare-and-set, exactly like Oak's
6//! `SegmentNodeStore`. Untouched subtrees, property values, and child
7//! maps are preserved by record identifier, so a commit writes only the
8//! spine of records from the changed nodes up to the super-root.
9//!
10//! Checkpoints follow `LockBasedScheduler.CPCreator` faithfully: creation
11//! first purges expired or corrupt checkpoints, then records `timestamp`
12//! (the expiry, clamped at the 64-bit maximum), `created`, a `properties`
13//! child with the caller's string metadata, and a `root` child that
14//! *shares* the current content root's record — an immutable snapshot at
15//! zero cost. Releasing a checkpoint is plain child removal.
16
17use std::collections::BTreeMap;
18
19use crate::content::provider::SegmentProvider;
20use crate::content::template::ChildNodeArity;
21use crate::error::{Error, Result};
22use crate::segment::record::RecordIdentifier;
23use crate::writer::record_writer::{
24    ChildNodesToWrite, PropertyToWrite, PropertyValuesToWrite, RecordWriter, SegmentSink,
25};
26use crate::writer::store_writer::WritableRepository;
27
28/// A named edit to one node's children: `Some` inserts or replaces the
29/// child with the given node record, `None` removes it.
30pub type ChildEdits = BTreeMap<String, Option<RecordIdentifier>>;
31
32/// The reusable parts of an existing node.
33struct NodeParts {
34    primary_type: Option<String>,
35    mixin_types: Vec<String>,
36    properties: Vec<PropertyToWrite>,
37    children: ExistingChildren,
38}
39
40/// The child structure of an existing node.
41enum ExistingChildren {
42    Zero,
43    One {
44        name: String,
45        node: RecordIdentifier,
46    },
47    Many {
48        map: RecordIdentifier,
49    },
50}
51
52/// Reads the parts of an existing node needed to rewrite it while
53/// preserving untouched records by identifier.
54fn read_node_parts(provider: &dyn SegmentProvider, node: RecordIdentifier) -> Result<NodeParts> {
55    let view = provider.segment(node.segment)?;
56    let template_identifier = view.read_record_identifier(node.record_number, 0, 1)?;
57    let template = provider.template(template_identifier)?;
58
59    let children = match &template.child_arity {
60        ChildNodeArity::Zero => ExistingChildren::Zero,
61        ChildNodeArity::One { child_name } => ExistingChildren::One {
62            name: child_name.clone(),
63            node: view.read_record_identifier(node.record_number, 0, 2)?,
64        },
65        ChildNodeArity::Many => ExistingChildren::Many {
66            map: view.read_record_identifier(node.record_number, 0, 2)?,
67        },
68    };
69
70    let mut properties = Vec::with_capacity(template.properties.len());
71    if !template.properties.is_empty() {
72        let slot = if matches!(template.child_arity, ChildNodeArity::Zero) {
73            2
74        } else {
75            3
76        };
77        let list_identifier = view.read_record_identifier(node.record_number, 0, slot)?;
78        for (index, property) in template.properties.iter().enumerate() {
79            let value_slot = crate::content::list::uncounted_list_entry(
80                provider,
81                list_identifier,
82                template.properties.len() as u64,
83                index as u64,
84            )?;
85            properties.push(PropertyToWrite {
86                name: property.name.clone(),
87                property_type: property.property_type,
88                values: PropertyValuesToWrite::PreservedSlot {
89                    value_slot,
90                    is_multiple: property.is_multiple,
91                },
92            });
93        }
94    }
95
96    Ok(NodeParts {
97        primary_type: template.primary_type.clone(),
98        mixin_types: template.mixin_types.clone(),
99        properties,
100        children,
101    })
102}
103
104/// Rewrites a node with edits to its children, preserving its primary
105/// type, mixin types, and every property slot. `base` of `None` starts
106/// from an empty node.
107#[allow(
108    clippy::missing_panics_doc,
109    reason = "the single-entry expect is guarded by the match arm's length"
110)]
111pub fn rewrite_node_with_child_edits<Sink: SegmentSink>(
112    provider: &dyn SegmentProvider,
113    writer: &mut RecordWriter<Sink>,
114    base: Option<RecordIdentifier>,
115    edits: &ChildEdits,
116) -> Result<RecordIdentifier> {
117    let parts = match base {
118        Some(node) => read_node_parts(provider, node)?,
119        None => NodeParts {
120            primary_type: None,
121            mixin_types: Vec::new(),
122            properties: Vec::new(),
123            children: ExistingChildren::Zero,
124        },
125    };
126
127    // Materialize the resulting child set. Unchanged many-child maps are
128    // preserved wholesale; any edit rebuilds the map from its entries
129    // (which themselves preserve the child records by identifier).
130    let children = if edits.is_empty() {
131        match parts.children {
132            ExistingChildren::Zero => ChildNodesToWrite::Zero,
133            ExistingChildren::One { name, node } => ChildNodesToWrite::One { name, node },
134            ExistingChildren::Many { map } => ChildNodesToWrite::ManyExistingMap(map),
135        }
136    } else {
137        let mut entries: BTreeMap<String, RecordIdentifier> = match parts.children {
138            ExistingChildren::Zero => BTreeMap::new(),
139            ExistingChildren::One { name, node } => {
140                let mut entries = BTreeMap::new();
141                entries.insert(name, node);
142                entries
143            }
144            ExistingChildren::Many { map } => crate::content::map::map_entries(provider, map)?
145                .into_iter()
146                .map(|entry| (entry.name, entry.value))
147                .collect(),
148        };
149        for (name, edit) in edits {
150            match edit {
151                Some(node) => {
152                    entries.insert(name.clone(), *node);
153                }
154                None => {
155                    entries.remove(name);
156                }
157            }
158        }
159        match entries.len() {
160            0 => ChildNodesToWrite::Zero,
161            1 => {
162                let (name, node) = entries.into_iter().next().expect("one entry");
163                ChildNodesToWrite::One { name, node }
164            }
165            _ => ChildNodesToWrite::Many(entries.into_iter().collect()),
166        }
167    };
168
169    writer.write_node(
170        parts.primary_type.as_deref(),
171        &parts.mixin_types,
172        &children,
173        &parts.properties,
174    )
175}
176
177/// The checkpoint metadata AEM reads back.
178#[derive(Clone, Debug, PartialEq, Eq)]
179pub struct CheckpointDescription {
180    /// The checkpoint's name (a random UUID string).
181    pub name: String,
182    /// Creation time in milliseconds since the epoch, when present.
183    pub created_milliseconds: Option<i64>,
184    /// Expiry time in milliseconds since the epoch, when present.
185    pub expires_milliseconds: Option<i64>,
186}
187
188/// Lists the checkpoints of the store's current head.
189pub fn list_checkpoints(store: &WritableRepository) -> Result<Vec<CheckpointDescription>> {
190    let head = store.head_node();
191    let Some(checkpoints) = head.child_node("checkpoints")? else {
192        return Ok(Vec::new());
193    };
194    let mut descriptions = Vec::new();
195    for (name, checkpoint) in checkpoints.child_node_entries()? {
196        descriptions.push(CheckpointDescription {
197            created_milliseconds: read_long_property(&checkpoint, "created")?,
198            expires_milliseconds: read_long_property(&checkpoint, "timestamp")?,
199            name,
200        });
201    }
202    Ok(descriptions)
203}
204
205fn read_long_property(
206    node: &crate::content::node::NodeState<'_>,
207    name: &str,
208) -> Result<Option<i64>> {
209    use crate::content::node::PropertyValues;
210    use crate::content::property::PropertyValue;
211    Ok(node
212        .property(name)?
213        .and_then(|property| match property.values {
214            PropertyValues::Single(PropertyValue::Long(value)) => Some(value),
215            _ => None,
216        }))
217}
218
219/// Creates a checkpoint with the given lifetime and string metadata,
220/// returning its generated name. Expired checkpoints are purged first,
221/// exactly like Oak's checkpoint creator.
222pub fn create_checkpoint(
223    store: &WritableRepository,
224    lifetime_milliseconds: i64,
225    properties: &[(String, String)],
226) -> Result<String> {
227    if lifetime_milliseconds <= 0 {
228        return Err(Error::InvalidFormat {
229            details: "checkpoint lifetime must be positive".to_owned(),
230        });
231    }
232    let now = current_time_milliseconds();
233    let name = random_checkpoint_name();
234    let expiry = if i64::MAX - now > lifetime_milliseconds {
235        now + lifetime_milliseconds
236    } else {
237        i64::MAX
238    };
239
240    let head = store.head();
241    let head_node = store.head_node();
242    let checkpoints_container = head_node.child_node("checkpoints")?;
243    let content_root = head_node
244        .child_node("root")?
245        .ok_or_else(|| Error::InvalidFormat {
246            details: "the super-root has no \"root\" child node".to_owned(),
247        })?;
248
249    let generation = store.writing_generation()?;
250    let mut writer = store.record_writer(generation);
251
252    // Purge expired or corrupt checkpoints while assembling the edits.
253    let mut edits: ChildEdits = ChildEdits::new();
254    if let Some(container) = &checkpoints_container {
255        for (existing_name, checkpoint) in container.child_node_entries()? {
256            let expires = read_long_property(&checkpoint, "timestamp")?;
257            if expires.is_none_or(|expires| now > expires) {
258                edits.insert(existing_name, None);
259            }
260        }
261    }
262
263    // The properties child holds the caller's string metadata.
264    let mut property_writes = Vec::with_capacity(properties.len());
265    for (key, value) in properties {
266        let value_identifier = writer.write_string(value)?;
267        property_writes.push(PropertyToWrite {
268            name: key.clone(),
269            property_type: crate::content::property::PropertyType::String,
270            values: PropertyValuesToWrite::Single(value_identifier),
271        });
272    }
273    crate::writer::record_writer::sort_properties_for_template(&mut property_writes);
274    let properties_node =
275        writer.write_node(None, &[], &ChildNodesToWrite::Zero, &property_writes)?;
276
277    // The checkpoint node: timestamp, created, properties, root snapshot.
278    let timestamp_value = writer.write_string(&expiry.to_string())?;
279    let created_value = writer.write_string(&now.to_string())?;
280    let mut checkpoint_properties = vec![
281        PropertyToWrite {
282            name: "timestamp".to_owned(),
283            property_type: crate::content::property::PropertyType::Long,
284            values: PropertyValuesToWrite::Single(timestamp_value),
285        },
286        PropertyToWrite {
287            name: "created".to_owned(),
288            property_type: crate::content::property::PropertyType::Long,
289            values: PropertyValuesToWrite::Single(created_value),
290        },
291    ];
292    crate::writer::record_writer::sort_properties_for_template(&mut checkpoint_properties);
293    let checkpoint_node = writer.write_node(
294        None,
295        &[],
296        &ChildNodesToWrite::Many(vec![
297            ("properties".to_owned(), properties_node),
298            ("root".to_owned(), content_root.record_identifier()),
299        ]),
300        &checkpoint_properties,
301    )?;
302    edits.insert(name.clone(), Some(checkpoint_node));
303
304    // Rebuild the checkpoints container and the super-root.
305    let container_base = checkpoints_container.map(|container| container.record_identifier());
306    let new_container = rewrite_node_with_child_edits(store, &mut writer, container_base, &edits)?;
307    let mut super_root_edits = ChildEdits::new();
308    super_root_edits.insert("checkpoints".to_owned(), Some(new_container));
309    let new_head =
310        rewrite_node_with_child_edits(store, &mut writer, Some(head), &super_root_edits)?;
311    writer.finish()?;
312
313    if !store.set_head(head, new_head) {
314        return Err(Error::InvalidFormat {
315            details: "the head moved while creating the checkpoint".to_owned(),
316        });
317    }
318    store.flush()?;
319    Ok(name)
320}
321
322/// Releases (removes) a checkpoint by name. Returns whether it existed.
323pub fn release_checkpoint(store: &WritableRepository, name: &str) -> Result<bool> {
324    let head = store.head();
325    let head_node = store.head_node();
326    let Some(container) = head_node.child_node("checkpoints")? else {
327        return Ok(false);
328    };
329    if container.child_node(name)?.is_none() {
330        return Ok(false);
331    }
332
333    let generation = store.writing_generation()?;
334    let mut writer = store.record_writer(generation);
335    let mut edits = ChildEdits::new();
336    edits.insert(name.to_owned(), None);
337    let new_container = rewrite_node_with_child_edits(
338        store,
339        &mut writer,
340        Some(container.record_identifier()),
341        &edits,
342    )?;
343    let mut super_root_edits = ChildEdits::new();
344    super_root_edits.insert("checkpoints".to_owned(), Some(new_container));
345    let new_head =
346        rewrite_node_with_child_edits(store, &mut writer, Some(head), &super_root_edits)?;
347    writer.finish()?;
348
349    if !store.set_head(head, new_head) {
350        return Err(Error::InvalidFormat {
351            details: "the head moved while releasing the checkpoint".to_owned(),
352        });
353    }
354    store.flush()?;
355    Ok(true)
356}
357
358/// Removes every checkpoint. Returns how many were removed.
359pub fn remove_all_checkpoints(store: &WritableRepository) -> Result<u64> {
360    let names: Vec<String> = list_checkpoints(store)?
361        .into_iter()
362        .map(|checkpoint| checkpoint.name)
363        .collect();
364    let mut removed = 0u64;
365    for name in names {
366        if release_checkpoint(store, &name)? {
367            removed += 1;
368        }
369    }
370    Ok(removed)
371}
372
373/// Removes every checkpoint not referenced from the asynchronous indexer
374/// state at `/:async`. Conservative: any string property value (or
375/// member of a multi-valued string property) of that node counts as a
376/// reference.
377pub fn remove_unreferenced_checkpoints(store: &WritableRepository) -> Result<u64> {
378    use crate::content::node::PropertyValues;
379    use crate::content::property::PropertyValue;
380
381    let mut referenced: std::collections::HashSet<String> = std::collections::HashSet::new();
382    let head_node = store.head_node();
383    if let Some(content_root) = head_node.child_node("root")?
384        && let Some(async_state) = content_root.child_node(":async")?
385    {
386        for property in async_state.properties()? {
387            match property.values {
388                PropertyValues::Single(PropertyValue::String(value)) => {
389                    referenced.insert(value);
390                }
391                PropertyValues::Multiple(values) => {
392                    for value in values {
393                        if let PropertyValue::String(value) = value {
394                            referenced.insert(value);
395                        }
396                    }
397                }
398                PropertyValues::Single(_) => {}
399            }
400        }
401    }
402
403    let names: Vec<String> = list_checkpoints(store)?
404        .into_iter()
405        .map(|checkpoint| checkpoint.name)
406        .filter(|name| !referenced.contains(name))
407        .collect();
408    let mut removed = 0u64;
409    for name in names {
410        if release_checkpoint(store, &name)? {
411            removed += 1;
412        }
413    }
414    Ok(removed)
415}
416
417/// Replaces the content root (`/root` of the super-root) with an already
418/// written node record and advances the head. The restore operation's
419/// final step.
420pub fn replace_content_root(store: &WritableRepository, new_root: RecordIdentifier) -> Result<()> {
421    let head = store.head();
422    let generation = store.writing_generation()?;
423    let mut writer = store.record_writer(generation);
424    let mut edits = ChildEdits::new();
425    edits.insert("root".to_owned(), Some(new_root));
426    let new_head = rewrite_node_with_child_edits(store, &mut writer, Some(head), &edits)?;
427    writer.finish()?;
428    if !store.set_head(head, new_head) {
429        return Err(Error::InvalidFormat {
430            details: "the head moved while replacing the content root".to_owned(),
431        });
432    }
433    store.flush()?;
434    Ok(())
435}
436
437/// The current wall clock in milliseconds since the Unix epoch.
438fn current_time_milliseconds() -> i64 {
439    std::time::SystemTime::now()
440        .duration_since(std::time::UNIX_EPOCH)
441        .map_or(0, |duration| duration.as_millis() as i64)
442}
443
444/// A random version 4 UUID string, the checkpoint naming scheme.
445fn random_checkpoint_name() -> String {
446    // Reuse the segment identifier entropy stream; a checkpoint name is
447    // an ordinary v4 UUID without the segment kind marker.
448    let identifier = crate::writer::identifier_generator::new_data_segment_identifier();
449    let most = identifier.most_significant_bits;
450    // Restore a proper random variant nibble (10xx) in place of the data
451    // segment marker.
452    let least = (identifier.least_significant_bits & 0x3FFF_FFFF_FFFF_FFFF) | 0x8000_0000_0000_0000;
453    format!(
454        "{:08x}-{:04x}-{:04x}-{:04x}-{:012x}",
455        most >> 32,
456        (most >> 16) & 0xFFFF,
457        most & 0xFFFF,
458        least >> 48,
459        least & 0xFFFF_FFFF_FFFF,
460    )
461}
462
463#[cfg(test)]
464mod tests {
465    use super::{create_checkpoint, list_checkpoints, release_checkpoint, remove_all_checkpoints};
466    use crate::store::Repository;
467    use crate::writer::store_writer::WritableRepository;
468
469    struct TestDirectory {
470        path: std::path::PathBuf,
471    }
472
473    impl TestDirectory {
474        fn new(name: &str) -> Self {
475            let path =
476                std::env::temp_dir().join(format!("froe-commit-{name}-{}", std::process::id()));
477            let _ = std::fs::remove_dir_all(&path);
478            Self { path }
479        }
480    }
481
482    impl Drop for TestDirectory {
483        fn drop(&mut self) {
484            let _ = std::fs::remove_dir_all(&self.path);
485        }
486    }
487
488    #[test]
489    fn checkpoints_are_created_shared_and_released() {
490        let directory = TestDirectory::new("lifecycle");
491        {
492            let store = WritableRepository::open(&directory.path).expect("bootstrap");
493            let name = create_checkpoint(
494                &store,
495                1_000_000,
496                &[("creator".to_owned(), "froe-test".to_owned())],
497            )
498            .expect("create");
499
500            let listed = list_checkpoints(&store).expect("list");
501            assert_eq!(listed.len(), 1);
502            assert_eq!(listed[0].name, name);
503            assert!(listed[0].created_milliseconds.is_some());
504            assert!(listed[0].expires_milliseconds.is_some());
505            store.close().expect("close");
506
507            // The reader sees the checkpoint with a shared root snapshot.
508            let repository = Repository::open(&directory.path).expect("reader");
509            let checkpoints = repository.checkpoints().expect("checkpoints");
510            assert_eq!(checkpoints.len(), 1);
511            let (_, checkpoint) = &checkpoints[0];
512            let snapshot = checkpoint
513                .child_node("root")
514                .expect("read")
515                .expect("snapshot present");
516            let live_root = repository.content_root().expect("content root");
517            assert_eq!(
518                snapshot.record_identifier(),
519                live_root.record_identifier(),
520                "the snapshot shares the live root's record"
521            );
522            let properties = checkpoint
523                .child_node("properties")
524                .expect("read")
525                .expect("properties present");
526            let creator = properties
527                .property("creator")
528                .expect("read")
529                .expect("present");
530            assert_eq!(
531                creator.values,
532                crate::content::node::PropertyValues::Single(
533                    crate::content::property::PropertyValue::String("froe-test".to_owned())
534                )
535            );
536        }
537        {
538            let store = WritableRepository::open(&directory.path).expect("reopen");
539            let listed = list_checkpoints(&store).expect("list");
540            assert_eq!(listed.len(), 1);
541            assert!(
542                release_checkpoint(&store, &listed[0].name).expect("release"),
543                "the checkpoint existed"
544            );
545            assert!(list_checkpoints(&store).expect("list").is_empty());
546            assert!(
547                !release_checkpoint(&store, "missing").expect("release"),
548                "absent checkpoints report false"
549            );
550            store.close().expect("close");
551        }
552    }
553
554    #[test]
555    fn expired_checkpoints_are_purged_on_create() {
556        let directory = TestDirectory::new("expiry");
557        let store = WritableRepository::open(&directory.path).expect("bootstrap");
558        // A one-millisecond lifetime expires immediately.
559        let expired = create_checkpoint(&store, 1, &[]).expect("create short lived");
560        std::thread::sleep(std::time::Duration::from_millis(5));
561        let fresh = create_checkpoint(&store, 1_000_000, &[]).expect("create fresh");
562
563        let names: Vec<String> = list_checkpoints(&store)
564            .expect("list")
565            .into_iter()
566            .map(|checkpoint| checkpoint.name)
567            .collect();
568        assert!(
569            !names.contains(&expired),
570            "the expired checkpoint is purged"
571        );
572        assert!(names.contains(&fresh));
573        store.close().expect("close");
574    }
575
576    #[test]
577    fn remove_all_reports_the_count() {
578        let directory = TestDirectory::new("remove-all");
579        let store = WritableRepository::open(&directory.path).expect("bootstrap");
580        create_checkpoint(&store, 1_000_000, &[]).expect("first");
581        create_checkpoint(&store, 1_000_000, &[]).expect("second");
582        assert_eq!(remove_all_checkpoints(&store).expect("remove"), 2);
583        assert!(list_checkpoints(&store).expect("list").is_empty());
584        store.close().expect("close");
585    }
586}