StormByte-Buffer 2.0.0
C++26 buffer module of the StormByte suite
 
Loading...
Searching...
No Matches
StormByte-Buffer

Platform C++26 CMake License: LGPL v3 or commercial CI Sponsor

This repository is StormByte Buffer: FIFO, SharedFIFO, Ring, Producer/Consumer, Hopper, Sink, Bridge, Pumper, pipelines and buffered I/O for the StormByte C++ suite.

It depends on StormByte-String 1.0.0 or newer, which vendors StormByte Base 2.0.0 or newer, StormByte-System 2.0.0 or newer, and optionally StormByte-Logger 2.0.0 or newer for pipeline pipes (Scope). Public headers live under StormByte/buffer/.

The suite is split on purpose. Base, Config, Crypto, Database, Logger, Multimedia, Network, String and System are other repositories. This one does not implement them.

Designed to interconnect

Pieces plug into each other through ReadOnly / WriteOnly and through IO leaves.

  • Producer yields a Consumer over the same Ring.
  • Pipeline transforms stream buffers only (ReadOnly / WriteOnly). Pipes do not take IO leaves.
  • Bridge moves bytes from a ReadOnly or an IO reader into a WriteOnly or an IO writer. IO tips are taken by move; in-memory tips stay referenced.
  • Pumper owns a Bridge and runs Passthrough until EoF or failure.
  • BufferedFileReader is a BufferedReader. BufferedFileWriter is a BufferedWriter. Leaves implement Origin*. Cache, prefetch, backpressure, delayed seek and telemetry live in the bases.

Typical wires:

  • Producer → Consumer (same ring).
  • Producer → Pipeline::Process → Consumer.
  • That Consumer → Bridge::Passthrough → file (or wrap the Bridge in a Pumper).
  • BufferedFileReader / BufferedFileWriter for seekable files with a page map.

See Pipeline, Bridge, Pumper, Telemetry, IO::BufferedReader and IO::BufferedWriter.

What this module does

  • BinaryData — octet payloads are StormByte::BinaryData (Base). Lengths of byte buffers are StormByte::ByteSize. Hopper and Sink count items with StormByte::Size.
  • FIFO — grow-on-demand byte buffer. Not thread-safe. Read / Peek keep data; Extract consumes it.
  • SharedFIFO — thread-safe FIFO. Read / Extract block until data or Close / SetError.
  • Ring — concurrent ring (many-to-many).
  • Producer / Consumer — write-only / read-only handles over a shared Ring.
  • Hopper / Sink — SPSC typed items and a keyed map of hoppers.
  • Pipeline — user leaves of Pipe. Stream buffers only. Add clones or moves.
  • Bridge — manual transfer. Passthrough(n, Operation) only. No worker.
  • Pumper — owns a Bridge and pumps until EoF or Cancel.
  • Telemetry — ReadTelemetry / WriteTelemetry as const StormByte::Shared<…>. MeanRate is caller rate, not disk rate.
  • IO — BufferedReader / BufferedWriter bases and file leaves. Nested Parameters and knobs.
  • Lifecycle — Close(), SetError(), EoF(), IsReadable(), IsWritable().

The rest of the suite

Module Role API
Base Exceptions, Expected, serialization, UUID, concepts /StormByte
Buffer This repository /StormByte-Buffer
Config Human-readable text and versioned binary documents /StormByte-Config
Crypto Hash, compress, encrypt, sign — Crypto++ stays private /StormByte-Crypto
Database One API over SQLite, PostgreSQL and MariaDB /StormByte-Database
Logger Stream logger with levels, headers, components and Scope /StormByte-Logger
Multimedia Decode, encode and containers without raw FFmpeg types /StormByte-Multimedia
Network Framed packets, Client/Server, IPv4/IPv6 TCP /StormByte-Network
String Owned UTF-8 / wide text that can cross a DLL boundary /StormByte-String
System Processes, pipes, Device, host and environment /StormByte-System

Table of Contents

Documentation

Installation

