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. See
182 /// [`WireFn`] for the argument roles.
183 pub wire: WireFn,
184}
185
186/// The wiring function pointer emitted by `#[toolkit::consumes]`. Registers the
187/// client into the `ClientHub`, given the endpoint resolver and the process's
188/// platform-plane credential source (threaded onto a directory-resolving remote
189/// client). Aliased to keep the fn-pointer type readable.
190pub type WireFn = fn(
191 &ClientHub,
192 Arc<dyn EndpointResolver>,
193 Option<&toolkit_contract::runtime::config::InternalTokenProvider>,
194) -> anyhow::Result<WireOutcome>;
195
196inventory::collect!(ConsumerRegistration);
197
198#[cfg(test)]
199#[cfg_attr(coverage_nightly, coverage(off))]
200mod tests {
201 use super::*;
202 use cf_system_sdks::directory::{RegisterInstanceInfo, ServiceEndpoint, ServiceInstanceInfo};
203 use std::sync::atomic::{AtomicUsize, Ordering};
204
205 /// Counts `resolve_rest_service` calls; all other methods are unused here.
206 struct CountingDirectory {
207 calls: AtomicUsize,
208 }
209
210 #[async_trait]
211 impl DirectoryClient for CountingDirectory {
212 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
213 unimplemented!()
214 }
215 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
216 self.calls.fetch_add(1, Ordering::SeqCst);
217 Ok(ServiceEndpoint::new(format!("http://{gear}.local")))
218 }
219 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
220 unimplemented!()
221 }
222 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
223 unimplemented!()
224 }
225 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
226 unimplemented!()
227 }
228 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
229 unimplemented!()
230 }
231 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
232 unimplemented!()
233 }
234 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
235 unimplemented!()
236 }
237 }
238
239 #[tokio::test]
240 async fn memoizes_successful_resolution_within_ttl() {
241 let dir = Arc::new(CountingDirectory {
242 calls: AtomicUsize::new(0),
243 });
244 let resolver = DirectoryEndpointResolver::new(dir.clone());
245
246 for _ in 0..3 {
247 let ep = resolver.resolve_endpoint("billing").await.unwrap();
248 assert_eq!(ep.as_deref(), Some("http://billing.local"));
249 }
250
251 assert_eq!(
252 dir.calls.load(Ordering::SeqCst),
253 1,
254 "successful resolution must be memoized within the TTL (one directory lookup)"
255 );
256 }
257
258 #[tokio::test]
259 async fn static_resolver_always_returns_fixed_endpoint() {
260 let resolver = StaticEndpointResolver::new("http://localhost:8081");
261 assert_eq!(
262 resolver
263 .resolve_endpoint("anything")
264 .await
265 .unwrap()
266 .as_deref(),
267 Some("http://localhost:8081")
268 );
269 assert_eq!(
270 resolver.resolve_endpoint("other").await.unwrap().as_deref(),
271 Some("http://localhost:8081")
272 );
273 }
274
275 #[tokio::test]
276 async fn null_resolver_never_resolves() {
277 assert_eq!(
278 NullEndpointResolver
279 .resolve_endpoint("billing")
280 .await
281 .unwrap(),
282 None
283 );
284 }
285
286 /// Answers every lookup with a caller-chosen failure, so the resolver's
287 /// two error branches can be told apart.
288 struct FailingDirectory {
289 make_error: fn(&str) -> anyhow::Error,
290 }
291
292 #[async_trait]
293 impl DirectoryClient for FailingDirectory {
294 async fn resolve_grpc_service(&self, _: &str) -> anyhow::Result<ServiceEndpoint> {
295 unimplemented!()
296 }
297 async fn resolve_rest_service(&self, gear: &str) -> anyhow::Result<ServiceEndpoint> {
298 Err((self.make_error)(gear))
299 }
300 async fn get_openapi_spec(&self, _: &str) -> anyhow::Result<String> {
301 unimplemented!()
302 }
303 async fn list_instances(&self, _: &str) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
304 unimplemented!()
305 }
306 async fn list_all_instances(&self) -> anyhow::Result<Vec<ServiceInstanceInfo>> {
307 unimplemented!()
308 }
309 async fn register_instance(&self, _: RegisterInstanceInfo) -> anyhow::Result<()> {
310 unimplemented!()
311 }
312 async fn deregister_instance(&self, _: &str, _: &str) -> anyhow::Result<()> {
313 unimplemented!()
314 }
315 async fn send_heartbeat(&self, _: &str, _: &str) -> anyhow::Result<()> {
316 unimplemented!()
317 }
318 }
319
320 #[tokio::test]
321 async fn not_found_sentinel_resolves_to_ok_none() {
322 // "The provider has not registered yet" is a routine startup condition,
323 // not a failure: it must come back as `Ok(None)` so the readiness probe
324 // logs it at `debug` and the resolving client reports `Unresolved`
325 // rather than warning on every call.
326 let dir = Arc::new(FailingDirectory {
327 make_error: |gear| DirectoryNotFound::new(format!("gear {gear}")).into(),
328 });
329 let resolver = DirectoryEndpointResolver::new(dir);
330
331 assert_eq!(resolver.resolve_endpoint("billing").await.unwrap(), None);
332 }
333
334 #[tokio::test]
335 async fn backend_failure_resolves_to_err() {
336 // Anything that is not the sentinel is a genuine directory outage and
337 // must stay an error, so a stuck 503 has a diagnostic trail.
338 let dir = Arc::new(FailingDirectory {
339 make_error: |_| anyhow::anyhow!("connection refused"),
340 });
341 let resolver = DirectoryEndpointResolver::new(dir);
342
343 assert!(resolver.resolve_endpoint("billing").await.is_err());
344 }
345
346 /// Mirrors the body `#[toolkit::consumes]` generates, so the wiring
347 /// contract is exercised without going through `inventory` (which is
348 /// link-time global and cannot hold two registrations built by a test).
349 #[allow(
350 clippy::unnecessary_wraps,
351 reason = "mirrors the fallible signature of ConsumerRegistration::wire so the test \
352 exercises the same shape the macro emits"
353 )]
354 fn generated_wire_body<C>(hub: &ClientHub, proxy: Arc<C>) -> anyhow::Result<WireOutcome>
355 where
356 C: ?Sized + Send + Sync + 'static,
357 {
358 if hub.try_get_local::<C>().is_some() {
359 return Ok(WireOutcome::Local);
360 }
361 if hub.has_remote_proxy::<C>() {
362 return Ok(WireOutcome::Remote);
363 }
364 hub.register_remote_proxy::<C>(proxy);
365 Ok(WireOutcome::Remote)
366 }
367
368 trait Payments: Send + Sync {}
369 struct PaymentsProxy;
370 impl Payments for PaymentsProxy {}
371 struct PaymentsLocal;
372 impl Payments for PaymentsLocal {}
373
374 #[test]
375 fn two_consumers_of_one_contract_both_report_remote() {
376 // Two gears in one process consuming the same contract from the same
377 // provider. Whichever wires second used to find the first one's proxy,
378 // conclude the dependency was local, and get it marked
379 // readiness-resolved — with both registrations sharing a `dep_gear`,
380 // that flipped the very entry the unresolved remote dep was gating on,
381 // and `/readyz` went green with no provider up.
382 let hub = ClientHub::new();
383
384 let first = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
385 let second = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
386
387 assert_eq!(first, WireOutcome::Remote);
388 assert_eq!(
389 second,
390 WireOutcome::Remote,
391 "the second consumer must still gate readiness, not short-circuit to Local"
392 );
393 }
394
395 #[test]
396 fn a_genuine_local_impl_still_wins() {
397 // The short-circuit must keep working for its actual purpose: a
398 // co-located provider registered during init means no HTTP hop and no
399 // readiness gate.
400 let hub = ClientHub::new();
401 hub.register::<dyn Payments>(Arc::new(PaymentsLocal));
402
403 let outcome = generated_wire_body::<dyn Payments>(&hub, Arc::new(PaymentsProxy)).unwrap();
404
405 assert_eq!(outcome, WireOutcome::Local);
406 }
407
408 #[tokio::test]
409 async fn local_directory_client_not_found_reaches_the_resolver_as_ok_none() {
410 // End-to-end over the real in-process client rather than a stub: this
411 // is the pairing that was broken (the client returned a bare `anyhow`,
412 // so the resolver could never produce `Ok(None)`).
413 let dir = Arc::new(crate::directory::LocalDirectoryClient::new(Arc::new(
414 crate::runtime::GearManager::new(),
415 )));
416 let resolver = DirectoryEndpointResolver::new(dir);
417
418 assert_eq!(
419 resolver.resolve_endpoint("never-registered").await.unwrap(),
420 None
421 );
422 }
423}