Skip to main content

JournalStorage

Struct JournalStorage 

pub struct JournalStorage { /* private fields */ }
Expand description

A study in one append-only file: grep-able, resumable, and shared between processes with no server.

This is the flagship backend (the four backends). A journal is JSON Lines: a versioned header record, then exactly one record per mutation. Each mutation stages the existing complete prefix plus one new record and atomically replaces the path, preserving the observable append-only format while giving the lease a real fencing point. That design buys all of:

  • resume is replay. Opening a journal folds its records into the same in-memory shape MemStorage keeps, so a resumed study is not similar to the one that was interrupted, it is the same one. Nothing has to be flushed on shutdown, because nothing is ever held back.
  • multi-process with no coordinator. A SLURM array of workers against one journal on a shared filesystem is the whole distribution story: zero servers, zero ports, zero database. Appends are serialised by an advisory lock built from the one primitive NFS actually guarantees (see JournalOptions::lock_grace_millis and the module source); readers never take it.
  • a study you can read with grep. Every record is one line of JSON tagged with its op: grep '"kind":"set-values"' study.jsonl is the whole result set, on a cluster login node, with no atune binary in sight.
  • crash tolerance by construction. A worker killed while staging leaves only an unreferenced ticket; a worker killed after publication leaves a complete image. A torn old tail is excluded from the next image, because a record that was never complete is a record that never happened.

§Files

PathWhat
<path>the log — the only authoritative file
<path>.lockthe renewable append-lease directory
<path>.lock.append-<generation>an ephemeral generation-bound publish ticket
<path>.snapshotan optimization: the replayed state as of op n

Deleting a study means deleting all three. A snapshot left next to a recreated journal is normally recognised as unrelated (it carries the journal’s header line) and ignored, but there is no reason to rely on that.

§Durability

Exactly one complete image is published per mutating call, with the append lock held. If the current image is B bytes, one mutation copies O(B) bytes before adding its record. Over N mutations, cumulative copied volume is Σ O(B_i), potentially O(N²) as the image grows. The journal_bytes_staged and journal_canonical_size metrics expose this work shape; they are not throughput measurements. Both durability settings sync the staged image before publication. Durability::Record additionally requests a containing-directory sync after the rename; Durability::Os leaves that metadata durability to the operating system. The publication is an atomic same-directory rename, so readers see either the previous complete image or the new one. Snapshots are likewise synced before they are published, and are either wholly there or wholly absent. These local sync calls are not evidence of power-loss behavior on a particular mount.

§Concurrency

Threads of one process share one journal handle and are serialised by an in-process mutex; separate processes are serialised by the lock file. Reads take the in-process mutex but never the file lock: they catch up on the tail they have not read yet, which is how a worker learns what the other workers did. (Because a write holds the mutex across acquiring the file lock, an in-process read can wait behind a concurrent writer — up to that writer’s lock-acquisition time — but it never contends for the file lock itself.) A Cursor is therefore durable and portable — it is a study-local sequence number derived from the log, identical in every process that replays the same file, not a memory offset that resets when a worker restarts.

use atune_core::prelude::*;
use atune_core::storage::JournalStorage;

let dir = std::env::temp_dir().join("atune-journal-doc-basic");
std::fs::create_dir_all(&dir).expect("create the study directory");
let path = dir.join("study.jsonl");

let storage = JournalStorage::open(&path)?;
let study = storage.create_study(StudyConfig::new("demo").with_seed(7))?;
let trial = storage.create_trial(study, None)?;
storage.try_transition(trial.id, TrialState::Waiting, TrialState::Running)?;
storage.set_values(trial.id, &[0.25])?;
storage.try_transition(trial.id, TrialState::Running, TrialState::Complete)?;
drop(storage);

// A second process — here, a second handle — resumes by replaying.
let resumed = JournalStorage::open(&path)?;
assert_eq!(resumed.get_trial(trial.id)?.objective_value(), Some(&[0.25][..]));
assert_eq!(resumed.study_config(study)?.seed(), 7);

Implementations§

§

impl JournalStorage

pub fn open(path: impl AsRef<Path>) -> Result<Self>

Opens (creating if necessary) the journal at path, with the defaults.

§Errors

Error::Storage if the file cannot be created or read, if it is not a journal, if its format version is newer than this build understands, or if it is corrupt in a way replay cannot explain (a gap in the record sequence, an identity that does not match the one replay derives). A torn final record is none of those: it is discarded. Error::ResourceLimit is returned when the finite byte, operation, snapshot, or record-entry limits are exceeded.

