Skip to main content

ab_farmer/
single_disk_farm.rs

1//! Primary [`Farm`] implementation that deals with hardware directly
2//!
3//! Single disk farm is an abstraction that contains an identity, associated plot with metadata and
4//! a small piece cache. It fully manages farming and plotting process, including listening to node
5//! notifications, producing solutions and sealing blocks.
6
7mod 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
89// Refuse to compile on non-64-bit platforms, offsets may fail on those when converting from u64 to
90// usize depending on chain parameters
91const {
92    assert!(size_of::<usize>() >= size_of::<u64>());
93}
94
95/// Reserve 1M of space for plot metadata (for potential future expansion)
96const RESERVED_PLOT_METADATA: u64 = 1024 * 1024;
97/// Reserve 1M of space for farm info (for potential future expansion)
98const RESERVED_FARM_INFO: u64 = 1024 * 1024;
99const NEW_SEGMENT_PROCESSING_DELAY: Duration = Duration::from_mins(10);
100
101/// Exclusive lock for single disk farm info file, ensuring no concurrent edits by cooperating
102/// processes is done
103#[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/// Important information about the contents of the `SingleDiskFarm`
110#[derive(Debug, Copy, Clone, Serialize, Deserialize)]
111#[serde(rename_all = "camelCase")]
112pub enum SingleDiskFarmInfo {
113    /// V0 of the info
114    #[serde(rename_all = "camelCase")]
115    V0 {
116        /// ID of the farm
117        id: FarmId,
118        /// Genesis root of the beacon chain used for farm creation
119        genesis_root: BlockRoot,
120        /// Public key of identity used for farm creation
121        public_key: Ed25519PublicKey,
122        /// Seed used for deriving shard commitments
123        shard_commitments_seed: Blake3Hash,
124        /// How many pieces does one sector contain.
125        pieces_in_sector: u16,
126        /// How much space in bytes is allocated for this farm
127        allocated_space: u64,
128    },
129}
130
131impl SingleDiskFarmInfo {
132    const FILE_NAME: &'static str = "single_disk_farm.json";
133
134    /// Create a new instance
135    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    /// Load `SingleDiskFarm` from path is supposed to be stored, `None` means no info file was
154    /// found, happens during first start.
155    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    /// Store `SingleDiskFarm` info to path, so it can be loaded again upon restart.
173    ///
174    /// Can optionally return a lock.
175    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    /// Try to acquire exclusive lock on the single disk farm info file, ensuring no concurrent
195    /// edits by cooperating processes is done
196    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    /// ID of the farm
204    #[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    /// Genesis hash of the chain used for farm creation
214    #[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    /// Public key of identity used for farm creation
224    #[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    /// Seed used for deriving shard commitments
234    #[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    /// How many pieces does one sector contain
247    #[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    /// How much space in bytes is allocated for this farm
260    #[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/// Summary of single disk farm for presentational purposes
274#[derive(Debug)]
275pub enum SingleDiskFarmSummary {
276    /// Farm was found and read successfully
277    Found {
278        /// Farm info
279        info: SingleDiskFarmInfo,
280        /// Path to directory where farm is stored.
281        directory: PathBuf,
282    },
283    /// Farm was not found
284    NotFound {
285        /// Path to directory where farm is stored.
286        directory: PathBuf,
287    },
288    /// Failed to open farm
289    Error {
290        /// Path to directory where farm is stored.
291        directory: PathBuf,
292        /// Error itself
293        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/// Options used to open single disk farm
316#[derive(Debug)]
317pub struct SingleDiskFarmOptions<'a, NC>
318where
319    NC: Clone,
320{
321    /// Path to directory where farm is stored.
322    pub directory: PathBuf,
323    /// Information necessary for farmer application
324    pub farmer_app_info: FarmerAppInfo,
325    /// The amount of space in bytes that was allocated
326    pub allocated_space: u64,
327    /// How many pieces one sector is supposed to contain (max)
328    pub max_pieces_in_sector: u16,
329    /// RPC client connected to the node
330    pub node_client: NC,
331    /// Address where farming rewards should go
332    pub reward_address: Address,
333    /// Plotter
334    pub plotter: Arc<dyn Plotter + Send + Sync>,
335    /// Erasure coding instance to use.
336    pub erasure_coding: ErasureCoding,
337    /// Percentage of allocated space dedicated for caching purposes
338    pub cache_percentage: u8,
339    /// Thread pool size used for farming (mostly for blocking I/O, but also for some
340    /// compute-intensive operations during proving)
341    pub farming_thread_pool_size: usize,
342    /// Notification for plotter to start, can be used to delay plotting until some initialization
343    /// has happened externally
344    pub plotting_delay: Option<oneshot::Receiver<()>>,
345    /// Global mutex that can restrict concurrency of resource-intensive operations and make sure
346    /// that those operations that are very sensitive (like proving) have all the resources
347    /// available to them for the highest probability of success
348    pub global_mutex: Arc<AsyncMutex<()>>,
349    /// How many sectors a will be plotted concurrently per farm
350    pub max_plotting_sectors_per_farm: NonZeroUsize,
351    /// Disable farm locking, for example if file system doesn't support it
352    pub disable_farm_locking: bool,
353    /// Prometheus registry
354    pub registry: Option<&'a Mutex<&'a mut Registry>>,
355    /// Whether to create a farm if it doesn't yet exist
356    pub create: bool,
357}
358
359/// Errors happening when trying to create/open single disk farm
360#[derive(Debug, Error)]
361pub enum SingleDiskFarmError {
362    /// Failed to open or create identity
363    #[error("Failed to open or create identity: {0}")]
364    FailedToOpenIdentity(#[from] IdentityError),
365    /// Farm is likely already in use, make sure no other farmer is using it
366    #[error("Farm is likely already in use, make sure no other farmer is using it: {0}")]
367    LikelyAlreadyInUse(io::Error),
368    /// I/O error occurred
369    #[error("Single disk farm I/O error: {0}")]
370    Io(#[from] io::Error),
371    /// Failed to spawn task for blocking thread
372    #[error("Failed to spawn task for blocking thread: {0}")]
373    TokioJoinError(#[from] task::JoinError),
374    /// Piece cache error
375    #[error("Piece cache error: {0}")]
376    PieceCacheError(#[from] DiskPieceCacheError),
377    /// Can't preallocate metadata file, probably not enough space on disk
378    #[error("Can't preallocate metadata file, probably not enough space on disk: {0}")]
379    CantPreallocateMetadataFile(io::Error),
380    /// Can't preallocate plot file, probably not enough space on disk
381    #[error("Can't preallocate plot file, probably not enough space on disk: {0}")]
382    CantPreallocatePlotFile(io::Error),
383    /// Wrong chain (genesis hash)
384    #[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        /// Farm ID
390        id: FarmId,
391        /// Hex-encoded genesis hash during farm creation
392        // TODO: Wrapper type with `Display` impl for genesis hash
393        correct_chain: String,
394        /// Hex-encoded current genesis hash
395        wrong_chain: String,
396    },
397    /// Public key in identity doesn't match metadata
398    #[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        /// Farm ID
404        id: FarmId,
405        /// Public key used during farm creation
406        correct_public_key: Ed25519PublicKey,
407        /// Current public key
408        wrong_public_key: Ed25519PublicKey,
409    },
410    /// Invalid number pieces in sector
411    #[error(
412        "Invalid number pieces in sector: max supported {max_supported}, farm initialized with \
413        {initialized_with}"
414    )]
415    InvalidPiecesInSector {
416        /// Farm ID
417        id: FarmId,
418        /// Max supported pieces in sector
419        max_supported: u16,
420        /// Number of pieces in sector farm is initialized with
421        initialized_with: u16,
422    },
423    /// Failed to decode metadata header
424    #[error("Failed to decode metadata header: {0}")]
425    FailedToDecodeMetadataHeader(parity_scale_codec::Error),
426    /// Unexpected metadata version
427    #[error("Unexpected metadata version {0}")]
428    UnexpectedMetadataVersion(u8),
429    /// Allocated space is not enough for one sector
430    #[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        /// Minimal allocated space
437        min_space: u64,
438        /// Current allocated space
439        allocated_space: u64,
440    },
441    /// Farm is too large
442    #[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
449        allocated_space: u64,
450        /// Allocated space in sectors
451        allocated_sectors: u64,
452        /// Max supported allocated space
453        max_space: u64,
454        /// Max supported allocated space in sectors
455        max_sectors: u16,
456    },
457    /// Failed to create thread pool
458    #[error("Failed to create thread pool: {0}")]
459    FailedToCreateThreadPool(ThreadPoolBuildError),
460}
461
462/// Errors happening during scrubbing
463#[derive(Debug, Error)]
464pub enum SingleDiskFarmScrubError {
465    /// Farm is likely already in use, make sure no other farmer is using it
466    #[error("Farm is likely already in use, make sure no other farmer is using it: {0}")]
467    LikelyAlreadyInUse(io::Error),
468    /// Failed to determine file size
469    #[error("Failed to file size of {file}: {error}")]
470    FailedToDetermineFileSize {
471        /// Affected file
472        file: PathBuf,
473        /// Low-level error
474        error: io::Error,
475    },
476    /// Failed to read bytes from file
477    #[error("Failed to read {size} bytes from {file} at offset {offset}: {error}")]
478    FailedToReadBytes {
479        /// Affected file
480        file: PathBuf,
481        /// Number of bytes to read
482        size: u64,
483        /// Offset in the file
484        offset: u64,
485        /// Low-level error
486        error: io::Error,
487    },
488    /// Failed to write bytes from file
489    #[error("Failed to write {size} bytes from {file} at offset {offset}: {error}")]
490    FailedToWriteBytes {
491        /// Affected file
492        file: PathBuf,
493        /// Number of bytes to read
494        size: u64,
495        /// Offset in the file
496        offset: u64,
497        /// Low-level error
498        error: io::Error,
499    },
500    /// Farm info file does not exist
501    #[error("Farm info file does not exist at {file}")]
502    FarmInfoFileDoesNotExist {
503        /// Info file
504        file: PathBuf,
505    },
506    /// Farm info can't be opened
507    #[error("Farm info at {file} can't be opened: {error}")]
508    FarmInfoCantBeOpened {
509        /// Info file
510        file: PathBuf,
511        /// Low-level error
512        error: io::Error,
513    },
514    /// Identity file does not exist
515    #[error("Identity file does not exist at {file}")]
516    IdentityFileDoesNotExist {
517        /// Identity file
518        file: PathBuf,
519    },
520    /// Identity can't be opened
521    #[error("Identity at {file} can't be opened: {error}")]
522    IdentityCantBeOpened {
523        /// Identity file
524        file: PathBuf,
525        /// Low-level error
526        error: IdentityError,
527    },
528    /// Identity public key doesn't match public key in the disk farm info
529    #[error("Identity public key {identity} doesn't match public key in the disk farm info {info}")]
530    PublicKeyMismatch {
531        /// Identity public key
532        identity: Ed25519PublicKey,
533        /// Disk farm info public key
534        info: Ed25519PublicKey,
535    },
536    /// Metadata file does not exist
537    #[error("Metadata file does not exist at {file}")]
538    MetadataFileDoesNotExist {
539        /// Metadata file
540        file: PathBuf,
541    },
542    /// Metadata can't be opened
543    #[error("Metadata at {file} can't be opened: {error}")]
544    MetadataCantBeOpened {
545        /// Metadata file
546        file: PathBuf,
547        /// Low-level error
548        error: io::Error,
549    },
550    /// Metadata file too small
551    #[error(
552        "Metadata file at {file} is too small: reserved size is {reserved_size} bytes, file size \
553        is {size}"
554    )]
555    MetadataFileTooSmall {
556        /// Metadata file
557        file: PathBuf,
558        /// Reserved size
559        reserved_size: u64,
560        /// File size
561        size: u64,
562    },
563    /// Failed to decode metadata header
564    #[error("Failed to decode metadata header: {0}")]
565    FailedToDecodeMetadataHeader(parity_scale_codec::Error),
566    /// Unexpected metadata version
567    #[error("Unexpected metadata version {0}")]
568    UnexpectedMetadataVersion(u8),
569    /// Cache can't be opened
570    #[error("Cache at {file} can't be opened: {error}")]
571    CacheCantBeOpened {
572        /// Cache file
573        file: PathBuf,
574        /// Low-level error
575        error: io::Error,
576    },
577}
578
579/// Errors that happen in background tasks
580#[derive(Debug, Error)]
581pub enum BackgroundTaskError {
582    /// Plotting error
583    #[error(transparent)]
584    Plotting(#[from] PlottingError),
585    /// Farming error
586    #[error(transparent)]
587    Farming(#[from] FarmingError),
588    /// Block sealing
589    #[error(transparent)]
590    BlockSealing(#[from] anyhow::Error),
591    /// Background task panicked
592    #[error("Background task {task} panicked")]
593    BackgroundTaskPanicked {
594        /// Name of the task
595        task: String,
596    },
597}
598
599type BackgroundTask = Pin<Box<dyn Future<Output = Result<(), BackgroundTaskError>> + Send>>;
600
601/// Scrub target
602#[derive(Debug, Copy, Clone)]
603pub enum ScrubTarget {
604    /// Scrub everything
605    All,
606    /// Scrub just metadata
607    Metadata,
608    /// Scrub metadata and corresponding plot
609    Plot,
610    /// Only scrub cache
611    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        // Fixed space usage regardless of plot size
679        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        // Calculate how many sectors can fit
684        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            // Do the rounding to make sure we have exactly as much space as fits whole number of
689            // sectors, account for disk sector size just in case
690            (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            // Cache must not be empty, ensure it contains at least one element even if
697            // percentage-wise it will use more space
698            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        // Align plot file size for disk sector size
712        let plot_file_size = plot_file_size.div_ceil(DISK_PAGE_SIZE as u64) * DISK_PAGE_SIZE as u64;
713
714        // Remaining space will be used for caching purposes
715        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                // We use `u16` for both count and index, hence index must not reach actual `MAX`
730                // (consensus doesn't care about this, just farmer implementation detail)
731                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/// Single disk farm abstraction is a container for everything necessary to plot/farm with a single
776/// disk.
777///
778/// Farm starts operating during creation and doesn't stop until dropped (or error happens).
779#[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    /// Metadata of all sectors plotted so far
786    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    /// Sender that will be used to signal to background threads that they should start
796    start_sender: Option<broadcast::Sender<()>>,
797    /// Sender that will be used to signal to background threads that they must stop
798    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        // Make background threads that are waiting to do something exit immediately
807        self.start_sender.take();
808        // Notify background tasks that they must stop
809        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    /// Name of the plot file
856    pub const PLOT_FILE: &'static str = "plot.bin";
857    /// Name of the metadata file
858    pub const METADATA_FILE: &'static str = "metadata.bin";
859    const SUPPORTED_PLOT_VERSION: u8 = 0;
860
861    /// Create new single disk farm instance
862    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            // Convert farm ID into cache ID for single disk farm
928            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        // Some sectors may already be plotted, skip them
986        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(&sectors_metadata);
1011            let handlers = Arc::clone(&handlers);
1012            let sectors_being_modified = Arc::clone(&sectors_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: &sectors_metadata,
1026                    sectors_being_modified: &sectors_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                        // Dropped before starting
1047                        return Ok(());
1048                    }
1049
1050                    if let Some(plotting_delay) = plotting_delay
1051                        && plotting_delay.await.is_err()
1052                    {
1053                        // Dropped before resolving
1054                        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                            // Nothing, just exit
1075                        }
1076                    }
1077                });
1078            }
1079        });
1080        let plotting_join_handle = AsyncJoinOnDrop::new(plotting_join_handle, false);
1081
1082        tasks.push(Box::pin(async move {
1083            // Panic will already be printed by now
1084            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(&sectors_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(&sectors_being_modified);
1127            let sectors_metadata = Arc::clone(&sectors_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                        // Dropped before starting
1140                        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                            // Nothing, just exit
1178                        }
1179                    }
1180                });
1181            }
1182        });
1183        let farming_join_handle = AsyncJoinOnDrop::new(farming_join_handle, false);
1184
1185        tasks.push(Box::pin(async move {
1186            // Panic will already be printed by now
1187            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(&sectors_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                            // Nothing, just exit
1214                        }
1215                        _ = stop_receiver.recv().fuse() => {
1216                            // Nothing, just exit
1217                        }
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            // Panic will already be printed by now
1227            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        // Align plot file size for disk sector size
1397        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                // Allocating the whole file (`set_len` below can create a sparse file, which will
1414                // cause writes to fail later)
1415                metadata_file
1416                    .preallocate(expected_metadata_size)
1417                    .map_err(SingleDiskFarmError::CantPreallocateMetadataFile)?;
1418                // Truncating file (if necessary)
1419                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            // Allocating the whole file (`set_len` below can create a sparse file, which will cause
1488            // writes to fail later)
1489            plot_file
1490                .preallocate(allocated_space_distribution.plot_file_size)
1491                .map_err(SingleDiskFarmError::CantPreallocatePlotFile)?;
1492            // Truncating file (if necessary)
1493            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            &sectors_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    /// Collect summary of single disk farm for presentational purposes
1520    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    /// Effective on-disk allocation of the files related to the farm (takes some buffer space
1538    /// into consideration).
1539    ///
1540    /// This is a helpful number in case some files were not allocated properly or were removed and
1541    /// do not correspond to allocated space in the farm info accurately.
1542    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                // No farm info, try to collect actual file sizes is any
1564                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                    // File is not stored on disk
1582                } 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                    // File is not stored on disk
1598                } 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                    // File is not stored on disk
1614                } else {
1615                    return Err(error.into());
1616                }
1617            }
1618        }
1619
1620        Ok(effective_disk_usage)
1621    }
1622
1623    /// Read all sectors metadata
1624    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    /// ID of this farm
1668    pub fn id(&self) -> &FarmId {
1669        self.single_disk_farm_info.id()
1670    }
1671
1672    /// Info of this farm
1673    pub fn info(&self) -> &SingleDiskFarmInfo {
1674        &self.single_disk_farm_info
1675    }
1676
1677    /// Number of sectors in this farm
1678    pub fn total_sectors_count(&self) -> u16 {
1679        self.total_sectors_count
1680    }
1681
1682    /// Read information about sectors plotted so far
1683    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    /// Get piece cache instance
1694    pub fn piece_cache(&self) -> SingleDiskPieceCache {
1695        self.piece_cache.clone()
1696    }
1697
1698    /// Get plot cache instance
1699    pub fn plot_cache(&self) -> DiskPlotCache {
1700        self.plot_cache.clone()
1701    }
1702
1703    /// Get piece reader to read plotted pieces later
1704    pub fn piece_reader(&self) -> DiskPieceReader {
1705        self.piece_reader.clone()
1706    }
1707
1708    /// Subscribe to sector updates
1709    pub fn on_sector_update(&self, callback: HandlerFn<(SectorIndex, SectorUpdate)>) -> HandlerId {
1710        self.handlers.sector_update.add(callback)
1711    }
1712
1713    /// Subscribe to farming notifications
1714    pub fn on_farming_notification(&self, callback: HandlerFn<FarmingNotification>) -> HandlerId {
1715        self.handlers.farming_notification.add(callback)
1716    }
1717
1718    /// Subscribe to new solution notification
1719    pub fn on_solution(&self, callback: HandlerFn<SolutionResponse>) -> HandlerId {
1720        self.handlers.solution.add(callback)
1721    }
1722
1723    /// Run and wait for background threads to exit or return an error
1724    pub async fn run(mut self) -> anyhow::Result<()> {
1725        if let Some(start_sender) = self.start_sender.take() {
1726            // Do not care if anyone is listening on the other side
1727            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    /// Wipe everything that belongs to this single disk farm
1738    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        // TODO: Identity should be able to wipe itself instead of assuming a specific file name
1773        //  here
1774        {
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    /// Check the farm for corruption and repair errors (caused by disk errors or something else),
1792    /// returns an error when irrecoverable errors occur.
1793    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                // Error doesn't matter here
1879                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                // Error doesn't matter here
1983                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                            // Read sector bytes and compute checksum
2129                            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                            // Verify checksum
2176                            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                                // Fill sector with zeroes and compute checksum
2199                                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                                // Write checksum
2223                                {
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        // Error doesn't matter here
2300        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}