Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 68 additions & 21 deletions src/api/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
/// Anthropic API client — port of services/api/claude.ts
pub mod ollama;
pub mod openai_compat;
pub mod retry;
pub mod types;

use anyhow::{Context, Result, anyhow};
Expand Down Expand Up @@ -69,6 +70,9 @@ pub struct ClaudeClient {
/// would send the field twice. Merging into the single per-request value
/// keeps exactly one.
credential_betas: Vec<String>,
/// Optional sink for retry notices, so a backoff sleep is visible rather
/// than looking like a hang. Set by the TUI and the headless runner.
retry_notifier: Option<retry::RetryNotifier>,
}

impl ClaudeClient {
Expand Down Expand Up @@ -119,9 +123,27 @@ impl ClaudeClient {
api_key,
base_url: ANTHROPIC_API_BASE.to_string(),
credential_betas,
retry_notifier: None,
})
}

/// Test-only: point the client at a local scripted server.
///
/// Deliberately `#[cfg(test)]`. A runtime base-URL override is a
/// credential-exfiltration vector, which is exactly why
/// `ANTHROPIC_BASE_URL` is excluded from the `.env` allowlist in main.rs.
/// This exists so the retry wiring can be proven without weakening that.
#[cfg(test)]
pub(crate) fn set_base_url_for_test(&mut self, url: impl Into<String>) {
self.base_url = url.into();
}

/// Install a sink for retry notices. Without one a backoff sleep is
/// invisible and a rate-limited turn looks like a hang.
pub fn set_retry_notifier(&mut self, n: retry::RetryNotifier) {
self.retry_notifier = Some(n);
}

