Skip to main content

subspace_farmer/single_disk_farm/
piece_reader.rs

1//! Piece reader for single disk farm
2
3use 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/// Wrapper data structure that can be used to read pieces from single disk farm
31#[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    /// Creates new piece reader instance and background future that handles reads internally.
50    ///
51    /// NOTE: Background future is async, but does blocking operations and should be running in
52    /// dedicated thread.
53    #[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    /// Read piece from sector by offset, `None` means input parameters are incorrect or piece
94    /// reader was shut down
95    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    // Keep a warm generator for each proof-of-space so old and new sectors can be read without
131    // re-allocating table caches per request.
132    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(&sector_index) {
149            // Skip sector that is being modified right now
150            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        // Sector must be plotted
174        if sector_index >= sector_count {
175            warn!(
176                %sector_index,
177                %piece_offset,
178                %sector_count,
179                "Incorrect sector offset"
180            );
181            // Doesn't matter if receiver still cares about it
182            let _ = response_sender.send(None);
183            continue;
184        }
185        // Piece must be within sector
186        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            // Doesn't matter if receiver still cares about it
194            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        // Take mutex briefly to make sure piece reading is allowed right now
202        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            &sector_metadata,
215            // TODO: Async
216            &ReadAt::from_sync(&sector),
217            &erasure_coding,
218            mode,
219            table_generator,
220        )
221        .await;
222
223        // Doesn't matter if receiver still cares about it
224        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        &sector_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            // The evaluation seed regenerates the exact proof-of-space table for this read, making
264            // a failure locally reproducible.
265            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}