chronon::sender::MultiProducerQueueAdapter
#include <MessageQueue.hpp>
Inherits from chronon::sender::IMessageQueue< T >
Public Functions
| Name | |
|---|---|
| bool | usesActiveFrontier() const |
| virtual std::optional< T > | tryPop(uint64_t current_cycle) override |
| uint64_t | transportOverflowEvents() const |
| bool | storageFullForThread(size_t queue_id) const |
| virtual size_t | storageCapacity() const override |
| size_t | stagedSize() const |
| virtual size_t | size() const override |
| size_t | sharedFifoHighWatermark() const |
| virtual void | setCapacity(size_t capacity) override |
| bool | pushFromThread(size_t queue_id, T data, uint64_t arrive_cycle, uint32_t sender_id =0) |
| virtual bool | push(T data, uint64_t arrive_cycle) override |
| void | prepareSharedFifo(uint64_t current_cycle) |
| virtual void | popAllInto(std::vector< T > & out, uint64_t current_cycle) override Reusable-buffer drain: avoids the fresh-vector allocation popAll() makes per call. |
| virtual std::vector< T > | popAll(uint64_t current_cycle) override |
| virtual std::optional< uint64_t > | minArrivalCycle() const override |
| virtual bool | hasReady(uint64_t current_cycle) const override |
| size_t | getQueueIdForThread(size_t thread_id) const |
| bool | fullForThread(size_t queue_id) const |
| virtual bool | full() const override |
| void | ensurePerThreadUsableCapacity(size_t min_usable_capacity) |
| virtual bool | empty() const override |
| template <typename Visitor > bool | consumeReady(uint64_t current_cycle, Visitor && visitor) |
| virtual void | clear() override |
| virtual size_t | capacity() const override |
| virtual size_t | available() const override |
| size_t | admissionOccupancyForThread(size_t queue_id, uint64_t send_cycle) const |
| std::optional< uint64_t > | admissionMinArrivalCycleForThread(size_t queue_id, uint64_t send_cycle) const |
| size_t | addProducerThread(size_t thread_id) |
| size_t | addProducerThread(size_t thread_id, bool track_admission) |
| MultiProducerQueueAdapter(size_t capacity =std::numeric_limits< size_t >::max(), size_t min_per_thread_usable_capacity =0) |
Public Attributes
| Name | |
|---|---|
| size_t | kLanesPerSignalWord |
| size_t | kFrontierLaneThreshold |
Additional inherited members
Public Functions inherited from chronon::sender::IMessageQueue< T >
| Name | |
|---|---|
| virtual | ~IMessageQueue() =default |
| virtual size_t | admissionOccupancy(uint64_t send_cycle) const |
| virtual std::optional< uint64_t > | admissionMinArrivalCycle(uint64_t send_cycle) const |
Detailed Description
template <typename T >
class chronon::sender::MultiProducerQueueAdapter;
MultiProducerQueueAdapter - deterministic MPSC built from independent SPSC lanes.
A Chronon Connection has exactly one producer and an InPort has exactly one consumer. Keeping one SPSC lane per Connection avoids the contended tail, slot-sequence CAS loops, and reclamation machinery of a general MPMC queue. The consumer merges lane heads by (arrive_cycle, sender_id, lane_id), which supplies the total order required for cycle-count reproducibility.
An unbounded InPort consumes selected lane slots in place. A bounded InPort adds a separately allocated, receiver-only shared FIFO: at each receiver cycle it moves at most capacity ready entries from the lanes into that preallocated ring. Capacity is therefore both aggregate destination depth and aggregate per-cycle admission without a contended producer tail, shared reservation RMW, lock, or receive-side allocation.
Small fan-in uses a branch-friendly linear scan. At larger fan-in, producers publish lane activity into cache-line-separated shards, and the consumer maintains a private min-heap containing at most one head per active lane. The activity bits are notifications, not ownership flags: exchange(0) coalesces repeated pushes, and a lane is reinserted after every pop while it remains non-empty. This avoids the producer-set versus consumer-clear lost-wakeup race of a persistent active-bit protocol.
Public Functions Documentation
function usesActiveFrontier
inline bool usesActiveFrontier() const
function tryPop
inline virtual std::optional< T > tryPop(
uint64_t current_cycle
) override
Reimplements: chronon::sender::IMessageQueue::tryPop
function transportOverflowEvents
inline uint64_t transportOverflowEvents() const
function storageFullForThread
inline bool storageFullForThread(
size_t queue_id
) const
function storageCapacity
inline virtual size_t storageCapacity() const override
Reimplements: chronon::sender::IMessageQueue::storageCapacity
function stagedSize
inline size_t stagedSize() const
Entries still resident in per-Connection transport lanes.
function size
inline virtual size_t size() const override
Reimplements: chronon::sender::IMessageQueue::size
function sharedFifoHighWatermark
inline size_t sharedFifoHighWatermark() const
Maximum observed aggregate destination FIFO occupancy.
function setCapacity
inline virtual void setCapacity(
size_t capacity
) override
Reimplements: chronon::sender::IMessageQueue::setCapacity
function pushFromThread
inline bool pushFromThread(
size_t queue_id,
T data,
uint64_t arrive_cycle,
uint32_t sender_id =0
)
function push
inline virtual bool push(
T data,
uint64_t arrive_cycle
) override
Reimplements: chronon::sender::IMessageQueue::push
function prepareSharedFifo
inline void prepareSharedFifo(
uint64_t current_cycle
)
Materialize deterministic ready ingress entries in the bounded shared destination FIFO. Consumer-only; Unit invokes it once before the model's tick, while receive queries may call it again to fill slots released in that tick. At most capacity entries cross the InPort boundary per cycle.
function popAllInto
inline virtual void popAllInto(
std::vector< T > & out,
uint64_t current_cycle
) override
Reusable-buffer drain: avoids the fresh-vector allocation popAll() makes per call.
Reimplements: chronon::sender::IMessageQueue::popAllInto
function popAll
inline virtual std::vector< T > popAll(
uint64_t current_cycle
) override
Reimplements: chronon::sender::IMessageQueue::popAll
function minArrivalCycle
inline virtual std::optional< uint64_t > minArrivalCycle() const override
Reimplements: chronon::sender::IMessageQueue::minArrivalCycle
function hasReady
inline virtual bool hasReady(
uint64_t current_cycle
) const override
Reimplements: chronon::sender::IMessageQueue::hasReady
function getQueueIdForThread
inline size_t getQueueIdForThread(
size_t thread_id
) const
function fullForThread
inline bool fullForThread(
size_t queue_id
) const
function full
inline virtual bool full() const override
Reimplements: chronon::sender::IMessageQueue::full
function ensurePerThreadUsableCapacity
inline void ensurePerThreadUsableCapacity(
size_t min_usable_capacity
)
function empty
inline virtual bool empty() const override
Reimplements: chronon::sender::IMessageQueue::empty
function consumeReady
template <typename Visitor >
inline bool consumeReady(
uint64_t current_cycle,
Visitor && visitor
)
Consumer-only path. Unbounded ports visit the selected lane head in place. Bounded ports visit the head of their receiver-owned shared FIFO. In both cases InPort applies receiver filtering without type erasure.
function clear
inline virtual void clear() override
Reimplements: chronon::sender::IMessageQueue::clear
function capacity
inline virtual size_t capacity() const override
Reimplements: chronon::sender::IMessageQueue::capacity
function available
inline virtual size_t available() const override
Reimplements: chronon::sender::IMessageQueue::available
function admissionOccupancyForThread
inline size_t admissionOccupancyForThread(
size_t queue_id,
uint64_t send_cycle
) const
function admissionMinArrivalCycleForThread
inline std::optional< uint64_t > admissionMinArrivalCycleForThread(
size_t queue_id,
uint64_t send_cycle
) const
function addProducerThread
inline size_t addProducerThread(
size_t thread_id
)
Register one stable producer key and return its lane id.
function addProducerThread
inline size_t addProducerThread(
size_t thread_id,
bool track_admission
)
function MultiProducerQueueAdapter
inline explicit MultiProducerQueueAdapter(
size_t capacity =std::numeric_limits< size_t >::max(),
size_t min_per_thread_usable_capacity =0
)
Public Attributes Documentation
variable kLanesPerSignalWord
static size_t kLanesPerSignalWord = 64;
variable kFrontierLaneThreshold
static size_t kFrontierLaneThreshold = 32;
Updated on 2026-07-23 at 16:24:33 +0000