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
52 changes: 27 additions & 25 deletions crates/ironclaw_turns/src/filesystem_store.rs
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
//! Filesystem-backed [`TurnStateStore`] implementation.
//!
//! Persists the lower-churn [`TurnPersistenceSnapshot`] as a JSON blob under
//! the `/turns` mount alias (alias-relative path: `/turns/state.json`) and
//! high-churn runner lease heartbeats as per-run CAS records under
//! `/turns/runner-leases`. Snapshot mutations read the snapshot, overlay
//! current runner leases, delegate to an [`InMemoryTurnStateStore`] in a
//! transient `apply` closure, and write the resulting snapshot back with
//! the `/turns` mount alias (alias-relative path: `/turns/state.json`).
//! High-churn runner lease heartbeats are memory-backed per
//! [`FilesystemTurnStateStore`] instance and are overlaid onto the durable
//! snapshot while the process is alive. Snapshot mutations read the snapshot,
//! overlay current runner leases, delegate to an [`InMemoryTurnStateStore`] in
//! a transient `apply` closure, and write the resulting snapshot back with
//! optimistic CAS + bounded retry. Reads load the snapshot, overlay current
//! runner leases, and project through the in-memory store without writing
//! back.
Expand All @@ -22,7 +23,6 @@
//!
//! ```text
//! /turns/state.json
//! /turns/runner-leases/{run_id}.json
//! ```
//!
//! Within-tenant scoping (agent/project/thread) is encoded inside the
Expand All @@ -33,13 +33,15 @@
//! path itself.

use std::{
collections::HashMap,
sync::{Arc, Mutex},
time::{Duration, Instant},
};

use async_trait::async_trait;
use ironclaw_filesystem::{CasExpectation, RecordVersion, RootFilesystem, ScopedFilesystem};
use ironclaw_host_api::{ResourceScope, ScopedPath, UserId};
use tokio::sync::RwLock;

use crate::{
AllowAllTurnAdmissionLimitProvider, CancelRunRequest, CancelRunResponse, EventCursor,
Expand Down Expand Up @@ -70,7 +72,7 @@ use io::{
put_with_cas, snapshot_entry, snapshot_path,
};
use profile_resolver::PreResolvedRunProfileResolver;
use runner_lease::{RunnerLeaseOverlay, RunnerLeaseRecord, RunnerLeaseSidecar};
use runner_lease::{RunnerLeaseMemory, RunnerLeaseOverlay, RunnerLeaseRecord, RunnerLeaseStore};

#[cfg(test)]
mod tests;
Expand Down Expand Up @@ -124,6 +126,7 @@ where
limits: InMemoryTurnStateStoreLimits,
admission_limit_provider: Arc<dyn TurnAdmissionLimitProvider>,
snapshot_cache: Mutex<Option<CachedSnapshot>>,
runner_leases: RunnerLeaseMemory,
apply_timeout: Duration,
}

Expand All @@ -137,6 +140,7 @@ where
limits: InMemoryTurnStateStoreLimits::default(),
admission_limit_provider: Arc::new(AllowAllTurnAdmissionLimitProvider),
snapshot_cache: Mutex::new(None),
runner_leases: Arc::new(RwLock::new(HashMap::new())),
apply_timeout: FILESYSTEM_APPLY_TIMEOUT,
}
}
Expand Down Expand Up @@ -225,38 +229,36 @@ where
snapshot: (TurnPersistenceSnapshot, Option<RecordVersion>),
overlay: RunnerLeaseOverlay,
) -> Result<(TurnPersistenceSnapshot, Option<RecordVersion>), TurnError> {
self.runner_lease_sidecar().overlay(snapshot, overlay).await
self.runner_lease_store().overlay(snapshot, overlay).await
}

async fn seed_runner_lease_from_snapshot_inner(
&self,
run_id: TurnRunId,
) -> Result<(), TurnError> {
let (snapshot, _version) = self.read_snapshot_from_filesystem().await?;
self.runner_lease_sidecar()
self.runner_lease_store()
.seed_from_snapshot(&snapshot, run_id)
.await?;
self.clear_snapshot_cache();
Ok(())
}

async fn cleanup_runner_lease_after_state(&self, result: &Result<TurnRunState, TurnError>) {
self.runner_lease_sidecar()
.cleanup_after_state(result)
.await;
self.runner_lease_store().cleanup_after_state(result).await;
self.clear_snapshot_cache();
}

async fn heartbeat_runner_lease(
&self,
request: HeartbeatRequest,
) -> Result<EventCursor, TurnError> {
let sidecar = self.runner_lease_sidecar();
let cursor = match sidecar.heartbeat(request.clone()).await {
let lease_store = self.runner_lease_store();
let cursor = match lease_store.heartbeat(request.clone()).await {
Err(TurnError::ScopeNotFound) => {
self.seed_missing_runner_lease_from_snapshot(request.run_id)
.await?;
self.runner_lease_sidecar().heartbeat(request).await?
self.runner_lease_store().heartbeat(request).await?
}
result => result?,
};
Expand All @@ -269,7 +271,7 @@ where
run_id: TurnRunId,
) -> Result<(), TurnError> {
let (snapshot, _version) = self.read_snapshot_from_filesystem().await?;
self.runner_lease_sidecar()
self.runner_lease_store()
.seed_from_snapshot_if_missing(&snapshot, run_id)
.await
}
Expand All @@ -292,7 +294,7 @@ where
) {
return Ok(None);
}
self.runner_lease_sidecar()
self.runner_lease_store()
.mark_cancel_requested_from_snapshot(&snapshot, request.run_id)
.await
}
Expand All @@ -305,7 +307,7 @@ where
retired_status: TurnStatus,
) -> Result<Option<RunnerLeaseRecord>, TurnError> {
let (snapshot, _version) = self.read_snapshot_from_filesystem().await?;
self.runner_lease_sidecar()
self.runner_lease_store()
.retire_runner_lease_from_snapshot(
&snapshot,
run_id,
Expand All @@ -324,15 +326,15 @@ where
let Some(previous) = previous else {
return;
};
self.runner_lease_sidecar()
self.runner_lease_store()
.restore_if_current_status(previous, current_status)
.await;
self.clear_snapshot_cache();
}

fn runner_lease_sidecar(&self) -> RunnerLeaseSidecar<F> {
RunnerLeaseSidecar::new(
Arc::clone(&self.filesystem),
fn runner_lease_store(&self) -> RunnerLeaseStore {
RunnerLeaseStore::new(
Arc::clone(&self.runner_leases),
self.limits.runner_lease_ttl,
self.apply_timeout,
)
Expand Down Expand Up @@ -503,7 +505,7 @@ where
tracing::debug!(
run_id = %run_id,
error = %error,
"failed to compensate turn claim after runner lease sidecar seed failed"
"failed to compensate turn claim after memory runner lease seed failed"
);
}
self.clear_snapshot_cache();
Expand Down Expand Up @@ -582,7 +584,7 @@ where
let response = result?;
match response.status {
status if status.is_terminal() => {
self.runner_lease_sidecar()
self.runner_lease_store()
.delete_best_effort(response.run_id)
.await;
}
Expand Down Expand Up @@ -789,7 +791,7 @@ where
.await;
if let Ok(response) = &result {
for state in &response.recovered {
self.runner_lease_sidecar()
self.runner_lease_store()
.delete_best_effort(state.run_id)
.await;
}
Expand Down
Loading
Loading