subspace_farmer/single_disk_farm/
piece_reader.rs1use crate::farm::{FarmError, PieceReader};
4use crate::single_disk_farm::direct_io_file::DirectIoFile;
5use async_lock::{Mutex as AsyncMutex, RwLock as AsyncRwLock};
6use async_trait::async_trait;
7use futures::channel::{mpsc, oneshot};
8use futures::{SinkExt, StreamExt};
9use std::collections::HashSet;
10use std::future::Future;
11use std::sync::Arc;
12use subspace_core_primitives::PublicKey;
13use subspace_core_primitives::pieces::{Piece, PieceOffset};
14use subspace_core_primitives::sectors::{SectorId, SectorIndex};
15use subspace_core_primitives::segments::HistorySize;
16use subspace_erasure_coding::ErasureCoding;
17use subspace_farmer_components::reading::ReadSectorRecordChunksMode;
18use subspace_farmer_components::sector::{SectorMetadataChecksummed, sector_size};
19use subspace_farmer_components::{ReadAt, ReadAtAsync, ReadAtSync, reading};
20use subspace_proof_of_space::Table;
21use tracing::{error, warn};
22
23#[derive(Debug)]
24struct ReadPieceRequest {
25 sector_index: SectorIndex,
26 piece_offset: PieceOffset,
27 response_sender: oneshot::Sender<Option<Piece>>,
28}
29
30#[derive(Debug, Clone)]
32pub struct DiskPieceReader {
33 read_piece_sender: mpsc::Sender<ReadPieceRequest>,
34}
35
36#[async_trait]
37impl PieceReader for DiskPieceReader {
38 #[inline]
39 async fn read_piece(
40 &self,
41 sector_index: SectorIndex,
42 piece_offset: PieceOffset,
43 ) -> Result<Option<Piece>, FarmError> {
44 Ok(self.read_piece(sector_index, piece_offset).await)
45 }
46}
47
48impl DiskPieceReader {
49 #[allow(clippy::too_many_arguments)]
54 pub(super) fn new<PosTable>(
55 public_key: PublicKey,
56 pieces_in_sector: u16,
57 plot_file: Arc<DirectIoFile>,
58 sectors_metadata: Arc<AsyncRwLock<Vec<SectorMetadataChecksummed>>>,
59 erasure_coding: ErasureCoding,
60 sectors_being_modified: Arc<AsyncRwLock<HashSet<SectorIndex>>>,
61 read_sector_record_chunks_mode: ReadSectorRecordChunksMode,
62 cutover: Option<HistorySize>,
63 global_mutex: Arc<AsyncMutex<()>>,
64 ) -> (Self, impl Future<Output = ()>)
65 where
66 PosTable: Table,
67 {
68 let (read_piece_sender, read_piece_receiver) = mpsc::channel(10);
69
70 let reading_fut = async move {
71 read_pieces::<PosTable, _>(
72 public_key,
73 pieces_in_sector,
74 &*plot_file,
75 sectors_metadata,
76 erasure_coding,
77 sectors_being_modified,
78 read_piece_receiver,
79 read_sector_record_chunks_mode,
80 cutover,
81 global_mutex,
82 )
83 .await
84 };
85
86 (Self { read_piece_sender }, reading_fut)
87 }
88
89 pub(super) fn close_all_readers(&mut self) {
90 self.read_piece_sender.close_channel();
91 }
92
93 pub async fn read_piece(
96 &self,
97 sector_index: SectorIndex,
98 piece_offset: PieceOffset,
99 ) -> Option<Piece> {
100 let (response_sender, response_receiver) = oneshot::channel();
101 self.read_piece_sender
102 .clone()
103 .send(ReadPieceRequest {
104 sector_index,
105 piece_offset,
106 response_sender,
107 })
108 .await
109 .ok()?;
110 response_receiver.await.ok()?
111 }
112}
113
114#[allow(clippy::too_many_arguments)]
115async fn read_pieces<PosTable, S>(
116 public_key: PublicKey,
117 pieces_in_sector: u16,
118 plot_file: S,
119 sectors_metadata: Arc<AsyncRwLock<Vec<SectorMetadataChecksummed>>>,
120 erasure_coding: ErasureCoding,
121 sectors_being_modified: Arc<AsyncRwLock<HashSet<SectorIndex>>>,
122 mut read_piece_receiver: mpsc::Receiver<ReadPieceRequest>,
123 mode: ReadSectorRecordChunksMode,
124 cutover: Option<HistorySize>,
125 global_mutex: Arc<AsyncMutex<()>>,
126) where
127 PosTable: Table,
128 S: ReadAtSync,
129{
130 let mut table_generator_old = PosTable::generator_for(false);
133 let mut table_generator_new = PosTable::generator_for(true);
134
135 while let Some(read_piece_request) = read_piece_receiver.next().await {
136 let ReadPieceRequest {
137 sector_index,
138 piece_offset,
139 response_sender,
140 } = read_piece_request;
141
142 if response_sender.is_canceled() {
143 continue;
144 }
145
146 let sectors_being_modified = &*sectors_being_modified.read().await;
147
148 if sectors_being_modified.contains(§or_index) {
149 continue;
151 }
152
153 let (sector_metadata, sector_count) = {
154 let sectors_metadata = sectors_metadata.read().await;
155
156 let sector_count = sectors_metadata.len() as SectorIndex;
157
158 let sector_metadata = match sectors_metadata.get(sector_index as usize) {
159 Some(sector_metadata) => sector_metadata.clone(),
160 None => {
161 error!(
162 %sector_index,
163 %sector_count,
164 "Tried to read piece from sector that is not yet plotted"
165 );
166 continue;
167 }
168 };
169
170 (sector_metadata, sector_count)
171 };
172
173 if sector_index >= sector_count {
175 warn!(
176 %sector_index,
177 %piece_offset,
178 %sector_count,
179 "Incorrect sector offset"
180 );
181 let _ = response_sender.send(None);
183 continue;
184 }
185 if u16::from(piece_offset) >= pieces_in_sector {
187 warn!(
188 %sector_index,
189 %piece_offset,
190 %sector_count,
191 "Incorrect piece offset"
192 );
193 let _ = response_sender.send(None);
195 continue;
196 }
197
198 let sector_size = sector_size(pieces_in_sector);
199 let sector = plot_file.offset(u64::from(sector_index) * sector_size as u64);
200
201 global_mutex.lock().await;
203
204 let is_post_cutover = super::is_post_cutover(cutover, sector_metadata.history_size);
205 let table_generator = if is_post_cutover {
206 &mut table_generator_new
207 } else {
208 &mut table_generator_old
209 };
210
211 let maybe_piece = read_piece::<PosTable, _, _>(
212 &public_key,
213 piece_offset,
214 §or_metadata,
215 &ReadAt::from_sync(§or),
217 &erasure_coding,
218 mode,
219 table_generator,
220 )
221 .await;
222
223 let _ = response_sender.send(maybe_piece);
225 }
226}
227
228async fn read_piece<PosTable, S, A>(
229 public_key: &PublicKey,
230 piece_offset: PieceOffset,
231 sector_metadata: &SectorMetadataChecksummed,
232 sector: &ReadAt<S, A>,
233 erasure_coding: &ErasureCoding,
234 mode: ReadSectorRecordChunksMode,
235 table_generator: &mut PosTable::Generator,
236) -> Option<Piece>
237where
238 PosTable: Table,
239 S: ReadAtSync,
240 A: ReadAtAsync,
241{
242 let sector_index = sector_metadata.sector_index;
243
244 let sector_id = SectorId::new(
245 public_key.hash(),
246 sector_index,
247 sector_metadata.history_size,
248 );
249
250 let piece = match reading::read_piece::<PosTable, _, _>(
251 piece_offset,
252 §or_id,
253 sector_metadata,
254 sector,
255 erasure_coding,
256 mode,
257 table_generator,
258 )
259 .await
260 {
261 Ok(piece) => piece,
262 Err(error) => {
263 let evaluation_seed = sector_id.derive_evaluation_seed(piece_offset);
266 error!(
267 %sector_index,
268 %piece_offset,
269 evaluation_seed = %hex::encode(*evaluation_seed),
270 %error,
271 "Failed to read piece from sector"
272 );
273 return None;
274 }
275 };
276
277 Some(piece)
278}