Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 41 additions & 1 deletion crates/busbar-kernel/src/egress/engine/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,27 @@ use hyper::body::Incoming;
use super::pool::{self, CheckedOut, ClientInner, EngineError, ErrorKind, PoolKey, PoolMap};
use super::{EngineConnector, H2KeepAlive};

/// THE CONNECTION'S RETURN, as a response extension: present on an HTTP/1.1 response whose
/// connection goes back to its pool only once hyper's dispatcher finishes the exchange (the body
/// was not drained at head time). [`ConnReturned::settled`] resolves once that return has happened,
/// or the connection died and is dropped. A caller that reads the body to its end and then reports
/// that end awaits it first, so a hop opened after that end is lent the connection back in the pool
/// instead of racing its return with a fresh dial (a second connection the second hop never rides).
#[derive(Clone, Debug)]
pub struct ConnReturned(tokio::sync::watch::Receiver<bool>);

impl ConnReturned {
/// The connection is back in its pool, or dropped (it died, or the pool is gone).
pub async fn settled(mut self) {
let _ = self.0.wait_for(|returned| *returned).await;
}
}

/// A delay before a watcher's return, per authority (`host:port`): the window between a body's end
/// and its connection's return, held open so a test meets the order a loaded host gives by chance.
#[cfg(test)]
pub(crate) static RETURN_DELAY_FOR_TESTS: Mutex<Option<(String, Duration)>> = Mutex::new(None);

