From 97bc20e2870a0ba938b07deb6e83bb928527b36f Mon Sep 17 00:00:00 2001 From: Daniel Kowalski Date: Thu, 9 Jul 2026 10:23:42 +0200 Subject: [PATCH 1/5] feat: add LatestFrameWriter/LatestFrameReader (producer-owned latest-frame, seqlock) --- .gitignore | 8 +- cpp/CMakeLists.txt | 1 + cpp/include/zerobuffer/latest_frame.h | 200 +++++++++++ cpp/src/latest_frame.cpp | 482 ++++++++++++++++++++++++++ cpp/tests/CMakeLists.txt | 59 +++- cpp/tests/test_latest_frame.cpp | 367 ++++++++++++++++++++ 6 files changed, 1108 insertions(+), 9 deletions(-) create mode 100644 cpp/include/zerobuffer/latest_frame.h create mode 100644 cpp/src/latest_frame.cpp create mode 100644 cpp/tests/test_latest_frame.cpp diff --git a/.gitignore b/.gitignore index 44cfa65..160e85b 100644 --- a/.gitignore +++ b/.gitignore @@ -42,6 +42,10 @@ Testing/ /cpp/tests/*_test /cpp/benchmarks/test_* /cpp/benchmarks/*_test +# ...but keep test/benchmark SOURCES tracked (the patterns above target built binaries) +!/cpp/tests/test_*.cpp +!/cpp/tests/test_*.h +!/cpp/benchmarks/test_*.cpp # C# build artifacts /csharp/.vs/ @@ -126,5 +130,5 @@ vcpkg_installed/ /python/ZeroBuffer.Python.Integration.Tests/bin/Release/net9.0 /modules/harmony/ModelingEvolution.Harmony.Shared.Tests/bin/Release/net9.0 /tools/harmony-cpp-gen/bin/Release/net9.0 -/csharp/ZeroBuffer/nupkg -/logs +/csharp/ZeroBuffer/nupkg +/logs diff --git a/cpp/CMakeLists.txt b/cpp/CMakeLists.txt index c6b1e54..eda28eb 100644 --- a/cpp/CMakeLists.txt +++ b/cpp/CMakeLists.txt @@ -46,6 +46,7 @@ endif() add_library(zerobuffer STATIC src/reader.cpp src/writer.cpp + src/latest_frame.cpp src/duplex_client.cpp src/immutable_duplex_server.cpp src/duplex_channel_factory.cpp diff --git a/cpp/include/zerobuffer/latest_frame.h b/cpp/include/zerobuffer/latest_frame.h new file mode 100644 index 0000000..2a8907f --- /dev/null +++ b/cpp/include/zerobuffer/latest_frame.h @@ -0,0 +1,200 @@ +#ifndef ZEROBUFFER_LATEST_FRAME_H +#define ZEROBUFFER_LATEST_FRAME_H + +// Producer-owned, consumer-optional, newest-frame-wins display buffer +// (Epic 020 / feature-002, ADR-1 + ADR-11). Purely additive: built on +// zerobuffer's platform layer (SharedMemory / process_exists) and does NOT +// touch the SPSC Reader/Writer, so the AI path's SPSC contract is unchanged. +// +// The writer creates and owns the segment and never waits on a reader; N +// readers map it read-only and read the newest published slot tear-free via a +// per-slot seqlock over a triple- (or wider) slot ring. Caps travel in the +// header as a length-prefixed JSON block consumed verbatim by native-player's +// ShmCaps::parse (4-byte little-endian length prefix + JSON). + +#include "zerobuffer/platform.h" +#include "zerobuffer/reader.h" // ZeroBufferException + +#include +#include +#include +#include +#include +#include + +namespace zerobuffer { + +// ----- on-segment layout (one definition shared by writer and reader) ------- + +constexpr uint32_t LATEST_FRAME_MAGIC = 0x3146524Cu; // 'LFR1' +constexpr uint32_t LATEST_FRAME_VERSION = 1u; +constexpr uint32_t LATEST_FRAME_MIN_SLOTS = 3u; // triple-slot minimum +constexpr size_t LATEST_FRAME_ALIGNMENT = 64u; // cache-line alignment +constexpr size_t LATEST_FRAME_METADATA_CAPACITY = 8192u; // caps JSON block bytes + +// Per-slot header. `seqlock` is the tear-free guard: even = stable, odd = write +// in progress. `frame_number` and `size` are plain fields protected by it. +struct LatestFrameSlotHeader { + uint64_t seqlock; + uint64_t frame_number; + uint64_t size; + uint64_t reserved; +}; +static_assert(sizeof(LatestFrameSlotHeader) == 32, "slot header must be 32 bytes"); + +// Segment header. POD; accessed cross-process. Concurrency-sensitive fields +// (seqlock, publish_index, heartbeat_ns, metadata_seq) are read/written through +// std::atomic_ref so the underlying bytes stay a plain cross-language layout. +struct LatestFrameSharedHeader { + uint32_t magic; + uint32_t version; + uint32_t slot_count; + uint32_t slot_header_size; // sizeof(LatestFrameSlotHeader) + uint64_t slot_size; // payload bytes per slot + uint64_t slot_stride; // aligned(slot_header_size + slot_size) + uint64_t slots_offset; // byte offset of slot 0 from segment start + uint64_t metadata_offset; // byte offset of the caps block + uint64_t metadata_capacity; // reserved bytes for [u32 len][json] + uint32_t metadata_seq; // seqlock over the caps block (even = stable) + uint32_t reserved0; + int64_t publish_index; // newest published slot; -1 = none published + uint64_t writer_pid; // owning writer pid (0 = none) + uint64_t writer_start_time; // writer process start time (pid-reuse guard) + uint64_t heartbeat_ns; // steady_clock tick stamped by the writer + uint64_t reserved1[4]; +}; +static_assert(sizeof(LatestFrameSharedHeader) == 128, "shared header must be 128 bytes"); +static_assert(sizeof(LatestFrameSharedHeader) % LATEST_FRAME_ALIGNMENT == 0, + "shared header must be cache-line aligned"); + +// ----- frame view returned by the reader ------------------------------------ + +// Zero-allocation view of one tear-free frame copy. `data()` points into the +// reader's reusable scratch buffer and is valid until the next read_latest(). +class LatestFrame { +public: + LatestFrame() = default; + + const void* data() const { return _data; } + size_t size() const { return _size; } + uint64_t sequence() const { return _sequence; } + bool valid() const { return _valid; } + +private: + friend class LatestFrameReader; + const void* _data = nullptr; + size_t _size = 0; + uint64_t _sequence = 0; + bool _valid = false; +}; + +// ----- writer (producer) ---------------------------------------------------- + +// Creates and OWNS the segment. Constructs successfully with no reader present +// and never blocks on or waits for a reader. Reclaims a stale segment left by a +// dead writer on create. Unlinks the segment on destruction. +class LatestFrameWriter { +public: + // slot_count defaults to the triple-slot minimum; values below it are + // raised to LATEST_FRAME_MIN_SLOTS. slot_size is the max payload per frame. + LatestFrameWriter(std::string name, size_t slot_size, + uint32_t slot_count = LATEST_FRAME_MIN_SLOTS); + ~LatestFrameWriter(); + + LatestFrameWriter(const LatestFrameWriter&) = delete; + LatestFrameWriter& operator=(const LatestFrameWriter&) = delete; + LatestFrameWriter(LatestFrameWriter&&) noexcept; + LatestFrameWriter& operator=(LatestFrameWriter&&) noexcept; + + // Write the length-prefixed caps block (4-byte LE length + JSON) into the + // header. Rewritable in place on a caps change (FR-3). + void set_metadata(const void* caps, size_t len); + + // Reserve the next writable slot (never the just-published slot) and mark it + // as being written. Returns a pointer to its payload region; the caller + // writes up to slot_size() bytes then calls publish(). + uint8_t* get_slot(); + + // Finalize the slot reserved by get_slot(): record (sequence, size), make it + // the newest published slot (seqlock release), and stamp the heartbeat. + // `size` above slot_size() is clamped. Never blocks on a reader. + void publish(uint64_t sequence, size_t size); + + // Stamp writer liveness without publishing a frame. + void heartbeat(); + + size_t slot_size() const { return _slot_size; } + uint32_t slot_count() const { return _slot_count; } + const std::string& name() const { return _name; } + +private: + void reclaim_stale_or_throw(); + void release(); + LatestFrameSlotHeader* slot_hdr(uint32_t i); + uint8_t* slot_payload(uint32_t i); + + std::string _name; + size_t _slot_size = 0; + uint32_t _slot_count = 0; + std::unique_ptr _shm; + LatestFrameSharedHeader* _header = nullptr; + uint8_t* _base = nullptr; + uint32_t _write_index = 0; + int _pending_slot = -1; // slot reserved by get_slot(), awaiting publish + uint64_t _pending_seq_even = 0; // seqlock value publish() must store +}; + +// ----- reader (consumer) ---------------------------------------------------- + +// Maps the segment read-only. Tolerates the segment being absent (the caller +// polls / retries). Never blocks or signals the writer. Supports 1..N readers. +class LatestFrameReader { +public: + explicit LatestFrameReader(std::string name); + ~LatestFrameReader(); + + LatestFrameReader(const LatestFrameReader&) = delete; + LatestFrameReader& operator=(const LatestFrameReader&) = delete; + LatestFrameReader(LatestFrameReader&&) noexcept; + LatestFrameReader& operator=(LatestFrameReader&&) noexcept; + + // Read the newest published frame tear-free (seqlock). Returns an invalid + // frame if the segment is absent or no NEW frame appears within `timeout`. + // Intermediate frames the writer overwrote are drops (sequence gap). + LatestFrame read_latest(std::chrono::milliseconds timeout); + + // Caps block as [4-byte LE len][JSON], consumed verbatim by ShmCaps::parse. + const void* get_metadata_raw(); + size_t get_metadata_size(); + + // False = producer gone (dead pid / reused pid), for RECONNECTING (FR-9). + bool is_writer_alive(); + + bool is_attached() const { return _header != nullptr; } + const std::string& name() const { return _name; } + +private: + bool try_attach(); + void detach(); + bool try_read_once(LatestFrame& out); + void refresh_metadata(); + bool writer_gone(); + const LatestFrameSlotHeader* slot_hdr(uint32_t i) const; + const uint8_t* slot_payload(uint32_t i) const; + + std::string _name; + std::unique_ptr _shm; + LatestFrameSharedHeader* _header = nullptr; + uint8_t* _base = nullptr; + std::vector _scratch; // reused frame copy (no per-frame alloc) + std::vector _meta_cache; // cached [u32 len][json] + size_t _meta_size = 0; + uint32_t _meta_seq_seen = 0; + bool _meta_valid = false; + uint64_t _last_sequence = 0; + bool _have_last = false; +}; + +} // namespace zerobuffer + +#endif // ZEROBUFFER_LATEST_FRAME_H diff --git a/cpp/src/latest_frame.cpp b/cpp/src/latest_frame.cpp new file mode 100644 index 0000000..977a9c4 --- /dev/null +++ b/cpp/src/latest_frame.cpp @@ -0,0 +1,482 @@ +#include "zerobuffer/latest_frame.h" + +#include +#include +#include +#include +#include + +// Producer-owned latest-frame primitive (ADR-11). The seqlock lives in each +// slot's `seqlock` field: the writer makes it odd before touching the payload +// and even after, so a reader that sees an even value both before and after its +// copy is guaranteed a tear-free frame. The triple- (or wider) slot ring means +// the slot a reader is copying is not the writer's next target, so a lap during +// a read is rare; the reader retries and lands on the newest slot. +// +// Cross-process atomicity uses std::atomic_ref over plain POD fields so the +// on-segment bytes stay a language-neutral layout. Ordering: +// - writer: payload store, then seqlock even (release), then publish_index +// (release) -> a reader that acquires publish_index sees the payload. +// - reader: publish_index (acquire), seqlock (acquire), copy, acquire fence, +// seqlock re-check. + +namespace zerobuffer { +namespace { + +constexpr int MAX_READ_RETRIES = 16; +constexpr auto READ_POLL_INTERVAL = std::chrono::microseconds(300); + +uint64_t now_ns() { + return static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()) + .count()); +} + +size_t align_up(size_t v, size_t a) { + return (v + a - 1) & ~(a - 1); +} + +// Force `seqlock` to the next odd value strictly greater than its current one, +// so a slot left odd by a dropped write still advances monotonically. +uint64_t begin_write_value(uint64_t current) { + return (current % 2 == 0) ? (current + 1) : (current + 2); +} + +} // namespace + +// ----------------------------- Writer --------------------------------------- + +LatestFrameWriter::LatestFrameWriter(std::string name, size_t slot_size, uint32_t slot_count) + : _name(std::move(name)), _slot_size(slot_size) { + if (_slot_size == 0) { + throw ZeroBufferException("LatestFrameWriter: slot_size must be non-zero"); + } + _slot_count = slot_count < LATEST_FRAME_MIN_SLOTS ? LATEST_FRAME_MIN_SLOTS : slot_count; + + const uint64_t slot_stride = align_up(sizeof(LatestFrameSlotHeader) + _slot_size, + LATEST_FRAME_ALIGNMENT); + const uint64_t metadata_offset = sizeof(LatestFrameSharedHeader); + const uint64_t slots_offset = align_up(metadata_offset + LATEST_FRAME_METADATA_CAPACITY, + LATEST_FRAME_ALIGNMENT); + const size_t total = static_cast(slots_offset) + + static_cast(_slot_count) * static_cast(slot_stride); + + try { + _shm = SharedMemory::create(_name, total); + } catch (const ZeroBufferException&) { + reclaim_stale_or_throw(); + _shm = SharedMemory::create(_name, total); // propagates on genuine failure + } + + _base = static_cast(_shm->data()); + _header = reinterpret_cast(_base); + + // Segment is zeroed by create(); fill everything but magic, then publish the + // magic with release so a reader only attaches once the layout is valid. + _header->version = LATEST_FRAME_VERSION; + _header->slot_count = _slot_count; + _header->slot_header_size = static_cast(sizeof(LatestFrameSlotHeader)); + _header->slot_size = _slot_size; + _header->slot_stride = slot_stride; + _header->slots_offset = slots_offset; + _header->metadata_offset = metadata_offset; + _header->metadata_capacity = LATEST_FRAME_METADATA_CAPACITY; + _header->metadata_seq = 0; + _header->publish_index = -1; + _header->writer_pid = platform::get_current_pid(); + _header->writer_start_time = platform::get_current_process_start_time(); + _header->heartbeat_ns = now_ns(); + + std::atomic_ref(_header->magic).store(LATEST_FRAME_MAGIC, std::memory_order_release); +} + +LatestFrameWriter::~LatestFrameWriter() { + release(); +} + +LatestFrameWriter::LatestFrameWriter(LatestFrameWriter&& o) noexcept + : _name(std::move(o._name)), _slot_size(o._slot_size), _slot_count(o._slot_count), + _shm(std::move(o._shm)), _header(o._header), _base(o._base), + _write_index(o._write_index), _pending_slot(o._pending_slot), + _pending_seq_even(o._pending_seq_even) { + o._header = nullptr; + o._base = nullptr; + o._pending_slot = -1; +} + +LatestFrameWriter& LatestFrameWriter::operator=(LatestFrameWriter&& o) noexcept { + if (this != &o) { + release(); + _name = std::move(o._name); + _slot_size = o._slot_size; + _slot_count = o._slot_count; + _shm = std::move(o._shm); + _header = o._header; + _base = o._base; + _write_index = o._write_index; + _pending_slot = o._pending_slot; + _pending_seq_even = o._pending_seq_even; + o._header = nullptr; + o._base = nullptr; + o._pending_slot = -1; + } + return *this; +} + +void LatestFrameWriter::reclaim_stale_or_throw() { + std::unique_ptr shm; + try { + shm = SharedMemory::open(_name); + } catch (const ZeroBufferException&) { + // Present but unmappable (e.g. a partial segment from a writer that died + // mid-create): clear it so the retry create() can own the name. + SharedMemory::remove(_name); + return; + } + + if (!shm || shm->size() < sizeof(LatestFrameSharedHeader)) { + SharedMemory::remove(_name); + return; + } + + auto* h = static_cast(shm->data()); + uint32_t magic = std::atomic_ref(h->magic).load(std::memory_order_acquire); + if (magic == LATEST_FRAME_MAGIC) { + uint64_t pid = h->writer_pid; + uint64_t start = h->writer_start_time; + bool alive = pid != 0 && platform::process_exists(pid) && + (start == 0 || platform::get_process_start_time(pid) == start); + if (alive) { + throw ZeroBufferException("LatestFrameWriter: segment '" + _name + + "' is owned by a live writer (pid " + + std::to_string(pid) + ")"); + } + } + // Stale (dead writer) or foreign/uninitialized content -> reclaim the name. + shm.reset(); + SharedMemory::remove(_name); +} + +void LatestFrameWriter::release() { + if (_shm && _header) { + // Invalidate the magic in the shared mapping BEFORE unlinking so any + // reader sharing these pages sees "producer closed" immediately and + // re-attaches to a fresh segment (works even in-process, where pid + // liveness cannot distinguish a restart). Then unlink the name. + std::atomic_ref(_header->magic).store(0, std::memory_order_release); + SharedMemory::remove(_name); + } + _shm.reset(); + _header = nullptr; + _base = nullptr; +} + +LatestFrameSlotHeader* LatestFrameWriter::slot_hdr(uint32_t i) { + return reinterpret_cast(_base + _header->slots_offset + + static_cast(i) * _header->slot_stride); +} + +uint8_t* LatestFrameWriter::slot_payload(uint32_t i) { + return reinterpret_cast(slot_hdr(i)) + sizeof(LatestFrameSlotHeader); +} + +void LatestFrameWriter::set_metadata(const void* caps, size_t len) { + if (4 + len > _header->metadata_capacity) { + throw ZeroBufferException("LatestFrameWriter: caps metadata (" + std::to_string(len) + + " bytes) exceeds capacity"); + } + uint8_t* mbase = _base + _header->metadata_offset; + std::atomic_ref mseq(_header->metadata_seq); + uint64_t begin = begin_write_value(mseq.load(std::memory_order_relaxed)); + mseq.store(static_cast(begin), std::memory_order_release); + std::atomic_thread_fence(std::memory_order_release); + + uint32_t l = static_cast(len); + mbase[0] = static_cast(l & 0xFF); + mbase[1] = static_cast((l >> 8) & 0xFF); + mbase[2] = static_cast((l >> 16) & 0xFF); + mbase[3] = static_cast((l >> 24) & 0xFF); + if (len > 0) { + std::memcpy(mbase + 4, caps, len); + } + + std::atomic_thread_fence(std::memory_order_release); + mseq.store(static_cast(begin + 1), std::memory_order_release); +} + +uint8_t* LatestFrameWriter::get_slot() { + int64_t pub = std::atomic_ref(_header->publish_index).load(std::memory_order_acquire); + uint32_t w = (_write_index + 1) % _slot_count; + if (static_cast(w) == pub) { + w = (w + 1) % _slot_count; + } + _write_index = w; + _pending_slot = static_cast(w); + + LatestFrameSlotHeader* sh = slot_hdr(w); + std::atomic_ref sref(sh->seqlock); + uint64_t begin = begin_write_value(sref.load(std::memory_order_relaxed)); + sref.store(begin, std::memory_order_release); + _pending_seq_even = begin + 1; + std::atomic_thread_fence(std::memory_order_release); + return slot_payload(w); +} + +void LatestFrameWriter::publish(uint64_t sequence, size_t size) { + if (_pending_slot < 0) { + return; + } + uint32_t w = static_cast(_pending_slot); + if (size > _slot_size) { + size = _slot_size; + } + LatestFrameSlotHeader* sh = slot_hdr(w); + sh->frame_number = sequence; + sh->size = size; + + std::atomic_thread_fence(std::memory_order_release); + std::atomic_ref(sh->seqlock).store(_pending_seq_even, std::memory_order_release); + std::atomic_ref(_header->publish_index).store(static_cast(w), + std::memory_order_release); + std::atomic_ref(_header->heartbeat_ns).store(now_ns(), std::memory_order_release); + _pending_slot = -1; +} + +void LatestFrameWriter::heartbeat() { + std::atomic_ref(_header->heartbeat_ns).store(now_ns(), std::memory_order_release); +} + +// ----------------------------- Reader --------------------------------------- + +LatestFrameReader::LatestFrameReader(std::string name) : _name(std::move(name)) { + try_attach(); +} + +LatestFrameReader::~LatestFrameReader() { + detach(); +} + +LatestFrameReader::LatestFrameReader(LatestFrameReader&& o) noexcept + : _name(std::move(o._name)), _shm(std::move(o._shm)), _header(o._header), _base(o._base), + _scratch(std::move(o._scratch)), _meta_cache(std::move(o._meta_cache)), + _meta_size(o._meta_size), _meta_seq_seen(o._meta_seq_seen), _meta_valid(o._meta_valid), + _last_sequence(o._last_sequence), _have_last(o._have_last) { + o._header = nullptr; + o._base = nullptr; +} + +LatestFrameReader& LatestFrameReader::operator=(LatestFrameReader&& o) noexcept { + if (this != &o) { + detach(); + _name = std::move(o._name); + _shm = std::move(o._shm); + _header = o._header; + _base = o._base; + _scratch = std::move(o._scratch); + _meta_cache = std::move(o._meta_cache); + _meta_size = o._meta_size; + _meta_seq_seen = o._meta_seq_seen; + _meta_valid = o._meta_valid; + _last_sequence = o._last_sequence; + _have_last = o._have_last; + o._header = nullptr; + o._base = nullptr; + } + return *this; +} + +bool LatestFrameReader::try_attach() { + if (_header) { + return true; + } + std::unique_ptr shm; + try { + shm = SharedMemory::open(_name); + } catch (const ZeroBufferException&) { + return false; // absent -> caller polls + } + if (!shm || shm->size() < sizeof(LatestFrameSharedHeader)) { + return false; + } + + auto* h = static_cast(shm->data()); + if (std::atomic_ref(h->magic).load(std::memory_order_acquire) != LATEST_FRAME_MAGIC) { + return false; // creator has not published the layout yet + } + if (h->version != LATEST_FRAME_VERSION || h->slot_count == 0 || h->slot_size == 0) { + return false; + } + uint64_t need = h->slots_offset + static_cast(h->slot_count) * h->slot_stride; + if (need > shm->size()) { + return false; + } + + _shm = std::move(shm); + _header = h; + _base = static_cast(_shm->data()); + _scratch.assign(static_cast(_header->slot_size), 0); + _meta_valid = false; + _meta_size = 0; + _meta_seq_seen = 0; + _have_last = false; + refresh_metadata(); + return true; +} + +void LatestFrameReader::detach() { + _shm.reset(); + _header = nullptr; + _base = nullptr; + _meta_valid = false; + _have_last = false; +} + +const LatestFrameSlotHeader* LatestFrameReader::slot_hdr(uint32_t i) const { + return reinterpret_cast( + _base + _header->slots_offset + static_cast(i) * _header->slot_stride); +} + +const uint8_t* LatestFrameReader::slot_payload(uint32_t i) const { + return reinterpret_cast(slot_hdr(i)) + sizeof(LatestFrameSlotHeader); +} + +bool LatestFrameReader::try_read_once(LatestFrame& out) { + for (int attempt = 0; attempt < MAX_READ_RETRIES; ++attempt) { + int64_t idx = + std::atomic_ref(_header->publish_index).load(std::memory_order_acquire); + if (idx < 0 || idx >= static_cast(_header->slot_count)) { + return false; // nothing published yet + } + uint32_t w = static_cast(idx); + const LatestFrameSlotHeader* sh = slot_hdr(w); + uint64_t& seq_field = const_cast(sh->seqlock); + + uint64_t s1 = std::atomic_ref(seq_field).load(std::memory_order_acquire); + if (s1 & 1) { + std::this_thread::yield(); // writer mid-write on this slot + continue; + } + uint64_t fn = sh->frame_number; + uint64_t sz = sh->size; + if (sz > _scratch.size()) { + sz = _scratch.size(); // defensive: never read past the slot + } + std::memcpy(_scratch.data(), slot_payload(w), sz); + std::atomic_thread_fence(std::memory_order_acquire); + + uint64_t s2 = std::atomic_ref(seq_field).load(std::memory_order_acquire); + if (s1 != s2) { + continue; // slot rewritten under us -> retry, lands on the newest + } + + if (_have_last && fn == _last_sequence) { + return false; // no new frame since the last read + } + _last_sequence = fn; + _have_last = true; + out._data = _scratch.data(); + out._size = static_cast(sz); + out._sequence = fn; + out._valid = true; + return true; + } + return false; +} + +LatestFrame LatestFrameReader::read_latest(std::chrono::milliseconds timeout) { + auto deadline = std::chrono::steady_clock::now() + timeout; + for (;;) { + if (!_header) { + try_attach(); + } + if (_header) { + if (std::atomic_ref(_header->magic).load(std::memory_order_acquire) != + LATEST_FRAME_MAGIC) { + detach(); // producer closed / restarted -> re-attach to a fresh segment + } else { + refresh_metadata(); + LatestFrame f; + if (try_read_once(f)) { + return f; + } + } + } + if (std::chrono::steady_clock::now() >= deadline) { + // Producer gone: drop the mapping so the next call re-attaches to a + // freshly created segment (writer restart recovery, ADR-5). + if (_header && writer_gone()) { + detach(); + } + return LatestFrame{}; + } + std::this_thread::sleep_for(READ_POLL_INTERVAL); + } +} + +void LatestFrameReader::refresh_metadata() { + if (!_header) { + return; + } + std::atomic_ref mseq(_header->metadata_seq); + uint32_t s1 = mseq.load(std::memory_order_acquire); + if (s1 & 1) { + return; // caps being rewritten -> keep the previous cache + } + if (_meta_valid && s1 == _meta_seq_seen) { + return; // unchanged + } + + const uint8_t* mbase = _base + _header->metadata_offset; + uint32_t len = static_cast(mbase[0]) | + (static_cast(mbase[1]) << 8) | + (static_cast(mbase[2]) << 16) | + (static_cast(mbase[3]) << 24); + if (4 + static_cast(len) > _header->metadata_capacity) { + return; // malformed -> keep the previous cache + } + _meta_cache.resize(4 + static_cast(len)); + std::memcpy(_meta_cache.data(), mbase, _meta_cache.size()); + std::atomic_thread_fence(std::memory_order_acquire); + + uint32_t s2 = mseq.load(std::memory_order_acquire); + if (s1 != s2) { + return; // torn -> refresh next time + } + _meta_seq_seen = s1; + _meta_size = _meta_cache.size(); + _meta_valid = true; +} + +const void* LatestFrameReader::get_metadata_raw() { + return _meta_valid ? _meta_cache.data() : nullptr; +} + +size_t LatestFrameReader::get_metadata_size() { + return _meta_valid ? _meta_size : 0; +} + +bool LatestFrameReader::is_writer_alive() { + if (!_header) { + return false; + } + if (std::atomic_ref(_header->magic).load(std::memory_order_acquire) != + LATEST_FRAME_MAGIC) { + return false; // producer closed the segment + } + uint64_t pid = _header->writer_pid; + if (pid == 0 || !platform::process_exists(pid)) { + return false; + } + uint64_t start = _header->writer_start_time; + if (start != 0 && platform::get_process_start_time(pid) != start) { + return false; // pid reused by a different process + } + return true; +} + +bool LatestFrameReader::writer_gone() { + return !is_writer_alive(); +} + +} // namespace zerobuffer diff --git a/cpp/tests/CMakeLists.txt b/cpp/tests/CMakeLists.txt index 2535e1d..02cf65d 100644 --- a/cpp/tests/CMakeLists.txt +++ b/cpp/tests/CMakeLists.txt @@ -1,5 +1,27 @@ cmake_minimum_required(VERSION 3.20) +# Ensure nlohmann_json is available for the generated tests, which link it but +# rely on the parent scope providing the target. The main CMakeLists only fetches +# it under BUILD_SERVE and AFTER add_subdirectory(tests), so hoist it here to make +# the test tree self-sufficient (breaks CMake configure otherwise). +if(NOT TARGET nlohmann_json::nlohmann_json) + find_package(nlohmann_json QUIET) + if(NOT nlohmann_json_FOUND) + include(FetchContent) + set(CMAKE_POLICY_DEFAULT_CMP0048 NEW) + cmake_policy(SET CMP0048 NEW) + FetchContent_Declare( + json + GIT_REPOSITORY https://github.com/nlohmann/json.git + GIT_TAG v3.11.3 + CMAKE_ARGS -DCMAKE_POLICY_DEFAULT_CMP0048=NEW + ) + set(JSON_BuildTests OFF CACHE INTERNAL "") + set(JSON_Install OFF CACHE INTERNAL "") + FetchContent_MakeAvailable(json) + endif() +endif() + # POC tests for testing infrastructure development # DISABLED: These are proof-of-concept tests, not needed for production # add_subdirectory(poc) @@ -34,21 +56,44 @@ cmake_minimum_required(VERSION 3.20) # ) # Frame RAII tests -add_executable(test_frame_raii - test_frame_raii.cpp +# NOTE: test_frame_raii.cpp is referenced here but is absent from the repo (not +# git-tracked on master), which breaks CMake configure with BUILD_TESTS=ON. +# Guarded so the source is only added when present; unrelated to the latest-frame +# primitive. Flagged to the team lead. +if(EXISTS "${CMAKE_CURRENT_SOURCE_DIR}/test_frame_raii.cpp") + add_executable(test_frame_raii + test_frame_raii.cpp + ) + + target_link_libraries(test_frame_raii + PRIVATE + zerobuffer + GTest::gtest + GTest::gtest_main + ) + + include(GoogleTest) + gtest_discover_tests(test_frame_raii + PROPERTIES TIMEOUT 30 + ) +endif() + +include(GoogleTest) + +# Latest-frame primitive tests (LatestFrameWriter / LatestFrameReader) +add_executable(test_latest_frame + test_latest_frame.cpp ) -target_link_libraries(test_frame_raii +target_link_libraries(test_latest_frame PRIVATE zerobuffer GTest::gtest GTest::gtest_main ) -# Add test discovery -include(GoogleTest) -gtest_discover_tests(test_frame_raii - PROPERTIES TIMEOUT 30 +gtest_discover_tests(test_latest_frame + PROPERTIES TIMEOUT 60 ) # Generated tests from feature files diff --git a/cpp/tests/test_latest_frame.cpp b/cpp/tests/test_latest_frame.cpp new file mode 100644 index 0000000..be69841 --- /dev/null +++ b/cpp/tests/test_latest_frame.cpp @@ -0,0 +1,367 @@ +#include + +#include "zerobuffer/latest_frame.h" +#include "zerobuffer/platform.h" + +#include +#include +#include +#include +#include +#include + +using namespace zerobuffer; +using namespace std::chrono_literals; + +namespace { + +// Unique buffer name per test so parallel/repeated runs never collide, and a +// pre-clean in case a previous crashed run left a segment behind. +std::string make_name(const char* suffix) { + std::string name = "/lf_test_" + std::to_string(platform::get_current_pid()) + "_" + suffix; + SharedMemory::remove(name); + return name; +} + +// A frame whose every byte equals (sequence & 0xFF). A tear-free read must find +// all bytes identical; a torn read (mixed laps) would show two values. +void fill_pattern(uint8_t* p, size_t size, uint64_t seq) { + std::memset(p, static_cast(seq & 0xFF), size); +} + +bool all_bytes_equal(const uint8_t* p, size_t size, uint8_t value) { + for (size_t i = 0; i < size; ++i) { + if (p[i] != value) return false; + } + return true; +} + +} // namespace + +// FR-7: the writer creates and publishes with NO reader present (the ADR-1 +// preroll bug must be impossible for this primitive). +TEST(LatestFrameTest, WriterCreatesAndPublishesWithNoReader) { + std::string name = make_name("no_reader"); + LatestFrameWriter writer(name, 4096, 3); + + for (uint64_t seq = 0; seq < 10; ++seq) { + uint8_t* slot = writer.get_slot(); + ASSERT_NE(slot, nullptr); + fill_pattern(slot, 4096, seq); + writer.publish(seq, 4096); // never blocks, no reader exists + } + SUCCEED(); +} + +// Slot count below the triple-slot minimum is raised, not accepted verbatim. +TEST(LatestFrameTest, SlotCountRaisedToMinimum) { + std::string name = make_name("min_slots"); + LatestFrameWriter writer(name, 1024, 1); + EXPECT_EQ(writer.slot_count(), LATEST_FRAME_MIN_SLOTS); +} + +TEST(LatestFrameTest, ZeroSlotSizeRejected) { + std::string name = make_name("zero_size"); + EXPECT_THROW(LatestFrameWriter(name, 0, 3), ZeroBufferException); +} + +// Reader attaching AFTER the writer gets the latest published frame. +TEST(LatestFrameTest, ReaderAttachesAfterWriterGetsLatest) { + std::string name = make_name("attach_after"); + LatestFrameWriter writer(name, 4096, 3); + for (uint64_t seq = 0; seq <= 5; ++seq) { + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 4096, seq); + writer.publish(seq, 4096); + } + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_EQ(f.sequence(), 5u); + EXPECT_EQ(f.size(), 4096u); + EXPECT_TRUE(all_bytes_equal(static_cast(f.data()), f.size(), 5 & 0xFF)); +} + +// Newest-frame-wins: under a burst, a late reader skips intermediates. +TEST(LatestFrameTest, NewestWinsSkipsIntermediates) { + std::string name = make_name("newest_wins"); + LatestFrameWriter writer(name, 1024, 3); + LatestFrameReader reader(name); + + for (uint64_t seq = 0; seq < 100; ++seq) { + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, seq); + writer.publish(seq, 1024); + } + + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_EQ(f.sequence(), 99u); // landed on the newest, not seq 0 +} + +// Segment absent for the reader -> read_latest returns invalid (poll/retry), +// never crashes or blocks past the timeout. +TEST(LatestFrameTest, ReaderToleratesAbsentSegment) { + std::string name = make_name("absent"); + LatestFrameReader reader(name); // nothing created it + EXPECT_FALSE(reader.is_attached()); + + auto start = std::chrono::steady_clock::now(); + LatestFrame f = reader.read_latest(100ms); + auto elapsed = std::chrono::steady_clock::now() - start; + EXPECT_FALSE(f.valid()); + EXPECT_GE(elapsed, 90ms); + EXPECT_LT(elapsed, 2000ms); +} + +// Reader detach + re-attach mid-stream; writer is unaffected and keeps cadence. +TEST(LatestFrameTest, ReaderDetachReattachWriterUnaffected) { + std::string name = make_name("reattach"); + LatestFrameWriter writer(name, 1024, 3); + + { + LatestFrameReader reader(name); // attach before any frame, then detach + EXPECT_FALSE(reader.read_latest(50ms).valid()); + } + uint64_t seq = 0; + auto publish_one = [&](uint64_t s) { + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, s); + writer.publish(s, 1024); + }; + publish_one(seq++); + + { + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + } + // Writer keeps publishing after the reader detached. + for (int i = 0; i < 20; ++i) publish_one(seq++); + + LatestFrameReader reader2(name); + LatestFrame f = reader2.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_EQ(f.sequence(), seq - 1); +} + +// Writer restart (segment recreated, sequence reset) while a reader polls -> +// reader recovers on the next publish. +TEST(LatestFrameTest, WriterRestartReaderRecovers) { + std::string name = make_name("restart"); + auto writer = std::make_unique(name, 1024, 3); + for (uint64_t seq = 0; seq < 50; ++seq) { + uint8_t* slot = writer->get_slot(); + fill_pattern(slot, 1024, seq); + writer->publish(seq, 1024); + } + + LatestFrameReader reader(name); + LatestFrame f1 = reader.read_latest(1000ms); + ASSERT_TRUE(f1.valid()); + EXPECT_EQ(f1.sequence(), 49u); + + // Restart: destroy the writer (unlinks) and create a fresh one (seq from 0). + writer.reset(); + writer = std::make_unique(name, 1024, 3); + + // Reader recovers and reads the new segment's frames. + bool recovered = false; + for (int i = 0; i < 100 && !recovered; ++i) { + uint8_t* slot = writer->get_slot(); + fill_pattern(slot, 1024, static_cast(i)); + writer->publish(static_cast(i), 1024); + LatestFrame f = reader.read_latest(200ms); + if (f.valid()) { + EXPECT_TRUE(all_bytes_equal(static_cast(f.data()), f.size(), + f.sequence() & 0xFF)); + recovered = true; + } + } + EXPECT_TRUE(recovered); +} + +// Writer gone / stale heartbeat -> reader reports producer-gone, no crash/block. +TEST(LatestFrameTest, WriterGoneReportedNoCrash) { + std::string name = make_name("gone"); + { + LatestFrameWriter writer(name, 1024, 3); + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, 1); + writer.publish(1, 1024); + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_TRUE(reader.is_writer_alive()); + // Writer destroyed at end of scope, before the reader below. + LatestFrameReader probe(name); + EXPECT_TRUE(probe.is_writer_alive()); + } + // Segment is now unlinked; a fresh reader cannot attach and reports gone. + LatestFrameReader reader(name); + EXPECT_FALSE(reader.is_writer_alive()); + LatestFrame f = reader.read_latest(100ms); + EXPECT_FALSE(f.valid()); +} + +// Oversized frame vs slot_size -> handled without out-of-bounds; the reader +// never sees more than slot_size bytes. +TEST(LatestFrameTest, OversizedFrameClampedNoOverflow) { + std::string name = make_name("oversize"); + LatestFrameWriter writer(name, 1024, 3); + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, 7); + writer.publish(7, 100000); // claim way more than slot_size + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_LE(f.size(), 1024u); // clamped to slot_size, no OOB +} + +// Metadata/caps set + rewritten on change -> reader re-reads (FR-3). Verifies +// the exact ShmCaps wire shape: 4-byte LE length prefix + JSON. +TEST(LatestFrameTest, MetadataSetAndRewrite) { + std::string name = make_name("metadata"); + LatestFrameWriter writer(name, 1024, 3); + + std::string caps1 = R"({"caps":"video/x-raw,format=GRAY8,width=640,height=480"})"; + writer.set_metadata(caps1.data(), caps1.size()); + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, 1); + writer.publish(1, 1024); + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + + const uint8_t* md = static_cast(reader.get_metadata_raw()); + size_t md_size = reader.get_metadata_size(); + ASSERT_NE(md, nullptr); + ASSERT_EQ(md_size, 4 + caps1.size()); + uint32_t prefix = static_cast(md[0]) | (static_cast(md[1]) << 8) | + (static_cast(md[2]) << 16) | (static_cast(md[3]) << 24); + EXPECT_EQ(prefix, caps1.size()); + EXPECT_EQ(0, std::memcmp(md + 4, caps1.data(), caps1.size())); + + // Rewrite caps on a change; the reader re-reads on the next poll. + std::string caps2 = R"({"caps":"video/x-raw,format=I420,width=1920,height=1080"})"; + writer.set_metadata(caps2.data(), caps2.size()); + slot = writer.get_slot(); + fill_pattern(slot, 1024, 2); + writer.publish(2, 1024); + + f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + md = static_cast(reader.get_metadata_raw()); + md_size = reader.get_metadata_size(); + ASSERT_EQ(md_size, 4 + caps2.size()); + EXPECT_EQ(0, std::memcmp(md + 4, caps2.data(), caps2.size())); +} + +// Tear-free read under concurrent write: a reader loop vs a fast writer over +// many iterations must never return a torn (mixed-lap) frame. +TEST(LatestFrameTest, TearFreeUnderConcurrentWrite) { + std::string name = make_name("tearfree"); + constexpr size_t FRAME = 256 * 1024; // large payload widens the tear window + constexpr uint64_t ITERATIONS = 20000; + + LatestFrameWriter writer(name, FRAME, 3); + std::atomic stop{false}; + std::atomic torn{0}; + std::atomic reads{0}; + + std::thread reader_thread([&]() { + LatestFrameReader reader(name); + while (!stop.load(std::memory_order_relaxed)) { + LatestFrame f = reader.read_latest(50ms); + if (!f.valid()) continue; + const uint8_t* p = static_cast(f.data()); + uint8_t expected = static_cast(f.sequence() & 0xFF); + if (!all_bytes_equal(p, f.size(), expected)) { + torn.fetch_add(1, std::memory_order_relaxed); + } + reads.fetch_add(1, std::memory_order_relaxed); + } + }); + + for (uint64_t seq = 0; seq < ITERATIONS; ++seq) { + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, FRAME, seq); + writer.publish(seq, FRAME); + } + // Let the reader drain a little, then stop. + std::this_thread::sleep_for(50ms); + stop.store(true, std::memory_order_relaxed); + reader_thread.join(); + + EXPECT_EQ(torn.load(), 0u) << "torn frames observed"; + EXPECT_GT(reads.load(), 0u) << "reader never observed a frame"; +} + +// Two readers on one segment both read tear-free (1..N readers, FR). +TEST(LatestFrameTest, MultipleReaders) { + std::string name = make_name("multi_reader"); + LatestFrameWriter writer(name, 4096, 4); + for (uint64_t seq = 0; seq <= 3; ++seq) { + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 4096, seq); + writer.publish(seq, 4096); + } + + LatestFrameReader r1(name); + LatestFrameReader r2(name); + LatestFrame f1 = r1.read_latest(1000ms); + LatestFrame f2 = r2.read_latest(1000ms); + ASSERT_TRUE(f1.valid()); + ASSERT_TRUE(f2.valid()); + EXPECT_EQ(f1.sequence(), 3u); + EXPECT_EQ(f2.sequence(), 3u); +} + +// A live writer owning the segment cannot be displaced by a second writer. +TEST(LatestFrameTest, SecondWriterOnLiveSegmentRejected) { + std::string name = make_name("double_writer"); + LatestFrameWriter writer(name, 1024, 3); + EXPECT_THROW(LatestFrameWriter(name, 1024, 3), ZeroBufferException); +} + +// A stale segment left by a dead writer is reclaimed on create. Build a genuine +// stale segment through the platform layer (magic + dead pid, never unlinked), +// so LatestFrameWriter::create() hits EEXIST and must reclaim it. +TEST(LatestFrameTest, StaleSegmentReclaimedOnCreate) { + std::string name = make_name("stale"); + + // Hand-craft a leftover segment: valid magic, but writer_pid = 0 (which the + // reclaim path treats as no live owner). It is NOT unlinked, mimicking a + // writer that crashed without cleanup. + { + auto shm = SharedMemory::create(name, 65536); + auto* h = static_cast(shm->data()); + std::memset(h, 0, sizeof(*h)); + h->version = LATEST_FRAME_VERSION; + h->slot_count = 3; + h->slot_size = 1024; + h->slot_stride = 1088; + h->slots_offset = 8320; + h->metadata_offset = sizeof(LatestFrameSharedHeader); + h->metadata_capacity = LATEST_FRAME_METADATA_CAPACITY; + h->publish_index = -1; + h->writer_pid = 0; // dead / no owner -> stale + std::atomic_ref(h->magic).store(LATEST_FRAME_MAGIC, std::memory_order_release); + // shm mapping dropped here WITHOUT shm_unlink -> the named segment persists. + } + + // A new writer must reclaim the stale name and create cleanly. + LatestFrameWriter fresh(name, 1024, 3); + uint8_t* slot = fresh.get_slot(); + fill_pattern(slot, 1024, 9); + fresh.publish(9, 1024); + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_EQ(f.sequence(), 9u); +} From a3ea38635fa24e874e1ac846fc45d769b82b24e6 Mon Sep 17 00:00:00 2001 From: Daniel Kowalski Date: Thu, 9 Jul 2026 10:40:58 +0200 Subject: [PATCH 2/5] =?UTF-8?q?refactor:=20address=20review=20=E2=80=94=20?= =?UTF-8?q?drop=20unused=20heartbeat,=20harden=20seqlock=20fields,=20gate?= =?UTF-8?q?=20empty=20caps?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 4 +--- cpp/include/zerobuffer/latest_frame.h | 12 ++++------ cpp/src/latest_frame.cpp | 33 ++++++++++++++------------- 3 files changed, 22 insertions(+), 27 deletions(-) diff --git a/.gitignore b/.gitignore index 160e85b..27b7c58 100644 --- a/.gitignore +++ b/.gitignore @@ -42,10 +42,8 @@ Testing/ /cpp/tests/*_test /cpp/benchmarks/test_* /cpp/benchmarks/*_test -# ...but keep test/benchmark SOURCES tracked (the patterns above target built binaries) +# ...but keep test SOURCES tracked (patterns above target built binaries) !/cpp/tests/test_*.cpp -!/cpp/tests/test_*.h -!/cpp/benchmarks/test_*.cpp # C# build artifacts /csharp/.vs/ diff --git a/cpp/include/zerobuffer/latest_frame.h b/cpp/include/zerobuffer/latest_frame.h index 2a8907f..25889ff 100644 --- a/cpp/include/zerobuffer/latest_frame.h +++ b/cpp/include/zerobuffer/latest_frame.h @@ -43,7 +43,7 @@ struct LatestFrameSlotHeader { static_assert(sizeof(LatestFrameSlotHeader) == 32, "slot header must be 32 bytes"); // Segment header. POD; accessed cross-process. Concurrency-sensitive fields -// (seqlock, publish_index, heartbeat_ns, metadata_seq) are read/written through +// (seqlock, publish_index, metadata_seq, magic) are read/written through // std::atomic_ref so the underlying bytes stay a plain cross-language layout. struct LatestFrameSharedHeader { uint32_t magic; @@ -60,8 +60,7 @@ struct LatestFrameSharedHeader { int64_t publish_index; // newest published slot; -1 = none published uint64_t writer_pid; // owning writer pid (0 = none) uint64_t writer_start_time; // writer process start time (pid-reuse guard) - uint64_t heartbeat_ns; // steady_clock tick stamped by the writer - uint64_t reserved1[4]; + uint64_t reserved1[5]; }; static_assert(sizeof(LatestFrameSharedHeader) == 128, "shared header must be 128 bytes"); static_assert(sizeof(LatestFrameSharedHeader) % LATEST_FRAME_ALIGNMENT == 0, @@ -115,14 +114,11 @@ class LatestFrameWriter { // writes up to slot_size() bytes then calls publish(). uint8_t* get_slot(); - // Finalize the slot reserved by get_slot(): record (sequence, size), make it - // the newest published slot (seqlock release), and stamp the heartbeat. + // Finalize the slot reserved by get_slot(): record (sequence, size) and make + // it the newest published slot (seqlock release). // `size` above slot_size() is clamped. Never blocks on a reader. void publish(uint64_t sequence, size_t size); - // Stamp writer liveness without publishing a frame. - void heartbeat(); - size_t slot_size() const { return _slot_size; } uint32_t slot_count() const { return _slot_count; } const std::string& name() const { return _name; } diff --git a/cpp/src/latest_frame.cpp b/cpp/src/latest_frame.cpp index 977a9c4..1421105 100644 --- a/cpp/src/latest_frame.cpp +++ b/cpp/src/latest_frame.cpp @@ -26,12 +26,9 @@ namespace { constexpr int MAX_READ_RETRIES = 16; constexpr auto READ_POLL_INTERVAL = std::chrono::microseconds(300); -uint64_t now_ns() { - return static_cast( - std::chrono::duration_cast( - std::chrono::steady_clock::now().time_since_epoch()) - .count()); -} +// A completed metadata write leaves metadata_seq even and >= 2 (0 = never +// written, 1 = first write in progress). Readers treat < 2 as "no caps yet". +constexpr uint32_t METADATA_WRITTEN = 2; size_t align_up(size_t v, size_t a) { return (v + a - 1) & ~(a - 1); @@ -86,7 +83,6 @@ LatestFrameWriter::LatestFrameWriter(std::string name, size_t slot_size, uint32_ _header->publish_index = -1; _header->writer_pid = platform::get_current_pid(); _header->writer_start_time = platform::get_current_process_start_time(); - _header->heartbeat_ns = now_ns(); std::atomic_ref(_header->magic).store(LATEST_FRAME_MAGIC, std::memory_order_release); } @@ -232,21 +228,19 @@ void LatestFrameWriter::publish(uint64_t sequence, size_t size) { size = _slot_size; } LatestFrameSlotHeader* sh = slot_hdr(w); - sh->frame_number = sequence; - sh->size = size; + // Relaxed atomic stores of the seqlock-guarded fields: the seqlock release + // below is the actual publication fence; relaxed keeps these free of any + // abstract-machine data race on the fields themselves. + std::atomic_ref(sh->frame_number).store(sequence, std::memory_order_relaxed); + std::atomic_ref(sh->size).store(size, std::memory_order_relaxed); std::atomic_thread_fence(std::memory_order_release); std::atomic_ref(sh->seqlock).store(_pending_seq_even, std::memory_order_release); std::atomic_ref(_header->publish_index).store(static_cast(w), std::memory_order_release); - std::atomic_ref(_header->heartbeat_ns).store(now_ns(), std::memory_order_release); _pending_slot = -1; } -void LatestFrameWriter::heartbeat() { - std::atomic_ref(_header->heartbeat_ns).store(now_ns(), std::memory_order_release); -} - // ----------------------------- Reader --------------------------------------- LatestFrameReader::LatestFrameReader(std::string name) : _name(std::move(name)) { @@ -357,11 +351,15 @@ bool LatestFrameReader::try_read_once(LatestFrame& out) { std::this_thread::yield(); // writer mid-write on this slot continue; } - uint64_t fn = sh->frame_number; - uint64_t sz = sh->size; + uint64_t fn = std::atomic_ref(const_cast(sh->frame_number)) + .load(std::memory_order_relaxed); + uint64_t sz = std::atomic_ref(const_cast(sh->size)) + .load(std::memory_order_relaxed); if (sz > _scratch.size()) { sz = _scratch.size(); // defensive: never read past the slot } + // Plain payload copy, validated by the seqlock re-check below: a byte + // race here is torn iff the seqlock changed, which the re-check catches. std::memcpy(_scratch.data(), slot_payload(w), sz); std::atomic_thread_fence(std::memory_order_acquire); @@ -420,6 +418,9 @@ void LatestFrameReader::refresh_metadata() { } std::atomic_ref mseq(_header->metadata_seq); uint32_t s1 = mseq.load(std::memory_order_acquire); + if (s1 < METADATA_WRITTEN) { + return; // no caps written yet -> stays "not available", not empty caps + } if (s1 & 1) { return; // caps being rewritten -> keep the previous cache } From 6f9dc523d9dc41be26b4644305c08438775860ac Mon Sep 17 00:00:00 2001 From: Daniel Kowalski Date: Thu, 9 Jul 2026 10:48:39 +0200 Subject: [PATCH 3/5] test: assert caps not-ready before first set_metadata --- cpp/tests/test_latest_frame.cpp | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/cpp/tests/test_latest_frame.cpp b/cpp/tests/test_latest_frame.cpp index be69841..b40369e 100644 --- a/cpp/tests/test_latest_frame.cpp +++ b/cpp/tests/test_latest_frame.cpp @@ -261,6 +261,22 @@ TEST(LatestFrameTest, MetadataSetAndRewrite) { EXPECT_EQ(0, std::memcmp(md + 4, caps2.data(), caps2.size())); } +// A frame published BEFORE the first set_metadata -> reader reports caps +// not-ready (null / size 0), never a valid-but-empty caps block (MINOR-2). +TEST(LatestFrameTest, MetadataNotReadyBeforeFirstSet) { + std::string name = make_name("meta_notready"); + LatestFrameWriter writer(name, 1024, 3); + uint8_t* slot = writer.get_slot(); + fill_pattern(slot, 1024, 1); + writer.publish(1, 1024); // published, but no set_metadata yet + + LatestFrameReader reader(name); + LatestFrame f = reader.read_latest(1000ms); + ASSERT_TRUE(f.valid()); + EXPECT_EQ(reader.get_metadata_raw(), nullptr); + EXPECT_EQ(reader.get_metadata_size(), 0u); +} + // Tear-free read under concurrent write: a reader loop vs a fast writer over // many iterations must never return a torn (mixed-lap) frame. TEST(LatestFrameTest, TearFreeUnderConcurrentWrite) { From 5f52872ad9d84f45f02541549136eb78ab1b1fa1 Mon Sep 17 00:00:00 2001 From: Daniel Kowalski Date: Thu, 9 Jul 2026 11:22:37 +0200 Subject: [PATCH 4/5] refactor: zero-copy seqlock read (read_latest_into), drop per-frame scratch copy --- cpp/include/zerobuffer/latest_frame.h | 132 +++++++++--- cpp/src/latest_frame.cpp | 102 +--------- cpp/tests/test_latest_frame.cpp | 279 ++++++++++++++------------ 3 files changed, 254 insertions(+), 259 deletions(-) diff --git a/cpp/include/zerobuffer/latest_frame.h b/cpp/include/zerobuffer/latest_frame.h index 25889ff..5f423ca 100644 --- a/cpp/include/zerobuffer/latest_frame.h +++ b/cpp/include/zerobuffer/latest_frame.h @@ -8,18 +8,23 @@ // // The writer creates and owns the segment and never waits on a reader; N // readers map it read-only and read the newest published slot tear-free via a -// per-slot seqlock over a triple- (or wider) slot ring. Caps travel in the -// header as a length-prefixed JSON block consumed verbatim by native-player's -// ShmCaps::parse (4-byte little-endian length prefix + JSON). +// per-slot seqlock over a triple- (or wider) slot ring. Reads are ZERO-COPY: +// LatestFrameReader::read_latest_into() hands the caller a pointer directly into +// the published slot and validates the seqlock around the caller's consume — no +// per-frame copy and no per-frame allocation (zerobuffer perf rule #1). Caps +// travel in the header as a length-prefixed JSON block consumed verbatim by +// native-player's ShmCaps::parse (4-byte little-endian length prefix + JSON). #include "zerobuffer/platform.h" #include "zerobuffer/reader.h" // ZeroBufferException +#include #include #include #include #include #include +#include #include namespace zerobuffer { @@ -31,6 +36,8 @@ constexpr uint32_t LATEST_FRAME_VERSION = 1u; constexpr uint32_t LATEST_FRAME_MIN_SLOTS = 3u; // triple-slot minimum constexpr size_t LATEST_FRAME_ALIGNMENT = 64u; // cache-line alignment constexpr size_t LATEST_FRAME_METADATA_CAPACITY = 8192u; // caps JSON block bytes +constexpr int LATEST_FRAME_READ_RETRIES = 16; // seqlock retries per read +constexpr std::chrono::microseconds LATEST_FRAME_READ_POLL{300}; // poll gap while waiting // Per-slot header. `seqlock` is the tear-free guard: even = stable, odd = write // in progress. `frame_number` and `size` are plain fields protected by it. @@ -66,27 +73,6 @@ static_assert(sizeof(LatestFrameSharedHeader) == 128, "shared header must be 128 static_assert(sizeof(LatestFrameSharedHeader) % LATEST_FRAME_ALIGNMENT == 0, "shared header must be cache-line aligned"); -// ----- frame view returned by the reader ------------------------------------ - -// Zero-allocation view of one tear-free frame copy. `data()` points into the -// reader's reusable scratch buffer and is valid until the next read_latest(). -class LatestFrame { -public: - LatestFrame() = default; - - const void* data() const { return _data; } - size_t size() const { return _size; } - uint64_t sequence() const { return _sequence; } - bool valid() const { return _valid; } - -private: - friend class LatestFrameReader; - const void* _data = nullptr; - size_t _size = 0; - uint64_t _sequence = 0; - bool _valid = false; -}; - // ----- writer (producer) ---------------------------------------------------- // Creates and OWNS the segment. Constructs successfully with no reader present @@ -154,10 +140,43 @@ class LatestFrameReader { LatestFrameReader(LatestFrameReader&&) noexcept; LatestFrameReader& operator=(LatestFrameReader&&) noexcept; - // Read the newest published frame tear-free (seqlock). Returns an invalid - // frame if the segment is absent or no NEW frame appears within `timeout`. - // Intermediate frames the writer overwrote are drops (sequence gap). - LatestFrame read_latest(std::chrono::milliseconds timeout); + // Zero-copy tear-free read of the newest published frame. Invokes + // consume(const uint8_t* src, size_t size, uint64_t sequence) + // with `src` pointing DIRECTLY into the published slot (no copy). The seqlock + // is validated AROUND the consume: if the writer laps the slot mid-read, + // consume is re-invoked on the newest slot (bounded retries), so the caller + // must treat its work as committed only when this returns true. Returns false + // if the segment is absent, no NEW frame appears within `timeout`, or all + // retries tore. Intermediate frames the writer overwrote are drops (sequence + // gap). `size` is clamped to slot_size. The primitive stays generic — it + // never copies and knows nothing about the caller's consume. + template + bool read_latest_into(Fn&& consume, std::chrono::milliseconds timeout) { + auto deadline = std::chrono::steady_clock::now() + timeout; + for (;;) { + if (!_header) { + try_attach(); + } + if (_header) { + if (std::atomic_ref(_header->magic).load(std::memory_order_acquire) != + LATEST_FRAME_MAGIC) { + detach(); // producer closed / restarted -> re-attach to a fresh segment + } else { + refresh_metadata(); + if (try_consume_newest(consume)) { + return true; + } + } + } + if (std::chrono::steady_clock::now() >= deadline) { + if (_header && writer_gone()) { + detach(); // producer gone -> next call re-attaches (restart recovery) + } + return false; + } + std::this_thread::sleep_for(LATEST_FRAME_READ_POLL); + } + } // Caps block as [4-byte LE len][JSON], consumed verbatim by ShmCaps::parse. const void* get_metadata_raw(); @@ -172,18 +191,65 @@ class LatestFrameReader { private: bool try_attach(); void detach(); - bool try_read_once(LatestFrame& out); void refresh_metadata(); bool writer_gone(); - const LatestFrameSlotHeader* slot_hdr(uint32_t i) const; - const uint8_t* slot_payload(uint32_t i) const; + + LatestFrameSlotHeader* slot_hdr(uint32_t i) { + return reinterpret_cast( + _base + _header->slots_offset + static_cast(i) * _header->slot_stride); + } + uint8_t* slot_payload(uint32_t i) { + return reinterpret_cast(slot_hdr(i)) + sizeof(LatestFrameSlotHeader); + } + + // Seqlock read of the newest published slot, invoking consume zero-copy and + // re-checking the seqlock; retries onto the newest slot on a mid-read lap. + template + bool try_consume_newest(Fn& consume) { + for (int attempt = 0; attempt < LATEST_FRAME_READ_RETRIES; ++attempt) { + int64_t idx = + std::atomic_ref(_header->publish_index).load(std::memory_order_acquire); + if (idx < 0 || idx >= static_cast(_header->slot_count)) { + return false; // nothing published yet + } + uint32_t w = static_cast(idx); + LatestFrameSlotHeader* sh = slot_hdr(w); + + uint64_t s1 = std::atomic_ref(sh->seqlock).load(std::memory_order_acquire); + if (s1 & 1) { + std::this_thread::yield(); // writer mid-write on this slot + continue; + } + uint64_t fn = std::atomic_ref(sh->frame_number).load(std::memory_order_relaxed); + uint64_t sz = std::atomic_ref(sh->size).load(std::memory_order_relaxed); + if (sz > _header->slot_size) { + sz = _header->slot_size; // bounds clamp: never hand out past the slot + } + if (_have_last && fn == _last_sequence) { + return false; // no new frame since the last read + } + + // Zero-copy: hand the slot payload straight to the caller. + consume(slot_payload(w), static_cast(sz), fn); + + std::atomic_thread_fence(std::memory_order_acquire); + uint64_t s2 = std::atomic_ref(sh->seqlock).load(std::memory_order_acquire); + if (s1 != s2) { + continue; // slot rewritten under us -> redo consume on the newest + } + + _last_sequence = fn; + _have_last = true; + return true; + } + return false; + } std::string _name; std::unique_ptr _shm; LatestFrameSharedHeader* _header = nullptr; uint8_t* _base = nullptr; - std::vector _scratch; // reused frame copy (no per-frame alloc) - std::vector _meta_cache; // cached [u32 len][json] + std::vector _meta_cache; // cached [u32 len][json] (on-change only) size_t _meta_size = 0; uint32_t _meta_seq_seen = 0; bool _meta_valid = false; diff --git a/cpp/src/latest_frame.cpp b/cpp/src/latest_frame.cpp index 1421105..d4c69e1 100644 --- a/cpp/src/latest_frame.cpp +++ b/cpp/src/latest_frame.cpp @@ -1,9 +1,7 @@ #include "zerobuffer/latest_frame.h" #include -#include #include -#include #include // Producer-owned latest-frame primitive (ADR-11). The seqlock lives in each @@ -17,15 +15,13 @@ // on-segment bytes stay a language-neutral layout. Ordering: // - writer: payload store, then seqlock even (release), then publish_index // (release) -> a reader that acquires publish_index sees the payload. -// - reader: publish_index (acquire), seqlock (acquire), copy, acquire fence, -// seqlock re-check. +// - reader: publish_index (acquire), seqlock (acquire), consume zero-copy, +// acquire fence, seqlock re-check. The zero-copy read loop lives in the +// header (templated on the caller's consume functor). namespace zerobuffer { namespace { -constexpr int MAX_READ_RETRIES = 16; -constexpr auto READ_POLL_INTERVAL = std::chrono::microseconds(300); - // A completed metadata write leaves metadata_seq even and >= 2 (0 = never // written, 1 = first write in progress). Readers treat < 2 as "no caps yet". constexpr uint32_t METADATA_WRITTEN = 2; @@ -253,8 +249,8 @@ LatestFrameReader::~LatestFrameReader() { LatestFrameReader::LatestFrameReader(LatestFrameReader&& o) noexcept : _name(std::move(o._name)), _shm(std::move(o._shm)), _header(o._header), _base(o._base), - _scratch(std::move(o._scratch)), _meta_cache(std::move(o._meta_cache)), - _meta_size(o._meta_size), _meta_seq_seen(o._meta_seq_seen), _meta_valid(o._meta_valid), + _meta_cache(std::move(o._meta_cache)), _meta_size(o._meta_size), + _meta_seq_seen(o._meta_seq_seen), _meta_valid(o._meta_valid), _last_sequence(o._last_sequence), _have_last(o._have_last) { o._header = nullptr; o._base = nullptr; @@ -267,7 +263,6 @@ LatestFrameReader& LatestFrameReader::operator=(LatestFrameReader&& o) noexcept _shm = std::move(o._shm); _header = o._header; _base = o._base; - _scratch = std::move(o._scratch); _meta_cache = std::move(o._meta_cache); _meta_size = o._meta_size; _meta_seq_seen = o._meta_seq_seen; @@ -309,7 +304,6 @@ bool LatestFrameReader::try_attach() { _shm = std::move(shm); _header = h; _base = static_cast(_shm->data()); - _scratch.assign(static_cast(_header->slot_size), 0); _meta_valid = false; _meta_size = 0; _meta_seq_seen = 0; @@ -326,92 +320,6 @@ void LatestFrameReader::detach() { _have_last = false; } -const LatestFrameSlotHeader* LatestFrameReader::slot_hdr(uint32_t i) const { - return reinterpret_cast( - _base + _header->slots_offset + static_cast(i) * _header->slot_stride); -} - -const uint8_t* LatestFrameReader::slot_payload(uint32_t i) const { - return reinterpret_cast(slot_hdr(i)) + sizeof(LatestFrameSlotHeader); -} - -bool LatestFrameReader::try_read_once(LatestFrame& out) { - for (int attempt = 0; attempt < MAX_READ_RETRIES; ++attempt) { - int64_t idx = - std::atomic_ref(_header->publish_index).load(std::memory_order_acquire); - if (idx < 0 || idx >= static_cast(_header->slot_count)) { - return false; // nothing published yet - } - uint32_t w = static_cast(idx); - const LatestFrameSlotHeader* sh = slot_hdr(w); - uint64_t& seq_field = const_cast(sh->seqlock); - - uint64_t s1 = std::atomic_ref(seq_field).load(std::memory_order_acquire); - if (s1 & 1) { - std::this_thread::yield(); // writer mid-write on this slot - continue; - } - uint64_t fn = std::atomic_ref(const_cast(sh->frame_number)) - .load(std::memory_order_relaxed); - uint64_t sz = std::atomic_ref(const_cast(sh->size)) - .load(std::memory_order_relaxed); - if (sz > _scratch.size()) { - sz = _scratch.size(); // defensive: never read past the slot - } - // Plain payload copy, validated by the seqlock re-check below: a byte - // race here is torn iff the seqlock changed, which the re-check catches. - std::memcpy(_scratch.data(), slot_payload(w), sz); - std::atomic_thread_fence(std::memory_order_acquire); - - uint64_t s2 = std::atomic_ref(seq_field).load(std::memory_order_acquire); - if (s1 != s2) { - continue; // slot rewritten under us -> retry, lands on the newest - } - - if (_have_last && fn == _last_sequence) { - return false; // no new frame since the last read - } - _last_sequence = fn; - _have_last = true; - out._data = _scratch.data(); - out._size = static_cast(sz); - out._sequence = fn; - out._valid = true; - return true; - } - return false; -} - -LatestFrame LatestFrameReader::read_latest(std::chrono::milliseconds timeout) { - auto deadline = std::chrono::steady_clock::now() + timeout; - for (;;) { - if (!_header) { - try_attach(); - } - if (_header) { - if (std::atomic_ref(_header->magic).load(std::memory_order_acquire) != - LATEST_FRAME_MAGIC) { - detach(); // producer closed / restarted -> re-attach to a fresh segment - } else { - refresh_metadata(); - LatestFrame f; - if (try_read_once(f)) { - return f; - } - } - } - if (std::chrono::steady_clock::now() >= deadline) { - // Producer gone: drop the mapping so the next call re-attaches to a - // freshly created segment (writer restart recovery, ADR-5). - if (_header && writer_gone()) { - detach(); - } - return LatestFrame{}; - } - std::this_thread::sleep_for(READ_POLL_INTERVAL); - } -} - void LatestFrameReader::refresh_metadata() { if (!_header) { return; diff --git a/cpp/tests/test_latest_frame.cpp b/cpp/tests/test_latest_frame.cpp index b40369e..28533fa 100644 --- a/cpp/tests/test_latest_frame.cpp +++ b/cpp/tests/test_latest_frame.cpp @@ -6,6 +6,7 @@ #include #include #include +#include #include #include #include @@ -36,6 +37,36 @@ bool all_bytes_equal(const uint8_t* p, size_t size, uint8_t value) { return true; } +// Test-side capture of one zero-copy read: records the slot pointer, size and +// sequence, and (optionally) copies the bytes out FOR ASSERTIONS ONLY — the +// primitive itself never copies. +struct Captured { + bool valid = false; + const void* ptr = nullptr; + size_t size = 0; + uint64_t sequence = 0; + std::vector bytes; +}; + +Captured read_capture(LatestFrameReader& r, std::chrono::milliseconds t, bool copy_bytes = true) { + Captured c; + c.valid = r.read_latest_into( + [&](const uint8_t* src, size_t size, uint64_t seq) { + c.ptr = src; + c.size = size; + c.sequence = seq; + if (copy_bytes) c.bytes.assign(src, src + size); + }, + t); + return c; +} + +void publish_pattern(LatestFrameWriter& w, uint64_t seq, size_t size) { + uint8_t* slot = w.get_slot(); + fill_pattern(slot, size, seq); + w.publish(seq, size); +} + } // namespace // FR-7: the writer creates and publishes with NO reader present (the ADR-1 @@ -43,7 +74,6 @@ bool all_bytes_equal(const uint8_t* p, size_t size, uint8_t value) { TEST(LatestFrameTest, WriterCreatesAndPublishesWithNoReader) { std::string name = make_name("no_reader"); LatestFrameWriter writer(name, 4096, 3); - for (uint64_t seq = 0; seq < 10; ++seq) { uint8_t* slot = writer.get_slot(); ASSERT_NE(slot, nullptr); @@ -53,7 +83,6 @@ TEST(LatestFrameTest, WriterCreatesAndPublishesWithNoReader) { SUCCEED(); } -// Slot count below the triple-slot minimum is raised, not accepted verbatim. TEST(LatestFrameTest, SlotCountRaisedToMinimum) { std::string name = make_name("min_slots"); LatestFrameWriter writer(name, 1024, 1); @@ -65,22 +94,18 @@ TEST(LatestFrameTest, ZeroSlotSizeRejected) { EXPECT_THROW(LatestFrameWriter(name, 0, 3), ZeroBufferException); } -// Reader attaching AFTER the writer gets the latest published frame. +// Reader attaching AFTER the writer gets the latest published frame (zero-copy). TEST(LatestFrameTest, ReaderAttachesAfterWriterGetsLatest) { std::string name = make_name("attach_after"); LatestFrameWriter writer(name, 4096, 3); - for (uint64_t seq = 0; seq <= 5; ++seq) { - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 4096, seq); - writer.publish(seq, 4096); - } + for (uint64_t seq = 0; seq <= 5; ++seq) publish_pattern(writer, seq, 4096); LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); - EXPECT_EQ(f.sequence(), 5u); - EXPECT_EQ(f.size(), 4096u); - EXPECT_TRUE(all_bytes_equal(static_cast(f.data()), f.size(), 5 & 0xFF)); + Captured f = read_capture(reader, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_EQ(f.sequence, 5u); + EXPECT_EQ(f.size, 4096u); + EXPECT_TRUE(all_bytes_equal(f.bytes.data(), f.bytes.size(), 5 & 0xFF)); } // Newest-frame-wins: under a burst, a late reader skips intermediates. @@ -88,29 +113,50 @@ TEST(LatestFrameTest, NewestWinsSkipsIntermediates) { std::string name = make_name("newest_wins"); LatestFrameWriter writer(name, 1024, 3); LatestFrameReader reader(name); + for (uint64_t seq = 0; seq < 100; ++seq) publish_pattern(writer, seq, 1024); - for (uint64_t seq = 0; seq < 100; ++seq) { - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 1024, seq); - writer.publish(seq, 1024); - } + Captured f = read_capture(reader, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_EQ(f.sequence, 99u); // landed on the newest, not seq 0 +} - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); - EXPECT_EQ(f.sequence(), 99u); // landed on the newest, not seq 0 +// Zero-copy: the pointer handed to consume moves through the slot ring (proving +// no per-frame copy into a single reused buffer) and is bounded by slot_count +// (proving no per-frame allocation). +TEST(LatestFrameTest, ZeroCopyPointerRotatesThroughSlots) { + std::string name = make_name("zerocopy"); + LatestFrameWriter writer(name, 4096, 3); + LatestFrameReader reader(name); + + std::set ptrs; + for (uint64_t seq = 0; seq < 20; ++seq) { + publish_pattern(writer, seq, 4096); + Captured f = read_capture(reader, 1000ms, /*copy_bytes=*/false); + ASSERT_TRUE(f.valid); + ASSERT_NE(f.ptr, nullptr); + EXPECT_EQ(f.sequence, seq); + ptrs.insert(f.ptr); + } + // >1 distinct pointer: not a single reused scratch buffer (that would be 1). + EXPECT_GT(ptrs.size(), 1u); + // <= slot_count: a fixed ring, no per-frame allocation. + EXPECT_LE(ptrs.size(), writer.slot_count()); } -// Segment absent for the reader -> read_latest returns invalid (poll/retry), -// never crashes or blocks past the timeout. +// Segment absent for the reader -> read returns false (poll/retry), consume is +// never invoked, never crashes or blocks past the timeout. TEST(LatestFrameTest, ReaderToleratesAbsentSegment) { std::string name = make_name("absent"); LatestFrameReader reader(name); // nothing created it EXPECT_FALSE(reader.is_attached()); + bool consumed = false; auto start = std::chrono::steady_clock::now(); - LatestFrame f = reader.read_latest(100ms); + bool ok = reader.read_latest_into([&](const uint8_t*, size_t, uint64_t) { consumed = true; }, + 100ms); auto elapsed = std::chrono::steady_clock::now() - start; - EXPECT_FALSE(f.valid()); + EXPECT_FALSE(ok); + EXPECT_FALSE(consumed); EXPECT_GE(elapsed, 90ms); EXPECT_LT(elapsed, 2000ms); } @@ -122,28 +168,21 @@ TEST(LatestFrameTest, ReaderDetachReattachWriterUnaffected) { { LatestFrameReader reader(name); // attach before any frame, then detach - EXPECT_FALSE(reader.read_latest(50ms).valid()); + EXPECT_FALSE(read_capture(reader, 50ms).valid); } uint64_t seq = 0; - auto publish_one = [&](uint64_t s) { - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 1024, s); - writer.publish(s, 1024); - }; - publish_one(seq++); + publish_pattern(writer, seq++, 1024); { LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); + EXPECT_TRUE(read_capture(reader, 1000ms).valid); } - // Writer keeps publishing after the reader detached. - for (int i = 0; i < 20; ++i) publish_one(seq++); + for (int i = 0; i < 20; ++i) publish_pattern(writer, seq++, 1024); LatestFrameReader reader2(name); - LatestFrame f = reader2.read_latest(1000ms); - ASSERT_TRUE(f.valid()); - EXPECT_EQ(f.sequence(), seq - 1); + Captured f = read_capture(reader2, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_EQ(f.sequence, seq - 1); } // Writer restart (segment recreated, sequence reset) while a reader polls -> @@ -151,62 +190,46 @@ TEST(LatestFrameTest, ReaderDetachReattachWriterUnaffected) { TEST(LatestFrameTest, WriterRestartReaderRecovers) { std::string name = make_name("restart"); auto writer = std::make_unique(name, 1024, 3); - for (uint64_t seq = 0; seq < 50; ++seq) { - uint8_t* slot = writer->get_slot(); - fill_pattern(slot, 1024, seq); - writer->publish(seq, 1024); - } + for (uint64_t seq = 0; seq < 50; ++seq) publish_pattern(*writer, seq, 1024); LatestFrameReader reader(name); - LatestFrame f1 = reader.read_latest(1000ms); - ASSERT_TRUE(f1.valid()); - EXPECT_EQ(f1.sequence(), 49u); + Captured f1 = read_capture(reader, 1000ms); + ASSERT_TRUE(f1.valid); + EXPECT_EQ(f1.sequence, 49u); - // Restart: destroy the writer (unlinks) and create a fresh one (seq from 0). - writer.reset(); - writer = std::make_unique(name, 1024, 3); + writer.reset(); // unlink (invalidates magic) + writer = std::make_unique(name, 1024, 3); // seq resets from 0 - // Reader recovers and reads the new segment's frames. bool recovered = false; for (int i = 0; i < 100 && !recovered; ++i) { - uint8_t* slot = writer->get_slot(); - fill_pattern(slot, 1024, static_cast(i)); - writer->publish(static_cast(i), 1024); - LatestFrame f = reader.read_latest(200ms); - if (f.valid()) { - EXPECT_TRUE(all_bytes_equal(static_cast(f.data()), f.size(), - f.sequence() & 0xFF)); + publish_pattern(*writer, static_cast(i), 1024); + Captured f = read_capture(reader, 200ms); + if (f.valid) { + EXPECT_TRUE(all_bytes_equal(f.bytes.data(), f.bytes.size(), f.sequence & 0xFF)); recovered = true; } } EXPECT_TRUE(recovered); } -// Writer gone / stale heartbeat -> reader reports producer-gone, no crash/block. +// Writer gone -> reader reports producer-gone, no crash/block. TEST(LatestFrameTest, WriterGoneReportedNoCrash) { std::string name = make_name("gone"); { LatestFrameWriter writer(name, 1024, 3); - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 1024, 1); - writer.publish(1, 1024); + publish_pattern(writer, 1, 1024); LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); EXPECT_TRUE(reader.is_writer_alive()); - // Writer destroyed at end of scope, before the reader below. - LatestFrameReader probe(name); - EXPECT_TRUE(probe.is_writer_alive()); } // Segment is now unlinked; a fresh reader cannot attach and reports gone. LatestFrameReader reader(name); EXPECT_FALSE(reader.is_writer_alive()); - LatestFrame f = reader.read_latest(100ms); - EXPECT_FALSE(f.valid()); + EXPECT_FALSE(read_capture(reader, 100ms).valid); } -// Oversized frame vs slot_size -> handled without out-of-bounds; the reader +// Oversized frame vs slot_size -> handled without out-of-bounds; the consume // never sees more than slot_size bytes. TEST(LatestFrameTest, OversizedFrameClampedNoOverflow) { std::string name = make_name("oversize"); @@ -216,9 +239,9 @@ TEST(LatestFrameTest, OversizedFrameClampedNoOverflow) { writer.publish(7, 100000); // claim way more than slot_size LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); - EXPECT_LE(f.size(), 1024u); // clamped to slot_size, no OOB + Captured f = read_capture(reader, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_LE(f.size, 1024u); // clamped to slot_size, no OOB } // Metadata/caps set + rewritten on change -> reader re-reads (FR-3). Verifies @@ -229,13 +252,10 @@ TEST(LatestFrameTest, MetadataSetAndRewrite) { std::string caps1 = R"({"caps":"video/x-raw,format=GRAY8,width=640,height=480"})"; writer.set_metadata(caps1.data(), caps1.size()); - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 1024, 1); - writer.publish(1, 1024); + publish_pattern(writer, 1, 1024); LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); const uint8_t* md = static_cast(reader.get_metadata_raw()); size_t md_size = reader.get_metadata_size(); @@ -246,15 +266,11 @@ TEST(LatestFrameTest, MetadataSetAndRewrite) { EXPECT_EQ(prefix, caps1.size()); EXPECT_EQ(0, std::memcmp(md + 4, caps1.data(), caps1.size())); - // Rewrite caps on a change; the reader re-reads on the next poll. std::string caps2 = R"({"caps":"video/x-raw,format=I420,width=1920,height=1080"})"; writer.set_metadata(caps2.data(), caps2.size()); - slot = writer.get_slot(); - fill_pattern(slot, 1024, 2); - writer.publish(2, 1024); + publish_pattern(writer, 2, 1024); - f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); md = static_cast(reader.get_metadata_raw()); md_size = reader.get_metadata_size(); ASSERT_EQ(md_size, 4 + caps2.size()); @@ -266,75 +282,87 @@ TEST(LatestFrameTest, MetadataSetAndRewrite) { TEST(LatestFrameTest, MetadataNotReadyBeforeFirstSet) { std::string name = make_name("meta_notready"); LatestFrameWriter writer(name, 1024, 3); - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 1024, 1); - writer.publish(1, 1024); // published, but no set_metadata yet + publish_pattern(writer, 1, 1024); // published, but no set_metadata yet LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); EXPECT_EQ(reader.get_metadata_raw(), nullptr); EXPECT_EQ(reader.get_metadata_size(), 0u); } -// Tear-free read under concurrent write: a reader loop vs a fast writer over -// many iterations must never return a torn (mixed-lap) frame. +// Tear-free under concurrent write: a reader in a tight read_latest_into loop vs +// a fast writer over many iterations. The COMMITTED frame (read returns true) is +// never torn; the seqlock retry path is exercised (the reader observes and +// rejects mid-lap slots). Zero-copy: consume validates the slot in place. TEST(LatestFrameTest, TearFreeUnderConcurrentWrite) { std::string name = make_name("tearfree"); constexpr size_t FRAME = 256 * 1024; // large payload widens the tear window constexpr uint64_t ITERATIONS = 20000; - LatestFrameWriter writer(name, FRAME, 3); + LatestFrameWriter writer(name, FRAME, 3); // min slots -> tightest safe window std::atomic stop{false}; - std::atomic torn{0}; + std::atomic committed_torn{0}; + std::atomic observed_torn{0}; // mid-lap reads caught by the seqlock std::atomic reads{0}; + std::atomic reader_ready{false}; std::thread reader_thread([&]() { LatestFrameReader reader(name); + reader_ready.store(true, std::memory_order_release); while (!stop.load(std::memory_order_relaxed)) { - LatestFrame f = reader.read_latest(50ms); - if (!f.valid()) continue; - const uint8_t* p = static_cast(f.data()); - uint8_t expected = static_cast(f.sequence() & 0xFF); - if (!all_bytes_equal(p, f.size(), expected)) { - torn.fetch_add(1, std::memory_order_relaxed); - } + bool consistent = true; + bool ok = reader.read_latest_into( + [&](const uint8_t* src, size_t size, uint64_t seq) { + uint8_t expected = static_cast(seq & 0xFF); + bool c = all_bytes_equal(src, size, expected); + // Widen the read window so a lapping writer can overwrite this + // slot mid-read, then re-scan: this deterministically exercises + // the seqlock retry path over many iterations. + volatile int sink = 0; + for (int spin = 0; spin < 4000; ++spin) { + sink = sink + 1; + } + (void)sink; + if (c && !all_bytes_equal(src, size, expected)) c = false; + if (!c) observed_torn.fetch_add(1, std::memory_order_relaxed); + consistent = c; + }, + 50ms); + if (!ok) continue; + if (!consistent) committed_torn.fetch_add(1, std::memory_order_relaxed); reads.fetch_add(1, std::memory_order_relaxed); } }); - for (uint64_t seq = 0; seq < ITERATIONS; ++seq) { - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, FRAME, seq); - writer.publish(seq, FRAME); + while (!reader_ready.load(std::memory_order_acquire)) { + std::this_thread::yield(); } - // Let the reader drain a little, then stop. + for (uint64_t seq = 0; seq < ITERATIONS; ++seq) publish_pattern(writer, seq, FRAME); std::this_thread::sleep_for(50ms); stop.store(true, std::memory_order_relaxed); reader_thread.join(); - EXPECT_EQ(torn.load(), 0u) << "torn frames observed"; + EXPECT_EQ(committed_torn.load(), 0u) << "a torn frame was committed"; EXPECT_GT(reads.load(), 0u) << "reader never observed a frame"; + EXPECT_GT(observed_torn.load(), 0u) << "seqlock retry path was never exercised"; } -// Two readers on one segment both read tear-free (1..N readers, FR). +// Two readers on one segment both read the newest tear-free (1..N readers). TEST(LatestFrameTest, MultipleReaders) { std::string name = make_name("multi_reader"); LatestFrameWriter writer(name, 4096, 4); - for (uint64_t seq = 0; seq <= 3; ++seq) { - uint8_t* slot = writer.get_slot(); - fill_pattern(slot, 4096, seq); - writer.publish(seq, 4096); - } + for (uint64_t seq = 0; seq <= 3; ++seq) publish_pattern(writer, seq, 4096); LatestFrameReader r1(name); LatestFrameReader r2(name); - LatestFrame f1 = r1.read_latest(1000ms); - LatestFrame f2 = r2.read_latest(1000ms); - ASSERT_TRUE(f1.valid()); - ASSERT_TRUE(f2.valid()); - EXPECT_EQ(f1.sequence(), 3u); - EXPECT_EQ(f2.sequence(), 3u); + Captured f1 = read_capture(r1, 1000ms); + Captured f2 = read_capture(r2, 1000ms); + ASSERT_TRUE(f1.valid); + ASSERT_TRUE(f2.valid); + EXPECT_EQ(f1.sequence, 3u); + EXPECT_EQ(f2.sequence, 3u); + EXPECT_TRUE(all_bytes_equal(f1.bytes.data(), f1.bytes.size(), 3)); + EXPECT_TRUE(all_bytes_equal(f2.bytes.data(), f2.bytes.size(), 3)); } // A live writer owning the segment cannot be displaced by a second writer. @@ -349,10 +377,6 @@ TEST(LatestFrameTest, SecondWriterOnLiveSegmentRejected) { // so LatestFrameWriter::create() hits EEXIST and must reclaim it. TEST(LatestFrameTest, StaleSegmentReclaimedOnCreate) { std::string name = make_name("stale"); - - // Hand-craft a leftover segment: valid magic, but writer_pid = 0 (which the - // reclaim path treats as no live owner). It is NOT unlinked, mimicking a - // writer that crashed without cleanup. { auto shm = SharedMemory::create(name, 65536); auto* h = static_cast(shm->data()); @@ -370,14 +394,11 @@ TEST(LatestFrameTest, StaleSegmentReclaimedOnCreate) { // shm mapping dropped here WITHOUT shm_unlink -> the named segment persists. } - // A new writer must reclaim the stale name and create cleanly. LatestFrameWriter fresh(name, 1024, 3); - uint8_t* slot = fresh.get_slot(); - fill_pattern(slot, 1024, 9); - fresh.publish(9, 1024); + publish_pattern(fresh, 9, 1024); LatestFrameReader reader(name); - LatestFrame f = reader.read_latest(1000ms); - ASSERT_TRUE(f.valid()); - EXPECT_EQ(f.sequence(), 9u); + Captured f = read_capture(reader, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_EQ(f.sequence, 9u); } From fb72cf3dfcce63b9ec043c6a90ff347489809c5b Mon Sep 17 00:00:00 2001 From: Daniel Kowalski Date: Thu, 9 Jul 2026 11:30:09 +0200 Subject: [PATCH 5/5] docs: spell out zero-copy consume contract in header --- cpp/include/zerobuffer/latest_frame.h | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/cpp/include/zerobuffer/latest_frame.h b/cpp/include/zerobuffer/latest_frame.h index 5f423ca..40403ab 100644 --- a/cpp/include/zerobuffer/latest_frame.h +++ b/cpp/include/zerobuffer/latest_frame.h @@ -143,13 +143,17 @@ class LatestFrameReader { // Zero-copy tear-free read of the newest published frame. Invokes // consume(const uint8_t* src, size_t size, uint64_t sequence) // with `src` pointing DIRECTLY into the published slot (no copy). The seqlock - // is validated AROUND the consume: if the writer laps the slot mid-read, - // consume is re-invoked on the newest slot (bounded retries), so the caller - // must treat its work as committed only when this returns true. Returns false - // if the segment is absent, no NEW frame appears within `timeout`, or all - // retries tore. Intermediate frames the writer overwrote are drops (sequence - // gap). `size` is clamped to slot_size. The primitive stays generic — it - // never copies and knows nothing about the caller's consume. + // is validated AROUND the consume, so `consume` must obey a stricter contract + // than a copy API: + // 1. it MAY be handed torn bytes (a writer lap mid-read) — do not trust the + // pixel values until this function returns true; + // 2. it MAY be invoked more than once per call (one per retry) — keep the + // work restartable/idempotent (e.g. repack into the same target slot); + // 3. commit / present only when this returns true. + // Returns false if the segment is absent, no NEW frame appears within + // `timeout`, or all retries tore. Intermediate frames the writer overwrote + // are drops (sequence gap). `size` is clamped to slot_size. The primitive + // stays generic — it never copies and knows nothing about the caller's consume. template bool read_latest_into(Fn&& consume, std::chrono::milliseconds timeout) { auto deadline = std::chrono::steady_clock::now() + timeout;