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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

109 changes: 94 additions & 15 deletions crates/agent/ARCHITECTURE.md

Large diffs are not rendered by default.

5 changes: 5 additions & 0 deletions crates/agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,11 @@ htmd = "0.5"
# Already resolved transitively at this exact version (tokio → signal-hook-registry → errno), so
# depending on it directly adds nothing to the build graph.
libc = "0.2"
# The session and manifest locks are open file description locks (`fcntl(F_OFD_SETLK)`, see
# `file_lock`): owned by the one `open` that took them, so no other descriptor of the file can release
# them, and released explicitly so a forked child's copy cannot keep them. Already in the build graph
# at this version (transitively); its safe `fcntl` keeps the `unsafe` out of this crate, which forbids it.
nix = { version = "0.31", default-features = false, features = ["fs"] }
# `export_html`'s markdown rendering (headers/bold/lists/code fences) for message text — pure Rust, no
# vendored JS. Deliberately not paired with a syntax-highlighting crate (e.g. `syntect`): that pulls in
# several MB of bundled syntax/theme data and would slow every build of this CLI, including `run`/
Expand Down
600 changes: 479 additions & 121 deletions crates/agent/src/file_lock.rs

Large diffs are not rendered by default.

24 changes: 12 additions & 12 deletions crates/agent/src/serve_ws.rs
Original file line number Diff line number Diff line change
Expand Up @@ -590,9 +590,9 @@ where
/// a replica that died mid-start) is left for an out-of-band sweep; see ARCHITECTURE.md.
///
/// `remove_dir` *is* the emptiness check, and an atomic one: `000001.jsonl` is created with `O_EXCL`
/// before anything else and never deleted, so a directory holding nothing but `lock` has never been
/// before anything else and never deleted, so a directory holding nothing but its lock files has never been
/// a session — and if a racing replica got one written in between, `ENOTEMPTY` leaves everything
/// alone. The `lock` file is unlinked first because `remove_dir` would otherwise always fail; the
/// alone. The lock files are unlinked first because `remove_dir` would otherwise always fail; the
/// window that opens between that unlink and this task's own release is harmless, since the lock is
/// liveness-only (correctness is the epoch fence) and this replica is already done with the session.
async fn release_session_lock(lock: Option<SessionLock>, path: Option<std::path::PathBuf>) {
Expand All @@ -601,16 +601,15 @@ async fn release_session_lock(lock: Option<SessionLock>, path: Option<std::path:
};
// Off the runtime thread: three network-filesystem calls and a descriptor close.
let _ = tokio::task::spawn_blocking(move || {
let only_lock = std::fs::read_dir(&path).is_ok_and(|entries| {
entries
.flatten()
.all(|e| e.file_name() == std::ffi::OsStr::new("lock"))
});
if only_lock {
let _ = std::fs::remove_file(path.join("lock"));
// The lock files are `file_lock`'s to name and unlink. Released (descriptors closed) before
// they are unlinked: an NFS client silly-renames a file it still holds open to `.nfs*`, and
// that leftover would keep the directory from ever being removed.
if crate::file_lock::dir_holds_only_lock_files(&path) {
lock.release_and_remove_files();
let _ = std::fs::remove_dir(&path);
} else {
drop(lock);
}
drop(lock);
})
.await;
}
Expand Down Expand Up @@ -3651,8 +3650,9 @@ mod tests {
"the segment must survive"
);
assert!(
path.join("lock").is_file(),
"so must the lock file it is locked through"
!crate::file_lock::dir_holds_only_lock_files(&path)
&& std::fs::read_dir(&path).unwrap().count() > 1,
"so must the lock files it is locked through"
);
// The lock itself is released, so the next owner can take it.
assert!(acquire_session_lock(&path).unwrap().is_some());
Expand Down
197 changes: 180 additions & 17 deletions crates/agent/src/session_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -862,7 +862,7 @@ fn make_journal_key(
}
let deadline = std::time::Instant::now() + JOURNAL_KEY_LOCK_WAIT;
let lock = loop {
match crate::file_lock::try_lock(&path) {
match crate::file_lock::try_lock(crate::file_lock::Target::Itself(&path)) {
Ok(Some(lock)) => break lock,
Ok(None) if std::time::Instant::now() < deadline => {
std::thread::sleep(std::time::Duration::from_millis(2));
Expand Down Expand Up @@ -6402,31 +6402,194 @@ pub type SessionLock = crate::file_lock::FileLock;
/// more. Correctness comes from the epoch fence, which needs no lock at all — so a lock lost to a
/// crash, a network partition or a stuck NFS client costs a retry, never history.
///
/// Three details that matter on EFS:
/// - the descriptor is opened **read+write**, because NFS emulates `flock` with POSIX record locks and
/// those need a writable descriptor;
/// - after locking, `fstat` on the held descriptor is compared with `stat` of the path, so a directory
/// renamed into `.trash/` between the open and the lock is caught rather than silently "locked";
/// - a session is opened **once per process**, because closing *any* descriptor to a POSIX-locked file
/// drops the lock — a second open in the same process would quietly release the first one's hold.
///
/// All three are [`crate::file_lock::try_lock`]'s; this only names the session's lock file.
/// It is an open file description lock ([`crate::file_lock`]): no other descriptor of the lock file
/// releases it, a second acquire in this process is refused, and it is released with an explicit
/// unlock — a forked child inherits the description, but cannot keep the lock once its holder lets
/// go. After locking, `fstat` on the held descriptor is compared with `stat` of the path, so a
/// directory renamed into `.trash/` between the open and the lock is caught rather than silently
/// "locked". This only says which lock file is the session's.
pub fn acquire_session_lock(session_path: &Path) -> std::io::Result<Option<SessionLock>> {
let lock_path = if session_path.is_dir() {
session_path.join("lock")
crate::file_lock::try_lock(if session_path.is_dir() {
crate::file_lock::Target::Dir(session_path)
} else {
let mut name = session_path.as_os_str().to_os_string();
name.push(".lock");
PathBuf::from(name)
};
crate::file_lock::try_lock(&lock_path)
crate::file_lock::Target::File(session_path)
})
}

#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;

/// The root cause of `serve_ws::tests::a_session_directory_with_a_segment_is_left_alone`'s
/// flakes under load: a child forked by *another* thread (a tool, an MCP server) shares the
/// lock's open file description until it execs, so a lock released by closing its descriptor was
/// still held about 6% of the time while processes were being spawned (180/3000 measured). The
/// OFD lock is inherited the same way, but it is released with an explicit unlock, which drops it
/// whatever copies a child holds: zero.
#[test]
fn a_released_session_lock_is_free_at_once_while_other_threads_spawn_processes() {
let stop = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let spawners: Vec<_> = (0..8)
.map(|_| {
let stop = stop.clone();
std::thread::spawn(move || {
while !stop.load(std::sync::atomic::Ordering::Relaxed) {
let _ = std::process::Command::new("true").status();
}
})
})
.collect();
let root = tempfile::tempdir().unwrap();
let mut held = 0;
for i in 0..3000 {
let p = root.path().join(format!("s{i}"));
std::fs::create_dir_all(&p).unwrap();
let lock = acquire_session_lock(&p).unwrap().unwrap();
drop(lock);
if acquire_session_lock(&p).unwrap().is_none() {
held += 1;
}
}
stop.store(true, std::sync::atomic::Ordering::Relaxed);
for s in spawners {
s.join().unwrap();
}
assert_eq!(
held, 0,
"a released lock was still held by a forked child {held}/3000 times"
);
}

/// The lock reached through another name — a symlink (`innocent.txt -> <session>/<lock file>`,
/// which the model can make with `bash`) or a hard link — by every in-process reader that does
/// not know it is a lock file: the file tools' read, write, edit and search, the memory tool's
/// view (under `<session>/memory`), and the `@file` read (`std::fs::read_to_string`). Each opens and
/// closes the lock file on its own descriptor, and none of them releases the lock, which belongs
/// to the open file description that took it.
#[cfg(unix)]
#[tokio::test]
async fn every_in_process_reader_reaching_the_lock_file_leaves_it_held() {
use agent_core::Tool as _;
let dir = tmpdir();
let session_dir = dir.path().join("s1");
std::fs::create_dir_all(session_dir.join("memory")).unwrap();
let lock = acquire_session_lock(&session_dir).unwrap().unwrap();
let target = crate::file_lock::Target::Dir(&session_dir);
let Some(true) = crate::file_lock::tests::held_for_another_process(target) else {
eprintln!("no python3: skipping");
return;
};
let record = crate::file_lock::tests::record_path(target);
let work = dir.path().join("work");
std::fs::create_dir_all(&work).unwrap();
let symlink = work.join("innocent.txt");
std::os::unix::fs::symlink(&record, &symlink).unwrap();
let hardlink = work.join("also-innocent.txt");
std::fs::hard_link(&record, &hardlink).unwrap();
std::os::unix::fs::symlink(&record, session_dir.join("memory").join("notes.md")).unwrap();

for link in [&symlink, &hardlink] {
let p = link.to_str().unwrap();
let _ = crate::tools::read::Read::new(&work)
.run(serde_json::json!({ "path": p }))
.await;
let _ = crate::tools::write::Write::new(&work)
.run(serde_json::json!({ "path": p, "content": "scribble" }))
.await;
let _ = crate::tools::edit::Edit::new(&work)
.run(serde_json::json!({ "path": p, "old_string": "", "new_string": "x" }))
.await;
// What `@file` does with the path it is given.
let _ = std::fs::read_to_string(link);
}
let _ = crate::tools::grep::Grep::new(&work)
.run(serde_json::json!({ "pattern": "x", "path": work.to_str().unwrap() }))
.await;
let memory = crate::tools::memory::Memory::new(std::sync::Arc::new(
crate::memory::file::FileBackend::session_at(session_dir.join("memory")),
));
let _ = memory
.run(serde_json::json!({ "command": "view", "path": "/session/notes.md" }))
.await;

assert_eq!(
crate::file_lock::tests::held_for_another_process(target),
Some(true),
"no reader of the lock file released it"
);
assert!(
acquire_session_lock(&session_dir).unwrap().is_none(),
"nor can this process take it a second time"
);
assert_eq!(
std::fs::metadata(&record).unwrap().len(),
0,
"and nothing replaced or wrote into the lock file"
);
drop(lock);
assert_eq!(
crate::file_lock::tests::held_for_another_process(target),
Some(false)
);
}

/// Every public path that reads a locked session directory — listing, opening, forking, the
/// file tools' read and search of every file in it, the lock file included — leaves the lock
/// held for another process: none of their descriptors is the one that holds it.
#[tokio::test]
async fn reading_a_locked_session_through_any_public_api_leaves_it_locked() {
use crate::tools::fs::FsBackend as _;
let dir = tmpdir();
let repo = SessionRepo::open_with(
dir.path(),
RepoOptions {
layout: Layout::Segmented { codec: None },
..RepoOptions::default()
},
)
.unwrap();
let mut store = repo.create(SessionMeta::new("/w", "m")).unwrap();
let mut session = Session::new();
session.user("hello lock");
store.append_new(&session.messages).unwrap();
let id = store.meta().id.clone();
let session_dir = store.path().to_path_buf();
drop(store);
let lock = acquire_session_lock(&session_dir).unwrap().unwrap();
let target = crate::file_lock::Target::Dir(&session_dir);
let Some(true) = crate::file_lock::tests::held_for_another_process(target) else {
eprintln!("no python3: skipping");
return;
};

let _ = repo.list().unwrap();
let _ = repo.open_id_read_only(&id).unwrap();
let _ = repo.fork(&id, usize::MAX).unwrap();
let fs = crate::tools::fs::local::LocalFs::new();
for entry in std::fs::read_dir(&session_dir).unwrap().flatten() {
// Read every file there by the tool path; a lock file is refused, never opened.
let _ = fs.read_bytes(&entry.path(), 0, 1024).await;
}
let _ = fs
.search(&crate::tools::fs::local::query(
"lock",
dir.path().to_path_buf(),
100,
))
.await;
assert_eq!(
crate::file_lock::tests::held_for_another_process(target),
Some(true),
"no public read of the session released its lock"
);
drop(lock);
assert_eq!(
crate::file_lock::tests::held_for_another_process(target),
Some(false)
);
}

fn tmpdir() -> tempfile::TempDir {
tempfile::tempdir().unwrap()
}
Expand Down
18 changes: 17 additions & 1 deletion crates/agent/src/tools/mcp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,9 @@ struct HttpDial {
/// The daemon's `mcp_http` client: ALPN, the `web` tool's SSRF resolver, and deliberately not
/// the gateway's h2c pool — a tenant's connector is not the gateway.
client: reqwest::Client,
/// The egress policy that client enforces, for the fresh client that ends this connection's
/// session on the way out of the process (`mcp_http_exit`).
policy: Arc<crate::tools::web::ssrf::EgressPolicy>,
headers: Vec<SecretHeader>,
}

Expand Down Expand Up @@ -453,6 +456,9 @@ impl ClientHandler for McpHandler {
}
self.skills_changed
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
// rmcp runs this on a task of its own, after (or before) the response that followed the
// notification on the wire; this line is how an operator — or a test — knows it has run.
tracing::debug!(server = %self.server_name, "resources/list_changed handled: view cache cleared");
}

async fn on_progress(
Expand Down Expand Up @@ -2112,6 +2118,7 @@ fn granted_jobs(
host: host.clone(),
http: Some(HttpDial {
client: egress.client.clone(),
policy: egress.policy.clone(),
headers,
}),
apps,
Expand Down Expand Up @@ -2511,13 +2518,21 @@ async fn connect_http(
// Wrapped so extension results survive rmcp's result decoding — see `mcp_wire` — so an MCP
// App view's read is refused over its cap without being read whole — see `mcp_view_http` — and
// so a rejected OAuth token is refreshed and the request retried once — see `mcp_oauth`.
// The session this connection opens is ended on the way out of the process if nothing ended
// it before — see `mcp_http_exit`.
let exit = crate::tools::mcp_http_exit::HttpSession::new(
url,
auth.clone(),
dial.http.as_ref().map(|h| h.policy.clone()),
);
let transport = StreamableHttpClientTransport::with_client(
crate::tools::mcp_oauth::OAuthHttp::new(
crate::tools::mcp_view_http::ViewCappedHttp::new(crate::tools::mcp_wire::HttpClient {
client,
}),
auth,
),
)
.ending_session_on_exit(exit),
transport_config,
);

Expand Down Expand Up @@ -3346,6 +3361,7 @@ mod tests {
agent_core::ensure_provider();
let dial = HttpDial {
client: reqwest::Client::new(),
policy: Arc::default(),
headers: vec![header],
};
assert!(!format!("{dial:?}").contains("tenant-token"));
Expand Down
Loading
Loading