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 { .. } => {}
}
}
}
impl ReplicaApplySink for ReplicaApplier {
fn apply<'a>(&'a self, payload: &'a [u8]) -> BoxFuture<'a, ()> {
Box::pin(async move { self.apply_op(payload).await })
}
}
#[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;
}
}