diff --git a/.gitignore b/.gitignore index 44cfa65..27b7c58 100644 --- a/.gitignore +++ b/.gitignore @@ -42,6 +42,8 @@ Testing/ /cpp/tests/*_test /cpp/benchmarks/test_* /cpp/benchmarks/*_test +# ...but keep test SOURCES tracked (patterns above target built binaries) +!/cpp/tests/test_*.cpp # C# build artifacts /csharp/.vs/ @@ -126,5 +128,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..40403ab --- /dev/null +++ b/cpp/include/zerobuffer/latest_frame.h @@ -0,0 +1,266 @@ +#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. 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 { + +// ----- 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 +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. +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, 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; + 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 reserved1[5]; +}; +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"); + +// ----- 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) 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); + + 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; + + // 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, 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; + 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(); + 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(); + void refresh_metadata(); + bool writer_gone(); + + 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 _meta_cache; // cached [u32 len][json] (on-change only) + 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..d4c69e1 --- /dev/null +++ b/cpp/src/latest_frame.cpp @@ -0,0 +1,391 @@ +#include "zerobuffer/latest_frame.h" + +#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), 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 { + +// 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); +} + +// 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(); + + 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); + // 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); + _pending_slot = -1; +} + +// ----------------------------- 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), + _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; + _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()); + _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; +} + +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 < METADATA_WRITTEN) { + return; // no caps written yet -> stays "not available", not empty caps + } + 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..28533fa --- /dev/null +++ b/cpp/tests/test_latest_frame.cpp @@ -0,0 +1,404 @@ +#include + +#include "zerobuffer/latest_frame.h" +#include "zerobuffer/platform.h" + +#include +#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; +} + +// 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 +// 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(); +} + +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 (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) publish_pattern(writer, seq, 4096); + + LatestFrameReader reader(name); + 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. +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); + + Captured f = read_capture(reader, 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 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(); + 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(ok); + EXPECT_FALSE(consumed); + 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(read_capture(reader, 50ms).valid); + } + uint64_t seq = 0; + publish_pattern(writer, seq++, 1024); + + { + LatestFrameReader reader(name); + EXPECT_TRUE(read_capture(reader, 1000ms).valid); + } + for (int i = 0; i < 20; ++i) publish_pattern(writer, seq++, 1024); + + LatestFrameReader reader2(name); + 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 -> +// 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) publish_pattern(*writer, seq, 1024); + + LatestFrameReader reader(name); + Captured f1 = read_capture(reader, 1000ms); + ASSERT_TRUE(f1.valid); + EXPECT_EQ(f1.sequence, 49u); + + writer.reset(); // unlink (invalidates magic) + writer = std::make_unique(name, 1024, 3); // seq resets from 0 + + bool recovered = false; + for (int i = 0; i < 100 && !recovered; ++i) { + 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 -> reader reports producer-gone, no crash/block. +TEST(LatestFrameTest, WriterGoneReportedNoCrash) { + std::string name = make_name("gone"); + { + LatestFrameWriter writer(name, 1024, 3); + publish_pattern(writer, 1, 1024); + + LatestFrameReader reader(name); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); + EXPECT_TRUE(reader.is_writer_alive()); + } + // Segment is now unlinked; a fresh reader cannot attach and reports gone. + LatestFrameReader reader(name); + EXPECT_FALSE(reader.is_writer_alive()); + EXPECT_FALSE(read_capture(reader, 100ms).valid); +} + +// 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"); + 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); + 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 +// 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()); + publish_pattern(writer, 1, 1024); + + LatestFrameReader reader(name); + 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(); + 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())); + + std::string caps2 = R"({"caps":"video/x-raw,format=I420,width=1920,height=1080"})"; + writer.set_metadata(caps2.data(), caps2.size()); + publish_pattern(writer, 2, 1024); + + 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()); + 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); + publish_pattern(writer, 1, 1024); // published, but no set_metadata yet + + LatestFrameReader reader(name); + ASSERT_TRUE(read_capture(reader, 1000ms).valid); + EXPECT_EQ(reader.get_metadata_raw(), nullptr); + EXPECT_EQ(reader.get_metadata_size(), 0u); +} + +// 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); // min slots -> tightest safe window + std::atomic stop{false}; + 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)) { + 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); + } + }); + + while (!reader_ready.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + 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(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 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) publish_pattern(writer, seq, 4096); + + LatestFrameReader r1(name); + LatestFrameReader r2(name); + 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. +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"); + { + 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. + } + + LatestFrameWriter fresh(name, 1024, 3); + publish_pattern(fresh, 9, 1024); + + LatestFrameReader reader(name); + Captured f = read_capture(reader, 1000ms); + ASSERT_TRUE(f.valid); + EXPECT_EQ(f.sequence, 9u); +}