2024-01-24 14:37:03 -08:00
|
|
|
//! [`tower::Service`] for zebra-scan.
|
|
|
|
|
2024-02-01 18:29:36 -08:00
|
|
|
use std::{collections::BTreeMap, future::Future, pin::Pin, task::Poll, time::Duration};
|
2024-01-24 14:37:03 -08:00
|
|
|
|
|
|
|
use futures::future::FutureExt;
|
2024-03-14 13:27:33 -07:00
|
|
|
use tower::{BoxError, Service};
|
2024-01-24 14:37:03 -08:00
|
|
|
|
2024-02-07 14:36:01 -08:00
|
|
|
use zebra_chain::{diagnostic::task::WaitForPanics, parameters::Network, transaction::Hash};
|
2024-02-01 18:29:36 -08:00
|
|
|
|
2024-01-24 14:37:03 -08:00
|
|
|
use zebra_state::ChainTipChange;
|
|
|
|
|
2024-02-06 10:41:00 -08:00
|
|
|
use crate::{scan, storage::Storage, Config, Request, Response};
|
2024-01-24 14:37:03 -08:00
|
|
|
|
2024-01-31 11:34:24 -08:00
|
|
|
#[cfg(test)]
|
|
|
|
mod tests;
|
|
|
|
|
2024-02-06 10:41:00 -08:00
|
|
|
pub mod scan_task;
|
|
|
|
|
|
|
|
pub use scan_task::{ScanTask, ScanTaskCommand};
|
|
|
|
|
|
|
|
#[cfg(any(test, feature = "proptest-impl"))]
|
2024-02-12 16:42:40 -08:00
|
|
|
use tokio::sync::mpsc::Receiver;
|
2024-02-06 10:41:00 -08:00
|
|
|
|
2024-01-24 14:37:03 -08:00
|
|
|
/// Zebra-scan [`tower::Service`]
|
|
|
|
#[derive(Debug)]
|
|
|
|
pub struct ScanService {
|
|
|
|
/// On-disk storage
|
2024-01-31 11:34:24 -08:00
|
|
|
pub db: Storage,
|
2024-01-24 14:37:03 -08:00
|
|
|
|
|
|
|
/// Handle to scan task that's responsible for writing results
|
|
|
|
scan_task: ScanTask,
|
|
|
|
}
|
|
|
|
|
2024-01-31 11:34:24 -08:00
|
|
|
/// A timeout applied to `DeleteKeys` requests.
|
2024-03-14 13:27:33 -07:00
|
|
|
///
|
|
|
|
/// This should be shorter than [`SCAN_SERVICE_TIMEOUT`](crate::init::SCAN_SERVICE_TIMEOUT) so the
|
|
|
|
/// request can try to delete entries from storage after the timeout before the future is dropped.
|
2024-01-31 11:34:24 -08:00
|
|
|
const DELETE_KEY_TIMEOUT: Duration = Duration::from_secs(15);
|
|
|
|
|
2024-01-24 14:37:03 -08:00
|
|
|
impl ScanService {
|
|
|
|
/// Create a new [`ScanService`].
|
2024-02-07 14:36:01 -08:00
|
|
|
pub async fn new(
|
2024-01-24 14:37:03 -08:00
|
|
|
config: &Config,
|
|
|
|
network: Network,
|
|
|
|
state: scan::State,
|
|
|
|
chain_tip_change: ChainTipChange,
|
|
|
|
) -> Self {
|
2024-02-07 14:36:01 -08:00
|
|
|
let config = config.clone();
|
|
|
|
let storage = tokio::task::spawn_blocking(move || Storage::new(&config, network, false))
|
|
|
|
.wait_for_panics()
|
|
|
|
.await;
|
|
|
|
|
2024-01-24 14:37:03 -08:00
|
|
|
Self {
|
2024-02-07 14:36:01 -08:00
|
|
|
scan_task: ScanTask::spawn(storage.clone(), state, chain_tip_change),
|
|
|
|
db: storage,
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
}
|
2024-01-25 17:29:37 -08:00
|
|
|
|
|
|
|
/// Create a new [`ScanService`] with a mock `ScanTask`
|
2024-02-06 10:41:00 -08:00
|
|
|
// TODO: Move this to tests behind `cfg(any(test, feature = "proptest-impl"))`
|
2024-01-31 11:34:24 -08:00
|
|
|
#[cfg(any(test, feature = "proptest-impl"))]
|
2024-02-06 10:41:00 -08:00
|
|
|
pub fn new_with_mock_scanner(db: Storage) -> (Self, Receiver<ScanTaskCommand>) {
|
2024-01-31 11:34:24 -08:00
|
|
|
let (scan_task, cmd_receiver) = ScanTask::mock();
|
|
|
|
(Self { db, scan_task }, cmd_receiver)
|
2024-01-25 17:29:37 -08:00
|
|
|
}
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
|
|
|
|
impl Service<Request> for ScanService {
|
|
|
|
type Response = Response;
|
2024-03-14 13:27:33 -07:00
|
|
|
type Error = BoxError;
|
2024-01-24 14:37:03 -08:00
|
|
|
type Future =
|
|
|
|
Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send + 'static>>;
|
|
|
|
|
|
|
|
fn poll_ready(&mut self, _cx: &mut std::task::Context<'_>) -> Poll<Result<(), Self::Error>> {
|
|
|
|
// TODO: If scan task returns an error, add error to the panic message
|
|
|
|
assert!(
|
|
|
|
!self.scan_task.handle.is_finished(),
|
|
|
|
"scan task finished unexpectedly"
|
|
|
|
);
|
|
|
|
|
|
|
|
self.db.check_for_panics();
|
|
|
|
|
|
|
|
Poll::Ready(Ok(()))
|
|
|
|
}
|
|
|
|
|
|
|
|
fn call(&mut self, req: Request) -> Self::Future {
|
2024-02-01 12:07:31 -08:00
|
|
|
if let Err(error) = req.check() {
|
|
|
|
return async move { Err(error) }.boxed();
|
|
|
|
}
|
|
|
|
|
2024-01-24 14:37:03 -08:00
|
|
|
match req {
|
2024-01-25 17:29:37 -08:00
|
|
|
Request::Info => {
|
|
|
|
let db = self.db.clone();
|
|
|
|
|
|
|
|
return async move {
|
|
|
|
Ok(Response::Info {
|
2024-02-09 07:23:19 -08:00
|
|
|
min_sapling_birthday_height: db.network().sapling_activation_height(),
|
2024-01-25 17:29:37 -08:00
|
|
|
})
|
|
|
|
}
|
|
|
|
.boxed();
|
|
|
|
}
|
|
|
|
|
2024-02-09 07:23:19 -08:00
|
|
|
Request::RegisterKeys(keys) => {
|
|
|
|
let mut scan_task = self.scan_task.clone();
|
|
|
|
|
|
|
|
return async move {
|
2024-03-14 13:27:33 -07:00
|
|
|
let newly_registered_keys = scan_task.register_keys(keys)?.await?;
|
|
|
|
if !newly_registered_keys.is_empty() {
|
|
|
|
Ok(Response::RegisteredKeys(newly_registered_keys))
|
|
|
|
} else {
|
|
|
|
Err("no keys were registered, check that keys are not already registered and \
|
|
|
|
are valid Sapling extended full viewing keys".into())
|
|
|
|
}
|
2024-02-09 07:23:19 -08:00
|
|
|
}
|
|
|
|
.boxed();
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
|
2024-01-31 11:34:24 -08:00
|
|
|
Request::DeleteKeys(keys) => {
|
|
|
|
let mut db = self.db.clone();
|
|
|
|
let mut scan_task = self.scan_task.clone();
|
|
|
|
|
|
|
|
return async move {
|
|
|
|
// Wait for a message to confirm that the scan task has removed the key up to `DELETE_KEY_TIMEOUT`
|
2024-03-04 13:25:50 -08:00
|
|
|
let remove_keys_result = tokio::time::timeout(
|
|
|
|
DELETE_KEY_TIMEOUT,
|
|
|
|
scan_task.remove_keys(keys.clone())?,
|
|
|
|
)
|
|
|
|
.await
|
2024-03-14 13:27:33 -07:00
|
|
|
.map_err(|_| "request timed out removing keys from scan task".to_string());
|
2024-01-31 11:34:24 -08:00
|
|
|
|
|
|
|
// Delete the key from the database after either confirmation that it's been removed from the scan task, or
|
|
|
|
// waiting `DELETE_KEY_TIMEOUT`.
|
|
|
|
let delete_key_task = tokio::task::spawn_blocking(move || {
|
2024-02-01 12:07:31 -08:00
|
|
|
db.delete_sapling_keys(keys);
|
2024-01-31 11:34:24 -08:00
|
|
|
});
|
|
|
|
|
|
|
|
// Return timeout errors or `RecvError`s, or wait for the key to be deleted from the database.
|
|
|
|
remove_keys_result??;
|
|
|
|
delete_key_task.await?;
|
|
|
|
|
|
|
|
Ok(Response::DeletedKeys)
|
|
|
|
}
|
|
|
|
.boxed();
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
|
2024-02-01 18:29:36 -08:00
|
|
|
Request::Results(keys) => {
|
|
|
|
let db = self.db.clone();
|
|
|
|
|
|
|
|
return async move {
|
|
|
|
let mut final_result = BTreeMap::new();
|
|
|
|
for key in keys {
|
|
|
|
let db = db.clone();
|
|
|
|
let mut heights_and_transactions = BTreeMap::new();
|
|
|
|
let txs = {
|
|
|
|
let key = key.clone();
|
|
|
|
tokio::task::spawn_blocking(move || db.sapling_results_for_key(&key))
|
|
|
|
}
|
|
|
|
.await?;
|
|
|
|
txs.iter().for_each(|(k, v)| {
|
|
|
|
heights_and_transactions
|
|
|
|
.entry(*k)
|
|
|
|
.or_insert_with(Vec::new)
|
|
|
|
.extend(v.iter().map(|x| Hash::from(*x)));
|
|
|
|
});
|
|
|
|
final_result.entry(key).or_insert(heights_and_transactions);
|
|
|
|
}
|
|
|
|
|
|
|
|
Ok(Response::Results(final_result))
|
|
|
|
}
|
|
|
|
.boxed();
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
|
2024-02-12 16:42:40 -08:00
|
|
|
Request::SubscribeResults(keys) => {
|
|
|
|
let mut scan_task = self.scan_task.clone();
|
|
|
|
|
|
|
|
return async move {
|
2024-03-14 13:27:33 -07:00
|
|
|
let results_receiver = scan_task.subscribe(keys)?.await.map_err(|_| {
|
|
|
|
"scan task dropped responder, check that keys are registered"
|
|
|
|
})?;
|
2024-02-12 16:42:40 -08:00
|
|
|
|
|
|
|
Ok(Response::SubscribeResults(results_receiver))
|
|
|
|
}
|
|
|
|
.boxed();
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
|
2024-02-01 12:07:31 -08:00
|
|
|
Request::ClearResults(keys) => {
|
|
|
|
let mut db = self.db.clone();
|
|
|
|
|
|
|
|
return async move {
|
|
|
|
// Clear results from db for the provided `keys`
|
|
|
|
tokio::task::spawn_blocking(move || {
|
|
|
|
db.delete_sapling_results(keys);
|
|
|
|
})
|
|
|
|
.await?;
|
|
|
|
|
|
|
|
Ok(Response::ClearedResults)
|
|
|
|
}
|
|
|
|
.boxed();
|
2024-01-24 14:37:03 -08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|