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