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