fraiseql-server 2.16.0

HTTP server for FraiseQL v2 GraphQL engine
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
//! GraphQL HTTP handlers and execution logic.

use std::{sync::atomic::Ordering, time::Instant};

use axum::{
    Json,
    extract::{Query, State},
    http::HeaderMap,
};
use fraiseql_core::{
    apq::{ApqMetrics, ApqStorage},
    security::SecurityContext,
};
use fraiseql_error::FraiseQLError;
use tracing::{debug, error, warn};

use super::{
    app_state::AppState,
    request::{GraphQLGetParams, GraphQLRequest, GraphQLResponse},
};
use crate::{
    error::{ErrorResponse, GraphQLError},
    extractors::{OptionalSecurityContext, PeerIp},
    tracing_utils,
};

/// GraphQL HTTP handler for POST requests.
///
/// Handles POST requests to the GraphQL endpoint:
/// 1. Extract W3C trace context from traceparent header (if present)
/// 2. Validate GraphQL request (depth, complexity)
/// 3. Parse GraphQL request body
/// 4. Execute query via Executor with optional `SecurityContext`
/// 5. Return GraphQL response with proper error formatting
///
/// Tracks execution timing and operation name for monitoring.
/// Provides GraphQL spec-compliant error responses.
/// Supports W3C Trace Context for distributed tracing.
/// Supports OIDC authentication for RLS policy evaluation.
///
/// # Errors
///
/// Returns appropriate HTTP status codes based on error type.
#[tracing::instrument(skip_all, fields(operation_name))]
#[doc(hidden)] // Internal-pub: axum route handler wired via Server::route; downstream uses Server::serve(), not this fn directly.
pub async fn graphql_handler(
    State(state): State<AppState>,
    headers: HeaderMap,
    PeerIp(peer_ip): PeerIp,
    OptionalSecurityContext(security_context): OptionalSecurityContext,
    // #958: an SSE `@stream` delivery outlives its request, so it re-checks the
    // principal per batch — which needs the decoded `jti`/`iat` the auth middleware
    // put in the extensions. Absent for anonymous and service-account callers, which
    // is exactly the signal that no revocation re-check applies.
    token_claims: Option<axum::Extension<crate::middleware::oidc_auth::SessionTokenClaims>>,
    Json(request): Json<GraphQLRequest>,
) -> Result<axum::response::Response, ErrorResponse> {
    // Extract trace context from W3C headers
    let trace_context = tracing_utils::extract_trace_context(&headers);
    if trace_context.is_some() {
        debug!("Extracted W3C trace context from incoming request");
    }

    if security_context.is_some() {
        debug!("Authenticated request with security context");
    }

    // GraphQL-over-SSE (#387): opt-in, negotiated by Accept. This branch runs
    // INSIDE the authenticated route, so it inherits the full middleware stack.
    if let Some(wire) = state
        .graphql_incremental_enabled
        .then(|| incremental::negotiate(&headers))
        .flatten()
    {
        return Box::pin(sse::handle_sse(
            state,
            wire,
            headers,
            peer_ip,
            security_context,
            token_claims.map(|axum::Extension(claims)| claims),
            request,
        ))
        .await;
    }

    execute_graphql_request(state, request, trace_context, security_context, &headers, &peer_ip)
        .await
        .map(axum::response::IntoResponse::into_response)
}

