Skip to main content

ab_farmer_components/
reading.rs

1//! Reading utilities
2//!
3//! This module contains utilities for extracting data from plots/sectors created by functions in
4//! [`plotting`](crate::plotting) module earlier. This is a relatively expensive operation and is
5//! only used for cold storage purposes or when there is a need to prove a solution to consensus.
6
7use crate::sector::{
8    RecordMetadata, SectorContentsMap, SectorContentsMapFromBytesError, SectorMetadataChecksummed,
9    sector_record_chunks_size,
10};
11use crate::{ReadAt, ReadAtAsync, ReadAtSync};
12use ab_core_primitives::hashes::Blake3Hash;
13use ab_core_primitives::pieces::{Piece, PieceOffset, Record, RecordChunk};
14use ab_core_primitives::sectors::{SBucket, SectorId};
15use ab_erasure_coding::{ErasureCoding, ErasureCodingError, ShardsPresent};
16use ab_proof_of_space::{PosProofs, Table, TableGenerator};
17use futures::StreamExt;
18use futures::stream::FuturesUnordered;
19use parity_scale_codec::Decode;
20use rayon::prelude::*;
21use std::io;
22use std::simd::Simd;
23use thiserror::Error;
24use tracing::debug;
25
26/// Errors that happen during reading
27#[derive(Debug, Error)]
28pub enum ReadingError {
29    /// Failed to read chunk.
30    ///
31    /// This is an implementation bug, most likely due to mismatch between sector contents map and
32    /// other farming parameters.
33    #[error("Failed to read chunk at location {chunk_location}: {error}")]
34    FailedToReadChunk {
35        /// Chunk location
36        chunk_location: u64,
37        /// Low-level error
38        error: io::Error,
39    },
40    /// Missing proof of space proof.
41    ///
42    /// This is either hardware issue or if happens for everyone all the time an implementation
43    /// bug.
44    #[error("Missing PoS proof for s-bucket {s_bucket}")]
45    MissingPosProof {
46        /// S-bucket
47        s_bucket: SBucket,
48    },
49    /// Failed to erasure-decode record
50    #[error("Failed to erasure-decode record at offset {piece_offset}: {error}")]
51    FailedToErasureDecodeRecord {
52        /// Piece offset
53        piece_offset: PieceOffset,
54        /// Lower-level error
55        error: ErasureCodingError,
56    },
57    /// Wrong record size after decoding
58    #[error("Wrong record size after decoding: expected {expected}, actual {actual}")]
59    WrongRecordSizeAfterDecoding {
60        /// Expected size in bytes
61        expected: usize,
62        /// Actual size in bytes
63        actual: usize,
64    },
65    /// Failed to decode sector contents map
66    #[error("Failed to decode sector contents map: {0}")]
67    FailedToDecodeSectorContentsMap(#[from] SectorContentsMapFromBytesError),
68    /// I/O error occurred
69    #[error("Reading I/O error: {0}")]
70    Io(#[from] io::Error),
71    /// Checksum mismatch
72    #[error("Checksum mismatch")]
73    ChecksumMismatch,
74}
75
76impl ReadingError {
77    /// Whether this error is fatal and renders farm unusable
78    pub fn is_fatal(&self) -> bool {
79        #[expect(
80            clippy::rest_pattern_accessible_field,
81            reason = "Do not care about fields"
82        )]
83        match self {
84            ReadingError::FailedToReadChunk { .. } => false,
85            ReadingError::MissingPosProof { .. } => false,
86            ReadingError::FailedToErasureDecodeRecord { .. } => false,
87            ReadingError::WrongRecordSizeAfterDecoding { .. } => false,
88            ReadingError::FailedToDecodeSectorContentsMap(_) => false,
89            ReadingError::Io(_) => true,
90            ReadingError::ChecksumMismatch => false,
91        }
92    }
93}
94
95/// Record chunks read from a sector.
96///
97/// Source and parity chunks are separate allocations so that a caller which only needs the source
98/// record can drop the parity one as soon as recovery is done. Contents of chunks that are not
99/// present are unspecified.
100#[derive(Debug)]
101pub struct SectorRecordChunks {
102    /// Source chunks
103    pub source: Box<Record>,
104    /// Parity chunks
105    pub parity: Box<Record>,
106    /// Which chunks are present
107    pub present: ShardsPresent<{ Record::NUM_CHUNKS }>,
108}
109
110impl SectorRecordChunks {
111    /// Returns the chunk at the given s-bucket
112    #[inline]
113    pub fn get_chunk(&self, s_bucket: SBucket) -> RecordChunk {
114        let s_bucket = usize::from(s_bucket);
115        RecordChunk::from(
116            *if let Some(parity_index) = s_bucket.checked_sub(Record::NUM_CHUNKS) {
117                self.parity
118                    .get(parity_index)
119                    .expect("Within correct range; qed")
120            } else {
121                self.source
122                    .get(s_bucket)
123                    .expect("Within correct range; qed")
124            },
125        )
126    }
127}
128
129/// Read sector record chunks, only plotted s-buckets are marked as present (in decoded form).
130///
131/// NOTE: This is an async function, but it also does CPU-intensive operation internally, while it
132/// is not very long, make sure it is okay to do so in your context.
133pub async fn read_sector_record_chunks<S, A>(
134    piece_offset: PieceOffset,
135    pieces_in_sector: u16,
136    s_bucket_offsets: &[u32; Record::NUM_S_BUCKETS],
137    sector_contents_map: &SectorContentsMap,
138    pos_proofs: &PosProofs,
139    sector: &ReadAt<S, A>,
140) -> Result<SectorRecordChunks, ReadingError>
141where
142    S: ReadAtSync,
143    A: ReadAtAsync,
144{
145    let mut source = Record::new_boxed();
146    let mut parity = Record::new_boxed();
147    let mut present = ShardsPresent::none();
148
149    let read_chunks_inputs = source
150        .par_iter_mut()
151        .chain(parity.par_iter_mut())
152        .zip(sector_contents_map.par_iter_record_chunk_to_plot(piece_offset))
153        .zip(s_bucket_offsets.par_iter())
154        .enumerate()
155        .map(
156            |(index, ((record_chunk, maybe_chunk_offset), &s_bucket_offset))| {
157                let chunk_offset = maybe_chunk_offset?;
158
159                let chunk_location = chunk_offset as u64 + u64::from(s_bucket_offset);
160
161                Some((index, record_chunk, chunk_location))
162            },
163        )
164        .flatten()
165        .collect::<Vec<_>>();
166
167    for &(index, _, _) in &read_chunks_inputs {
168        if let Some(parity_index) = index.checked_sub(Record::NUM_CHUNKS) {
169            present.parity.set(parity_index);
170        } else {
171            present.source.set(index);
172        }
173    }
174
175    let sector_contents_map_size = SectorContentsMap::encoded_size(pieces_in_sector) as u64;
176    match sector {
177        ReadAt::Sync(sector) => {
178            read_chunks_inputs
179                .into_par_iter()
180                .zip(&pos_proofs.proofs)
181                .try_for_each(|((_index, output_chunk, chunk_location), pos_proof)| {
182                    let mut record_chunk = [0; RecordChunk::SIZE];
183                    sector
184                        .read_at(
185                            &mut record_chunk,
186                            sector_contents_map_size + chunk_location * RecordChunk::SIZE as u64,
187                        )
188                        .map_err(|error| ReadingError::FailedToReadChunk {
189                            chunk_location,
190                            error,
191                        })?;
192
193                    // TODO: Use SIMD for hashing
194                    record_chunk =
195                        Simd::to_array(Simd::from(record_chunk) ^ Simd::from(*pos_proof.hash()));
196
197                    *output_chunk = record_chunk;
198
199                    Ok::<_, ReadingError>(())
200                })?;
201        }
202        ReadAt::Async(sector) => {
203            let processing_chunks = read_chunks_inputs
204                .into_iter()
205                .zip(&pos_proofs.proofs)
206                .map(
207                    |((_index, output_chunk, chunk_location), pos_proof)| async move {
208                        let mut record_chunk = [0; RecordChunk::SIZE];
209                        record_chunk.copy_from_slice(
210                            &sector
211                                .read_at(
212                                    vec![0; RecordChunk::SIZE],
213                                    sector_contents_map_size
214                                        + chunk_location * RecordChunk::SIZE as u64,
215                                )
216                                .await
217                                .map_err(|error| ReadingError::FailedToReadChunk {
218                                    chunk_location,
219                                    error,
220                                })?,
221                        );
222
223                        // TODO: Use SIMD for hashing
224                        record_chunk = Simd::to_array(
225                            Simd::from(record_chunk) ^ Simd::from(*pos_proof.hash()),
226                        );
227
228                        *output_chunk = record_chunk;
229
230                        Ok::<_, ReadingError>(())
231                    },
232                )
233                .collect::<FuturesUnordered<_>>()
234                .filter_map(|result| async move { result.err() });
235
236            std::pin::pin!(processing_chunks)
237                .next()
238                .await
239                .map_or(Ok(()), Err)?;
240        }
241    }
242
243    Ok(SectorRecordChunks {
244        source,
245        parity,
246        present,
247    })
248}
249
250/// Given sector record chunks recover the source record
251pub fn recover_source_record(
252    sector_record_chunks: SectorRecordChunks,
253    piece_offset: PieceOffset,
254    erasure_coding: &ErasureCoding,
255) -> Result<Box<Record>, ReadingError> {
256    let SectorRecordChunks {
257        mut source,
258        parity,
259        present,
260    } = sector_record_chunks;
261
262    erasure_coding
263        .recover_source(&mut source, &parity, &present)
264        .map_err(|error| ReadingError::FailedToErasureDecodeRecord {
265            piece_offset,
266            error,
267        })?;
268
269    // Parity chunks are no longer needed, dropping them here keeps peak memory usage down
270    drop(parity);
271
272    Ok(source)
273}
274
275/// Read metadata (roots and proof) for record
276pub(crate) async fn read_record_metadata<S, A>(
277    piece_offset: PieceOffset,
278    pieces_in_sector: u16,
279    sector: &ReadAt<S, A>,
280) -> Result<RecordMetadata, ReadingError>
281where
282    S: ReadAtSync,
283    A: ReadAtAsync,
284{
285    let sector_metadata_start = SectorContentsMap::encoded_size(pieces_in_sector) as u64
286        + sector_record_chunks_size(pieces_in_sector) as u64;
287    // Move to the beginning of the root and proof we care about
288    let record_metadata_offset =
289        sector_metadata_start + RecordMetadata::encoded_size() as u64 * u64::from(piece_offset);
290
291    let mut record_metadata_bytes = vec![0; RecordMetadata::encoded_size()];
292    match sector {
293        ReadAt::Sync(sector) => {
294            sector.read_at(&mut record_metadata_bytes, record_metadata_offset)?;
295        }
296        ReadAt::Async(sector) => {
297            record_metadata_bytes = sector
298                .read_at(record_metadata_bytes, record_metadata_offset)
299                .await?;
300        }
301    }
302    let record_metadata = RecordMetadata::decode(&mut record_metadata_bytes.as_ref())
303        .expect("Length is correct, contents doesn't have specific structure to it; qed");
304
305    Ok(record_metadata)
306}
307
308/// Read piece from sector.
309///
310/// NOTE: Even though this function is async, proof of time table generation is expensive and should
311/// be done in a dedicated thread where blocking is allowed.
312pub async fn read_piece<PosTable, S, A>(
313    piece_offset: PieceOffset,
314    sector_id: &SectorId,
315    sector_metadata: &SectorMetadataChecksummed,
316    sector: &ReadAt<S, A>,
317    erasure_coding: &ErasureCoding,
318    table_generator: &PosTable::Generator,
319) -> Result<Piece, ReadingError>
320where
321    PosTable: Table,
322    S: ReadAtSync,
323    A: ReadAtAsync,
324{
325    let pieces_in_sector = sector_metadata.pieces_in_sector;
326
327    let sector_contents_map = {
328        let mut sector_contents_map_bytes =
329            vec![0; SectorContentsMap::encoded_size(pieces_in_sector)];
330        match sector {
331            ReadAt::Sync(sector) => {
332                sector.read_at(&mut sector_contents_map_bytes, 0)?;
333            }
334            ReadAt::Async(sector) => {
335                sector_contents_map_bytes = sector.read_at(sector_contents_map_bytes, 0).await?;
336            }
337        }
338
339        SectorContentsMap::from_bytes(&sector_contents_map_bytes, pieces_in_sector)?
340    };
341
342    let sector_record_chunks = read_sector_record_chunks(
343        piece_offset,
344        pieces_in_sector,
345        &sector_metadata.s_bucket_offsets(),
346        &sector_contents_map,
347        &table_generator.create_proofs(&sector_id.derive_evaluation_seed(piece_offset)),
348        sector,
349    )
350    .await?;
351    // Restore source record scalars
352    let record = recover_source_record(sector_record_chunks, piece_offset, erasure_coding)?;
353
354    let RecordMetadata {
355        piece_header,
356        piece_checksum,
357    } = read_record_metadata(piece_offset, pieces_in_sector, sector).await?;
358
359    let mut piece = Piece::default();
360
361    piece.header = piece_header;
362    // Fancy way to insert value to avoid going through stack (if naive dereferencing is used)
363    // and potentially causing stack overflow as the result
364    piece.record.copy_from_slice(&**record);
365
366    // Verify checksum
367    let actual_checksum = Blake3Hash::from(blake3::hash(piece.as_ref()));
368    if actual_checksum != piece_checksum {
369        debug!(
370            ?sector_id,
371            %piece_offset,
372            %actual_checksum,
373            expected_checksum = %piece_checksum,
374            "Hash doesn't match, plotted piece is corrupted"
375        );
376
377        return Err(ReadingError::ChecksumMismatch);
378    }
379
380    Ok(piece.to_shared())
381}