Skip to content

Commit 1a4a2bb

Browse files
committed
fix(examples): split expanded content guard output within the chunk limit
Redacting a term shorter than the marker makes STREAM output longer than its input, so one output chunk could exceed the max_chunk_bytes the preflight offered and fail the message closed. Split every output chunk within that limit. Signed-off-by: Piotr Mlocek <pmlocek@nvidia.com>
1 parent c2684e8 commit 1a4a2bb

2 files changed

Lines changed: 106 additions & 44 deletions

File tree

‎examples/supervisor-middleware-content-guard/README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,7 @@ The echoed JSON body contains `[FILTERED]` instead of the configured term.
9999

100100
## HTTP behavior
101101

102-
At preflight, the guard selects its configured `body_mode`, BUFFERED by default, and falls back to the other mode when OpenShell offers only that one. BUFFERED inspects one complete body of at most 256 KiB. STREAM has no size limit: it releases every complete line at once and withholds only the bytes that may begin a term split across chunks, so line-oriented streams such as server-sent events keep flowing.
102+
At preflight, the guard selects its configured `body_mode`, BUFFERED by default, and falls back to the other mode when OpenShell offers only that one. BUFFERED inspects one complete body of at most 256 KiB. STREAM has no size limit: it releases every complete line at once and withholds only the bytes that may begin a term split across chunks, so line-oriented streams such as server-sent events keep flowing. Redaction can make the output longer than the input, so the guard splits it into chunks within the limit OpenShell offers.
103103

104104
- A request whose path or query contains a configured term is rejected at preflight with reason code `content_match`, before its body is read or the upstream is contacted.
105105
- A message with an empty or absent body continues without inspection.

‎examples/supervisor-middleware-content-guard/src/http.rs‎

Lines changed: 105 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -67,9 +67,10 @@ where
6767

6868
enum Stage {
6969
Preflight,
70-
Begin(GuardConfig, BodyMode),
70+
Begin(GuardConfig, Selected),
7171
Buffered(GuardConfig),
72-
Stream(GuardConfig, StreamScanner),
72+
/// STREAM, with the largest output chunk OpenShell accepts.
73+
Stream(GuardConfig, StreamScanner, usize),
7374
Done,
7475
}
7576

