#![cfg(test)]
use serde_json::json;
use crate::generation::manifest::{
CatalogManifest, ManifestProjection, ManifestStoreOption, ManifestTable,
};
use crate::metrics::{MetricsRecorder, NoopMetrics};
use crate::runtime::projection::{ProjectionEngine, ProjectionPlan};
fn opt(key: &str, value: &str) -> ManifestStoreOption {
ManifestStoreOption {
key: key.to_string(),
value: value.to_string(),
}
}
fn projection_for(
message: &str,
kind: &str,
backend: &str,
resource: &str,
options: Vec<ManifestStoreOption>,
) -> ManifestProjection {
ManifestProjection {
message_type: message.to_string(),
projection_kind: kind.to_string(),
backend: backend.to_string(),
instance: "default".to_string(),
resource_name: resource.to_string(),
write_policy: "primary".to_string(),
fanout_policy: "outbox".to_string(),
options,
..ManifestProjection::default()
}
}
fn full_multi_backend_manifest() -> CatalogManifest {
let table = ManifestTable {
message_name: "acme.billing.v1.Customer".to_string(),
schema: "public".to_string(),
table: "customers".to_string(),
primary_key: vec!["id".to_string()],
projections: vec![
projection_for(
"acme.billing.v1.Customer",
"document",
"mongodb",
"customers",
vec![],
),
projection_for(
"acme.billing.v1.Customer",
"vector",
"qdrant",
"customers_vectors",
vec![opt("vector_size", "384")],
),
projection_for(
"acme.billing.v1.Customer",
"graph",
"neo4j",
"Customer",
vec![],
),
projection_for(
"acme.billing.v1.Customer",
"analytics",
"clickhouse",
"customers_events",
vec![],
),
projection_for(
"acme.billing.v1.Customer",
"cache",
"redis",
"udb:{tenant}:Customer:{id}",
vec![opt("ttl_seconds", "300")],
),
projection_for(
"acme.billing.v1.Customer",
"object",
"s3",
"customer-snapshots",
vec![opt("key_prefix", "v1/customers")],
),
],
..Default::default()
};
CatalogManifest {
checksum_sha256: "full-multi-backend-v1".to_string(),
tables: vec![table],
..Default::default()
}
}
#[test]
fn full_manifest_produces_one_target_per_backend() {
let plans = ProjectionPlan::from_manifest(&full_multi_backend_manifest());
assert_eq!(plans.len(), 1, "one source message type");
let plan = &plans[0];
let backends: Vec<&str> = plan.targets.iter().map(|t| t.backend.as_str()).collect();
assert_eq!(
backends,
vec!["mongodb", "qdrant", "neo4j", "clickhouse", "redis", "s3"],
"every U3-named backend must appear as a target"
);
let by_backend: std::collections::BTreeMap<&str, &str> = plan
.targets
.iter()
.map(|t| (t.backend.as_str(), t.resource_name.as_str()))
.collect();
assert_eq!(by_backend["mongodb"], "customers");
assert_eq!(by_backend["qdrant"], "customers_vectors");
assert_eq!(by_backend["clickhouse"], "customers_events");
}
#[test]
fn idempotency_key_is_stable_across_replays() {
let row = json!({ "id": "cust-1" });
let key_1 = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v1",
"source-checksum-1",
);
let key_2 = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v1",
"source-checksum-1",
);
assert_eq!(
key_1, key_2,
"same logical write must produce the same idempotency key"
);
}
#[test]
fn idempotency_key_changes_at_each_correctness_boundary() {
let row = json!({ "id": "cust-1" });
let base = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v1",
"source-checksum-1",
);
let changed_source = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v1",
"source-checksum-2",
);
assert_ne!(
base, changed_source,
"source payload change must invalidate idempotency key"
);
let migrated = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v2",
"source-checksum-1",
);
assert_ne!(
base, migrated,
"manifest migration must invalidate idempotency key"
);
let other_backend = ProjectionEngine::idempotency_key(
"billing",
"customers",
&row,
"upsert",
"qdrant",
"default",
"manifest-v1",
"source-checksum-1",
);
assert_ne!(
base, other_backend,
"different target backend must produce a different task"
);
let other_project = ProjectionEngine::idempotency_key(
"analytics",
"customers",
&row,
"upsert",
"mongodb",
"default",
"manifest-v1",
"source-checksum-1",
);
assert_ne!(
base, other_project,
"different project must produce a different task"
);
}
#[test]
fn extract_row_key_handles_compound_keys() {
let table = ManifestTable {
message_name: "acme.billing.v1.OrderItem".to_string(),
schema: "billing".to_string(),
table: "order_items".to_string(),
primary_key: vec!["order_id".to_string(), "item_id".to_string()],
projections: vec![projection_for(
"acme.billing.v1.OrderItem",
"analytics",
"clickhouse",
"order_items_facts",
vec![],
)],
..Default::default()
};
let manifest = CatalogManifest {
checksum_sha256: "compound-pk-v1".into(),
tables: vec![table],
..Default::default()
};
let plan = &ProjectionPlan::from_manifest(&manifest)[0];
let payload = json!({
"order_id": "ord-1",
"item_id": "sku-9",
"quantity": 3,
"price_cents": 4_999,
});
let key = plan.extract_row_key(&payload);
let expected = json!({ "order_id": "ord-1", "item_id": "sku-9" });
assert_eq!(
key, expected,
"compound key extraction must drop non-PK fields"
);
}
#[test]
fn source_checksum_is_payload_sensitive() {
let a = json!({ "id": "cust-1", "email": "alice@acme.com" });
let b = json!({ "id": "cust-1", "email": "alice@acme.com" });
let c = json!({ "id": "cust-1", "email": "alice+changed@acme.com" });
let nul = json!({ "id": "cust-1", "email": "alice@acme.com", "note": "ok\u{0}then" });
let stripped = json!({ "id": "cust-1", "email": "alice@acme.com", "note": "okthen" });
assert_eq!(
ProjectionEngine::source_checksum(&a),
ProjectionEngine::source_checksum(&b),
"identical payloads → identical checksums"
);
assert_ne!(
ProjectionEngine::source_checksum(&a),
ProjectionEngine::source_checksum(&c),
"changed payload → new checksum"
);
assert_eq!(
ProjectionEngine::source_checksum(&nul),
ProjectionEngine::source_checksum(&stripped),
"NUL bytes are stripped before projection JSONB persistence"
);
}
#[test]
fn oldest_pending_age_metric_is_reachable_with_full_label_set() {
let metrics: std::sync::Arc<dyn MetricsRecorder> = std::sync::Arc::new(NoopMetrics);
metrics.set_projection_oldest_pending_age_seconds(
"billing", "mongodb", "default", "document", 42.5,
);
metrics.set_projection_oldest_pending_age_seconds("billing", "redis", "default", "cache", 1.0);
}