Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
Aleph::CA::Async_Checkpoint_Writer< Engine > Class Template Reference

Background writer pumping snapshots onto a dedicated thread. More...

#include <ca-checkpoint.H>

Collaboration diagram for Aleph::CA::Async_Checkpoint_Writer< Engine >:
[legend]

Classes

struct  Task
 

Public Types

enum class  Queue_Policy { Block , Drop_Oldest }
 Behaviour when the internal queue is full. More...
 

Public Member Functions

 Async_Checkpoint_Writer (const std::size_t capacity=8, const Queue_Policy policy=Queue_Policy::Block)
 Construct the writer and spawn its worker thread.
 
 Async_Checkpoint_Writer (const Async_Checkpoint_Writer &)=delete
 
Async_Checkpoint_Writer & operator= (const Async_Checkpoint_Writer &)=delete
 
 ~Async_Checkpoint_Writer () noexcept
 Destructor: drains the queue and joins the worker.
 
void submit (Engine &engine, std::filesystem::path path, Checkpoint_Options options={})
 Synchronously capture the engine state and enqueue a write.
 
void flush ()
 Block until the queue is drained.
 
std::size_t pending () const
 Snapshot of pending queue depth (informational).
 
void set_before_write_hook (std::function< void()> hook)
 Install a hook called by the worker just before each write.
 
std::size_t total_written () const noexcept
 Total number of files successfully written by the worker.
 
std::size_t total_dropped () const noexcept
 Total number of pending tasks discarded due to Drop_Oldest.
 
std::size_t capacity () const noexcept
 Configured queue capacity.
 
Queue_Policy policy () const noexcept
 Configured back-pressure policy.
 

Private Member Functions

void worker_loop () noexcept
 

Private Attributes

std::mutex mu_
 
std::condition_variable not_empty_
 
std::condition_variable not_full_
 
std::condition_variable drained_
 
std::deque< Task > queue_
 
std::size_t capacity_
 
Queue_Policy policy_
 
bool stop_ = false
 
bool in_flight_ = false
 
std::optional< std::string > worker_error_
 
std::atomic< std::size_t > total_written_ {0}
 
std::atomic< std::size_t > total_dropped_ {0}
 
std::function< void()> before_write_hook_
 
std::thread worker_
 

Detailed Description

template<typename Engine>
class Aleph::CA::Async_Checkpoint_Writer< Engine >

Background writer pumping snapshots onto a dedicated thread.

Decouples engine.step() from disk I/O: submit() synchronously copies the current frame into an internal task buffer (a fast memcpy) and returns immediately; a worker thread later runs optional DEFLATE compression and the atomic-rename write.

Queue policy:

  • Block — submit() blocks until queue depth < capacity.
  • Drop_Oldest — when the queue is full, the oldest pending task is discarded and the new one enqueued.

The destructor blocks on flush() so no snapshot is silently lost when the writer goes out of scope.

Thread safety: submit() and flush() are safe to call from a single producer thread (typically the engine thread). The writer is not designed for multi-producer use — wrap it in your own synchronisation if multiple engines share the same writer.

Template Parameters
Engineengine type satisfying Checkpointable_Engine.

Definition at line 1289 of file ca-checkpoint.H.

Member Enumeration Documentation

◆ Queue_Policy

Behaviour when the internal queue is full.

Enumerator
Block 
Drop_Oldest 

Definition at line 1293 of file ca-checkpoint.H.

Constructor & Destructor Documentation

◆ Async_Checkpoint_Writer() [1/2]

template<typename Engine >
Aleph::CA::Async_Checkpoint_Writer< Engine >::Async_Checkpoint_Writer ( const std::size_t  capacity = 8,
const Queue_Policy  policy = Queue_Policy::Block 
)
inlineexplicit

Construct the writer and spawn its worker thread.

Parameters
[in]capacitymaximum number of pending tasks. Must be >= 1.
[in]policywhat to do when capacity is reached.
Exceptions
std::domain_errorwhen capacity == 0.

Definition at line 1305 of file ca-checkpoint.H.

References ah_domain_error_if.

◆ Async_Checkpoint_Writer() [2/2]

template<typename Engine >
Aleph::CA::Async_Checkpoint_Writer< Engine >::Async_Checkpoint_Writer ( const Async_Checkpoint_Writer< Engine > &  )
delete

◆ ~Async_Checkpoint_Writer()

Destructor: drains the queue and joins the worker.

Definition at line 1319 of file ca-checkpoint.H.

Member Function Documentation

◆ capacity()

template<typename Engine >
std::size_t Aleph::CA::Async_Checkpoint_Writer< Engine >::capacity ( ) const
inlinenoexcept

Configured queue capacity.

Definition at line 1427 of file ca-checkpoint.H.

◆ flush()

Block until the queue is drained.

Returns once every previously-submitted task has been written (or dropped). Subsequent submit() calls are still accepted.

Exceptions
std::runtime_errorrethrown from the worker if it captured a write failure.

Definition at line 1380 of file ca-checkpoint.H.

