claude_codex/providers/kimi/
client.rs1use 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 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 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}