praxis-proxy-protocol 0.7.3

HTTP, TCP, and protocol adapters for Praxis
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
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2024 Praxis Contributors

//! Conversions between Pingora types and Praxis transport-agnostic types.
//!
//! This module bridges Pingora's request/response lifecycle to Praxis's
//! transport-agnostic filter pipeline. The conversion functions run on every
//! request and are marked inline to minimize overhead at this boundary.
//!
//! The rejection path handles both plain HTTP responses and gRPC trailers-only
//! error mapping, dispatching based on whether the pipeline armed gRPC mode for
//! a given request.

use pingora_proxy::Session;
use praxis_filter::{Rejection, Request, Response};
use tracing::debug;

// -----------------------------------------------------------------------------
// Pingora - Request / Response Conversion
// -----------------------------------------------------------------------------

/// Build a transport-agnostic [`Request`] from a Pingora session.
///
/// ```ignore
/// // Requires a `pingora_proxy::Session` which cannot be constructed
/// // outside of a Pingora request lifecycle.
/// use praxis_protocol::http::pingora::convert::request_header_from_session;
///
/// let req = request_header_from_session(&mut session);
/// assert!(!req.method.is_safe());
/// ```
///
/// [`Request`]: praxis_filter::Request
// Hot path: called per-request, cross-crate boundary.
#[inline]
pub(crate) fn request_header_from_session(session: &mut Session) -> Request {
    let req = session.req_header_mut();
    let method = req.method.clone();
    let uri = req.uri.clone();
    let headers = req.headers.clone();

    Request { method, uri, headers }
}

/// Build a transport-agnostic [`Response`] by copying a Pingora response's
/// status and headers.
///
/// Clones the [`HeaderMap`]: Pingora 0.9.0 removed `DerefMut` on
/// `ResponseHeader`, so the map can no longer be moved out in place, and the
/// original must stay intact regardless so its status stays readable for the
/// retry and commit checks that run before write-back.
///
/// ```ignore
/// // Requires `pingora_http::ResponseHeader` from Pingora internals.
/// use praxis_protocol::http::pingora::convert::response_header_from_pingora;
///
/// let upstream = pingora_http::ResponseHeader::build(200, None).unwrap();
/// let resp = response_header_from_pingora(&upstream);
/// assert_eq!(resp.status.as_u16(), 200);
/// ```
///
/// [`Response`]: praxis_filter::Response
/// [`HeaderMap`]: http::HeaderMap
// Hot path: called per-request, cross-crate boundary.
#[inline]
pub(crate) fn response_header_from_pingora(upstream: &pingora_http::ResponseHeader) -> Response {
    Response {
        status: upstream.status,
        headers: upstream.headers.clone(),
    }
}

// -----------------------------------------------------------------------------
// Pingora - Rejection
// -----------------------------------------------------------------------------

/// Send a rejection, rendering it in gRPC's shape when the pipeline
/// armed one for this request.
///
/// Follows a three-path dispatch:
/// 1. If the rejection carries a `grpc-status` header (gRPC filter built it), send it as a gRPC rejection with framing
///    but no status remapping.
/// 2. If a `GrpcErrorMapping` is armed and the status is an error (>= 400), map the HTTP status to a gRPC code and send
///    trailers-only.
/// 3. Otherwise, send as a plain HTTP response.
///
/// Only error statuses become gRPC statuses: a successful short-circuit
/// (a CORS preflight `204`, a `static_response` `200`) is a real HTTP
/// response the client asked for, not a proxy error.
pub(crate) async fn send_rejection_for(
    session: &mut Session,
    rejection: Rejection,
    ctx: &crate::http::pingora::context::PingoraRequestCtx,
) {
    // Path 1: gRPC filter already set grpc-status, just frame it.
    if rejection
        .headers
        .iter()
        .any(|(name, _value)| name.eq_ignore_ascii_case("grpc-status"))
    {
        crate::http::pingora::grpc_trailers::send_grpc_rejection(session, &rejection).await;
        return;
    }

    // Path 2: GrpcErrorMapping armed and status is an error, map to gRPC.
    let mapping = ctx.extensions.get::<praxis_filter::GrpcErrorMapping>();
    if let Some(mapping) = mapping.filter(|_mapping| rejection.status >= 400) {
        let message = grpc_error_message(&rejection);
        crate::http::pingora::grpc_trailers::send_trailers_only(session, mapping, rejection.status, message).await;
        return;
    }
    // Path 3: plain HTTP rejection.
    send_rejection(session, rejection).await;
}

