diff --git a/experimental/silverwood/crates/silverwood-core/src/docstore.rs b/experimental/silverwood/crates/silverwood-core/src/docstore.rs index 65b8bec..f170e31 100644 --- a/experimental/silverwood/crates/silverwood-core/src/docstore.rs +++ b/experimental/silverwood/crates/silverwood-core/src/docstore.rs @@ -1,6 +1,8 @@ use std::fs; +use std::fs::{File, OpenOptions}; use std::io; use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; use crate::error::{Error, Result}; use crate::id::{IdScheme, WorkstreamId}; @@ -8,6 +10,9 @@ use crate::id::{IdScheme, WorkstreamId}; /// File extension for a persisted workstream document. const DOC_EXT: &str = "loro"; +/// Per-process counter making each save's temp file name unique (see [`FilesDocStore::save`]). +static SAVE_SEQ: AtomicU64 = AtomicU64::new(0); + /// Where workstream documents are persisted. Each document is keyed by its /// [`WorkstreamId`]; the store is oblivious to document contents (opaque bytes). /// @@ -22,8 +27,27 @@ pub trait DocStore { /// Enumerate the ids of all documents currently present. fn list_ids(&self) -> Result>; + + /// Acquire an exclusive, cross-process write lock for document `id`, held until the + /// returned guard is dropped. The forest holds it across an entire load-modify-save so + /// two processes writing the same document — e.g. papyrus and Codex's separate + /// SessionStart-hook process (`__register-codex-session`) — take turns instead of racing + /// and losing each other's write. This is the one guarantee a snapshot-overwrite store + /// needs to be correct under multiple writers. A backend with no cross-process sharing + /// may keep the default no-op guard. + fn write_lock(&self, _id: WorkstreamId) -> Result> { + Ok(Box::new(())) + } } +/// An RAII handle for a [`DocStore::write_lock`]: the lock is held for as long as the guard +/// lives and released when it drops. The unit type is a valid no-op guard for backends that +/// need no locking; [`FilesDocStore`] returns the locked [`File`] itself (dropping it closes +/// the fd and releases the `flock`). +pub trait DocWriteGuard {} +impl DocWriteGuard for () {} +impl DocWriteGuard for File {} + /// A [`DocStore`] backed by one file per document under a directory. pub struct FilesDocStore { dir: PathBuf, @@ -59,6 +83,13 @@ impl FilesDocStore { } self.bare_path(id).filter(|p| p.exists()) } + + /// The sibling advisory-lock file for `id` (see [`DocStore::write_lock`]). Always keyed + /// by the scheme-explicit stem, so the deprecated bare and canonical doc names still map + /// to one lock; the `.lock` extension is never `.loro`, so `list_ids` ignores it. + fn lock_path(&self, id: WorkstreamId) -> PathBuf { + self.dir.join(format!("{}.lock", id.storage_key())) + } } impl DocStore for FilesDocStore { @@ -80,12 +111,18 @@ impl DocStore for FilesDocStore { // previous document intact rather than a truncated one. Reuse the existing // on-disk name if the document already exists — a pre-scheme forest's bare // name is updated in place, never renamed — and mint the canonical - // scheme-explicit name only for a brand-new document. The `.tmp` extension - // is not `.loro`, so `list_ids` ignores any leftover. + // scheme-explicit name only for a brand-new document. let path = self .existing_path(id) .unwrap_or_else(|| self.explicit_path(id)); - let tmp = path.with_extension("tmp"); + // The temp name is unique per (process, save): a fixed `.tmp` would be shared by + // two writers racing the same document from different processes — e.g. papyrus and + // Codex's separate SessionStart-hook process (`__register-codex-session`) — so one's + // `rename` could find the temp already moved and fail with `ENOENT`. pid keeps it + // distinct across processes, the counter across concurrent saves within one. Any + // extension here ends in `.tmp` (never `.loro`), so `list_ids` ignores leftovers. + let seq = SAVE_SEQ.fetch_add(1, Ordering::Relaxed); + let tmp = path.with_extension(format!("{}.{seq}.tmp", std::process::id())); fs::write(&tmp, bytes).map_err(|e| Error::io(&tmp, e))?; fs::rename(&tmp, &path).map_err(|e| Error::io(&path, e)) } @@ -118,6 +155,21 @@ impl DocStore for FilesDocStore { } Ok(ids) } + + fn write_lock(&self, id: WorkstreamId) -> Result> { + // Advisory `flock` on a per-document lock file. Blocking: a concurrent writer to the + // same `id` (in this or another process) waits here until the holder's guard drops. + // `flock` is per-open-file-description, so this must never be re-acquired for an `id` + // already locked on this call stack — the forest locks once, at the top of a mutation. + let path = self.lock_path(id); + let file = OpenOptions::new() + .create(true) + .write(true) + .open(&path) + .map_err(|e| Error::io(&path, e))?; + file.lock().map_err(|e| Error::io(&path, e))?; + Ok(Box::new(file)) + } } /// Whether `path` names a document file (`*.loro`). @@ -182,6 +234,63 @@ mod tests { assert_eq!(store.list_ids().unwrap(), vec![id]); } + #[test] + fn concurrent_saves_of_one_doc_never_collide_on_temp() { + use std::sync::Arc; + use std::thread; + + let dir = TempDir::new("docstore-concurrent"); + let store = Arc::new(FilesDocStore::new(&dir.0)); + let id = WorkstreamId::generate(); + store.save(id, b"seed").unwrap(); + + // Many writers hammering one document must never trip the temp-file race: a fixed + // `.tmp` shared across savers lets one's `rename` lose the temp to another's and + // fail with ENOENT. Each save must return Ok and the doc stays loadable throughout. + let handles: Vec<_> = (0..8) + .map(|t| { + let store = Arc::clone(&store); + thread::spawn(move || { + for i in 0..50 { + store.save(id, format!("t{t}-{i}").as_bytes()).unwrap(); + assert!(store.load(id).unwrap().is_some()); + } + }) + }) + .collect(); + for h in handles { + h.join().unwrap(); + } + } + + #[test] + fn write_lock_is_exclusive_per_document() { + let dir = TempDir::new("docstore-lock"); + let store = FilesDocStore::new(&dir.0); + let id = WorkstreamId::generate(); + + let held = store.write_lock(id).unwrap(); + // A second acquisition of the same doc's lock (a fresh descriptor) must not be + // grantable while the first is held — `flock` contends across descriptors even within + // one process, which is what serializes the codex hook against papyrus. + let probe = std::fs::OpenOptions::new() + .create(true) + .write(true) + .open(store.lock_path(id)) + .unwrap(); + assert!( + probe.try_lock().is_err(), + "a second holder must be blocked while the lock is held" + ); + + // Once the guard drops, the lock is grantable again. + drop(held); + assert!( + probe.try_lock().is_ok(), + "the lock must be grantable after release" + ); + } + #[test] fn new_document_is_written_with_explicit_scheme() { let dir = TempDir::new("docstore-new-explicit"); diff --git a/experimental/silverwood/crates/silverwood-core/src/forest.rs b/experimental/silverwood/crates/silverwood-core/src/forest.rs index 84ba81d..62994c2 100644 --- a/experimental/silverwood/crates/silverwood-core/src/forest.rs +++ b/experimental/silverwood/crates/silverwood-core/src/forest.rs @@ -245,6 +245,7 @@ impl Forest { /// [`Error::NotAwaitingCheckout`] (leaving the document untouched) if it is already /// checked out, mid-provision, or previously failed. pub fn checkout_workstream(&self, id: WorkstreamId) -> Result { + let _guard = self.docs.write_lock(id)?; let bytes = self.docs.load(id)?.ok_or(Error::NotFound(id))?; let ws = doc::hydrate(id, &bytes)?; @@ -322,6 +323,7 @@ impl Forest { /// Archive a workstream (tombstone; the document and checkout are retained). pub fn archive(&self, id: WorkstreamId) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::set_status(&doc, Status::Archived)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -352,6 +354,9 @@ impl Forest { } // Tombstone first (the sync-relevant record of truth), then discard the directory. + // The lock is taken here, after the (VCS-touching) removability check, so a slow `jj` + // probe doesn't hold it — only the doc read-modify-save is serialized. + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::set_status(&doc, Status::Deleted)?; self.docs.save(id, &doc::snapshot(&doc)?)?; @@ -366,6 +371,7 @@ impl Forest { /// Rename a workstream (overwrite its `name`). pub fn rename(&self, id: WorkstreamId, name: &str) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::set_name(&doc, name)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -376,6 +382,7 @@ impl Forest { /// Namespaces under the core-reserved prefix (e.g. sessions) are rejected. pub fn set_kv(&self, id: WorkstreamId, namespace: &str, key: &str, value: &str) -> Result<()> { reject_reserved(namespace)?; + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::set_kv(&doc, namespace, key, value)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -385,6 +392,7 @@ impl Forest { /// are rejected (they are core-owned; use the session API). pub fn unset_kv(&self, id: WorkstreamId, namespace: &str, key: &str) -> Result<()> { reject_reserved(namespace)?; + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::unset_kv(&doc, namespace, key)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -421,6 +429,7 @@ impl Forest { agent_kind: SessionKind, name: &str, ) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::create_session(&doc, session_id, agent_kind, name, &now_rfc3339())?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -428,6 +437,7 @@ impl Forest { /// Rename a session (preserving its kind + created_at). Errors if absent. pub fn rename_session(&self, id: WorkstreamId, session_id: &str, name: &str) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::rename_session(&doc, session_id, name)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -435,6 +445,7 @@ impl Forest { /// Remove a session from a workstream (no-op if absent). pub fn remove_session(&self, id: WorkstreamId, session_id: &str) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; doc::remove_session(&doc, session_id)?; self.docs.save(id, &doc::snapshot(&doc)?) @@ -468,7 +479,7 @@ impl Forest { claude_config_dir: &Path, codex_home: &Path, ) -> Result { - let doc = self.load_doc(id)?; + let doc = self.load_doc_ephemeral(id)?; let session = doc::get_session(&doc, session_id)? .ok_or_else(|| Error::SessionNotFound(session_id.to_string()))?; let conversation_exists = @@ -514,7 +525,7 @@ impl Forest { let ws = self.get(id)?; let cwd = spawn_cwd(&ws)?; - let session = doc::get_session(&self.load_doc(id)?, session_id)? + let session = doc::get_session(&self.load_doc_ephemeral(id)?, session_id)? .ok_or_else(|| Error::SessionNotFound(session_id.to_string()))?; // First-run vs resume from ground truth: does Claude's transcript exist on disk? @@ -574,6 +585,7 @@ impl Forest { holder: &str, force: bool, ) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; let current = doc::get_session(&doc, session_id)? .ok_or_else(|| Error::SessionNotFound(session_id.to_string()))? @@ -608,6 +620,7 @@ impl Forest { holder: Option<&str>, force: bool, ) -> Result<()> { + let _guard = self.docs.write_lock(id)?; let doc = self.load_doc(id)?; let current = doc::get_session(&doc, session_id)? .ok_or_else(|| Error::SessionNotFound(session_id.to_string()))? @@ -643,6 +656,16 @@ impl Forest { doc::load(self.peer_id(), &bytes) } + /// Like [`Forest::load_doc`] but **never persists** the lazy schema migration, so a + /// read-only caller performs no write and needs no [`DocStore::write_lock`]. An + /// unpersisted upgrade is simply re-derived on the next read and written out by the next + /// mutation (or an explicit [`Forest::upgrade_all`]). + fn load_doc_ephemeral(&self, id: WorkstreamId) -> Result { + let bytes = self.docs.load(id)?.ok_or(Error::NotFound(id))?; + let bytes = doc::migrate_bytes(id, &bytes, self.peer_id())?.unwrap_or(bytes); + doc::load(self.peer_id(), &bytes) + } + /// Upgrade every document in the forest to the latest schema version, /// returning a per-document report sorted by id. With `dry_run`, inspects and /// reports without writing. Idempotent — documents already at the latest are @@ -652,6 +675,13 @@ impl Forest { let latest = migrate::DOC_SCHEMA_VERSION; let mut reports = Vec::new(); for id in self.docs.list_ids()? { + // Serialize each doc's read-migrate-save against other writers. A dry run only + // reads (it never saves), so it takes no lock. + let _guard = if dry_run { + None + } else { + Some(self.docs.write_lock(id)?) + }; let bytes = self.docs.load(id)?.ok_or(Error::NotFound(id))?; let from = doc::peek_version(id, &bytes)?; if from > latest {