/// GraphQL HTTP handler for GET requests.
///
/// Handles GET requests to the GraphQL endpoint per the GraphQL over HTTP spec.
/// Query parameters:
/// - `query`: Required, the GraphQL query string (URL-encoded)
/// - `variables`: Optional, JSON-encoded variables object (URL-encoded)
/// - `operationName`: Optional, name of the operation to execute
///
/// Supports W3C Trace Context via traceparent header for distributed tracing.
///
/// Example:
/// ```text
/// GET /graphql?query={users{id,name}}&variables={"limit":10}
/// ```
///
/// # Errors
///
/// Returns `413 Payload Too Large` (via `ErrorResponse`) when the query string
/// exceeds `AppState::max_get_query_bytes` (default 100 `KiB`, configurable via
/// `ServerConfig::max_get_query_bytes`). Returns other HTTP status codes for
/// additional error conditions.
///
/// # Note
///
/// Per GraphQL over HTTP spec, GET requests should only be used for queries,
/// not mutations (which should use POST). This handler does not enforce that
/// restriction but logs a warning for mutation-like queries.
#[tracing::instrument(skip_all, fields(operation_name))]
#[doc(hidden)] // Internal-pub: axum route handler wired via Server::route; downstream uses Server::serve(), not this fn directly.
pub async fn graphql_get_handler(
    State(state): State<AppState>,
    headers: HeaderMap,
    PeerIp(peer_ip): PeerIp,
    OptionalSecurityContext(security_context): OptionalSecurityContext,
    // #958: see `graphql_handler` — the SSE branch needs the token's revocation claims.
    token_claims: Option<axum::Extension<crate::middleware::oidc_auth::SessionTokenClaims>>,
    Query(params): Query<GraphQLGetParams>,
) -> Result<axum::response::Response, ErrorResponse> {
    // Reject oversized GET queries early to prevent DoS via query parsing.
    let max_get_bytes = state.max_get_query_bytes;
    if params.query.len() > max_get_bytes {
        return Err(ErrorResponse::from_error(GraphQLError::payload_too_large(format!(
            "GET query string exceeds maximum allowed length ({max_get_bytes} bytes)"
        ))));
    }

    // Parse variables from JSON string.
    // Apply the same size cap as the query string — the URL-length limit imposed
    // by reverse proxies/OS is real but not enforced by axum itself, so we guard
    // explicitly to prevent parser DoS from a very large variables value.
    let variables = if let Some(vars_str) = params.variables {
        if vars_str.len() > max_get_bytes {
            return Err(ErrorResponse::from_error(GraphQLError::payload_too_large(format!(
                "GET variables string exceeds maximum allowed length ({max_get_bytes} bytes)"
            ))));
        }
        match serde_json::from_str::<serde_json::Value>(&vars_str) {
            Ok(v) => Some(v),
            Err(e) => {
                // Log the fault, never the payload. `variables` is client-supplied
                // and may hold PII or bearer tokens; emitting up to
                // `max_get_query_bytes` of it at `warn!` put that into every log
                // sink the deployment ships to (#730). The serde error already
                // carries the line/column, which is what a diagnosis needs.
                warn!(
                    error = %e,
                    variables_bytes = vars_str.len(),
                    "Failed to parse variables JSON in GET request"
                );
                return Err(ErrorResponse::from_error(GraphQLError::request(format!(
                    "Invalid variables JSON: {e}"
                ))));
            },
        }
    } else {
        None
    };

    // Reject mutations over GET with 405 per the GraphQL-over-HTTP spec: GET is for
    // queries only, and allowing mutations sidesteps the POST-only CSRF posture
    // (M-get-mutations). Detection parses the operation (reliable) rather than matching a
    // `mutation` string prefix (which a leading comment or named query defeats).
    if detect_mutation_name(&params.query).is_some() {
        warn!(
            operation_name = ?params.operation_name,
            "Mutation sent via GET request — rejected (use POST)"
        );
        return Err(ErrorResponse::from_error(GraphQLError::method_not_allowed(
            "Mutations must be sent over POST, not GET",
        )));
    }

    let trace_context = tracing_utils::extract_trace_context(&headers);
    if trace_context.is_some() {
        debug!("Extracted W3C trace context from incoming request");
    }

    let request = GraphQLRequest {
        query: Some(params.query),
        variables,
        operation_name: params.operation_name,
        extensions: None,
        document_id: None,
    };

    if security_context.is_some() {
        debug!("Authenticated GET request with security context");
    }

    // GraphQL-over-SSE (#387): the GET arm is what an `EventSource` client can
    // reach (it cannot POST). Mutations were already rejected above.
    if let Some(wire) = state
        .graphql_incremental_enabled
        .then(|| incremental::negotiate(&headers))
        .flatten()
    {
        return Box::pin(sse::handle_sse(
            state,
            wire,
            headers,
            peer_ip,
            security_context,
            token_claims.map(|axum::Extension(claims)| claims),
            request,
        ))
        .await;
    }

    execute_graphql_request(state, request, trace_context, security_context, &headers, &peer_ip)
        .await
        .map(axum::response::IntoResponse::into_response)
}

