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::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}