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;