/// The HTTP `QUERY` method (RFC 10008), as a byte string.
///
/// **This is the single seam to swap for native support.** Neither `http` 1.4.x nor
/// axum 0.8.x exposes `Method::QUERY` / `MethodFilter::QUERY` yet
/// (hyperium/http#798, tokio-rs/axum#3799). When they ship, replace this constant
/// and the `MethodRouter::fallback` wiring in `server/routing/graphql.rs` with the
/// typed filter — those two places are the whole workaround.
pub(crate) const HTTP_QUERY_METHOD: &str = "QUERY";

/// `QUERY /graphql` — RFC 10008 (#508).
///
/// Mounted only when `enable_http_query = true`, as a `MethodRouter` fallback, so
/// `GET` and `POST` dispatch is untouched. The body is parsed exactly like a `POST`,
/// then the operation is **restricted to queries**.
///
/// # Why the restriction is mandatory
///
/// `QUERY` is defined as safe and idempotent, which entitles any intermediary —
/// proxy, CDN, retry layer — to replay it. A mutation carried over it could
/// therefore be silently executed more than once. Refusing non-query operations is
/// what makes accepting the method safe at all.
///
/// # Errors
///
/// - `405` when the request method is not `QUERY` (any other unmatched method).
/// - `405` when the document is a `mutation` or `subscription`.
/// - `400` when the body is not a valid `GraphQLRequest`.
pub async fn graphql_query_method_handler(
    State(state): State<AppState>,
    method: axum::http::Method,
    headers: HeaderMap,
    PeerIp(peer_ip): PeerIp,
    OptionalSecurityContext(security_context): OptionalSecurityContext,
    body: axum::body::Bytes,
) -> Result<GraphQLResponse, ErrorResponse> {
    // This handler is a method-router fallback, so it also catches PUT/DELETE/…
    // Anything that is not QUERY gets the same 405 the router would have produced.
    if method.as_str() != HTTP_QUERY_METHOD {
        return Err(ErrorResponse::from_error(GraphQLError::method_not_allowed(
            "Method not allowed on the GraphQL endpoint",
        )));
    }

    let request: GraphQLRequest = serde_json::from_slice(&body).map_err(|e| {
        ErrorResponse::from_error(GraphQLError::request(format!("Invalid QUERY body: {e}")))
    })?;

    // Queries-only. `detect_operation_type` uses the same parser the executor uses,
    // so the gate cannot disagree with what would actually execute. A document that
    // does not parse falls through (like the GET path) and surfaces a GraphQL parse
    // error from the executor — it executes nothing, so no state can change.
    if let Some(op) = non_query_operation(request.query.as_deref().unwrap_or_default()) {
        warn!(operation = op, "Non-query operation sent via QUERY — rejected (use POST)");
        return Err(ErrorResponse::from_error(GraphQLError::method_not_allowed(
            "Only query operations may be sent over the QUERY method; use POST",
        )));
    }

    let trace_context = tracing_utils::extract_trace_context(&headers);

    execute_graphql_request(state, request, trace_context, security_context, &headers, &peer_ip)
        .await
}

