toolkit/runtime/oop_registration.rs
1//! Background self-registration (directory presence) for `OoP` gears
2//! (`cpt-cf-component-oop-bootstrap`).
3//!
4//! **Presence** ([`presence_loop`]) is the single owner of this instance's
5//! `DirectoryService` presence. It registers the instance's gRPC/REST endpoints
6//! and `OpenAPI` spec (retrying with exponential backoff, 100ms → 30s cap), then
7//! runs one loop that both sends periodic **heartbeats** (the steady-state
8//! liveness signal the directory uses to keep the instance routable) and
9//! periodically **re-registers** to self-heal after a `DirectoryService` restart
10//! / connection loss. Consolidating both into one task (owned by the `OoP` serve
11//! lifecycle) avoids the split-brain where an independent heartbeat task and a
12//! re-registration task raced to write conflicting liveness state onto the same
13//! directory record.
14//!
15//! Dependency resolution is **not** here: consumers are wired by the
16//! proxy-wiring phase via typed `#[toolkit::consumes]` directory-resolving
17//! clients, which feed the shared
18//! [`DependencyChecker`](super::readiness::DependencyChecker) that gates
19//! `/readyz`.
20
21use std::sync::Arc;
22use std::time::Duration;
23
24use tokio_util::sync::CancellationToken;
25
26use cf_system_sdks::directory::{
27 DirectoryClient, DirectoryInvalidArgument, DirectoryPermissionDenied, RegisterInstanceInfo,
28};
29
30/// Initial retry backoff for registration and dependency polling.
31const INITIAL_BACKOFF: Duration = Duration::from_millis(100);
32/// Maximum retry backoff (cap).
33const MAX_BACKOFF: Duration = Duration::from_secs(30);
34/// Interval at which a successfully-registered instance re-registers to
35/// self-heal after a directory restart / connection loss.
36const RE_REGISTER_INTERVAL: Duration = Duration::from_secs(30);
37
38/// Next backoff in the exponential schedule (doubles, capped at [`MAX_BACKOFF`]).
39fn next_backoff(current: Duration) -> Duration {
40 current.saturating_mul(2).min(MAX_BACKOFF)
41}
42
43/// Sleep for `dur`, returning early (`false`) if `cancel` fires first.
44async fn sleep_or_cancel(dur: Duration, cancel: &CancellationToken) -> bool {
45 tokio::select! {
46 () = cancel.cancelled() => false,
47 () = tokio::time::sleep(dur) => true,
48 }
49}
50
51/// Classification of a failed registration attempt.
52///
53/// Permanence drives the retry decision and is *independent* of any log text, so
54/// a permanent rejection with no remediation hint still stops the loop rather
55/// than being retried forever.
56enum RejectionKind {
57 /// Retrying the identical request can never succeed — stop the backoff and
58 /// log loudly. The optional string is an operator-facing "what to check"
59 /// hint (the rejection *category* is already in the sentinel's `Display`).
60 Permanent(Option<&'static str>),
61 /// Transient or unrecognised — keep retrying with backoff.
62 Transient,
63}
64
65/// Classify a registration error for the retry decision.
66fn classify_rejection(err: &anyhow::Error) -> RejectionKind {
67 if err.is::<DirectoryInvalidArgument>() {
68 return RejectionKind::Permanent(Some("Check its configured labels and endpoint URIs"));
69 }
70 if err.is::<DirectoryPermissionDenied>() {
71 return RejectionKind::Permanent(Some(
72 "Check that this peer is authorized to register this gear (its ServiceAccount \
73 namespace / trust domain allowlist)",
74 ));
75 }
76 RejectionKind::Transient
77}
78
79/// Register `info` with the directory, retrying with exponential backoff.
80///
81/// Returns `true` if the presence loop should keep running (the directory
82/// accepted the registration, or [permanently rejected](classify_rejection)
83/// it), and `false` only if the task was cancelled. A permanent rejection still
84/// returns `true`, so the caller keeps the presence loop alive; see
85/// [`presence_loop`].
86///
87/// **Not cancel-safe under `abort()`:** `cancel` is checked only between attempts,
88/// so the token exits the loop cleanly between RPCs but an abort can drop this
89/// mid-`register_instance`. Shutdown aborts the task then deregisters, so an abort
90/// landing mid-RPC can leave a registration committed after the deregister — a
91/// stale instance until its heartbeat TTL evicts it.
92async fn register_once_with_backoff(
93 directory: &Arc<dyn DirectoryClient>,
94 info: &RegisterInstanceInfo,
95 cancel: &CancellationToken,
96) -> bool {
97 let mut backoff = INITIAL_BACKOFF;
98 loop {
99 if cancel.is_cancelled() {
100 return false;
101 }
102 match directory.register_instance(info.clone()).await {
103 Ok(()) => {
104 tracing::info!(gear = %info.gear, instance = %info.instance_id, "registered with DirectoryService");
105 return true;
106 }
107 Err(e) => {
108 if let RejectionKind::Permanent(guidance) = classify_rejection(&e) {
109 tracing::error!(
110 gear = %info.gear,
111 instance = %info.instance_id,
112 error = %e,
113 guidance = guidance.unwrap_or("(no remediation hint available)"),
114 "registration permanently rejected by the directory; this instance \
115 may not be discoverable"
116 );
117 return true;
118 }
119 tracing::warn!(
120 gear = %info.gear,
121 error = %e,
122 backoff_ms = backoff.as_millis(),
123 "registration attempt failed; retrying"
124 );
125 if !sleep_or_cancel(backoff, cancel).await {
126 return false;
127 }
128 backoff = next_backoff(backoff);
129 }
130 }
131 }
132}
133
134/// Single directory-presence loop for an `OoP` instance.
135///
136/// Registers `info` (with backoff), then owns **both** liveness signals in one
137/// task until `cancel` fires:
138///
139/// - a **heartbeat** every `heartbeat_interval` — the steady-state signal the
140/// directory uses to keep the instance `Healthy`/routable (and to avoid
141/// heartbeat-timeout eviction);
142/// - an idempotent **re-registration** every [`RE_REGISTER_INTERVAL`] — self-heals
143/// after a `DirectoryService` restart / connection loss (a heartbeat alone cannot,
144/// since the directory silently ignores heartbeats for an instance it has
145/// forgotten). Re-registration preserves the directory-side liveness state
146/// (see `GearManager::register_instance`), so it does not disturb the
147/// heartbeat-maintained `Healthy` state.
148///
149/// Registering *before* the first heartbeat is deliberate: a heartbeat for an
150/// unregistered instance is a no-op on the directory.
151///
152/// A [permanent rejection](classify_rejection) (a malformed request or
153/// an authorization denial) never tears the presence loop down: it is logged
154/// loudly at every attempt, but the loop keeps sending heartbeats (so an
155/// already-registered instance is not evicted by a transient directory-side
156/// policy / validation change) and keeps retrying the idempotent
157/// re-registration (so it self-heals if the directory later accepts it). A
158/// *recoverable* rejection — e.g. a gRPC service-name conflict that clears once
159/// the current owner deregisters — is retried with backoff like any transient
160/// error. Only cancellation stops the loop.
161pub(super) async fn presence_loop(
162 directory: Arc<dyn DirectoryClient>,
163 info: RegisterInstanceInfo,
164 heartbeat_interval: Duration,
165 cancel: CancellationToken,
166) {
167 if !register_once_with_backoff(&directory, &info, &cancel).await {
168 return;
169 }
170
171 // `tokio::time::interval` panics on a zero period; clamp defensively.
172 let heartbeat_interval = heartbeat_interval.max(Duration::from_secs(1));
173 let mut heartbeat = tokio::time::interval(heartbeat_interval);
174 heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
175 heartbeat.tick().await; // consume the immediate first tick
176
177 let mut reregister = tokio::time::interval(RE_REGISTER_INTERVAL);
178 reregister.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
179 reregister.tick().await; // consume the immediate first tick
180
181 loop {
182 tokio::select! {
183 () = cancel.cancelled() => return,
184 _ = heartbeat.tick() => {
185 if let Err(e) = directory.send_heartbeat(&info.gear, &info.instance_id).await {
186 tracing::warn!(
187 gear = %info.gear,
188 instance = %info.instance_id,
189 error = %e,
190 "heartbeat failed; re-registering to self-heal"
191 );
192 if !register_once_with_backoff(&directory, &info, &cancel).await {
193 return;
194 }
195 } else {
196 tracing::trace!(gear = %info.gear, "heartbeat sent");
197 }
198 }
199 _ = reregister.tick() => {
200 if !register_once_with_backoff(&directory, &info, &cancel).await {
201 return;
202 }
203 }
204 }
205 }
206}
207
208#[cfg(test)]
209#[cfg_attr(coverage_nightly, coverage(off))]
210#[path = "oop_registration_tests.rs"]
211mod tests;