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
MemStoragekeeps, 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_millisand 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.jsonlis the whole result set, on a cluster login node, with noatunebinary 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
| Path | What |
|---|---|
<path> | the log — the only authoritative file |
<path>.lock | the renewable append-lease directory |
<path>.lock.append-<generation> | an ephemeral generation-bound publish ticket |
<path>.snapshot | an 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
impl JournalStorage
pub fn open(path: impl AsRef<Path>) -> Result<Self>
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>
pub fn with_options( path: impl AsRef<Path>, options: JournalOptions, ) -> Result<Self>
pub fn open_read_only(path: impl AsRef<Path>) -> Result<Self>
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>
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>
pub fn open_read_only_with_options( path: impl AsRef<Path>, options: JournalOptions, ) -> Result<Self>
pub fn snapshot_path(&self) -> &Path
pub fn snapshot_path(&self) -> &Path
The snapshot side-file’s path.
pub fn ops(&self) -> Result<u64>
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<()>
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>>
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.
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>
fn import_batch(&self, batch: &ImportBatch) -> Result<ImportBatchResult>
§impl Debug for JournalStorage
impl Debug for JournalStorage
§impl Storage for JournalStorage
impl Storage for JournalStorage
§fn backing_paths(&self) -> Option<Vec<PathBuf>>
fn backing_paths(&self) -> Option<Vec<PathBuf>>
§fn create_study(&self, cfg: StudyConfig) -> Result<StudyId>
fn create_study(&self, cfg: StudyConfig) -> Result<StudyId>
§fn study_config(&self, study: StudyId) -> Result<StudyConfig>
fn study_config(&self, study: StudyId) -> Result<StudyConfig>
§fn study_user_attrs(&self, study: StudyId) -> Result<UserAttrs>
fn study_user_attrs(&self, study: StudyId) -> Result<UserAttrs>
§fn set_study_user_attr(
&self,
study: StudyId,
key: &str,
value: Value,
) -> Result<()>
fn set_study_user_attr( &self, study: StudyId, key: &str, value: Value, ) -> Result<()>
§fn create_trial(
&self,
study: StudyId,
template: Option<TrialTemplate>,
) -> Result<TrialRecord>
fn create_trial( &self, study: StudyId, template: Option<TrialTemplate>, ) -> Result<TrialRecord>
§fn try_transition(
&self,
trial: TrialId,
from: TrialState,
to: TrialState,
) -> Result<bool>
fn try_transition( &self, trial: TrialId, from: TrialState, to: TrialState, ) -> Result<bool>
§fn write_params(
&self,
trial: TrialId,
batch: &[(String, Distribution, ParamValue)],
) -> Result<()>
fn write_params( &self, trial: TrialId, batch: &[(String, Distribution, ParamValue)], ) -> Result<()>
§fn report(&self, trial: TrialId, step: u64, values: &[f64]) -> Result<()>
fn report(&self, trial: TrialId, step: u64, values: &[f64]) -> Result<()>
§fn set_values(&self, trial: TrialId, values: &[f64]) -> Result<()>
fn set_values(&self, trial: TrialId, values: &[f64]) -> Result<()>
§fn set_error(&self, trial: TrialId, message: &str) -> Result<()>
fn set_error(&self, trial: TrialId, message: &str) -> Result<()>
§fn heartbeat(&self, trial: TrialId, now_millis: u64) -> Result<()>
fn heartbeat(&self, trial: TrialId, now_millis: u64) -> Result<()>
now_millis. Read more§fn put_state(&self, scope: Scope, blob: StateBlob) -> Result<()>
fn put_state(&self, scope: Scope, blob: StateBlob) -> Result<()>
scope, replacing any previous one. Read more§fn trial_lifecycle(&self) -> Option<&dyn TrialLifecycleStorage>
fn trial_lifecycle(&self) -> Option<&dyn TrialLifecycleStorage>
§fn atomic_import(&self) -> Option<&dyn AtomicImportStorage>
fn atomic_import(&self) -> Option<&dyn AtomicImportStorage>
§fn study_catalog(&self) -> Option<&dyn StudyCatalog>
fn study_catalog(&self) -> Option<&dyn StudyCatalog>
§fn get_state(&self, scope: Scope) -> Result<Option<StateBlob>>
fn get_state(&self, scope: Scope) -> Result<Option<StateBlob>>
scope, if one was ever written. Read more§fn read_snapshot(&self, request: &SnapshotRequest) -> Result<StorageSnapshot>
fn read_snapshot(&self, request: &SnapshotRequest) -> Result<StorageSnapshot>
§impl StudyCatalog for JournalStorage
impl StudyCatalog for JournalStorage
§fn create_if_empty(&self, cfg: StudyConfig) -> Result<StudyCreateOutcome>
fn create_if_empty(&self, cfg: StudyConfig) -> Result<StudyCreateOutcome>
cfg only when the catalog has no studies. Read more§fn list_studies(&self) -> Result<Vec<StudySummary>>
fn list_studies(&self) -> Result<Vec<StudySummary>>
§fn find_by_name(&self, name: &str) -> Result<StudyLoadOutcome>
fn find_by_name(&self, name: &str) -> Result<StudyLoadOutcome>
§fn list(&self) -> Result<Vec<StudySummary>>
fn list(&self) -> Result<Vec<StudySummary>>
list_studies. Read more§impl TrialLifecycleStorage for JournalStorage
impl TrialLifecycleStorage for JournalStorage
§fn observe_finalization(
&self,
study: StudyId,
) -> Result<FinalizationObservation>
fn observe_finalization( &self, study: StudyId, ) -> Result<FinalizationObservation>
§fn acquire_finalization(
&self,
request: FinalizationAcquireRequest,
) -> Result<FinalizationToken>
fn acquire_finalization( &self, request: FinalizationAcquireRequest, ) -> Result<FinalizationToken>
§fn renew_finalization(
&self,
token: FinalizationToken,
) -> Result<FinalizationObservation>
fn renew_finalization( &self, token: FinalizationToken, ) -> Result<FinalizationObservation>
§fn release_finalization(&self, token: FinalizationToken) -> Result<()>
fn release_finalization(&self, token: FinalizationToken) -> Result<()>
§fn reserve_and_create_if_budget(
&self,
request: ReserveAndCreateRequest,
) -> Result<ReservationOutcome>
fn reserve_and_create_if_budget( &self, request: ReserveAndCreateRequest, ) -> Result<ReservationOutcome>
Created consumes a trial number; Exhausted consumes nothing. Read more