use anyhow::{Context, Result};
use std::io::{BufRead, BufReader};
use std::time::Duration;
const BASE_URL: &str = "http://127.0.0.1:11434";
fn agent() -> ureq::Agent {
let config = ureq::Agent::config_builder()
.timeout_connect(Some(Duration::from_secs(5)))
.timeout_recv_response(Some(Duration::from_secs(10)))
.http_status_as_error(false)
.build();
ureq::Agent::new_with_config(config)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OllamaModel {
pub name: String,
pub size_bytes: Option<u64>,
pub parameter_size: Option<String>,
pub quantization: Option<String>,
}
impl OllamaModel {
pub fn display_row(&self) -> String {
let mut parts = vec![self.name.clone()];
if let Some(bytes) = self.size_bytes {
parts.push(format_size(bytes));
}
if let Some(quant) = &self.quantization {
parts.push(quant.clone());
}
parts.join(" ")
}
}
fn format_size(bytes: u64) -> String {
const GB: f64 = 1e9;
const MB: f64 = 1e6;
let bytes = bytes as f64;
if bytes >= GB {
format!("{:.1}GB", bytes / GB)
} else {
format!("{:.0}MB", bytes / MB)
}
}
fn parse_tags_response(body: &str) -> Result<Vec<OllamaModel>> {
let parsed: serde_json::Value =
serde_json::from_str(body).context("parsing /api/tags response as JSON")?;
let models_value = parsed
.get("models")
.ok_or_else(|| anyhow::anyhow!("malformed /api/tags response: missing \"models\" field"))?;
let models = models_value.as_array().ok_or_else(|| {
anyhow::anyhow!(
"malformed /api/tags response: \"models\" is not an array (got {})",
models_value
)
})?;
Ok(models
.iter()
.filter_map(|entry| {
let name = entry.get("name")?.as_str()?.to_string();
let size_bytes = entry.get("size").and_then(serde_json::Value::as_u64);
let details = entry.get("details");
let parameter_size = details
.and_then(|d| d.get("parameter_size"))
.and_then(serde_json::Value::as_str)
.map(str::to_string);
let quantization = details
.and_then(|d| d.get("quantization_level"))
.and_then(serde_json::Value::as_str)
.map(str::to_string);
Some(OllamaModel {
name,
size_bytes,
parameter_size,
quantization,
})
})
.collect())
}
pub fn list_installed() -> Result<Vec<OllamaModel>> {
list_installed_from(BASE_URL)
}
fn list_installed_from(base_url: &str) -> Result<Vec<OllamaModel>> {
let mut response = agent()
.get(format!("{base_url}/api/tags"))
.call()
.context("is the Ollama daemon running? (`ollama serve`)")?;
if !response.status().is_success() {
let detail = super::provider::read_error_body(&mut response);
anyhow::bail!(
"GET /api/tags failed ({}): {detail}",
response.status().as_u16()
);
}
let body = response
.body_mut()
.read_to_string()
.context("reading /api/tags response body")?;
parse_tags_response(&body)
}
fn parse_ps_response(body: &str) -> Result<Vec<String>> {
let parsed: serde_json::Value =
serde_json::from_str(body).context("parsing /api/ps response as JSON")?;
let models_value = parsed
.get("models")
.ok_or_else(|| anyhow::anyhow!("malformed /api/ps response: missing \"models\" field"))?;
let models = models_value.as_array().ok_or_else(|| {
anyhow::anyhow!(
"malformed /api/ps response: \"models\" is not an array (got {})",
models_value
)
})?;
Ok(models
.iter()
.filter_map(|entry| entry.get("name")?.as_str().map(str::to_string))
.collect())
}
pub fn running_models() -> Result<Vec<String>> {
running_models_from(BASE_URL)
}
fn running_models_from(base_url: &str) -> Result<Vec<String>> {
let mut response = agent()
.get(format!("{base_url}/api/ps"))
.call()
.context("is the Ollama daemon running? (`ollama serve`)")?;
if !response.status().is_success() {
let detail = super::provider::read_error_body(&mut response);
anyhow::bail!(
"GET /api/ps failed ({}): {detail}",
response.status().as_u16()
);
}
let body = response
.body_mut()
.read_to_string()
.context("reading /api/ps response body")?;
parse_ps_response(&body)
}
#[derive(Debug, Clone)]
pub enum PullEvent {
Status(String),
Progress { status: String, percent: u8 },
Done,
Error(String),
}
fn parse_pull_line(line: &str) -> Option<PullEvent> {
let value: serde_json::Value = serde_json::from_str(line).ok()?;
if let Some(error) = value.get("error").and_then(serde_json::Value::as_str) {
return Some(PullEvent::Error(error.to_string()));
}
let status = value.get("status").and_then(serde_json::Value::as_str)?;
if status == "success" {
return Some(PullEvent::Done);
}
let completed = value.get("completed").and_then(serde_json::Value::as_u64);
let total = value.get("total").and_then(serde_json::Value::as_u64);
match (completed, total) {
(Some(completed), Some(total)) if total > 0 => Some(PullEvent::Progress {
status: status.to_string(),
percent: ((completed as f64 / total as f64) * 100.0).round() as u8,
}),
_ => Some(PullEvent::Status(status.to_string())),
}
}
pub fn pull(name: &str, mut on_event: impl FnMut(PullEvent)) -> Result<()> {
let body = serde_json::json!({ "name": name });
let pull_agent_config = ureq::Agent::config_builder()
.timeout_connect(Some(Duration::from_secs(5)))
.timeout_recv_response(Some(Duration::from_secs(30)))
.http_status_as_error(false)
.build();
let response = ureq::Agent::new_with_config(pull_agent_config)
.post(format!("{BASE_URL}/api/pull"))
.send_json(&body)
.context("is the Ollama daemon running? (`ollama serve`)")?;
let reader = BufReader::new(response.into_body().into_reader());
let mut saw_done = false;
let mut last_error: Option<String> = None;
for line in reader.lines() {
let line = line.context("reading /api/pull response stream")?;
let line = line.trim();
if line.is_empty() {
continue;
}
let Some(event) = parse_pull_line(line) else {
continue;
};
match &event {
PullEvent::Done => saw_done = true,
PullEvent::Error(message) => last_error = Some(message.clone()),
_ => {}
}
on_event(event);
}
if let Some(message) = last_error {
anyhow::bail!("pull failed: {message}");
}
if !saw_done {
anyhow::bail!("pull stream ended without a success status — connection dropped?");
}
Ok(())
}
pub fn delete(name: &str) -> Result<()> {
let body = serde_json::json!({ "name": name });
let mut response = agent()
.delete(format!("{BASE_URL}/api/delete"))
.force_send_body()
.send_json(&body)
.context("is the Ollama daemon running? (`ollama serve`)")?;
if !response.status().is_success() {
let detail = super::provider::read_error_body(&mut response);
anyhow::bail!("delete failed ({}): {detail}", response.status().as_u16());
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_tags_response_with_full_and_partial_details() {
let body = serde_json::json!({
"models": [
{
"name": "llama3.2:latest",
"size": 2_019_393_189u64,
"details": { "parameter_size": "3.2B", "quantization_level": "Q4_K_M" }
},
{ "name": "no-details-model" }
]
})
.to_string();
let models = parse_tags_response(&body).unwrap();
assert_eq!(models.len(), 2);
assert_eq!(models[0].name, "llama3.2:latest");
assert_eq!(models[0].size_bytes, Some(2_019_393_189));
assert_eq!(models[0].quantization.as_deref(), Some("Q4_K_M"));
assert_eq!(models[1].name, "no-details-model");
assert_eq!(models[1].size_bytes, None);
assert_eq!(models[1].quantization, None);
}
#[test]
fn empty_tags_response_is_an_empty_list_not_an_error() {
let body = serde_json::json!({ "models": [] }).to_string();
assert_eq!(parse_tags_response(&body).unwrap(), vec![]);
}
#[test]
fn missing_models_key_is_malformed_not_empty() {
let error_body = serde_json::json!({
"error": "model \"ghost\" not found, try pulling it first"
})
.to_string();
let error = parse_tags_response(&error_body)
.expect_err("a body with no \"models\" key must be Err, not Ok(vec![])");
assert!(error.to_string().contains("missing"));
}
#[test]
fn wrong_type_models_field_is_malformed_not_empty() {
let body = serde_json::json!({ "models": {} }).to_string();
let error = parse_tags_response(&body)
.expect_err("\"models\": {} must be Err (wrong type), not Ok(vec![])");
assert!(error.to_string().contains("not an array"));
}
#[test]
fn wrong_type_models_field_is_malformed_not_empty_ps() {
let body = serde_json::json!({ "models": "not-an-array" }).to_string();
let error = parse_ps_response(&body)
.expect_err("\"models\": \"...\" must be Err (wrong type), not Ok(vec![])");
assert!(error.to_string().contains("not an array"));
}
#[test]
fn missing_models_key_is_malformed_not_empty_ps() {
let body = serde_json::json!({ "unexpected": [] }).to_string();
let error = parse_ps_response(&body)
.expect_err("a body with no \"models\" key must be Err, not Ok(vec![])");
assert!(error.to_string().contains("missing"));
}
#[test]
fn display_row_omits_missing_fields_without_stray_whitespace() {
let full = OllamaModel {
name: "llama3.2:latest".to_string(),
size_bytes: Some(2_000_000_000),
parameter_size: Some("3.2B".to_string()),
quantization: Some("Q4_K_M".to_string()),
};
assert_eq!(full.display_row(), "llama3.2:latest 2.0GB Q4_K_M");
let bare = OllamaModel {
name: "no-details-model".to_string(),
size_bytes: None,
parameter_size: None,
quantization: None,
};
assert_eq!(bare.display_row(), "no-details-model");
assert_eq!(
bare.display_row().split_whitespace().next(),
Some("no-details-model")
);
}
#[test]
fn parses_ps_response() {
let body = serde_json::json!({
"models": [ { "name": "llama3.2:latest" }, { "name": "qwen2.5:7b" } ]
})
.to_string();
assert_eq!(
parse_ps_response(&body).unwrap(),
vec!["llama3.2:latest".to_string(), "qwen2.5:7b".to_string()]
);
}
#[test]
fn empty_ps_response_is_an_empty_list() {
let body = serde_json::json!({ "models": [] }).to_string();
assert_eq!(parse_ps_response(&body).unwrap(), Vec::<String>::new());
}
#[test]
fn parses_pull_progress_and_terminal_events() {
assert!(matches!(
parse_pull_line(r#"{"status":"pulling manifest"}"#),
Some(PullEvent::Status(s)) if s == "pulling manifest"
));
match parse_pull_line(r#"{"status":"downloading","completed":50,"total":200}"#) {
Some(PullEvent::Progress { status, percent }) => {
assert_eq!(status, "downloading");
assert_eq!(percent, 25);
}
other => panic!("expected Progress, got {other:?}"),
}
assert!(matches!(
parse_pull_line(r#"{"status":"success"}"#),
Some(PullEvent::Done)
));
assert!(matches!(
parse_pull_line(r#"{"error":"model not found"}"#),
Some(PullEvent::Error(e)) if e == "model not found"
));
assert!(parse_pull_line("not json at all").is_none());
}
#[test]
fn zero_total_falls_back_to_status_not_a_division_by_zero() {
assert!(matches!(
parse_pull_line(r#"{"status":"verifying digest","completed":0,"total":0}"#),
Some(PullEvent::Status(s)) if s == "verifying digest"
));
}
use std::io::{BufRead, BufReader, Write};
use std::net::TcpListener;
fn spawn_one_shot_response(
status_line: &'static str,
body: String,
) -> (String, std::thread::JoinHandle<()>) {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock ollama listener");
let port = listener.local_addr().expect("local addr").port();
let handle = std::thread::spawn(move || {
let (stream, _) = listener.accept().expect("accept connection");
let mut reader = BufReader::new(&stream);
loop {
let mut line = String::new();
reader.read_line(&mut line).expect("read header line");
if line == "\r\n" || line.is_empty() {
break;
}
}
let mut writer = &stream;
write!(
writer,
"{status_line}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
)
.expect("write mock response");
writer.flush().ok();
});
(format!("http://127.0.0.1:{port}"), handle)
}
#[test]
fn list_installed_from_surfaces_a_real_500_as_an_error_not_empty_ok() {
let error_body = serde_json::json!({
"error": "model \"ghost\" not found, try pulling it first"
})
.to_string();
let (base_url, handle) =
spawn_one_shot_response("HTTP/1.1 500 Internal Server Error", error_body);
let result = list_installed_from(&base_url);
handle.join().expect("mock server thread panicked");
let error = result.expect_err(
"a live 500 with a real Ollama error body must surface as Err, \
not silently become Ok(vec![]) as it did before the Section 8 fix",
);
let message = error.to_string();
assert!(
message.contains("500"),
"error should mention the real HTTP status: {message}"
);
}
#[test]
fn list_installed_from_returns_models_on_a_real_200() {
let body = serde_json::json!({
"models": [{"name": "llama3.2:latest"}]
})
.to_string();
let (base_url, handle) = spawn_one_shot_response("HTTP/1.1 200 OK", body);
let models = list_installed_from(&base_url).expect("200 with valid body must succeed");
handle.join().expect("mock server thread panicked");
assert_eq!(models.len(), 1);
assert_eq!(models[0].name, "llama3.2:latest");
}
#[test]
fn list_installed_from_surfaces_a_real_200_with_malformed_body_as_an_error() {
let body = serde_json::json!({ "unexpected_shape": true }).to_string();
let (base_url, handle) = spawn_one_shot_response("HTTP/1.1 200 OK", body);
let error = list_installed_from(&base_url)
.expect_err("a 200 with no \"models\" key must be Err, not Ok(vec![])");
handle.join().expect("mock server thread panicked");
assert!(error.to_string().contains("missing"));
}
#[test]
fn running_models_from_surfaces_a_real_404_as_an_error_not_empty_ok() {
let error_body = serde_json::json!({ "error": "not found" }).to_string();
let (base_url, handle) = spawn_one_shot_response("HTTP/1.1 404 Not Found", error_body);
let result = running_models_from(&base_url);
handle.join().expect("mock server thread panicked");
let error = result.expect_err("a live 404 must surface as Err, not Ok(vec![])");
assert!(error.to_string().contains("404"));
}
}