Skip to content

Commit 039b265

Browse files
authored
feat(middleware): define HTTP response pre-return interface (#3073)
* feat(middleware): define HTTP response pre-return interface Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): clarify HTTP response interface Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): align response result actions Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): expose response reason codes Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): share session end reasons Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * feat(middleware)!: finalize HTTP response pre-return contract Replace the separate body_end event with HttpResponseBodyUnit.end_of_stream. Every body-inspecting stage receives exactly one flagged unit, which may be empty; a zero-byte body is one empty flagged unit and OpenShell never reads ahead to set the flag. Defer response trailers from V1 and reserve their field numbers. HTTP/1.0 clients and Content-Length bodies cannot carry trailers and that behavior was undefined. Add HttpResponsePreflight.permitted_body_modes, computed once from the original upstream head so every stage sees the same list, and make an unlisted selection a failure rather than a downgrade. Add the block_delivery preflight action as a successful decision enforced regardless of on_error. Expose Content-Length, Content-Encoding, and Content-Range read-only in preflight. Cap STREAM_BYTES input units at half of max_payload_bytes and permit deferring bytes across replacements only for fail_closed bindings, surfaced as deferral_permitted. Split PEER_DISCONNECT into DOWNSTREAM_DISCONNECT and UPSTREAM_DISCONNECT and attribute WebSocket relay failures by direction instead of a generic peer error. Compile the content-guard example in lint and branch checks so proto renames cannot break it silently. BREAKING CHANGE: WebSocketSessionEndReason and WebSocketSessionEnd are replaced by the shared MiddlewareSessionEndReason and MiddlewareSessionEnd. NORMAL_CLOSE is now NORMAL, UPSTREAM_REJECTED is now UPSTREAM_FAILURE, and PEER_DISCONNECT is split into DOWNSTREAM_DISCONNECT and UPSTREAM_DISCONNECT. Enum numbers are unchanged so binary wire compatibility is preserved; generated symbols and JSON names change. Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): describe skip as opting out of inspection Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * feat(middleware): add body-phase block_delivery and skip_remaining actions Body results may now stop delivery or opt out of inspecting the rest of the response after a prefix. One HttpResponseBlockDelivery message is shared by preflight and body results and documents the difference between blocking before and after head commitment. Drop the field reservations, since nothing in this contract has shipped, and renumber session_end to close the gap. Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): share HTTP body leaf messages across directions HttpBodyUnit, HttpBodyPassThrough, HttpBodyTransform, HttpBodySkipRemaining, and HttpBodyMode carry no response-specific semantics, so name them for reuse by the streaming request hook. Envelopes, results, preflight, and block_delivery stay response-specific because commitment semantics differ. Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): keep HTTP body leaf messages response-specific Reverts the shared HttpBody* naming. A direction-specific payload such as a response-only semantic mode would otherwise add unreachable variants to the other direction or force a source-breaking fork after 0.1.0. The streaming request hook defines its own HttpRequestBody* messages and copies the shape; SDKs present a direction-neutral body handler over both. Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): simplify response proto comments Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * fix(middleware): reject undispatched response bindings Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): simplify phase field comment Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * refactor(middleware): rename HTTP response preflight result Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * feat(middleware): add HTTP response trailer results Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): define response block delivery Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): define stage-local response body modes Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): define final response body units Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): define streaming response deferral Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): defer whole-body accumulation timeout Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): align response result diagnostics Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * test(middleware): cover upstream WebSocket disconnect Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * fix(middleware): keep response streams unit-local Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> * docs(middleware): trim disconnect compatibility note Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com> --------- Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com>
1 parent c96b9bf commit 039b265

14 files changed

Lines changed: 852 additions & 180 deletions

File tree

‎.github/workflows/branch-checks.yml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -158,12 +158,14 @@ jobs:
158158
cargo fmt --all -- --check
159159
cargo fmt --manifest-path e2e/rust/Cargo.toml --all -- --check
160160
cargo fmt --manifest-path examples/governance-interceptor/Cargo.toml --all -- --check
161+
cargo fmt --manifest-path examples/supervisor-middleware-content-guard/Cargo.toml --all -- --check
161162
162163
- name: Lint
163164
run: |
164165
cargo clippy --workspace --all-targets -- -D warnings
165166
cargo clippy --manifest-path e2e/rust/Cargo.toml --all-targets -- -D warnings
166167
cargo check --manifest-path examples/governance-interceptor/Cargo.toml --all-targets
168+
cargo check --manifest-path examples/supervisor-middleware-content-guard/Cargo.toml --all-targets
167169
168170
- name: Test
169171
env:

