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.
Pieces plug into each other through ReadOnly / WriteOnly and through IO leaves.
Produceryields aConsumerover the sameRing.Pipelinetransforms stream buffers only (ReadOnly/WriteOnly). Pipes do not take IO leaves.Bridgemoves bytes from aReadOnlyor an IO reader into aWriteOnlyor an IO writer. IO tips are taken by move; in-memory tips stay referenced.Pumperowns aBridgeand runsPassthroughuntil EoF or failure.BufferedFileReaderis aBufferedReader.BufferedFileWriteris aBufferedWriter. Leaves implementOrigin*. 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 aPumper). BufferedFileReader/BufferedFileWriterfor seekable files with a page map.
See Pipeline, Bridge, Pumper, Telemetry, IO::BufferedReader and IO::BufferedWriter.
- BinaryData — octet payloads are
StormByte::BinaryData(Base). Lengths of byte buffers areStormByte::ByteSize.HopperandSinkcount items withStormByte::Size. - FIFO — grow-on-demand byte buffer. Not thread-safe.
Read/Peekkeep data;Extractconsumes it. - SharedFIFO — thread-safe FIFO.
Read/Extractblock until data orClose/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.Addclones or moves. - Bridge — manual transfer.
Passthrough(n, Operation)only. No worker. - Pumper — owns a Bridge and pumps until EoF or
Cancel. - Telemetry —
ReadTelemetry/WriteTelemetryasconst StormByte::Shared<…>.MeanRateis caller rate, not disk rate. - IO —
BufferedReader/BufferedWriterbases and file leaves. NestedParametersand knobs. - Lifecycle —
Close(),SetError(),EoF(),IsReadable(),IsWritable().
| 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 |
- Designed to interconnect
- What this module does
- The rest of the suite
- Documentation
- Installation
- Usage
- Support
- Contributing
- License
- This README: how to build, ownership, examples.
- Doxygen class reference: https://dev.stormbyte.org/StormByte-Buffer/.
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 buildShared 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>.
Namespace root is StormByte::Buffer. I/O types live in StormByte::Buffer::IO. Octet payloads use StormByte::BinaryData. Byte lengths use StormByte::ByteSize.
#include <StormByte/buffer/fifo.hxx>
using StormByte::BinaryData;
using StormByte::Buffer::FIFO;
using StormByte::Buffer::Position;
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"
}FIFO is not thread-safe. Concurrent ends use SharedFIFO or Ring.
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 <StormByte/buffer/consumer.hxx>
#include <StormByte/buffer/producer.hxx>
#include <thread>
using StormByte::BinaryData;
using StormByte::Buffer::Consumer;
using StormByte::Buffer::Producer;
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();
}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 <StormByte/buffer/hopper.hxx>
#include <memory>
#include <thread>
using StormByte::Buffer::Hopper;
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();
}Sink<T> maps integer keys to Hopper<T> buckets. Wire with To(key) / >> / <<.
#include <StormByte/buffer/sink.hxx>
#include <memory>
#include <string>
#include <thread>
using StormByte::Buffer::Sink;
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();
}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.
using StormByte::Buffer::IO::MaxMemory;
using StormByte::Buffer::IO::ReadAhead;
using StormByte::Buffer::IO::BufferedFileReader;
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);Explicit 0 is 0. It does not probe.
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 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.
#include <StormByte/buffer/pipe.hxx>
#include <StormByte/buffer/pipeline.hxx>
#include <StormByte/buffer/producer.hxx>
#include <StormByte/safe_pointers.hxx>
using StormByte::Buffer::Consumer;
using StormByte::Buffer::ExecutionMode;
using StormByte::Buffer::Pipe;
using StormByte::Buffer::Pipeline;
using StormByte::Buffer::Producer;
using StormByte::Buffer::ReadOnly;
using StormByte::Buffer::WriteOnly;
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 {
StormByte::BinaryData raw;
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 {
return StormByte::Unique<Pipe>::MakePointer<StripCrPipe>(*this);
}
PointerType Move() noexcept override {
return StormByte::Unique<Pipe>::MakePointer<StripCrPipe>(std::move(*this));
}
};
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;
}Sync runs on the caller thread. Async returns at once. Parallel is one thread per pipe. Combine with |.
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. WriteTryAgainis retried until that call completes.Operationapplies to the read tip only:Blockingwaits fornor EoF;NonBlockingtakes what is available now, up ton.n == 0is current contents (Available()on IO does not touch the origin).StateisOpen,ClosedorFailed.Failed()is only a real tip fault. A moved-from Bridge isClosed, notFailed.Close()ends the session. Owned adapters and stolen IO leaves are released. AfterClose, a consumed source EoF, or a move-from, the instance isClosedand cannot be re-armed. Construct a new Bridge to transfer again.Close()is idempotent and does not setFailed.- A NonBlocking
Passthroughthat returns0is not the end of the session. - Telemetry: the Bridge always keeps a
Sharedcopy 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. AfterClosethe same handles remain valid. ASharedyou 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 <StormByte/buffer/bridge.hxx>
#include <StormByte/buffer/fifo.hxx>
#include <StormByte/buffer/io/buffered_file_writer.hxx>
#include <filesystem>
using StormByte::Buffer::Bridge;
using StormByte::Buffer::FIFO;
using StormByte::Buffer::IO::BufferedFileWriter;
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");
}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.
#include <StormByte/buffer/bridge.hxx>
#include <StormByte/buffer/io/buffered_file_writer.hxx>
#include <StormByte/buffer/pipeline.hxx>
#include <StormByte/buffer/producer.hxx>
using StormByte::Buffer::Bridge;
using StormByte::Buffer::Consumer;
using StormByte::Buffer::ExecutionMode;
using StormByte::Buffer::Pipeline;
using StormByte::Buffer::Producer;
using StormByte::Buffer::IO::BufferedFileWriter;
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 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 ofPassthrougheach cycle.0is automatic chunking, not Bridge “current contents”.HighWater— input cap only. Omitted:0if the source is IO, otherwise a backend default (constexpr in the PIMPL). Explicit0: no Pumper cap. Use0when the IO source already limits itself. Non-IO sources are unbounded by design.Togglepauses and resumes. No-op ifFailedorCanceled.Cancelis terminal (Canceled, no restart). It does not setFailed. ItClose()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), notFailedand notCanceled.- Telemetry handles are copied from the Bridge at construction. They stay valid after
Cancel. ASharedyou 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.
#include <StormByte/buffer/bridge.hxx>
#include <StormByte/buffer/pumper.hxx>
#include <StormByte/buffer/io/buffered_file_reader.hxx>
#include <StormByte/buffer/io/buffered_file_writer.hxx>
using StormByte::Buffer::Bridge;
using StormByte::Buffer::Chunk;
using StormByte::Buffer::HighWater;
using StormByte::Buffer::Pumper;
using StormByte::Buffer::IO::BufferedFileReader;
using StormByte::Buffer::IO::BufferedFileWriter;
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
}finish = !pump.Failed() && !pump.Canceled() && pump.EoF().
active = !pump.Failed() && !pump.Canceled() && !pump.EoF().
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.
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.
Final file leaves. Path-only constructors probe the device at Setup. Explicit knobs stay explicit.
#include <StormByte/buffer/io/buffered_file_reader.hxx>
#include <StormByte/buffer/fifo.hxx>
using StormByte::Buffer::FIFO;
using StormByte::Buffer::IO::BufferedFileReader;
using StormByte::Buffer::IO::MaxMemory;
using StormByte::Buffer::IO::ReadAhead;
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();
}Questions and bugs: GitHub issues on this repository. Sponsorship: github.com/sponsors/StormBytePP.
Open an issue before large work. Conventional Commits. Public headers need Doxygen. Do not send patches that reintroduce raw new for StormByte pointer types.
Dual license: GNU Lesser General Public License v3.0 or later, or a commercial license from the copyright holder. See LICENSE, COPYING.LGPLv3 and https://www.gnu.org/licenses/lgpl-3.0.html. Third-party trees under thirdparty/ keep their own licenses.