use super::{
Provider, ProviderEvent, ProviderRequest,
anthropic::parse_captured_sse,
claude_admission::{
CLAUDE_MAX_REQUEST_DURATION, CLAUDE_STARTUP_TIMEOUT, ClaudeAdmission, claude_error_hint,
claude_error_with_hint,
},
claude_projection::ClaudeProjection,
};
use crate::{cancellation::AgentCancellation, tools::process::terminate_child_tree_and_wait};
use anyhow::{Context, Result, bail};
use serde_json::{Value, json};
use std::{
io::{BufRead, BufReader, Read, Write},
process::{Child, Command, Stdio},
sync::mpsc,
thread,
time::{Duration, Instant},
};
const MAX_NATIVE_LINE: u64 = 1024 * 1024;
const NATIVE_PROGRESS_TIMEOUT: Duration = Duration::from_secs(300);
pub(super) struct ClaudeSubscriptionProvider {
pub(super) cwd: std::path::PathBuf,
}
fn conflicting_native_auth_override() -> Option<&'static str> {
[
"ANTHROPIC_API_KEY",
"ANTHROPIC_AUTH_TOKEN",
"ANTHROPIC_FOUNDRY_API_KEY",
"CLAUDE_CODE_USE_BEDROCK",
"CLAUDE_CODE_USE_VERTEX",
"CLAUDE_CODE_USE_FOUNDRY",
]
.into_iter()
.find(|key| std::env::var_os(key).is_some_and(|value| !value.is_empty()))
}
fn configure_cli_effort(command: &mut Command, thinking_level: crate::thinking::ThinkingLevel) {
if let Some(effort) = thinking_level.explicit_effort() {
command.args(["--effort", effort]);
}
}
fn claude_loopback_no_proxy(existing: &str) -> String {
if existing.is_empty() {
"127.0.0.1,localhost".to_string()
} else {
format!("127.0.0.1,localhost,{existing}")
}
}
pub(crate) fn cli_subscription_ready() -> bool {
cli_subscription_status().is_some()
}
pub(crate) fn cli_subscription_status() -> Option<Value> {
probe_cli_subscription().ok()
}
fn claude_cli_program() -> std::path::PathBuf {
let on_path = std::env::var_os("PATH")
.is_some_and(|path| std::env::split_paths(&path).any(|dir| dir.join("claude").is_file()));
let native_install = dirs::home_dir().map(|home| home.join(".local/bin/claude"));
match native_install {
Some(path) if !on_path && path.is_file() => path,
_ => "claude".into(),
}
}
pub(crate) fn probe_cli_subscription() -> Result<Value, String> {
if !cfg!(unix) {
return Err("Claude subscription requires Unix (macOS, Linux, or WSL)".into());
}
if let Some(key) = conflicting_native_auth_override() {
return Err(format!("unset {key}; it overrides claude.ai login"));
}
use crate::tools::process::{BoundedChildProcessLimits, run_bounded_child_process};
let mut command = Command::new(claude_cli_program());
command
.args(["auth", "status", "--json"])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.process_group(0);
}
let child = command.spawn().map_err(|error| {
format!("cannot run `claude` ({error}); install Claude Code and put `claude` on PATH")
})?;
let output = run_bounded_child_process(
child,
BoundedChildProcessLimits {
stdout_max_bytes: 8192,
stderr_max_bytes: 16 * 1024,
timeout: Duration::from_secs(10),
poll_interval: Duration::from_millis(20),
},
&AgentCancellation::default(),
)
.map_err(|error| format!("`claude auth status --json` failed: {error}"))?;
if output.timed_out {
return Err("`claude auth status --json` timed out".into());
}
if output.stdout_truncated || output.stderr_truncated {
return Err("`claude auth status --json` output too large".into());
}
if let Some(warning) = output.cleanup_warning {
return Err(format!("`claude auth status --json` {warning}"));
}
if !output.status.is_some_and(|status| status.success()) {
return Err(
"`claude auth status --json` exited unsuccessfully; run `claude auth login`".into(),
);
}
let status = serde_json::from_str::<Value>(&output.stdout)
.map_err(|_| "`claude auth status --json` returned invalid JSON; update Claude Code")?;
if status.get("loggedIn").and_then(Value::as_bool) != Some(true) {
return Err("Claude CLI not signed in; run `claude auth login`".into());
}
if status.get("authMethod").and_then(Value::as_str) != Some("claude.ai") {
return Err(
"Claude CLI not signed in with a claude.ai subscription; run `claude auth login`"
.into(),
);
}
Ok(status)
}
impl Provider for ClaudeSubscriptionProvider {
fn stream_cancellable(
&self,
request: ProviderRequest,
cancellation: &AgentCancellation,
on_event: &mut dyn FnMut(ProviderEvent) -> Result<()>,
) -> Result<()> {
cancellation.check()?;
if cfg!(not(unix)) {
bail!("Claude subscription process-tree cleanup requires Unix");
}
let estimated_tokens = crate::context::project_provider_request_input_tokens(
super::CLAUDE_SUBSCRIPTION_PROVIDER,
&request,
)
.tokens;
let ceiling = super::CLAUDE_SUBSCRIPTION_MAX_PROMPT_TOKENS;
if estimated_tokens > ceiling {
bail!(
"Claude CLI prompt is estimated at {estimated_tokens} tokens, exceeding its {ceiling}-token limit; compact the session, shorten the input, or use another provider; no CLI request was sent"
);
}
let projection = ClaudeProjection::from_request(&request)?;
let root = tempfile::tempdir().context("creating Claude request directory")?;
let system = root.path().join("system.md");
let manifest = root.path().join("tools.json");
let settings = root.path().join("settings.json");
std::fs::write(&system, &projection.system)?;
std::fs::write(&manifest, serde_json::to_vec(&projection.manifest)?)?;
std::fs::write(
&settings,
serde_json::to_vec(&json!({"env":{
"CLAUDE_CODE_EXTRA_BODY": serde_json::to_string(&projection.extra_body)?
}}))?,
)?;
let host_messages = projection
.frames
.iter()
.map(|frame| frame["message"].clone())
.collect();
let relay = ClaudeAdmission::start(cancellation, host_messages)?;
let executable = std::env::current_exe()?;
let mcp = json!({"mcpServers":{"magi":{"command":executable,"env":{
"MAGI_CLAUDE_INERT_MCP_INVENTORY":manifest
}}}});
let no_proxy = std::env::var("NO_PROXY")
.or_else(|_| std::env::var("no_proxy"))
.unwrap_or_default();
let mut command = Command::new(claude_cli_program());
command
.args([
"-p",
"--model",
&request.model,
"--input-format",
"stream-json",
"--output-format",
"stream-json",
"--verbose",
"--include-partial-messages",
"--tools",
"",
"--system-prompt-file",
])
.arg(&system)
.args(["--settings"])
.arg(&settings)
.args([
"--setting-sources",
"",
"--strict-mcp-config",
"--disable-slash-commands",
"--max-turns",
"1",
"--permission-mode",
"dontAsk",
"--no-session-persistence",
"--mcp-config",
&mcp.to_string(),
])
.current_dir(&self.cwd)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.env("ANTHROPIC_BASE_URL", &relay.url)
.env("NO_PROXY", claude_loopback_no_proxy(&no_proxy))
.env("CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC", "1")
.env("CLAUDE_CODE_MAX_RETRIES", "0")
.env("DISABLE_AUTO_COMPACT", "1")
.env("DISABLE_COMPACT", "1")
.env("ENABLE_TOOL_SEARCH", "false")
.env("CLAUDE_CODE_TOTAL_TOKENS_REMINDER", "off")
.env_remove("CLAUDE_CODE_EXTRA_BODY");
if let Some(key) = conflicting_native_auth_override() {
bail!(
"Claude subscription refuses conflicting native authentication or backend override: {key}"
);
}
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.process_group(0);
}
configure_cli_effort(&mut command, request.thinking_level);
let child = command
.spawn()
.context("Claude CLI not installed; install and sign in to Claude Code")?;
let mut child = OwnedChild(Some(child));
let stdout = child
.0
.as_mut()
.context("Claude child missing")?
.stdout
.take()
.context("Claude CLI missing stdout")?;
let stderr = child
.0
.as_mut()
.context("Claude child missing")?
.stderr
.take()
.context("Claude CLI missing stderr")?;
let stderr_reader = thread::spawn(move || {
let mut source = stderr;
let mut captured = Vec::new();
let mut buffer = [0u8; 4096];
while let Ok(count) = source.read(&mut buffer) {
if count == 0 {
break;
}
let remaining = 4096usize.saturating_sub(captured.len());
captured.extend_from_slice(&buffer[..count.min(remaining)]);
}
captured
});
let (sender, receiver) = mpsc::sync_channel(16);
let reader = thread::spawn(move || {
let mut source = BufReader::new(stdout);
loop {
let mut line = Vec::new();
let next = source
.by_ref()
.take(MAX_NATIVE_LINE + 1)
.read_until(b'\n', &mut line);
let finished =
matches!(next, Ok(0) | Err(_)) || line.len() as u64 > MAX_NATIVE_LINE;
if sender
.send(if finished { None } else { Some(line) })
.is_err()
|| finished
{
break;
}
}
});
let outcome = run_native(
child.0.as_mut().context("Claude child missing")?,
&receiver,
&projection.frames,
cancellation,
NativeDeadline::new(
CLAUDE_STARTUP_TIMEOUT,
request
.semantic_progress_timeout()
.unwrap_or(NATIVE_PROGRESS_TIMEOUT),
CLAUDE_MAX_REQUEST_DURATION,
),
);
let cleanup = child
.0
.take()
.map(|mut process| terminate_child_tree_and_wait(&mut process))
.context("Claude child missing")?;
drop(receiver);
let _ = reader.join();
let stderr = stderr_reader.join().unwrap_or_default();
let stderr_hint = claude_error_hint(&stderr);
cancellation.check()?;
let cleanup = cleanup?;
if let Some(warning) = cleanup.cleanup_warning {
bail!("Claude process cleanup failed: {warning}");
}
let completion = outcome.map_err(|error| claude_error_with_hint(error, stderr_hint))?;
let hint = completion.error_hint.or(stderr_hint);
let response = relay
.finish()
.map_err(|error| claude_error_with_hint(error, hint))?;
if !(completion.assistant && completion.result && completion.stop) {
return Err(claude_error_with_hint(
anyhow::anyhow!(
"Claude CLI missing assistant, result, or message_stop; check Claude CLI version and login"
),
hint,
));
}
let _ = response.denied;
let mut events = parse_captured_sse(&response.body)?;
let tool_boundary = completion.subtype == "error_max_turns"
&& events
.iter()
.any(|event| matches!(event, ProviderEvent::ToolCall(_)));
if completion.is_error && !tool_boundary {
return Err(claude_error_with_hint(
anyhow::anyhow!("Claude CLI returned an error; check Claude CLI version and login"),
hint,
));
}
if completion.subtype != "success" && !tool_boundary {
bail!("Claude CLI did not complete successfully");
}
if !cleanup
.status
.is_some_and(|status| status.success() || (tool_boundary && status.code() == Some(1)))
{
bail!("Claude CLI exited unsuccessfully");
}
let names: std::collections::HashSet<&str> = projection
.manifest
.iter()
.filter_map(|tool| tool.get("name").and_then(Value::as_str))
.collect();
for event in &mut events {
match event {
ProviderEvent::ToolCall(call) => {
call.name = host_tool_name(&call.name, &names)?.to_owned();
}
ProviderEvent::ResponseItem(item)
if item.get("type").and_then(Value::as_str) == Some("function_call") =>
{
let name = item.get("name").and_then(Value::as_str).unwrap_or("");
item["name"] = json!(host_tool_name(name, &names)?);
}
_ => {}
}
}
for event in events {
cancellation.check()?;
on_event(event)?;
}
Ok(())
}
}
fn host_tool_name<'a>(
name: &'a str,
inventory: &std::collections::HashSet<&str>,
) -> Result<&'a str> {
const PREFIX: &str = "mcp__magi__";
let printable = if name.len() <= PREFIX.len() + 50
&& name
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_' || byte == b'-')
{
name
} else {
"[invalid tool name]"
};
let Some(host_name) = name.strip_prefix(PREFIX) else {
bail!(
"Claude CLI returned tool with unexpected name `{printable}` (expected {PREFIX} prefix)"
);
};
if !inventory.contains(host_name) {
bail!(
"Claude CLI returned tool outside host inventory: `{printable}` (not provided for this request)"
);
}
Ok(host_name)
}
struct OwnedChild(Option<Child>);
impl Drop for OwnedChild {
fn drop(&mut self) {
if let Some(mut child) = self.0.take() {
let _ = terminate_child_tree_and_wait(&mut child);
}
}
}
struct NativeDeadline {
startup_deadline: Instant,
progress_deadline: Option<Instant>,
overall_deadline: Instant,
progress_timeout: Duration,
}
impl NativeDeadline {
fn new(startup: Duration, progress: Duration, overall: Duration) -> Self {
let now = Instant::now();
Self {
startup_deadline: now + startup,
progress_deadline: None,
overall_deadline: now + overall,
progress_timeout: progress,
}
}
fn check(&self) -> Result<()> {
let now = Instant::now();
if now >= self.overall_deadline {
bail!("Claude CLI exceeded overall time limit");
}
if now >= self.progress_deadline.unwrap_or(self.startup_deadline) {
if self.progress_deadline.is_some() {
bail!("Claude CLI timed out waiting for progress");
}
bail!("Claude CLI startup timed out waiting for output");
}
Ok(())
}
fn observe(&mut self, event: &Value) {
if native_event_has_progress(event) {
self.progress_deadline = Some(Instant::now() + self.progress_timeout);
}
}
}
fn native_event_has_progress(event: &Value) -> bool {
if event.get("type").and_then(Value::as_str) != Some("stream_event") {
return false;
}
let event = &event["event"];
let nonempty = |value: &Value| value.as_str().is_some_and(|text| !text.is_empty());
match event.get("type").and_then(Value::as_str) {
Some("content_block_start") => {
let block = &event["content_block"];
match block.get("type").and_then(Value::as_str) {
Some("tool_use") => nonempty(&block["id"]) && nonempty(&block["name"]),
Some("text") => nonempty(&block["text"]),
Some("thinking") => nonempty(&block["thinking"]),
Some("redacted_thinking") => nonempty(&block["data"]),
_ => false,
}
}
Some("content_block_delta") => {
let delta = &event["delta"];
let field = match delta.get("type").and_then(Value::as_str) {
Some("text_delta") => "text",
Some("thinking_delta") => "thinking",
Some("signature_delta") => "signature",
Some("redacted_thinking_delta") => "data",
Some("input_json_delta") => "partial_json",
_ => return false,
};
nonempty(&delta[field])
}
_ => false,
}
}
fn receive(
receiver: &mpsc::Receiver<Option<Vec<u8>>>,
cancel: &AgentCancellation,
deadline: &NativeDeadline,
) -> Result<Option<Value>> {
loop {
cancel.check()?;
deadline.check()?;
let wait = deadline
.progress_deadline
.unwrap_or(deadline.startup_deadline)
.min(deadline.overall_deadline)
.saturating_duration_since(Instant::now())
.min(Duration::from_millis(100));
let received = receiver.recv_timeout(wait);
cancel.check()?;
deadline.check()?;
match received {
Ok(Some(line)) => {
return serde_json::from_slice(&line)
.map(Some)
.context("invalid Claude CLI stream-json output");
}
Ok(None) | Err(mpsc::RecvTimeoutError::Disconnected) => return Ok(None),
Err(mpsc::RecvTimeoutError::Timeout) => {}
}
}
}
fn replay_acknowledged(event: &Value) -> Result<bool> {
if event.get("type").and_then(Value::as_str) != Some("result") {
return Ok(false);
}
if event.get("num_turns").and_then(Value::as_u64) != Some(0)
|| event.get("is_error").and_then(Value::as_bool) == Some(true)
{
let hint = event
.get("result")
.and_then(Value::as_str)
.and_then(|message| claude_error_hint(message.as_bytes()));
return Err(claude_error_with_hint(
anyhow::anyhow!(
"Claude CLI does not support zero-turn history replay; check Claude CLI version"
),
hint,
));
}
Ok(true)
}
struct NativeCompletion {
assistant: bool,
result: bool,
stop: bool,
subtype: String,
is_error: bool,
error_hint: Option<&'static str>,
}
fn run_native(
child: &mut Child,
receiver: &mpsc::Receiver<Option<Vec<u8>>>,
frames: &[Value],
cancellation: &AgentCancellation,
mut deadline: NativeDeadline,
) -> Result<NativeCompletion> {
let mut stdin = child.stdin.take().context("Claude CLI missing stdin")?;
#[cfg(unix)]
{
use std::os::fd::AsRawFd;
crate::tools::process::set_nonblocking(stdin.as_raw_fd())?;
}
for frame in frames {
let mut encoded = serde_json::to_vec(frame)?;
encoded.push(b'\n');
let mut remaining = encoded.as_slice();
while !remaining.is_empty() {
cancellation.check()?;
deadline.check().context("writing Claude CLI stdin")?;
match stdin.write(remaining) {
Ok(0) => bail!("Claude CLI stdin closed"),
Ok(count) => remaining = &remaining[count..],
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(10));
}
Err(_) => bail!("Claude CLI stdin write failed"),
}
}
if frame.get("shouldQuery") == Some(&Value::Bool(false)) {
loop {
let event = receive(receiver, cancellation, &deadline)?
.context("Claude CLI exited before replay acknowledgment")?;
if replay_acknowledged(&event)? {
break;
}
}
}
}
drop(stdin);
let mut completion = NativeCompletion {
assistant: false,
result: false,
stop: false,
subtype: String::new(),
is_error: false,
error_hint: None,
};
while let Some(event) = receive(receiver, cancellation, &deadline)? {
deadline.observe(&event);
match event.get("type").and_then(Value::as_str) {
Some("assistant") => {
if event.get("error").is_some() || event.pointer("/message/error").is_some() {
completion.is_error = true;
}
completion.assistant = true;
}
Some("result") => {
if completion.result {
bail!("Claude CLI returned multiple results");
}
completion.result = true;
completion.is_error |= event.get("is_error").and_then(Value::as_bool) == Some(true);
if completion.is_error {
completion.error_hint = event
.get("result")
.and_then(Value::as_str)
.and_then(|message| claude_error_hint(message.as_bytes()));
}
completion.subtype = event
.get("subtype")
.and_then(Value::as_str)
.unwrap_or("")
.to_owned();
}
Some("stream_event")
if event.pointer("/event/type").and_then(Value::as_str) == Some("message_stop") =>
{
completion.stop = true;
}
_ => {}
}
}
Ok(completion)
}