‎crates/openshell-core/src/middleware.rs‎

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,11 +11,17 @@ use tokio::sync::mpsc;
1111
use tonic::{Request, Response, Status};
1212

1313
use crate::proto::{
14-
HttpHeader, HttpRequestEvaluation, HttpRequestResult, HttpRequestTarget, MiddlewareManifest,
15-
RequestContext, SupervisorMiddlewarePhase, ValidateConfigRequest, ValidateConfigResponse,
16-
WebSocketSessionEvent, WebSocketSessionEventResult,
14+
HttpHeader, HttpRequestEvaluation, HttpRequestResult, HttpRequestTarget, HttpResponseEvent,
15+
HttpResponseEventResult, MiddlewareManifest, RequestContext, SupervisorMiddlewarePhase,
16+
ValidateConfigRequest, ValidateConfigResponse, WebSocketSessionEvent,
17+
WebSocketSessionEventResult,
1718
};
1819

20+
/// Transport-neutral result stream for one HTTP response middleware stage.
21+
pub type HttpResponseResultStream = Pin<
22+
Box<dyn tokio_stream::Stream<Item = Result<HttpResponseEventResult, Status>> + Send + 'static>,
23+
>;
24+
1925
/// Transport-neutral response stream for one WebSocket middleware stage.
2026
pub type WebSocketResponseStream = Pin<
2127
Box<
@@ -47,6 +53,15 @@ pub trait SupervisorMiddlewareEndpoint: Send + Sync {
4753
&self,
4854
requests: mpsc::Receiver<WebSocketSessionEvent>,
4955
) -> Result<WebSocketResponseStream, Status>;
56+
57+
async fn open_http_response_pre_return(
58+
&self,
59+
_requests: mpsc::Receiver<HttpResponseEvent>,
60+
) -> Result<HttpResponseResultStream, Status> {
61+
Err(Status::unimplemented(
62+
"middleware does not implement HTTP response pre-return evaluation",
63+
))
64+
}
5065
}
5166

5267
/// Borrowed request state exposed to one in-process middleware invocation.
@@ -242,6 +257,18 @@ pub trait InProcessMiddleware: Send + Sync {
242257
"middleware does not implement WebSocket sessions",
243258
))
244259
}
260+
261+
/// Open one HTTP response pre-return stream.
262+
///
263+
/// Request-only implementations may keep the default unsupported response.
264+
async fn open_http_response_pre_return(
265+
&self,
266+
_requests: mpsc::Receiver<HttpResponseEvent>,
267+
) -> std::result::Result<HttpResponseResultStream, Status> {
268+
Err(Status::unimplemented(
269+
"middleware does not implement HTTP response pre-return evaluation",
270+
))
271+
}
245272
}
246273

247274
/// Default timeout for one supervisor middleware RPC.

‎crates/openshell-supervisor-middleware/src/headers.rs‎

Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ pub const MAX_HEADER_MUTATION_BYTES: usize = 32 * 1024;
1616
pub enum HeaderAuthority {
1717
Request,
1818
Response,
19+
ResponseTrailers,
1920
}
2021

2122
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -30,6 +31,7 @@ pub enum HeaderMutationError {
3031
InvalidExistingAction,
3132
MissingExistingAction { name: String },
3233
UnsupportedExistingAction,
34+
AbsentTrailerName { name: String },
3335
Empty,
3436
}
3537

