Skip to main content

zebra_state/
service.rs

1//! [`tower::Service`]s for Zebra's cached chain state.
2//!
3//! Zebra provides cached state access via two main services:
4//! - [`StateService`]: a read-write service that writes blocks to the state,
5//!   and redirects most read requests to the [`ReadStateService`].
6//! - [`ReadStateService`]: a read-only service that answers from the most
7//!   recent committed block.
8//!
9//! Most users should prefer [`ReadStateService`], unless they need to write blocks to the state.
10//!
11//! Zebra also provides access to the best chain tip via:
12//! - [`LatestChainTip`]: a read-only channel that contains the latest committed
13//!   tip.
14//! - [`ChainTipChange`]: a read-only channel that can asynchronously await
15//!   chain tip changes.
16
17use std::{
18    collections::HashMap,
19    future::Future,
20    pin::Pin,
21    sync::Arc,
22    task::{Context, Poll},
23    time::{Duration, Instant},
24};
25
26use futures::future::FutureExt;
27use tokio::sync::oneshot;
28use tower::{util::BoxService, Service, ServiceExt};
29use tracing::{instrument, Instrument, Span};
30
31#[cfg(any(test, feature = "proptest-impl"))]
32use tower::buffer::Buffer;
33
34use zebra_chain::{
35    block::{self, CountedHeader, HeightDiff},
36    diagnostic::CodeTimer,
37    parameters::{Network, NetworkUpgrade},
38    serialization::ZcashSerialize,
39    subtree::NoteCommitmentSubtreeIndex,
40};
41
42use crate::{
43    constants::{
44        MAX_FIND_BLOCK_HASHES_RESULTS, MAX_FIND_BLOCK_HEADERS_RESULTS, MAX_LEGACY_CHAIN_BLOCKS,
45    },
46    error::{CommitBlockError, CommitCheckpointVerifiedError, InvalidateError, ReconsiderError},
47    request::TimedSpan,
48    response::NonFinalizedBlocksListener,
49    service::{
50        block_iter::any_ancestor_blocks,
51        chain_tip::{ChainTipBlock, ChainTipChange, ChainTipSender, LatestChainTip},
52        finalized_state::{FinalizedState, ZebraDb},
53        non_finalized_state::{Chain, NonFinalizedState},
54        pending_utxos::PendingUtxos,
55        queued_blocks::QueuedBlocks,
56        read::find,
57        watch_receiver::WatchReceiver,
58    },
59    BoxError, CheckpointVerifiedBlock, CommitSemanticallyVerifiedError, Config, KnownBlock,
60    ReadRequest, ReadResponse, Request, Response, SemanticallyVerifiedBlock, StateInitError,
61};
62
63pub mod block_iter;
64pub mod chain_tip;
65pub mod watch_receiver;
66
67pub mod check;
68
69pub(crate) mod finalized_state;
70pub(crate) mod non_finalized_state;
71mod pending_utxos;
72mod queued_blocks;
73pub(crate) mod read;
74mod traits;
75mod write;
76
77#[cfg(any(test, feature = "proptest-impl"))]
78pub mod arbitrary;
79
80#[cfg(test)]
81mod tests;
82
83pub use finalized_state::{OutputLocation, TransactionIndex, TransactionLocation};
84use write::NonFinalizedWriteMessage;
85
86use self::queued_blocks::{QueuedCheckpointVerified, QueuedSemanticallyVerified, SentHashes};
87
88pub use self::traits::{ReadState, State};
89
90/// A read-write service for Zebra's cached blockchain state.
91///
92/// This service modifies and provides access to:
93/// - the non-finalized state: the most recent blocks, up to
94///   [`MAX_BLOCK_REORG_HEIGHT`](crate::MAX_BLOCK_REORG_HEIGHT) of them.
95///   Zebra allows chain forks in the non-finalized state,
96///   stores it in memory, and re-downloads it when restarted.
97/// - the finalized state: older blocks that have many confirmations.
98///   Zebra stores the single best chain in the finalized state,
99///   and re-loads it from disk when restarted.
100///
101/// Read requests to this service are buffered, then processed concurrently.
102/// Block write requests are buffered, then queued, then processed in order by a separate task.
103///
104/// Most state users can get faster read responses using the [`ReadStateService`],
105/// because its requests do not share a [`tower::buffer::Buffer`] with block write requests.
106///
107/// To quickly get the latest block, use [`LatestChainTip`] or [`ChainTipChange`].
108/// They can read the latest block directly, without queueing any requests.
109#[derive(Debug)]
110pub(crate) struct StateService {
111    // Configuration
112    //
113    /// The configured Zcash network.
114    network: Network,
115
116    /// The height that we start storing UTXOs from finalized blocks.
117    ///
118    /// This height should be lower than the last few checkpoints,
119    /// so the full verifier can verify UTXO spends from those blocks,
120    /// even if they haven't been committed to the finalized state yet.
121    full_verifier_utxo_lookahead: block::Height,
122
123    // Queued Blocks
124    //
125    /// Queued blocks for the [`NonFinalizedState`] that arrived out of order.
126    /// These blocks are awaiting their parent blocks before they can do contextual verification.
127    non_finalized_state_queued_blocks: QueuedBlocks,
128
129    /// Queued blocks for the [`FinalizedState`] that arrived out of order.
130    /// These blocks are awaiting their parent blocks before they can do contextual verification.
131    ///
132    /// Indexed by their parent block hash.
133    finalized_state_queued_blocks: HashMap<block::Hash, QueuedCheckpointVerified>,
134
135    /// Channels to send blocks to the block write task.
136    block_write_sender: write::BlockWriteSender,
137
138    /// The [`block::Hash`] of the most recent block sent on
139    /// `finalized_block_write_sender` or `non_finalized_block_write_sender`.
140    ///
141    /// On startup, this is:
142    /// - the finalized tip, if there are stored blocks, or
143    /// - the genesis block's parent hash, if the database is empty.
144    ///
145    /// If `invalid_block_write_reset_receiver` gets a reset, this is:
146    /// - the hash of the last valid committed block (the parent of the invalid block).
147    finalized_block_write_last_sent_hash: block::Hash,
148
149    /// A set of block hashes that have been sent to the block write task.
150    /// Hashes of blocks below the finalized tip height are periodically pruned.
151    non_finalized_block_write_sent_hashes: SentHashes,
152
153    /// If an invalid block is sent on `finalized_block_write_sender`
154    /// or `non_finalized_block_write_sender`,
155    /// this channel gets the [`block::Hash`] of the valid tip.
156    //
157    // TODO: add tests for finalized and non-finalized resets (#2654)
158    invalid_block_write_reset_receiver: tokio::sync::mpsc::UnboundedReceiver<block::Hash>,
159
160    /// Receives the hash of every non-finalized block that the write task
161    /// rejected, so the corresponding entry can be removed from
162    /// `non_finalized_block_write_sent_hashes`.
163    ///
164    /// Without this, a rejected same-hash block locks out a later honest
165    /// re-delivery of a block at the same hash as a "duplicate" until restart
166    /// or reorg.
167    non_finalized_rejected_receiver: tokio::sync::mpsc::UnboundedReceiver<block::Hash>,
168
169    // Pending UTXO Request Tracking
170    //
171    /// The set of outpoints with pending requests for their associated transparent::Output.
172    pending_utxos: PendingUtxos,
173
174    /// Instant tracking the last time `pending_utxos` was pruned.
175    last_prune: Instant,
176
177    // Updating Concurrently Readable State
178    //
179    /// A cloneable [`ReadStateService`], used to answer concurrent read requests.
180    ///
181    /// TODO: move users of read [`Request`]s to [`ReadStateService`], and remove `read_service`.
182    read_service: ReadStateService,
183
184    // Metrics
185    //
186    /// A metric tracking the maximum height that's currently in `finalized_state_queued_blocks`
187    ///
188    /// Set to `f64::NAN` if `finalized_state_queued_blocks` is empty, because grafana shows NaNs
189    /// as a break in the graph.
190    max_finalized_queue_height: f64,
191}
192
193/// A read-only service for accessing Zebra's cached blockchain state.
194///
195/// This service provides read-only access to:
196/// - the non-finalized state: the most recent blocks, up to
197///   [`MAX_BLOCK_REORG_HEIGHT`](crate::MAX_BLOCK_REORG_HEIGHT) of them.
198/// - the finalized state: older blocks that have many confirmations.
199///
200/// Requests to this service are processed in parallel,
201/// ignoring any blocks queued by the read-write [`StateService`].
202///
203/// This quick response behavior is better for most state users.
204/// It allows other async tasks to make progress while concurrently reading data from disk.
205#[derive(Clone, Debug)]
206pub struct ReadStateService {
207    // Configuration
208    //
209    /// The configured Zcash network.
210    network: Network,
211
212    // Shared Concurrently Readable State
213    //
214    /// A watch channel with a cached copy of the [`NonFinalizedState`].
215    ///
216    /// This state is only updated between requests,
217    /// so it might include some block data that is also on `disk`.
218    non_finalized_state_receiver: WatchReceiver<NonFinalizedState>,
219
220    /// The shared inner on-disk database for the finalized state.
221    ///
222    /// RocksDB allows reads and writes via a shared reference,
223    /// but [`ZebraDb`] doesn't expose any write methods or types.
224    ///
225    /// This chain is updated concurrently with requests,
226    /// so it might include some block data that is also in `best_mem`.
227    db: ZebraDb,
228
229    /// A shared handle to a task that writes blocks to the [`NonFinalizedState`] or [`FinalizedState`],
230    /// once the queues have received all their parent blocks.
231    ///
232    /// Used to check for panics when writing blocks.
233    block_write_task: Option<Arc<std::thread::JoinHandle<()>>>,
234}
235
236impl Drop for StateService {
237    fn drop(&mut self) {
238        // The state service owns the state, tasks, and channels,
239        // so dropping it should shut down everything.
240
241        // Close the channels (non-blocking)
242        // This makes the block write thread exit the next time it checks the channels.
243        // We want to do this here so we get any errors or panics from the block write task before it shuts down.
244        self.invalid_block_write_reset_receiver.close();
245        self.non_finalized_rejected_receiver.close();
246
247        std::mem::drop(self.block_write_sender.finalized.take());
248        std::mem::drop(self.block_write_sender.non_finalized.take());
249
250        self.clear_finalized_block_queue(CommitBlockError::WriteTaskExited);
251        self.clear_non_finalized_block_queue(CommitBlockError::WriteTaskExited);
252
253        // Log database metrics before shutting down
254        info!("dropping the state: logging database metrics");
255        self.log_db_metrics();
256
257        // Then drop self.read_service, which checks the block write task for panics,
258        // and tries to shut down the database.
259    }
260}
261
262impl Drop for ReadStateService {
263    fn drop(&mut self) {
264        // The read state service shares the state,
265        // so dropping it should check if we can shut down.
266
267        // TODO: move this into a try_shutdown() method
268        if let Some(block_write_task) = self.block_write_task.take() {
269            if let Some(block_write_task_handle) = Arc::into_inner(block_write_task) {
270                // We're the last database user, so we can tell it to shut down (blocking):
271                // - flushes the database to disk, and
272                // - drops the database, which cleans up any database tasks correctly.
273                self.db.shutdown(true);
274
275                // We are the last state with a reference to this thread, so we can
276                // wait until the block write task finishes, then check for panics (blocking).
277                // (We'd also like to abort the thread, but std::thread::JoinHandle can't do that.)
278
279                // This log is verbose during tests.
280                #[cfg(not(test))]
281                info!("waiting for the block write task to finish");
282                #[cfg(test)]
283                debug!("waiting for the block write task to finish");
284
285                // TODO: move this into a check_for_panics() method
286                if let Err(thread_panic) = block_write_task_handle.join() {
287                    std::panic::resume_unwind(thread_panic);
288                } else {
289                    debug!("shutting down the state because the block write task has finished");
290                }
291            }
292        } else {
293            // Even if we're not the last database user, try shutting it down.
294            //
295            // TODO: rename this to try_shutdown()?
296            self.db.shutdown(false);
297        }
298    }
299}
300
301impl StateService {
302    const PRUNE_INTERVAL: Duration = Duration::from_secs(30);
303
304    /// Creates a new state service for the state `config` and `network`.
305    ///
306    /// Uses the `max_checkpoint_height` and `checkpoint_verify_concurrency_limit`
307    /// to work out when it is near the final checkpoint.
308    ///
309    /// Returns the read-write and read-only state services,
310    /// and read-only watch channels for its best chain tip.
311    pub async fn new(
312        config: Config,
313        network: &Network,
314        max_checkpoint_height: block::Height,
315        checkpoint_verify_concurrency_limit: usize,
316    ) -> (Self, ReadStateService, LatestChainTip, ChainTipChange) {
317        let (finalized_state, finalized_tip, timer) = {
318            let config = config.clone();
319            let network = network.clone();
320            tokio::task::spawn_blocking(move || {
321                let timer = CodeTimer::start();
322                let finalized_state = FinalizedState::new(
323                    &config,
324                    &network,
325                    #[cfg(feature = "elasticsearch")]
326                    true,
327                )
328                .expect(
329                    "opening the read-write finalized state database failed; check that the \
330                     state cache directory is writable and not locked by another Zebra instance, \
331                     and that there is free disk space",
332                );
333                timer.finish_desc("opening finalized state database");
334
335                let timer = CodeTimer::start();
336                let finalized_tip = finalized_state.db.tip_block();
337
338                (finalized_state, finalized_tip, timer)
339            })
340            .await
341            .expect("failed to join blocking task")
342        };
343
344        // # Correctness
345        //
346        // The state service must set the finalized block write sender to `None`
347        // if there are blocks in the restored non-finalized state that are above
348        // the max checkpoint height so that non-finalized blocks can be written, otherwise,
349        // Zebra will be unable to commit semantically verified blocks, and its chain sync will stall.
350        //
351        // The state service must not set the finalized block write sender to `None` if there
352        // aren't blocks in the restored non-finalized state that are above the max checkpoint height,
353        // otherwise, unless checkpoint sync is disabled in the zebra-consensus configuration,
354        // Zebra will be unable to commit checkpoint verified blocks, and its chain sync will stall.
355        let is_finalized_tip_past_max_checkpoint = if let Some(tip) = &finalized_tip {
356            tip.coinbase_height().expect("valid block must have height") >= max_checkpoint_height
357        } else {
358            false
359        };
360        let backup_dir_path = config.non_finalized_state_backup_dir(network);
361        let skip_backup_task = config.debug_skip_non_finalized_state_backup_task;
362        let (non_finalized_state, non_finalized_state_sender, non_finalized_state_receiver) =
363            NonFinalizedState::new(network)
364                .with_backup(
365                    backup_dir_path.clone(),
366                    &finalized_state.db,
367                    is_finalized_tip_past_max_checkpoint,
368                    config.debug_skip_non_finalized_state_backup_task,
369                )
370                .await;
371
372        let non_finalized_block_write_sent_hashes = SentHashes::new(&non_finalized_state);
373        let initial_tip = non_finalized_state
374            .best_tip_block()
375            .map(|cv_block| cv_block.block.clone())
376            .or(finalized_tip)
377            .map(CheckpointVerifiedBlock::from)
378            .map(ChainTipBlock::from);
379
380        tracing::info!(chain_tip = ?initial_tip.as_ref().map(|tip| (tip.hash, tip.height)), "loaded Zebra state cache");
381
382        let (chain_tip_sender, latest_chain_tip, chain_tip_change) =
383            ChainTipSender::new(initial_tip, network);
384
385        let finalized_state_for_writing = finalized_state.clone();
386        let should_use_finalized_block_write_sender = non_finalized_state.is_chain_set_empty();
387        let sync_backup_dir_path = backup_dir_path.filter(|_| skip_backup_task);
388        let (
389            block_write_sender,
390            invalid_block_write_reset_receiver,
391            non_finalized_rejected_receiver,
392            block_write_task,
393        ) = write::BlockWriteSender::spawn(
394            finalized_state_for_writing,
395            non_finalized_state,
396            chain_tip_sender,
397            non_finalized_state_sender,
398            should_use_finalized_block_write_sender,
399            sync_backup_dir_path,
400        );
401
402        let read_service = ReadStateService::new(
403            &finalized_state,
404            block_write_task,
405            non_finalized_state_receiver,
406        );
407
408        let full_verifier_utxo_lookahead = max_checkpoint_height
409            - HeightDiff::try_from(checkpoint_verify_concurrency_limit)
410                .expect("fits in HeightDiff");
411        let full_verifier_utxo_lookahead =
412            full_verifier_utxo_lookahead.unwrap_or(block::Height::MIN);
413        let non_finalized_state_queued_blocks = QueuedBlocks::default();
414        let pending_utxos = PendingUtxos::default();
415
416        let finalized_block_write_last_sent_hash =
417            tokio::task::spawn_blocking(move || finalized_state.db.finalized_tip_hash())
418                .await
419                .expect("failed to join blocking task");
420
421        let state = Self {
422            network: network.clone(),
423            full_verifier_utxo_lookahead,
424            non_finalized_state_queued_blocks,
425            finalized_state_queued_blocks: HashMap::new(),
426            block_write_sender,
427            finalized_block_write_last_sent_hash,
428            non_finalized_block_write_sent_hashes,
429            invalid_block_write_reset_receiver,
430            non_finalized_rejected_receiver,
431            pending_utxos,
432            last_prune: Instant::now(),
433            read_service: read_service.clone(),
434            max_finalized_queue_height: f64::NAN,
435        };
436        timer.finish_desc("initializing state service");
437
438        tracing::info!("starting legacy chain check");
439        let timer = CodeTimer::start();
440
441        if let (Some(tip), Some(nu5_activation_height)) = (
442            {
443                let read_state = state.read_service.clone();
444                tokio::task::spawn_blocking(move || read_state.best_tip())
445                    .await
446                    .expect("task should not panic")
447            },
448            NetworkUpgrade::Nu5.activation_height(network),
449        ) {
450            if let Err(error) = check::legacy_chain(
451                nu5_activation_height,
452                any_ancestor_blocks(
453                    &state.read_service.latest_non_finalized_state(),
454                    &state.read_service.db,
455                    tip.1,
456                ),
457                &state.network,
458                MAX_LEGACY_CHAIN_BLOCKS,
459            ) {
460                let legacy_db_path = state.read_service.db.path().to_path_buf();
461                panic!(
462                    "Cached state contains a legacy chain.\n\
463                     An outdated Zebra version did not know about a recent network upgrade,\n\
464                     so it followed a legacy chain using outdated consensus branch rules.\n\
465                     Hint: Delete your database, and restart Zebra to do a full sync.\n\
466                     Database path: {legacy_db_path:?}\n\
467                     Error: {error:?}",
468                );
469            }
470        }
471
472        tracing::info!("cached state consensus branch is valid: no legacy chain found");
473        timer.finish_desc("legacy chain check");
474
475        // Spawn a background task to periodically export RocksDB metrics to Prometheus
476        let db_for_metrics = read_service.db.clone();
477        tokio::spawn(async move {
478            let mut interval = tokio::time::interval(Duration::from_secs(30));
479            loop {
480                interval.tick().await;
481                db_for_metrics.export_metrics();
482            }
483        });
484
485        (state, read_service, latest_chain_tip, chain_tip_change)
486    }
487
488    /// Call read only state service to log rocksdb database metrics.
489    pub fn log_db_metrics(&self) {
490        self.read_service.db.print_db_metrics();
491    }
492
493    /// Queue a checkpoint verified block for verification and storage in the finalized state.
494    ///
495    /// Returns a channel receiver that provides the result of the block commit.
496    fn queue_and_commit_to_finalized_state(
497        &mut self,
498        checkpoint_verified: CheckpointVerifiedBlock,
499    ) -> oneshot::Receiver<Result<block::Hash, CommitCheckpointVerifiedError>> {
500        // # Correctness & Performance
501        //
502        // This method must not block, access the database, or perform CPU-intensive tasks,
503        // because it is called directly from the tokio executor's Future threads.
504
505        let queued_prev_hash = checkpoint_verified.block.header.previous_block_hash;
506        let queued_height = checkpoint_verified.height;
507
508        // If we're close to the final checkpoint, make the block's UTXOs available for
509        // semantic block verification, even when it is in the channel.
510        if self.is_close_to_final_checkpoint(queued_height) {
511            self.non_finalized_block_write_sent_hashes
512                .add_finalized(&checkpoint_verified)
513        }
514
515        let (rsp_tx, rsp_rx) = oneshot::channel();
516        let queued = (checkpoint_verified, rsp_tx);
517
518        if self.block_write_sender.finalized.is_some() {
519            // We're still committing checkpoint verified blocks
520            if let Some(duplicate_queued) = self
521                .finalized_state_queued_blocks
522                .insert(queued_prev_hash, queued)
523            {
524                Self::send_checkpoint_verified_block_error(
525                    duplicate_queued,
526                    CommitBlockError::new_duplicate(
527                        Some(queued_prev_hash.into()),
528                        KnownBlock::Queue,
529                    ),
530                );
531            }
532
533            self.drain_finalized_queue_and_commit();
534        } else {
535            // We've finished committing checkpoint verified blocks to the finalized state,
536            // so drop any repeated queued blocks, and return an error.
537            //
538            // TODO: track the latest sent height, and drop any blocks under that height
539            //       every time we send some blocks (like QueuedSemanticallyVerifiedBlocks)
540            Self::send_checkpoint_verified_block_error(
541                queued,
542                CommitBlockError::new_duplicate(None, KnownBlock::Finalized),
543            );
544
545            self.clear_finalized_block_queue(CommitBlockError::new_duplicate(
546                None,
547                KnownBlock::Finalized,
548            ));
549        }
550
551        if self.finalized_state_queued_blocks.is_empty() {
552            self.max_finalized_queue_height = f64::NAN;
553        } else if self.max_finalized_queue_height.is_nan()
554            || self.max_finalized_queue_height < queued_height.0 as f64
555        {
556            // if there are still blocks in the queue, then either:
557            //   - the new block was lower than the old maximum, and there was a gap before it,
558            //     so the maximum is still the same (and we skip this code), or
559            //   - the new block is higher than the old maximum, and there is at least one gap
560            //     between the finalized tip and the new maximum
561            self.max_finalized_queue_height = queued_height.0 as f64;
562        }
563
564        metrics::gauge!("state.checkpoint.queued.max.height").set(self.max_finalized_queue_height);
565        metrics::gauge!("state.checkpoint.queued.block.count")
566            .set(self.finalized_state_queued_blocks.len() as f64);
567
568        rsp_rx
569    }
570
571    /// Finds finalized state queue blocks to be committed to the state in order,
572    /// removes them from the queue, and sends them to the block commit task.
573    ///
574    /// After queueing a finalized block, this method checks whether the newly
575    /// queued block (and any of its descendants) can be committed to the state.
576    ///
577    /// Returns an error if the block commit channel has been closed.
578    pub fn drain_finalized_queue_and_commit(&mut self) {
579        use tokio::sync::mpsc::error::{SendError, TryRecvError};
580
581        // # Correctness & Performance
582        //
583        // This method must not block, access the database, or perform CPU-intensive tasks,
584        // because it is called directly from the tokio executor's Future threads.
585
586        // If a block failed, we need to start again from a valid tip.
587        match self.invalid_block_write_reset_receiver.try_recv() {
588            Ok(reset_tip_hash) => self.finalized_block_write_last_sent_hash = reset_tip_hash,
589            Err(TryRecvError::Disconnected) => {
590                info!("Block commit task closed the block reset channel. Is Zebra shutting down?");
591                return;
592            }
593            // There are no errors, so we can just use the last block hash we sent
594            Err(TryRecvError::Empty) => {}
595        }
596
597        while let Some(queued_block) = self
598            .finalized_state_queued_blocks
599            .remove(&self.finalized_block_write_last_sent_hash)
600        {
601            let last_sent_finalized_block_height = queued_block.0.height;
602
603            self.finalized_block_write_last_sent_hash = queued_block.0.hash;
604
605            // If we've finished sending finalized blocks, ignore any repeated blocks.
606            // (Blocks can be repeated after a syncer reset.)
607            if let Some(finalized_block_write_sender) = &self.block_write_sender.finalized {
608                let send_result = finalized_block_write_sender.send(queued_block);
609
610                // If the receiver is closed, we can't send any more blocks.
611                if let Err(SendError(queued)) = send_result {
612                    // If Zebra is shutting down, drop blocks and return an error.
613                    Self::send_checkpoint_verified_block_error(
614                        queued,
615                        CommitBlockError::WriteTaskExited,
616                    );
617
618                    self.clear_finalized_block_queue(CommitBlockError::WriteTaskExited);
619                } else {
620                    metrics::gauge!("state.checkpoint.sent.block.height")
621                        .set(last_sent_finalized_block_height.0 as f64);
622                };
623            }
624        }
625    }
626
627    /// Drains every hash queued on `non_finalized_rejected_receiver` and
628    /// removes it from `non_finalized_block_write_sent_hashes`.
629    ///
630    /// This closes the lockout window where a rejected block keeps its hash
631    /// recorded as "sent", so a subsequent honest re-delivery of a block at
632    /// the same hash is not short-circuited as a false "duplicate".
633    ///
634    /// # Correctness & Performance
635    ///
636    /// Like the other drain methods on `StateService`, this must not block,
637    /// access the database, or perform CPU-intensive work, because it is
638    /// called directly from the tokio executor's Future threads.
639    fn drain_non_finalized_rejected_hashes(&mut self) {
640        use tokio::sync::mpsc::error::TryRecvError;
641
642        loop {
643            match self.non_finalized_rejected_receiver.try_recv() {
644                Ok(hash) => {
645                    self.non_finalized_block_write_sent_hashes.remove(&hash);
646                }
647                Err(TryRecvError::Empty) => break,
648                Err(TryRecvError::Disconnected) => {
649                    info!(
650                        "Block commit task closed the non-finalized rejected hash channel. \
651                         Is Zebra shutting down?"
652                    );
653                    break;
654                }
655            }
656        }
657    }
658
659    /// Drops all finalized state queue blocks, and sends an error on their result channels.
660    fn clear_finalized_block_queue(
661        &mut self,
662        error: impl Into<CommitCheckpointVerifiedError> + Clone,
663    ) {
664        for (_hash, queued) in self.finalized_state_queued_blocks.drain() {
665            Self::send_checkpoint_verified_block_error(queued, error.clone());
666        }
667    }
668
669    /// Send an error on a `QueuedCheckpointVerified` block's result channel, and drop the block
670    fn send_checkpoint_verified_block_error(
671        queued: QueuedCheckpointVerified,
672        error: impl Into<CommitCheckpointVerifiedError>,
673    ) {
674        let (finalized, rsp_tx) = queued;
675
676        // The block sender might have already given up on this block,
677        // so ignore any channel send errors.
678        let _ = rsp_tx.send(Err(error.into()));
679        std::mem::drop(finalized);
680    }
681
682    /// Drops all non-finalized state queue blocks, and sends an error on their result channels.
683    fn clear_non_finalized_block_queue(
684        &mut self,
685        error: impl Into<CommitSemanticallyVerifiedError> + Clone,
686    ) {
687        for (_hash, queued) in self.non_finalized_state_queued_blocks.drain() {
688            Self::send_semantically_verified_block_error(queued, error.clone());
689        }
690    }
691
692    /// Send an error on a `QueuedSemanticallyVerified` block's result channel, and drop the block
693    fn send_semantically_verified_block_error(
694        queued: QueuedSemanticallyVerified,
695        error: impl Into<CommitSemanticallyVerifiedError>,
696    ) {
697        let (finalized, rsp_tx) = queued;
698
699        // The block sender might have already given up on this block,
700        // so ignore any channel send errors.
701        let _ = rsp_tx.send(Err(error.into()));
702        std::mem::drop(finalized);
703    }
704
705    /// Queue a semantically verified block for contextual verification and check if any queued
706    /// blocks are ready to be verified and committed to the state.
707    ///
708    /// This function encodes the logic for [committing non-finalized blocks][1]
709    /// in RFC0005.
710    ///
711    /// [1]: https://zebra.zfnd.org/dev/rfcs/0005-state-updates.html#committing-non-finalized-blocks
712    #[instrument(level = "debug", skip(self, semantically_verified))]
713    fn queue_and_commit_to_non_finalized_state(
714        &mut self,
715        semantically_verified: SemanticallyVerifiedBlock,
716    ) -> oneshot::Receiver<Result<block::Hash, CommitSemanticallyVerifiedError>> {
717        tracing::debug!(block = %semantically_verified.block, "queueing block for contextual verification");
718        let parent_hash = semantically_verified.block.header.previous_block_hash;
719
720        // Drop hashes of any blocks the write task has rejected before checking
721        // the SentHashes membership below. Without this, a rejected same-hash
722        // block would lock out a later honest re-delivery of a block at the
723        // same hash as a false "duplicate".
724        self.drain_non_finalized_rejected_hashes();
725
726        if self
727            .non_finalized_block_write_sent_hashes
728            .contains(&semantically_verified.hash)
729        {
730            let (rsp_tx, rsp_rx) = oneshot::channel();
731            let _ = rsp_tx.send(Err(CommitBlockError::new_duplicate(
732                Some(semantically_verified.hash.into()),
733                KnownBlock::WriteChannel,
734            )
735            .into()));
736            return rsp_rx;
737        }
738
739        if self
740            .read_service
741            .db
742            .contains_height(semantically_verified.height)
743        {
744            let (rsp_tx, rsp_rx) = oneshot::channel();
745            let _ = rsp_tx.send(Err(CommitBlockError::new_duplicate(
746                Some(semantically_verified.height.into()),
747                KnownBlock::Finalized,
748            )
749            .into()));
750            return rsp_rx;
751        }
752
753        // [`Request::CommitSemanticallyVerifiedBlock`] contract: a request to commit a block which
754        // has been queued but not yet committed to the state fails the older request and replaces
755        // it with the newer request.
756        let rsp_rx = if let Some((_, old_rsp_tx)) = self
757            .non_finalized_state_queued_blocks
758            .get_mut(&semantically_verified.hash)
759        {
760            tracing::debug!("replacing older queued request with new request");
761            let (mut rsp_tx, rsp_rx) = oneshot::channel();
762            std::mem::swap(old_rsp_tx, &mut rsp_tx);
763            let _ = rsp_tx.send(Err(CommitBlockError::new_duplicate(
764                Some(semantically_verified.hash.into()),
765                KnownBlock::Queue,
766            )
767            .into()));
768            rsp_rx
769        } else {
770            let (rsp_tx, rsp_rx) = oneshot::channel();
771            self.non_finalized_state_queued_blocks
772                .queue((semantically_verified, rsp_tx));
773            rsp_rx
774        };
775
776        // We've finished sending checkpoint verified blocks when:
777        // - we've sent the verified block for the last checkpoint, and
778        // - it has been successfully written to disk.
779        //
780        // We detect the last checkpoint by looking for non-finalized blocks
781        // that are a child of the last block we sent.
782        //
783        // TODO: configure the state with the last checkpoint hash instead?
784        if self.block_write_sender.finalized.is_some()
785            && self
786                .non_finalized_state_queued_blocks
787                .has_queued_children(self.finalized_block_write_last_sent_hash)
788            && self.read_service.db.finalized_tip_hash()
789                == self.finalized_block_write_last_sent_hash
790        {
791            // Tell the block write task to stop committing checkpoint verified blocks to the finalized state,
792            // and move on to committing semantically verified blocks to the non-finalized state.
793            std::mem::drop(self.block_write_sender.finalized.take());
794            // Remove any checkpoint-verified block hashes from `non_finalized_block_write_sent_hashes`.
795            self.non_finalized_block_write_sent_hashes = SentHashes::default();
796            // Mark `SentHashes` as usable by the `can_fork_chain_at()` method.
797            self.non_finalized_block_write_sent_hashes
798                .can_fork_chain_at_hashes = true;
799            // Send blocks from non-finalized queue
800            self.send_ready_non_finalized_queued(self.finalized_block_write_last_sent_hash);
801            // We've finished committing checkpoint verified blocks to finalized state, so drop any repeated queued blocks.
802            self.clear_finalized_block_queue(CommitBlockError::new_duplicate(
803                None,
804                KnownBlock::Finalized,
805            ));
806        } else if !self.can_fork_chain_at(&parent_hash) {
807            tracing::trace!("unready to verify, returning early");
808        } else if self.block_write_sender.finalized.is_none() {
809            // Wait until block commit task is ready to write non-finalized blocks before dequeuing them
810            self.send_ready_non_finalized_queued(parent_hash);
811
812            let finalized_tip_height = self.read_service.db.finalized_tip_height().expect(
813                "Finalized state must have at least one block before committing non-finalized state",
814            );
815
816            self.non_finalized_state_queued_blocks
817                .prune_by_height(finalized_tip_height);
818
819            self.non_finalized_block_write_sent_hashes
820                .prune_by_height(finalized_tip_height);
821        }
822
823        rsp_rx
824    }
825
826    /// Returns `true` if `hash` is a valid previous block hash for new non-finalized blocks.
827    fn can_fork_chain_at(&self, hash: &block::Hash) -> bool {
828        self.non_finalized_block_write_sent_hashes
829            .can_fork_chain_at(hash)
830            || &self.read_service.db.finalized_tip_hash() == hash
831    }
832
833    /// Returns `true` if `queued_height` is near the final checkpoint.
834    ///
835    /// The semantic block verifier needs access to UTXOs from checkpoint verified blocks
836    /// near the final checkpoint, so that it can verify blocks that spend those UTXOs.
837    ///
838    /// If it doesn't have the required UTXOs, some blocks will time out,
839    /// but succeed after a syncer restart.
840    fn is_close_to_final_checkpoint(&self, queued_height: block::Height) -> bool {
841        queued_height >= self.full_verifier_utxo_lookahead
842    }
843
844    /// Sends all queued blocks whose parents have recently arrived starting from `new_parent`
845    /// in breadth-first ordering to the block write task which will attempt to validate and commit them
846    #[tracing::instrument(level = "debug", skip(self, new_parent))]
847    fn send_ready_non_finalized_queued(&mut self, new_parent: block::Hash) {
848        use tokio::sync::mpsc::error::SendError;
849        if let Some(non_finalized_block_write_sender) = &self.block_write_sender.non_finalized {
850            let mut new_parents: Vec<block::Hash> = vec![new_parent];
851
852            while let Some(parent_hash) = new_parents.pop() {
853                let queued_children = self
854                    .non_finalized_state_queued_blocks
855                    .dequeue_children(parent_hash);
856
857                for queued_child in queued_children {
858                    let (SemanticallyVerifiedBlock { hash, .. }, _) = queued_child;
859
860                    self.non_finalized_block_write_sent_hashes
861                        .add(&queued_child.0);
862                    let send_result = non_finalized_block_write_sender.send(queued_child.into());
863
864                    if let Err(SendError(NonFinalizedWriteMessage::Commit(queued))) = send_result {
865                        // If Zebra is shutting down, drop blocks and return an error.
866                        Self::send_semantically_verified_block_error(
867                            queued,
868                            CommitBlockError::WriteTaskExited,
869                        );
870
871                        self.clear_non_finalized_block_queue(CommitBlockError::WriteTaskExited);
872
873                        return;
874                    };
875
876                    new_parents.push(hash);
877                }
878            }
879
880            self.non_finalized_block_write_sent_hashes.finish_batch();
881        };
882    }
883
884    /// Return the tip of the current best chain.
885    pub fn best_tip(&self) -> Option<(block::Height, block::Hash)> {
886        self.read_service.best_tip()
887    }
888
889    fn send_invalidate_block(
890        &self,
891        hash: block::Hash,
892    ) -> oneshot::Receiver<Result<block::Hash, InvalidateError>> {
893        let (rsp_tx, rsp_rx) = oneshot::channel();
894
895        let Some(sender) = &self.block_write_sender.non_finalized else {
896            let _ = rsp_tx.send(Err(InvalidateError::ProcessingCheckpointedBlocks));
897            return rsp_rx;
898        };
899
900        if let Err(tokio::sync::mpsc::error::SendError(error)) =
901            sender.send(NonFinalizedWriteMessage::Invalidate { hash, rsp_tx })
902        {
903            let NonFinalizedWriteMessage::Invalidate { rsp_tx, .. } = error else {
904                unreachable!("should return the same Invalidate message could not be sent");
905            };
906
907            let _ = rsp_tx.send(Err(InvalidateError::SendInvalidateRequestFailed));
908        }
909
910        rsp_rx
911    }
912
913    fn send_reconsider_block(
914        &self,
915        hash: block::Hash,
916    ) -> oneshot::Receiver<Result<Vec<block::Hash>, ReconsiderError>> {
917        let (rsp_tx, rsp_rx) = oneshot::channel();
918
919        let Some(sender) = &self.block_write_sender.non_finalized else {
920            let _ = rsp_tx.send(Err(ReconsiderError::CheckpointCommitInProgress));
921            return rsp_rx;
922        };
923
924        if let Err(tokio::sync::mpsc::error::SendError(error)) =
925            sender.send(NonFinalizedWriteMessage::Reconsider { hash, rsp_tx })
926        {
927            let NonFinalizedWriteMessage::Reconsider { rsp_tx, .. } = error else {
928                unreachable!("should return the same Reconsider message could not be sent");
929            };
930
931            let _ = rsp_tx.send(Err(ReconsiderError::ReconsiderSendFailed));
932        }
933
934        rsp_rx
935    }
936
937    /// Assert some assumptions about the semantically verified `block` before it is queued.
938    fn assert_block_can_be_validated(&self, block: &SemanticallyVerifiedBlock) {
939        // required by `Request::CommitSemanticallyVerifiedBlock` call
940        assert!(
941            block.height > self.network.mandatory_checkpoint_height(),
942            "invalid semantically verified block height: the canopy checkpoint is mandatory, pre-canopy \
943            blocks, and the canopy activation block, must be committed to the state as finalized \
944            blocks"
945        );
946    }
947
948    fn known_sent_hash(&self, hash: &block::Hash) -> Option<KnownBlock> {
949        self.non_finalized_block_write_sent_hashes
950            .contains(hash)
951            .then_some(KnownBlock::WriteChannel)
952    }
953}
954
955impl ReadStateService {
956    /// Creates a new read-only state service, using the provided finalized state and
957    /// block write task handle.
958    ///
959    /// Returns the newly created service,
960    /// and a watch channel for updating the shared recent non-finalized chain.
961    pub(crate) fn new(
962        finalized_state: &FinalizedState,
963        block_write_task: Option<Arc<std::thread::JoinHandle<()>>>,
964        non_finalized_state_receiver: WatchReceiver<NonFinalizedState>,
965    ) -> Self {
966        let read_service = Self {
967            network: finalized_state.network(),
968            db: finalized_state.db.clone(),
969            non_finalized_state_receiver,
970            block_write_task,
971        };
972
973        tracing::debug!("created new read-only state service");
974
975        read_service
976    }
977
978    /// Return the tip of the current best chain.
979    pub fn best_tip(&self) -> Option<(block::Height, block::Hash)> {
980        read::best_tip(&self.latest_non_finalized_state(), &self.db)
981    }
982
983    /// Gets a clone of the latest non-finalized state from the `non_finalized_state_receiver`
984    fn latest_non_finalized_state(&self) -> NonFinalizedState {
985        self.non_finalized_state_receiver.cloned_watch_data()
986    }
987
988    /// Gets a clone of the latest, best non-finalized chain from the `non_finalized_state_receiver`
989    fn latest_best_chain(&self) -> Option<Arc<Chain>> {
990        self.non_finalized_state_receiver
991            .borrow_mapped(|non_finalized_state| non_finalized_state.best_chain().cloned())
992    }
993
994    /// Test-only access to the inner database.
995    /// Can be used to modify the database without doing any consensus checks.
996    #[cfg(any(test, feature = "proptest-impl"))]
997    pub fn db(&self) -> &ZebraDb {
998        &self.db
999    }
1000
1001    /// Logs rocksdb metrics using the read only state service.
1002    pub fn log_db_metrics(&self) {
1003        self.db.print_db_metrics();
1004    }
1005}
1006
1007impl Service<Request> for StateService {
1008    type Response = Response;
1009    type Error = BoxError;
1010    type Future =
1011        Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send + 'static>>;
1012
1013    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1014        // Check for panics in the block write task
1015        let poll = self.read_service.poll_ready(cx);
1016
1017        // Prune outdated UTXO requests
1018        let now = Instant::now();
1019
1020        if self.last_prune + Self::PRUNE_INTERVAL < now {
1021            let tip = self.best_tip();
1022            let old_len = self.pending_utxos.len();
1023
1024            self.pending_utxos.prune();
1025            self.last_prune = now;
1026
1027            let new_len = self.pending_utxos.len();
1028            let prune_count = old_len
1029                .checked_sub(new_len)
1030                .expect("prune does not add any utxo requests");
1031            if prune_count > 0 {
1032                tracing::debug!(
1033                    ?old_len,
1034                    ?new_len,
1035                    ?prune_count,
1036                    ?tip,
1037                    "pruned utxo requests"
1038                );
1039            } else {
1040                tracing::debug!(len = ?old_len, ?tip, "no utxo requests needed pruning");
1041            }
1042        }
1043
1044        poll
1045    }
1046
1047    #[instrument(name = "state", skip(self, req))]
1048    fn call(&mut self, req: Request) -> Self::Future {
1049        req.count_metric();
1050        let span = Span::current();
1051
1052        match req {
1053            // Uses non_finalized_state_queued_blocks and pending_utxos in the StateService
1054            // Accesses shared writeable state in the StateService, NonFinalizedState, and ZebraDb.
1055            //
1056            // The expected error type for this request is `CommitSemanticallyVerifiedError`.
1057            Request::CommitSemanticallyVerifiedBlock(semantically_verified) => {
1058                let timer = CodeTimer::start();
1059                self.assert_block_can_be_validated(&semantically_verified);
1060
1061                self.pending_utxos
1062                    .check_against_ordered(&semantically_verified.new_outputs);
1063
1064                // # Performance
1065                //
1066                // Allow other async tasks to make progress while blocks are being verified
1067                // and written to disk. But wait for the blocks to finish committing,
1068                // so that `StateService` multi-block queries always observe a consistent state.
1069                //
1070                // Since each block is spawned into its own task,
1071                // there shouldn't be any other code running in the same task,
1072                // so we don't need to worry about blocking it:
1073                // https://docs.rs/tokio/latest/tokio/task/fn.block_in_place.html
1074
1075                let rsp_rx = tokio::task::block_in_place(move || {
1076                    span.in_scope(|| {
1077                        self.queue_and_commit_to_non_finalized_state(semantically_verified)
1078                    })
1079                });
1080
1081                // TODO:
1082                //   - check for panics in the block write task here,
1083                //     as well as in poll_ready()
1084
1085                // The work is all done, the future just waits on a channel for the result
1086                timer.finish_desc("CommitSemanticallyVerifiedBlock");
1087
1088                // Await the channel response, flatten the result, map receive errors to
1089                // `CommitSemanticallyVerifiedError::WriteTaskExited`.
1090                // Then flatten the nested Result and convert any errors to a BoxError.
1091                let span = Span::current();
1092                async move {
1093                    rsp_rx
1094                        .await
1095                        .map_err(|_recv_error| CommitBlockError::WriteTaskExited.into())
1096                        .and_then(|result| result)
1097                        .map_err(BoxError::from)
1098                        .map(Response::Committed)
1099                }
1100                .instrument(span)
1101                .boxed()
1102            }
1103
1104            // Uses finalized_state_queued_blocks and pending_utxos in the StateService.
1105            // Accesses shared writeable state in the StateService.
1106            //
1107            // The expected error type for this request is `CommitCheckpointVerifiedError`.
1108            Request::CommitCheckpointVerifiedBlock(finalized) => {
1109                let timer = CodeTimer::start();
1110                // # Consensus
1111                //
1112                // A semantic block verification could have called AwaitUtxo
1113                // before this checkpoint verified block arrived in the state.
1114                // So we need to check for pending UTXO requests sent by running
1115                // semantic block verifications.
1116                //
1117                // This check is redundant for most checkpoint verified blocks,
1118                // because semantic verification can only succeed near the final
1119                // checkpoint, when all the UTXOs are available for the verifying block.
1120                //
1121                // (Checkpoint block UTXOs are verified using block hash checkpoints
1122                // and transaction merkle tree block header commitments.)
1123                self.pending_utxos
1124                    .check_against_ordered(&finalized.new_outputs);
1125
1126                // # Performance
1127                //
1128                // This method doesn't block, access the database, or perform CPU-intensive tasks,
1129                // so we can run it directly in the tokio executor's Future threads.
1130                let rsp_rx = self.queue_and_commit_to_finalized_state(finalized);
1131
1132                // TODO:
1133                //   - check for panics in the block write task here,
1134                //     as well as in poll_ready()
1135
1136                // The work is all done, the future just waits on a channel for the result
1137                timer.finish_desc("CommitCheckpointVerifiedBlock");
1138
1139                // Await the channel response, flatten the result, map receive errors to
1140                // `CommitCheckpointVerifiedError::WriteTaskExited`.
1141                // Then flatten the nested Result and convert any errors to a BoxError.
1142                async move {
1143                    rsp_rx
1144                        .await
1145                        .map_err(|_recv_error| CommitBlockError::WriteTaskExited.into())
1146                        .and_then(|result| result)
1147                        .map_err(BoxError::from)
1148                        .map(Response::Committed)
1149                }
1150                .instrument(span)
1151                .boxed()
1152            }
1153
1154            // Uses pending_utxos and non_finalized_state_queued_blocks in the StateService.
1155            // If the UTXO isn't in the queued blocks, runs concurrently using the ReadStateService.
1156            Request::AwaitUtxo(outpoint) => {
1157                let timer = CodeTimer::start();
1158                // Prepare the AwaitUtxo future from PendingUxtos.
1159                let response_fut = self.pending_utxos.queue(outpoint);
1160                // Only instrument `response_fut`, the ReadStateService already
1161                // instruments its requests with the same span.
1162
1163                let response_fut = response_fut.instrument(span).boxed();
1164
1165                // Check the non-finalized block queue outside the returned future,
1166                // so we can access mutable state fields.
1167                if let Some(utxo) = self.non_finalized_state_queued_blocks.utxo(&outpoint) {
1168                    self.pending_utxos.respond(&outpoint, utxo);
1169
1170                    // We're finished, the returned future gets the UTXO from the respond() channel.
1171                    timer.finish_desc("AwaitUtxo/queued-non-finalized");
1172
1173                    return response_fut;
1174                }
1175
1176                // Check the sent non-finalized blocks
1177                self.drain_non_finalized_rejected_hashes();
1178
1179                if let Some(utxo) = self.non_finalized_block_write_sent_hashes.utxo(&outpoint) {
1180                    self.pending_utxos.respond(&outpoint, utxo);
1181
1182                    // We're finished, the returned future gets the UTXO from the respond() channel.
1183                    timer.finish_desc("AwaitUtxo/sent-non-finalized");
1184
1185                    return response_fut;
1186                }
1187
1188                // We ignore any UTXOs in FinalizedState.finalized_state_queued_blocks,
1189                // because it is only used during checkpoint verification.
1190                //
1191                // This creates a rare race condition, but it doesn't seem to happen much in practice.
1192                // See #5126 for details.
1193
1194                // Manually send a request to the ReadStateService,
1195                // to get UTXOs from any non-finalized chain or the finalized chain.
1196                let read_service = self.read_service.clone();
1197
1198                // Run the request in an async block, so we can await the response.
1199                async move {
1200                    let req = ReadRequest::AnyChainUtxo(outpoint);
1201
1202                    let rsp = read_service.oneshot(req).await?;
1203
1204                    // Optional TODO:
1205                    //  - make pending_utxos.respond() async using a channel,
1206                    //    so we can respond to all waiting requests here
1207                    //
1208                    // This change is not required for correctness, because:
1209                    // - any waiting requests should have returned when the block was sent to the state
1210                    // - otherwise, the request returns immediately if:
1211                    //   - the block is in the non-finalized queue, or
1212                    //   - the block is in any non-finalized chain or the finalized state
1213                    //
1214                    // And if the block is in the finalized queue,
1215                    // that's rare enough that a retry is ok.
1216                    if let ReadResponse::AnyChainUtxo(Some(utxo)) = rsp {
1217                        // We got a UTXO, so we replace the response future with the result own.
1218                        timer.finish_desc("AwaitUtxo/any-chain");
1219
1220                        return Ok(Response::Utxo(utxo));
1221                    }
1222
1223                    // We're finished, but the returned future is waiting on the respond() channel.
1224                    timer.finish_desc("AwaitUtxo/waiting");
1225
1226                    response_fut.await
1227                }
1228                .boxed()
1229            }
1230
1231            // Used by sync, inbound, and block verifier to check if a block is already in the state
1232            // before downloading or validating it.
1233            Request::KnownBlock(hash) => {
1234                let timer = CodeTimer::start();
1235
1236                self.drain_non_finalized_rejected_hashes();
1237
1238                let sent_hash_response = self.known_sent_hash(&hash);
1239                let read_service = self.read_service.clone();
1240
1241                async move {
1242                    if sent_hash_response.is_some() {
1243                        return Ok(Response::KnownBlock(sent_hash_response));
1244                    };
1245
1246                    let response = read::non_finalized_state_contains_block_hash(
1247                        &read_service.latest_non_finalized_state(),
1248                        hash,
1249                    )
1250                    // TODO: Move this to a blocking task, perhaps by moving some of this logic to the ReadStateService.
1251                    .or_else(|| read::finalized_state_contains_block_hash(&read_service.db, hash));
1252
1253                    timer.finish_desc("Request::KnownBlock");
1254
1255                    Ok(Response::KnownBlock(response))
1256                }
1257                .boxed()
1258            }
1259
1260            // The expected error type for this request is `InvalidateError`
1261            Request::InvalidateBlock(block_hash) => {
1262                let rsp_rx = tokio::task::block_in_place(move || {
1263                    span.in_scope(|| self.send_invalidate_block(block_hash))
1264                });
1265
1266                // Await the channel response, flatten the result, map receive errors to
1267                // `InvalidateError::InvalidateRequestDropped`.
1268                // Then flatten the nested Result and convert any errors to a BoxError.
1269                let span = Span::current();
1270                async move {
1271                    rsp_rx
1272                        .await
1273                        .map_err(|_recv_error| InvalidateError::InvalidateRequestDropped)
1274                        .and_then(|result| result)
1275                        .map_err(BoxError::from)
1276                        .map(Response::Invalidated)
1277                }
1278                .instrument(span)
1279                .boxed()
1280            }
1281
1282            // The expected error type for this request is `ReconsiderError`
1283            Request::ReconsiderBlock(block_hash) => {
1284                let rsp_rx = tokio::task::block_in_place(move || {
1285                    span.in_scope(|| self.send_reconsider_block(block_hash))
1286                });
1287
1288                // Await the channel response, flatten the result, map receive errors to
1289                // `ReconsiderError::ReconsiderResponseDropped`.
1290                // Then flatten the nested Result and convert any errors to a BoxError.
1291                let span = Span::current();
1292                async move {
1293                    rsp_rx
1294                        .await
1295                        .map_err(|_recv_error| ReconsiderError::ReconsiderResponseDropped)
1296                        .and_then(|result| result)
1297                        .map_err(BoxError::from)
1298                        .map(Response::Reconsidered)
1299                }
1300                .instrument(span)
1301                .boxed()
1302            }
1303
1304            // Runs concurrently using the ReadStateService
1305            Request::Tip
1306            | Request::Depth(_)
1307            | Request::BestChainNextMedianTimePast
1308            | Request::BestChainBlockHash(_)
1309            | Request::BlockLocator
1310            | Request::Transaction(_)
1311            | Request::AnyChainTransaction(_)
1312            | Request::UnspentBestChainUtxo(_)
1313            | Request::Block(_)
1314            | Request::AnyChainBlock(_)
1315            | Request::BlockAndSize(_)
1316            | Request::BlockHeader(_)
1317            | Request::FindBlockHashes { .. }
1318            | Request::FindBlockHeaders { .. }
1319            | Request::CheckBestChainTipNullifiersAndAnchors(_)
1320            | Request::CheckBlockProposalValidity(_) => {
1321                // Redirect the request to the concurrent ReadStateService
1322                let read_service = self.read_service.clone();
1323
1324                async move {
1325                    let req = req
1326                        .try_into()
1327                        .expect("ReadRequest conversion should not fail");
1328
1329                    let rsp = read_service.oneshot(req).await?;
1330                    let rsp = rsp.try_into().expect("Response conversion should not fail");
1331
1332                    Ok(rsp)
1333                }
1334                .boxed()
1335            }
1336        }
1337    }
1338}
1339
1340impl Service<ReadRequest> for ReadStateService {
1341    type Response = ReadResponse;
1342    type Error = BoxError;
1343    type Future =
1344        Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send + 'static>>;
1345
1346    fn poll_ready(&mut self, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1347        // Check for panics in the block write task
1348        //
1349        // TODO: move into a check_for_panics() method
1350        let block_write_task = self.block_write_task.take();
1351
1352        if let Some(block_write_task) = block_write_task {
1353            if block_write_task.is_finished() {
1354                if let Some(block_write_task) = Arc::into_inner(block_write_task) {
1355                    // We are the last state with a reference to this task, so we can propagate any panics
1356                    if let Err(thread_panic) = block_write_task.join() {
1357                        std::panic::resume_unwind(thread_panic);
1358                    }
1359                }
1360            } else {
1361                // It hasn't finished, so we need to put it back
1362                self.block_write_task = Some(block_write_task);
1363            }
1364        }
1365
1366        self.db.check_for_panics();
1367
1368        Poll::Ready(Ok(()))
1369    }
1370
1371    #[instrument(name = "read_state", skip(self, req))]
1372    fn call(&mut self, req: ReadRequest) -> Self::Future {
1373        req.count_metric();
1374        let timer = CodeTimer::start_desc(req.variant_name());
1375        let span = Span::current();
1376        let timed_span = TimedSpan::new(timer, span);
1377        let state = self.clone();
1378
1379        if let ReadRequest::NonFinalizedBlocksListener { known_chain_tips } = req {
1380            // The non-finalized blocks listener is used to notify the state service
1381            // about new blocks that have been added to the non-finalized state.
1382            let non_finalized_blocks_listener = NonFinalizedBlocksListener::spawn(
1383                self.non_finalized_state_receiver.clone(),
1384                known_chain_tips,
1385            );
1386
1387            return async move {
1388                Ok(ReadResponse::NonFinalizedBlocksListener(
1389                    non_finalized_blocks_listener,
1390                ))
1391            }
1392            .boxed();
1393        };
1394
1395        let request_handler = move || match req {
1396            // Used by the `getblockchaininfo` RPC.
1397            ReadRequest::UsageInfo => Ok(ReadResponse::UsageInfo(state.db.size())),
1398
1399            // Used by the StateService.
1400            ReadRequest::Tip => Ok(ReadResponse::Tip(read::tip(
1401                state.latest_best_chain(),
1402                &state.db,
1403            ))),
1404
1405            // Used by `getblockchaininfo` RPC method.
1406            ReadRequest::TipPoolValues => {
1407                let (tip_height, tip_hash, value_balance) =
1408                    read::tip_with_value_balance(state.latest_best_chain(), &state.db)?
1409                        .ok_or(BoxError::from("no chain tip available yet"))?;
1410
1411                Ok(ReadResponse::TipPoolValues {
1412                    tip_height,
1413                    tip_hash,
1414                    value_balance,
1415                })
1416            }
1417
1418            // Used by getblock
1419            ReadRequest::BlockInfo(hash_or_height) => Ok(ReadResponse::BlockInfo(
1420                read::block_info(state.latest_best_chain(), &state.db, hash_or_height),
1421            )),
1422
1423            // Used by the StateService.
1424            ReadRequest::Depth(hash) => Ok(ReadResponse::Depth(read::depth(
1425                state.latest_best_chain(),
1426                &state.db,
1427                hash,
1428            ))),
1429
1430            // Used by the StateService.
1431            ReadRequest::BestChainNextMedianTimePast => {
1432                Ok(ReadResponse::BestChainNextMedianTimePast(
1433                    read::next_median_time_past(&state.latest_non_finalized_state(), &state.db)?,
1434                ))
1435            }
1436
1437            // Used by the get_block (raw) RPC and the StateService.
1438            ReadRequest::Block(hash_or_height) => Ok(ReadResponse::Block(read::block(
1439                state.latest_best_chain(),
1440                &state.db,
1441                hash_or_height,
1442            ))),
1443
1444            ReadRequest::AnyChainBlock(hash_or_height) => Ok(ReadResponse::Block(read::any_block(
1445                state.latest_non_finalized_state().chain_iter(),
1446                &state.db,
1447                hash_or_height,
1448            ))),
1449
1450            // Used by the get_block (raw) RPC and the StateService.
1451            ReadRequest::BlockAndSize(hash_or_height) => Ok(ReadResponse::BlockAndSize(
1452                read::block_and_size(state.latest_best_chain(), &state.db, hash_or_height),
1453            )),
1454
1455            // Used by the get_block (verbose) RPC and the StateService.
1456            ReadRequest::BlockHeader(hash_or_height) => {
1457                let best_chain = state.latest_best_chain();
1458
1459                let height = hash_or_height
1460                    .height_or_else(|hash| {
1461                        read::find::height_by_hash(best_chain.clone(), &state.db, hash)
1462                    })
1463                    .ok_or_else(|| BoxError::from("block hash or height not found"))?;
1464
1465                let hash = hash_or_height
1466                    .hash_or_else(|height| {
1467                        read::find::hash_by_height(best_chain.clone(), &state.db, height)
1468                    })
1469                    .ok_or_else(|| BoxError::from("block hash or height not found"))?;
1470
1471                let next_height = height.next()?;
1472                let next_block_hash =
1473                    read::find::hash_by_height(best_chain.clone(), &state.db, next_height);
1474
1475                let header = read::block_header(best_chain, &state.db, height.into())
1476                    .ok_or_else(|| BoxError::from("block hash or height not found"))?;
1477
1478                Ok(ReadResponse::BlockHeader {
1479                    header,
1480                    hash,
1481                    height,
1482                    next_block_hash,
1483                })
1484            }
1485
1486            // For the get_raw_transaction RPC and the StateService.
1487            ReadRequest::Transaction(hash) => Ok(ReadResponse::Transaction(
1488                read::mined_transaction(state.latest_best_chain(), &state.db, hash),
1489            )),
1490
1491            ReadRequest::AnyChainTransaction(hash) => {
1492                Ok(ReadResponse::AnyChainTransaction(read::any_transaction(
1493                    state.latest_non_finalized_state().chain_iter(),
1494                    &state.db,
1495                    hash,
1496                )))
1497            }
1498
1499            // Used by the getblock (verbose) RPC.
1500            ReadRequest::TransactionIdsForBlock(hash_or_height) => Ok(
1501                ReadResponse::TransactionIdsForBlock(read::transaction_hashes_for_block(
1502                    state.latest_best_chain(),
1503                    &state.db,
1504                    hash_or_height,
1505                )),
1506            ),
1507
1508            ReadRequest::AnyChainTransactionIdsForBlock(hash_or_height) => {
1509                Ok(ReadResponse::AnyChainTransactionIdsForBlock(
1510                    read::transaction_hashes_for_any_block(
1511                        state.latest_non_finalized_state().chain_iter(),
1512                        &state.db,
1513                        hash_or_height,
1514                    ),
1515                ))
1516            }
1517
1518            #[cfg(feature = "indexer")]
1519            ReadRequest::SpendingTransactionId(spend) => Ok(ReadResponse::TransactionId(
1520                read::spending_transaction_hash(state.latest_best_chain(), &state.db, spend),
1521            )),
1522
1523            ReadRequest::UnspentBestChainUtxo(outpoint) => Ok(ReadResponse::UnspentBestChainUtxo(
1524                read::unspent_utxo(state.latest_best_chain(), &state.db, outpoint),
1525            )),
1526
1527            // Manually used by the StateService to implement part of AwaitUtxo.
1528            ReadRequest::AnyChainUtxo(outpoint) => Ok(ReadResponse::AnyChainUtxo(read::any_utxo(
1529                state.latest_non_finalized_state(),
1530                &state.db,
1531                outpoint,
1532            ))),
1533
1534            // Used by the StateService.
1535            ReadRequest::BlockLocator => Ok(ReadResponse::BlockLocator(
1536                read::block_locator(state.latest_best_chain(), &state.db).unwrap_or_default(),
1537            )),
1538
1539            // Used by the StateService.
1540            ReadRequest::FindBlockHashes { known_blocks, stop } => {
1541                Ok(ReadResponse::BlockHashes(read::find_chain_hashes(
1542                    state.latest_best_chain(),
1543                    &state.db,
1544                    known_blocks,
1545                    stop,
1546                    MAX_FIND_BLOCK_HASHES_RESULTS,
1547                )))
1548            }
1549
1550            // Used by the StateService.
1551            ReadRequest::FindBlockHeaders { known_blocks, stop } => Ok(ReadResponse::BlockHeaders(
1552                read::find_chain_headers(
1553                    state.latest_best_chain(),
1554                    &state.db,
1555                    known_blocks,
1556                    stop,
1557                    MAX_FIND_BLOCK_HEADERS_RESULTS,
1558                )
1559                .into_iter()
1560                .map(|header| CountedHeader { header })
1561                .collect(),
1562            )),
1563
1564            ReadRequest::FindForkPoint { known_blocks } => {
1565                // Reject over-long locators before doing any work, so an untrusted
1566                // caller can't force unbounded lookups.
1567                let locator_len: u64 = known_blocks
1568                    .len()
1569                    .try_into()
1570                    .expect("usize always fits in u64 on supported (<=64-bit) platforms");
1571                if locator_len > block::MAX_BLOCK_LOCATOR_LENGTH {
1572                    return Err(BoxError::from(format!(
1573                        "FindForkPoint locator length {locator_len} exceeds \
1574                         MAX_BLOCK_LOCATOR_LENGTH ({})",
1575                        block::MAX_BLOCK_LOCATOR_LENGTH,
1576                    )));
1577                }
1578
1579                Ok(ReadResponse::ForkPoint(read::find_fork_point(
1580                    state.latest_best_chain(),
1581                    &state.db,
1582                    known_blocks,
1583                )))
1584            }
1585
1586            ReadRequest::SaplingTree(hash_or_height) => Ok(ReadResponse::SaplingTree(
1587                read::sapling_tree(state.latest_best_chain(), &state.db, hash_or_height),
1588            )),
1589
1590            ReadRequest::OrchardTree(hash_or_height) => Ok(ReadResponse::OrchardTree(
1591                read::orchard_tree(state.latest_best_chain(), &state.db, hash_or_height),
1592            )),
1593
1594            ReadRequest::IronwoodTree(hash_or_height) => Ok(ReadResponse::IronwoodTree(
1595                read::ironwood_tree(state.latest_best_chain(), &state.db, hash_or_height),
1596            )),
1597
1598            ReadRequest::SaplingSubtrees { start_index, limit } => {
1599                let end_index = limit
1600                    .and_then(|limit| start_index.0.checked_add(limit.0))
1601                    .map(NoteCommitmentSubtreeIndex);
1602
1603                let best_chain = state.latest_best_chain();
1604                let sapling_subtrees = if let Some(end_index) = end_index {
1605                    read::sapling_subtrees(best_chain, &state.db, start_index..end_index)
1606                } else {
1607                    // If there is no end bound, just return all the trees.
1608                    // If the end bound would overflow, just returns all the trees, because that's what
1609                    // `zcashd` does. (It never calculates an end bound, so it just keeps iterating until
1610                    // the trees run out.)
1611                    read::sapling_subtrees(best_chain, &state.db, start_index..)
1612                };
1613
1614                Ok(ReadResponse::SaplingSubtrees(sapling_subtrees))
1615            }
1616
1617            ReadRequest::OrchardSubtrees { start_index, limit } => {
1618                let end_index = limit
1619                    .and_then(|limit| start_index.0.checked_add(limit.0))
1620                    .map(NoteCommitmentSubtreeIndex);
1621
1622                let best_chain = state.latest_best_chain();
1623                let orchard_subtrees = if let Some(end_index) = end_index {
1624                    read::orchard_subtrees(best_chain, &state.db, start_index..end_index)
1625                } else {
1626                    // If there is no end bound, just return all the trees.
1627                    // If the end bound would overflow, just returns all the trees, because that's what
1628                    // `zcashd` does. (It never calculates an end bound, so it just keeps iterating until
1629                    // the trees run out.)
1630                    read::orchard_subtrees(best_chain, &state.db, start_index..)
1631                };
1632
1633                Ok(ReadResponse::OrchardSubtrees(orchard_subtrees))
1634            }
1635
1636            ReadRequest::IronwoodSubtrees { start_index, limit } => {
1637                let end_index = limit
1638                    .and_then(|limit| start_index.0.checked_add(limit.0))
1639                    .map(NoteCommitmentSubtreeIndex);
1640
1641                let best_chain = state.latest_best_chain();
1642                let ironwood_subtrees = if let Some(end_index) = end_index {
1643                    read::ironwood_subtrees(best_chain, &state.db, start_index..end_index)
1644                } else {
1645                    // If there is no end bound, just return all the trees.
1646                    // If the end bound would overflow, just returns all the trees, because that's what
1647                    // `zcashd` does. (It never calculates an end bound, so it just keeps iterating until
1648                    // the trees run out.)
1649                    read::ironwood_subtrees(best_chain, &state.db, start_index..)
1650                };
1651
1652                Ok(ReadResponse::IronwoodSubtrees(ironwood_subtrees))
1653            }
1654
1655            // For the get_address_balance RPC.
1656            ReadRequest::AddressBalance(addresses) => {
1657                let (balance, received) =
1658                    read::transparent_balance(state.latest_best_chain(), &state.db, addresses)?;
1659                Ok(ReadResponse::AddressBalance { balance, received })
1660            }
1661
1662            // For the get_address_tx_ids RPC.
1663            ReadRequest::TransactionIdsByAddresses {
1664                addresses,
1665                height_range,
1666            } => read::transparent_tx_ids(
1667                state.latest_best_chain(),
1668                &state.db,
1669                addresses,
1670                height_range,
1671            )
1672            .map(ReadResponse::AddressesTransactionIds),
1673
1674            // For the get_address_utxos RPC.
1675            ReadRequest::UtxosByAddresses(addresses) => read::address_utxos(
1676                &state.network,
1677                state.latest_best_chain(),
1678                &state.db,
1679                addresses,
1680            )
1681            .map(ReadResponse::AddressUtxos),
1682
1683            ReadRequest::CheckBestChainTipNullifiersAndAnchors(unmined_tx) => {
1684                let latest_non_finalized_best_chain = state.latest_best_chain();
1685
1686                check::nullifier::tx_no_duplicates_in_chain(
1687                    &state.db,
1688                    latest_non_finalized_best_chain.as_ref(),
1689                    &unmined_tx.transaction,
1690                )?;
1691
1692                check::anchors::tx_anchors_refer_to_final_treestates(
1693                    &state.db,
1694                    latest_non_finalized_best_chain.as_ref(),
1695                    &unmined_tx,
1696                )?;
1697
1698                Ok(ReadResponse::ValidBestChainTipNullifiersAndAnchors)
1699            }
1700
1701            // Used by the get_block and get_block_hash RPCs.
1702            ReadRequest::BestChainBlockHash(height) => Ok(ReadResponse::BlockHash(
1703                read::hash_by_height(state.latest_best_chain(), &state.db, height),
1704            )),
1705
1706            // Used by get_block_template and getblockchaininfo RPCs.
1707            ReadRequest::ChainInfo => {
1708                // # Correctness
1709                //
1710                // It is ok to do these lookups using multiple database calls. Finalized state updates
1711                // can only add overlapping blocks, and block hashes are unique across all chain forks.
1712                //
1713                // If there is a large overlap between the non-finalized and finalized states,
1714                // where the finalized tip is above the non-finalized tip,
1715                // Zebra is receiving a lot of blocks, or this request has been delayed for a long time.
1716                //
1717                // In that case, the `getblocktemplate` RPC will return an error because Zebra
1718                // is not synced to the tip. That check happens before the RPC makes this request.
1719                read::difficulty::get_block_template_chain_info(
1720                    &state.latest_non_finalized_state(),
1721                    &state.db,
1722                    &state.network,
1723                )
1724                .map(ReadResponse::ChainInfo)
1725            }
1726
1727            // Used by getmininginfo, getnetworksolps, and getnetworkhashps RPCs.
1728            ReadRequest::SolutionRate { num_blocks, height } => {
1729                let latest_non_finalized_state = state.latest_non_finalized_state();
1730                // # Correctness
1731                //
1732                // It is ok to do these lookups using multiple database calls. Finalized state updates
1733                // can only add overlapping blocks, and block hashes are unique across all chain forks.
1734                //
1735                // The worst that can happen here is that the default `start_hash` will be below
1736                // the chain tip.
1737                let (tip_height, tip_hash) =
1738                    match read::tip(latest_non_finalized_state.best_chain(), &state.db) {
1739                        Some(tip_hash) => tip_hash,
1740                        None => return Ok(ReadResponse::SolutionRate(None)),
1741                    };
1742
1743                let start_hash = match height {
1744                    Some(height) if height < tip_height => read::hash_by_height(
1745                        latest_non_finalized_state.best_chain(),
1746                        &state.db,
1747                        height,
1748                    ),
1749                    // use the chain tip hash if height is above it or not provided.
1750                    _ => Some(tip_hash),
1751                };
1752
1753                let solution_rate = start_hash.and_then(|start_hash| {
1754                    read::difficulty::solution_rate(
1755                        &latest_non_finalized_state,
1756                        &state.db,
1757                        num_blocks,
1758                        start_hash,
1759                    )
1760                });
1761
1762                Ok(ReadResponse::SolutionRate(solution_rate))
1763            }
1764
1765            ReadRequest::CheckBlockProposalValidity(semantically_verified) => {
1766                tracing::debug!(
1767                    "attempting to validate and commit block proposal \
1768                         onto a cloned non-finalized state"
1769                );
1770                let mut latest_non_finalized_state = state.latest_non_finalized_state();
1771
1772                // The previous block of a valid proposal must be on the best chain tip.
1773                let Some((_best_tip_height, best_tip_hash)) =
1774                    read::best_tip(&latest_non_finalized_state, &state.db)
1775                else {
1776                    return Err(
1777                        "state is empty: wait for Zebra to sync before submitting a proposal"
1778                            .into(),
1779                    );
1780                };
1781
1782                if semantically_verified.block.header.previous_block_hash != best_tip_hash {
1783                    return Err("proposal is not based on the current best chain tip: \
1784                                    previous block hash must be the best chain tip"
1785                        .into());
1786                }
1787
1788                // This clone of the non-finalized state is dropped when this closure returns.
1789                // The non-finalized state that's used in the rest of the state (including finalizing
1790                // blocks into the db) is not mutated here.
1791                //
1792                // TODO: Convert `CommitSemanticallyVerifiedError` to a new `ValidateProposalError`?
1793                latest_non_finalized_state.disable_metrics();
1794
1795                write::validate_and_commit_non_finalized(
1796                    &state.db,
1797                    &mut latest_non_finalized_state,
1798                    semantically_verified,
1799                )?;
1800
1801                Ok(ReadResponse::ValidBlockProposal)
1802            }
1803
1804            ReadRequest::TipBlockSize => {
1805                // Respond with the length of the obtained block if any.
1806                Ok(ReadResponse::TipBlockSize(
1807                    state
1808                        .best_tip()
1809                        .and_then(|(tip_height, _)| {
1810                            read::block_info(
1811                                state.latest_best_chain(),
1812                                &state.db,
1813                                tip_height.into(),
1814                            )
1815                        })
1816                        .map(|info| info.size().try_into().expect("u32 should fit in usize"))
1817                        .or_else(|| {
1818                            find::tip_block(state.latest_best_chain(), &state.db)
1819                                .map(|b| b.zcash_serialized_size())
1820                        }),
1821                ))
1822            }
1823
1824            ReadRequest::NonFinalizedBlocksListener { .. } => {
1825                unreachable!("should return early");
1826            }
1827
1828            // Used by `gettxout` RPC method.
1829            ReadRequest::IsTransparentOutputSpent(outpoint) => {
1830                let is_spent = read::unspent_utxo(state.latest_best_chain(), &state.db, outpoint);
1831                Ok(ReadResponse::IsTransparentOutputSpent(is_spent.is_none()))
1832            }
1833        };
1834
1835        timed_span.spawn_blocking(request_handler)
1836    }
1837}
1838
1839/// Initialize a state service from the provided [`Config`].
1840/// Returns a boxed state service, a read-only state service,
1841/// and receivers for state chain tip updates.
1842///
1843/// Each `network` has its own separate on-disk database.
1844///
1845/// The state uses the `max_checkpoint_height` and `checkpoint_verify_concurrency_limit`
1846/// to work out when it is near the final checkpoint.
1847///
1848/// To share access to the state, wrap the returned service in a `Buffer`,
1849/// or clone the returned [`ReadStateService`].
1850///
1851/// It's possible to construct multiple state services in the same application (as
1852/// long as they, e.g., use different storage locations), but doing so is
1853/// probably not what you want.
1854pub async fn init(
1855    config: Config,
1856    network: &Network,
1857    max_checkpoint_height: block::Height,
1858    checkpoint_verify_concurrency_limit: usize,
1859) -> (
1860    BoxService<Request, Response, BoxError>,
1861    ReadStateService,
1862    LatestChainTip,
1863    ChainTipChange,
1864) {
1865    let (state_service, read_only_state_service, latest_chain_tip, chain_tip_change) =
1866        StateService::new(
1867            config,
1868            network,
1869            max_checkpoint_height,
1870            checkpoint_verify_concurrency_limit,
1871        )
1872        .await;
1873
1874    (
1875        BoxService::new(state_service),
1876        read_only_state_service,
1877        latest_chain_tip,
1878        chain_tip_change,
1879    )
1880}
1881
1882/// Initialize a read state service from the provided [`Config`].
1883/// Returns a read-only state service,
1884///
1885/// Each `network` has its own separate on-disk database.
1886///
1887/// To share access to the state, clone the returned [`ReadStateService`].
1888pub fn init_read_only(
1889    config: Config,
1890    network: &Network,
1891) -> Result<
1892    (
1893        ReadStateService,
1894        ZebraDb,
1895        tokio::sync::watch::Sender<NonFinalizedState>,
1896    ),
1897    StateInitError,
1898> {
1899    let finalized_state = FinalizedState::new_with_debug(
1900        &config,
1901        network,
1902        true,
1903        #[cfg(feature = "elasticsearch")]
1904        false,
1905        true,
1906    )?;
1907    let (non_finalized_state_sender, non_finalized_state_receiver) =
1908        tokio::sync::watch::channel(NonFinalizedState::new(network));
1909
1910    Ok((
1911        ReadStateService::new(
1912            &finalized_state,
1913            None,
1914            WatchReceiver::new(non_finalized_state_receiver),
1915        ),
1916        finalized_state.db.clone(),
1917        non_finalized_state_sender,
1918    ))
1919}
1920
1921/// Calls [`init_read_only`] with the provided [`Config`] and [`Network`] from a blocking task.
1922///
1923/// Returns a [`tokio::task::JoinHandle`] whose output is a [`Result`]: awaiting it yields a
1924/// [`JoinError`](tokio::task::JoinError) if the blocking task panicked or was cancelled, and
1925/// otherwise an `Err(`[`StateInitError`]`)` if the read-only state could not be opened (for
1926/// example, a missing read-only database).
1927pub fn spawn_init_read_only(
1928    config: Config,
1929    network: &Network,
1930) -> tokio::task::JoinHandle<
1931    Result<
1932        (
1933            ReadStateService,
1934            ZebraDb,
1935            tokio::sync::watch::Sender<NonFinalizedState>,
1936        ),
1937        StateInitError,
1938    >,
1939> {
1940    let network = network.clone();
1941    tokio::task::spawn_blocking(move || init_read_only(config, &network))
1942}
1943
1944/// Returns a [`StateService`] with an ephemeral [`Config`] and a buffer with a single slot.
1945///
1946/// This can be used to create a state service for testing. See also [`init`].
1947#[cfg(any(test, feature = "proptest-impl"))]
1948pub async fn init_test(
1949    network: &Network,
1950) -> Buffer<BoxService<Request, Response, BoxError>, Request> {
1951    // TODO: pass max_checkpoint_height and checkpoint_verify_concurrency limit
1952    //       if we ever need to test final checkpoint sent UTXO queries
1953    let (state_service, _, _, _) =
1954        StateService::new(Config::ephemeral(), network, block::Height::MAX, 0).await;
1955
1956    Buffer::new(BoxService::new(state_service), 1)
1957}
1958
1959/// Initializes a state service with an ephemeral [`Config`] and a buffer with a single slot,
1960/// then returns the read-write service, read-only service, and tip watch channels.
1961///
1962/// This can be used to create a state service for testing. See also [`init`].
1963#[cfg(any(test, feature = "proptest-impl"))]
1964pub async fn init_test_services(
1965    network: &Network,
1966) -> (
1967    Buffer<BoxService<Request, Response, BoxError>, Request>,
1968    ReadStateService,
1969    LatestChainTip,
1970    ChainTipChange,
1971) {
1972    // TODO: pass max_checkpoint_height and checkpoint_verify_concurrency limit
1973    //       if we ever need to test final checkpoint sent UTXO queries
1974    let (state_service, read_state_service, latest_chain_tip, chain_tip_change) =
1975        StateService::new(Config::ephemeral(), network, block::Height::MAX, 0).await;
1976
1977    let state_service = Buffer::new(BoxService::new(state_service), 1);
1978
1979    (
1980        state_service,
1981        read_state_service,
1982        latest_chain_tip,
1983        chain_tip_change,
1984    )
1985}