use std::sync::Arc;
use dynomite::embed::hooks::DatastoreError;
use dynomite::embed::Datastore;
use dynomite::net::client::BoxFuture;
use dynomite::net::ReplicaApplySink;
use crate::proto::http::object::HttpObject;
use crate::proto::replica_wire::decode_peer_op;
use crate::router::PeerOp;
#[derive(Clone)]
pub struct ReplicaApplier {
datastore: Arc<dyn Datastore>,
}
impl std::fmt::Debug for ReplicaApplier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ReplicaApplier").finish_non_exhaustive()
}
}
impl ReplicaApplier {
#[must_use]
pub fn new(datastore: Arc<dyn Datastore>) -> Self {
Self { datastore }
}
async fn apply_op(&self, payload: &[u8]) {
let op = match decode_peer_op(payload) {
Ok(op) => op,
Err(e) => {
tracing::warn!(error = %e, "riak replica: undecodable peer op dropped");
return;
}
};
match op {
PeerOp::Put {
bucket, key, value, ..
} => {
let envelope = HttpObject {
value,
content_type: None,
indexes: Vec::new(),
links: Vec::new(),
};
let storage = envelope.to_storage_bytes();
match self.datastore.riak_put(&bucket, &key, &storage, &[]).await {
Ok(()) | Err(DatastoreError::Unsupported(_)) => {}
Err(e) => {
tracing::warn!(error = %e, "riak replica: local put failed");
}
}
}
PeerOp::Del { bucket, key, .. } => {
match self.datastore.riak_delete(&bucket, &key).await {
Ok(_) | Err(DatastoreError::Unsupported(_)) => {}
Err(e) => {
tracing::warn!(error = %e, "riak replica: local delete failed");
}
}
}
PeerOp::Get { .. } | PeerOp::DtFetch { .. } => {}
PeerOp::DtUpdate {
bucket, key, op, ..
} => {
let store =
crate::crdt_store::CrdtStore::new(std::sync::Arc::clone(&self.datastore));
let res = match op.split_first() {
Some((&crate::crdt_store::DT_WIRE_STATE, body)) => {
store.merge_state(&bucket, &key, body).await
}
Some((&crate::crdt_store::DT_WIRE_OP, body)) => {
match crate::crdt_store::CrdtOp::from_bytes(body) {
Ok(parsed) => store.apply(&bucket, &key, &parsed).await.map(|_| ()),
Err(e) => {
tracing::warn!(error = %e, "riak replica: undecodable dt op");
return;
}
}
}
_ => {
store.merge_state(&bucket, &key, &op).await
}
};
if let Err(e) = res {
if !matches!(
e,
crate::crdt_store::CrdtStoreError::Datastore(DatastoreError::Unsupported(
_
))
) {
tracing::warn!(error = %e, "riak replica: local dt apply/merge failed");
}
}
}
}
}
}
impl ReplicaApplySink for ReplicaApplier {
fn apply<'a>(&'a self, payload: &'a [u8]) -> BoxFuture<'a, ()> {
Box::pin(async move { self.apply_op(payload).await })
}
fn apply_query<'a>(&'a self, payload: &'a [u8]) -> BoxFuture<'a, Option<Vec<u8>>> {
Box::pin(async move {
let op = decode_peer_op(payload).ok()?;
let PeerOp::DtFetch { bucket, key, .. } = op else {
return None;
};
match self.datastore.riak_get(&bucket, &key).await {
Ok(Some(state)) => Some(state),
Ok(None) => Some(Vec::new()),
Err(e) => {
tracing::debug!(error = %e, "riak replica: dt fetch read failed");
Some(Vec::new())
}
}
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::proto::replica_wire::encode_peer_op;
#[tokio::test]
async fn undecodable_payload_is_dropped_without_panic() {
let ds: Arc<dyn Datastore> = Arc::new(dynomite::embed::MemoryDatastore::new());
let applier = ReplicaApplier::new(ds);
applier.apply(&[0xFF, 0x00, 0x01]).await;
let op = PeerOp::Put {
bucket_type: b"default".to_vec(),
bucket: b"b".to_vec(),
key: b"k".to_vec(),
value: b"v".to_vec(),
};
applier.apply(&encode_peer_op(&op)).await;
}
}