#[cfg(feature = "tracer")]
use opentelemetry::trace::Tracer;
use crate::cluster::EtcdNode;
use crate::etcdpb::etcdserverpb::kv_server::Kv;
use crate::etcdpb::etcdserverpb::{CompactionRequest, CompactionResponse, DeleteRangeRequest, DeleteRangeResponse, PutRequest, PutResponse, RangeRequest, RangeResponse, TxnRequest, TxnResponse};
use tonic::{async_trait, Request, Response, Status};
impl EtcdNode {
async fn deny_remote(&self, key: &[u8], remote: Option<std::net::SocketAddr>)
-> Result<(), Status>
{
if self.policy.read().await.may_serve(key, remote) {
return Ok(());
}
Err(Status::permission_denied(
"this key prefix is served only on a local connection"))
}
}
#[async_trait]
impl Kv for EtcdNode {
async fn range(&self, request: Request<RangeRequest>) -> Result<Response<RangeResponse>, Status> {
self.deny_remote(&request.get_ref().key, request.remote_addr()).await?;
#[cfg(feature = "tracer")]
let _s = self.tracer.read().await.as_ref().map(|t| t.start("get"));
let result = self.get_impl(request).await;
result
}
async fn put(&self, request: Request<PutRequest>) -> Result<Response<PutResponse>, Status> {
self.deny_remote(&request.get_ref().key, request.remote_addr()).await?;
#[cfg(feature = "tracer")]
let _s = self.tracer.read().await.as_ref().map(|t| t.start("put"));
let result = self.put_impl(request).await;
result
}
async fn delete_range(&self, request: Request<DeleteRangeRequest>) -> Result<Response<DeleteRangeResponse>, Status> {
self.deny_remote(&request.get_ref().key, request.remote_addr()).await?;
#[cfg(feature = "tracer")]
let _s = self.tracer.read().await.as_ref().map(|t| t.start("delete"));
let result = self.delete_impl(request).await;
result
}
async fn txn(&self, request: Request<TxnRequest>) -> Result<Response<TxnResponse>, Status> {
let remote = request.remote_addr();
for k in crate::kv::txn_keys(request.get_ref()) {
self.deny_remote(&k, remote).await?;
}
#[cfg(feature = "tracer")]
let _s = self.tracer.read().await.as_ref().map(|t| t.start("txn"));
self.txn_impl(request).await
}
async fn compact(&self, _request: Request<CompactionRequest>) -> Result<Response<CompactionResponse>, Status> {
Ok(Response::new(CompactionResponse { header: None }))
}
}