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
18 changes: 18 additions & 0 deletions crates/ironclaw_filesystem/src/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,24 @@ impl RootFilesystem for CompositeRootFilesystem {
.await
}

async fn head_seq(
&self,
path: &VirtualPath,
from: SeqNo,
) -> Result<Option<SeqNo>, FilesystemError> {
self.matching_mount(path)?
.backend
.head_seq(path, from)
.await
}

async fn reserve_sequence(&self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
self.matching_mount(path)?
.backend
.reserve_sequence(path)
.await
}

// ── Legacy bytes plane ──

async fn read_file(&self, path: &VirtualPath) -> Result<Vec<u8>, FilesystemError> {
Expand Down
121 changes: 121 additions & 0 deletions crates/ironclaw_filesystem/src/in_memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ struct State {
entries: HashMap<VirtualPath, StoredEntry>,
indexes: HashMap<String, Vec<IndexSpec>>,
event_logs: HashMap<String, Vec<EventRecord>>,
sequences: HashMap<String, SeqNo>,
}

/// In-memory backend serving the full unified [`RootFilesystem`] surface.
Expand All @@ -64,6 +65,7 @@ impl InMemoryBackend {
entries: HashMap::new(),
indexes: HashMap::new(),
event_logs: HashMap::new(),
sequences: HashMap::new(),
}),
}
}
Expand Down Expand Up @@ -145,13 +147,17 @@ impl RootFilesystem for InMemoryBackend {
state
.entries
.retain(|key, _| !key.as_str().starts_with(&prefix));
clear_event_logs_under(&mut state.event_logs, path.as_str(), &prefix);
clear_sequences_under(&mut state.sequences, path.as_str(), &prefix);
return Ok(());
}
let prefix = with_trailing_slash(path.as_str());
let before = state.entries.len();
state
.entries
.retain(|key, _| !key.as_str().starts_with(&prefix));
clear_event_logs_under(&mut state.event_logs, path.as_str(), &prefix);
clear_sequences_under(&mut state.sequences, path.as_str(), &prefix);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if state.entries.len() == before {
return Err(FilesystemError::NotFound {
path: path.clone(),
Expand Down Expand Up @@ -432,6 +438,17 @@ impl RootFilesystem for InMemoryBackend {
.collect())
}

async fn reserve_sequence(&self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
let mut state = self.state.lock().await;
let next = state
.sequences
.entry(path.as_str().to_string())
.or_insert_with(|| SeqNo::ZERO.next());
let reserved = *next;
*next = next.next();
Ok(reserved)
}

// Legacy bytes ops — default impls in the trait route them through put/get
// and use our native implementations. The only one needing an explicit
// impl is the required-method `list_dir`, which we already overrode above.
Expand Down Expand Up @@ -588,6 +605,28 @@ fn with_trailing_slash(s: &str) -> String {
}
}

/// Drop any reserved sequence counter for `exact` and every path under
/// `prefix` (the trailing-slash form of `exact`). Mirrors the entries-delete
/// subtree semantics so a delete/recreate of the same path restarts its
/// sequence from 1 instead of resuming stale state. The exact match is kept
/// separate from the prefix scan so a sibling sharing a string prefix
/// (`/a/b` vs `/a/bc`) is not swept.
fn clear_sequences_under(sequences: &mut HashMap<String, SeqNo>, exact: &str, prefix: &str) {
sequences.retain(|key, _| key.as_str() != exact && !key.starts_with(prefix));
}

/// Drop the append-event log for `exact` and every log under `prefix`. Append-only
/// finalized assistant messages live in `event_logs`, so a delete/recreate of the
/// same thread path would otherwise rehydrate stale append-log history. Mirrors
/// the entries/sequences subtree semantics.
fn clear_event_logs_under(
event_logs: &mut HashMap<String, Vec<EventRecord>>,
exact: &str,
prefix: &str,
) {
event_logs.retain(|key, _| key.as_str() != exact && !key.starts_with(prefix));
}

fn first_segment(s: &str) -> (&str, bool) {
match s.find('/') {
Some(idx) => (&s[..idx], true),
Expand Down Expand Up @@ -667,6 +706,46 @@ mod tests {
assert!(fs.append_batch(&log, Vec::new()).await.unwrap().is_empty());
}

#[tokio::test]
async fn reserve_sequence_assigns_path_local_monotonic_values() {
let fs = InMemoryBackend::new();
let thread_a = vpath("/threads/a/message_sequence");
let thread_b = vpath("/threads/b/message_sequence");

assert_eq!(fs.reserve_sequence(&thread_a).await.unwrap().get(), 1);
assert_eq!(fs.reserve_sequence(&thread_a).await.unwrap().get(), 2);
assert_eq!(fs.reserve_sequence(&thread_b).await.unwrap().get(), 1);
assert_eq!(fs.reserve_sequence(&thread_a).await.unwrap().get(), 3);
}

#[tokio::test]
async fn deleting_a_subtree_clears_its_reserved_sequences() {
// Counters are now path-scoped in a side table rather than embedded in
// the thread record, so delete must sweep them or a delete/recreate of
// the same path would resume from stale state. A sibling that merely
// shares a string prefix must NOT be swept.
let fs = InMemoryBackend::new();
let seq_path = vpath("/threads/a/message_sequence");
let sibling = vpath("/threads/ab/message_sequence");
let thread_dir = vpath("/threads/a");

assert_eq!(fs.reserve_sequence(&seq_path).await.unwrap().get(), 1);
assert_eq!(fs.reserve_sequence(&seq_path).await.unwrap().get(), 2);
assert_eq!(fs.reserve_sequence(&sibling).await.unwrap().get(), 1);

// Materialize an entry so the directory delete has something to remove,
// then delete the whole /threads/a subtree.
fs.put(&seq_path, Entry::bytes(vec![1]), CasExpectation::Any)
.await
.unwrap();
fs.delete(&thread_dir).await.unwrap();

// The deleted subtree's counter restarts from 1.
assert_eq!(fs.reserve_sequence(&seq_path).await.unwrap().get(), 1);
// The string-prefix sibling /threads/ab is untouched.
assert_eq!(fs.reserve_sequence(&sibling).await.unwrap().get(), 2);
}

#[tokio::test]
async fn cas_absent_rejects_when_present() {
let fs = InMemoryBackend::new();
Expand Down Expand Up @@ -899,6 +978,48 @@ mod tests {
assert!(matches!(err, FilesystemError::NotFound { .. }));
}

#[tokio::test]
async fn delete_sweeps_append_event_logs_under_path() {
// Append-only finalized assistant messages live in the per-thread
// append log. Deleting a thread must clear that log so a later
// recreate of the same path does not replay stale history. This covers
// both delete branches: an exact append-log path, and a thread-root
// directory delete whose log lives under the subtree.
let fs = InMemoryBackend::new();

// Exact-entry branch: an entry and an append log share the deleted
// path. Deleting the entry must also sweep the co-located log.
let exact = vpath("/threads/exact");
fs.put(&exact, Entry::bytes(vec![1]), CasExpectation::Absent)
.await
.unwrap();
fs.append(&exact, b"finalized-a".to_vec()).await.unwrap();
assert_eq!(fs.tail(&exact, SeqNo::ZERO).await.unwrap().len(), 1);
fs.delete(&exact).await.unwrap();
assert!(fs.tail(&exact, SeqNo::ZERO).await.unwrap().is_empty());

// Directory branch: the thread root has no exact entry, but a child
// entry and an append log live under it. Deleting the root must sweep
// the log too.
let thread_root = vpath("/threads/t1");
let thread_doc = vpath("/threads/t1/thread.json");
let message_log = vpath("/threads/t1/messages_append.log");
fs.put(&thread_doc, Entry::bytes(vec![1]), CasExpectation::Absent)
.await
.unwrap();
fs.append(&message_log, b"finalized-b".to_vec())
.await
.unwrap();
assert_eq!(fs.tail(&message_log, SeqNo::ZERO).await.unwrap().len(), 1);

fs.delete(&thread_root).await.unwrap();

// History is gone, and recreating the same thread starts from an empty
// log rather than resurrecting the finalized message.
assert!(fs.get(&thread_doc).await.unwrap().is_none());
assert!(fs.tail(&message_log, SeqNo::ZERO).await.unwrap().is_empty());
}

#[tokio::test]
async fn list_dir_upgrades_to_directory_on_later_child_discovery() {
// PR #3659 reviewer fix: with `or_insert`, the first discovery
Expand Down
70 changes: 70 additions & 0 deletions crates/ironclaw_filesystem/src/libsql.rs
Original file line number Diff line number Diff line change
Expand Up @@ -960,6 +960,25 @@ impl RootFilesystem for LibSqlRootFilesystem {
if deleted == 0 {
return Err(not_found(path.clone(), FilesystemOperation::Delete));
}
// Sweep the append-event log for this path and its subtree. Append-only
// finalized assistant messages live in `root_filesystem_events`, so a
// delete/recreate of the same thread would otherwise replay stale
// history from the old log. Mirrors the entries-delete predicate above.
conn.execute(
"DELETE FROM root_filesystem_events WHERE path = ?1 OR path LIKE ?2 ESCAPE '!'",
libsql::params![path.as_str(), child_path_like_pattern(path)],
)
.await
.map_err(|error| libsql_db_error(path.clone(), FilesystemOperation::Delete, error))?;
// Sweep any reserved sequence counter for this path and its subtree so
// a delete/recreate restarts sequences from 1 rather than resuming
// stale state. Mirrors the entries-delete predicate above.
conn.execute(
"DELETE FROM root_filesystem_sequences WHERE path = ?1 OR path LIKE ?2 ESCAPE '!'",
libsql::params![path.as_str(), child_path_like_pattern(path)],
)
.await
.map_err(|error| libsql_db_error(path.clone(), FilesystemOperation::Delete, error))?;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Ok(())
}

Expand Down Expand Up @@ -1157,6 +1176,39 @@ impl RootFilesystem for LibSqlRootFilesystem {
}
}

async fn reserve_sequence(&self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
let conn = self.connect().await?;
let mut rows = conn
.query(
r#"
INSERT INTO root_filesystem_sequences (path, next_seq, updated_at)
VALUES (?1, 2, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
ON CONFLICT(path) DO UPDATE SET
next_seq = root_filesystem_sequences.next_seq + 1,
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
RETURNING next_seq - 1
"#,
libsql::params![path.as_str()],
)
.await
.map_err(|error| {
libsql_db_error(path.clone(), FilesystemOperation::ReserveSeq, error)
})?;
let row = rows
.next()
.await
.map_err(|error| libsql_db_error(path.clone(), FilesystemOperation::ReserveSeq, error))?
.ok_or_else(|| FilesystemError::Backend {
path: path.clone(),
operation: FilesystemOperation::ReserveSeq,
reason: "sequence reservation returned no row".to_string(),
})?;
let seq_raw: i64 = row.get(0).map_err(|error| {
libsql_db_error(path.clone(), FilesystemOperation::ReserveSeq, error)
})?;
seq_no_from_i64(path, seq_raw, FilesystemOperation::ReserveSeq)
}

async fn create_dir_all(&self, path: &VirtualPath) -> Result<(), FilesystemError> {
let conn = self.connect().await?;
let transaction = conn.transaction().await.map_err(|error| {
Expand Down Expand Up @@ -1218,6 +1270,7 @@ async fn run_libsql_migrations_inner(conn: &libsql::Connection) -> Result<(), Fi
ensure_libsql_records_columns(conn).await?;
ensure_libsql_index_specs_table(conn).await?;
ensure_libsql_events_table(conn).await?;
ensure_libsql_sequences_table(conn).await?;
Ok(())
}

Expand Down Expand Up @@ -1634,6 +1687,14 @@ async fn ensure_libsql_events_table(conn: &libsql::Connection) -> Result<(), Fil
Ok(())
}

#[cfg(feature = "libsql")]
async fn ensure_libsql_sequences_table(conn: &libsql::Connection) -> Result<(), FilesystemError> {
conn.execute_batch(LIBSQL_SEQUENCES_SCHEMA)
.await
.map_err(|error| infrastructure_libsql_error(FilesystemOperation::ReserveSeq, error))?;
Ok(())
}

#[cfg(feature = "libsql")]
fn seq_no_from_i64(
path: &VirtualPath,
Expand Down Expand Up @@ -1931,6 +1992,15 @@ CREATE INDEX IF NOT EXISTS idx_root_filesystem_events_path_seq
ON root_filesystem_events(path, seq);
"#;

#[cfg(feature = "libsql")]
const LIBSQL_SEQUENCES_SCHEMA: &str = r#"
CREATE TABLE IF NOT EXISTS root_filesystem_sequences (
path TEXT PRIMARY KEY,
next_seq INTEGER NOT NULL CHECK (next_seq > 0),
updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
);
"#;

#[cfg(test)]
mod tests {
//! Deterministic regression tests for libSQL behaviours that aren't
Expand Down
43 changes: 43 additions & 0 deletions crates/ironclaw_filesystem/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -748,6 +748,26 @@ impl RootFilesystem for PostgresRootFilesystem {
}
}

async fn reserve_sequence(&self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
let client = self.client().await?;
let row = cached_query_one(
&client,
r#"
INSERT INTO root_filesystem_sequences (path, next_seq, updated_at)
VALUES ($1, 2, NOW())
ON CONFLICT (path) DO UPDATE SET
next_seq = root_filesystem_sequences.next_seq + 1,
updated_at = NOW()
RETURNING next_seq - 1 AS reserved
"#,
&[&path.as_str()],
)
.await
.map_err(|error| db_error(path.clone(), FilesystemOperation::ReserveSeq, error))?;
let reserved: i64 = row.get("reserved");
seq_no_from_i64(path, reserved, FilesystemOperation::ReserveSeq)
}

async fn create_dir_all(&self, path: &VirtualPath) -> Result<(), FilesystemError> {
let mut client = self.client().await?;
let transaction = client
Expand Down Expand Up @@ -1395,6 +1415,27 @@ async fn postgres_delete_with_client(
if deleted == 0 {
return Err(not_found(path.clone(), FilesystemOperation::Delete));
}
// Sweep the append-event log for this path and its subtree. Append-only
// finalized assistant messages live in `root_filesystem_events`, so a
// delete/recreate of the same thread would otherwise replay stale history
// from the old log. Mirrors the entries-delete predicate above.
cached_execute(
client,
"DELETE FROM root_filesystem_events WHERE path = $1 OR (path >= $2 AND path < $3)",
&[&path.as_str(), &prefix_lower, &prefix_upper],
)
.await
.map_err(|error| db_error(path.clone(), FilesystemOperation::Delete, error))?;
// Sweep any reserved sequence counter for this path and its subtree so a
// delete/recreate restarts sequences from 1 rather than resuming stale
// state. Mirrors the entries-delete predicate above.
cached_execute(
client,
"DELETE FROM root_filesystem_sequences WHERE path = $1 OR (path >= $2 AND path < $3)",
&[&path.as_str(), &prefix_lower, &prefix_upper],
)
.await
.map_err(|error| db_error(path.clone(), FilesystemOperation::Delete, error))?;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Ok(())
}

Expand Down Expand Up @@ -1712,6 +1753,8 @@ const POSTGRES_ROOT_FILESYSTEM_SCHEMA: &str = concat!(
include_str!("../../../migrations/V30__root_filesystem_events.sql"),
"\n",
include_str!("../../../migrations/V31__root_filesystem_path_collation.sql"),
"\n",
include_str!("../../../migrations/V32__root_filesystem_sequences.sql"),
);

#[cfg(all(test, feature = "postgres"))]
Expand Down
11 changes: 11 additions & 0 deletions crates/ironclaw_filesystem/src/root.rs
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,17 @@ pub trait RootFilesystem: Send + Sync {
Ok(records.into_iter().map(|record| record.seq).max())
}

/// Reserve and return the next monotonic sequence number for `path`.
///
/// Unlike [`append`](Self::append), this sequence is scoped to `path`
/// rather than the backend's global event table. Consumers use it to
/// assign row-native ordering keys without rewriting a shared metadata
/// record under CAS. A failed follow-up write may leave a gap; callers must
/// rely on monotonicity, not contiguity.
async fn reserve_sequence(&self, path: &VirtualPath) -> Result<SeqNo, FilesystemError> {
unsupported(path, FilesystemOperation::ReserveSeq)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

// ─── Legacy bytes plane (DEPRECATED — removed after consumer migration) ─
//
// The methods below predate the unified [`put`]/[`get`] surface and exist
Expand Down
Loading
Loading