Skip to main content

toolkit_contract/runtime/
resolving.rs

1//! Directory-resolving REST transport.
2//!
3//! [`DirectoryResolvingClient`] is a thin, self-healing wrapper around a
4//! generated `<Trait>RestClient`. Instead of binding to a single static base
5//! URL at construction time, it resolves the provider's endpoint from the
6//! service directory on every call and rebuilds the underlying client only
7//! when the resolved endpoint changes.
8//!
9//! This is what makes the consumer side tolerant of *eventual readiness* and
10//! runtime churn:
11//! - **Not ready yet** — the provider hasn't registered → resolution yields
12//!   `None` → calls fail with [`TransportError::Unresolved`] (a transient
13//!   error that maps to `service_unavailable`), never a panic.
14//! - **Provider moved / pod replaced** — the directory returns a new endpoint
15//!   → the wrapper rebuilds the client against it on the next call.
16//! - **Provider vanished** — all instances evicted → resolution yields `None`
17//!   again → `Unresolved`; the wrapper recovers automatically once a live
18//!   instance reappears.
19//!
20//! The `Arc<dyn Trait>` registered in the `ClientHub` is the long-lived
21//! resolving wrapper, so the hub entry is wired once and never replaced; all
22//! churn handling lives inside this transport.
23
24use std::sync::{Arc, RwLock};
25
26use async_trait::async_trait;
27use tracing::{debug, warn};
28
29use crate::runtime::config::ClientConfig;
30use crate::runtime::transport_error::TransportError;
31use crate::wiring::ClientTuning;
32
33/// Directory-lookup failure (the directory backend itself could not answer —
34/// e.g. the gRPC directory is unreachable), as opposed to a successful lookup
35/// that found no live instance (`Ok(None)`). Keeping the two distinct lets
36/// callers and observability tell a not-ready provider apart from a directory
37/// outage.
38#[derive(Debug, thiserror::Error)]
39#[error("directory lookup for gear `{gear}` failed: {source}")]
40pub struct ResolveError {
41    /// Gear name whose lookup failed.
42    pub gear: String,
43    /// Underlying directory transport/backend error.
44    #[source]
45    pub source: Box<dyn std::error::Error + Send + Sync + 'static>,
46}
47
48impl ResolveError {
49    /// Build a [`ResolveError`] from any boxable source error.
50    pub fn new<E>(gear: impl Into<String>, source: E) -> Self
51    where
52        E: Into<Box<dyn std::error::Error + Send + Sync + 'static>>,
53    {
54        Self {
55            gear: gear.into(),
56            source: source.into(),
57        }
58    }
59}
60
61/// Resolves a logical gear name to a live base endpoint URI.
62///
63/// Kept as a minimal trait here so `toolkit-contract` does not depend on the
64/// service-directory SDK. The host layer (`toolkit`) provides an adapter over
65/// its `DirectoryClient`; tests can implement it directly.
66///
67/// Return contract:
68/// - `Ok(Some(uri))` — a live instance was resolved.
69/// - `Ok(None)` — the lookup succeeded but no live instance is registered
70///   (provider not ready yet, or every instance evicted).
71/// - `Err(ResolveError)` — the directory backend itself failed to answer.
72///
73/// Both `Ok(None)` and `Err(..)` surface to the caller as a retryable
74/// [`TransportError::Unresolved`]; the distinction exists for logging/metrics.
75///
76/// [`DirectoryResolvingClient`] calls this **on every request** so provider
77/// churn is observed immediately. Implementations whose lookup is expensive
78/// (e.g. an out-of-process gRPC directory) SHOULD memoize internally (short
79/// TTL) rather than performing a network round-trip per call.
80#[async_trait]
81pub trait EndpointResolver: Send + Sync {
82    /// Resolve `gear` to a base endpoint URI (e.g. `http://billing:8080`).
83    ///
84    /// # Errors
85    /// Returns [`ResolveError`] only when the directory backend could not be
86    /// queried; a successful query with no live instance is `Ok(None)`.
87    async fn resolve_endpoint(&self, gear: &str) -> Result<Option<String>, ResolveError>;
88}
89
90type ClientBuilder<C> = dyn Fn(ClientConfig) -> Result<C, TransportError> + Send + Sync;
91
92/// Self-healing REST transport that resolves its endpoint from the directory.
93///
94/// Generic over the concrete generated client `C` (e.g. `PaymentApiRestClient`).
95/// The generated `<Trait>ResolvingClient` wraps one of these and delegates each
96/// trait method through [`DirectoryResolvingClient::resolved`].
97pub struct DirectoryResolvingClient<C> {
98    resolver: Arc<dyn EndpointResolver>,
99    from_gear: String,
100    tuning: ClientTuning,
101    build: Box<ClientBuilder<C>>,
102    /// Cached `(endpoint, client)` reused while the resolved endpoint is stable.
103    cache: RwLock<Option<(String, Arc<C>)>>,
104    /// Serializes the build step. Without this, N callers that concurrently
105    /// observe the same stale/absent cache entry (first call, or immediately
106    /// after an endpoint change) would each construct a redundant client —
107    /// only the guard is held, never across an `.await` (the directory lookup
108    /// has already completed by the time this lock is taken; `build` itself
109    /// is synchronous).
110    build_lock: std::sync::Mutex<()>,
111}
112
113impl<C: Send + Sync + 'static> DirectoryResolvingClient<C> {
114    /// Construct a resolving client.
115    ///
116    /// - `resolver` — resolves the provider gear name to a live endpoint.
117    /// - `from_gear` — the logical provider gear name to resolve.
118    /// - `tuning` — timeout/retry/reconnect overrides applied to each built client.
119    /// - `build` — constructs the concrete client `C` from a [`ClientConfig`]
120    ///   (the generated `C::new(config)`, mapping its error into
121    ///   [`TransportError`]).
122    pub fn new(
123        resolver: Arc<dyn EndpointResolver>,
124        from_gear: impl Into<String>,
125        tuning: ClientTuning,
126        build: impl Fn(ClientConfig) -> Result<C, TransportError> + Send + Sync + 'static,
127    ) -> Self {
128        Self {
129            resolver,
130            from_gear: from_gear.into(),
131            tuning,
132            build: Box::new(build),
133            cache: RwLock::new(None),
134            build_lock: std::sync::Mutex::new(()),
135        }
136    }
137
138    /// The logical provider gear name this client resolves against.
139    #[must_use]
140    pub fn from_gear(&self) -> &str {
141        &self.from_gear
142    }
143
144    /// Resolve (or reuse) the underlying client for the current endpoint.
145    ///
146    /// Re-resolves the directory on every call so a moved or vanished provider
147    /// is observed immediately; the built client is cached and reused while the
148    /// endpoint is unchanged. Returns [`TransportError::Unresolved`] when no
149    /// live instance is registered.
150    ///
151    /// # Errors
152    /// - [`TransportError::Unresolved`] if the directory has no live endpoint.
153    /// - Whatever `build` returns if constructing the client fails.
154    pub async fn resolved(&self) -> Result<Arc<C>, TransportError> {
155        let uri = match self.resolver.resolve_endpoint(&self.from_gear).await {
156            Ok(Some(uri)) => uri,
157            Ok(None) => {
158                debug!(gear = %self.from_gear, "no live instance registered; provider not ready");
159                self.invalidate();
160                return Err(TransportError::unresolved(&self.from_gear));
161            }
162            Err(err) => {
163                warn!(gear = %self.from_gear, error = %err, "directory lookup failed");
164                self.invalidate();
165                return Err(TransportError::unresolved(&self.from_gear));
166            }
167        };
168
169        // Fast path: endpoint unchanged → reuse the cached client.
170        if let Some(client) = self.cached_for(&uri) {
171            return Ok(client);
172        }
173
174        // First call, or the endpoint changed: build and cache. Single-flight
175        // via `build_lock` so a thundering herd of concurrent callers that all
176        // missed the fast path above builds the client ONCE, not N times.
177        let _build_guard = self
178            .build_lock
179            .lock()
180            .unwrap_or_else(std::sync::PoisonError::into_inner);
181        // Re-check under the lock: another caller may have just finished
182        // building for this exact endpoint while we were waiting for it.
183        if let Some(client) = self.cached_for(&uri) {
184            return Ok(client);
185        }
186        debug!(gear = %self.from_gear, endpoint = %uri, "building REST client for resolved endpoint");
187        let cfg = self.tuning.apply_to(&uri);
188        let client = Arc::new((self.build)(cfg)?);
189        if let Ok(mut w) = self.cache.write() {
190            *w = Some((uri, Arc::clone(&client)));
191        }
192        Ok(client)
193    }
194
195    /// Returns the cached client if it is present and was built for `uri`.
196    fn cached_for(&self, uri: &str) -> Option<Arc<C>> {
197        let r = self.cache.read().ok()?;
198        let (cached_uri, client) = r.as_ref()?;
199        (cached_uri == uri).then(|| Arc::clone(client))
200    }
201
202    /// Drop any cached client so a now-absent or moved provider isn't masked
203    /// on the next call.
204    fn invalidate(&self) {
205        if let Ok(mut w) = self.cache.write() {
206            *w = None;
207        }
208    }
209}
210
211#[cfg(test)]
212#[cfg_attr(coverage_nightly, coverage(off))]
213mod tests {
214    use super::*;
215    use std::sync::atomic::{AtomicUsize, Ordering};
216
217    /// One scripted resolver outcome.
218    #[derive(Clone)]
219    enum Step {
220        /// `Ok(Some(uri))`.
221        Found(String),
222        /// `Ok(None)` — looked up, nothing registered.
223        Empty,
224        /// `Err(..)` — directory backend failure.
225        Fail,
226    }
227
228    /// Resolver returning a scripted sequence of outcomes.
229    struct ScriptResolver {
230        steps: Vec<Step>,
231        idx: AtomicUsize,
232    }
233
234    #[async_trait]
235    impl EndpointResolver for ScriptResolver {
236        async fn resolve_endpoint(&self, gear: &str) -> Result<Option<String>, ResolveError> {
237            let i = self.idx.fetch_add(1, Ordering::SeqCst);
238            match self.steps.get(i).cloned() {
239                Some(Step::Found(uri)) => Ok(Some(uri)),
240                Some(Step::Empty) | None => Ok(None),
241                Some(Step::Fail) => Err(ResolveError::new(gear, "directory down")),
242            }
243        }
244    }
245
246    /// Dummy "client" carrying the endpoint it was built against.
247    #[derive(Debug)]
248    struct DummyClient {
249        base_url: String,
250    }
251
252    /// Build a resolving client with a **per-instance** build counter (returned
253    /// alongside) so parallel tests don't race on a shared global.
254    fn resolving(steps: Vec<Step>) -> (DirectoryResolvingClient<DummyClient>, Arc<AtomicUsize>) {
255        let builds = Arc::new(AtomicUsize::new(0));
256        let builds_for_closure = Arc::clone(&builds);
257        let client = DirectoryResolvingClient::new(
258            Arc::new(ScriptResolver {
259                steps,
260                idx: AtomicUsize::new(0),
261            }),
262            "billing",
263            ClientTuning::default(),
264            move |cfg| {
265                builds_for_closure.fetch_add(1, Ordering::SeqCst);
266                Ok(DummyClient {
267                    base_url: cfg.base_url,
268                })
269            },
270        );
271        (client, builds)
272    }
273
274    #[tokio::test]
275    async fn unresolved_when_no_endpoint() {
276        let (c, _builds) = resolving(vec![Step::Empty]);
277        let err = c.resolved().await.unwrap_err();
278        assert!(matches!(err, TransportError::Unresolved { .. }));
279    }
280
281    #[tokio::test]
282    async fn directory_failure_surfaces_as_unresolved() {
283        let (c, builds) = resolving(vec![Step::Fail]);
284        let err = c.resolved().await.unwrap_err();
285        assert!(matches!(err, TransportError::Unresolved { .. }));
286        assert_eq!(
287            builds.load(Ordering::SeqCst),
288            0,
289            "must not build a client on directory failure"
290        );
291    }
292
293    #[tokio::test]
294    async fn caches_while_endpoint_stable_and_rebuilds_on_change() {
295        let (c, builds) = resolving(vec![
296            Step::Found("http://a:8080".into()),
297            Step::Found("http://a:8080".into()),
298            Step::Found("http://b:9090".into()),
299        ]);
300
301        let c1 = c.resolved().await.unwrap();
302        assert_eq!(c1.base_url, "http://a:8080");
303        let c2 = c.resolved().await.unwrap();
304        assert_eq!(c2.base_url, "http://a:8080");
305        assert_eq!(
306            builds.load(Ordering::SeqCst),
307            1,
308            "stable endpoint reuses client"
309        );
310
311        let c3 = c.resolved().await.unwrap();
312        assert_eq!(c3.base_url, "http://b:9090");
313        assert_eq!(
314            builds.load(Ordering::SeqCst),
315            2,
316            "endpoint change rebuilds client"
317        );
318    }
319
320    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
321    async fn concurrent_resolves_build_client_only_once() {
322        // Real OS-thread concurrency (not single-threaded cooperative
323        // interleaving): N callers all resolving to the same first-seen
324        // endpoint race to build. Without single-flight (`build_lock`), each
325        // one that misses the cache before any write lands would construct
326        // its own redundant client.
327        const N: usize = 16;
328        let (c, builds) = resolving(vec![Step::Found("http://a:8080".into()); N]);
329        let c = Arc::new(c);
330
331        let handles: Vec<_> = (0..N)
332            .map(|_| {
333                let c = Arc::clone(&c);
334                tokio::spawn(async move { c.resolved().await.unwrap() })
335            })
336            .collect();
337
338        for h in handles {
339            let client = h.await.unwrap();
340            assert_eq!(client.base_url, "http://a:8080");
341        }
342        assert_eq!(
343            builds.load(Ordering::SeqCst),
344            1,
345            "single-flight: exactly one build for {N} concurrent resolves to the same endpoint"
346        );
347    }
348
349    #[tokio::test]
350    async fn recovers_after_provider_vanishes_and_returns() {
351        let (c, _builds) = resolving(vec![
352            Step::Found("http://a:8080".into()),
353            Step::Empty,
354            Step::Found("http://c:7070".into()),
355        ]);
356
357        assert_eq!(c.resolved().await.unwrap().base_url, "http://a:8080");
358        // Provider vanished → Unresolved, stale cache dropped.
359        assert!(matches!(
360            c.resolved().await.unwrap_err(),
361            TransportError::Unresolved { .. }
362        ));
363        // New instance appears → resolves again, self-healed.
364        assert_eq!(c.resolved().await.unwrap().base_url, "http://c:7070");
365    }
366}