use std::{
collections::{HashMap, HashSet},
str,
sync::Arc,
time::{Duration, SystemTime},
};
use flume::{Receiver, Sender};
use futures::{pin_mut, select, FutureExt};
use tokio::{sync::RwLock, time::interval};
use zenoh::key_expr::keyexpr;
use zenoh_backend_traits::config::{ReplicaConfig, StorageConfig};
use crate::{backends_mgt::StoreIntercept, storages_mgt::StorageMessage};
pub mod align_queryable;
pub mod aligner;
pub mod digest;
pub mod snapshotter;
pub mod storage;
pub use align_queryable::AlignQueryable;
pub use aligner::Aligner;
pub use digest::{Digest, DigestConfig, EraType, LogEntry};
pub use snapshotter::Snapshotter;
pub use storage::{ReplicationService, StorageService};
use zenoh::{key_expr::OwnedKeyExpr, sample::Locality, time::Timestamp, Session};
const ERA: &str = "era";
const INTERVALS: &str = "intervals";
const SUBINTERVALS: &str = "subintervals";
const CONTENTS: &str = "contents";
pub const EPOCH_START: SystemTime = SystemTime::UNIX_EPOCH;
pub const SUBINTERVAL_CHUNKS: usize = 10;
lazy_static::lazy_static!(
static ref KE_PREFIX_DIGEST: &'static keyexpr = unsafe { keyexpr::from_str_unchecked("@-digest") };
);
pub struct Replica {
name: String, session: Arc<Session>,
key_expr: OwnedKeyExpr,
replica_config: ReplicaConfig,
digests_published: RwLock<HashSet<u64>>, }
impl Replica {
pub async fn start(
session: Arc<Session>,
store_intercept: StoreIntercept,
storage_config: StorageConfig,
name: &str,
rx: Receiver<StorageMessage>,
) {
tracing::trace!("[REPLICA] Opening session...");
let startup_entries = match store_intercept.storage.get_all_entries().await {
Ok(entries) => {
let mut result = Vec::new();
for entry in entries {
if entry.0.is_none() {
if let Some(prefix) = storage_config.clone().strip_prefix {
result.push((prefix, entry.1));
} else {
tracing::error!("Empty key found with timestamp `{}`", entry.1);
}
} else {
result.push((
crate::prefix(storage_config.strip_prefix.as_ref(), &entry.0.unwrap()),
entry.1,
));
}
}
result
}
Err(e) => {
tracing::error!("[REPLICA] Error fetching entries from storage: {}", e);
return;
}
};
let replica = Replica {
name: name.to_string(),
session: session.clone(),
key_expr: storage_config.key_expr.clone(),
replica_config: storage_config.replica_config.clone().unwrap(),
digests_published: RwLock::new(HashSet::new()),
};
let (tx_digest, rx_digest) = flume::unbounded();
let (tx_sample, rx_sample) = flume::unbounded();
let (tx_log, rx_log) = flume::unbounded();
let config = replica.replica_config.clone();
let snapshotter =
Arc::new(Snapshotter::new(session, rx_log, &startup_entries, &config).await);
let digest_sub = replica.start_digest_sub(tx_digest).fuse();
let digest_key = Replica::get_digest_key(&replica.key_expr);
let align_q = AlignQueryable::start_align_queryable(
replica.session.clone(),
digest_key.clone(),
&replica.name,
snapshotter.clone(),
)
.fuse();
let aligner = Aligner::start_aligner(
replica.session.clone(),
digest_key,
rx_digest,
tx_sample,
snapshotter.clone(),
)
.fuse();
let digest_pub = replica.start_digest_pub(snapshotter.clone()).fuse();
let snapshot_task = snapshotter.start().fuse();
let replication = ReplicationService {
empty_start: startup_entries.is_empty(),
aligner_updates: rx_sample,
log_propagation: tx_log,
};
let storage_task = StorageService::start(
replica.session.clone(),
storage_config,
&replica.name,
store_intercept,
rx,
Some(replication),
)
.fuse();
pin_mut!(
digest_sub,
align_q,
aligner,
digest_pub,
snapshot_task,
storage_task
);
select!(
() = digest_sub => tracing::trace!("[REPLICA] Exiting digest subscriber"),
() = align_q => tracing::trace!("[REPLICA] Exiting align queryable"),
() = aligner => tracing::trace!("[REPLICA] Exiting aligner"),
() = digest_pub => tracing::trace!("[REPLICA] Exiting digest publisher"),
() = snapshot_task => tracing::trace!("[REPLICA] Exiting snapshot task"),
() = storage_task => tracing::trace!("[REPLICA] Exiting storage task"),
)
}
pub async fn start_digest_sub(&self, tx: Sender<(String, Digest)>) {
let mut received = HashMap::<String, Timestamp>::new();
let digest_key = Replica::get_digest_key(&self.key_expr).join("**").unwrap();
tracing::debug!(
"[DIGEST_SUB] Declaring Subscriber named {} on '{}'",
self.name,
digest_key
);
let subscriber = self
.session
.declare_subscriber(&digest_key)
.allowed_origin(Locality::Remote)
.await
.unwrap();
loop {
let sample = match subscriber.recv_async().await {
Ok(sample) => sample,
Err(e) => {
tracing::error!("[DIGEST_SUB] Error receiving sample: {}", e);
continue;
}
};
let from =
&sample.key_expr().as_str()[Replica::get_digest_key(&self.key_expr).len() + 1..];
let digest: Digest = match serde_json::from_reader(sample.payload().reader()) {
Ok(digest) => digest,
Err(e) => {
tracing::error!("[DIGEST_SUB] Error in decoding the digest: {}", e);
continue;
}
};
tracing::trace!(
"[DIGEST_SUB] From {} Received {} ('{}': '{:?}')",
from,
sample.kind(),
sample.key_expr().as_str(),
digest,
);
let ts = digest.timestamp;
let to_be_processed = self
.processing_needed(
from,
digest.timestamp,
digest.checksum,
received.clone(),
digest.config.clone(),
)
.await;
if to_be_processed {
tracing::trace!("[DIGEST_SUB] sending {} to aligner", digest.checksum);
match tx.send_async((from.to_string(), digest)).await {
Ok(()) => {}
Err(e) => {
tracing::error!("[DIGEST_SUB] Error sending digest to aligner: {}", e)
}
}
}
received.insert(from.to_string(), ts);
}
}
pub async fn start_digest_pub(&self, snapshotter: Arc<Snapshotter>) {
let digest_key = Replica::get_digest_key(&self.key_expr)
.join(&self.name)
.unwrap();
tracing::debug!("[DIGEST_PUB] Declaring Publisher on '{}'...", digest_key);
let publisher = self.session.declare_publisher(digest_key).await.unwrap();
let mut interval = interval(self.replica_config.publication_interval);
loop {
let _ = interval.tick().await;
let digest = snapshotter.get_digest().await;
let digest = digest.compress();
let digest_json = serde_json::to_string(&digest).unwrap();
let mut digests_published = self.digests_published.write().await;
digests_published.insert(digest.checksum);
drop(digests_published);
drop(digest);
tracing::trace!("[DIGEST_PUB] Putting Digest: {} ...", digest_json);
match publisher.put(digest_json).await {
Ok(()) => {}
Err(e) => tracing::error!("[DIGEST_PUB] Digest publication failed: {}", e),
}
}
}
async fn processing_needed(
&self,
from: &str,
ts: Timestamp,
checksum: u64,
received: HashMap<String, Timestamp>,
config: DigestConfig,
) -> bool {
if checksum == 0 {
return false;
}
let digests_published = self.digests_published.read().await;
if digests_published.contains(&checksum) {
tracing::trace!("[DIGEST_SUB] Dropping since matching digest already seen");
return false;
}
if received.contains_key(from) && *received.get(from).unwrap() > ts {
tracing::trace!("[DIGEST_SUB] Dropping older digest at {} from {}", ts, from);
return false;
}
if config.delta != self.replica_config.delta
|| config.hot
!= Replica::get_hot_interval_number(
self.replica_config.publication_interval,
self.replica_config.delta,
)
{
tracing::error!("[DIGEST_SUB] Mismatching digest configs, cannot be aligned");
return false;
}
true
}
fn get_digest_key(key_expr: &keyexpr) -> OwnedKeyExpr {
*KE_PREFIX_DIGEST / key_expr
}
pub fn get_hot_interval_number(publication_interval: Duration, delta: Duration) -> usize {
((publication_interval.as_nanos() / delta.as_nanos()) as usize) + 1
}
pub fn get_warm_interval_number(publication_interval: Duration, delta: Duration) -> usize {
Replica::get_hot_interval_number(publication_interval, delta) * 5
}
}