/// The text to carry as `grpc-message` for a rejection.
///
/// Prefers the rejection's own body when present: filters set it to name the
/// rule that fired (for example "rate limit exceeded"), giving clients a more
/// specific diagnostic than the HTTP status's canonical reason phrase.
/// Falls back through: body text → status reason phrase → "proxy error".
fn grpc_error_message(rejection: &Rejection) -> &str {
    rejection
        .body
        .as_deref()
        .and_then(|body| str::from_utf8(body).ok())
        .filter(|text| !text.is_empty())
        .or_else(|| {
            http::StatusCode::from_u16(rejection.status)
                .ok()
                .and_then(|status| status.canonical_reason())
        })
        .unwrap_or("proxy error")
}

/// Send a rejection response to the client, including any headers and body from the [`Rejection`].
///
/// Disables downstream keep-alive by default so the connection closes after
/// a short-circuit response. Complete responses may explicitly preserve it.
///
/// ```ignore
/// // Requires an active `pingora_proxy::Session` from a live request.
/// use praxis_protocol::http::pingora::convert::send_rejection;
///
/// let rejection = praxis_filter::Rejection::status(403);
/// send_rejection(&mut session, rejection).await;
/// ```
///
/// [`Rejection`]: praxis_filter::Rejection
pub(crate) async fn send_rejection(session: &mut Session, rejection: Rejection) {
    debug!(status = rejection.status, "sending rejection response");
    if !rejection.preserve_keepalive {
        session.set_keepalive(None);
    }

    let mut header = build_rejection_header(&rejection);
    let has_body = rejection.body.is_some();
    if let Some(body) = &rejection.body {
        let _insert = header.insert_header("content-length", body.len().to_string());
    }
    if let Err(e) = session.write_response_header(Box::new(header), !has_body).await {
        debug!(error = %e, "failed to write rejection response header");
        return;
    }
    if let Some(body) = rejection.body
        && let Err(e) = session.write_response_body(Some(body), true).await
    {
        debug!(error = %e, "failed to write rejection response body");
    }
}

/// Build a Pingora [`ResponseHeader`] from a [`Rejection`], falling back
/// to 500 if the status code is not a final status.
///
/// The 500 fallback handles rejections with out-of-range status codes
/// (for example, manually constructed `Rejection` with an invalid value)
/// and informational `1xx` codes, which per [RFC 9110 Section 15.2] cannot
/// complete a response. This ensures the proxy always sends a valid final
/// HTTP response even if a filter produced a malformed rejection.
///
/// [RFC 9110 Section 15.2]: https://datatracker.ietf.org/doc/html/rfc9110#section-15.2
/// [`ResponseHeader`]: pingora_http::ResponseHeader
/// [`Rejection`]: praxis_filter::Rejection
fn build_rejection_header(rejection: &Rejection) -> pingora_http::ResponseHeader {
    let header_count = Some(
        rejection
            .headers
            .len()
            .saturating_add(rejection.header_map.as_ref().map_or(0, |headers| headers.len())),
    );
    let status = if (200..=599).contains(&rejection.status) {
        rejection.status
    } else {
        tracing::error!(status = rejection.status, "non-final rejection status; using 500");
        500
    };
    let mut header = match pingora_http::ResponseHeader::build(status, header_count) {
        Ok(h) => h,
        Err(e) => {
            tracing::error!(status = rejection.status, error = %e, "invalid rejection status; using 500");
            #[expect(clippy::expect_used, reason = "500 is a valid status code")]
            pingora_http::ResponseHeader::build(500, header_count).expect("500 is a valid status code")
        },
    };
    append_rejection_headers(&mut header, rejection);
    header
}

