use std::sync::Arc;
use std::time::{Duration, Instant};
use tonic::metadata::MetadataValue;
use tonic::{Request, Status};
use crate::proto::udb::core::embedding::services::v1 as embedding_pb;
use crate::proto::udb::core::embedding::services::v1::embedding_service_server::EmbeddingService;
use crate::proto::{ErrorDetail, ErrorKind};
use crate::runtime::catalog::CatalogManager;
use crate::runtime::executor_utils::ERROR_DETAIL_METADATA_KEY;
use super::EmbeddingServiceImpl;
use super::config::{EMBEDDING_BACKFILL_PAGE_LIMIT, STATUS_ACTIVE};
use super::errors::{
embedding_capability_status, embedding_source_not_found_status, require_source_tenant_column,
validate_reported_vector,
};
use super::events::{bound_embedding_text, build_work_event_payload};
use super::handlers::{parse_grpc_timeout, remaining_before_deadline, retrieve_hit_payload_json};
use super::model::{StoredSource, build_embedding_point, merge_retrieve_filter};
use super::workers::{
backfill_read_context, backfill_select_request, embedding_teardown_jobs_sql,
embedding_teardown_point_ids_sql, embedding_work_jobs_sql, extract_source_text,
parse_source_text_fields, source_change_is_delete, source_change_row, source_change_row_pk,
};
fn decode_detail(status: &Status) -> ErrorDetail {
let raw = status
.metadata()
.get_bin(ERROR_DETAIL_METADATA_KEY)
.expect("typed detail metadata");
crate::runtime::executor_utils::decode_error_detail_from_raw(&raw)
}
fn assert_single_field_violation(status: &Status, field: &str, description: &str) {
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 1);
assert_eq!(detail.field_violations[0].field, field);
assert_eq!(detail.field_violations[0].description, description);
}
fn assert_schema_not_found_detail(status: &Status, operation: &str) {
assert_eq!(status.code(), tonic::Code::NotFound);
assert_eq!(status.message(), "embedding source not found");
let detail = decode_detail(status);
assert_eq!(detail.kind, ErrorKind::Schema as i32);
assert_eq!(detail.backend, "embedding");
assert_eq!(detail.operation, operation);
assert_eq!(detail.capability_required, "embedding_source_not_found");
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
}
#[test]
fn work_event_payload_has_no_credentials() {
let payload = build_work_event_payload(
"acme",
"contacts",
"row-1",
"Ada Lovelace, London",
"text-embedding-3-small",
"acme_contacts_vectors",
8000,
);
let object = payload.as_object().expect("object payload");
for key in [
"tenant_id",
"source",
"row_pk",
"text",
"model_id",
"target_collection",
] {
assert!(object.contains_key(key), "missing routing key {key}");
}
assert_eq!(object.get("row_pk").and_then(|v| v.as_str()), Some("row-1"));
assert_eq!(
object.get("text").and_then(|v| v.as_str()),
Some("Ada Lovelace, London")
);
for forbidden in [
"api_key",
"apikey",
"secret",
"token",
"password",
"credential",
"credentials",
"authorization",
"key",
] {
assert!(
!object.contains_key(forbidden),
"work payload leaked credential-shaped key {forbidden}"
);
}
}
#[test]
fn work_event_text_is_bounded() {
assert_eq!(bound_embedding_text("hello world", 100), "hello world");
assert_eq!(bound_embedding_text("abcde", 5), "abcde");
assert_eq!(
bound_embedding_text("alpha beta gamma delta", 12),
"alpha beta"
);
let hard = bound_embedding_text("abcdefghijklmno", 5);
assert_eq!(hard, "abcde");
assert_eq!(hard.chars().count(), 5);
assert!(
bound_embedding_text("café ☕ ambulance 🚑 dispatch", 8)
.chars()
.count()
<= 8
);
assert_eq!(bound_embedding_text("ééééééé", 3).chars().count(), 3);
assert_eq!(bound_embedding_text("anything", 0), "");
let payload = build_work_event_payload("t", "s", "r", "alpha beta gamma delta", "m", "c", 12);
assert_eq!(
payload.get("text").and_then(|v| v.as_str()),
Some("alpha beta")
);
}
#[tokio::test]
async fn report_embedding_rejects_cross_tenant_body() {
let svc = EmbeddingServiceImpl::new(); let mut request = Request::new(embedding_pb::ReportEmbeddingRequest {
tenant_id: "tenant-b".to_string(),
source_name: "contacts".to_string(),
row_pk: "row-1".to_string(),
vector: vec![0.1, 0.2, 0.3],
model: "m".to_string(),
dims: 3,
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.report_embedding(request)
.await
.expect_err("cross-tenant body must be rejected");
assert_eq!(err.code(), tonic::Code::PermissionDenied);
}
#[tokio::test]
async fn register_source_missing_required_fields_carries_field_violations() {
let svc = EmbeddingServiceImpl::new(); let mut request = Request::new(embedding_pb::RegisterSourceRequest {
tenant_id: "tenant-a".to_string(),
source_name: " ".to_string(),
source_message_type: String::new(),
text_fields: Vec::new(),
target_collection: "contacts_vec".to_string(),
model_id: String::new(),
metadata_json: String::new(),
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.register_source(request)
.await
.expect_err("missing source identity must fail");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"source_name and source_message_type are required"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 2);
assert_eq!(detail.field_violations[0].field, "source_name");
assert_eq!(
detail.field_violations[0].description,
"must be a non-empty embedding source name"
);
assert_eq!(detail.field_violations[1].field, "source_message_type");
assert_eq!(
detail.field_violations[1].description,
"must be a non-empty source message type"
);
}
#[tokio::test]
async fn report_embedding_missing_identity_carries_field_violations() {
let svc = EmbeddingServiceImpl::new(); let mut request = Request::new(embedding_pb::ReportEmbeddingRequest {
tenant_id: "tenant-a".to_string(),
source_name: String::new(),
row_pk: " ".to_string(),
vector: vec![0.1, 0.2],
model: "m".to_string(),
dims: 2,
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.report_embedding(request)
.await
.expect_err("missing report identity must fail");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(err.message(), "source_name and row_pk are required");
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Validation as i32);
assert_eq!(detail.field_violations.len(), 2);
assert_eq!(detail.field_violations[0].field, "source_name");
assert_eq!(detail.field_violations[1].field, "row_pk");
}
#[tokio::test]
async fn retrieve_missing_query_vector_carries_field_violation() {
let svc = EmbeddingServiceImpl::new(); let mut request = Request::new(embedding_pb::RetrieveRequest {
tenant_id: "tenant-a".to_string(),
source_name: "contacts".to_string(),
query_text: String::new(),
query_vector: Vec::new(),
top_k: 10,
filter_json: String::new(),
score_threshold: 0.0,
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.retrieve(request)
.await
.expect_err("missing query vector must fail");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"query_vector is required (the broker does not embed queries; supply a vector)"
);
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, "query_vector");
assert_eq!(
detail.field_violations[0].description,
"must contain at least one embedding dimension"
);
}
#[test]
fn source_message_type_missing_from_catalog_carries_field_violation() {
let catalog = Arc::new(CatalogManager::new(
crate::generation::CatalogManifest::default(),
));
let svc = EmbeddingServiceImpl::new().with_catalog(Some(catalog));
let err = svc
.resolve_source_tenant_column("default", "acme.crm.entity.v1.Contact")
.expect_err("unknown source message type must fail before store/vector access");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"source_message_type 'acme.crm.entity.v1.Contact' is not present in the active catalog manifest"
);
assert_single_field_violation(
&err,
"source_message_type",
"must be present in the active catalog manifest",
);
}
#[test]
fn embedding_missing_setup_capabilities_carry_typed_detail() {
for (operation, capability, message) in [
(
"native_entity_dispatch",
"runtime_native_entity_dispatch",
"embedding service requires runtime native-entity dispatch (no runtime configured)",
),
(
"catalog_lookup",
"active_catalog",
"embedding service requires the active catalog (no catalog configured)",
),
] {
let err = embedding_capability_status(operation, capability, message);
assert_eq!(err.code(), tonic::Code::FailedPrecondition);
assert_eq!(err.message(), message);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Capability as i32);
assert_eq!(detail.backend, "embedding");
assert_eq!(detail.operation, operation);
assert_eq!(detail.capability_required, capability);
assert!(!detail.retryable);
}
}
#[test]
fn embedding_source_not_found_statuses_carry_schema_detail() {
for operation in ["backfill", "report_embedding", "retrieve"] {
assert_schema_not_found_detail(&embedding_source_not_found_status(operation), operation);
}
}
#[test]
fn embedding_point_is_tenant_tagged_no_fail_open() {
let err = build_embedding_point("row-1", vec![0.1, 0.2], " ", "orders")
.expect_err("empty tenant must fail closed");
assert_eq!(err.code(), tonic::Code::PermissionDenied);
assert_eq!(
err.message(),
"embedding upsert requires a verified tenant; refusing to store an unscoped vector \
(no fail-open)"
);
let detail = decode_detail(&err);
assert_eq!(detail.kind, ErrorKind::Policy as i32);
assert_eq!(detail.operation, "embedding_vector_upsert");
assert_eq!(detail.policy_decision_id, "verified_tenant_required");
assert!(!detail.retryable);
assert_eq!(detail.retry_after_ms, 0);
let point = build_embedding_point("row-1", vec![0.1, 0.2], "acme", "orders")
.expect("verified tenant ok");
assert_eq!(point.id, "row-1");
let payload = point.payload.expect("tenant-tag payload present");
assert!(
payload.fields.contains_key("_tenant_id"),
"stored vector is not tenant-tagged"
);
assert!(
payload.fields.contains_key("_source"),
"stored vector is not source-tagged (teardown filter-delete would miss it)"
);
let pjson = crate::runtime::executor_utils::struct_to_json(&payload);
assert_eq!(pjson["_parent_pk"], "row-1");
assert_eq!(pjson["_chunk_seq"].as_f64(), Some(0.0));
let chunked = build_embedding_point("row-1#chunk:4", vec![0.1, 0.2], "acme", "orders")
.expect("verified tenant ok");
let cjson =
crate::runtime::executor_utils::struct_to_json(&chunked.payload.expect("payload present"));
assert_eq!(
cjson["_parent_pk"], "row-1",
"chunk point must tag its PARENT row pk for row-scoped teardown"
);
assert_eq!(cjson["_chunk_seq"].as_f64(), Some(4.0));
let untagged =
build_embedding_point("row-2", vec![0.1, 0.2], "acme", " ").expect("verified tenant ok");
assert!(
!untagged
.payload
.expect("payload present")
.fields
.contains_key("_source"),
"empty source must not write an empty _source tag"
);
}
#[test]
fn source_teardown_filter_is_fully_scoped_or_none() {
use super::model::source_teardown_filter;
let filter = source_teardown_filter(" acme ", " orders ").expect("scoped filter");
let must = filter
.get("must")
.and_then(|m| m.as_array())
.expect("must array");
assert_eq!(must.len(), 2, "filter must AND tenant + source");
assert_eq!(must[0]["key"], "_tenant_id");
assert_eq!(must[0]["match"]["value"], "acme");
assert_eq!(must[1]["key"], "_source");
assert_eq!(must[1]["match"]["value"], "orders");
assert!(source_teardown_filter("", "orders").is_none());
assert!(source_teardown_filter("acme", " ").is_none());
assert!(source_teardown_filter(" ", " ").is_none());
}
#[test]
fn row_teardown_filter_scopes_to_a_single_parent() {
use super::model::row_teardown_filter;
let filter = row_teardown_filter("acme", "orders", "row-7").expect("scoped filter");
let must = filter.get("must").and_then(|m| m.as_array()).expect("must");
assert_eq!(must.len(), 3, "tenant + source + parent");
assert_eq!(must[2]["key"], "_parent_pk");
assert_eq!(must[2]["match"]["value"], "row-7");
assert!(row_teardown_filter("acme", "orders", " ").is_none());
assert!(row_teardown_filter("", "orders", "row-7").is_none());
}
#[test]
fn chunk_point_id_round_trips_and_is_backward_compatible() {
use super::chunking::{chunk_point_id, parse_chunk_point_id};
assert_eq!(chunk_point_id("row-1", 0, 1), "row-1");
assert_eq!(parse_chunk_point_id("row-1"), ("row-1", 0));
assert_eq!(chunk_point_id("row-1", 3, 5), "row-1#chunk:3");
assert_eq!(parse_chunk_point_id("row-1#chunk:3"), ("row-1", 3));
assert_eq!(parse_chunk_point_id("a#chunk:b"), ("a#chunk:b", 0));
assert_eq!(parse_chunk_point_id("a#chunk:x#chunk:2"), ("a#chunk:x", 2));
}
#[test]
fn chunk_text_splits_with_overlap_and_bounds() {
use super::chunking::chunk_text;
let one = chunk_text(" hello world ", 100, 10, 8);
assert_eq!(one.len(), 1);
assert_eq!(one[0].seq, 0);
assert_eq!(one[0].text, "hello world");
assert!(chunk_text(" ", 100, 10, 8).is_empty());
assert!(chunk_text("", 100, 10, 8).is_empty());
let text = "alpha beta gamma delta epsilon zeta eta theta iota kappa";
let chunks = chunk_text(text, 20, 5, 100);
assert!(chunks.len() >= 2, "long text must split");
for (i, c) in chunks.iter().enumerate() {
assert_eq!(c.seq as usize, i, "seqs are sequential from 0");
assert!(c.text.chars().count() <= 20, "chunk within size");
assert!(
!c.text.starts_with(' ') && !c.text.ends_with(' '),
"trimmed"
);
}
let joined: String = chunks
.iter()
.map(|c| c.text.clone())
.collect::<Vec<_>>()
.join(" ");
for word in text.split_whitespace() {
assert!(joined.contains(word), "word '{word}' lost across chunks");
}
let capped = chunk_text(text, 10, 2, 2);
assert!(capped.len() <= 2, "max_chunks caps the count");
let mb = chunk_text(
"café ☕ ambulance 🚑 dispatch café ☕ ambulance 🚑 dispatch",
8,
2,
100,
);
for c in &mb {
assert!(c.text.chars().count() <= 8);
}
}
#[test]
fn retrieve_filter_merges_under_authoritative_tenant_clause() {
let tenant_only =
serde_json::json!({ "must": [{ "key": "_tenant_id", "match": { "value": "t1" } }] });
assert_eq!(merge_retrieve_filter("t1", "").unwrap(), tenant_only);
assert_eq!(merge_retrieve_filter("t1", " ").unwrap(), tenant_only);
let merged = merge_retrieve_filter(
"t1",
r#"{"must":[{"key":"doc_type","match":{"value":"invoice"}}]}"#,
)
.unwrap();
let must = merged.get("must").unwrap().as_array().unwrap();
assert_eq!(must.len(), 2);
assert_eq!(must[0].get("key").unwrap(), "_tenant_id");
assert_eq!(must[1].get("key").unwrap(), "doc_type");
let merged = merge_retrieve_filter(
"t1",
r#"{"should":[{"key":"tag","match":{"value":"x"}}],"must_not":[{"key":"archived","match":{"value":true}}]}"#,
)
.unwrap();
assert!(merged.get("should").is_some());
assert!(merged.get("must_not").is_some());
assert_eq!(
merged.get("must").unwrap().as_array().unwrap()[0]
.get("key")
.unwrap(),
"_tenant_id"
);
assert!(
merge_retrieve_filter(
"t1",
r#"{"must":[{"key":"_tenant_id","match":{"value":"OTHER"}}]}"#
)
.is_err(),
"must reject an attempt to override the tenant clause"
);
assert!(
merge_retrieve_filter(
"t1",
r#"{"must":[{"key":"_project","match":{"value":"p"}}]}"#
)
.is_err()
);
assert!(
merge_retrieve_filter(
"t1",
r#"{"must":[{"should":[{"key":"_tenant_id","match":{"value":"x"}}]}]}"#
)
.is_err(),
"must reject a nested internal-key reference"
);
assert!(merge_retrieve_filter("t1", "not json").is_err());
assert!(merge_retrieve_filter("t1", "[1,2,3]").is_err());
assert!(merge_retrieve_filter("t1", r#"{"whatever":[]}"#).is_err());
assert!(merge_retrieve_filter("t1", r#"{"must":{}}"#).is_err());
}
#[test]
fn register_source_fails_closed_without_source_tenant_column() {
let err = require_source_tenant_column(None, "acme.crm.entity.v1.Contact")
.expect_err("missing tenant column must fail closed");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"source entity 'acme.crm.entity.v1.Contact' has no resolvable tenant column; refusing to register a tenant-scoped embedding source over it (fail closed)"
);
assert_single_field_violation(
&err,
"source_message_type",
"must resolve to a tenant-scoped source entity",
);
let blank = require_source_tenant_column(Some(" ".to_string()), "X")
.expect_err("blank column must fail closed");
assert_eq!(blank.code(), tonic::Code::InvalidArgument);
assert_single_field_violation(
&blank,
"source_message_type",
"must resolve to a tenant-scoped source entity",
);
assert_eq!(
require_source_tenant_column(Some("tenant_id".to_string()), "X").unwrap(),
"tenant_id"
);
}
#[test]
fn retrieve_past_deadline_returns_deadline_exceeded() {
let base = Instant::now();
let now = base + Duration::from_secs(1);
let err = remaining_before_deadline(Some(base), now)
.expect_err("past deadline must be deadline_exceeded");
assert_eq!(err.code(), tonic::Code::DeadlineExceeded);
let future = base + Duration::from_secs(30);
let remaining = remaining_before_deadline(Some(future), base)
.expect("future deadline ok")
.expect("some remaining budget");
assert!(remaining > Duration::from_secs(20));
assert!(
remaining_before_deadline(None, base)
.expect("no deadline ok")
.is_none()
);
}
#[test]
fn parses_grpc_timeout_units() {
assert_eq!(parse_grpc_timeout("5S"), Some(Duration::from_secs(5)));
assert_eq!(parse_grpc_timeout("250m"), Some(Duration::from_millis(250)));
assert_eq!(parse_grpc_timeout("1M"), Some(Duration::from_secs(60)));
assert_eq!(parse_grpc_timeout(""), None);
assert_eq!(parse_grpc_timeout("abc"), None);
}
#[test]
fn source_change_helpers_accept_cdc_envelope() {
let payload = serde_json::json!({
"tenant_id": "tenant-a",
"project_id": "project-a",
"operation": "upsert",
"document_id": "invoice-1",
"payload": {
"email": "ada@example.com",
"status": "OPEN"
}
});
assert_eq!(source_change_row_pk(&payload), "invoice-1");
assert!(!source_change_is_delete(&payload));
let row = source_change_row(&payload).expect("source row");
assert_eq!(
extract_source_text(row, &["email".to_string(), "status".to_string()]),
"ada@example.com OPEN"
);
let deleted = serde_json::json!({
"tenant_id": "tenant-a",
"operation": "delete",
"document_id": "invoice-1",
"payload": { "email": "ada@example.com" }
});
assert!(source_change_is_delete(&deleted));
}
#[test]
fn source_text_fields_parse_jsonb_and_string_wrapped_json() {
assert_eq!(
parse_source_text_fields(r#"["email","status"]"#),
vec!["email".to_string(), "status".to_string()]
);
assert_eq!(
parse_source_text_fields(r#""[\"email\",\"status\"]""#),
vec!["email".to_string(), "status".to_string()]
);
assert!(parse_source_text_fields("not-json").is_empty());
}
#[test]
fn backfill_select_request_is_tenant_scoped_and_cursor_bounded() {
let source = StoredSource {
source_id: "src-1".to_string(),
tenant_id: "tenant-a".to_string(),
source_name: "contacts".to_string(),
source_message_type: "acme.crm.v1.Contact".to_string(),
text_fields_json: r#"["email","status"]"#.to_string(),
target_collection: "contacts_vec".to_string(),
model_id: "model-a".to_string(),
tenant_column: "tenant_id".to_string(),
source_cdc_topic: "acme.contacts.cdc".to_string(),
status: STATUS_ACTIVE.to_string(),
};
let request = backfill_select_request(
&source,
"contact_id",
&["email".to_string(), "status".to_string()],
Some("contact-7"),
None,
"",
);
assert_eq!(request.message_type, "acme.crm.v1.Contact");
assert_eq!(request.fields, ["contact_id", "email", "status"]);
assert_eq!(request.limit, EMBEDDING_BACKFILL_PAGE_LIMIT);
assert_eq!(request.sort[0].field, "contact_id");
let filter = request.filter.as_ref().expect("tenant filter");
let json = crate::runtime::executor_utils::struct_to_json(filter);
assert_eq!(json["$and"][0]["tenant_id"], "tenant-a");
assert_eq!(json["$and"][1]["contact_id"]["$gt"], "contact-7");
let scoped = backfill_select_request(
&source,
"contact_id",
&["email".to_string()],
None,
Some("project_id"),
"proj-9",
);
let scoped_json =
crate::runtime::executor_utils::struct_to_json(scoped.filter.as_ref().unwrap());
assert_eq!(scoped_json["$and"][0]["tenant_id"], "tenant-a");
assert_eq!(scoped_json["$and"][1]["project_id"], "proj-9");
}
#[test]
fn backfill_read_context_uses_served_read_gate() {
let context = backfill_read_context("tenant-a", "project-a");
assert_eq!(context.tenant_id, "tenant-a");
assert_eq!(context.project_id, "project-a");
assert_eq!(context.purpose, "embedding_backfill");
assert_eq!(context.scopes, vec!["udb:read".to_string()]);
assert_eq!(context.service_identity, "udb.embedding.backfill");
}
#[test]
fn reported_vector_shape_gate_rejects_empty_and_mismatched_dims() {
let err = validate_reported_vector(0, &[]).expect_err("empty vector must fail");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(err.message(), "vector is required");
assert_single_field_violation(
&err,
"vector",
"must contain at least one embedding dimension",
);
let err = validate_reported_vector(5, &[0.1, 0.2, 0.3]).expect_err("dims mismatch must fail");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_eq!(
err.message(),
"dims (5) does not match the reported vector length (3)"
);
assert_single_field_violation(
&err,
"dims",
"must equal the reported vector length when set",
);
validate_reported_vector(0, &[0.1]).expect("unset dims must pass");
validate_reported_vector(-1, &[0.1]).expect("negative dims treated as unset");
validate_reported_vector(3, &[0.1, 0.2, 0.3]).expect("matching dims must pass");
}
#[tokio::test]
async fn report_embedding_dims_mismatch_carries_field_violation() {
let svc = EmbeddingServiceImpl::new(); let mut request = Request::new(embedding_pb::ReportEmbeddingRequest {
tenant_id: "tenant-a".to_string(),
source_name: "contacts".to_string(),
row_pk: "row-1".to_string(),
vector: vec![0.1, 0.2, 0.3],
model: "m".to_string(),
dims: 5,
});
request
.metadata_mut()
.insert("x-tenant-id", MetadataValue::from_static("tenant-a"));
let err = svc
.report_embedding(request)
.await
.expect_err("dims mismatch must be rejected");
assert_eq!(err.code(), tonic::Code::InvalidArgument);
assert_single_field_violation(
&err,
"dims",
"must equal the reported vector length when set",
);
}
#[test]
fn work_loader_sql_reads_nested_journal_envelope() {
let sql = embedding_work_jobs_sql("udb_cdc_event_journal", "udb_outbox");
for fragment in [
"COALESCE(j.payload->>'tenant_id', j.payload->'payload'->>'tenant_id', '')",
"COALESCE(j.payload->>'project_id', j.payload->'payload'->>'project_id', '')",
"COALESCE(o.payload->>'source_event_id', o.payload->'payload'->>'source_event_id')",
"COALESCE(done.payload->>'source_event_id', done.payload->'payload'->>'source_event_id')",
] {
assert!(
sql.contains(fragment),
"work loader SQL lost the nested-envelope fallback: missing {fragment}"
);
}
assert!(
!sql.contains("COALESCE(j.payload->>'tenant_id', '')"),
"work loader SQL still reads tenant_id top-level only"
);
}
#[test]
fn teardown_sqls_read_nested_journal_envelope() {
let jobs_sql = embedding_teardown_jobs_sql("udb_cdc_event_journal", "udb_outbox");
for fragment in [
"COALESCE(j.payload->>'tenant_id', j.payload->'payload'->>'tenant_id', '')",
"COALESCE(j.payload->>'source', j.payload->'payload'->>'source', '')",
"COALESCE(j.payload->>'target_collection', j.payload->'payload'->>'target_collection', '')",
"COALESCE(o.payload->>'teardown_event_id', o.payload->'payload'->>'teardown_event_id')",
"COALESCE(done.payload->>'teardown_event_id', done.payload->'payload'->>'teardown_event_id')",
] {
assert!(
jobs_sql.contains(fragment),
"teardown loader SQL missing nested-envelope fallback: {fragment}"
);
}
let ids_sql = embedding_teardown_point_ids_sql("udb_cdc_event_journal", "udb_outbox");
for fragment in [
"COALESCE(j.payload->>'row_pk', j.payload->'payload'->>'row_pk', '')",
"COALESCE(o.payload->>'row_pk', o.payload->'payload'->>'row_pk', '')",
"COALESCE(j.payload->>'tenant_id', j.payload->'payload'->>'tenant_id', '')",
"COALESCE(o.payload->>'source', o.payload->'payload'->>'source', '')",
] {
assert!(
ids_sql.contains(fragment),
"teardown point-id SQL missing nested-envelope fallback: {fragment}"
);
}
assert!(ids_sql.contains("LIMIT $6"));
assert!(ids_sql.contains("(ids.target_collection, ids.row_pk) > ($4, $5)"));
}
#[test]
fn retrieve_hit_payload_strips_internal_tenant_tag() {
let payload = crate::runtime::executor_utils::json_to_struct(&serde_json::json!({
"_tenant_id": "acme",
"title": "Q3 report",
"chunk": 4,
}))
.expect("struct payload");
let json = retrieve_hit_payload_json(Some(&payload));
let value: serde_json::Value = serde_json::from_str(&json).expect("payload_json is valid JSON");
assert_eq!(value["title"], "Q3 report");
assert_eq!(value["chunk"], 4.0);
assert!(
value.get("_tenant_id").is_none(),
"internal tenant tag leaked into RetrieveHit payload_json"
);
let tag_only = crate::runtime::executor_utils::json_to_struct(&serde_json::json!({
"_tenant_id": "acme",
}))
.expect("struct payload");
assert_eq!(retrieve_hit_payload_json(Some(&tag_only)), "");
assert_eq!(retrieve_hit_payload_json(None), "");
}