Files
ai-app/server/src/session/llama.rs
T
irisandClaude Opus 5 3deeffd1e7 Run GGUF models through llama-server, and stop orphaning them
The second half of the llama.cpp work: a session can now name a downloaded
model and talk to it. `llama-server` is spawned through the same transport
as any other driver, polled until the model is loaded, then driven over its
OpenAI-compatible streaming endpoint and translated into the same events
the Claude driver emits -- so the transcript, the SSE stream and the phone
need to know nothing new.

**The conversation is rebuilt from the transcript, not held in the driver.**
llama-server is stateless between requests, so the whole history goes with
every one, and the obvious place to keep it is a Vec in the driver. That
fails the requirement: memory in a driver is invisible to a second device
and gone on restart, and this app is meant to work across devices. Reading
it back also means the model is prompted with exactly what the phone was
shown -- including a reply that was interrupted half way, which is in the
transcript because the deltas were already emitted.

That leaves the Claude driver as the odd one out rather than this one: the
CLI's memory of a conversation is a cache in front of the same transcript,
not a second truth. Said so at the top of llama.rs, because it is the sort
of inconsistency that gets "fixed" in the wrong direction.

Session settings arrive as a driver-interpreted `params` map rather than
new typed fields, so the shared schema does not grow one dialect's
vocabulary. Context size, gpu layers and threads become server flags;
temperature and the rest ride on each request, so changing them need not
reload a model.

**Also fixes an orphan this feature would have created.** Drivers set
kill_on_drop, which covers a session being deleted -- but nothing drops on
the way out of a SIGTERM, so signalling the server left its children
running. For the Claude CLI that is untidy; for a llama-server holding a
model it is gigabytes belonging to nobody. The server now stops its
sessions on SIGTERM and SIGINT. Found by killing a test server and noticing
two 600 MB processes still resident.

Remote llama sessions are refused rather than half-working: the model is
reached over HTTP, and forwarding that port to an ssh host is the "reach
this port" operation the transport does not have yet.

Verified end to end against a real model: downloaded Qwen3-0.6B Q8_0
through the app's own download route, spawned a session on it, and held a
two-turn conversation -- "my favourite colour is teal" then "what is my
favourite colour?", answered "teal", which is the transcript replay doing
its job. Token counts arrive. An earlier attempt with the IQ2_XXS quant
produced fluent nonsense, which turned out to be the quantisation rather
than the pipeline: llama-cli produces the same from that file directly.
Four unit tests cover the fold and the path guard.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017xn8nHw1tw1R6PtiY1eEtw
2026-08-28 05:23:12 -04:00

543 lines
20 KiB
Rust

