Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 61 additions & 3 deletions aimdb-mqtt-connector/src/native.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ pub(crate) fn build<'a>(
})?;

let futures: Vec<BoxFuture> = vec![
Box::pin(run_event_loop(event_loop, inbound, broker_url.to_string())),
Box::pin(run_event_loop(event_loop, inbound, redacted(broker_url))),
Box::pin(run_publish_loop(client, outbound, opts)),
];
Ok(futures)
Expand Down Expand Up @@ -91,8 +91,10 @@ impl MqttConnectorImpl {
url = format!("{}/dummy", url.trim_end_matches('/'));
}

let connector_url =
ConnectorUrl::parse(&url).map_err(|e| format!("Invalid MQTT URL: {}", e))?;
let connector_url = ConnectorUrl::parse(&url).map_err(|e| {
// The parse error quotes its input, password included.
format!("Invalid MQTT URL: {}", e).replace(&url, &redacted(&url))
})?;

let host = connector_url.host.clone();
let port = connector_url.port.unwrap_or_else(|| {
Expand Down Expand Up @@ -254,6 +256,26 @@ async fn run_publish_loop(
log_info!("MQTT publish loop: every outbound route has closed");
}

/// `url` with the password in its authority masked, as `ConnectorUrl`'s
/// `Display` masks it: `mqtt://user:****@host:1883`. Safe for logs and errors.
/// A URL missing its scheme (the one input `ConnectorUrl::parse` rejects, and
/// quotes) is masked the same way.
fn redacted(url: &str) -> String {
let (prefix, rest) = match url.split_once("://") {
Some((_, rest)) => (url.len() - rest.len(), rest),
None => (0, url),
};
let authority_end = rest.find('/').unwrap_or(rest.len());
let (authority, path) = rest.split_at(authority_end);
let Some((userinfo, host)) = authority.rsplit_once('@') else {
return url.to_string();
};
match userinfo.split_once(':') {
Some((user, _password)) => format!("{}{user}:****@{host}{path}", &url[..prefix]),
None => url.to_string(),
}
}

/// The TLS configuration for `mqtts://`, from whichever backend this build
/// selected.
#[cfg(feature = "tokio-native-tls")]
Expand Down Expand Up @@ -420,6 +442,42 @@ mod tests {
}
}

#[test]
fn redacted_masks_only_the_password() {
assert_eq!(
redacted("mqtt://aimdb:s3cret@localhost:1884"),
"mqtt://aimdb:****@localhost:1884"
);
assert_eq!(
redacted("mqtts://user:p@ss:w0rd@broker.example.com:8883/x"),
"mqtts://user:****@broker.example.com:8883/x"
);
// Without a scheme, as a mistyped URL arrives.
assert_eq!(
redacted("user:s3cret@localhost:1884"),
"user:****@localhost:1884"
);
// Nothing to mask: unchanged.
for url in [
"mqtt://localhost:1883",
"mqtt://user@localhost:1883",
"not-a-url",
] {
assert_eq!(redacted(url), url);
}
}

#[tokio::test]
async fn an_invalid_url_error_does_not_repeat_the_password() {
let Err(err) =
MqttConnectorImpl::build_internal("user:s3cret@localhost:1884", None, None, 60, &[])
.await
else {
panic!("a URL without a scheme must not build");
};
assert!(!err.contains("s3cret"), "{err}");
}

/// The plain scheme is unaffected by which backend, if any, is selected.
#[tokio::test]
async fn test_connector_mqtt_url_needs_no_tls_backend() {
Expand Down