Skip to main content

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;