b231eb4f04
* Remove AsyncRead/Write traits from native client - moving them to stream/ * Substream model first push * Update / add examples * Update lockfile * Clippy * clippy examples * remove codecs * Remove unused bincode setup * Revert a lot of changes when SDK client itself implemented AsyncRead/Write * Remove unnecessary mut * Use local PollSender in MixnetStream instead of client_input.input_sender Now that client-core's input_sender is back to mpsc::Sender (reverted PollSender migration), MixnetStream creates its own PollSender wrapper for the AsyncWrite impl's poll_ready/start_send calls. * Remove now-unnecessary parameter * Clippy * Cleanup more stragglers from previous setup (Async traits on MixnetClient) * Rename files (remove module inception) * - Shrink StreamId from 16 bytes to u64, add version byte to wire format - Introduce MixStreamHeader/MixStreamFrame structs for decode - Replace StreamMap type alias with struct using tokio::sync::Mutex - Add StreamMap helper methods, eliminate lock().expect() panics - Rename stream.rs -> mixnet_stream.rs to avoid module inception - Document irrevocable stream mode, add LP integration TODO * - Remove dummy channel - Add err variant for reciever alredy taken - Remove panics * add timeout to stream * clippy
68 lines
2.2 KiB
Rust
68 lines
2.2 KiB
Rust
// Copyright 2022 - Nym Technologies SA <contact@nymtech.net>
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
#![allow(deprecated)] // silences clippy warning: use of deprecated associated function `nym_crypto::generic_array::GenericArray::<T, N>::from_exact_iter`: please upgrade to generic-array 1.x - TODO
|
|
pub use backend::*;
|
|
pub use combined::CombinedReplyStorage;
|
|
pub use key_storage::SentReplyKeys;
|
|
pub use surb_storage::{ReceivedReplySurb, ReceivedReplySurbsMap, RetrievedReplySurb};
|
|
pub use tag_storage::UsedSenderTags;
|
|
use time::OffsetDateTime;
|
|
|
|
mod backend;
|
|
mod combined;
|
|
mod key_storage;
|
|
mod surb_storage;
|
|
mod tag_storage;
|
|
|
|
// only really exists to get information about shutdown and save data to the backing storage
|
|
pub struct PersistentReplyStorage<T = backend::Empty>
|
|
where
|
|
T: ReplyStorageBackend,
|
|
{
|
|
backend: T,
|
|
}
|
|
|
|
impl<T> PersistentReplyStorage<T>
|
|
where
|
|
T: ReplyStorageBackend + Send + Sync,
|
|
{
|
|
pub fn new(backend: T) -> Self {
|
|
PersistentReplyStorage { backend }
|
|
}
|
|
|
|
pub async fn load_state_from_backend(
|
|
&self,
|
|
surb_freshness_cutoff: OffsetDateTime,
|
|
) -> Result<CombinedReplyStorage, T::StorageError> {
|
|
self.backend.load_surb_storage(surb_freshness_cutoff).await
|
|
}
|
|
|
|
pub async fn flush_on_shutdown(
|
|
mut self,
|
|
mem_state: CombinedReplyStorage,
|
|
shutdown: nym_task::ShutdownToken,
|
|
) {
|
|
use tracing::{debug, error, info};
|
|
|
|
debug!("Started PersistentReplyStorage");
|
|
if let Err(err) = self.backend.start_storage_session().await {
|
|
error!("failed to start the storage session - {err}");
|
|
return;
|
|
}
|
|
|
|
shutdown.cancelled().await;
|
|
|
|
info!("PersistentReplyStorage is flushing all reply-related data to underlying storage");
|
|
if let Err(err) = self.backend.flush_surb_storage(&mem_state).await {
|
|
error!("failed to flush our reply-related data to the persistent storage: {err}")
|
|
} else {
|
|
info!("Data flush is complete")
|
|
}
|
|
|
|
if let Err(err) = self.backend.stop_storage_session().await {
|
|
error!("failed to properly stop the storage session - {err}. We might not be able to smoothly restore it")
|
|
}
|
|
}
|
|
}
|