@@ -13,11 +13,65 @@ use crate::envelope::{RequestEnvelope, ResponseEnvelope, ResponseMetadata};
1313use crate :: internal:: { dispatch_and_split, dispatch_parts, to_response_envelope_text} ;
1414use crate :: registry:: resolve_app_router;
1515use crate :: wire:: {
16- WIRE_HEADER_RESERVE , WIRE_VERSION , error_wire , header_capacity_estimate , parse_wire_header ,
17- split_wire_borrowed, split_wire_request, to_wire_bytes, write_wire_header_into_slice ,
18- write_wire_header_into_vec,
16+ WIRE_HEADER_RESERVE , WIRE_VERSION , WireRequestHeader , error_wire , header_capacity_estimate ,
17+ parse_wire_header , split_wire_borrowed, split_wire_request, to_wire_bytes,
18+ write_wire_header_into_slice , write_wire_header_into_vec,
1919} ;
2020
21+ // ── Shared wire prelude (used by every wire entry point) ─────────────
22+
23+ /// Ingress-cap guard shared by the **buffered** wire entry points
24+ /// (`dispatch_from_bytes_async`, `dispatch_into_async`,
25+ /// `dispatch_into_async_borrowed`, and the response-streaming pair).
26+ /// Returns the `413` wire bytes when the request exceeds the configured
27+ /// maximum, else `None`. Centralizing the message keeps the cap identical
28+ /// across entry points; **bidirectional** streaming is intentionally exempt
29+ /// (it is `O(chunk)` RAM) and so does not call this.
30+ #[ inline]
31+ pub fn check_ingress_cap ( len : usize ) -> Option < Vec < u8 > > {
32+ if crate :: config:: request_exceeds_limit ( len) {
33+ Some ( error_wire (
34+ 413 ,
35+ & format ! (
36+ "request size {len} bytes exceeds configured maximum of {} bytes" ,
37+ crate :: config:: max_request_bytes( )
38+ ) ,
39+ ) )
40+ } else {
41+ None
42+ }
43+ }
44+
45+ /// Wire-prelude shared by **every** wire entry point (buffered,
46+ /// direct-write, and streaming): parse the header, enforce the protocol
47+ /// [`WIRE_VERSION`], and resolve the target app [`Router`]. Centralizing
48+ /// this keeps the security-sensitive version check + app resolution
49+ /// byte-identical across all dispatchers — the previous per-entry-point
50+ /// copies were a drift hazard.
51+ ///
52+ /// `header_bytes` is the wire header-JSON region; the returned
53+ /// [`WireRequestHeader`] borrows from it, so the caller MUST keep it alive
54+ /// for as long as the header is used. On failure the `Err` carries the
55+ /// exact wire error bytes to deliver in the caller's shape (`400` for a
56+ /// parse error or version mismatch, `400`/`404` from app resolution).
57+ #[ inline]
58+ pub fn parse_validate_resolve (
59+ header_bytes : & [ u8 ] ,
60+ ) -> Result < ( WireRequestHeader < ' _ > , Router ) , Vec < u8 > > {
61+ let header = parse_wire_header ( header_bytes) . map_err ( |msg| error_wire ( 400 , & msg) ) ?;
62+ if header. v != WIRE_VERSION {
63+ return Err ( error_wire (
64+ 400 ,
65+ & format ! (
66+ "unsupported wire version: got {}, expected {WIRE_VERSION}" ,
67+ header. v
68+ ) ,
69+ ) ) ;
70+ }
71+ let router = resolve_app_router ( & header) ?;
72+ Ok ( ( header, router) )
73+ }
74+
2175// ── Dispatch (direct API — backward compatible) ──────────────────────
2276
2377/// Dispatch a [`RequestEnvelope`] through an axum [`Router`] and
@@ -118,40 +172,20 @@ pub fn dispatch_from_bytes(input: Vec<u8>, runtime: &tokio::runtime::Runtime) ->
118172/// guarantees as [`dispatch_from_bytes`]), including `404` when no app
119173/// is registered under the requested name.
120174pub async fn dispatch_from_bytes_async ( input : Vec < u8 > ) -> Vec < u8 > {
121- // Ingress cap (defense-in-depth): reject an oversized buffered
122- // request with 413 before doing any further work. Unlimited by
123- // default (see `max_request_bytes`); streaming paths are exempt.
124- if crate :: config:: request_exceeds_limit ( input. len ( ) ) {
125- return error_wire (
126- 413 ,
127- & format ! (
128- "request size {} bytes exceeds configured maximum of {} bytes" ,
129- input. len( ) ,
130- crate :: config:: max_request_bytes( )
131- ) ,
132- ) ;
175+ // Ingress cap (defense-in-depth): reject an oversized buffered request
176+ // with 413 before any further work. Unlimited by default; bidirectional
177+ // streaming is exempt. See [`check_ingress_cap`].
178+ if let Some ( err) = check_ingress_cap ( input. len ( ) ) {
179+ return err;
133180 }
134- // Wire-level checks next: malformed input must report parse
135- // errors regardless of whether an app is registered .
181+ // Malformed input must report parse errors regardless of whether an app
182+ // is registered, so split first, then the shared parse/version/resolve .
136183 let ( header_bytes, body_bytes) = match split_wire_request ( input) {
137184 Ok ( parts) => parts,
138185 Err ( msg) => return error_wire ( 400 , & msg) ,
139186 } ;
140- let header = match parse_wire_header ( & header_bytes) {
141- Ok ( h) => h,
142- Err ( msg) => return error_wire ( 400 , & msg) ,
143- } ;
144- if header. v != WIRE_VERSION {
145- return error_wire (
146- 400 ,
147- & format ! (
148- "unsupported wire version: got {}, expected {WIRE_VERSION}" ,
149- header. v
150- ) ,
151- ) ;
152- }
153- let router = match resolve_app_router ( & header) {
154- Ok ( r) => r,
187+ let ( header, router) = match parse_validate_resolve ( & header_bytes) {
188+ Ok ( parts) => parts,
155189 Err ( wire) => return wire,
156190 } ;
157191
@@ -296,41 +330,15 @@ pub fn dispatch_into(
296330pub async fn dispatch_into_async ( input : Vec < u8 > , out : & mut [ u8 ] ) -> DirectWriteResult {
297331 // Ingress cap (defense-in-depth) — same policy as
298332 // `dispatch_from_bytes_async`; 413 written into the caller buffer.
299- if crate :: config:: request_exceeds_limit ( input. len ( ) ) {
300- return write_wire_into (
301- out,
302- & error_wire (
303- 413 ,
304- & format ! (
305- "request size {} bytes exceeds configured maximum of {} bytes" ,
306- input. len( ) ,
307- crate :: config:: max_request_bytes( )
308- ) ,
309- ) ,
310- ) ;
333+ if let Some ( err) = check_ingress_cap ( input. len ( ) ) {
334+ return write_wire_into ( out, & err) ;
311335 }
312336 let ( header_bytes, body_bytes) = match split_wire_request ( input) {
313337 Ok ( parts) => parts,
314338 Err ( msg) => return write_wire_into ( out, & error_wire ( 400 , & msg) ) ,
315339 } ;
316- let header = match parse_wire_header ( & header_bytes) {
317- Ok ( h) => h,
318- Err ( msg) => return write_wire_into ( out, & error_wire ( 400 , & msg) ) ,
319- } ;
320- if header. v != WIRE_VERSION {
321- return write_wire_into (
322- out,
323- & error_wire (
324- 400 ,
325- & format ! (
326- "unsupported wire version: got {}, expected {WIRE_VERSION}" ,
327- header. v
328- ) ,
329- ) ,
330- ) ;
331- }
332- let router = match resolve_app_router ( & header) {
333- Ok ( r) => r,
340+ let ( header, router) = match parse_validate_resolve ( & header_bytes) {
341+ Ok ( parts) => parts,
334342 Err ( wire) => return write_wire_into ( out, & wire) ,
335343 } ;
336344
@@ -381,41 +389,15 @@ pub async fn dispatch_into_async(input: Vec<u8>, out: &mut [u8]) -> DirectWriteR
381389/// the same error / `422` / overflow semantics apply.
382390pub async fn dispatch_into_async_borrowed ( input : & [ u8 ] , out : & mut [ u8 ] ) -> DirectWriteResult {
383391 // Ingress cap (defense-in-depth) — same policy as `dispatch_into_async`.
384- if crate :: config:: request_exceeds_limit ( input. len ( ) ) {
385- return write_wire_into (
386- out,
387- & error_wire (
388- 413 ,
389- & format ! (
390- "request size {} bytes exceeds configured maximum of {} bytes" ,
391- input. len( ) ,
392- crate :: config:: max_request_bytes( )
393- ) ,
394- ) ,
395- ) ;
392+ if let Some ( err) = check_ingress_cap ( input. len ( ) ) {
393+ return write_wire_into ( out, & err) ;
396394 }
397395 let ( header_bytes, body_bytes) = match split_wire_borrowed ( input) {
398396 Ok ( parts) => parts,
399397 Err ( msg) => return write_wire_into ( out, & error_wire ( 400 , & msg) ) ,
400398 } ;
401- let header = match parse_wire_header ( header_bytes) {
402- Ok ( h) => h,
403- Err ( msg) => return write_wire_into ( out, & error_wire ( 400 , & msg) ) ,
404- } ;
405- if header. v != WIRE_VERSION {
406- return write_wire_into (
407- out,
408- & error_wire (
409- 400 ,
410- & format ! (
411- "unsupported wire version: got {}, expected {WIRE_VERSION}" ,
412- header. v
413- ) ,
414- ) ,
415- ) ;
416- }
417- let router = match resolve_app_router ( & header) {
418- Ok ( r) => r,
399+ let ( header, router) = match parse_validate_resolve ( header_bytes) {
400+ Ok ( parts) => parts,
419401 Err ( wire) => return write_wire_into ( out, & wire) ,
420402 } ;
421403
0 commit comments