pub fn with_options( path: impl AsRef<Path>, options: JournalOptions, ) -> Result<Self>

Opens the journal at path with options.

§Errors

As open.

pub fn open_read_only(path: impl AsRef<Path>) -> Result<Self>

Opens the journal at path for read-only queries.

This is the path a pure reader — atune best, atune trials, any command that only inspects a study — should take. Unlike open it neither creates the file nor takes the append lock, so it will not leave an empty journal behind on a mistyped path and it succeeds in a directory the process cannot write (where taking the .lock would fail with a permission error).

The returned handle serves every read of the Storage trait — study_config, sync, get_trial, get_state. A mutating call will fail because this handle is read-only; use open for anything that writes. This is an inherent method, not part of the frozen Storage trait.

§Errors

Error::Storage if the file does not exist or cannot be read, if it is not a journal, or if its format version is newer than this build understands. Error::ResourceLimit is returned when the finite byte, operation, snapshot, or record-entry limits are exceeded. A missing file is reported rather than created — the caller distinguishes “no study yet” by checking existence first if it wants a friendlier message.

pub fn open_read_only_with( path: impl AsRef<Path>, options: JournalOptions, ) -> Result<Self>

Opens the journal read-only with explicit finite admission options.

Unlike with_options, this never creates the journal or takes its append lock. The options still govern every bounded read and snapshot admission.

§Errors

Returns Error::Storage when the path or its journal image cannot be opened, read, or validated under the supplied options, or Error::ResourceLimit when a finite option cap is exceeded.

pub fn open_read_only_with_options( path: impl AsRef<Path>, options: JournalOptions, ) -> Result<Self>

Alias for open_read_only_with.

§Errors

Returns the same errors as open_read_only_with.

pub fn path(&self) -> &Path

The journal’s path.

pub fn lock_path(&self) -> &Path

The append lock’s path.

pub fn snapshot_path(&self) -> &Path

The snapshot side-file’s path.

pub const fn worker(&self) -> WorkerId

This handle’s worker identity.

pub fn ops(&self) -> Result<u64>

How many op records this handle has replayed.

Catches up first, so it counts what is in the file, including what other processes have written.

§Errors

Error::Storage if the tail cannot be read.

pub fn snapshot(&self) -> Result<()>

Writes a snapshot now, whatever the configured cadence.

Useful before a long pause, and the way to find out why automatic snapshots are failing: the periodic ones are best-effort and silent (a mutation that succeeded must not be reported as failed because an optimization could not be written), this one reports.

§Errors

Error::Storage if the tail cannot be read or the snapshot cannot be written, synced or published. Error::ResourceLimit is returned when the finite snapshot byte or collection-entry cap is exceeded.

pub fn fail_over_stale( &self, study: StudyId, now_millis: u64, ) -> Result<Vec<FailedOver>>

Fails over every trial of study whose heartbeat has gone quiet.

A Running trial whose last heartbeat is older than the grace period (JournalOptions::stale_after_millis, by default twice the heartbeat interval) belongs to a worker that is not coming back: the node was preempted, the job hit its wall clock, somebody kill -9ed it. Any worker may take it over — that is what “workers coordinate only through storage” means when a worker dies.

The trial is failed with an explanatory message, and — if JournalOptions::requeue_on_failover is on — its parameters are re-enqueued as a fresh trial, so the configuration is still evaluated. A trial that never heartbeat at all is judged on its creation time, because “started long ago and never said anything” is the same evidence as “said something long ago”.

The shared core-side mechanism is kept out of any one backend: its implementation uses the public TrialLifecycleStorage capability to observe and conditionally fail the exact owner epoch. This handle keeps the unchanged-observation proof across calls; now_millis therefore must come from one monotonic sweeper clock. Owner heartbeat wall-clock values are not compared with it.

§Races

Two workers may decide to fail the same trial over at the same moment. Exactly one wins the Running -> Failed transition; the loser reports nothing for that trial, so a caller can treat the returned list as “trials I failed over”.

§Errors

Error::NotFound if the study does not exist; Error::Storage on a backend fault.

Trait Implementations§

§

impl AtomicImportStorage for JournalStorage

Applies one import in a shadow journal while the source generation lease is held, then publishes the shadow’s complete image exactly once.

§

fn import_batch(&self, batch: &ImportBatch) -> Result<ImportBatchResult>

Publishes one completely staged batch at one backend commit point. Read more
§

impl Debug for JournalStorage

§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
§

impl Storage for JournalStorage

§

fn backing_paths(&self) -> Option<Vec<PathBuf>>

