Skip to main content

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}