toolkit/discovery.rs
1//! Consumer-side discovery wiring for eventual readiness.
2//!
3//! Bridges the service [`DirectoryClient`](crate::DirectoryClient) into the
4//! contract layer's [`EndpointResolver`], and carries the
5//! [`ConsumerRegistration`]s emitted by `#[toolkit::consumes]` that the
6//! runtime's proxy-wiring phase replays at startup.
7//!
8//! Unconditional: consuming another gear over a discovered REST transport is
9//! the normal case, and `#[toolkit::consumes]` emits its registration
10//! unconditionally too. Putting this behind a feature only created a way to
11//! build a gear whose declared dependencies were silently never wired.
12
13use std::sync::Arc;
14
15use async_trait::async_trait;
16
17use cf_system_sdks::directory::DirectoryNotFound;
18
19use crate::DirectoryClient;
20use crate::client_hub::ClientHub;
21
22/// Re-export so the `#[toolkit::consumes]`-generated wire fn can name the
23/// resolver type as `toolkit::discovery::EndpointResolver` without the consuming
24/// gear crate needing a direct `toolkit-contract` dependency.
25pub use toolkit_contract::runtime::resolving::{EndpointResolver, ResolveError};
26
27/// Adapts the service [`DirectoryClient`] into the contract-layer
28/// [`EndpointResolver`] consumed by `DirectoryResolvingClient`.
29///
30/// Keeps `toolkit-contract` free of a dependency on the directory SDK: the
31/// adapter lives here in `toolkit`, which already owns `DirectoryClient`.
32///
33/// Successful resolutions are memoized for a short TTL ([`Self::TTL`]) so a
34/// per-call resolving client does not incur a directory round-trip on every
35/// business call (the `OoP` directory is a network hop). Only `Ok(Some(uri))` is
36/// cached; `Ok(None)` / `Err` always re-resolve. The TTL is deliberately short
37/// so endpoint churn self-heals within ~TTL.
38pub struct DirectoryEndpointResolver {
39 client: Arc<dyn DirectoryClient>,
40 cache: parking_lot::Mutex<std::collections::HashMap<String, (String, std::time::Instant)>>,
41}
42
43impl DirectoryEndpointResolver {
44 /// Memoization window for successful resolutions.
45 ///
46 /// Public because it bounds observable behaviour: a provider that moves is
47 /// followed within roughly this window, not instantly, since the cache
48 /// holds only successes and is not invalidated by a failed call.
49 pub const TTL: std::time::Duration = std::time::Duration::from_millis(1500);
50
51 /// Wrap a [`DirectoryClient`] with a short-TTL resolution cache.
52 #[must_use]
53 pub fn new(client: Arc<dyn DirectoryClient>) -> Self {
54 Self {
55 client,
56 cache: parking_lot::Mutex::new(std::collections::HashMap::new()),
57 }
58 }
59
60 fn cached(&self, gear: &str) -> Option<String> {
61 self.cache
62 .lock()
63 .get(gear)
64 .filter(|(_, at)| at.elapsed() < Self::TTL)
65 .map(|(uri, _)| uri.clone())
66 }
67
68 fn store(&self, gear: &str, uri: &str) {
69 self.cache
70 .lock()
71 .insert(gear.to_owned(), (uri.to_owned(), std::time::Instant::now()));
72 }
73}
74
75#[async_trait]
76impl EndpointResolver for DirectoryEndpointResolver {
77 async fn resolve_endpoint(&self, gear: &str) -> Result<Option<String>, ResolveError> {
78 // Fresh cache hit avoids the directory round-trip.
79 if let Some(uri) = self.cached(gear) {
80 return Ok(Some(uri));
81 }
82 // Preserve the trait's `Ok(None)` (no live instance / not ready) vs
83 // `Err` (directory backend unreachable) distinction so a real directory
84 // outage surfaces at `warn` rather than being silently treated as
85 // "provider not ready". The directory returns the typed
86 // `DirectoryNotFound` sentinel for the not-found case; anything else is
87 // a genuine backend/transport failure (relevant for the OoP gRPC
88 // directory client).
89 match self.client.resolve_rest_service(gear).await {
90 Ok(ep) => {
91 self.store(gear, &ep.uri);
92 Ok(Some(ep.uri))
93 }
94 Err(e) if e.downcast_ref::<DirectoryNotFound>().is_some() => Ok(None),
95 Err(e) => Err(ResolveError::new(gear, e)),
96 }
97 }
98}
99
100/// An [`EndpointResolver`] that never resolves any gear (`Ok(None)` for every
101/// lookup). Used by the proxy-wiring phase when no [`DirectoryClient`] is present
102/// in the [`ClientHub`]: co-located (local) consumers still short-circuit and
103/// wire successfully, while remote consumers register as *unresolved* readiness
104/// gates (so `/readyz` correctly stays 503) rather than being silently skipped.
105pub struct NullEndpointResolver;
106
107#[async_trait]
108impl EndpointResolver for NullEndpointResolver {
109 async fn resolve_endpoint(&self, _gear: &str) -> Result<Option<String>, ResolveError> {
110 Ok(None)
111 }
112}
113
114/// An [`EndpointResolver`] that always resolves to a single, fixed endpoint —
115/// the ADR-0004 static-endpoint override escape hatch for local development and
116/// integration testing. It bypasses the service directory: a consumer configured
117/// with `gears.<owner>.config.consumer_wiring.<dep> = "http://host:port"` wires a
118/// resolving client whose endpoint is this constant. Not for production use.
119pub struct StaticEndpointResolver {
120 endpoint: String,
121}
122
123impl StaticEndpointResolver {
124 /// Wrap a fixed base endpoint URI.
125 #[must_use]
126 pub fn new(endpoint: impl Into<String>) -> Self {
127 Self {
128 endpoint: endpoint.into(),
129 }
130 }
131}
132
133#[async_trait]
134impl EndpointResolver for StaticEndpointResolver {
135 async fn resolve_endpoint(&self, _gear: &str) -> Result<Option<String>, ResolveError> {
136 Ok(Some(self.endpoint.clone()))
137 }
138}
139
140/// Outcome of a [`ConsumerRegistration::wire`] call: whether the dependency was
141/// satisfied by a co-located compile-time impl (no directory dependency) or by
142/// a directory-resolving REST client that still needs endpoint resolution.
143///
144/// The runtime uses this to gate readiness: `Local` deps are readiness-resolved
145/// immediately (no `/readyz` gating, no probe loop), while `Remote` deps gate
146/// readiness until the directory resolves them (ADR-0007).
147#[derive(Debug, Clone, Copy, PartialEq, Eq)]
148pub enum WireOutcome {
149 /// A compile-time (in-process) impl was already registered — it won the
150 /// `hub.try_get_local` short-circuit; no remote binding was created.
151 Local,
152 /// A directory-resolving REST client was registered; its endpoint is
153 /// resolved lazily/at readiness-probe time.
154 Remote,
155}
156
157/// A consumer→provider contract wiring, registered by `#[toolkit::consumes]`
158/// via `inventory::submit!` and replayed by the runtime's proxy-wiring phase.
159///
160/// `wire` is a non-capturing `fn` (the provider gear name is inlined into the
161/// generated body). It must:
162/// 1. short-circuit returning [`WireOutcome::Local`] if a compile-time (local)
163/// impl is already registered (`hub.try_get_local::<dyn Trait>().is_some()`),
164/// so in-process providers win;
165/// 2. otherwise register a directory-resolving REST client under the contract
166/// trait via [`ClientHub::register_remote_proxy`] (using the supplied
167/// [`EndpointResolver`]) and return [`WireOutcome::Remote`].
168///
169/// Step 1 must use `try_get_local`, not `try_get`. The hub is keyed by type, so
170/// when two consumers in one process depend on the same contract the second to
171/// wire would find the first one's proxy under `try_get`, report `Local`, and
172/// have its dependency marked readiness-resolved — flipping `/readyz` to 200
173/// while the provider is still absent. `register_remote_proxy` is what makes
174/// that distinction visible to the hub.
175pub struct ConsumerRegistration {
176 /// Gear that declares the dependency (diagnostics).
177 pub owner_gear: &'static str,
178 /// Provider gear name the contract resolves against (diagnostics).
179 pub dep_gear: &'static str,
180 /// Wiring action: register the client into the hub; reports whether the
181 /// dependency was satisfied locally or bound to a remote client.
182 pub wire: fn(&ClientHub, Arc<dyn EndpointResolver>) -> anyhow::Result<WireOutcome>,
183}
184
185inventory::collect!(ConsumerRegistration);
186
187#[cfg(test)]
188#[cfg_attr(coverage_nightly, coverage(off))]
189mod tests {
190 use super::*;
191 use cf_system_sdks::directory::{RegisterInstanceInfo, ServiceEndpoint, ServiceInstanceInfo};
192 use std::sync::atomic::{AtomicUsize, Ordering};
193
194 /// Counts `resolve_rest_service` calls; all other methods are unused here.
195 struct CountingDirectory {
196 calls: AtomicUsize,
197 }
198
199 #[async_trait]
200 impl DirectoryClient for CountingDirectory {
201 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
202 unimplemented!()
203 }
204 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
205 self.calls.fetch_add(1, Ordering::SeqCst);
206 Ok(ServiceEndpoint::new(format!("http://{gear}.local")))
207 }
208 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
209 unimplemented!()
210 }
211 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
212 unimplemented!()
213 }
214 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
215 unimplemented!()
216 }
217 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
218 unimplemented!()
219 }
220 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
221 unimplemented!()
222 }
223 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
224 unimplemented!()
225 }
226 }
227
228 #[tokio::test]
229 async fn memoizes_successful_resolution_within_ttl() {
230 let dir = Arc::new(CountingDirectory {
231 calls: AtomicUsize::new(0),
232 });
233 let resolver = DirectoryEndpointResolver::new(dir.clone());
234
235 for _ in 0..3 {
236 let ep = resolver.resolve_endpoint("billing").await.unwrap();
237 assert_eq!(ep.as_deref(), Some("http://billing.local"));
238 }
239
240 assert_eq!(
241 dir.calls.load(Ordering::SeqCst),
242 1,
243 "successful resolution must be memoized within the TTL (one directory lookup)"
244 );
245 }
246
247 #[tokio::test]
248 async fn static_resolver_always_returns_fixed_endpoint() {
249 let resolver = StaticEndpointResolver::new("http://localhost:8081");
250 assert_eq!(
251 resolver
252 .resolve_endpoint("anything")
253 .await
254 .unwrap()
255 .as_deref(),
256 Some("http://localhost:8081")
257 );
258 assert_eq!(
259 resolver.resolve_endpoint("other").await.unwrap().as_deref(),
260 Some("http://localhost:8081")
261 );
262 }
263
264 #[tokio::test]
265 async fn null_resolver_never_resolves() {
266 assert_eq!(
267 NullEndpointResolver
268 .resolve_endpoint("billing")
269 .await
270 .unwrap(),
271 None
272 );
273 }
274
275 /// Answers every lookup with a caller-chosen failure, so the resolver's
276 /// two error branches can be told apart.
277 struct FailingDirectory {
278 make_error: fn(&str) -> anyhow::Error,
279 }
280
281 #[async_trait]
282 impl DirectoryClient for FailingDirectory {
283 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
284 unimplemented!()
285 }
286 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
287 Err((self.make_error)(gear))
288 }
289 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
290 unimplemented!()
291 }
292 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
293 unimplemented!()
294 }
295 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
296 unimplemented!()
297 }
298 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
299 unimplemented!()
300 }
301 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
302 unimplemented!()
303 }
304 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
305 unimplemented!()
306 }
307 }
308
309 #[tokio::test]
310 async fn not_found_sentinel_resolves_to_ok_none() {
311 // "The provider has not registered yet" is a routine startup condition,
312 // not a failure: it must come back as `Ok(None)` so the readiness probe
313 // logs it at `debug` and the resolving client reports `Unresolved`
314 // rather than warning on every call.
315 let dir = Arc::new(FailingDirectory {
316 make_error: |gear| DirectoryNotFound::new(format!("gear {gear}")).into(),
317 });
318 let resolver = DirectoryEndpointResolver::new(dir);
319
320 assert_eq!(resolver.resolve_endpoint("billing").await.unwrap(), None);
321 }
322
323 #[tokio::test]
324 async fn backend_failure_resolves_to_err() {
325 // Anything that is not the sentinel is a genuine directory outage and
326 // must stay an error, so a stuck 503 has a diagnostic trail.
327 let dir = Arc::new(FailingDirectory {
328 make_error: |_| anyhow::anyhow!("connection refused"),
329 });
330 let resolver = DirectoryEndpointResolver::new(dir);
331
332 assert!(resolver.resolve_endpoint("billing").await.is_err());
333 }
334
335 /// Mirrors the body `#[toolkit::consumes]` generates, so the wiring
336 /// contract is exercised without going through `inventory` (which is
337 /// link-time global and cannot hold two registrations built by a test).
338 #[allow(
339 clippy::unnecessary_wraps,
340 reason = "mirrors the fallible signature of ConsumerRegistration::wire so the test \
341 exercises the same shape the macro emits"
342 )]
343 fn generated_wire_body<C>(hub: &ClientHub, proxy: Arc<C>) -> anyhow::Result<WireOutcome>
344 where
345 C: ?Sized + Send + Sync + 'static,
346 {
347 if hub.try_get_local::<C>().is_some() {
348 return Ok(WireOutcome::Local);
349 }
350 if hub.has_remote_proxy::<C>() {
351 return Ok(WireOutcome::Remote);
352 }
353 hub.register_remote_proxy::<C>(proxy);
354 Ok(WireOutcome::Remote)
355 }
356
357 trait Payments: Send + Sync {}
358 struct PaymentsProxy;
359 impl Payments for PaymentsProxy {}
360 struct PaymentsLocal;
361 impl Payments for PaymentsLocal {}
362
363 #[test]
364 fn two_consumers_of_one_contract_both_report_remote() {
365 // Two gears in one process consuming the same contract from the same
366 // provider. Whichever wires second used to find the first one's proxy,
367 // conclude the dependency was local, and get it marked
368 // readiness-resolved — with both registrations sharing a `dep_gear`,
369 // that flipped the very entry the unresolved remote dep was gating on,
370 // and `/readyz` went green with no provider up.
371 let hub = ClientHub::new();
372
373 let first = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
374 let second = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
375
376 assert_eq!(first, WireOutcome::Remote);
377 assert_eq!(
378 second,
379 WireOutcome::Remote,
380 "the second consumer must still gate readiness, not short-circuit to Local"
381 );
382 }
383
384 #[test]
385 fn a_genuine_local_impl_still_wins() {
386 // The short-circuit must keep working for its actual purpose: a
387 // co-located provider registered during init means no HTTP hop and no
388 // readiness gate.
389 let hub = ClientHub::new();
390 hub.register::<dyn Payments>(Arc::new(PaymentsLocal));
391
392 let outcome = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
393
394 assert_eq!(outcome, WireOutcome::Local);
395 }
396
397 #[tokio::test]
398 async fn local_directory_client_not_found_reaches_the_resolver_as_ok_none() {
399 // End-to-end over the real in-process client rather than a stub: this
400 // is the pairing that was broken (the client returned a bare `anyhow`,
401 // so the resolver could never produce `Ok(None)`).
402 let dir = Arc::new(crate::directory::LocalDirectoryClient::new(Arc::new(
403 crate::runtime::GearManager::new(),
404 )));
405 let resolver = DirectoryEndpointResolver::new(dir);
406
407 assert_eq!(
408 resolver.resolve_endpoint("never-registered").await.unwrap(),
409 None
410 );
411 }
412}