@@ -47,6 +49,7 @@ impl HeaderMutationError {
4749
Self::InvalidExistingAction => "header_mutation_invalid_existing_action",
4850
Self::MissingExistingAction { .. } => "header_mutation_missing_existing_action",
4951
Self::UnsupportedExistingAction => "header_mutation_unsupported_existing_action",
52+
Self::AbsentTrailerName { .. } => "trailer_mutation_absent_name",
5053
Self::Empty => "header_mutation_empty",
5154
}
5255
}
@@ -104,6 +107,10 @@ impl std::fmt::Display for HeaderMutationError {
104107
"middleware returned unsupported on_existing action"
105108
)
106109
}
110+
Self::AbsentTrailerName { name } => write!(
111+
formatter,
112+
"middleware cannot create absent response trailer '{name}'"
113+
),
107114
Self::Empty => write!(formatter, "middleware returned an empty header mutation"),
108115
}
109116
}
@@ -133,6 +140,15 @@ pub fn apply(
133140
Some(header_mutation::Operation::Write(write)) => {
134141
let name = validate_name(&write.name)?;
135142
validate_authority(authority, MutationKind::Write, &write.name, &name)?;
143+
if authority == HeaderAuthority::ResponseTrailers
144+
&& !existing_headers
145+
.iter()
146+
.any(|existing| existing.name.eq_ignore_ascii_case(&name))
147+
{
148+
return Err(HeaderMutationError::AbsentTrailerName {
149+
name: write.name.clone(),
150+
});
151+
}
136152
if is_connection_nominated(connection_nominated_headers, &name) {
137153
return Err(HeaderMutationError::HopByHop {
138154
name: write.name.clone(),
@@ -229,6 +245,7 @@ fn validate_authority(
229245
is_response_protected(normalized_name)
230246
|| (kind == MutationKind::Write && is_response_remove_only(normalized_name))
231247
}
248+
HeaderAuthority::ResponseTrailers => is_response_protected(normalized_name),
232249
};
233250
if protected {
234251
return Err(HeaderMutationError::Protected {
@@ -650,4 +667,115 @@ mod tests {
650667
}
651668
}
652669
}
670+
671+
#[test]
672+
fn empty_response_trailers_accept_pass_through() {
673+
assert_eq!(
674+
apply(HeaderAuthority::ResponseTrailers, &[], &[], &[]),
675+
Ok(Vec::new())
676+
);
677+
}
678+
679+
#[test]
680+
fn response_trailer_pass_through_preserves_fields_and_order() {
681+
let existing = [
682+
header("x-checksum", "one"),
683+
header("x-trace", "middle"),
684+
header("x-checksum", "two"),
685+
];
686+
687+
let updated = apply(HeaderAuthority::ResponseTrailers, &existing, &[], &[])
688+
.expect("empty mutation list");
689+
690+
assert_eq!(updated, existing);
691+
}
692+
693+
#[test]
694+
fn response_trailer_mutations_modify_remove_and_preserve_order() {
695+
let existing = [
696+
header("x-checksum", "one"),
697+
header("x-remove", "gone"),
698+
header("x-trace", "middle"),
699+
header("x-checksum", "two"),
700+
];
701+
702+
let updated = apply(
703+
HeaderAuthority::ResponseTrailers,
704+
&existing,
705+
&[],
706+
&[
707+
write("X-Checksum", "replacement", ExistingHeaderAction::Overwrite),
708+
remove("X-Remove"),
709+
write("X-Trace", "last", ExistingHeaderAction::Append),
710+
],
711+
)
712+
.expect("permitted response trailer mutations");
713+
714+
assert_eq!(
715+
updated,
716+
vec![
717+
header("x-trace", "middle"),
718+
header("x-checksum", "replacement"),
719+
header("x-trace", "last"),
720+
]
721+
);
722+
}
723+
724+
#[test]
725+
fn response_trailer_write_cannot_introduce_an_absent_name() {
726+
let existing = [header("x-checksum", "one")];
727+
let error = apply(
728+
HeaderAuthority::ResponseTrailers,
729+
&existing,
730+
&[],
731+
&[write(
732+
"X-New-Trailer",
733+
"value",
734+
ExistingHeaderAction::Overwrite,
735+
)],
736+
)
737+
.expect_err("absent response trailer name");
738+
739+
assert_eq!(
740+
error,
741+
HeaderMutationError::AbsentTrailerName {
742+
name: "X-New-Trailer".into()
743+
}
744+
);
745+
}
746+
747+
#[test]
748+
fn response_trailer_removal_of_an_absent_name_is_a_noop() {
749+
let existing = [header("x-checksum", "one")];
750+
let updated = apply(
751+
HeaderAuthority::ResponseTrailers,
752+
&existing,
753+
&[],
754+
&[remove("X-Missing")],
755+
)
756+
.expect("absent response trailer removal");
757+
758+
assert_eq!(updated, existing);
759+
}
760+
761+
#[test]
762+
fn response_trailer_protected_fields_cannot_be_mutated() {
763+
for mutation in [
764+
write("Content-Length", "10", ExistingHeaderAction::Overwrite),
765+
remove("Set-Cookie"),
766+
] {
767+
let error = apply(
768+
HeaderAuthority::ResponseTrailers,
769+
&[
770+
header("content-length", "5"),
771+
header("set-cookie", "session=upstream"),
772+
],
773+
&[],
774+
&[mutation],
775+
)
776+
.expect_err("protected response trailer mutation");
777+
778+
assert!(matches!(error, HeaderMutationError::Protected { .. }));
779+
}
780+
}
653781
}

0 commit comments

Comments
 (0)