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;