@@ -92,70 +93,69 @@ impl Stage {
9293
diagnostics: Some(diagnostics),
9394
}))])
9495
}
95-
Decision::Inspect(config, mode) => {
96-
let selected = match mode {
96+
Decision::Inspect(config, selected) => {
97+
let mode = match selected {
9798
Selected::Buffered(max_body_bytes) => {
98-
*self = Self::Begin(config, BodyMode::Buffered);
9999
http_inspect::Mode::Buffered(HttpBufferedMode { max_body_bytes })
100100
}
101-
Selected::Stream => {
102-
*self = Self::Begin(config, BodyMode::Stream);
103-
http_inspect::Mode::Stream(HttpStreamMode {})
104-
}
101+
Selected::Stream(_) => http_inspect::Mode::Stream(HttpStreamMode {}),
105102
};
103+
*self = Self::Begin(config, selected);
106104
Ok(vec![result(http_result::Result::PreflightResult(
107105
HttpPreflightResult {
108106
decision: Some(http_preflight_result::Decision::Inspect(
109-
HttpInspect {
110-
mode: Some(selected),
111-
},
107+
HttpInspect { mode: Some(mode) },
112108
)),
113109
..Default::default()
114110
},
115111
))])
116112
}
117113
}
118114
}
119-
(Self::Begin(config, BodyMode::Buffered), Some(http_event::Event::Begin(_))) => {
115+
(Self::Begin(config, Selected::Buffered(_)), Some(http_event::Event::Begin(_))) => {
120116
*self = Self::Buffered(config);
121117
Ok(Vec::new())
122118
}
123-
(Self::Begin(config, BodyMode::Stream), Some(http_event::Event::Begin(_))) => {
119+
(
120+
Self::Begin(config, Selected::Stream(max_chunk_bytes)),
121+
Some(http_event::Event::Begin(_)),
122+
) => {
124123
let scanner = StreamScanner::new(&config);
125-
*self = Self::Stream(config, scanner);
124+
*self = Self::Stream(config, scanner, max_chunk_bytes);
126125
Ok(vec![result(http_result::Result::OutputStart(
127126
HttpOutputStart::default(),
128127
))])
129128
}
130129
(Self::Buffered(config), Some(http_event::Event::BufferedBody(body))) => {
131130
Ok(vec![buffered_result(inspect(&config, &body.data))])
132131
}
133-
(Self::Stream(config, mut scanner), Some(http_event::Event::InputChunk(chunk))) => {
134-
match scanner.push(&chunk.data, config.mode == Mode::Deny) {
135-
Ok(output) => {
136-
*self = Self::Stream(config, scanner);
137-
Ok(output_chunk(output).into_iter().collect())
138-
}
139-
Err(_) => Ok(vec![reject(&config, &scanner)]),
132+
(
133+
Self::Stream(config, mut scanner, max_chunk_bytes),
134+
Some(http_event::Event::InputChunk(chunk)),
135+
) => match scanner.push(&chunk.data, config.mode == Mode::Deny) {
136+
Ok(output) => {
137+
*self = Self::Stream(config, scanner, max_chunk_bytes);
138+
Ok(output_chunks(&output, max_chunk_bytes))
140139
}
141-
}
142-
(Self::Stream(config, mut scanner), Some(http_event::Event::InputEnd(_))) => {
143-
match scanner.finish(config.mode == Mode::Deny) {
144-
Ok(output) => {
145-
let (match_count, matched_term_count) = scanner.counts();
146-
let diagnostics = (match_count > 0).then(|| {
147-
diagnostics(outcome(&config, match_count, matched_term_count))
148-
});
149-
let mut results: Vec<_> = output_chunk(output).into_iter().collect();
150-
results.push(result(http_result::Result::Finish(HttpFinish {
151-
trailer_mutations: Vec::new(),
152-
diagnostics,
153-
})));
154-
Ok(results)
155-
}
156-
Err(_) => Ok(vec![reject(&config, &scanner)]),
140+
Err(_) => Ok(vec![reject(&config, &scanner)]),
141+
},
142+
(
143+
Self::Stream(config, mut scanner, max_chunk_bytes),
144+
Some(http_event::Event::InputEnd(_)),
145+
) => match scanner.finish(config.mode == Mode::Deny) {
146+
Ok(output) => {
147+
let (match_count, matched_term_count) = scanner.counts();
148+
let diagnostics = (match_count > 0)
149+
.then(|| diagnostics(outcome(&config, match_count, matched_term_count)));
150+
let mut results = output_chunks(&output, max_chunk_bytes);
151+
results.push(result(http_result::Result::Finish(HttpFinish {
152+
trailer_mutations: Vec::new(),
153+
diagnostics,
154+
})));
155+
Ok(results)
157156
}
158-
}
157+
Err(_) => Ok(vec![reject(&config, &scanner)]),
158+
},
159159
(_, None) => Err(Status::invalid_argument("HTTP event is required")),
160160
// Not FAILED_PRECONDITION, which OpenShell reads as "cannot
161161
// inspect".
@@ -166,9 +166,11 @@ impl Stage {
166166
}
167167
}
168168

169+
#[derive(Clone, Copy)]
169170
enum Selected {
170171
Buffered(u64),
171-
Stream,
172+
/// The largest output chunk OpenShell accepts.
173+
Stream(usize),
172174
}
173175

174176
enum Decision {
@@ -223,7 +225,10 @@ fn preflight_decision(preflight: HttpPreflight) -> Result<Decision, Status> {
223225
.map_or(0, |limits| limits.max_buffered_body_bytes);
224226
Selected::Buffered(offered.min(MAX_PAYLOAD_BYTES))
225227
});
226-
let stream = permitted(HttpBodyMode::Stream).then_some(Selected::Stream);
228+
let stream = permitted(HttpBodyMode::Stream).then(|| {
229+
let max_chunk_bytes = preflight.limits.map_or(0, |limits| limits.max_chunk_bytes);
230+
Selected::Stream(usize::try_from(max_chunk_bytes).unwrap_or(usize::MAX))
231+
});
227232
let selected = match config.body_mode {
228233
BodyMode::Buffered => buffered.or(stream),
229234
BodyMode::Stream => stream.or(buffered),
@@ -232,6 +237,9 @@ fn preflight_decision(preflight: HttpPreflight) -> Result<Decision, Status> {
232237
Some(Selected::Buffered(0)) => Err(Status::invalid_argument(
233238
"BUFFERED requires a positive max_buffered_body_bytes",
234239
)),
240+
Some(Selected::Stream(0)) => Err(Status::invalid_argument(
241+
"STREAM requires a positive max_chunk_bytes",
242+
)),
235243
Some(selected) => Ok(Decision::Inspect(config, selected)),
236244
// A message without a body has nothing to guard. The guard never
237245
// passes a body it could not inspect.
@@ -283,8 +291,16 @@ fn buffered_result(outcome: GuardOutcome) -> HttpResult {
283291
}))
284292
}
285293

286-
fn output_chunk(data: Vec<u8>) -> Option<HttpResult> {
287-
(!data.is_empty()).then(|| result(http_result::Result::OutputChunk(HttpOutputChunk { data })))
294+
/// Output results for `data`, each within OpenShell's chunk limit: redacting
295+
/// a short term makes the output longer than the input.
296+
fn output_chunks(data: &[u8], max_chunk_bytes: usize) -> Vec<HttpResult> {
297+
data.chunks(max_chunk_bytes)
298+
.map(|data| {
299+
result(http_result::Result::OutputChunk(HttpOutputChunk {
300+
data: data.to_vec(),
301+
}))
302+
})
303+
.collect()
288304
}
289305

290306
fn reject(config: &GuardConfig, scanner: &StreamScanner) -> HttpResult {
@@ -320,6 +336,7 @@ mod tests {
320336
config: Some(config(mode, &["prototype-secret", "秘密"], extra)),
321337
permitted_body_modes: permitted.iter().map(|mode| *mode as i32).collect(),
322338
limits: Some(HttpBodyLimits {
339+
max_chunk_bytes: 64 * 1024,
323340
max_buffered_body_bytes: 1024 * 1024,
324341
..Default::default()
325342
}),
@@ -471,6 +488,51 @@ mod tests {
471488
assert_eq!(finish.diagnostics.unwrap().findings[0].count, 1);
472489
}
473490

491+
#[test]
492+
fn stream_stage_splits_expanded_output_within_the_chunk_limit() {
493+
let mut selected = preflight(
494+
response(),
495+
"redact",
496+
&[("body_mode", "stream")],
497+
&[HttpBodyMode::Stream],
498+
);
499+
let Some(http_event::Event::Preflight(preflight_event)) = selected.event.as_mut() else {
500+
unreachable!()
501+
};
502+
preflight_event.limits = Some(HttpBodyLimits {
503+
max_chunk_bytes: 16,
504+
..Default::default()
505+
});
506+
let mut stage = Stage::Preflight;
507+
stage.handle(selected).unwrap();
508+
stage
509+
.handle(event(http_event::Event::Begin(HttpBegin::default())))
510+
.unwrap();
511+
// Each 6-byte term becomes the 10-byte marker.
512+
let mut results = stage
513+
.handle(event(http_event::Event::InputChunk(HttpInputChunk {
514+
data: "秘密".repeat(8).into_bytes(),
515+
})))
516+
.unwrap();
517+
results.extend(
518+
stage
519+
.handle(event(http_event::Event::InputEnd(HttpInputEnd::default())))
520+
.unwrap(),
521+
);
522+
let mut output = Vec::new();
523+
for result in results {
524+
match result.result {
525+
Some(http_result::Result::OutputChunk(chunk)) => {
526+
assert!((1..=16).contains(&chunk.data.len()), "{chunk:?}");
527+
output.extend(chunk.data);
528+
}
529+
Some(http_result::Result::Finish(_)) => {}
530+
result => panic!("unexpected result: {result:?}"),
531+
}
532+
}
533+
assert_eq!(output, "[REDACTED]".repeat(8).as_bytes());
534+
}
535+
474536
#[test]
475537
fn preflight_continues_bodyless_and_empty_messages() {
476538
let mut bodyless = preflight(response(), "redact", &[], &[]);

0 commit comments

Comments
 (0)