|
Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
|
Background writer pumping snapshots onto a dedicated thread. More...
#include <ca-checkpoint.H>
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_ |
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.
| Engine | engine type satisfying Checkpointable_Engine. |
Definition at line 1289 of file ca-checkpoint.H.
Behaviour when the internal queue is full.
| Enumerator | |
|---|---|
| Block | |
| Drop_Oldest | |
Definition at line 1293 of file ca-checkpoint.H.
|
inlineexplicit |
Construct the writer and spawn its worker thread.
| [in] | capacity | maximum number of pending tasks. Must be >= 1. |
| [in] | policy | what to do when capacity is reached. |
| std::domain_error | when capacity == 0. |
Definition at line 1305 of file ca-checkpoint.H.
References ah_domain_error_if.
|
delete |
|
inlinenoexcept |
Destructor: drains the queue and joins the worker.
Definition at line 1319 of file ca-checkpoint.H.
|
inlinenoexcept |
Configured queue capacity.
Definition at line 1427 of file ca-checkpoint.H.
|
inline |
Block until the queue is drained.
Returns once every previously-submitted task has been written (or dropped). Subsequent submit() calls are still accepted.
| std::runtime_error | rethrown from the worker if it captured a write failure. |
Definition at line 1380 of file ca-checkpoint.H.
References ah_runtime_error_if.
|
delete |
|
inline |
Snapshot of pending queue depth (informational).
Definition at line 1393 of file ca-checkpoint.H.
|
inlinenoexcept |
Configured back-pressure policy.
Definition at line 1430 of file ca-checkpoint.H.
|
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.
submit() calls; call only before the first submit() or after flush(). Definition at line 1409 of file ca-checkpoint.H.
|
inline |
Synchronously capture the engine state and enqueue a write.
| [in] | engine | engine to snapshot. |
| [in] | path | destination path. |
| [in] | options | compression / dir-fsync knobs. |
Definition at line 1345 of file ca-checkpoint.H.
Referenced by Aleph::CA::Periodic_Checkpoint_Observer< Engine >::on_step_end().
|
inlinenoexcept |
Total number of pending tasks discarded due to Drop_Oldest.
Definition at line 1421 of file ca-checkpoint.H.
|
inlinenoexcept |
Total number of files successfully written by the worker.
Definition at line 1415 of file ca-checkpoint.H.
|
inlineprivatenoexcept |
|
private |
Definition at line 1512 of file ca-checkpoint.H.
|
private |
Definition at line 1505 of file ca-checkpoint.H.
|
private |
Definition at line 1503 of file ca-checkpoint.H.
|
private |
Definition at line 1508 of file ca-checkpoint.H.
|
mutableprivate |
Definition at line 1500 of file ca-checkpoint.H.
|
private |
Definition at line 1501 of file ca-checkpoint.H.
|
private |
Definition at line 1502 of file ca-checkpoint.H.
|
private |
Definition at line 1506 of file ca-checkpoint.H.
|
private |
Definition at line 1504 of file ca-checkpoint.H.
Definition at line 1507 of file ca-checkpoint.H.
|
private |
Definition at line 1511 of file ca-checkpoint.H.
|
private |
Definition at line 1510 of file ca-checkpoint.H.
|
private |
Definition at line 1513 of file ca-checkpoint.H.
|
private |
Definition at line 1509 of file ca-checkpoint.H.