Skip to main content

claude_codex/providers/kimi/
client.rs

1use std::time::{SystemTime, UNIX_EPOCH};
2
3use crate::providers::kimi::auth::constants::api_base_url;
4use crate::providers::kimi::auth::headers::common_headers;
5use crate::providers::kimi::auth::manager::KimiAuthManager;
6use crate::providers::kimi::auth::token_store::{StoredAuth, file_store};
7use crate::providers::kimi::translate::request::KimiChatRequest;
8use crate::retry::{MAX_RATE_LIMIT_RETRIES, compute_backoff_delay};
9
10#[derive(Debug)]
11pub struct KimiError {
12    pub status: u16,
13    pub message: String,
14    pub detail: Option<String>,
15    pub retry_after: Option<String>,
16}
17
18pub struct KimiResponse {
19    pub body: Vec<u8>,
20    pub status: u16,
21    pub request_start_time: u64,
22}
23
24pub struct KimiHttpClient {
25    client: reqwest::blocking::Client,
26    auth_manager: KimiAuthManager<crate::auth::FileAuthStore<StoredAuth>>,
27}
28
29impl Default for KimiHttpClient {
30    fn default() -> Self {
31        Self::new()
32    }
33}
34
35impl KimiHttpClient {
36    pub fn new() -> Self {
37        Self {
38            client: reqwest::blocking::Client::builder()
39                .timeout(std::time::Duration::from_secs(120))
40                .build()
41                .expect("failed to create HTTP client"),
42            auth_manager: KimiAuthManager::new(file_store()),
43        }
44    }
45
46    pub fn auth_manager(&self) -> &KimiAuthManager<crate::auth::FileAuthStore<StoredAuth>> {
47        &self.auth_manager
48    }
49
50    pub fn post_kimi(&self, body: &KimiChatRequest) -> Result<KimiResponse, KimiError> {
51        let mut auth = self.auth_manager.get_auth().map_err(|e| KimiError {
52            status: 401,
53            message: "Auth error".to_string(),
54            detail: Some(e.to_string()),
55            retry_after: None,
56        })?;
57
58        let mut attempt = 0u32;
59        loop {
60            let result = self.attempt_post(&auth.access, body);
61
62            match result {
63                Ok(response) if response.status == 401 && attempt == 0 => {
64                    // First 401: try refresh
65                    match self.auth_manager.force_refresh() {
66                        Ok(new_auth) => {
67                            auth = new_auth;
68                            attempt += 1;
69                            continue;
70                        }
71                        Err(e) => {
72                            return Err(KimiError {
73                                status: 401,
74                                message: "Unauthorized".to_string(),
75                                detail: Some(e.to_string()),
76                                retry_after: None,
77                            });
78                        }
79                    }
80                }
81                Ok(response) => return Ok(response),
82                Err(err @ KimiError { status: 429, .. }) => {
83                    if attempt < MAX_RATE_LIMIT_RETRIES {
84                        let delay = compute_backoff_delay(attempt, err.retry_after.as_deref());
85                        std::thread::sleep(std::time::Duration::from_millis(delay.wait_ms));
86                        attempt += 1;
87                        continue;
88                    }
89                    return Err(err);
90                }
91                Err(err) => return Err(err),
92            }
93        }
94    }
95
96    fn attempt_post(
97        &self,
98        access_token: &str,
99        body: &KimiChatRequest,
100    ) -> Result<KimiResponse, KimiError> {
101        let headers = common_headers().map_err(|e| KimiError {
102            status: 500,
103            message: "Failed to build headers".to_string(),
104            detail: Some(e.to_string()),
105            retry_after: None,
106        })?;
107
108        let url = format!("{}/chat/completions", api_base_url());
109        let body_json = serde_json::to_string(body).map_err(|e| KimiError {
110            status: 500,
111            message: "Failed to serialize request".to_string(),
112            detail: Some(e.to_string()),
113            retry_after: None,
114        })?;
115
116        let request_start_time = now_ms();
117
118        let mut req_builder = self
119            .client
120            .post(&url)
121            .header("Content-Type", "application/json")
122            .header("Accept", "application/json")
123            .header("Authorization", format!("Bearer {access_token}"));
124
125        // Add common headers
126        for (k, v) in &headers {
127            if let Ok(name) = reqwest::header::HeaderName::from_bytes(k.as_bytes())
128                && let Ok(value) = reqwest::header::HeaderValue::from_str(v)
129            {
130                req_builder = req_builder.header(name, value);
131            }
132        }
133
134        let resp = match req_builder.body(body_json).send() {
135            Ok(r) => r,
136            Err(e) => {
137                return Err(KimiError {
138                    status: 0,
139                    message: "Network error".to_string(),
140                    detail: Some(e.to_string()),
141                    retry_after: None,
142                });
143            }
144        };
145
146        let status = resp.status().as_u16();
147
148        if status == 429 {
149            let retry_after = resp
150                .headers()
151                .get("retry-after")
152                .and_then(|v| v.to_str().ok())
153                .map(|s| s.to_string());
154            let text = resp.text().unwrap_or_default();
155            return Err(KimiError {
156                status: 429,
157                message: "Rate limited".to_string(),
158                detail: if text.is_empty() { None } else { Some(text) },
159                retry_after,
160            });
161        }
162
163        if status == 401 || status == 403 {
164            let text = resp.text().unwrap_or_default();
165            return Err(KimiError {
166                status,
167                message: if status == 401 {
168                    "Unauthorized"
169                } else {
170                    "Forbidden"
171                }
172                .to_string(),
173                detail: if text.is_empty() { None } else { Some(text) },
174                retry_after: None,
175            });
176        }
177
178        if !resp.status().is_success() {
179            let text = resp.text().unwrap_or_default();
180            return Err(KimiError {
181                status,
182                message: "Upstream error".to_string(),
183                detail: if text.is_empty() { None } else { Some(text) },
184                retry_after: None,
185            });
186        }
187
188        let body_bytes = resp.bytes().map(|b| b.to_vec()).unwrap_or_default();
189
190        Ok(KimiResponse {
191            body: body_bytes,
192            status,
193            request_start_time,
194        })
195    }
196}
197
198fn now_ms() -> u64 {
199    SystemTime::now()
200        .duration_since(UNIX_EPOCH)
201        .unwrap_or_default()
202        .as_millis() as u64
203}