1#![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
55struct 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 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 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 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 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 fn finalize_obj(&mut self) -> Result<()> {
148 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 fn finalize_ref(&mut self) -> Result<()> {
173 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 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 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 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 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
270pub struct StateSnapshotWriterV1 {
273 local_staging_dir: PathBuf,
274 file_compression: FileCompression,
275 remote_object_store: Arc<DynObjectStore>,
276 local_staging_store: Arc<DynObjectStore>,
277 checkpoint_store: Arc<CheckpointStore>,
280 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 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 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 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 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 write_handler
380 .await?
381 .context(format!("Failed to write state snapshot for epoch: {epoch}"))?;
382
383 upload_handle.await?.context(format!(
385 "Failed to upload state snapshot for epoch: {epoch}"
386 ))?;
387
388 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 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 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 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 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 for (_, writer) in object_writers.into_iter() {
486 files.extend(writer.done()?);
487 }
488 let epoch_info_metadata = self.write_epoch_info(epoch, &local_staging_dir_path, &sender)?;
493 files.push(epoch_info_metadata);
494 self.write_manifest(epoch, files)?;
496 Ok(())
497 }
498
499 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 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 for epoch_id in 0..=epoch {
540 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 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 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 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 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 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 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 1u32
656 }
657
658 fn epoch_dir(&self, epoch: u64) -> Path {
659 Path::from(format!("epoch_{epoch}"))
660 }
661
662 async fn setup_epoch_dir(&self, epoch: u64) -> Result<()> {
665 let epoch_dir = self.epoch_dir(epoch);
666 delete_recursively(&epoch_dir, &self.remote_object_store, self.concurrency).await?;
668 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 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}