Local files that an exporter must never replace, including fixed sidecars. Read more
§

fn create_study(&self, cfg: StudyConfig) -> Result<StudyId>

Creates a study and returns its identity. Read more
§

fn study_config(&self, study: StudyId) -> Result<StudyConfig>

Reads a study’s configuration back. Read more
§

fn study_user_attrs(&self, study: StudyId) -> Result<UserAttrs>

Returns all user-owned JSON metadata attached to a study. Read more
§

fn set_study_user_attr( &self, study: StudyId, key: &str, value: Value, ) -> Result<()>

Inserts or replaces one study user attribute. Read more
§

fn create_trial( &self, study: StudyId, template: Option<TrialTemplate>, ) -> Result<TrialRecord>

Creates a trial, optionally pre-loaded from a template. Read more
§

fn try_transition( &self, trial: TrialId, from: TrialState, to: TrialState, ) -> Result<bool>

Atomically moves a trial from from to to. Read more
§

fn write_params( &self, trial: TrialId, batch: &[(String, Distribution, ParamValue)], ) -> Result<()>

Records sampled parameters together with the distributions they came from. Read more
§

fn report(&self, trial: TrialId, step: u64, values: &[f64]) -> Result<()>

Appends an intermediate report. Read more
§

fn set_values(&self, trial: TrialId, values: &[f64]) -> Result<()>

Records a trial’s final objective values. Read more
§

fn set_error(&self, trial: TrialId, message: &str) -> Result<()>

Records why a trial failed. Read more
§

fn heartbeat(&self, trial: TrialId, now_millis: u64) -> Result<()>

Records that the trial’s owner is still alive at now_millis. Read more
§

fn set_trial_user_attr( &self, trial: TrialId, key: &str, value: Value, ) -> Result<()>

Inserts or replaces one user-owned JSON attribute on an unfinished trial and marks that trial changed for sync. Read more
§

fn put_state(&self, scope: Scope, blob: StateBlob) -> Result<()>

Writes the state blob of scope, replacing any previous one. Read more
§

fn trial_lifecycle(&self) -> Option<&dyn TrialLifecycleStorage>

Returns the backend’s atomic trial-lifecycle capability, if available. Read more
§

fn atomic_import(&self) -> Option<&dyn AtomicImportStorage>

Returns the backend’s whole-import atomic-publication capability, if available. Read more
§

fn study_catalog(&self) -> Option<&dyn StudyCatalog>

Returns the backend’s study-catalog capability, if available. Read more
§

fn get_state(&self, scope: Scope) -> Result<Option<StateBlob>>

Reads the state blob of scope, if one was ever written. Read more
§

fn read_snapshot(&self, request: &SnapshotRequest) -> Result<StorageSnapshot>

Reads trial deltas, requested side state, and optional finalization records for one observation. Read more
§

fn sync( &self, study: StudyId, cursor: Cursor, ) -> Result<(Vec<TrialDelta>, Cursor)>

Returns everything that happened in study after cursor. Read more
§

fn get_trial(&self, trial: TrialId) -> Result<FrozenTrial>

Reads one trial’s snapshot. Read more
§

fn sync_limited( &self, study: StudyId, cursor: Cursor, max_trials: usize, ) -> Result<SyncPage>

sync, bounded to at most max_trials deltas. Read more
§

impl StudyCatalog for JournalStorage

§

fn create_if_empty(&self, cfg: StudyConfig) -> Result<StudyCreateOutcome>

Atomically creates cfg only when the catalog has no studies. Read more
§

fn list_studies(&self) -> Result<Vec<StudySummary>>

Lists study identities in deterministic ascending StudyId order. Read more
§

fn find_by_name(&self, name: &str) -> Result<StudyLoadOutcome>

Resolves an exact name without selecting arbitrarily among duplicates. Read more
§

fn list(&self) -> Result<Vec<StudySummary>>

Short alias for list_studies. Read more
§

impl TrialLifecycleStorage for JournalStorage

§

fn observe_finalization( &self, study: StudyId, ) -> Result<FinalizationObservation>

Reads the study-wide finalization ownership observation. Read more
§

fn acquire_finalization( &self, request: FinalizationAcquireRequest, ) -> Result<FinalizationToken>

Atomically acquires or reclaims the study-wide finalization claim. Read more
§

fn renew_finalization( &self, token: FinalizationToken, ) -> Result<FinalizationObservation>

Atomically renews a held study-wide finalization claim. Read more
§

fn release_finalization(&self, token: FinalizationToken) -> Result<()>