/// Returns the operation type when `query` is **not** a plain query — i.e.
/// `Some("mutation")` or `Some("subscription")`. `None` for queries and for
/// documents that fail to parse (which execute nothing).
pub(crate) fn non_query_operation(query: &str) -> Option<&'static str> {
    let parsed = fraiseql_core::graphql::parse_query(query).ok()?;
    match parsed.operation_type.as_str() {
        "mutation" => Some("mutation"),
        "subscription" => Some("subscription"),
        _ => None,
    }
}

/// Extract the mutation name from a GraphQL query string, if the operation is a mutation.
///
/// Returns `Some(root_field_name)` when the query parses successfully and the operation
/// type is `"mutation"`. Returns `None` for queries, subscriptions, or parse errors.
///
/// Used to look up before-mutation hooks: a single `HashMap::get` on the trigger
/// registry — O(1) and allocation-free when no hooks are registered.
pub(crate) fn detect_mutation_name(query: &str) -> Option<String> {
    let parsed = fraiseql_core::graphql::parse_query(query).ok()?;
    if parsed.operation_type == "mutation" {
        Some(parsed.root_field)
    } else {
        None
    }
}

/// Extract client IP address from headers.
///
/// # Security
///
/// Does NOT trust X-Forwarded-For or X-Real-IP headers, as these are trivially
/// spoofable by attackers to bypass rate limiting. Returns "unknown" as a safe
/// fallback — callers requiring real IPs should use `ConnectInfo<SocketAddr>`
/// or `ProxyConfig::extract_client_ip()` with validated proxy chains.
#[cfg(feature = "auth")]
#[allow(dead_code)] // Reason: used only in tests that verify spoofable headers are ignored
pub(crate) fn extract_ip_from_headers(_headers: &HeaderMap) -> String {
    // SECURITY: Spoofable headers removed. Use ConnectInfo<SocketAddr> or
    // ProxyConfig::extract_client_ip() for validated IP extraction.
    "unknown".to_string()
}

/// Extract the APQ SHA-256 hash from the `extensions.persistedQuery` field, if present.
pub(crate) fn extract_apq_hash(extensions: Option<&serde_json::Value>) -> Option<&str> {
    extensions?.get("persistedQuery")?.get("sha256Hash")?.as_str()
}

/// Extract a trusted document ID from the request.
///
/// Supports three formats:
/// 1. `documentId` (GraphQL over HTTP spec)
/// 2. `extensions.persistedQuery.sha256Hash` (Apollo APQ format)
/// 3. `extensions.doc_id` (Relay format)
fn extract_document_id(request: &GraphQLRequest) -> Option<String> {
    // 1. Top-level documentId field (GraphQL over HTTP spec)
    if let Some(ref doc_id) = request.document_id {
        return Some(doc_id.clone());
    }
    // 2. Extensions-based formats
    if let Some(ext) = request.extensions.as_ref() {
        // Relay format: extensions.doc_id
        if let Some(doc_id) = ext.get("doc_id").and_then(|v| v.as_str()) {
            return Some(doc_id.to_string());
        }
        // Apollo APQ format: extensions.persistedQuery.sha256Hash (also used for APQ)
        if let Some(hash) = ext
            .get("persistedQuery")
            .and_then(|pq| pq.get("sha256Hash"))
            .and_then(|h| h.as_str())
        {
            return Some(hash.to_string());
        }
    }
    None
}

