Skip to main content

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;