/// Merge the request's betas with any the credential requires.
/// Returns `None` when there are none, so the header is omitted entirely.
fn beta_header(&self, request_betas: &[String]) -> Option<String> {
Expand All @@ -148,15 +170,22 @@ impl ClaudeClient {
let url = format!("{}/v1/messages", self.base_url);
debug!("POST {url} model={}", request.model);

let mut builder = self.client.post(&url).json(&request);
if let Some(betas) = self.beta_header(&request.betas) {
builder = builder.header("anthropic-beta", betas);
}
if let Some(ref sid) = request.session_id {
builder = builder.header("X-Claude-Code-Session-Id", sid.as_str());
}

let resp = builder.send().await.context("API request failed")?;
let betas = self.beta_header(&request.betas);
let resp = retry::send_with_retry(
|| {
let mut builder = self.client.post(&url).json(&request);
if let Some(ref b) = betas {
builder = builder.header("anthropic-beta", b.as_str());
}
if let Some(ref sid) = request.session_id {
builder = builder.header("X-Claude-Code-Session-Id", sid.as_str());
}
builder
},
self.retry_notifier.as_ref(),
"API request failed",
)
.await?;

let status = resp.status();
if !status.is_success() {
Expand All @@ -181,18 +210,24 @@ impl ClaudeClient {
let url = format!("{}/v1/messages", self.base_url);
debug!("POST {url} stream=true model={}", request.model);

let mut builder = self.client.post(&url).json(&request);
if let Some(betas) = self.beta_header(&request.betas) {
builder = builder.header("anthropic-beta", betas);
}
if let Some(ref sid) = request.session_id {
builder = builder.header("X-Claude-Code-Session-Id", sid.as_str());
}

let resp = builder
.send()
.await
.context("Streaming API request failed")?;
// Retrying is safe here and only here: nothing has been handed to
// `on_text` yet, so a retry cannot duplicate text the user has seen.
let betas = self.beta_header(&request.betas);
let resp = retry::send_with_retry(
|| {
let mut builder = self.client.post(&url).json(&request);
if let Some(ref b) = betas {
builder = builder.header("anthropic-beta", b.as_str());
}
if let Some(ref sid) = request.session_id {
builder = builder.header("X-Claude-Code-Session-Id", sid.as_str());
}
builder
},
self.retry_notifier.as_ref(),
"Streaming API request failed",
)
.await?;

let status = resp.status();
if !status.is_success() {
Expand Down Expand Up @@ -452,6 +487,18 @@ impl ApiBackend {
}
}

/// Route retry notices to the UI. Applies to every backend that talks to
/// a remote provider.
pub fn set_retry_notifier(&mut self, n: retry::RetryNotifier) {
match self {
Self::Anthropic(c) => c.set_retry_notifier(n),
Self::OpenAiCompat(c) => c.set_retry_notifier(n),
// Ollama is a local daemon: it does not rate limit, and a
// connection failure there means the server is down, not busy.
Self::Ollama(_) => {}
}
}

/// Returns true the first time called after tools are disabled — for a one-time user notice.
pub fn take_tools_notice(&self) -> bool {
match self {
Expand Down
62 changes: 36 additions & 26 deletions src/api/openai_compat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -529,6 +529,9 @@ pub struct OpenAiCompatClient {
/// Set to true after the first 400 "does not support tools" error.
no_tools: Arc<AtomicBool>,
tools_notice_sent: Arc<AtomicBool>,
/// See `ClaudeClient::retry_notifier`. Rate limiting is far more common on
/// these providers than on Anthropic — Groq and OpenRouter throttle hard.
retry_notifier: Option<super::retry::RetryNotifier>,
}

impl OpenAiCompatClient {
Expand Down Expand Up @@ -595,6 +598,7 @@ impl OpenAiCompatClient {
.context("Failed to build HTTP client")?;

Ok(Self {
retry_notifier: None,
client,
base_url,
api_key,
Expand All @@ -610,6 +614,11 @@ impl OpenAiCompatClient {
self.no_tools.load(Ordering::Relaxed)
}

/// See `ClaudeClient::set_retry_notifier`.
pub fn set_retry_notifier(&mut self, n: super::retry::RetryNotifier) {
self.retry_notifier = Some(n);
}

pub fn take_tools_notice(&self) -> bool {
self.no_tools.load(Ordering::Relaxed)
&& !self.tools_notice_sent.swap(true, Ordering::Relaxed)
Expand Down Expand Up @@ -658,22 +667,28 @@ impl OpenAiCompatClient {
}),
};

let mut builder = self.client.post(&url).json(&oai_request);

// Auth header
if !self.api_key.is_empty() {
builder = builder.bearer_auth(&self.api_key);
}

// Provider-specific headers (e.g. OpenRouter's HTTP-Referer)
for (k, v) in &self.extra_headers {
builder = builder.header(k.as_str(), v.as_str());
}
// Nothing has reached `on_text` yet, so retrying cannot duplicate
// output. `send_with_retry` hands back the final response even on an
// error status, which keeps the "does not support tools" sniff below
// working exactly as before.
let build = |req: &OaiRequest| {
let mut builder = self.client.post(&url).json(req);
if !self.api_key.is_empty() {
builder = builder.bearer_auth(&self.api_key);
}
// Provider-specific headers (e.g. OpenRouter's HTTP-Referer)
for (k, v) in &self.extra_headers {
builder = builder.header(k.as_str(), v.as_str());
}
builder
};

let resp = builder
.send()
.await
.context(format!("{} request failed", self.provider_name))?;
let resp = super::retry::send_with_retry(
|| build(&oai_request),
self.retry_notifier.as_ref(),
&format!("{} request failed", self.provider_name),
)
.await?;

let status = resp.status();
let resp = if !status.is_success() {
Expand All @@ -687,17 +702,12 @@ impl OpenAiCompatClient {
oai_request.messages = translate_messages(&patched_system, &request.messages);
oai_request.tools = vec![];

let mut retry = self.client.post(&url).json(&oai_request);
if !self.api_key.is_empty() {
retry = retry.bearer_auth(&self.api_key);
}
for (k, v) in &self.extra_headers {
retry = retry.header(k.as_str(), v.as_str());
}
retry
.send()
.await
.context(format!("{} request failed", self.provider_name))?
super::retry::send_with_retry(
|| build(&oai_request),
self.retry_notifier.as_ref(),
&format!("{} request failed", self.provider_name),
)
.await?
} else {
return Err(anyhow!("{} error {status}: {body}", self.provider_name));
}
Expand Down
Loading
Loading