Merge remote-tracking branch 'origin/rustify' into worktree-agent-a27094a7db775552a
# Conflicts: # AGENTS.md # server/src/session/driver.rs
This commit is contained in:
commit
88631f5e8b
216 files changed
+52342
-461
No files matched your search
Generated
+9
@@ -26,6 +26,7 @@ dependencies = [
|
||||
"axum-server",
|
||||
"base64 0.23.1",
|
||||
"clap",
|
||||
"event-model",
|
||||
"libc",
|
||||
"rand",
|
||||
"ron",
|
||||
@@ -544,6 +545,14 @@ dependencies = [
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "event-model"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fastrand"
|
||||
version = "2.5.0"
|
||||
|
||||
@@ -13,6 +13,10 @@ path = "src/main.rs"
|
||||
# the RON house rules. Extracted from the two copies that had drifted --
|
||||
# see that repo's README for the evidence and the bug the extraction found.
|
||||
wg-app-link = { path = "../wg-app-link/server" }
|
||||
# The event model, shared with `client-core` so a Rust client and this
|
||||
# server read the same `Event`/`SeqEvent` rather than the app hand-mirroring
|
||||
# it the way `Events.kt` used to.
|
||||
event-model = { path = "../event-model" }
|
||||
axum = { version = "0.8", features = ["json", "multipart"] }
|
||||
axum-server = { version = "0.8", features = ["tls-rustls"] }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "sync", "time", "process", "io-util", "signal"] }
|
||||
|
||||
+11
-373
@@ -5,364 +5,20 @@
|
||||
//! and accepts the small inbound vocabulary below. The transcript, the SSE
|
||||
//! stream, and the phone UI work purely in this model; nothing downstream
|
||||
//! of a driver may branch on the session kind.
|
||||
//!
|
||||
//! The event model itself -- [`Event`], [`QuestionOption`], [`SessionStatus`],
|
||||
//! [`AttachmentRef`], `ImageRef`, [`context_tokens`] and [`context_after`] --
|
||||
//! moved to the `event-model` crate on 2026-09-04, so `client-core` can share
|
||||
//! one definition with this server instead of a hand-kept Kotlin mirror.
|
||||
//! Re-exported here so nothing downstream of this module had to change; what
|
||||
//! stayed behind is the *driver* abstraction, which is how this server runs
|
||||
//! a session rather than part of what a client reads off the wire.
|
||||
pub use event_model::{
|
||||
AttachmentRef, Event, QuestionOption, SessionStatus, context_after, context_tokens,
|
||||
};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
/// The name a session's image is stored and served under -- minted for an
|
||||
/// upload or for one a tool produced, and fetched back from
|
||||
/// `/sessions/{id}/files/{ref}`. One id both directions, so the transcript
|
||||
/// renders them identically.
|
||||
pub type ImageRef = String;
|
||||
|
||||
/// The name an upload is stored and served under: an image is
|
||||
/// `<hex>.<extension>` and is an [`ImageRef`] like any other; any other file
|
||||
/// keeps its own name after the hex, `<hex>-<name>`, because the name is what
|
||||
/// the reader attached and what the session is told. Told apart by
|
||||
/// `crate::media::media_type_for`.
|
||||
pub type AttachmentRef = String;
|
||||
|
||||
/// One choice offered in answer to a [`Event::Question`]. More than a label
|
||||
/// because the reader is deciding rather than confirming: what an option
|
||||
/// means, and what picking it would produce, are what decide it.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct QuestionOption {
|
||||
pub label: String,
|
||||
/// A sentence about what this option means.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub description: Option<String>,
|
||||
/// A block to show as written -- a mockup, a diff, a config file.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub preview: Option<String>,
|
||||
}
|
||||
|
||||
impl QuestionOption {
|
||||
pub fn plain(label: impl Into<String>) -> Self {
|
||||
Self {
|
||||
label: label.into(),
|
||||
description: None,
|
||||
preview: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Everything a session can tell the outside world. Every event is appended
|
||||
/// to the transcript with a sequence number, then fanned out to SSE
|
||||
/// subscribers, so reconnecting is just "events after seq N" -- no separate
|
||||
/// history path to drift from the live one.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
// `rename_all` renames the variants; `rename_all_fields` renames what is
|
||||
// inside them. Both are needed and only the first is obvious: every field
|
||||
// here was one lowercase word until `pre_tokens` arrived, so a multi-word
|
||||
// field went out as snake_case, the app looked for camelCase and found
|
||||
// nothing, and the event still rendered -- as the "no counts reported" case,
|
||||
// which is a state it is allowed to be in.
|
||||
#[serde(
|
||||
tag = "type",
|
||||
rename_all = "camelCase",
|
||||
rename_all_fields = "camelCase"
|
||||
)]
|
||||
pub enum Event {
|
||||
/// What the user sent, written into the transcript by the manager (not by
|
||||
/// drivers) so every device renders the conversation from one stream.
|
||||
/// Recorded when the session reads it, which is what `MessageTaken` reports.
|
||||
UserMessage {
|
||||
/// The [`Event::MessageQueued`] this resolves, when it waited. The
|
||||
/// phone has a bubble on screen for the waiting message and needs to
|
||||
/// know *which* one this is, rather than matching on the text and
|
||||
/// clearing the wrong one when the same thing was sent twice.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
id: Option<String>,
|
||||
text: String,
|
||||
/// What was attached, by the ref the files route serves. On the
|
||||
/// message rather than beside it: these used to be their own `Image`
|
||||
/// events just before, which left the phone deciding from adjacency
|
||||
/// which message an image belonged to. `images` on disk until
|
||||
/// 2026-09-03, when files joined them; the alias reads the older rows.
|
||||
#[serde(default, alias = "images", skip_serializing_if = "Vec::is_empty")]
|
||||
attachments: Vec<AttachmentRef>,
|
||||
},
|
||||
/// A message accepted from the phone that the session cannot read yet.
|
||||
///
|
||||
/// Recorded, unlike the message itself, and that difference is the point:
|
||||
/// the message belongs in the transcript where the session read it, but
|
||||
/// something has to say it is waiting, and it has to be the server. The
|
||||
/// phone used to remember its own outgoing messages, so leaving the
|
||||
/// screen showed nothing pending when something was.
|
||||
///
|
||||
/// Carries no row of its own; resolved by the `UserMessage` bearing the
|
||||
/// same id, as `CommandQueued` is resolved by `CommandSent`.
|
||||
MessageQueued {
|
||||
id: String,
|
||||
text: String,
|
||||
/// Carried for the same reason [`Event::UserMessage`] carries it,
|
||||
/// and it matters more here: a waiting message is on screen for as
|
||||
/// long as the turn runs, so its attachment has nowhere else to be.
|
||||
#[serde(default, alias = "images", skip_serializing_if = "Vec::is_empty")]
|
||||
attachments: Vec<AttachmentRef>,
|
||||
},
|
||||
/// A message taken out of the queue before the session read it.
|
||||
///
|
||||
/// Recorded for the same reason `MessageQueued` is: the queue is the
|
||||
/// server's, so what is waiting has to be answerable from the transcript
|
||||
/// alone. Without it a phone that reconnects replays the `MessageQueued`
|
||||
/// and puts back a bubble nothing will ever resolve -- the `UserMessage`
|
||||
/// that normally does is exactly what is not coming.
|
||||
///
|
||||
/// Only ever sent for a message that had not been handed over; see
|
||||
/// [`Unqueued::AlreadySent`].
|
||||
MessageDropped {
|
||||
id: String,
|
||||
},
|
||||
/// A driver has taken one of the user's messages and started reading it.
|
||||
/// The manager turns this into the `UserMessage` above, so it never
|
||||
/// reaches a phone itself.
|
||||
///
|
||||
/// It exists because sending and being read are not the same moment. A
|
||||
/// message sent into a running turn waits, and recording it among things
|
||||
/// already read puts it in the transcript above output that predates it.
|
||||
MessageTaken {
|
||||
/// The `MessageQueued` this answers, or `None` when it never waited.
|
||||
/// Carried through onto the `UserMessage`.
|
||||
id: Option<String>,
|
||||
text: String,
|
||||
#[serde(default, alias = "images", skip_serializing_if = "Vec::is_empty")]
|
||||
attachments: Vec<AttachmentRef>,
|
||||
},
|
||||
/// Streaming assistant text; the phone renders the concatenation as
|
||||
/// markdown.
|
||||
AssistantText {
|
||||
delta: String,
|
||||
},
|
||||
ToolStart {
|
||||
id: String,
|
||||
tool: String,
|
||||
input: serde_json::Value,
|
||||
},
|
||||
ToolUpdate {
|
||||
id: String,
|
||||
output: String,
|
||||
},
|
||||
ToolEnd {
|
||||
id: String,
|
||||
output: String,
|
||||
},
|
||||
/// An image the session produced or was sent, saved under the session
|
||||
/// dir and referenced by id; the phone fetches it by URL.
|
||||
Image {
|
||||
#[serde(rename = "ref")]
|
||||
image: ImageRef,
|
||||
/// The tool call whose result carried it, when one did. A screenshot
|
||||
/// belongs under the call that took it, not floating beside it -- the
|
||||
/// reader has to pair them by position otherwise, and position is
|
||||
/// exactly what a page boundary breaks.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
about: Option<String>,
|
||||
},
|
||||
/// Anything the session needs a human for: AskUserQuestion, and
|
||||
/// permission requests, are the same shape with different options.
|
||||
Question {
|
||||
id: String,
|
||||
prompt: String,
|
||||
/// A few words naming what the question is about, when the asker
|
||||
/// offered one. `None` for a permission, which is about the call
|
||||
/// above it.
|
||||
header: Option<String>,
|
||||
options: Vec<QuestionOption>,
|
||||
/// Whether several options may be chosen at once. Here rather than
|
||||
/// left for a phone to work out from the dialect underneath: how many
|
||||
/// answers a question takes is a fact about the question, and the
|
||||
/// alternative was Claude Code's tool-input schema written out a
|
||||
/// second time in Kotlin, where no other dialect could reach it.
|
||||
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
|
||||
multi_select: bool,
|
||||
/// The tool call this is permission for, when it is one, so a phone
|
||||
/// can draw the ask on the tool's own row rather than as a second
|
||||
/// card repeating its input. `None` for anything not about a tool.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
about: Option<String>,
|
||||
},
|
||||
/// A message another agent sent this session.
|
||||
///
|
||||
/// Its own kind rather than a `UserMessage`, because it is not something
|
||||
/// the reader said and a transcript that renders it in their voice is
|
||||
/// claiming they did. It also explains what would otherwise be
|
||||
/// inexplicable: a session working on something nobody here asked for.
|
||||
PeerMessage {
|
||||
/// The sending session's own name, which is what the reader
|
||||
/// recognises it by -- the socket path it came from is not.
|
||||
from: String,
|
||||
text: String,
|
||||
/// The seq of the `Status::Running` that opened the turn this message
|
||||
/// started, so a reader can draw it above that turn.
|
||||
///
|
||||
/// The CLI says nothing about a peer message until the turn's
|
||||
/// `result`, so the event is appended after everything it caused, and
|
||||
/// an append-only transcript cannot go back and insert it. Carrying
|
||||
/// the position instead keeps one order on the wire and one on screen.
|
||||
///
|
||||
/// Filled in by the pump, the only place that knows a seq, and only
|
||||
/// where a turn was open: `None` for a message replayed by `import`,
|
||||
/// which already has it in the right place.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
turn_start: Option<u64>,
|
||||
},
|
||||
/// The manager's record of a question being answered, so a rendered
|
||||
/// question card resolves on every device rather than only the one that
|
||||
/// answered.
|
||||
///
|
||||
/// A list because a question can take several answers, and one that took
|
||||
/// one is the list of length one rather than a different shape.
|
||||
Answered {
|
||||
id: String,
|
||||
answers: Vec<String>,
|
||||
},
|
||||
Status {
|
||||
state: SessionStatus,
|
||||
},
|
||||
/// What the session is set to, as the session itself reports it.
|
||||
///
|
||||
/// Asking for a change and having one are different things, and only this
|
||||
/// is a measurement: a model name the dialect does not know, a mode it
|
||||
/// refuses, or a driver whose model is fixed at startup all leave a
|
||||
/// request that was sent and nothing that changed. Reporting from the
|
||||
/// request put the answer on the phone before the question was answered.
|
||||
///
|
||||
/// Either field alone, because the two are confirmed separately.
|
||||
Settings {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
model: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
permission_mode: Option<String>,
|
||||
},
|
||||
/// Per-turn token counts, where the dialect reports them.
|
||||
UsageDelta {
|
||||
/// What this turn cost: the tokens it was charged for.
|
||||
tokens: u64,
|
||||
/// What the model was holding when the turn ended -- see
|
||||
/// [`context_tokens`].
|
||||
///
|
||||
/// Carried rather than summed by whoever is reading, because it is
|
||||
/// not a sum: context goes *down* at a compaction and a clear, so
|
||||
/// adding turns up would report a figure the session stopped being
|
||||
/// true of long ago.
|
||||
///
|
||||
/// `None` where the dialect did not say, which every reader has to be
|
||||
/// able to draw.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
context: Option<u64>,
|
||||
},
|
||||
/// A compaction that finished, and how much context it recovered.
|
||||
///
|
||||
/// The counts are the point, and a spinner is not. They are optional
|
||||
/// because the record has shipped without them, and "the compaction
|
||||
/// happened, we don't know by how much" is a state this has to be able to
|
||||
/// say -- a plausible number would be indistinguishable from a counted one.
|
||||
Compacted {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pre_tokens: Option<u64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
post_tokens: Option<u64>,
|
||||
/// What asked for it, in the dialect's own word -- `auto` when the
|
||||
/// session compacted on its own. Carried rather than reduced to a bool
|
||||
/// so an unrecognised trigger stays unrecognised: an automatic
|
||||
/// compaction is the one worth naming, because it explains a wait
|
||||
/// nobody asked for.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
trigger: Option<String>,
|
||||
},
|
||||
/// A command the session was asked to run on itself, held because it
|
||||
/// cannot run yet. These are not messages: `/compact` and `/rename` are
|
||||
/// instructions about the session, and a session mid-turn reads a line
|
||||
/// written to it as something the model should see. So they wait, and
|
||||
/// this is what a phone draws while they do.
|
||||
CommandQueued {
|
||||
id: String,
|
||||
text: String,
|
||||
},
|
||||
/// The same command, now handed to the session. Its [`CommandQueued`]
|
||||
/// stops being pending when this arrives, matched by `id`; a command
|
||||
/// that ran immediately has only this.
|
||||
CommandSent {
|
||||
id: String,
|
||||
text: String,
|
||||
},
|
||||
/// The conversation was cleared: everything above this is still in the
|
||||
/// record but is no longer in the session's context.
|
||||
///
|
||||
/// Nothing is deleted. A transcript is the thing a person scrolls back
|
||||
/// through, so this is a divider, not a truncation.
|
||||
///
|
||||
/// **Load-bearing, not decorative.** For any driver that rebuilds its
|
||||
/// conversation from the transcript, this marker decides what the model
|
||||
/// is given -- dropping it, or treating it as something only the phone
|
||||
/// draws, silently puts a cleared conversation back in front of the model
|
||||
/// at full cost. Today `llama::conversation` is the only fold that reads
|
||||
/// it, which is why this is written down rather than left to be inferred
|
||||
/// from a second example that does not exist.
|
||||
Cleared,
|
||||
/// The account behind this session has no quota left, so the turn stopped
|
||||
/// without finishing.
|
||||
///
|
||||
/// Its own event rather than an [`Event::Error`] carrying the dialect's
|
||||
/// sentence, because two things act on it that cannot read English: the
|
||||
/// transcript draws it as a state the session is in rather than as a
|
||||
/// failure of something it did, and `crate::resume` schedules the message
|
||||
/// that picks the work back up. Recognising it belongs to the driver, which
|
||||
/// is the only layer that knows its dialect's wording -- above here nothing
|
||||
/// matches on strings.
|
||||
///
|
||||
/// `resets_at` is epoch seconds, and `None` is a real state: the dialect
|
||||
/// said the limit was hit without saying when it lifts. Nothing here
|
||||
/// invents one -- what the wait is actually decided against is the usage
|
||||
/// endpoint, and this is the hint that starts the waiting.
|
||||
LimitReached {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
resets_at: Option<f64>,
|
||||
},
|
||||
Error {
|
||||
message: String,
|
||||
},
|
||||
}
|
||||
|
||||
/// How much the model was holding, from the three figures a turn reports:
|
||||
/// the input side only, prompt plus both cache figures. A cached token is
|
||||
/// cheaper but it is still one the model was given; output is what the turn
|
||||
/// produced rather than what continuing has to carry.
|
||||
///
|
||||
/// One function so the definition cannot drift, because it is extracted two
|
||||
/// quite different ways -- the live translators have the usage object parsed,
|
||||
/// and `import::context_tokens` scans it out of a raw line without parsing.
|
||||
pub fn context_tokens(input: u64, cache_creation: u64, cache_read: u64) -> u64 {
|
||||
input + cache_creation + cache_read
|
||||
}
|
||||
|
||||
/// The context after `event`, given what it was before.
|
||||
///
|
||||
/// The whole rule in one place, because three readers need the same answer:
|
||||
/// the pump keeping a live session's figure, the transcript seeding it at
|
||||
/// startup, and the phone folding the same events into what it draws.
|
||||
///
|
||||
/// The two that *lower* it are the point. A clear takes the conversation away
|
||||
/// and a compaction replaces it with a summary, so a figure measured before
|
||||
/// either stopped being true at that moment -- and carrying it forward is how
|
||||
/// a session that had just been cleared went on reporting the context it no
|
||||
/// longer had.
|
||||
///
|
||||
/// `None` is "we don't know", which each of them can reach.
|
||||
pub fn context_after(current: Option<u64>, event: &Event) -> Option<u64> {
|
||||
match event {
|
||||
// `or`, so a turn the dialect reported no usage for leaves the last
|
||||
// measurement standing: stale by a turn, which every context figure
|
||||
// is, rather than wrong.
|
||||
Event::UsageDelta { context, .. } => context.or(current),
|
||||
Event::Compacted { post_tokens, .. } => *post_tokens,
|
||||
Event::Cleared => None,
|
||||
_ => current,
|
||||
}
|
||||
}
|
||||
|
||||
/// Something a session can be asked to do to itself.
|
||||
///
|
||||
/// A closed set rather than a string, because the two that are not
|
||||
@@ -401,24 +57,6 @@ impl SessionCommand {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub enum SessionStatus {
|
||||
Idle,
|
||||
Running,
|
||||
AwaitingInput,
|
||||
Compacting,
|
||||
Exited,
|
||||
/// There is a process recorded for this session and the machine will not
|
||||
/// say whether it is still running.
|
||||
///
|
||||
/// Its own state rather than the nearest of the others, because both
|
||||
/// neighbours are lies with consequences: `Exited` invites starting a
|
||||
/// second process against a conversation that may already have one, and
|
||||
/// `Idle` claims a session is waiting for you when nobody has checked.
|
||||
Unknown,
|
||||
}
|
||||
|
||||
/// What became of a request to take a queued message back.
|
||||
///
|
||||
/// Three states rather than a bool because the two failures are not the same
|
||||
|
||||
@@ -13,20 +13,14 @@ use std::os::unix::fs::OpenOptionsExt;
|
||||
use std::path::Path;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::driver::{Event, SessionStatus, context_after};
|
||||
|
||||
/// One transcript line: an [`Event`] plus its position and time. The event
|
||||
/// is flattened so the wire shape stays one flat object.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct SeqEvent {
|
||||
pub seq: u64,
|
||||
/// Epoch seconds.
|
||||
pub ts: f64,
|
||||
#[serde(flatten)]
|
||||
pub event: Event,
|
||||
}
|
||||
// `SeqEvent` moved to `event-model` on 2026-09-04 along with the rest of the
|
||||
// event model, so `client-core` can read the same wire shape; re-exported
|
||||
// here since every caller in this crate reaches it through this module.
|
||||
pub use event_model::SeqEvent;
|
||||
|
||||
pub struct Transcript {
|
||||
file: File,
|
||||
|
||||
Reference in new issue
Block a user