use std::sync::Arc;
use serde_json::Value as JsonValue;
use crate::runtime::backend_context::{
AppliedContext, BackendContextEnforcer, ContextEffect, enforce_with_mechanism,
};
use crate::runtime::executor_utils::{build_probe, executor_timeout_duration};
use crate::runtime::executors::{
BackendExecutor, BackendHealth, BackendProbe, MutationExecutor, ObjectExecutor, QueryExecutor,
ResourceAdminExecutor, SearchExecutor,
};
#[derive(Clone)]
pub struct MemcachedClient {
inner: Arc<memcache::Client>,
}
impl std::fmt::Debug for MemcachedClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("MemcachedClient").finish()
}
}
impl MemcachedClient {
pub fn connect(url: &str) -> Result<Self, String> {
let client =
memcache::connect(url).map_err(|e| format!("memcached connect failed: {e}"))?;
Ok(Self {
inner: Arc::new(client),
})
}
}
fn validate_memcached_key(key: &str) -> Result<(), tonic::Status> {
if key.is_empty() || key.len() > 250 {
return Err(tonic::Status::invalid_argument(format!(
"memcached key must be 1–250 bytes (got {})",
key.len()
)));
}
if key.chars().any(|c| c.is_whitespace() || c.is_control()) {
return Err(tonic::Status::invalid_argument(
"memcached key may not contain whitespace or control characters",
));
}
Ok(())
}
#[derive(Clone)]
pub struct MemcachedExecutor {
client: MemcachedClient,
}
impl std::fmt::Debug for MemcachedExecutor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("MemcachedExecutor").finish()
}
}
impl MemcachedExecutor {
pub fn new(client: MemcachedClient) -> Self {
Self { client }
}
async fn blocking<F, R>(&self, f: F) -> Result<R, tonic::Status>
where
F: FnOnce(Arc<memcache::Client>) -> Result<R, String> + Send + 'static,
R: Send + 'static,
{
let inner = self.client.inner.clone();
tokio::time::timeout(
executor_timeout_duration(),
tokio::task::spawn_blocking(move || f(inner)),
)
.await
.map_err(|_| tonic::Status::deadline_exceeded("memcached operation deadline exceeded"))?
.map_err(|e| tonic::Status::internal(format!("memcached blocking task join: {e}")))?
.map_err(|e| tonic::Status::unavailable(format!("memcached: {e}")))
}
}
impl BackendContextEnforcer for MemcachedExecutor {
fn backend_label(&self) -> &str {
"memcached"
}
fn enforce(&self, ctx: &AppliedContext) -> ContextEffect {
enforce_with_mechanism(ctx, "key namespace prefix udb:{project}:{tenant}:")
}
}
impl BackendHealth for MemcachedExecutor {
async fn ping(&self) -> Result<(), String> {
self.blocking(|c| c.version().map(|_| ()).map_err(|e| e.to_string()))
.await
.map_err(|s| s.message().to_string())
}
}
fn parse_kv_request(
request_json: &str,
) -> Result<(String, String, Option<Vec<u8>>, u32), tonic::Status> {
let v: JsonValue = serde_json::from_str(request_json)
.map_err(|e| tonic::Status::invalid_argument(format!("invalid dispatch JSON: {e}")))?;
let op = v
.get("op")
.or_else(|| v.get("operation"))
.and_then(|x| x.as_str())
.ok_or_else(|| {
tonic::Status::invalid_argument("missing `op`/`operation` in dispatch request")
})?
.to_string();
let key = if matches!(op.as_str(), "scan" | "cache_scan") {
String::new()
} else {
let key = v
.get("key")
.and_then(|x| x.as_str())
.ok_or_else(|| tonic::Status::invalid_argument("missing `key` in dispatch request"))?
.to_string();
validate_memcached_key(&key)?;
key
};
let value = match v.get("value") {
Some(JsonValue::String(s)) => {
if let Some(b64) = s.strip_prefix("base64:") {
use base64::Engine as _;
Some(
base64::engine::general_purpose::STANDARD
.decode(b64)
.map_err(|e| {
tonic::Status::invalid_argument(format!("bad base64 value: {e}"))
})?,
)
} else {
Some(s.as_bytes().to_vec())
}
}
Some(JsonValue::Array(arr)) => Some(
arr.iter()
.filter_map(|n| n.as_u64().map(|u| u as u8))
.collect(),
),
Some(JsonValue::Null) | None => None,
Some(other) => Some(other.to_string().into_bytes()),
};
let ttl = v
.get("ttl_seconds")
.or_else(|| v.get("ttl"))
.and_then(|x| x.as_u64())
.unwrap_or(0) as u32;
Ok((op, key, value, ttl))
}
impl QueryExecutor for MemcachedExecutor {
async fn query(&self, request_json: &str) -> Result<String, tonic::Status> {
let (op, key, _, _) = parse_kv_request(request_json)?;
if op == "scan" || op == "cache_scan" {
return Ok(serde_json::json!({
"entries": [],
"next_page_token": ""
})
.to_string());
}
if op != "get" && op != "cache_get" {
return Err(tonic::Status::invalid_argument(format!(
"memcached query expects op=\"get\", got '{op}'"
)));
}
let value = self
.blocking(move |c| c.get::<Vec<u8>>(&key).map_err(|e| e.to_string()))
.await?;
let resp = match value {
Some(bytes) => {
use base64::Engine as _;
let b64 = base64::engine::general_purpose::STANDARD.encode(&bytes);
serde_json::json!({ "found": true, "hit": true, "value": format!("base64:{b64}") })
}
None => serde_json::json!({ "found": false, "hit": false }),
};
Ok(resp.to_string())
}
}
impl MutationExecutor for MemcachedExecutor {
async fn mutate(&self, request_json: &str) -> Result<String, tonic::Status> {
let (op, key, value, ttl) = parse_kv_request(request_json)?;
match op.as_str() {
"set" | "cache_set" => {
let bytes = value
.ok_or_else(|| tonic::Status::invalid_argument("set op requires `value`"))?;
self.blocking(move |c| {
c.set(&key, bytes.as_slice(), ttl)
.map_err(|e| e.to_string())
})
.await?;
Ok(serde_json::json!({ "ok": true, "op": "set" }).to_string())
}
"delete" | "del" | "cache_delete" => self
.blocking(move |c| {
c.delete(&key)
.map(|deleted| deleted)
.map_err(|e| e.to_string())
})
.await
.map(|deleted| {
serde_json::json!({ "ok": true, "op": "delete", "deleted": deleted })
.to_string()
}),
other => Err(tonic::Status::invalid_argument(format!(
"memcached mutate op '{other}' is not supported"
))),
}
}
}
impl SearchExecutor for MemcachedExecutor {
async fn search(&self, _request_json: &str) -> Result<String, tonic::Status> {
Err(tonic::Status::failed_precondition(
"UDB_UNSUPPORTED_OPERATION: Memcached has no search surface; route to a vector / text backend",
))
}
}
impl ObjectExecutor for MemcachedExecutor {
async fn get_object(&self, _request_json: &str) -> Result<Vec<u8>, tonic::Status> {
Err(tonic::Status::failed_precondition(
"UDB_UNSUPPORTED_OPERATION: Memcached is not an object store; route to S3/MinIO",
))
}
async fn put_object(
&self,
_request_json: &str,
_bytes: Vec<u8>,
) -> Result<String, tonic::Status> {
Err(tonic::Status::failed_precondition(
"UDB_UNSUPPORTED_OPERATION: Memcached is not an object store; route to S3/MinIO",
))
}
}
impl ResourceAdminExecutor for MemcachedExecutor {
async fn ensure_resource(
&self,
_resource_name: &str,
_spec_json: &str,
) -> Result<(), tonic::Status> {
Ok(())
}
async fn drop_resource(&self, _resource_name: &str) -> Result<(), tonic::Status> {
Err(tonic::Status::failed_precondition(
"UDB_UNSUPPORTED_OPERATION: Memcached has no per-resource drop; use flush_all via operator console",
))
}
async fn list_resources(&self) -> Result<Vec<String>, tonic::Status> {
Ok(Vec::new())
}
}
impl BackendExecutor for MemcachedExecutor {
async fn transaction(&self, _request_json: &str) -> Result<String, tonic::Status> {
Err(tonic::Status::failed_precondition(
"UDB_UNSUPPORTED_OPERATION: Memcached has no multi-key transaction primitive",
))
}
async fn probe(&self) -> Result<BackendProbe, tonic::Status> {
Ok(build_probe("memcached", self.ping().await))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn validate_key_rejects_empty() {
assert!(validate_memcached_key("").is_err());
}
#[test]
fn validate_key_rejects_too_long() {
let key = "a".repeat(251);
assert!(validate_memcached_key(&key).is_err());
}
#[test]
fn validate_key_rejects_whitespace() {
assert!(validate_memcached_key("has space").is_err());
assert!(validate_memcached_key("has\nnewline").is_err());
assert!(validate_memcached_key("has\ttab").is_err());
}
#[test]
fn validate_key_accepts_canonical_template() {
assert!(validate_memcached_key("udb:proj:tenant:msg:pk").is_ok());
}
#[test]
fn parse_kv_get_request() {
let req = r#"{"op":"get","key":"udb:a:b:msg:pk"}"#;
let (op, key, value, ttl) = parse_kv_request(req).unwrap();
assert_eq!(op, "get");
assert_eq!(key, "udb:a:b:msg:pk");
assert!(value.is_none());
assert_eq!(ttl, 0);
}
#[test]
fn parse_kv_set_request_with_base64_value() {
let req = r#"{"op":"set","key":"k","value":"base64:aGVsbG8=","ttl_seconds":60}"#;
let (op, _, value, ttl) = parse_kv_request(req).unwrap();
assert_eq!(op, "set");
assert_eq!(value.unwrap(), b"hello");
assert_eq!(ttl, 60);
}
#[test]
fn parse_kv_set_request_with_raw_string() {
let req = r#"{"op":"set","key":"k","value":"hello world"}"#;
let (_, _, value, _) = parse_kv_request(req).unwrap();
assert_eq!(value.unwrap(), b"hello world");
}
#[test]
fn parse_kv_rejects_invalid_key() {
let req = r#"{"op":"get","key":"has space"}"#;
let err = parse_kv_request(req).unwrap_err();
assert_eq!(err.code(), tonic::Code::InvalidArgument);
}
#[test]
fn enforce_returns_enforced_with_namespace_mechanism() {
let ctx = AppliedContext {
tenant_id: "acme".into(),
..Default::default()
};
let effect = if ctx.is_empty() {
ContextEffect::Advisory {
recorded_in: "no_context".into(),
}
} else {
ContextEffect::Enforced {
mechanism: "key namespace prefix udb:{project}:{tenant}:".into(),
}
};
assert!(matches!(effect, ContextEffect::Enforced { .. }));
}
}