mirror of
https://github.com/telemt/telemt.git
synced 2026-09-30 22:45:58 +03:00
Compare commits
2 Commits
326c0ecdb9
...
38eabc50e3
| Author | SHA1 | Date | |
|---|---|---|---|
| 38eabc50e3 | |||
| 5feded2919 |
@@ -447,6 +447,9 @@ fn deep_merge(base: &mut Toml, patch: &Toml) {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "config_edit/base_path_tests.rs"]
|
||||
mod base_path_tests;
|
||||
#[cfg(test)]
|
||||
#[path = "config_edit/tests.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
use super::*;
|
||||
|
||||
fn web_config() -> &'static str {
|
||||
r#"
|
||||
[access.users]
|
||||
alice = "000102030405060708090a0b0c0d0e0f"
|
||||
|
||||
[[server.listeners]]
|
||||
ip = "127.0.0.1"
|
||||
port = 18080
|
||||
transport = "web"
|
||||
proxy_protocol = false
|
||||
web_client_ip_source = "x_forwarded_for"
|
||||
web_trusted_proxy_cidrs = ["127.0.0.1/32"]
|
||||
|
||||
[web]
|
||||
enabled = true
|
||||
|
||||
[[web.vhosts]]
|
||||
host = "proxy.example.com"
|
||||
public_addr = "203.0.113.10:443"
|
||||
|
||||
[web.vhosts.decoy]
|
||||
mode = "http_upstream"
|
||||
upstream = "http://127.0.0.1:18081"
|
||||
|
||||
[[web.vhosts.profiles]]
|
||||
user = "alice"
|
||||
secret_mode = "plain"
|
||||
"#
|
||||
}
|
||||
|
||||
fn vhosts_patch(base_path: &str) -> Json {
|
||||
serde_json::json!({
|
||||
"web": {
|
||||
"vhosts": [{
|
||||
"host": "proxy.example.com",
|
||||
"base_path": base_path,
|
||||
"public_addr": "203.0.113.10:443",
|
||||
"decoy": {
|
||||
"mode": "http_upstream",
|
||||
"upstream": "http://127.0.0.1:18081"
|
||||
},
|
||||
"profiles": [{
|
||||
"user": "alice",
|
||||
"secret_mode": "plain"
|
||||
}]
|
||||
}]
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn config_api_applies_valid_base_path_and_preserves_source_on_invalid_patch() {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let path = directory.path().join("config.toml");
|
||||
std::fs::write(&path, web_config()).unwrap();
|
||||
let active = ProxyConfig::load(&path).unwrap();
|
||||
|
||||
let mut response = apply_patch_to_path(&path, &vhosts_patch("MixedCase/path"), None)
|
||||
.await
|
||||
.unwrap();
|
||||
let desired = ProxyConfig::load(&path).unwrap();
|
||||
let resolved = reconcile_runtime_effect(&mut response, &active, &desired).unwrap();
|
||||
assert!(!response.restart_required);
|
||||
assert!(response.runtime_reload_required);
|
||||
assert!(!response.process_restart_required);
|
||||
assert!(response.deferred_process_fields.is_empty());
|
||||
assert!(resolved.runtime_changed);
|
||||
assert_eq!(desired.web.vhosts[0].base_path, "MixedCase/path");
|
||||
assert_eq!(
|
||||
resolved.effective.web.runtime.as_ref().unwrap().vhosts["proxy.example.com"].base,
|
||||
"/MixedCase/path/"
|
||||
);
|
||||
|
||||
let (managed, _revision) = read_managed_config(&path).await.unwrap();
|
||||
let vhosts = managed["web"]["vhosts"].as_array().unwrap();
|
||||
assert_eq!(vhosts[0]["base_path"].as_str(), Some("MixedCase/path"));
|
||||
assert!(managed["web"].get("runtime").is_none());
|
||||
assert!(!managed.as_table().unwrap().contains_key("access"));
|
||||
|
||||
let before_invalid = std::fs::read(&path).unwrap();
|
||||
let error = apply_patch_to_path(&path, &vhosts_patch("/invalid"), None)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert_eq!(error.status, hyper::StatusCode::BAD_REQUEST);
|
||||
assert_eq!(std::fs::read(&path).unwrap(), before_invalid);
|
||||
}
|
||||
@@ -62,5 +62,8 @@ use reporting::log_changes;
|
||||
#[cfg(test)]
|
||||
use watcher::{ReloadState, reload_config};
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "hot_reload/base_path_tests.rs"]
|
||||
mod base_path_tests;
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
use base64::Engine as _;
|
||||
|
||||
use super::*;
|
||||
|
||||
fn write_base_path_config(path: &Path, base_path: &str) {
|
||||
let base_path = if base_path.is_empty() {
|
||||
String::new()
|
||||
} else {
|
||||
format!("base_path = \"{base_path}\"\n")
|
||||
};
|
||||
let config = format!(
|
||||
r#"
|
||||
[access.users]
|
||||
alice = "000102030405060708090a0b0c0d0e0f"
|
||||
|
||||
[[server.listeners]]
|
||||
ip = "127.0.0.1"
|
||||
port = 18080
|
||||
transport = "web"
|
||||
proxy_protocol = false
|
||||
web_client_ip_source = "x_forwarded_for"
|
||||
web_trusted_proxy_cidrs = ["127.0.0.1/32"]
|
||||
|
||||
[web]
|
||||
enabled = true
|
||||
|
||||
[[web.vhosts]]
|
||||
host = "proxy.example.com"
|
||||
{base_path}public_addr = "203.0.113.10:443"
|
||||
|
||||
[web.vhosts.decoy]
|
||||
mode = "http_upstream"
|
||||
upstream = "http://127.0.0.1:18081"
|
||||
|
||||
[[web.vhosts.profiles]]
|
||||
user = "alice"
|
||||
secret_mode = "plain"
|
||||
"#,
|
||||
);
|
||||
std::fs::write(path, config).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reload_rejects_invalid_base_then_publishes_route_identity_together() {
|
||||
let directory = tempfile::tempdir().unwrap();
|
||||
let path = directory.path().join("config.toml");
|
||||
write_base_path_config(&path, "");
|
||||
let initial = Arc::new(ProxyConfig::load(&path).unwrap());
|
||||
let initial_hash = ProxyConfig::load_with_metadata(&path)
|
||||
.unwrap()
|
||||
.rendered_hash;
|
||||
let initial_capability = initial.web.runtime.as_ref().unwrap().capabilities[0];
|
||||
let (config_tx, _config_rx) = watch::channel(Arc::clone(&initial));
|
||||
let (log_tx, _log_rx) = watch::channel(initial.general.log_level.clone());
|
||||
let mut reload_state = ReloadState::new(Some(initial_hash));
|
||||
|
||||
write_base_path_config(&path, "/invalid");
|
||||
reload_config(&path, &config_tx, &log_tx, None, None, &mut reload_state);
|
||||
let unchanged = config_tx.borrow().clone();
|
||||
assert!(Arc::ptr_eq(&unchanged, &initial));
|
||||
assert_eq!(unchanged.web.vhosts[0].base_path, "");
|
||||
assert_eq!(
|
||||
unchanged.web.runtime.as_ref().unwrap().capabilities[0],
|
||||
initial_capability
|
||||
);
|
||||
|
||||
write_base_path_config(&path, "dobry-cola-super-app");
|
||||
reload_config(&path, &config_tx, &log_tx, None, None, &mut reload_state);
|
||||
let applied = config_tx.borrow().clone();
|
||||
let runtime = applied.web.runtime.as_ref().unwrap();
|
||||
let vhost = &runtime.vhosts["proxy.example.com"];
|
||||
assert_eq!(applied.web.vhosts[0].base_path, "dobry-cola-super-app");
|
||||
assert_eq!(vhost.base, "/dobry-cola-super-app/");
|
||||
assert_eq!(vhost.capabilities[0], vhost.profiles[0].capability);
|
||||
assert_eq!(runtime.capabilities.as_ref(), vhost.capabilities.as_ref());
|
||||
assert!(!runtime.capabilities.contains(&initial_capability));
|
||||
assert_eq!(
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(vhost.capabilities[0]),
|
||||
"hHz99Xs93EN1j91G9gpNepXwGNNt5YdAFkEVk_LlqdQ"
|
||||
);
|
||||
}
|
||||
@@ -30,7 +30,8 @@ use crate::util::secure_fs::open_dir_nofollow;
|
||||
#[cfg(not(unix))]
|
||||
mod static_site_fallback;
|
||||
|
||||
const WEB_CAPABILITY_CONTEXT: &[u8] = b"tdesktop-web-proxy-bridge-v1\n";
|
||||
const WEB_CAPABILITY_CONTEXT_V1: &[u8] = b"tdesktop-web-proxy-bridge-v1\n";
|
||||
const WEB_CAPABILITY_CONTEXT_V2: &[u8] = b"tdesktop-web-proxy-bridge-v2\n";
|
||||
const WEB_DEBUG_FINGERPRINT_CONTEXT: &[u8] = b"telemt-web-debug-key-fingerprint-v1\0";
|
||||
const MAX_WEB_STATIC_DEPTH: usize = 64;
|
||||
|
||||
@@ -41,6 +42,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
|
||||
})?;
|
||||
let mut runtime_vhosts = BTreeMap::new();
|
||||
let mut runtime_profiles = Vec::new();
|
||||
let mut runtime_capabilities = Vec::new();
|
||||
let mut static_files = 0usize;
|
||||
let mut static_bytes = 0usize;
|
||||
|
||||
@@ -67,8 +69,11 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
|
||||
})?;
|
||||
let (client_secret, client_secret_len) =
|
||||
client_secret(auth_entry.secret, profile.secret_mode);
|
||||
let capability =
|
||||
derive_web_capability(&client_secret[..client_secret_len], vhost.host.as_bytes())?;
|
||||
let capability = derive_web_capability(
|
||||
&client_secret[..client_secret_len],
|
||||
vhost.host.as_bytes(),
|
||||
vhost.base_path.as_bytes(),
|
||||
)?;
|
||||
let key_fingerprint = debug_key_fingerprint(&client_secret[..client_secret_len]);
|
||||
if !capabilities.insert(capability) {
|
||||
return Err(ProxyError::Config(format!(
|
||||
@@ -104,6 +109,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
|
||||
.unwrap_or(config.web.limits.max_streams_per_session),
|
||||
});
|
||||
capability_table.push(capability);
|
||||
runtime_capabilities.push(capability);
|
||||
profiles.push(Arc::clone(&runtime_profile));
|
||||
runtime_profiles.push(runtime_profile);
|
||||
}
|
||||
@@ -111,6 +117,11 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
|
||||
vhost.host.clone(),
|
||||
Arc::new(WebRuntimeVhost {
|
||||
host: vhost.host.clone(),
|
||||
base: if vhost.base_path.is_empty() {
|
||||
"/".to_string()
|
||||
} else {
|
||||
format!("/{}/", vhost.base_path)
|
||||
},
|
||||
decoy_fasttrack_mode: config.web.decoy_fasttrack_mode,
|
||||
decoy,
|
||||
decoy_header_secs: config.web.timeouts.decoy_header_secs,
|
||||
@@ -123,6 +134,7 @@ pub(super) fn rebuild(config: &mut ProxyConfig) -> Result<()> {
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts: runtime_vhosts,
|
||||
profiles: runtime_profiles,
|
||||
capabilities: runtime_capabilities.into_boxed_slice(),
|
||||
}));
|
||||
Ok(())
|
||||
}
|
||||
@@ -134,12 +146,23 @@ fn debug_key_fingerprint(secret: &[u8]) -> String {
|
||||
hex::encode(&digest.finalize()[..8])
|
||||
}
|
||||
|
||||
/// Derives the Telegram Desktop WEB capability for one exact secret and host.
|
||||
pub(crate) fn derive_web_capability(secret: &[u8], host: &[u8]) -> Result<[u8; 32]> {
|
||||
/// Derives the Telegram Desktop WEB capability for one exact secret, host, and base path.
|
||||
pub(crate) fn derive_web_capability(
|
||||
secret: &[u8],
|
||||
host: &[u8],
|
||||
base_path: &[u8],
|
||||
) -> Result<[u8; 32]> {
|
||||
let mut mac = Hmac::<Sha256>::new_from_slice(secret)
|
||||
.map_err(|_| ProxyError::Config("WEB capability secret must not be empty".to_string()))?;
|
||||
mac.update(WEB_CAPABILITY_CONTEXT);
|
||||
mac.update(host);
|
||||
if base_path.is_empty() {
|
||||
mac.update(WEB_CAPABILITY_CONTEXT_V1);
|
||||
mac.update(host);
|
||||
} else {
|
||||
mac.update(WEB_CAPABILITY_CONTEXT_V2);
|
||||
mac.update(host);
|
||||
mac.update(b"\n");
|
||||
mac.update(base_path);
|
||||
}
|
||||
Ok(mac.finalize().into_bytes().into())
|
||||
}
|
||||
|
||||
@@ -482,62 +505,7 @@ fn static_content_type(path: &Path) -> &'static str {
|
||||
}
|
||||
}
|
||||
|
||||
// Runtime WEB construction tests remain separate from the production loader.
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use base64::Engine as _;
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn capability_matches_reference_vectors() {
|
||||
let secret = hex::decode("000102030405060708090a0b0c0d0e0f").unwrap();
|
||||
let plain = derive_web_capability(&secret, b"proxy.example.com").unwrap();
|
||||
assert_eq!(
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(plain),
|
||||
"MHLEY5PmW1GWqJkSrlmJpvJUiLhBH_QKy6yKg8a0JPk"
|
||||
);
|
||||
let mut dd_secret = vec![0xdd];
|
||||
dd_secret.extend_from_slice(&secret);
|
||||
let dd = derive_web_capability(&dd_secret, b"proxy.example.com").unwrap();
|
||||
assert_eq!(
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(dd),
|
||||
"IpJrt3e7sKtzPyoXy6w-Zj6GGEvsvclN66JzQEfPYLA"
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn static_snapshot_remains_anchored_after_root_path_replacement() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let root = temp.path().join("site");
|
||||
let detached = temp.path().join("detached");
|
||||
let replacement = temp.path().join("replacement");
|
||||
fs::create_dir(&root).unwrap();
|
||||
fs::write(root.join("index.html"), b"original").unwrap();
|
||||
fs::create_dir(&replacement).unwrap();
|
||||
fs::write(replacement.join("index.html"), b"replacement").unwrap();
|
||||
|
||||
let directory = open_static_root(&root).unwrap();
|
||||
fs::rename(&root, &detached).unwrap();
|
||||
symlink(&replacement, &root).unwrap();
|
||||
|
||||
let mut assets = BTreeMap::new();
|
||||
let mut total_files = 0;
|
||||
let mut total_bytes = 0;
|
||||
load_static_directory(
|
||||
directory,
|
||||
Path::new(""),
|
||||
&root,
|
||||
&mut assets,
|
||||
&mut total_files,
|
||||
&mut total_bytes,
|
||||
&WebLimitsConfig::default(),
|
||||
0,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(assets["/index.html"].body.as_ref(), b"original");
|
||||
}
|
||||
}
|
||||
#[path = "runtime_web/tests.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
use base64::Engine as _;
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn capability_matches_reference_vectors() {
|
||||
let secret = hex::decode("000102030405060708090a0b0c0d0e0f").unwrap();
|
||||
let mut dd_secret = vec![0xdd];
|
||||
dd_secret.extend_from_slice(&secret);
|
||||
for (client_secret, base_path, expected) in [
|
||||
(
|
||||
secret.as_slice(),
|
||||
b"".as_slice(),
|
||||
"MHLEY5PmW1GWqJkSrlmJpvJUiLhBH_QKy6yKg8a0JPk",
|
||||
),
|
||||
(
|
||||
dd_secret.as_slice(),
|
||||
b"".as_slice(),
|
||||
"IpJrt3e7sKtzPyoXy6w-Zj6GGEvsvclN66JzQEfPYLA",
|
||||
),
|
||||
(
|
||||
secret.as_slice(),
|
||||
b"dobry-cola-super-app".as_slice(),
|
||||
"hHz99Xs93EN1j91G9gpNepXwGNNt5YdAFkEVk_LlqdQ",
|
||||
),
|
||||
(
|
||||
dd_secret.as_slice(),
|
||||
b"dobry-cola-super-app".as_slice(),
|
||||
"TGUkZaevsavLbHvlNWipnRoYxgzZ51ioWvbxgGT3wHo",
|
||||
),
|
||||
] {
|
||||
let capability =
|
||||
derive_web_capability(client_secret, b"proxy.example.com", base_path).unwrap();
|
||||
assert_eq!(
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability),
|
||||
expected
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn capability_binds_the_exact_host_and_base_path_identity() {
|
||||
let secret = hex::decode("000102030405060708090a0b0c0d0e0f").unwrap();
|
||||
let root = derive_web_capability(&secret, b"proxy.example.com", b"").unwrap();
|
||||
let mixed = derive_web_capability(&secret, b"proxy.example.com", b"MixedCase/path").unwrap();
|
||||
let lower = derive_web_capability(&secret, b"proxy.example.com", b"mixedcase/path").unwrap();
|
||||
let other_path =
|
||||
derive_web_capability(&secret, b"proxy.example.com", b"MixedCase/other").unwrap();
|
||||
let other_host =
|
||||
derive_web_capability(&secret, b"other.example.com", b"MixedCase/path").unwrap();
|
||||
|
||||
let identities = [root, mixed, lower, other_path, other_host]
|
||||
.into_iter()
|
||||
.collect::<std::collections::HashSet<_>>();
|
||||
assert_eq!(identities.len(), 5);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn static_snapshot_remains_anchored_after_root_path_replacement() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let root = temp.path().join("site");
|
||||
let detached = temp.path().join("detached");
|
||||
let replacement = temp.path().join("replacement");
|
||||
fs::create_dir(&root).unwrap();
|
||||
fs::write(root.join("index.html"), b"original").unwrap();
|
||||
fs::create_dir(&replacement).unwrap();
|
||||
fs::write(replacement.join("index.html"), b"replacement").unwrap();
|
||||
|
||||
let directory = open_static_root(&root).unwrap();
|
||||
fs::rename(&root, &detached).unwrap();
|
||||
symlink(&replacement, &root).unwrap();
|
||||
|
||||
let mut assets = BTreeMap::new();
|
||||
let mut total_files = 0;
|
||||
let mut total_bytes = 0;
|
||||
load_static_directory(
|
||||
directory,
|
||||
Path::new(""),
|
||||
&root,
|
||||
&mut assets,
|
||||
&mut total_files,
|
||||
&mut total_bytes,
|
||||
&WebLimitsConfig::default(),
|
||||
0,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(assets["/index.html"].body.as_ref(), b"original");
|
||||
}
|
||||
@@ -365,7 +365,7 @@ const WEB_TIMEOUTS_CONFIG_KEYS: &[&str] = &[
|
||||
"decoy_header_secs",
|
||||
];
|
||||
|
||||
const WEB_VHOST_CONFIG_KEYS: &[&str] = &["host", "public_addr", "decoy", "profiles"];
|
||||
const WEB_VHOST_CONFIG_KEYS: &[&str] = &["host", "base_path", "public_addr", "decoy", "profiles"];
|
||||
const WEB_DECOY_CONFIG_KEYS: &[&str] = &["mode", "upstream", "directory", "index"];
|
||||
const WEB_PROFILE_CONFIG_KEYS: &[&str] = &[
|
||||
"user",
|
||||
|
||||
@@ -6,7 +6,8 @@ const WEB_DEBUG_GROUP_SCRATCH_BYTES: usize = 4 * 1024 * 1024;
|
||||
const WEB_CARRIER_LEARNING_ENTRY_BYTES: usize = 512;
|
||||
const WEB_LANE_STATE_BYTES: usize = 512;
|
||||
const WEB_OVERLOAD_CONNECTION_BYTES: usize = 4 * 1024;
|
||||
const WEB_CAPABILITY_INDEX_ENTRY_BYTES: usize = 32;
|
||||
// Each profile capability is stored in its vhost and in the global containment table.
|
||||
const WEB_CAPABILITY_INDEX_ENTRY_BYTES: usize = 64;
|
||||
|
||||
/// Validates process-wide body, header, queue, static, and debug reservations.
|
||||
pub(super) fn validate(limits: &WebLimitsConfig) -> Result<()> {
|
||||
|
||||
@@ -9,6 +9,10 @@ pub(super) fn validate_vhosts(config: &mut ProxyConfig) -> Result<()> {
|
||||
let mut profile_count = 0usize;
|
||||
for (vhost_idx, vhost) in config.web.vhosts.iter_mut().enumerate() {
|
||||
vhost.host = normalize_web_host(&vhost.host, &format!("web.vhosts[{vhost_idx}].host"))?;
|
||||
validate_web_base_path(
|
||||
&vhost.base_path,
|
||||
&format!("web.vhosts[{vhost_idx}].base_path"),
|
||||
)?;
|
||||
if !hosts.insert(vhost.host.clone()) {
|
||||
return config_error(&format!("duplicate WEB vhost host `{}`", vhost.host));
|
||||
}
|
||||
@@ -75,6 +79,25 @@ pub(super) fn validate_vhosts(config: &mut ProxyConfig) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn validate_web_base_path(value: &str, field: &str) -> Result<()> {
|
||||
let valid = value.len() <= 128
|
||||
&& !value.starts_with('/')
|
||||
&& !value.ends_with('/')
|
||||
&& value.split('/').all(|segment| {
|
||||
let mut bytes = segment.bytes();
|
||||
bytes
|
||||
.next()
|
||||
.is_some_and(|byte| byte.is_ascii_alphanumeric())
|
||||
&& bytes.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_'))
|
||||
});
|
||||
if value.is_empty() || valid {
|
||||
return Ok(());
|
||||
}
|
||||
config_error(&format!(
|
||||
"{field} must be empty or contain at most 128 ASCII bytes in slash-separated [A-Za-z0-9][A-Za-z0-9_-]* segments"
|
||||
))
|
||||
}
|
||||
|
||||
pub(super) fn normalize_web_host(value: &str, field: &str) -> Result<String> {
|
||||
let input = value.trim();
|
||||
if input.is_empty()
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
use super::*;
|
||||
|
||||
#[path = "web_tests/base_path_tests.rs"]
|
||||
mod base_path_tests;
|
||||
|
||||
const WEB_CONFIG: &str = r#"
|
||||
[access.users]
|
||||
alice = "000102030405060708090a0b0c0d0e0f"
|
||||
@@ -40,9 +43,11 @@ fn web_config_builds_canonical_runtime_snapshot() {
|
||||
.vhosts
|
||||
.get("proxy.example.com")
|
||||
.expect("canonical WEB vhost");
|
||||
assert_eq!(vhost.base, "/");
|
||||
assert_eq!(vhost.profiles.len(), 1);
|
||||
assert_eq!(vhost.capabilities.len(), vhost.profiles.len());
|
||||
assert_eq!(vhost.capabilities[0], vhost.profiles[0].capability);
|
||||
assert_eq!(runtime.capabilities.as_ref(), vhost.capabilities.as_ref());
|
||||
assert_eq!(vhost.decoy_fasttrack_mode, WebDecoyFastTrackMode::Off);
|
||||
assert_eq!(vhost.profiles[0].user, "alice");
|
||||
assert_eq!(vhost.profiles[0].secret_mode, WebSecretMode::Dd);
|
||||
@@ -59,6 +64,66 @@ fn web_config_builds_canonical_runtime_snapshot() {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn web_base_path_is_canonical_and_precomputed() {
|
||||
let maximum = "a".repeat(128);
|
||||
for base_path in ["a", "a/b/c9_x-y", "Dobry-Cola/super_app", maximum.as_str()] {
|
||||
let configured = WEB_CONFIG.replace(
|
||||
"host = \"Proxy.Example.COM\"",
|
||||
&format!("host = \"Proxy.Example.COM\"\nbase_path = \"{base_path}\""),
|
||||
);
|
||||
let config = load_config_from_temp_toml(&configured);
|
||||
assert_eq!(config.web.vhosts[0].base_path, base_path);
|
||||
assert_eq!(
|
||||
config.web.runtime.as_ref().unwrap().vhosts["proxy.example.com"].base,
|
||||
format!("/{base_path}/")
|
||||
);
|
||||
}
|
||||
|
||||
let strict = format!(
|
||||
"[general]\nconfig_strict = true\n{}",
|
||||
WEB_CONFIG.replace(
|
||||
"host = \"Proxy.Example.COM\"",
|
||||
"host = \"Proxy.Example.COM\"\nbase_path = \"relay\"",
|
||||
)
|
||||
);
|
||||
assert_eq!(
|
||||
load_config_from_temp_toml(&strict).web.vhosts[0].base_path,
|
||||
"relay"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn web_base_path_rejects_noncanonical_forms() {
|
||||
let oversized = "a".repeat(129);
|
||||
for base_path in [
|
||||
"/relay",
|
||||
"relay/",
|
||||
"relay//nested",
|
||||
"-relay",
|
||||
"_relay",
|
||||
"a/-lead",
|
||||
"a/_lead",
|
||||
"dot.ted",
|
||||
"..",
|
||||
"/",
|
||||
"relay/.hidden",
|
||||
"relay/%2fhidden",
|
||||
"relay path",
|
||||
"relay/тест",
|
||||
oversized.as_str(),
|
||||
] {
|
||||
let invalid = WEB_CONFIG.replace(
|
||||
"host = \"Proxy.Example.COM\"",
|
||||
&format!("host = \"Proxy.Example.COM\"\nbase_path = \"{base_path}\""),
|
||||
);
|
||||
assert!(
|
||||
load_config_error_from_temp_toml(&invalid).contains("web.vhosts[0].base_path"),
|
||||
"base path {base_path:?} was accepted"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn web_decoy_fasttrack_mode_is_typed_and_defaults_off() {
|
||||
let defaults = ProxyConfig::default();
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn web_runtime_collects_every_vhost_capability() {
|
||||
let configured = format!(
|
||||
"{WEB_CONFIG}\n{}",
|
||||
r#"
|
||||
[[web.vhosts]]
|
||||
host = "Other.Example.COM"
|
||||
base_path = "other/path"
|
||||
public_addr = "203.0.113.11:443"
|
||||
|
||||
[web.vhosts.decoy]
|
||||
mode = "http_upstream"
|
||||
upstream = "http://127.0.0.1:18082"
|
||||
|
||||
[[web.vhosts.profiles]]
|
||||
user = "alice"
|
||||
secret_mode = "dd"
|
||||
"#
|
||||
);
|
||||
let config = load_config_from_temp_toml(&configured);
|
||||
let runtime = config.web.runtime.as_ref().unwrap();
|
||||
let first = &runtime.vhosts["proxy.example.com"];
|
||||
let second = &runtime.vhosts["other.example.com"];
|
||||
|
||||
assert_eq!(runtime.capabilities.len(), 2);
|
||||
assert_eq!(first.capabilities[0], first.profiles[0].capability);
|
||||
assert_eq!(second.capabilities[0], second.profiles[0].capability);
|
||||
assert_ne!(first.capabilities[0], second.capabilities[0]);
|
||||
assert!(runtime.capabilities.contains(&first.capabilities[0]));
|
||||
assert!(runtime.capabilities.contains(&second.capabilities[0]));
|
||||
}
|
||||
@@ -71,6 +71,9 @@ pub enum WebDecoyConfig {
|
||||
pub struct WebVhostConfig {
|
||||
/// Canonical lowercase ACE hostname used by Telegram Desktop.
|
||||
pub host: String,
|
||||
/// Optional canonical WEB endpoint prefix without surrounding slashes.
|
||||
#[serde(default)]
|
||||
pub base_path: String,
|
||||
/// Stable public destination tuple used by inner relay routing and KDF metadata.
|
||||
pub public_addr: SocketAddr,
|
||||
/// Ordinary-site fallback for this hostname.
|
||||
|
||||
@@ -7,6 +7,8 @@ pub(crate) struct WebRuntimeConfig {
|
||||
pub(crate) vhosts: BTreeMap<String, Arc<WebRuntimeVhost>>,
|
||||
/// Flat profile inventory used by startup link emission.
|
||||
pub(crate) profiles: Vec<Arc<WebRuntimeProfile>>,
|
||||
/// Complete active capability table used to contain misplaced credentials.
|
||||
pub(crate) capabilities: Box<[[u8; 32]]>,
|
||||
}
|
||||
|
||||
/// Precomputed immutable virtual-host data.
|
||||
@@ -14,6 +16,8 @@ pub(crate) struct WebRuntimeConfig {
|
||||
pub(crate) struct WebRuntimeVhost {
|
||||
/// Canonical lowercase ACE hostname.
|
||||
pub(crate) host: String,
|
||||
/// Exact slash-delimited endpoint base, including the trailing slash.
|
||||
pub(crate) base: String,
|
||||
/// Restart-frozen decoy capability-scan policy.
|
||||
pub(crate) decoy_fasttrack_mode: WebDecoyFastTrackMode,
|
||||
/// Immutable ordinary-site fallback snapshot.
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
use std::time::Duration;
|
||||
|
||||
use base64::Engine as _;
|
||||
use tokio::sync::watch;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
@@ -80,19 +81,52 @@ pub(crate) fn print_web_proxy_links(config: &ProxyConfig) {
|
||||
let Some(secret) = config.access.users.get(&profile.user) else {
|
||||
continue;
|
||||
};
|
||||
let prefix = match profile.secret_mode {
|
||||
crate::config::WebSecretMode::Plain => "",
|
||||
crate::config::WebSecretMode::Dd => "dd",
|
||||
let Some(vhost) = config
|
||||
.web
|
||||
.vhosts
|
||||
.iter()
|
||||
.find(|vhost| vhost.host == profile.host)
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
print_maestro_line(format!(
|
||||
"User: {} ({:?})",
|
||||
profile.user, profile.secret_mode
|
||||
));
|
||||
print_maestro_line(format!(
|
||||
"WEB: tg://webproxy?server={}&secret={prefix}{secret}",
|
||||
profile.host,
|
||||
if let Some(link) =
|
||||
format_web_proxy_link(&profile.host, &vhost.base_path, secret, profile.secret_mode)
|
||||
{
|
||||
print_maestro_line(format!("WEB: {link}"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn format_web_proxy_link(
|
||||
host: &str,
|
||||
base_path: &str,
|
||||
secret: &str,
|
||||
mode: crate::config::WebSecretMode,
|
||||
) -> Option<String> {
|
||||
if base_path.is_empty() {
|
||||
let prefix = match mode {
|
||||
crate::config::WebSecretMode::Plain => "",
|
||||
crate::config::WebSecretMode::Dd => "dd",
|
||||
};
|
||||
return Some(format!(
|
||||
"tg://webproxy?server={host}&secret={prefix}{secret}"
|
||||
));
|
||||
}
|
||||
let decoded = hex::decode(secret).ok()?;
|
||||
let mut marked = Vec::with_capacity(decoded.len() + 2);
|
||||
marked.push(0x70);
|
||||
if mode == crate::config::WebSecretMode::Dd {
|
||||
marked.push(0xdd);
|
||||
}
|
||||
marked.extend_from_slice(&decoded);
|
||||
let server = url::form_urlencoded::byte_serialize(format!("{host}/{base_path}").as_bytes())
|
||||
.collect::<String>();
|
||||
let marked = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(marked);
|
||||
Some(format!("tg://webproxy?server={server}&secret={marked}"))
|
||||
}
|
||||
|
||||
/// Durably replaces one Beobachten snapshot without following Unix symlinks.
|
||||
@@ -367,3 +401,76 @@ pub(crate) async fn load_startup_proxy_config_snapshot(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::WebSecretMode;
|
||||
|
||||
const SECRET: &str = "000102030405060708090a0b0c0d0e0f";
|
||||
|
||||
#[test]
|
||||
fn root_web_proxy_links_keep_the_legacy_secret_form() {
|
||||
assert_eq!(
|
||||
format_web_proxy_link("proxy.example.com", "", SECRET, WebSecretMode::Plain),
|
||||
Some(format!(
|
||||
"tg://webproxy?server=proxy.example.com&secret={SECRET}"
|
||||
))
|
||||
);
|
||||
assert_eq!(
|
||||
format_web_proxy_link("proxy.example.com", "", SECRET, WebSecretMode::Dd),
|
||||
Some(format!(
|
||||
"tg://webproxy?server=proxy.example.com&secret=dd{SECRET}"
|
||||
))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn path_web_proxy_links_use_the_tdesktop_marker() {
|
||||
assert_eq!(
|
||||
format_web_proxy_link(
|
||||
"proxy.example.com",
|
||||
"dobry-cola/super_app",
|
||||
SECRET,
|
||||
WebSecretMode::Plain,
|
||||
),
|
||||
Some("tg://webproxy?server=proxy.example.com%2Fdobry-cola%2Fsuper_app&secret=cAABAgMEBQYHCAkKCwwNDg8".to_string())
|
||||
);
|
||||
assert_eq!(
|
||||
format_web_proxy_link(
|
||||
"proxy.example.com",
|
||||
"dobry-cola/super_app",
|
||||
SECRET,
|
||||
WebSecretMode::Dd,
|
||||
),
|
||||
Some("tg://webproxy?server=proxy.example.com%2Fdobry-cola%2Fsuper_app&secret=cN0AAQIDBAUGBwgJCgsMDQ4P".to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn path_web_proxy_link_round_trips_through_the_tdesktop_grammar() {
|
||||
for (mode, expected_secret) in [
|
||||
(WebSecretMode::Plain, hex::decode(SECRET).unwrap()),
|
||||
(
|
||||
WebSecretMode::Dd,
|
||||
[vec![0xdd], hex::decode(SECRET).unwrap()].concat(),
|
||||
),
|
||||
] {
|
||||
let link = format_web_proxy_link("proxy.example.com", "MixedCase/a_b-9", SECRET, mode)
|
||||
.unwrap();
|
||||
let parsed = url::Url::parse(&link).unwrap();
|
||||
let query = parsed
|
||||
.query_pairs()
|
||||
.collect::<std::collections::BTreeMap<_, _>>();
|
||||
assert_eq!(
|
||||
query["server"].as_ref(),
|
||||
"proxy.example.com/MixedCase/a_b-9"
|
||||
);
|
||||
let marked = base64::engine::general_purpose::URL_SAFE_NO_PAD
|
||||
.decode(query["secret"].as_bytes())
|
||||
.unwrap();
|
||||
assert_eq!(marked.first(), Some(&0x70));
|
||||
assert_eq!(&marked[1..], expected_secret.as_slice());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -295,6 +295,28 @@ fn web_decoy_fasttrack_mode_is_deferred_without_runtime_publication() {
|
||||
assert!(!resolved.runtime_changed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn base_path_change_is_runtime_owned_and_rebuilds_route_identity() {
|
||||
let old = web_config_with_fasttrack("off");
|
||||
let old_runtime = old.web.runtime.as_ref().unwrap();
|
||||
let old_capability = old_runtime.vhosts["proxy.example.com"].capabilities[0];
|
||||
let mut desired = old.clone();
|
||||
desired.web.vhosts[0].base_path = "MixedCase/path".to_string();
|
||||
|
||||
let resolved = resolve_reload_config(&old, &desired).unwrap();
|
||||
|
||||
assert!(resolved.deferred_process_fields.is_empty());
|
||||
assert!(resolved.runtime_changed);
|
||||
assert_eq!(resolved.effective.web.vhosts[0].base_path, "MixedCase/path");
|
||||
let runtime = resolved.effective.web.runtime.as_ref().unwrap();
|
||||
let vhost = &runtime.vhosts["proxy.example.com"];
|
||||
assert_eq!(vhost.base, "/MixedCase/path/");
|
||||
assert_ne!(vhost.capabilities[0], old_capability);
|
||||
assert_eq!(vhost.capabilities[0], vhost.profiles[0].capability);
|
||||
assert_eq!(runtime.capabilities.as_ref(), vhost.capabilities.as_ref());
|
||||
assert!(!runtime.capabilities.contains(&old_capability));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn enabling_learning_is_deferred_when_retained_capacity_is_too_small() {
|
||||
let mut old = ProxyConfig::default();
|
||||
|
||||
@@ -17,6 +17,7 @@ pub(crate) struct BridgePage {
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) fn render(
|
||||
host: &str,
|
||||
base: &str,
|
||||
bootstrap: &str,
|
||||
batch_limit: usize,
|
||||
queue_limit: usize,
|
||||
@@ -55,6 +56,7 @@ pub(crate) fn render(
|
||||
} else {
|
||||
String::new()
|
||||
};
|
||||
let base_prefix = base.strip_suffix('/').unwrap_or(base);
|
||||
let body = DOCUMENT
|
||||
.replace("__DIAGNOSTIC_RUNTIME__\n", &diagnostic_script)
|
||||
.replace("__RESPONSE_RUNTIME__", RESPONSE_RUNTIME)
|
||||
@@ -120,6 +122,7 @@ pub(crate) fn render(
|
||||
)
|
||||
.replace("__NONCE__", &nonce)
|
||||
.replace("__HOST__", host)
|
||||
.replace("__BASE_PREFIX__", base_prefix)
|
||||
.replace("__BOOTSTRAP__", bootstrap)
|
||||
.replace("__BATCH_LIMIT__", &batch_limit.to_string())
|
||||
.replace("__QUEUE_LIMIT__", &queue_limit.to_string())
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
(()=>{'use strict';
|
||||
let bootstrap="__BOOTSTRAP__",hello=false,emitted=0,boundaryTimer=null;
|
||||
const relayOrigin='https://__HOST__',requestMs=__BRIDGE_REQUEST_SECS__*1000;
|
||||
const relayOrigin='https://__HOST__',relayBase=relayOrigin+'__BASE_PREFIX__',requestMs=__BRIDGE_REQUEST_SECS__*1000;
|
||||
function eventBit(event){
|
||||
if(event==='runtime_started')return 1;if(event==='status_posted')return 2;if(event==='hello_received')return 4;if(event==='boundary_timeout')return 8;
|
||||
if(event==='hello_timeout')return 16;if(event==='client_close_before_hello')return 32;if(event==='document_unloaded_before_hello')return 64;
|
||||
@@ -12,7 +12,7 @@ function report(event){
|
||||
let timer=null;
|
||||
try{
|
||||
const controller=new AbortController(),body=JSON.stringify({v:1,event});timer=setTimeout(()=>controller.abort(),requestMs);
|
||||
fetch(relayOrigin+'/api/v1/diagnostic',{method:'POST',body,signal:controller.signal,keepalive:true,mode:'same-origin',credentials:'omit',cache:'no-store',redirect:'error',referrerPolicy:'no-referrer',headers:{Authorization:'Bearer '+bootstrap,'Content-Type':'application/json'}})
|
||||
fetch(relayBase+'/api/v1/diagnostic',{method:'POST',body,signal:controller.signal,keepalive:true,mode:'same-origin',credentials:'omit',cache:'no-store',redirect:'error',referrerPolicy:'no-referrer',headers:{Authorization:'Bearer '+bootstrap,'Content-Type':'application/json'}})
|
||||
.then(discard,()=>{}).then(()=>clearTimeout(timer),()=>clearTimeout(timer));
|
||||
}catch(error){if(timer)clearTimeout(timer)}
|
||||
}
|
||||
|
||||
@@ -40,7 +40,7 @@ function create(settings){
|
||||
const timer=setTimeout(()=>{timedOut=true;controller.abort()},Math.max(1,Math.min(attemptLimit,remaining)));
|
||||
let response=null,wait=0;
|
||||
try{
|
||||
const fetched=await fetch(settings.origin()+path,requestOptions);
|
||||
const fetched=await fetch(settings.base()+path,requestOptions);
|
||||
if(retryableStatus(fetched.status)){
|
||||
lastReason='http';wait=retryAfterMs(fetched);settings.cancel(fetched);
|
||||
}else{
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
(()=>{'use strict';
|
||||
let bootstrap="__BOOTSTRAP__";
|
||||
const relayOrigin='https://__HOST__',carrierCapabilities='https,https-lanes,websocket,websocket-lanes';
|
||||
const relayOrigin='https://__HOST__',relayBase=relayOrigin+'__BASE_PREFIX__',carrierCapabilities='https,https-lanes,websocket,websocket-lanes';
|
||||
__DIAGNOSTIC_BINDING__;
|
||||
__DIAGNOSTIC_RUNTIME_STARTED__;
|
||||
const responseBody=globalThis.TelemtBridgeResponse;if(!responseBody)throw new Error('missing response runtime');
|
||||
@@ -26,9 +26,9 @@ const canonicalFailures=['timeout','network','upgrade','http','protocol'];
|
||||
const failure=(reason,message)=>Object.assign(new Error(message||reason),{telemtReason:reason});
|
||||
const failureReason=(error,fallback)=>error&&canonicalFailures.includes(error.telemtReason)?error.telemtReason:fallback;
|
||||
const status=__STATUS_FUNCTION__;
|
||||
const socketURL=()=>relayOrigin.replace(/^https:/,'wss:')+'/api/v1/ws';
|
||||
const socketURL=()=>relayBase.replace(/^https:/,'wss:')+'/api/v1/ws';
|
||||
const requestClient=requestSupport.create({
|
||||
origin:()=>relayOrigin,closed:()=>closed,retryMs:()=>bridgeRetryMs,longPollMs:()=>longPollMs,requestMs:()=>bridgeRequestMs,
|
||||
base:()=>relayBase,closed:()=>closed,retryMs:()=>bridgeRetryMs,longPollMs:()=>longPollMs,requestMs:()=>bridgeRequestMs,
|
||||
batchLimit:()=>batchLimit,read:(response,limit,exact,signal)=>responseBody.read(response,limit,exact,signal),cancel:responseBody.cancel,
|
||||
failure,reason:failureReason,retrying:()=>status('reconnecting')
|
||||
});
|
||||
@@ -478,7 +478,7 @@ async function pollLane(lane){
|
||||
}
|
||||
function deleteSession(){
|
||||
const token=cleanupToken||sessionToken,headers=canonicalFailures.includes(terminalFailure)?{'X-Carrier-Failure':terminalFailure}:null;
|
||||
if(token)fetch(relayOrigin+'/api/v1/session',options('DELETE',token,null,headers,undefined,true)).catch(()=>{});
|
||||
if(token)fetch(relayBase+'/api/v1/session',options('DELETE',token,null,headers,undefined,true)).catch(()=>{});
|
||||
}
|
||||
function close(notifyServer){
|
||||
if(closed)return;closed=true;if(recoveryController)recoveryController.cancel();rejectRecoveryCommit(failure('network','bridge closed'));if(helloTimer)clearTimeout(helloTimer);helloTimer=null;if(carrierTimer)clearTimeout(carrierTimer);clearProbeTimer();if(schedulerTimer)clearTimeout(schedulerTimer);schedulerTimer=null;if(attemptController)attemptController.abort();if(pollController)pollController.abort();
|
||||
|
||||
+48
-1
@@ -3,6 +3,7 @@ use super::*;
|
||||
fn render_page(bootstrap: &str, candidate_count: usize) -> BridgePage {
|
||||
render(
|
||||
"proxy.example.com",
|
||||
"/",
|
||||
bootstrap,
|
||||
2 * 1024 * 1024,
|
||||
32 * 1024 * 1024,
|
||||
@@ -26,6 +27,7 @@ fn render_page(bootstrap: &str, candidate_count: usize) -> BridgePage {
|
||||
fn render_diagnostic_page(bootstrap: &str) -> BridgePage {
|
||||
render(
|
||||
"proxy.example.com",
|
||||
"/",
|
||||
bootstrap,
|
||||
2 * 1024 * 1024,
|
||||
32 * 1024 * 1024,
|
||||
@@ -72,6 +74,49 @@ fn rendered_page_contains_bounded_negotiation_contract() {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rendered_page_resolves_carriers_against_the_exact_base_path() {
|
||||
let page = render(
|
||||
"proxy.example.com",
|
||||
"/Dobry-Cola/super_app/",
|
||||
"AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
|
||||
2 * 1024 * 1024,
|
||||
32 * 1024 * 1024,
|
||||
16 * 1024,
|
||||
1024,
|
||||
true,
|
||||
4,
|
||||
[3, 5, 8, 12],
|
||||
25,
|
||||
10,
|
||||
90,
|
||||
15,
|
||||
15,
|
||||
120,
|
||||
0,
|
||||
true,
|
||||
&SecureRandom::new(),
|
||||
);
|
||||
|
||||
assert!(
|
||||
page.body
|
||||
.contains("relayBase=relayOrigin+'/Dobry-Cola/super_app'")
|
||||
);
|
||||
assert!(page.body.contains("fetch(settings.base()+path"));
|
||||
assert!(
|
||||
page.body
|
||||
.contains("relayBase.replace(/^https:/,'wss:')+'/api/v1/ws'")
|
||||
);
|
||||
assert!(page.body.contains("fetch(relayBase+'/api/v1/diagnostic'"));
|
||||
assert!(page.body.contains("url:()=>relayOrigin+recoveryPath"));
|
||||
assert!(
|
||||
!page
|
||||
.body
|
||||
.contains("/Dobry-Cola/super_app/Dobry-Cola/super_app")
|
||||
);
|
||||
assert!(!page.body.contains("__BASE_PREFIX__"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rendered_page_preserves_the_ios_bootstrap_literal() {
|
||||
let bootstrap = "BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB";
|
||||
@@ -86,6 +131,7 @@ fn rendered_page_preserves_the_ios_bootstrap_literal() {
|
||||
fn rendered_page_embeds_the_configured_bridge_timing_policy() {
|
||||
let page = render(
|
||||
"proxy.example.com",
|
||||
"/",
|
||||
"GGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGGG",
|
||||
2 * 1024 * 1024,
|
||||
32 * 1024 * 1024,
|
||||
@@ -138,6 +184,7 @@ fn effective_deadline_formula_uses_the_final_checkpoint() {
|
||||
fn disabled_negotiation_does_not_arm_a_carrier_deadline() {
|
||||
let page = render(
|
||||
"proxy.example.com",
|
||||
"/",
|
||||
"DDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDDD",
|
||||
2 * 1024 * 1024,
|
||||
32 * 1024 * 1024,
|
||||
@@ -284,7 +331,7 @@ fn enabled_bridge_diagnostics_use_the_https_sideband_only() {
|
||||
let page = render_diagnostic_page("JJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJJ");
|
||||
|
||||
assert!(!page.body.contains("__"));
|
||||
assert!(page.body.contains("fetch(relayOrigin+'/api/v1/diagnostic'"));
|
||||
assert!(page.body.contains("fetch(relayBase+'/api/v1/diagnostic'"));
|
||||
assert!(page.body.contains("JSON.stringify({v:1,event})"));
|
||||
assert!(page.body.contains("'Content-Type':'application/json'"));
|
||||
assert!(page.body.contains("keepalive:true"));
|
||||
|
||||
+21
-15
@@ -28,6 +28,8 @@ mod activity;
|
||||
mod body;
|
||||
// Canonical capability parsing and complete scans remain isolated from HTTP routing.
|
||||
mod capability;
|
||||
// Authentic credential containment stays independent from carrier routing.
|
||||
mod secrets;
|
||||
// Decoy routing and upstream proxying are isolated from carrier authentication.
|
||||
mod decoy;
|
||||
// Authenticated generated-bridge diagnostics remain outside carrier framing.
|
||||
@@ -71,13 +73,13 @@ type BoxError = Box<dyn Error + Send + Sync>;
|
||||
type HttpBody = UnsyncBoxBody<Bytes, BoxError>;
|
||||
type HttpResponse = Response<HttpBody>;
|
||||
|
||||
const TRANSPORT_PATHS: [&str; 4] = [
|
||||
"/api/v1/session",
|
||||
"/api/v1/up",
|
||||
"/api/v1/down",
|
||||
"/api/v1/diagnostic",
|
||||
const TRANSPORT_SUFFIXES: [&str; 4] = [
|
||||
"api/v1/session",
|
||||
"api/v1/up",
|
||||
"api/v1/down",
|
||||
"api/v1/diagnostic",
|
||||
];
|
||||
const WEBSOCKET_PATH: &str = "/api/v1/ws";
|
||||
const WEBSOCKET_SUFFIX: &str = "api/v1/ws";
|
||||
|
||||
/// Serves one bounded HTTP/1.1 connection accepted from an external TLS terminator.
|
||||
pub(crate) async fn serve_connection(
|
||||
@@ -185,8 +187,9 @@ async fn handle_request(
|
||||
let Some(vhost) = web_runtime.vhosts.get(host).cloned() else {
|
||||
return generic_not_found();
|
||||
};
|
||||
let path = request.uri().path();
|
||||
if path == WEBSOCKET_PATH {
|
||||
secrets::mark_internal_credential(&mut request, web_runtime, &runtime);
|
||||
let suffix = request.uri().path().strip_prefix(&vhost.base);
|
||||
if suffix == Some(WEBSOCKET_SUFFIX) {
|
||||
return websocket::handle(
|
||||
request,
|
||||
peer,
|
||||
@@ -197,7 +200,7 @@ async fn handle_request(
|
||||
)
|
||||
.await;
|
||||
}
|
||||
if TRANSPORT_PATHS.contains(&path) {
|
||||
if suffix.is_some_and(|suffix| TRANSPORT_SUFFIXES.contains(&suffix)) {
|
||||
return handle_api(
|
||||
request,
|
||||
peer,
|
||||
@@ -208,7 +211,7 @@ async fn handle_request(
|
||||
)
|
||||
.await;
|
||||
}
|
||||
if path == "/" && matches!(*request.method(), Method::GET | Method::HEAD) {
|
||||
if suffix == Some("") && matches!(*request.method(), Method::GET | Method::HEAD) {
|
||||
return handle_root(
|
||||
request,
|
||||
peer,
|
||||
@@ -369,6 +372,7 @@ async fn handle_root(
|
||||
}
|
||||
let page = bridge::render(
|
||||
&vhost.host,
|
||||
&vhost.base,
|
||||
&bootstrap.token,
|
||||
config.web.limits.carrier_batch_bytes,
|
||||
config.web.limits.pending_bytes_per_session,
|
||||
@@ -440,11 +444,13 @@ async fn handle_api(
|
||||
let Some(token_hash) = bearer_token_hash(&request) else {
|
||||
return serve_decoy(request, vhost, true, &runtime).await;
|
||||
};
|
||||
match request.uri().path() {
|
||||
"/api/v1/session" => handle_session(request, runtime, vhost, token_hash, client_ip).await,
|
||||
"/api/v1/up" => handle_up(request, runtime, vhost, token_hash).await,
|
||||
"/api/v1/down" => handle_down(request, runtime, vhost, token_hash).await,
|
||||
"/api/v1/diagnostic" => {
|
||||
match request.uri().path().strip_prefix(&vhost.base) {
|
||||
Some("api/v1/session") => {
|
||||
handle_session(request, runtime, vhost, token_hash, client_ip).await
|
||||
}
|
||||
Some("api/v1/up") => handle_up(request, runtime, vhost, token_hash).await,
|
||||
Some("api/v1/down") => handle_down(request, runtime, vhost, token_hash).await,
|
||||
Some("api/v1/diagnostic") => {
|
||||
diagnostic::handle(request, runtime, vhost, token_hash, client_ip).await
|
||||
}
|
||||
_ => serve_decoy(request, vhost, true, &runtime).await,
|
||||
|
||||
@@ -0,0 +1,432 @@
|
||||
use super::*;
|
||||
|
||||
#[path = "base_path_tests/credentials.rs"]
|
||||
mod credentials;
|
||||
#[path = "base_path_tests/reload.rs"]
|
||||
mod reload;
|
||||
#[path = "base_path_tests/routing.rs"]
|
||||
mod routing;
|
||||
|
||||
fn bridge_request(path: &str) -> Vec<u8> {
|
||||
format!(
|
||||
"GET {path} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes()
|
||||
}
|
||||
|
||||
fn bridge_token(body: &[u8]) -> String {
|
||||
std::str::from_utf8(body)
|
||||
.unwrap()
|
||||
.split_once("bootstrap=\"")
|
||||
.and_then(|(_, suffix)| suffix.split_once('"'))
|
||||
.map(|(token, _)| token.to_string())
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn request_without_body(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
request_head: &[u8],
|
||||
) -> Vec<u8> {
|
||||
let addr = listener.local_addr().unwrap();
|
||||
let (accepted, client) = tokio::join!(listener.accept(), TcpStream::connect(addr));
|
||||
let (server, peer) = accepted.unwrap();
|
||||
let mut client = client.unwrap();
|
||||
let permit = runtime.try_http_connection().unwrap();
|
||||
let task = tokio::spawn(serve_connection(
|
||||
server,
|
||||
peer,
|
||||
WebClientIpSource::XForwardedFor,
|
||||
Arc::from(["127.0.0.1/32".parse().unwrap()]),
|
||||
Arc::clone(runtime),
|
||||
CancellationToken::new(),
|
||||
permit,
|
||||
));
|
||||
client.write_all(request_head).await.unwrap();
|
||||
let mut response = Vec::new();
|
||||
tokio::time::timeout(
|
||||
std::time::Duration::from_secs(2),
|
||||
client.read_to_end(&mut response),
|
||||
)
|
||||
.await
|
||||
.expect("private rejection waited for the request body")
|
||||
.unwrap();
|
||||
task.await.unwrap();
|
||||
response
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn base_path_routes_only_the_exact_prefixed_contract() {
|
||||
let capability = [21u8; 32];
|
||||
let config = runtime_config_with_base(capability, WebCarrier::Https, "/dobry-cola/");
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
|
||||
let response = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/dobry-cola/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(
|
||||
std::str::from_utf8(body)
|
||||
.unwrap()
|
||||
.contains("relayBase=relayOrigin+'/dobry-cola'")
|
||||
);
|
||||
|
||||
for path in [
|
||||
format!("/?bridge={encoded}"),
|
||||
format!("/dobry-cola?bridge={encoded}"),
|
||||
format!("/dobry-cola/nested/?bridge={encoded}"),
|
||||
] {
|
||||
let response = request(&listener, &runtime, bridge_request(&path)).await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
for path in [
|
||||
"/",
|
||||
"/api/v1/session",
|
||||
"/dobry-cola/unknown?q=1",
|
||||
"/dobry-cola//api/v1/up",
|
||||
"/dobry-cola%2Fapi/v1/up",
|
||||
] {
|
||||
let response = request(&listener, &runtime, bridge_request(path)).await;
|
||||
let (_, body) = split_response(&response);
|
||||
assert_eq!(body, b"<!doctype html><title>decoy</title>");
|
||||
}
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prefixed_https_carrier_creates_uses_and_closes_a_session() {
|
||||
let capability = [22u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/nested/");
|
||||
config.web.timeouts.long_poll_secs = 1;
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
|
||||
let bridge = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/relay/nested/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let bootstrap = bridge_token(split_response(&bridge).1);
|
||||
let hello = frame::encode(FrameType::Hello, 0, &[1]);
|
||||
let mut create = format!(
|
||||
"POST /relay/nested/api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
hello.len()
|
||||
)
|
||||
.into_bytes();
|
||||
create.extend_from_slice(&hello);
|
||||
let created = request(&listener, &runtime, create).await;
|
||||
let (headers, _) = split_response(&created);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
let session = response_header(headers, "x-session-token").to_string();
|
||||
|
||||
let misplaced = format!(
|
||||
"GET /wrong HTTP/1.1\r\nHost: proxy.example.com\r\nAuthorization: Bearer {session}\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes();
|
||||
let misplaced = request(&listener, &runtime, misplaced).await;
|
||||
let (headers, body) = split_response(&misplaced);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
|
||||
let pong = frame::encode(FrameType::Pong, 0, &[]);
|
||||
let mut uplink = format!(
|
||||
"POST /relay/nested/api/v1/up HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nContent-Type: application/octet-stream\r\nX-Up-Seq: 1\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
pong.len()
|
||||
)
|
||||
.into_bytes();
|
||||
uplink.extend_from_slice(&pong);
|
||||
assert!(
|
||||
request(&listener, &runtime, uplink)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
let downlink = format!(
|
||||
"POST /relay/nested/api/v1/down HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nX-Down-Cursor: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes();
|
||||
assert!(
|
||||
request(&listener, &runtime, downlink)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
let close = format!(
|
||||
"DELETE /relay/nested/api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes();
|
||||
assert!(
|
||||
request(&listener, &runtime, close)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn authentic_credentials_never_reach_the_decoy() {
|
||||
let capability = [23u8; 32];
|
||||
let config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let bridge = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/relay/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let bootstrap = bridge_token(split_response(&bridge).1);
|
||||
let escaped = format!("%{:02X}{}", bootstrap.as_bytes()[0], &bootstrap[1..]);
|
||||
|
||||
for raw in [
|
||||
format!(
|
||||
"GET /wrong/{bootstrap} HTTP/1.1\r\nHost: proxy.example.com\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /wrong/{escaped} HTTP/1.1\r\nHost: proxy.example.com\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /wrong HTTP/1.1\r\nHost: proxy.example.com\r\nCookie: opaque={bootstrap}\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /wrong HTTP/1.1\r\nHost: proxy.example.com\r\nReferer: https://example.invalid/{encoded}\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
] {
|
||||
let response = request(&listener, &runtime, raw.into_bytes()).await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
let mut forged = bootstrap.into_bytes();
|
||||
forged[0] = if forged[0] == b'A' { b'B' } else { b'A' };
|
||||
let forged = String::from_utf8(forged).unwrap();
|
||||
let response = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/wrong/{forged}")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
split_response(&response).1,
|
||||
b"<!doctype html><title>decoy</title>"
|
||||
);
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn misplaced_process_token_is_rejected_without_reading_the_body() {
|
||||
let capability = [24u8; 32];
|
||||
let config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let bridge = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/relay/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let bootstrap = bridge_token(split_response(&bridge).1);
|
||||
let head = format!(
|
||||
"POST /wrong HTTP/1.1\r\nHost: proxy.example.com\r\nAuthorization: Bearer {bootstrap}\r\nContent-Length: 1048576\r\nConnection: close\r\n\r\n"
|
||||
);
|
||||
|
||||
let response = request_without_body(&listener, &runtime, head.as_bytes()).await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn process_token_provenance_survives_registry_expiry() {
|
||||
let capability = [25u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
config.web.timeouts.bootstrap_lifetime_secs = 1;
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let bridge = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/relay/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let bootstrap = bridge_token(split_response(&bridge).1);
|
||||
|
||||
tokio::time::timeout(std::time::Duration::from_secs(4), async {
|
||||
loop {
|
||||
let status = serde_json::to_value(runtime.try_status()).unwrap();
|
||||
if status["manager"]["bootstraps"] == 0 {
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("bootstrap registry entry did not expire");
|
||||
|
||||
let response = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/wrong/{bootstrap}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn generation_swap_switches_base_path_and_capability_together() {
|
||||
let initial_capability = [26u8; 32];
|
||||
let mut initial = runtime_config_with_base(initial_capability, WebCarrier::Https, "/old-path/");
|
||||
initial.web.limits.max_bootstraps_per_ip = 2;
|
||||
let generation = test_runtime_generation(1, initial);
|
||||
let active = Arc::new(ArcSwap::from(Arc::clone(&generation)));
|
||||
let runtime = WebProcessRuntime::start(Arc::clone(&active));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let initial_encoded =
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(initial_capability);
|
||||
|
||||
let old_bridge = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/old-path/?bridge={initial_encoded}")),
|
||||
)
|
||||
.await;
|
||||
assert!(old_bridge.starts_with(b"HTTP/1.1 200"));
|
||||
|
||||
let replacement_capability = [27u8; 32];
|
||||
let replacement = test_runtime_generation(
|
||||
2,
|
||||
runtime_config_with_base(replacement_capability, WebCarrier::Https, "/new-path/"),
|
||||
);
|
||||
active.store(Arc::clone(&replacement));
|
||||
let stale = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/old-path/?bridge={initial_encoded}")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
split_response(&stale).1,
|
||||
b"<!doctype html><title>decoy</title>"
|
||||
);
|
||||
|
||||
let replacement_encoded =
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(replacement_capability);
|
||||
let current = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/new-path/?bridge={replacement_encoded}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(¤t);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(
|
||||
std::str::from_utf8(body)
|
||||
.unwrap()
|
||||
.contains("relayBase=relayOrigin+'/new-path'")
|
||||
);
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
replacement.stop_sessions().await;
|
||||
replacement.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prefixed_decoy_request_keeps_its_original_path_and_query() {
|
||||
let site = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let site_addr = site.local_addr().unwrap();
|
||||
let site_task = tokio::spawn(async move {
|
||||
let (mut stream, _) = site.accept().await.unwrap();
|
||||
let mut request = vec![0; 4096];
|
||||
let read = stream.read(&mut request).await.unwrap();
|
||||
request.truncate(read);
|
||||
stream
|
||||
.write_all(
|
||||
b"HTTP/1.1 404 Not Found\r\nContent-Length: 4\r\nConnection: close\r\n\r\nsite",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
request
|
||||
});
|
||||
|
||||
let capability = [28u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
let profile = Arc::clone(&config.web.runtime.as_ref().unwrap().profiles[0]);
|
||||
let vhost = Arc::new(WebRuntimeVhost {
|
||||
host: "proxy.example.com".to_string(),
|
||||
base: "/relay/".to_string(),
|
||||
decoy_fasttrack_mode: WebDecoyFastTrackMode::Off,
|
||||
decoy: WebRuntimeDecoy::HttpUpstream {
|
||||
addr: site_addr,
|
||||
authority: "decoy.internal".to_string(),
|
||||
},
|
||||
decoy_header_secs: 1,
|
||||
profiles: vec![Arc::clone(&profile)],
|
||||
capabilities: vec![capability].into_boxed_slice(),
|
||||
});
|
||||
let mut vhosts = BTreeMap::new();
|
||||
vhosts.insert("proxy.example.com".to_string(), vhost);
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts,
|
||||
profiles: vec![profile],
|
||||
capabilities: vec![capability].into_boxed_slice(),
|
||||
}));
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
|
||||
let response = request(&listener, &runtime, bridge_request("/relay/ordinary?q=1")).await;
|
||||
assert!(response.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(split_response(&response).1, b"site");
|
||||
let forwarded = site_task.await.unwrap();
|
||||
assert!(forwarded.starts_with(b"GET /relay/ordinary?q=1 HTTP/1.1\r\n"));
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
@@ -0,0 +1,162 @@
|
||||
use super::*;
|
||||
|
||||
fn assert_private_not_found(response: &[u8]) {
|
||||
let (headers, body) = split_response(response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
fn percent_encode(value: &str) -> String {
|
||||
value
|
||||
.bytes()
|
||||
.enumerate()
|
||||
.map(|(index, byte)| {
|
||||
if index % 2 == 0 {
|
||||
format!("%{byte:02x}")
|
||||
} else {
|
||||
format!("%{byte:02X}")
|
||||
}
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn issue_session_credentials(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
capability: [u8; 32],
|
||||
) -> (String, String) {
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let bridge = request(
|
||||
listener,
|
||||
runtime,
|
||||
bridge_request(&format!("/relay/?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let bootstrap = bridge_token(split_response(&bridge).1);
|
||||
let hello = frame::encode(FrameType::Hello, 0, &[1]);
|
||||
let mut create = format!(
|
||||
"POST /relay/api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
hello.len()
|
||||
)
|
||||
.into_bytes();
|
||||
create.extend_from_slice(&hello);
|
||||
let created = request(listener, runtime, create).await;
|
||||
let (headers, _) = split_response(&created);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
(
|
||||
bootstrap,
|
||||
response_header(headers, "x-session-token").to_string(),
|
||||
)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn every_authentic_credential_placement_stays_out_of_the_upstream() {
|
||||
let site = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let site_addr = site.local_addr().unwrap();
|
||||
let capability = [91u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
config.web.limits.max_bootstraps_per_ip = 2;
|
||||
let runtime_config = config.web.runtime.as_ref().unwrap();
|
||||
let profile = runtime_config.profiles[0].clone();
|
||||
let mut vhosts = BTreeMap::new();
|
||||
vhosts.insert(
|
||||
"proxy.example.com".to_string(),
|
||||
Arc::new(WebRuntimeVhost {
|
||||
host: "proxy.example.com".to_string(),
|
||||
base: "/relay/".to_string(),
|
||||
decoy_fasttrack_mode: WebDecoyFastTrackMode::Off,
|
||||
decoy: WebRuntimeDecoy::HttpUpstream {
|
||||
addr: site_addr,
|
||||
authority: "decoy.internal".to_string(),
|
||||
},
|
||||
decoy_header_secs: 1,
|
||||
profiles: vec![profile.clone()],
|
||||
capabilities: vec![capability].into_boxed_slice(),
|
||||
}),
|
||||
);
|
||||
vhosts.insert(
|
||||
"other.example.com".to_string(),
|
||||
Arc::new(WebRuntimeVhost {
|
||||
host: "other.example.com".to_string(),
|
||||
base: "/other/".to_string(),
|
||||
decoy_fasttrack_mode: WebDecoyFastTrackMode::Off,
|
||||
decoy: WebRuntimeDecoy::HttpUpstream {
|
||||
addr: site_addr,
|
||||
authority: "decoy.internal".to_string(),
|
||||
},
|
||||
decoy_header_secs: 1,
|
||||
profiles: Vec::new(),
|
||||
capabilities: Vec::new().into_boxed_slice(),
|
||||
}),
|
||||
);
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts,
|
||||
profiles: vec![profile],
|
||||
capabilities: vec![capability].into_boxed_slice(),
|
||||
}));
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded_capability = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let (bootstrap, session) = issue_session_credentials(&listener, &runtime, capability).await;
|
||||
|
||||
for secret in [&encoded_capability, &bootstrap, &session] {
|
||||
for raw in [
|
||||
format!(
|
||||
"GET /wrong/{secret} HTTP/1.1\r\nHost: proxy.example.com\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /wrong/{} HTTP/1.1\r\nHost: proxy.example.com\r\nConnection: close\r\n\r\n",
|
||||
percent_encode(secret)
|
||||
),
|
||||
format!(
|
||||
"GET /?bridge={secret}&extra=1 HTTP/1.1\r\nHost: proxy.example.com\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"POST /wrong HTTP/1.1\r\nHost: proxy.example.com\r\nAuthorization: Bearer {secret}\r\nCookie: opaque={secret}\r\nReferer: https://example.invalid/{secret}\r\nX-Unexpected: random,{secret}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /api/v1/ws HTTP/1.1\r\nHost: proxy.example.com\r\nSec-WebSocket-Protocol: chat, tproxy-v1.{secret}\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
format!(
|
||||
"GET /other/wrong HTTP/1.1\r\nHost: other.example.com\r\nX-Unexpected: {secret}\r\nConnection: close\r\n\r\n"
|
||||
),
|
||||
] {
|
||||
assert_private_not_found(&request(&listener, &runtime, raw.into_bytes()).await);
|
||||
}
|
||||
}
|
||||
|
||||
assert!(
|
||||
tokio::time::timeout(std::time::Duration::from_millis(100), site.accept())
|
||||
.await
|
||||
.is_err()
|
||||
);
|
||||
|
||||
let site_task = tokio::spawn(async move {
|
||||
let (mut stream, _) = site.accept().await.unwrap();
|
||||
let mut received = [0; 1024];
|
||||
let _ = stream.read(&mut received).await.unwrap();
|
||||
stream
|
||||
.write_all(
|
||||
b"HTTP/1.1 404 Not Found\r\nContent-Length: 4\r\nConnection: close\r\n\r\nsite",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
let mut forged = bootstrap.into_bytes();
|
||||
forged[0] = if forged[0] == b'A' { b'B' } else { b'A' };
|
||||
let forged = String::from_utf8(forged).unwrap();
|
||||
let response = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/wrong/{forged}")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(split_response(&response).1, b"site");
|
||||
site_task.await.unwrap();
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
@@ -0,0 +1,283 @@
|
||||
use super::*;
|
||||
|
||||
async fn issue_bootstrap(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
base: &str,
|
||||
capability: [u8; 32],
|
||||
) -> String {
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let response = request(
|
||||
listener,
|
||||
runtime,
|
||||
bridge_request(&format!("{base}?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
bridge_token(body)
|
||||
}
|
||||
|
||||
async fn create_session(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
base: &str,
|
||||
bootstrap: &str,
|
||||
) -> String {
|
||||
let hello = frame::encode(FrameType::Hello, 0, &[1]);
|
||||
let mut request_bytes = format!(
|
||||
"POST {base}api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
hello.len()
|
||||
)
|
||||
.into_bytes();
|
||||
request_bytes.extend_from_slice(&hello);
|
||||
let response = request(listener, runtime, request_bytes).await;
|
||||
let (headers, _) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
response_header(headers, "x-session-token").to_string()
|
||||
}
|
||||
|
||||
fn session_request(base: &str, session: &str) -> Vec<u8> {
|
||||
format!(
|
||||
"POST {base}api/v1/down HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nX-Down-Cursor: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes()
|
||||
}
|
||||
|
||||
fn assert_private_not_found(response: &[u8]) {
|
||||
let (headers, body) = split_response(response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
fn assert_decoy(response: &[u8]) {
|
||||
assert_eq!(
|
||||
split_response(response).1,
|
||||
b"<!doctype html><title>decoy</title>"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn generation_swap_preserves_process_tokens_and_replaces_route_identity() {
|
||||
let old_capability = [101u8; 32];
|
||||
let mut old_config = runtime_config_with_base(old_capability, WebCarrier::Https, "/old/");
|
||||
old_config.web.limits.max_bootstraps_per_ip = 8;
|
||||
old_config.web.timeouts.long_poll_secs = 1;
|
||||
let old_generation = test_runtime_generation(1, old_config);
|
||||
let active = Arc::new(ArcSwap::from(Arc::clone(&old_generation)));
|
||||
let runtime = WebProcessRuntime::start(Arc::clone(&active));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let unused_bootstrap = issue_bootstrap(&listener, &runtime, "/old/", old_capability).await;
|
||||
let used_bootstrap = issue_bootstrap(&listener, &runtime, "/old/", old_capability).await;
|
||||
let session = create_session(&listener, &runtime, "/old/", &used_bootstrap).await;
|
||||
|
||||
let new_capability = [102u8; 32];
|
||||
let mut new_config = runtime_config_with_base(new_capability, WebCarrier::Https, "/new/");
|
||||
new_config.web.limits.max_bootstraps_per_ip = 8;
|
||||
new_config.web.timeouts.long_poll_secs = 1;
|
||||
let new_generation = test_runtime_generation(2, new_config);
|
||||
active.store(Arc::clone(&new_generation));
|
||||
let old_encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(old_capability);
|
||||
let new_encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(new_capability);
|
||||
|
||||
assert_decoy(
|
||||
&request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/old/?bridge={old_encoded}")),
|
||||
)
|
||||
.await,
|
||||
);
|
||||
assert_private_not_found(
|
||||
&request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/old/?bridge={new_encoded}")),
|
||||
)
|
||||
.await,
|
||||
);
|
||||
assert_decoy(
|
||||
&request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/new/?bridge={old_encoded}")),
|
||||
)
|
||||
.await,
|
||||
);
|
||||
|
||||
let old_bootstrap_path = format!(
|
||||
"GET /old/wrong HTTP/1.1\r\nHost: proxy.example.com\r\nAuthorization: Bearer {unused_bootstrap}\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes();
|
||||
assert_private_not_found(&request(&listener, &runtime, old_bootstrap_path).await);
|
||||
assert_private_not_found(
|
||||
&request(&listener, &runtime, session_request("/old/", &session)).await,
|
||||
);
|
||||
assert!(
|
||||
request(&listener, &runtime, session_request("/new/", &session),)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
let hello = frame::encode(FrameType::Hello, 0, &[1]);
|
||||
let mut stale_bootstrap_request = format!(
|
||||
"POST /new/api/v1/session HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {unused_bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
hello.len()
|
||||
)
|
||||
.into_bytes();
|
||||
stale_bootstrap_request.extend_from_slice(&hello);
|
||||
assert_private_not_found(&request(&listener, &runtime, stale_bootstrap_request).await);
|
||||
|
||||
let current = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/new/?bridge={new_encoded}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(¤t);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
let new_bootstrap = bridge_token(body);
|
||||
create_session(&listener, &runtime, "/new/", &new_bootstrap).await;
|
||||
|
||||
runtime.shutdown().await;
|
||||
old_generation.stop_sessions().await;
|
||||
old_generation.stop_background_tasks().await;
|
||||
new_generation.stop_sessions().await;
|
||||
new_generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn routed_uplink_finishes_under_its_acquisition_generation() {
|
||||
let capability = [103u8; 32];
|
||||
let mut old_config = runtime_config_with_base(capability, WebCarrier::Https, "/old/");
|
||||
old_config.web.timeouts.long_poll_secs = 1;
|
||||
let old_generation = test_runtime_generation(1, old_config);
|
||||
let active = Arc::new(ArcSwap::from(Arc::clone(&old_generation)));
|
||||
let runtime = WebProcessRuntime::start(Arc::clone(&active));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let bootstrap = issue_bootstrap(&listener, &runtime, "/old/", capability).await;
|
||||
let session = create_session(&listener, &runtime, "/old/", &bootstrap).await;
|
||||
|
||||
let address = listener.local_addr().unwrap();
|
||||
let (accepted, client) = tokio::join!(listener.accept(), TcpStream::connect(address));
|
||||
let (server, peer) = accepted.unwrap();
|
||||
let mut client = client.unwrap();
|
||||
let permit = runtime.try_http_connection().unwrap();
|
||||
let task = tokio::spawn(serve_connection(
|
||||
server,
|
||||
peer,
|
||||
WebClientIpSource::XForwardedFor,
|
||||
Arc::from(["127.0.0.1/32".parse().unwrap()]),
|
||||
Arc::clone(&runtime),
|
||||
CancellationToken::new(),
|
||||
permit,
|
||||
));
|
||||
let pong = frame::encode(FrameType::Pong, 0, &[]);
|
||||
let head = format!(
|
||||
"POST /old/api/v1/up HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nContent-Type: application/octet-stream\r\nX-Up-Seq: 1\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
pong.len()
|
||||
);
|
||||
client.write_all(head.as_bytes()).await.unwrap();
|
||||
tokio::time::timeout(std::time::Duration::from_secs(2), async {
|
||||
loop {
|
||||
let body_readers = runtime
|
||||
.capacity_snapshot()
|
||||
.resources
|
||||
.into_iter()
|
||||
.find(|resource| resource.resource == "body_readers")
|
||||
.unwrap();
|
||||
if body_readers.used == 1 {
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("uplink never entered body collection");
|
||||
|
||||
let mut new_config = runtime_config_with_base([104u8; 32], WebCarrier::Https, "/new/");
|
||||
new_config.web.timeouts.long_poll_secs = 1;
|
||||
let new_generation = test_runtime_generation(2, new_config);
|
||||
active.store(Arc::clone(&new_generation));
|
||||
client.write_all(&pong).await.unwrap();
|
||||
let mut response = Vec::new();
|
||||
client.read_to_end(&mut response).await.unwrap();
|
||||
task.await.unwrap();
|
||||
assert!(response.starts_with(b"HTTP/1.1 204"));
|
||||
|
||||
assert_private_not_found(
|
||||
&request(&listener, &runtime, session_request("/old/", &session)).await,
|
||||
);
|
||||
assert!(
|
||||
request(&listener, &runtime, session_request("/new/", &session),)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
runtime.shutdown().await;
|
||||
old_generation.stop_sessions().await;
|
||||
old_generation.stop_background_tasks().await;
|
||||
new_generation.stop_sessions().await;
|
||||
new_generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn generation_transition_burst_never_authenticates_a_torn_route_identity() {
|
||||
let old_capability = [105u8; 32];
|
||||
let new_capability = [106u8; 32];
|
||||
let old_generation = test_runtime_generation(
|
||||
1,
|
||||
runtime_config_with_base(old_capability, WebCarrier::Https, "/old/"),
|
||||
);
|
||||
let new_generation = test_runtime_generation(
|
||||
2,
|
||||
runtime_config_with_base(new_capability, WebCarrier::Https, "/new/"),
|
||||
);
|
||||
let active = Arc::new(ArcSwap::from(Arc::clone(&old_generation)));
|
||||
let runtime = WebProcessRuntime::start(Arc::clone(&active));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let barrier = Arc::new(tokio::sync::Barrier::new(2));
|
||||
let swap_barrier = Arc::clone(&barrier);
|
||||
let swap_active = Arc::clone(&active);
|
||||
let swap_old = Arc::clone(&old_generation);
|
||||
let swap_new = Arc::clone(&new_generation);
|
||||
let swaps = tokio::spawn(async move {
|
||||
swap_barrier.wait().await;
|
||||
for index in 0..10_000 {
|
||||
if index % 2 == 0 {
|
||||
swap_active.store(Arc::clone(&swap_new));
|
||||
} else {
|
||||
swap_active.store(Arc::clone(&swap_old));
|
||||
}
|
||||
if index % 8 == 0 {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
}
|
||||
});
|
||||
barrier.wait().await;
|
||||
let old_encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(old_capability);
|
||||
let new_encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(new_capability);
|
||||
|
||||
for index in 0..256 {
|
||||
let target = if index % 2 == 0 {
|
||||
format!("/old/?bridge={new_encoded}")
|
||||
} else {
|
||||
format!("/new/?bridge={old_encoded}")
|
||||
};
|
||||
let response = request(&listener, &runtime, bridge_request(&target)).await;
|
||||
assert!(!response.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(
|
||||
!std::str::from_utf8(split_response(&response).1)
|
||||
.unwrap()
|
||||
.contains("bootstrap=\"")
|
||||
);
|
||||
}
|
||||
swaps.await.unwrap();
|
||||
|
||||
runtime.shutdown().await;
|
||||
old_generation.stop_sessions().await;
|
||||
old_generation.stop_background_tasks().await;
|
||||
new_generation.stop_sessions().await;
|
||||
new_generation.stop_background_tasks().await;
|
||||
}
|
||||
@@ -0,0 +1,291 @@
|
||||
use super::*;
|
||||
|
||||
const RECOVERY_TYPE: &str = "application/vnd.telemt.web-recovery+json";
|
||||
|
||||
async fn issue_bootstrap(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
path: &str,
|
||||
capability: [u8; 32],
|
||||
) -> String {
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let response = request(
|
||||
listener,
|
||||
runtime,
|
||||
bridge_request(&format!("{path}?bridge={encoded}")),
|
||||
)
|
||||
.await;
|
||||
let (headers, body) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
bridge_token(body)
|
||||
}
|
||||
|
||||
async fn create_session(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
path: &str,
|
||||
bootstrap: &str,
|
||||
) -> String {
|
||||
let hello = frame::encode(FrameType::Hello, 0, &[1]);
|
||||
let mut request_bytes = format!(
|
||||
"POST {path} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/octet-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
hello.len()
|
||||
)
|
||||
.into_bytes();
|
||||
request_bytes.extend_from_slice(&hello);
|
||||
let response = request(listener, runtime, request_bytes).await;
|
||||
let (headers, _) = split_response(&response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
response_header(headers, "x-session-token").to_string()
|
||||
}
|
||||
|
||||
fn assert_private_not_found(response: &[u8]) {
|
||||
let (headers, body) = split_response(response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(headers, "cache-control"), "no-store");
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
fn down_request(path: &str, session: &str) -> Vec<u8> {
|
||||
format!(
|
||||
"POST {path} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {session}\r\nX-Down-Cursor: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn session_token_distinguishes_the_exact_route_from_every_path_alias() {
|
||||
let capability = [81u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
config.web.timeouts.long_poll_secs = 1;
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let bootstrap = issue_bootstrap(&listener, &runtime, "/relay/", capability).await;
|
||||
let session = create_session(&listener, &runtime, "/relay/api/v1/session", &bootstrap).await;
|
||||
|
||||
let exact = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
down_request("/relay/api/v1/down", &session),
|
||||
)
|
||||
.await;
|
||||
assert!(exact.starts_with(b"HTTP/1.1 204"));
|
||||
|
||||
for alias in [
|
||||
"/api/v1/down",
|
||||
"/Relay/api/v1/down",
|
||||
"/relayx/api/v1/down",
|
||||
"/relay//api/v1/down",
|
||||
"/relay%2Fapi/v1/down",
|
||||
"/relay/api%2Fv1/down",
|
||||
"/relay/api/v1/down/",
|
||||
"/relay/api/v1/down?q=1",
|
||||
] {
|
||||
let response = request(&listener, &runtime, down_request(alias, &session)).await;
|
||||
assert_private_not_found(&response);
|
||||
}
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn inactive_capabilities_remain_decoy_and_the_active_path_is_case_sensitive() {
|
||||
let active_capability = [82u8; 32];
|
||||
let inactive_capability = [83u8; 32];
|
||||
let generation = test_runtime_generation(
|
||||
1,
|
||||
runtime_config_with_base(
|
||||
active_capability,
|
||||
WebCarrier::Https,
|
||||
"/Dobry-Cola/super_app/",
|
||||
),
|
||||
);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let active = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(active_capability);
|
||||
let inactive = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(inactive_capability);
|
||||
|
||||
let exact = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/Dobry-Cola/super_app/?bridge={active}")),
|
||||
)
|
||||
.await;
|
||||
assert!(exact.starts_with(b"HTTP/1.1 200"));
|
||||
|
||||
for path in [
|
||||
format!("/dobry-cola/super_app/?bridge={active}"),
|
||||
format!("/Dobry-Cola/super_app?bridge={active}"),
|
||||
format!("/Dobry-Cola/super_app/nested?bridge={active}"),
|
||||
] {
|
||||
assert_private_not_found(&request(&listener, &runtime, bridge_request(&path)).await);
|
||||
}
|
||||
|
||||
let decoy = request(
|
||||
&listener,
|
||||
&runtime,
|
||||
bridge_request(&format!("/Dobry-Cola/super_app/?bridge={inactive}")),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
split_response(&decoy).1,
|
||||
b"<!doctype html><title>decoy</title>"
|
||||
);
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn decoy_forwarding_preserves_every_reference_request_target() {
|
||||
let targets = [
|
||||
"/relay/",
|
||||
"/relay/whatever?q=1",
|
||||
"/relay/api/v1/session",
|
||||
"/relay",
|
||||
"/api/v1/ws",
|
||||
"/relay//api/v1/up",
|
||||
"/relay%2Fapi/v1/up",
|
||||
"/Relay/api/v1/down",
|
||||
];
|
||||
let site = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let site_addr = site.local_addr().unwrap();
|
||||
let site_task = tokio::spawn(async move {
|
||||
let mut request_targets = Vec::new();
|
||||
for _ in 0..targets.len() {
|
||||
let (mut stream, _) = site.accept().await.unwrap();
|
||||
let mut received = Vec::new();
|
||||
loop {
|
||||
let mut chunk = [0; 1024];
|
||||
let read = stream.read(&mut chunk).await.unwrap();
|
||||
if read == 0 {
|
||||
break;
|
||||
}
|
||||
received.extend_from_slice(&chunk[..read]);
|
||||
if received.windows(4).any(|window| window == b"\r\n\r\n") {
|
||||
break;
|
||||
}
|
||||
}
|
||||
let line = std::str::from_utf8(&received)
|
||||
.unwrap()
|
||||
.split("\r\n")
|
||||
.next()
|
||||
.unwrap()
|
||||
.to_string();
|
||||
request_targets.push(line);
|
||||
stream
|
||||
.write_all(
|
||||
b"HTTP/1.1 404 Not Found\r\nContent-Length: 4\r\nConnection: close\r\n\r\nsite",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
request_targets
|
||||
});
|
||||
|
||||
let capability = [84u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
let runtime_config = config.web.runtime.as_ref().unwrap();
|
||||
let mut vhosts = runtime_config.vhosts.clone();
|
||||
let previous = &runtime_config.vhosts["proxy.example.com"];
|
||||
vhosts.insert(
|
||||
"proxy.example.com".to_string(),
|
||||
Arc::new(WebRuntimeVhost {
|
||||
host: previous.host.clone(),
|
||||
base: previous.base.clone(),
|
||||
decoy_fasttrack_mode: previous.decoy_fasttrack_mode,
|
||||
decoy: WebRuntimeDecoy::HttpUpstream {
|
||||
addr: site_addr,
|
||||
authority: "decoy.internal".to_string(),
|
||||
},
|
||||
decoy_header_secs: 1,
|
||||
profiles: previous.profiles.clone(),
|
||||
capabilities: previous.capabilities.clone(),
|
||||
}),
|
||||
);
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts,
|
||||
profiles: runtime_config.profiles.clone(),
|
||||
capabilities: runtime_config.capabilities.clone(),
|
||||
}));
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
|
||||
for target in targets {
|
||||
let response = request(&listener, &runtime, bridge_request(target)).await;
|
||||
assert_eq!(split_response(&response).1, b"site");
|
||||
}
|
||||
let forwarded = site_task.await.unwrap();
|
||||
let expected = targets
|
||||
.into_iter()
|
||||
.map(|target| format!("GET {target} HTTP/1.1"))
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(forwarded, expected);
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn diagnostic_and_recovery_use_only_the_exact_prefixed_root() {
|
||||
let capability = [85u8; 32];
|
||||
let mut config = runtime_config_with_base(capability, WebCarrier::Https, "/relay/");
|
||||
config.web.debug.enabled = true;
|
||||
config.web.debug.sideband = true;
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
let bootstrap = issue_bootstrap(&listener, &runtime, "/relay/", capability).await;
|
||||
|
||||
let body = br#"{"v":1,"event":"runtime_started"}"#;
|
||||
let mut diagnostic = format!(
|
||||
"POST /relay/api/v1/diagnostic HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAuthorization: Bearer {bootstrap}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
|
||||
body.len()
|
||||
)
|
||||
.into_bytes();
|
||||
diagnostic.extend_from_slice(body);
|
||||
assert!(
|
||||
request(&listener, &runtime, diagnostic)
|
||||
.await
|
||||
.starts_with(b"HTTP/1.1 204")
|
||||
);
|
||||
|
||||
let root_diagnostic = format!(
|
||||
"POST /api/v1/diagnostic HTTP/1.1\r\nHost: proxy.example.com\r\nAuthorization: Bearer {bootstrap}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes();
|
||||
assert_private_not_found(&request(&listener, &runtime, root_diagnostic).await);
|
||||
|
||||
let session = create_session(&listener, &runtime, "/relay/api/v1/session", &bootstrap).await;
|
||||
let recovery_request = |path: &str| {
|
||||
format!(
|
||||
"GET {path}?bridge={encoded} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nAccept: {RECOVERY_TYPE}\r\nAuthorization: Bearer {session}\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
.into_bytes()
|
||||
};
|
||||
assert_private_not_found(&request(&listener, &runtime, recovery_request("/")).await);
|
||||
let recovered = request(&listener, &runtime, recovery_request("/relay/")).await;
|
||||
let (headers, body) = split_response(&recovered);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert_eq!(response_header(headers, "content-type"), RECOVERY_TYPE);
|
||||
let document: serde_json::Value = serde_json::from_slice(body).unwrap();
|
||||
let recovery_bootstrap = document["bootstrap"].as_str().unwrap();
|
||||
create_session(
|
||||
&listener,
|
||||
&runtime,
|
||||
"/relay/api/v1/session",
|
||||
recovery_bootstrap,
|
||||
)
|
||||
.await;
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
@@ -32,28 +32,35 @@ pub(super) fn bridge_candidate(query: Option<&str>) -> BridgeCandidate {
|
||||
let Some(value) = query.and_then(|query| query.strip_prefix("bridge=")) else {
|
||||
return BridgeCandidate::NonCanonical;
|
||||
};
|
||||
canonical_credential(value.as_bytes())
|
||||
.map(BridgeCandidate::Canonical)
|
||||
.unwrap_or(BridgeCandidate::NonCanonical)
|
||||
}
|
||||
|
||||
/// Decodes one exact canonical 32-byte base64url credential.
|
||||
pub(super) fn canonical_credential(value: &[u8]) -> Option<[u8; 32]> {
|
||||
if value.len() != 43 {
|
||||
return BridgeCandidate::NonCanonical;
|
||||
return None;
|
||||
}
|
||||
let mut decoded = [0u8; 32];
|
||||
let Ok(decoded_len) =
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.decode_slice(value, &mut decoded)
|
||||
else {
|
||||
return BridgeCandidate::NonCanonical;
|
||||
return None;
|
||||
};
|
||||
let mut canonical = [0u8; 43];
|
||||
let Ok(encoded_len) =
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode_slice(decoded, &mut canonical)
|
||||
else {
|
||||
return BridgeCandidate::NonCanonical;
|
||||
return None;
|
||||
};
|
||||
if decoded_len != decoded.len()
|
||||
|| encoded_len != canonical.len()
|
||||
|| !bool::from(canonical.ct_eq(value.as_bytes()))
|
||||
|| !bool::from(canonical.ct_eq(value))
|
||||
{
|
||||
return BridgeCandidate::NonCanonical;
|
||||
return None;
|
||||
}
|
||||
BridgeCandidate::Canonical(decoded)
|
||||
Some(decoded)
|
||||
}
|
||||
|
||||
/// Internal result of one complete capability-table scan.
|
||||
|
||||
@@ -29,6 +29,9 @@ where
|
||||
B: hyper::body::Body<Data = Bytes> + Send + 'static,
|
||||
B::Error: Error + Send + Sync + 'static,
|
||||
{
|
||||
if super::secrets::has_internal_credential(&request) {
|
||||
return super::response::private_not_found();
|
||||
}
|
||||
super::set_trace_route(&request, crate::web::trace::TraceRoute::Decoy);
|
||||
if sanitize_transport {
|
||||
sanitize_transport_request(&mut request);
|
||||
|
||||
@@ -153,12 +153,9 @@ async fn pause_preserves_decoy_retry_and_exact_session_replay() {
|
||||
.into_bytes();
|
||||
let decoy = request(&listener, &runtime, decoy).await;
|
||||
let (decoy_headers, decoy_body) = split_response(&decoy);
|
||||
assert!(decoy_headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(
|
||||
!decoy_body
|
||||
.windows(11)
|
||||
.any(|window| window == b"bootstrap=\"")
|
||||
);
|
||||
assert!(decoy_headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(response_header(decoy_headers, "cache-control"), "no-store");
|
||||
assert_eq!(decoy_body, b"not found\n");
|
||||
|
||||
runtime.resume_operator().await.unwrap();
|
||||
let created = request(&listener, &runtime, create_request(&bootstrap, &hello)).await;
|
||||
|
||||
@@ -195,12 +195,12 @@ async fn malformed_or_over_capacity_recovery_is_indistinguishable_from_decoy() {
|
||||
|
||||
let malformed = recover(&listener, &runtime, &encoded, "Bearer malformed").await;
|
||||
let (malformed_headers, malformed_body) = split_response(&malformed);
|
||||
assert!(malformed_headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(malformed_headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(
|
||||
response_header(malformed_headers, "cache-control"),
|
||||
"no-store"
|
||||
);
|
||||
assert_eq!(malformed_body, b"<!doctype html><title>decoy</title>");
|
||||
assert_eq!(malformed_body, b"not found\n");
|
||||
|
||||
let invalid_capability = recover(
|
||||
&listener,
|
||||
@@ -209,7 +209,13 @@ async fn malformed_or_over_capacity_recovery_is_indistinguishable_from_decoy() {
|
||||
&format!("Bearer {}", "U".repeat(43)),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(invalid_capability, malformed);
|
||||
let (invalid_headers, invalid_body) = split_response(&invalid_capability);
|
||||
assert!(invalid_headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert_eq!(
|
||||
response_header(invalid_headers, "cache-control"),
|
||||
"no-store"
|
||||
);
|
||||
assert_eq!(invalid_body, b"<!doctype html><title>decoy</title>");
|
||||
|
||||
let malformed_accept = request(
|
||||
&listener,
|
||||
@@ -231,12 +237,12 @@ async fn malformed_or_over_capacity_recovery_is_indistinguishable_from_decoy() {
|
||||
)
|
||||
.await;
|
||||
let (capacity_headers, capacity_body) = split_response(&over_capacity);
|
||||
assert!(capacity_headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(capacity_headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(
|
||||
response_header(capacity_headers, "cache-control"),
|
||||
"no-store"
|
||||
);
|
||||
assert_eq!(capacity_body, b"<!doctype html><title>decoy</title>");
|
||||
assert_eq!(capacity_body, b"not found\n");
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
|
||||
@@ -62,6 +62,15 @@ pub(super) fn generic_not_found() -> HttpResponse {
|
||||
full_response(StatusCode::NOT_FOUND, Bytes::from_static(b"not found\n"))
|
||||
}
|
||||
|
||||
/// Builds a non-cacheable local rejection for misplaced internal credentials.
|
||||
pub(super) fn private_not_found() -> HttpResponse {
|
||||
let mut response = generic_not_found();
|
||||
response
|
||||
.headers_mut()
|
||||
.insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store"));
|
||||
response
|
||||
}
|
||||
|
||||
/// Builds one in-memory response with an exact content length.
|
||||
pub(super) fn full_response(status: StatusCode, body: Bytes) -> HttpResponse {
|
||||
let length = body.len();
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
use super::*;
|
||||
|
||||
/// Builds a static-decoy runtime with one explicit WEB endpoint base.
|
||||
pub(in crate::web::http) fn runtime_config_with_base(
|
||||
capability: [u8; 32],
|
||||
carrier: WebCarrier,
|
||||
base: &str,
|
||||
) -> ProxyConfig {
|
||||
runtime_config_with_carriers_and_deadlines(
|
||||
capability,
|
||||
carrier,
|
||||
false,
|
||||
true,
|
||||
Arc::from([carrier]),
|
||||
TEST_CARRIER_DEADLINES_SECS,
|
||||
base,
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
use hyper::Request;
|
||||
|
||||
use super::capability::{canonical_credential, scan_capabilities};
|
||||
use crate::config::WebRuntimeConfig;
|
||||
use crate::web::manager::WebProcessRuntime;
|
||||
|
||||
/// Request extension proving that metadata contains an authentic internal credential.
|
||||
#[derive(Clone, Copy)]
|
||||
struct InternalCredential;
|
||||
|
||||
/// Marks requests whose metadata contains a capability or process token.
|
||||
pub(super) fn mark_internal_credential<B>(
|
||||
request: &mut Request<B>,
|
||||
config: &WebRuntimeConfig,
|
||||
runtime: &WebProcessRuntime,
|
||||
) {
|
||||
let uri = request.uri();
|
||||
let uri_contains = uri
|
||||
.authority()
|
||||
.is_some_and(|authority| contains_secret(authority.as_str().as_bytes(), config, runtime))
|
||||
|| uri
|
||||
.path_and_query()
|
||||
.is_some_and(|path| contains_secret(path.as_str().as_bytes(), config, runtime));
|
||||
let headers_contain = request.headers().iter().any(|(name, value)| {
|
||||
contains_secret(name.as_str().as_bytes(), config, runtime)
|
||||
|| contains_secret(value.as_bytes(), config, runtime)
|
||||
});
|
||||
if uri_contains || headers_contain {
|
||||
request.extensions_mut().insert(InternalCredential);
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns whether request metadata was authenticated before routing.
|
||||
pub(super) fn has_internal_credential<B>(request: &Request<B>) -> bool {
|
||||
request.extensions().get::<InternalCredential>().is_some()
|
||||
}
|
||||
|
||||
fn contains_secret(text: &[u8], config: &WebRuntimeConfig, runtime: &WebProcessRuntime) -> bool {
|
||||
let mut window = [0; 43];
|
||||
let mut run_len = 0usize;
|
||||
let mut offset = 0usize;
|
||||
while offset < text.len() {
|
||||
let (byte, consumed) = if text[offset] == b'%' && offset + 2 < text.len() {
|
||||
match (hex_value(text[offset + 1]), hex_value(text[offset + 2])) {
|
||||
(Some(high), Some(low)) => ((high << 4) | low, 3),
|
||||
_ => (text[offset], 1),
|
||||
}
|
||||
} else {
|
||||
(text[offset], 1)
|
||||
};
|
||||
offset += consumed;
|
||||
if base64url_byte(byte) {
|
||||
if run_len < window.len() {
|
||||
window[run_len] = byte;
|
||||
run_len += 1;
|
||||
} else {
|
||||
window.copy_within(1.., 0);
|
||||
window[42] = byte;
|
||||
}
|
||||
if run_len >= window.len()
|
||||
&& canonical_credential(&window).is_some_and(|candidate| {
|
||||
bool::from(scan_capabilities(&config.capabilities, &candidate).matched)
|
||||
|| runtime.authentic_token(&candidate)
|
||||
})
|
||||
{
|
||||
return true;
|
||||
}
|
||||
} else {
|
||||
run_len = 0;
|
||||
}
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
fn base64url_byte(value: u8) -> bool {
|
||||
value.is_ascii_alphanumeric() || matches!(value, b'-' | b'_')
|
||||
}
|
||||
|
||||
fn hex_value(value: u8) -> Option<u8> {
|
||||
match value {
|
||||
b'0'..=b'9' => Some(value - b'0'),
|
||||
b'a'..=b'f' => Some(value - b'a' + 10),
|
||||
b'A'..=b'F' => Some(value - b'A' + 10),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "secrets/tests.rs"]
|
||||
mod tests;
|
||||
@@ -0,0 +1,76 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use arc_swap::ArcSwap;
|
||||
use base64::Engine as _;
|
||||
|
||||
use super::*;
|
||||
use crate::config::WebCarrier;
|
||||
use crate::maestro::generation::test_runtime_generation;
|
||||
use crate::web::http::tests::runtime_config_with_base;
|
||||
|
||||
fn percent_encode_byte(value: u8, lowercase: bool) -> String {
|
||||
if lowercase {
|
||||
format!("%{value:02x}")
|
||||
} else {
|
||||
format!("%{value:02X}")
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanner_recognizes_exact_credentials_across_encoding_boundaries() {
|
||||
let capability = [71u8; 32];
|
||||
let generation = test_runtime_generation(
|
||||
1,
|
||||
runtime_config_with_base(capability, WebCarrier::Https, "/relay/"),
|
||||
);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let generation_config = generation.config();
|
||||
let config = generation_config.web.runtime.as_ref().unwrap();
|
||||
let encoded = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(capability);
|
||||
|
||||
assert!(contains_secret(encoded.as_bytes(), config, &runtime));
|
||||
assert!(contains_secret(
|
||||
format!("prefix{encoded}suffix").as_bytes(),
|
||||
config,
|
||||
&runtime,
|
||||
));
|
||||
for index in 0..encoded.len() {
|
||||
let mut escaped = String::with_capacity(encoded.len() + 2);
|
||||
escaped.push_str(&encoded[..index]);
|
||||
escaped.push_str(&percent_encode_byte(
|
||||
encoded.as_bytes()[index],
|
||||
index % 2 == 0,
|
||||
));
|
||||
escaped.push_str(&encoded[index + 1..]);
|
||||
assert!(
|
||||
contains_secret(escaped.as_bytes(), config, &runtime),
|
||||
"percent-encoded byte {index} was not recognized"
|
||||
);
|
||||
}
|
||||
|
||||
let profile = config.profiles[0].clone();
|
||||
let bootstrap = runtime
|
||||
.issue_bootstrap(profile, "192.0.2.10".parse().unwrap())
|
||||
.unwrap()
|
||||
.token;
|
||||
let escaped_bootstrap = bootstrap
|
||||
.bytes()
|
||||
.enumerate()
|
||||
.map(|(index, byte)| percent_encode_byte(byte, index % 2 == 0))
|
||||
.collect::<String>();
|
||||
assert!(contains_secret(bootstrap.as_bytes(), config, &runtime));
|
||||
assert!(contains_secret(
|
||||
escaped_bootstrap.as_bytes(),
|
||||
config,
|
||||
&runtime,
|
||||
));
|
||||
|
||||
let split = format!("{}%zz{}", &encoded[..20], &encoded[20..]);
|
||||
assert!(!contains_secret(split.as_bytes(), config, &runtime));
|
||||
let inactive = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode([72u8; 32]);
|
||||
assert!(!contains_secret(inactive.as_bytes(), config, &runtime));
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
+15
-2
@@ -49,11 +49,18 @@ mod recovery_tests;
|
||||
// Decoy fast-track routing and telemetry remain isolated from carrier protocol scenarios.
|
||||
#[path = "decoy_fasttrack_tests.rs"]
|
||||
mod decoy_fasttrack_tests;
|
||||
// Base-path routing and credential containment share reference-contract coverage.
|
||||
#[path = "base_path_tests.rs"]
|
||||
mod base_path_tests;
|
||||
// Raw response parsing helpers are shared by the HTTP integration test modules.
|
||||
#[path = "response_test_support.rs"]
|
||||
mod response_test_support;
|
||||
// Alternate runtime fixtures remain separate from the main integration scenarios.
|
||||
#[path = "runtime_test_support.rs"]
|
||||
mod runtime_test_support;
|
||||
|
||||
pub(super) use response_test_support::{response_header, split_response};
|
||||
pub(super) use runtime_test_support::runtime_config_with_base;
|
||||
|
||||
const TEST_CARRIER_DEADLINES_SECS: [u64; 4] = [3, 5, 8, 12];
|
||||
|
||||
@@ -86,6 +93,7 @@ fn runtime_config_with_carriers(
|
||||
carrier_learning,
|
||||
carriers,
|
||||
TEST_CARRIER_DEADLINES_SECS,
|
||||
"/",
|
||||
)
|
||||
}
|
||||
|
||||
@@ -104,6 +112,7 @@ pub(super) fn negotiation_runtime_config_with_deadlines(
|
||||
carrier_learning,
|
||||
carriers,
|
||||
carrier_negotiation_deadlines_secs,
|
||||
"/",
|
||||
)
|
||||
}
|
||||
|
||||
@@ -114,6 +123,7 @@ fn runtime_config_with_carriers_and_deadlines(
|
||||
carrier_learning: bool,
|
||||
carriers: Arc<[WebCarrier]>,
|
||||
carrier_negotiation_deadlines_secs: [u64; 4],
|
||||
base: &str,
|
||||
) -> ProxyConfig {
|
||||
let profile = Arc::new(WebRuntimeProfile {
|
||||
host: "proxy.example.com".to_string(),
|
||||
@@ -147,6 +157,7 @@ fn runtime_config_with_carriers_and_deadlines(
|
||||
});
|
||||
let vhost = Arc::new(WebRuntimeVhost {
|
||||
host: "proxy.example.com".to_string(),
|
||||
base: base.to_string(),
|
||||
decoy_fasttrack_mode: WebDecoyFastTrackMode::Off,
|
||||
decoy: WebRuntimeDecoy::StaticDirectory(Arc::clone(&site)),
|
||||
decoy_header_secs: 1,
|
||||
@@ -159,6 +170,7 @@ fn runtime_config_with_carriers_and_deadlines(
|
||||
"other.example.com".to_string(),
|
||||
Arc::new(WebRuntimeVhost {
|
||||
host: "other.example.com".to_string(),
|
||||
base: "/".to_string(),
|
||||
decoy_fasttrack_mode: WebDecoyFastTrackMode::Off,
|
||||
decoy: WebRuntimeDecoy::StaticDirectory(site),
|
||||
decoy_header_secs: 1,
|
||||
@@ -181,6 +193,7 @@ fn runtime_config_with_carriers_and_deadlines(
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts,
|
||||
profiles: vec![profile],
|
||||
capabilities: vec![capability].into_boxed_slice(),
|
||||
}));
|
||||
config
|
||||
}
|
||||
@@ -328,12 +341,12 @@ async fn rejected_bridge_bootstrap_falls_back_to_uncacheable_static_index() {
|
||||
|
||||
let fallback_response = request(&listener, &runtime, bridge_request()).await;
|
||||
let (fallback_headers, fallback_body) = split_response(&fallback_response);
|
||||
assert!(fallback_headers.starts_with(b"HTTP/1.1 200"));
|
||||
assert!(fallback_headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(
|
||||
response_header(fallback_headers, "cache-control"),
|
||||
"no-store"
|
||||
);
|
||||
assert_eq!(fallback_body, b"<!doctype html><title>decoy</title>");
|
||||
assert_eq!(fallback_body, b"not found\n");
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
|
||||
@@ -14,11 +14,16 @@ use tokio_util::sync::CancellationToken;
|
||||
|
||||
use crate::maestro::generation::{RuntimeGeneration, test_runtime_generation};
|
||||
use crate::web::frame::{self, FrameType};
|
||||
use crate::web::http::tests::{negotiation_runtime_config, runtime_config};
|
||||
use crate::web::http::tests::{
|
||||
negotiation_runtime_config, runtime_config, runtime_config_with_base,
|
||||
};
|
||||
use crate::web::manager::{
|
||||
CarrierCapabilities, CarrierClientClass, CarrierFailure, CarrierRequest, WebProcessRuntime,
|
||||
};
|
||||
|
||||
#[path = "tests/base_path.rs"]
|
||||
mod base_path;
|
||||
|
||||
fn request(protocol: &str) -> Request<()> {
|
||||
Request::builder()
|
||||
.method(Method::GET)
|
||||
@@ -210,6 +215,15 @@ async fn upgrade(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
protocol: &str,
|
||||
) -> WebSocketStream<TcpStream> {
|
||||
upgrade_at(listener, runtime, "/api/v1/ws", protocol).await
|
||||
}
|
||||
|
||||
async fn upgrade_at(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
path: &str,
|
||||
protocol: &str,
|
||||
) -> WebSocketStream<TcpStream> {
|
||||
let addr = listener.local_addr().unwrap();
|
||||
let (accepted, client) = tokio::join!(listener.accept(), TcpStream::connect(addr));
|
||||
@@ -226,7 +240,7 @@ async fn upgrade(
|
||||
permit,
|
||||
));
|
||||
let request = format!(
|
||||
"GET /api/v1/ws HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nConnection: Upgrade\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Protocol: {protocol}\r\nCookie: browser-state=allowed\r\n\r\n"
|
||||
"GET {path} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nConnection: Upgrade\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Protocol: {protocol}\r\nCookie: browser-state=allowed\r\n\r\n"
|
||||
);
|
||||
client.write_all(request.as_bytes()).await.unwrap();
|
||||
let mut response = Vec::new();
|
||||
@@ -244,6 +258,38 @@ async fn upgrade(
|
||||
WebSocketStream::from_raw_socket(client, Role::Client, None).await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn prefixed_websocket_route_upgrades_at_the_exact_base() {
|
||||
let live = live_runtime_from_config(
|
||||
runtime_config_with_base([31; 32], WebCarrier::Websocket, "/relay/nested/"),
|
||||
1,
|
||||
);
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let (session, _) = create_session(&live.runtime);
|
||||
let protocol = format!("tproxy-v1.{session}");
|
||||
let mut socket = upgrade_at(
|
||||
&listener,
|
||||
&live.runtime,
|
||||
"/relay/nested/api/v1/ws",
|
||||
&protocol,
|
||||
)
|
||||
.await;
|
||||
|
||||
socket
|
||||
.send(Message::Ping(Bytes::from_static(b"prefixed")))
|
||||
.await
|
||||
.unwrap();
|
||||
let response = tokio::time::timeout(Duration::from_secs(2), socket.next())
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(response, Message::Pong(Bytes::from_static(b"prefixed")));
|
||||
|
||||
let _ = socket.close(None).await;
|
||||
live.shutdown().await;
|
||||
}
|
||||
|
||||
fn masked_message(opcode: u8, payload: &[u8], mask: [u8; 4]) -> Vec<u8> {
|
||||
masked_frame(true, opcode, payload, mask)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,107 @@
|
||||
use super::*;
|
||||
|
||||
async fn rejected_upgrade(
|
||||
listener: &TcpListener,
|
||||
runtime: &Arc<WebProcessRuntime>,
|
||||
path: &str,
|
||||
protocol: &str,
|
||||
) -> Vec<u8> {
|
||||
let address = listener.local_addr().unwrap();
|
||||
let (accepted, client) = tokio::join!(listener.accept(), TcpStream::connect(address));
|
||||
let (server, peer) = accepted.unwrap();
|
||||
let mut client = client.unwrap();
|
||||
let permit = runtime.try_http_connection().unwrap();
|
||||
let task = tokio::spawn(super::super::super::serve_connection(
|
||||
server,
|
||||
peer,
|
||||
WebClientIpSource::XForwardedFor,
|
||||
Arc::from(["127.0.0.1/32".parse().unwrap()]),
|
||||
Arc::clone(runtime),
|
||||
CancellationToken::new(),
|
||||
permit,
|
||||
));
|
||||
let request = format!(
|
||||
"GET {path} HTTP/1.1\r\nHost: proxy.example.com\r\nX-Forwarded-For: 192.0.2.10\r\nConnection: close, Upgrade\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Protocol: {protocol}\r\n\r\n"
|
||||
);
|
||||
client.write_all(request.as_bytes()).await.unwrap();
|
||||
let mut response = Vec::new();
|
||||
client.read_to_end(&mut response).await.unwrap();
|
||||
task.await.unwrap();
|
||||
response
|
||||
}
|
||||
|
||||
fn assert_private_not_found(response: &[u8]) {
|
||||
let (headers, body) = crate::web::http::tests::split_response(response);
|
||||
assert!(headers.starts_with(b"HTTP/1.1 404"));
|
||||
assert_eq!(
|
||||
crate::web::http::tests::response_header(headers, "cache-control"),
|
||||
"no-store"
|
||||
);
|
||||
assert_eq!(body, b"not found\n");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn authentic_websocket_protocol_rejects_every_base_path_alias_locally() {
|
||||
let config = runtime_config_with_base([32; 32], WebCarrier::Websocket, "/relay/");
|
||||
let generation = test_runtime_generation(1, config);
|
||||
let runtime = WebProcessRuntime::start(Arc::new(ArcSwap::from(Arc::clone(&generation))));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let (session, _) = create_session(&runtime);
|
||||
let protocol = format!("tproxy-v1.{session}");
|
||||
|
||||
for path in [
|
||||
"/api/v1/ws",
|
||||
"/Relay/api/v1/ws",
|
||||
"/relayx/api/v1/ws",
|
||||
"/relay//api/v1/ws",
|
||||
"/relay%2Fapi/v1/ws",
|
||||
"/relay/api%2Fv1/ws",
|
||||
"/relay/api/v1/ws/",
|
||||
] {
|
||||
assert_private_not_found(&rejected_upgrade(&listener, &runtime, path, &protocol).await);
|
||||
}
|
||||
|
||||
runtime.shutdown().await;
|
||||
generation.stop_sessions().await;
|
||||
generation.stop_background_tasks().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upgraded_socket_survives_base_path_generation_swap() {
|
||||
let old_generation = test_runtime_generation(
|
||||
1,
|
||||
runtime_config_with_base([33; 32], WebCarrier::Websocket, "/old/"),
|
||||
);
|
||||
let active = Arc::new(ArcSwap::from(Arc::clone(&old_generation)));
|
||||
let runtime = WebProcessRuntime::start(Arc::clone(&active));
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let (session, _) = create_session(&runtime);
|
||||
let protocol = format!("tproxy-v1.{session}");
|
||||
let mut socket = upgrade_at(&listener, &runtime, "/old/api/v1/ws", &protocol).await;
|
||||
|
||||
let new_generation = test_runtime_generation(
|
||||
2,
|
||||
runtime_config_with_base([34; 32], WebCarrier::Websocket, "/new/"),
|
||||
);
|
||||
active.store(Arc::clone(&new_generation));
|
||||
socket
|
||||
.send(Message::Ping(Bytes::from_static(b"after-reload")))
|
||||
.await
|
||||
.unwrap();
|
||||
let response = tokio::time::timeout(Duration::from_secs(2), socket.next())
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(response, Message::Pong(Bytes::from_static(b"after-reload")));
|
||||
assert_private_not_found(
|
||||
&rejected_upgrade(&listener, &runtime, "/old/api/v1/ws", &protocol).await,
|
||||
);
|
||||
|
||||
let _ = socket.close(None).await;
|
||||
runtime.shutdown().await;
|
||||
old_generation.stop_sessions().await;
|
||||
old_generation.stop_background_tasks().await;
|
||||
new_generation.stop_sessions().await;
|
||||
new_generation.stop_background_tasks().await;
|
||||
}
|
||||
@@ -5,12 +5,17 @@ use std::sync::atomic::AtomicU64;
|
||||
use std::time::Duration;
|
||||
|
||||
use arc_swap::ArcSwap;
|
||||
use hmac::{Hmac, Mac};
|
||||
use parking_lot::Mutex;
|
||||
use sha2::Sha256;
|
||||
use subtle::ConstantTimeEq;
|
||||
use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore, TryAcquireError};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tokio_util::task::TaskTracker;
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
use crate::config::{WebCarrier, WebLimitsConfig};
|
||||
use crate::crypto::SecureRandom;
|
||||
use crate::maestro::generation::RuntimeGeneration;
|
||||
use crate::web::telemetry::{WebRejectionReason, WebTelemetry};
|
||||
use crate::web::trace::WebTraceStore;
|
||||
@@ -67,6 +72,57 @@ pub(crate) use websocket::{WebSocketConnection, WebSocketKind};
|
||||
|
||||
const TOKEN_BYTES: usize = 32;
|
||||
const CLEANUP_INTERVAL: Duration = Duration::from_secs(1);
|
||||
const TOKEN_NONCE_BYTES: usize = 16;
|
||||
const BOOTSTRAP_TOKEN_CONTEXT: &[u8] = b"telemt-web-bootstrap-token-v1\0";
|
||||
const SESSION_TOKEN_CONTEXT: &[u8] = b"telemt-web-session-token-v1\0";
|
||||
|
||||
/// Distinguishes process-authenticated WEB credential domains.
|
||||
#[derive(Clone, Copy)]
|
||||
enum TokenKind {
|
||||
Bootstrap,
|
||||
Session,
|
||||
}
|
||||
|
||||
struct TokenAuthenticator {
|
||||
key: Zeroizing<[u8; 32]>,
|
||||
}
|
||||
|
||||
impl TokenAuthenticator {
|
||||
fn new(rng: &SecureRandom) -> Self {
|
||||
let mut key = Zeroizing::new([0; 32]);
|
||||
rng.fill(key.as_mut());
|
||||
Self { key }
|
||||
}
|
||||
|
||||
fn issue(&self, kind: TokenKind, nonce: [u8; TOKEN_NONCE_BYTES]) -> [u8; TOKEN_BYTES] {
|
||||
let mut token = [0; TOKEN_BYTES];
|
||||
token[..TOKEN_NONCE_BYTES].copy_from_slice(&nonce);
|
||||
let tag = self.tag(kind, &nonce);
|
||||
token[TOKEN_NONCE_BYTES..].copy_from_slice(&tag[..TOKEN_NONCE_BYTES]);
|
||||
token
|
||||
}
|
||||
|
||||
fn authentic(&self, token: &[u8; TOKEN_BYTES]) -> bool {
|
||||
let nonce = &token[..TOKEN_NONCE_BYTES];
|
||||
let tag = &token[TOKEN_NONCE_BYTES..];
|
||||
let bootstrap = self.tag(TokenKind::Bootstrap, nonce);
|
||||
let session = self.tag(TokenKind::Session, nonce);
|
||||
bool::from(
|
||||
bootstrap[..TOKEN_NONCE_BYTES].ct_eq(tag) | session[..TOKEN_NONCE_BYTES].ct_eq(tag),
|
||||
)
|
||||
}
|
||||
|
||||
fn tag(&self, kind: TokenKind, nonce: &[u8]) -> [u8; 32] {
|
||||
let mut mac = Hmac::<Sha256>::new_from_slice(self.key.as_ref())
|
||||
.expect("HMAC accepts every WEB token key length");
|
||||
mac.update(match kind {
|
||||
TokenKind::Bootstrap => BOOTSTRAP_TOKEN_CONTEXT,
|
||||
TokenKind::Session => SESSION_TOKEN_CONTEXT,
|
||||
});
|
||||
mac.update(nonce);
|
||||
mac.finalize().into_bytes().into()
|
||||
}
|
||||
}
|
||||
|
||||
/// Stable hash key used for bootstrap and session credentials.
|
||||
pub(crate) type TokenHash = [u8; TOKEN_BYTES];
|
||||
@@ -162,6 +218,7 @@ pub(crate) struct WebProcessRuntime {
|
||||
runtime_instance: Arc<str>,
|
||||
active_runtime: Arc<ArcSwap<RuntimeGeneration>>,
|
||||
trace: Arc<WebTraceStore>,
|
||||
token_authenticator: TokenAuthenticator,
|
||||
limits: WebLimitsConfig,
|
||||
state: Mutex<ManagerState>,
|
||||
stream_admission: Mutex<StreamAdmissionState>,
|
||||
@@ -205,6 +262,7 @@ impl WebProcessRuntime {
|
||||
) -> Arc<Self> {
|
||||
let initial_generation = active_runtime.load_full();
|
||||
let config = initial_generation.config();
|
||||
let token_authenticator = TokenAuthenticator::new(&initial_generation.rng);
|
||||
trace.apply_policy(initial_generation.id, &config.web.debug);
|
||||
let limits = config.web.limits.clone();
|
||||
let learning_capacity = limits.max_carrier_learning_entries;
|
||||
@@ -231,6 +289,7 @@ impl WebProcessRuntime {
|
||||
runtime_instance,
|
||||
active_runtime,
|
||||
trace,
|
||||
token_authenticator,
|
||||
http_connections: Arc::new(Semaphore::new(limits.max_http_connections)),
|
||||
http_overload_connections: Arc::new(Semaphore::new(
|
||||
limits.max_http_overload_connections,
|
||||
@@ -302,6 +361,11 @@ impl WebProcessRuntime {
|
||||
&self.telemetry
|
||||
}
|
||||
|
||||
/// Returns whether one canonical raw credential was minted by this process.
|
||||
pub(crate) fn authentic_token(&self, token: &[u8; TOKEN_BYTES]) -> bool {
|
||||
self.token_authenticator.authentic(token)
|
||||
}
|
||||
|
||||
/// Returns whether terminal process shutdown has started.
|
||||
pub(crate) fn is_shutdown(&self) -> bool {
|
||||
self.shutdown.is_cancelled()
|
||||
@@ -448,3 +512,7 @@ impl WebProcessRuntime {
|
||||
self.telemetry.record_limit_hit();
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "manager/token_authenticator_tests.rs"]
|
||||
mod token_authenticator_tests;
|
||||
|
||||
@@ -9,7 +9,7 @@ use super::state::{
|
||||
Bootstrap, CarrierChainPhase, allow_rate, evict_oldest_unused_bootstrap, matching_profile,
|
||||
new_unique_token, profile_key, remove_expired_locked,
|
||||
};
|
||||
use super::{BootstrapResult, ManagerError, TOKEN_BYTES, TokenHash, WebProcessRuntime};
|
||||
use super::{BootstrapResult, ManagerError, TOKEN_BYTES, TokenHash, TokenKind, WebProcessRuntime};
|
||||
use crate::config::WebRuntimeProfile;
|
||||
use crate::maestro::generation::RuntimeGeneration;
|
||||
use crate::web::session::{SessionCloseReason, WebSession};
|
||||
@@ -123,7 +123,12 @@ impl WebProcessRuntime {
|
||||
.record_rejection(WebRejectionReason::BootstrapCapacity);
|
||||
return Err(ManagerError::Limit);
|
||||
}
|
||||
let Some((token, hash)) = new_unique_token(generation, &state) else {
|
||||
let Some((token, hash)) = new_unique_token(
|
||||
generation,
|
||||
&state,
|
||||
&self.token_authenticator,
|
||||
TokenKind::Bootstrap,
|
||||
) else {
|
||||
self.record_limit_hit();
|
||||
self.telemetry
|
||||
.record_rejection(WebRejectionReason::BootstrapCapacity);
|
||||
|
||||
@@ -13,7 +13,7 @@ use super::state::{
|
||||
profile_key, remember_closed_token_locked, remove_expired_locked,
|
||||
};
|
||||
use super::{
|
||||
CarrierLearningContext, CarrierRequest, CreateResult, ManagerError, TokenHash,
|
||||
CarrierLearningContext, CarrierRequest, CreateResult, ManagerError, TokenHash, TokenKind,
|
||||
WebProcessRuntime,
|
||||
};
|
||||
use crate::config::{WebCarrier, WebRuntimeProfile};
|
||||
@@ -316,7 +316,12 @@ impl WebProcessRuntime {
|
||||
if !admit_initial(self, &mut state, now, client_ip, profile_key, &profile) {
|
||||
return Err(ManagerError::Limit);
|
||||
}
|
||||
let Some((session_token, session_hash)) = new_unique_token(&generation, &state) else {
|
||||
let Some((session_token, session_hash)) = new_unique_token(
|
||||
&generation,
|
||||
&state,
|
||||
&self.token_authenticator,
|
||||
TokenKind::Session,
|
||||
) else {
|
||||
self.record_limit_hit();
|
||||
self.telemetry
|
||||
.record_rejection(crate::web::telemetry::WebRejectionReason::SessionCapacity);
|
||||
|
||||
@@ -60,7 +60,12 @@ impl WebProcessRuntime {
|
||||
self.cancel_replacement(bootstrap_hash, &replacement.old_session);
|
||||
return Err(ManagerError::Closed);
|
||||
};
|
||||
let Some((session_token, session_hash)) = new_unique_token(&generation, &state) else {
|
||||
let Some((session_token, session_hash)) = new_unique_token(
|
||||
&generation,
|
||||
&state,
|
||||
&self.token_authenticator,
|
||||
TokenKind::Session,
|
||||
) else {
|
||||
self.record_limit_hit();
|
||||
self.telemetry
|
||||
.record_rejection(crate::web::telemetry::WebRejectionReason::SessionCapacity);
|
||||
|
||||
@@ -7,7 +7,7 @@ use base64::Engine as _;
|
||||
use sha2::{Digest, Sha256};
|
||||
use zeroize::Zeroizing;
|
||||
|
||||
use super::{CarrierRequest, ProfileKey, TOKEN_BYTES, TokenHash};
|
||||
use super::{CarrierRequest, ProfileKey, TokenAuthenticator, TokenHash, TokenKind};
|
||||
use crate::config::{WebCarrier, WebRuntimeConfig, WebRuntimeProfile, WebTimeoutsConfig};
|
||||
use crate::maestro::generation::RuntimeGeneration;
|
||||
use crate::proxy::user_admission::UserSessionRegistration;
|
||||
@@ -233,10 +233,13 @@ impl Default for ManagerState {
|
||||
pub(super) fn new_unique_token(
|
||||
generation: &RuntimeGeneration,
|
||||
state: &ManagerState,
|
||||
authenticator: &TokenAuthenticator,
|
||||
kind: TokenKind,
|
||||
) -> Option<(String, TokenHash)> {
|
||||
for _ in 0..8 {
|
||||
let mut raw = [0u8; TOKEN_BYTES];
|
||||
generation.rng.fill(&mut raw);
|
||||
let mut nonce = [0u8; 16];
|
||||
generation.rng.fill(&mut nonce);
|
||||
let raw = authenticator.issue(kind, nonce);
|
||||
let token = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(raw);
|
||||
let hash = Sha256::digest(raw).into();
|
||||
if !state.bootstraps.contains_key(&hash)
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
use super::*;
|
||||
|
||||
fn authenticator(key: u8) -> TokenAuthenticator {
|
||||
TokenAuthenticator {
|
||||
key: Zeroizing::new([key; 32]),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn process_tokens_are_domain_separated_and_fail_closed_after_mutation() {
|
||||
let issuer = authenticator(0x11);
|
||||
let other_process = authenticator(0x22);
|
||||
let nonce = [0x33; TOKEN_NONCE_BYTES];
|
||||
let bootstrap = issuer.issue(TokenKind::Bootstrap, nonce);
|
||||
let session = issuer.issue(TokenKind::Session, nonce);
|
||||
|
||||
assert_ne!(bootstrap, session);
|
||||
assert!(issuer.authentic(&bootstrap));
|
||||
assert!(issuer.authentic(&session));
|
||||
assert!(!other_process.authentic(&bootstrap));
|
||||
assert!(!other_process.authentic(&session));
|
||||
|
||||
for token in [bootstrap, session] {
|
||||
for index in 0..TOKEN_BYTES {
|
||||
let mut mutated = token;
|
||||
mutated[index] ^= 1;
|
||||
assert!(!issuer.authentic(&mutated), "mutation at byte {index}");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -111,6 +111,7 @@ fn test_runtime_with_dc(
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts: BTreeMap::new(),
|
||||
profiles: vec![Arc::clone(&profile)],
|
||||
capabilities: vec![[7; 32]].into_boxed_slice(),
|
||||
}));
|
||||
config.rebuild_runtime_user_auth().unwrap();
|
||||
let limits = config.web.limits.clone();
|
||||
|
||||
@@ -70,6 +70,7 @@ fn runtime(admission: bool) -> TestRuntime {
|
||||
config.web.runtime = Some(Arc::new(WebRuntimeConfig {
|
||||
vhosts: BTreeMap::new(),
|
||||
profiles: vec![Arc::clone(&profile)],
|
||||
capabilities: vec![[7; 32]].into_boxed_slice(),
|
||||
}));
|
||||
config.rebuild_runtime_user_auth().unwrap();
|
||||
let limits = config.web.limits.clone();
|
||||
|
||||
Reference in New Issue
Block a user