Skip to main content

fakecloud_core/
dispatch.rs

1use axum::body::Body;
2use axum::extract::{ConnectInfo, Extension, Query};
3use axum::http::{Request, StatusCode};
4use axum::response::Response;
5use bytes::Bytes;
6use std::collections::HashMap;
7use std::net::SocketAddr;
8use std::sync::Arc;
9
10use crate::auth::{
11    is_root_bypass, ConditionContext, CredentialResolver, IamMode, IamPolicyEvaluator,
12    InternalCaller, Principal, PrincipalType, ResourcePolicyProvider,
13};
14use crate::protocol::{self, AwsProtocol};
15use crate::registry::ServiceRegistry;
16use crate::service::{AwsRequest, ResponseBody};
17
18/// Pins an in-process request to one REST service; see [`dispatch_to_service`].
19/// Private, so only this module can attach it.
20#[derive(Clone, Copy, Debug)]
21struct PinnedService(&'static str);
22
23/// Dispatch an in-process request straight to one REST-protocol service
24/// (`"s3"`), bypassing the HTTP router and service detection.
25///
26/// The CloudFront data plane fetches an S3 origin this way: the request is by
27/// construction an S3 request addressed to the origin bucket (via its `Host`),
28/// so no viewer path (`/_fakecloud/*`, `/latest/*`, ...) can reach one of
29/// fakecloud's own routes, and no viewer header (`X-Amz-Target`, an
30/// `Authorization` scoped to another service, `?Action=`) can steer it to
31/// another service. Authentication and IAM enforcement run exactly as for any
32/// request, including an [`InternalCaller`] extension. The source address is
33/// the request's `ConnectInfo<SocketAddr>` extension, or loopback.
34pub async fn dispatch_to_service(
35    service: &'static str,
36    registry: Arc<ServiceRegistry>,
37    config: Arc<DispatchConfig>,
38    mut request: Request<Body>,
39) -> Response<Body> {
40    let remote_addr = request
41        .extensions()
42        .get::<ConnectInfo<SocketAddr>>()
43        .map(|c| c.0)
44        .unwrap_or_else(|| SocketAddr::from(([127, 0, 0, 1], 0)));
45    let query = match Query::<HashMap<String, String>>::try_from_uri(request.uri()) {
46        Ok(q) => q,
47        Err(e) => {
48            return build_error_response(
49                StatusCode::BAD_REQUEST,
50                "InvalidArgument",
51                &format!("Invalid query string: {e}"),
52                &uuid::Uuid::new_v4().to_string(),
53                AwsProtocol::Rest,
54            )
55        }
56    };
57    request.extensions_mut().insert(PinnedService(service));
58    dispatch(
59        ConnectInfo(remote_addr),
60        Extension(registry),
61        Extension(config),
62        query,
63        request,
64    )
65    .await
66}
67
68/// Services whose Smithy model marks operations `@requestCompression`, so a
69/// client may send a gzip `Content-Encoding` body that the service decodes.
70/// Elsewhere `Content-Encoding` is left alone: S3 stores it as object metadata
71/// and API Gateway forwards it to the backend.
72const REQUEST_COMPRESSION_SERVICES: &[&str] = &["monitoring"];
73
74/// Decompressed bodies are capped like buffered ones, so a small gzip bomb
75/// can't expand past the request size limit.
76fn decode_request_compression(
77    headers: &http::HeaderMap,
78    rpc_v2_cbor: Option<&protocol::DetectedRequest>,
79    body: Bytes,
80) -> Result<Bytes, (String, AwsProtocol)> {
81    let gzipped = headers
82        .get_all(http::header::CONTENT_ENCODING)
83        .iter()
84        .filter_map(|v| v.to_str().ok())
85        .flat_map(|v| v.split(','))
86        .any(|enc| enc.trim().eq_ignore_ascii_case("gzip"));
87    if !gzipped || body.is_empty() {
88        return Ok(body);
89    }
90    let target = headers
91        .get("x-amz-target")
92        .and_then(|v| v.to_str().ok())
93        .and_then(protocol::parse_amz_target);
94    let (service, protocol) = if let Some(d) = rpc_v2_cbor {
95        (Some(d.service.clone()), AwsProtocol::RpcV2Cbor)
96    } else if let Some(d) = target {
97        (Some(d.service), AwsProtocol::Json)
98    } else {
99        (
100            protocol::extract_service_from_auth(headers),
101            AwsProtocol::Query,
102        )
103    };
104    if !service.is_some_and(|s| REQUEST_COMPRESSION_SERVICES.contains(&s.as_str())) {
105        return Ok(body);
106    }
107    use std::io::Read;
108    let limit = max_request_body_bytes() as u64;
109    let mut out = Vec::new();
110    let read = flate2::read::MultiGzDecoder::new(body.as_ref())
111        .take(limit + 1)
112        .read_to_end(&mut out);
113    match read {
114        Ok(_) if out.len() as u64 > limit => {
115            Err(("Decompressed request body too large".to_string(), protocol))
116        }
117        Ok(_) => Ok(Bytes::from(out)),
118        Err(e) => Err((
119            format!("Unable to decompress gzip request body: {e}"),
120            protocol,
121        )),
122    }
123}
124
125/// awsJson content type of protocol version 1.0.
126const AWS_JSON_1_0: &str = "application/x-amz-json-1.0";
127/// awsJson content type of protocol version 1.1, the services' default.
128const AWS_JSON_1_1: &str = "application/x-amz-json-1.1";
129
130/// The main dispatch handler. All HTTP requests come through here.
131pub async fn dispatch(
132    ConnectInfo(remote_addr): ConnectInfo<SocketAddr>,
133    Extension(registry): Extension<Arc<ServiceRegistry>>,
134    Extension(config): Extension<Arc<DispatchConfig>>,
135    Query(query_params): Query<HashMap<String, String>>,
136    request: Request<Body>,
137) -> Response<Body> {
138    let json_1_0 = request
139        .headers()
140        .get(http::header::CONTENT_TYPE)
141        .and_then(|v| v.to_str().ok())
142        .is_some_and(|ct| ct.trim().eq_ignore_ascii_case(AWS_JSON_1_0));
143    let mut response = dispatch_inner(remote_addr, registry, config, query_params, request).await;
144    if json_1_0 {
145        answer_in_json_1_0(&mut response);
146    }
147    response
148}
149
150/// An awsJson 1.0 service (DynamoDB, SQS, CloudWatch, Step Functions, ...)
151/// answers in the content type of its protocol version, which its clients
152/// send on the request. Handlers render JSON as 1.1, so relabel the response
153/// (success or error) for a 1.0 caller.
154fn answer_in_json_1_0(response: &mut Response<Body>) {
155    let is_1_1 = response
156        .headers()
157        .get(http::header::CONTENT_TYPE)
158        .and_then(|v| v.to_str().ok())
159        .is_some_and(|ct| ct.eq_ignore_ascii_case(AWS_JSON_1_1));
160    if is_1_1 {
161        response.headers_mut().insert(
162            http::header::CONTENT_TYPE,
163            http::HeaderValue::from_static(AWS_JSON_1_0),
164        );
165    }
166}
167
168async fn dispatch_inner(
169    remote_addr: SocketAddr,
170    registry: Arc<ServiceRegistry>,
171    config: Arc<DispatchConfig>,
172    query_params: HashMap<String, String>,
173    request: Request<Body>,
174) -> Response<Body> {
175    let remote_addr = Some(remote_addr);
176    let request_id = uuid::Uuid::new_v4().to_string();
177
178    let (mut parts, body) = request.into_parts();
179    // Set only by [`dispatch_to_service`]: the request is for this service,
180    // whatever its headers or query say.
181    let pinned = parts.extensions.get::<PinnedService>().map(|p| p.0);
182
183    // Streaming opt-in: if the route is a known large-body S3 / ECR
184    // upload, we skip the buffered `to_bytes` step entirely and hand
185    // the raw body to the service handler. The handler spills it to
186    // disk on the fly. Header-only detection covers every streaming
187    // candidate (none of them rely on form-body sniffing).
188    // Smithy RPC v2 CBOR is recognized by its `smithy-protocol` header and
189    // `/service/{Service}/operation/{Op}` path before any other detection:
190    // its SigV4 scope alone would otherwise route it as the service's Query
191    // protocol.
192    let rpc_v2_cbor = if pinned.is_none() {
193        protocol::detect_rpc_v2_cbor(&parts.headers, parts.uri.path())
194    } else {
195        None
196    };
197    let stream_route = streaming_route(
198        &parts.method,
199        parts.uri.path(),
200        &parts.headers,
201        &query_params,
202    );
203    let header_only = protocol::detect_service_headers_only(&parts.headers, &query_params);
204    let stream_dispatch = match (&stream_route, &header_only) {
205        // A pinned request is always buffered: its caller already holds the
206        // whole body, and header detection must not pick its service.
207        _ if pinned.is_some() => None,
208        _ if rpc_v2_cbor.is_some() => None,
209        // Header-only detection agrees with the URL match — covers S3
210        // PUT object (SigV4 service=s3 in Authorization).
211        (Some(sr), Some(detected)) if sr.0 == detected.service => Some(detected.clone()),
212        // ECR OCI v2 blob upload has no AWS auth header; the path
213        // alone (`/v2/.../blobs/uploads/...`) tells us the route is
214        // ECR. Synthesize a DetectedRequest so dispatch picks the
215        // streaming path. Same special-case the buffered branch
216        // applies on detect_service None (see below).
217        (Some((service, _)), None) if *service == "ecr" => Some(protocol::DetectedRequest {
218            service: "ecr".to_string(),
219            action: String::new(),
220            protocol: AwsProtocol::Rest,
221        }),
222        _ => None,
223    };
224
225    let (body_bytes, body_stream) = if stream_dispatch.is_some() {
226        (Bytes::new(), Some(body))
227    } else {
228        // Buffered path: materialize the body into memory under the
229        // configured cap. `FAKECLOUD_MAX_REQUEST_BODY_BYTES` (default
230        // 1 GiB) caps non-streaming requests; streaming routes have no
231        // cap because nothing materializes the entire body in RAM.
232        let max_body_bytes = max_request_body_bytes();
233        match axum::body::to_bytes(body, max_body_bytes).await {
234            Ok(b) => (b, None),
235            Err(_) => {
236                return build_error_response(
237                    StatusCode::PAYLOAD_TOO_LARGE,
238                    "RequestEntityTooLarge",
239                    "Request body too large",
240                    &request_id,
241                    AwsProtocol::Query,
242                );
243            }
244        }
245    };
246
247    // `@requestCompression`: a client may gzip the body of an operation whose
248    // model allows it. Decode it before anything parses the body (Query
249    // detection reads the form body), but keep the wire bytes for SigV4,
250    // which signed the compressed payload.
251    let wire_body = body_bytes.clone();
252    let body_bytes = if pinned.is_none() && stream_dispatch.is_none() {
253        match decode_request_compression(&parts.headers, rpc_v2_cbor.as_ref(), body_bytes) {
254            Ok(b) => b,
255            Err((message, protocol)) => {
256                return build_error_response(
257                    StatusCode::BAD_REQUEST,
258                    "SerializationException",
259                    &message,
260                    &request_id,
261                    protocol,
262                );
263            }
264        }
265    } else {
266        body_bytes
267    };
268
269    // Detect service and action
270    let detected = if let Some(service) = pinned {
271        protocol::DetectedRequest {
272            service: service.to_string(),
273            action: String::new(),
274            protocol: AwsProtocol::Rest,
275        }
276    } else if let Some(d) = rpc_v2_cbor {
277        d
278    } else if let Some(d) = stream_dispatch {
279        d
280    } else {
281        match protocol::detect_service(&parts.headers, &query_params, &body_bytes) {
282            Some(d) => d,
283            None => {
284                // A request carrying X-Amz-Target is unambiguously an awsJson
285                // call whose operation we couldn't map to a known service. AWS
286                // answers these with UnknownOperationException; routing them to
287                // the apigateway catch-all below would 404 with a misleading
288                // "Stage not found" instead.
289                if let Some(target) = parts
290                    .headers
291                    .get("x-amz-target")
292                    .and_then(|v| v.to_str().ok())
293                {
294                    return build_error_response(
295                        StatusCode::BAD_REQUEST,
296                        "UnknownOperationException",
297                        &format!("The operation {target} is not recognized."),
298                        &request_id,
299                        AwsProtocol::Json,
300                    );
301                }
302                // OPTIONS requests (CORS preflight) don't carry Authorization headers.
303                // One addressed to an API's `{api-id}.execute-api.` host is that
304                // API's preflight (its OPTIONS method, e.g. a MOCK CORS
305                // integration); any other unsigned OPTIONS goes to S3, the other
306                // REST service that answers CORS preflights.
307                if parts.method == http::Method::OPTIONS && is_execute_api_host(&parts.headers) {
308                    protocol::DetectedRequest {
309                        service: "apigateway".to_string(),
310                        action: String::new(),
311                        protocol: AwsProtocol::RestJson,
312                    }
313                } else if parts.method == http::Method::OPTIONS {
314                    protocol::DetectedRequest {
315                        service: "s3".to_string(),
316                        action: String::new(),
317                        protocol: AwsProtocol::Rest,
318                    }
319                } else if parts.uri.path() == "/v2" || parts.uri.path().starts_with("/v2/") {
320                    // OCI Distribution v2 protocol. Docker CLI / OCI clients
321                    // use Basic auth (not SigV4) and GET /v2/ with no body,
322                    // so this must be matched before the apigateway fallback.
323                    protocol::DetectedRequest {
324                        service: "ecr".to_string(),
325                        action: String::new(),
326                        protocol: AwsProtocol::Rest,
327                    }
328                } else if let Some(bucket) = anonymous_s3_bucket(&parts.uri, &config) {
329                    // Unsigned request whose first path segment names an
330                    // existing S3 bucket: an anonymous path-style S3 access
331                    // (e.g. serving a public-read object to a browser). Without
332                    // this it would fall through to the apigateway catch-all
333                    // below and 404 with "Stage not found" (#1707). Authorization
334                    // for the anonymous caller still runs in the IAM block.
335                    tracing::debug!(bucket = %bucket, "routing unsigned request to S3 (existing bucket)");
336                    protocol::DetectedRequest {
337                        service: "s3".to_string(),
338                        action: String::new(),
339                        protocol: AwsProtocol::Rest,
340                    }
341                } else if !parts.uri.path().starts_with("/_")
342                    || parts.uri.path().starts_with("/_aws/execute-api/")
343                {
344                    // Requests without AWS auth that don't match any service might be
345                    // API Gateway execute API calls (plain HTTP without signatures).
346                    // Route them to apigateway service which will validate if a matching
347                    // API/stage exists. Skip special FakeCloud endpoints (/_*),
348                    // except LocalStack's path-style execute-api invocation URL
349                    // `/_aws/execute-api/{api-id}/{stage}/{path}`.
350                    protocol::DetectedRequest {
351                        service: "apigateway".to_string(),
352                        action: String::new(),
353                        protocol: AwsProtocol::RestJson,
354                    }
355                } else {
356                    return build_error_response(
357                        StatusCode::BAD_REQUEST,
358                        "MissingAction",
359                        "Could not determine target service or action from request",
360                        &request_id,
361                        AwsProtocol::Query,
362                    );
363                }
364            }
365        }
366    };
367
368    // Bedrock-agent and bedrock-runtime both send `bedrock` in the SigV4
369    // credential scope, but bedrock-agent has its own service handler.
370    // Disambiguate based on the request path.
371    let detected = if detected.service == "bedrock" {
372        match bedrock_agent_service_for(&parts.method, parts.uri.path()) {
373            Some(service) => protocol::DetectedRequest {
374                service: service.to_string(),
375                ..detected
376            },
377            None => detected,
378        }
379    } else {
380        detected
381    };
382
383    // Amazon DocumentDB shares RDS's Query wire protocol, `rds` SigV4
384    // signing scope, and `rds.<region>.amazonaws.com` endpoint, so a real
385    // `aws-sdk-docdb` request is indistinguishable from `aws-sdk-rds` by
386    // signing name or host alone. The DocumentDB SDK does stamp an
387    // `api/docdb` token into its `user-agent`; use it to route the request
388    // to the dedicated `docdb` handler. (The conformance probe signs the
389    // `docdb` scope directly, so it never reaches this branch.) Requests
390    // without the token stay on `rds`.
391    let detected = if detected.service == "rds" && user_agent_indicates_docdb(&parts.headers) {
392        protocol::DetectedRequest {
393            service: "docdb".to_string(),
394            ..detected
395        }
396    } else {
397        detected
398    };
399
400    // Amazon Neptune shares RDS's Query wire protocol, `rds` SigV4 signing
401    // scope, and `rds.<region>.amazonaws.com` endpoint, so a real
402    // `aws-sdk-neptune` request is indistinguishable from `aws-sdk-rds` by
403    // signing name or host alone. The Neptune SDK does stamp an
404    // `api/neptune` token into its `user-agent`; use it to route the request
405    // to the dedicated `neptune` handler. (The conformance probe signs the
406    // `neptune` scope directly, so it never reaches this branch.) Requests
407    // without the token stay on `rds`.
408    let detected = if detected.service == "rds" && user_agent_indicates_neptune(&parts.headers) {
409        protocol::DetectedRequest {
410            service: "neptune".to_string(),
411            ..detected
412        }
413    } else {
414        detected
415    };
416
417    // Look up service
418    let service = match registry.get(&detected.service) {
419        Some(s) => s,
420        None => {
421            return build_error_response(
422                detected.protocol.error_status(),
423                "UnknownService",
424                &format!("Service '{}' is not available", detected.service),
425                &request_id,
426                ErrorEnvelope::for_request(&detected, &parts.headers),
427            );
428        }
429    };
430
431    // Extract region and access key from auth header (or presigned query).
432    let auth_header = parts
433        .headers
434        .get("authorization")
435        .and_then(|v| v.to_str().ok())
436        .unwrap_or("");
437    let header_info = fakecloud_aws::sigv4::parse_sigv4(auth_header);
438    let presigned_info = if header_info.is_none() {
439        // Presigned URL: credentials live in the query string.
440        fakecloud_aws::sigv4::parse_sigv4_presigned(&query_params).map(|p| p.as_info())
441    } else {
442        None
443    };
444    let sigv4_info = header_info.or(presigned_info);
445    // SigV2 presigned URLs (`AWSAccessKeyId` + `Signature` + `Expires` query
446    // parameters) carry the access key outside the SigV4 grammar, so the SigV4
447    // parsers above return None. Recover the key here so a SigV2-presigned
448    // request is attributed to its caller instead of being treated as
449    // anonymous (which would deny it under the object read auth gate).
450    let access_key_id = sigv4_info
451        .as_ref()
452        .map(|info| info.access_key.clone())
453        .or_else(|| sigv2_presigned_access_key(&query_params));
454
455    // Host-header routing hint: LocalStack-shaped
456    // `<svc>.<region>.localhost.localstack.cloud[:port]`, real-AWS
457    // `<svc>.<region>.amazonaws.com`, and every S3 virtual-hosted variant
458    // of both. Secondary region source and carries the bucket for
459    // virtual-hosted S3 path rewrite.
460    let host_info = protocol::parse_routing_host_from_headers(&parts.headers);
461
462    let region = sigv4_info
463        .map(|info| info.region)
464        .or_else(|| host_info.as_ref().map(|h| h.region.clone()))
465        .or_else(|| extract_region_from_user_agent(&parts.headers))
466        .unwrap_or_else(|| config.region.clone());
467
468    // Resolve the caller's principal up front so both SigV4 verification
469    // (which needs the secret) and the service handler (which needs the
470    // identity for GetCallerIdentity and IAM enforcement) share a single
471    // lookup. The root-bypass AKID skips resolution entirely — `test`
472    // credentials have no backing identity and must always pass.
473    let caller_akid = access_key_id.as_deref().unwrap_or("");
474    let resolved = if !caller_akid.is_empty() && !is_root_bypass(caller_akid) {
475        config
476            .credential_resolver
477            .as_ref()
478            .and_then(|r| r.resolve(caller_akid))
479    } else {
480        None
481    };
482    let caller_principal = resolved.as_ref().map(|r| r.principal.clone());
483    let caller_session_policies = resolved
484        .as_ref()
485        .map(|r| r.session_policies.clone())
486        .unwrap_or_default();
487
488    // Opt-in SigV4 cryptographic verification. Runs before the service
489    // handler so a failing signature never reaches business logic. The
490    // reserved `test*` root identity short-circuits verification to keep
491    // local-dev workflows frictionless.
492    //
493    // A fully anonymous request — no `Authorization` header AND no presigned
494    // credential (SigV4 `X-Amz-Credential` or SigV2 `AWSAccessKeyId`) — carries
495    // no signature to verify, so it must NOT be hard-403'd here even under
496    // `--verify-sigv4`. AWS treats an unsigned request as the anonymous
497    // principal and lets a resource policy / public-read ACL authorize it
498    // (e.g. a public S3 object GET). Let it fall through to the
499    // anonymous-authorization path below, which allows a public resource and
500    // otherwise denies. A request that DID present a credential but failed to
501    // parse is NOT anonymous (`Authorization` non-empty or a presign query
502    // present) and still returns the malformed-signature error below.
503    let is_fully_anonymous = auth_header.is_empty()
504        && !query_params.contains_key("X-Amz-Credential")
505        && sigv2_presigned_access_key(&query_params).is_none();
506    if config.verify_sigv4
507        && !is_fully_anonymous
508        && !is_root_bypass(caller_akid)
509        && config.credential_resolver.is_some()
510    {
511        let amz_date = parts
512            .headers
513            .get("x-amz-date")
514            .and_then(|v| v.to_str().ok());
515        let parsed = fakecloud_aws::sigv4::parse_sigv4_header(auth_header, amz_date)
516            .or_else(|| fakecloud_aws::sigv4::parse_sigv4_presigned(&query_params));
517        let parsed = match parsed {
518            Some(p) => p,
519            None => {
520                return build_error_response(
521                    StatusCode::FORBIDDEN,
522                    "IncompleteSignature",
523                    "Request is missing or has a malformed AWS Signature",
524                    &request_id,
525                    ErrorEnvelope::for_request(&detected, &parts.headers),
526                );
527            }
528        };
529        let resolved_for_verify = match resolved.as_ref() {
530            Some(r) => r,
531            None => {
532                return unresolved_credential_response(
533                    &config,
534                    caller_akid,
535                    &request_id,
536                    ErrorEnvelope::for_request(&detected, &parts.headers),
537                );
538            }
539        };
540        let headers_vec = fakecloud_aws::sigv4::headers_from_http(&parts.headers);
541        let raw_query_for_verify = parts.uri.query().unwrap_or("").to_string();
542        let verify_req = fakecloud_aws::sigv4::VerifyRequest {
543            method: parts.method.as_str(),
544            path: parts.uri.path(),
545            query: &raw_query_for_verify,
546            headers: &headers_vec,
547            body: &wire_body,
548        };
549        match fakecloud_aws::sigv4::verify(
550            &parsed,
551            &verify_req,
552            &resolved_for_verify.secret_access_key,
553            chrono::Utc::now(),
554        ) {
555            Ok(()) => {
556                // Bind the buffered request body to the signed
557                // `x-amz-content-sha256`. The sigv4 canonical builder uses that
558                // header value verbatim (it is a signed header) and deliberately
559                // does NOT re-hash the body: for streaming / aws-chunked routes
560                // `body_bytes` is empty at that layer, so re-hashing there would
561                // reject legitimate signed requests (see the
562                // `feedback_sigv4_body_hash_wrong_layer` caveat). S3 rebinds the
563                // body hash in its own write path (`XAmzContentSHA256Mismatch`);
564                // every other service is bound HERE, where the full buffered body
565                // is available. Only a genuine 64-char lowercase-hex digest is
566                // checked — `UNSIGNED-PAYLOAD`, `STREAMING-*`, and presigned
567                // (UNSIGNED-PAYLOAD) requests are skipped, so correct clients are
568                // unaffected. A mismatch means the body was altered after signing;
569                // AWS would then compute a different canonical request, so return
570                // `SignatureDoesNotMatch`.
571                if !parsed.is_presigned && detected.service != "s3" {
572                    if let Some(signed_hash) = parts
573                        .headers
574                        .get("x-amz-content-sha256")
575                        .and_then(|v| v.to_str().ok())
576                        .filter(|h| is_hex_sha256(h))
577                    {
578                        if sha256_hex_lower(&wire_body) != signed_hash {
579                            return build_error_response(
580                                StatusCode::FORBIDDEN,
581                                "SignatureDoesNotMatch",
582                                "The request signature we calculated does not match the signature you provided",
583                                &request_id,
584                                ErrorEnvelope::for_request(&detected, &parts.headers),
585                            );
586                        }
587                    }
588                }
589            }
590            Err(fakecloud_aws::sigv4::SigV4Error::RequestTimeTooSkewed { .. }) => {
591                return build_error_response(
592                    StatusCode::FORBIDDEN,
593                    "RequestTimeTooSkewed",
594                    "The difference between the request time and the current time is too large",
595                    &request_id,
596                    ErrorEnvelope::for_request(&detected, &parts.headers),
597                );
598            }
599            Err(fakecloud_aws::sigv4::SigV4Error::InvalidDate(msg)) => {
600                return build_error_response(
601                    StatusCode::FORBIDDEN,
602                    "IncompleteSignature",
603                    &format!("Invalid x-amz-date: {msg}"),
604                    &request_id,
605                    ErrorEnvelope::for_request(&detected, &parts.headers),
606                );
607            }
608            Err(fakecloud_aws::sigv4::SigV4Error::Malformed(msg)) => {
609                return build_error_response(
610                    StatusCode::FORBIDDEN,
611                    "IncompleteSignature",
612                    &format!("Malformed SigV4 signature: {msg}"),
613                    &request_id,
614                    ErrorEnvelope::for_request(&detected, &parts.headers),
615                );
616            }
617            Err(fakecloud_aws::sigv4::SigV4Error::SignatureMismatch) => {
618                return build_error_response(
619                    StatusCode::FORBIDDEN,
620                    "SignatureDoesNotMatch",
621                    "The request signature we calculated does not match the signature you provided",
622                    &request_id,
623                    ErrorEnvelope::for_request(&detected, &parts.headers),
624                );
625            }
626            Err(fakecloud_aws::sigv4::SigV4Error::PresignedUrlExpired { .. }) => {
627                return build_error_response(
628                    StatusCode::FORBIDDEN,
629                    "AccessDenied",
630                    "Request has expired",
631                    &request_id,
632                    ErrorEnvelope::for_request(&detected, &parts.headers),
633                );
634            }
635            Err(fakecloud_aws::sigv4::SigV4Error::InvalidPresignExpires(_)) => {
636                return build_error_response(
637                    StatusCode::BAD_REQUEST,
638                    "AuthorizationQueryParametersError",
639                    "X-Amz-Expires must be a number between 1 and 604800 seconds",
640                    &request_id,
641                    ErrorEnvelope::for_request(&detected, &parts.headers),
642                );
643            }
644        }
645    }
646
647    // A SigV4 presigned S3 URL can carry request headers as query parameters
648    // (the JS v3 presigner hoists every `x-amz-*` header this way). S3 applies
649    // them as if they were headers, so lift them into the header map now that
650    // the signature over the raw query has been verified.
651    if detected.service == "s3" && query_params.contains_key("X-Amz-Credential") {
652        hoist_presigned_query_headers(&mut parts.headers, &query_params);
653    }
654
655    // Build path segments. For S3 virtual-hosted-style requests the bucket
656    // lives in the Host header, not the path — prepend it so the S3 handler
657    // sees a uniform path-style request. SigV4 verification above already
658    // ran against the wire path, so this rewrite is signature-safe.
659    let wire_path = parts.uri.path();
660    let path = if detected.service == "s3" {
661        s3_routing_path(
662            wire_path,
663            host_info.as_ref().and_then(|h| h.bucket.as_deref()),
664        )
665    } else {
666        wire_path.to_string()
667    };
668    let raw_query = parts.uri.query().unwrap_or("").to_string();
669    // Split on `/` first, then percent-decode each segment once, so handlers
670    // see decoded `@httpLabel` values and an encoded `/` (`%2F`) stays inside
671    // its label. `raw_path` keeps the undecoded wire form.
672    let path_segments = crate::path::split_path_segments(&path);
673
674    // rpcv2Cbor: the signature above covered the CBOR wire body; from here on
675    // the request is handled on the service's JSON path, so swap in the
676    // equivalent awsJson document.
677    let body_bytes = if detected.protocol == AwsProtocol::RpcV2Cbor {
678        match crate::cbor::decode_to_json(&body_bytes) {
679            Ok(json) => Bytes::from(json.to_string()),
680            Err(e) => {
681                return build_error_response(
682                    StatusCode::BAD_REQUEST,
683                    "SerializationException",
684                    &format!("Unable to decode CBOR request body: {e}"),
685                    &request_id,
686                    AwsProtocol::RpcV2Cbor,
687                );
688            }
689        }
690    } else {
691        body_bytes
692    };
693
694    // For JSON protocol, validate that non-empty bodies are valid JSON
695    if detected.protocol == AwsProtocol::Json
696        && !body_bytes.is_empty()
697        && serde_json::from_slice::<serde_json::Value>(&body_bytes).is_err()
698    {
699        return build_error_response(
700            StatusCode::BAD_REQUEST,
701            "SerializationException",
702            "Start of structure or map found where not expected",
703            &request_id,
704            AwsProtocol::Json,
705        );
706    }
707
708    // Merge query params with form body params for the Query family (awsQuery
709    // and ec2Query share identical form-encoded request bodies).
710    let mut all_params = query_params;
711    if matches!(
712        detected.protocol,
713        AwsProtocol::Query | AwsProtocol::Ec2Query
714    ) {
715        let body_params = protocol::parse_query_body(&body_bytes);
716        for (k, v) in body_params {
717            all_params.entry(k).or_insert(v);
718        }
719    }
720
721    // CloudWatch (`monitoring`) advertises awsJson1_0 alongside awsQuery. Its
722    // handlers all read the flat awsQuery param map, so when a client uses the
723    // JSON protocol we flatten the JSON body into that same map, leaving the
724    // handlers unchanged. The handler emits a JSON response for JSON callers.
725    if matches!(
726        detected.protocol,
727        AwsProtocol::Json | AwsProtocol::RpcV2Cbor
728    ) && detected.service == "monitoring"
729    {
730        let body_params = protocol::flatten_json_to_query(&body_bytes);
731        for (k, v) in body_params {
732            all_params.entry(k).or_insert(v);
733        }
734    }
735
736    // An in-process request fakecloud issues as an AWS-owned principal (e.g.
737    // CloudFront fetching an S3 origin through an origin access control). The
738    // identity rides a request extension that only in-process code can set,
739    // and is honored only when the request presents no credentials of its
740    // own: a request that does is authorized as that caller instead.
741    let internal_caller = if is_fully_anonymous && access_key_id.is_none() {
742        parts.extensions.get::<InternalCaller>().cloned()
743    } else {
744        None
745    };
746    let caller_principal =
747        caller_principal.or_else(|| internal_caller.as_ref().map(|c| c.principal()));
748
749    let error_envelope = ErrorEnvelope::for_request(&detected, &parts.headers);
750    let aws_request = AwsRequest {
751        service: detected.service.clone(),
752        action: detected.action.clone(),
753        region,
754        // A service acting for a customer works on the target resource in its
755        // owner's account: an S3 origin bucket may belong to another account
756        // than the distribution fetching it.
757        account_id: match internal_caller.as_ref() {
758            Some(caller) => {
759                internal_caller_account(caller, &detected.service, &path_segments, &config)
760            }
761            None => caller_principal
762                .as_ref()
763                .map(|p| p.account_id.clone())
764                .unwrap_or_else(|| config.account_id.clone()),
765        },
766        request_id: request_id.clone(),
767        headers: parts.headers,
768        query_params: all_params,
769        body: body_bytes,
770        body_stream: parking_lot::Mutex::new(body_stream),
771        path_segments,
772        raw_path: path,
773        raw_query,
774        method: parts.method,
775        is_query_protocol: matches!(
776            detected.protocol,
777            AwsProtocol::Query | AwsProtocol::Ec2Query
778        ),
779        access_key_id,
780        principal: caller_principal,
781    };
782
783    tracing::info!(
784        service = %aws_request.service,
785        action = %aws_request.action,
786        request_id = %aws_request.request_id,
787        "handling request"
788    );
789
790    // Opt-in IAM identity-policy enforcement. Runs before the service
791    // handler so a deny never reaches business logic. Root principals
792    // (both `test*` bypass AKIDs and the account's IAM root) are exempt,
793    // matching AWS behavior. Services that haven't opted in via
794    // `iam_enforceable()` are transparently skipped — the startup log
795    // lists which services are under enforcement so users always know.
796    if config.iam_mode.is_enabled()
797        && service.iam_enforceable()
798        && !is_root_bypass(aws_request.access_key_id.as_deref().unwrap_or(""))
799    {
800        if let Some(evaluator) = config.policy_evaluator.as_ref() {
801            if let Some(caller) = internal_caller.as_ref() {
802                if let Some(denied) = authorize_internal_caller(
803                    caller,
804                    service.as_ref(),
805                    &aws_request,
806                    evaluator.as_ref(),
807                    &config,
808                    &detected,
809                    &request_id,
810                ) {
811                    return denied;
812                }
813            } else if let Some(principal) = aws_request.principal.as_ref() {
814                if principal.is_root() {
815                    if let Some(denied) = authorize_root_cross_account(
816                        principal,
817                        service.as_ref(),
818                        &aws_request,
819                        evaluator.as_ref(),
820                        &config,
821                        &detected,
822                        &request_id,
823                        remote_addr,
824                        resolved.as_ref(),
825                    ) {
826                        return denied;
827                    }
828                } else {
829                    // A request can need several authorizations -- one per
830                    // table in a batch, say -- and every one must allow it.
831                    let iam_actions = service.iam_actions_for(&aws_request);
832                    if !iam_actions.is_empty() {
833                        for iam_action in &iam_actions {
834                            let mut condition_context = principal_condition_context(
835                                principal,
836                                resolved.as_ref(),
837                                service.as_ref(),
838                                &aws_request,
839                                iam_action,
840                                &detected,
841                                remote_addr,
842                            );
843                            let service_resource = !iam_action.is_pass_role();
844
845                            // Phase 2: fetch the resource-based policy (if
846                            // any) attached to the target resource and
847                            // pass it to the evaluator alongside the
848                            // principal's identity policies. The resource's
849                            // owning account is parsed from the ARN (#381
850                            // multi-account alignment); S3 ARNs have an
851                            // empty account field, so we fall back to the
852                            // server's configured account ID in that case.
853                            // A request creating its resource in the caller's
854                            // account (S3 CreateBucket) is same-account, whatever
855                            // account holds that name today.
856                            let in_caller_account =
857                                service.iam_resource_in_caller_account(&aws_request);
858                            let resource_policy_json = config
859                                .resource_policy_provider
860                                .as_ref()
861                                .filter(|_| service_resource && !in_caller_account)
862                                .and_then(|p| {
863                                    p.resource_policy(&detected.service, &iam_action.resource)
864                                });
865                            // Derive the resource-owning account. Prefer a provider
866                            // lookup (S3 ARNs carry no account, so the bucket's
867                            // owner is resolved from state — without this, account
868                            // A reaching account B's bucket would be mis-read as
869                            // same-account and skip B's bucket-policy requirement,
870                            // bug-audit 2026-05-28, 5.3), then fall back to the
871                            // account embedded in the ARN (SQS/SNS/Lambda/…), then
872                            // to the caller's account for wildcard / unscoped
873                            // actions (ListQueues, GetCallerIdentity).
874                            let resource_account_id = config
875                                .resource_policy_provider
876                                .as_ref()
877                                .filter(|_| service_resource && !in_caller_account)
878                                .and_then(|p| {
879                                    p.resource_owner_account(
880                                        &detected.service,
881                                        &iam_action.resource,
882                                    )
883                                })
884                                .or_else(|| parse_account_from_arn(&iam_action.resource))
885                                .unwrap_or_else(|| principal.account_id.clone());
886                            // SCP ceiling: resolve the inherited SCP chain
887                            // for this principal (management accounts and
888                            // service-linked roles come back as `None`, in
889                            // which case the evaluator treats the layer as
890                            // absent). Audit breadcrumbs emitted by the
891                            // resolver itself, not here.
892                            let scps = config
893                                .scp_resolver
894                                .as_ref()
895                                .and_then(|r| r.scps_for(principal));
896                            // Global keys AWS sets on every request, which
897                            // perimeter guardrails gate on (`Deny ...
898                            // StringNotEquals aws:PrincipalOrgID` /
899                            // `aws:ResourceAccount`): absent, those negated
900                            // conditions would deny everyone.
901                            add_global_request_keys(
902                                &mut condition_context,
903                                principal,
904                                &resource_account_id,
905                                config.scp_resolver.as_deref(),
906                            );
907                            let decision = evaluator.evaluate_with_resource_policy(
908                                principal,
909                                iam_action,
910                                &condition_context,
911                                resource_policy_json.as_deref(),
912                                &resource_account_id,
913                                &caller_session_policies,
914                                scps.as_deref(),
915                            );
916                            if !decision.is_allow() {
917                                tracing::warn!(
918                                    target: "fakecloud::iam::audit",
919                                    service = %detected.service,
920                                    action = %iam_action.action_string(),
921                                    resource = %iam_action.resource,
922                                    principal = %principal.arn,
923                                    resource_policy_present = resource_policy_json.is_some(),
924                                    decision = ?decision,
925                                    mode = %config.iam_mode,
926                                    request_id = %request_id,
927                                    "IAM policy evaluation denied request"
928                                );
929                                if config.iam_mode.is_strict() {
930                                    // Real AWS includes an "Encoded
931                                    // authorization failure message" suffix
932                                    // on AccessDeniedException — an opaque
933                                    // base64+zlib JSON blob that the caller
934                                    // can pass to STS
935                                    // `DecodeAuthorizationMessage` to
936                                    // recover the structured deny reason
937                                    // (action, principal, matched
938                                    // statements, condition context). We
939                                    // produce the same blob inline so
940                                    // existing tooling that decodes deny
941                                    // reasons works against fakecloud.
942                                    let context_summary = serde_json::json!({
943                                        "aws:PrincipalArn": principal.arn,
944                                        "aws:PrincipalAccount": principal.account_id,
945                                        "aws:RequestedRegion": condition_context
946                                            .aws_requested_region
947                                            .clone()
948                                            .unwrap_or_default(),
949                                        "aws:SecureTransport": condition_context
950                                            .aws_secure_transport
951                                            .unwrap_or(false),
952                                        "aws:Action": iam_action.action_string(),
953                                        "aws:Resource": iam_action.resource,
954                                        "decision": format!("{:?}", decision),
955                                    });
956                                    let action_string = iam_action.action_string();
957                                    let encoded = crate::auth_message::encode_deny(
958                                        matches!(decision, crate::auth::IamDecision::ExplicitDeny),
959                                        Some(&action_string),
960                                        Some(&principal.arn),
961                                        Vec::new(),
962                                        Some(context_summary),
963                                    );
964                                    return build_error_response(
965                                    StatusCode::FORBIDDEN,
966                                    "AccessDeniedException",
967                                    &format!(
968                                        "User: {} is not authorized to perform: {} on resource: {} Encoded authorization failure message: {}",
969                                        principal.arn,
970                                        iam_action.action_string(),
971                                        iam_action.resource,
972                                        encoded,
973                                    ),
974                                    &request_id,
975                                    error_envelope,
976                                );
977                                }
978                                // Soft mode: audit log already emitted; fall
979                                // through to the handler.
980                            }
981                        }
982                    } else {
983                        // Service opted in via `iam_enforceable()` but its
984                        // `iam_action_for` returned no `IamAction` for this
985                        // specific operation (e.g. S3's `s3_detect_action`
986                        // has `_ => return None` arms for unrecognized
987                        // sub-resources). Under strict enforcement that must
988                        // fail closed: an operation we cannot map to an IAM
989                        // action cannot be authorized, so serving it would be
990                        // a fail-open bypass of `--iam strict`. Deny by
991                        // default. In soft mode we preserve the historical
992                        // warn-and-allow so an incomplete mapping surfaces
993                        // during rollout without blocking traffic.
994                        tracing::warn!(
995                            target: "fakecloud::iam::audit",
996                            service = %detected.service,
997                            action = %aws_request.action,
998                            mode = %config.iam_mode,
999                            request_id = %request_id,
1000                            "service is iam_enforceable but has no IamAction mapping for this action; denying under strict, allowing under soft"
1001                        );
1002                        if config.iam_mode.is_strict() {
1003                            return build_error_response(
1004                                StatusCode::FORBIDDEN,
1005                                "AccessDeniedException",
1006                                &format!(
1007                                    "User: {} is not authorized to perform: {}: no IAM action mapping exists for this operation, so it cannot be authorized under strict IAM enforcement",
1008                                    principal.arn, aws_request.action,
1009                                ),
1010                                &request_id,
1011                                error_envelope,
1012                            );
1013                        }
1014                        // Soft mode: audit log emitted; fall through to the
1015                        // handler.
1016                    }
1017                }
1018            } else if let Some(akid) = aws_request
1019                .access_key_id
1020                .as_deref()
1021                .filter(|_| aws_request.principal.is_none())
1022            {
1023                // The request presented an access key that resolves to no
1024                // identity: unknown, deactivated (Inactive), or an expired /
1025                // revoked STS session. AWS rejects such a request before
1026                // authorization (`InvalidClientTokenId` / `ExpiredToken`), so
1027                // under strict enforcement it must not run unchecked -- nor
1028                // reach STS, which would otherwise read the missing principal
1029                // as the account root and mint real credentials. The
1030                // explicit exemptions stay: the reserved `test*` root
1031                // identity (excluded above) and truly anonymous requests
1032                // (next branch). Soft mode logs and lets it through.
1033                tracing::warn!(
1034                    target: "fakecloud::iam::audit",
1035                    service = %detected.service,
1036                    action = %aws_request.action,
1037                    mode = %config.iam_mode,
1038                    request_id = %request_id,
1039                    "request credential does not resolve to an identity; denying under strict, allowing under soft"
1040                );
1041                if config.iam_mode.is_strict() {
1042                    return unresolved_credential_response(
1043                        &config,
1044                        akid,
1045                        &request_id,
1046                        ErrorEnvelope::for_request(&detected, &aws_request.headers),
1047                    );
1048                }
1049            } else if aws_request.access_key_id.is_none() {
1050                // Truly anonymous (unsigned) caller — no Authorization header at
1051                // all. No identity policies exist, so authorization rests
1052                // entirely on the resource policy (a bucket policy granting
1053                // `Principal:"*"`) and public-read ACLs — mirroring AWS, which
1054                // denies anonymous requests unless the resource is explicitly
1055                // made public. Without this an anonymous request that reached an
1056                // iam_enforceable service in enforcement mode would bypass
1057                // authorization entirely.
1058                //
1059                // A request whose credential did not resolve (principal `None`
1060                // with `access_key_id` `Some`) is handled by the branch above.
1061                let iam_actions = service.iam_actions_for(&aws_request);
1062                if !iam_actions.is_empty() {
1063                    for iam_action in &iam_actions {
1064                        let now = chrono::Utc::now();
1065                        let mut condition_context = ConditionContext {
1066                            aws_source_ip: remote_addr.map(|sa| sa.ip()),
1067                            aws_current_time: Some(now),
1068                            aws_epoch_time: Some(now.timestamp()),
1069                            aws_secure_transport: Some(is_secure_transport(&aws_request.headers)),
1070                            aws_requested_region: Some(aws_request.region.clone()),
1071                            ..Default::default()
1072                        };
1073                        condition_context.service_keys =
1074                            service.iam_condition_keys_for(&aws_request, iam_action);
1075                        // `aws:ResourceAccount` is set for anonymous requests
1076                        // too: the target resource's owner.
1077                        if let Some(owner) = config
1078                            .resource_policy_provider
1079                            .as_ref()
1080                            .and_then(|p| {
1081                                p.resource_owner_account(&detected.service, &iam_action.resource)
1082                            })
1083                            .or_else(|| parse_account_from_arn(&iam_action.resource))
1084                        {
1085                            condition_context
1086                                .service_keys
1087                                .entry("aws:resourceaccount".to_string())
1088                                .or_insert_with(|| vec![owner]);
1089                        }
1090                        let resource_policy_json =
1091                            config.resource_policy_provider.as_ref().and_then(|p| {
1092                                p.resource_policy(&detected.service, &iam_action.resource)
1093                            });
1094                        let policy_decision = evaluator.evaluate_anonymous(
1095                            iam_action,
1096                            &condition_context,
1097                            resource_policy_json.as_deref(),
1098                        );
1099                        let policy_allows = policy_decision.is_allow();
1100                        // An explicit Deny in the resource policy always wins, even
1101                        // over a public-read ACL — matching AWS's Deny-overrides
1102                        // precedence. Collapsing the decision to a bool and ORing the
1103                        // ACL let a public ACL override an explicit anonymous Deny.
1104                        let policy_explicit_deny =
1105                            matches!(policy_decision, crate::auth::IamDecision::ExplicitDeny);
1106                        let acl_allows = !policy_explicit_deny
1107                            && config.resource_policy_provider.as_ref().is_some_and(|p| {
1108                                p.public_acl_allows(
1109                                    &detected.service,
1110                                    &iam_action.resource,
1111                                    iam_action.action,
1112                                )
1113                            });
1114                        if !policy_allows && !acl_allows {
1115                            tracing::warn!(
1116                                target: "fakecloud::iam::audit",
1117                                service = %detected.service,
1118                                action = %iam_action.action_string(),
1119                                resource = %iam_action.resource,
1120                                resource_policy_present = resource_policy_json.is_some(),
1121                                mode = %config.iam_mode,
1122                                request_id = %request_id,
1123                                "anonymous request denied: no public bucket policy or ACL grants the action"
1124                            );
1125                            if config.iam_mode.is_strict() {
1126                                return build_error_response(
1127                                    StatusCode::FORBIDDEN,
1128                                    "AccessDenied",
1129                                    "Access Denied",
1130                                    &request_id,
1131                                    error_envelope,
1132                                );
1133                            }
1134                            // Soft mode: audit log emitted; fall through to the handler.
1135                        }
1136                    }
1137                } else {
1138                    // Anonymous request to an iam_enforceable service whose
1139                    // operation has no IamAction mapping. Mirror the signed
1140                    // branch above: an operation we cannot map to an IAM action
1141                    // cannot be authorized, so serving it would be a fail-open
1142                    // bypass of `--iam strict`. Deny by default under strict;
1143                    // soft mode warns and falls through.
1144                    tracing::warn!(
1145                        target: "fakecloud::iam::audit",
1146                        service = %detected.service,
1147                        action = %aws_request.action,
1148                        mode = %config.iam_mode,
1149                        request_id = %request_id,
1150                        "anonymous request to iam_enforceable service has no IamAction mapping; denying under strict, allowing under soft"
1151                    );
1152                    if config.iam_mode.is_strict() {
1153                        return build_error_response(
1154                            StatusCode::FORBIDDEN,
1155                            "AccessDenied",
1156                            "Access Denied",
1157                            &request_id,
1158                            error_envelope,
1159                        );
1160                    }
1161                }
1162            }
1163        }
1164    }
1165
1166    match service.handle(aws_request).await {
1167        Ok(resp) => {
1168            let resp = if detected.protocol == AwsProtocol::RpcV2Cbor {
1169                rpc_v2_cbor_response(resp)
1170            } else {
1171                resp
1172            };
1173            let mut builder = Response::builder()
1174                .status(resp.status)
1175                .header("x-amzn-requestid", &request_id)
1176                .header("x-amz-request-id", &request_id);
1177
1178            if !resp.content_type.is_empty() {
1179                builder = builder.header("content-type", &resp.content_type);
1180            }
1181
1182            let has_content_length = resp
1183                .headers
1184                .iter()
1185                .any(|(k, _)| k.as_str().eq_ignore_ascii_case("content-length"));
1186
1187            for (k, v) in &resp.headers {
1188                builder = builder.header(k, v);
1189            }
1190
1191            match resp.body {
1192                ResponseBody::Bytes(b) => builder.body(Body::from(b)).unwrap(),
1193                ResponseBody::File { file, size } => {
1194                    let stream = tokio_util::io::ReaderStream::new(file);
1195                    let body = Body::from_stream(stream);
1196                    if !has_content_length {
1197                        builder = builder.header("content-length", size.to_string());
1198                    }
1199                    builder.body(body).unwrap()
1200                }
1201            }
1202        }
1203        Err(err) => {
1204            tracing::warn!(
1205                service = %detected.service,
1206                action = %detected.action,
1207                error = %err,
1208                "request failed"
1209            );
1210            let error_headers = err.response_headers().to_vec();
1211            let mut resp = build_error_response_with_fields(
1212                err.status(),
1213                err.code(),
1214                &err.message(),
1215                &request_id,
1216                error_envelope,
1217                err.extra_fields(),
1218            );
1219            for (k, v) in &error_headers {
1220                if let (Ok(name), Ok(val)) = (
1221                    k.parse::<http::header::HeaderName>(),
1222                    v.parse::<http::header::HeaderValue>(),
1223                ) {
1224                    // `Vary` combines, so two entries for it must both survive
1225                    // (S3 adds its CORS value to an error that may carry its
1226                    // own) — but an identical value is not worth repeating.
1227                    // Everything else replaces, so a service can still override
1228                    // a header this builder already set.
1229                    if name == http::header::VARY {
1230                        let already = resp
1231                            .headers()
1232                            .get_all(&name)
1233                            .iter()
1234                            .any(|existing| existing == val);
1235                        if !already {
1236                            resp.headers_mut().append(name, val);
1237                        }
1238                    } else {
1239                        resp.headers_mut().insert(name, val);
1240                    }
1241                }
1242            }
1243            resp
1244        }
1245    }
1246}
1247
1248/// Frame a service response for an rpcv2Cbor caller. A service that renders
1249/// its own CBOR (content type `application/cbor`) passes through; a JSON body
1250/// is transcoded schema-lessly. Either way the response carries the
1251/// `smithy-protocol` header the client checks.
1252fn rpc_v2_cbor_response(mut resp: crate::service::AwsResponse) -> crate::service::AwsResponse {
1253    if resp.content_type != crate::cbor::CBOR_CONTENT_TYPE {
1254        if let ResponseBody::Bytes(bytes) = &resp.body {
1255            let json = if bytes.is_empty() {
1256                serde_json::Value::Object(serde_json::Map::new())
1257            } else {
1258                serde_json::from_slice(bytes).unwrap_or_else(|e| {
1259                    tracing::warn!(error = %e, "non-JSON response body for an rpcv2Cbor request");
1260                    serde_json::Value::Object(serde_json::Map::new())
1261                })
1262            };
1263            resp.body = ResponseBody::Bytes(Bytes::from(crate::cbor::encode(
1264                &crate::cbor::json_to_cbor(&json),
1265            )));
1266            resp.content_type = crate::cbor::CBOR_CONTENT_TYPE.to_string();
1267        }
1268    }
1269    resp.headers.insert(
1270        http::HeaderName::from_static(crate::cbor::SMITHY_PROTOCOL_HEADER),
1271        http::HeaderValue::from_static(crate::cbor::RPC_V2_CBOR),
1272    );
1273    resp
1274}
1275
1276/// Configuration passed to the dispatch handler.
1277#[derive(Clone)]
1278pub struct DispatchConfig {
1279    pub region: String,
1280    pub account_id: String,
1281    /// Whether to cryptographically verify SigV4 signatures on incoming
1282    /// requests. Wired through from `--verify-sigv4` /
1283    /// `FAKECLOUD_VERIFY_SIGV4`. Off by default.
1284    pub verify_sigv4: bool,
1285    /// IAM policy evaluation mode. Wired through from `--iam` /
1286    /// `FAKECLOUD_IAM`. Defaults to [`IamMode::Off`]. Actual evaluation is
1287    /// added in a later batch; today this field is plumbed but never
1288    /// consulted.
1289    pub iam_mode: IamMode,
1290    /// Resolves access key IDs to their secrets and owning principals.
1291    /// Required when `verify_sigv4` or `iam_mode != Off`. When `None`, both
1292    /// features gracefully degrade to off-by-default behavior.
1293    pub credential_resolver: Option<Arc<dyn CredentialResolver>>,
1294    /// Evaluates IAM identity policies for a resolved principal + action.
1295    /// Required when `iam_mode != Off`. When `None`, enforcement silently
1296    /// degrades to off even if `iam_mode` is set.
1297    pub policy_evaluator: Option<Arc<dyn IamPolicyEvaluator>>,
1298    /// Resolves resource-based policies (S3 bucket policies in the
1299    /// initial rollout) to hand to the evaluator alongside the
1300    /// principal's identity policies. `None` means the server was
1301    /// started without any resource-policy-owning service registered;
1302    /// dispatch then behaves as if no resource policy is attached to
1303    /// any resource, identical to the Phase 1 behavior.
1304    pub resource_policy_provider: Option<Arc<dyn ResourcePolicyProvider>>,
1305    /// Resolves the ordered SCP chain that applies to a principal's
1306    /// account (root-OU first, account-direct last). `None` means no
1307    /// organizations resolver has been registered — SCPs never gate
1308    /// any request in that case. Off-by-default matches the Batch 4
1309    /// contract: zero behavior change until a user calls
1310    /// `CreateOrganization` and the resolver is wired.
1311    pub scp_resolver: Option<Arc<dyn crate::auth::ScpResolver>>,
1312}
1313
1314impl std::fmt::Debug for DispatchConfig {
1315    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1316        f.debug_struct("DispatchConfig")
1317            .field("region", &self.region)
1318            .field("account_id", &self.account_id)
1319            .field("verify_sigv4", &self.verify_sigv4)
1320            .field("iam_mode", &self.iam_mode)
1321            .field(
1322                "credential_resolver",
1323                &self
1324                    .credential_resolver
1325                    .as_ref()
1326                    .map(|_| "<CredentialResolver>"),
1327            )
1328            .field(
1329                "policy_evaluator",
1330                &self
1331                    .policy_evaluator
1332                    .as_ref()
1333                    .map(|_| "<IamPolicyEvaluator>"),
1334            )
1335            .field(
1336                "resource_policy_provider",
1337                &self
1338                    .resource_policy_provider
1339                    .as_ref()
1340                    .map(|_| "<ResourcePolicyProvider>"),
1341            )
1342            .field(
1343                "scp_resolver",
1344                &self.scp_resolver.as_ref().map(|_| "<ScpResolver>"),
1345            )
1346            .finish()
1347    }
1348}
1349
1350impl DispatchConfig {
1351    /// Minimal constructor for tests and call sites that don't care about the
1352    /// opt-in security features.
1353    pub fn new(region: impl Into<String>, account_id: impl Into<String>) -> Self {
1354        Self {
1355            region: region.into(),
1356            account_id: account_id.into(),
1357            verify_sigv4: false,
1358            iam_mode: IamMode::Off,
1359            credential_resolver: None,
1360            policy_evaluator: None,
1361            resource_policy_provider: None,
1362            scp_resolver: None,
1363        }
1364    }
1365}
1366
1367/// The path an S3 request is routed on: for virtual-hosted-style the bucket
1368/// comes from the Host header and is prepended, unless the client already put
1369/// it in the path (`PUT /<bucket>` / `PUT /<bucket>/key` against a
1370/// virtual-hosted host). Path-style requests pass through unchanged.
1371///
1372/// Shared with [`streaming_route`] on purpose. The streaming gate has to agree
1373/// with routing about where the bucket ends and the key begins; when the two
1374/// derived it separately they drifted, and a request the router read as a
1375/// bucket-level operation was dispatched unbuffered, so the handler and IAM
1376/// enforcement saw an empty body.
1377///
1378/// Under virtual-hosted addressing the whole wire path is the object key, as
1379/// on real S3, so the bucket is always prefixed: on
1380/// `docs.s3.<region>.amazonaws.com`, `GET /docs/intro.html` addresses the key
1381/// `docs/intro.html`, not `intro.html`.
1382fn s3_routing_path(wire_path: &str, host_bucket: Option<&str>) -> String {
1383    let Some(bucket) = host_bucket else {
1384        return wire_path.to_string();
1385    };
1386    if wire_path == "/" || wire_path.is_empty() {
1387        format!("/{bucket}")
1388    } else {
1389        format!("/{bucket}{wire_path}")
1390    }
1391}
1392
1393/// Extract the 12-digit account ID segment from an AWS ARN.
1394///
1395/// ARNs follow `arn:<partition>:<service>:<region>:<account>:<resource>`.
1396/// Identifies routes that opt into streaming request bodies. Returns
1397/// `Some((service, action_hint))` when the dispatch path should hand
1398/// the raw body to the service handler unbuffered, otherwise `None`
1399/// for the default buffered path. The handler reads the stream via
1400/// [`crate::service::AwsRequest::take_body_stream`].
1401///
1402/// Streaming-eligible routes today:
1403///
1404/// * `s3` PUT object — `PUT /<bucket>/<key>` with a SigV4 (or
1405///   presigned) auth header. Covers PutObject, UploadPart, and
1406///   UploadPartCopy. The S3 service spills to disk via
1407///   [`fakecloud_persistence::BodySource::File`] when the stream is
1408///   present.
1409/// * `ecr` OCI Distribution v2 blob upload — `PATCH` and `PUT` on
1410///   `/v2/{name}/blobs/uploads/{uuid}`. The ECR service spools the
1411///   stream into a per-upload temp file before computing the digest.
1412fn streaming_route(
1413    method: &http::Method,
1414    path: &str,
1415    headers: &http::HeaderMap,
1416    query_params: &HashMap<String, String>,
1417) -> Option<(&'static str, &'static str)> {
1418    // ECR OCI v2 blob upload (PATCH chunk + final PUT).
1419    if (method == http::Method::PATCH || method == http::Method::PUT)
1420        && path.starts_with("/v2/")
1421        && path.contains("/blobs/uploads/")
1422    {
1423        return Some(("ecr", ""));
1424    }
1425
1426    // S3 PutObject / UploadPart / UploadPartCopy. Detect either via
1427    // SigV4 service field in the Authorization header OR via a SigV4
1428    // presigned URL (X-Amz-Credential .../s3/...) OR a SigV2 presigned
1429    // URL (AWSAccessKeyId + Signature + Expires query parameters).
1430    if method == http::Method::PUT {
1431        // An object upload needs a bucket AND a key. Resolve the request the
1432        // way routing will (bucket from the Host header for virtual-hosted,
1433        // straight from the path otherwise), then split it the way the request
1434        // builder builds `path_segments` -- on '/', dropping empty segments.
1435        // Anything that leaves no key is a bucket-level operation:
1436        // `PUT /<bucket>`, `PUT /<bucket>/` (the form the AWS SDKs actually
1437        // send for CreateBucket), `PUT /<bucket>//`, and the virtual-hosted
1438        // equivalents. Streaming those left the body unbuffered, so everything
1439        // that runs before the handler -- IAM enforcement above all -- saw an
1440        // empty body and could not read `CreateBucketConfiguration` (no
1441        // `aws:RequestTag` from a create-time tag set, for one).
1442        let host_bucket = protocol::parse_routing_host_from_headers(headers)
1443            .filter(|h| h.service == "s3")
1444            .and_then(|h| h.bucket);
1445        let routed = s3_routing_path(path, host_bucket.as_deref());
1446        let has_key = routed.split('/').filter(|seg| !seg.is_empty()).count() >= 2;
1447        if !has_key {
1448            return None;
1449        }
1450        let header_s3 = headers
1451            .get("authorization")
1452            .and_then(|v| v.to_str().ok())
1453            .and_then(fakecloud_aws::sigv4::parse_sigv4)
1454            .map(|info| info.service == "s3")
1455            .unwrap_or(false);
1456        let presigned_v4_s3 = query_params
1457            .get("X-Amz-Credential")
1458            .and_then(|c| c.split('/').nth(3).map(|s| s.to_string()))
1459            .map(|service| service == "s3")
1460            .unwrap_or(false);
1461        let presigned_v2 = query_params.contains_key("AWSAccessKeyId")
1462            && query_params.contains_key("Signature")
1463            && query_params.contains_key("Expires");
1464        if header_s3 || presigned_v4_s3 || presigned_v2 {
1465            return Some(("s3", ""));
1466        }
1467    }
1468
1469    None
1470}
1471
1472/// The SigV4 query-authentication parameters themselves. They authenticate
1473/// a presigned request and are not request headers.
1474const PRESIGN_AUTH_PARAMS: &[&str] = &[
1475    "x-amz-algorithm",
1476    "x-amz-credential",
1477    "x-amz-date",
1478    "x-amz-expires",
1479    "x-amz-signedheaders",
1480    "x-amz-signature",
1481    "x-amz-security-token",
1482];
1483
1484/// Copy the `x-amz-*` request headers a presigned URL carries in its query
1485/// string into `headers`, skipping the SigV4 auth parameters. A header the
1486/// client also sent directly wins over the query parameter. A non-ASCII user
1487/// metadata value is RFC 2047-encoded so S3 stores it as the UTF-8 the
1488/// percent-encoded query string meant, not as ISO-8859-1 header bytes.
1489fn hoist_presigned_query_headers(
1490    headers: &mut http::HeaderMap,
1491    query_params: &HashMap<String, String>,
1492) {
1493    for (key, value) in query_params {
1494        let lower = key.to_ascii_lowercase();
1495        if !lower.starts_with("x-amz-") || PRESIGN_AUTH_PARAMS.contains(&lower.as_str()) {
1496            continue;
1497        }
1498        let Ok(name) = http::HeaderName::from_bytes(lower.as_bytes()) else {
1499            continue;
1500        };
1501        if headers.contains_key(&name) {
1502            continue;
1503        }
1504        let value = if lower.starts_with("x-amz-meta-") {
1505            crate::rfc2047::encode(value)
1506        } else {
1507            value.clone()
1508        };
1509        if let Ok(value) = http::HeaderValue::from_str(&value) {
1510            headers.insert(name, value);
1511        }
1512    }
1513}
1514
1515/// Default request-body buffering cap. fakecloud reads the entire
1516/// request body into memory before handing it to a service handler,
1517/// so this ceiling caps RAM usage per in-flight request.
1518///
1519/// Default 1 GiB — comfortably above legitimate single S3 PutObject
1520/// payloads (AWS recommends multipart above ~100 MiB) and each
1521/// multipart part dispatches through here separately. Override with
1522/// `FAKECLOUD_MAX_REQUEST_BODY_BYTES` (decimal bytes) when running
1523/// stress tests that push past the default.
1524const DEFAULT_MAX_REQUEST_BODY_BYTES: usize = 1024 * 1024 * 1024;
1525
1526/// The server-wide buffered request body cap in bytes, from
1527/// `FAKECLOUD_MAX_REQUEST_BODY_BYTES` (default 1 GiB). Public so buffered
1528/// sub-proxies (e.g. the CloudFront viewer data plane) apply the SAME cap as
1529/// direct traffic rather than an inconsistent limit of their own.
1530pub fn max_request_body_bytes() -> usize {
1531    static CACHED: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
1532    *CACHED.get_or_init(|| {
1533        std::env::var("FAKECLOUD_MAX_REQUEST_BODY_BYTES")
1534            .ok()
1535            .and_then(|s| s.parse::<usize>().ok())
1536            .filter(|&n| n > 0)
1537            .unwrap_or(DEFAULT_MAX_REQUEST_BODY_BYTES)
1538    })
1539}
1540
1541/// For the cross-account decision in IAM enforcement, the "resource
1542/// account" is the ARN's account segment. Some services (notably S3)
1543/// produce ARNs with an empty account field — for those we return
1544/// `None` and let the caller fall back to the server's configured
1545/// account ID. Malformed or non-ARN strings also return `None`.
1546fn parse_account_from_arn(arn: &str) -> Option<String> {
1547    let mut parts = arn.splitn(6, ':');
1548    if parts.next()? != "arn" {
1549        return None;
1550    }
1551    let _partition = parts.next()?;
1552    let _service = parts.next()?;
1553    let _region = parts.next()?;
1554    let account = parts.next()?;
1555    // Resource segment must exist (parts.next().is_some()) for the ARN
1556    // to be well-formed, but we don't consume its value here.
1557    parts.next()?;
1558    if account.is_empty() {
1559        None
1560    } else {
1561        Some(account.to_string())
1562    }
1563}
1564
1565/// Whether the request's `user-agent` (or `x-amz-user-agent`) carries the
1566/// `api/neptune` SDK-metadata token the `aws-sdk-neptune` client stamps in.
1567/// Used to disambiguate Neptune from RDS, which share the `rds` SigV4 scope
1568/// and endpoint. Matches the `api/neptune` token that AWS SDKs emit (a
1569/// leading `api/neptune` optionally followed by `#`/`/` and a version).
1570fn user_agent_indicates_neptune(headers: &http::HeaderMap) -> bool {
1571    for name in ["user-agent", "x-amz-user-agent"] {
1572        if let Some(ua) = headers.get(name).and_then(|v| v.to_str().ok()) {
1573            for part in ua.split_whitespace() {
1574                if let Some(rest) = part.strip_prefix("api/neptune") {
1575                    if rest.is_empty() || rest.starts_with('#') || rest.starts_with('/') {
1576                        return true;
1577                    }
1578                }
1579            }
1580        }
1581    }
1582    false
1583}
1584
1585/// Extract region from User-Agent header suffix `region/<region>`.
1586/// Whether the request's `user-agent` (or `x-amz-user-agent`) carries the
1587/// `api/docdb` SDK-metadata token the `aws-sdk-docdb` client stamps in. Used
1588/// to disambiguate DocumentDB from RDS, which share the `rds` SigV4 scope
1589/// and endpoint. Matches the `api/docdb` token that AWS SDKs emit (a leading
1590/// `api/docdb` optionally followed by `#`/`/` and a version).
1591fn user_agent_indicates_docdb(headers: &http::HeaderMap) -> bool {
1592    for name in ["user-agent", "x-amz-user-agent"] {
1593        if let Some(ua) = headers.get(name).and_then(|v| v.to_str().ok()) {
1594            for part in ua.split_whitespace() {
1595                if let Some(rest) = part.strip_prefix("api/docdb") {
1596                    if rest.is_empty() || rest.starts_with('#') || rest.starts_with('/') {
1597                        return true;
1598                    }
1599                }
1600            }
1601        }
1602    }
1603    false
1604}
1605
1606fn extract_region_from_user_agent(headers: &http::HeaderMap) -> Option<String> {
1607    let ua = headers.get("user-agent")?.to_str().ok()?;
1608    for part in ua.split_whitespace() {
1609        if let Some(region) = part.strip_prefix("region/") {
1610            if !region.is_empty() {
1611                return Some(region.to_string());
1612            }
1613        }
1614    }
1615    None
1616}
1617
1618/// The wire shape of an error body. It follows the protocol, except within
1619/// REST-XML: S3 answers with a bare `<Error>` document, while CloudFront and
1620/// Route 53 wrap theirs in a namespaced `<ErrorResponse>`. The AWS SDKs parse
1621/// each service's errors strictly against its own shape.
1622#[derive(Clone, Copy, Debug)]
1623struct ErrorEnvelope {
1624    protocol: AwsProtocol,
1625    /// `Some(xmlns)` selects the `<ErrorResponse>` REST-XML shape.
1626    rest_xml_namespace: Option<&'static str>,
1627    /// An S3 Control request (`s3-control` host, served by the `s3`
1628    /// handler): its errors use S3 Control's un-namespaced `<ErrorResponse>`.
1629    s3_control: bool,
1630}
1631
1632impl ErrorEnvelope {
1633    /// The envelope for an error answering a request routed to `detected`.
1634    fn for_request(detected: &protocol::DetectedRequest, headers: &http::HeaderMap) -> Self {
1635        let rest_xml_namespace = if detected.protocol == AwsProtocol::Rest {
1636            fakecloud_aws::error::rest_xml_error_namespace(&detected.service)
1637        } else {
1638            None
1639        };
1640        Self {
1641            protocol: detected.protocol,
1642            rest_xml_namespace,
1643            s3_control: detected.protocol == AwsProtocol::Rest
1644                && detected.service == "s3"
1645                && protocol::is_s3_control_host(headers),
1646        }
1647    }
1648}
1649
1650impl From<AwsProtocol> for ErrorEnvelope {
1651    fn from(protocol: AwsProtocol) -> Self {
1652        Self {
1653            protocol,
1654            rest_xml_namespace: None,
1655            s3_control: false,
1656        }
1657    }
1658}
1659
1660fn build_error_response(
1661    status: StatusCode,
1662    code: &str,
1663    message: &str,
1664    request_id: &str,
1665    envelope: impl Into<ErrorEnvelope>,
1666) -> Response<Body> {
1667    build_error_response_with_fields(status, code, message, request_id, envelope, &[])
1668}
1669
1670fn build_error_response_with_fields(
1671    status: StatusCode,
1672    code: &str,
1673    message: &str,
1674    request_id: &str,
1675    envelope: impl Into<ErrorEnvelope>,
1676    extra_fields: &[(String, String)],
1677) -> Response<Body> {
1678    let envelope = envelope.into();
1679    let (status, content_type, body) = match (envelope.protocol, envelope.rest_xml_namespace) {
1680        // awsQuery services (SQS, SNS, IAM, STS, RDS, ELBv2, CloudWatch,
1681        // AutoScaling, ...) share the `<ErrorResponse>` envelope.
1682        (AwsProtocol::Query, _) => {
1683            fakecloud_aws::error::xml_error_response(status, code, message, request_id)
1684        }
1685        // EC2 uses the distinct ec2Query error envelope
1686        // (`<Response><Errors><Error>...</Errors><RequestID>`). Only EC2 is
1687        // classified `Ec2Query`, so the other Query-protocol services above are
1688        // untouched.
1689        (AwsProtocol::Ec2Query, _) => {
1690            fakecloud_aws::ec2query::ec2_error_response(status, code, message, request_id)
1691        }
1692        // CloudFront and Route 53 wrap the error in a namespaced
1693        // `<ErrorResponse>`; S3 (and the REST-XML fallbacks) use a bare `<Error>`.
1694        (AwsProtocol::Rest, Some(namespace)) => fakecloud_aws::error::rest_xml_error_response(
1695            status, code, message, request_id, namespace,
1696        ),
1697        (AwsProtocol::Rest, None) if envelope.s3_control => {
1698            fakecloud_aws::error::s3_control_xml_error_response(
1699                status,
1700                code,
1701                message,
1702                request_id,
1703                extra_fields,
1704            )
1705        }
1706        (AwsProtocol::Rest, None) => fakecloud_aws::error::s3_xml_error_response_with_fields(
1707            status,
1708            code,
1709            message,
1710            request_id,
1711            extra_fields,
1712        ),
1713        (AwsProtocol::Json | AwsProtocol::RestJson, _) => {
1714            fakecloud_aws::error::json_error_response_with_fields(
1715                status,
1716                code,
1717                message,
1718                extra_fields,
1719            )
1720        }
1721        (AwsProtocol::RpcV2Cbor, _) => (
1722            status,
1723            crate::cbor::CBOR_CONTENT_TYPE.to_string(),
1724            Bytes::from(crate::cbor::error_body(code, message, extra_fields)),
1725        ),
1726    };
1727
1728    // S3 (and other REST-XML services) place the error code in
1729    // `x-amz-error-code` so HEAD responses — which HTTP forbids from
1730    // carrying a body — still surface the code. AWS SDKs read this header
1731    // when the body is empty. Emit it on every error response so HEAD,
1732    // OPTIONS, and any client that strips the body still see the code.
1733    // Backend errors regularly include newlines (multi-line stderr from
1734    // docker/podman/etc.); HTTP header values reject control characters,
1735    // so sanitize before insertion or the builder rejects the response
1736    // and the connection drops.
1737    let safe_code = sanitize_header_value(code);
1738    let safe_message = sanitize_header_value(message);
1739    let mut builder = Response::builder()
1740        .status(status)
1741        .header("content-type", content_type)
1742        .header("x-amzn-requestid", request_id)
1743        .header("x-amz-request-id", request_id);
1744    if let Ok(v) = http::HeaderValue::from_str(&safe_code) {
1745        builder = builder.header("x-amz-error-code", v);
1746    }
1747    if let Ok(v) = http::HeaderValue::from_str(&safe_message) {
1748        builder = builder.header("x-amz-error-message", v);
1749    }
1750    if envelope.protocol == AwsProtocol::RpcV2Cbor {
1751        builder = builder.header(
1752            crate::cbor::SMITHY_PROTOCOL_HEADER,
1753            crate::cbor::RPC_V2_CBOR,
1754        );
1755    }
1756    builder.body(Body::from(body)).unwrap_or_else(|_| {
1757        // Builder only fails if a header is invalid; we sanitized the two
1758        // we control, so the remaining ones (content-type, request id) are
1759        // ASCII and safe. This fallback exists purely so we never panic.
1760        Response::new(Body::empty())
1761    })
1762}
1763
1764/// Strip characters that HTTP header values reject (control bytes, CR/LF/TAB)
1765/// and truncate to a length that AWS SDKs handle cleanly. Backend tools
1766/// (docker, podman, kubectl, …) emit multi-line stderr, and forwarding that
1767/// raw into `x-amz-error-message` previously panicked the dispatcher.
1768fn sanitize_header_value(s: &str) -> String {
1769    const MAX_LEN: usize = 1024;
1770    let mut out = String::with_capacity(s.len().min(MAX_LEN));
1771    for ch in s.chars() {
1772        if out.len() >= MAX_LEN {
1773            break;
1774        }
1775        // Header values forbid CR, LF, and other control bytes (RFC 9110).
1776        // Replace with a single space so multi-line messages stay readable.
1777        if ch.is_control() {
1778            if !out.ends_with(' ') {
1779                out.push(' ');
1780            }
1781        } else {
1782            out.push(ch);
1783        }
1784    }
1785    out.trim().to_string()
1786}
1787
1788/// Build the [`ConditionContext`] passed to the IAM evaluator for one
1789/// request. Populates the 10 global condition keys from the resolved
1790/// principal + the HTTP request. Service-specific keys are deferred to
1791/// a follow-up batch and left empty.
1792/// For an unsigned request that no other detection rule claimed, return the
1793/// bucket name when the first path segment names an existing S3 bucket.
1794///
1795/// fakecloud serves every service from one endpoint, so an anonymous
1796/// path-style S3 request (`GET /bucket/key`, no SigV4) is indistinguishable
1797/// from an API Gateway execute-api call by headers alone. Bucket existence is
1798/// the disambiguator: if the segment is a real bucket, route to S3; otherwise
1799/// fall through to the apigateway catch-all. Uses the already-wired
1800/// `resource_policy_provider`, which resolves S3 bucket ownership from state
1801/// (`Some` => the bucket exists). Returns `None` when no provider is wired.
1802/// Recover the access key from a SigV2 presigned URL. AWS SigV2 presigning
1803/// puts the key in the `AWSAccessKeyId` query parameter alongside `Signature`
1804/// and `Expires`; all three must be present for the URL to be a SigV2 presign.
1805/// Returns None for SigV4 presigns (which use `X-Amz-Credential`) or unsigned
1806/// requests.
1807fn sigv2_presigned_access_key(query_params: &HashMap<String, String>) -> Option<String> {
1808    if query_params.contains_key("Signature") && query_params.contains_key("Expires") {
1809        query_params.get("AWSAccessKeyId").cloned()
1810    } else {
1811        None
1812    }
1813}
1814
1815/// True when `s` is a 64-character lowercase-hex SHA-256 digest — the form a
1816/// signed `x-amz-content-sha256` payload hash takes. Distinguishes a real body
1817/// hash from the `UNSIGNED-PAYLOAD` / `STREAMING-*` markers (shorter and
1818/// containing non-hex characters), so only genuine hashes are re-bound to the
1819/// buffered body.
1820fn is_hex_sha256(s: &str) -> bool {
1821    s.len() == 64 && s.bytes().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
1822}
1823
1824/// Lowercase-hex SHA-256 of `bytes`, used to re-bind a signed
1825/// `x-amz-content-sha256` to the buffered request body for non-S3 services.
1826fn sha256_hex_lower(bytes: &[u8]) -> String {
1827    use sha2::{Digest, Sha256};
1828    let digest = Sha256::digest(bytes);
1829    const HEX: &[u8] = b"0123456789abcdef";
1830    let mut out = String::with_capacity(64);
1831    for b in digest {
1832        out.push(HEX[(b >> 4) as usize] as char);
1833        out.push(HEX[(b & 0x0f) as usize] as char);
1834    }
1835    out
1836}
1837
1838/// The condition context for evaluating `iam_action` as `principal`: the
1839/// principal's global keys, the session keys its credential carries (MFA,
1840/// token issue time, federated provider), the service's own condition keys,
1841/// and the resource / request / principal tags for ABAC. Callers add the
1842/// keys that depend on the resource's owner (see [`add_global_request_keys`]).
1843#[allow(clippy::too_many_arguments)]
1844fn principal_condition_context(
1845    principal: &Principal,
1846    resolved: Option<&crate::auth::ResolvedCredential>,
1847    service: &dyn crate::service::AwsService,
1848    aws_request: &AwsRequest,
1849    iam_action: &crate::auth::IamAction,
1850    detected: &protocol::DetectedRequest,
1851    remote_addr: Option<SocketAddr>,
1852) -> ConditionContext {
1853    let mut ctx = build_condition_context(
1854        principal,
1855        remote_addr,
1856        &aws_request.region,
1857        is_secure_transport(&aws_request.headers),
1858    );
1859    // F3 keys riding on the resolved credential. STS
1860    // populates these at mint time so subsequent
1861    // requests under the credential can be evaluated
1862    // against `aws:MultiFactorAuthPresent`,
1863    // `aws:MultiFactorAuthAge`, `aws:TokenIssueTime`,
1864    // and `aws:FederatedProvider`. IAM user access
1865    // keys carry none of these, matching AWS.
1866    if let Some(rc) = resolved {
1867        ctx.aws_mfa_present = Some(rc.mfa_present);
1868        ctx.aws_token_issue_time = rc.token_issued_at;
1869        ctx.aws_federated_provider = rc.federated_provider.clone();
1870        // `aws:MultiFactorAuthAge` is "seconds since
1871        // MFA was asserted" — computed at evaluation
1872        // time from the token issue moment so the
1873        // value increases monotonically as the session
1874        // ages. Only set when the session was actually
1875        // minted with MFA; otherwise the key is
1876        // absent, matching AWS.
1877        if rc.mfa_present {
1878            if let Some(issued) = rc.token_issued_at {
1879                let age = chrono::Utc::now()
1880                    .signed_duration_since(issued)
1881                    .num_seconds()
1882                    .max(0);
1883                ctx.aws_mfa_age_seconds = Some(age);
1884            }
1885        }
1886    }
1887    ctx.service_keys = service.iam_condition_keys_for(aws_request, iam_action);
1888
1889    // ABAC: populate tag-based condition keys.
1890    // aws:ResourceTag/*. An `iam:PassRole` resource is
1891    // an IAM role, not one of this service's resources,
1892    // so the service has no tags (or resource policy)
1893    // for it.
1894    let service_resource = !iam_action.is_pass_role();
1895    match service_resource
1896        .then(|| service.resource_tags_for(&iam_action.resource))
1897        .flatten()
1898    {
1899        Some(tags) => ctx.resource_tags = Some(tags),
1900        None => tracing::debug!(
1901            target: "fakecloud::iam::audit",
1902            service = %detected.service,
1903            resource = %iam_action.resource,
1904            "service does not expose resource tags for ABAC; skipping aws:ResourceTag/* evaluation"
1905        ),
1906    }
1907    // aws:RequestTag/* + aws:TagKeys
1908    match service.request_tags_from(aws_request, iam_action.action) {
1909        Some(tags) => ctx.request_tags = Some(tags),
1910        None => tracing::debug!(
1911            target: "fakecloud::iam::audit",
1912            service = %detected.service,
1913            action = %iam_action.action_string(),
1914            "service does not expose request tags for ABAC; skipping aws:RequestTag/* / aws:TagKeys evaluation"
1915        ),
1916    }
1917    // aws:PrincipalTag/*
1918    ctx.principal_tags = principal.tags.clone();
1919    ctx
1920}
1921
1922/// Populate `aws:ResourceAccount` (the account owning the target resource)
1923/// and, when the principal's account is in an organization,
1924/// `aws:PrincipalOrgID` / `aws:PrincipalOrgPaths`. A value the service
1925/// already supplied wins.
1926fn add_global_request_keys(
1927    ctx: &mut ConditionContext,
1928    principal: &Principal,
1929    resource_account_id: &str,
1930    scp_resolver: Option<&dyn crate::auth::ScpResolver>,
1931) {
1932    if !resource_account_id.is_empty() {
1933        ctx.service_keys
1934            .entry("aws:resourceaccount".to_string())
1935            .or_insert_with(|| vec![resource_account_id.to_string()]);
1936    }
1937    if let Some((org_id, path)) = scp_resolver.and_then(|r| r.principal_org(&principal.account_id))
1938    {
1939        ctx.service_keys
1940            .entry("aws:principalorgid".to_string())
1941            .or_insert_with(|| vec![org_id]);
1942        ctx.service_keys
1943            .entry("aws:principalorgpaths".to_string())
1944            .or_insert_with(|| vec![path]);
1945    }
1946}
1947
1948/// The AWS error for a request whose access key resolves to no identity:
1949/// `ExpiredToken` for a temporary credential that has expired, otherwise
1950/// `InvalidClientTokenId` (unknown or deactivated key, revoked session).
1951fn unresolved_credential_response(
1952    config: &DispatchConfig,
1953    access_key_id: &str,
1954    request_id: &str,
1955    envelope: ErrorEnvelope,
1956) -> Response<Body> {
1957    let expired = config
1958        .credential_resolver
1959        .as_ref()
1960        .is_some_and(|r| r.is_expired(access_key_id));
1961    let (code, message) = if expired {
1962        (
1963            "ExpiredToken",
1964            "The security token included in the request is expired",
1965        )
1966    } else {
1967        (
1968            "InvalidClientTokenId",
1969            "The security token included in the request is invalid",
1970        )
1971    };
1972    build_error_response(StatusCode::FORBIDDEN, code, message, request_id, envelope)
1973}
1974
1975fn anonymous_s3_bucket(uri: &http::Uri, config: &DispatchConfig) -> Option<String> {
1976    let provider = config.resource_policy_provider.as_ref()?;
1977    let segment = uri.path().split('/').find(|s| !s.is_empty())?.to_string();
1978    let arn = fakecloud_aws::arn::Arn::s3(&segment).to_string();
1979    provider.resource_owner_account("s3", &arn).map(|_| segment)
1980}
1981
1982/// The account an [`InternalCaller`] request works in: for S3, the owner of
1983/// the addressed bucket (bucket names are global, and a bucket policy can
1984/// grant a service acting for another account's resource); otherwise, or
1985/// when the bucket does not exist, the account the caller acts for.
1986fn internal_caller_account(
1987    caller: &InternalCaller,
1988    service: &str,
1989    path_segments: &[String],
1990    config: &DispatchConfig,
1991) -> String {
1992    let bucket_owner = (service == "s3")
1993        .then(|| path_segments.first())
1994        .flatten()
1995        .and_then(|bucket| {
1996            let arn = fakecloud_aws::arn::Arn::s3(bucket).to_string();
1997            config
1998                .resource_policy_provider
1999                .as_ref()?
2000                .resource_owner_account("s3", &arn)
2001        });
2002    bucket_owner.unwrap_or_else(|| caller.acting_account().to_string())
2003}
2004
2005/// Authorize a request fakecloud makes in-process as an AWS-owned principal
2006/// (see [`InternalCaller`]). Such a principal has no identity policies and no
2007/// account, so -- as in AWS -- only the resource policy (S3 bucket policy) can
2008/// grant it the action, with a public-read ACL honored the way it is for any
2009/// caller. The request context carries the caller's keys (`aws:SourceArn`,
2010/// `aws:SourceAccount`, ...) so a confused-deputy condition scoped to the
2011/// acting resource matches. `aws:SourceIp` is absent, as AWS omits it for a
2012/// request a service makes.
2013///
2014/// Returns the error response to send when the request is denied under
2015/// strict mode; soft mode only logs.
2016fn authorize_internal_caller(
2017    caller: &InternalCaller,
2018    service: &dyn crate::service::AwsService,
2019    aws_request: &AwsRequest,
2020    evaluator: &dyn IamPolicyEvaluator,
2021    config: &DispatchConfig,
2022    detected: &protocol::DetectedRequest,
2023    request_id: &str,
2024) -> Option<Response<Body>> {
2025    let principal = caller.principal();
2026    let denied = || {
2027        config.iam_mode.is_strict().then(|| {
2028            build_error_response(
2029                StatusCode::FORBIDDEN,
2030                "AccessDenied",
2031                "Access Denied",
2032                request_id,
2033                ErrorEnvelope::for_request(detected, &aws_request.headers),
2034            )
2035        })
2036    };
2037    let iam_actions = service.iam_actions_for(aws_request);
2038    if iam_actions.is_empty() {
2039        tracing::warn!(
2040            target: "fakecloud::iam::audit",
2041            service = %detected.service,
2042            action = %aws_request.action,
2043            principal = %principal.arn,
2044            mode = %config.iam_mode,
2045            request_id = %request_id,
2046            "service-principal request has no IamAction mapping; denying under strict, allowing under soft"
2047        );
2048        return denied();
2049    }
2050    // A service acting on AWS's behalf is not subject to `iam:PassRole`:
2051    // that check binds the customer principal handing a role to a service.
2052    for iam_action in iam_actions.iter().filter(|a| !a.is_pass_role()) {
2053        let now = chrono::Utc::now();
2054        let mut context = ConditionContext {
2055            aws_principal_arn: Some(principal.arn.clone()),
2056            aws_current_time: Some(now),
2057            aws_epoch_time: Some(now.timestamp()),
2058            aws_secure_transport: Some(is_secure_transport(&aws_request.headers)),
2059            aws_requested_region: Some(aws_request.region.clone()),
2060            ..Default::default()
2061        };
2062        context.service_keys = service.iam_condition_keys_for(aws_request, iam_action);
2063        context.service_keys.extend(caller.condition_keys());
2064        // `aws:ResourceAccount`: the account the service works in for this
2065        // request (for S3, the addressed bucket's owner).
2066        context
2067            .service_keys
2068            .entry("aws:resourceaccount".to_string())
2069            .or_insert_with(|| vec![aws_request.account_id.clone()]);
2070        let resource_policy_json = config
2071            .resource_policy_provider
2072            .as_ref()
2073            .and_then(|p| p.resource_policy(&detected.service, &iam_action.resource));
2074        let decision = evaluator.evaluate_resource_policy_only(
2075            &principal,
2076            iam_action,
2077            &context,
2078            resource_policy_json.as_deref(),
2079        );
2080        let explicit_deny = matches!(decision, crate::auth::IamDecision::ExplicitDeny);
2081        let acl_allows = !explicit_deny
2082            && config.resource_policy_provider.as_ref().is_some_and(|p| {
2083                p.public_acl_allows(&detected.service, &iam_action.resource, iam_action.action)
2084            });
2085        if !decision.is_allow() && !acl_allows {
2086            tracing::warn!(
2087                target: "fakecloud::iam::audit",
2088                service = %detected.service,
2089                action = %iam_action.action_string(),
2090                resource = %iam_action.resource,
2091                principal = %principal.arn,
2092                resource_policy_present = resource_policy_json.is_some(),
2093                decision = ?decision,
2094                mode = %config.iam_mode,
2095                request_id = %request_id,
2096                "service-principal request denied: the resource policy does not grant the action"
2097            );
2098            if let Some(resp) = denied() {
2099                return Some(resp);
2100            }
2101        }
2102    }
2103    None
2104}
2105
2106/// Authorize an account root reaching a resource ANOTHER account owns.
2107///
2108/// Root needs no identity policy in its own account, which is why dispatch
2109/// exempts it from identity evaluation. That exemption ends at the account
2110/// boundary: on AWS, root of account B reaching account A's bucket (or table)
2111/// is a cross-account request like any other, allowed only when A's resource
2112/// policy grants it. This applies where the resource policy provider resolves
2113/// the owner of a resource (S3 buckets, whose names are global and whose ARNs
2114/// carry no account, and DynamoDB tables and streams) -- the resource-policy
2115/// governed data services whose handlers serve a request in the owner's
2116/// account. A same-account resource, or one no provider claims, keeps the
2117/// root exemption.
2118///
2119/// Returns the error response to send when the request is denied under
2120/// strict mode; soft mode only logs.
2121#[allow(clippy::too_many_arguments)]
2122fn authorize_root_cross_account(
2123    principal: &Principal,
2124    service: &dyn crate::service::AwsService,
2125    aws_request: &AwsRequest,
2126    evaluator: &dyn IamPolicyEvaluator,
2127    config: &DispatchConfig,
2128    detected: &protocol::DetectedRequest,
2129    request_id: &str,
2130    remote_addr: Option<SocketAddr>,
2131    resolved: Option<&crate::auth::ResolvedCredential>,
2132) -> Option<Response<Body>> {
2133    if service.iam_resource_in_caller_account(aws_request) {
2134        return None;
2135    }
2136    let provider = config.resource_policy_provider.as_ref()?;
2137    for iam_action in service.iam_actions_for(aws_request) {
2138        let Some(owner) = provider.resource_owner_account(&detected.service, &iam_action.resource)
2139        else {
2140            continue;
2141        };
2142        if owner == principal.account_id {
2143            continue;
2144        }
2145        let mut context = principal_condition_context(
2146            principal,
2147            resolved,
2148            service,
2149            aws_request,
2150            &iam_action,
2151            detected,
2152            remote_addr,
2153        );
2154        add_global_request_keys(
2155            &mut context,
2156            principal,
2157            &owner,
2158            config.scp_resolver.as_deref(),
2159        );
2160        let resource_policy_json =
2161            provider.resource_policy(&detected.service, &iam_action.resource);
2162        let decision = evaluator.evaluate_resource_policy_only(
2163            principal,
2164            &iam_action,
2165            &context,
2166            resource_policy_json.as_deref(),
2167        );
2168        let explicit_deny = matches!(decision, crate::auth::IamDecision::ExplicitDeny);
2169        let acl_allows = !explicit_deny
2170            && provider.public_acl_allows(
2171                &detected.service,
2172                &iam_action.resource,
2173                iam_action.action,
2174            );
2175        if decision.is_allow() || acl_allows {
2176            continue;
2177        }
2178        tracing::warn!(
2179            target: "fakecloud::iam::audit",
2180            service = %detected.service,
2181            action = %iam_action.action_string(),
2182            resource = %iam_action.resource,
2183            principal = %principal.arn,
2184            resource_account = %owner,
2185            resource_policy_present = resource_policy_json.is_some(),
2186            decision = ?decision,
2187            mode = %config.iam_mode,
2188            request_id = %request_id,
2189            "cross-account root request denied: the resource policy does not grant the action"
2190        );
2191        if config.iam_mode.is_strict() {
2192            return Some(build_error_response(
2193                StatusCode::FORBIDDEN,
2194                "AccessDeniedException",
2195                &format!(
2196                    "User: {} is not authorized to perform: {} on resource: {} because no resource-based policy allows the {} action",
2197                    principal.arn,
2198                    iam_action.action_string(),
2199                    iam_action.resource,
2200                    iam_action.action_string(),
2201                ),
2202                request_id,
2203                ErrorEnvelope::for_request(detected, &aws_request.headers),
2204            ));
2205        }
2206    }
2207    None
2208}
2209
2210fn build_condition_context(
2211    principal: &Principal,
2212    remote_addr: Option<SocketAddr>,
2213    region: &str,
2214    secure_transport: bool,
2215) -> ConditionContext {
2216    let now = chrono::Utc::now();
2217    ConditionContext {
2218        aws_username: aws_username_from_principal(principal),
2219        aws_userid: Some(principal.user_id.clone()),
2220        aws_principal_arn: Some(principal.arn.clone()),
2221        aws_principal_account: Some(principal.account_id.clone()),
2222        aws_principal_type: Some(principal_type_label(principal.principal_type).to_string()),
2223        aws_source_ip: remote_addr.map(|sa| sa.ip()),
2224        aws_current_time: Some(now),
2225        aws_epoch_time: Some(now.timestamp()),
2226        aws_secure_transport: Some(secure_transport),
2227        aws_requested_region: Some(region.to_string()),
2228        // F3 keys: populated from the caller's session context when STS
2229        // mints credentials with MFA / SAML / OIDC / VPC-endpoint hints.
2230        // Default-None here so tests/dispatch sites that don't set them
2231        // safe-fail any policy referencing them — matching AWS for keys
2232        // that aren't asserted.
2233        aws_mfa_present: None,
2234        aws_mfa_age_seconds: None,
2235        aws_called_via: Vec::new(),
2236        aws_source_vpce: None,
2237        aws_source_vpc: None,
2238        aws_vpc_source_ip: None,
2239        aws_federated_provider: None,
2240        aws_token_issue_time: None,
2241        service_keys: Default::default(),
2242        resource_tags: None,
2243        request_tags: None,
2244        principal_tags: None,
2245    }
2246}
2247
2248/// `aws:username` is only set for IAM users, matching AWS. For assumed
2249/// roles, federated users, root, and unknown principals the key is
2250/// absent — operators that reference it without `IfExists` safe-fail.
2251fn aws_username_from_principal(principal: &Principal) -> Option<String> {
2252    if principal.principal_type != PrincipalType::User {
2253        return None;
2254    }
2255    let after = principal.arn.rsplit_once(":user/").map(|(_, s)| s)?;
2256    // Strip any IAM path prefix; bare username is the last segment.
2257    Some(after.rsplit('/').next().unwrap_or(after).to_string())
2258}
2259
2260/// AWS's `aws:PrincipalType` uses PascalCase identifiers, distinct from
2261/// the lowercase ones [`PrincipalType::as_str`] returns for ARNs.
2262fn principal_type_label(t: PrincipalType) -> &'static str {
2263    match t {
2264        PrincipalType::User => "User",
2265        PrincipalType::AssumedRole => "AssumedRole",
2266        PrincipalType::FederatedUser => "FederatedUser",
2267        PrincipalType::Root => "Account",
2268        PrincipalType::Unknown => "Unknown",
2269        PrincipalType::Service => "Service",
2270    }
2271}
2272
2273/// Best-effort detection of TLS-terminated requests. Direct HTTPS
2274/// connections are not yet supported by the fakecloud server (it speaks
2275/// plain HTTP), so the only signal is an `x-forwarded-proto: https`
2276/// header set by an upstream proxy. Anything else evaluates to `false`,
2277/// which matches the typical local-dev setup.
2278fn is_secure_transport(headers: &http::HeaderMap) -> bool {
2279    headers
2280        .get("x-forwarded-proto")
2281        .and_then(|v| v.to_str().ok())
2282        .map(|s| s.eq_ignore_ascii_case("https"))
2283        .unwrap_or(false)
2284}
2285
2286trait ProtocolExt {
2287    fn error_status(&self) -> StatusCode;
2288}
2289
2290impl ProtocolExt for AwsProtocol {
2291    fn error_status(&self) -> StatusCode {
2292        StatusCode::BAD_REQUEST
2293    }
2294}
2295
2296/// Whether a (possibly percent-encoded) `/tags/{resourceArn}` label is the ARN
2297/// of a Bedrock Agents runtime session (`...:session/<id>`).
2298fn names_bedrock_session(label: &str) -> bool {
2299    let decoded = label
2300        .to_ascii_lowercase()
2301        .replace("%3a", ":")
2302        .replace("%2f", "/");
2303    decoded.starts_with("arn:") && decoded.contains(":bedrock:") && decoded.contains(":session/")
2304}
2305
2306/// Which Bedrock Agents handler a `bedrock`-scoped request belongs to:
2307/// `bedrock-agent-runtime` or `bedrock-agent` for their path families, `None`
2308/// for any other Bedrock (bedrock-runtime / control plane) path. Runtime and
2309/// control plane share `/agents/...` and `/flows/...` prefixes, so the split
2310/// is by path shape and, where the shapes coincide, by method.
2311fn bedrock_agent_service_for(method: &http::Method, path: &str) -> Option<&'static str> {
2312    let first_seg = path.split('/').nth(1);
2313    if !matches!(
2314        first_seg,
2315        Some(
2316            "agents"
2317                | "knowledgebases"
2318                | "flows"
2319                | "prompts"
2320                | "tags"
2321                | "retrieveAndGenerate"
2322                | "retrieveAndGenerateStream"
2323                | "optimize-prompt"
2324                | "sessions"
2325                | "invocations"
2326                | "generate-query"
2327                | "rerank"
2328        )
2329    ) {
2330        return None;
2331    }
2332    let segs: Vec<&str> = path.split('/').collect();
2333    let is_runtime = matches!(
2334        segs.as_slice(),
2335        ["", "agents", _, "agentAliases", _, ..]  // InvokeAgent
2336            | ["", "flows", _, "executions"] // ListFlowExecutions
2337            | ["", "flows", _, "aliases", _, "executions", ..] // flow executions
2338            | ["", "knowledgebases", _, "retrieve"] // Retrieve
2339            | ["", "retrieveAndGenerate"]
2340            | ["", "retrieveAndGenerateStream"]
2341            | ["", "optimize-prompt"]
2342            | ["", "sessions", ..]
2343            | ["", "invocations", ..]
2344            | ["", "generate-query"]
2345            | ["", "rerank"]
2346    ) || (
2347        // InvokeFlow is POST /flows/{flow}/aliases/{alias}; the same path
2348        // under GET / PUT / DELETE is the control plane's flow-alias CRUD.
2349        *method == http::Method::POST && matches!(segs.as_slice(), ["", "flows", _, "aliases", _])
2350    ) || (
2351        // InvokeInlineAgent is POST /agents/{sessionId}; the control plane's
2352        // PrepareAgent is POST /agents/{agentId}/ (trailing slash).
2353        *method == http::Method::POST
2354            && matches!(segs.as_slice(), ["", "agents", id] if !id.is_empty())
2355    ) || path
2356        // Tagging is shared: a runtime resource (a session) is tagged through
2357        // the runtime, an agent/flow/knowledge-base resource through the
2358        // control plane.
2359        .strip_prefix("/tags/")
2360        .is_some_and(names_bedrock_session);
2361    Some(if is_runtime {
2362        "bedrock-agent-runtime"
2363    } else {
2364        "bedrock-agent"
2365    })
2366}
2367
2368#[cfg(test)]
2369mod tests {
2370
2371    fn gzip(data: &[u8]) -> Bytes {
2372        use std::io::Write;
2373        let mut enc = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
2374        enc.write_all(data).unwrap();
2375        Bytes::from(enc.finish().unwrap())
2376    }
2377
2378    fn gzip_headers(extra: &[(&'static str, &str)]) -> http::HeaderMap {
2379        let mut h = http::HeaderMap::new();
2380        h.insert("content-encoding", "gzip".parse().unwrap());
2381        for (k, v) in extra {
2382            h.insert(*k, v.parse().unwrap());
2383        }
2384        h
2385    }
2386
2387    #[test]
2388    fn request_compression_decodes_gzip_for_cloudwatch() {
2389        let body = br#"{"Namespace":"App"}"#;
2390        // awsJson (X-Amz-Target).
2391        let h = gzip_headers(&[(
2392            "x-amz-target",
2393            "GraniteServiceVersion20100801.PutMetricData",
2394        )]);
2395        assert_eq!(
2396            decode_request_compression(&h, None, gzip(body)).unwrap(),
2397            Bytes::from_static(body)
2398        );
2399        // awsQuery (SigV4 scope only).
2400        let h = gzip_headers(&[(
2401            "authorization",
2402            "AWS4-HMAC-SHA256 Credential=test/20240101/us-east-1/monitoring/aws4_request, SignedHeaders=host, Signature=0",
2403        )]);
2404        let form = b"Action=PutMetricData&Namespace=App";
2405        assert_eq!(
2406            decode_request_compression(&h, None, gzip(form)).unwrap(),
2407            Bytes::from_static(form)
2408        );
2409        // rpcv2Cbor.
2410        let detected = protocol::DetectedRequest {
2411            service: "monitoring".to_string(),
2412            action: "PutMetricData".to_string(),
2413            protocol: AwsProtocol::RpcV2Cbor,
2414        };
2415        let h = gzip_headers(&[]);
2416        assert_eq!(
2417            decode_request_compression(&h, Some(&detected), gzip(&[0xa0])).unwrap(),
2418            Bytes::from_static(&[0xa0])
2419        );
2420        // Corrupt gzip is a serialization error in the caller's protocol.
2421        let err = decode_request_compression(&h, Some(&detected), Bytes::from_static(b"nope"))
2422            .unwrap_err();
2423        assert_eq!(err.1, AwsProtocol::RpcV2Cbor);
2424    }
2425
2426    #[test]
2427    fn request_compression_leaves_other_services_alone() {
2428        // S3 keeps Content-Encoding as object metadata: the body is stored as sent.
2429        let h = gzip_headers(&[(
2430            "authorization",
2431            "AWS4-HMAC-SHA256 Credential=test/20240101/us-east-1/s3/aws4_request, SignedHeaders=host, Signature=0",
2432        )]);
2433        let body = gzip(b"object bytes");
2434        assert_eq!(
2435            decode_request_compression(&h, None, body.clone()).unwrap(),
2436            body
2437        );
2438        // No Content-Encoding: untouched.
2439        let h = http::HeaderMap::new();
2440        assert_eq!(
2441            decode_request_compression(&h, None, Bytes::from_static(b"x")).unwrap(),
2442            Bytes::from_static(b"x")
2443        );
2444    }
2445
2446    #[test]
2447    fn request_compression_services_match_the_models() {
2448        // Every vendored model with an `@requestCompression` operation must be
2449        // listed (by registry name) in REQUEST_COMPRESSION_SERVICES.
2450        let dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../aws-models");
2451        let mut with_trait: Vec<String> = std::fs::read_dir(&dir)
2452            .unwrap_or_else(|e| panic!("read {}: {e}", dir.display()))
2453            .filter_map(|e| e.ok())
2454            .filter(|e| {
2455                std::fs::read_to_string(e.path())
2456                    .is_ok_and(|s| s.contains("\"smithy.api#requestCompression\""))
2457            })
2458            .map(|e| e.file_name().to_string_lossy().into_owned())
2459            .collect();
2460        with_trait.sort();
2461        assert_eq!(
2462            with_trait,
2463            vec!["cloudwatch.json".to_string()],
2464            "update REQUEST_COMPRESSION_SERVICES for new @requestCompression models"
2465        );
2466        assert_eq!(REQUEST_COMPRESSION_SERVICES, &["monitoring"]);
2467    }
2468    #[test]
2469    fn bedrock_agent_paths_split_between_runtime_and_control_plane() {
2470        use http::Method;
2471        let runtime = Some("bedrock-agent-runtime");
2472        let agent = Some("bedrock-agent");
2473        for (method, path, want) in [
2474            (Method::POST, "/flows/F/aliases/A", runtime),
2475            (Method::GET, "/flows/F/aliases/A", agent),
2476            (Method::PUT, "/flows/F/aliases/A", agent),
2477            (Method::DELETE, "/flows/F/aliases/A", agent),
2478            (Method::POST, "/flows/F/aliases/A/executions", runtime),
2479            (Method::GET, "/flows/F/aliases/A/executions/E", runtime),
2480            (
2481                Method::POST,
2482                "/flows/F/aliases/A/executions/E/stop",
2483                runtime,
2484            ),
2485            (
2486                Method::GET,
2487                "/flows/F/aliases/A/executions/E/events",
2488                runtime,
2489            ),
2490            (
2491                Method::GET,
2492                "/flows/F/aliases/A/executions/E/flowsnapshot",
2493                runtime,
2494            ),
2495            (Method::GET, "/flows/F/executions", runtime),
2496            (Method::GET, "/flows/F/aliases", agent),
2497            (Method::POST, "/flows/F/versions", agent),
2498            (Method::GET, "/flows/F", agent),
2499            (
2500                Method::POST,
2501                "/agents/X/agentAliases/Y/sessions/S/text",
2502                runtime,
2503            ),
2504            (Method::GET, "/agents/X", agent),
2505            (Method::POST, "/agents/session-1", runtime),
2506            (Method::POST, "/agents/X/", agent),
2507            (Method::GET, "/agents/X/", agent),
2508            (Method::PUT, "/agents/X/", agent),
2509            (
2510                Method::POST,
2511                "/tags/arn%3Aaws%3Abedrock%3Aus-east-1%3A123456789012%3Asession%2F0f1e2d3c-4b5a-6978-8a9b-0c1d2e3f4a5b",
2512                runtime,
2513            ),
2514            (
2515                Method::GET,
2516                "/tags/arn:aws:bedrock:us-east-1:123456789012:session/0f1e2d3c-4b5a-6978-8a9b-0c1d2e3f4a5b",
2517                runtime,
2518            ),
2519            (
2520                Method::DELETE,
2521                "/tags/arn%3aaws%3abedrock%3aus-east-1%3a123456789012%3asession%2fabc",
2522                runtime,
2523            ),
2524            (
2525                Method::POST,
2526                "/tags/arn%3Aaws%3Abedrock%3Aus-east-1%3A123456789012%3Aagent%2FAGENT12345",
2527                agent,
2528            ),
2529            (
2530                Method::GET,
2531                "/tags/arn%3Aaws%3Abedrock%3Aus-east-1%3A123456789012%3Aflow%2FFLOW123456",
2532                agent,
2533            ),
2534            (Method::POST, "/model/m/invoke", None),
2535        ] {
2536            assert_eq!(
2537                bedrock_agent_service_for(&method, path),
2538                want,
2539                "{method} {path}"
2540            );
2541        }
2542    }
2543
2544    use super::*;
2545
2546    #[test]
2547    fn default_max_request_body_bytes_is_one_gib() {
2548        // Without the env override, the cap defaults to 1 GiB. The
2549        // public function caches via OnceLock so only the first call
2550        // in the process matters; we assert the constant directly.
2551        assert_eq!(DEFAULT_MAX_REQUEST_BODY_BYTES, 1024 * 1024 * 1024);
2552    }
2553
2554    #[test]
2555    fn sigv2_presigned_access_key_extracted_with_signature_and_expires() {
2556        let mut q = HashMap::new();
2557        q.insert("AWSAccessKeyId".to_string(), "AKIAEXAMPLE".to_string());
2558        q.insert("Signature".to_string(), "abc%2Bdef".to_string());
2559        q.insert("Expires".to_string(), "1700000000".to_string());
2560        assert_eq!(
2561            sigv2_presigned_access_key(&q).as_deref(),
2562            Some("AKIAEXAMPLE")
2563        );
2564    }
2565
2566    #[test]
2567    fn sigv2_presigned_access_key_none_without_signature_or_expires() {
2568        // AWSAccessKeyId alone (e.g. a stray query param) is not a SigV2
2569        // presign and must not be treated as a credential.
2570        let mut q = HashMap::new();
2571        q.insert("AWSAccessKeyId".to_string(), "AKIAEXAMPLE".to_string());
2572        assert_eq!(sigv2_presigned_access_key(&q), None);
2573
2574        q.insert("Expires".to_string(), "1700000000".to_string());
2575        assert_eq!(
2576            sigv2_presigned_access_key(&q),
2577            None,
2578            "missing Signature must not qualify"
2579        );
2580    }
2581
2582    #[test]
2583    fn sigv2_presigned_access_key_none_for_unsigned_request() {
2584        assert_eq!(sigv2_presigned_access_key(&HashMap::new()), None);
2585    }
2586
2587    #[test]
2588    fn is_hex_sha256_accepts_real_digest_rejects_markers() {
2589        // A genuine 64-char lowercase-hex digest is a bindable body hash.
2590        assert!(is_hex_sha256(&sha256_hex_lower(b"hello")));
2591        assert!(is_hex_sha256(
2592            "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
2593        ));
2594        // The SigV4 payload markers are NOT hashes and must be skipped.
2595        assert!(!is_hex_sha256("UNSIGNED-PAYLOAD"));
2596        assert!(!is_hex_sha256("STREAMING-AWS4-HMAC-SHA256-PAYLOAD"));
2597        assert!(!is_hex_sha256("STREAMING-UNSIGNED-PAYLOAD-TRAILER"));
2598        // Wrong length or non-lowercase-hex characters are rejected.
2599        assert!(!is_hex_sha256("abc123"));
2600        assert!(!is_hex_sha256(
2601            "E3B0C44298FC1C149AFBF4C8996FB92427AE41E4649B934CA495991B7852B855"
2602        ));
2603    }
2604
2605    #[test]
2606    fn sha256_hex_lower_matches_known_vectors() {
2607        // Empty input -> the well-known SHA-256 of "".
2608        assert_eq!(
2609            sha256_hex_lower(b""),
2610            "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
2611        );
2612        assert_eq!(
2613            sha256_hex_lower(b"abc"),
2614            "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad"
2615        );
2616        assert_eq!(sha256_hex_lower(b"abc").len(), 64);
2617    }
2618
2619    #[test]
2620    fn global_request_keys_cover_resource_account_and_org() {
2621        struct Org;
2622        impl crate::auth::ScpResolver for Org {
2623            fn scps_for(&self, _: &Principal) -> Option<Vec<String>> {
2624                None
2625            }
2626            fn principal_org(&self, account: &str) -> Option<(String, String)> {
2627                (account == "111111111111")
2628                    .then(|| ("o-abc".to_string(), "o-abc/r-1/ou-2/".to_string()))
2629            }
2630        }
2631        let principal = |account: &str| Principal {
2632            arn: format!("arn:aws:iam::{account}:user/u"),
2633            user_id: "AIDAU".into(),
2634            account_id: account.into(),
2635            principal_type: PrincipalType::User,
2636            source_identity: None,
2637            tags: None,
2638        };
2639        let mut ctx = ConditionContext::default();
2640        add_global_request_keys(
2641            &mut ctx,
2642            &principal("111111111111"),
2643            "222222222222",
2644            Some(&Org),
2645        );
2646        assert_eq!(
2647            ctx.lookup("aws:ResourceAccount"),
2648            Some(vec!["222222222222".into()])
2649        );
2650        assert_eq!(ctx.lookup("aws:PrincipalOrgID"), Some(vec!["o-abc".into()]));
2651        assert_eq!(
2652            ctx.lookup("aws:PrincipalOrgPaths"),
2653            Some(vec!["o-abc/r-1/ou-2/".into()])
2654        );
2655        // A principal in no organization gets no org keys (AWS omits them).
2656        let mut ctx = ConditionContext::default();
2657        add_global_request_keys(
2658            &mut ctx,
2659            &principal("333333333333"),
2660            "333333333333",
2661            Some(&Org),
2662        );
2663        assert_eq!(ctx.lookup("aws:PrincipalOrgID"), None);
2664        assert_eq!(
2665            ctx.lookup("aws:ResourceAccount"),
2666            Some(vec!["333333333333".into()])
2667        );
2668    }
2669
2670    #[test]
2671    fn unresolved_credential_distinguishes_expired_from_invalid() {
2672        struct Expired;
2673        impl CredentialResolver for Expired {
2674            fn resolve(&self, _: &str) -> Option<crate::auth::ResolvedCredential> {
2675                None
2676            }
2677            fn is_expired(&self, akid: &str) -> bool {
2678                akid == "ASIAEXPIRED"
2679            }
2680        }
2681        let mut cfg = DispatchConfig::new("us-east-1", "123456789012");
2682        cfg.credential_resolver = Some(Arc::new(Expired));
2683        let detected = protocol::DetectedRequest {
2684            service: "sts".into(),
2685            action: "GetCallerIdentity".into(),
2686            protocol: AwsProtocol::Query,
2687        };
2688        let code = |akid: &str| {
2689            let resp = unresolved_credential_response(
2690                &cfg,
2691                akid,
2692                "r",
2693                ErrorEnvelope::for_request(&detected, &http::HeaderMap::new()),
2694            );
2695            assert_eq!(resp.status(), StatusCode::FORBIDDEN);
2696            resp.headers()["x-amz-error-code"]
2697                .to_str()
2698                .unwrap()
2699                .to_string()
2700        };
2701        assert_eq!(code("ASIAEXPIRED"), "ExpiredToken");
2702        assert_eq!(code("AKIAUNKNOWN"), "InvalidClientTokenId");
2703    }
2704
2705    #[test]
2706    fn dispatch_config_new_defaults_to_off() {
2707        let cfg = DispatchConfig::new("us-east-1", "123456789012");
2708        assert_eq!(cfg.region, "us-east-1");
2709        assert_eq!(cfg.account_id, "123456789012");
2710        assert!(!cfg.verify_sigv4);
2711        assert_eq!(cfg.iam_mode, IamMode::Off);
2712    }
2713
2714    #[test]
2715    fn aws_username_strips_iam_path_for_users() {
2716        let p = Principal {
2717            arn: "arn:aws:iam::123456789012:user/engineering/alice".into(),
2718            user_id: "AIDAALICE".into(),
2719            account_id: "123456789012".into(),
2720            principal_type: PrincipalType::User,
2721            source_identity: None,
2722            tags: None,
2723        };
2724        assert_eq!(aws_username_from_principal(&p), Some("alice".into()));
2725    }
2726
2727    #[test]
2728    fn aws_username_unset_for_assumed_role() {
2729        let p = Principal {
2730            arn: "arn:aws:sts::123456789012:assumed-role/ops/session".into(),
2731            user_id: "AROAOPS:session".into(),
2732            account_id: "123456789012".into(),
2733            principal_type: PrincipalType::AssumedRole,
2734            source_identity: None,
2735            tags: None,
2736        };
2737        assert_eq!(aws_username_from_principal(&p), None);
2738    }
2739
2740    #[test]
2741    fn principal_type_label_matches_aws_casing() {
2742        assert_eq!(principal_type_label(PrincipalType::User), "User");
2743        assert_eq!(
2744            principal_type_label(PrincipalType::AssumedRole),
2745            "AssumedRole"
2746        );
2747        assert_eq!(principal_type_label(PrincipalType::Root), "Account");
2748    }
2749
2750    #[test]
2751    fn build_condition_context_populates_global_keys() {
2752        let p = Principal {
2753            arn: "arn:aws:iam::123456789012:user/alice".into(),
2754            user_id: "AIDAALICE".into(),
2755            account_id: "123456789012".into(),
2756            principal_type: PrincipalType::User,
2757            source_identity: None,
2758            tags: None,
2759        };
2760        let addr: SocketAddr = "10.0.0.1:54321".parse().unwrap();
2761        let ctx = build_condition_context(&p, Some(addr), "us-east-1", false);
2762        assert_eq!(ctx.aws_username.as_deref(), Some("alice"));
2763        assert_eq!(ctx.aws_userid.as_deref(), Some("AIDAALICE"));
2764        assert_eq!(
2765            ctx.aws_principal_arn.as_deref(),
2766            Some("arn:aws:iam::123456789012:user/alice")
2767        );
2768        assert_eq!(ctx.aws_principal_account.as_deref(), Some("123456789012"));
2769        assert_eq!(ctx.aws_principal_type.as_deref(), Some("User"));
2770        assert_eq!(
2771            ctx.aws_source_ip.map(|i| i.to_string()).as_deref(),
2772            Some("10.0.0.1")
2773        );
2774        assert_eq!(ctx.aws_requested_region.as_deref(), Some("us-east-1"));
2775        assert_eq!(ctx.aws_secure_transport, Some(false));
2776        assert!(ctx.aws_current_time.is_some());
2777        assert!(ctx.aws_epoch_time.is_some());
2778    }
2779
2780    #[test]
2781    fn is_secure_transport_reads_x_forwarded_proto() {
2782        let mut headers = http::HeaderMap::new();
2783        headers.insert("x-forwarded-proto", "https".parse().unwrap());
2784        assert!(is_secure_transport(&headers));
2785        headers.insert("x-forwarded-proto", "http".parse().unwrap());
2786        assert!(!is_secure_transport(&headers));
2787        let empty = http::HeaderMap::new();
2788        assert!(!is_secure_transport(&empty));
2789    }
2790
2791    #[test]
2792    fn parse_account_from_arn_extracts_standard_shapes() {
2793        assert_eq!(
2794            parse_account_from_arn("arn:aws:sqs:us-east-1:123456789012:queue"),
2795            Some("123456789012".to_string())
2796        );
2797        assert_eq!(
2798            parse_account_from_arn("arn:aws:iam::123456789012:user/alice"),
2799            Some("123456789012".to_string())
2800        );
2801    }
2802
2803    #[test]
2804    fn parse_account_from_arn_returns_none_for_s3_empty_account() {
2805        // S3 ARNs have both region and account empty.
2806        assert_eq!(parse_account_from_arn("arn:aws:s3:::my-bucket"), None);
2807        assert_eq!(
2808            parse_account_from_arn("arn:aws:s3:::my-bucket/path/to/key"),
2809            None
2810        );
2811    }
2812
2813    #[test]
2814    fn parse_account_from_arn_returns_none_for_malformed() {
2815        assert_eq!(parse_account_from_arn(""), None);
2816        assert_eq!(parse_account_from_arn("not-an-arn"), None);
2817        assert_eq!(parse_account_from_arn("arn:aws:sqs:us-east-1"), None);
2818        assert_eq!(parse_account_from_arn("arn:aws:sqs"), None);
2819    }
2820
2821    #[test]
2822    fn extract_region_from_user_agent_finds_region_segment() {
2823        let mut headers = http::HeaderMap::new();
2824        headers.insert(
2825            "user-agent",
2826            "aws-sdk-rust/1.0 os/linux region/eu-central-1"
2827                .parse()
2828                .unwrap(),
2829        );
2830        assert_eq!(
2831            extract_region_from_user_agent(&headers),
2832            Some("eu-central-1".to_string())
2833        );
2834    }
2835
2836    #[test]
2837    fn extract_region_from_user_agent_none_without_header() {
2838        let headers = http::HeaderMap::new();
2839        assert_eq!(extract_region_from_user_agent(&headers), None);
2840    }
2841
2842    #[test]
2843    fn extract_region_from_user_agent_ignores_empty_region() {
2844        let mut headers = http::HeaderMap::new();
2845        headers.insert("user-agent", "aws-sdk-java region/".parse().unwrap());
2846        assert_eq!(extract_region_from_user_agent(&headers), None);
2847    }
2848
2849    #[test]
2850    fn extract_region_from_user_agent_none_when_no_region_marker() {
2851        let mut headers = http::HeaderMap::new();
2852        headers.insert("user-agent", "curl/7.79.1".parse().unwrap());
2853        assert_eq!(extract_region_from_user_agent(&headers), None);
2854    }
2855
2856    #[test]
2857    fn aws_username_none_for_root() {
2858        let p = Principal {
2859            arn: "arn:aws:iam::123456789012:root".into(),
2860            user_id: "123456789012".into(),
2861            account_id: "123456789012".into(),
2862            principal_type: PrincipalType::Root,
2863            source_identity: None,
2864            tags: None,
2865        };
2866        assert_eq!(aws_username_from_principal(&p), None);
2867    }
2868
2869    #[test]
2870    fn aws_username_bare_no_path() {
2871        let p = Principal {
2872            arn: "arn:aws:iam::123456789012:user/bob".into(),
2873            user_id: "AIDABOB".into(),
2874            account_id: "123456789012".into(),
2875            principal_type: PrincipalType::User,
2876            source_identity: None,
2877            tags: None,
2878        };
2879        assert_eq!(aws_username_from_principal(&p), Some("bob".into()));
2880    }
2881
2882    #[test]
2883    fn principal_type_label_covers_federated_and_unknown() {
2884        assert_eq!(
2885            principal_type_label(PrincipalType::FederatedUser),
2886            "FederatedUser"
2887        );
2888        assert_eq!(principal_type_label(PrincipalType::Unknown), "Unknown");
2889    }
2890
2891    #[test]
2892    fn build_condition_context_marks_secure_when_flag_set() {
2893        let p = Principal {
2894            arn: "arn:aws:iam::123456789012:user/alice".into(),
2895            user_id: "AIDAALICE".into(),
2896            account_id: "123456789012".into(),
2897            principal_type: PrincipalType::User,
2898            source_identity: None,
2899            tags: None,
2900        };
2901        let ctx = build_condition_context(&p, None, "us-west-2", true);
2902        assert_eq!(ctx.aws_secure_transport, Some(true));
2903        assert!(ctx.aws_source_ip.is_none());
2904        assert_eq!(ctx.aws_requested_region.as_deref(), Some("us-west-2"));
2905    }
2906
2907    #[test]
2908    fn is_secure_transport_case_insensitive() {
2909        let mut headers = http::HeaderMap::new();
2910        headers.insert("x-forwarded-proto", "HTTPS".parse().unwrap());
2911        assert!(is_secure_transport(&headers));
2912    }
2913
2914    #[test]
2915    fn is_secure_transport_non_ascii_bytes_false() {
2916        let mut headers = http::HeaderMap::new();
2917        headers.insert(
2918            "x-forwarded-proto",
2919            http::HeaderValue::from_bytes(&[0xFF, 0xFE]).unwrap(),
2920        );
2921        assert!(!is_secure_transport(&headers));
2922    }
2923
2924    #[test]
2925    fn protocol_ext_error_status_is_bad_request() {
2926        assert_eq!(AwsProtocol::Query.error_status(), StatusCode::BAD_REQUEST);
2927        assert_eq!(AwsProtocol::Json.error_status(), StatusCode::BAD_REQUEST);
2928        assert_eq!(AwsProtocol::Rest.error_status(), StatusCode::BAD_REQUEST);
2929        assert_eq!(
2930            AwsProtocol::RestJson.error_status(),
2931            StatusCode::BAD_REQUEST
2932        );
2933    }
2934
2935    #[test]
2936    fn build_error_response_json_has_json_content_type() {
2937        let resp = build_error_response(
2938            StatusCode::BAD_REQUEST,
2939            "TestCode",
2940            "test msg",
2941            "req-1",
2942            AwsProtocol::Json,
2943        );
2944        assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
2945        let ct = resp
2946            .headers()
2947            .get("content-type")
2948            .unwrap()
2949            .to_str()
2950            .unwrap();
2951        assert!(ct.contains("json"));
2952        let rid = resp
2953            .headers()
2954            .get("x-amzn-requestid")
2955            .unwrap()
2956            .to_str()
2957            .unwrap();
2958        assert_eq!(rid, "req-1");
2959    }
2960
2961    #[test]
2962    fn build_error_response_rest_returns_xml_content_type() {
2963        let resp = build_error_response(
2964            StatusCode::NOT_FOUND,
2965            "NoSuchBucket",
2966            "bucket missing",
2967            "req-2",
2968            AwsProtocol::Rest,
2969        );
2970        assert_eq!(resp.status(), StatusCode::NOT_FOUND);
2971        let ct = resp
2972            .headers()
2973            .get("content-type")
2974            .unwrap()
2975            .to_str()
2976            .unwrap();
2977        assert!(ct.contains("xml"));
2978    }
2979
2980    fn rest_detected(service: &str) -> protocol::DetectedRequest {
2981        protocol::DetectedRequest {
2982            service: service.to_string(),
2983            action: String::new(),
2984            protocol: AwsProtocol::Rest,
2985        }
2986    }
2987
2988    async fn body_string(resp: Response<Body>) -> String {
2989        let bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2990            .await
2991            .unwrap();
2992        String::from_utf8(bytes.to_vec()).unwrap()
2993    }
2994
2995    #[tokio::test]
2996    async fn cloudfront_and_route53_errors_use_error_response_wrapper() {
2997        for (service, ns) in [
2998            (
2999                "cloudfront",
3000                "http://cloudfront.amazonaws.com/doc/2020-05-31/",
3001            ),
3002            ("route53", "https://route53.amazonaws.com/doc/2013-04-01/"),
3003        ] {
3004            let resp = build_error_response(
3005                StatusCode::NOT_FOUND,
3006                "NoSuchThing",
3007                "missing",
3008                "req-w",
3009                ErrorEnvelope::for_request(&rest_detected(service), &http::HeaderMap::new()),
3010            );
3011            assert_eq!(resp.status(), StatusCode::NOT_FOUND);
3012            assert_eq!(
3013                resp.headers().get("x-amz-error-code").unwrap(),
3014                "NoSuchThing"
3015            );
3016            let body = body_string(resp).await;
3017            assert!(
3018                body.contains(&format!(
3019                    "<ErrorResponse xmlns=\"{ns}\"><Error><Type>Sender</Type>\
3020                     <Code>NoSuchThing</Code><Message>missing</Message></Error>\
3021                     <RequestId>req-w</RequestId></ErrorResponse>"
3022                )),
3023                "{service}: {body}"
3024            );
3025        }
3026    }
3027
3028    #[tokio::test]
3029    async fn s3_errors_keep_bare_error_document() {
3030        let resp = build_error_response(
3031            StatusCode::NOT_FOUND,
3032            "NoSuchBucket",
3033            "missing",
3034            "req-s3",
3035            ErrorEnvelope::for_request(&rest_detected("s3"), &http::HeaderMap::new()),
3036        );
3037        let body = body_string(resp).await;
3038        assert!(!body.contains("<ErrorResponse"), "{body}");
3039        assert!(body.contains("<Error>"), "{body}");
3040        assert!(body.contains("<Code>NoSuchBucket</Code>"), "{body}");
3041    }
3042
3043    #[tokio::test]
3044    async fn s3_control_errors_use_error_response_wrapper() {
3045        let mut headers = http::HeaderMap::new();
3046        headers.insert(
3047            "host",
3048            "000000000000.s3-control.us-east-1.amazonaws.com"
3049                .parse()
3050                .unwrap(),
3051        );
3052        let resp = build_error_response(
3053            StatusCode::NOT_FOUND,
3054            "NoSuchAccessPoint",
3055            "missing",
3056            "req-ctl",
3057            ErrorEnvelope::for_request(&rest_detected("s3"), &headers),
3058        );
3059        assert_eq!(
3060            resp.headers().get("x-amz-error-code").unwrap(),
3061            "NoSuchAccessPoint"
3062        );
3063        let body = body_string(resp).await;
3064        assert!(
3065            body.contains(
3066                "<ErrorResponse><Error><Code>NoSuchAccessPoint</Code>\
3067                 <Message>missing</Message></Error>\
3068                 <RequestId>req-ctl</RequestId></ErrorResponse>"
3069            ),
3070            "{body}"
3071        );
3072    }
3073
3074    #[test]
3075    fn build_error_response_query_returns_xml() {
3076        let resp = build_error_response(
3077            StatusCode::BAD_REQUEST,
3078            "InvalidParameter",
3079            "bad param",
3080            "req-3",
3081            AwsProtocol::Query,
3082        );
3083        let ct = resp
3084            .headers()
3085            .get("content-type")
3086            .unwrap()
3087            .to_str()
3088            .unwrap();
3089        assert!(ct.contains("xml"));
3090    }
3091
3092    /// Regression for issue #1539: multi-line backend errors (e.g. podman
3093    /// stderr) used to panic the dispatcher when stuffed into the
3094    /// `x-amz-error-message` HTTP header. The response must build cleanly
3095    /// and the header value must not contain control characters.
3096    #[test]
3097    fn build_error_response_with_multiline_message_does_not_panic() {
3098        let resp = build_error_response(
3099            StatusCode::INTERNAL_SERVER_ERROR,
3100            "ServiceException",
3101            "Lambda execution failed: container failed to start: docker start failed: \
3102             Error: unable to start container \"abc\": \
3103             failed to create new hosts file:\nhost-gateway is empty\n",
3104            "req-multi",
3105            AwsProtocol::Json,
3106        );
3107        assert_eq!(resp.status(), StatusCode::INTERNAL_SERVER_ERROR);
3108        let msg = resp
3109            .headers()
3110            .get("x-amz-error-message")
3111            .expect("x-amz-error-message must be set even when input contains newlines")
3112            .to_str()
3113            .unwrap();
3114        assert!(!msg.contains('\n'));
3115        assert!(!msg.contains('\r'));
3116        assert!(msg.contains("Lambda execution failed"));
3117        assert!(msg.contains("host-gateway is empty"));
3118    }
3119
3120    #[test]
3121    fn build_error_response_with_control_chars_strips_them() {
3122        let resp = build_error_response(
3123            StatusCode::BAD_REQUEST,
3124            "Code\twith\ttabs",
3125            "msg\x00with\x01nulls",
3126            "req-ctrl",
3127            AwsProtocol::Json,
3128        );
3129        let code = resp
3130            .headers()
3131            .get("x-amz-error-code")
3132            .unwrap()
3133            .to_str()
3134            .unwrap();
3135        let msg = resp
3136            .headers()
3137            .get("x-amz-error-message")
3138            .unwrap()
3139            .to_str()
3140            .unwrap();
3141        assert!(!code.contains('\t'));
3142        assert!(!msg.contains('\x00'));
3143        assert!(!msg.contains('\x01'));
3144    }
3145
3146    #[test]
3147    fn sanitize_header_value_truncates_long_input() {
3148        let huge = "x".repeat(5_000);
3149        let out = sanitize_header_value(&huge);
3150        assert!(out.len() <= 1024);
3151    }
3152
3153    #[test]
3154    fn sanitize_header_value_collapses_consecutive_control_runs() {
3155        let out = sanitize_header_value("a\n\n\n\rb");
3156        assert_eq!(out, "a b");
3157    }
3158
3159    #[test]
3160    fn anonymous_s3_probe_finds_a_bucket_on_a_china_server() {
3161        // Answers only for ARNs the S3 policy provider's bucket parser reads
3162        // (`arn:aws:s3:::<bucket>`), so a probe the real provider would not
3163        // understand misses here too.
3164        struct RecordingProvider(parking_lot::Mutex<Vec<String>>);
3165        impl crate::auth::ResourcePolicyProvider for RecordingProvider {
3166            fn resource_policy(&self, _service: &str, _resource_arn: &str) -> Option<String> {
3167                None
3168            }
3169            fn resource_owner_account(&self, _service: &str, resource_arn: &str) -> Option<String> {
3170                self.0.lock().push(resource_arn.to_string());
3171                resource_arn
3172                    .strip_prefix("arn:aws:s3:::")
3173                    .filter(|bucket| *bucket == "my-bucket")
3174                    .map(|_| "000000000000".to_string())
3175            }
3176        }
3177        let provider = Arc::new(RecordingProvider(parking_lot::Mutex::new(Vec::new())));
3178        let mut cfg = DispatchConfig::new("cn-north-1", "000000000000");
3179        cfg.resource_policy_provider = Some(provider.clone());
3180        let uri: http::Uri = "/my-bucket/key.txt".parse().unwrap();
3181        assert_eq!(
3182            anonymous_s3_bucket(&uri, &cfg),
3183            Some("my-bucket".to_string())
3184        );
3185        assert_eq!(
3186            *provider.0.lock(),
3187            vec!["arn:aws:s3:::my-bucket".to_string()]
3188        );
3189    }
3190
3191    #[test]
3192    fn dispatch_config_carries_opt_in_flags() {
3193        let cfg = DispatchConfig {
3194            region: "eu-west-1".to_string(),
3195            account_id: "000000000000".to_string(),
3196            verify_sigv4: true,
3197            iam_mode: IamMode::Strict,
3198            credential_resolver: None,
3199            policy_evaluator: None,
3200            resource_policy_provider: None,
3201            scp_resolver: None,
3202        };
3203        assert!(cfg.verify_sigv4);
3204        assert!(cfg.iam_mode.is_strict());
3205        assert!(cfg.resource_policy_provider.is_none());
3206        assert!(cfg.scp_resolver.is_none());
3207    }
3208
3209    fn s3_sigv4_headers() -> http::HeaderMap {
3210        let mut headers = http::HeaderMap::new();
3211        headers.insert(
3212            "authorization",
3213            "AWS4-HMAC-SHA256 Credential=test/20240101/us-east-1/s3/aws4_request, \
3214             SignedHeaders=host, Signature=fake"
3215                .parse()
3216                .unwrap(),
3217        );
3218        headers
3219    }
3220
3221    #[test]
3222    fn streaming_route_path_style_s3_put_object() {
3223        let headers = s3_sigv4_headers();
3224        assert_eq!(
3225            streaming_route(
3226                &http::Method::PUT,
3227                "/my-bucket/key.txt",
3228                &headers,
3229                &HashMap::new(),
3230            ),
3231            Some(("s3", "")),
3232        );
3233    }
3234
3235    #[test]
3236    fn streaming_route_path_style_create_bucket_skipped() {
3237        // `PUT /bucket` (no trailing key) is CreateBucket — must NOT
3238        // hit the streaming path.
3239        let headers = s3_sigv4_headers();
3240        assert_eq!(
3241            streaming_route(&http::Method::PUT, "/my-bucket", &headers, &HashMap::new(),),
3242            None,
3243        );
3244    }
3245
3246    #[test]
3247    fn s3_routing_path_prefixes_the_host_bucket() {
3248        assert_eq!(s3_routing_path("/", Some("b")), "/b");
3249        assert_eq!(s3_routing_path("", Some("b")), "/b");
3250        assert_eq!(s3_routing_path("/k.txt", Some("b")), "/b/k.txt");
3251        assert_eq!(s3_routing_path("/dir/k", Some("a.b")), "/a.b/dir/k");
3252        assert_eq!(s3_routing_path("/b/k", None), "/b/k");
3253    }
3254
3255    #[test]
3256    fn s3_routing_path_keeps_a_key_that_starts_with_the_bucket_name() {
3257        // The wire path is the whole key on a virtual-hosted request, so a key
3258        // whose first segment equals the bucket name must keep it.
3259        assert_eq!(
3260            s3_routing_path("/docs/intro.html", Some("docs")),
3261            "/docs/docs/intro.html"
3262        );
3263        assert_eq!(s3_routing_path("/docs", Some("docs")), "/docs/docs");
3264    }
3265
3266    #[test]
3267    fn streaming_route_path_style_create_bucket_with_trailing_slash_skipped() {
3268        // The AWS SDKs send CreateBucket as `PUT /<bucket>/`. The trailing
3269        // slash must not make it look like an object upload: the body carries
3270        // `CreateBucketConfiguration` (location constraint, tag set) and has
3271        // to be buffered before IAM enforcement reads it.
3272        let headers = s3_sigv4_headers();
3273        assert_eq!(
3274            streaming_route(&http::Method::PUT, "/my-bucket/", &headers, &HashMap::new(),),
3275            None,
3276        );
3277    }
3278
3279    #[test]
3280    fn streaming_route_path_style_doubled_slash_skipped() {
3281        // Routing filters empty path segments, so `PUT /<bucket>//` is a
3282        // bucket-level operation, not an object upload with a "/" key. The
3283        // streaming gate has to agree, or the body goes unbuffered for a
3284        // request the handler reads as CreateBucket / PutBucketTagging.
3285        let headers = s3_sigv4_headers();
3286        assert_eq!(
3287            streaming_route(
3288                &http::Method::PUT,
3289                "/my-bucket//",
3290                &headers,
3291                &HashMap::new()
3292            ),
3293            None,
3294        );
3295    }
3296
3297    #[test]
3298    fn streaming_route_path_style_key_with_trailing_slash_streams() {
3299        // A real object key that ends in '/' (a directory marker) still
3300        // streams -- the gate rejects an absent key, not a trailing slash.
3301        let headers = s3_sigv4_headers();
3302        assert_eq!(
3303            streaming_route(
3304                &http::Method::PUT,
3305                "/my-bucket/folder/",
3306                &headers,
3307                &HashMap::new(),
3308            ),
3309            Some(("s3", "")),
3310        );
3311    }
3312
3313    #[test]
3314    fn streaming_route_virtual_hosted_s3_put_object() {
3315        let mut headers = s3_sigv4_headers();
3316        headers.insert(
3317            "host",
3318            "vhost-bucket.s3.us-east-1.localhost.localstack.cloud:4566"
3319                .parse()
3320                .unwrap(),
3321        );
3322        // Virtual-hosted PUT has no bucket in the URL path (`/<key>`),
3323        // so the slash check used for path-style would reject it. The
3324        // Host parser confirms this is virtual-hosted S3 and the key
3325        // flows through the streaming dispatch.
3326        assert_eq!(
3327            streaming_route(&http::Method::PUT, "/hello.txt", &headers, &HashMap::new(),),
3328            Some(("s3", "")),
3329        );
3330    }
3331
3332    #[test]
3333    fn streaming_route_virtual_hosted_path_naming_the_bucket_streams() {
3334        // Under virtual-hosted addressing the whole path is the key, as on
3335        // real S3: `PUT /my-bucket` on `Host: my-bucket.s3...` uploads the
3336        // object `my-bucket`, not a CreateBucket, so every one of these
3337        // streams.
3338        let mut headers = s3_sigv4_headers();
3339        headers.insert(
3340            "host",
3341            "my-bucket.s3.us-east-1.amazonaws.com".parse().unwrap(),
3342        );
3343        for path in ["/my-bucket", "/my-bucket/", "/my-bucket/key.txt"] {
3344            assert_eq!(
3345                streaming_route(&http::Method::PUT, path, &headers, &HashMap::new()),
3346                Some(("s3", "")),
3347                "{path}",
3348            );
3349        }
3350    }
3351
3352    #[test]
3353    fn streaming_route_virtual_hosted_s3_root_skipped() {
3354        // `PUT /` against a virtual-hosted Host = CreateBucket, which
3355        // is handled buffered. Empty path-after-slash must short-circuit.
3356        let mut headers = s3_sigv4_headers();
3357        headers.insert(
3358            "host",
3359            "vhost-bucket.s3.us-east-1.localhost.localstack.cloud:4566"
3360                .parse()
3361                .unwrap(),
3362        );
3363        assert_eq!(
3364            streaming_route(&http::Method::PUT, "/", &headers, &HashMap::new()),
3365            None,
3366        );
3367    }
3368
3369    #[test]
3370    fn streaming_route_ecr_blob_upload() {
3371        let headers = http::HeaderMap::new();
3372        assert_eq!(
3373            streaming_route(
3374                &http::Method::PATCH,
3375                "/v2/my-repo/blobs/uploads/abcd1234",
3376                &headers,
3377                &HashMap::new(),
3378            ),
3379            Some(("ecr", "")),
3380        );
3381        assert_eq!(
3382            streaming_route(
3383                &http::Method::PUT,
3384                "/v2/my-repo/blobs/uploads/abcd1234",
3385                &headers,
3386                &HashMap::new(),
3387            ),
3388            Some(("ecr", "")),
3389        );
3390    }
3391
3392    #[test]
3393    fn hoist_presigned_query_headers_skips_auth_params_and_keeps_direct_headers() {
3394        let mut headers = http::HeaderMap::new();
3395        headers.insert("x-amz-meta-color", "red".parse().unwrap());
3396        let query: HashMap<String, String> = [
3397            (
3398                "X-Amz-Credential",
3399                "AKID/20260101/us-east-1/s3/aws4_request",
3400            ),
3401            ("X-Amz-Signature", "00"),
3402            ("X-Amz-Security-Token", "tok"),
3403            ("x-amz-meta-color", "blue"),
3404            ("X-Amz-Meta-Shape", "round"),
3405            ("x-amz-meta-name", "café"),
3406            ("x-amz-tagging", "env=test"),
3407            ("response-content-type", "text/plain"),
3408            ("partNumber", "1"),
3409        ]
3410        .into_iter()
3411        .map(|(k, v)| (k.to_string(), v.to_string()))
3412        .collect();
3413        hoist_presigned_query_headers(&mut headers, &query);
3414
3415        assert_eq!(headers["x-amz-meta-color"], "red");
3416        assert_eq!(headers["x-amz-meta-shape"], "round");
3417        assert_eq!(headers["x-amz-meta-name"], "=?UTF-8?B?Y2Fmw6k=?=");
3418        assert_eq!(headers["x-amz-tagging"], "env=test");
3419        for absent in [
3420            "x-amz-credential",
3421            "x-amz-signature",
3422            "x-amz-security-token",
3423            "response-content-type",
3424            "partnumber",
3425        ] {
3426            assert!(headers.get(absent).is_none(), "{absent}");
3427        }
3428    }
3429
3430    #[test]
3431    fn streaming_route_presigned_v4_s3_put() {
3432        let headers = http::HeaderMap::new();
3433        let mut query_params = HashMap::new();
3434        query_params.insert(
3435            "X-Amz-Credential".to_string(),
3436            "test/20240101/us-east-1/s3/aws4_request".to_string(),
3437        );
3438        assert_eq!(
3439            streaming_route(
3440                &http::Method::PUT,
3441                "/my-bucket/key.txt",
3442                &headers,
3443                &query_params,
3444            ),
3445            Some(("s3", "")),
3446        );
3447    }
3448
3449    #[test]
3450    fn streaming_route_non_s3_auth_header_skipped() {
3451        // Same path shape but the SigV4 service is lambda — must not
3452        // wire the streaming dispatch (Lambda has its own buffered path).
3453        let mut headers = http::HeaderMap::new();
3454        headers.insert(
3455            "authorization",
3456            "AWS4-HMAC-SHA256 Credential=test/20240101/us-east-1/lambda/aws4_request, \
3457             SignedHeaders=host, Signature=fake"
3458                .parse()
3459                .unwrap(),
3460        );
3461        assert_eq!(
3462            streaming_route(
3463                &http::Method::PUT,
3464                "/my-bucket/key.txt",
3465                &headers,
3466                &HashMap::new(),
3467            ),
3468            None,
3469        );
3470    }
3471
3472    #[test]
3473    fn streaming_route_get_skipped() {
3474        let headers = s3_sigv4_headers();
3475        assert_eq!(
3476            streaming_route(
3477                &http::Method::GET,
3478                "/my-bucket/key.txt",
3479                &headers,
3480                &HashMap::new(),
3481            ),
3482            None,
3483        );
3484    }
3485
3486    /// Root of account B reaching a bucket account A owns: allowed only when
3487    /// A's resource policy grants it, as for any cross-account principal.
3488    /// Root in its own account, and resources no provider claims, keep the
3489    /// root exemption.
3490    #[test]
3491    fn root_cross_account_needs_the_resource_policy() {
3492        use crate::auth::IamAction;
3493        use crate::service::{AwsResponse, AwsServiceError};
3494        struct OwnerProvider;
3495        impl crate::auth::ResourcePolicyProvider for OwnerProvider {
3496            fn resource_policy(&self, _service: &str, resource_arn: &str) -> Option<String> {
3497                resource_arn
3498                    .starts_with("arn:aws:s3:::granted")
3499                    .then(|| "granting-policy".to_string())
3500            }
3501            fn resource_owner_account(&self, _service: &str, resource_arn: &str) -> Option<String> {
3502                let bucket = resource_arn.strip_prefix("arn:aws:s3:::")?;
3503                let bucket = bucket.split('/').next()?;
3504                match bucket {
3505                    "own" => Some("222222222222".to_string()),
3506                    "granted" | "denied" => Some("111111111111".to_string()),
3507                    _ => None,
3508                }
3509            }
3510        }
3511        struct PolicyEvaluator(parking_lot::Mutex<Vec<ConditionContext>>);
3512        impl IamPolicyEvaluator for PolicyEvaluator {
3513            fn evaluate(
3514                &self,
3515                _: &Principal,
3516                _: &IamAction,
3517                _: &ConditionContext,
3518                _: &[String],
3519                _: Option<&[String]>,
3520            ) -> crate::auth::IamDecision {
3521                crate::auth::IamDecision::ImplicitDeny
3522            }
3523            fn evaluate_with_resource_policy(
3524                &self,
3525                _: &Principal,
3526                _: &IamAction,
3527                _: &ConditionContext,
3528                _: Option<&str>,
3529                _: &str,
3530                _: &[String],
3531                _: Option<&[String]>,
3532            ) -> crate::auth::IamDecision {
3533                crate::auth::IamDecision::ImplicitDeny
3534            }
3535            fn evaluate_resource_policy_only(
3536                &self,
3537                _: &Principal,
3538                _: &IamAction,
3539                context: &ConditionContext,
3540                policy: Option<&str>,
3541            ) -> crate::auth::IamDecision {
3542                self.0.lock().push(context.clone());
3543                if policy == Some("granting-policy") {
3544                    crate::auth::IamDecision::Allow
3545                } else {
3546                    crate::auth::IamDecision::ImplicitDeny
3547                }
3548            }
3549        }
3550        struct BucketService;
3551        #[async_trait::async_trait]
3552        impl crate::service::AwsService for BucketService {
3553            fn service_name(&self) -> &str {
3554                "s3"
3555            }
3556            async fn handle(&self, _: AwsRequest) -> Result<AwsResponse, AwsServiceError> {
3557                unreachable!()
3558            }
3559            fn supported_actions(&self) -> &[&str] {
3560                &[]
3561            }
3562            fn iam_action_for(&self, request: &AwsRequest) -> Option<IamAction> {
3563                Some(IamAction {
3564                    service: "s3",
3565                    action: "GetObject",
3566                    resource: format!("arn:aws:s3:::{}", request.path_segments.join("/")),
3567                })
3568            }
3569        }
3570        let root = Principal {
3571            arn: "arn:aws:iam::222222222222:root".to_string(),
3572            user_id: "222222222222".to_string(),
3573            account_id: "222222222222".to_string(),
3574            principal_type: PrincipalType::Root,
3575            source_identity: None,
3576            tags: None,
3577        };
3578        let detected = protocol::DetectedRequest {
3579            service: "s3".to_string(),
3580            action: String::new(),
3581            protocol: AwsProtocol::Rest,
3582        };
3583        let request = |bucket: &str| AwsRequest {
3584            service: "s3".to_string(),
3585            action: String::new(),
3586            region: "us-east-1".to_string(),
3587            account_id: "222222222222".to_string(),
3588            request_id: "req".to_string(),
3589            headers: http::HeaderMap::new(),
3590            query_params: HashMap::new(),
3591            body: Bytes::new(),
3592            body_stream: parking_lot::Mutex::new(None),
3593            path_segments: vec![bucket.to_string(), "key".to_string()],
3594            raw_path: format!("/{bucket}/key"),
3595            raw_query: String::new(),
3596            method: http::Method::GET,
3597            is_query_protocol: false,
3598            access_key_id: Some("AKIAROOTB".to_string()),
3599            principal: Some(root.clone()),
3600        };
3601        struct Org;
3602        impl crate::auth::ScpResolver for Org {
3603            fn scps_for(&self, _: &Principal) -> Option<Vec<String>> {
3604                None
3605            }
3606            fn principal_org(&self, account: &str) -> Option<(String, String)> {
3607                (account == "222222222222").then(|| ("o-abc".to_string(), "o-abc/r-1/".to_string()))
3608            }
3609        }
3610        let issued = chrono::Utc::now();
3611        let credential = crate::auth::ResolvedCredential {
3612            secret_access_key: "secret".to_string(),
3613            session_token: Some("token".to_string()),
3614            principal: root.clone(),
3615            session_policies: Vec::new(),
3616            mfa_present: true,
3617            token_issued_at: Some(issued),
3618            federated_provider: None,
3619        };
3620        let evaluator = PolicyEvaluator(parking_lot::Mutex::new(Vec::new()));
3621        let run = |bucket: &str, mode: IamMode| {
3622            let mut cfg = DispatchConfig::new("us-east-1", "111111111111");
3623            cfg.iam_mode = mode;
3624            cfg.resource_policy_provider = Some(Arc::new(OwnerProvider));
3625            cfg.scp_resolver = Some(Arc::new(Org));
3626            authorize_root_cross_account(
3627                &root,
3628                &BucketService,
3629                &request(bucket),
3630                &evaluator,
3631                &cfg,
3632                &detected,
3633                "req",
3634                None,
3635                Some(&credential),
3636            )
3637        };
3638        let denied = run("denied", IamMode::Strict).expect("no grant must deny under strict");
3639        assert_eq!(denied.status(), StatusCode::FORBIDDEN);
3640        assert!(run("granted", IamMode::Strict).is_none());
3641        assert!(run("own", IamMode::Strict).is_none());
3642        assert!(run("unclaimed", IamMode::Strict).is_none());
3643        assert!(run("denied", IamMode::Soft).is_none());
3644
3645        // The root path evaluates with the same condition context as any
3646        // other principal: the owner's account, the principal's organization
3647        // and the session keys its credential carries.
3648        let contexts = evaluator.0.lock();
3649        assert!(!contexts.is_empty());
3650        for ctx in contexts.iter() {
3651            assert_eq!(
3652                ctx.lookup("aws:ResourceAccount"),
3653                Some(vec!["111111111111".into()])
3654            );
3655            assert_eq!(ctx.lookup("aws:PrincipalOrgID"), Some(vec!["o-abc".into()]));
3656            assert_eq!(
3657                ctx.lookup("aws:PrincipalOrgPaths"),
3658                Some(vec!["o-abc/r-1/".into()])
3659            );
3660            assert_eq!(ctx.aws_mfa_present, Some(true));
3661            assert_eq!(ctx.aws_token_issue_time, Some(issued));
3662            assert_eq!(ctx.aws_principal_type.as_deref(), Some("Account"));
3663        }
3664    }
3665}
3666
3667/// Whether the request is addressed to an API Gateway execute-api host
3668/// (`{api-id}.execute-api.{region}.amazonaws.com`, or a local alias of it).
3669fn is_execute_api_host(headers: &http::HeaderMap) -> bool {
3670    headers
3671        .get(http::header::HOST)
3672        .and_then(|v| v.to_str().ok())
3673        .is_some_and(|host| host.contains(".execute-api."))
3674}
3675
3676#[cfg(test)]
3677mod execute_api_host_tests {
3678    use super::is_execute_api_host;
3679
3680    #[test]
3681    fn recognizes_execute_api_hosts_only() {
3682        let mut h = http::HeaderMap::new();
3683        h.insert(
3684            http::header::HOST,
3685            "abc123.execute-api.us-east-1.amazonaws.com"
3686                .parse()
3687                .unwrap(),
3688        );
3689        assert!(is_execute_api_host(&h));
3690        h.insert(http::header::HOST, "localhost:4566".parse().unwrap());
3691        assert!(!is_execute_api_host(&h));
3692        assert!(!is_execute_api_host(&http::HeaderMap::new()));
3693    }
3694}