/// Resolve an APQ request: look up or register a persisted query.
///
/// Returns the resolved query body, or an error if the query is not found and no body was
/// provided (the client should resend with the full body).
///
/// # Errors
///
/// Returns [`ErrorResponse`] if the hash doesn't match the body, or if the
/// hash is unknown and no query body was provided (client must retry with full body).
pub(crate) async fn resolve_apq(
    apq_store: &dyn ApqStorage,
    apq_metrics: &ApqMetrics,
    hash: &str,
    query_body: Option<&str>,
) -> Result<String, ErrorResponse> {
    if let Some(body) = query_body {
        // Hash + body present: verify and register.
        if !fraiseql_core::apq::verify_hash(body, hash) {
            apq_metrics.record_error();
            return Err(ErrorResponse::from_error(GraphQLError::persisted_query_mismatch()));
        }
        // Store the query (best-effort; log on failure).
        if let Err(e) = apq_store.set(hash.to_owned(), body.to_owned()).await {
            warn!(error = %e, "Failed to store APQ query — proceeding without caching");
            apq_metrics.record_error();
        } else {
            apq_metrics.record_store();
        }
        Ok(body.to_owned())
    } else {
        // Hash only: look up.
        match apq_store.get(hash).await {
            Ok(Some(stored)) => {
                apq_metrics.record_hit();
                Ok(stored)
            },
            Ok(None) => {
                apq_metrics.record_miss();
                Err(ErrorResponse::from_error(GraphQLError::persisted_query_not_found()))
            },
            Err(e) => {
                warn!(error = %e, "APQ store lookup failed — treating as miss");
                apq_metrics.record_error();
                Err(ErrorResponse::from_error(GraphQLError::persisted_query_not_found()))
            },
        }
    }
}

