Skip to main content

iota_snapshot/
writer.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5#![allow(dead_code)]
6
7use std::{
8    collections::{HashMap, hash_map::Entry::Vacant},
9    fs,
10    fs::{File, OpenOptions},
11    io::{BufWriter, Seek, SeekFrom, Write},
12    num::NonZeroUsize,
13    path::PathBuf,
14    sync::Arc,
15};
16
17use anyhow::{Context, Result, anyhow};
18use byteorder::{BigEndian, ByteOrder};
19use fastcrypto::hash::MultisetHash;
20use futures::StreamExt;
21use integer_encoding::VarInt;
22use iota_config::object_storage_config::ObjectStoreConfig;
23use iota_core::{
24    authority::authority_store_tables::{AuthorityPerpetualTables, LiveObject, SnapshotLiveObject},
25    checkpoints::CheckpointStore,
26    global_state_hasher::GlobalStateHasher,
27};
28use iota_sdk_types::{ObjectId, ObjectReference};
29use iota_storage::{
30    blob::{BLOB_ENCODING_BYTES, Blob, BlobEncoding},
31    object_store::util::{copy_file, delete_recursively, path_to_filesystem},
32};
33use iota_types::{
34    digests::ChainIdentifier, global_state_hash::GlobalStateHash,
35    messages_checkpoint::ECMHLiveObjectSetDigest,
36};
37use object_store::{DynObjectStore, path::Path};
38use tokio::{
39    sync::{
40        mpsc,
41        mpsc::{Receiver, Sender},
42    },
43    task::JoinHandle,
44};
45use tokio_stream::wrappers::ReceiverStream;
46use tracing::debug;
47
48use crate::{
49    EPOCH_INFO_FILE_MAGIC, EpochInfo, EpochInfoV1, FILE_MAX_BYTES, FileCompression, FileMetadata,
50    FileType, MAGIC_BYTES, MANIFEST_FILE_MAGIC, Manifest, ManifestV2, OBJECT_FILE_MAGIC,
51    OBJECT_REF_BYTES, REFERENCE_FILE_MAGIC, SEQUENCE_NUM_BYTES, compute_sha3_checksum,
52    create_file_metadata,
53};
54
55/// LiveObjectSetWriterV1 writes live object set. It creates multiple *.obj
56/// files and *.ref file
57struct LiveObjectSetWriterV1 {
58    dir_path: PathBuf,
59    bucket_num: u32,
60    current_part_num: u32,
61    obj_wbuf: BufWriter<File>,
62    ref_wbuf: BufWriter<File>,
63    object_file_size: usize,
64    files: Vec<FileMetadata>,
65    sender: Option<Sender<FileMetadata>>,
66    file_compression: FileCompression,
67}
68
69impl LiveObjectSetWriterV1 {
70    fn new(
71        dir_path: PathBuf,
72        bucket_num: u32,
73        file_compression: FileCompression,
74        sender: Sender<FileMetadata>,
75    ) -> Result<Self> {
76        let part_num = 1;
77        let (n, obj_file) = Self::object_file(dir_path.clone(), bucket_num, part_num)?;
78        let ref_file = Self::ref_file(dir_path.clone(), bucket_num, part_num)?;
79        Ok(LiveObjectSetWriterV1 {
80            dir_path,
81            bucket_num,
82            current_part_num: part_num,
83            obj_wbuf: BufWriter::new(obj_file),
84            ref_wbuf: BufWriter::new(ref_file),
85            object_file_size: n,
86            files: vec![],
87            sender: Some(sender),
88            file_compression,
89        })
90    }
91
92    /// Writes a live object to the object file and the reference to the
93    /// reference file.
94    pub fn write(&mut self, live_object: &LiveObject) -> Result<()> {
95        let object_reference = live_object.object_reference();
96        self.write_object(live_object)?;
97        self.write_object_ref(&object_reference)?;
98        Ok(())
99    }
100
101    /// Finalizes the object and reference files and returns the FileMetadata of
102    /// the files.
103    pub fn done(mut self) -> Result<Vec<FileMetadata>> {
104        self.finalize_obj()?;
105        self.finalize_ref()?;
106        self.sender = None;
107        Ok(self.files.clone())
108    }
109
110    /// Creates a new object file for the provided bucket number and part
111    /// number, and returns the file and the number of bytes written to it.
112    fn object_file(dir_path: PathBuf, bucket_num: u32, part_num: u32) -> Result<(usize, File)> {
113        let next_part_file_path = dir_path.join(format!("{bucket_num}_{part_num}.obj"));
114        let next_part_file_tmp_path = dir_path.join(format!("{bucket_num}_{part_num}.obj.tmp"));
115        let mut f = File::create(next_part_file_tmp_path.clone())?;
116        let mut metab = [0u8; MAGIC_BYTES];
117        BigEndian::write_u32(&mut metab, OBJECT_FILE_MAGIC);
118        f.rewind()?;
119        let n = f.write(&metab)?;
120        drop(f);
121        fs::rename(next_part_file_tmp_path, next_part_file_path.clone())?;
122        let mut f = OpenOptions::new().append(true).open(next_part_file_path)?;
123        f.seek(SeekFrom::Start(n as u64))?;
124        Ok((n, f))
125    }
126
127    /// Creates a new reference file for the provided bucket number and part
128    /// number, and returns the file and the number of bytes written to the
129    /// file.
130    fn ref_file(dir_path: PathBuf, bucket_num: u32, part_num: u32) -> Result<File> {
131        let ref_path = dir_path.join(format!("{bucket_num}_{part_num}.ref"));
132        let ref_tmp_path = dir_path.join(format!("{bucket_num}_{part_num}.ref.tmp"));
133        let mut f = File::create(ref_tmp_path.clone())?;
134        f.rewind()?;
135        let mut metab = [0u8; MAGIC_BYTES];
136        BigEndian::write_u32(&mut metab, REFERENCE_FILE_MAGIC);
137        let n = f.write(&metab)?;
138        drop(f);
139        fs::rename(ref_tmp_path, ref_path.clone())?;
140        let mut f = OpenOptions::new().append(true).open(ref_path)?;
141        f.seek(SeekFrom::Start(n as u64))?;
142        Ok(f)
143    }
144
145    /// Finalizes the object file by flushing the buffer to disk and sends its
146    /// FileMetadata to the channel.
147    fn finalize_obj(&mut self) -> Result<()> {
148        // Flushes the buffer and sync the data to disk
149        self.obj_wbuf.flush()?;
150        self.obj_wbuf.get_ref().sync_data()?;
151        let off = self.obj_wbuf.get_ref().stream_position()?;
152        self.obj_wbuf.get_ref().set_len(off)?;
153        let file_path = self
154            .dir_path
155            .join(format!("{}_{}.obj", self.bucket_num, self.current_part_num));
156        let file_metadata = create_file_metadata(
157            &file_path,
158            self.file_compression,
159            FileType::Object,
160            self.bucket_num,
161            self.current_part_num,
162        )?;
163        self.files.push(file_metadata.clone());
164        if let Some(sender) = &self.sender {
165            sender.blocking_send(file_metadata)?;
166        }
167        Ok(())
168    }
169
170    /// Finalizes the reference file by flushing the buffer to disk and sends
171    /// its FileMetadata to the channel.
172    fn finalize_ref(&mut self) -> Result<()> {
173        // Flushes the buffer and sync the data to disk
174        self.ref_wbuf.flush()?;
175        self.ref_wbuf.get_ref().sync_data()?;
176        let off = self.ref_wbuf.get_ref().stream_position()?;
177        self.ref_wbuf.get_ref().set_len(off)?;
178        let file_path = self
179            .dir_path
180            .join(format!("{}_{}.ref", self.bucket_num, self.current_part_num));
181        let file_metadata = create_file_metadata(
182            &file_path,
183            self.file_compression,
184            FileType::Reference,
185            self.bucket_num,
186            self.current_part_num,
187        )?;
188        self.files.push(file_metadata.clone());
189        if let Some(sender) = &self.sender {
190            sender.blocking_send(file_metadata)?;
191        }
192        Ok(())
193    }
194
195    /// Finalizes the object file of current partition and creates a new one for
196    /// the next partition.
197    fn cut(&mut self) -> Result<()> {
198        self.finalize_obj()?;
199        let (n, f) = Self::object_file(
200            self.dir_path.clone(),
201            self.bucket_num,
202            self.current_part_num + 1,
203        )?;
204        self.object_file_size = n;
205        self.obj_wbuf = BufWriter::new(f);
206        Ok(())
207    }
208
209    /// Finalizes the reference file of current partition and creates a new one
210    /// for the next partition.
211    fn cut_reference_file(&mut self) -> Result<()> {
212        self.finalize_ref()?;
213        let f = Self::ref_file(
214            self.dir_path.clone(),
215            self.bucket_num,
216            self.current_part_num + 1,
217        )?;
218        self.ref_wbuf = BufWriter::new(f);
219        Ok(())
220    }
221
222    /// Writes a live object to the object file. Creates a new partition
223    /// (new object file and reference file) if it exceeds the maximum size.
224    fn write_object(&mut self, live_object: &LiveObject) -> Result<()> {
225        let previous_transaction_checkpoint =
226            live_object.previous_transaction_checkpoint.ok_or_else(|| {
227                anyhow!(
228                    "Snapshot V2 writer: live object {:?} (version {:?}) was lifted from a \
229                     pre-V2 store row and has no `previous_transaction_checkpoint`. This node \
230                     cannot publish V2 snapshots without re-syncing from genesis under V2 or \
231                     starting from a valid V2 snapshot so the entire perpetual store is in V2 \
232                     format.",
233                    live_object.object.id(),
234                    live_object.object.version(),
235                )
236            })?;
237        let snapshot_live_object = SnapshotLiveObject {
238            object: live_object.object.clone(),
239            previous_transaction_checkpoint,
240        };
241        let blob = Blob::encode(&snapshot_live_object, BlobEncoding::Bcs)?;
242        let mut blob_size = blob.data.len().required_space();
243        blob_size += BLOB_ENCODING_BYTES;
244        blob_size += blob.data.len();
245        let cut_new_part_file = (self.object_file_size + blob_size) > FILE_MAX_BYTES;
246        if cut_new_part_file {
247            self.cut()?;
248            self.cut_reference_file()?;
249            self.current_part_num += 1;
250        }
251        self.object_file_size += blob.write(&mut self.obj_wbuf)?;
252        Ok(())
253    }
254
255    /// Writes an object reference to the reference file.
256    fn write_object_ref(&mut self, object_ref: &ObjectReference) -> Result<()> {
257        let mut buf = [0u8; OBJECT_REF_BYTES];
258        buf[0..ObjectId::LENGTH].copy_from_slice(object_ref.object_id.as_ref());
259        BigEndian::write_u64(
260            &mut buf[ObjectId::LENGTH..OBJECT_REF_BYTES],
261            object_ref.version.as_u64(),
262        );
263        buf[ObjectId::LENGTH + SEQUENCE_NUM_BYTES..OBJECT_REF_BYTES]
264            .copy_from_slice(object_ref.digest.as_ref());
265        self.ref_wbuf.write_all(&buf)?;
266        Ok(())
267    }
268}
269
270/// StateSnapshotWriterV1 writes snapshot files to a local staging dir and
271/// simultaneously uploads them to a remote object store
272pub struct StateSnapshotWriterV1 {
273    local_staging_dir: PathBuf,
274    file_compression: FileCompression,
275    remote_object_store: Arc<DynObjectStore>,
276    local_staging_store: Arc<DynObjectStore>,
277    /// Source of `EPOCH_INFO` data for the snapshot: the CheckpointStore's
278    /// `epoch_info` table.
279    checkpoint_store: Arc<CheckpointStore>,
280    /// Chain identifier written into the `ManifestV2`.
281    chain_id: ChainIdentifier,
282    concurrency: NonZeroUsize,
283}
284
285impl StateSnapshotWriterV1 {
286    pub async fn new_from_store(
287        local_staging_path: &std::path::Path,
288        local_staging_store: &Arc<DynObjectStore>,
289        remote_object_store: &Arc<DynObjectStore>,
290        checkpoint_store: Arc<CheckpointStore>,
291        chain_id: ChainIdentifier,
292        file_compression: FileCompression,
293        concurrency: NonZeroUsize,
294    ) -> Result<Self> {
295        Ok(StateSnapshotWriterV1 {
296            file_compression,
297            local_staging_dir: local_staging_path.to_path_buf(),
298            remote_object_store: remote_object_store.clone(),
299            local_staging_store: local_staging_store.clone(),
300            checkpoint_store,
301            chain_id,
302            concurrency,
303        })
304    }
305
306    pub async fn new(
307        local_store_config: &ObjectStoreConfig,
308        remote_store_config: &ObjectStoreConfig,
309        checkpoint_store: Arc<CheckpointStore>,
310        chain_id: ChainIdentifier,
311        file_compression: FileCompression,
312        concurrency: NonZeroUsize,
313    ) -> Result<Self> {
314        let remote_object_store = remote_store_config.make()?;
315        let local_staging_store = local_store_config.make()?;
316        let local_staging_dir = local_store_config
317            .directory
318            .as_ref()
319            .context("No local directory specified")?
320            .clone();
321        Ok(StateSnapshotWriterV1 {
322            local_staging_dir,
323            file_compression,
324            remote_object_store,
325            local_staging_store,
326            checkpoint_store,
327            chain_id,
328            concurrency,
329        })
330    }
331
332    /// Retrieves the system state object from the perpetual database, writes
333    /// the state snapshot for the specified epoch to the local staging
334    /// directory, and uploads it to the remote store.
335    pub async fn write(
336        self,
337        epoch: u64,
338        perpetual_db: Arc<AuthorityPerpetualTables>,
339        root_state_hash: ECMHLiveObjectSetDigest,
340    ) -> Result<()> {
341        self.write_internal(epoch, perpetual_db, root_state_hash)
342            .await
343    }
344
345    /// Writes the state snapshot for the provided epoch to the local staging
346    /// directory and uploads it to the remote store.
347    pub(crate) async fn write_internal(
348        mut self,
349        epoch: u64,
350        perpetual_db: Arc<AuthorityPerpetualTables>,
351        root_state_hash: ECMHLiveObjectSetDigest,
352    ) -> Result<()> {
353        // Fail fast on the epoch-info completeness precondition so a node with
354        // an incomplete epoch chain does not perform a full live-object scan
355        // (tens of GiB on mainnet-sized DBs) before failing.
356        self.check_epoch_watermark(epoch)?;
357
358        self.setup_epoch_dir(epoch).await?;
359
360        let manifest_file_path = self.epoch_dir(epoch).child("MANIFEST");
361        let local_staging_dir = self.local_staging_dir.clone();
362        let local_object_store = self.local_staging_store.clone();
363        let remote_object_store = self.remote_object_store.clone();
364
365        let (sender, receiver) = mpsc::channel::<FileMetadata>(1000);
366        // Starts the upload loop, which listens on the receiver for FileMetadata
367        let upload_handle = self.start_upload(epoch, receiver)?;
368        let write_handler = tokio::task::spawn_blocking(move || {
369            self.write_live_object_set(
370                epoch,
371                perpetual_db,
372                sender,
373                Self::bucket_func,
374                root_state_hash,
375            )
376        });
377        // Awaits the object and reference files to be written to the local staging
378        // directory and informs the upload loop
379        write_handler
380            .await?
381            .context(format!("Failed to write state snapshot for epoch: {epoch}"))?;
382
383        // Awaits the upload loop to finish
384        upload_handle.await?.context(format!(
385            "Failed to upload state snapshot for epoch: {epoch}"
386        ))?;
387
388        // Syncs the manifest file to the remote store
389        Self::sync_file_to_remote(
390            local_staging_dir,
391            manifest_file_path,
392            local_object_store,
393            remote_object_store,
394        )
395        .await?;
396        Ok(())
397    }
398
399    /// Starts listening on the receiver for FileMetadata and uploads the files
400    /// to the remote store in parallel.
401    fn start_upload(
402        &self,
403        epoch: u64,
404        receiver: Receiver<FileMetadata>,
405    ) -> Result<JoinHandle<Result<Vec<()>, anyhow::Error>>> {
406        let remote_object_store = self.remote_object_store.clone();
407        let local_staging_store = self.local_staging_store.clone();
408        let local_dir_path = self.local_staging_dir.clone();
409        let epoch_dir = self.epoch_dir(epoch);
410        let concurrency = self.concurrency;
411        let join_handle = tokio::spawn(async move {
412            // Uploads the files to the remote store in parallel for each received
413            // FileMetadata
414            let results: Vec<Result<(), anyhow::Error>> = ReceiverStream::new(receiver)
415                .map(|file_metadata| {
416                    let file_path = file_metadata.file_path(&epoch_dir);
417                    let remote_object_store = remote_object_store.clone();
418                    let local_object_store = local_staging_store.clone();
419                    let local_dir_path = local_dir_path.clone();
420                    async move {
421                        Self::sync_file_to_remote(
422                            local_dir_path.clone(),
423                            file_path.clone(),
424                            local_object_store.clone(),
425                            remote_object_store.clone(),
426                        )
427                        .await?;
428                        Ok(())
429                    }
430                })
431                .boxed()
432                .buffer_unordered(concurrency.into())
433                .collect()
434                .await;
435            results
436                .into_iter()
437                .collect::<Result<Vec<()>, anyhow::Error>>()
438        });
439        Ok(join_handle)
440    }
441
442    /// Writes the provided live object set in the form of reference files,
443    /// object files, EPOCH_INFO, and MANIFEST. These files are stored in the
444    /// local staging directory and the FileMetadata is sent to the channel.
445    fn write_live_object_set<F>(
446        &mut self,
447        epoch: u64,
448        perpetual_db: Arc<AuthorityPerpetualTables>,
449        sender: Sender<FileMetadata>,
450        bucket_func: F,
451        root_state_hash: ECMHLiveObjectSetDigest,
452    ) -> Result<()>
453    where
454        F: Fn(&LiveObject) -> u32,
455    {
456        let mut object_writers: HashMap<u32, LiveObjectSetWriterV1> = HashMap::new();
457        let local_staging_dir_path =
458            path_to_filesystem(self.local_staging_dir.clone(), &self.epoch_dir(epoch))?;
459        let mut acc = GlobalStateHash::default();
460        for live_object in perpetual_db.iter_live_object_set() {
461            GlobalStateHasher::accumulate_live_object(&mut acc, &live_object);
462            let bucket_num = bucket_func(&live_object);
463            // Creates a new LiveObjectSetWriterV1 for the bucket if it does not exist
464            if let Vacant(slot) = object_writers.entry(bucket_num) {
465                slot.insert(LiveObjectSetWriterV1::new(
466                    local_staging_dir_path.clone(),
467                    bucket_num,
468                    self.file_compression,
469                    sender.clone(),
470                )?);
471            }
472            let writer = object_writers
473                .get_mut(&bucket_num)
474                .context("Unexpected missing bucket writer")?;
475            writer.write(&live_object)?;
476        }
477        assert_eq!(
478            ECMHLiveObjectSetDigest::from(acc.digest()),
479            root_state_hash,
480            "Root state hash mismatch!"
481        );
482        let mut files = vec![];
483        // Flushes the object and reference files to disk, informs the file channel of
484        // flushed files and get the FileMetadata
485        for (_, writer) in object_writers.into_iter() {
486            files.extend(writer.done()?);
487        }
488        // Emit the EPOCH_INFO file alongside the bucket files. It must go through
489        // the same upload channel as `.obj`/`.ref` files so the existing
490        // upload-MANIFEST-last invariant continues to imply all referenced
491        // files are present.
492        let epoch_info_metadata = self.write_epoch_info(epoch, &local_staging_dir_path, &sender)?;
493        files.push(epoch_info_metadata);
494        // Write the manifest file for the epoch(bucket)
495        self.write_manifest(epoch, files)?;
496        Ok(())
497    }
498
499    /// Verifies that every epoch in `[0, epoch]` is finalized in the
500    /// `epoch_info` table, failing fast before any disk work. `None` and
501    /// `Some(h) where h < epoch` are distinct failure modes with distinct
502    /// remediations — keep them as separate messages.
503    fn check_epoch_watermark(&self, epoch: u64) -> Result<()> {
504        match self.checkpoint_store.highest_indexed_epoch()? {
505            None => Err(anyhow!(
506                "Snapshot V2 writer: the epoch_info completeness watermark is \
507                 absent — no epoch_info rows have been finalized on this node \
508                 yet. Wait until at least epoch 0 closes under live indexing, \
509                 or restore this node from a formal snapshot."
510            )),
511            Some(h) if h < epoch => Err(anyhow!(
512                "Snapshot V2 writer: the epoch_info completeness watermark is at \
513                 epoch {h}, but snapshot_epoch is {epoch}. The chain is \
514                 incomplete; restore this node from a formal snapshot or resync \
515                 it from genesis before publishing."
516            )),
517            Some(_) => Ok(()),
518        }
519    }
520
521    /// Writes the per-snapshot `EPOCH_INFO` file, one entry per epoch in
522    /// `[0, epoch]`. Callers must have run
523    /// [`Self::check_epoch_watermark`] first; this function trusts the
524    /// precondition and panics on any unfinalized row.
525    ///
526    /// File layout: 4-byte magic | bcs(EpochInfo). Integrity is anchored
527    /// by `FileMetadata::sha3_digest` in the MANIFEST — no in-file sha3
528    /// trailer.
529    fn write_epoch_info(
530        &self,
531        epoch: u64,
532        local_staging_dir_path: &std::path::Path,
533        sender: &Sender<FileMetadata>,
534    ) -> Result<FileMetadata> {
535        let mut entries = Vec::with_capacity((epoch + 1) as usize);
536        // O(epochs) point lookups. Cheap relative to writing the live-object
537        // set (millions of rows) and to the snapshot upload, so the simple
538        // loop is fine; a range scan would be a micro-optimization.
539        for epoch_id in 0..=epoch {
540            // The watermark precondition above guarantees every entry in
541            // `[0, epoch]` is present and finalized; the panics below turn any
542            // watermark/row inconsistency into a loud failure rather than a
543            // silently truncated snapshot.
544            // `panic!` is deliberate: this runs inside `spawn_blocking`,
545            // so the panic surfaces as `JoinError` and fails only the
546            // snapshot task — exactly the desired blast radius.
547            let epoch_info = self
548                .checkpoint_store
549                .get_epoch_info(epoch_id)?
550                .unwrap_or_else(|| {
551                    panic!(
552                        "epoch_info[{epoch_id}] is absent despite the completeness \
553                         watermark covering it — watermark/row inconsistency"
554                    )
555                });
556            // The boundary's whole entry is committed in one atomic batch, so a
557            // row the watermark covers must be finalized; a missing entry means
558            // the watermark advanced over an unfinalized row. Panics for the
559            // same reason as above.
560            let entry = epoch_info.epoch_close_proof.unwrap_or_else(|| {
561                panic!(
562                    "epoch_info[{epoch_id}] is not finalized despite the completeness \
563                     watermark covering it — the watermark must never cover an \
564                     unfinalized row"
565                )
566            });
567            // Turn a silent miswrite (row stored under the wrong epoch key)
568            // into a loud panic at snapshot time.
569            assert_eq!(
570                entry.last_checkpoint_summary.epoch(),
571                epoch_id,
572                "epoch_info[{epoch_id}] is populated with an entry for epoch {}; the \
573                 snapshot would silently misattribute checkpoints",
574                entry.last_checkpoint_summary.epoch(),
575            );
576
577            entries.push(entry);
578        }
579        let epoch_info = EpochInfo::V1(EpochInfoV1 { entries });
580        let serialized = bcs::to_bytes(&epoch_info)?;
581
582        let file_path = local_staging_dir_path.join("EPOCH_INFO");
583        let mut metab = [0u8; MAGIC_BYTES];
584        BigEndian::write_u32(&mut metab, EPOCH_INFO_FILE_MAGIC);
585
586        let mut f = File::create(&file_path)?;
587        f.write_all(&metab)?;
588        f.write_all(&serialized)?;
589        f.sync_data()?;
590        drop(f);
591
592        // Use bucket_num/part_num 0; EPOCH_INFO is a singleton per snapshot
593        // and the filename does not include them.
594        let file_metadata =
595            create_file_metadata(&file_path, self.file_compression, FileType::EpochInfo, 0, 0)?;
596        sender.blocking_send(file_metadata.clone())?;
597        Ok(file_metadata)
598    }
599
600    /// Writes the manifest file for the provided FileMetadata of an epoch and
601    /// its sha3 checksum.
602    fn write_manifest(&mut self, epoch: u64, file_metadata: Vec<FileMetadata>) -> Result<()> {
603        let (f, manifest_file_path) = self.manifest_file(epoch)?;
604        let mut wbuf = BufWriter::new(f);
605        let manifest: Manifest = Manifest::V2(ManifestV2 {
606            snapshot_version: 2,
607            address_length: ObjectId::LENGTH as u64,
608            file_metadata,
609            epoch,
610            chain_id: self.chain_id,
611        });
612        let serialized_manifest = bcs::to_bytes(&manifest)?;
613        wbuf.write_all(&serialized_manifest)?;
614        wbuf.flush()?;
615        wbuf.get_ref().sync_data()?;
616        // Computes the sha3 checksum of the manifest file and write it to the end of
617        // the file
618        let sha3_digest = compute_sha3_checksum(&manifest_file_path)?;
619        wbuf.write_all(&sha3_digest)?;
620        wbuf.flush()?;
621        wbuf.get_ref().sync_data()?;
622        let off = wbuf.get_ref().stream_position()?;
623        wbuf.get_ref().set_len(off)?;
624        Ok(())
625    }
626
627    /// Creates a new manifest file for the provided epoch and returns the file
628    /// and the path to the file.
629    fn manifest_file(&mut self, epoch: u64) -> Result<(File, PathBuf)> {
630        let manifest_file_path = path_to_filesystem(
631            self.local_staging_dir.clone(),
632            &self.epoch_dir(epoch).child("MANIFEST"),
633        )?;
634        let manifest_file_tmp_path = path_to_filesystem(
635            self.local_staging_dir.clone(),
636            &self.epoch_dir(epoch).child("MANIFEST.tmp"),
637        )?;
638        let mut f = File::create(manifest_file_tmp_path.clone())?;
639        let mut metab = vec![0u8; MAGIC_BYTES];
640        BigEndian::write_u32(&mut metab, MANIFEST_FILE_MAGIC);
641        f.rewind()?;
642        f.write_all(&metab)?;
643        drop(f);
644        fs::rename(manifest_file_tmp_path, manifest_file_path.clone())?;
645        let mut f = OpenOptions::new()
646            .append(true)
647            .open(manifest_file_path.clone())?;
648        f.seek(SeekFrom::Start(MAGIC_BYTES as u64))?;
649        Ok((f, manifest_file_path))
650    }
651
652    fn bucket_func(_live_object: &LiveObject) -> u32 {
653        // TODO: Use the hash bucketing function used for accumulator tree if there is
654        // one
655        1u32
656    }
657
658    fn epoch_dir(&self, epoch: u64) -> Path {
659        Path::from(format!("epoch_{epoch}"))
660    }
661
662    /// Creates a new epoch directory and a new staging directory for the epoch
663    /// in the local store. Deletes the old ones if they exist.
664    async fn setup_epoch_dir(&self, epoch: u64) -> Result<()> {
665        let epoch_dir = self.epoch_dir(epoch);
666        // Deletes remote epoch dir if it exists
667        delete_recursively(&epoch_dir, &self.remote_object_store, self.concurrency).await?;
668        // Deletes local staging epoch dir if it exists
669        let local_epoch_dir_path = self.local_staging_dir.join(format!("epoch_{epoch}"));
670        if local_epoch_dir_path.exists() {
671            fs::remove_dir_all(&local_epoch_dir_path)?;
672        }
673        fs::create_dir_all(&local_epoch_dir_path)?;
674        Ok(())
675    }
676
677    /// Syncs a file from local store to remote store and removes the local file
678    async fn sync_file_to_remote(
679        local_path: PathBuf,
680        path: Path,
681        from: Arc<DynObjectStore>,
682        to: Arc<DynObjectStore>,
683    ) -> Result<()> {
684        debug!("Syncing snapshot file to remote: {:?}", path);
685        copy_file(&path, &path, &from, &to).await?;
686        fs::remove_file(path_to_filesystem(local_path, &path)?)?;
687        Ok(())
688    }
689}