Skip to main content

conjure_runtime/
builder.rs

1// Copyright 2020 Palantir Technologies, Inc.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14//! The client builder.
15#[cfg(not(target_arch = "wasm32"))]
16use crate::blocking;
17use crate::client::ClientState;
18use crate::config::{ProxyConfig, SecurityConfig, ServiceConfig};
19use crate::weak_cache::Cached;
20use crate::{Client, HostMetricsRegistry, UserAgent};
21use arc_swap::ArcSwap;
22use conjure_error::Error;
23use conjure_http::client::ConjureRuntime;
24use std::sync::Arc;
25use std::time::Duration;
26#[cfg(not(target_arch = "wasm32"))]
27use tokio::runtime::Handle;
28use url::Url;
29use witchcraft_metrics::MetricRegistry;
30
31/// A builder to construct [`Client`]s and [`blocking::Client`]s.
32pub struct Builder<T = Complete>(T);
33
34/// The service builder stage.
35pub struct ServiceStage(());
36
37/// The user agent builder stage.
38pub struct UserAgentStage {
39    service: String,
40}
41
42#[derive(Clone, PartialEq, Eq, Hash)]
43pub(crate) struct CachedConfig {
44    service: String,
45    user_agent: UserAgent,
46    uris: Vec<Url>,
47    security: SecurityConfig,
48    proxy: ProxyConfig,
49    connect_timeout: Duration,
50    read_timeout: Duration,
51    write_timeout: Duration,
52    backoff_slot_size: Duration,
53    max_num_retries: u32,
54    client_qos: ClientQos,
55    server_qos: ServerQos,
56    service_error: ServiceError,
57    idempotency: Idempotency,
58    node_selection_strategy: NodeSelectionStrategy,
59    override_host_index: Option<usize>,
60}
61
62#[derive(Clone)]
63pub(crate) struct UncachedConfig {
64    pub(crate) metrics: Option<Arc<MetricRegistry>>,
65    pub(crate) host_metrics: Option<Arc<HostMetricsRegistry>>,
66    #[cfg(not(target_arch = "wasm32"))]
67    pub(crate) blocking_handle: Option<Handle>,
68    pub(crate) conjure_runtime: Arc<ConjureRuntime>,
69}
70
71/// The complete builder stage.
72pub struct Complete {
73    cached: CachedConfig,
74    uncached: UncachedConfig,
75}
76
77impl Default for Builder<ServiceStage> {
78    #[inline]
79    fn default() -> Self {
80        Builder::new()
81    }
82}
83
84impl Builder<ServiceStage> {
85    /// Creates a new builder with default settings.
86    #[inline]
87    pub fn new() -> Self {
88        Builder(ServiceStage(()))
89    }
90
91    /// Sets the name of the service this client will communicate with.
92    ///
93    /// This is used in logging and metrics to allow differentiation between different clients.
94    #[inline]
95    pub fn service(self, service: &str) -> Builder<UserAgentStage> {
96        Builder(UserAgentStage {
97            service: service.to_string(),
98        })
99    }
100}
101
102impl Builder<UserAgentStage> {
103    /// Sets the user agent sent by this client.
104    #[inline]
105    pub fn user_agent(self, user_agent: UserAgent) -> Builder {
106        Builder(Complete {
107            cached: CachedConfig {
108                service: self.0.service,
109                user_agent,
110                uris: vec![],
111                security: SecurityConfig::builder().build(),
112                proxy: ProxyConfig::Direct,
113                connect_timeout: Duration::from_secs(10),
114                read_timeout: Duration::from_secs(5 * 60),
115                write_timeout: Duration::from_secs(5 * 60),
116                backoff_slot_size: Duration::from_millis(250),
117                max_num_retries: 4,
118                client_qos: ClientQos::Enabled,
119                server_qos: ServerQos::AutomaticRetry,
120                service_error: ServiceError::WrapInNewError,
121                idempotency: Idempotency::ByMethod,
122                node_selection_strategy: NodeSelectionStrategy::PinUntilError,
123                override_host_index: None,
124            },
125            uncached: UncachedConfig {
126                metrics: None,
127                host_metrics: None,
128                #[cfg(not(target_arch = "wasm32"))]
129                blocking_handle: None,
130                conjure_runtime: Arc::new(ConjureRuntime::new()),
131            },
132        })
133    }
134}
135
136#[cfg(test)]
137impl Builder {
138    pub(crate) fn for_test() -> Self {
139        use crate::Agent;
140
141        Builder::new()
142            .service("test")
143            .user_agent(UserAgent::new(Agent::new("test", "0.0.0")))
144    }
145}
146
147impl Builder<Complete> {
148    pub(crate) fn cached_config(&self) -> &CachedConfig {
149        &self.0.cached
150    }
151
152    /// Applies configuration settings from a `ServiceConfig` to the builder.
153    #[inline]
154    pub fn from_config(mut self, config: &ServiceConfig) -> Self {
155        self = self.uris(config.uris().to_vec());
156
157        if let Some(security) = config.security() {
158            self = self.security(security.clone());
159        }
160
161        if let Some(proxy) = config.proxy() {
162            self = self.proxy(proxy.clone());
163        }
164
165        if let Some(connect_timeout) = config.connect_timeout() {
166            self = self.connect_timeout(connect_timeout);
167        }
168
169        if let Some(read_timeout) = config.read_timeout() {
170            self = self.read_timeout(read_timeout);
171        }
172
173        if let Some(write_timeout) = config.write_timeout() {
174            self = self.write_timeout(write_timeout);
175        }
176
177        if let Some(backoff_slot_size) = config.backoff_slot_size() {
178            self = self.backoff_slot_size(backoff_slot_size);
179        }
180
181        if let Some(max_num_retries) = config.max_num_retries() {
182            self = self.max_num_retries(max_num_retries);
183        }
184
185        self
186    }
187
188    /// Returns the builder's configured service name.
189    #[inline]
190    pub fn get_service(&self) -> &str {
191        &self.0.cached.service
192    }
193
194    /// Returns the builder's configured user agent.
195    #[inline]
196    pub fn get_user_agent(&self) -> &UserAgent {
197        &self.0.cached.user_agent
198    }
199
200    /// Appends a URI to the URIs list.
201    ///
202    /// Defaults to an empty list.
203    #[inline]
204    pub fn uri(mut self, uri: Url) -> Self {
205        self.0.cached.uris.push(uri);
206        self
207    }
208
209    /// Sets the URIs list.
210    ///
211    /// Defaults to an empty list.
212    #[inline]
213    pub fn uris(mut self, uris: Vec<Url>) -> Self {
214        self.0.cached.uris = uris;
215        self
216    }
217
218    /// Returns the builder's configured URIs list.
219    #[inline]
220    pub fn get_uris(&self) -> &[Url] {
221        &self.0.cached.uris
222    }
223
224    /// Sets the security configuration.
225    ///
226    /// Defaults to an empty configuration.
227    #[inline]
228    pub fn security(mut self, security: SecurityConfig) -> Self {
229        self.0.cached.security = security;
230        self
231    }
232
233    /// Returns the builder's configured security configuration.
234    #[inline]
235    pub fn get_security(&self) -> &SecurityConfig {
236        &self.0.cached.security
237    }
238
239    /// Sets the proxy configuration.
240    ///
241    /// Defaults to `ProxyConfig::Direct` (i.e. no proxy).
242    #[inline]
243    pub fn proxy(mut self, proxy: ProxyConfig) -> Self {
244        self.0.cached.proxy = proxy;
245        self
246    }
247
248    /// Returns the builder's configured proxy configuration.
249    #[inline]
250    pub fn get_proxy(&self) -> &ProxyConfig {
251        &self.0.cached.proxy
252    }
253
254    /// Sets the connect timeout.
255    ///
256    /// Defaults to 10 seconds.
257    #[inline]
258    pub fn connect_timeout(mut self, connect_timeout: Duration) -> Self {
259        self.0.cached.connect_timeout = connect_timeout;
260        self
261    }
262
263    /// Returns the builder's configured connect timeout.
264    #[inline]
265    pub fn get_connect_timeout(&self) -> Duration {
266        self.0.cached.connect_timeout
267    }
268
269    /// Sets the read timeout.
270    ///
271    /// This timeout applies to socket-level read attempts.
272    ///
273    /// Defaults to 5 minutes.
274    #[inline]
275    pub fn read_timeout(mut self, read_timeout: Duration) -> Self {
276        self.0.cached.read_timeout = read_timeout;
277        self
278    }
279
280    /// Returns the builder's configured read timeout.
281    #[inline]
282    pub fn get_read_timeout(&self) -> Duration {
283        self.0.cached.read_timeout
284    }
285
286    /// Sets the write timeout.
287    ///
288    /// This timeout applies to socket-level write attempts.
289    ///
290    /// Defaults to 5 minutes.
291    #[inline]
292    pub fn write_timeout(mut self, write_timeout: Duration) -> Self {
293        self.0.cached.write_timeout = write_timeout;
294        self
295    }
296
297    /// Returns the builder's configured write timeout.
298    #[inline]
299    pub fn get_write_timeout(&self) -> Duration {
300        self.0.cached.write_timeout
301    }
302
303    /// Sets the backoff slot size.
304    ///
305    /// This is the upper bound on the initial delay before retrying a request. It grows exponentially as additional
306    /// attempts are made for a given request.
307    ///
308    /// Defaults to 250 milliseconds.
309    #[inline]
310    pub fn backoff_slot_size(mut self, backoff_slot_size: Duration) -> Self {
311        self.0.cached.backoff_slot_size = backoff_slot_size;
312        self
313    }
314
315    /// Returns the builder's configured backoff slot size.
316    #[inline]
317    pub fn get_backoff_slot_size(&self) -> Duration {
318        self.0.cached.backoff_slot_size
319    }
320
321    /// Sets the maximum number of times a request attempt will be retried before giving up.
322    ///
323    /// Defaults to 4.
324    #[inline]
325    pub fn max_num_retries(mut self, max_num_retries: u32) -> Self {
326        self.0.cached.max_num_retries = max_num_retries;
327        self
328    }
329
330    /// Returns the builder's configured maximum number of retries.
331    #[inline]
332    pub fn get_max_num_retries(&self) -> u32 {
333        self.0.cached.max_num_retries
334    }
335
336    /// Sets the client's internal rate limiting behavior.
337    ///
338    /// Defaults to `ClientQos::Enabled`.
339    #[inline]
340    pub fn client_qos(mut self, client_qos: ClientQos) -> Self {
341        self.0.cached.client_qos = client_qos;
342        self
343    }
344
345    /// Returns the builder's configured internal rate limiting behavior.
346    #[inline]
347    pub fn get_client_qos(&self) -> ClientQos {
348        self.0.cached.client_qos
349    }
350
351    /// Sets the client's behavior in response to a QoS error from the server.
352    ///
353    /// Defaults to `ServerQos::AutomaticRetry`.
354    #[inline]
355    pub fn server_qos(mut self, server_qos: ServerQos) -> Self {
356        self.0.cached.server_qos = server_qos;
357        self
358    }
359
360    /// Returns the builder's configured server QoS behavior.
361    #[inline]
362    pub fn get_server_qos(&self) -> ServerQos {
363        self.0.cached.server_qos
364    }
365
366    /// Sets the client's behavior in response to a service error from the server.
367    ///
368    /// Defaults to `ServiceError::WrapInNewError`.
369    #[inline]
370    pub fn service_error(mut self, service_error: ServiceError) -> Self {
371        self.0.cached.service_error = service_error;
372        self
373    }
374
375    /// Returns the builder's configured service error handling behavior.
376    #[inline]
377    pub fn get_service_error(&self) -> ServiceError {
378        self.0.cached.service_error
379    }
380
381    /// Sets the client's behavior to determine if a request is idempotent or not.
382    ///
383    /// Only idempotent requests will be retried.
384    ///
385    /// Defaults to `Idempotency::ByMethod`.
386    #[inline]
387    pub fn idempotency(mut self, idempotency: Idempotency) -> Self {
388        self.0.cached.idempotency = idempotency;
389        self
390    }
391
392    /// Returns the builder's configured idempotency handling behavior.
393    #[inline]
394    pub fn get_idempotency(&self) -> Idempotency {
395        self.0.cached.idempotency
396    }
397
398    /// Sets the client's strategy for selecting a node for a request.
399    ///
400    /// Defaults to `NodeSelectionStrategy::PinUntilError`.
401    #[inline]
402    pub fn node_selection_strategy(
403        mut self,
404        node_selection_strategy: NodeSelectionStrategy,
405    ) -> Self {
406        self.0.cached.node_selection_strategy = node_selection_strategy;
407        self
408    }
409
410    /// Returns the builder's configured node selection strategy.
411    #[inline]
412    pub fn get_node_selection_strategy(&self) -> NodeSelectionStrategy {
413        self.0.cached.node_selection_strategy
414    }
415
416    /// Sets the metric registry used to register client metrics.
417    ///
418    /// Defaults to no registry.
419    #[inline]
420    pub fn metrics(mut self, metrics: Arc<MetricRegistry>) -> Self {
421        self.0.uncached.metrics = Some(metrics);
422        self
423    }
424
425    /// Returns the builder's configured metric registry.
426    #[inline]
427    pub fn get_metrics(&self) -> Option<&Arc<MetricRegistry>> {
428        self.0.uncached.metrics.as_ref()
429    }
430
431    /// Sets the host metrics registry used to track host performance.
432    ///
433    /// Defaults to no registry.
434    #[inline]
435    pub fn host_metrics(mut self, host_metrics: Arc<HostMetricsRegistry>) -> Self {
436        self.0.uncached.host_metrics = Some(host_metrics);
437        self
438    }
439
440    /// Returns the builder's configured host metrics registry.
441    #[inline]
442    pub fn get_host_metrics(&self) -> Option<&Arc<HostMetricsRegistry>> {
443        self.0.uncached.host_metrics.as_ref()
444    }
445
446    /// Sets the Conjure runtime used to configure request and response encodings.
447    ///
448    /// Defaults to `ConjureRuntime::default()`.
449    #[inline]
450    pub fn conjure_runtime(mut self, conjure_runtime: Arc<ConjureRuntime>) -> Self {
451        self.0.uncached.conjure_runtime = conjure_runtime;
452        self
453    }
454
455    /// Returns the configured Conjure runtime.
456    pub fn get_conjure_runtime(&self) -> &Arc<ConjureRuntime> {
457        &self.0.uncached.conjure_runtime
458    }
459
460    /// Overrides the `hostIndex` field included in metrics.
461    #[inline]
462    pub fn override_host_index(mut self, override_host_index: usize) -> Self {
463        self.0.cached.override_host_index = Some(override_host_index);
464        self
465    }
466
467    /// Returns the builder's `hostIndex` override.
468    #[inline]
469    pub fn get_override_host_index(&self) -> Option<usize> {
470        self.0.cached.override_host_index
471    }
472
473    /// Creates a new `Client`.
474    pub fn build(&self) -> Result<Client, Error> {
475        let state = ClientState::new(self)?;
476        Ok(Client::new(
477            Arc::new(ArcSwap::new(Arc::new(Cached::uncached(state)))),
478            None,
479        ))
480    }
481}
482
483#[cfg(not(target_arch = "wasm32"))]
484impl Builder<Complete> {
485    /// Returns the `Handle` to the tokio `Runtime` to be used by blocking clients.
486    ///
487    /// This has no effect on async clients.
488    ///
489    /// Defaults to a `conjure-runtime` internal `Runtime`.
490    #[inline]
491    pub fn blocking_handle(mut self, blocking_handle: Handle) -> Self {
492        self.0.uncached.blocking_handle = Some(blocking_handle);
493        self
494    }
495
496    /// Returns the builder's configured blocking handle.
497    #[inline]
498    pub fn get_blocking_handle(&self) -> Option<&Handle> {
499        self.0.uncached.blocking_handle.as_ref()
500    }
501
502    /// Creates a new `blocking::Client`.
503    pub fn build_blocking(&self) -> Result<blocking::Client, Error> {
504        self.build().map(|client| blocking::Client {
505            client,
506            handle: self.0.uncached.blocking_handle.clone(),
507        })
508    }
509}
510
511/// Specifies the beahavior of client-side sympathetic rate limiting.
512#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
513#[non_exhaustive]
514pub enum ClientQos {
515    /// Enable client side rate limiting.
516    ///
517    /// This is the default behavior.
518    Enabled,
519
520    /// Disables client-side rate limiting.
521    ///
522    /// This should only be used when there are known issues with the interaction between a service's rate limiting
523    /// implementation and the client's.
524    DangerousDisableSympatheticClientQos,
525}
526
527/// Specifies the behavior of a client in response to a `QoS` error from a server.
528///
529/// QoS errors have status codes 429 or 503.
530#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
531#[non_exhaustive]
532pub enum ServerQos {
533    /// The client will automatically retry the request when possible in response to a QoS error.
534    ///
535    /// This is the default behavior.
536    AutomaticRetry,
537
538    /// The client will transparently propagate the QoS error without retrying.
539    ///
540    /// This is designed for use when an upstream service has better context on how to handle a QoS error. Propagating
541    /// the error upstream to that service without retrying allows it to handle retry logic internally.
542    Propagate429And503ToCaller,
543}
544
545/// Specifies the behavior of the client in response to a service error from a server.
546///
547/// Service errors are encoded as responses with a 4xx or 5xx response code and a body containing a `SerializableError`.
548#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
549#[non_exhaustive]
550pub enum ServiceError {
551    /// The service error will be propagated as a new internal service error.
552    ///
553    /// The error's cause will contain the information about the received service error, but the error constructed by
554    /// the client will have a different error instance ID, type, etc.
555    ///
556    /// This is the default behavior.
557    WrapInNewError,
558
559    /// The service error will be transparently propagated without change.
560    ///
561    /// This is designed for use when proxying a request to another node, commonly of the same service. By preserving
562    /// the original error's instance ID, type, etc, the upstream service will be able to process the error properly.
563    PropagateToCaller,
564}
565
566/// Specifies the manner in which the client decides if a request is idempotent or not.
567#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
568#[non_exhaustive]
569pub enum Idempotency {
570    /// All requests are assumed to be idempotent.
571    Always,
572
573    /// Only requests with HTTP methods defined as idempotent (GET, HEAD, OPTIONS, TRACE, PUT, and DELETE) are assumed
574    /// to be idempotent.
575    ///
576    /// This is the default behavior.
577    ByMethod,
578
579    /// No requests are assumed to be idempotent.
580    Never,
581}
582
583/// Specifies the strategy used to select a node of a service to use for a request attempt.
584#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
585#[non_exhaustive]
586pub enum NodeSelectionStrategy {
587    /// Pin to a single host as long as it continues to successfully respond to requests.
588    ///
589    /// If the pinned node fails to successfully respond, the client will rotate through the other nodes until it finds
590    /// one that can successfully respond and then pin to that new node. The pinned node will also be randomly rotated
591    /// periodically to help spread load across the cluster.
592    ///
593    /// This is the default behavior.
594    PinUntilError,
595
596    /// Like `PinUntilError` except that the pinned node is never randomly shuffled.
597    PinUntilErrorWithoutReshuffle,
598
599    /// For each new request, select the "next" node (in some unspecified order).
600    Balanced,
601}