Skip to main content

typesafe_ai_rs/
blocking.rs

1//! Synchronous client, enabled by the `blocking` Cargo feature.
2//!
3//! Use outside async runtimes, or inside `tokio::task::spawn_blocking`.
4use crate::{
5    config::{config_methods, Config, ConfigBuilder},
6    transport, Error, ListModelsResponse, RawResponse, RequestOptions, SystemOneRequest,
7    SystemOneResponse,
8};
9use reqwest::Method;
10use serde_json::Value;
11use std::time::Instant;
12
13/// Synchronous TypeSafe client with pooled HTTP connections.
14#[derive(Clone)]
15pub struct Client {
16    config: Config,
17    http: reqwest::blocking::Client,
18}
19
20/// Configure a blocking client using the same settings as the async client.
21#[derive(Default)]
22pub struct ClientBuilder {
23    config: ConfigBuilder,
24    http: Option<reqwest::blocking::Client>,
25}
26
27impl std::fmt::Debug for Client {
28    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
29        f.debug_struct("Client")
30            .field("config", &self.config)
31            .finish_non_exhaustive()
32    }
33}
34
35impl std::fmt::Debug for ClientBuilder {
36    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
37        f.debug_struct("ClientBuilder").finish_non_exhaustive()
38    }
39}
40
41impl ClientBuilder {
42    config_methods!();
43    /// Supply an HTTP client configured with custom proxy, TLS, or pool settings.
44    pub fn http_client(mut self, value: reqwest::blocking::Client) -> Self {
45        self.http = Some(value);
46        self
47    }
48    /// Build a client. Call outside an async runtime.
49    pub fn build(self) -> Result<Client, Error> {
50        let config = self.config.resolve()?;
51        let http = match self.http {
52            Some(http) => http,
53            None => reqwest::blocking::Client::builder()
54                .redirect(reqwest::redirect::Policy::none())
55                .build()
56                .map_err(Error::Connection)?,
57        };
58        Ok(Client { config, http })
59    }
60}
61impl Client {
62    /// Construct from `TYPESAFE_*` environment variables.
63    pub fn new() -> Result<Self, Error> {
64        Self::builder().build()
65    }
66    /// Configure a blocking client.
67    pub fn builder() -> ClientBuilder {
68        ClientBuilder::default()
69    }
70    /// Resolved default model.
71    pub fn default_model(&self) -> &str {
72        &self.config.model
73    }
74    /// Resolved API root.
75    pub fn base_url(&self) -> &str {
76        &self.config.base_url
77    }
78    /// Evaluate named questions about text or structured state.
79    pub fn system_one(&self, request: SystemOneRequest) -> Result<SystemOneResponse, Error> {
80        self.system_one_with_options(request, RequestOptions::default())
81    }
82    /// Evaluate questions with per-call overrides.
83    pub fn system_one_with_options(
84        &self,
85        request: SystemOneRequest,
86        options: RequestOptions,
87    ) -> Result<SystemOneResponse, Error> {
88        let body = request.prepare(&self.config.model)?;
89        self.send(Method::POST, "/v1/systemone", Some(body), options, |raw| {
90            SystemOneResponse::from_raw_with_log_level(raw, self.config.log_level)
91        })
92    }
93    /// Access model discovery.
94    pub fn models(&self) -> Models<'_> {
95        Models(self)
96    }
97
98    fn send<T>(
99        &self,
100        method: Method,
101        path: &str,
102        body: Option<Value>,
103        options: RequestOptions,
104        decode: impl Fn(RawResponse) -> Result<T, Error>,
105    ) -> Result<T, Error> {
106        if options.cancellation_token.is_some() {
107            return Err(Error::Configuration(
108                "Cancellation requires the async client".into(),
109            ));
110        }
111        let request = transport::prepare(&self.config, method, path, body, &options)?;
112        let started = Instant::now();
113        let mut attempt = 0;
114        loop {
115            let headers = request.headers(attempt);
116            let attempt_started = Instant::now();
117            request.log_request(&self.config, &headers, attempt);
118            let mut builder = self
119                .http
120                .request(request.method.clone(), &request.url)
121                .headers(headers)
122                .timeout(request.timeout);
123            if let Some(body) = &request.body {
124                builder = builder.body(body.clone());
125            }
126            let result = (|| {
127                let response = builder.send().map_err(|e| request.map_error(e))?;
128                let status = response.status();
129                let headers = response.headers().clone();
130                let body = response.bytes().map_err(|e| request.map_error(e))?;
131                let raw = RawResponse {
132                    status,
133                    headers,
134                    body,
135                };
136                request.log_response(&self.config, &raw, attempt_started);
137                request.finish(raw, &decode)
138            })();
139            match result {
140                Ok(response) => return Ok(response),
141                Err(error) => {
142                    let Some(delay) = request.retry_delay(&error, attempt, started) else {
143                        return Err(error);
144                    };
145                    if self.config.log_level >= log::LevelFilter::Info {
146                        log::info!(target: "typesafe_ai_rs", "retry={} delay_ms={}", attempt + 1, delay.as_millis());
147                    }
148                    std::thread::sleep(delay);
149                    attempt += 1;
150                }
151            }
152        }
153    }
154}
155
156/// Synchronous model discovery resource.
157#[derive(Clone, Copy, Debug)]
158pub struct Models<'a>(&'a Client);
159impl Models<'_> {
160    /// List model cards and response metadata.
161    pub fn list(&self) -> Result<ListModelsResponse, Error> {
162        self.list_with_options(RequestOptions::default())
163    }
164    /// List models with per-call overrides.
165    pub fn list_with_options(&self, options: RequestOptions) -> Result<ListModelsResponse, Error> {
166        self.0.send(
167            Method::GET,
168            "/v1/models",
169            None,
170            options,
171            ListModelsResponse::from_raw,
172        )
173    }
174}