/// Shared GraphQL execution logic for both GET and POST handlers.
#[tracing::instrument(skip_all, fields(operation_name = request.operation_name.as_deref().unwrap_or("anonymous")))]
async fn execute_graphql_request(
    state: AppState,
    mut request: GraphQLRequest,
    #[cfg(feature = "federation")] _trace_context: Option<
        fraiseql_core::federation::FederationTraceContext,
    >,
    #[cfg(not(feature = "federation"))] _trace_context: Option<()>,
    security_context: Option<SecurityContext>,
    headers: &HeaderMap,
    peer_ip: &str,
) -> Result<GraphQLResponse, ErrorResponse> {
    // ── Who is asking ────────────────────────────────────────────────────────
    //
    // Every `await` on a stage is `Box::pin`ned. An awaited `async fn` embeds its
    // whole state machine in the caller's future, so the nine stages compose into
    // a type deep enough that rustc 1.97 refuses to compute the handler's layout
    // ("queries overflow the depth limit"). Boxing also keeps the request future
    // small enough for `clippy::large_futures`, which this crate has tripped
    // repeatedly on `Server`-shaped futures. The cost is one allocation per stage
    // per request, against a database round-trip.
    let mut security_context =
        Box::pin(stages::authenticate(&state, headers, security_context)).await?;
    security_context = stages::stamp_trace_context(headers, security_context);
    #[cfg(feature = "auth")]
    Box::pin(stages::enrich_identity(&state, &mut security_context)).await?;

    // ── What they are asking ─────────────────────────────────────────────────
    let query = Box::pin(stages::resolve_query_body(&state, &mut request)).await?;

    let start_time = Instant::now();
    let metrics = &state.metrics;
    metrics.queries_total.fetch_add(1, Ordering::Relaxed);

    // Reason (F041): per-request execution log moved to `debug!`. At >100 RPS
    // this event drowns the operator's `info!`-level signal-to-noise ratio.
    // `info!` is reserved for startup/shutdown/schema-reload events.
    debug!(
        query_length = query.len(),
        has_variables = request.variables.is_some(),
        operation_name = ?request.operation_name,
        "Executing GraphQL query"
    );

    // ── Whether they may ─────────────────────────────────────────────────────
    stages::enforce_introspection_policy(&state, &query, security_context.as_ref())?;
    stages::validate_request(&state, &query, &request, peer_ip)?;

    #[cfg(feature = "federation")]
    let cb_entity_types =
        stages::check_federation_circuit_breakers(&state, &query, request.variables.as_ref())?;

    // Resolve the tenant key — the token's tenant for an authenticated caller, the
    // client hints otherwise — through the same seam the MCP transport uses (#858). A
    // refusal keeps its own class: a header naming a tenant the token is not bound to
    // is FORBIDDEN, a malformed header a validation error.
    let tenant_key =
        super::tenant_dispatch::resolve_tenant_key(&state, security_context.as_ref(), headers)
            .map_err(|e| ErrorResponse::from_error(GraphQLError::from_fraiseql_error(&e)))?;

    // ── Idempotency (#747) ───────────────────────────────────────────────────
    // A mutation carrying an `Idempotency-Key` header executes at most once per
    // key: a repeat with the same body replays the stored response, a repeat
    // with a different body is a 409 conflict. This is the receiving half of
    // the saga at-least-once dispatch contract — a peer coordinator re-sends a
    // step's mutation under the same key after an ambiguous failure (timeout,
    // connection reset after send) or a crash-recovery replay, and this check
    // is what turns those re-sends into one logical effect. Scoped by tenant and
    // by principal, so a key never replays across tenants nor to a caller other
    // than the one whose gates produced the stored response; mutations only
    // arrive via POST (the GET handler rejects them).
    let idempotency_key = if detect_mutation_name(&query).is_some() {
        headers.get("idempotency-key").and_then(|v| v.to_str().ok()).map(|client_key| {
            let scope = crate::routes::idempotency::IdempotencyScope {
                tenant:    tenant_key.clone(),
                principal: security_context.as_ref().map(|c| c.user_id.to_string()),
                method:    "POST".to_string(),
                path:      "/graphql".to_string(),
            };
            let body_hash = crate::routes::idempotency::hash_body(&serde_json::json!({
                "query": query,
                "variables": request.variables,
                "operationName": request.operation_name,
            }));
            (scope.key(client_key), body_hash)
        })
    } else {
        None
    };
    if let Some((ref key, body_hash)) = idempotency_key {
        match state.idempotency_store.check(key, body_hash).await {
            crate::routes::idempotency::IdempotencyCheck::Replay(stored) => {
                debug!("Replaying stored response for repeated Idempotency-Key mutation");
                return Ok(GraphQLResponse {
                    body: stored.body.unwrap_or(serde_json::Value::Null),
                });
            },
            crate::routes::idempotency::IdempotencyCheck::Conflict => {
                return Err(ErrorResponse::from_error(GraphQLError::idempotency_conflict()));
            },
            crate::routes::idempotency::IdempotencyCheck::New => {},
        }
    }

    // The `before:mutation` chain used to run here, once per request, keyed on the
    // first root field and handed this `variables` map. That shape was bypassable
    // three ways (#1327), so enforcement moved into the engine: the chain now runs
    // from `execute_mutation_impl` — per executed root, in document order, with the
    // arguments the write binds from — on every transport. It is installed on the
    // executor's `RuntimeConfig` at serve time by `prepare_functions_runtime`.
    let variables = request.variables;

    // ── Execution ────────────────────────────────────────────────────────────
    // Dispatch, the suspended-tenant gate and the per-tenant quotas all live in
    // the shared seam so the MCP transport enforces the identical policy (#858).
    // `dispatch` holds the concurrency permit for the rest of this scope, which
    // is why it is not extracted into a stage.
    let dispatch = super::tenant_dispatch::dispatch_to_tenant(&state, tenant_key.as_deref())
        .map_err(|e| ErrorResponse::from_error(tenant_dispatch_error(&e)))?;
    let executor = &dispatch.executor;

    // M-quotas (cost): one estimate serves both budget enforcement and
    // observability (#379). Rejection surfaces at the same chokepoint as the
    // other per-tenant quotas, in the shared seam for the same reason (#858).
    let estimated_cost =
        super::tenant_dispatch::estimate_request_cost(&query, variables.as_ref(), executor);
    super::tenant_dispatch::charge_cost_budget(
        &state,
        tenant_key.as_deref(),
        security_context.as_ref(),
        estimated_cost,
    )
    .map_err(|e| ErrorResponse::from_error(tenant_dispatch_error(&e)))?;

    // Preserve subject for audit logging before security_context is consumed.
    #[cfg(feature = "auth")]
    let audit_subject = security_context.as_ref().map(|ctx| ctx.user_id.to_string());
    // Error propagation is deferred so the circuit-breaker outcome is recorded first.
    // GraphQL § 6.1 — the request's `operationName` selects which operation runs.
    // Before this was threaded, a document carrying two operations always ran the
    // first one, whatever the client named.
    let operation_name = request.operation_name.as_deref();
    let exec_result = if let Some(sec_ctx) = security_context {
        executor
            .execute_operation_with_security(&query, variables.as_ref(), &sec_ctx, operation_name)
            .await
    } else {
        executor.execute_operation(&query, variables.as_ref(), operation_name).await
    };

    // Record circuit breaker outcome for federation entity queries
    #[cfg(feature = "federation")]
    if !cb_entity_types.is_empty() {
        if let Some(ref cb_manager) = state.circuit_breaker {
            if exec_result.is_ok() {
                for entity_type in &cb_entity_types {
                    cb_manager.record_success(entity_type);
                }
            } else {
                for entity_type in &cb_entity_types {
                    cb_manager.record_failure(entity_type);
                }
            }
        }
    }

    // Propagate execution errors with metrics
    let op_name = request.operation_name.as_deref().unwrap_or("");
    let result = exec_result.map_err(|e| {
        let elapsed = start_time.elapsed();
        #[allow(clippy::cast_possible_truncation)]
        // Reason: microsecond counter cannot exceed u64 in any practical uptime
        let elapsed_us = elapsed.as_micros() as u64;
        error!(
            error = %e,
            elapsed_ms = elapsed.as_millis(),
            operation_name = ?request.operation_name,
            "Query execution failed"
        );
        metrics.queries_error.fetch_add(1, Ordering::Relaxed);
        metrics.execution_errors_total.fetch_add(1, Ordering::Relaxed);
        // Record duration even for failed queries
        metrics.queries_duration_us.fetch_add(elapsed_us, Ordering::Relaxed);
        metrics.operation_metrics.record(op_name, elapsed_us, true);

        // S46: emit AuthorizationDenied audit event for compliance (SOC 2).
        // Must be emitted before error sanitization so we log the real reason.
        #[cfg(feature = "auth")]
        if matches!(e, fraiseql_core::FraiseQLError::Authorization { .. }) {
            use fraiseql_auth::audit::logger::{
                AuditEntry, AuditEventType, SecretType, get_audit_logger,
            };
            let resource =
                if let fraiseql_core::FraiseQLError::Authorization { ref resource, .. } = e {
                    resource.clone().unwrap_or_else(|| op_name.to_string())
                } else {
                    op_name.to_string()
                };
            get_audit_logger().log_entry(AuditEntry {
                event_type:    AuditEventType::AuthorizationDenied,
                secret_type:   SecretType::JwtToken,
                subject:       audit_subject.clone(),
                operation:     op_name.to_string(),
                success:       false,
                error_message: Some(resource),
                context:       Some(format!("peer_ip={peer_ip}")),
                chain_hash:    None,
            });
        }

        let err = state.error_sanitizer.sanitize(GraphQLError::from_fraiseql_error(&e));
        ErrorResponse::from_error(err)
    })?;

    let elapsed = start_time.elapsed();
    #[allow(clippy::cast_possible_truncation)]
    // Reason: microsecond counter cannot exceed u64 in any practical uptime
    let elapsed_us = elapsed.as_micros() as u64;

    // Record successful query metrics
    metrics.queries_success.fetch_add(1, Ordering::Relaxed);
    metrics.queries_duration_us.fetch_add(elapsed_us, Ordering::Relaxed);
    metrics.db_queries_total.fetch_add(1, Ordering::Relaxed);
    metrics.db_queries_duration_us.fetch_add(elapsed_us, Ordering::Relaxed);
    metrics.operation_metrics.record(op_name, elapsed_us, false);

    // #379 acceptance: the audit trail records the estimated cost alongside
    // tenant and operation for every executed request, and the running sum is
    // exported at /metrics so budgets can be sized from observed traffic.
    // Per-request, but on a dedicated target (the `mutation_audit` pattern) so
    // operators subscribe to it explicitly instead of it flooding the default
    // `info!` stream.
    if let Some(cost) = estimated_cost {
        metrics.queries_cost_total.fetch_add(cost, Ordering::Relaxed);
        tracing::info!(
            target: "fraiseql::cost_audit",
            cost,
            tenant = tenant_key.as_deref().unwrap_or(""),
            operation = %op_name,
            "operation cost"
        );
    }

    // Record federation-specific metrics for federation queries
    #[cfg(feature = "federation")]
    if fraiseql_core::federation::is_federation_query(&query) {
        metrics.record_entity_resolution(elapsed_us, true);
    }

    debug!(
        elapsed_ms = elapsed.as_millis(),
        operation_name = ?request.operation_name,
        "Query executed successfully"
    );

    // ── Post-processing ──────────────────────────────────────────────────────
    #[allow(unused_mut)]
    // Reason: mut is required by decrypt_response_fields(&mut ...) under the secrets feature
    let mut response_json = result;

    #[cfg(feature = "secrets")]
    Box::pin(stages::decrypt_response_fields(&state, &mut response_json)).await?;

    // Idempotency (#747): persist the successful response so a re-send of the
    // same mutation under the same key replays it instead of executing again.
    // Only success is stored — a failed mutation stays retryable under its key.
    if let Some((key, body_hash)) = idempotency_key {
        state
            .idempotency_store
            .store(
                key,
                body_hash,
                crate::routes::idempotency::StoredResponse {
                    status:  200,
                    headers: Vec::new(),
                    body:    Some(response_json.clone()),
                },
            )
            .await;
    }

    Ok(GraphQLResponse {
        body: response_json,
    })
}

