ipcq 0.12.0
Loading...
Searching...
No Matches
README

ipcq

ipcq is not appropriate for all use-cases. Design trade-offs were made to achieve the following:

  • Single producer, multiple consumer (SPMC).
  • Templated element type that satisfies the TriviallyCopyable requirements.
  • No error propagation from reader to writer or other readers (the queue is optionally lock free for this reason). This includes starvation due to holding locks.
  • Support write frequencies in the kHz range with element size in the 0.5 - 1 MiB range (note that there is no strict limit that constrain user defined types).
  • Reader can copy a subset of the full element type.
Note
ipcq use atomic operations and memory barriers for synchronization, which due to cache coherency protocols is almost always guaranteed not to be optimal. The choice was made to be able to guarantee that readers cannot block the writer.

The ipcq::Writer is located in ipcq/src/include/ipcq/writer.hpp. The ipcq::Reader is located in ipcq/src/include/ipcq/reader.hpp.

The details around shared memory management is located in ipcq/src/include/ipcq/detail/shm.hpp.

Dependencies

Library dependencies:

  • boost (mainly boost::interprocess)
  • NUMA++

Basic Usage

See also demo application in separate "tools" project.

The following shows the most basic use:

/**
* @file
* @copyright
* SPDX-FileCopyrightText: 2023-2023 European Southern Observatory (ESO)
*
* SPDX-License-Identifier: LGPL-3.0-only
*/
#include <array>
#include <cstdint>
#include <ipcq/reader.hpp>
#include <ipcq/writer.hpp>
struct MyElement {
enum class Status {
};
std::size_t size;
std::array<float, 64 * 64> data;
};
int main() {
auto capacity = 64u;
ipcq::Writer<MyElement> writer("topic", capacity);
ipcq::Reader<MyElement> reader("topic");
// Write
{
auto ec = writer.Write(data, ipcq::Notify::All);
if (ec) {
std::cerr << "Failed to write to queue: " << ec.message() << std::endl;
return 1;
}
}
// Read
{
using namespace std::chrono_literals;
MyElement out;
auto [ec, count] =
reader.Read([&out](MyElement const& in) noexcept { out = in; }, 1u, 10ms);
if (ec) {
std::cerr << "Failed to read from queue: " << ec.message() << std::endl;
return 1;
} else {
// out is OK
// ...
}
}
}
BasicReader< T, detail::BoostConditionPolicy > Reader
Convenience alias with reasonable default condition policy.
Definition reader.hpp:489
BasicWriter< T, BoostConditionPolicy > Writer
Convenience alias with reasonable default condition policy.
Definition writer.hpp:230
@ All
Notify all waiting readers.
Definition writer.hpp:31
Contains definition of ipcq::BasicReader.
Status status
std::array< float, 64 *64 > data
std::size_t size
Contains declarations for ipcq::BasicWriter.

To read a batch consider the following (incomplete) example:

/**
* @file
* @copyright
* SPDX-FileCopyrightText: 2023-2023 European Southern Observatory (ESO)
*
* SPDX-License-Identifier: LGPL-3.0-only
*/
#include <chrono>
#include <ipcq/adapter.hpp>
#include <ipcq/reader.hpp>
#include <ipcq/writer.hpp>
struct MyElement {
std::array<uint8_t, 128> data;
};
/** Writer side */
void Writer() {
auto capacity = 64u;
ipcq::Writer<MyElement> writer("myqueue", capacity);
MyElement element = {};
// Prepare element
auto err = writer.Write(element, ipcq::Notify::All);
if (err) {
// handle error
}
}
/** Reader side */
void Reader() {
using namespace std::chrono_literals;
// Create reader, waiting up to 30s for the topic to be created if it does not exist.
// If the topic is not created after that it will throw `std::system_error` with error
// code `ipcq::Error::Timeout`
auto reader = ipcq::Reader<MyElement>::MakeReader("myqueue", 30s);
std::array<MyElement, 40> elements;
// Read *up to* 40 elements, using the adapter utility that performs assignment
// for each element into the array.
// If there is no data at all yet we wait 2s before giving up (returning `ipcq::Error::Timeout`)
auto [err, count] =
reader.Read(ipcq::OutputIteratorAdapter(std::begin(elements)), elements.size(), 2s);
if (err) {
// Handle errors.
// Particularly ipcq::Error::InconsistentState indicates that writer has broken the
// sequential reads by `reader`. This is recovered by resetting the internal state to read
// the next written element by invoking `ipcq::BasicReader::Reset()`.
} else {
// Success
// `count` elements have been copied.
// process result ...
}
}
Provides iterator adapter for ipcq::BasicReader.
static BasicReader MakeReader(char const *topic_name, std::chrono::duration< Rep, Period > timeout)
Definition reader.hpp:190
Adapts either LegacyOutputIterator or Iterator type (e.g.
Definition adapter.hpp:52

Debugging Tool

The debug tool ipcq-spy is available, which can attach to any ipcq topic, irrespective of the contained types. It cannot be notified of availabillity of new data and will simply poll at a given interval. Nevertheless it is possible to tell if data is being written, together with other additional information.

Example output when spying on the ipcq-demo topic from the ipcq-demo application:

$ ipcq-spy
topic-name ipcq-demo
owner-pid 30258
owner-cmdline ipcq-demo writer
shm-size 37.61Mib
(39435624b)
shm-capacity 700
element-size 55.02kib
(56336b)
element-type-name GenericTopic<14080ul>
(12GenericTopicILm14080EE)
condition-type-name ipcq::detail::BoostConditionPolicy
(N4ipcq6detail20BoostConditionPolicyE)
shm-numa-map (ctl) 7f0e12e24000 default file=/dev/shm/ipcq-ipcq-demo dirty=1 mapmax=4 active=0 N0=1
kernelpagesize_kB=4
shm-numa-map (data) 7f0e0efbe000 default file=/dev/shm/ipcq-ipcq-demo dirty=193 mapped=9628 mapmax=4
active=0 N0=193 N1=9435 kernelpagesize_kB=4
stop spying with CTRL-C
[/] last internal counter value: 709, avg. freq: 631.6