Needs a C++26 compiler, CMake 3.28 or newer, StormByte-String 1.0.0 or newer (vendors StormByte Base 2.0.0), StormByte-System 2.0.0 or newer, and optionally StormByte-Logger 2.0.0 when pipeline pipes take a logger.

git clone --recursive https://github.com/StormBytePP/StormByte-Buffer.git
cd StormByte-Buffer
cmake -S . -B build
cmake --build build

Shared vs static follows CMake BUILD_SHARED_LIBS (declared in lib/, default ON). A plain configure builds the shared library. -DBUILD_SHARED_LIBS=OFF builds a static archive; on Windows the headers then do not use dllimport. Vendored StormByte pins follow the same mode and are configured with ENABLE_TEST=OFF.

A shared build keeps this library as its own .so / .dll. Under the LGPL that is usually the simpler way to ship: the user can replace that file. A static archive is folded into your binary. The LGPL still applies to this code; you must give the recipient a way to relink your product with a different build of this library. If that does not fit how you distribute the final product, a commercial license is available from the copyright holder (see License).

Link StormByte-Buffer. Headers: #include <StormByte/buffer/….hxx>.

Usage

Namespace root is StormByte::Buffer. I/O types live in StormByte::Buffer::IO. Octet payloads use StormByte::BinaryData. Byte lengths use StormByte::ByteSize.

FIFO

int main() {
FIFO fifo;
fifo.Write("Hello World");
BinaryData data;
fifo.Read(5, data); // "Hello", still in the buffer
fifo.Seek(6, Position::Absolute);
BinaryData extracted;
fifo.Extract(5, extracted); // "World"
}
Byte-oriented FIFO buffer with grow-on-demand and close/error support.
Definition fifo.hxx:78
bool Write(const StormByte::ByteSize &count, const StormByte::BinaryData &data) noexcept override
Append bytes from a StormByte::BinaryData (copy).
Position
Forward declaration of the WriteOnly interface.
Definition typedefs.hxx:73

FIFO is not thread-safe. Concurrent ends use SharedFIFO or Ring.

Producer and Consumer

A Producer yields a Consumer over the same ring. The Consumer is a ReadOnly; the Producer is a WriteOnly. Either tip can be passed to a Bridge.

#include <thread>
int main() {
Producer producer;
Consumer consumer = producer.Consumer();
std::thread writer([producer]() mutable {
producer.Write("Data chunk 1");
producer.Write("Data chunk 2");
producer.Close();
});
std::thread reader([consumer]() mutable {
while (!consumer.EoF()) {
BinaryData data;
if (consumer.Extract(0, data) && !data.empty()) {
// process
}
}
});
writer.join();
reader.join();
}
Read-oriented handle over a shared Ring.
Definition consumer.hxx:84
Write-only handle over a shared Ring.
Definition producer.hxx:66
class Consumer Consumer()
Consumer that shares this Producer’s Ring.
Definition producer.hxx:262

Hopper

Hopper<T> is an SPSC queue for StormByte::Type::MoveConstructible items. Capacity 0 is unbounded. Push blocks when a bounded hopper is full. Eof() ends production. Smart-pointer items discard nulls.

Notify(cv) stores a pointer the Hopper does not own. Call Unnotify() before that CV dies.

#include <memory>
#include <thread>
int main() {
Hopper<std::unique_ptr<int>> hopper(5);
std::thread producer([&hopper]() {
for (int i = 0; i < 10; ++i)
hopper << std::make_unique<int>(i);
hopper.Eof();
});
std::thread consumer([&hopper]() {
while (!hopper.Empty() || !hopper.EoF()) {
auto item = hopper.Pop();
(void)item;
}
});
producer.join();
consumer.join();
}
Single-producer single-consumer (SPSC) typed item queue.
Definition hopper.hxx:76

Sink

Sink<T> maps integer keys to Hopper<T> buckets. Wire with To(key) / >> / <<.

