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::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;