Skip to main content

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}