Skip to main content

toolkit/runtime/
oop_registration.rs

1//! Background self-registration and dependency resolution for `OoP` gears
2//! (`cpt-cf-component-oop-bootstrap`).
3//!
4//! Both run as non-blocking background tasks so the HTTP server and its probes
5//! are up immediately:
6//!
7//! - **Presence** ([`presence_loop`]) is the single owner of this instance's
8//!   `DirectoryService` presence. It registers the instance's gRPC/REST
9//!   endpoints and `OpenAPI` spec (retrying with exponential backoff,
10//!   100ms → 30s cap), then runs one loop that both sends periodic **heartbeats**
11//!   (the steady-state liveness signal the directory uses to keep the instance
12//!   routable) and periodically **re-registers** to self-heal after a
13//!   `DirectoryService` restart / connection loss. Consolidating both into one task (owned
14//!   by the `OoP` serve lifecycle) avoids the split-brain where an independent
15//!   heartbeat task and a re-registration task raced to write conflicting
16//!   liveness state onto the same directory record.
17//! - **Dependency resolution** ([`resolve_deps`]) polls
18//!   `DirectoryService.resolve_rest_service` for each declared dependency, wires
19//!   the resolved base URL into the [`ClientHub`] via [`ResolvedRestEndpoints`],
20//!   and marks the dep resolved so `/readyz` can flip to `200`.
21//!
22//! In Profile 1 (in-process) dependency resolution is a no-op — the topo-sorted
23//! `HostRuntime` already guarantees deps are initialized.
24
25use std::sync::Arc;
26use std::time::Duration;
27
28use dashmap::DashMap;
29use tokio_util::sync::CancellationToken;
30
31use cf_system_sdks::directory::{DirectoryClient, RegisterInstanceInfo};
32
33use super::readiness::ReadinessState;
34
35/// Initial retry backoff for registration and dependency polling.
36const INITIAL_BACKOFF: Duration = Duration::from_millis(100);
37/// Maximum retry backoff (cap).
38const MAX_BACKOFF: Duration = Duration::from_secs(30);
39/// Interval at which a successfully-registered instance re-registers to
40/// self-heal after a directory restart / connection loss.
41const RE_REGISTER_INTERVAL: Duration = Duration::from_secs(30);
42
43/// Next backoff in the exponential schedule (doubles, capped at [`MAX_BACKOFF`]).
44fn next_backoff(current: Duration) -> Duration {
45    current.saturating_mul(2).min(MAX_BACKOFF)
46}
47
48/// Sleep for `dur`, returning early (`false`) if `cancel` fires first.
49async fn sleep_or_cancel(dur: Duration, cancel: &CancellationToken) -> bool {
50    tokio::select! {
51        () = cancel.cancelled() => false,
52        () = tokio::time::sleep(dur) => true,
53    }
54}
55
56/// Resolved REST base URLs for the gear's dependencies, keyed by gear name.
57///
58/// Registered into the [`ClientHub`] by the `OoP` bootstrap and populated as
59/// dependencies resolve. Until typed REST client codegen lands, gears look up a
60/// dependency's base URL here (`hub.get::<ResolvedRestEndpoints>()`).
61#[derive(Debug, Default)]
62pub struct ResolvedRestEndpoints {
63    inner: DashMap<String, String>,
64}
65
66impl ResolvedRestEndpoints {
67    /// Create an empty registry.
68    #[must_use]
69    pub fn new() -> Self {
70        Self::default()
71    }
72
73    /// Record a resolved base URL for `gear`.
74    pub fn set(&self, gear: impl Into<String>, uri: impl Into<String>) {
75        self.inner.insert(gear.into(), uri.into());
76    }
77
78    /// Look up the resolved base URL for `gear`, if known.
79    #[must_use]
80    pub fn get(&self, gear: &str) -> Option<String> {
81        self.inner.get(gear).map(|v| v.value().clone())
82    }
83
84    /// Number of resolved dependencies.
85    #[must_use]
86    pub fn len(&self) -> usize {
87        self.inner.len()
88    }
89
90    /// Returns `true` if no dependencies have been resolved yet.
91    #[must_use]
92    pub fn is_empty(&self) -> bool {
93        self.inner.is_empty()
94    }
95}
96
97/// Register `info` with the directory, retrying with exponential backoff until
98/// success or cancellation. Returns `true` on success, `false` if cancelled.
99async fn register_once_with_backoff(
100    directory: &Arc<dyn DirectoryClient>,
101    info: &RegisterInstanceInfo,
102    cancel: &CancellationToken,
103) -> bool {
104    let mut backoff = INITIAL_BACKOFF;
105    loop {
106        if cancel.is_cancelled() {
107            return false;
108        }
109        match directory.register_instance(info.clone()).await {
110            Ok(()) => {
111                tracing::info!(gear = %info.gear, instance = %info.instance_id, "registered with DirectoryService");
112                return true;
113            }
114            Err(e) => {
115                tracing::warn!(
116                    gear = %info.gear,
117                    error = %e,
118                    backoff_ms = backoff.as_millis(),
119                    "registration attempt failed; retrying"
120                );
121                if !sleep_or_cancel(backoff, cancel).await {
122                    return false;
123                }
124                backoff = next_backoff(backoff);
125            }
126        }
127    }
128}
129
130/// Single directory-presence loop for an `OoP` instance.
131///
132/// Registers `info` (with backoff), then owns **both** liveness signals in one
133/// task until `cancel` fires:
134///
135/// - a **heartbeat** every `heartbeat_interval` — the steady-state signal the
136///   directory uses to keep the instance `Healthy`/routable (and to avoid
137///   heartbeat-timeout eviction);
138/// - an idempotent **re-registration** every [`RE_REGISTER_INTERVAL`] — self-heals
139///   after a `DirectoryService` restart / connection loss (a heartbeat alone cannot,
140///   since the directory silently ignores heartbeats for an instance it has
141///   forgotten). Re-registration preserves the directory-side liveness state
142///   (see `GearManager::register_instance`), so it does not disturb the
143///   heartbeat-maintained `Healthy` state.
144///
145/// Registering *before* the first heartbeat is deliberate: a heartbeat for an
146/// unregistered instance is a no-op on the directory.
147pub(super) async fn presence_loop(
148    directory: Arc<dyn DirectoryClient>,
149    info: RegisterInstanceInfo,
150    heartbeat_interval: Duration,
151    cancel: CancellationToken,
152) {
153    if !register_once_with_backoff(&directory, &info, &cancel).await {
154        return;
155    }
156
157    // `tokio::time::interval` panics on a zero period; clamp defensively.
158    let heartbeat_interval = heartbeat_interval.max(Duration::from_secs(1));
159    let mut heartbeat = tokio::time::interval(heartbeat_interval);
160    heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
161    heartbeat.tick().await; // consume the immediate first tick
162
163    let mut reregister = tokio::time::interval(RE_REGISTER_INTERVAL);
164    reregister.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
165    reregister.tick().await; // consume the immediate first tick
166
167    loop {
168        tokio::select! {
169            () = cancel.cancelled() => return,
170            _ = heartbeat.tick() => {
171                if let Err(e) = directory.send_heartbeat(&info.gear, &info.instance_id).await {
172                    tracing::warn!(
173                        gear = %info.gear,
174                        instance = %info.instance_id,
175                        error = %e,
176                        "heartbeat failed; re-registering to self-heal"
177                    );
178                    if !register_once_with_backoff(&directory, &info, &cancel).await {
179                        return;
180                    }
181                } else {
182                    tracing::trace!(gear = %info.gear, "heartbeat sent");
183                }
184            }
185            _ = reregister.tick() => {
186                if !register_once_with_backoff(&directory, &info, &cancel).await {
187                    return;
188                }
189            }
190        }
191    }
192}
193
194/// Poll for a single dependency until resolved (or cancelled), then wire it into
195/// the resolved-endpoints registry and mark it resolved in readiness.
196async fn resolve_one_dep(
197    directory: Arc<dyn DirectoryClient>,
198    dep: String,
199    readiness: Arc<ReadinessState>,
200    resolved: Arc<ResolvedRestEndpoints>,
201    cancel: CancellationToken,
202) {
203    let mut backoff = INITIAL_BACKOFF;
204    loop {
205        if cancel.is_cancelled() {
206            return;
207        }
208        match directory.resolve_rest_service(&dep).await {
209            Ok(endpoint) => {
210                tracing::info!(dep = %dep, endpoint = %endpoint.uri, "resolved REST dependency");
211                resolved.set(dep.clone(), endpoint.uri);
212                readiness.mark_dep_resolved(&dep);
213                return;
214            }
215            Err(e) => {
216                tracing::debug!(dep = %dep, error = %e, "dependency not yet resolvable; retrying");
217                if !sleep_or_cancel(backoff, &cancel).await {
218                    return;
219                }
220                backoff = next_backoff(backoff);
221            }
222        }
223    }
224}
225
226/// Spawn one resolution task per dependency. Each resolves independently and
227/// gates `/readyz` via `readiness`. A no-op when `deps` is empty (Profile 1).
228pub(super) fn resolve_deps(
229    directory: &Arc<dyn DirectoryClient>,
230    deps: Vec<String>,
231    readiness: &Arc<ReadinessState>,
232    resolved: &Arc<ResolvedRestEndpoints>,
233    cancel: &CancellationToken,
234) {
235    for dep in deps {
236        tokio::spawn(resolve_one_dep(
237            Arc::clone(directory),
238            dep,
239            Arc::clone(readiness),
240            Arc::clone(resolved),
241            cancel.clone(),
242        ));
243    }
244}
245
246#[cfg(test)]
247#[cfg_attr(coverage_nightly, coverage(off))]
248#[path = "oop_registration_tests.rs"]
249mod tests;