/// Append a rejection's filter-supplied headers, dropping reserved and
/// hop-by-hop ones.
///
/// Reserved internal (x-praxis-* / x-ext-*) headers never reach the client:
/// the upstream-response and terminal-response paths both strip them, and a
/// rejection built from filter-supplied headers must hold the same invariant.
/// Response hop-by-hop headers ([RFC 9110 Section 7.6.1]) are dropped too:
/// framing and connection management belong to the proxy, so a filter-set
/// `transfer-encoding` or `connection` cannot contradict the
/// `content-length` and keep-alive state Praxis applies itself.
///
/// [RFC 9110 Section 7.6.1]: https://datatracker.ietf.org/doc/html/rfc9110#section-7.6.1
fn append_rejection_headers(header: &mut pingora_http::ResponseHeader, rejection: &Rejection) {
    for (name, value) in &rejection.headers {
        if is_dropped_rejection_header(name) {
            debug!(header = %name, "dropping reserved or hop-by-hop header from rejection response");
            continue;
        }
        let _append = header.append_header(name.clone(), value.clone());
    }
    if let Some(headers) = &rejection.header_map {
        for (name, value) in headers.iter() {
            if is_dropped_rejection_header(name.as_str()) {
                debug!(header = %name, "dropping reserved or hop-by-hop header from rejection response");
                continue;
            }
            let _append = header.append_header(name.clone(), value.clone());
        }
    }
}

/// Whether a filter-supplied rejection header must be withheld from the
/// client: reserved internal headers and response hop-by-hop headers.
fn is_dropped_rejection_header(name: &str) -> bool {
    praxis_core::reserved_headers::is_reserved(name)
        || praxis_core::reserved_headers::RESPONSE_HOP_BY_HOP_HEADERS
            .iter()
            .any(|hop| name.eq_ignore_ascii_case(hop))
}

// -----------------------------------------------------------------------------
// Tests
// -----------------------------------------------------------------------------

#[cfg(test)]
#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
#[allow(
    clippy::unwrap_used,
    clippy::expect_used,
    clippy::indexing_slicing,
    clippy::too_many_lines,
    reason = "tests"
)]
mod tests {
    use http::StatusCode;

    use super::*;

    #[test]
    fn response_header_preserves_status() {
        let upstream = pingora_http::ResponseHeader::build(200, None).unwrap();
        let resp = response_header_from_pingora(&upstream);
        assert_eq!(resp.status, StatusCode::OK, "status should be 200 OK");
    }

    #[test]
    fn response_header_preserves_headers() {
        let mut upstream = pingora_http::ResponseHeader::build(200, Some(2)).unwrap();
        let _insert1 = upstream.insert_header("x-custom", "value");
        let _insert2 = upstream.insert_header("content-type", "text/plain");

        let resp = response_header_from_pingora(&upstream);
        assert_eq!(
            resp.headers.get("x-custom").unwrap(),
            "value",
            "x-custom header should be preserved"
        );
        assert_eq!(
            resp.headers.get("content-type").unwrap(),
            "text/plain",
            "content-type header should be preserved"
        );
    }

    #[test]
    fn response_header_copies_headers_from_upstream() {
        let mut upstream = pingora_http::ResponseHeader::build(200, Some(1)).unwrap();
        let _insert = upstream.insert_header("x-test", "taken");

        let resp = response_header_from_pingora(&upstream);
        assert_eq!(
            resp.headers.get("x-test").unwrap(),
            "taken",
            "header should be in response"
        );
        assert!(
            !upstream.headers.is_empty(),
            "upstream headers stay intact after the copy"
        );
    }

    #[test]
    fn response_header_empty_headers() {
        let upstream = pingora_http::ResponseHeader::build(404, None).unwrap();
        let resp = response_header_from_pingora(&upstream);
        assert_eq!(resp.status, StatusCode::NOT_FOUND, "status should be 404 Not Found");
        assert!(
            resp.headers.is_empty(),
            "headers should be empty when upstream has none"
        );
    }

