1use 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#[derive(Debug, Error)]
28pub enum ReadingError {
29 #[error("Failed to read chunk at location {chunk_location}: {error}")]
34 FailedToReadChunk {
35 chunk_location: u64,
37 error: io::Error,
39 },
40 #[error("Missing PoS proof for s-bucket {s_bucket}")]
45 MissingPosProof {
46 s_bucket: SBucket,
48 },
49 #[error("Failed to erasure-decode record at offset {piece_offset}: {error}")]
51 FailedToErasureDecodeRecord {
52 piece_offset: PieceOffset,
54 error: ErasureCodingError,
56 },
57 #[error("Wrong record size after decoding: expected {expected}, actual {actual}")]
59 WrongRecordSizeAfterDecoding {
60 expected: usize,
62 actual: usize,
64 },
65 #[error("Failed to decode sector contents map: {0}")]
67 FailedToDecodeSectorContentsMap(#[from] SectorContentsMapFromBytesError),
68 #[error("Reading I/O error: {0}")]
70 Io(#[from] io::Error),
71 #[error("Checksum mismatch")]
73 ChecksumMismatch,
74}
75
76impl ReadingError {
77 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#[derive(Debug)]
101pub struct SectorRecordChunks {
102 pub source: Box<Record>,
104 pub parity: Box<Record>,
106 pub present: ShardsPresent<{ Record::NUM_CHUNKS }>,
108}
109
110impl SectorRecordChunks {
111 #[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
129pub 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 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 §or
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 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
250pub 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 drop(parity);
271
272 Ok(source)
273}
274
275pub(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 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
308pub 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(§or_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 §or_metadata.s_bucket_offsets(),
346 §or_contents_map,
347 &table_generator.create_proofs(§or_id.derive_evaluation_seed(piece_offset)),
348 sector,
349 )
350 .await?;
351 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 piece.record.copy_from_slice(&**record);
365
366 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}