use super::*;
use crate::runtime::executor_utils::{invalid_argument_fields, json_into_struct, struct_to_json};
fn store_rpc_invalid_fields<I, F, D>(message: impl Into<String>, fields: I) -> Status
where
I: IntoIterator<Item = (F, D)>,
F: Into<String>,
D: Into<String>,
{
invalid_argument_fields(message, fields)
}
fn require_resource_backend(
resource: Option<&crate::proto::StoreResource>,
) -> Result<String, Status> {
resource
.map(|r| r.backend.clone())
.filter(|b| !b.trim().is_empty())
.ok_or_else(|| {
store_rpc_invalid_fields(
"resource.backend is required",
[("resource.backend", "must be a non-empty backend name")],
)
})
}
impl DataBrokerService {
async fn run_store_op(
&self,
security: &SecurityContext,
resource: Option<&crate::proto::StoreResource>,
write: bool,
method: &str,
spec: serde_json::Value,
) -> Result<String, Status> {
let backend = require_resource_backend(resource)?;
let resource_name = resource
.map(|r| r.resource_name.clone())
.unwrap_or_default();
self.execute_backend_operation(
&security.request_context(),
&backend,
write,
method.to_string(),
resource_name,
spec.to_string(),
)
.await
}
pub(crate) async fn cache_get_inner(
&self,
request: Request<crate::proto::CacheGetRequest>,
) -> Result<Response<crate::proto::CacheGetResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("CacheGet", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "cache.get")
.await
{
return self.record_grpc("CacheGet", started, Err(e));
}
let spec = serde_json::json!({ "operation": "get", "key": req.key });
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| {
let v = parse_json(&json);
let found = v
.get("hit")
.and_then(serde_json::Value::as_bool)
.unwrap_or_else(|| v.get("value").map(|x| !x.is_null()).unwrap_or(false));
let value = v
.get("value")
.and_then(serde_json::Value::as_str)
.map(cache_value_bytes)
.unwrap_or_default();
Response::new(crate::proto::CacheGetResponse {
found,
value,
..Default::default()
})
});
self.record_grpc("CacheGet", started, out)
}
pub(crate) async fn cache_set_inner(
&self,
request: Request<crate::proto::CacheSetRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("CacheSet", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "cache.set")
.await
{
return self.record_grpc("CacheSet", started, Err(e));
}
let mut spec = serde_json::json!({
"operation": "set",
"key": req.key,
"value": String::from_utf8_lossy(&req.value),
});
if req.ttl_seconds > 0 {
spec["ttl"] = serde_json::json!(req.ttl_seconds);
}
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("CacheSet", started, out)
}
pub(crate) async fn cache_delete_inner(
&self,
request: Request<crate::proto::CacheDeleteRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("CacheDelete", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "cache.delete")
.await
{
return self.record_grpc("CacheDelete", started, Err(e));
}
let spec = serde_json::json!({ "operation": "delete", "key": req.key });
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("CacheDelete", started, out)
}
pub(crate) async fn cache_scan_inner(
&self,
request: Request<crate::proto::CacheScanRequest>,
) -> Result<Response<crate::proto::CacheScanResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("CacheScan", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "cache.scan")
.await
{
return self.record_grpc("CacheScan", started, Err(e));
}
let spec = serde_json::json!({
"operation": "scan",
"pattern": req.key_pattern,
"limit": req.limit,
"cursor": req.page_token,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| {
let v = parse_json(&json);
let entries = v
.get("entries")
.and_then(serde_json::Value::as_array)
.map(|items| {
items
.iter()
.map(|e| crate::proto::CacheEntry {
key: e
.get("key")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string(),
value: e
.get("value")
.and_then(serde_json::Value::as_str)
.map(cache_value_bytes)
.unwrap_or_default(),
..Default::default()
})
.collect()
})
.unwrap_or_default();
let next = v
.get("next_page_token")
.or_else(|| v.get("cursor"))
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string();
Response::new(crate::proto::CacheScanResponse {
entries,
next_page_token: next,
..Default::default()
})
});
self.record_grpc("CacheScan", started, out)
}
pub(crate) async fn document_get_inner(
&self,
request: Request<crate::proto::DocumentGetRequest>,
) -> Result<Response<crate::proto::DocumentSet>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("DocumentGet", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "document.get")
.await
{
return self.record_grpc("DocumentGet", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("DocumentGet", started, Err(e));
}
let spec = serde_json::json!({
"collection": collection_of(&req.resource),
"filter": { "_id": req.document_id },
"limit": 1,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| Response::new(document_set_from_json(&json)));
self.record_grpc("DocumentGet", started, out)
}
pub(crate) async fn document_find_inner(
&self,
request: Request<crate::proto::DocumentFindRequest>,
) -> Result<Response<crate::proto::DocumentSet>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("DocumentFind", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "document.find")
.await
{
return self.record_grpc("DocumentFind", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("DocumentFind", started, Err(e));
}
let spec = serde_json::json!({
"collection": collection_of(&req.resource),
"filter": struct_field(&req.filter),
"limit": req.limit,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| Response::new(document_set_from_json(&json)));
self.record_grpc("DocumentFind", started, out)
}
pub(crate) async fn document_upsert_inner(
&self,
request: Request<crate::proto::DocumentUpsertRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("DocumentUpsert", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "document.upsert")
.await
{
return self.record_grpc("DocumentUpsert", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("DocumentUpsert", started, Err(e));
}
let spec = serde_json::json!({
"collection": collection_of(&req.resource),
"operation": if req.replace { "update" } else { "upsert" },
"filter": { "_id": req.document_id },
"update": struct_field(&req.document),
});
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("DocumentUpsert", started, out)
}
pub(crate) async fn document_delete_inner(
&self,
request: Request<crate::proto::DocumentDeleteRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("DocumentDelete", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "document.delete")
.await
{
return self.record_grpc("DocumentDelete", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("DocumentDelete", started, Err(e));
}
let filter = if req.document_id.is_empty() {
struct_field(&req.filter)
} else {
serde_json::json!({ "_id": req.document_id })
};
let spec = serde_json::json!({
"collection": collection_of(&req.resource),
"operation": "delete",
"filter": filter,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("DocumentDelete", started, out)
}
pub(crate) async fn graph_query_inner(
&self,
request: Request<crate::proto::GraphQueryRequest>,
) -> Result<Response<crate::proto::GraphResultSet>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("GraphQuery", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "graph.query")
.await
{
return self.record_grpc("GraphQuery", started, Err(e));
}
let spec = serde_json::json!({
"cypher": req.query,
"parameters": struct_field(&req.parameters),
"limit": req.limit,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| {
Response::new(crate::proto::GraphResultSet {
records: structs_from_result(&json),
..Default::default()
})
});
self.record_grpc("GraphQuery", started, out)
}
pub(crate) async fn graph_mutate_inner(
&self,
request: Request<crate::proto::GraphMutationRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("GraphMutate", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "graph.mutate")
.await
{
return self.record_grpc("GraphMutate", started, Err(e));
}
let spec = serde_json::json!({
"operation": "cypher",
"cypher": req.query,
"parameters": struct_field(&req.parameters),
});
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("GraphMutate", started, out)
}
pub(crate) async fn time_series_write_inner(
&self,
request: Request<crate::proto::TimeSeriesWriteRequest>,
) -> Result<Response<MutationResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("TimeSeriesWrite", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "timeseries.write")
.await
{
return self.record_grpc("TimeSeriesWrite", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("TimeSeriesWrite", started, Err(e));
}
let rows: Vec<serde_json::Value> = req
.points
.iter()
.map(|p| {
let mut row = serde_json::Map::new();
for (k, v) in &p.tags {
row.insert(k.clone(), serde_json::json!(v));
}
for (k, v) in &p.values {
row.insert(k.clone(), serde_json::json!(v));
}
if let serde_json::Value::Object(map) = struct_field(&p.fields) {
row.extend(map);
}
serde_json::Value::Object(row)
})
.collect();
let spec = serde_json::json!({
"table": collection_of(&req.resource),
"rows": rows,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), true, "mutate", spec)
.await
.map(|json| Response::new(mutation_from_json(&json)));
self.record_grpc("TimeSeriesWrite", started, out)
}
pub(crate) async fn time_series_query_inner(
&self,
request: Request<crate::proto::TimeSeriesQueryRequest>,
) -> Result<Response<crate::proto::TimeSeriesQueryResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("TimeSeriesQuery", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "timeseries.query")
.await
{
return self.record_grpc("TimeSeriesQuery", started, Err(e));
}
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("TimeSeriesQuery", started, Err(e));
}
let spec = serde_json::json!({
"table": collection_of(&req.resource),
"filter": struct_field(&req.filter),
"limit": req.limit,
});
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| {
let points = structs_from_result(&json)
.into_iter()
.map(|fields| crate::proto::TimeSeriesPoint {
fields: Some(fields),
..Default::default()
})
.collect();
Response::new(crate::proto::TimeSeriesQueryResponse {
points,
..Default::default()
})
});
self.record_grpc("TimeSeriesQuery", started, out)
}
pub(crate) async fn analytical_query_inner(
&self,
request: Request<crate::proto::AnalyticalQueryRequest>,
) -> Result<Response<crate::proto::AnalyticalQueryResponse>, Status> {
let started = Instant::now();
let security = match security_from_request(&request) {
Ok(s) => s,
Err(e) => return self.record_grpc("AnalyticalQuery", started, Err(e)),
};
let req = request.into_inner();
if let Err(e) = self
.authorize(&security, msg_type(&req.resource), "analytical.query")
.await
{
return self.record_grpc("AnalyticalQuery", started, Err(e));
}
if req.query.trim().is_empty() {
if let Err(e) = require_collection(&req.resource) {
return self.record_grpc("AnalyticalQuery", started, Err(e));
}
}
let spec = if req.query.trim().is_empty() {
serde_json::json!({ "table": collection_of(&req.resource), "limit": req.limit })
} else {
serde_json::json!({ "sql": req.query })
};
let out = self
.run_store_op(&security, req.resource.as_ref(), false, "query", spec)
.await
.map(|json| {
let rows = structs_from_result(&json)
.into_iter()
.map(|s| crate::proto::Row {
fields: s.fields.into_iter().collect(),
})
.collect();
Response::new(crate::proto::AnalyticalQueryResponse {
rows,
..Default::default()
})
});
self.record_grpc("AnalyticalQuery", started, out)
}
}
fn parse_json(s: &str) -> serde_json::Value {
serde_json::from_str(s).unwrap_or(serde_json::Value::Null)
}
fn msg_type(resource: &Option<crate::proto::StoreResource>) -> &str {
resource
.as_ref()
.map(|r| r.message_type.as_str())
.unwrap_or("")
}
fn collection_of(resource: &Option<crate::proto::StoreResource>) -> String {
resource
.as_ref()
.map(|r| {
if r.resource_name.is_empty() {
r.message_type.clone()
} else {
r.resource_name.clone()
}
})
.unwrap_or_default()
}
fn require_collection(resource: &Option<crate::proto::StoreResource>) -> Result<(), Status> {
if collection_of(resource).trim().is_empty() {
return Err(store_rpc_invalid_fields(
"resource.resource_name (or resource.message_type) is required",
[
(
"resource.resource_name",
"must be non-empty when resource.message_type is empty",
),
(
"resource.message_type",
"must be non-empty when resource.resource_name is empty",
),
],
));
}
Ok(())
}
fn cache_value_bytes(value: &str) -> Vec<u8> {
if let Some(encoded) = value.strip_prefix("base64:") {
use base64::Engine as _;
return base64::engine::general_purpose::STANDARD
.decode(encoded)
.unwrap_or_default();
}
value.as_bytes().to_vec()
}
fn struct_field(value: &Option<prost_types::Struct>) -> serde_json::Value {
value
.as_ref()
.map(struct_to_json)
.unwrap_or_else(|| serde_json::json!({}))
}
fn mutation_from_json(json: &str) -> MutationResponse {
let v = parse_json(json);
let affected_rows = v
.get("affected_rows")
.and_then(serde_json::Value::as_i64)
.or_else(|| v.get("inserted_id").map(|_| 1))
.unwrap_or(0);
MutationResponse {
affected_rows,
..Default::default()
}
}
fn document_set_from_json(json: &str) -> crate::proto::DocumentSet {
crate::proto::DocumentSet {
documents: structs_from_result(json),
..Default::default()
}
}
fn structs_from_result(json: &str) -> Vec<prost_types::Struct> {
let mut v = parse_json(json);
let array = match &mut v {
serde_json::Value::Array(items) => Some(std::mem::take(items)),
serde_json::Value::Object(map) => {
let slot = if map.contains_key("rows") {
map.get_mut("rows")
} else if map.contains_key("records") {
map.get_mut("records")
} else {
map.get_mut("documents")
};
slot.and_then(serde_json::Value::as_array_mut)
.map(std::mem::take)
}
_ => None,
};
array
.map(|items| items.into_iter().filter_map(json_into_struct).collect())
.unwrap_or_default()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::proto::{ErrorDetail, ErrorKind};
use crate::runtime::executor_utils::ERROR_DETAIL_METADATA_KEY;
use prost::Message as _;
use tonic::{Code, Status};
fn decode_detail(status: &Status) -> ErrorDetail {
let raw = status
.metadata()
.get_bin(ERROR_DETAIL_METADATA_KEY)
.expect("typed detail trailer is present");
crate::runtime::executor_utils::decode_error_detail_from_raw(&raw)
}
fn assert_validation_fields(status: &Status, expected: &[(&str, &str)]) {
assert_eq!(status.code(), Code::InvalidArgument);
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), expected.len());
for (actual, (field, description)) in detail.field_violations.iter().zip(expected) {
assert_eq!(actual.field, *field);
assert_eq!(actual.description, *description);
}
}
#[test]
fn store_rpc_missing_backend_carries_field_violation() {
let err = require_resource_backend(Some(&crate::proto::StoreResource {
backend: " ".to_string(),
..Default::default()
}))
.expect_err("missing resource.backend must fail before backend dispatch");
assert_eq!(err.message(), "resource.backend is required");
assert_validation_fields(
&err,
&[("resource.backend", "must be a non-empty backend name")],
);
}
#[test]
fn store_rpc_missing_collection_carries_field_violations() {
let err = require_collection(&Some(crate::proto::StoreResource {
resource_name: " ".to_string(),
message_type: " ".to_string(),
..Default::default()
}))
.expect_err("missing collection identifier must fail before backend dispatch");
assert_eq!(
err.message(),
"resource.resource_name (or resource.message_type) is required"
);
assert_validation_fields(
&err,
&[
(
"resource.resource_name",
"must be non-empty when resource.message_type is empty",
),
(
"resource.message_type",
"must be non-empty when resource.resource_name is empty",
),
],
);
}
}