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::{DirectoryInvalidArgument, 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) if e.downcast_ref::<DirectoryInvalidArgument>().is_some() => {
96 // A malformed provider name is a permanent configuration error,
97 // not a transient outage: re-resolving can never succeed. Log
98 // it distinctly from the generic backend-failure path (at
99 // `error`, not the caller's transient `warn`) so the actionable
100 // "fix your config" message is not mistaken for a directory
101 // outage, and still return `Err` so the call fails.
102 tracing::error!(
103 gear,
104 error = %e,
105 "provider name permanently rejected by the directory (invalid \
106 argument); this is a static configuration error, not a directory \
107 outage, and will never resolve. Check the consumed gear name"
108 );
109 Err(ResolveError::new(gear, e))
110 }
111 Err(e) => Err(ResolveError::new(gear, e)),
112 }
113 }
114}
115
116/// An [`EndpointResolver`] that never resolves any gear (`Ok(None)` for every
117/// lookup). Used by the proxy-wiring phase when no [`DirectoryClient`] is present
118/// in the [`ClientHub`]: co-located (local) consumers still short-circuit and
119/// wire successfully, while remote consumers register as *unresolved* readiness
120/// gates (so `/readyz` correctly stays 503) rather than being silently skipped.
121pub struct NullEndpointResolver;
122
123#[async_trait]
124impl EndpointResolver for NullEndpointResolver {
125 async fn resolve_endpoint(&self, _gear: &str) -> Result<Option<String>, ResolveError> {
126 Ok(None)
127 }
128}
129
130/// An [`EndpointResolver`] that always resolves to a single, fixed endpoint —
131/// the ADR-0004 static-endpoint override escape hatch for local development and
132/// integration testing. It bypasses the service directory: a consumer configured
133/// with `gears.<owner>.config.consumer_wiring.<dep> = "http://host:port"` wires a
134/// resolving client whose endpoint is this constant. Not for production use.
135pub struct StaticEndpointResolver {
136 endpoint: String,
137}
138
139impl StaticEndpointResolver {
140 /// Wrap a fixed base endpoint URI.
141 #[must_use]
142 pub fn new(endpoint: impl Into<String>) -> Self {
143 Self {
144 endpoint: endpoint.into(),
145 }
146 }
147}
148
149#[async_trait]
150impl EndpointResolver for StaticEndpointResolver {
151 async fn resolve_endpoint(&self, _gear: &str) -> Result<Option<String>, ResolveError> {
152 Ok(Some(self.endpoint.clone()))
153 }
154}
155
156/// Outcome of a [`ConsumerRegistration::wire`] call: whether the dependency was
157/// satisfied by a co-located compile-time impl (no directory dependency) or by
158/// a directory-resolving REST client that still needs endpoint resolution.
159///
160/// The runtime uses this to gate readiness: `Local` deps are readiness-resolved
161/// immediately (no `/readyz` gating, no probe loop), while `Remote` deps gate
162/// readiness until the directory resolves them (ADR-0007).
163#[derive(Debug, Clone, Copy, PartialEq, Eq)]
164pub enum WireOutcome {
165 /// A compile-time (in-process) impl was already registered — it won the
166 /// `hub.try_get_local` short-circuit; no remote binding was created.
167 Local,
168 /// A directory-resolving REST client was registered; its endpoint is
169 /// resolved lazily/at readiness-probe time.
170 Remote,
171}
172
173/// A consumer→provider contract wiring, registered by `#[toolkit::consumes]`
174/// via `inventory::submit!` and replayed by the runtime's proxy-wiring phase.
175///
176/// `wire` is a non-capturing `fn` (the provider gear name is inlined into the
177/// generated body). It must:
178/// 1. short-circuit returning [`WireOutcome::Local`] if a compile-time (local)
179/// impl is already registered (`hub.try_get_local::<dyn Trait>().is_some()`),
180/// so in-process providers win;
181/// 2. otherwise register a directory-resolving REST client under the contract
182/// trait via [`ClientHub::register_remote_proxy`] (using the supplied
183/// [`EndpointResolver`]) and return [`WireOutcome::Remote`].
184///
185/// Step 1 must use `try_get_local`, not `try_get`. The hub is keyed by type, so
186/// when two consumers in one process depend on the same contract the second to
187/// wire would find the first one's proxy under `try_get`, report `Local`, and
188/// have its dependency marked readiness-resolved — flipping `/readyz` to 200
189/// while the provider is still absent. `register_remote_proxy` is what makes
190/// that distinction visible to the hub.
191pub struct ConsumerRegistration {
192 /// Gear that declares the dependency (diagnostics).
193 pub owner_gear: &'static str,
194 /// Provider gear name the contract resolves against (diagnostics).
195 pub dep_gear: &'static str,
196 /// Wiring action: register the client into the hub; reports whether the
197 /// dependency was satisfied locally or bound to a remote client. See
198 /// [`WireFn`] for the argument roles.
199 pub wire: WireFn,
200}
201
202/// The wiring function pointer emitted by `#[toolkit::consumes]`. Registers the
203/// client into the `ClientHub`, given the endpoint resolver and the process's
204/// platform-plane credential source (threaded onto a directory-resolving remote
205/// client). Aliased to keep the fn-pointer type readable.
206pub type WireFn = fn(
207 &ClientHub,
208 Arc<dyn EndpointResolver>,
209 Option<&toolkit_contract::runtime::config::InternalTokenProvider>,
210) -> anyhow::Result<WireOutcome>;
211
212inventory::collect!(ConsumerRegistration);
213
214#[cfg(test)]
215#[cfg_attr(coverage_nightly, coverage(off))]
216mod tests {
217 use super::*;
218 use cf_system_sdks::directory::{RegisterInstanceInfo, ServiceEndpoint, ServiceInstanceInfo};
219 use std::sync::atomic::{AtomicUsize, Ordering};
220
221 /// Counts `resolve_rest_service` calls; all other methods are unused here.
222 struct CountingDirectory {
223 calls: AtomicUsize,
224 }
225
226 #[async_trait]
227 impl DirectoryClient for CountingDirectory {
228 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
229 unimplemented!()
230 }
231 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
232 self.calls.fetch_add(1, Ordering::SeqCst);
233 Ok(ServiceEndpoint::new(format!("http://{gear}.local")))
234 }
235 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
236 unimplemented!()
237 }
238 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
239 unimplemented!()
240 }
241 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
242 unimplemented!()
243 }
244 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
245 unimplemented!()
246 }
247 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
248 unimplemented!()
249 }
250 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
251 unimplemented!()
252 }
253 }
254
255 #[tokio::test]
256 async fn memoizes_successful_resolution_within_ttl() {
257 let dir = Arc::new(CountingDirectory {
258 calls: AtomicUsize::new(0),
259 });
260 let resolver = DirectoryEndpointResolver::new(dir.clone());
261
262 for _ in 0..3 {
263 let ep = resolver.resolve_endpoint("billing").await.unwrap();
264 assert_eq!(ep.as_deref(), Some("http://billing.local"));
265 }
266
267 assert_eq!(
268 dir.calls.load(Ordering::SeqCst),
269 1,
270 "successful resolution must be memoized within the TTL (one directory lookup)"
271 );
272 }
273
274 #[tokio::test]
275 async fn static_resolver_always_returns_fixed_endpoint() {
276 let resolver = StaticEndpointResolver::new("http://localhost:8081");
277 assert_eq!(
278 resolver
279 .resolve_endpoint("anything")
280 .await
281 .unwrap()
282 .as_deref(),
283 Some("http://localhost:8081")
284 );
285 assert_eq!(
286 resolver.resolve_endpoint("other").await.unwrap().as_deref(),
287 Some("http://localhost:8081")
288 );
289 }
290
291 #[tokio::test]
292 async fn null_resolver_never_resolves() {
293 assert_eq!(
294 NullEndpointResolver
295 .resolve_endpoint("billing")
296 .await
297 .unwrap(),
298 None
299 );
300 }
301
302 /// Answers every lookup with a caller-chosen failure, so the resolver's
303 /// two error branches can be told apart.
304 struct FailingDirectory {
305 make_error: fn(&str) -> anyhow::Error,
306 }
307
308 #[async_trait]
309 impl DirectoryClient for FailingDirectory {
310 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
311 unimplemented!()
312 }
313 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
314 Err((self.make_error)(gear))
315 }
316 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
317 unimplemented!()
318 }
319 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
320 unimplemented!()
321 }
322 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
323 unimplemented!()
324 }
325 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
326 unimplemented!()
327 }
328 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
329 unimplemented!()
330 }
331 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
332 unimplemented!()
333 }
334 }
335
336 #[tokio::test]
337 async fn not_found_sentinel_resolves_to_ok_none() {
338 // "The provider has not registered yet" is a routine startup condition,
339 // not a failure: it must come back as `Ok(None)` so the readiness probe
340 // logs it at `debug` and the resolving client reports `Unresolved`
341 // rather than warning on every call.
342 let dir = Arc::new(FailingDirectory {
343 make_error: |gear| DirectoryNotFound::new(format!("gear {gear}")).into(),
344 });
345 let resolver = DirectoryEndpointResolver::new(dir);
346
347 assert_eq!(resolver.resolve_endpoint("billing").await.unwrap(), None);
348 }
349
350 #[tokio::test]
351 async fn backend_failure_resolves_to_err() {
352 // Anything that is not the sentinel is a genuine directory outage and
353 // must stay an error, so a stuck 503 has a diagnostic trail.
354 let dir = Arc::new(FailingDirectory {
355 make_error: |_| anyhow::anyhow!("connection refused"),
356 });
357 let resolver = DirectoryEndpointResolver::new(dir);
358
359 assert!(resolver.resolve_endpoint("billing").await.is_err());
360 }
361
362 #[tokio::test]
363 #[tracing_test::traced_test]
364 async fn invalid_argument_resolves_to_err_and_logs_config_error() {
365 // A malformed provider name is a permanent configuration error, handled
366 // by the dedicated arm. Unlike the not-found sentinel (`Ok(None)`), it
367 // must stay an `Err` so the call fails rather than being treated as a
368 // provider that is merely not ready yet. It must also be logged
369 // distinctly at `error` (a static configuration error) so it is not
370 // folded into the generic backend-failure path, which emits no log of
371 // its own at the resolver.
372 let dir = Arc::new(FailingDirectory {
373 make_error: |gear| DirectoryInvalidArgument::new(format!("bad name {gear}")).into(),
374 });
375 let resolver = DirectoryEndpointResolver::new(dir);
376
377 assert!(resolver.resolve_endpoint("Bad Name").await.is_err());
378
379 logs_assert(|lines: &[&str]| {
380 match lines
381 .iter()
382 .find(|line| line.contains("permanently rejected by the directory"))
383 {
384 Some(line) if line.contains("ERROR") => Ok(()),
385 Some(_) => Err("invalid-argument arm logged at the wrong level".to_owned()),
386 None => {
387 Err("invalid-argument arm must emit a distinct config-error log".to_owned())
388 }
389 }
390 });
391 }
392
393 /// Mirrors the body `#[toolkit::consumes]` generates, so the wiring
394 /// contract is exercised without going through `inventory` (which is
395 /// link-time global and cannot hold two registrations built by a test).
396 #[allow(
397 clippy::unnecessary_wraps,
398 reason = "mirrors the fallible signature of ConsumerRegistration::wire so the test \
399 exercises the same shape the macro emits"
400 )]
401 fn generated_wire_body<C>(hub: &ClientHub, proxy: Arc<C>) -> anyhow::Result<WireOutcome>
402 where
403 C: ?Sized + Send + Sync + 'static,
404 {
405 if hub.try_get_local::<C>().is_some() {
406 return Ok(WireOutcome::Local);
407 }
408 if hub.has_remote_proxy::<C>() {
409 return Ok(WireOutcome::Remote);
410 }
411 hub.register_remote_proxy::<C>(proxy);
412 Ok(WireOutcome::Remote)
413 }
414
415 trait Payments: Send + Sync {}
416 struct PaymentsProxy;
417 impl Payments for PaymentsProxy {}
418 struct PaymentsLocal;
419 impl Payments for PaymentsLocal {}
420
421 #[test]
422 fn two_consumers_of_one_contract_both_report_remote() {
423 // Two gears in one process consuming the same contract from the same
424 // provider. Whichever wires second used to find the first one's proxy,
425 // conclude the dependency was local, and get it marked
426 // readiness-resolved — with both registrations sharing a `dep_gear`,
427 // that flipped the very entry the unresolved remote dep was gating on,
428 // and `/readyz` went green with no provider up.
429 let hub = ClientHub::new();
430
431 let first = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
432 let second = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
433
434 assert_eq!(first, WireOutcome::Remote);
435 assert_eq!(
436 second,
437 WireOutcome::Remote,
438 "the second consumer must still gate readiness, not short-circuit to Local"
439 );
440 }
441
442 #[test]
443 fn a_genuine_local_impl_still_wins() {
444 // The short-circuit must keep working for its actual purpose: a
445 // co-located provider registered during init means no HTTP hop and no
446 // readiness gate.
447 let hub = ClientHub::new();
448 hub.register::<dyn Payments>(Arc::new(PaymentsLocal));
449
450 let outcome = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
451
452 assert_eq!(outcome, WireOutcome::Local);
453 }
454
455 #[tokio::test]
456 async fn local_directory_client_not_found_reaches_the_resolver_as_ok_none() {
457 // End-to-end over the real in-process client rather than a stub: this
458 // is the pairing that was broken (the client returned a bare `anyhow`,
459 // so the resolver could never produce `Ok(None)`).
460 let dir = Arc::new(crate::directory::LocalDirectoryClient::new(Arc::new(
461 crate::runtime::GearManager::new(),
462 )));
463 let resolver = DirectoryEndpointResolver::new(dir);
464
465 assert_eq!(
466 resolver.resolve_endpoint("never-registered").await.unwrap(),
467 None
468 );
469 }
470}