Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
ca-checkpoint.H
Go to the documentation of this file.
1/*
2 Aleph_w
3
4 Data structures & Algorithms
5 version 2.0.0b
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
84#ifndef CA_CHECKPOINT_H
85#define CA_CHECKPOINT_H
86
87#include <array>
88#include <atomic>
89#include <cerrno>
90#include <chrono>
91#include <concepts>
92#include <condition_variable>
93#include <cstddef>
94#include <cstdint>
95#include <cstdio>
96#include <cstring>
97#include <deque>
98#include <filesystem>
99#include <fstream>
100#include <functional>
101#include <iosfwd>
102#include <mutex>
103#include <optional>
104#include <random>
105#include <string>
106#include <string_view>
107#include <system_error>
108#include <thread>
109#include <type_traits>
110#include <typeinfo>
111#include <utility>
112#include <vector>
113
114#if defined(_WIN32)
115# include <io.h>
116# include <windows.h>
117#else
118# include <fcntl.h>
119# include <sys/stat.h>
120# include <sys/types.h>
121# include <unistd.h>
122#endif
123
124#include <ah-errors.H>
125
127#include <ca-rng.H>
128#include <ca-traits.H>
129#include <tpl_ca_concepts.H>
130
131namespace Aleph {
132namespace CA {
133
134namespace ca_checkpoint_detail {
135
139inline constexpr std::array<char, 8> magic_bytes
140 = {{'A', 'L', 'E', 'P', 'H', 'C', 'A', '1'}};
141
145inline constexpr std::uint32_t format_version = 2;
146
149inline constexpr std::uint32_t format_version_min_read = 1;
150
152inline constexpr std::streamoff header_bytes_v1
153 = 8 + 4 + 8 + 4 + 4 + 4 * 8 + 8 + 1 + 7 + 8 + 8; // 92
154
157inline constexpr std::streamoff header_bytes_v2 = header_bytes_v1 + 24; // 116
158
163[[nodiscard]] inline constexpr std::uint64_t fnv1a_64(std::string_view s) noexcept
164{
165 std::uint64_t h = 0xcbf29ce484222325ull;
166 for (const char c : s)
167 {
168 h ^= static_cast<std::uint64_t>(static_cast<unsigned char>(c));
169 h *= 0x100000001b3ull;
170 }
171 return h;
172}
173
178template <typename T>
179[[nodiscard]] inline std::uint64_t type_hash() noexcept
180{
181 const std::uint64_t name_hash = fnv1a_64(typeid(T).name());
182 return splitmix64(name_hash ^ static_cast<std::uint64_t>(sizeof(T)));
183}
184
185// --- Little-endian byte serialisation --------------------------------
186
187template <typename T>
188inline void write_le(std::ostream &out, const T &value)
189{
190 static_assert(std::is_trivially_copyable_v<T>,
191 "write_le requires a trivially copyable type");
192 // We always emit the host's byte representation; for the integral
193 // sizes we use it is identical to little-endian on every platform
194 // Aleph-w supports.
195 std::array<char, sizeof(T)> bytes{};
196 std::memcpy(bytes.data(), &value, sizeof(T));
197 out.write(bytes.data(), sizeof(T));
198}
199
200template <typename T>
201[[nodiscard]] inline T read_le(std::istream &in)
202{
203 static_assert(std::is_trivially_copyable_v<T>,
204 "read_le requires a trivially copyable type");
205 std::array<char, sizeof(T)> bytes{};
206 in.read(bytes.data(), sizeof(T));
207 T value{};
208 std::memcpy(&value, bytes.data(), sizeof(T));
209 return value;
210}
211
212// --- Atomic-write primitives (Phase 17) ------------------------------
213
222[[nodiscard]] inline std::filesystem::path
223make_tmp_path(const std::filesystem::path &target)
224{
225 std::mt19937_64 rng(static_cast<std::uint64_t>(
226 std::chrono::steady_clock::now().time_since_epoch().count())
227 ^ static_cast<std::uint64_t>(
228 reinterpret_cast<std::uintptr_t>(&target)));
229 const std::uint64_t r = rng();
230#if defined(_WIN32)
231 const long pid = static_cast<long>(::GetCurrentProcessId());
232#else
233 const long pid = static_cast<long>(::getpid());
234#endif
235 std::string suffix = ".tmp." + std::to_string(pid) + "." + std::to_string(r);
236 std::filesystem::path tmp = target;
237 tmp += suffix;
238 return tmp;
239}
240
251inline void sync_file(const std::filesystem::path &path)
252{
253#if defined(_WIN32)
254 HANDLE h = ::CreateFileW(path.wstring().c_str(),
257 nullptr,
260 nullptr);
262 << "sync_file: CreateFileW failed for '" << path.string() << "'";
263 const BOOL ok = ::FlushFileBuffers(h);
264 const DWORD err = ok ? 0 : ::GetLastError();
267 << "sync_file: FlushFileBuffers failed for '" << path.string()
268 << "' (error " << err << ")";
269#else
270 const int fd = ::open(path.c_str(), O_WRONLY);
272 << "sync_file: open() failed for '" << path.string()
273 << "' (errno=" << errno << ")";
274 const int rc = ::fsync(fd);
275 const int saved = errno;
276 ::close(fd);
278 << "sync_file: fsync() failed for '" << path.string()
279 << "' (errno=" << saved << ")";
280#endif
281}
282
291inline void sync_directory(const std::filesystem::path &dir) noexcept
292{
293#if defined(_WIN32)
294 (void) dir;
295#else
296 const int fd = ::open(dir.c_str(), O_RDONLY | O_DIRECTORY);
297 if (fd < 0)
298 return;
299 (void) ::fsync(fd);
300 ::close(fd);
301#endif
302}
303
310{
311 std::filesystem::path path_;
312 bool committed_ = false;
313
314public:
315 explicit Tmp_File_Guard(std::filesystem::path p) noexcept
316 : path_(std::move(p))
317 {}
318
321
323 {
324 if (not committed_)
325 {
326 std::error_code ec;
327 std::filesystem::remove(path_, ec);
328 }
329 }
330
332 void commit() noexcept { committed_ = true; }
333
335 [[nodiscard]] const std::filesystem::path &path() const noexcept { return path_; }
336};
337
338} // namespace ca_checkpoint_detail
339
340// =====================================================================
341// Public types
342// =====================================================================
343
348namespace ca_checkpoint_flags {
349
351inline constexpr std::uint32_t compressed = 0x1;
352
355inline constexpr std::uint32_t delta = 0x2;
356
357} // namespace ca_checkpoint_flags
358
369{
371 std::array<char, 8> magic{};
373 std::uint32_t format_version = 0;
375 std::uint64_t engine_type_hash = 0;
377 std::uint32_t state_type_size = 0;
379 std::uint32_t rank = 0;
381 std::array<std::uint64_t, 4> extents{};
383 std::uint64_t step_count = 0;
385 std::uint8_t has_rng = 0;
387 std::uint64_t master_seed = 0;
389 std::uint64_t cell_count = 0;
390
391 // --- v2-only trailer fields (default to 0 when reading v1) --------
392
394 std::uint32_t flags = 0;
396 std::uint32_t compression_level = 0;
400 std::uint64_t payload_size = 0;
403 std::uint64_t delta_base_step = 0;
404};
405
413{
415 bool compress = false;
417 int level = 6;
420 bool sync_dir = true;
421};
422
432template <typename R>
433concept Has_Master_Seed = requires(R &r, std::uint64_t s) {
434 { r.master_seed() } -> std::convertible_to<std::uint64_t>;
435 r.set_master_seed(s);
436};
437
446template <typename L>
448 = LatticeLike<L> and requires(const L &l) {
449 { l.extents() } -> std::convertible_to<typename L::extents_type>;
450 };
451
471template <typename E>
474 and requires { typename E::rule_type; }
475 and requires(E &e, std::size_t s) {
476 { e.frame() } -> std::convertible_to<const typename E::lattice_type &>;
477 { e.steps_run() } -> std::convertible_to<std::size_t>;
478 { e.rule() } -> std::convertible_to<typename E::rule_type &>;
479 { e.write_lattice() } -> std::convertible_to<typename E::lattice_type &>;
480 e.set_step_count(s);
481 };
482
491{
493 std::filesystem::path source_path;
496};
497
498// =====================================================================
499// inspect_checkpoint
500// =====================================================================
501
514[[nodiscard]] inline Checkpoint_Header
515inspect_checkpoint(const std::filesystem::path &path)
516{
517 using namespace ca_checkpoint_detail;
518 std::ifstream in(path, std::ios::binary);
520 << "inspect_checkpoint: cannot open '" << path.string() << "'";
521
523 in.read(h.magic.data(), h.magic.size());
524 ah_runtime_error_if(not in or h.magic != magic_bytes)
525 << "inspect_checkpoint: bad magic in '" << path.string() << "'";
526
528 ah_runtime_error_if(h.format_version < format_version_min_read
529 or h.format_version > format_version)
530 << "inspect_checkpoint: unsupported format version " << h.format_version
531 << " (supported [" << format_version_min_read << ", " << format_version
532 << "]) in '" << path.string() << "'";
533
534 h.engine_type_hash = read_le<std::uint64_t>(in);
535 h.state_type_size = read_le<std::uint32_t>(in);
537 // save_checkpoint and load_checkpoint_into currently implement
538 // rank-1, rank-2 and rank-3 serialisation branches; rank-4 is
539 // intentionally rejected here so that a corrupted or future-version
540 // header cannot bypass the save/load symmetry.
541 ah_runtime_error_if(h.rank == 0 or h.rank > 3)
542 << "inspect_checkpoint: rank " << h.rank << " out of range in '"
543 << path.string() << "'";
544 for (std::size_t d = 0; d < h.extents.size(); ++d)
545 h.extents[d] = read_le<std::uint64_t>(in);
546 h.step_count = read_le<std::uint64_t>(in);
547 h.has_rng = read_le<std::uint8_t>(in);
548 // 7 bytes of padding to keep the next field 8-aligned in the file.
549 for (int i = 0; i < 7; ++i)
551 h.master_seed = read_le<std::uint64_t>(in);
552 h.cell_count = read_le<std::uint64_t>(in);
553
554 if (h.format_version >= 2)
555 {
556 h.flags = read_le<std::uint32_t>(in);
557 h.compression_level = read_le<std::uint32_t>(in);
558 h.payload_size = read_le<std::uint64_t>(in);
559 h.delta_base_step = read_le<std::uint64_t>(in);
560 }
561 else
562 {
563 h.flags = 0;
564 h.compression_level = 0;
565 h.payload_size = h.cell_count * h.state_type_size;
566 h.delta_base_step = 0;
567 }
568
570 << "inspect_checkpoint: short read while parsing header of '"
571 << path.string() << "'";
572
573 // --- Hardening against hostile/corrupt headers (CWE-789, CWE-190) ---
574 // Every field above is attacker-controlled. Validate internal
575 // consistency here, once, so no downstream consumer can be driven into
576 // an out-of-bounds or excessively large allocation:
577 // * cell_count must equal the product of the declared extents
578 // (this is an invariant of every file save_checkpoint writes), and
579 // that product must not overflow;
580 // * cell_count * state_type_size (the raw frame size) must not
581 // overflow size_t.
582 std::uint64_t expected_cells = 1;
583 for (std::uint32_t d = 0; d < h.rank; ++d)
584 {
585 const std::uint64_t e = h.extents[d];
587 << "inspect_checkpoint: extent product overflows in '"
588 << path.string() << "'";
589 expected_cells *= e;
590 }
592 << "inspect_checkpoint: cell_count (" << h.cell_count
593 << ") inconsistent with the product of extents (" << expected_cells
594 << ") in '" << path.string() << "'";
595 ah_runtime_error_if(h.state_type_size != 0
596 and h.cell_count > UINT64_MAX / h.state_type_size)
597 << "inspect_checkpoint: cell_count * state_type_size overflows in '"
598 << path.string() << "'";
599
600 return h;
601}
602
603// =====================================================================
604// Frame serialisation helpers (shared by sync / async / delta paths)
605// =====================================================================
606
607namespace ca_checkpoint_detail {
608
620template <typename Engine>
621[[nodiscard]] std::vector<std::uint8_t> snapshot_frame_bytes(const Engine &engine)
623{
624 using Lattice = typename Engine::lattice_type;
625 using state_t = typename Lattice::state_type;
626 using coord_t = typename Lattice::coord_type;
627 static_assert(std::is_trivially_copyable_v<state_t>,
628 "snapshot_frame_bytes requires a trivially copyable state_type");
629 static_assert(Lattice::rank <= 3,
630 "snapshot_frame_bytes currently supports rank <= 3");
631
632 const Lattice &frame = engine.frame();
633 std::uint64_t cell_count = 1;
634 for (std::size_t d = 0; d < Lattice::rank; ++d)
635 cell_count *= static_cast<std::uint64_t>(frame.size(d));
636
637 std::vector<std::uint8_t> bytes(static_cast<std::size_t>(cell_count) * sizeof(state_t));
638 std::size_t off = 0;
639 auto append = [&](const state_t &v)
640 {
641 std::memcpy(bytes.data() + off, &v, sizeof(state_t));
642 off += sizeof(state_t);
643 };
644
645 if constexpr (Lattice::rank == 1)
646 {
647 for (ca_size_t i = 0; i < frame.size(0); ++i)
648 append(frame.at(coord_t{static_cast<ca_index_t>(i)}));
649 }
650 else if constexpr (Lattice::rank == 2)
651 {
652 for (ca_size_t i = 0; i < frame.size(0); ++i)
653 for (ca_size_t j = 0; j < frame.size(1); ++j)
654 append(frame.at(coord_t{static_cast<ca_index_t>(i),
655 static_cast<ca_index_t>(j)}));
656 }
657 else if constexpr (Lattice::rank == 3)
658 {
659 for (ca_size_t i = 0; i < frame.size(0); ++i)
660 for (ca_size_t j = 0; j < frame.size(1); ++j)
661 for (ca_size_t k = 0; k < frame.size(2); ++k)
662 append(frame.at(coord_t{static_cast<ca_index_t>(i),
663 static_cast<ca_index_t>(j),
664 static_cast<ca_index_t>(k)}));
665 }
666 return bytes;
667}
668
671template <typename Engine>
672void restore_frame_from_bytes(Engine &engine, const std::vector<std::uint8_t> &bytes)
674{
675 using Lattice = typename Engine::lattice_type;
676 using state_t = typename Lattice::state_type;
677 using coord_t = typename Lattice::coord_type;
678 static_assert(std::is_trivially_copyable_v<state_t>,
679 "restore_frame_from_bytes requires a trivially copyable state_type");
680 static_assert(Lattice::rank <= 3,
681 "restore_frame_from_bytes currently supports rank <= 3");
682
683 const Lattice &frame = engine.frame();
684 std::uint64_t cell_count = 1;
685 for (std::size_t d = 0; d < Lattice::rank; ++d)
686 cell_count *= static_cast<std::uint64_t>(frame.size(d));
687 ah_runtime_error_if(bytes.size() != static_cast<std::size_t>(cell_count) * sizeof(state_t))
688 << "restore_frame_from_bytes: size mismatch (expected "
689 << (cell_count * sizeof(state_t)) << ", got " << bytes.size() << ")";
690
691 std::size_t off = 0;
692 auto next = [&]() -> state_t
693 {
694 state_t v{};
695 std::memcpy(&v, bytes.data() + off, sizeof(state_t));
696 off += sizeof(state_t);
697 return v;
698 };
699
700 if constexpr (Lattice::rank == 1)
701 {
702 for (ca_size_t i = 0; i < frame.size(0); ++i)
703 engine.write_lattice().set(coord_t{static_cast<ca_index_t>(i)}, next());
704 }
705 else if constexpr (Lattice::rank == 2)
706 {
707 for (ca_size_t i = 0; i < frame.size(0); ++i)
708 for (ca_size_t j = 0; j < frame.size(1); ++j)
709 engine.write_lattice().set(
710 coord_t{static_cast<ca_index_t>(i), static_cast<ca_index_t>(j)},
711 next());
712 }
713 else if constexpr (Lattice::rank == 3)
714 {
715 for (ca_size_t i = 0; i < frame.size(0); ++i)
716 for (ca_size_t j = 0; j < frame.size(1); ++j)
717 for (ca_size_t k = 0; k < frame.size(2); ++k)
718 engine.write_lattice().set(coord_t{static_cast<ca_index_t>(i),
719 static_cast<ca_index_t>(j),
720 static_cast<ca_index_t>(k)},
721 next());
722 }
723}
724
733{
734 std::uint64_t engine_type_hash = 0;
735 std::uint32_t state_type_size = 0;
736 std::uint32_t rank = 0;
737 std::array<std::uint64_t, 4> extents{};
738 std::uint64_t step_count = 0;
739 std::uint8_t has_rng = 0;
740 std::uint64_t master_seed = 0;
741 std::uint64_t cell_count = 0;
742};
743
754template <typename Engine>
755[[nodiscard]] Header_Snapshot capture_header(const Engine &engine)
757{
758 using Lattice = typename Engine::lattice_type;
759 using state_t = typename Lattice::state_type;
761 hs.engine_type_hash = type_hash<Engine>();
762 hs.state_type_size = static_cast<std::uint32_t>(sizeof(state_t));
763 hs.rank = static_cast<std::uint32_t>(Lattice::rank);
764 const auto &ext = engine.frame().extents();
765 for (std::size_t d = 0; d < Lattice::rank; ++d)
766 hs.extents[d] = static_cast<std::uint64_t>(ext[d]);
767 hs.step_count = static_cast<std::uint64_t>(engine.steps_run());
768 std::uint64_t n = 1;
769 for (std::size_t d = 0; d < Lattice::rank; ++d)
770 n *= hs.extents[d];
771 hs.cell_count = n;
773 {
774 hs.has_rng = 1;
775 hs.master_seed = static_cast<std::uint64_t>(engine.rule().master_seed());
776 }
777 return hs;
778}
779
781inline void write_header_bytes(std::ostream &out,
782 const Header_Snapshot &hs,
783 const std::uint32_t flags,
784 const std::uint32_t compression_level,
785 const std::uint64_t payload_size,
786 const std::uint64_t delta_base_step)
787{
788 out.write(magic_bytes.data(), magic_bytes.size());
789 write_le<std::uint32_t>(out, format_version);
790 write_le<std::uint64_t>(out, hs.engine_type_hash);
791 write_le<std::uint32_t>(out, hs.state_type_size);
792 write_le<std::uint32_t>(out, hs.rank);
793 for (std::size_t d = 0; d < hs.extents.size(); ++d)
794 write_le<std::uint64_t>(out, hs.extents[d]);
795 write_le<std::uint64_t>(out, hs.step_count);
796 write_le<std::uint8_t>(out, hs.has_rng);
797 for (int i = 0; i < 7; ++i)
798 write_le<std::uint8_t>(out, std::uint8_t{0});
799 write_le<std::uint64_t>(out, hs.master_seed);
800 write_le<std::uint64_t>(out, hs.cell_count);
801 write_le<std::uint32_t>(out, flags);
802 write_le<std::uint32_t>(out, compression_level);
803 write_le<std::uint64_t>(out, payload_size);
804 write_le<std::uint64_t>(out, delta_base_step);
805}
806
814inline void atomic_write_file(const std::filesystem::path &path,
815 const Header_Snapshot &hs,
816 const std::uint32_t flags,
817 const std::uint32_t compression_level,
818 const std::vector<std::uint8_t> &payload,
819 const std::uint64_t delta_base_step,
821{
822 const std::filesystem::path tmp = make_tmp_path(path);
823 Tmp_File_Guard guard(tmp);
824 {
825 std::ofstream out(tmp, std::ios::binary | std::ios::trunc);
827 << "atomic_write_file: cannot open '" << tmp.string() << "' for write";
829 hs,
830 flags,
831 compression_level,
832 static_cast<std::uint64_t>(payload.size()),
833 delta_base_step);
834 if (not payload.empty())
835 out.write(reinterpret_cast<const char *>(payload.data()),
836 static_cast<std::streamsize>(payload.size()));
837 out.flush();
839 << "atomic_write_file: write failed for '" << tmp.string() << "'";
840 }
841 sync_file(tmp);
842
843 std::error_code ec;
844 std::filesystem::rename(tmp, path, ec);
846 << "atomic_write_file: rename '" << tmp.string() << "' -> '"
847 << path.string() << "' failed: " << ec.message();
848 guard.commit();
849
850 if (options.sync_dir)
851 {
852 std::filesystem::path parent = path.parent_path();
853 if (parent.empty())
854 parent = std::filesystem::path(".");
855 sync_directory(parent);
856 }
857}
858
861inline std::vector<std::uint8_t>
862build_payload(const std::vector<std::uint8_t> &raw,
864 std::uint32_t &flags_out,
865 std::uint32_t &level_out)
866{
867 if (not options.compress or raw.empty())
868 {
869 flags_out = 0;
870 level_out = 0;
871 return raw;
872 }
873 const int level = (options.level <= 0) ? 6 : options.level;
874 std::vector<std::uint8_t> compressed = deflate_bytes(raw.data(), raw.size(), level);
875 flags_out = ca_checkpoint_flags::compressed;
876 level_out = static_cast<std::uint32_t>(level);
877 return compressed;
878}
879
882inline std::vector<std::uint8_t>
883read_raw_payload(const std::filesystem::path &path, const Checkpoint_Header &h)
884{
885 std::ifstream in(path, std::ios::binary);
886 ah_runtime_error_if(not in)
887 << "read_raw_payload: cannot reopen '" << path.string() << "'";
888 const std::streamoff header_bytes
889 = (h.format_version >= 2) ? header_bytes_v2 : header_bytes_v1;
890
891 // Hardening (CWE-789): the on-disk payload occupies exactly the bytes
892 // after the header. A hostile header that claims a payload larger than
893 // the file would otherwise drive a multi-gigabyte allocation below.
894 const std::uint64_t file_size
895 = static_cast<std::uint64_t>(std::filesystem::file_size(path));
896 const std::uint64_t hdr = static_cast<std::uint64_t>(header_bytes);
897 ah_runtime_error_if(file_size < hdr or h.payload_size > file_size - hdr)
898 << "read_raw_payload: declared payload_size " << h.payload_size
899 << " exceeds the file body ("
900 << (file_size >= hdr ? file_size - hdr : 0)
901 << " bytes) in '" << path.string() << "'";
902
903 in.seekg(header_bytes, std::ios::beg);
904 std::vector<std::uint8_t> payload(static_cast<std::size_t>(h.payload_size));
905 if (not payload.empty())
906 {
907 in.read(reinterpret_cast<char *>(payload.data()),
908 static_cast<std::streamsize>(payload.size()));
909 ah_runtime_error_if(not in)
910 << "read_raw_payload: short read of payload in '" << path.string() << "'";
911 }
912 if (h.flags & ca_checkpoint_flags::compressed)
913 {
914 const std::size_t expected
915 = static_cast<std::size_t>(h.cell_count) * h.state_type_size;
916 return inflate_bytes(payload.data(), payload.size(), expected);
917 }
918 return payload;
919}
920
921} // namespace ca_checkpoint_detail
922
923// =====================================================================
924// save_checkpoint
925// =====================================================================
926
954template <typename Engine>
956 const std::filesystem::path &path,
957 const Checkpoint_Options &options = {})
958 requires Checkpointable_Engine<Engine>
959{
960 using namespace ca_checkpoint_detail;
961 const Header_Snapshot hs = capture_header(engine);
962 std::vector<std::uint8_t> raw = snapshot_frame_bytes(engine);
963 std::uint32_t flags = 0;
964 std::uint32_t level = 0;
965 std::vector<std::uint8_t> payload = build_payload(raw, options, flags, level);
966 atomic_write_file(path, hs, flags, level, payload, /*delta_base_step=*/0, options);
967}
968
969// =====================================================================
970// load_checkpoint_into
971// =====================================================================
972
997template <typename Engine>
999 const std::filesystem::path &path)
1001{
1002 using namespace ca_checkpoint_detail;
1003 using Lattice = typename Engine::lattice_type;
1004 using state_t = typename Lattice::state_type;
1005 static_assert(std::is_trivially_copyable_v<state_t>,
1006 "load_checkpoint_into requires a trivially copyable state_type");
1007
1008 Resume_Token token;
1009 token.source_path = path;
1010 token.header = inspect_checkpoint(path);
1011
1012 ah_runtime_error_if(token.header.engine_type_hash != type_hash<Engine>())
1013 << "load_checkpoint_into: type hash mismatch in '" << path.string()
1014 << "' (file=" << token.header.engine_type_hash
1015 << " engine=" << type_hash<Engine>() << ")";
1016 ah_runtime_error_if(token.header.state_type_size != sizeof(state_t))
1017 << "load_checkpoint_into: state size mismatch in '" << path.string()
1018 << "' (file=" << token.header.state_type_size << " engine=" << sizeof(state_t) << ")";
1019 ah_runtime_error_if(token.header.rank != Lattice::rank)
1020 << "load_checkpoint_into: rank mismatch in '" << path.string()
1021 << "' (file=" << token.header.rank << " engine=" << Lattice::rank << ")";
1022 ah_runtime_error_if(token.header.flags & ca_checkpoint_flags::delta)
1023 << "load_checkpoint_into: '" << path.string() << "' is a delta snapshot; "
1024 << "use apply_delta_checkpoint() after loading the baseline";
1025
1026 const auto &ext = engine.frame().extents();
1027 for (std::size_t d = 0; d < Lattice::rank; ++d)
1028 ah_runtime_error_if(token.header.extents[d] != static_cast<std::uint64_t>(ext[d]))
1029 << "load_checkpoint_into: extent[" << d << "] mismatch in '" << path.string()
1030 << "' (file=" << token.header.extents[d] << " engine=" << ext[d] << ")";
1031
1032 std::vector<std::uint8_t> raw = read_raw_payload(path, token.header);
1033 restore_frame_from_bytes(engine, raw);
1034
1035 // Restore step counter.
1036 engine.set_step_count(token.header.step_count);
1037
1038 // Restore master RNG seed when both sides agree on the contract.
1040 {
1041 if (token.header.has_rng != 0)
1042 engine.rule().set_master_seed(token.header.master_seed);
1043 }
1044 return token;
1045}
1046
1047// =====================================================================
1048// Delta snapshots (best-effort)
1049// =====================================================================
1050
1081template <typename Engine>
1083 const std::vector<std::uint8_t> &baseline,
1084 const std::uint64_t base_step,
1085 const std::filesystem::path &path,
1086 const Checkpoint_Options &options = {})
1087 requires Checkpointable_Engine<Engine>
1088{
1089 using namespace ca_checkpoint_detail;
1090 using Lattice = typename Engine::lattice_type;
1091 using state_t = typename Lattice::state_type;
1092 static_assert(std::is_trivially_copyable_v<state_t>,
1093 "save_delta_checkpoint requires a trivially copyable state_type");
1094
1095 const Header_Snapshot hs = capture_header(engine);
1096 const std::size_t cell_bytes = sizeof(state_t);
1097 const std::size_t total_cells = static_cast<std::size_t>(hs.cell_count);
1098 ah_runtime_error_if(baseline.size() != total_cells * cell_bytes)
1099 << "save_delta_checkpoint: baseline size mismatch (expected "
1100 << (total_cells * cell_bytes) << ", got " << baseline.size() << ")";
1101
1102 std::vector<std::uint8_t> current = snapshot_frame_bytes(engine);
1103 std::vector<std::uint8_t> diff;
1104 diff.reserve(64); // arbitrary small reserve to avoid tiny reallocs.
1105 for (std::size_t i = 0; i < total_cells; ++i)
1106 {
1107 const std::size_t off = i * cell_bytes;
1108 if (std::memcmp(current.data() + off,
1109 baseline.data() + off,
1110 cell_bytes)
1111 != 0)
1112 {
1113 const std::uint64_t idx = static_cast<std::uint64_t>(i);
1114 const std::size_t head = diff.size();
1115 diff.resize(head + sizeof(std::uint64_t) + cell_bytes);
1116 std::memcpy(diff.data() + head, &idx, sizeof(std::uint64_t));
1117 std::memcpy(diff.data() + head + sizeof(std::uint64_t),
1118 current.data() + off,
1119 cell_bytes);
1120 }
1121 }
1122
1123 std::uint32_t flags = ca_checkpoint_flags::delta;
1124 std::uint32_t level = 0;
1125 std::vector<std::uint8_t> payload;
1126 if (options.compress and not diff.empty())
1127 {
1128 const int lvl = (options.level <= 0) ? 6 : options.level;
1129 payload = deflate_bytes(diff.data(), diff.size(), lvl);
1130 flags |= ca_checkpoint_flags::compressed;
1131 level = static_cast<std::uint32_t>(lvl);
1132 }
1133 else
1134 payload = std::move(diff);
1135
1136 atomic_write_file(path, hs, flags, level, payload, base_step, options);
1137}
1138
1154template <typename Engine>
1155[[nodiscard]] Resume_Token
1156apply_delta_checkpoint(Engine &engine, const std::filesystem::path &path)
1158{
1159 using namespace ca_checkpoint_detail;
1160 using Lattice = typename Engine::lattice_type;
1161 using state_t = typename Lattice::state_type;
1162 using coord_t = typename Lattice::coord_type;
1163 static_assert(std::is_trivially_copyable_v<state_t>,
1164 "apply_delta_checkpoint requires a trivially copyable state_type");
1165 static_assert(Lattice::rank <= 3,
1166 "apply_delta_checkpoint currently supports rank <= 3");
1167
1168 Resume_Token token;
1169 token.source_path = path;
1170 token.header = inspect_checkpoint(path);
1171
1172 ah_runtime_error_if(not (token.header.flags & ca_checkpoint_flags::delta))
1173 << "apply_delta_checkpoint: '" << path.string() << "' is not a delta snapshot";
1174 ah_runtime_error_if(token.header.engine_type_hash != type_hash<Engine>())
1175 << "apply_delta_checkpoint: type hash mismatch in '" << path.string() << "'";
1176 ah_runtime_error_if(token.header.state_type_size != sizeof(state_t))
1177 << "apply_delta_checkpoint: state size mismatch in '" << path.string() << "'";
1178 ah_runtime_error_if(token.header.rank != Lattice::rank)
1179 << "apply_delta_checkpoint: rank mismatch in '" << path.string() << "'";
1180
1181 std::ifstream in(path, std::ios::binary);
1182 ah_runtime_error_if(not in)
1183 << "apply_delta_checkpoint: cannot reopen '" << path.string() << "'";
1184 in.seekg(header_bytes_v2, std::ios::beg);
1185 std::vector<std::uint8_t> payload(static_cast<std::size_t>(token.header.payload_size));
1186 if (not payload.empty())
1187 {
1188 in.read(reinterpret_cast<char *>(payload.data()),
1189 static_cast<std::streamsize>(payload.size()));
1190 ah_runtime_error_if(not in)
1191 << "apply_delta_checkpoint: short read in '" << path.string() << "'";
1192 }
1193
1194 std::vector<std::uint8_t> diff;
1195 if (token.header.flags & ca_checkpoint_flags::compressed)
1196 {
1197 // Upper bound on the uncompressed delta size: at most every cell
1198 // is mutated, producing `cell_count * (8 + state_size)` bytes.
1199 const std::size_t bound
1200 = static_cast<std::size_t>(token.header.cell_count)
1201 * (sizeof(std::uint64_t) + sizeof(state_t));
1202 diff = inflate_bytes_up_to(payload.data(), payload.size(), bound);
1203 }
1204 else
1205 diff = std::move(payload);
1206
1207 const std::size_t cell_bytes = sizeof(state_t);
1208 const std::size_t entry_bytes = sizeof(std::uint64_t) + cell_bytes;
1209 ah_runtime_error_if(diff.size() % entry_bytes != 0)
1210 << "apply_delta_checkpoint: delta payload size " << diff.size()
1211 << " is not a multiple of (" << entry_bytes << ")";
1212
1213 auto extents_from_header = [&](std::size_t d) -> ca_size_t
1214 { return static_cast<ca_size_t>(token.header.extents[d]); };
1215
1216 const std::size_t entries = diff.size() / entry_bytes;
1217 for (std::size_t e = 0; e < entries; ++e)
1218 {
1219 std::uint64_t idx = 0;
1220 state_t value{};
1221 std::memcpy(&idx, diff.data() + e * entry_bytes, sizeof(std::uint64_t));
1222 std::memcpy(&value,
1223 diff.data() + e * entry_bytes + sizeof(std::uint64_t),
1224 cell_bytes);
1226 << "apply_delta_checkpoint: linear index " << idx
1227 << " out of range (cell_count=" << token.header.cell_count << ")";
1228
1229 if constexpr (Lattice::rank == 1)
1230 {
1231 engine.write_lattice().set(
1232 coord_t{static_cast<ca_index_t>(idx)},
1233 value);
1234 }
1235 else if constexpr (Lattice::rank == 2)
1236 {
1237 const auto j = static_cast<ca_index_t>(idx % extents_from_header(1));
1238 const auto i = static_cast<ca_index_t>(idx / extents_from_header(1));
1239 engine.write_lattice().set(coord_t{i, j}, value);
1240 }
1241 else if constexpr (Lattice::rank == 3)
1242 {
1243 const auto e2 = extents_from_header(2);
1244 const auto e1 = extents_from_header(1);
1245 const auto k = static_cast<ca_index_t>(idx % e2);
1246 const auto j = static_cast<ca_index_t>((idx / e2) % e1);
1247 const auto i = static_cast<ca_index_t>(idx / (e1 * e2));
1248 engine.write_lattice().set(coord_t{i, j, k}, value);
1249 }
1250 }
1251
1252 engine.set_step_count(token.header.step_count);
1254 {
1255 if (token.header.has_rng != 0)
1256 engine.rule().set_master_seed(token.header.master_seed);
1257 }
1258 return token;
1259}
1260
1261// =====================================================================
1262// Async_Checkpoint_Writer
1263// =====================================================================
1264
1288template <typename Engine>
1290{
1291public:
1293 enum class Queue_Policy
1294 {
1295 Block,
1296 Drop_Oldest
1297 };
1298
1305 explicit Async_Checkpoint_Writer(const std::size_t capacity = 8,
1306 const Queue_Policy policy = Queue_Policy::Block)
1307 : capacity_(capacity)
1308 , policy_(policy)
1309 {
1310 ah_domain_error_if(capacity == 0)
1311 << "Async_Checkpoint_Writer: capacity must be >= 1";
1312 worker_ = std::thread([this] { this->worker_loop(); });
1313 }
1314
1317
1320 {
1321 try
1322 {
1323 flush();
1324 }
1325 catch (...)
1326 {
1327 // Swallow flush errors during destruction; the worker thread
1328 // is still joined below.
1329 }
1330 {
1331 std::unique_lock<std::mutex> lk(mu_);
1332 stop_ = true;
1333 not_empty_.notify_all();
1334 }
1335 if (worker_.joinable())
1336 worker_.join();
1337 }
1338
1345 void submit(Engine &engine,
1346 std::filesystem::path path,
1348 {
1349 Task task;
1350 task.header = ca_checkpoint_detail::capture_header(engine);
1351 task.raw = ca_checkpoint_detail::snapshot_frame_bytes(engine);
1352 task.path = std::move(path);
1353 task.options = options;
1354
1355 std::unique_lock<std::mutex> lk(mu_);
1356 if (policy_ == Queue_Policy::Drop_Oldest)
1357 {
1358 while (queue_.size() >= capacity_)
1359 {
1360 queue_.pop_front();
1361 total_dropped_.fetch_add(1, std::memory_order_relaxed);
1362 }
1363 }
1364 else
1365 {
1366 not_full_.wait(lk, [this] { return queue_.size() < capacity_ or stop_; });
1367 }
1368 queue_.push_back(std::move(task));
1369 not_empty_.notify_one();
1370 }
1371
1380 void flush()
1381 {
1382 std::unique_lock<std::mutex> lk(mu_);
1383 drained_.wait(lk, [this] { return queue_.empty() and not in_flight_; });
1384 if (worker_error_)
1385 {
1386 const std::string msg = std::move(*worker_error_);
1387 worker_error_.reset();
1388 ah_runtime_error_if(true) << msg;
1389 }
1390 }
1391
1393 [[nodiscard]] std::size_t pending() const
1394 {
1395 std::lock_guard<std::mutex> lk(mu_);
1396 return queue_.size();
1397 }
1398
1409 void set_before_write_hook(std::function<void()> hook)
1410 {
1411 before_write_hook_ = std::move(hook);
1412 }
1413
1415 [[nodiscard]] std::size_t total_written() const noexcept
1416 {
1417 return total_written_.load(std::memory_order_relaxed);
1418 }
1419
1421 [[nodiscard]] std::size_t total_dropped() const noexcept
1422 {
1423 return total_dropped_.load(std::memory_order_relaxed);
1424 }
1425
1427 [[nodiscard]] std::size_t capacity() const noexcept { return capacity_; }
1428
1430 [[nodiscard]] Queue_Policy policy() const noexcept { return policy_; }
1431
1432private:
1433 struct Task
1434 {
1436 std::vector<std::uint8_t> raw;
1437 std::filesystem::path path;
1439 };
1440
1441 void worker_loop() noexcept
1442 {
1443 for (;;)
1444 {
1445 Task task;
1446 {
1447 std::unique_lock<std::mutex> lk(mu_);
1448 not_empty_.wait(lk, [this] { return not queue_.empty() or stop_; });
1449 if (queue_.empty() and stop_)
1450 return;
1451 task = std::move(queue_.front());
1452 queue_.pop_front();
1453 in_flight_ = true;
1454 not_full_.notify_one();
1455 }
1456
1457 if (before_write_hook_)
1458 before_write_hook_();
1459
1460 bool ok = true;
1461 std::string err;
1462 try
1463 {
1464 std::uint32_t flags = 0;
1465 std::uint32_t level = 0;
1466 std::vector<std::uint8_t> payload
1467 = ca_checkpoint_detail::build_payload(task.raw, task.options, flags, level);
1468 ca_checkpoint_detail::atomic_write_file(task.path,
1469 task.header,
1470 flags,
1471 level,
1472 payload,
1473 /*delta_base_step=*/0,
1474 task.options);
1475 }
1476 catch (const std::exception &e)
1477 {
1478 ok = false;
1479 err = e.what();
1480 }
1481 catch (...)
1482 {
1483 ok = false;
1484 err = "Async_Checkpoint_Writer: unknown exception";
1485 }
1486
1487 {
1488 std::unique_lock<std::mutex> lk(mu_);
1489 in_flight_ = false;
1490 if (ok)
1491 total_written_.fetch_add(1, std::memory_order_relaxed);
1492 else if (not worker_error_)
1493 worker_error_ = std::move(err);
1494 if (queue_.empty())
1495 drained_.notify_all();
1496 }
1497 }
1498 }
1499
1500 mutable std::mutex mu_;
1501 std::condition_variable not_empty_;
1502 std::condition_variable not_full_;
1503 std::condition_variable drained_;
1504 std::deque<Task> queue_;
1505 std::size_t capacity_;
1507 bool stop_ = false;
1508 bool in_flight_ = false;
1509 std::optional<std::string> worker_error_;
1510 std::atomic<std::size_t> total_written_{0};
1511 std::atomic<std::size_t> total_dropped_{0};
1512 std::function<void()> before_write_hook_;
1513 std::thread worker_;
1514};
1515
1516// =====================================================================
1517// Periodic_Checkpoint_Observer
1518// =====================================================================
1519
1562template <typename Engine>
1564{
1565 Engine *engine_;
1566 std::size_t every_;
1567 std::string pattern_;
1568 std::size_t zero_pad_;
1569 std::filesystem::path last_path_;
1570 Async_Checkpoint_Writer<Engine> *async_writer_ = nullptr;
1572
1574 [[nodiscard]] std::filesystem::path format_path(const std::size_t step) const
1575 {
1576 constexpr std::string_view placeholder = "{step}";
1577 const std::size_t pos = pattern_.find(placeholder);
1578 std::string padded = std::to_string(step);
1579 if (padded.size() < zero_pad_)
1580 padded.insert(0, zero_pad_ - padded.size(), '0');
1581 if (pos == std::string::npos)
1582 return std::filesystem::path(pattern_ + "." + padded);
1583 return std::filesystem::path(pattern_.substr(0, pos) + padded
1584 + pattern_.substr(pos + placeholder.size()));
1585 }
1586
1587public:
1598 const std::size_t every,
1599 std::string path_pattern,
1600 const std::size_t zero_pad = 6)
1601 : engine_(&engine), every_(every), pattern_(std::move(path_pattern)), zero_pad_(zero_pad)
1602 {
1603 ah_domain_error_if(every == 0)
1604 << "Periodic_Checkpoint_Observer: every must be >= 1";
1605 }
1606
1620 const std::size_t every,
1621 std::string path_pattern,
1624 const std::size_t zero_pad = 6)
1625 : engine_(&engine)
1626 , every_(every)
1627 , pattern_(std::move(path_pattern))
1628 , zero_pad_(zero_pad)
1629 , async_writer_(&writer)
1630 , options_(options)
1631 {
1632 ah_domain_error_if(every == 0)
1633 << "Periodic_Checkpoint_Observer: every must be >= 1";
1634 }
1635
1647 template <typename Lattice>
1648 void on_step_begin(const std::size_t step, const Lattice &frame) noexcept
1649 {
1650 (void) step;
1651 (void) frame;
1652 }
1653
1663 template <typename Lattice>
1664 void on_step_end(const std::size_t step, const Lattice &)
1665 {
1666 if (step == 0 or (step % every_) != 0)
1667 return;
1668 last_path_ = format_path(step);
1669 if (async_writer_ != nullptr)
1670 async_writer_->submit(*engine_, last_path_, options_);
1671 else
1672 save_checkpoint(*engine_, last_path_, options_);
1673 }
1674
1688 [[nodiscard]] const std::filesystem::path &last_path() const noexcept
1689 {
1690 return last_path_;
1691 }
1692};
1693
1694} // namespace CA
1695} // namespace Aleph
1696
1697#endif // CA_CHECKPOINT_H
Exception handling system with formatted messages for Aleph-w.
#define ah_domain_error_if(C)
Throws std::domain_error if condition holds.
Definition ah-errors.H:527
#define ah_runtime_error_if(C)
Throws std::runtime_error if condition holds.
Definition ah-errors.H:271
long double h
Definition btreepic.C:154
size_t size_t int32_t value
Definition ca-c-api.h:116
size_t size_t int32_t * out
Definition ca-c-api.h:120
Internal compress / decompress helpers used by ca-checkpoint.H (Phase 17).
Reproducible random-number support for stochastic CA rules (Phase 8).
Common typedefs and tag types for the Cellular Automata module.
Background writer pumping snapshots onto a dedicated thread.
std::optional< std::string > worker_error_
std::size_t pending() const
Snapshot of pending queue depth (informational).
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() noexcept
Destructor: drains the queue and joins the worker.
std::condition_variable not_empty_
std::size_t capacity() const noexcept
Configured queue capacity.
std::condition_variable drained_
std::size_t total_dropped() const noexcept
Total number of pending tasks discarded due to Drop_Oldest.
Async_Checkpoint_Writer(const Async_Checkpoint_Writer &)=delete
Async_Checkpoint_Writer & operator=(const Async_Checkpoint_Writer &)=delete
std::size_t total_written() const noexcept
Total number of files successfully written by the worker.
void set_before_write_hook(std::function< void()> hook)
Install a hook called by the worker just before each write.
void submit(Engine &engine, std::filesystem::path path, Checkpoint_Options options={})
Synchronously capture the engine state and enqueue a write.
std::function< void()> before_write_hook_
Queue_Policy policy() const noexcept
Configured back-pressure policy.
Queue_Policy
Behaviour when the internal queue is full.
std::condition_variable not_full_
void flush()
Block until the queue is drained.
Lattice that adds boundary-aware access on top of a storage.
typename Storage::state_type state_type
typename Storage::coord_type coord_type
static constexpr std::size_t rank
ca_size_t size() const noexcept
state_type at(const coord_type &c) const
Strict access: throws if c is out of range.
Observer that auto-saves the engine state every every steps to a templated path.
const std::filesystem::path & last_path() const noexcept
Path of the most recent checkpoint scheduled by this observer.
std::filesystem::path format_path(const std::size_t step) const
Replace the first {step} placeholder with step.
void on_step_end(const std::size_t step, const Lattice &)
Snapshot when step % every == 0.
Periodic_Checkpoint_Observer(Engine &engine, const std::size_t every, std::string path_pattern, Async_Checkpoint_Writer< Engine > &writer, Checkpoint_Options options={}, const std::size_t zero_pad=6)
Build a periodic observer routing writes through an Async_Checkpoint_Writer.
void on_step_begin(const std::size_t step, const Lattice &frame) noexcept
Pre-step hook.
Periodic_Checkpoint_Observer(Engine &engine, const std::size_t every, std::string path_pattern, const std::size_t zero_pad=6)
Build a synchronous periodic checkpoint observer.
RAII guard that removes a temporary file unless committed.
Tmp_File_Guard(std::filesystem::path p) noexcept
Tmp_File_Guard & operator=(const Tmp_File_Guard &)=delete
const std::filesystem::path & path() const noexcept
Access the underlying path.
Tmp_File_Guard(const Tmp_File_Guard &)=delete
void commit() noexcept
Mark the temporary file as successfully renamed; suppress cleanup.
Minimal std::expected-style result type for C++20.
Shape (per-axis sizes) of an mdspan, mixing compile-time and run-time extents.
Definition ah-mdspan.H:303
Concept: engine that exposes the minimum surface needed by checkpoint save/load.
Concept: lattice exposing the standard CA layout interface.
Concept: rule type that exposes a master RNG seed.
Storage + topology that carries the cell values.
static mt19937 rng
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
std::uint64_t type_hash() noexcept
Stable type hash of T.
std::vector< std::uint8_t > deflate_bytes(const std::uint8_t *src, const std::size_t len, const int level)
Compress a raw byte buffer with DEFLATE (miniz backend).
constexpr std::array< char, 8 > magic_bytes
Magic prefix identifying an Aleph CA checkpoint stream.
Header_Snapshot capture_header(const Engine &engine)
Capture all header fields a checkpoint file needs.
std::vector< std::uint8_t > build_payload(const std::vector< std::uint8_t > &raw, const Checkpoint_Options &options, std::uint32_t &flags_out, std::uint32_t &level_out)
Build the (possibly compressed) frame payload from a flat byte buffer, honouring options....
constexpr std::streamoff header_bytes_v1
Size in bytes of the v1 header (Phase 15 layout).
constexpr std::uint64_t fnv1a_64(std::string_view s) noexcept
FNV-1a 64-bit hash.
constexpr std::streamoff header_bytes_v2
Size in bytes of the v2 header (Phase 17 layout: v1 + 4 trailer fields of 4 + 4 + 8 + 8 = 24 bytes).
void sync_file(const std::filesystem::path &path)
Force a file's contents to be durable on the underlying storage device.
constexpr std::uint32_t format_version
Current format version.
std::vector< std::uint8_t > inflate_bytes(const std::uint8_t *src, const std::size_t len, const std::size_t expected_uncompressed)
Decompress a DEFLATE byte buffer back to its raw form.
void restore_frame_from_bytes(Engine &engine, const std::vector< std::uint8_t > &bytes)
Inverse of snapshot_frame_bytes: spray the buffer back into the engine's write lattice in row-major o...
void atomic_write_file(const std::filesystem::path &path, const Header_Snapshot &hs, const std::uint32_t flags, const std::uint32_t compression_level, const std::vector< std::uint8_t > &payload, const std::uint64_t delta_base_step, const Checkpoint_Options &options)
Write header + payload atomically to path.
std::filesystem::path make_tmp_path(const std::filesystem::path &target)
Generate a temporary sibling path for atomic writes.
void write_header_bytes(std::ostream &out, const Header_Snapshot &hs, const std::uint32_t flags, const std::uint32_t compression_level, const std::uint64_t payload_size, const std::uint64_t delta_base_step)
Serialise the v2 header bytes for a checkpoint file.
constexpr std::uint32_t format_version_min_read
Oldest format version this file knows how to read.
void write_le(std::ostream &out, const T &value)
std::vector< std::uint8_t > read_raw_payload(const std::filesystem::path &path, const Checkpoint_Header &h)
Read the raw, uncompressed payload from a checkpoint file given a validated header.
void sync_directory(const std::filesystem::path &dir) noexcept
Force a directory's metadata to be durable.
std::vector< std::uint8_t > snapshot_frame_bytes(const Engine &engine)
Flatten an engine frame into a row-major byte buffer.
constexpr std::uint32_t compressed
Frame payload was written through DEFLATE (miniz).
constexpr std::uint32_t delta
File contains only the cells that changed relative to a baseline.
@ R
Recovered (and immune).
Checkpoint_Header inspect_checkpoint(const std::filesystem::path &path)
Read just the header from a checkpoint file.
std::ptrdiff_t ca_index_t
Signed coordinate component used by lattices and neighborhoods.
Definition ca-traits.H:60
Resume_Token apply_delta_checkpoint(Engine &engine, const std::filesystem::path &path)
Apply a delta checkpoint to the engine's current state.
constexpr std::uint64_t splitmix64(std::uint64_t x) noexcept
64-bit SplitMix hash.
Definition ca-rng.H:82
void save_delta_checkpoint(Engine &engine, const std::vector< std::uint8_t > &baseline, const std::uint64_t base_step, const std::filesystem::path &path, const Checkpoint_Options &options={})
Write a delta checkpoint relative to a previous payload.
std::size_t ca_size_t
Unsigned size component used for extents and counts.
Definition ca-traits.H:63
Resume_Token load_checkpoint_into(Engine &engine, const std::filesystem::path &path)
Restore an engine's state in-place from a checkpoint file.
void save_checkpoint(Engine &engine, const std::filesystem::path &path, const Checkpoint_Options &options={})
Write a complete engine snapshot to disk atomically.
Main namespace for Aleph-w library functions.
Definition ah-arena.H:89
static void suffix(Node *root, DynList< Node * > &acc)
@ Block
Scoped sequence of statements.
and
Check uniqueness with explicit hash + equality functors.
std::decay_t< typename HeadC::Item_Type > T
Definition ah-zip.H:105
bool diff(const C1 &c1, const C2 &c2, Eq e=Eq())
Check if two containers differ.
void next()
Advance all underlying iterators (bounds-checked).
Definition ah-zip.H:171
Itor::difference_type count(const Itor &beg, const Itor &end, const T &value)
Count elements equal to a value.
Definition ahAlgo.H:127
STL namespace.
static struct argp_option options[]
Definition ntreepic.C:1886
ca_checkpoint_detail::Header_Snapshot header
Public on-disk header surfaced by inspect_checkpoint.
std::uint32_t flags
Bitmask of ca_checkpoint_flags::* values.
std::uint64_t cell_count
Number of cells stored in the frame payload (uncompressed view).
std::uint32_t format_version
Format version (1 = Phase 15, 2 = Phase 17).
std::uint64_t step_count
Number of completed steps at save time.
std::uint64_t delta_base_step
Step number of the baseline snapshot this delta references (meaningful only when flags & delta).
std::uint64_t master_seed
Recorded master RNG seed (only valid when has_rng == 1).
std::uint64_t engine_type_hash
Stable type hash of the engine type that produced this snapshot.
std::uint32_t rank
Lattice rank (1, 2 or 3 in the bundled lattices).
std::array< char, 8 > magic
Stream magic ("ALEPHCA1").
std::uint32_t compression_level
DEFLATE compression level used when flags & compressed.
std::uint32_t state_type_size
sizeof(state_type) recorded at save time, for sanity checks.
std::uint64_t payload_size
Actual payload byte count on disk (compressed size when flags & compressed, or (cell_count * state_ty...
std::array< std::uint64_t, 4 > extents
Per-axis extents (only the first rank entries are meaningful).
std::uint8_t has_rng
1 when the snapshot carries a master RNG seed.
User-tunable knobs for save_checkpoint.
bool compress
Enable DEFLATE compression of the frame payload.
bool sync_dir
Call fsync on the parent directory after the atomic rename.
int level
DEFLATE level in [1, 9]. 0 falls back to 6.
Resume handle returned by load_checkpoint_into.
std::filesystem::path source_path
Path of the checkpoint file consumed by load.
Checkpoint_Header header
Validated header read from source_path.
Header bundle pre-computed from an engine.
Task with priority for job scheduling.
static int * k
static mt19937 engine
gsl_rng * r
C++20 concepts for the Cellular Automata module.
DynList< int > l