use serde_json::json;
use tonic::metadata::MetadataValue;
use tonic::{Request, Status};
use crate::ir::{ComparisonOp, LogicalFilter, LogicalValue};
use crate::proto::udb::core::livequery::services::v1 as lq_pb;
use crate::proto::udb::core::livequery::services::v1::live_query_service_server::LiveQueryService;
use crate::proto::{ErrorDetail, ErrorKind};
use crate::runtime::executor_utils::ERROR_DETAIL_METADATA_KEY;
use super::LiveQueryServiceImpl;
use super::budget::{active_streams, try_acquire_stream_slot};
use super::config::max_streams_per_tenant;
use super::errors::{livequery_backpressure_status, livequery_capability_status};
use super::predicate::{event_matches_tenant_scope, filter_matches_row, resolve_source};
fn decode_detail(status: &Status) -> ErrorDetail {
let raw = status
.metadata()
.get_bin(ERROR_DETAIL_METADATA_KEY)
.expect("error-detail trailer present")
.to_bytes()
.expect("trailer decodes to bytes");
crate::runtime::executor_utils::decode_error_detail_from_raw(&raw)
}
#[tokio::test]
async fn subscribe_rejects_cross_tenant_body() {
let svc = LiveQueryServiceImpl::new();
let mut request = Request::new(lq_pb::SubscribeRequest {
tenant_id: "tenant-b".to_string(),
message_type: "udb.core.lock.entity.v1.Lock".to_string(),
..Default::default()
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.subscribe(request)
.await
.err()
.expect("cross-tenant body must be rejected");
assert_eq!(err.code(), tonic::Code::PermissionDenied);
}
#[tokio::test]
async fn subscribe_missing_message_type_carries_field_violation() {
let svc = LiveQueryServiceImpl::new(); let mut request = Request::new(lq_pb::SubscribeRequest {
tenant_id: "tenant-a".to_string(),
message_type: " ".to_string(),
..Default::default()
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.subscribe(request)
.await
.err()
.expect("missing message_type must be rejected");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(err.message(), "message_type is required");
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, "message_type");
assert_eq!(
detail.field_violations[0].description,
"must be a non-empty source message type"
);
}
#[tokio::test]
async fn subscribe_empty_predicate_field_carries_field_violation() {
let svc = LiveQueryServiceImpl::new(); let mut request = Request::new(lq_pb::SubscribeRequest {
tenant_id: "tenant-a".to_string(),
message_type: "udb.core.lock.entity.v1.Lock".to_string(),
filters: vec![lq_pb::LiveQueryPredicate {
field: " ".to_string(),
op: lq_pb::LiveQueryComparison::Eq as i32,
value: "HELD".to_string(),
}],
..Default::default()
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.subscribe(request)
.await
.err()
.expect("empty predicate field must be rejected");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"live query predicate field must not be empty"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, "filters.field");
assert_eq!(
detail.field_violations[0].description,
"must be a non-empty live query predicate field"
);
}
#[tokio::test]
async fn subscribe_unspecified_predicate_op_carries_field_violation() {
let svc = LiveQueryServiceImpl::new(); let mut request = Request::new(lq_pb::SubscribeRequest {
tenant_id: "tenant-a".to_string(),
message_type: "udb.core.lock.entity.v1.Lock".to_string(),
filters: vec![lq_pb::LiveQueryPredicate {
field: "status".to_string(),
op: lq_pb::LiveQueryComparison::Unspecified as i32,
value: "HELD".to_string(),
}],
..Default::default()
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.subscribe(request)
.await
.err()
.expect("unspecified predicate op must be rejected");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"live query predicate comparison op is unspecified"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, "filters.op");
assert_eq!(
detail.field_violations[0].description,
"must specify a live query predicate comparison operator"
);
}
#[test]
fn per_event_scope_drops_foreign_and_tenantless_keeps_match() {
let topic = "udb.lock.lock.acquired.v1";
assert!(event_matches_tenant_scope(
topic,
&json!({"tenant_id": "tenant-a", "project_id": "project-a"}),
"tenant-a",
"project-a",
));
assert!(!event_matches_tenant_scope(
topic,
&json!({"tenant_id": "tenant-b"}),
"tenant-a",
"",
));
assert!(!event_matches_tenant_scope(
topic,
&json!({"event_id": "e1"}),
"tenant-a",
"",
));
assert!(!event_matches_tenant_scope(
"app.customer.changed",
&json!({"tenant_id": "tenant-a"}),
"tenant-a",
"",
));
assert!(!event_matches_tenant_scope(
topic,
&json!({"tenant_id": "tenant-a"}),
"",
"",
));
assert!(!event_matches_tenant_scope(
topic,
&json!({"tenant_id": "tenant-a", "project_id": "project-b"}),
"tenant-a",
"project-a",
));
}
#[test]
fn unknown_source_fails_closed() {
let err = resolve_source("not.a.real.Entity")
.err()
.expect("unknown source must fail closed");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"live query source 'not.a.real.Entity' is not a known UDB entity"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, "message_type");
assert_eq!(
detail.field_violations[0].description,
"must name exactly one known tenant-scoped UDB entity"
);
}
#[test]
fn backpressure_status_carries_typed_quota_detail() {
let status = livequery_backpressure_status(
"subscriber_channel",
"live query subscriber too slow; stream closed",
);
assert_eq!(status.code(), tonic::Code::ResourceExhausted);
let detail = decode_detail(&status);
assert_eq!(detail.kind, ErrorKind::Quota as i32);
assert!(detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
assert_eq!(detail.backend, "livequery");
assert_eq!(detail.operation, "subscriber_channel");
}
#[test]
fn stream_budget_refuses_over_budget_and_reuses_freed_slot() {
let tenant = "tenant-stream-budget-a";
let other_tenant = "tenant-stream-budget-b";
let budget = max_streams_per_tenant();
let mut slots = Vec::with_capacity(budget);
for _ in 0..budget {
slots.push(try_acquire_stream_slot(tenant).expect("within budget"));
}
let err = try_acquire_stream_slot(tenant)
.err()
.expect("over-budget subscribe must be refused");
assert_eq!(err.code(), tonic::Code::ResourceExhausted);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Quota as i32);
assert_eq!(detail.backend, "livequery");
assert_eq!(detail.operation, "tenant_stream_budget");
assert!(!detail.retryable);
let other_slot = try_acquire_stream_slot(other_tenant).expect("other tenant unaffected");
drop(other_slot);
drop(slots.pop());
let reacquired = try_acquire_stream_slot(tenant).expect("freed slot must be reusable");
drop(reacquired);
drop(slots);
let map = active_streams()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
assert!(!map.contains_key(tenant));
assert!(!map.contains_key(other_tenant));
}
#[test]
fn livequery_missing_runtime_capability_carries_typed_detail() {
let err = livequery_capability_status(
"native_entity_dispatch",
"runtime_native_entity_dispatch",
"live query service requires runtime native-entity dispatch (no runtime configured)",
);
assert_eq!(err.code(), tonic::Code::FailedPrecondition);
assert_eq!(
err.message(),
"live query service requires runtime native-entity dispatch (no runtime configured)"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Capability as i32);
assert_eq!(detail.backend, "livequery");
assert_eq!(detail.operation, "native_entity_dispatch");
assert_eq!(detail.capability_required, "runtime_native_entity_dispatch");
assert!(!detail.retryable);
}
#[test]
fn single_row_evaluator_matches_predicate() {
let filter = LogicalFilter::And(vec![
LogicalFilter::Comparison {
field: "status".to_string(),
op: ComparisonOp::Eq,
value: LogicalValue::String("HELD".to_string()),
},
LogicalFilter::Comparison {
field: "fencing_token".to_string(),
op: ComparisonOp::Ge,
value: LogicalValue::Int(5),
},
]);
assert!(filter_matches_row(
&filter,
&json!({"status": "HELD", "fencing_token": 7})
));
assert!(!filter_matches_row(
&filter,
&json!({"status": "HELD", "fencing_token": 3})
));
assert!(!filter_matches_row(
&filter,
&json!({"status": "RELEASED", "fencing_token": 9})
));
assert!(!filter_matches_row(&filter, &json!({"status": "HELD"})));
}