#include <memory>
#include <string>
#include <thread>
int main() {
Sink<std::shared_ptr<std::string>> producer_sink;
Sink<std::shared_ptr<std::string>> consumer_sink;
producer_sink.To(1, consumer_sink);
producer_sink.To(2, consumer_sink);
std::thread writer([&producer_sink]() {
producer_sink.Push(1, std::make_shared<std::string>("ch1"));
producer_sink.Push(2, std::make_shared<std::string>("ch2"));
producer_sink.Eof();
});
std::thread reader([&consumer_sink]() {
while (!consumer_sink.EoF())
(void)consumer_sink.Pop();
});
writer.join();
reader.join();
}
Set of Hopper buckets keyed by an integer.
Definition sink.hxx:79
Lane To(int key) noexcept
Redirect of one hopper key.

Parameters

IO leaves and Pumper take a nested Parameters object. Omitted knobs keep the office default. Brace-init and a named bag are both valid. The variadic list is applied in your TU (STORMBYTE_FORCE_INLINE); the DLL sees numbers.

BufferedFileReader in("in.bin"); // probe at Open
BufferedFileReader in2("in.bin", { ReadAhead{1 << 20} }); // one knob
BufferedFileReader::Parameters p{ ReadAhead{1 << 20}, MaxMemory{8 << 20} };
BufferedFileReader in3("in.bin", p);
Final BufferedLocationReader over a filesystem file.
Definition buffered_file_reader.hxx:98
Cache / dirty-page budget knob.
Definition parameters.hxx:102
Prefetch length knob.
Definition parameters.hxx:77

Explicit 0 is 0. It does not probe.

Telemetry

Every IO office and every Bridge exposes const StormByte::Shared<ReadTelemetry> / WriteTelemetry. The handle is the same object for the life of the office. IO types add cache / origin / seek / wait counters. Non-IO Bridge tips use the basic type.

MeanRate is octets per second of requested user operations, including cache hits. It is not a disk benchmark. A cached write can look like GiB/s. Worker, GC and internal flushes enter the rate only when they delay the caller. Explicit Flush / Close Flush pull it back.

Flatten with operator StormByte::String::String or operator std::string() (the latter is FORCE_INLINE so the std::string lives in your TU):

auto tel = reader.Telemetry();
if (tel)
log << Level::Info << *tel << std::endl;

Pipeline

Pipeline transforms stream buffers (ReadOnly / WriteOnly). It does not take IO leaves. A file or device is attached later with a Bridge.

A pipe is a user leaf of Pipe. Implement Run, Clone and Move. Add(const Pipe&) clones onto Base's heap and does not touch the caller object. Add(Pipe&&) takes Move().

Pipe::Run(ReadOnly&, WriteOnly&, const Shared<Logger::Log>&). Close or SetError the sink before return. Pipeline::Process(buffer, log, mode) — mode last. When a logger is set, each pipe receives a scoped handle.

class StripCrPipe final: public Pipe {
public:
StripCrPipe() = default;
StripCrPipe(const StripCrPipe&) = default;
StripCrPipe(StripCrPipe&&) noexcept = default;
~StripCrPipe() noexcept override = default;
StripCrPipe& operator=(const StripCrPipe&) = default;
StripCrPipe& operator=(StripCrPipe&&) noexcept = default;
void Run(ReadOnly& in, WriteOnly& out,
const StormByte::Shared<StormByte::Logger::Log>&) override {
in.Extract(0, raw);
StormByte::BinaryData unix_newlines;
unix_newlines.reserve(raw.size());
for (const std::byte octet : raw) {
if (octet != std::byte{'\r'})
unix_newlines.push_back(octet);
}
out.Write(unix_newlines.size(), std::move(unix_newlines));
out.Close();
}
PointerType Clone() const noexcept override {
}
PointerType Move() noexcept override {
}
};
int main() {
Producer src;
src.Write("a\r\nb\r\n");
src.Close();
StripCrPipe crlf;
Pipeline pipe;
pipe.Add(crlf);
Consumer result = pipe.Process(src.Consumer(), {}, ExecutionMode::Sync);
(void)result;
}
void reserve(const StormByte::ByteSize &new_cap)
void push_back(std::byte value)
StormByte::ByteSize size() const noexcept
One transformation in a Pipeline.
Definition pipe.hxx:68
Multi-pipe byte transformation over stream buffers.
Definition pipeline.hxx:92
Buffer that can be read and not written.
Definition generic.hxx:228
Buffer that can be written and not read.
Definition generic.hxx:505
ExecutionMode
Bitmask controlling how Pipeline::Process schedules work.
Definition typedefs.hxx:105
Root namespace of the StormByte C++ suite.