References ah_runtime_error_if.

◆ operator=()

◆ pending()

template<typename Engine >
std::size_t Aleph::CA::Async_Checkpoint_Writer< Engine >::pending ( ) const
inline

Snapshot of pending queue depth (informational).

Definition at line 1393 of file ca-checkpoint.H.

◆ policy()

template<typename Engine >
Queue_Policy Aleph::CA::Async_Checkpoint_Writer< Engine >::policy ( ) const
inlinenoexcept

Configured back-pressure policy.

Definition at line 1430 of file ca-checkpoint.H.

◆ set_before_write_hook()

template<typename Engine >
void Aleph::CA::Async_Checkpoint_Writer< Engine >::set_before_write_hook ( std::function< void()>  hook)
inline

Install a hook called by the worker just before each write.

Intended for unit tests that need to throttle the worker to make back-pressure behaviour deterministic. The hook runs on the worker thread, outside the mutex. Passing an empty std::function clears any previously installed hook.

Note
Not thread-safe with respect to concurrent submit() calls; call only before the first submit() or after flush().

Definition at line 1409 of file ca-checkpoint.H.

◆ submit()

template<typename Engine >
void Aleph::CA::Async_Checkpoint_Writer< Engine >::submit ( Engine &  engine,
std::filesystem::path  path,
Checkpoint_Options  options = {} 
)
inline

Synchronously capture the engine state and enqueue a write.

Parameters
[in]engineengine to snapshot.
[in]pathdestination path.
[in]optionscompression / dir-fsync knobs.

Definition at line 1345 of file ca-checkpoint.H.

Referenced by Aleph::CA::Periodic_Checkpoint_Observer< Engine >::on_step_end().

◆ total_dropped()

template<typename Engine >
std::size_t Aleph::CA::Async_Checkpoint_Writer< Engine >::total_dropped ( ) const
inlinenoexcept

Total number of pending tasks discarded due to Drop_Oldest.

Definition at line 1421 of file ca-checkpoint.H.

◆ total_written()

template<typename Engine >
std::size_t Aleph::CA::Async_Checkpoint_Writer< Engine >::total_written ( ) const
inlinenoexcept

Total number of files successfully written by the worker.

Definition at line 1415 of file ca-checkpoint.H.

◆ worker_loop()

Member Data Documentation

◆ before_write_hook_

template<typename Engine >
std::function<void()> Aleph::CA::Async_Checkpoint_Writer< Engine >::before_write_hook_
private

Definition at line 1512 of file ca-checkpoint.H.

◆ capacity_

template<typename Engine >
std::size_t Aleph::CA::Async_Checkpoint_Writer< Engine >::capacity_
private

Definition at line 1505 of file ca-checkpoint.H.

◆ drained_

template<typename Engine >
std::condition_variable Aleph::CA::Async_Checkpoint_Writer< Engine >::drained_
private

Definition at line 1503 of file ca-checkpoint.H.

◆ in_flight_

template<typename Engine >
bool Aleph::CA::Async_Checkpoint_Writer< Engine >::in_flight_ = false
private

Definition at line 1508 of file ca-checkpoint.H.

◆ mu_

template<typename Engine >
std::mutex Aleph::CA::Async_Checkpoint_Writer< Engine >::mu_
mutableprivate

Definition at line 1500 of file ca-checkpoint.H.

◆ not_empty_

template<typename Engine >
std::condition_variable Aleph::CA::Async_Checkpoint_Writer< Engine >::not_empty_
private

Definition at line 1501 of file ca-checkpoint.H.

◆ not_full_

template<typename Engine >
std::condition_variable Aleph::CA::Async_Checkpoint_Writer< Engine >::not_full_
private

Definition at line 1502 of file ca-checkpoint.H.

◆ policy_

Definition at line 1506 of file ca-checkpoint.H.

◆ queue_

template<typename Engine >
std::deque<Task> Aleph::CA::Async_Checkpoint_Writer< Engine >::queue_
private

Definition at line 1504 of file ca-checkpoint.H.

◆ stop_

Definition at line 1507 of file ca-checkpoint.H.

◆ total_dropped_

template<typename Engine >
std::atomic<std::size_t> Aleph::CA::Async_Checkpoint_Writer< Engine >::total_dropped_ {0}
private

Definition at line 1511 of file ca-checkpoint.H.

◆ total_written_

template<typename Engine >
std::atomic<std::size_t> Aleph::CA::Async_Checkpoint_Writer< Engine >::total_written_ {0}
private

Definition at line 1510 of file ca-checkpoint.H.

◆ worker_

template<typename Engine >
std::thread Aleph::CA::Async_Checkpoint_Writer< Engine >::worker_
private

Definition at line 1513 of file ca-checkpoint.H.

◆ worker_error_

template<typename Engine >
std::optional<std::string> Aleph::CA::Async_Checkpoint_Writer< Engine >::worker_error_
private

Definition at line 1509 of file ca-checkpoint.H.


The documentation for this class was generated from the following file: