Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 58 additions & 6 deletions crates/agent/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -3559,20 +3559,35 @@ the task, not the connection or the process, as the unit of work.
and the owning `sessionId`. Journal writes come from the run's event sink and from the background
resumer. They are drained by the idle loop, the busy loop and at run end, and a write whose
`sessionId` is not the session being persisted to is dropped. The journal sits beside the
pre-dispatch checkpoint, which already persisted the assistant's `tool_use`.
pre-dispatch checkpoint, which already persisted the assistant's `tool_use`. Headless `run`
journals the same entries with its persisted session (`RunJournal`, from its event sink).
- **Pending.** A journaled task is pending when three things hold: its record belongs to this session,
its `tool_use` is anywhere on the active path, and nothing answers it, neither a `tool_result` nor
a journaled `mcp_task_result`. Requiring this session's record means only the session that created
a task may resume it. A fork copies messages, not custom entries, so it inherits no journal; even a
copied journal under another session id is ignored. The call then gets the generic "interrupted"
repair.
a task may resume it. A fork copies messages, not custom entries, so it inherits no journal; a
journal that does travel under another session id (a session file copied or restored under a new
id) is ignored. The call then gets the generic "interrupted" repair. A record must also describe the
call it names (`mcp__<server>__<tool>` equal to the call's tool) and name a configured server, and a
client's `append_custom` cannot write the journal's kinds (`mcp_task`, `mcp_task_result`): a forged
record would otherwise have the resumer poll a task id of the client's choosing.
- **Resume, in the background.** When a session is loaded, and the moment a command switches to
another one (`switch_session`, `new_session`, `fork`, `clone`), `mcp_resume::Resumer` polls each of
that session's pending tasks to a terminal status through its `McpCatalog`. The previous session's
resumer is dropped. It never cancels a task, and from that moment nothing it still has in flight
reaches the clients.
- The `ttlMs` backstop counts from `createdAtMs`. A TTL that ran out while the agent was down
resolves as expired without contacting the server, even if the server is gone.
- A server that cannot be reached is not an answer. The resume redials a few times (250 ms
doubling, about 5 s); if it still cannot connect, or the connection is lost for good mid-poll,
the call gets a placeholder `tool_result` reading `[MCP task result pending] …` (the model is
told the result is pending), no result is journaled, and an `mcp_task_placeholder` entry marks
that call's `tool_result` as a placeholder, so the task stays pending. A placeholder is known by
that mark, never by its text: a real result that happens to begin the same way is an answer. A
later prompt (a fresh resumer, when the last one finished with a server unreachable) or start
(`serve`, or `run --continue`) tries again, and the real result replaces the placeholder where it
sits. A configured server that could not be dialed at startup is kept dormant in the catalog for
this.
Only a terminal status, a JSON-RPC error such as `-32602`, or an expired TTL resolves a call.
- Answered keys are seeded from the record, so the user is never asked twice.
- In-task input reaches the session's host.
- Its progress streams as `tool_progress`, starting with a "resuming" notice. That notice is not
Expand All @@ -3594,9 +3609,46 @@ the task, not the connection or the process, as the unit of work.
The session store applies every journaled result on the active path when it materializes
messages (`session_store::materialize`), so `get_messages`, the HTML export and `run --continue`
show exactly what `serve` sent the model. The file stays append-only and ids stay stable.
- **`run --continue`.** `run` resumes too: before its turn, it polls each pending task of the session it
opened (journaled by `run` or `serve`) to its result, in turn, with a line on stderr (it has no
client to stream progress to; Ctrl-C ends the wait as it ends the run), journals each result, and
splices it into the turn it sends.
- **Privacy.** Custom entries never reach the model. A task id can be a bearer token for the
server's stored state. The session file keeps it in the clear because resume needs it, and the file
is as private as the conversation it holds. The HTML export withholds `mcp_task*` entries.
server's stored state. Resume needs it, so the session keeps it, exactly as private as the
conversation it holds. A local session file is created `0600`, and every append tightens an
existing file whose mode is looser (an older version's, a restore's, a `chmod`) before writing.
In service mode the journal goes through the tenant-keyed sealed segments with the rest of the
transcript (`tests/session_segments_sealing.rs`). The HTML export withholds `mcp_task*` entries.
- **Journal authentication** (`mcp_resume::JournalAuth`, held by the session's store, which seals
what it journals and replays only what passes). The key goes where the session goes, never with
the machine, so a session resumes on another replica, another machine, a fresh `$HOME` or after
an upgrade.
- Service mode (segments sealed under the tenant key): no MAC at all. The storage already
authenticates every line, which is strictly stronger, and any replica holding the tenant key
reads the journal (`tests/mcp_tasks_service.rs`).
- Local: every entry carries a `mac`, an HMAC-SHA256 over its kind and content, keyed by 32 random
bytes in a `0600` sidecar beside the session (`<session>.jsonl.mcp-task-journal.json`: a suffix
on the whole file name, so `work.1` and `work.2` never share a key; or inside a segmented
session's directory), made by the session's first journal write, so a session that never
journals a task gets no extra file. Every first writer (threads, other processes) agrees on one
key: it is made under an exclusive lock on the sidecar and re-read once the lock is held, and
every journal write adopts the key on disk. A key under the earlier `with_extension` name
(`<session>.mcp-task-journal.json`) is read, and carried over by the next journal write. It is one of the session's sidecars, so it moves,
trashes and restores with it. On replay an entry without a valid `mac` is ignored: a line appended to the
`.jsonl` by a model with `write`/`edit` neither causes a poll nor is delivered as a result.
- A session without a key (from before per-session keys, or that has not journaled yet) is read
as before: its entries are accepted, and opening it writes nothing. Its next journal write makes
the key and records the MACs of the entries already there as accepted, so a task that resumed
before still resumes. Entries written after that must carry a `mac`.
- The residual risk, honestly: locally the key is a file the agent's user can read, and a session
that never journaled has no key yet, so a model that can also read files, or that targets such a
session, can forge an entry. A model with `write` can also delete the key or corrupt it
(overwrite it with anything that does not parse): the session is then keyless (what is there is
accepted), and its next journal write makes a new key that accepts every line already in the
file, planted ones included; and the transcript itself is unauthenticated (a model
that can write the file can plant an ordinary `tool_result`). Locally the seal turns "append a
line" into a deliberate read-then-forge; it does not close forgery. In service mode the store is
out of the tools' reach (they run in the tenant's sandbox), and forgery is closed.

**MCP Apps.** A service session's apps pool (`ServiceSession::mcp_apps_pool`) dials the same grant
connectors a second time, with `io.modelcontextprotocol/ui` advertised, the first time the session's
Expand Down
12 changes: 11 additions & 1 deletion crates/agent/src/bin/mcp_apps_fixture_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -430,7 +430,7 @@ async fn main() {
&& let Ok(marker) = std::env::var("MCP_APPS_FIXTURE_FAIL_FIRST_UI")
&& !std::path::Path::new(&marker).exists()
{
let _ = std::fs::write(&marker, "failed once");
write_atomically(&marker, "failed once");
std::process::exit(3);
}
let changes_view = method == "tools/call"
Expand Down Expand Up @@ -459,3 +459,13 @@ async fn main() {
});
}
}