Sync runs on the caller thread. Async returns at once. Parallel is one thread per pipe. Combine with |.

Bridge

Bridge is a manual transfer. Bytes move only when you call Bridge::Passthrough. There is no worker and no occupancy cap here. Continuous transfer is Pumper.

  • In-memory tips: ReadOnly& / WriteOnly&. Those buffers must outlive the Bridge.
  • IO tips: stolen by move as the concrete leaf (BufferedFileReader, BufferedFileWriter, …).
  • Passthrough(n, Operation) is atomic. Write TryAgain is retried until that call completes.
  • Operation applies to the read tip only: Blocking waits for n or EoF; NonBlocking takes what is available now, up to n.
  • n == 0 is current contents (Available() on IO does not touch the origin).
  • State is Open, Closed or Failed. Failed() is only a real tip fault. A moved-from Bridge is Closed, not Failed.
  • Close() ends the session. Owned adapters and stolen IO leaves are released. After Close, a consumed source EoF, or a move-from, the instance is Closed and cannot be re-armed. Construct a new Bridge to transfer again. Close() is idempotent and does not set Failed.
  • A NonBlocking Passthrough that returns 0 is not the end of the session.
  • Telemetry: the Bridge always keeps a Shared copy of the read and write counters. Non-IO tips use a basic telemetry owned and updated by the Bridge. IO tips donate the leaf handle at attach; the leaf updates that object. After Close the same handles remain valid. A Shared you already copied stays valid when the Bridge dies.

On Windows an IO leaf keeps the origin handle until that leaf is destroyed. While the Bridge still owns a BufferedFileReader or a BufferedFileWriter, the path stays locked: DeleteFile / std::filesystem::remove fail. In-memory tips do not lock a file. Call Bridge::Close() when you are done so the stolen IO origin is closed and the path can be unlinked. The destructor does the same work; Close is for when *this must stay alive.

#include <filesystem>
int main() {
FIFO in;
in.Write("payload");
in.Close();
BufferedFileWriter out("out.bin");
out.Open();
Bridge bridge(in, std::move(out));
while (!bridge.EoF() && !bridge.Failed()) {
if (bridge.Passthrough(64 * 1024) == 0 && !bridge.EoF())
break;
}
auto read = bridge.ReadTelemetry();
bridge.Close();
(void)read; // still valid
std::filesystem::remove("out.bin");
}
Private pump for StormByte::Buffer::Bridge.
Manual bridge between two StormByte-Buffer ends.
Definition bridge.hxx:156
Final BufferedLocationWriter over a filesystem file.
Definition buffered_file_writer.hxx:100
virtual Result Write(const FIFO &src) final
Write every unread byte of src.

A Pipeline is not a Bridge tip. Pipeline::Process is. It returns a Consumer, and a Consumer is ReadOnly, so it can be the source of a Bridge. The leaf StripCrPipe is the one from Pipeline.

int main() {
Producer capture;
capture.Write("line 1\r\nline 2\r\n");
capture.Close();
StripCrPipe crlf;
Pipeline pipe;
pipe.Add(crlf);
Consumer normalized = pipe.Process(capture.Consumer(), {}, ExecutionMode::Sync);
BufferedFileWriter out("session.log");
out.Open();
Bridge to_disk(normalized, std::move(out));
while (!to_disk.EoF() && !to_disk.Failed()) {
if (to_disk.Passthrough(4 * 1024) == 0 && !to_disk.EoF())
break;
}
to_disk.Close();
}

