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;