/// Write a file a test reads, all at once: to a temporary sibling, then `rename` it into place. A
/// reader polling for the file (or its content) can otherwise see it created but still empty,
/// between `write`'s create and its write.
fn write_atomically(path: &str, contents: &str) {
let tmp = format!("{path}.tmp-{}", std::process::id());
if std::fs::write(&tmp, contents).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
}
17 changes: 15 additions & 2 deletions crates/agent/src/bin/mcp_fixture_events_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1758,7 +1758,10 @@ async fn main() {
if let Ok(pidfile) = std::env::var("MCP_FIXTURE_ORPHAN_PIDFILE") {
let _ = tokio::process::Command::new("sh")
.arg("-c")
.arg(format!("sleep 600 & echo $! > {pidfile}"))
// Written to a temporary name and renamed, so a test never reads a half-written pid.
.arg(format!(
"sleep 600 & echo $! > {pidfile}.tmp && mv -f {pidfile}.tmp {pidfile}"
))
.status()
.await;
}
Expand All @@ -1767,7 +1770,7 @@ async fn main() {
// to "clean up" (as a server closing a browser would), then record that we got to.
if let Ok(path) = std::env::var("MCP_FIXTURE_EXIT_MARKER") {
tokio::time::sleep(Duration::from_millis(500)).await;
let _ = std::fs::write(path, "clean exit");
write_atomically(&path, "clean exit");
}
} else {
println!(
Expand All @@ -1777,3 +1780,13 @@ async fn main() {
let _ = accept.await;
}
}

/// Write a file a test reads, all at once: to a temporary sibling, then `rename` it into place. A
/// reader polling for the file (or its content) can otherwise see it created but still empty,
/// between `write`'s create and its write.
fn write_atomically(path: &str, contents: &str) {
let tmp = format!("{path}.tmp-{}", std::process::id());
if std::fs::write(&tmp, contents).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
}
16 changes: 14 additions & 2 deletions crates/agent/src/bin/mcp_fixture_stdio_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,11 +69,21 @@ fn tasks() -> &'static Mutex<HashMap<String, FixtureTask>> {
TASKS.get_or_init(|| Mutex::new(HashMap::new()))
}

/// Write a file a test reads, all at once: to a temporary sibling, then `rename` it into place. A
/// reader polling for the file (or its content) can otherwise see it created but still empty,
/// between `write`'s create and its write.
fn write_atomically(path: &str, contents: &str) {
let tmp = format!("{path}.tmp-{}", std::process::id());
if std::fs::write(&tmp, contents).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
}

fn record_cancel_flag(task_id: &str) {
let Ok(path) = std::env::var("MCP_FIXTURE_CANCEL_FLAG") else {
return;
};
let _ = std::fs::write(path, task_id);
write_atomically(&path, task_id);
}

fn capabilities() -> Value {
Expand All @@ -97,7 +107,9 @@ async fn main() {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
}
if let Ok(pidfile) = std::env::var("MCP_FIXTURE_ORPHAN_PIDFILE") {
let script = format!("sleep 600 & echo $! > {pidfile}");
// Written to a temporary name and renamed, so a test never reads a half-written pid.
let script =
format!("sleep 600 & echo $! > {pidfile}.tmp && mv -f {pidfile}.tmp {pidfile}");
let _ = tokio::process::Command::new("sh")
.arg("-c")
.arg(script)
Expand Down
54 changes: 50 additions & 4 deletions crates/agent/src/bin/mcp_fixture_tasks_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@
//! - `ttl_task`: `ttlMs: 400`, never leaves `working`.
//! - `blip_task`: over HTTP, polls 2 and 3 have their connection dropped without a response (a
//! network blip); completes on poll 5 with `blip-done`.
//! - `gated_task`: `working` until the file at `MCP_TASKS_FIXTURE_GATE` exists, then `gated-done`.
//! - `gated_task`: `working` until the file at `MCP_TASKS_FIXTURE_GATE` exists, then `gated-done`
//! (or `MCP_TASKS_FIXTURE_GATED_TEXT`).
//! - `sample_task`: one in-task `sampling/createMessage` (`draft`); completes with `sampled:<text>`.
//! - `ttl_shift_task`: created with `ttlMs: null`; every poll then says `ttlMs: 300`, never finishing.
//! Only a client honouring the *latest* TTL stops (after 60 polls it completes `ttl-ignored`).
Expand All @@ -51,7 +52,7 @@
//! Resource `fixture-tasks://doc`: `resources/read` answers a `CreateTaskResult`, which the client
//! MUST treat as an invalid response (tasks are defined for `tools/call` only).
//!
//! Env: `MCP_TASKS_FIXTURE_LOG` (path), `MCP_TASKS_FIXTURE_LEGACY=1` (answer `server/discover` with
//! Env: `MCP_TASKS_FIXTURE_KEEP` (path: gated tasks survive a restart), `MCP_TASKS_FIXTURE_LOG` (path), `MCP_TASKS_FIXTURE_LEGACY=1` (answer `server/discover` with
//! `-32601`, forcing the client onto legacy `initialize`), `MCP_TASKS_FIXTURE_GATE` (path),
//! `MCP_TASKS_FIXTURE_HTTP_PORT_FILE` (serve Streamable HTTP on `127.0.0.1:<ephemeral>/mcp` instead
//! of stdio, writing the port to this file; each HTTP log line also records the `Mcp-Method` /
Expand Down Expand Up @@ -264,6 +265,9 @@ impl Server {
if kind == Kind::CrashOnce {
save_state(&task_id);
}
if kind == Kind::Gated {
keep_gated(&self.tasks);
}
// A long interval for `ttl_task`, so only a client that caps its wait by the TTL notices
// the TTL on time.
let interval = if kind == Kind::Ttl { 10_000 } else { 60 };
Expand Down Expand Up @@ -324,7 +328,8 @@ impl Server {
Kind::Blip => completed(task_id, ttl, "blip-done"),
Kind::Gated => {
if gate_open() {
completed(task_id, ttl, "gated-done")
let text = std::env::var("MCP_TASKS_FIXTURE_GATED_TEXT");
completed(task_id, ttl, text.as_deref().unwrap_or("gated-done"))
} else {
task_json(task_id, "working", ttl, 50)
}
Expand Down Expand Up @@ -477,14 +482,54 @@ fn ask(task_id: &str, ttl: Value, key: &str, message: &str, interval: u64) -> Va
v
}

/// Write a file a test reads, all at once: to a temporary sibling, then `rename` it into place. A
/// reader polling for the file (or its content) can otherwise see it created but still empty,
/// between `write`'s create and its write.
fn write_atomically(path: &str, contents: &str) {
let tmp = format!("{path}.tmp-{}", std::process::id());
if std::fs::write(&tmp, contents).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
}

/// `crash_once_task`: remember the task across a process restart.
fn save_state(task_id: &str) {
if let Ok(path) = std::env::var("MCP_TASKS_FIXTURE_STATE") {
let _ = std::fs::write(path, task_id);
write_atomically(&path, task_id);
}
}

/// The `crash_once_task` a previous process saved, already past its crash.
/// `MCP_TASKS_FIXTURE_KEEP=<path>`: the gated tasks this server holds, one id per line, so a
/// restarted HTTP server still knows them (a server whose tasks outlive its process).
fn keep_gated(tasks: &HashMap<String, Task>) {
if let Ok(path) = std::env::var("MCP_TASKS_FIXTURE_KEEP") {
let ids: Vec<&str> = tasks
.iter()
.filter(|(_, t)| t.kind == Kind::Gated)
.map(|(id, _)| id.as_str())
.collect();
write_atomically(&path, &ids.join("\n"));
}
}

fn load_kept(tasks: &mut HashMap<String, Task>) {
let Ok(path) = std::env::var("MCP_TASKS_FIXTURE_KEEP") else {
return;
};
for id in std::fs::read_to_string(path).unwrap_or_default().lines() {
tasks.insert(
id.to_owned(),
Task {
kind: Kind::Gated,
polls: 0,
answers: HashMap::new(),
updates: 0,
},
);
}
}

fn load_state(tasks: &mut HashMap<String, Task>) -> bool {
let Ok(path) = std::env::var("MCP_TASKS_FIXTURE_STATE") else {
return false;
Expand Down Expand Up @@ -630,6 +675,7 @@ fn envelope(server: &mut Server, id: Value, method: &str, params: &Value) -> Opt
#[tokio::main(flavor = "current_thread")]
async fn main() {
let mut tasks = HashMap::new();
load_kept(&mut tasks);
if load_state(&mut tasks)
&& let Some(ms) = std::env::var("MCP_TASKS_FIXTURE_RESTART_DELAY_MS")
.ok()
Expand Down
12 changes: 11 additions & 1 deletion crates/agent/src/bin/mcp_skills_fixture_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -469,7 +469,7 @@ fn handle(method: &str, params: &Value) -> Result<Value, (i64, String)> {
}] })),
"tools/call" if params.get("name").and_then(Value::as_str) == Some("publish_late") => {
if let Ok(path) = std::env::var("MCP_SKILLS_FIXTURE_LATE_FLAG") {
let _ = std::fs::write(path, "published");
write_atomically(&path, "published");
}
Ok(
json!({ "content": [{ "type": "text", "text": "LATE-PUBLISHED" }], "isError": false }),
Expand Down Expand Up @@ -821,3 +821,13 @@ async fn serve_http() {
() = eof => {}
}
}

/// Write a file a test reads, all at once: to a temporary sibling, then `rename` it into place. A
/// reader polling for the file (or its content) can otherwise see it created but still empty,
/// between `write`'s create and its write.
fn write_atomically(path: &str, contents: &str) {
let tmp = format!("{path}.tmp-{}", std::process::id());
if std::fs::write(&tmp, contents).is_ok() {
let _ = std::fs::rename(&tmp, path);
}
}
12 changes: 10 additions & 2 deletions crates/agent/src/file_lock.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,18 @@ const RETRIES: usize = 5;
/// A held lock. Dropping it closes the descriptor (releasing the kernel lock) and frees the
/// registration.
pub struct FileLock {
_file: File,
file: File,
_registration: Registration,
}

impl FileLock {
/// The locked file, for a holder that writes the lock file's own content. Write through this
/// descriptor, never a second open: over NFS, closing *any* descriptor to the file drops the lock.
pub fn file(&self) -> &File {
&self.file
}
}

/// "This process holds the lock at this path."
struct Registration(PathBuf);

Expand Down Expand Up @@ -80,7 +88,7 @@ pub fn try_lock(lock_path: &Path) -> std::io::Result<Option<FileLock>> {
}
if same_file(&file, lock_path)? {
return Ok(Some(FileLock {
_file: file,
file,
_registration: registration,
}));
}
Expand Down
Loading
Loading