session.log holds line 1\nline 2\n. For a long capture, wrap to_disk in a Pumper instead of the Passthrough loop. The file stays locked while that Pumper is alive; ~Pumper joins and then releases the owned Bridge.

Pumper

Pumper takes a Bridge by move and starts a worker immediately. The destructor joins and may block until the current cycle ends. It does not Cancel: it lets the worker finish.

  • Chunk — bytes asked of Passthrough each cycle. 0 is automatic chunking, not Bridge “current contents”.
  • HighWater — input cap only. Omitted: 0 if the source is IO, otherwise a backend default (constexpr in the PIMPL). Explicit 0: no Pumper cap. Use 0 when the IO source already limits itself. Non-IO sources are unbounded by design.
  • Toggle pauses and resumes. No-op if Failed or Canceled.
  • Cancel is terminal (Canceled, no restart). It does not set Failed. It Close()s the owned Bridge so Windows can unlink an IO path.
  • Failed() is only a real Bridge fault. A moved-from Pumper is empty (EoF), not Failed and not Canceled.
  • Telemetry handles are copied from the Bridge at construction. They stay valid after Cancel. A Shared you already copied stays valid when the Pumper dies.

An IO path stays locked while the Pumper is alive. After Cancel (or after ~Pumper joins and destroys the Bridge) the handle is released.

int main() {
BufferedFileReader in("in.bin");
BufferedFileWriter out("out.bin");
in.Open();
out.Open();
Pumper pump(Bridge(std::move(in), std::move(out)), { Chunk{1 << 20} });
while (!pump.EoF() && !pump.Failed() && !pump.Canceled()) {
// work elsewhere; pump runs on its thread
}
auto read = pump.ReadTelemetry();
(void)read; // still valid after ~Pumper
// ~Pumper joins
}
Private worker for StormByte::Buffer::Pumper.
Pumper cycle size.
Definition pumper.hxx:80
Input occupancy cap for Pumper.
Definition pumper.hxx:109
Owns a Bridge and moves bytes until EoF, failure or Cancel.
Definition pumper.hxx:168

finish = !pump.Failed() && !pump.Canceled() && pump.EoF().
active = !pump.Failed() && !pump.Canceled() && !pump.EoF().

IO::BufferedReader

Public base for a binary origin. Leaves implement OriginOpen, OriginClose, OriginPull, OriginCanSeek, OriginSeek, OriginHasSize, OriginSize. Construction is Unavailable; a successful Open is Idle.

Available() is cached bytes at Tell. It does not call the origin.

Seek is logical. A cache hit does not move the device. Tell never lies.

IO::BufferedWriter

Public base for a binary sink. Leaves implement OriginOpen, OriginClose, OriginPush, OriginFlush, OriginTruncate. Writes are lazy up to MaxMemory. Public Flush and the Flush inside Close count toward MeanRate. Internal drains do not.

BufferedFileReader / BufferedFileWriter

Final file leaves. Path-only constructors probe the device at Setup. Explicit knobs stay explicit.

int main() {
BufferedFileReader in("in.bin", { ReadAhead{1 << 20}, MaxMemory{8 << 20} });
in.Open();
FIFO dest;
(void)in.Read(64 * 1024, dest);
auto tel = in.Telemetry();
in.Close();
}

Support

Questions and bugs: GitHub issues on this repository. Sponsorship: github.com/sponsors/StormBytePP.

Contributing

Open an issue before large work. Conventional Commits. Public headers need Doxygen. Do not send patches that reintroduce raw new for StormByte pointer types.

License

Dual license: GNU Lesser General Public License v3.0 or later, or a commercial license from the copyright holder. See [LICENSE](LICENSE), COPYING.LGPLv3 and https://www.gnu.org/licenses/lgpl-3.0.html. Third-party trees under thirdparty/ keep their own licenses.