Atomically releases a held study-wide finalization claim. Read more
§

fn reserve_and_create_if_budget( &self, request: ReserveAndCreateRequest, ) -> Result<ReservationOutcome>

Reserves and creates a waiting trial if the durable-record budget allows it. Created consumes a trial number; Exhausted consumes nothing. Read more
§

fn abandon(&self, request: AbandonRequest) -> Result<bool>

Conditionally abandons an as-yet-unclaimed reservation. Read more
§

fn observe(&self, trial: TrialId) -> Result<LeaseObservation>

Reads the complete current lease observation for a trial. Read more
§

fn claim(&self, request: ClaimRequest) -> Result<TrialOwnerToken>

Atomically claims a reservation or paused trial, writes its initial parameters/state, and starts a new owner epoch. Read more
§

fn renew(&self, request: RenewRequest) -> Result<LeaseObservation>

Renews a running owner and returns the new observation. Read more
§

fn fail_stale(&self, request: FailStaleRequest) -> Result<bool>

Conditionally fails a running trial whose complete observed lease the caller proved unchanged for the request’s grace interval. Backends must compare the complete observation atomically, but must not compare the sweeper timestamp to the owner’s heartbeat timestamp: those clocks may belong to different nodes. Returns true for the sole race winner. The winner atomically commits the failed outcome and an Initial outbox record keyed by the observed trial and epoch, including the request’s finalization intent. A race loss commits neither. Read more
§

fn fenced_mutate(&self, request: FencedMutationRequest) -> Result<()>

Applies one mutation only if the owner/epoch token is still exact. Read more
§

fn mutate_owned( &self, owner: FinalizationToken, request: FencedMutationRequest, ) -> Result<()>

Applies one lifecycle mutation under a study-wide finalization owner. Read more
§

fn terminal(&self, request: TerminalRequest) -> Result<()>

Validates and commits a terminal outcome only if the owner/epoch token is still exact. Read more
§

fn terminal_with_intent( &self, request: TerminalRequest, intent: FinalizationIntent, ) -> Result<OutboxKey>

Commits a terminal outcome together with its durable initial outbox intent. Read more
§

fn terminal_with_intent_owned( &self, owner: FinalizationToken, request: TerminalRequest, intent: FinalizationIntent, ) -> Result<OutboxKey>

Commits a terminal result and initial outbox intent under a study-wide finalization owner. Read more
§

fn advance_outbox(&self, request: OutboxAdvanceRequest) -> Result<OutboxPhase>

Atomically persists callback state and discovered intents while moving one outbox record through its explicit phases. Read more
§

fn advance_outbox_owned( &self, owner: FinalizationToken, request: OutboxAdvanceRequest, _operation: OperationId, ) -> Result<OutboxPhase>

Advances one outbox phase under a study-wide finalization owner. Read more
§

fn get_outbox(&self, key: OutboxKey) -> Result<Option<OutboxRecord>>

Reads one durable outbox record, including its replay cursor and completed effect keys. Read more
§

fn list_outbox(&self) -> Result<Vec<OutboxRecord>>

Lists durable outbox records in key order, including fully-done records. Read more
§

fn list_outbox_page(&self, request: OutboxPageRequest) -> Result<OutboxPage>

Reads one bounded page of study-scoped durable outbox records. Read more
§

fn complete_outbox_effect(&self, request: OutboxEffectCompletion) -> Result<()>

Idempotently acknowledges an externally executed effect. Read more
§

fn create_replacement( &self, effect: OutboxEffectKey, ) -> Result<ReplacementOutcome>

Idempotently creates the replacement described by a durable effect. Read more
§

fn put_finalized_state(&self, request: FinalizedStateRequest) -> Result<()>

Replaces one state blob on a terminal lifecycle-managed trial. Read more
§

fn advance_outbox_with_operation( &self, request: OutboxAdvanceRequest, _operation: OperationId, ) -> Result<OutboxPhase>

Advances one outbox phase with the caller’s lifecycle operation correlation. Read more
§

fn complete_outbox_effect_with_operation( &self, request: OutboxEffectCompletion, _operation: OperationId, ) -> Result<()>

Acknowledges one externally executed effect with caller correlation. Read more
§

fn create_replacement_with_operation( &self, effect: OutboxEffectKey, _operation: OperationId, ) -> Result<ReplacementOutcome>

Creates one replacement with caller correlation. Read more
§

fn put_finalized_state_with_operation( &self, request: FinalizedStateRequest, _operation: OperationId, ) -> Result<()>

Replaces finalized state with caller correlation. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more