/// Map a tenant-dispatch error from [`AppState::executor_for_tenant`] to the
/// correct GraphQL error code (#332).
///
/// `executor_for_tenant` returns [`FraiseQLError::Authorization`] for an unknown
/// tenant key (→ 403 Forbidden) and [`FraiseQLError::ServiceUnavailable`] for a
/// suspended tenant (→ 503 with a `Retry-After` header carrying `retry_after`).
/// Previously both collapsed to 403, discarding the variant and the retry hint.
pub(super) fn tenant_dispatch_error(error: &FraiseQLError) -> GraphQLError {
    match error {
        FraiseQLError::ServiceUnavailable { retry_after, .. } => {
            GraphQLError::service_unavailable(error.to_string(), *retry_after)
        },
        // Per-tenant concurrency limit reached (M-quotas) → 429 Too Many Requests.
        FraiseQLError::RateLimited { .. } => GraphQLError::rate_limited(error.to_string()),
        // Cost rejections carry their own codes (#379): a per-request ceiling
        // is permanent for the operation (200 + errors[], no retry invitation),
        // an exhausted rolling window is retryable (429 + Retry-After).
        FraiseQLError::CostExceeded {
            retry_after_secs, ..
        } => match retry_after_secs {
            Some(secs) => GraphQLError::cost_budget_exhausted(error.to_string(), *secs),
            None => GraphQLError::operation_cost_exceeded(error.to_string()),
        },
        // Unknown tenant key (Authorization) and any other dispatch error stay
        // 403 Forbidden, preserving the prior behaviour.
        _ => GraphQLError::new(error.to_string(), crate::error::ErrorCode::Forbidden),
    }
}

mod incremental;
mod sse;
mod stages;

#[cfg(test)]
mod tests;