1mod block_sealing;
8pub mod direct_io_file_wrapper;
9pub mod farming;
10pub mod identity;
11mod metrics;
12pub mod piece_cache;
13pub mod piece_reader;
14pub mod plot_cache;
15mod plotted_sectors;
16mod plotting;
17
18use crate::disk_piece_cache::{DiskPieceCache, DiskPieceCacheError};
19use crate::farm::{
20 Farm, FarmId, FarmingError, FarmingNotification, HandlerFn, PieceCacheId, PieceReader,
21 PlottedSectors, SectorUpdate,
22};
23use crate::node_client::NodeClient;
24use crate::plotter::Plotter;
25use crate::single_disk_farm::block_sealing::block_sealing;
26use crate::single_disk_farm::direct_io_file_wrapper::{DISK_PAGE_SIZE, DirectIoFileWrapper};
27use crate::single_disk_farm::farming::rayon_files::RayonFiles;
28use crate::single_disk_farm::farming::{
29 FarmingOptions, PlotAudit, farming, slot_notification_forwarder,
30};
31use crate::single_disk_farm::identity::{Identity, IdentityError};
32use crate::single_disk_farm::metrics::SingleDiskFarmMetrics;
33use crate::single_disk_farm::piece_cache::SingleDiskPieceCache;
34use crate::single_disk_farm::piece_reader::DiskPieceReader;
35use crate::single_disk_farm::plot_cache::DiskPlotCache;
36use crate::single_disk_farm::plotted_sectors::SingleDiskPlottedSectors;
37pub use crate::single_disk_farm::plotting::PlottingError;
38use crate::single_disk_farm::plotting::{
39 PlottingOptions, PlottingSchedulerOptions, SectorPlottingOptions, plotting, plotting_scheduler,
40};
41use crate::utils::{AsyncJoinOnDrop, tokio_rayon_spawn_handler};
42use crate::{KNOWN_PEERS_CACHE_SIZE, farm};
43use ab_core_primitives::address::Address;
44use ab_core_primitives::block::BlockRoot;
45use ab_core_primitives::ed25519::Ed25519PublicKey;
46use ab_core_primitives::hashes::Blake3Hash;
47use ab_core_primitives::pieces::Record;
48use ab_core_primitives::sectors::SectorIndex;
49use ab_core_primitives::segments::{HistorySize, SegmentIndex};
50use ab_erasure_coding::ErasureCoding;
51use ab_farmer_components::FarmerProtocolInfo;
52use ab_farmer_components::file_ext::FileExt;
53use ab_farmer_components::sector::{SectorMetadata, SectorMetadataChecksummed, sector_size};
54use ab_farmer_components::shard_commitment::ShardCommitmentsRootsCache;
55use ab_farmer_rpc_primitives::{FarmerAppInfo, SolutionResponse};
56use ab_networking::KnownPeersManager;
57use ab_proof_of_space::Table;
58use async_lock::{Mutex as AsyncMutex, RwLock as AsyncRwLock};
59use async_trait::async_trait;
60use bytesize::ByteSize;
61use event_listener_primitives::{Bag, HandlerId};
62use futures::channel::{mpsc, oneshot};
63use futures::stream::FuturesUnordered;
64use futures::{FutureExt, StreamExt, select};
65use parity_scale_codec::{Decode, Encode};
66use parking_lot::Mutex;
67use prometheus_client::registry::Registry;
68use rayon::prelude::*;
69use rayon::{ThreadPoolBuildError, ThreadPoolBuilder};
70use serde::{Deserialize, Serialize};
71use std::collections::HashSet;
72use std::fs::{File, OpenOptions};
73use std::future::Future;
74use std::io::Write;
75use std::num::{NonZeroU32, NonZeroUsize};
76use std::path::{Path, PathBuf};
77use std::pin::Pin;
78use std::str::FromStr;
79use std::sync::Arc;
80use std::sync::atomic::{AtomicUsize, Ordering};
81use std::time::Duration;
82use std::{fmt, fs, io};
83use thiserror::Error;
84use tokio::runtime::Handle;
85use tokio::sync::broadcast;
86use tokio::task;
87use tracing::{Instrument, Span, error, info, trace, warn};
88
89const {
92 assert!(size_of::<usize>() >= size_of::<u64>());
93}
94
95const RESERVED_PLOT_METADATA: u64 = 1024 * 1024;
97const RESERVED_FARM_INFO: u64 = 1024 * 1024;
99const NEW_SEGMENT_PROCESSING_DELAY: Duration = Duration::from_mins(10);
100
101#[derive(Debug)]
104#[must_use = "Lock file must be kept around or as long as farm is used"]
105pub struct SingleDiskFarmInfoLock {
106 _file: File,
107}
108
109#[derive(Debug, Copy, Clone, Serialize, Deserialize)]
111#[serde(rename_all = "camelCase")]
112pub enum SingleDiskFarmInfo {
113 #[serde(rename_all = "camelCase")]
115 V0 {
116 id: FarmId,
118 genesis_root: BlockRoot,
120 public_key: Ed25519PublicKey,
122 shard_commitments_seed: Blake3Hash,
124 pieces_in_sector: u16,
126 allocated_space: u64,
128 },
129}
130
131impl SingleDiskFarmInfo {
132 const FILE_NAME: &'static str = "single_disk_farm.json";
133
134 pub fn new(
136 id: FarmId,
137 genesis_root: BlockRoot,
138 public_key: Ed25519PublicKey,
139 shard_commitments_seed: Blake3Hash,
140 pieces_in_sector: u16,
141 allocated_space: u64,
142 ) -> Self {
143 Self::V0 {
144 id,
145 genesis_root,
146 public_key,
147 shard_commitments_seed,
148 pieces_in_sector,
149 allocated_space,
150 }
151 }
152
153 pub fn load_from(directory: &Path) -> io::Result<Option<Self>> {
156 let bytes = match fs::read(directory.join(Self::FILE_NAME)) {
157 Ok(bytes) => bytes,
158 Err(error) => {
159 return if error.kind() == io::ErrorKind::NotFound {
160 Ok(None)
161 } else {
162 Err(error)
163 };
164 }
165 };
166
167 serde_json::from_slice(&bytes)
168 .map(Some)
169 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))
170 }
171
172 pub fn store_to(
176 &self,
177 directory: &Path,
178 lock: bool,
179 ) -> io::Result<Option<SingleDiskFarmInfoLock>> {
180 let mut file = OpenOptions::new()
181 .write(true)
182 .create(true)
183 .truncate(false)
184 .open(directory.join(Self::FILE_NAME))?;
185 if lock {
186 fs4::FileExt::try_lock(&file)?;
187 }
188 file.set_len(0)?;
189 file.write_all(&serde_json::to_vec(self).expect("Info serialization never fails; qed"))?;
190
191 Ok(lock.then_some(SingleDiskFarmInfoLock { _file: file }))
192 }
193
194 pub fn try_lock(directory: &Path) -> io::Result<SingleDiskFarmInfoLock> {
197 let file = File::open(directory.join(Self::FILE_NAME))?;
198 fs4::FileExt::try_lock(&file)?;
199
200 Ok(SingleDiskFarmInfoLock { _file: file })
201 }
202
203 #[expect(
205 clippy::rest_pattern_accessible_field,
206 reason = "Do not need other fields"
207 )]
208 pub fn id(&self) -> &FarmId {
209 let Self::V0 { id, .. } = self;
210 id
211 }
212
213 #[expect(
215 clippy::rest_pattern_accessible_field,
216 reason = "Do not need other fields"
217 )]
218 pub fn genesis_root(&self) -> &BlockRoot {
219 let Self::V0 { genesis_root, .. } = self;
220 genesis_root
221 }
222
223 #[expect(
225 clippy::rest_pattern_accessible_field,
226 reason = "Do not need other fields"
227 )]
228 pub fn public_key(&self) -> &Ed25519PublicKey {
229 let Self::V0 { public_key, .. } = self;
230 public_key
231 }
232
233 #[expect(
235 clippy::rest_pattern_accessible_field,
236 reason = "Do not need other fields"
237 )]
238 pub fn shard_commitments_seed(&self) -> &Blake3Hash {
239 let Self::V0 {
240 shard_commitments_seed,
241 ..
242 } = self;
243 shard_commitments_seed
244 }
245
246 #[expect(
248 clippy::rest_pattern_accessible_field,
249 reason = "Do not need other fields"
250 )]
251 pub fn pieces_in_sector(&self) -> u16 {
252 match self {
253 SingleDiskFarmInfo::V0 {
254 pieces_in_sector, ..
255 } => *pieces_in_sector,
256 }
257 }
258
259 #[expect(
261 clippy::rest_pattern_accessible_field,
262 reason = "Do not need other fields"
263 )]
264 pub fn allocated_space(&self) -> u64 {
265 match self {
266 SingleDiskFarmInfo::V0 {
267 allocated_space, ..
268 } => *allocated_space,
269 }
270 }
271}
272
273#[derive(Debug)]
275pub enum SingleDiskFarmSummary {
276 Found {
278 info: SingleDiskFarmInfo,
280 directory: PathBuf,
282 },
283 NotFound {
285 directory: PathBuf,
287 },
288 Error {
290 directory: PathBuf,
292 error: io::Error,
294 },
295}
296
297#[derive(Debug, Encode, Decode)]
298struct PlotMetadataHeader {
299 version: u8,
300 plotted_sector_count: u16,
301}
302
303impl PlotMetadataHeader {
304 #[inline]
305 fn encoded_size() -> usize {
306 let default = PlotMetadataHeader {
307 version: 0,
308 plotted_sector_count: 0,
309 };
310
311 default.encoded_size()
312 }
313}
314
315#[derive(Debug)]
317pub struct SingleDiskFarmOptions<'a, NC>
318where
319 NC: Clone,
320{
321 pub directory: PathBuf,
323 pub farmer_app_info: FarmerAppInfo,
325 pub allocated_space: u64,
327 pub max_pieces_in_sector: u16,
329 pub node_client: NC,
331 pub reward_address: Address,
333 pub plotter: Arc<dyn Plotter + Send + Sync>,
335 pub erasure_coding: ErasureCoding,
337 pub cache_percentage: u8,
339 pub farming_thread_pool_size: usize,
342 pub plotting_delay: Option<oneshot::Receiver<()>>,
345 pub global_mutex: Arc<AsyncMutex<()>>,
349 pub max_plotting_sectors_per_farm: NonZeroUsize,
351 pub disable_farm_locking: bool,
353 pub registry: Option<&'a Mutex<&'a mut Registry>>,
355 pub create: bool,
357}
358
359#[derive(Debug, Error)]
361pub enum SingleDiskFarmError {
362 #[error("Failed to open or create identity: {0}")]
364 FailedToOpenIdentity(#[from] IdentityError),
365 #[error("Farm is likely already in use, make sure no other farmer is using it: {0}")]
367 LikelyAlreadyInUse(io::Error),
368 #[error("Single disk farm I/O error: {0}")]
370 Io(#[from] io::Error),
371 #[error("Failed to spawn task for blocking thread: {0}")]
373 TokioJoinError(#[from] task::JoinError),
374 #[error("Piece cache error: {0}")]
376 PieceCacheError(#[from] DiskPieceCacheError),
377 #[error("Can't preallocate metadata file, probably not enough space on disk: {0}")]
379 CantPreallocateMetadataFile(io::Error),
380 #[error("Can't preallocate plot file, probably not enough space on disk: {0}")]
382 CantPreallocatePlotFile(io::Error),
383 #[error(
385 "Genesis hash of farm {id} {wrong_chain} is different from {correct_chain} when farm was \
386 created, it is not possible to use farm on a different chain"
387 )]
388 WrongChain {
389 id: FarmId,
391 correct_chain: String,
394 wrong_chain: String,
396 },
397 #[error(
399 "Public key of farm {id} {wrong_public_key} is different from {correct_public_key} when \
400 farm was created, something went wrong, likely due to manual edits"
401 )]
402 IdentityMismatch {
403 id: FarmId,
405 correct_public_key: Ed25519PublicKey,
407 wrong_public_key: Ed25519PublicKey,
409 },
410 #[error(
412 "Invalid number pieces in sector: max supported {max_supported}, farm initialized with \
413 {initialized_with}"
414 )]
415 InvalidPiecesInSector {
416 id: FarmId,
418 max_supported: u16,
420 initialized_with: u16,
422 },
423 #[error("Failed to decode metadata header: {0}")]
425 FailedToDecodeMetadataHeader(parity_scale_codec::Error),
426 #[error("Unexpected metadata version {0}")]
428 UnexpectedMetadataVersion(u8),
429 #[error(
431 "Allocated space is not enough for one sector. \
432 The lowest acceptable value for allocated space is {min_space} bytes, \
433 provided {allocated_space} bytes."
434 )]
435 InsufficientAllocatedSpace {
436 min_space: u64,
438 allocated_space: u64,
440 },
441 #[error(
443 "Farm is too large: allocated {allocated_sectors} sectors ({allocated_space} bytes), max \
444 supported is {max_sectors} ({max_space} bytes). Consider creating multiple smaller farms \
445 instead."
446 )]
447 FarmTooLarge {
448 allocated_space: u64,
450 allocated_sectors: u64,
452 max_space: u64,
454 max_sectors: u16,
456 },
457 #[error("Failed to create thread pool: {0}")]
459 FailedToCreateThreadPool(ThreadPoolBuildError),
460}
461
462#[derive(Debug, Error)]
464pub enum SingleDiskFarmScrubError {
465 #[error("Farm is likely already in use, make sure no other farmer is using it: {0}")]
467 LikelyAlreadyInUse(io::Error),
468 #[error("Failed to file size of {file}: {error}")]
470 FailedToDetermineFileSize {
471 file: PathBuf,
473 error: io::Error,
475 },
476 #[error("Failed to read {size} bytes from {file} at offset {offset}: {error}")]
478 FailedToReadBytes {
479 file: PathBuf,
481 size: u64,
483 offset: u64,
485 error: io::Error,
487 },
488 #[error("Failed to write {size} bytes from {file} at offset {offset}: {error}")]
490 FailedToWriteBytes {
491 file: PathBuf,
493 size: u64,
495 offset: u64,
497 error: io::Error,
499 },
500 #[error("Farm info file does not exist at {file}")]
502 FarmInfoFileDoesNotExist {
503 file: PathBuf,
505 },
506 #[error("Farm info at {file} can't be opened: {error}")]
508 FarmInfoCantBeOpened {
509 file: PathBuf,
511 error: io::Error,
513 },
514 #[error("Identity file does not exist at {file}")]
516 IdentityFileDoesNotExist {
517 file: PathBuf,
519 },
520 #[error("Identity at {file} can't be opened: {error}")]
522 IdentityCantBeOpened {
523 file: PathBuf,
525 error: IdentityError,
527 },
528 #[error("Identity public key {identity} doesn't match public key in the disk farm info {info}")]
530 PublicKeyMismatch {
531 identity: Ed25519PublicKey,
533 info: Ed25519PublicKey,
535 },
536 #[error("Metadata file does not exist at {file}")]
538 MetadataFileDoesNotExist {
539 file: PathBuf,
541 },
542 #[error("Metadata at {file} can't be opened: {error}")]
544 MetadataCantBeOpened {
545 file: PathBuf,
547 error: io::Error,
549 },
550 #[error(
552 "Metadata file at {file} is too small: reserved size is {reserved_size} bytes, file size \
553 is {size}"
554 )]
555 MetadataFileTooSmall {
556 file: PathBuf,
558 reserved_size: u64,
560 size: u64,
562 },
563 #[error("Failed to decode metadata header: {0}")]
565 FailedToDecodeMetadataHeader(parity_scale_codec::Error),
566 #[error("Unexpected metadata version {0}")]
568 UnexpectedMetadataVersion(u8),
569 #[error("Cache at {file} can't be opened: {error}")]
571 CacheCantBeOpened {
572 file: PathBuf,
574 error: io::Error,
576 },
577}
578
579#[derive(Debug, Error)]
581pub enum BackgroundTaskError {
582 #[error(transparent)]
584 Plotting(#[from] PlottingError),
585 #[error(transparent)]
587 Farming(#[from] FarmingError),
588 #[error(transparent)]
590 BlockSealing(#[from] anyhow::Error),
591 #[error("Background task {task} panicked")]
593 BackgroundTaskPanicked {
594 task: String,
596 },
597}
598
599type BackgroundTask = Pin<Box<dyn Future<Output = Result<(), BackgroundTaskError>> + Send>>;
600
601#[derive(Debug, Copy, Clone)]
603pub enum ScrubTarget {
604 All,
606 Metadata,
608 Plot,
610 Cache,
612}
613
614impl fmt::Display for ScrubTarget {
615 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
616 match self {
617 Self::All => f.write_str("all"),
618 Self::Metadata => f.write_str("metadata"),
619 Self::Plot => f.write_str("plot"),
620 Self::Cache => f.write_str("cache"),
621 }
622 }
623}
624
625impl FromStr for ScrubTarget {
626 type Err = String;
627
628 fn from_str(s: &str) -> Result<Self, Self::Err> {
629 match s {
630 "all" => Ok(Self::All),
631 "metadata" => Ok(Self::Metadata),
632 "plot" => Ok(Self::Plot),
633 "cache" => Ok(Self::Cache),
634 s => Err(format!("Can't parse {s} as `ScrubTarget`")),
635 }
636 }
637}
638
639impl ScrubTarget {
640 fn metadata(self) -> bool {
641 match self {
642 Self::All | Self::Metadata | Self::Plot => true,
643 Self::Cache => false,
644 }
645 }
646
647 fn plot(self) -> bool {
648 match self {
649 Self::All | Self::Plot => true,
650 Self::Metadata | Self::Cache => false,
651 }
652 }
653
654 fn cache(self) -> bool {
655 match self {
656 Self::All | Self::Cache => true,
657 Self::Metadata | Self::Plot => false,
658 }
659 }
660}
661
662struct AllocatedSpaceDistribution {
663 piece_cache_file_size: u64,
664 piece_cache_capacity: u32,
665 plot_file_size: u64,
666 target_sector_count: u16,
667 metadata_file_size: u64,
668}
669
670impl AllocatedSpaceDistribution {
671 fn new(
672 allocated_space: u64,
673 sector_size: u64,
674 cache_percentage: u8,
675 sector_metadata_size: u64,
676 ) -> Result<Self, SingleDiskFarmError> {
677 let single_sector_overhead = sector_size + sector_metadata_size;
678 let fixed_space_usage = RESERVED_PLOT_METADATA
680 + RESERVED_FARM_INFO
681 + Identity::file_size() as u64
682 + KnownPeersManager::file_size(KNOWN_PEERS_CACHE_SIZE) as u64;
683 let target_sector_count = {
685 let potentially_plottable_space = allocated_space.saturating_sub(fixed_space_usage)
686 / 100
687 * (100 - u64::from(cache_percentage));
688 (potentially_plottable_space - DISK_PAGE_SIZE as u64) / single_sector_overhead
691 };
692
693 if target_sector_count == 0 {
694 let mut single_plot_with_cache_space =
695 single_sector_overhead.div_ceil(100 - u64::from(cache_percentage)) * 100;
696 if single_plot_with_cache_space - single_sector_overhead
699 < u64::from(DiskPieceCache::element_size())
700 {
701 single_plot_with_cache_space =
702 single_sector_overhead + u64::from(DiskPieceCache::element_size());
703 }
704
705 return Err(SingleDiskFarmError::InsufficientAllocatedSpace {
706 min_space: fixed_space_usage + single_plot_with_cache_space,
707 allocated_space,
708 });
709 }
710 let plot_file_size = target_sector_count * sector_size;
711 let plot_file_size = plot_file_size.div_ceil(DISK_PAGE_SIZE as u64) * DISK_PAGE_SIZE as u64;
713
714 let piece_cache_capacity = if cache_percentage > 0 {
716 let cache_space = allocated_space
717 - fixed_space_usage
718 - plot_file_size
719 - (sector_metadata_size * target_sector_count);
720 (cache_space / u64::from(DiskPieceCache::element_size())) as u32
721 } else {
722 0
723 };
724 let target_sector_count = match u16::try_from(target_sector_count).map(SectorIndex::from) {
725 Ok(target_sector_count) if target_sector_count < SectorIndex::MAX => {
726 u16::from(target_sector_count)
727 }
728 _ => {
729 let max_sectors = u16::from(SectorIndex::MAX) - 1;
732 return Err(SingleDiskFarmError::FarmTooLarge {
733 allocated_space: target_sector_count * sector_size,
734 allocated_sectors: target_sector_count,
735 max_space: u64::from(max_sectors) * sector_size,
736 max_sectors,
737 });
738 }
739 };
740
741 Ok(Self {
742 piece_cache_file_size: u64::from(piece_cache_capacity)
743 * u64::from(DiskPieceCache::element_size()),
744 piece_cache_capacity,
745 plot_file_size,
746 target_sector_count,
747 metadata_file_size: RESERVED_PLOT_METADATA
748 + sector_metadata_size * u64::from(target_sector_count),
749 })
750 }
751}
752
753type Handler<A> = Bag<HandlerFn<A>, A>;
754
755#[derive(Default, Debug)]
756struct Handlers {
757 sector_update: Handler<(SectorIndex, SectorUpdate)>,
758 farming_notification: Handler<FarmingNotification>,
759 solution: Handler<SolutionResponse>,
760}
761
762struct SingleDiskFarmInit {
763 identity: Identity,
764 single_disk_farm_info: SingleDiskFarmInfo,
765 single_disk_farm_info_lock: Option<SingleDiskFarmInfoLock>,
766 plot_file: Arc<DirectIoFileWrapper>,
767 metadata_file: DirectIoFileWrapper,
768 metadata_header: PlotMetadataHeader,
769 target_sector_count: u16,
770 sectors_metadata: Arc<AsyncRwLock<Vec<SectorMetadataChecksummed>>>,
771 piece_cache_capacity: u32,
772 plot_cache: DiskPlotCache,
773}
774
775#[derive(Debug)]
780#[must_use = "Plot does not function properly unless run() method is called"]
781pub struct SingleDiskFarm {
782 farmer_protocol_info: FarmerProtocolInfo,
783 single_disk_farm_info: SingleDiskFarmInfo,
784 shard_commitments_roots_cache: ShardCommitmentsRootsCache,
785 sectors_metadata: Arc<AsyncRwLock<Vec<SectorMetadataChecksummed>>>,
787 pieces_in_sector: u16,
788 total_sectors_count: u16,
789 span: Span,
790 tasks: FuturesUnordered<BackgroundTask>,
791 handlers: Arc<Handlers>,
792 piece_cache: SingleDiskPieceCache,
793 plot_cache: DiskPlotCache,
794 piece_reader: DiskPieceReader,
795 start_sender: Option<broadcast::Sender<()>>,
797 stop_sender: Option<broadcast::Sender<()>>,
799 _single_disk_farm_info_lock: Option<SingleDiskFarmInfoLock>,
800}
801
802impl Drop for SingleDiskFarm {
803 #[inline]
804 fn drop(&mut self) {
805 self.piece_reader.close_all_readers();
806 self.start_sender.take();
808 self.stop_sender.take();
810 }
811}
812
813#[async_trait(?Send)]
814impl Farm for SingleDiskFarm {
815 fn id(&self) -> &FarmId {
816 self.id()
817 }
818
819 fn total_sectors_count(&self) -> u16 {
820 self.total_sectors_count
821 }
822
823 fn plotted_sectors(&self) -> Arc<dyn PlottedSectors + 'static> {
824 Arc::new(self.plotted_sectors())
825 }
826
827 fn piece_reader(&self) -> Arc<dyn PieceReader + 'static> {
828 Arc::new(self.piece_reader())
829 }
830
831 fn on_sector_update(
832 &self,
833 callback: HandlerFn<(SectorIndex, SectorUpdate)>,
834 ) -> Box<dyn farm::HandlerId> {
835 Box::new(self.on_sector_update(callback))
836 }
837
838 fn on_farming_notification(
839 &self,
840 callback: HandlerFn<FarmingNotification>,
841 ) -> Box<dyn farm::HandlerId> {
842 Box::new(self.on_farming_notification(callback))
843 }
844
845 fn on_solution(&self, callback: HandlerFn<SolutionResponse>) -> Box<dyn farm::HandlerId> {
846 Box::new(self.on_solution(callback))
847 }
848
849 fn run(self: Box<Self>) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send>> {
850 Box::pin((*self).run())
851 }
852}
853
854impl SingleDiskFarm {
855 pub const PLOT_FILE: &'static str = "plot.bin";
857 pub const METADATA_FILE: &'static str = "metadata.bin";
859 const SUPPORTED_PLOT_VERSION: u8 = 0;
860
861 pub async fn new<NC, PosTable>(
863 options: SingleDiskFarmOptions<'_, NC>,
864 farm_index: usize,
865 ) -> Result<Self, SingleDiskFarmError>
866 where
867 NC: NodeClient + Clone,
868 PosTable: Table,
869 {
870 let span = Span::current();
871
872 let SingleDiskFarmOptions {
873 directory,
874 farmer_app_info,
875 allocated_space,
876 max_pieces_in_sector,
877 node_client,
878 reward_address,
879 plotter,
880 erasure_coding,
881 cache_percentage,
882 farming_thread_pool_size,
883 plotting_delay,
884 global_mutex,
885 max_plotting_sectors_per_farm,
886 disable_farm_locking,
887 registry,
888 create,
889 } = options;
890
891 let single_disk_farm_init_fut = task::spawn_blocking({
892 let directory = directory.clone();
893 let farmer_app_info = farmer_app_info.clone();
894 let span = span.clone();
895
896 move || {
897 let _span_guard = span.enter();
898 Self::init(
899 &directory,
900 &farmer_app_info,
901 allocated_space,
902 max_pieces_in_sector,
903 cache_percentage,
904 disable_farm_locking,
905 create,
906 )
907 }
908 });
909
910 let single_disk_farm_init =
911 AsyncJoinOnDrop::new(single_disk_farm_init_fut, false).await??;
912
913 let SingleDiskFarmInit {
914 identity,
915 single_disk_farm_info,
916 single_disk_farm_info_lock,
917 plot_file,
918 metadata_file,
919 metadata_header,
920 target_sector_count,
921 sectors_metadata,
922 piece_cache_capacity,
923 plot_cache,
924 } = single_disk_farm_init;
925
926 let piece_cache = {
927 let FarmId::Ulid(id) = *single_disk_farm_info.id();
929 let id = PieceCacheId::Ulid(id);
930
931 SingleDiskPieceCache::new(
932 id,
933 if let Some(piece_cache_capacity) = NonZeroU32::new(piece_cache_capacity) {
934 Some(task::block_in_place(|| {
935 if let Some(registry) = registry {
936 DiskPieceCache::open(
937 &directory,
938 piece_cache_capacity,
939 Some(id),
940 Some(*registry.lock()),
941 )
942 } else {
943 DiskPieceCache::open(&directory, piece_cache_capacity, Some(id), None)
944 }
945 })?)
946 } else {
947 None
948 },
949 )
950 };
951
952 let public_key = *single_disk_farm_info.public_key();
953 let shard_commitments_roots_cache =
954 ShardCommitmentsRootsCache::new(*single_disk_farm_info.shard_commitments_seed());
955 let pieces_in_sector = single_disk_farm_info.pieces_in_sector();
956 let sector_size = sector_size(pieces_in_sector);
957
958 let metrics = registry.map(|registry| {
959 Arc::new(SingleDiskFarmMetrics::new(
960 *registry.lock(),
961 single_disk_farm_info.id(),
962 target_sector_count,
963 sectors_metadata.read_blocking().len() as u16,
964 ))
965 });
966
967 let (error_sender, error_receiver) = oneshot::channel();
968 let error_sender = Arc::new(Mutex::new(Some(error_sender)));
969
970 let tasks = FuturesUnordered::<BackgroundTask>::new();
971
972 tasks.push(Box::pin(async move {
973 if let Ok(error) = error_receiver.await {
974 return Err(error);
975 }
976
977 Ok(())
978 }));
979
980 let handlers = Arc::<Handlers>::default();
981 let (start_sender, mut start_receiver) = broadcast::channel::<()>(1);
982 let (stop_sender, mut stop_receiver) = broadcast::channel::<()>(1);
983 let sectors_being_modified = Arc::<AsyncRwLock<HashSet<SectorIndex>>>::default();
984 let (sectors_to_plot_sender, sectors_to_plot_receiver) = mpsc::channel(1);
985 let sectors_indices_left_to_plot = SectorIndex::from(metadata_header.plotted_sector_count)
987 ..SectorIndex::from(target_sector_count);
988
989 let farming_thread_pool = ThreadPoolBuilder::new()
990 .thread_name(move |thread_index| format!("farming-{farm_index}.{thread_index}"))
991 .num_threads(farming_thread_pool_size)
992 .spawn_handler(tokio_rayon_spawn_handler())
993 .build()
994 .map_err(SingleDiskFarmError::FailedToCreateThreadPool)?;
995 let farming_plot_fut = task::spawn_blocking(|| {
996 farming_thread_pool
997 .install(move || {
998 RayonFiles::open_with(directory.join(Self::PLOT_FILE), |path| {
999 DirectIoFileWrapper::open(path)
1000 })
1001 })
1002 .map(|farming_plot| (farming_plot, farming_thread_pool))
1003 });
1004
1005 let (farming_plot, farming_thread_pool) =
1006 AsyncJoinOnDrop::new(farming_plot_fut, false).await??;
1007
1008 let plotting_join_handle = task::spawn_blocking({
1009 let shard_commitments_roots_cache = shard_commitments_roots_cache.clone();
1010 let sectors_metadata = Arc::clone(§ors_metadata);
1011 let handlers = Arc::clone(&handlers);
1012 let sectors_being_modified = Arc::clone(§ors_being_modified);
1013 let node_client = node_client.clone();
1014 let plot_file = Arc::clone(&plot_file);
1015 let error_sender = Arc::clone(&error_sender);
1016 let span = span.clone();
1017 let global_mutex = Arc::clone(&global_mutex);
1018 let metrics = metrics.clone();
1019
1020 move || {
1021 let _span_guard = span.enter();
1022
1023 let plotting_options = PlottingOptions {
1024 metadata_header,
1025 sectors_metadata: §ors_metadata,
1026 sectors_being_modified: §ors_being_modified,
1027 sectors_to_plot_receiver,
1028 sector_plotting_options: SectorPlottingOptions {
1029 public_key,
1030 shard_commitments_roots_cache,
1031 node_client: &node_client,
1032 pieces_in_sector,
1033 sector_size,
1034 plot_file,
1035 metadata_file: Arc::new(metadata_file),
1036 handlers: &handlers,
1037 global_mutex: &global_mutex,
1038 plotter,
1039 metrics,
1040 },
1041 max_plotting_sectors_per_farm,
1042 };
1043
1044 let plotting_fut = async {
1045 if start_receiver.recv().await.is_err() {
1046 return Ok(());
1048 }
1049
1050 if let Some(plotting_delay) = plotting_delay
1051 && plotting_delay.await.is_err()
1052 {
1053 return Ok(());
1055 }
1056
1057 plotting(plotting_options).await
1058 };
1059
1060 Handle::current().block_on(async {
1061 select! {
1062 plotting_result = plotting_fut.fuse() => {
1063 if let Err(error) = plotting_result
1064 && let Some(error_sender) = error_sender.lock().take()
1065 && let Err(error) = error_sender.send(error.into())
1066 {
1067 error!(
1068 %error,
1069 "Plotting failed to send error to background task"
1070 );
1071 }
1072 }
1073 _ = stop_receiver.recv().fuse() => {
1074 }
1076 }
1077 });
1078 }
1079 });
1080 let plotting_join_handle = AsyncJoinOnDrop::new(plotting_join_handle, false);
1081
1082 tasks.push(Box::pin(async move {
1083 plotting_join_handle.await.map_err(|_error| {
1085 BackgroundTaskError::BackgroundTaskPanicked {
1086 task: format!("plotting-{farm_index}"),
1087 }
1088 })
1089 }));
1090
1091 let public_key_hash = public_key.hash();
1092
1093 let plotting_scheduler_options = PlottingSchedulerOptions {
1094 public_key_hash,
1095 shard_commitments_roots_cache: shard_commitments_roots_cache.clone(),
1096 sectors_indices_left_to_plot,
1097 target_sector_count,
1098 last_archived_segment_index: farmer_app_info.protocol_info.history_size.segment_index(),
1099 min_sector_lifetime: farmer_app_info.protocol_info.min_sector_lifetime,
1100 node_client: node_client.clone(),
1101 handlers: Arc::clone(&handlers),
1102 sectors_metadata: Arc::clone(§ors_metadata),
1103 sectors_to_plot_sender,
1104 new_segment_processing_delay: NEW_SEGMENT_PROCESSING_DELAY,
1105 metrics: metrics.clone(),
1106 };
1107 tasks.push(Box::pin(plotting_scheduler(plotting_scheduler_options)));
1108
1109 let (slot_info_forwarder_sender, slot_info_forwarder_receiver) = mpsc::channel(0);
1110
1111 tasks.push(Box::pin({
1112 let node_client = node_client.clone();
1113 let metrics = metrics.clone();
1114
1115 async move {
1116 slot_notification_forwarder(&node_client, slot_info_forwarder_sender, metrics)
1117 .await
1118 .map_err(BackgroundTaskError::Farming)
1119 }
1120 }));
1121
1122 let farming_join_handle = task::spawn_blocking({
1123 let shard_commitments_roots_cache = shard_commitments_roots_cache.clone();
1124 let erasure_coding = erasure_coding.clone();
1125 let handlers = Arc::clone(&handlers);
1126 let sectors_being_modified = Arc::clone(§ors_being_modified);
1127 let sectors_metadata = Arc::clone(§ors_metadata);
1128 let mut start_receiver = start_sender.subscribe();
1129 let mut stop_receiver = stop_sender.subscribe();
1130 let node_client = node_client.clone();
1131 let span = span.clone();
1132 let global_mutex = Arc::clone(&global_mutex);
1133
1134 move || {
1135 let _span_guard = span.enter();
1136
1137 let farming_fut = async move {
1138 if start_receiver.recv().await.is_err() {
1139 return Ok(());
1141 }
1142
1143 let plot_audit = PlotAudit::new(&farming_plot);
1144
1145 let farming_options = FarmingOptions {
1146 public_key_hash,
1147 shard_commitments_roots_cache,
1148 reward_address,
1149 node_client,
1150 plot_audit,
1151 sectors_metadata,
1152 erasure_coding,
1153 handlers,
1154 sectors_being_modified,
1155 slot_info_notifications: slot_info_forwarder_receiver,
1156 thread_pool: farming_thread_pool,
1157 global_mutex,
1158 metrics,
1159 };
1160 farming::<PosTable, _, _>(farming_options).await
1161 };
1162
1163 Handle::current().block_on(async {
1164 select! {
1165 farming_result = farming_fut.fuse() => {
1166 if let Err(error) = farming_result
1167 && let Some(error_sender) = error_sender.lock().take()
1168 && let Err(error) = error_sender.send(error.into())
1169 {
1170 error!(
1171 %error,
1172 "Farming failed to send error to background task",
1173 );
1174 }
1175 }
1176 _ = stop_receiver.recv().fuse() => {
1177 }
1179 }
1180 });
1181 }
1182 });
1183 let farming_join_handle = AsyncJoinOnDrop::new(farming_join_handle, false);
1184
1185 tasks.push(Box::pin(async move {
1186 farming_join_handle.await.map_err(|_error| {
1188 BackgroundTaskError::BackgroundTaskPanicked {
1189 task: format!("farming-{farm_index}"),
1190 }
1191 })
1192 }));
1193
1194 let (piece_reader, reading_fut) = DiskPieceReader::new::<PosTable>(
1195 public_key_hash,
1196 shard_commitments_roots_cache.clone(),
1197 pieces_in_sector,
1198 plot_file,
1199 Arc::clone(§ors_metadata),
1200 erasure_coding,
1201 sectors_being_modified,
1202 global_mutex,
1203 );
1204
1205 let reading_join_handle = task::spawn_blocking({
1206 let mut stop_receiver = stop_sender.subscribe();
1207 let reading_fut = reading_fut.instrument(span.clone());
1208
1209 move || {
1210 Handle::current().block_on(async {
1211 select! {
1212 () = reading_fut.fuse() => {
1213 }
1215 _ = stop_receiver.recv().fuse() => {
1216 }
1218 }
1219 });
1220 }
1221 });
1222
1223 let reading_join_handle = AsyncJoinOnDrop::new(reading_join_handle, false);
1224
1225 tasks.push(Box::pin(async move {
1226 reading_join_handle.await.map_err(|_error| {
1228 BackgroundTaskError::BackgroundTaskPanicked {
1229 task: format!("reading-{farm_index}"),
1230 }
1231 })
1232 }));
1233
1234 tasks.push(Box::pin(async move {
1235 match block_sealing(node_client, identity).await {
1236 Ok(block_sealing_fut) => {
1237 block_sealing_fut.await;
1238 }
1239 Err(error) => {
1240 return Err(BackgroundTaskError::BlockSealing(anyhow::anyhow!(
1241 "Failed to subscribe to block sealing notifications: {error}"
1242 )));
1243 }
1244 }
1245
1246 Ok(())
1247 }));
1248
1249 let farm = Self {
1250 farmer_protocol_info: farmer_app_info.protocol_info,
1251 single_disk_farm_info,
1252 shard_commitments_roots_cache,
1253 sectors_metadata,
1254 pieces_in_sector,
1255 total_sectors_count: target_sector_count,
1256 span,
1257 tasks,
1258 handlers,
1259 piece_cache,
1260 plot_cache,
1261 piece_reader,
1262 start_sender: Some(start_sender),
1263 stop_sender: Some(stop_sender),
1264 _single_disk_farm_info_lock: single_disk_farm_info_lock,
1265 };
1266 Ok(farm)
1267 }
1268
1269 fn init(
1270 directory: &PathBuf,
1271 farmer_app_info: &FarmerAppInfo,
1272 allocated_space: u64,
1273 max_pieces_in_sector: u16,
1274 cache_percentage: u8,
1275 disable_farm_locking: bool,
1276 create: bool,
1277 ) -> Result<SingleDiskFarmInit, SingleDiskFarmError> {
1278 fs::create_dir_all(directory)?;
1279
1280 let identity = if create {
1281 Identity::open_or_create(directory)?
1282 } else {
1283 Identity::open(directory)?.ok_or_else(|| {
1284 IdentityError::Io(io::Error::new(
1285 io::ErrorKind::NotFound,
1286 "Farm does not exist and creation was explicitly disabled",
1287 ))
1288 })?
1289 };
1290 let public_key = identity.public_key();
1291
1292 let (single_disk_farm_info, single_disk_farm_info_lock) = if let Some(
1293 mut single_disk_farm_info,
1294 ) =
1295 SingleDiskFarmInfo::load_from(directory)?
1296 {
1297 if &farmer_app_info.genesis_root != single_disk_farm_info.genesis_root() {
1298 return Err(SingleDiskFarmError::WrongChain {
1299 id: *single_disk_farm_info.id(),
1300 correct_chain: hex::encode(single_disk_farm_info.genesis_root()),
1301 wrong_chain: hex::encode(farmer_app_info.genesis_root),
1302 });
1303 }
1304
1305 if &public_key != single_disk_farm_info.public_key() {
1306 return Err(SingleDiskFarmError::IdentityMismatch {
1307 id: *single_disk_farm_info.id(),
1308 correct_public_key: *single_disk_farm_info.public_key(),
1309 wrong_public_key: public_key,
1310 });
1311 }
1312
1313 let pieces_in_sector = single_disk_farm_info.pieces_in_sector();
1314
1315 if max_pieces_in_sector < pieces_in_sector {
1316 return Err(SingleDiskFarmError::InvalidPiecesInSector {
1317 id: *single_disk_farm_info.id(),
1318 max_supported: max_pieces_in_sector,
1319 initialized_with: pieces_in_sector,
1320 });
1321 }
1322
1323 if max_pieces_in_sector > pieces_in_sector {
1324 info!(
1325 pieces_in_sector,
1326 max_pieces_in_sector,
1327 "Farm initialized with smaller number of pieces in sector, farm needs \
1328 to be re-created for increase"
1329 );
1330 }
1331
1332 let mut single_disk_farm_info_lock = None;
1333
1334 if allocated_space != single_disk_farm_info.allocated_space() {
1335 info!(
1336 old_space = %ByteSize::b(single_disk_farm_info.allocated_space()).display().iec(),
1337 new_space = %ByteSize::b(allocated_space).display().iec(),
1338 "Farm size has changed"
1339 );
1340
1341 let new_allocated_space = allocated_space;
1342 #[expect(
1343 clippy::rest_pattern_accessible_field,
1344 reason = "Do not need other fields"
1345 )]
1346 match &mut single_disk_farm_info {
1347 SingleDiskFarmInfo::V0 {
1348 allocated_space, ..
1349 } => {
1350 *allocated_space = new_allocated_space;
1351 }
1352 }
1353
1354 single_disk_farm_info_lock =
1355 single_disk_farm_info.store_to(directory, !disable_farm_locking)?;
1356 } else if !disable_farm_locking {
1357 single_disk_farm_info_lock = Some(
1358 SingleDiskFarmInfo::try_lock(directory)
1359 .map_err(SingleDiskFarmError::LikelyAlreadyInUse)?,
1360 );
1361 }
1362
1363 (single_disk_farm_info, single_disk_farm_info_lock)
1364 } else {
1365 let single_disk_farm_info = SingleDiskFarmInfo::new(
1366 FarmId::new(),
1367 farmer_app_info.genesis_root,
1368 public_key,
1369 identity.shard_commitments_seed(),
1370 max_pieces_in_sector,
1371 allocated_space,
1372 );
1373
1374 let single_disk_farm_info_lock =
1375 single_disk_farm_info.store_to(directory, !disable_farm_locking)?;
1376
1377 (single_disk_farm_info, single_disk_farm_info_lock)
1378 };
1379
1380 let pieces_in_sector = single_disk_farm_info.pieces_in_sector();
1381 let sector_size = sector_size(pieces_in_sector) as u64;
1382 let sector_metadata_size = SectorMetadataChecksummed::encoded_size();
1383 let allocated_space_distribution = AllocatedSpaceDistribution::new(
1384 allocated_space,
1385 sector_size,
1386 cache_percentage,
1387 sector_metadata_size as u64,
1388 )?;
1389 let target_sector_count = allocated_space_distribution.target_sector_count;
1390
1391 let metadata_file_path = directory.join(Self::METADATA_FILE);
1392 let metadata_file = DirectIoFileWrapper::open(&metadata_file_path)?;
1393
1394 let metadata_size = metadata_file.size()?;
1395 let expected_metadata_size = allocated_space_distribution.metadata_file_size;
1396 let expected_metadata_size =
1398 expected_metadata_size.div_ceil(DISK_PAGE_SIZE as u64) * DISK_PAGE_SIZE as u64;
1399 let metadata_header = if metadata_size == 0 {
1400 let metadata_header = PlotMetadataHeader {
1401 version: SingleDiskFarm::SUPPORTED_PLOT_VERSION,
1402 plotted_sector_count: 0,
1403 };
1404
1405 metadata_file
1406 .preallocate(expected_metadata_size)
1407 .map_err(SingleDiskFarmError::CantPreallocateMetadataFile)?;
1408 metadata_file.write_all_at(metadata_header.encode().as_slice(), 0)?;
1409
1410 metadata_header
1411 } else {
1412 if metadata_size != expected_metadata_size {
1413 metadata_file
1416 .preallocate(expected_metadata_size)
1417 .map_err(SingleDiskFarmError::CantPreallocateMetadataFile)?;
1418 metadata_file.set_len(expected_metadata_size)?;
1420 }
1421
1422 let mut metadata_header_bytes = vec![0; PlotMetadataHeader::encoded_size()];
1423 metadata_file.read_exact_at(&mut metadata_header_bytes, 0)?;
1424
1425 let mut metadata_header =
1426 PlotMetadataHeader::decode(&mut metadata_header_bytes.as_ref())
1427 .map_err(SingleDiskFarmError::FailedToDecodeMetadataHeader)?;
1428
1429 if metadata_header.version != SingleDiskFarm::SUPPORTED_PLOT_VERSION {
1430 return Err(SingleDiskFarmError::UnexpectedMetadataVersion(
1431 metadata_header.version,
1432 ));
1433 }
1434
1435 if metadata_header.plotted_sector_count > target_sector_count {
1436 metadata_header.plotted_sector_count = target_sector_count;
1437 metadata_file.write_all_at(&metadata_header.encode(), 0)?;
1438 }
1439
1440 metadata_header
1441 };
1442
1443 let sectors_metadata = {
1444 let mut sectors_metadata =
1445 Vec::<SectorMetadataChecksummed>::with_capacity(usize::from(target_sector_count));
1446
1447 let mut sector_metadata_bytes = vec![0; sector_metadata_size];
1448 for sector_index in
1449 SectorIndex::ZERO..SectorIndex::from(metadata_header.plotted_sector_count)
1450 {
1451 let sector_offset =
1452 RESERVED_PLOT_METADATA + sector_metadata_size as u64 * u64::from(sector_index);
1453 metadata_file.read_exact_at(&mut sector_metadata_bytes, sector_offset)?;
1454
1455 let sector_metadata =
1456 match SectorMetadataChecksummed::decode(&mut sector_metadata_bytes.as_ref()) {
1457 Ok(sector_metadata) => sector_metadata,
1458 Err(error) => {
1459 warn!(
1460 path = %metadata_file_path.display(),
1461 %error,
1462 %sector_index,
1463 "Failed to decode sector metadata, replacing with dummy expired \
1464 sector metadata"
1465 );
1466
1467 let dummy_sector = SectorMetadataChecksummed::from(SectorMetadata {
1468 sector_index,
1469 pieces_in_sector,
1470 s_bucket_sizes: Box::new([0; Record::NUM_S_BUCKETS]),
1471 history_size: HistorySize::from(SegmentIndex::ZERO),
1472 });
1473 metadata_file.write_all_at(&dummy_sector.encode(), sector_offset)?;
1474
1475 dummy_sector
1476 }
1477 };
1478 sectors_metadata.push(sector_metadata);
1479 }
1480
1481 Arc::new(AsyncRwLock::new(sectors_metadata))
1482 };
1483
1484 let plot_file = DirectIoFileWrapper::open(directory.join(Self::PLOT_FILE))?;
1485
1486 if plot_file.size()? != allocated_space_distribution.plot_file_size {
1487 plot_file
1490 .preallocate(allocated_space_distribution.plot_file_size)
1491 .map_err(SingleDiskFarmError::CantPreallocatePlotFile)?;
1492 plot_file.set_len(allocated_space_distribution.plot_file_size)?;
1494 }
1495
1496 let plot_file = Arc::new(plot_file);
1497
1498 let plot_cache = DiskPlotCache::new(
1499 &plot_file,
1500 §ors_metadata,
1501 target_sector_count,
1502 sector_size,
1503 );
1504
1505 Ok(SingleDiskFarmInit {
1506 identity,
1507 single_disk_farm_info,
1508 single_disk_farm_info_lock,
1509 plot_file,
1510 metadata_file,
1511 metadata_header,
1512 target_sector_count,
1513 sectors_metadata,
1514 piece_cache_capacity: allocated_space_distribution.piece_cache_capacity,
1515 plot_cache,
1516 })
1517 }
1518
1519 pub fn collect_summary(directory: PathBuf) -> SingleDiskFarmSummary {
1521 let single_disk_farm_info = match SingleDiskFarmInfo::load_from(&directory) {
1522 Ok(Some(single_disk_farm_info)) => single_disk_farm_info,
1523 Ok(None) => {
1524 return SingleDiskFarmSummary::NotFound { directory };
1525 }
1526 Err(error) => {
1527 return SingleDiskFarmSummary::Error { directory, error };
1528 }
1529 };
1530
1531 SingleDiskFarmSummary::Found {
1532 info: single_disk_farm_info,
1533 directory,
1534 }
1535 }
1536
1537 pub fn effective_disk_usage(
1543 directory: &Path,
1544 cache_percentage: u8,
1545 ) -> Result<u64, SingleDiskFarmError> {
1546 let mut effective_disk_usage;
1547 match SingleDiskFarmInfo::load_from(directory)? {
1548 Some(single_disk_farm_info) => {
1549 let allocated_space_distribution = AllocatedSpaceDistribution::new(
1550 single_disk_farm_info.allocated_space(),
1551 sector_size(single_disk_farm_info.pieces_in_sector()) as u64,
1552 cache_percentage,
1553 SectorMetadataChecksummed::encoded_size() as u64,
1554 )?;
1555
1556 effective_disk_usage = single_disk_farm_info.allocated_space();
1557 effective_disk_usage -= Identity::file_size() as u64;
1558 effective_disk_usage -= allocated_space_distribution.metadata_file_size;
1559 effective_disk_usage -= allocated_space_distribution.plot_file_size;
1560 effective_disk_usage -= allocated_space_distribution.piece_cache_file_size;
1561 }
1562 None => {
1563 effective_disk_usage = 0;
1565 }
1566 }
1567
1568 if Identity::open(directory)?.is_some() {
1569 effective_disk_usage += Identity::file_size() as u64;
1570 }
1571
1572 match OpenOptions::new()
1573 .read(true)
1574 .open(directory.join(Self::METADATA_FILE))
1575 {
1576 Ok(metadata_file) => {
1577 effective_disk_usage += metadata_file.size()?;
1578 }
1579 Err(error) => {
1580 if error.kind() == io::ErrorKind::NotFound {
1581 } else {
1583 return Err(error.into());
1584 }
1585 }
1586 }
1587
1588 match OpenOptions::new()
1589 .read(true)
1590 .open(directory.join(Self::PLOT_FILE))
1591 {
1592 Ok(plot_file) => {
1593 effective_disk_usage += plot_file.size()?;
1594 }
1595 Err(error) => {
1596 if error.kind() == io::ErrorKind::NotFound {
1597 } else {
1599 return Err(error.into());
1600 }
1601 }
1602 }
1603
1604 match OpenOptions::new()
1605 .read(true)
1606 .open(directory.join(DiskPieceCache::FILE_NAME))
1607 {
1608 Ok(piece_cache) => {
1609 effective_disk_usage += piece_cache.size()?;
1610 }
1611 Err(error) => {
1612 if error.kind() == io::ErrorKind::NotFound {
1613 } else {
1615 return Err(error.into());
1616 }
1617 }
1618 }
1619
1620 Ok(effective_disk_usage)
1621 }
1622
1623 pub fn read_all_sectors_metadata(
1625 directory: &Path,
1626 ) -> io::Result<Vec<SectorMetadataChecksummed>> {
1627 let metadata_file = DirectIoFileWrapper::open(directory.join(Self::METADATA_FILE))?;
1628
1629 let metadata_size = metadata_file.size()?;
1630 let sector_metadata_size = SectorMetadataChecksummed::encoded_size();
1631
1632 let mut metadata_header_bytes = vec![0; PlotMetadataHeader::encoded_size()];
1633 metadata_file.read_exact_at(&mut metadata_header_bytes, 0)?;
1634
1635 let metadata_header = PlotMetadataHeader::decode(&mut metadata_header_bytes.as_ref())
1636 .map_err(|error| {
1637 io::Error::other(format!("Failed to decode metadata header: {error}"))
1638 })?;
1639
1640 if metadata_header.version != SingleDiskFarm::SUPPORTED_PLOT_VERSION {
1641 return Err(io::Error::other(format!(
1642 "Unsupported metadata version {}",
1643 metadata_header.version
1644 )));
1645 }
1646
1647 let mut sectors_metadata = Vec::<SectorMetadataChecksummed>::with_capacity(
1648 ((metadata_size - RESERVED_PLOT_METADATA) / sector_metadata_size as u64) as usize,
1649 );
1650
1651 let mut sector_metadata_bytes = vec![0; sector_metadata_size];
1652 for sector_index in 0..metadata_header.plotted_sector_count {
1653 metadata_file.read_exact_at(
1654 &mut sector_metadata_bytes,
1655 RESERVED_PLOT_METADATA + sector_metadata_size as u64 * u64::from(sector_index),
1656 )?;
1657 sectors_metadata.push(
1658 SectorMetadataChecksummed::decode(&mut sector_metadata_bytes.as_ref()).map_err(
1659 |error| io::Error::other(format!("Failed to decode sector metadata: {error}")),
1660 )?,
1661 );
1662 }
1663
1664 Ok(sectors_metadata)
1665 }
1666
1667 pub fn id(&self) -> &FarmId {
1669 self.single_disk_farm_info.id()
1670 }
1671
1672 pub fn info(&self) -> &SingleDiskFarmInfo {
1674 &self.single_disk_farm_info
1675 }
1676
1677 pub fn total_sectors_count(&self) -> u16 {
1679 self.total_sectors_count
1680 }
1681
1682 pub fn plotted_sectors(&self) -> SingleDiskPlottedSectors {
1684 SingleDiskPlottedSectors {
1685 public_key_hash: self.single_disk_farm_info.public_key().hash(),
1686 shard_commitments_roots_cache: self.shard_commitments_roots_cache.clone(),
1687 pieces_in_sector: self.pieces_in_sector,
1688 farmer_protocol_info: self.farmer_protocol_info,
1689 sectors_metadata: Arc::clone(&self.sectors_metadata),
1690 }
1691 }
1692
1693 pub fn piece_cache(&self) -> SingleDiskPieceCache {
1695 self.piece_cache.clone()
1696 }
1697
1698 pub fn plot_cache(&self) -> DiskPlotCache {
1700 self.plot_cache.clone()
1701 }
1702
1703 pub fn piece_reader(&self) -> DiskPieceReader {
1705 self.piece_reader.clone()
1706 }
1707
1708 pub fn on_sector_update(&self, callback: HandlerFn<(SectorIndex, SectorUpdate)>) -> HandlerId {
1710 self.handlers.sector_update.add(callback)
1711 }
1712
1713 pub fn on_farming_notification(&self, callback: HandlerFn<FarmingNotification>) -> HandlerId {
1715 self.handlers.farming_notification.add(callback)
1716 }
1717
1718 pub fn on_solution(&self, callback: HandlerFn<SolutionResponse>) -> HandlerId {
1720 self.handlers.solution.add(callback)
1721 }
1722
1723 pub async fn run(mut self) -> anyhow::Result<()> {
1725 if let Some(start_sender) = self.start_sender.take() {
1726 let _: Result<_, _> = start_sender.send(());
1728 }
1729
1730 while let Some(result) = self.tasks.next().instrument(self.span.clone()).await {
1731 result?;
1732 }
1733
1734 Ok(())
1735 }
1736
1737 pub fn wipe(directory: &Path) -> io::Result<()> {
1739 let single_disk_info_info_path = directory.join(SingleDiskFarmInfo::FILE_NAME);
1740 match SingleDiskFarmInfo::load_from(directory) {
1741 Ok(Some(single_disk_farm_info)) => {
1742 info!("Found single disk farm {}", single_disk_farm_info.id());
1743 }
1744 Ok(None) => {
1745 return Err(io::Error::new(
1746 io::ErrorKind::NotFound,
1747 format!(
1748 "Single disk farm info not found at {}",
1749 single_disk_info_info_path.display()
1750 ),
1751 ));
1752 }
1753 Err(error) => {
1754 warn!("Found unknown single disk farm: {}", error);
1755 }
1756 }
1757
1758 {
1759 let plot = directory.join(Self::PLOT_FILE);
1760 if plot.exists() {
1761 info!("Deleting plot file at {}", plot.display());
1762 fs::remove_file(plot)?;
1763 }
1764 }
1765 {
1766 let metadata = directory.join(Self::METADATA_FILE);
1767 if metadata.exists() {
1768 info!("Deleting metadata file at {}", metadata.display());
1769 fs::remove_file(metadata)?;
1770 }
1771 }
1772 {
1775 let identity = directory.join("identity.bin");
1776 if identity.exists() {
1777 info!("Deleting identity file at {}", identity.display());
1778 fs::remove_file(identity)?;
1779 }
1780 }
1781
1782 DiskPieceCache::wipe(directory)?;
1783
1784 info!(
1785 "Deleting info file at {}",
1786 single_disk_info_info_path.display()
1787 );
1788 fs::remove_file(single_disk_info_info_path)
1789 }
1790
1791 pub fn scrub(
1794 directory: &Path,
1795 disable_farm_locking: bool,
1796 target: ScrubTarget,
1797 dry_run: bool,
1798 ) -> Result<(), SingleDiskFarmScrubError> {
1799 let span = Span::current();
1800
1801 if dry_run {
1802 info!("Dry run is used, no changes will be written to disk");
1803 }
1804
1805 if target.metadata() || target.plot() {
1806 let info = {
1807 let file = directory.join(SingleDiskFarmInfo::FILE_NAME);
1808 info!(path = %file.display(), "Checking info file");
1809
1810 match SingleDiskFarmInfo::load_from(directory) {
1811 Ok(Some(info)) => info,
1812 Ok(None) => {
1813 return Err(SingleDiskFarmScrubError::FarmInfoFileDoesNotExist { file });
1814 }
1815 Err(error) => {
1816 return Err(SingleDiskFarmScrubError::FarmInfoCantBeOpened { file, error });
1817 }
1818 }
1819 };
1820
1821 let _single_disk_farm_info_lock = if disable_farm_locking {
1822 None
1823 } else {
1824 Some(
1825 SingleDiskFarmInfo::try_lock(directory)
1826 .map_err(SingleDiskFarmScrubError::LikelyAlreadyInUse)?,
1827 )
1828 };
1829
1830 let identity = {
1831 let file = directory.join(Identity::FILE_NAME);
1832 info!(path = %file.display(), "Checking identity file");
1833
1834 match Identity::open(directory) {
1835 Ok(Some(identity)) => identity,
1836 Ok(None) => {
1837 return Err(SingleDiskFarmScrubError::IdentityFileDoesNotExist { file });
1838 }
1839 Err(error) => {
1840 return Err(SingleDiskFarmScrubError::IdentityCantBeOpened { file, error });
1841 }
1842 }
1843 };
1844
1845 if &identity.public_key() != info.public_key() {
1846 return Err(SingleDiskFarmScrubError::PublicKeyMismatch {
1847 identity: identity.public_key(),
1848 info: *info.public_key(),
1849 });
1850 }
1851
1852 let sector_metadata_size = SectorMetadataChecksummed::encoded_size();
1853
1854 let metadata_file_path = directory.join(Self::METADATA_FILE);
1855 let (metadata_file, mut metadata_header) = {
1856 info!(path = %metadata_file_path.display(), "Checking metadata file");
1857
1858 let metadata_file = match OpenOptions::new()
1859 .read(true)
1860 .write(!dry_run)
1861 .open(&metadata_file_path)
1862 {
1863 Ok(metadata_file) => metadata_file,
1864 Err(error) => {
1865 return Err(if error.kind() == io::ErrorKind::NotFound {
1866 SingleDiskFarmScrubError::MetadataFileDoesNotExist {
1867 file: metadata_file_path,
1868 }
1869 } else {
1870 SingleDiskFarmScrubError::MetadataCantBeOpened {
1871 file: metadata_file_path,
1872 error,
1873 }
1874 });
1875 }
1876 };
1877
1878 let _: Result<(), _> = metadata_file.advise_sequential_access();
1880
1881 let metadata_size = match metadata_file.size() {
1882 Ok(metadata_size) => metadata_size,
1883 Err(error) => {
1884 return Err(SingleDiskFarmScrubError::FailedToDetermineFileSize {
1885 file: metadata_file_path,
1886 error,
1887 });
1888 }
1889 };
1890
1891 if metadata_size < RESERVED_PLOT_METADATA {
1892 return Err(SingleDiskFarmScrubError::MetadataFileTooSmall {
1893 file: metadata_file_path,
1894 reserved_size: RESERVED_PLOT_METADATA,
1895 size: metadata_size,
1896 });
1897 }
1898
1899 let mut metadata_header = {
1900 let mut reserved_metadata = vec![0; RESERVED_PLOT_METADATA as usize];
1901
1902 if let Err(error) = metadata_file.read_exact_at(&mut reserved_metadata, 0) {
1903 return Err(SingleDiskFarmScrubError::FailedToReadBytes {
1904 file: metadata_file_path,
1905 size: RESERVED_PLOT_METADATA,
1906 offset: 0,
1907 error,
1908 });
1909 }
1910
1911 PlotMetadataHeader::decode(&mut reserved_metadata.as_slice())
1912 .map_err(SingleDiskFarmScrubError::FailedToDecodeMetadataHeader)?
1913 };
1914
1915 if metadata_header.version != SingleDiskFarm::SUPPORTED_PLOT_VERSION {
1916 return Err(SingleDiskFarmScrubError::UnexpectedMetadataVersion(
1917 metadata_header.version,
1918 ));
1919 }
1920
1921 let plotted_sector_count = metadata_header.plotted_sector_count;
1922
1923 let expected_metadata_size = RESERVED_PLOT_METADATA
1924 + sector_metadata_size as u64 * u64::from(plotted_sector_count);
1925
1926 if metadata_size < expected_metadata_size {
1927 warn!(
1928 %metadata_size,
1929 %expected_metadata_size,
1930 "Metadata file size is smaller than expected, shrinking number of plotted \
1931 sectors to correct value"
1932 );
1933
1934 metadata_header.plotted_sector_count =
1935 ((metadata_size - RESERVED_PLOT_METADATA) / sector_metadata_size as u64)
1936 as u16;
1937 let metadata_header_bytes = metadata_header.encode();
1938
1939 if !dry_run
1940 && let Err(error) = metadata_file.write_all_at(&metadata_header_bytes, 0)
1941 {
1942 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
1943 file: metadata_file_path,
1944 size: metadata_header_bytes.len() as u64,
1945 offset: 0,
1946 error,
1947 });
1948 }
1949 }
1950
1951 (metadata_file, metadata_header)
1952 };
1953
1954 let pieces_in_sector = info.pieces_in_sector();
1955 let sector_size = sector_size(pieces_in_sector) as u64;
1956
1957 let plot_file_path = directory.join(Self::PLOT_FILE);
1958 let plot_file = {
1959 let plot_file_path = directory.join(Self::PLOT_FILE);
1960 info!(path = %plot_file_path.display(), "Checking plot file");
1961
1962 let plot_file = match OpenOptions::new()
1963 .read(true)
1964 .write(!dry_run)
1965 .open(&plot_file_path)
1966 {
1967 Ok(plot_file) => plot_file,
1968 Err(error) => {
1969 return Err(if error.kind() == io::ErrorKind::NotFound {
1970 SingleDiskFarmScrubError::MetadataFileDoesNotExist {
1971 file: plot_file_path,
1972 }
1973 } else {
1974 SingleDiskFarmScrubError::MetadataCantBeOpened {
1975 file: plot_file_path,
1976 error,
1977 }
1978 });
1979 }
1980 };
1981
1982 let _: Result<(), _> = plot_file.advise_sequential_access();
1984
1985 let plot_size = match plot_file.size() {
1986 Ok(metadata_size) => metadata_size,
1987 Err(error) => {
1988 return Err(SingleDiskFarmScrubError::FailedToDetermineFileSize {
1989 file: plot_file_path,
1990 error,
1991 });
1992 }
1993 };
1994
1995 let min_expected_plot_size =
1996 u64::from(metadata_header.plotted_sector_count) * sector_size;
1997 if plot_size < min_expected_plot_size {
1998 warn!(
1999 %plot_size,
2000 %min_expected_plot_size,
2001 "Plot file size is smaller than expected, shrinking number of plotted \
2002 sectors to correct value"
2003 );
2004
2005 metadata_header.plotted_sector_count = (plot_size / sector_size) as u16;
2006 let metadata_header_bytes = metadata_header.encode();
2007
2008 if !dry_run
2009 && let Err(error) = metadata_file.write_all_at(&metadata_header_bytes, 0)
2010 {
2011 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
2012 file: plot_file_path,
2013 size: metadata_header_bytes.len() as u64,
2014 offset: 0,
2015 error,
2016 });
2017 }
2018 }
2019
2020 plot_file
2021 };
2022
2023 let sector_bytes_range = 0..(sector_size as usize - Blake3Hash::SIZE);
2024
2025 info!("Checking sectors and corresponding metadata");
2026 (0..metadata_header.plotted_sector_count)
2027 .into_par_iter()
2028 .map(SectorIndex::from)
2029 .map_init(
2030 || vec![0u8; Record::SIZE],
2031 |scratch_buffer, sector_index| {
2032 let _span_guard = span.enter();
2033
2034 let offset = RESERVED_PLOT_METADATA
2035 + u64::from(sector_index) * sector_metadata_size as u64;
2036 if let Err(error) = metadata_file
2037 .read_exact_at(&mut scratch_buffer[..sector_metadata_size], offset)
2038 {
2039 warn!(
2040 path = %metadata_file_path.display(),
2041 %error,
2042 %offset,
2043 size = %sector_metadata_size,
2044 %sector_index,
2045 "Failed to read sector metadata, replacing with dummy expired \
2046 sector metadata"
2047 );
2048
2049 if !dry_run {
2050 write_dummy_sector_metadata(
2051 &metadata_file,
2052 &metadata_file_path,
2053 sector_index,
2054 pieces_in_sector,
2055 )?;
2056 }
2057 return Ok(());
2058 }
2059
2060 let sector_metadata = match SectorMetadataChecksummed::decode(
2061 &mut &scratch_buffer[..sector_metadata_size],
2062 ) {
2063 Ok(sector_metadata) => sector_metadata,
2064 Err(error) => {
2065 warn!(
2066 path = %metadata_file_path.display(),
2067 %error,
2068 %sector_index,
2069 "Failed to decode sector metadata, replacing with dummy \
2070 expired sector metadata"
2071 );
2072
2073 if !dry_run {
2074 write_dummy_sector_metadata(
2075 &metadata_file,
2076 &metadata_file_path,
2077 sector_index,
2078 pieces_in_sector,
2079 )?;
2080 }
2081 return Ok(());
2082 }
2083 };
2084
2085 if sector_metadata.sector_index != sector_index {
2086 warn!(
2087 path = %metadata_file_path.display(),
2088 %sector_index,
2089 found_sector_index = %sector_metadata.sector_index,
2090 "Sector index mismatch, replacing with dummy expired sector \
2091 metadata"
2092 );
2093
2094 if !dry_run {
2095 write_dummy_sector_metadata(
2096 &metadata_file,
2097 &metadata_file_path,
2098 sector_index,
2099 pieces_in_sector,
2100 )?;
2101 }
2102 return Ok(());
2103 }
2104
2105 if sector_metadata.pieces_in_sector != pieces_in_sector {
2106 warn!(
2107 path = %metadata_file_path.display(),
2108 %sector_index,
2109 %pieces_in_sector,
2110 found_pieces_in_sector = sector_metadata.pieces_in_sector,
2111 "Pieces in sector mismatch, replacing with dummy expired sector \
2112 metadata"
2113 );
2114
2115 if !dry_run {
2116 write_dummy_sector_metadata(
2117 &metadata_file,
2118 &metadata_file_path,
2119 sector_index,
2120 pieces_in_sector,
2121 )?;
2122 }
2123 return Ok(());
2124 }
2125
2126 if target.plot() {
2127 let mut hasher = blake3::Hasher::new();
2128 for offset_in_sector in
2130 sector_bytes_range.clone().step_by(scratch_buffer.len())
2131 {
2132 let offset =
2133 u64::from(sector_index) * sector_size + offset_in_sector as u64;
2134 let bytes_to_read = (offset_in_sector + scratch_buffer.len())
2135 .min(sector_bytes_range.end)
2136 - offset_in_sector;
2137
2138 let bytes = &mut scratch_buffer[..bytes_to_read];
2139
2140 if let Err(error) = plot_file.read_exact_at(bytes, offset) {
2141 warn!(
2142 path = %plot_file_path.display(),
2143 %error,
2144 %sector_index,
2145 %offset,
2146 size = %bytes.len() as u64,
2147 "Failed to read sector bytes"
2148 );
2149
2150 continue;
2151 }
2152
2153 hasher.update(bytes);
2154 }
2155
2156 let actual_checksum = *hasher.finalize().as_bytes();
2157 let mut expected_checksum = [0; Blake3Hash::SIZE];
2158 {
2159 let offset = u64::from(sector_index) * sector_size
2160 + sector_bytes_range.end as u64;
2161 if let Err(error) =
2162 plot_file.read_exact_at(&mut expected_checksum, offset)
2163 {
2164 warn!(
2165 path = %plot_file_path.display(),
2166 %error,
2167 %sector_index,
2168 %offset,
2169 size = %expected_checksum.len() as u64,
2170 "Failed to read sector checksum bytes"
2171 );
2172 }
2173 }
2174
2175 if actual_checksum != expected_checksum {
2177 warn!(
2178 path = %plot_file_path.display(),
2179 %sector_index,
2180 actual_checksum = %hex::encode(actual_checksum),
2181 expected_checksum = %hex::encode(expected_checksum),
2182 "Plotted sector checksum mismatch, replacing with dummy \
2183 expired sector"
2184 );
2185
2186 if !dry_run {
2187 write_dummy_sector_metadata(
2188 &metadata_file,
2189 &metadata_file_path,
2190 sector_index,
2191 pieces_in_sector,
2192 )?;
2193 }
2194
2195 scratch_buffer.fill(0);
2196
2197 hasher.reset();
2198 for offset_in_sector in
2200 sector_bytes_range.clone().step_by(scratch_buffer.len())
2201 {
2202 let offset = u64::from(sector_index) * sector_size
2203 + offset_in_sector as u64;
2204 let bytes_to_write = (offset_in_sector + scratch_buffer.len())
2205 .min(sector_bytes_range.end)
2206 - offset_in_sector;
2207 let bytes = &mut scratch_buffer[..bytes_to_write];
2208
2209 if !dry_run
2210 && let Err(error) = plot_file.write_all_at(bytes, offset)
2211 {
2212 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
2213 file: plot_file_path.clone(),
2214 size: scratch_buffer.len() as u64,
2215 offset,
2216 error,
2217 });
2218 }
2219
2220 hasher.update(bytes);
2221 }
2222 {
2224 let checksum = *hasher.finalize().as_bytes();
2225 let offset = u64::from(sector_index) * sector_size
2226 + sector_bytes_range.end as u64;
2227 if !dry_run
2228 && let Err(error) =
2229 plot_file.write_all_at(&checksum, offset)
2230 {
2231 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
2232 file: plot_file_path.clone(),
2233 size: checksum.len() as u64,
2234 offset,
2235 error,
2236 });
2237 }
2238 }
2239
2240 return Ok(());
2241 }
2242 }
2243
2244 trace!(%sector_index, "Sector is in good shape");
2245
2246 Ok(())
2247 },
2248 )
2249 .try_for_each({
2250 let span = &span;
2251 let checked_sectors = AtomicUsize::new(0);
2252
2253 move |result| {
2254 let _span_guard = span.enter();
2255
2256 let checked_sectors = checked_sectors.fetch_add(1, Ordering::Relaxed);
2257 if checked_sectors > 1 && checked_sectors.is_multiple_of(10) {
2258 info!(
2259 "Checked {}/{} sectors",
2260 checked_sectors, metadata_header.plotted_sector_count
2261 );
2262 }
2263
2264 result
2265 }
2266 })?;
2267 }
2268
2269 if target.cache() {
2270 Self::scrub_cache(directory, dry_run)?;
2271 }
2272
2273 info!("Farm check completed");
2274
2275 Ok(())
2276 }
2277
2278 fn scrub_cache(directory: &Path, dry_run: bool) -> Result<(), SingleDiskFarmScrubError> {
2279 let span = Span::current();
2280
2281 let file = directory.join(DiskPieceCache::FILE_NAME);
2282 info!(path = %file.display(), "Checking cache file");
2283
2284 let cache_file = match OpenOptions::new().read(true).write(!dry_run).open(&file) {
2285 Ok(plot_file) => plot_file,
2286 Err(error) => {
2287 return if error.kind() == io::ErrorKind::NotFound {
2288 warn!(
2289 file = %file.display(),
2290 "Cache file does not exist, this is expected in farming cluster"
2291 );
2292 Ok(())
2293 } else {
2294 Err(SingleDiskFarmScrubError::CacheCantBeOpened { file, error })
2295 };
2296 }
2297 };
2298
2299 let _: Result<(), _> = cache_file.advise_sequential_access();
2301
2302 let cache_size = match cache_file.size() {
2303 Ok(cache_size) => cache_size,
2304 Err(error) => {
2305 return Err(SingleDiskFarmScrubError::FailedToDetermineFileSize { file, error });
2306 }
2307 };
2308
2309 let element_size = DiskPieceCache::element_size();
2310 let number_of_cached_elements = cache_size / u64::from(element_size);
2311 let dummy_element = vec![0; element_size as usize];
2312 (0..number_of_cached_elements)
2313 .into_par_iter()
2314 .map_with(vec![0; element_size as usize], |element, cache_offset| {
2315 let _span_guard = span.enter();
2316
2317 let offset = cache_offset * u64::from(element_size);
2318 if let Err(error) = cache_file.read_exact_at(element, offset) {
2319 warn!(
2320 path = %file.display(),
2321 %cache_offset,
2322 size = %element.len() as u64,
2323 %offset,
2324 %error,
2325 "Failed to read cached piece, replacing with dummy element"
2326 );
2327
2328 if !dry_run && let Err(error) = cache_file.write_all_at(&dummy_element, offset)
2329 {
2330 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
2331 file: file.clone(),
2332 size: u64::from(element_size),
2333 offset,
2334 error,
2335 });
2336 }
2337
2338 return Ok(());
2339 }
2340
2341 let (index_and_piece_bytes, expected_checksum) =
2342 element.split_at(element_size as usize - Blake3Hash::SIZE);
2343 let actual_checksum = *blake3::hash(index_and_piece_bytes).as_bytes();
2344 if actual_checksum != expected_checksum && element != &dummy_element {
2345 warn!(
2346 %cache_offset,
2347 actual_checksum = %hex::encode(actual_checksum),
2348 expected_checksum = %hex::encode(expected_checksum),
2349 "Cached piece checksum mismatch, replacing with dummy element"
2350 );
2351
2352 if !dry_run && let Err(error) = cache_file.write_all_at(&dummy_element, offset)
2353 {
2354 return Err(SingleDiskFarmScrubError::FailedToWriteBytes {
2355 file: file.clone(),
2356 size: u64::from(element_size),
2357 offset,
2358 error,
2359 });
2360 }
2361
2362 return Ok(());
2363 }
2364
2365 Ok(())
2366 })
2367 .try_for_each({
2368 let span = &span;
2369 let checked_elements = AtomicUsize::new(0);
2370
2371 move |result| {
2372 let _span_guard = span.enter();
2373
2374 let checked_elements = checked_elements.fetch_add(1, Ordering::Relaxed);
2375 if checked_elements > 1 && checked_elements.is_multiple_of(1000) {
2376 info!(
2377 "Checked {}/{} cache elements",
2378 checked_elements, number_of_cached_elements
2379 );
2380 }
2381
2382 result
2383 }
2384 })?;
2385
2386 Ok(())
2387 }
2388}
2389
2390fn write_dummy_sector_metadata(
2391 metadata_file: &File,
2392 metadata_file_path: &Path,
2393 sector_index: SectorIndex,
2394 pieces_in_sector: u16,
2395) -> Result<(), SingleDiskFarmScrubError> {
2396 let dummy_sector_bytes = SectorMetadataChecksummed::from(SectorMetadata {
2397 sector_index,
2398 pieces_in_sector,
2399 s_bucket_sizes: Box::new([0; Record::NUM_S_BUCKETS]),
2400 history_size: HistorySize::from(SegmentIndex::ZERO),
2401 })
2402 .encode();
2403 let sector_offset = RESERVED_PLOT_METADATA
2404 + u64::from(sector_index) * SectorMetadataChecksummed::encoded_size() as u64;
2405 metadata_file
2406 .write_all_at(&dummy_sector_bytes, sector_offset)
2407 .map_err(|error| SingleDiskFarmScrubError::FailedToWriteBytes {
2408 file: metadata_file_path.to_path_buf(),
2409 size: dummy_sector_bytes.len() as u64,
2410 offset: sector_offset,
2411 error,
2412 })
2413}