    #[test]
    fn rejection_header_strips_reserved_internal_headers() {
        let mut header_map = http::HeaderMap::new();
        header_map.insert("x-ext-protocol-task", http::HeaderValue::from_static("meta"));
        header_map.insert("x-request-id", http::HeaderValue::from_static("abc-123"));
        let mut rejection = Rejection::status(403)
            .with_header("x-praxis-route", "internal-cluster")
            .with_header("content-type", "text/plain");
        rejection.header_map = Some(Box::new(header_map));

        let header = build_rejection_header(&rejection);
        assert!(
            header.headers.get("x-praxis-route").is_none(),
            "reserved x-praxis-* header must never reach the client via a rejection"
        );
        assert!(
            header.headers.get("x-ext-protocol-task").is_none(),
            "reserved x-ext-* header must never reach the client via a rejection header_map"
        );
        assert_eq!(
            header.headers.get("content-type").map(http::HeaderValue::as_bytes),
            Some(b"text/plain".as_slice()),
            "non-reserved rejection headers must be preserved"
        );
        assert_eq!(
            header.headers.get("x-request-id").map(http::HeaderValue::as_bytes),
            Some(b"abc-123".as_slice()),
            "non-reserved header_map entries must be preserved"
        );
    }

    #[test]
    fn rejection_header_strips_mixed_case_reserved_from_string_list() {
        let rejection = Rejection::status(403)
            .with_header("X-Praxis-Route", "internal-cluster")
            .with_header("X-Ext-Agent-Task", "meta")
            .with_header("X-Custom", "keep");

        let header = build_rejection_header(&rejection);
        assert!(
            header.headers.get("x-praxis-route").is_none(),
            "a mixed-case reserved x-praxis-* header must be dropped from a rejection"
        );
        assert!(
            header.headers.get("x-ext-agent-task").is_none(),
            "a mixed-case reserved x-ext-agent-* header must be dropped from a rejection"
        );
        assert_eq!(
            header.headers.get("x-custom").map(http::HeaderValue::as_bytes),
            Some(b"keep".as_slice()),
            "non-reserved rejection headers must still be preserved"
        );
    }

    #[test]
    fn rejection_header_preserves_duplicate_values() {
        let rejection = Rejection::status(200)
            .with_header("set-cookie", "first=1")
            .with_header("set-cookie", "second=2");

        let header = build_rejection_header(&rejection);
        let values: Vec<_> = header
            .headers
            .get_all("set-cookie")
            .iter()
            .map(|value| value.to_str().unwrap())
            .collect();

        assert_eq!(values, ["first=1", "second=2"]);
    }

    #[test]
    fn rejection_header_preserves_opaque_values() {
        let mut rejection = Rejection::status(200);
        rejection
            .header_map
            .get_or_insert_with(Default::default)
            .append("x-opaque", http::HeaderValue::from_bytes(&[b'a', 0x80, b'z']).unwrap());

        let header = build_rejection_header(&rejection);

        assert_eq!(header.headers["x-opaque"].as_bytes(), &[b'a', 0x80, b'z']);
    }

    #[test]
    fn invalid_rejection_status_falls_back_to_500() {
        let rejection = Rejection {
            body: None,
            headers: Vec::new(),
            header_map: None,
            preserve_keepalive: false,
            status: 99,
        };
        let header = build_rejection_header(&rejection);
        assert_eq!(header.status.as_u16(), 500, "invalid status codes must map to 500");
    }

    #[test]
    fn rejection_header_drops_hop_by_hop_headers() {
        let rejection = Rejection::status(403)
            .with_header("Transfer-Encoding", "chunked")
            .with_header("connection", "close")
            .with_header("content-type", "text/plain");

        let header = build_rejection_header(&rejection);

        assert!(
            header.headers.get("transfer-encoding").is_none(),
            "filter-supplied transfer-encoding must not reach the client"
        );
        assert!(
            header.headers.get("connection").is_none(),
            "filter-supplied connection must not reach the client"
        );
        assert_eq!(
            header.headers.get("content-type").map(http::HeaderValue::as_bytes),
            Some(b"text/plain".as_slice()),
            "end-to-end rejection headers must be preserved"
        );
    }

    #[test]
    fn rejection_header_informational_status_becomes_500() {
        let header = build_rejection_header(&Rejection::status(101));

        assert_eq!(
            header.status.as_u16(),
            500,
            "1xx rejection status must fall back to 500"
        );
    }

    #[test]
    fn rejection_header_final_status_unchanged() {
        let header = build_rejection_header(&Rejection::status(204));

        assert_eq!(header.status.as_u16(), 204, "final rejection status must be kept");
    }
}