use serde_json::Value;
fn fnv1a_feed(h: &mut u64, bytes: &[u8]) {
for &b in bytes {
*h ^= b as u64;
*h = h.wrapping_mul(0x100000001b3);
}
}
const FNV1A_SEED: u64 = 0xcbf29ce484222325;
fn fnv1a64(bytes: &[u8]) -> u64 {
let mut h = FNV1A_SEED;
fnv1a_feed(&mut h, bytes);
h
}
fn first_forwarded_value<'a>(mut get: impl FnMut(&str) -> Option<&'a str>) -> Option<&'a str> {
if let Some(xff) = get("x-forwarded-for")
&& let Some(first) = xff.split(',').next().map(str::trim)
&& !first.is_empty()
{
return Some(first);
}
get("x-real-ip").map(str::trim).filter(|v| !v.is_empty())
}
pub fn rollout_identity<'a>(metadata: &'a Value, sticky_header: &str) -> Option<&'a str> {
let headers = metadata.get("headers")?.as_object()?;
if !sticky_header.is_empty()
&& let Some((_, v)) = headers
.iter()
.find(|(k, _)| k.eq_ignore_ascii_case(sticky_header))
&& let Some(v) = v.as_str()
&& !v.is_empty()
{
return Some(v);
}
first_forwarded_value(|name| headers.get(name).and_then(|v| v.as_str()))
}
pub fn serialize_task_trace_capped(
trace: Option<&dataflow_rs::ExecutionTrace>,
max_bytes: usize,
context: &str,
) -> Option<String> {
let json = serde_json::to_string(trace?).ok()?;
if max_bytes > 0 && json.len() > max_bytes {
crate::metrics::record_error("task_trace_size_exceeded");
tracing::warn!(
context = %context,
task_trace_bytes = json.len(),
limit_bytes = max_bytes,
"task_trace_json exceeds queue.max_result_size_bytes; dropping task detail"
);
return None;
}
Some(json)
}
pub fn rollout_bucket_for_identity(identity: Option<&str>) -> u8 {
match identity {
Some(id) if !id.is_empty() => (fnv1a64(id.as_bytes()) % 100) as u8,
_ => (rand::random::<u32>() % 100) as u8,
}
}
pub const CREDENTIAL_HEADERS: [&str; 4] = [
"authorization",
"cookie",
"proxy-authorization",
"x-api-key",
];
pub fn is_credential_header(name: &str) -> bool {
CREDENTIAL_HEADERS.contains(&name)
}
const STRING_MAP_KEYS: [&str; 4] = ["headers", "params", "query", "cookies"];
pub fn prepare_offline_metadata(metadata: Value) -> Result<Value, String> {
let mut metadata = match metadata {
Value::Null => return Ok(serde_json::json!({})),
Value::Object(_) => metadata,
other => {
return Err(format!(
"'metadata' must be a JSON object, got {}",
json_kind(&other)
));
}
};
for key in STRING_MAP_KEYS {
let Some(value) = metadata.get(key) else {
continue;
};
let Some(map) = value.as_object() else {
return Err(format!(
"'metadata.{key}' must be an object of strings, got {}",
json_kind(value)
));
};
if let Some((name, bad)) = map.iter().find(|(_, v)| !v.is_string()) {
return Err(format!(
"'metadata.{key}.{name}' must be a string, got {}",
json_kind(bad)
));
}
}
for key in ["channel", "http_method"] {
if let Some(value) = metadata.get(key)
&& !value.is_string()
{
return Err(format!(
"'metadata.{key}' must be a string, got {}",
json_kind(value)
));
}
}
if let Some(auth) = metadata.get("auth") {
let Some(map) = auth.as_object() else {
return Err(format!(
"'metadata.auth' must be an object, got {}",
json_kind(auth)
));
};
if let Some(unknown) = map.keys().find(|k| k.as_str() != "claims") {
return Err(format!(
"'metadata.auth.{unknown}' is not settable — the request path builds \
'auth' as {{\"claims\": …}} and nothing else reaches a workflow"
));
}
}
crate::engine::clear_error_context(&mut metadata);
if let Some(headers) = metadata.get("headers").and_then(Value::as_object) {
let normalized: serde_json::Map<String, Value> = headers
.iter()
.map(|(name, value)| {
let name = name.to_ascii_lowercase();
let value = if is_credential_header(&name) {
Value::String(crate::connector::MASK.to_string())
} else {
value.clone()
};
(name, value)
})
.collect();
metadata["headers"] = Value::Object(normalized);
}
Ok(metadata)
}
pub fn json_kind(value: &Value) -> &'static str {
match value {
Value::Null => "null",
Value::Bool(_) => "a boolean",
Value::Number(_) => "a number",
Value::String(_) => "a string",
Value::Array(_) => "an array",
Value::Object(_) => "an object",
}
}
#[cfg(test)]
mod tests {
use super::*;
use dataflow_rs::Message;
use serde_json::json;
fn ingress_message(payload: Value, metadata: Value, identity: Option<&str>) -> Message {
Message::builder()
.payload_json(&payload)
.metadata_json(&metadata)
.routing_bucket(rollout_bucket_for_identity(identity))
.build()
}
#[test]
fn test_task_trace_cap_drops_oversized_detail() {
let trace = dataflow_rs::ExecutionTrace::new();
assert!(serialize_task_trace_capped(Some(&trace), 1024, "t").is_some());
assert!(serialize_task_trace_capped(Some(&trace), 1, "t").is_none());
assert!(serialize_task_trace_capped(Some(&trace), 0, "t").is_some());
assert!(serialize_task_trace_capped(None, 1024, "t").is_none());
}
#[test]
fn ingress_seeds_metadata_with_literal_keys() {
let msg = ingress_message(json!({}), json!({"source": "test", "a.b": 2}), None);
assert_eq!(
msg.metadata().get("source").and_then(|v| v.as_str()),
Some("test")
);
assert_eq!(msg.metadata().get("a.b").and_then(|v| v.as_i64()), Some(2));
assert!(
msg.metadata().get("a").is_none(),
"a dotted metadata key must not be re-read as a path"
);
}
#[test]
fn test_rollout_bucket_is_sticky_per_identity() {
let a1 = rollout_bucket_for_identity(Some("10.0.0.7"));
let a2 = rollout_bucket_for_identity(Some("10.0.0.7"));
assert_eq!(a1, a2, "same identity must map to the same bucket");
assert!(a1 < 100);
let buckets: std::collections::HashSet<u8> = (0..50)
.map(|i| rollout_bucket_for_identity(Some(&format!("user-{i}"))))
.collect();
assert!(buckets.len() > 10, "expected spread, got {buckets:?}");
}
#[test]
fn test_rollout_identity_prefers_sticky_header() {
let metadata = json!({
"method": "POST",
"headers": {
"x-user-id": "user-7",
"x-forwarded-for": "10.1.1.1, 10.1.1.2"
}
});
assert_eq!(rollout_identity(&metadata, "X-User-Id"), Some("user-7"));
assert_eq!(rollout_identity(&metadata, ""), Some("10.1.1.1"));
let metadata = json!({"headers": {"x-real-ip": "10.2.2.2"}});
assert_eq!(rollout_identity(&metadata, "x-user-id"), Some("10.2.2.2"));
assert_eq!(rollout_identity(&json!({}), "x-user-id"), None);
}
#[test]
fn test_rollout_bucket_empty_identity_falls_back_to_random() {
let buckets: std::collections::HashSet<u8> = (0..100)
.map(|_| rollout_bucket_for_identity(Some("")))
.collect();
assert!(buckets.len() > 1, "empty identity should randomize");
}
#[test]
fn the_routing_bucket_is_not_part_of_the_message_body() {
let msg = ingress_message(json!({"order_id": 7}), json!({}), Some("caller-1"));
assert!(
msg.routing_bucket().is_some(),
"precondition: the ingress set a bucket"
);
let body: Value = msg.data().into();
assert_eq!(body, json!({}), "routing must not write into `data`");
let serialized = serde_json::to_string(&msg).expect("message serializes");
assert!(
!serialized.contains("_rollout_bucket") && !serialized.contains("routing_bucket"),
"the bucket must not reach the persisted message: {serialized}"
);
}
}
#[cfg(test)]
mod offline_metadata_tests {
use super::*;
use serde_json::json;
#[test]
fn case_metadata_is_normalized_the_way_the_ingress_builds_it() {
let out = prepare_offline_metadata(json!({
"headers": { "DeviceId": "device-abc", "Authorization": "Bearer secret" },
"auth": { "claims": { "sub": "asha@example.com" } },
"custom": { "anything": true },
}))
.expect("a well-formed metadata object is accepted");
assert_eq!(
out["headers"]["deviceid"], "device-abc",
"header keys lowercase, because axum yields lowercase names"
);
assert!(
out["headers"].get("DeviceId").is_none(),
"the original casing must not survive alongside it"
);
assert_eq!(
out["headers"]["authorization"],
crate::connector::MASK,
"a credential header is masked here exactly as at ingress"
);
assert_eq!(out["auth"]["claims"]["sub"], "asha@example.com");
assert_eq!(
out["custom"]["anything"], true,
"keys outside the reserved set pass through: the envelope merges \
arbitrary caller metadata, so a closed set would be wrong"
);
}
#[test]
fn the_engine_owned_error_context_is_cleared() {
let out = prepare_offline_metadata(json!({"_orion_errors": [{"code": "FAKE"}]}))
.expect("accepted");
assert!(out.get("_orion_errors").is_none());
}
#[test]
fn shapes_the_ingress_cannot_produce_are_refused() {
let err = prepare_offline_metadata(json!({"headers": ["a", "b"]}))
.expect_err("headers must be an object");
assert!(err.contains("metadata.headers"), "{err}");
let err = prepare_offline_metadata(json!({"query": {"page": 2}}))
.expect_err("query values must be strings");
assert!(err.contains("metadata.query.page"), "{err}");
let err =
prepare_offline_metadata(json!({"channel": 7})).expect_err("channel must be a string");
assert!(err.contains("metadata.channel"), "{err}");
let err = prepare_offline_metadata(json!({"auth": {"token": "t"}}))
.expect_err("only auth.claims is reachable");
assert!(err.contains("metadata.auth.token"), "{err}");
let err = prepare_offline_metadata(json!("nope")).expect_err("root must be an object");
assert!(err.contains("must be a JSON object"), "{err}");
}
#[test]
fn the_credential_header_set_is_the_ingress_set() {
for name in CREDENTIAL_HEADERS {
assert!(is_credential_header(name), "{name} must be masked");
}
assert!(!is_credential_header("deviceid"));
assert!(!is_credential_header("content-type"));
}
}