use crate::http::{Request, send};
use std::io::{BufRead, BufReader};
use std::process::{Command, Stdio};
use std::sync::mpsc::{Receiver, Sender, channel};
const BETA: &str = "managed-agents-2026-04-01";
const VERSION: &str = "2023-06-01";
#[derive(Debug, Clone)]
pub enum Backend {
FirstParty { api_key: String },
ClaudePlatformAwsKey {
api_key: String,
region: String,
workspace_id: String,
},
ClaudePlatformAwsSigV4 {
region: String,
workspace_id: String,
},
}
impl Backend {
#[allow(dead_code)] pub fn label(&self) -> &'static str {
match self {
Backend::FirstParty { .. } => "first-party Claude API",
Backend::ClaudePlatformAwsKey { .. } => "Claude Platform on AWS (API key)",
Backend::ClaudePlatformAwsSigV4 { .. } => "Claude Platform on AWS (SigV4)",
}
}
pub fn base(&self) -> String {
match self {
Backend::FirstParty { .. } => "https://api.anthropic.com".to_string(),
Backend::ClaudePlatformAwsKey { region, .. }
| Backend::ClaudePlatformAwsSigV4 { region, .. } => {
format!("https://aws-external-anthropic.{region}.api.aws")
}
}
}
pub fn headers(&self) -> Vec<(String, String)> {
let mut out = vec![
("anthropic-version".to_string(), VERSION.to_string()),
("anthropic-beta".to_string(), BETA.to_string()),
("content-type".to_string(), "application/json".to_string()),
];
match self {
Backend::FirstParty { api_key } => {
out.push(("x-api-key".to_string(), api_key.clone()));
}
Backend::ClaudePlatformAwsKey {
api_key,
workspace_id,
..
} => {
out.push(("x-api-key".to_string(), api_key.clone()));
out.push(("anthropic-workspace-id".to_string(), workspace_id.clone()));
}
Backend::ClaudePlatformAwsSigV4 { workspace_id, .. } => {
out.push(("anthropic-workspace-id".to_string(), workspace_id.clone()));
}
}
out
}
fn is_sigv4(&self) -> bool {
matches!(self, Backend::ClaudePlatformAwsSigV4 { .. })
}
}
fn aws_export_credentials() -> Result<Vec<(String, String)>, String> {
let out = Command::new("aws")
.args(["configure", "export-credentials", "--format", "env"])
.output()
.map_err(|e| format!("spawn `aws configure export-credentials`: {e}"))?;
if !out.status.success() {
return Err(format!(
"`aws configure export-credentials` failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
let text = String::from_utf8_lossy(&out.stdout);
let mut creds = Vec::new();
for line in text.lines() {
let line = line.trim();
let rest = line.strip_prefix("export ").unwrap_or(line);
let Some((k, v)) = rest.split_once('=') else {
continue;
};
let v = v.trim_matches(|c| c == '\'' || c == '"');
creds.push((k.to_string(), v.to_string()));
}
if !creds.iter().any(|(k, _)| k == "AWS_ACCESS_KEY_ID") {
return Err(
"no AWS_ACCESS_KEY_ID in aws export — credentials not available (run `aws sso login`?)"
.to_string(),
);
}
Ok(creds)
}
fn cred(creds: &[(String, String)], key: &str) -> String {
creds
.iter()
.find(|(k, _)| k == key)
.map(|(_, v)| v.clone())
.unwrap_or_default()
}
fn curl_sigv4(
region: &str,
method: &str,
url: &str,
headers: &[(String, String)],
body: Option<&str>,
) -> Result<(u16, String), String> {
let creds = aws_export_credentials()?;
let access = cred(&creds, "AWS_ACCESS_KEY_ID");
let secret = cred(&creds, "AWS_SECRET_ACCESS_KEY");
let session_token = cred(&creds, "AWS_SESSION_TOKEN");
if access.is_empty() || secret.is_empty() {
return Err("aws export-credentials returned no access key / secret".to_string());
}
let mut cmd = Command::new("curl");
cmd.args([
"-sS",
"-X",
method,
"-w",
"\n__HTTP_STATUS__%{http_code}",
"--max-time",
"30",
"--aws-sigv4",
&format!("aws:amz:{region}:aws-external-anthropic"),
"--user",
]);
cmd.arg(format!("{access}:{secret}"));
if !session_token.is_empty() {
cmd.arg("-H");
cmd.arg(format!("x-amz-security-token: {session_token}"));
}
for (k, v) in headers {
cmd.arg("-H");
cmd.arg(format!("{k}: {v}"));
}
if let Some(b) = body {
cmd.arg("--data-binary");
cmd.arg(b);
}
cmd.arg(url);
let out = cmd.output().map_err(|e| format!("spawn curl: {e}"))?;
if !out.status.success() {
return Err(format!(
"curl failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
let raw = String::from_utf8_lossy(&out.stdout).to_string();
let (body, status) = match raw.rsplit_once("\n__HTTP_STATUS__") {
Some((b, s)) => (b.to_string(), s.parse::<u16>().unwrap_or(0)),
None => (raw, 0),
};
Ok((status, body))
}
fn dispatch(
backend: &Backend,
method: &str,
path: &str,
body: Option<String>,
) -> Result<(u16, String), String> {
let url = format!("{}{}", backend.base(), path);
if backend.is_sigv4() {
let region = match backend {
Backend::ClaudePlatformAwsSigV4 { region, .. } => region.clone(),
_ => unreachable!(),
};
return curl_sigv4(®ion, method, &url, &backend.headers(), body.as_deref());
}
let req = Request {
method: method.to_string(),
url,
headers: backend.headers(),
body,
insecure: false,
};
let resp = send(&req).map_err(|e| format!("send: {e}"))?;
Ok((resp.status, resp.body))
}
pub fn detect_backend() -> Result<Backend, String> {
let aws_key = std::env::var("ANTHROPIC_AWS_API_KEY").ok();
let aws_region = std::env::var("AWS_REGION")
.ok()
.or_else(|| std::env::var("AWS_DEFAULT_REGION").ok());
let aws_workspace = std::env::var("ANTHROPIC_AWS_WORKSPACE_ID").ok();
let region_ok = aws_region
.as_deref()
.map(|s| !s.is_empty())
.unwrap_or(false);
let workspace_ok = aws_workspace
.as_deref()
.map(|s| !s.is_empty())
.unwrap_or(false);
let key_set = aws_key.as_deref().map(|s| !s.is_empty()).unwrap_or(false);
if region_ok && workspace_ok {
let region = aws_region.unwrap();
let workspace_id = aws_workspace.unwrap();
if key_set {
return Ok(Backend::ClaudePlatformAwsKey {
api_key: aws_key.unwrap(),
region,
workspace_id,
});
}
return Ok(Backend::ClaudePlatformAwsSigV4 {
region,
workspace_id,
});
}
let first = std::env::var("ANTHROPIC_API_KEY")
.map_err(|_| "no managed-agents auth found — set ANTHROPIC_API_KEY (first-party), OR AWS_REGION + ANTHROPIC_AWS_WORKSPACE_ID (Claude Platform on AWS via SigV4), optionally + ANTHROPIC_AWS_API_KEY (AWS API key auth)".to_string())?;
if first.is_empty() {
return Err("ANTHROPIC_API_KEY is empty".to_string());
}
Ok(Backend::FirstParty { api_key: first })
}
#[derive(Debug)]
pub struct Created {
pub id: String,
}
#[derive(Debug)]
pub struct CreatedSession {
pub id: String,
#[allow(dead_code)]
pub agent_id: String,
#[allow(dead_code)]
pub environment_id: String,
}
fn extract_id(body: &str) -> Result<String, String> {
let v: serde_json::Value =
serde_json::from_str(body).map_err(|e| format!("response JSON parse: {e}"))?;
v.get("id")
.and_then(|x| x.as_str())
.map(str::to_string)
.ok_or_else(|| format!("response missing `id` field: {body}"))
}
pub fn create_agent(
backend: &Backend,
name: &str,
model: &str,
system: &str,
) -> Result<Created, String> {
let body = serde_json::json!({
"name": name,
"model": model,
"system": system,
"tools": [{"type": "agent_toolset_20260401"}],
})
.to_string();
let (status, body) = dispatch(backend, "POST", "/v1/agents", Some(body))
.map_err(|e| format!("create_agent: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("create_agent HTTP {status}: {body}"));
}
Ok(Created {
id: extract_id(&body)?,
})
}
pub fn create_environment(
backend: &Backend,
name: &str,
config_kind: &str,
) -> Result<Created, String> {
let config = match config_kind {
"cloud" => serde_json::json!({
"type": "cloud",
"networking": {"type": "unrestricted"},
}),
"self_hosted" => serde_json::json!({"type": "self_hosted"}),
other => return Err(format!("unknown environment kind: {other}")),
};
let body = serde_json::json!({"name": name, "config": config}).to_string();
let (status, body) = dispatch(backend, "POST", "/v1/environments", Some(body))
.map_err(|e| format!("create_environment: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("create_environment HTTP {status}: {body}"));
}
Ok(Created {
id: extract_id(&body)?,
})
}
pub fn create_session(
backend: &Backend,
agent_id: &str,
environment_id: &str,
title: &str,
) -> Result<CreatedSession, String> {
let body = serde_json::json!({
"agent": agent_id,
"environment_id": environment_id,
"title": title,
})
.to_string();
let (status, body) = dispatch(backend, "POST", "/v1/sessions", Some(body))
.map_err(|e| format!("create_session: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("create_session HTTP {status}: {body}"));
}
Ok(CreatedSession {
id: extract_id(&body)?,
agent_id: agent_id.to_string(),
environment_id: environment_id.to_string(),
})
}
pub fn stop_session(backend: &Backend, session_id: &str) -> Result<(), String> {
let path = format!("/v1/sessions/{session_id}/stop");
let (status, body) = dispatch(backend, "POST", &path, Some("{}".to_string()))
.map_err(|e| format!("stop_session: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("stop_session HTTP {status}: {body}"));
}
Ok(())
}
pub fn send_user_message(backend: &Backend, session_id: &str, text: &str) -> Result<(), String> {
let body = serde_json::json!({
"events": [{
"type": "user.message",
"content": [{"type": "text", "text": text}],
}],
})
.to_string();
let path = format!("/v1/sessions/{session_id}/events");
let (status, body) = dispatch(backend, "POST", &path, Some(body))
.map_err(|e| format!("send_user_message: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("send_user_message HTTP {status}: {body}"));
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct SessionSummary {
pub id: String,
pub title: Option<String>,
pub status: String,
#[allow(dead_code)]
pub created_at: Option<String>,
#[allow(dead_code)]
pub agent_id: Option<String>,
#[allow(dead_code)]
pub environment_id: Option<String>,
}
pub fn list_sessions(backend: &Backend) -> Result<Vec<SessionSummary>, String> {
let (status, body) = dispatch(backend, "GET", "/v1/sessions?limit=50", None)
.map_err(|e| format!("list_sessions: {e}"))?;
if !(200..300).contains(&status) {
return Err(format!("list_sessions HTTP {status}: {body}"));
}
let v: serde_json::Value =
serde_json::from_str(&body).map_err(|e| format!("list_sessions JSON: {e}"))?;
let arr = v.get("data").and_then(|d| d.as_array());
let Some(arr) = arr else {
return Err(format!("list_sessions missing `data`: {body}"));
};
let mut out = Vec::with_capacity(arr.len());
for item in arr {
let Some(id) = item.get("id").and_then(|x| x.as_str()) else {
continue;
};
let status = item
.get("status")
.and_then(|x| x.as_str())
.unwrap_or("unknown")
.to_string();
let title = item.get("title").and_then(|x| x.as_str()).map(String::from);
let created_at = item
.get("created_at")
.and_then(|x| x.as_str())
.map(String::from);
let agent_id = item
.get("agent")
.and_then(|x| x.get("id"))
.and_then(|x| x.as_str())
.or_else(|| item.get("agent").and_then(|x| x.as_str()))
.map(String::from);
let environment_id = item
.get("environment_id")
.and_then(|x| x.as_str())
.map(String::from);
out.push(SessionSummary {
id: id.to_string(),
title,
status,
created_at,
agent_id,
environment_id,
});
}
Ok(out)
}
#[derive(Debug, Clone)]
pub enum SessionStreamEvent {
Line(String),
Done,
Error(String),
}
pub fn spawn_session_event_stream(session_id: String) -> Receiver<SessionStreamEvent> {
let (tx, rx) = channel::<SessionStreamEvent>();
std::thread::spawn(move || {
let backend = match detect_backend() {
Ok(b) => b,
Err(e) => {
let _ = tx.send(SessionStreamEvent::Error(format!("backend: {e}")));
return;
}
};
let url = format!("{}/v1/sessions/{}/stream", backend.base(), session_id);
let mut args: Vec<String> = vec![
"-sS".to_string(),
"-N".to_string(),
"-H".to_string(),
"Accept: text/event-stream".to_string(),
];
if let Backend::ClaudePlatformAwsSigV4 { region, .. } = &backend {
let creds = match aws_export_credentials() {
Ok(c) => c,
Err(e) => {
let _ = tx.send(SessionStreamEvent::Error(format!(
"aws export-credentials: {e}"
)));
return;
}
};
let access = cred(&creds, "AWS_ACCESS_KEY_ID");
let secret = cred(&creds, "AWS_SECRET_ACCESS_KEY");
let session_token = cred(&creds, "AWS_SESSION_TOKEN");
if access.is_empty() || secret.is_empty() {
let _ = tx.send(SessionStreamEvent::Error(
"aws export-credentials: missing access key / secret".to_string(),
));
return;
}
args.push("--aws-sigv4".to_string());
args.push(format!("aws:amz:{region}:aws-external-anthropic"));
args.push("--user".to_string());
args.push(format!("{access}:{secret}"));
if !session_token.is_empty() {
args.push("-H".to_string());
args.push(format!("x-amz-security-token: {session_token}"));
}
}
for (k, v) in backend.headers() {
args.push("-H".to_string());
args.push(format!("{k}: {v}"));
}
args.push(url);
let mut cmd = Command::new("curl");
cmd.args(&args).stdout(Stdio::piped()).stdin(Stdio::null());
match cmd.spawn() {
Ok(mut child) => {
let stdout = match child.stdout.take() {
Some(s) => s,
None => {
let _ = tx.send(SessionStreamEvent::Error(
"curl produced no stdout".to_string(),
));
return;
}
};
let reader = BufReader::new(stdout);
let _ = run_sse_reader(reader, &tx);
let _ = child.wait();
let _ = tx.send(SessionStreamEvent::Done);
}
Err(e) => {
let _ = tx.send(SessionStreamEvent::Error(format!("spawn curl: {e}")));
}
}
});
rx
}
fn run_sse_reader<R: BufRead>(reader: R, tx: &Sender<SessionStreamEvent>) -> std::io::Result<()> {
for line in reader.lines() {
let line = line?;
let Some(payload) = line.strip_prefix("data: ") else {
continue;
};
let v: serde_json::Value = match serde_json::from_str(payload) {
Ok(v) => v,
Err(_) => continue,
};
let rendered = render_stream_event(&v);
if !rendered.is_empty() && tx.send(SessionStreamEvent::Line(rendered)).is_err() {
break;
}
}
Ok(())
}
fn render_stream_event(v: &serde_json::Value) -> String {
let ty = v.get("type").and_then(|x| x.as_str()).unwrap_or("");
match ty {
"agent.message" => v
.get("content")
.and_then(|c| c.as_array())
.map(|arr| {
arr.iter()
.filter_map(|b| b.get("text").and_then(|t| t.as_str()))
.collect::<Vec<_>>()
.join(" ")
})
.filter(|s| !s.is_empty())
.unwrap_or_default(),
"agent.tool_use" => {
let name = v.get("name").and_then(|x| x.as_str()).unwrap_or("?");
format!("[tool {name}]")
}
"agent.tool_result" => {
let ok = v
.get("is_error")
.and_then(|x| x.as_bool())
.map(|b| !b)
.unwrap_or(true);
if ok {
"[tool ok]".to_string()
} else {
"[tool error]".to_string()
}
}
"user.message" => v
.get("content")
.and_then(|c| c.as_array())
.map(|arr| {
arr.iter()
.filter_map(|b| b.get("text").and_then(|t| t.as_str()))
.collect::<Vec<_>>()
.join(" ")
})
.filter(|s| !s.is_empty())
.map(|s| format!("[user] {s}"))
.unwrap_or_default(),
"session.status_idle" => "[idle]".to_string(),
"session.status_run_started" => "[run started]".to_string(),
"session.status_run_completed" => "[run completed]".to_string(),
"session.status_run_failed" => "[run FAILED]".to_string(),
"" => String::new(),
other => format!("[{other}]"),
}
}
fn short_id(id: &str) -> String {
let n = id.chars().count();
if n <= 14 {
return id.to_string();
}
let prefix: String = id.chars().take(4).collect();
let suffix: String = id.chars().skip(n.saturating_sub(8)).collect();
format!("{prefix}…{suffix}")
}
pub fn collect_managed_agent_rows() -> Vec<crate::claude_agents::AgentRow> {
use crate::claude_agents::{AgentRow, AgentSource, AgentState};
use std::path::PathBuf;
let backend = match detect_backend() {
Ok(b) => b,
Err(_) => return Vec::new(),
};
let sessions = match list_sessions(&backend) {
Ok(s) => s,
Err(e) => {
let _ = e;
return Vec::new();
}
};
sessions
.into_iter()
.map(|s| {
let state = match s.status.as_str() {
"in_progress" => AgentState::Streaming,
"pending" => AgentState::ToolCall,
"idle" => AgentState::Idle,
"completed" | "failed" | "cancelled" => AgentState::Ended,
_ => AgentState::Idle,
};
let workspace = s
.environment_id
.clone()
.filter(|e| !e.is_empty())
.map(|e| short_id(&e))
.unwrap_or_else(|| "managed".to_string());
AgentRow {
source: AgentSource::AnthropicManaged,
transcript_path: PathBuf::from(format!("/dev/null/managed/{}", s.id)),
session_id: s.id,
workspace,
cwd: None,
git_branch: None,
model: None,
last_activity: None,
tokens: 0,
input_tokens: 0,
output_tokens: 0,
cache_create_tokens: 0,
cache_read_tokens: 0,
cost_usd: 0.0,
event_count: 0,
last_user_msg: None,
last_assistant_msg: s.title.or_else(|| Some(s.status.clone())),
pid: None,
state,
current_tool: None,
todos: Vec::new(),
recent_bash: Vec::new(),
recent_files: Vec::new(),
recent_subagents: Vec::new(),
pending_tool_uses: 0,
tokens_per_min: None,
}
})
.collect()
}