1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
//! The one HTTP-client constructor every loopback/daemon caller routes through
//! (#4392).
//!
//! Why: reqwest 0.12 sends `127.0.0.1` through `HTTP_PROXY` / `http_proxy` /
//! `ALL_PROXY` when one of them is exported. hyper-util's proxy matcher has no
//! runtime loopback exemption — the only bypass is an explicit `NO_PROXY` entry
//! the operator is not required to have — so on a machine with a corporate proxy
//! configured, every tm↔daemon call is routed to the proxy and fails. The daemon
//! is up; the caller reports it down. `.no_proxy()` on the builder is the fix,
//! and it has to live at ONE entry point: the workspace held ~133 bare
//! `reqwest::Client::builder()` sites and adding the call per-site guarantees the
//! next site forgets it.
//!
//! What:
//! - [`loopback_client_builder`](crate::http_client::loopback_client_builder) —
//! the primitive. A `ClientBuilder` with proxies disabled and NO timeout
//! policy, for callers that own their own bounds (an SSE stream must not carry
//! a whole-request timeout at all).
//! - [`loopback_client`](crate::http_client::loopback_client) — that builder plus
//! the standard
//! [`LOOPBACK_CONNECT_TIMEOUT`](crate::http_client::LOOPBACK_CONNECT_TIMEOUT) /
//! [`LOOPBACK_REQUEST_TIMEOUT`](crate::http_client::LOOPBACK_REQUEST_TIMEOUT)
//! bounds, for callers with no bespoke timing requirement.
//! - `blocking_loopback_client_builder` — the `reqwest::blocking` counterpart,
//! behind the `blocking-http` feature. Deliberately NOT a link: the item does
//! not exist when that feature is off, and a link that resolves only under one
//! feature fails `#![deny(rustdoc::broken_intra_doc_links)]` for anyone who
//! documents this crate without it.
//!
//! Paths here are crate-absolute on purpose. This module carries docs in two
//! places — these `//!` lines and the `///` block on `pub mod http_client;` in
//! `lib.rs` — and rustdoc merges both into one doc string whose link-resolution
//! scope comes from the FIRST fragment, which is the `lib.rs` one, so a bare
//! `loopback_client` here resolves against the crate root and is not found
//! (#6027).
//!
//! Scope: LOOPBACK callers only. A client that genuinely talks to the public
//! internet — crates.io in `update`, an inference provider in `inference` /
//! `chat`, the GitHub API — must keep honouring the operator's proxy, and must
//! NOT be routed through here.
//!
//! Test: `tests::loopback_client_ignores_exported_http_proxy` sets `HTTP_PROXY`
//! to a dead address and proves both halves — a bare builder is diverted, this
//! one still reaches a loopback stub.
use std::time::Duration;
/// TCP connect bound for a loopback call.
///
/// Why: a closed loopback port refuses instantly, but a firewalled or wedged one
/// can hang for the OS default (~75s). Bounding connect separately from the whole
/// request keeps a dead peer cheap.
pub const LOOPBACK_CONNECT_TIMEOUT: Duration = Duration::from_secs(2);
/// Whole-request bound (connect + headers + body) for a loopback call.
///
/// Why: a live daemon answers in single-digit milliseconds; 5s is generous for a
/// loaded one and short enough that a CLI never appears to freeze.
pub const LOOPBACK_REQUEST_TIMEOUT: Duration = Duration::from_secs(5);
/// Start a `reqwest::Client` for a loopback/daemon target — proxies off, no
/// timeout policy applied.
///
/// Why: `.no_proxy()` is the load-bearing call (see the module docs); leaving the
/// timeouts to the caller is what lets every existing site adopt this entry point
/// without changing its own timing contract, including the SSE readers that must
/// have no whole-request bound.
/// What: `reqwest::Client::builder().no_proxy()`. Chain the caller's own
/// `.timeout()` / `.connect_timeout()` and `.build()`.
/// Test: `tests::loopback_client_ignores_exported_http_proxy`.
pub fn loopback_client_builder() -> reqwest::ClientBuilder {
reqwest::Client::builder().no_proxy()
}
/// Build a ready-to-use loopback client with the standard bounds.
///
/// Why: most callers want "talk to the daemon, fail fast, never via a proxy" and
/// should not restate the two timeouts.
/// What: [`loopback_client_builder`] plus [`LOOPBACK_CONNECT_TIMEOUT`] and
/// [`LOOPBACK_REQUEST_TIMEOUT`].
/// Test: `tests::loopback_client_ignores_exported_http_proxy`,
/// `tests::loopback_client_builds`.
pub fn loopback_client() -> reqwest::Result<reqwest::Client> {
loopback_client_builder()
.connect_timeout(LOOPBACK_CONNECT_TIMEOUT)
.timeout(LOOPBACK_REQUEST_TIMEOUT)
.build()
}
/// [`loopback_client_builder`] for the synchronous `reqwest::blocking` API.
///
/// Why: the index-registration and search-readiness paths run on detached
/// threads with no reactor, so they use the blocking client. They are loopback
/// callers and inherit the same defect.
/// What: `reqwest::blocking::Client::builder().no_proxy()`, with the timeout
/// policy left to the caller for the same reason as the async form.
/// Test: `tests::blocking_loopback_client_ignores_exported_http_proxy`.
#[cfg(feature = "blocking-http")]
pub fn blocking_loopback_client_builder() -> reqwest::blocking::ClientBuilder {
reqwest::blocking::Client::builder().no_proxy()
}
#[cfg(test)]
mod tests {
use super::*;
use serial_test::serial;
/// Upper bound on the request head the stub reads before it answers.
const STUB_HEAD_LIMIT: usize = 8 * 1024;
/// A loopback stub that answers one `200 OK` per connection, forever.
///
/// Returns the bound `host:port`. Modelled on
/// `daemon_guard::tests::spin_until_ready_returns_ok_for_live_server` — a
/// real listener rather than a mocked client, so the proxy behaviour under
/// test is the real transport's.
///
/// Each connection reads the request head (to `\r\n\r\n`, at most
/// [`STUB_HEAD_LIMIT`] bytes) before it writes the response, then sends
/// `Connection: close` and shuts down its write half.
/// Test: `stub_server_answers_only_after_the_request_arrives`.
async fn stub_server() -> String {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind loopback stub");
let addr = listener.local_addr().expect("stub addr").to_string();
tokio::spawn(async move {
while let Ok((mut stream, _)) = listener.accept().await {
tokio::spawn(async move {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
// #6575: answering before the request arrived let hyper read
// the 200 on an idle connection (`unexpected message`), and
// closing with the request unread sends RST, not FIN.
let mut head = Vec::with_capacity(512);
let mut chunk = [0u8; 512];
while !head.windows(4).any(|w| w == b"\r\n\r\n") {
match stream.read(&mut chunk).await {
Ok(0) | Err(_) => return,
Ok(n) => head.extend_from_slice(&chunk[..n]),
}
if head.len() >= STUB_HEAD_LIMIT {
break;
}
}
let _ = stream
.write_all(
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
)
.await;
let _ = stream.shutdown().await;
});
}
});
addr
}
/// Why (#6575): the proxy tests below flaked on the GUARDED client, which a
/// proxy cannot touch. The stub answered at `accept` time, before the
/// client had written its request, so the 200 could reach hyper on an idle
/// connection and the stub could close with the request unread.
/// What: connects with a raw socket, proves no byte arrives before a
/// request is sent, then sends a GET and reads to a clean EOF.
/// Test: This is the test.
#[tokio::test]
async fn stub_server_answers_only_after_the_request_arrives() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let addr = stub_server().await;
let mut stream = tokio::net::TcpStream::connect(&addr)
.await
.expect("connect to stub");
let mut early = [0u8; 64];
let premature =
tokio::time::timeout(Duration::from_millis(300), stream.read(&mut early)).await;
assert!(
premature.is_err(),
"the stub must not answer before the request arrives; got {premature:?}"
);
stream
.write_all(b"GET /health HTTP/1.1\r\nHost: stub\r\n\r\n")
.await
.expect("send request");
let mut response = Vec::new();
let read =
tokio::time::timeout(LOOPBACK_REQUEST_TIMEOUT, stream.read_to_end(&mut response))
.await
.expect("stub answers within the request bound");
assert!(
read.is_ok(),
"the stub must end with FIN, not RST: {read:?}"
);
let text = String::from_utf8_lossy(&response);
assert!(
text.starts_with("HTTP/1.1 200 OK\r\n") && text.ends_with("\r\n\r\nok"),
"the client must get the whole 200; got {text:?}"
);
}
/// An address with nothing listening: bind port 0, read it, release it.
fn dead_addr() -> String {
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind to free a port");
let addr = listener.local_addr().expect("dead addr").to_string();
drop(listener);
addr
}
/// Set `HTTP_PROXY`, run `body`, restore the previous value, and hold
/// `ENV_LOCK` for the whole call (#6575).
///
/// Why more than one domain: this binary serializes environment mutation in
/// three domains that do not exclude each other — `#[serial]`'s default
/// group, the named `#[serial(dotenv_credential_env)]` group
/// (`credentials`, `inference`, `memory_core`), and `data_dir::ENV_LOCK`, a
/// plain mutex that excludes only the tests that take it
/// (`daemon_addr.rs`'s own comment records that gap). `reqwest` reads
/// `HTTP_PROXY` / `http_proxy` / `ALL_PROXY` through `std::env::var_os` at
/// every client build
/// (`hyper_util::client::proxy::matcher::Builder::from_env`, no cache), so
/// a proxy test holding one domain still builds its clients while a test in
/// another domain is inside `setenv`/`unsetenv` — and
/// `credentials::dotenv`'s `load_env_from_path` republishes arbitrary keys
/// in bulk. The callers moved to the `dotenv_credential_env` key, which is
/// where every bulk env writer lives, and this takes `ENV_LOCK`. The default
/// group is what they gave up, and it costs nothing: its env writers
/// (`bm25`, `catchup`, `daemon_token`) each set one named variable of their
/// own, none of them a proxy variable.
///
/// SAFETY: the caller holds every domain above, so no other test in this
/// binary is reading or writing the environment concurrently.
fn with_http_proxy<T>(value: &str, body: impl FnOnce() -> T) -> T {
let _env = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let previous = std::env::var("HTTP_PROXY").ok();
unsafe { std::env::set_var("HTTP_PROXY", value) };
let out = body();
unsafe {
match previous {
Some(v) => std::env::set_var("HTTP_PROXY", v),
None => std::env::remove_var("HTTP_PROXY"),
}
}
out
}
/// Why (#4392): this is the whole defect. reqwest reads the proxy
/// environment at client-BUILD time and routes `127.0.0.1` through it, so a
/// developer with a corporate proxy exported sees every loopback daemon call
/// fail while the daemon is demonstrably up.
///
/// Both halves are asserted so that deleting `.no_proxy()` as "hygiene"
/// fails loudly: a bare builder must be diverted to the dead proxy, and
/// [`loopback_client`] must still reach the stub.
/// What: exports `HTTP_PROXY` pointing at a released ephemeral port, builds
/// both clients under it, and issues the same loopback GET with each.
/// Serialized because `HTTP_PROXY` is process-global — both serial keys,
/// per [`with_http_proxy`] (#6575).
/// Test: This is the test.
///
/// If a future reqwest exempts loopback from proxies, the first assertion
/// becomes obsolete and should be deleted — `.no_proxy()` itself must stay,
/// because this crate cannot pin every consumer's reqwest patch level.
#[tokio::test]
#[serial(dotenv_credential_env)]
async fn loopback_client_ignores_exported_http_proxy() {
let addr = stub_server().await;
let url = format!("http://{addr}/health");
let proxy = format!("http://{}", dead_addr());
let (leaky, guarded) = with_http_proxy(&proxy, || {
let leaky = reqwest::Client::builder()
.connect_timeout(Duration::from_millis(500))
.timeout(Duration::from_millis(500))
.build()
.expect("bare client builds");
let guarded = loopback_client().expect("loopback client builds");
(leaky, guarded)
});
let leaked = leaky.get(&url).send().await;
let reached = guarded.get(&url).send().await;
assert!(
leaked.is_err(),
"a client WITHOUT .no_proxy() must be diverted through HTTP_PROXY — \
that diversion IS the #4392 mechanism; got {leaked:?}"
);
assert!(
reached.is_ok_and(|r| r.status().is_success()),
"loopback_client() must reach a loopback peer with HTTP_PROXY exported"
);
}
/// Why (#4392): the blocking half of the entry point serves the
/// index-registration and readiness paths, which run off-reactor. It carries
/// the same defect and needs the same proof.
/// What: the async test's shape, on `reqwest::blocking`, driven from a
/// `spawn_blocking` so the blocking client never runs on a reactor thread.
/// Test: This is the test.
#[cfg(feature = "blocking-http")]
#[tokio::test]
#[serial(dotenv_credential_env)]
async fn blocking_loopback_client_ignores_exported_http_proxy() {
let addr = stub_server().await;
let url = format!("http://{addr}/health");
let proxy = format!("http://{}", dead_addr());
let (leaked, reached) = tokio::task::spawn_blocking(move || {
with_http_proxy(&proxy, || {
let leaky = reqwest::blocking::Client::builder()
.connect_timeout(Duration::from_millis(500))
.timeout(Duration::from_millis(500))
.build()
.expect("bare blocking client builds");
let guarded = blocking_loopback_client_builder()
.connect_timeout(LOOPBACK_CONNECT_TIMEOUT)
.timeout(LOOPBACK_REQUEST_TIMEOUT)
.build()
.expect("blocking loopback client builds");
(
leaky.get(&url).send().is_err(),
guarded
.get(&url)
.send()
.is_ok_and(|r| r.status().is_success()),
)
})
})
.await
.expect("blocking probe thread");
assert!(
leaked,
"a blocking client WITHOUT .no_proxy() must be diverted through HTTP_PROXY"
);
assert!(
reached,
"the blocking loopback builder must reach a loopback peer with HTTP_PROXY exported"
);
}
/// Why: the standard-bounds convenience form must construct.
/// What: builds it and drops it.
/// Test: This is the test.
#[test]
fn loopback_client_builds() {
drop(loopback_client().expect("loopback client builds"));
}
}