agentplane/netguard/intake.rs
1//! How much an answer from somewhere else may cost this process.
2//!
3//! The address rule next door answers *where may we connect*, and
4//! [`Egress`](crate::core::Egress) answers *which hosts may we reach*. Neither
5//! answers the question a counterparty decides on its own: **how many bytes come
6//! back**. Every outbound call here carries a timeout, so the time an answer may
7//! take is bounded; a timeout says nothing about how much may arrive inside it,
8//! and a fast endpoint delivers a gigabyte long before one fires. The model, the
9//! tools and the peers are untrusted here, so a tool server answering a
10//! directory listing with a gigabyte and a hostile one are the same case.
11//!
12//! The ceiling is applied twice per call. Once against a declared
13//! `Content-Length`, which refuses the honest large answer before a byte is
14//! read; then against the accumulated bytes, because the header is a claim by
15//! the party under suspicion and a check that stopped there would be a control
16//! the attacker configures.
17//!
18//! [`ANSWER`] bounds a reply that carries work — a model completion, a peer's
19//! answer. [`METADATA`] bounds a description of something: an Agent Card, a
20//! checkpoint note, a wrapped key, a body read only to say why a call failed.
21//! Both are constants rather than knobs, because a ceiling nobody can raise is a
22//! ceiling nobody quietly raises to whatever the last failure needed. Governed
23//! media keeps its own configurable limit, since there the payload size is the
24//! subject rather than the overhead.
25//!
26//! # Two responses this crate never holds
27//!
28//! Named rather than left to be inferred from a missing call, for the reason
29//! [`Egress`](crate::core::Egress) names its own exception: an uncovered surface
30//! that says it is uncovered gets revisited, and one that says nothing does not.
31//!
32//! * **MCP over stdio or streamable HTTP**, whose framing belongs to `rmcp` —
33//! this crate hands that transport a process or a URL and never sees the bytes.
34//! * **Bedrock**, for the same reason it takes no `Egress`: the AWS SDK owns the
35//! response.
36//!
37//! # Why this is public
38//!
39//! The shipped drivers are not the only drivers, and the version of this control
40//! that gets written by hand is the unbounded one. The cost is stated rather
41//! than hidden: these signatures name `reqwest` types, so this crate's HTTP
42//! client is part of its public surface on the features that link one.
43
44use futures_util::StreamExt as _;
45
46/// The most an answer carrying work may cost: 16 MiB.
47///
48/// Roughly twenty times the largest answer any shipped driver can legitimately
49/// produce — two hundred thousand tokens of text is under a megabyte, plus the
50/// reasoning blobs and tool arguments beside it. Wide enough that no real call
51/// meets it, narrow enough that meeting it is a fault.
52///
53/// A peer's reply is held to the same figure: an A2A artifact past it is a file,
54/// and a file belongs in a blob store addressed by digest rather than inlined
55/// into a response about to become a journal record.
56pub const ANSWER: usize = 16 * 1024 * 1024;
57
58/// The most a description of something may cost: 1 MiB.
59///
60/// An Agent Card, a `tlog-checkpoint` note, a wrapped data key, or the body of
61/// a failed response read only to say why it failed. Every one of these is
62/// small by construction, and the largest of them — a card advertising a few
63/// hundred skills — is three orders of magnitude under this.
64pub const METADATA: usize = 1024 * 1024;
65
66/// Why an answer was not read.
67///
68/// The three arms are kept apart because they call for different next steps and
69/// because two of them are the counterparty's fault while the third may be the
70/// network's. Collapsing them would file *this peer is misbehaving* under *try
71/// again later*, which is the retry loop that never ends.
72#[derive(Debug)]
73pub enum IntakeError {
74 /// The response declared a body past the ceiling. Refused before reading.
75 Declared { limit: usize, declared: u64 },
76 /// The body grew past the ceiling while it was being read.
77 Exceeded { limit: usize },
78 /// The stream failed. Not a size refusal — the counterparty may be blameless.
79 Transport(reqwest::Error),
80}
81
82impl std::fmt::Display for IntakeError {
83 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
84 match self {
85 Self::Declared { limit, declared } => write!(
86 f,
87 "the answer declared {declared} bytes and this plane reads at most {limit}"
88 ),
89 Self::Exceeded { limit } => write!(
90 f,
91 "the answer grew past {limit} bytes, which is the most this plane reads"
92 ),
93 Self::Transport(e) => write!(
94 f,
95 "the answer could not be read: {}",
96 super::transport_text(e)
97 ),
98 }
99 }
100}
101
102impl std::error::Error for IntakeError {
103 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
104 match self {
105 Self::Transport(e) => Some(e),
106 _ => None,
107 }
108 }
109}
110
111impl IntakeError {
112 /// Whether this was a size refusal rather than a failure to read.
113 ///
114 /// The distinction a caller needs in order to classify: a size refusal is
115 /// permanent and repeating the call reaches the same wall, while a transport
116 /// failure is the ordinary interrupted-call case every driver already has a
117 /// ladder for.
118 #[must_use]
119 pub const fn is_refusal(&self) -> bool {
120 matches!(self, Self::Declared { .. } | Self::Exceeded { .. })
121 }
122
123 /// The transport failure inside, when that is what this is.
124 #[must_use]
125 pub const fn transport(&self) -> Option<&reqwest::Error> {
126 match self {
127 Self::Transport(e) => Some(e),
128 _ => None,
129 }
130 }
131}
132
133/// A running byte budget for one answer.
134///
135/// Separate from [`read`] because a streamed answer is consumed event by event
136/// by a decoder that has to see each chunk — there is no whole body to hand
137/// back. The same ceiling applies to the same bytes; only the consumer differs.
138///
139/// What it bounds on that path is the **number** of events, not their size: a
140/// per-event ceiling stops the unterminated line and lets well-formed hundred-
141/// byte deltas accumulate until the process dies.
142#[derive(Debug)]
143pub struct Meter {
144 limit: usize,
145 seen: usize,
146}
147
148impl Meter {
149 #[must_use]
150 pub const fn new(limit: usize) -> Self {
151 Self { limit, seen: 0 }
152 }
153
154 /// Charge a chunk against the budget.
155 ///
156 /// # Errors
157 ///
158 /// [`IntakeError::Exceeded`] once the answer has cost more than the budget.
159 /// Checked *before* the caller keeps the chunk, so the ceiling bounds what
160 /// is held rather than what has already been held.
161 pub fn charge(&mut self, bytes: usize) -> Result<(), IntakeError> {
162 // Saturating rather than wrapping: an answer big enough to overflow a
163 // `usize` is one this refuses either way, and a wrap would refuse
164 // nothing.
165 self.seen = self.seen.saturating_add(bytes);
166 if self.seen > self.limit {
167 return Err(IntakeError::Exceeded { limit: self.limit });
168 }
169 Ok(())
170 }
171}
172
173/// Read a whole response body, refusing one that will not fit.
174///
175/// # Errors
176///
177/// [`IntakeError`], which distinguishes a size refusal from a stream that
178/// failed.
179pub async fn read(response: reqwest::Response, limit: usize) -> Result<Vec<u8>, IntakeError> {
180 // The header, not `Response::content_length()`. What this check is about is
181 // the *claim* the counterparty made, and the accessor answers from the body
182 // when there is no claim at all — which would make the cheap refusal depend
183 // on how the body happened to be framed. This crate links no decompression
184 // feature, so the two agree on every real response anyway.
185 let declared = response
186 .headers()
187 .get(reqwest::header::CONTENT_LENGTH)
188 .and_then(|value| value.to_str().ok())
189 .and_then(|value| value.parse::<u64>().ok());
190 if let Some(declared) = declared
191 && declared > limit as u64
192 {
193 return Err(IntakeError::Declared { limit, declared });
194 }
195
196 let mut meter = Meter::new(limit);
197 // With capacity from the declared length when there is one, capped at the
198 // ceiling: a body that lies about being small still cannot make this
199 // allocate more than the ceiling up front.
200 let mut body = Vec::with_capacity(
201 declared
202 .and_then(|n| usize::try_from(n).ok())
203 .unwrap_or(0)
204 .min(limit),
205 );
206 let mut stream = response.bytes_stream();
207 while let Some(chunk) = stream.next().await {
208 let chunk = chunk.map_err(IntakeError::Transport)?;
209 meter.charge(chunk.len())?;
210 body.extend_from_slice(&chunk);
211 }
212 Ok(body)
213}
214
215/// Read a whole response body as text, refusing one that will not fit.
216///
217/// Lossy on invalid UTF-8, matching `Response::text`. Every caller of this is
218/// reading a body to *explain* something — a provider's error message, a
219/// service's refusal — and a decode failure there would replace the explanation
220/// with a second error about the explanation.
221///
222/// # Errors
223///
224/// [`IntakeError`], as [`read`].
225pub async fn read_text(response: reqwest::Response, limit: usize) -> Result<String, IntakeError> {
226 let bytes = read(response, limit).await?;
227 Ok(String::from_utf8_lossy(&bytes).into_owned())
228}
229
230#[cfg(test)]
231mod tests {
232 use super::*;
233
234 #[test]
235 fn a_budget_refuses_the_chunk_that_crosses_it() {
236 let mut meter = Meter::new(10);
237 assert!(meter.charge(6).is_ok());
238 // The chunk that crosses is refused rather than accepted and reported
239 // afterwards: the ceiling bounds what is held.
240 let err = meter.charge(6).expect_err("11 bytes is past 10");
241 assert!(err.is_refusal());
242 assert!(matches!(err, IntakeError::Exceeded { limit: 10 }));
243 }
244
245 #[test]
246 fn a_budget_admits_an_answer_of_exactly_the_ceiling() {
247 // The boundary is inclusive on the permitted side. An answer that fits
248 // exactly is an answer that fits, and an off-by-one here would refuse
249 // the honest call that sized itself to the documented limit.
250 let mut meter = Meter::new(10);
251 assert!(meter.charge(10).is_ok());
252 assert!(meter.charge(1).is_err());
253 }
254
255 /// Built from an `http::Response`, so both halves of the rule can be probed
256 /// without a socket. The streamed half needs a real server and has one in
257 /// `tests/wire/drivers.rs`; this is the half a stub can state exactly.
258 fn response(body: &'static str, declared: Option<u64>) -> reqwest::Response {
259 let mut builder = http::Response::builder();
260 if let Some(n) = declared {
261 builder = builder.header(http::header::CONTENT_LENGTH, n);
262 }
263 reqwest::Response::from(builder.body(body).expect("a response"))
264 }
265
266 #[tokio::test]
267 async fn a_declared_oversize_is_refused_before_a_byte_is_read() {
268 // The cheap half: an honest counterparty says how big the answer is, and
269 // there is no reason to spend a single allocation finding out.
270 let err = read(response("hello", Some(9_000)), 10)
271 .await
272 .expect_err("9000 declared against a ceiling of 10");
273 assert!(matches!(
274 err,
275 IntakeError::Declared {
276 limit: 10,
277 declared: 9_000
278 }
279 ));
280 assert!(err.is_refusal());
281 }
282
283 #[tokio::test]
284 async fn a_body_that_understates_itself_is_still_refused() {
285 // The half that matters. `Content-Length` is a claim by the party under
286 // suspicion, so a check that stopped at the header would be a control the
287 // attacker configures: declare five, send fifty.
288 let err = read(response("hello world", Some(2)), 4)
289 .await
290 .expect_err("eleven bytes past a ceiling of four");
291 assert!(
292 matches!(err, IntakeError::Exceeded { limit: 4 }),
293 "the header said it would fit; the bytes are what decides"
294 );
295 }
296
297 #[tokio::test]
298 async fn an_answer_within_the_ceiling_is_returned_whole() {
299 // The positive half, which a refuse-everything implementation would
300 // otherwise satisfy.
301 let body = read(response("hello", Some(5)), 16).await.expect("it fits");
302 assert_eq!(body, b"hello");
303 }
304
305 #[test]
306 fn an_answer_big_enough_to_overflow_is_still_refused() {
307 // Wrapping addition would carry `seen` back under the limit and admit
308 // it, which is the arithmetic that turns a ceiling into nothing.
309 let mut meter = Meter::new(10);
310 assert!(meter.charge(usize::MAX).is_err());
311 }
312}