toolkit/runtime/oop_serve.rs
1//! Out-of-process (`OoP`) HTTP server: probes, two-plane auth, and graceful
2//! drain (`cpt-cf-component-oop-bootstrap`).
3//!
4//! This module owns the edge of an `OoP` gear's HTTP surface:
5//!
6//! - **Framework probes** — `/healthz` (liveness), `/readyz` (readiness, gated
7//! on dependency resolution + custom checks), `/health` (diagnostics, full
8//! healthcheck report), and `/.well-known/openapi.json` (the canonical
9//! discovery path; `cpt-cf-binding-constraint-openapi-well-known`).
10//! - **Two-plane auth** — platform-plane ([`internal_auth_middleware`]) runs
11//! *before* tenant-plane ([`security_context_middleware`]) per
12//! `cpt-cf-adr-two-plane-auth`,
13//! but only when the corresponding authenticator is injected. Because both
14//! authenticator traits use return-position `impl Trait` (not `dyn`-safe),
15//! callers inject the object-safe [`DynBearerAuthenticator`] /
16//! [`DynInternalAuthenticator`] adapters.
17//! - **Graceful drain** — a drain guard rejects new gear-route requests with
18//! `503 + Retry-After` once draining begins, while in-flight requests finish.
19//! - **Framework middleware** — a canonical-error layer normalizes gear/auth
20//! error responses into `problem+json` (`trace_id` / `instance` filled),
21//! matching the in-process (`api-gateway`) path. The drain guard also catches
22//! panics so a panic can neither drop the client connection nor leak the
23//! in-flight counter (which would otherwise stall graceful drain).
24//!
25//! The server is driven by [`OopHttpServer`], invoked from the `HostRuntime`
26//! `OoP` path: it binds and serves the probes **before** the gear's `start()`
27//! phase (so `/healthz` is up during a slow startup), then swaps in the composed
28//! gear routes. Self-registration and dependency resolution live in
29//! [`super::oop_registration`].
30
31use std::collections::BTreeMap;
32use std::panic::AssertUnwindSafe;
33use std::sync::Arc;
34use std::sync::atomic::{AtomicUsize, Ordering};
35use std::time::Duration;
36
37use axum::{
38 Json, Router,
39 extract::{Request, State},
40 http::{StatusCode, header::RETRY_AFTER},
41 middleware::{Next, from_fn, from_fn_with_state},
42 response::{IntoResponse, Response},
43 routing::get,
44};
45use futures_util::FutureExt as _;
46use tokio_util::sync::CancellationToken;
47use tower::ServiceExt as _;
48use url::Url;
49
50use cf_system_sdks::directory::{DirectoryClient, RegisterInstanceInfo, ServiceEndpoint};
51use toolkit_canonical_errors::CanonicalError;
52use toolkit_http_middleware::{internal_auth_middleware, security_context_middleware};
53use toolkit_security::{DynBearerAuthenticator, DynInternalAuthenticator};
54
55use super::readiness::ReadinessState;
56use crate::api::canonical_error_middleware;
57
58/// `Retry-After` (seconds) advertised while the gear is draining.
59const DRAIN_RETRY_AFTER_SECONDS: u64 = 5;
60
61/// `Retry-After` (seconds) advertised for gear routes before the gear has
62/// finished starting (probes are already live; routes come up shortly).
63const STARTING_RETRY_AFTER_SECONDS: u64 = 1;
64
65// ---------------------------------------------------------------------------
66// Serve options
67// ---------------------------------------------------------------------------
68
69/// Configuration and collaborators for the `OoP` HTTP runtime, assembled by the
70/// bootstrap layer and consumed by `HostRuntime`'s `OoP` serving path.
71///
72/// `#[non_exhaustive]` so new fields (like `labels`) stay additive. Construct via
73/// [`OopServeOptions::new`] + the `with_*` builders rather than a struct literal,
74/// exactly as [`RunOptions`](super::RunOptions) is built: the five arguments to
75/// `new` are the fields with no meaningful default, and everything else defaults
76/// to the same values the bootstrap config layer uses. (A `Default` impl is still
77/// not provided — the `directory` collaborator is a required dependency.)
78///
79/// The builders are what keep [`run_oop_serving`](super::run_oop_serving)
80/// callable: it is `pub` and takes this type by value, so an out-of-crate caller
81/// — an integration test driving the real serve path, for instance — needs some
82/// way to build one.
83#[non_exhaustive]
84pub struct OopServeOptions {
85 /// Logical gear name (used for registration + `OpenAPI` title).
86 pub gear_name: String,
87 /// Process instance id (used for registration/deregistration).
88 pub instance_id: String,
89 /// Optional gear version (used for registration + `OpenAPI` version).
90 pub version: Option<String>,
91 /// Base URL other services use to reach this instance's REST endpoint
92 /// (e.g. `http://billing.default.svc.cluster.local:8080`). Registered with
93 /// `DirectoryService` as the instance's `rest_endpoint`.
94 pub advertise_uri: String,
95 /// Address the main HTTP server binds to (gear routes + probes).
96 pub listen_addr: std::net::SocketAddr,
97 /// Optional separate address for probe endpoints (sidecar port). When set,
98 /// probes are served here *and* on the main listener.
99 pub probe_bind_addr: Option<std::net::SocketAddr>,
100 /// Maximum time to wait for in-flight requests to drain on shutdown.
101 pub drain_timeout: Duration,
102 /// Interval between `DirectoryService` heartbeats. The single presence task
103 /// ([`presence_loop`](super::oop_registration)) sends a heartbeat every
104 /// interval to keep the instance `Healthy`/routable and avoid
105 /// heartbeat-timeout eviction. The presence loop clamps this to a 1s minimum.
106 pub heartbeat_interval: Duration,
107 /// Per-check timeout for readiness healthchecks (`/readyz`). A gear's
108 /// [`Healthcheck`](crate::healthcheck::Healthcheck) that exceeds this is
109 /// reported `Unhealthy`. Mirrors the `api-gateway` `healthcheck_timeout_ms`.
110 pub healthcheck_timeout: Duration,
111 /// Directory client used for self-registration and dependency resolution.
112 pub directory: Arc<dyn DirectoryClient>,
113 /// Tenant-plane authenticator; when `Some`, `security_context_middleware`
114 /// is installed on gear routes.
115 pub bearer_authenticator: Option<DynBearerAuthenticator>,
116 /// Platform-plane authenticator; when `Some`, `internal_auth_middleware` is
117 /// installed on gear routes (runs before the tenant plane).
118 pub internal_authenticator: Option<DynInternalAuthenticator>,
119 /// Stable addressing labels (k8s `matchLabels` style) advertised with this
120 /// instance's directory registration for label-based selection.
121 pub labels: BTreeMap<String, String>,
122}
123
124impl OopServeOptions {
125 /// Create `OopServeOptions` with the fields that have no meaningful default;
126 /// every optional field defaults to the bootstrap config layer's own default
127 /// (`drain_timeout` 30s, `heartbeat_interval` 5s, `healthcheck_timeout`
128 /// 500ms, no probe sidecar port, no authenticators, no labels). Layer the
129 /// rest on with the `with_*` builders.
130 #[must_use]
131 pub fn new(
132 gear_name: String,
133 instance_id: String,
134 advertise_uri: String,
135 listen_addr: std::net::SocketAddr,
136 directory: Arc<dyn DirectoryClient>,
137 ) -> Self {
138 Self {
139 gear_name,
140 instance_id,
141 version: None,
142 advertise_uri,
143 listen_addr,
144 probe_bind_addr: None,
145 drain_timeout: Duration::from_secs(30),
146 heartbeat_interval: Duration::from_secs(5),
147 healthcheck_timeout: Duration::from_millis(500),
148 directory,
149 bearer_authenticator: None,
150 internal_authenticator: None,
151 labels: BTreeMap::new(),
152 }
153 }
154
155 /// Set the gear version (used for registration + `OpenAPI` version).
156 #[must_use]
157 pub fn with_version(mut self, version: Option<String>) -> Self {
158 self.version = version;
159 self
160 }
161
162 /// Serve probes on a separate sidecar port *in addition to* the main listener.
163 #[must_use]
164 pub fn with_probe_bind_addr(mut self, addr: Option<std::net::SocketAddr>) -> Self {
165 self.probe_bind_addr = addr;
166 self
167 }
168
169 /// Override how long in-flight requests may drain on shutdown.
170 #[must_use]
171 pub fn with_drain_timeout(mut self, timeout: Duration) -> Self {
172 self.drain_timeout = timeout;
173 self
174 }
175
176 /// Override the `DirectoryService` heartbeat interval.
177 #[must_use]
178 pub fn with_heartbeat_interval(mut self, interval: Duration) -> Self {
179 self.heartbeat_interval = interval;
180 self
181 }
182
183 /// Override the per-check readiness healthcheck timeout.
184 #[must_use]
185 pub fn with_healthcheck_timeout(mut self, timeout: Duration) -> Self {
186 self.healthcheck_timeout = timeout;
187 self
188 }
189
190 /// Install the tenant-plane authenticator (`security_context_middleware`).
191 #[must_use]
192 pub fn with_bearer_authenticator(mut self, auth: Option<DynBearerAuthenticator>) -> Self {
193 self.bearer_authenticator = auth;
194 self
195 }
196
197 /// Install the platform-plane authenticator (`internal_auth_middleware`).
198 #[must_use]
199 pub fn with_internal_authenticator(mut self, auth: Option<DynInternalAuthenticator>) -> Self {
200 self.internal_authenticator = auth;
201 self
202 }
203
204 /// Advertise stable addressing labels for label-based selection.
205 #[must_use]
206 pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
207 self.labels = labels;
208 self
209 }
210}
211
212impl std::fmt::Debug for OopServeOptions {
213 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
214 f.debug_struct("OopServeOptions")
215 .field("gear_name", &self.gear_name)
216 .field("instance_id", &self.instance_id)
217 .field("version", &self.version)
218 .field("advertise_uri", &self.advertise_uri)
219 .field("listen_addr", &self.listen_addr)
220 .field("probe_bind_addr", &self.probe_bind_addr)
221 .field("drain_timeout", &self.drain_timeout)
222 .field("heartbeat_interval", &self.heartbeat_interval)
223 .field("healthcheck_timeout", &self.healthcheck_timeout)
224 .field("bearer_authenticator", &self.bearer_authenticator.is_some())
225 .field(
226 "internal_authenticator",
227 &self.internal_authenticator.is_some(),
228 )
229 .field("labels", &self.labels)
230 .finish_non_exhaustive()
231 }
232}
233
234// ---------------------------------------------------------------------------
235// Probe router
236// ---------------------------------------------------------------------------
237
238/// Composed gear routes + `OpenAPI` document, published *after* the gear's
239/// `start()` phase completes.
240///
241/// Held behind an [`arc_swap::ArcSwapOption`] so the probe listener can bind and
242/// serve `/healthz` immediately (while a slow `start()` runs), then have the gear
243/// routes swapped in atomically — no port rebind, so liveness never blips.
244#[derive(Clone, Default)]
245struct LateRoutes {
246 inner: Arc<arc_swap::ArcSwapOption<LateRoutesInner>>,
247}
248
249struct LateRoutesInner {
250 /// Fully-layered gear router (auth planes + drain guard + canonical errors).
251 gear: Router,
252 /// Serialized `OpenAPI` document served at the well-known path.
253 openapi: Arc<str>,
254}
255
256impl LateRoutes {
257 /// Publish the composed gear routes; subsequent requests are served by them.
258 fn publish(&self, gear: Router, openapi: Arc<str>) {
259 self.inner
260 .store(Some(Arc::new(LateRoutesInner { gear, openapi })));
261 }
262
263 /// The current gear router, if the gear has finished starting.
264 fn gear(&self) -> Option<Router> {
265 self.inner.load_full().map(|i| i.gear.clone())
266 }
267
268 /// The current `OpenAPI` document, if published.
269 fn openapi(&self) -> Option<Arc<str>> {
270 self.inner.load_full().map(|i| Arc::clone(&i.openapi))
271 }
272}
273
274/// `503 + Retry-After` returned for gear routes before the gear has started.
275fn starting_response() -> Response {
276 (
277 StatusCode::SERVICE_UNAVAILABLE,
278 [(RETRY_AFTER, STARTING_RETRY_AFTER_SECONDS.to_string())],
279 "starting",
280 )
281 .into_response()
282}
283
284/// Shared state for the probe endpoints.
285#[derive(Clone)]
286struct ProbeState {
287 readiness: Arc<ReadinessState>,
288 late: LateRoutes,
289}
290
291/// Build the outer router bound at startup: framework probes plus a gear-route
292/// fallback that is `503 starting` until [`LateRoutes::publish`] swaps in the
293/// composed gear router. Probes are live as soon as the listener binds.
294fn build_outer_router(readiness: Arc<ReadinessState>, late: LateRoutes) -> Router {
295 let fallback_late = late.clone();
296 build_probe_router(readiness, late).fallback_service(tower::util::service_fn(
297 move |req: Request| {
298 let late = fallback_late.clone();
299 async move {
300 match late.gear() {
301 // `Router`'s service error is `Infallible`, hence `match e {}`.
302 Some(router) => Ok::<_, std::convert::Infallible>(
303 router.oneshot(req).await.unwrap_or_else(|e| match e {}),
304 ),
305 None => Ok::<_, std::convert::Infallible>(starting_response()),
306 }
307 }
308 },
309 ))
310}
311
312/// Build the framework probe router: `/healthz`, `/readyz`, `/health`, and the
313/// well-known `OpenAPI` discovery path.
314///
315/// These routes carry no tenant JWT and are never subject to the drain guard or
316/// the auth middlewares (they must respond during startup and drain).
317fn build_probe_router(readiness: Arc<ReadinessState>, late: LateRoutes) -> Router {
318 let state = ProbeState { readiness, late };
319
320 Router::new()
321 .route("/healthz", get(healthz))
322 .route("/readyz", get(readyz))
323 .route("/health", get(health))
324 .route("/.well-known/openapi.json", get(openapi))
325 .route("/openapi.json", get(openapi))
326 .with_state(state)
327}
328
329/// Liveness: always `200` with plain body `ok` once the server is listening.
330async fn healthz() -> &'static str {
331 "ok"
332}
333
334/// Readiness: `200` when ready, `503` with the unresolved-deps / failing-checks
335/// body otherwise. `Degraded` checks keep the response `200`.
336async fn readyz(State(state): State<ProbeState>) -> Response {
337 let report = state.readiness.evaluate().await;
338 let status = if report.ready {
339 StatusCode::OK
340 } else {
341 StatusCode::SERVICE_UNAVAILABLE
342 };
343 (status, Json(report)).into_response()
344}
345
346/// Health: returns the full `HealthcheckReport` (`status` + per-component
347/// detail). `200` for `Healthy`/`Degraded`, `503` for `Unhealthy`.
348async fn health(State(state): State<ProbeState>) -> impl IntoResponse {
349 let report = state.readiness.health_report().await;
350 let status = if report.is_ready() {
351 StatusCode::OK
352 } else {
353 StatusCode::SERVICE_UNAVAILABLE
354 };
355 (status, Json(report))
356}
357
358/// Serve the gear's generated `OpenAPI` document once published (after
359/// `start()`); `503 starting` beforehand.
360async fn openapi(State(state): State<ProbeState>) -> Response {
361 match state.late.openapi() {
362 Some(spec) => (
363 StatusCode::OK,
364 [(axum::http::header::CONTENT_TYPE, "application/json")],
365 spec.to_string(),
366 )
367 .into_response(),
368 None => starting_response(),
369 }
370}
371
372// ---------------------------------------------------------------------------
373// Drain guard
374// ---------------------------------------------------------------------------
375
376/// Tracks in-flight gear-route requests and enforces the drain state.
377#[derive(Clone)]
378struct DrainGuard {
379 readiness: Arc<ReadinessState>,
380 in_flight: Arc<AtomicUsize>,
381}
382
383impl DrainGuard {
384 fn new(readiness: Arc<ReadinessState>) -> Self {
385 Self {
386 readiness,
387 in_flight: Arc::new(AtomicUsize::new(0)),
388 }
389 }
390
391 /// Current number of in-flight gear-route requests.
392 fn in_flight(&self) -> usize {
393 self.in_flight.load(Ordering::SeqCst)
394 }
395
396 /// Begin draining: flip readiness to `503` so upstreams stop routing new
397 /// traffic and the guard starts rejecting new requests (drain step 1/2).
398 fn begin_drain(&self) {
399 self.readiness.set_draining(true);
400 }
401}
402
403/// RAII guard that increments the in-flight request counter on creation and
404/// decrements it on drop, so cancellation or panic cannot leak a count.
405struct InFlightGuard(Arc<AtomicUsize>);
406
407impl InFlightGuard {
408 fn new(counter: Arc<AtomicUsize>) -> Self {
409 counter.fetch_add(1, Ordering::SeqCst);
410 Self(counter)
411 }
412}
413
414impl Drop for InFlightGuard {
415 fn drop(&mut self) {
416 self.0.fetch_sub(1, Ordering::SeqCst);
417 }
418}
419
420/// Middleware: reject new requests with `503 + Retry-After` while draining;
421/// otherwise track the request as in-flight until it completes.
422async fn drain_guard_middleware(
423 State(guard): State<DrainGuard>,
424 request: axum::extract::Request,
425 next: Next,
426) -> Response {
427 if guard.readiness.is_draining() {
428 return (
429 StatusCode::SERVICE_UNAVAILABLE,
430 [(RETRY_AFTER, DRAIN_RETRY_AFTER_SECONDS.to_string())],
431 "draining",
432 )
433 .into_response();
434 }
435
436 // The in-flight counter is tracked by an RAII guard so it is decremented
437 // whether the handler completes normally, panics, or the future is cancelled.
438 // Panics are caught and converted to a canonical 500 (enriched by the outer
439 // canonical-error layer) instead of dropping the connection — the OoP
440 // equivalent of `tower_http`'s `CatchPanicLayer`.
441 let _in_flight = InFlightGuard::new(guard.in_flight.clone());
442 let outcome = AssertUnwindSafe(next.run(request)).catch_unwind().await;
443
444 match outcome {
445 Ok(response) => response,
446 Err(panic) => {
447 let detail = panic_message(panic.as_ref());
448 tracing::error!(panic = %detail, "OoP gear handler panicked");
449 // `detail` is a server-side diagnostic only; `CanonicalError::internal`
450 // keeps it off the wire (`#[serde(skip)]`) per the error contract.
451 CanonicalError::internal(format!("gear handler panicked: {detail}"))
452 .create()
453 .into_response()
454 }
455 }
456}
457
458/// Extract a human-readable message from a caught panic payload.
459fn panic_message(panic: &(dyn std::any::Any + Send)) -> String {
460 panic
461 .downcast_ref::<&str>()
462 .map(|s| (*s).to_owned())
463 .or_else(|| panic.downcast_ref::<String>().cloned())
464 .unwrap_or_else(|| "unknown panic".to_owned())
465}
466
467// ---------------------------------------------------------------------------
468// Router assembly
469// ---------------------------------------------------------------------------
470
471/// Apply the framework middleware stack to the composed gear router: auth planes
472/// (when injected), the drain guard, and the canonical-error layer.
473///
474/// Probes are NOT layered here — they live on the outer router (see
475/// [`build_outer_router`]) so they stay unguarded and answerable during startup
476/// and drain.
477fn layer_gear_router(
478 gear_router: Router,
479 drain_guard: DrainGuard,
480 options: &OopServeOptions,
481) -> Router {
482 let mut gear = gear_router;
483
484 // Auth planes (installed only when injected). Add tenant plane first so it
485 // is *inner* to the platform plane — `internal_auth_middleware` must run
486 // BEFORE `security_context_middleware` (`cpt-cf-adr-two-plane-auth`).
487 if let Some(bearer) = options.bearer_authenticator.clone() {
488 gear = gear.layer(from_fn_with_state(
489 Arc::new(bearer),
490 security_context_middleware::<DynBearerAuthenticator>,
491 ));
492 }
493 if let Some(internal) = options.internal_authenticator.clone() {
494 gear = gear.layer(from_fn_with_state(
495 Arc::new(internal),
496 internal_auth_middleware::<DynInternalAuthenticator>,
497 ));
498 }
499
500 // Drain guard sits just inside the canonical-error layer: it rejects new
501 // requests while draining and tracks the in-flight count (catching handler
502 // panics so the counter can never leak and stall the drain).
503 gear = gear.layer(from_fn_with_state(drain_guard, drain_guard_middleware));
504
505 // Canonical-error middleware is the outermost gear layer so it post-processes
506 // every gear/auth response into `problem+json` with `trace_id` / `instance`
507 // filled — matching the in-process (`api-gateway`) path.
508 gear.layer(from_fn(canonical_error_middleware))
509}
510
511// ---------------------------------------------------------------------------
512// Serve
513// ---------------------------------------------------------------------------
514
515/// Serve the pre-bound `OoP` listener(s) until `cancel` fires, then drain.
516///
517/// Binding happens earlier in [`OopHttpServer::start`] so `/healthz` is up
518/// before the gear's (possibly slow) `start()` phase; this loop only drives the
519/// accept + graceful-shutdown machinery:
520/// 1. Serve with graceful shutdown wired to `cancel`.
521/// 2. On `cancel`: flip readiness to draining (via the caller-owned
522/// `ReadinessState`), wait up to `drain_timeout` for in-flight to reach zero,
523/// then let `axum` close the listener.
524///
525/// The `DirectoryService` deregistration is orchestrated by [`OopHttpServer::join`]
526/// so the full drain sequence (`cpt-cf-component-oop-bootstrap`) is honored.
527///
528/// # Errors
529/// Returns an error if the server task fails.
530async fn serve_loop(
531 listener: tokio::net::TcpListener,
532 router: Router,
533 drain_guard: DrainGuard,
534 sidecar: Option<(tokio::net::TcpListener, Router)>,
535 drain_timeout: Duration,
536 cancel: CancellationToken,
537) -> anyhow::Result<()> {
538 // Optional sidecar probe listener (bound by the caller).
539 let sidecar = match sidecar {
540 Some((probe_listener, probe_router)) => {
541 let shutdown = {
542 let cancel = cancel.clone();
543 async move { cancel.cancelled().await }
544 };
545 Some(tokio::spawn(async move {
546 if let Err(e) = axum::serve(probe_listener, probe_router)
547 .with_graceful_shutdown(shutdown)
548 .await
549 {
550 tracing::warn!(error = %e, "OoP probe sidecar server error");
551 }
552 }))
553 }
554 None => None,
555 };
556
557 // Main server: graceful shutdown waits for cancellation, then drains.
558 let shutdown = {
559 let cancel = cancel.clone();
560 let guard = drain_guard.clone();
561 async move {
562 cancel.cancelled().await;
563 // Drain step 1/2: flip readiness to 503 and start rejecting new work.
564 guard.begin_drain();
565 tracing::info!("OoP HTTP server draining (graceful shutdown)");
566 // Drain step 3: wait for in-flight requests to complete.
567 drain_in_flight(&guard, drain_timeout).await;
568 }
569 };
570
571 let result = axum::serve(
572 listener,
573 router.into_make_service_with_connect_info::<std::net::SocketAddr>(),
574 )
575 .with_graceful_shutdown(shutdown)
576 .await
577 .map_err(anyhow::Error::from);
578
579 if let Some(handle) = sidecar {
580 handle.abort();
581 if let Err(e) = handle.await
582 && !e.is_cancelled()
583 {
584 tracing::warn!(error = %e, "OoP probe sidecar task join error");
585 }
586 }
587
588 result
589}
590
591/// Owns an `OoP` gear's HTTP surface across the startup boundary.
592///
593/// [`start`](Self::start) binds the listener and serves the framework probes
594/// **immediately** — `/healthz` → `200`, `/readyz` → `starting`/`503`, and gear
595/// routes → `503 starting` — so the kubelet's liveness probe passes while the
596/// gear's (possibly slow) `start()` phase runs. Once the gear router is composed,
597/// [`attach`](Self::attach) swaps it in atomically (no port rebind) and starts
598/// background directory presence + dependency resolution. [`join`](Self::join)
599/// waits for the graceful drain on shutdown, then deregisters from
600/// `DirectoryService` (drain step 4, `cpt-cf-component-oop-bootstrap`).
601///
602/// Steps 5–7 of the drain order (reverse-dependency wait, stopping runtime
603/// services, process exit) are the caller's / operator's responsibility;
604/// the presence task stops when `cancel` fires.
605pub(super) struct OopHttpServer {
606 readiness: Arc<ReadinessState>,
607 late: LateRoutes,
608 drain_guard: DrainGuard,
609 options: OopServeOptions,
610 cancel: CancellationToken,
611 registration_task: Option<tokio::task::JoinHandle<()>>,
612 serve: tokio::task::JoinHandle<anyhow::Result<()>>,
613}
614
615impl OopHttpServer {
616 /// Bind the listener(s) and start serving probes immediately; gear routes
617 /// reply `503 starting` until [`attach`](Self::attach) publishes them.
618 ///
619 /// # Errors
620 /// Returns an error if the main (or sidecar) listener cannot be bound.
621 pub(super) async fn start(
622 readiness: Arc<ReadinessState>,
623 options: OopServeOptions,
624 cancel: CancellationToken,
625 ) -> anyhow::Result<Self> {
626 let late = LateRoutes::default();
627 let drain_guard = DrainGuard::new(Arc::clone(&readiness));
628 let mut options = options;
629
630 let listener = tokio::net::TcpListener::bind(options.listen_addr).await?;
631 let bound = listener.local_addr()?;
632 let configured_port = options.listen_addr.port();
633 if bound.port() != configured_port {
634 if let Some(rewritten) =
635 rewrite_advertise_uri_port(&options.advertise_uri, bound.port(), configured_port)
636 {
637 tracing::info!(
638 old = %options.advertise_uri,
639 new = %rewritten,
640 "advertise_uri rewritten to use actually bound port"
641 );
642 options.advertise_uri = rewritten;
643 } else {
644 tracing::info!(
645 advertise_uri = %options.advertise_uri,
646 bound_port = bound.port(),
647 configured_port,
648 "bound port differs from configured; leaving advertise_uri unchanged"
649 );
650 }
651 }
652 options.listen_addr = bound;
653 tracing::info!(
654 addr = %options.listen_addr,
655 "OoP HTTP server bound (probes live; gear routes attach after start)"
656 );
657
658 let outer = build_outer_router(Arc::clone(&readiness), late.clone());
659
660 let sidecar = if let Some(addr) = options.probe_bind_addr {
661 let probe_listener = tokio::net::TcpListener::bind(addr).await?;
662 let bound_probe = probe_listener.local_addr()?;
663 options.probe_bind_addr = Some(bound_probe);
664 tracing::info!(addr = %bound_probe, "OoP probe sidecar bound");
665 Some((
666 probe_listener,
667 build_probe_router(Arc::clone(&readiness), late.clone()),
668 ))
669 } else {
670 None
671 };
672
673 let serve = tokio::spawn(serve_loop(
674 listener,
675 outer,
676 drain_guard.clone(),
677 sidecar,
678 options.drain_timeout,
679 cancel.clone(),
680 ));
681
682 Ok(Self {
683 readiness,
684 late,
685 drain_guard,
686 options,
687 cancel,
688 registration_task: None,
689 serve,
690 })
691 }
692
693 /// The serve options (used by the caller to compose the `OpenAPI` document).
694 pub(super) fn options(&self) -> &OopServeOptions {
695 &self.options
696 }
697
698 /// Resolve the tenant-plane authenticator from the populated `ClientHub`.
699 ///
700 /// Called by the `OoP` serving path after the gear lifecycle's `start`
701 /// phase and before [`attach`](Self::attach) layers the middleware. When an
702 /// in-process authn stack (e.g. the `authn-resolver` gear) is linked into
703 /// the binary it registers a [`DynBearerAuthenticator`] bridge during
704 /// `init`; picking it up here installs `security_context_middleware` on the
705 /// gear routes. A no-op if an authenticator was already supplied directly or
706 /// none is registered in the hub.
707 pub(super) fn resolve_bearer_authenticator(&mut self, hub: &crate::ClientHub) {
708 if self.options.bearer_authenticator.is_some() {
709 return;
710 }
711 if let Ok(auth) = hub.get::<DynBearerAuthenticator>() {
712 self.options.bearer_authenticator = Some((*auth).clone());
713 tracing::info!(
714 gear = %self.options.gear_name,
715 "tenant-plane authenticator installed (security_context_middleware enabled)"
716 );
717 } else {
718 tracing::warn!(
719 gear = %self.options.gear_name,
720 "no tenant-plane authenticator registered in ClientHub; tenant plane not installed"
721 );
722 }
723 }
724
725 /// Publish the composed gear routes (they go live atomically) and start
726 /// background directory presence + dependency resolution.
727 ///
728 /// Presence (registration + heartbeat) and dep resolution start here — only
729 /// once the gear can actually serve — so the directory never advertises a
730 /// not-yet-serving instance.
731 pub(super) fn attach(&mut self, gear_router: Router, openapi_json: String) {
732 let openapi_arc: Arc<str> = Arc::from(openapi_json);
733 let layered = layer_gear_router(gear_router, self.drain_guard.clone(), &self.options);
734 self.late.publish(layered, Arc::clone(&openapi_arc));
735 self.readiness.mark_startup_complete();
736 tracing::info!(gear = %self.options.gear_name, "OoP gear routes attached (now serving)");
737
738 // Single directory-presence task: registration + heartbeat + self-heal.
739 // Dependency resolution is handled by the proxy-wiring phase (typed
740 // `#[toolkit::consumes]` clients), which runs before serving.
741 let mut registration_info = RegisterInstanceInfo::new(
742 self.options.gear_name.clone(),
743 self.options.instance_id.clone(),
744 )
745 .with_rest_endpoint(ServiceEndpoint::new(self.options.advertise_uri.clone()))
746 .with_openapi_spec(openapi_arc.to_string())
747 .with_labels(self.options.labels.clone());
748 if let Some(version) = self.options.version.clone() {
749 registration_info = registration_info.with_version(version);
750 }
751 self.registration_task = Some(tokio::spawn(super::oop_registration::presence_loop(
752 Arc::clone(&self.options.directory),
753 registration_info,
754 self.options.heartbeat_interval,
755 self.cancel.clone(),
756 )));
757 }
758
759 /// Wait for the server to drain on shutdown, then deregister from
760 /// `DirectoryService` (drain step 4 — after in-flight requests finish so
761 /// consumers don't see stale "ready" state).
762 ///
763 /// # Errors
764 /// Propagates a serve error; deregistration failures are logged only.
765 pub(super) async fn join(mut self) -> anyhow::Result<()> {
766 let serve_result = match self.serve.await {
767 Ok(r) => r,
768 Err(e) => Err(anyhow::anyhow!("OoP serve task join error: {e}")),
769 };
770
771 if let Some(task) = self.registration_task.take() {
772 task.abort();
773 }
774 if let Err(e) = self
775 .options
776 .directory
777 .deregister_instance(&self.options.gear_name, &self.options.instance_id)
778 .await
779 {
780 tracing::warn!(
781 gear = %self.options.gear_name,
782 error = %e,
783 "deregistration from DirectoryService failed on shutdown"
784 );
785 } else {
786 tracing::info!(gear = %self.options.gear_name, "deregistered from DirectoryService");
787 }
788
789 serve_result
790 }
791}
792
793/// Rewrite `advertise_uri` to use the actually bound port, but only if the
794/// URI currently uses the `configured_port` (i.e. the port the caller asked
795/// the server to bind to). This preserves user-provided load-balancer URLs
796/// that intentionally advertise a different port from the local socket.
797fn rewrite_advertise_uri_port(
798 advertise_uri: &str,
799 bound_port: u16,
800 configured_port: u16,
801) -> Option<String> {
802 let mut url = Url::parse(advertise_uri).ok()?;
803 if url.port_or_known_default() != Some(configured_port) {
804 return None;
805 }
806 url.set_port(Some(bound_port)).ok()?;
807 Some(url.to_string())
808}
809
810/// Wait up to `timeout` for the in-flight counter to reach zero.
811async fn drain_in_flight(guard: &DrainGuard, timeout: Duration) {
812 let deadline = tokio::time::Instant::now() + timeout;
813 loop {
814 let in_flight = guard.in_flight();
815 if in_flight == 0 {
816 tracing::info!("OoP drain complete: no in-flight requests");
817 return;
818 }
819 if tokio::time::Instant::now() >= deadline {
820 tracing::warn!(
821 in_flight,
822 timeout_secs = timeout.as_secs(),
823 "OoP drain timed out with in-flight requests remaining"
824 );
825 return;
826 }
827 tokio::time::sleep(Duration::from_millis(50)).await;
828 }
829}
830
831#[cfg(test)]
832#[cfg_attr(coverage_nightly, coverage(off))]
833#[path = "oop_serve_tests.rs"]
834mod tests;