use serde_json::Value;
use std::collections::HashMap;
use crate::storage::models::ConnectorResponse;
pub(crate) const MASK: &str = "******";
const SECRET_KEY_SUBSTRINGS: &[&str] = &[
"password",
"passwd",
"passphrase",
"secret",
"token",
"credential",
"api_key",
"apikey",
"private_key",
"privatekey",
"access_key",
"accesskey",
"signature",
"bearer",
"dsn",
"webhook",
];
const SECRET_KEY_EXACT: &[&str] = &[
"key",
"pwd",
"cert",
"certificate",
"ssl_key",
"sasl_key",
"keyfile",
"keytab",
"authorization",
"connection_string",
"pat",
"sig",
];
fn is_secret_key(key: &str) -> bool {
let lower = key.to_ascii_lowercase();
SECRET_KEY_EXACT.contains(&lower.as_str())
|| SECRET_KEY_SUBSTRINGS.iter().any(|p| lower.contains(p))
}
const READABLE_KEYS: &[&str] = &[
"type",
"driver",
"backend",
"method",
"username",
"header",
"sasl_mechanism",
"sasl_username",
"url",
"brokers",
"topic",
"max_connections",
"connect_timeout_ms",
"query_timeout_ms",
"request_timeout_ms",
"max_response_size",
"default_ttl_secs",
"max_retries",
"retry_delay_ms",
"retry_non_idempotent",
"allow_private_urls",
"read",
"insert",
"update",
"delete",
"upsert",
"raw_write",
"write",
"publish",
"methods",
"require_schema",
"allowed_entities",
];
fn is_readable_key(key: &str) -> bool {
READABLE_KEYS.contains(&key.to_ascii_lowercase().as_str())
}
fn mask_connector_in_place(value: &mut Value, key: Option<&str>) {
match value {
Value::Object(map) => {
for (k, v) in map.iter_mut() {
mask_connector_in_place(v, Some(k.as_str()));
}
}
Value::Array(items) => {
for item in items.iter_mut() {
mask_connector_in_place(item, key);
}
}
Value::Null => {}
Value::String(s) => {
if crate::connector::secrets::is_resolvable_reference(s) {
return;
}
if key.is_some_and(is_readable_key) {
if let Some(redacted) = redact_url_secrets(s) {
*s = redacted;
}
} else {
*s = MASK.to_string();
}
}
other => {
if !key.is_some_and(is_readable_key) {
*other = Value::String(MASK.to_string());
}
}
}
}
pub fn mask_channel_config(config: &mut Value) {
let Some(auth) = config.get_mut("auth") else {
return;
};
if let Some(keys) = auth.get_mut("keys")
&& let Value::Array(items) = keys
{
for item in items.iter_mut() {
mask_channel_secret_leaf(item);
}
}
if let Some(secret) = auth.get_mut("secret") {
mask_channel_secret_leaf(secret);
}
}
fn mask_channel_secret_leaf(value: &mut Value) {
match value {
Value::String(s) if crate::connector::secrets::is_resolvable_reference(s) => {}
Value::Null => {}
other => *other = Value::String(MASK.to_string()),
}
}
pub fn unmask_channel_config(incoming: &mut Value, stored: &Value) {
let mut masked_stored = stored.clone();
mask_channel_config(&mut masked_stored);
restore_in_place(incoming, stored, &masked_stored);
}
struct SplitUrl<'a> {
scheme: &'a str,
user: Option<&'a str>,
password: Option<&'a str>,
host_path: &'a str,
query: Option<Vec<QueryPair<'a>>>,
fragment: Option<&'a str>,
}
struct QueryPair<'a> {
name: &'a str,
value: Option<&'a str>,
}
fn split_url(s: &str) -> Option<SplitUrl<'_>> {
let (scheme, rest) = s.split_once("://")?;
if scheme.is_empty()
|| !scheme
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '+' | '-' | '.' | '_'))
{
return None;
}
let authority_end = rest.find(['/', '?', '#']).unwrap_or(rest.len());
let (user, password, after_userinfo) = match rest[..authority_end].rfind('@') {
Some(at) => {
let (user, password) = match rest[..at].split_once(':') {
Some((user, password)) => (user, Some(password)),
None => (&rest[..at], None),
};
(Some(user), password, &rest[at + 1..])
}
None => (None, None, rest),
};
let (main, fragment) = match after_userinfo.split_once('#') {
Some((main, fragment)) => (main, Some(fragment)),
None => (after_userinfo, None),
};
let (host_path, query) = match main.split_once('?') {
Some((host_path, query)) => {
let pairs = query
.split('&')
.map(|pair| match pair.split_once('=') {
Some((name, value)) => QueryPair {
name,
value: Some(value),
},
None => QueryPair {
name: pair,
value: None,
},
})
.collect();
(host_path, Some(pairs))
}
None => (main, None),
};
Some(SplitUrl {
scheme,
user,
password,
host_path,
query,
fragment,
})
}
fn assemble_url(url: &SplitUrl<'_>, password: Option<&str>, query: Option<&[String]>) -> String {
let mut out = format!("{}://", url.scheme);
if let Some(user) = url.user {
out.push_str(user);
if let Some(password) = password {
out.push(':');
out.push_str(password);
}
out.push('@');
}
out.push_str(url.host_path);
if let Some(pairs) = query {
out.push('?');
out.push_str(&pairs.join("&"));
}
if let Some(fragment) = url.fragment {
out.push('#');
out.push_str(fragment);
}
out
}
fn render_pair(name: &str, value: Option<&str>) -> String {
match value {
Some(value) => format!("{name}={value}"),
None => name.to_string(),
}
}
pub fn redact_url_secrets(s: &str) -> Option<String> {
let url = split_url(s)?;
let mut changed = false;
let password = url.password.map(|_| {
changed = true;
MASK
});
let query: Option<Vec<String>> = url.query.as_ref().map(|pairs| {
pairs
.iter()
.map(|pair| match pair.value {
Some(_) if is_secret_key(pair.name) => {
changed = true;
render_pair(pair.name, Some(MASK))
}
value => render_pair(pair.name, value),
})
.collect()
});
changed.then(|| assemble_url(&url, password, query.as_deref()))
}
pub fn redact_url_secrets_or_raw(s: &str) -> String {
redact_url_secrets(s).unwrap_or_else(|| s.to_string())
}
fn restore_url_secrets(incoming: &str, stored: &str) -> Option<String> {
let inc = split_url(incoming)?;
let st = split_url(stored)?;
let mut restored = false;
let password = match inc.password {
Some(p) if p == MASK => match st.password {
Some(stored_password) => {
restored = true;
Some(stored_password)
}
None => Some(p),
},
other => other,
};
let stored_pairs = st.query.unwrap_or_default();
let mut occurrences: HashMap<&str, usize> = HashMap::new();
let query: Option<Vec<String>> = inc.query.as_ref().map(|pairs| {
pairs
.iter()
.map(|pair| {
let occurrence = occurrences.entry(pair.name).or_insert(0);
let counterpart = stored_pairs
.iter()
.filter(|sp| sp.name == pair.name)
.nth(*occurrence)
.and_then(|sp| sp.value);
*occurrence += 1;
match pair.value {
Some(v) if v == MASK && is_secret_key(pair.name) && counterpart.is_some() => {
restored = true;
render_pair(pair.name, counterpart)
}
value => render_pair(pair.name, value),
}
})
.collect()
});
restored.then(|| assemble_url(&inc, password, query.as_deref()))
}
fn url_carries_mask(s: &str) -> bool {
let Some(url) = split_url(s) else {
return false;
};
url.password == Some(MASK)
|| url
.query
.is_some_and(|pairs| pairs.iter().any(|pair| pair.value == Some(MASK)))
}
fn mask_in_place(value: &mut Value, key: Option<&str>, force: bool) {
let secret = force || key.is_some_and(is_secret_key);
match value {
Value::Object(map) => {
for (k, v) in map.iter_mut() {
mask_in_place(v, Some(k.as_str()), secret);
}
}
Value::Array(items) => {
for item in items.iter_mut() {
mask_in_place(item, key, secret);
}
}
Value::Null => {}
Value::String(s) => {
if secret && !crate::connector::secrets::is_resolvable_reference(s) {
*s = MASK.to_string();
} else if !secret && let Some(redacted) = redact_url_secrets(s) {
*s = redacted;
}
}
other => {
if secret {
*other = Value::String(MASK.to_string());
}
}
}
}
pub fn mask_secrets(value: &mut Value) {
mask_in_place(value, None, false);
}
pub fn mask_connector_secrets(config_json: &str) -> String {
let Ok(mut val) = serde_json::from_str::<Value>(config_json) else {
return config_json.to_string();
};
mask_connector_in_place(&mut val, None);
serde_json::to_string(&val).unwrap_or_else(|_| config_json.to_string())
}
pub fn mask_connector(connector: &crate::storage::models::Connector) -> ConnectorResponse {
let mut masked = ConnectorResponse::from(connector);
masked.config_json = mask_connector_secrets(&masked.config_json);
masked.config = serde_json::from_str(&masked.config_json).unwrap_or(Value::Null);
masked
}
pub fn unmask_config(incoming: &mut Value, stored: &Value) {
let mut masked_stored = stored.clone();
mask_connector_in_place(&mut masked_stored, None);
restore_in_place(incoming, stored, &masked_stored);
}
fn restore_in_place(incoming: &mut Value, stored: &Value, masked: &Value) {
match (incoming, stored, masked) {
(Value::Object(inc), Value::Object(st), Value::Object(mk)) => {
for (key, value) in inc.iter_mut() {
if let (Some(stored_value), Some(masked_value)) = (st.get(key), mk.get(key)) {
restore_in_place(value, stored_value, masked_value);
}
}
}
(Value::Array(inc), Value::Array(st), Value::Array(mk)) => {
for ((value, stored_value), masked_value) in inc.iter_mut().zip(st).zip(mk) {
restore_in_place(value, stored_value, masked_value);
}
}
(inc, stored_value, masked_value) => {
if inc == masked_value && inc != stored_value {
*inc = stored_value.clone();
} else if let (Value::String(inc_s), Value::String(stored_s)) = (inc, stored_value)
&& let Some(restored) = restore_url_secrets(inc_s, stored_s)
{
*inc_s = restored;
}
}
}
}
pub fn find_masked_value(value: &Value) -> Option<String> {
fn walk(value: &Value, path: &str) -> Option<String> {
match value {
Value::Object(map) => map.iter().find_map(|(key, child)| {
let child_path = if path.is_empty() {
key.clone()
} else {
format!("{path}.{key}")
};
walk(child, &child_path)
}),
Value::Array(items) => items
.iter()
.enumerate()
.find_map(|(i, item)| walk(item, &format!("{path}[{i}]"))),
Value::String(s) if s == MASK || url_carries_mask(s) => Some(path.to_string()),
_ => None,
}
}
walk(value, "")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn mask_connector_populates_the_parsed_config_from_the_masked_string() {
let row = crate::storage::models::Connector {
id: "c-1".to_string(),
name: "api".to_string(),
connector_type: "http".to_string(),
config_json:
r#"{"type":"http","url":"https://api.example.com","auth":{"type":"bearer","token":"secret123"}}"#
.to_string(),
enabled: true,
tags_json: "[]".to_string(),
created_at: Default::default(),
updated_at: Default::default(),
};
let masked = mask_connector(&row);
assert_eq!(masked.config["auth"]["token"], "******");
assert!(
!masked.config.to_string().contains("secret123"),
"the parsed config leaked a secret the string form masked: {}",
masked.config
);
let from_string: Value = serde_json::from_str(&masked.config_json).expect("test");
assert_eq!(masked.config, from_string);
}
#[test]
fn an_unparseable_config_yields_a_null_parsed_config() {
let row = crate::storage::models::Connector {
id: "c-2".to_string(),
name: "broken".to_string(),
connector_type: "http".to_string(),
config_json: "{not json".to_string(),
enabled: true,
tags_json: "[]".to_string(),
created_at: Default::default(),
updated_at: Default::default(),
};
assert_eq!(mask_connector(&row).config, Value::Null);
}
#[test]
fn test_mask_connector_secrets_bearer_token() {
let config = r#"{"type":"http","url":"https://api.example.com","auth":{"type":"bearer","token":"secret123"}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["auth"]["token"], "******");
}
#[test]
fn test_mask_connector_secrets_basic_password() {
let config = r#"{"type":"http","url":"https://api.example.com","auth":{"type":"basic","username":"user","password":"secret"}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["auth"]["password"], "******");
assert_eq!(val["auth"]["username"], "user");
}
#[test]
fn test_mask_connector_secrets_api_key() {
let config = r#"{"type":"http","url":"https://api.example.com","auth":{"type":"apikey","key":"mysecretkey"}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["auth"]["key"], "******");
}
#[test]
fn test_mask_connector_secrets_top_level_fields() {
let config = r#"{"type":"http","url":"https://api.example.com","password":"top_secret","api_key":"ak123","token":"tk456","secret":"shhh"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["password"], "******");
assert_eq!(val["api_key"], "******");
assert_eq!(val["token"], "******");
assert_eq!(val["secret"], "******");
assert_eq!(val["url"], "https://api.example.com");
}
#[test]
fn url_userinfo_password_is_redacted() {
let config = r#"{"type":"es","url":"https://elastic:s3cret@es.internal:9200"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["url"], "https://elastic:******@es.internal:9200");
}
#[test]
fn url_userinfo_password_is_redacted_without_a_username() {
let config =
r#"{"type":"cache","backend":"redis","url":"redis://:hunter2@cache.internal:6379/0"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["url"], "redis://:******@cache.internal:6379/0");
assert_eq!(val["backend"], "redis");
}
#[test]
fn broker_list_entries_are_redacted() {
let config = r#"{"type":"kafka","brokers":["sasl_ssl://svc:tok3n@b1.example:9093","b2.example:9092"],"topic":"orders"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["brokers"][0], "sasl_ssl://svc:******@b1.example:9093");
assert_eq!(val["brokers"][1], "b2.example:9092");
assert_eq!(val["topic"], "orders");
}
#[test]
fn nested_secrets_below_the_first_level_are_masked() {
let config = r#"{"type":"http","url":"https://api.example.com","headers":{"authorization":"Bearer abc","x-tenant":"acme"},"extra":{"deep":{"client_secret":"cs123"}}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["headers"]["authorization"], "******");
assert_eq!(val["headers"]["x-tenant"], "******");
assert!(
val["headers"]
.as_object()
.expect("test")
.contains_key("x-tenant"),
"header names stay visible; values do not"
);
assert_eq!(val["extra"]["deep"]["client_secret"], "******");
}
#[test]
fn credential_bundles_mask_every_child() {
let config = r#"{"type":"db","credentials":{"id":"AKIAEXAMPLE","value":"wJalrXUtn"},"database":"assets"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["credentials"]["id"], "******");
assert_eq!(val["credentials"]["value"], "******");
assert_eq!(val["database"], "******");
}
#[test]
fn kafka_sasl_password_shape_is_masked() {
let config = r#"{"type":"kafka","brokers":["b:9092"],"topic":"t","auth":{"sasl_mechanism":"PLAIN","sasl_username":"svc","sasl_password":"p@ss"}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["auth"]["sasl_password"], "******");
assert_eq!(val["auth"]["sasl_username"], "svc");
assert_eq!(val["auth"]["sasl_mechanism"], "PLAIN");
}
#[test]
fn innocuous_keys_are_left_alone() {
let config = r#"{"type":"db","driver":"postgres","max_connections":10,"query_timeout_ms":5000,"cache_key_fields":["tenant"],"operations":{"read":true,"delete":false},"retry":{"max_retries":3}}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["driver"], "postgres");
assert_eq!(val["max_connections"], 10);
assert_eq!(val["query_timeout_ms"], 5000);
assert_eq!(val["cache_key_fields"][0], "******");
assert_eq!(val["operations"]["read"], true);
assert_eq!(val["operations"]["delete"], false);
assert_eq!(val["retry"]["max_retries"], 3);
}
#[test]
fn a_secret_under_an_unanticipated_key_is_masked() {
let config = r#"{"type":"http","url":"https://api.example.com","signing_cert_pem":"-----BEGIN...","tenant_code":"t-42"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["signing_cert_pem"], "******");
assert_eq!(val["tenant_code"], "******");
assert_eq!(val["url"], "https://api.example.com");
}
#[test]
fn non_string_secret_values_are_masked() {
let config = r#"{"type":"http","token":12345}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["token"], "******");
}
#[test]
fn redact_url_secrets_ignores_non_urls_and_credential_free_urls() {
assert_eq!(redact_url_secrets("plain text"), None);
assert_eq!(redact_url_secrets("https://example.com/a?b=c"), None);
assert_eq!(
redact_url_secrets("https://example.com/search?q=orion&page=2"),
None
);
assert_eq!(redact_url_secrets("postgres://user@host/db"), None);
assert_eq!(redact_url_secrets("https://host/path@here"), None);
}
#[test]
fn query_string_secret_is_redacted() {
let config = r#"{"type":"http","url":"https://api.example.com/v1?api_key=SECRET&page=2"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(
val["url"],
"https://api.example.com/v1?api_key=******&page=2"
);
}
#[test]
fn query_parameter_names_use_the_same_predicate_as_keys() {
for (url, expected) in [
(
"https://acct.blob.example.com/c/b?se=2026-01-01&sig=fzo7Ax8u",
"https://acct.blob.example.com/c/b?se=2026-01-01&sig=******",
),
(
"https://s3.example.com/k?X-Amz-Signature=abc123",
"https://s3.example.com/k?X-Amz-Signature=******",
),
(
"https://api.example.com/cb?access_token=tok&state=xyz",
"https://api.example.com/cb?access_token=******&state=xyz",
),
] {
assert_eq!(redact_url_secrets(url).as_deref(), Some(expected));
}
}
#[test]
fn userinfo_and_query_secrets_are_both_redacted_in_one_url() {
assert_eq!(
redact_url_secrets("https://svc:hunter2@api.example.com/v1?token=t0k3n#frag")
.as_deref(),
Some("https://svc:******@api.example.com/v1?token=******#frag")
);
}
#[test]
fn fragment_is_not_treated_as_query() {
assert_eq!(
redact_url_secrets("https://example.com/docs?page=2#api_key=example"),
None
);
}
#[test]
fn path_embedded_tokens_rely_on_the_key_name_rule() {
assert_eq!(
redact_url_secrets("https://hooks.slack.com/services/T00/B00/XXXX"),
None
);
let config = r#"{"type":"http","webhook_url":"https://hooks.slack.com/services/T00/B00/XXXX","url":"https://api.example.com"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["webhook_url"], "******");
assert_eq!(val["url"], "https://api.example.com");
}
#[test]
fn extended_denylist_terms_are_masked() {
let config = r#"{"type":"http","bearer":"b","sentry_dsn":"https://k@sentry.example/1","webhook":"https://hooks.example/T/B/X","pat":"ghp_abc","sig":"xyz","pattern":"/orders/{id}","path":"/v1"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
for key in [
"bearer",
"sentry_dsn",
"webhook",
"pat",
"sig",
"pattern",
"path",
] {
assert_eq!(val[key], "******", "{key}");
}
}
#[test]
fn test_unmask_restores_query_secret() {
let stored: Value = serde_json::from_str(
r#"{"type":"http","url":"https://api.example.com/v1?api_key=real-key&page=2"}"#,
)
.expect("test");
let mut incoming: Value =
serde_json::from_str(&mask_connector_secrets(&stored.to_string())).expect("test");
assert_eq!(
incoming["url"],
"https://api.example.com/v1?api_key=******&page=2"
);
unmask_config(&mut incoming, &stored);
assert_eq!(
incoming["url"],
"https://api.example.com/v1?api_key=real-key&page=2"
);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_find_masked_value_detects_the_query_form() {
let config: Value =
serde_json::from_str(r#"{"url":"https://api.example.com/v1?api_key=******"}"#)
.expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("url"));
}
#[test]
fn test_mask_connector_secrets_no_auth() {
let config = r#"{"type":"http","url":"https://api.example.com"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["url"], "https://api.example.com");
}
#[test]
fn test_mask_connector_secrets_invalid_json() {
let config = "not valid json";
let masked = mask_connector_secrets(config);
assert_eq!(masked, config);
}
#[test]
fn test_mask_connector_model() {
use chrono::NaiveDate;
let connector = crate::storage::models::Connector {
tags_json: "[]".to_string(),
id: "c1".to_string(),
name: "test".to_string(),
connector_type: "http".to_string(),
config_json: r#"{"type":"http","url":"https://api.example.com","auth":{"type":"bearer","token":"secret"}}"#.to_string(),
enabled: true,
created_at: NaiveDate::from_ymd_opt(2025, 1, 1)
.expect("test")
.and_hms_opt(0, 0, 0)
.expect("test"),
updated_at: NaiveDate::from_ymd_opt(2025, 1, 1)
.expect("test")
.and_hms_opt(0, 0, 0)
.expect("test"),
};
let masked = mask_connector(&connector);
assert_eq!(masked.id, "c1");
let val: serde_json::Value = serde_json::from_str(&masked.config_json).expect("test");
assert_eq!(val["auth"]["token"], "******");
}
#[test]
fn test_mask_connector_secrets_connection_string() {
let config = r#"{"type":"db","connection_string":"postgres://user:pass@host/db","driver":"postgres"}"#;
let masked = mask_connector_secrets(config);
let val: serde_json::Value = serde_json::from_str(&masked).expect("test");
assert_eq!(val["connection_string"], "******");
assert_eq!(val["driver"], "postgres");
}
#[test]
fn test_unmask_restores_round_tripped_secret() {
let stored: Value = serde_json::from_str(
r#"{"type":"http","url":"https://api.example.com","timeout_secs":5,
"auth":{"type":"bearer","token":"real-secret"}}"#,
)
.expect("test");
let mut incoming: Value =
serde_json::from_str(&mask_connector_secrets(&stored.to_string())).expect("test");
incoming["timeout_secs"] = serde_json::json!(30);
unmask_config(&mut incoming, &stored);
assert_eq!(incoming["auth"]["token"], "real-secret");
assert_eq!(incoming["timeout_secs"], 30);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_unmask_restores_url_password() {
let stored: Value =
serde_json::from_str(r#"{"type":"cache","url":"redis://admin:hunter2@redis:6379"}"#)
.expect("test");
let mut incoming: Value =
serde_json::from_str(&mask_connector_secrets(&stored.to_string())).expect("test");
assert_eq!(incoming["url"], "redis://admin:******@redis:6379");
unmask_config(&mut incoming, &stored);
assert_eq!(incoming["url"], "redis://admin:hunter2@redis:6379");
}
#[test]
fn test_unmask_restores_array_elements() {
let stored: Value = serde_json::from_str(
r#"{"type":"kafka","brokers":["SASL_SSL://u:p1@b1:9093","plaintext://b2:9092"]}"#,
)
.expect("test");
let mut incoming: Value =
serde_json::from_str(&mask_connector_secrets(&stored.to_string())).expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(incoming["brokers"][0], "SASL_SSL://u:p1@b1:9093");
assert_eq!(incoming["brokers"][1], "plaintext://b2:9092");
}
#[test]
fn test_unmask_leaves_a_genuine_new_secret_alone() {
let stored: Value =
serde_json::from_str(r#"{"auth":{"token":"old-secret"}}"#).expect("test");
let mut incoming: Value =
serde_json::from_str(r#"{"auth":{"token":"new-secret"}}"#).expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(incoming["auth"]["token"], "new-secret");
}
#[test]
fn test_unmask_leaves_unmatched_mask_for_rejection() {
let stored: Value = serde_json::from_str(r#"{"auth":{"token":"old"}}"#).expect("test");
let mut incoming: Value =
serde_json::from_str(r#"{"auth":{"token":"old"},"password":"******"}"#).expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(find_masked_value(&incoming).as_deref(), Some("password"));
}
#[test]
fn test_unmask_restores_both_url_positions() {
let stored: Value = serde_json::from_str(
r#"{"type":"http","url":"https://svc:hunter2@api.example.com/v1?api_key=real-key&page=2"}"#,
)
.expect("test");
let mut incoming: Value =
serde_json::from_str(&mask_connector_secrets(&stored.to_string())).expect("test");
assert_eq!(
incoming["url"],
"https://svc:******@api.example.com/v1?api_key=******&page=2"
);
unmask_config(&mut incoming, &stored);
assert_eq!(
incoming["url"],
"https://svc:hunter2@api.example.com/v1?api_key=real-key&page=2"
);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_unmask_restores_query_secret_while_password_rotates() {
let stored: Value = serde_json::from_str(
r#"{"url":"https://svc:oldpass@api.example.com/v1?api_key=real-key&page=2"}"#,
)
.expect("test");
let mut incoming: Value = serde_json::from_str(
r#"{"url":"https://svc:newpass@api.example.com/v1?api_key=******&page=2"}"#,
)
.expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(
incoming["url"],
"https://svc:newpass@api.example.com/v1?api_key=real-key&page=2"
);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_unmask_restores_password_while_query_secret_rotates() {
let stored: Value = serde_json::from_str(
r#"{"url":"https://svc:oldpass@api.example.com/v1?api_key=old-key"}"#,
)
.expect("test");
let mut incoming: Value = serde_json::from_str(
r#"{"url":"https://svc:******@api.example.com/v1?api_key=new-key"}"#,
)
.expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(
incoming["url"],
"https://svc:oldpass@api.example.com/v1?api_key=new-key"
);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_unmask_restores_two_query_secrets_independently() {
let stored: Value = serde_json::from_str(
r#"{"url":"https://api.example.com/v1?api_key=k1&page=2&access_token=k2"}"#,
)
.expect("test");
let mut incoming: Value = serde_json::from_str(
r#"{"url":"https://api.example.com/v1?api_key=******&page=3&access_token=******"}"#,
)
.expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(
incoming["url"],
"https://api.example.com/v1?api_key=k1&page=3&access_token=k2"
);
assert_eq!(find_masked_value(&incoming), None);
}
#[test]
fn test_unmask_leaves_unmatched_query_mask_for_rejection() {
let stored: Value =
serde_json::from_str(r#"{"url":"https://api.example.com/v1?page=2"}"#).expect("test");
let mut incoming: Value =
serde_json::from_str(r#"{"url":"https://api.example.com/v1?page=2&api_key=******"}"#)
.expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(find_masked_value(&incoming).as_deref(), Some("url"));
}
#[test]
fn test_unmask_leaves_unmatched_url_password_for_rejection() {
let stored: Value =
serde_json::from_str(r#"{"url":"https://api.example.com/v1"}"#).expect("test");
let mut incoming: Value =
serde_json::from_str(r#"{"url":"https://svc:******@api.example.com/v1"}"#)
.expect("test");
unmask_config(&mut incoming, &stored);
assert_eq!(find_masked_value(&incoming).as_deref(), Some("url"));
}
#[test]
fn test_find_masked_value_detects_partially_masked_urls() {
let config: Value = serde_json::from_str(
r#"{"url":"https://svc:newpass@api.example.com/v1?api_key=******"}"#,
)
.expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("url"));
let config: Value = serde_json::from_str(
r#"{"url":"https://svc:******@api.example.com/v1?api_key=fresh"}"#,
)
.expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("url"));
}
#[test]
fn test_find_masked_value_detects_the_sentinel_under_any_query_name() {
let config: Value =
serde_json::from_str(r#"{"url":"https://api.example.com/v1?page=******"}"#)
.expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("url"));
}
#[test]
fn test_find_masked_value_reports_a_dotted_path() {
let config: Value = serde_json::from_str(r#"{"a":{"b":[{"c":"******"}]}}"#).expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("a.b[0].c"));
}
#[test]
fn test_find_masked_value_detects_the_url_form() {
let config: Value =
serde_json::from_str(r#"{"url":"redis://admin:******@redis:6379"}"#).expect("test");
assert_eq!(find_masked_value(&config).as_deref(), Some("url"));
}
#[test]
fn test_find_masked_value_ignores_a_clean_config() {
let config: Value = serde_json::from_str(
r#"{"url":"https://api.example.com","auth":{"token":"real"},"n":6}"#,
)
.expect("test");
assert_eq!(find_masked_value(&config), None);
}
use serde_json::json;
#[test]
fn a_resolvable_reference_survives_masking() {
let mut v = json!({"auth": {"token": "env://STRIPE_KEY"}});
mask_secrets(&mut v);
assert_eq!(v["auth"]["token"], "env://STRIPE_KEY");
}
#[test]
fn a_reserved_scheme_reference_survives_masking() {
let mut v = json!({"password": "vault://kv/data/db#password"});
mask_secrets(&mut v);
assert_eq!(v["password"], "vault://kv/data/db#password");
}
#[test]
fn a_database_url_is_not_a_reference_and_is_still_masked() {
let mut v = json!({"connection_string": "postgres://user:hunter2@db/orders"});
mask_secrets(&mut v);
assert_eq!(
v["connection_string"], MASK,
"a URL is not a secret reference; masking it is the whole policy"
);
}
#[test]
fn a_malformed_reference_is_still_masked() {
for value in ["env://", "environment://X", "env:/NAME", "env-var://NAME"] {
let mut v = json!({ "token": value });
mask_secrets(&mut v);
assert_eq!(v["token"], MASK, "{value} must not read as a reference");
}
}
#[test]
fn url_redaction_under_a_non_secret_key_is_unchanged() {
let mut v = json!({"url": "https://user:hunter2@es:9200"});
mask_secrets(&mut v);
assert_eq!(v["url"], format!("https://user:{MASK}@es:9200"));
}
#[test]
fn every_connector_struct_field_is_classified() {
const EXPECTED_MASKED: &[&str] = &["connection_string", "auth"];
fn walk(v: &Value, out: &mut Vec<String>) {
if let Value::Object(map) = v {
for (k, child) in map {
match child {
Value::Object(_) => walk(child, out),
_ => out.push(k.clone()),
}
}
}
}
for sample in [
r#"{"type":"http","url":"https://x"}"#,
r#"{"type":"kafka","brokers":["b:9092"],"topic":"t"}"#,
r#"{"type":"db","connection_string":"postgres://h/db"}"#,
r#"{"type":"cache","backend":"memory"}"#,
r#"{"type":"es","url":"https://es:9200"}"#,
] {
let parsed: crate::connector::ConnectorConfig =
serde_json::from_str(sample).expect("sample parses");
let serialized = serde_json::to_value(&parsed).expect("serializes");
let mut leaf_keys = Vec::new();
walk(&serialized, &mut leaf_keys);
assert!(!leaf_keys.is_empty());
for key in leaf_keys {
assert!(
is_readable_key(&key) || EXPECTED_MASKED.contains(&key.as_str()),
"connector field '{key}' is unclassified: add it to READABLE_KEYS \
(safe to serve readable) or to EXPECTED_MASKED in this test \
(a credential, masked by default)"
);
}
}
}
#[test]
fn channel_auth_material_is_masked_and_nothing_else() {
let mut config = json!({
"auth": {
"mode": "api_key",
"header": "x-api-key",
"keys": ["sk-live-1", "sk-live-2"],
"secret": "whsec_abc",
"signature_prefix": "sha256="
},
"rate_limit": {"requests_per_second": 5},
"cache_key_fields": ["data.user_id"]
});
mask_channel_config(&mut config);
assert_eq!(config["auth"]["keys"][0], MASK);
assert_eq!(config["auth"]["keys"][1], MASK);
assert_eq!(config["auth"]["secret"], MASK);
assert_eq!(config["auth"]["header"], "x-api-key");
assert_eq!(config["auth"]["signature_prefix"], "sha256=");
assert_eq!(config["rate_limit"]["requests_per_second"], 5);
assert_eq!(config["cache_key_fields"][0], "data.user_id");
}
#[test]
fn channel_auth_references_survive_masking() {
let mut config = json!({"auth": {"mode": "hmac", "secret": "env://WEBHOOK_SECRET"}});
mask_channel_config(&mut config);
assert_eq!(config["auth"]["secret"], "env://WEBHOOK_SECRET");
}
#[test]
fn channel_config_without_auth_is_untouched() {
let mut config = json!({"dedup": {"window_secs": 60}});
let before = config.clone();
mask_channel_config(&mut config);
assert_eq!(config, before);
}
#[tokio::test]
async fn channel_unmask_restores_round_tripped_auth() {
let stored = json!({"auth": {"mode": "api_key", "keys": ["real-1", "real-2"]}});
let mut incoming = stored.clone();
mask_channel_config(&mut incoming);
assert_eq!(incoming["auth"]["keys"][0], MASK);
unmask_channel_config(&mut incoming, &stored);
assert_eq!(incoming["auth"]["keys"][0], "real-1");
assert_eq!(incoming["auth"]["keys"][1], "real-2");
assert_eq!(find_masked_value(&incoming), None);
let mut incoming = json!({"auth": {"mode": "api_key", "keys": ["******", "fresh-key"]}});
unmask_channel_config(&mut incoming, &stored);
assert_eq!(incoming["auth"]["keys"][0], "real-1");
assert_eq!(incoming["auth"]["keys"][1], "fresh-key");
}
#[test]
fn channel_unmask_leaves_unmatched_sentinel_for_rejection() {
let stored = json!({"auth": {"mode": "api_key", "keys": ["real-1"]}});
let mut incoming = json!({"auth": {"mode": "hmac", "secret": "******"}});
unmask_channel_config(&mut incoming, &stored);
assert_eq!(find_masked_value(&incoming).as_deref(), Some("auth.secret"));
}
}