use std::borrow::Cow;
use std::collections::{BTreeMap, BTreeSet};
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use async_trait::async_trait;
use axum::body::{to_bytes, Bytes};
use axum::extract::ws::WebSocketUpgrade;
use axum::extract::{
ConnectInfo, DefaultBodyLimit, FromRequestParts, MatchedPath, Query, RawPathParams, Request,
State,
};
use axum::http::{HeaderMap, Method, StatusCode};
use axum::response::{IntoResponse, Response};
use axum::routing::any;
use axum::{Json, Router};
use base64::Engine;
use harn_vm::TenantId;
use ipnet::IpNet;
use serde_json::{json, Map, Value};
use crate::adapter::DispatchRuntime;
use crate::auth::AuthRequest;
use crate::http_codec::{
axum_response_from_call, axum_response_from_dispatch_error, classify_ws_upgrade,
fresh_request_id,
};
use crate::tls::HttpTlsConfig;
use crate::ws::{ws_accept, WsConfig, WsMessage, WsSession};
use crate::{
CallArguments, CallRequest, DispatchCore, DispatchError, ExportCatalog, RouteSpec,
TransportConfig, DEFAULT_HTTP_BODY_LIMIT_BYTES,
};
const SITE_ADAPTER: &str = "site";
#[async_trait]
pub trait SiteAuth: Send + Sync {
async fn authenticate(
&self,
parts: &axum::http::request::Parts,
route: &RouteSpec,
) -> SiteAuthOutcome;
}
pub enum SiteAuthOutcome {
Allow(SiteAuthContext),
Deny(Box<Response>),
}
#[derive(Clone, Debug, Default)]
pub struct SiteAuthContext {
pub tenant_id: Option<TenantId>,
pub scopes: BTreeSet<String>,
pub subject: Option<String>,
pub scheme: Option<String>,
pub kind: Option<String>,
pub context: Option<Value>,
}
impl SiteAuthContext {
pub(crate) fn principal(&self) -> harn_vm::AuthPrincipal {
harn_vm::AuthPrincipal {
subject: self.subject.clone().unwrap_or_default(),
scheme: self.scheme.clone().unwrap_or_default(),
scopes: self.scopes.clone(),
kind: self.kind.clone(),
}
}
}
impl SiteAuthOutcome {
pub fn deny(response: Response) -> Self {
Self::Deny(Box::new(response))
}
}
#[async_trait]
pub trait SiteStreamProvider: Send + Sync {
async fn open(
&self,
route: &RouteSpec,
auth: Option<&SiteAuthContext>,
request: Value,
body: Option<Bytes>,
) -> Response;
async fn upgrade(
&self,
route: &RouteSpec,
_auth: Option<&SiteAuthContext>,
_ws: WebSocketUpgrade,
_request: Value,
) -> Response {
(
StatusCode::UPGRADE_REQUIRED,
Json(json!({
"code": "ws_upgrade_unsupported",
"message": format!(
"route {} {} is marked @ws but this stream provider does not implement \
SiteStreamProvider::upgrade",
route.method, route.path
),
})),
)
.into_response()
}
}
#[derive(Clone, Debug)]
pub struct SiteHttpServeOptions {
pub bind: SocketAddr,
pub public_url: Option<String>,
pub tls: HttpTlsConfig,
}
pub struct SiteServerConfig {
pub core: DispatchCore,
pub transport: TransportConfig,
pub trusted_proxies: Vec<IpNet>,
pub auth: Option<Arc<dyn SiteAuth>>,
pub stream_provider: Option<Arc<dyn SiteStreamProvider>>,
}
impl SiteServerConfig {
pub fn new(core: DispatchCore) -> Self {
Self {
core,
transport: TransportConfig::default_enabled(),
trusted_proxies: Vec::new(),
auth: None,
stream_provider: None,
}
}
pub fn with_transport(mut self, transport: TransportConfig) -> Self {
self.transport = transport;
self
}
pub fn with_trusted_proxies(mut self, trusted_proxies: Vec<IpNet>) -> Self {
self.trusted_proxies = trusted_proxies;
self
}
pub fn with_auth(mut self, auth: Arc<dyn SiteAuth>) -> Self {
self.auth = Some(auth);
self
}
pub fn with_stream_provider(mut self, provider: Arc<dyn SiteStreamProvider>) -> Self {
self.stream_provider = Some(provider);
self
}
}
pub struct SiteServer {
config: SiteServerConfig,
}
impl SiteServer {
pub fn new(config: SiteServerConfig) -> Self {
Self { config }
}
pub fn router(self) -> Result<Router, String> {
let SiteServerConfig {
core,
transport,
trusted_proxies,
auth,
stream_provider,
} = self.config;
let catalog = Arc::new(core.catalog().clone());
let runtime = Arc::new(DispatchRuntime::start("SITE", Arc::new(core)));
build_site_router(
&catalog,
runtime,
&transport,
trusted_proxies,
auth,
stream_provider,
)
}
pub async fn run_http(self, options: SiteHttpServeOptions) -> Result<(), String> {
let tls = options.tls.clone();
let router = self.router()?;
let router = crate::tls::apply_security_headers(router, &tls);
let listener = crate::tls::bind_listener(options.bind)?;
let local_addr = listener
.local_addr()
.map_err(|error| format!("failed to read local addr: {error}"))?;
let advertised = options
.public_url
.clone()
.unwrap_or_else(|| format!("{}://{local_addr}", tls.listener_scheme()));
eprintln!("[harn] Site server ready on {advertised}");
crate::tls::serve_router_from_tcp(listener, router, &tls)
.await
.map_err(|error| format!("Site server failed: {error}"))
}
}
#[derive(Clone)]
struct SiteRoute {
function: String,
spec: RouteSpec,
required_scopes: BTreeSet<String>,
method_scopes: BTreeMap<String, BTreeSet<String>>,
allowed_kinds: BTreeSet<String>,
stream: bool,
raw: bool,
ws: bool,
}
impl SiteRoute {
fn scopes_for(&self, method: &Method) -> Cow<'_, BTreeSet<String>> {
match self.method_scopes.get(method.as_str()) {
None => Cow::Borrowed(&self.required_scopes),
Some(extra) if extra.is_empty() => Cow::Borrowed(&self.required_scopes),
Some(extra) if self.required_scopes.is_empty() => Cow::Borrowed(extra),
Some(extra) => Cow::Owned(self.required_scopes.union(extra).cloned().collect()),
}
}
}
#[derive(Clone)]
struct SiteState {
runtime: Arc<DispatchRuntime>,
routes: Arc<BTreeMap<String, BTreeMap<String, SiteRoute>>>,
trusted_proxies: Arc<Vec<IpNet>>,
auth: Option<Arc<dyn SiteAuth>>,
stream_provider: Option<Arc<dyn SiteStreamProvider>>,
}
impl SiteState {
fn resolve(&self, path: &str, method: &Method) -> Option<&SiteRoute> {
let methods = self.routes.get(path)?;
methods.get(method.as_str()).or_else(|| methods.get("*"))
}
}
fn build_site_router(
catalog: &ExportCatalog,
runtime: Arc<DispatchRuntime>,
transport: &TransportConfig,
trusted_proxies: Vec<IpNet>,
auth: Option<Arc<dyn SiteAuth>>,
stream_provider: Option<Arc<dyn SiteStreamProvider>>,
) -> Result<Router, String> {
let mut routes: BTreeMap<String, BTreeMap<String, SiteRoute>> = BTreeMap::new();
for function in catalog.functions.values() {
let Some(spec @ RouteSpec { method, path }) = function.route.as_ref() else {
continue;
};
if (function.stream || function.raw || function.ws) && stream_provider.is_none() {
let marker = if function.stream {
"@stream"
} else if function.raw {
"@raw"
} else {
"@ws"
};
return Err(format!(
"route {method} {path} (`{}`) is marked {marker} but no stream provider is \
configured; install one with SiteServerConfig::with_stream_provider(...)",
function.name
));
}
let by_method = routes.entry(path.clone()).or_default();
let entry = SiteRoute {
function: function.name.clone(),
spec: spec.clone(),
required_scopes: function.required_scopes.clone(),
method_scopes: function.method_scopes.clone(),
allowed_kinds: function
.policy
.as_ref()
.map(|policy| policy.allowed_kinds.clone())
.unwrap_or_default(),
stream: function.stream,
raw: function.raw,
ws: function.ws,
};
if let Some(existing) = by_method.insert(method.clone(), entry) {
return Err(format!(
"route conflict: {method} {path} is claimed by both `{}` and \
`{}`; give one of them a distinct @route(...)",
existing.function, function.name
));
}
}
if routes.is_empty() {
return Err(
"no HTTP routes found: export a `pub fn handler_*` or annotate a function with \
@route(\"METHOD\", \"/path\") to serve it"
.to_string(),
);
}
let state = SiteState {
runtime,
routes: Arc::new(routes),
trusted_proxies: Arc::new(trusted_proxies),
auth,
stream_provider,
};
let mut router: Router<SiteState> = Router::new();
for path in state.routes.keys() {
router = router.route(path, any(site_dispatch));
}
let router = router
.layer(DefaultBodyLimit::max(DEFAULT_HTTP_BODY_LIMIT_BYTES))
.with_state(state);
Ok(crate::apply_transport_layers(router, transport))
}
fn provider_scope_backstop(
auth_context: Option<&SiteAuthContext>,
required_scopes: &BTreeSet<String>,
request_id: &str,
) -> Option<Response> {
if auth_context.is_none() && !required_scopes.is_empty() {
return Some(axum_response_from_dispatch_error(
DispatchError::Forbidden {
required: required_scopes.clone(),
granted: BTreeSet::new(),
},
request_id,
));
}
None
}
async fn site_dispatch(State(state): State<SiteState>, request: Request) -> Response {
let (mut parts, body) = request.into_parts();
let matched_path = MatchedPath::from_request_parts(&mut parts, &())
.await
.ok()
.map(|m| m.as_str().to_string());
let raw_params = RawPathParams::from_request_parts(&mut parts, &())
.await
.ok();
let query = Query::<BTreeMap<String, String>>::from_request_parts(&mut parts, &())
.await
.map(|q| q.0)
.unwrap_or_default();
let peer = parts
.extensions
.get::<ConnectInfo<SocketAddr>>()
.map(|ConnectInfo(addr)| *addr);
let method = parts.method.clone();
let route_template = matched_path.unwrap_or_else(|| parts.uri.path().to_string());
let Some(route) = state.resolve(&route_template, &method).cloned() else {
return method_not_allowed(&state, &route_template);
};
let function = route.function.clone();
let request_id = fresh_request_id();
let required_scopes = route.scopes_for(&method);
let auth_context = match state.auth.as_ref() {
None => None,
Some(hook) => match hook.authenticate(&parts, &route.spec).await {
SiteAuthOutcome::Deny(response) => return *response,
SiteAuthOutcome::Allow(context) => {
if !required_scopes.is_subset(&context.scopes) {
return axum_response_from_dispatch_error(
DispatchError::Forbidden {
required: required_scopes.into_owned(),
granted: context.scopes,
},
&request_id,
);
}
if !route.allowed_kinds.is_empty()
&& !context
.kind
.as_deref()
.is_some_and(|kind| route.allowed_kinds.contains(kind))
{
return axum_response_from_dispatch_error(
DispatchError::ForbiddenPrincipalKind {
allowed: route.allowed_kinds.clone(),
},
&request_id,
);
}
Some(context)
}
},
};
if route.ws && (is_websocket_upgrade(&parts.headers) || !route.stream) {
if let Some(response) =
provider_scope_backstop(auth_context.as_ref(), &required_scopes, &request_id)
{
return response;
}
let Some(provider) = state.stream_provider.as_ref() else {
return axum_response_from_dispatch_error(
DispatchError::Execution(format!(
"provider route {route_template} has no stream provider"
)),
&request_id,
);
};
let upgrade = match WebSocketUpgrade::from_request_parts(&mut parts, &()).await {
Ok(upgrade) => upgrade,
Err(rejection) => return rejection.into_response(),
};
let request = build_request_value(
&method,
&parts.uri,
&route_template,
raw_params.as_ref(),
&query,
&parts.headers,
&[],
peer,
&state.trusted_proxies,
);
return provider
.upgrade(&route.spec, auth_context.as_ref(), upgrade, request)
.await;
}
if route.stream || route.raw {
if let Some(response) =
provider_scope_backstop(auth_context.as_ref(), &required_scopes, &request_id)
{
return response;
}
let Some(provider) = state.stream_provider.as_ref() else {
return axum_response_from_dispatch_error(
DispatchError::Execution(format!(
"provider route {route_template} has no stream provider"
)),
&request_id,
);
};
let raw_body = if route.raw {
match buffer_request_body(body).await {
Ok(bytes) => Some(bytes),
Err(response) => return response,
}
} else {
None
};
let request = build_request_value(
&method,
&parts.uri,
&route_template,
raw_params.as_ref(),
&query,
&parts.headers,
&[],
peer,
&state.trusted_proxies,
);
return provider
.open(&route.spec, auth_context.as_ref(), request, raw_body)
.await;
}
let wants_upgrade = is_websocket_upgrade(&parts.headers);
let body_bytes = if wants_upgrade {
Bytes::new()
} else {
match buffer_request_body(body).await {
Ok(bytes) => bytes,
Err(response) => return response,
}
};
let req_value = build_request_value(
&method,
&parts.uri,
&route_template,
raw_params.as_ref(),
&query,
&parts.headers,
&body_bytes,
peer,
&state.trusted_proxies,
);
let mut auth = AuthRequest::from_http(
&method,
parts.uri.path(),
body_bytes.to_vec(),
&parts.headers,
);
if let Some(context) = auth_context.as_ref() {
auth.granted_scopes = context.scopes.clone();
}
let call = CallRequest {
adapter: SITE_ADAPTER.to_string(),
function: function.clone(),
arguments: CallArguments::Positional(vec![req_value]),
auth,
caller: SITE_ADAPTER.to_string(),
replay_key: None,
trace_id: None,
parent_span_id: None,
metadata: BTreeMap::new(),
cancel_token: None,
agent_session_id: None,
agent_event_sink: None,
actor_chain: None,
actor_chain_hop: None,
progress: None,
tenant_id: auth_context
.as_ref()
.and_then(|context| context.tenant_id.clone()),
request_id: Some(request_id.clone()),
auth_context: auth_context
.as_ref()
.and_then(|context| context.context.clone()),
auth_principal: auth_context.as_ref().map(SiteAuthContext::principal),
};
let response = match state.runtime.call(call).await {
Ok(response) => response,
Err(error) => return axum_response_from_dispatch_error(error, &request_id),
};
if wants_upgrade {
if let Some(spec) = classify_ws_upgrade(&response) {
return upgrade_websocket(&mut parts, state.runtime, spec, auth_context).await;
}
}
axum_response_from_call(response, &request_id)
}
async fn buffer_request_body(body: axum::body::Body) -> Result<Bytes, Response> {
to_bytes(body, DEFAULT_HTTP_BODY_LIMIT_BYTES)
.await
.map_err(|_| {
(
StatusCode::PAYLOAD_TOO_LARGE,
Json(json!({
"code": "request_body_too_large",
"message": format!(
"request body exceeds the {DEFAULT_HTTP_BODY_LIMIT_BYTES}-byte limit"
),
})),
)
.into_response()
})
}
fn method_not_allowed(state: &SiteState, path: &str) -> Response {
let allow = state
.routes
.get(path)
.map(|methods| {
methods
.keys()
.filter(|method| method.as_str() != "*")
.cloned()
.collect::<Vec<_>>()
.join(", ")
})
.unwrap_or_default();
let mut response = (
StatusCode::METHOD_NOT_ALLOWED,
Json(json!({
"code": "method_not_allowed",
"message": format!("no handler for this method on {path}"),
})),
)
.into_response();
if !allow.is_empty() {
if let Ok(value) = allow.parse() {
response
.headers_mut()
.insert(axum::http::header::ALLOW, value);
}
}
response
}
fn is_websocket_upgrade(headers: &HeaderMap) -> bool {
let has_upgrade = headers
.get(axum::http::header::UPGRADE)
.and_then(|value| value.to_str().ok())
.map(|value| value.eq_ignore_ascii_case("websocket"))
.unwrap_or(false);
has_upgrade && headers.contains_key("sec-websocket-key")
}
async fn upgrade_websocket(
parts: &mut axum::http::request::Parts,
runtime: Arc<DispatchRuntime>,
spec: harn_vm::WsUpgradeSpec,
auth_context: Option<SiteAuthContext>,
) -> Response {
let config = ws_config_from_spec(&spec);
let on_message = spec.on_message.clone();
let headers = parts.headers.clone();
let upgrade = match WebSocketUpgrade::from_request_parts(parts, &()).await {
Ok(upgrade) => upgrade,
Err(rejection) => return rejection.into_response(),
};
ws_accept(config, headers, upgrade, move |session| {
let runtime = runtime.clone();
let on_message = on_message.clone();
let auth_context = auth_context.clone();
async move {
drive_ws_session(session, runtime, on_message, auth_context).await;
}
})
.await
}
fn ws_config_from_spec(spec: &harn_vm::WsUpgradeSpec) -> WsConfig {
let mut config = WsConfig {
subprotocols: spec.offered.clone(),
..WsConfig::default()
};
if let Some(ms) = spec.idle_ping_ms {
config.idle_ping = Some(std::time::Duration::from_millis(ms));
}
if let Some(bytes) = spec.max_message_bytes {
config.max_message_bytes = bytes as usize;
}
config
}
async fn drive_ws_session(
session: WsSession,
runtime: Arc<DispatchRuntime>,
on_message: Option<String>,
auth_context: Option<SiteAuthContext>,
) {
let Some(handler) = on_message else {
while let Ok(Some(_)) = session.recv().await {}
return;
};
while let Ok(Some(message)) = session.recv().await {
let message_value = match &message {
WsMessage::Text(text) => json!({ "type": "text", "data": text }),
WsMessage::Binary(bytes) => json!({
"type": "binary",
"data_base64": base64::engine::general_purpose::STANDARD.encode(bytes),
}),
};
let auth = AuthRequest {
granted_scopes: auth_context
.as_ref()
.map(|context| context.scopes.clone())
.unwrap_or_default(),
..AuthRequest::default()
};
let call = CallRequest {
adapter: SITE_ADAPTER.to_string(),
function: handler.clone(),
arguments: CallArguments::Positional(vec![message_value]),
auth,
caller: SITE_ADAPTER.to_string(),
replay_key: None,
trace_id: None,
parent_span_id: None,
metadata: BTreeMap::new(),
cancel_token: None,
agent_session_id: None,
agent_event_sink: None,
actor_chain: None,
actor_chain_hop: None,
progress: None,
tenant_id: auth_context
.as_ref()
.and_then(|context| context.tenant_id.clone()),
request_id: Some(fresh_request_id()),
auth_context: auth_context
.as_ref()
.and_then(|context| context.context.clone()),
auth_principal: auth_context.as_ref().map(SiteAuthContext::principal),
};
match runtime.call(call).await {
Ok(response) => {
if !send_ws_reply(&session, response.value).await {
break;
}
}
Err(_) => {
let _ = session.close(1011, "handler error").await;
break;
}
}
}
}
async fn send_ws_reply(session: &WsSession, value: Value) -> bool {
match value {
Value::Null => true,
Value::String(text) => session.send(text).await.is_ok(),
other => match serde_json::to_string(&other) {
Ok(text) => session.send(text).await.is_ok(),
Err(_) => true,
},
}
}
#[allow(clippy::too_many_arguments)]
fn build_request_value(
method: &Method,
uri: &axum::http::Uri,
route_template: &str,
raw_params: Option<&RawPathParams>,
query: &BTreeMap<String, String>,
headers: &HeaderMap,
body: &[u8],
peer: Option<SocketAddr>,
trusted_proxies: &[IpNet],
) -> Value {
let mut path_params = Map::new();
if let Some(params) = raw_params {
for (name, value) in params {
path_params.insert(name.to_string(), Value::String(value.to_string()));
}
}
let query = query
.iter()
.map(|(key, value)| (key.clone(), Value::String(value.clone())))
.collect::<Map<_, _>>();
let mut header_map = Map::new();
for (name, value) in headers.iter() {
if let Ok(text) = value.to_str() {
header_map.insert(
name.as_str().to_ascii_lowercase(),
Value::String(text.to_string()),
);
}
}
let path_params = Value::Object(path_params);
let mut request = Map::new();
request.insert("method".into(), Value::String(method.as_str().to_string()));
request.insert("path".into(), Value::String(uri.path().to_string()));
request.insert("route".into(), Value::String(route_template.to_string()));
request.insert("params".into(), path_params.clone());
request.insert("path_params".into(), path_params);
request.insert("query".into(), Value::Object(query));
request.insert("headers".into(), Value::Object(header_map));
match std::str::from_utf8(body) {
Ok(text) => {
request.insert("body".into(), Value::String(text.to_string()));
request.insert("body_kind".into(), Value::String("text".to_string()));
}
Err(_) => {
request.insert("body".into(), Value::Null);
request.insert("body_kind".into(), Value::String("base64".to_string()));
}
}
request.insert(
"body_base64".into(),
Value::String(base64::engine::general_purpose::STANDARD.encode(body)),
);
request.insert("content_length".into(), Value::from(body.len()));
request.insert(
"client_ip".into(),
resolve_client_ip(peer.map(|addr| addr.ip()), headers, trusted_proxies)
.map(|ip| Value::String(ip.to_string()))
.unwrap_or(Value::Null),
);
request.insert(
"remote_addr".into(),
peer.map(|addr| Value::String(addr.to_string()))
.unwrap_or(Value::Null),
);
Value::Object(request)
}
fn resolve_client_ip(
peer: Option<IpAddr>,
headers: &HeaderMap,
trusted_proxies: &[IpNet],
) -> Option<IpAddr> {
let peer = peer?;
let is_trusted = |ip: &IpAddr| trusted_proxies.iter().any(|net| net.contains(ip));
if trusted_proxies.is_empty() || !is_trusted(&peer) {
return Some(peer);
}
let forwarded: Vec<IpAddr> = headers
.get("x-forwarded-for")
.and_then(|value| value.to_str().ok())
.map(|value| {
value
.split(',')
.filter_map(|hop| hop.trim().parse::<IpAddr>().ok())
.collect()
})
.unwrap_or_default();
if let Some(client) = forwarded.iter().rev().find(|ip| !is_trusted(ip)) {
return Some(*client);
}
if let Some(originator) = forwarded.first() {
return Some(*originator);
}
let real_ip = headers
.get("x-real-ip")
.and_then(|value| value.to_str().ok())
.and_then(|value| value.trim().parse::<IpAddr>().ok());
Some(real_ip.unwrap_or(peer))
}
#[cfg(test)]
mod tests {
use super::*;
fn nets(cidrs: &[&str]) -> Vec<IpNet> {
cidrs.iter().map(|c| c.parse().unwrap()).collect()
}
fn headers(pairs: &[(&str, &str)]) -> HeaderMap {
let mut map = HeaderMap::new();
for (name, value) in pairs {
map.insert(
axum::http::HeaderName::from_bytes(name.as_bytes()).unwrap(),
value.parse().unwrap(),
);
}
map
}
fn ip(value: &str) -> IpAddr {
value.parse().unwrap()
}
#[test]
fn no_peer_resolves_to_none() {
assert_eq!(
resolve_client_ip(None, &headers(&[]), &nets(&["10.0.0.0/8"])),
None
);
}
#[test]
fn untrusted_config_returns_peer_ignoring_headers() {
let got = resolve_client_ip(
Some(ip("203.0.113.1")),
&headers(&[("x-forwarded-for", "1.2.3.4"), ("x-real-ip", "5.6.7.8")]),
&[],
);
assert_eq!(got, Some(ip("203.0.113.1")));
}
#[test]
fn real_ip_is_the_fallback_when_no_forwarded_for() {
let got = resolve_client_ip(
Some(ip("10.0.0.5")),
&headers(&[("x-real-ip", "9.9.9.9")]),
&nets(&["10.0.0.0/8"]),
);
assert_eq!(got, Some(ip("9.9.9.9")));
}
#[test]
fn all_trusted_chain_falls_back_to_leftmost_originator() {
let got = resolve_client_ip(
Some(ip("10.0.0.5")),
&headers(&[("x-forwarded-for", "10.0.0.1, 10.0.0.2")]),
&nets(&["10.0.0.0/8"]),
);
assert_eq!(got, Some(ip("10.0.0.1")));
}
#[test]
fn ipv6_proxy_and_client_are_supported() {
let got = resolve_client_ip(
Some(ip("2001:db8::1")),
&headers(&[("x-forwarded-for", "2606:4700::1234")]),
&nets(&["2001:db8::/32"]),
);
assert_eq!(got, Some(ip("2606:4700::1234")));
}
#[test]
fn garbage_forwarded_entries_are_skipped() {
let got = resolve_client_ip(
Some(ip("10.0.0.5")),
&headers(&[("x-forwarded-for", "not-an-ip, 1.2.3.4, also-bad")]),
&nets(&["10.0.0.0/8"]),
);
assert_eq!(got, Some(ip("1.2.3.4")));
}
}