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