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