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}