pub struct RunStore { /* private fields */ }Expand description
JSONL run-trace store rooted at <car_dir>/runs/.
Stateless across calls — each append opens, writes, and closes the run’s file. Flush points are sparse (turn granularity), so there is no long-lived file handle to manage, and a concurrently-restarting daemon always sees a consistent on-disk tail.
Implementations§
Source§impl RunStore
impl RunStore
Sourcepub fn retention(&self) -> RetentionConfig
pub fn retention(&self) -> RetentionConfig
Construct a store rooted at runs_root (the runs/ dir itself).
Use RunStore::from_journal_dir from the daemon, which derives
the root from the configured journal dir; this constructor is the
test/embedder seam.
The retention policy in force, for a caller reporting what GC did.
pub fn new(runs_root: PathBuf, retention: RetentionConfig) -> Self
pub fn with_failure_injector(self, failures: RunStoreFailureInjector) -> Self
pub fn with_summary_read_gate(self, gate: RunStoreSummaryReadGate) -> Self
pub fn with_append_gate(self, gate: RunStoreAppendGate) -> Self
pub fn with_lookup_gate(self, gate: RunStoreLookupGate) -> Self
Sourcepub fn with_private_path_failure_injector(
self,
failures: PrivatePathDurabilityFailureInjector,
) -> Self
pub fn with_private_path_failure_injector( self, failures: PrivatePathDurabilityFailureInjector, ) -> Self
Test/embedder seam for deterministic first-use directory-entry faults.
Sourcepub fn from_journal_dir(journal_dir: &Path) -> Self
pub fn from_journal_dir(journal_dir: &Path) -> Self
Derive the store from the daemon’s journal dir. The journal lives
at ~/.car/journals, so the run store is its sibling
~/.car/runs; retention is read from ~/.car/config.toml. When
the journal dir has no parent (a bare relative path), the store
falls back to journal_dir/../runs resolved lexically.
Sourcepub fn run_started(&self, run_id: &str) -> Result<Option<RunStarted>>
pub fn run_started(&self, run_id: &str) -> Result<Option<RunStarted>>
Read the exact durable RunStarted boundary for run_id without
collapsing absence into corruption. A zero-byte file left by an open
that failed before its first write is still absent; any non-empty file
without one valid, unambiguous RunStarted is invalid durable state.
Sourcepub fn rollback_empty_run_start(&self, started: &RunStarted) -> Result<()>
pub fn rollback_empty_run_start(&self, started: &RunStarted) -> Result<()>
Remove only the zero-byte trace created when runs.start opened its
destination but failed before writing RunStarted. The full expected
identity is accepted so callers cannot use this as a broad delete
primitive. Any bytes or parsed rows make rollback fail closed.
pub fn pending_provenance( &self, pending: &PendingProposalFinalization, ) -> Result<(RunStarted, ProposalExecutionMarker)>
Sourcepub fn write_execution_marker(
&self,
marker: &ProposalExecutionMarker,
) -> Result<()>
pub fn write_execution_marker( &self, marker: &ProposalExecutionMarker, ) -> Result<()>
Write the durable execution_in_progress marker before tool dispatch.
pub fn execution_marker( &self, run_id: &str, ) -> Result<Option<ProposalExecutionMarker>>
Sourcepub fn write_proposal_retry_rollback(
&self,
rollback: &ProposalRetryRollback,
) -> Result<()>
pub fn write_proposal_retry_rollback( &self, rollback: &ProposalRetryRollback, ) -> Result<()>
Persist the pre-execution rollback authority before acquiring the content-addressed retry owner. Exact rewrites are idempotent; a conflicting preimage for the same run fails closed.
pub fn proposal_retry_rollback( &self, run_id: &str, ) -> Result<Option<ProposalRetryRollback>>
pub fn clear_proposal_retry_rollback( &self, expected: &ProposalRetryRollback, ) -> Result<()>
Sourcepub fn reconcile_proposal_retry_rollbacks(&self) -> Result<usize>
pub fn reconcile_proposal_retry_rollbacks(&self) -> Result<usize>
Resolve every crash-left rollback intent. The intent is removed last: a second crash at any earlier boundary leaves enough authority for the next startup to repeat the exact cleanup.
Sourcepub fn claim_proposal_id(
&self,
run_id: &str,
client_id: &str,
proposal_id: &str,
original_submission: &Value,
) -> Result<ProposalIdClaimOutcome>
pub fn claim_proposal_id( &self, run_id: &str, client_id: &str, proposal_id: &str, original_submission: &Value, ) -> Result<ProposalIdClaimOutcome>
Claim a proposal identifier inside one authenticated run before any dispatch. Exact-submission retries are idempotent; reusing the same id for changed bytes is rejected durably, including after restart.
pub fn clear_execution_marker( &self, expected: &ProposalExecutionMarker, ) -> Result<()>
Sourcepub fn write_pending_proposal(
&self,
pending: &PendingProposalFinalization,
) -> Result<()>
pub fn write_pending_proposal( &self, pending: &PendingProposalFinalization, ) -> Result<()>
Persist one exact proposal-finalization transaction with private permissions and an fsynced atomic replacement. Exact retries are idempotent; a conflicting preimage for the same run is rejected.
pub fn pending_proposal( &self, run_id: &str, ) -> Result<Option<PendingProposalFinalization>>
pub fn all_pending_proposals(&self) -> Result<Vec<PendingProposalFinalization>>
Sourcepub fn clear_pending_proposal(
&self,
expected: &PendingProposalFinalization,
) -> Result<()>
pub fn clear_pending_proposal( &self, expected: &PendingProposalFinalization, ) -> Result<()>
Remove only the exact, provenance-validated transaction whose terminal was acknowledged. The caller clears this before its matching execution marker so provenance remains available for the check.
Sourcepub fn reserve_proposal_retry_owner(
&self,
run_id: &str,
client_id: &str,
requested_policy_session_id: Option<&str>,
original_submission: &Value,
) -> Result<ProposalRetryReservation>
pub fn reserve_proposal_retry_owner( &self, run_id: &str, client_id: &str, requested_policy_session_id: Option<&str>, original_submission: &Value, ) -> Result<ProposalRetryReservation>
Atomically reserve the process-wide retry tuple before any execution marker or runtime invocation. Direct exclusive creation is the linear point: one run/client wins, every concurrent or restarted contender observes that durable owner, and a partial claim fails closed.
Sourcepub fn release_proposal_retry_owner(
&self,
run_id: &str,
client_id: &str,
requested_policy_session_id: Option<&str>,
original_submission: &Value,
) -> Result<()>
pub fn release_proposal_retry_owner( &self, run_id: &str, client_id: &str, requested_policy_session_id: Option<&str>, original_submission: &Value, ) -> Result<()>
Release only the exact pre-dispatch retry owner. This is used when execution never became admissible (marker durability or policy bind failed), so a later scheduled run may safely submit the same proposal. The content-addressed owner is validated before unlink and the owner directory is fsynced before success is reported.
Sourcepub fn write_completed_proposal(
&self,
pending: &PendingProposalFinalization,
) -> Result<CompletedProposalResponse>
pub fn write_completed_proposal( &self, pending: &PendingProposalFinalization, ) -> Result<CompletedProposalResponse>
Persist an exact typed response only after the matching critical
proposal_completed append acknowledged. The finalization outbox and
execution marker must still be present and provenance-valid here, so a
forged receipt cannot become a cleanup authority.
Sourcepub fn completed_proposal(
&self,
run_id: &str,
client_id: &str,
requested_policy_session_id: Option<&str>,
original_submission: &Value,
) -> Result<Option<CompletedProposalResponse>>
pub fn completed_proposal( &self, run_id: &str, client_id: &str, requested_policy_session_id: Option<&str>, original_submission: &Value, ) -> Result<Option<CompletedProposalResponse>>
Load only the receipt whose immutable retry tuple exactly matches this bound run/client/policy/raw submission. Absence permits a new proposal; a present corrupt or self-inconsistent receipt fails closed.
Sourcepub fn completed_proposal_for_resumed_owner(
&self,
run_id: &str,
durable_client_id: &str,
original_submission: &Value,
) -> Result<Option<CompletedProposalResponse>>
pub fn completed_proposal_for_resumed_owner( &self, run_id: &str, durable_client_id: &str, original_submission: &Value, ) -> Result<Option<CompletedProposalResponse>>
Recover the sole retained response for an exact proposal after an authenticated run owner reconnects with a newly minted policy session.
The lookup is deliberately bounded to one run and one durable client. It never treats the replacement policy-session id as durable identity, and ambiguity fails closed instead of selecting an arbitrary receipt.
pub fn all_completed_proposals(&self) -> Result<Vec<CompletedProposalResponse>>
Sourcepub fn reconcile_completed_proposal_migration(&self) -> Result<()>
pub fn reconcile_completed_proposal_migration(&self) -> Result<()>
Backfill the content-addressed retry-owner index exactly once for receipts written before that index existed. Permanent receipts grow without bound by design, so normal daemon startup must never enumerate them. After this versioned checkpoint is durable, crash cleanup is driven by the bounded pending-finalization outbox instead.
pub fn completed_proposal_retry_owner( &self, requested_policy_session_id: Option<&str>, original_submission: &Value, ) -> Result<Option<(String, String)>>
Sourcepub fn cleanup_completed_proposal_guards(
&self,
receipt: &CompletedProposalResponse,
) -> Result<()>
pub fn cleanup_completed_proposal_guards( &self, receipt: &CompletedProposalResponse, ) -> Result<()>
Remove leftover pre-response guards only under a completed receipt’s exact typed authority. Each deletion is independently durable; a fault leaves the receipt and any remaining guard for safe retry/restart.
pub fn completed_proposal_response_value( &self, receipt: &CompletedProposalResponse, ) -> Result<Value>
Sourcepub fn prepare_storage(&self) -> Result<()>
pub fn prepare_storage(&self) -> Result<()>
Prepare the CAR-owned run tree for daemon startup. This is intentionally fallible: the daemon must refuse adoption/listening when the root or backup marker cannot be proven owner-private.
Sourcepub fn append_records(
&self,
agent_id: &str,
run_id: &str,
records: &[RunRecord],
) -> Result<()>
pub fn append_records( &self, agent_id: &str, run_id: &str, records: &[RunRecord], ) -> Result<()>
Append one or more RunRecords to a run’s file, creating it 0600
on first write. Records are written one JSONL line each, in order.
This is the single low-level flush primitive the wiring calls at
each boundary: RunStarted on runs.start, RunTurns as the
recorder produces them, and the terminal RunEnded/Incomplete on
runs.complete/disconnect.
Sourcepub fn write_started(&self, started: &RunStarted) -> Result<()>
pub fn write_started(&self, started: &RunStarted) -> Result<()>
Append the RunStarted line + create the run file (runs.start).
Sourcepub fn append_turns(
&self,
agent_id: &str,
run_id: &str,
turns: &[RunRecord],
) -> Result<()>
pub fn append_turns( &self, agent_id: &str, run_id: &str, turns: &[RunRecord], ) -> Result<()>
Append RunTurn records (the recorder’s per-proposal output).
Sourcepub fn run_trace_corruption_for(
&self,
agent_id: &str,
run_id: &str,
) -> Result<Option<RunTraceCorruption>>
pub fn run_trace_corruption_for( &self, agent_id: &str, run_id: &str, ) -> Result<Option<RunTraceCorruption>>
Read the durable summary’s corruption marker without materializing the trace. Active subscribe reads this off the async lock path, then revalidates live state before it trusts a snapshot or registers.
Sourcepub fn ensure_proposal_turns(
&self,
pending: &PendingProposalFinalization,
) -> Result<ProposalTraceEnsure>
pub fn ensure_proposal_turns( &self, pending: &PendingProposalFinalization, ) -> Result<ProposalTraceEnsure>
Ensure the exact action trace authenticated by a durable proposal finalization exists once before the terminal journal/receipt can advance. Existing exact rows make retries idempotent; partial or conflicting rows fail closed.
Sourcepub fn write_ended(&self, ended: &RunEnded) -> Result<()>
pub fn write_ended(&self, ended: &RunEnded) -> Result<()>
Append the terminal RunEnded line (runs.complete or the
disconnect-Incomplete path).
Sourcepub fn write_cancellation_requested(
&self,
agent_id: &str,
requested: &RunCancellationRequested,
) -> Result<()>
pub fn write_cancellation_requested( &self, agent_id: &str, requested: &RunCancellationRequested, ) -> Result<()>
Durably append the body-free cancellation request before attempting control. The exact key/preimage is idempotent; another key cannot race the active request.
Sourcepub fn write_cancellation_result(
&self,
agent_id: &str,
result: &RunCancelResponse,
) -> Result<()>
pub fn write_cancellation_result( &self, agent_id: &str, result: &RunCancelResponse, ) -> Result<()>
Persist an unconfirmed cancellation receipt. Confirmed cancellation is
represented by the terminal RunEnded::Cancelled row instead.
Sourcepub fn get_run_trace(&self, run_id: &str) -> Option<Vec<RunRecord>>
pub fn get_run_trace(&self, run_id: &str) -> Option<Vec<RunRecord>>
Load a run’s full ordered trace from disk by run_id — the U5
runs.get_trace read path. Works after a restart when memory is
empty. Resolves run_id -> agent_id by scanning the tree, then
reads the JSONL. A malformed unterminated final crash tail is ignored;
committed corruption makes this compatibility wrapper return None.
Returns None when no file exists for the run_id.
Sourcepub fn get_run_trace_checked(
&self,
run_id: &str,
) -> Result<Option<Vec<RunRecord>>>
pub fn get_run_trace_checked( &self, run_id: &str, ) -> Result<Option<Vec<RunRecord>>>
Strict replay read used by trust-bearing RPC surfaces. A malformed newline-terminated row rejects the trace; only a torn final row is ignored as an uncommitted crash tail.
Sourcepub fn get_run_trace_for(
&self,
agent_id: &str,
run_id: &str,
) -> Option<Vec<RunRecord>>
pub fn get_run_trace_for( &self, agent_id: &str, run_id: &str, ) -> Option<Vec<RunRecord>>
Load a run’s trace given both keys (cheaper — no tree scan). Used
by list_runs internally and available to callers that already
know the owning agent.
pub fn get_run_trace_for_checked( &self, agent_id: &str, run_id: &str, ) -> Result<Option<Vec<RunRecord>>>
Sourcepub fn get_run_trace_page_for(
&self,
agent_id: &str,
run_id: &str,
cursor: usize,
limit: usize,
) -> Result<Option<(Vec<RunRecord>, Option<usize>)>>
pub fn get_run_trace_page_for( &self, agent_id: &str, run_id: &str, cursor: usize, limit: usize, ) -> Result<Option<(Vec<RunRecord>, Option<usize>)>>
Stream one bounded page from a run’s JSONL trace. limit + 1
records are retained only to determine whether a continuation exists;
accumulated history is never materialized by the dashboard read path.
Sourcepub fn get_run_turn_page_for(
&self,
agent_id: &str,
run_id: &str,
cursor: usize,
limit: usize,
) -> Result<Option<(Vec<RunRecord>, Option<usize>, usize, RunStatus)>>
pub fn get_run_turn_page_for( &self, agent_id: &str, run_id: &str, cursor: usize, limit: usize, ) -> Result<Option<(Vec<RunRecord>, Option<usize>, usize, RunStatus)>>
Stream a bounded turns-only page for runs.subscribe. The exact total
comes from the durable summary sidecar, so returning live_cursor never
requires materializing the trace.
Sourcepub fn list_runs_page(
&self,
agent_id: &str,
cursor: usize,
limit: usize,
) -> Result<(Vec<RunSummary>, Option<usize>)>
pub fn list_runs_page( &self, agent_id: &str, cursor: usize, limit: usize, ) -> Result<(Vec<RunSummary>, Option<usize>)>
Read a newest-first run-summary page through the durable fixed-record
order index. The page performs at most limit + 1 sidecar reads and no
run-directory enumeration or trace scan.
Sourcepub fn list_runs(&self, agent_id: &str) -> Vec<RunSummary>
pub fn list_runs(&self, agent_id: &str) -> Vec<RunSummary>
List an agent’s runs newest-first — the U5 runs.list read path.
Each summary is built from the run file’s records (head for
RunStarted, tail for the terminal record, count of Turns).
Returns an empty Vec for an agent with no runs (the empty-state).
Sourcepub fn visit_run_boundaries<F>(&self, visitor: F)where
F: FnMut(RunStarted, Option<RunEnded>, Option<RunCancellationRequested>, Option<RunCancelResponse>),
pub fn visit_run_boundaries<F>(&self, visitor: F)where
F: FnMut(RunStarted, Option<RunEnded>, Option<RunCancellationRequested>, Option<RunCancelResponse>),
Visit only the durable lifecycle boundaries for one trace at a time. Startup reconciliation needs neither turns nor a process-wide snapshot; streaming each file prevents accumulated newsroom history from being materialized before the daemon can listen and schedule.
Sourcepub fn agent_for_run(&self, run_id: &str) -> Option<String>
pub fn agent_for_run(&self, run_id: &str) -> Option<String>
Resolve the owning agent_id for a run_id from disk. Mirrors
Self::resolve_run_file but returns the agent dir name — U5’s
authorization check (KTD10) needs run_id -> agent_id to verify
ownership before serving a trace.
Sourcepub fn gc(&self) -> usize
pub fn gc(&self) -> usize
Retention GC (R6) — call on daemon boot. Per agent: keep the most
recent max_per_agent completed runs and drop any completed
run older than max_age_days, whichever is more restrictive. An
in-progress run (no terminal record) is NEVER evicted. Returns the
number of run files removed.
Sourcepub fn adopt_orphans(&self) -> usize
pub fn adopt_orphans(&self) -> usize
Adopt crash-orphaned runs at boot (FIX 4). A daemon crash mid-run
leaves an on-disk run with RunStarted (+ Turns) but no terminal
RunEnded, so it reads InProgress forever and the GC — which never
evicts an in-progress run — can never reclaim it. The file leaks
across every crash.
This runs at store construction/boot, BEFORE gc(). At that moment
the in-memory runs map is always empty, so any run that is
InProgress on disk cannot have a live harness writing to it — it is
necessarily a crashed prior process. We append an Incomplete
terminal RunEnded marker to adopt it, making it terminal and thus
age-GC-eligible (so a later gc() in the same boot can reclaim it).
Returns the number of runs adopted. Best-effort: an unwritable file is skipped rather than failing startup.
Trait Implementations§
Auto Trait Implementations§
impl Freeze for RunStore
impl RefUnwindSafe for RunStore
impl Send for RunStore
impl Sync for RunStore
impl Unpin for RunStore
impl UnsafeUnpin for RunStore
impl UnwindSafe for RunStore
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
impl<S, T> Duplex<S> for Twhere
T: FromSample<S> + ToSample<S>,
impl<T> ErasedDestructor for Twhere
T: 'static,
Source§impl<S> FromSample<S> for S
impl<S> FromSample<S> for S
fn from_sample_(s: S) -> S
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more