//! The llama.cpp driver: a `llama-server` process per session, spoken to
//! over its OpenAI-compatible HTTP API and translated into the common
//! event model.
//!
//! Two things make this shaped differently from the Claude driver, and
//! both are worth knowing before changing anything here.
//!
//! **It is spawned but not spoken to over stdio.** The process is started
//! through the same [`Transport`] as any other, and then reached over
//! HTTP on a loopback port. That is the case the transport's doc comment
//! flags: a remote llama-server would need its port forwarded as well as
//! its command wrapped, which is not built, so a session on an ssh host
//! is refused rather than silently talking to the wrong machine.
//!
//! **The server is stateless between requests**, so the whole
//! conversation goes with every one. It is rebuilt from the session's
//! transcript rather than kept in this struct, which is not tidiness: a
//! copy in driver memory is invisible to a second device and gone when
//! this process restarts, and the app is meant to work across devices.
//! The transcript is already the source of truth for everything else, and
//! this makes it the source of truth for the prompt too.
//!
//! That leaves the Claude driver as the odd one out rather than this one:
//! the CLI's own memory of a conversation is a cache in front of the same
//! transcript, not a second truth. Anyone tempted to "fix" the
//! inconsistency should resolve it in this direction.
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use anyhow::{Context, Result, bail};
use serde::{Deserialize, Serialize};
use serde_json::json;
use super::driver::{Driver, Event, EventSink, ImageRef, SessionStatus};
use super::transport::{Launch, Transport};
use crate::config::{ProviderConfig, SessionConfig};
/// How long to wait for a model to load before giving up on it. Loading
/// is mostly disk, and a large quantised model on a cold cache is
/// genuinely slow, so this is generous -- the failure it exists for is a
/// server that will never answer, not one that is taking its time.
const READY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300);
/// One turn in the conversation this driver keeps on the server's behalf.
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Message {
role: String,
content: String,
}
pub struct LlamaDriver {
sink: EventSink,
/// Where this session's own llama-server answers.
endpoint: String,
/// Where the conversation is read back from, one line per event.
transcript: PathBuf,
/// Sampling settings chosen at spawn, sent with every request.
sampling: serde_json::Map<String, serde_json::Value>,
/// Set by [`Driver::interrupt`]; the streaming loop checks it between
/// chunks and stops, leaving what was generated in the transcript.
cancel: Arc<AtomicBool>,
/// Taken by shutdown to stop the server.
kill: Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
}
impl LlamaDriver {
pub fn spawn(
meta: &SessionConfig,
provider: &ProviderConfig,
transport: &Transport,
models_dir: &Path,
transcript: &Path,
sink: EventSink,
) -> Result<Self> {
if !matches!(transport, Transport::Here) {
bail!(
"llama.cpp sessions can only run on this machine for now: the model is served \
over HTTP, and forwarding that port to another host isn't built yet."
);
}
let model = meta.model.as_deref().context(
"a llama.cpp session needs a model -- one of the downloaded ones, by its key",
)?;
let path = model_path(models_dir, model)?;
let port = free_port().context("finding a port for llama-server")?;
let mut args: Vec<String> = vec![
"-m".into(),
path.to_string_lossy().into_owned(),
"--host".into(),
"127.0.0.1".into(),
"--port".into(),
port.to_string(),
];
// Settings that belong to the server because they decide how the
// model is loaded; the sampling ones ride on each request instead,
// so changing them later needn't reload anything.
for (key, flag) in [
("contextSize", "-c"),
("gpuLayers", "-ngl"),
("threads", "-t"),
] {
if let Some(value) = meta.params.get(key) {
args.push(flag.to_string());
args.push(value.clone());
}
}
let program = provider.command.as_deref().unwrap_or("llama-server");
let launch = Launch::new(program, args, meta.cwd.as_deref());
let mut child = transport.spawn(&launch)?;
tracing::info!(
"session {} running {program} for {model} on 127.0.0.1:{port}",
meta.id
);
let endpoint = format!("http://127.0.0.1:{port}");
let (kill_tx, kill_rx) = tokio::sync::oneshot::channel::<()>();
{
let sink = sink.clone();
let label = format!("{} ({model})", provider.name);
tokio::spawn(async move {
let status = tokio::select! {
status = child.wait() => status.ok(),
_ = kill_rx => {
let _ = child.kill().await;
None
}
};
if let Some(status) = status
&& !status.success()
{
let _ = sink.send(Event::Error {
message: format!("{label} exited: {status}"),
});
}
let _ = sink.send(Event::Status {
state: SessionStatus::Exited,
});
});
}
// Loading is slow enough to be worth saying so: the session shows
// as running until the model is in memory, then goes idle, rather
// than looking ready and refusing the first message.
let _ = sink.send(Event::Status {
state: SessionStatus::Running,
});
{
let sink = sink.clone();
let endpoint = endpoint.clone();
let model = model.to_string();
std::thread::spawn(move || match wait_until_ready(&endpoint) {
Ok(()) => {
tracing::info!("{model} loaded and answering at {endpoint}");
let _ = sink.send(Event::Status {
state: SessionStatus::Idle,
});
}
Err(err) => {
let _ = sink.send(Event::Error {
message: format!("{model} never became ready: {err:#}"),
});
let _ = sink.send(Event::Status {
state: SessionStatus::Exited,
});
}
});
}
let mut sampling = serde_json::Map::new();
for (key, field) in [
("temperature", "temperature"),
("topP", "top_p"),
("topK", "top_k"),
("maxTokens", "max_tokens"),
] {
if let Some(raw) = meta.params.get(key)
&& let Ok(number) = raw.parse::<f64>()
{
sampling.insert(field.to_string(), json!(number));
}
}
Ok(Self {
sink,
endpoint,
transcript: transcript.to_path_buf(),
sampling,
cancel: Arc::new(AtomicBool::new(false)),
kill: Mutex::new(Some(kill_tx)),
})
}
}
impl Driver for LlamaDriver {
fn send_user_message(&self, text: String, images: Vec<ImageRef>) {
if !images.is_empty() {
let _ = self.sink.send(Event::Error {
message: "this model can't be sent images".to_string(),
});
}
let sink = self.sink.clone();
let endpoint = self.endpoint.clone();
let transcript = self.transcript.clone();
let sampling = self.sampling.clone();
let cancel = Arc::clone(&self.cancel);
cancel.store(false, Ordering::Relaxed);
// Its own thread: the request blocks for as long as the model
// takes to generate, which is the whole point of streaming it.
std::thread::spawn(move || {
let _ = sink.send(Event::Status {
state: SessionStatus::Running,
});
// Everything before this message, plus this message. Read
// rather than remembered, and `text` is appended here rather
// than waited for, because the manager's UserMessage event is
// still on its way to the transcript when this runs.
let mut messages = conversation(&transcript);
messages.push(Message {
role: "user".into(),
content: text,
});
// The reply is not stored: the deltas below are the durable
// record, so the next turn reads back exactly what the phone
// was shown -- including a partial one that was interrupted.
if let Err(err) = generate(&endpoint, &messages, &sampling, &cancel, &sink) {
let _ = sink.send(Event::Error {
message: format!("{err:#}"),
});
}
let _ = sink.send(Event::Status {
state: SessionStatus::Idle,
});
});
}
fn answer_question(&self, _id: &str, _answer: &str) {
// Nothing here asks questions: this driver has no tools, so no
// permission prompts and no AskUserQuestion.
}
fn interrupt(&self) {
self.cancel.store(true, Ordering::Relaxed);
}
fn set_model(&self, _model: &str) {
let _ = self.sink.send(Event::Error {
message: "a llama.cpp session's model is fixed when it starts, because the server \
loads one model into memory. Spawn another session to use a different one."
.to_string(),
});
}
fn compact(&self) {
let _ = self.sink.send(Event::Error {
message: "llama.cpp has no compaction. When the context fills, start a new session."
.to_string(),
});
}
fn shutdown(&self) {
self.cancel.store(true, Ordering::Relaxed);
if let Some(kill) = self.kill.lock().unwrap().take() {
let _ = kill.send(());
}
}
}
/// The conversation so far, folded out of the transcript.
///
/// Consecutive `AssistantText` deltas are one assistant turn, closed by
/// the next user message -- which is also what makes an interrupted reply
/// come back as the partial text the phone actually saw, rather than
/// vanishing or being invented.
///
/// This must stay a pure function of the transcript and must never
/// re-render earlier turns. llama.cpp caches the prompt prefix, so a
/// growing conversation reprocesses almost nothing -- but only while
/// every turn is byte-identical to last time. Changing how an old turn is
/// rendered silently reprocesses the whole history on every message.
fn conversation(path: &Path) -> Vec<Message> {
let Ok(events) = crate::session::transcript::read_after(path, 0) else {
return Vec::new();
};
let mut messages: Vec<Message> = Vec::new();
let mut pending = String::new();
for event in events {
match event.event {
Event::UserMessage { text } => {
if !pending.is_empty() {
messages.push(Message {
role: "assistant".into(),
content: std::mem::take(&mut pending),
});
}
messages.push(Message {
role: "user".into(),
content: text,
});
}
Event::AssistantText { delta } => pending.push_str(&delta),
_ => {}
}
}
if !pending.is_empty() {
messages.push(Message {
role: "assistant".into(),
content: pending,
});
}
messages
}
/// Where a model key resolves to on disk, refusing anything that climbs
/// out of the models directory -- the key arrives from a phone.
fn model_path(models_dir: &Path, key: &str) -> Result<PathBuf> {
let mut path = models_dir.to_path_buf();
for part in key.split('/') {
if part.is_empty() || part == "." || part == ".." {
bail!("\"{key}\" is not a model key this can resolve");
}
path.push(part);
}
if !path.is_file() {
bail!("no downloaded model called \"{key}\" -- download it first");
}
Ok(path)
}
/// An unused loopback port, by asking the OS for one and letting it go.
///
/// Racy in principle: something else could take it between here and
/// llama-server binding. In practice nothing on this machine is hunting
/// for ports, and the alternative -- parsing the port back out of the
/// server's log -- couples us to its output format for no real gain.
fn free_port() -> Result<u16> {
let listener = std::net::TcpListener::bind("127.0.0.1:0")?;
Ok(listener.local_addr()?.port())
}
/// Polls until the server says it is ready, or gives up.
fn wait_until_ready(endpoint: &str) -> Result<()> {
let deadline = std::time::Instant::now() + READY_TIMEOUT;
let url = format!("{endpoint}/health");
loop {
if let Ok(response) = ureq::get(&url).call()
&& response.status() == 200
{
return Ok(());
}
if std::time::Instant::now() > deadline {
bail!("gave up after {}s", READY_TIMEOUT.as_secs());
}
std::thread::sleep(std::time::Duration::from_millis(250));
}
}
/// One streamed completion: posts the conversation, emits each delta as it
/// arrives. Emits rather than returns: the transcript those events land
/// in is what the next turn reads back, so there is nothing to hand up.
fn generate(
endpoint: &str,
messages: &[Message],
sampling: &serde_json::Map<String, serde_json::Value>,
cancel: &AtomicBool,
sink: &EventSink,
) -> Result<()> {
let mut body = json!({
"messages": messages,
"stream": true,
"stream_options": {"include_usage": true},
});
let map = body.as_object_mut().expect("built as an object");
for (key, value) in sampling {
map.insert(key.clone(), value.clone());
}
let mut response = ureq::post(format!("{endpoint}/v1/chat/completions"))
.header("Content-Type", "application/json")
.send_json(&body)
.context("asking llama-server to generate")?;
let reader = std::io::BufReader::new(response.body_mut().as_reader());
let mut tokens = 0u64;
for line in std::io::BufRead::lines(reader) {
if cancel.load(Ordering::Relaxed) {
break;
}
let line = line.context("reading the generation stream")?;
// Server-sent events: the payload lines are the ones that matter,
// and blank lines separate events.
let Some(payload) = line.strip_prefix("data: ") else {
continue;
};
if payload.trim() == "[DONE]" {
break;
}
let Ok(chunk) = serde_json::from_str::<serde_json::Value>(payload) else {
continue;
};
if let Some(usage) = chunk.get("usage").and_then(|u| u.get("total_tokens"))
&& let Some(total) = usage.as_u64()
{
tokens = total;
}
let delta = chunk
.get("choices")
.and_then(|c| c.get(0))
.and_then(|c| c.get("delta"))
.and_then(|d| d.get("content"))
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
if !delta.is_empty() {
let _ = sink.send(Event::AssistantText {
delta: delta.to_string(),
});
}
}
if tokens > 0 {
let _ = sink.send(Event::UsageDelta { tokens });
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::session::transcript::Transcript;
/// Writes a transcript the way the pump does, so the fold is tested
/// against the real file format rather than a hand-built vector.
fn transcript_with(events: &[Event]) -> (tempfile::TempDir, PathBuf) {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("transcript.jsonl");
let mut transcript = Transcript::open(&path).expect("open");
for event in events {
transcript.append(event.clone(), 0.0).expect("append");
}
(dir, path)
}
#[test]
fn deltas_between_user_messages_are_one_assistant_turn() {
let (_dir, path) = transcript_with(&[
Event::UserMessage {
text: "hello".into(),
},
Event::AssistantText {
delta: "hi ".into(),
},
Event::AssistantText {
delta: "there".into(),
},
Event::Status {
state: SessionStatus::Idle,
},
Event::UserMessage {
text: "again".into(),
},
Event::AssistantText {
delta: "yes".into(),
},
]);
let messages = conversation(&path);
assert_eq!(
messages
.iter()
.map(|m| (m.role.as_str(), m.content.as_str()))
.collect::<Vec<_>>(),
[
("user", "hello"),
("assistant", "hi there"),
("user", "again"),
("assistant", "yes")
],
);
}
#[test]
/// The interrupted case, which decides what a resumed conversation is
/// built from: whatever the phone was shown. The deltas that arrived
/// before the stop are in the transcript, so they are in the prompt --
/// the model is never told it said something the user did not see, and
/// never has a turn silently dropped from under it.
fn an_interrupted_reply_stays_in_the_conversation() {
let (_dir, path) = transcript_with(&[
Event::UserMessage {
text: "count".into(),
},
Event::AssistantText {
delta: "one two".into(),
},
Event::Status {
state: SessionStatus::Idle,
},
]);
let messages = conversation(&path);
assert_eq!(messages.len(), 2);
assert_eq!(messages[1].content, "one two");
}
#[test]
/// Events this driver does not produce must not disturb the fold: a
/// transcript can carry errors and status changes from a session that
/// was, say, relaunched.
fn other_events_are_not_part_of_the_conversation() {
let (_dir, path) = transcript_with(&[
Event::Status {
state: SessionStatus::Running,
},
Event::UserMessage {
text: "hello".into(),
},
Event::Error {
message: "something went wrong".into(),
},
Event::AssistantText {
delta: "still here".into(),
},
Event::UsageDelta { tokens: 12 },
]);
let messages = conversation(&path);
assert_eq!(messages.len(), 2);
assert_eq!(messages[0].content, "hello");
assert_eq!(messages[1].content, "still here");
}
#[test]
fn a_model_key_cannot_climb_out_of_the_models_directory() {
let dir = tempfile::tempdir().expect("tempdir");
for attempt in ["../../etc/passwd", "unsloth/../../escape.gguf", ""] {
assert!(
model_path(dir.path(), attempt).is_err(),
"{attempt:?} should have been refused",
);
}
}
}