/// The pooled egress client — the owned struct behind the seam every plane builds from
/// (`build_client`). Cheap to clone: clones share one pool; dropping the last clone releases the
/// pool and closes its idle sockets.
Expand Down Expand Up @@ -204,20 +225,39 @@ async fn send_request(
// drained at head time), else the per-exchange watcher — hyper's h1
// `poll_ready` resolves only when the dispatcher finishes the exchange,
// which is gated on the caller draining `Incoming`.
if sender.is_ready() {
// A test's lag on this authority's return (see `RETURN_DELAY_FOR_TESTS`)
// takes the watcher path even when the exchange already finished.
#[cfg(test)]
let lag = RETURN_DELAY_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
.filter(|(at, _)| *at == pool_key.1.as_str())
.map(|(_, d)| *d);
#[cfg(not(test))]
let lag: Option<Duration> = None;
if sender.is_ready() && lag.is_none() {
pool::return_h1_conn(&inner, &pool_key, sender, extras);
} else {
let weak = Arc::downgrade(&inner);
let key = pool_key.clone();
// The response carries the return's settling, for a caller that
// reports the body's end only once its connection is back.
let (returned, settled) = tokio::sync::watch::channel(false);
resp.extensions_mut().insert(ConnReturned(settled));
tokio::spawn(async move {
let ready = std::future::poll_fn(|cx| sender.poll_ready(cx)).await;
if let Some(lag) = lag {
tokio::time::sleep(lag).await;
}
// An Err means the conn died during the body read: drop it —
// never returned, nothing delivered, no counter touched.
if ready.is_ok() {
if let Some(inner) = weak.upgrade() {
pool::return_h1_conn(&inner, &key, sender, extras);
}
}
let _ = returned.send(true);
});
}
return Ok(resp);
Expand Down
4 changes: 3 additions & 1 deletion crates/busbar-kernel/src/egress/engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,10 +84,12 @@ pub use tls::{ClientIdentity, Trust};
pub type EngineConnector =
KeyPinObserve<ConnectDeadline<https::HttpsConnector<tunnel::TunnelConnector>>>;

#[cfg(test)]
pub(crate) use client::RETURN_DELAY_FOR_TESTS;
/// The pooled egress client — the OWNED pool with dial coalescing (`client.rs`/`pool.rs`),
/// behind the same `request()` surface the `hyper_util::client::legacy::Client` alias had.
/// `Full<Bytes>`: every engine egress body is one owned buffer.
pub use client::EngineClient;
pub use client::{ConnReturned, EngineClient};

/// The error type `EngineClient::request` yields — named so a consumer's transport-error
/// classification arms read as prose. Owned (`pool::EngineError`): carries its cause chain as
Expand Down
10 changes: 10 additions & 0 deletions crates/busbar-kernel/src/plane_host/egress.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1127,6 +1127,13 @@ fn run_http_stream(
// (its per-request `.timeout()` kept ticking through the body and failed the stream at
// the instant — so does this).
use http_body_util::BodyExt;
// Its connection's return to the pool, when that waits on the exchange's end: the body's
// end is reported only once the connection is back, so a hop opened after that end is lent
// it rather than dialling a second connection it never rides.
let mut returned = resp
.extensions()
.get::<busbar_kernel::egress::engine::ConnReturned>()
.cloned();
let mut body = resp.into_body();
loop {
tokio::select! {
Expand Down Expand Up @@ -1161,6 +1168,9 @@ fn run_http_stream(
break;
}
None => {
if let Some(returned) = returned.take() {
let _ = tokio::time::timeout_at(deadline, returned.settled()).await;
}
let _ = chunk_tx.send(ChunkMsg::End);
break;
}
Expand Down
50 changes: 50 additions & 0 deletions crates/busbar-kernel/src/plane_host/egress_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -683,6 +683,56 @@ fn a_repeat_hop_reuses_the_pooled_connection_instead_of_redialing() {
});
}

/// STEP-6 REUSE when the connection's return LAGS its body's end — the order a loaded host gives
/// [`a_repeat_hop_reuses_the_pooled_connection_instead_of_redialing`] by chance (its CI red: a
/// second connection carrying 0 requests): the first hop's connection goes back to the pool only
/// once hyper's dispatcher finishes the exchange, on a task of its own, and here that task is held
/// back 300 ms past the body's end. The body's end is reported to the plane only once the
/// connection is back, so the repeat hop is lent it every time instead of parking for a fresh dial
/// the returning connection then serves (leaving the dial's connection idle and unused).
#[test]
fn a_repeat_hop_reuses_the_connection_whose_return_lags_its_body() {
use busbar_kernel::egress::fixtures::{spawn_http, CannedResponse};
let fixture = spawn_http(CannedResponse::ok("warm"), 8);
*crate::egress::engine::RETURN_DELAY_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some((
fixture.addr.to_string(),
std::time::Duration::from_millis(300),
));
let url = format!("http://{}/hop", fixture.addr);
let desc = http_desc(url.as_bytes());
let app = crate::test_support::TestApp::new().build();
with_dispatch_scope(&app, |host, vt| {
for round in 1..=2 {
let mut out = std::mem::MaybeUninit::<EgressOpen>::uninit();
let class = host_authored_open(host, &desc, &mut out);
assert_eq!(class, StatusClass::Ok, "open {round} must succeed");
// SAFETY: Ok ⇒ the out-param is initialized.
let open = unsafe { out.assume_init() };
assert_eq!(
drain(vt, host, open.id),
b"warm",
"round {round} streams the body"
);
assert_eq!((vt.egress_close.unwrap())(host, open.id), StatusClass::Ok);
}
let records = fixture.records();
assert_eq!(
records.len(),
1,
"both hops must ride ONE connection — the repeat hop redialed: {records:?}"
);
assert_eq!(
records[0].requests, 2,
"the one connection served both requests"
);
});
*crate::egress::engine::RETURN_DELAY_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
}

/// FFI-F1 (SSRF pin bypass): a plane-supplied PINNED address gets NO trust — it is judged by the SAME
/// host-side address rule as a resolved one BEFORE connecting. A pinned cloud-metadata address is
/// refused EVEN under a fully permissive host scope (metadata is the guard, not a policy a scope can
Expand Down
37 changes: 36 additions & 1 deletion crates/plugin-loader/src/hook_door.rs
Original file line number Diff line number Diff line change
Expand Up @@ -355,7 +355,19 @@ impl Inner {
// as the ops holding them return — and fails `TimedOut` on that budget, never at
// once and never past it. A refusal with a unit free (one came back since) is tried
// again once at once; any other refusal is the plugin's answer.
if done.outcome == Outcome::Refused && !plugin.is_faulted() {
#[cfg(test)]
if done.outcome == Outcome::Refused {
hold_refused_until_faulted_for_tests(&self.settings, &plugin, deadline).await;
}
// REFUSED, AND THE INSTANCE FAULTED SINCE (the watchdog quarantined it between the
// refusal and this look): a refusal is never a crossing, so this call was never
// made; like one that met the fault before it was submitted (above), it waits for
// the trial window within its own budget (R2), never answering the refusal.
if done.outcome == Outcome::Refused && plugin.is_faulted() {
guard.answered = true;
continue 'route;
}
if done.outcome == Outcome::Refused {
if plugin.inflight() < plugin.max_inflight() {
if !retried_free {
retried_free = true;
Expand Down Expand Up @@ -411,6 +423,29 @@ impl Inner {
}
}

/// A mark in the settings of the instances whose refused calls wait, before they are judged, until the
/// instance is faulted (bounded by the call's deadline): the order a loaded host gives by chance
/// (a call refused at the cap, the watchdog faulting the wedged instance before the refusal is
/// looked at), held so a test meets it every time.
#[cfg(test)]
pub(crate) static HOLD_REFUSED_UNTIL_FAULTED_FOR_TESTS: Mutex<Option<Vec<u8>>> = Mutex::new(None);

#[cfg(test)]
async fn hold_refused_until_faulted_for_tests(
settings: &[u8],
plugin: &Plugin<Hook>,
deadline: Instant,
) {
let held = HOLD_REFUSED_UNTIL_FAULTED_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_deref()
.is_some_and(|mark| settings.windows(mark.len()).any(|w| w == mark));
while held && !plugin.is_faulted() && Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(5)).await;
}
}

/// What one submitted op answered (READY or FAILED, with its `out`).
struct Submitted<O> {
plugin: Plugin<Hook>,
Expand Down
53 changes: 53 additions & 0 deletions crates/plugin-loader/src/tests/hook_door_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -370,6 +370,59 @@ async fn a_slow_hook_is_cut_off_at_its_budget_through_the_axis() {
);
}

/// PB-81, THE QUARANTINE RACE: a call refused at the saturated cap whose wedged instance the
/// watchdog faults before the refusal is looked at (the order a loaded host gave
/// [`the_inflight_cap_saturates_and_fails_on_the_caller_deadline_through_the_axis`] by chance: its
/// freed-slot call answered `broken (hook ... answered Refused)` at 1.00 s). That call was never
/// made — a refusal is no crossing — so it waits for the trial window within its own budget, as a
/// call meeting the fault before it was submitted does, and is answered by the fresh instance.
/// The refused call here is held until the watchdog has faulted the instance, so the order is met
/// every time.
#[tokio::test]
async fn a_call_refused_at_the_cap_as_the_watchdog_faults_the_instance_waits_for_its_trial() {
let axis = rows(hook_door_plugin::conforming::door, "hook_door", Way::Linked).expect("linked");
// `sleep_ms` past the Call class budget: the watchdog faults the wedged crossings; the odd
// figure marks this instance's settings for the hold.
let wedged: Arc<dyn busbar_contract::hook_calls::HookCalls> = axis
.open(
NAME,
"hooks.wedged-race",
&json!({"reject_over_messages": 3, "sleep_ms": 1_501}),
BUDGET,
)
.expect("the wedged gate opens");
*crate::hook_door::HOLD_REFUSED_UNTIL_FAULTED_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(b"\"sleep_ms\":1501".to_vec());
let mut inflight = Vec::with_capacity(MAX_INFLIGHT as usize);
for _ in 0..MAX_INFLIGHT {
let wedged = Arc::clone(&wedged);
inflight.push(tokio::spawn(async move {
let _ = wedged.decide(frame(2), Duration::from_millis(50)).await;
}));
}
for h in inflight {
let _ = h.await;
}
// Every unit is held by a wedged crossing (the callers gave up; the crossings run on), so this
// call is refused at the cap, held until the watchdog faults the instance, then judged.
let answered = wedged.decide(frame(2), BUDGET).await;
*crate::hook_door::HOLD_REFUSED_UNTIL_FAULTED_FOR_TESTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
assert!(
matches!(
answered,
Answered::Answer {
outcome: busbar_contract::abi::mechanism::call::Outcome::Ready,
..
}
),
"a call the cap refused as the instance was quarantined waits for its trial: {}",
decided(2, &answered)
);
}

/// PB-81 (`max_inflight` per loaded hook; 1.5.5 pinned `MAX_INFLIGHT_HOOK_CALLS = 64`, a 1.6.0 hook
/// states its own in its Statement): one hook is saturated with `max_inflight` calls that never
/// return inside the test's budget. A further call must fail CLOSED on the caller's own deadline —
Expand Down
Loading