Queue¶
Coming from lockfreequeues v5?
See From lockfreequeues v5 for the package rename, import path changes, and the static-thread-affinity endpoint API.
Queue[T, ccProd, ccCons, ST, S, MaxThreads] is the unified unbounded,
lock-free queue exposed by lockfree v0.1.0. A single generic type
covers all four producer/consumer cardinality combinations — SPSC,
MPSC, SPMC, and MPMC — selected at compile time through ccProd and
ccCons. It replaces the family-prefixed unbounded types
(UnboundedSipsic, UnboundedSipmuc, UnboundedMupsic,
UnboundedMupmuc) shipped by the predecessor package lockfreequeues
through v4.x; the consolidated umbrella lockfree absorbs that surface
verbatim under a single generic.
Overview¶
- Capacity: Dynamic — grows by linking fixed-size segments
(
Sitems per segment). Push never fails for capacity reasons. - Push: Wait-free (SPSC) or lock-free (multi-producer)
- Pop: Wait-free (SPSC) or lock-free (multi-consumer)
Body layout splits on (ccProd, ccCons) is (ccSingle, ccSingle):
- SPSC (
ccSingle × ccSingle): no memory-reclamation manager. The linked-segment, committed-flag-free SPSC protocol (formerly the standaloneUnboundedSpsctype) is absorbed verbatim; retired segments are freed inline by the consumer-side advance. - MPSC / SPMC / MPMC: DEBRA+ epoch-based memory reclamation
(Brown 2015) via
the in-tree nebr SMR engine (
lockfree/smr/nebr). The queue owns or borrows aDebraManager(nebr.Manager), and each operating thread holds a per-thread handle for the pin/retire cycle.
Supported T¶
Per Path-C of the lockfree v0.1.0 element-type story, Queue[T, …]
supports the full Nim type vocabulary on the SPSC, MPSC, and SPMC
shapes — including ref T, string, and seq[T] — via the
ManagedRef / ManagedSlice wrappers documented in
Managed Ref and
Managed Slice. Plain copyable T flow
through unmodified; GC-managed T flow through the managed wrappers
while preserving the lock-free fast path. The unbounded MPMC shape
carries an additional T-width constraint documented below.
Type Parameters¶
T— Item typeccProd: static PinScopeCardinality— Producer cardinality (ccSingleorccMulti)ccCons: static PinScopeCardinality— Consumer cardinality (ccSingleorccMulti)ST: static DeallocationStrategy— Reclamation policy (stManualorstEager). Inert for the debra-free SPSC shape.S: static int— Segment size in items (must be> 0; power of 2, multiple of the cache line, recommended)MaxThreads: static int— DEBRA thread-registry capacity (must be> 0). A type-uniform phantom on the SPSC shape, which consumes no registry slots.
The parameter order is load-bearing: T, ccProd, ccCons, ST, S, MaxThreads.
lockfree v0.1.0 — unbounded MPMC T constraint
The unbounded MPMC shape (Queue[T, ccMulti, ccMulti, …]) requires
supportsCopyMem(T) AND sizeof(T) <= 8 (8 bytes on 64-bit;
4 bytes on 32-bit). Each cell packs a (seq, payload) pair into
a single DWCAS word per the LCRQ paper §4 close-CAS-on-empty
progress rule.
Violations fail at compile time with a {.error.} overload that
cites the migration path. For wider or move-only T, switch to
BQueue[T, ccMulti, ccMulti, …] (bounded MPMC, Vyukov per-slot
seq) — BQueue preserves general T support. To keep unbounded
MPMC, wrap as ptr T; see
From lockfreequeues v5
for the ptr T recipe and examples/job_scheduler.nim. The other
three unbounded shapes (SPSC / SPMC / MPSC) are unaffected and
accept ref T / string / seq[T] via the Path-C managed
wrappers.
Constructors¶
newQueue(Queue[T, ccProd, ccCons, ST, S, MaxThreads]) is the
canonical generic smart constructor. The typedesc-only overload
auto-creates a private DebraManager (and sets ownsManager = true)
for the reclaiming shapes, or skips manager allocation entirely for the
SPSC shape. Manager-borrowed overloads accept an existing
DebraManager (and optionally a ThreadHandle) and set
ownsManager = false.
Family-named thin wrappers are retained for ergonomic continuity with
the lockfreequeues v3.x/v4.x naming; all compile to the same Queue
type:
newUnboundedSpscQueue—ccSingle × ccSingle(SPSC)newUnboundedMpscQueue—ccMulti × ccSingle(MPSC)newUnboundedSpmcQueue—ccSingle × ccMulti(SPMC)newUnboundedMpmcQueue—ccMulti × ccMulti(MPMC)
Attach-Time Thread Registration¶
DEBRA thread registration is thread-affine: it stamps the calling thread and installs a signal handler on that OS thread. Registering on the wrong thread mis-routes the handle. Therefore, for the reclaiming shapes:
- No thread is registered at construction.
getProducer()/getConsumer()return an unregistered view: the calling thread reserves an index but does not register.- Each operating thread calls
attach()on its view on the thread that will subsequentlypush()/pop()through it.attach()performs the debra registration and may raiseDebraRegistrationError.
This differs from the pre-consolidation lockfreequeues API, where
get*() registered at get-time.
Usage¶
import lockfree
# SPSC: debra-free; no manager, no attach needed.
var spsc = newQueue(Queue[int, ccSingle, ccSingle, stEager, 64, 1])
var spscProducer = spsc.getProducer()
spscProducer.push(42) # never fails — grows as needed
let a = spsc.pop() # some(42)
# MPMC: each operating thread attaches before its first push/pop.
var mpmc = newQueue(Queue[int, ccMulti, ccMulti, stEager, 64, 8])
var producer = mpmc.getProducer()
discard producer.bindToThread() # registers this thread (may raise
# DebraRegistrationError)
producer.push(99)
var consumer = mpmc.getConsumer()
discard consumer.bindToThread()
let b = consumer.pop() # some(99)
# MPMC: when the calling thread is also the operating thread,
# `getProducerHere` / `getConsumerHere` are sugar for getX() + attach().
# Prefer this same-thread form; use the explicit getX() + attach()
# pair above when the view is handed off to a worker thread that does
# the push/pop (the attach() must run on that worker thread).
var mpmc2 = newQueue(Queue[int, ccMulti, ccMulti, stEager, 64, 8])
var producer2 = mpmc2.getProducerHere() # registers on current thread
producer2.push(99)
var consumer2 = mpmc2.getConsumerHere()
let c = consumer2.pop() # some(99)
Calling Convention by Cardinality¶
Queue always pushes through a Bound[T, Tag, Queue[...]] endpoint view — even for
ccProd == ccSingle. The view is obtained from queue.getProducer(),
and the single-producer arm doesn't require .bindToThread(). Pop is
asymmetric: ccCons == ccSingle arms expose a bare queue.pop(),
while ccCons == ccMulti requires a Bound[T, Tag, Queue[...]] endpoint view from
getConsumer() and per-thread .bindToThread(). Direct push on a
Queue or direct pop on a multi-consumer Queue is a
compile-time error whose diagnostic names only the user-visible
Bound[T, Tag, Queue[...]] endpoint / Bound[T, Tag, Queue[...]] endpoint aliases.
For the common same-thread case (the calling thread is also the
thread that will push/pop through the returned view),
getProducerHere() / getConsumerHere() are templates that combine
getX() + attach() in one call. They are pure sugar: identical
runtime behavior, identical diagnostics. Use the explicit
getX() + attach() pair when the view is handed off to a worker
thread that does the push/pop, so the attach() registers the
correct thread.
Typestate Notes¶
Queue carries a Lifecycle typestate (QueueInit -> QueueDestroyed)
driven by =destroy; the Bound[T, Tag, Queue[...]] endpoint / Bound[T, Tag, Queue[...]] endpoint views carry
a Claim-state typestate (QCUnclaimed -> QCBothClaimed). All
push/pop/attach/detach operations are state-preserving; only the
destructor moves a value to its terminal state. Use-after-destroy is a
documented limitation; see the lockfree CHANGELOG [0.1.0] entry.
See also¶
- Strict-LCRQ MPMC queue — the
unbounded MPMC engine and its DWCAS /
T-width constraints. - Legacy unbounded shapes — the
pre-consolidation SPSC / MPSC / SPMC unbounded protocols absorbed
into
Queue. - Managed Ref — Path-C wrapper for
ref T. - Managed Slice — Path-C wrapper for
stringandseq[T]. - Memory Management — DEBRA managers, deallocation strategies, and the attach/detach lifecycle.
- Safety Model — happens-before guarantees.
- Bounded vs Unbounded — choosing
between
QueueandBQueue.
queue
¶
Unbounded Queue generic.
Queue[T, ccProd, ccCons, ST, S, MaxThreads]
Param order is LOAD-BEARING: T, ccProd, ccCons, ST, S, MaxThreads
The bounded surface lives in bqueue.nim (BQueue[T, ccProd, ccCons,
N, P, C]). The (ccSingle, ccSingle) branch of the object body
carries no debra integration (no manager, no ownsManager, no
pin/retire wrappers) and uses the committed-flag-free linked-segment
protocol from the legacy standalone unbounded-SPSC module. The other
three cardinality combos (MPSC, SPMC, MPMC) carry debra integration.
Cardinality-illegal direct-on-queue calls (multi-producer push
or multi-consumer pop against Queue directly rather than via
QueueProducer / QueueConsumer) are gated by compile-time
{.error.} overloads. The error messages reference the user-visible
alias type names — no *Multi/*Single leakage.
Queueable[T] concept hookup. Bare
Queue[T, ...] does NOT satisfy the Queueable[T] concept defined
in ./typestates/with_bound — unbounded push/pop always routes
through a Bound[T, Tag, Queue[T, ...]] endpoint (no direct push/pop
on bare Queue, even for SPSC). Users wanting an unbounded queue
with the Queueable-style ergonomic surface use the typestate API
(getProducer().bindToThread()) or the withBoundProducer /
withBoundConsumer RAII templates in ./typestates/with_bound.
CLOSED_BIT
¶
const CLOSED_BIT = 1'u shl (sizeof(uint) * 8 - 1)
Strict-LCRQ close sentinel. A cell with seq == CLOSED_BIT
is permanently closed: no producer can publish into it, no consumer can claim it.
The sentinel occupies the high bit of the platform-native uint
so that LCRQCell[T] stays at native double-word width on every
target debra supports (16 bytes on 64-bit, 8 bytes on 32-bit).
On 64-bit (uint == uint64) the value is identical to 1'u64 shl 63.
LCRQCell
¶
type LCRQCell[T] = Atomic[Pair[uint, T]]
Strict-LCRQ cell: native double-word DWCAS-able pair of seq counter
(Pair.first) and payload (Pair.second). Transparent alias
for Atomic[Pair[uint, T]] — assigning a LCRQCell[T] to /
from the spelled-out type requires no conversion. Using
platform-native uint (rather than hardcoded uint64) keeps
the cell at the platform's native DWCAS width, preserving the
lock-free guarantee on every debra-supported target.
tryPublish inline ¶
proc tryPublish(cell: var LCRQCell[T]; expectedSeq: uint; value: T): bool
Producer publish via DWCAS into an empty cell.
Returns true on success (cell now (expectedSeq+1, value)).
Returns false if the cell is already filled, closed, or at a
different epoch.
Precondition for nullable T (ptr, ref, pointer, cstring,
proc, closures): value MUST NOT be nil. The std/options
transport used by tryClaim (some(val)) asserts not val.isNil
at runtime for nullable types; forbidding nil here surfaces the
contract violation at the producer rather than as a delayed
AssertionDefect inside an unrelated consumer's tryClaim call.
doAssert (not assert) so the guard
survives -d:danger builds. when compiles(value.isNil) covers
every nullable type Nim exposes (broader than T is ptr or ref).
Parameters
-
cell(var LCRQCell[T]) -
expectedSeq(uint) -
value(T)
Returns
bool
tryClaim inline ¶
proc tryClaim(cell: var LCRQCell[T]; expectedSeq: uint): Option[T]
Consumer claim via DWCAS.
CONTRACT: NEVER inspect observed.second. The CAS on the seq
encoding is the sole authority on cell state. A filled cell with
payload default(T) (e.g. q.push(0), q.push(nil)) is a
legitimate publish and MUST be returned via some(observed.second).
Parameters
-
cell(var LCRQCell[T]) -
expectedSeq(uint)
Returns
Option[T]
tryCloseOnEmpty inline ¶
proc tryCloseOnEmpty(cell: var LCRQCell[T]; expectedSeq: uint): bool
Consumer close-on-empty via DWCAS. Atomically sets
CLOSED_BIT on an empty cell so no producer can later publish
into it. Returns false if the cell is already filled or closed.
Parameters
-
cell(var LCRQCell[T]) -
expectedSeq(uint)
Returns
bool
QueueLifecycleCtx
¶
type QueueLifecycleCtx[T; ccProd, ccCons: static PinScopeCardinality; ST: static DeallocationStrategy;
S, MaxThreads: static int] = object
QueueInit
¶
type QueueInit[T; ccProd, ccCons: static PinScopeCardinality; ST: static DeallocationStrategy;
S, MaxThreads: static int] = distinct QueueLifecycleCtx[T, ccProd, ccCons, ST, S, MaxThreads]
Initial Lifecycle state for an unbounded Queue.
QueueDestroyed
¶
type QueueDestroyed[T; ccProd, ccCons: static PinScopeCardinality; ST: static DeallocationStrategy;
S, MaxThreads: static int] = distinct QueueLifecycleCtx[T, ccProd, ccCons, ST, S, MaxThreads]
Terminal Lifecycle state for an unbounded Queue.
Segment
¶
type Segment[T; ccProd, ccCons: static PinScopeCardinality; S: static int] = object
Unbounded-queue segment. One linked-segment payload, parameterized
by (ccProd, ccCons) so each cardinality variant's field set
matches its per-family analogue.
Field set:
- data: array[S, T] — slot storage (non-MPMC variants).
- cells: array[S, LCRQCell[T]] — strict-LCRQ cells (MPMC only).
Replaces committed + data on the ccMulti × ccMulti arm.
- next: Atomic[ptr Segment[...]] — linked-list pointer.
- tail: Atomic[int] — producer write index. Atomic for
multi-producer coordination and for spsc-equiv (publish
via release).
- head: int — single-consumer non-atomic read position.
Present on (ccProd × ccSingle) shapes (mpsc-equiv and
the absorbed spsc-equiv). Only the single consumer ever
writes it.
- committed: array[S, Atomic[bool]] — multi-producer
publication flags. Present on ccProd == ccMulti and
ccCons == ccSingle (MPSC only — MPMC migrated to cells).
- prevConsumerIdx: Atomic[int] — multi-consumer CAS slot.
Present on ccCons == ccMulti.
validateQueueParams ¶
proc validateQueueParams(_: typedesc[Queue[T, ccProd, ccCons, ST, S, MaxThreads]])
Compile-time entry point for the 2 param-coherence guards. Has no
runtime cost.
Parameters
-
_(typedesc[Queue[T, ccProd, ccCons, ST, S, MaxThreads]])
retireOnCAS discardable ¶
proc retireOnCAS(q: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]; scope: var PinnedScope[MaxThreads, CC]; atomic: var Atomic[U]; expected: var U; desired: U; dtor: Destructor): bool
Per-queue retireOnCAS wrapper. Delegates to nebr's
pinned_scope.retireOnCAS.
Parameters
-
q(var Queue[T, ccProd, ccCons, ST, S, MaxThreads]) -
scope(var PinnedScope[MaxThreads, CC]) -
atomic(var Atomic[U]) -
expected(var U) -
desired(U) -
dtor(Destructor)
Returns
bool
retireOnPublish ¶
proc retireOnPublish(q: var Queue[T, ccProd, ccSingle, ST, S, MaxThreads]; scope: var PinnedScope[MaxThreads, CC]; atomic: var Atomic[U]; desired: U; dtor: Destructor)
Per-queue retireOnPublish wrapper. **FOOT-GUN — single-writer
required (DR-S4).** Delegates to nebr's
pinned_scope.retireOnPublish.
Parameters
-
q(var Queue[T, ccProd, ccSingle, ST, S, MaxThreads]) -
scope(var PinnedScope[MaxThreads, CC]) -
atomic(var Atomic[U]) -
desired(U) -
dtor(Destructor)
newQueue ¶
proc newQueue(_: typedesc[Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, nebr.ccSingle]; handle: ThreadHandle[MaxThreads, nebr.ccSingle]): Queue[T, ccProd, ccSingle, ST, S, MaxThreads]
Manager-borrowed unbounded newQueue overload — ccCons == ccSingle
variants (mpsc-equiv only; the spsc-absorbed (ccSingle,
ccSingle) shape is debra-free and uses a separate {.error.}
overload below).
Caller owns the DebraManager. Sets ownsManager = false. The
handle is consumed by mpsc-equiv (ccProd == ccMulti) and stored
on the queue. ccProd-ccSingle spsc-equiv would be type-uniformly
constructable here, but is excluded by the dedicated spsc
{.error.} overload further below.
Parameters
-
_(typedesc[Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, nebr.ccSingle]) -
handle(ThreadHandle[MaxThreads, nebr.ccSingle])
Returns
Queue[T, ccProd, ccSingle, ST, S, MaxThreads]
newQueue ¶
proc newQueue(_: typedesc[Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, nebr.ccSingle]): Queue[T, ccProd, ccSingle, ST, S, MaxThreads]
Handle-free manager-borrowed unbounded newQueue overload for
ccCons == ccSingle (mpsc-equiv). The single consumer's debra
handle is NOT registered here; the consumer thread registers itself
via attachConsumer() before its first pop. The handle-carrying
overload above remains the escape hatch for callers who register
on the consumer thread and supply the handle directly.
Parameters
-
_(typedesc[Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, nebr.ccSingle])
Returns
Queue[T, ccProd, ccSingle, ST, S, MaxThreads]
newQueue ¶
proc newQueue(_: typedesc[Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, nebr.ccMulti]; handle: ThreadHandle[MaxThreads, nebr.ccMulti]): Queue[T, ccProd, ccMulti, ST, S, MaxThreads]
Manager-borrowed unbounded newQueue overload — ccCons == ccMulti
variants (spmc-equiv + mpmc-equiv).
Parameters
-
_(typedesc[Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, nebr.ccMulti]) -
handle(ThreadHandle[MaxThreads, nebr.ccMulti])
Returns
Queue[T, ccProd, ccMulti, ST, S, MaxThreads]
newQueue ¶
proc newQueue(_: typedesc[Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, nebr.ccMulti]): Queue[T, ccProd, ccMulti, ST, S, MaxThreads]
Handle-free manager-borrowed unbounded newQueue overload for
ccCons == ccMulti (spmc-equiv and mpmc-equiv).
Parameters
-
_(typedesc[Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, nebr.ccMulti])
Returns
Queue[T, ccProd, ccMulti, ST, S, MaxThreads]
newQueue error ¶
proc newQueue(_: typedesc[Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, CC]; handle: ThreadHandle[MaxThreads, CC]): Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
Parameters
-
_(typedesc[Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, CC]) -
handle(ThreadHandle[MaxThreads, CC])
Returns
Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
newQueue error ¶
proc newQueue(_: typedesc[Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]]; manager: ptr DebraManager[MaxThreads, CC]): Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
Parameters
-
_(typedesc[Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]]) -
manager(ptr DebraManager[MaxThreads, CC])
Returns
Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
newQueue ¶
proc newQueue(_: typedesc[Queue[T, ccProd, ccCons, ST, S, MaxThreads]]): Queue[T, ccProd, ccCons, ST, S, MaxThreads]
Auto-create unbounded newQueue overload. For the spsc-absorbed
(ccSingle, ccSingle) branch this skips manager allocation
entirely. For the other three cardinality combos, allocates a
private DebraManager[MaxThreads, ...] and sets ownsManager =
true.
No thread is registered at construction. registerThread is
thread-affine (stamps the calling thread + installs a signal
handler on the calling OS thread), so registering here would
mis-route the handle when the queue is later operated on a
different thread. Each operating thread registers itself at
attach-time: producers/consumers via getProducer().attach() /
getConsumer().attach(), and the mpsc-equiv single consumer via
attachConsumer() before its first pop.
Parameters
-
_(typedesc[Queue[T, ccProd, ccCons, ST, S, MaxThreads]])
Returns
Queue[T, ccProd, ccCons, ST, S, MaxThreads]
len ¶
proc len(self: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]): int
Number of items currently in the queue (atomic snapshot).
Parameters
-
self(var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Returns
int
segmentCount ¶
proc segmentCount(self: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]): int
Number of segments currently allocated (atomic snapshot).
Parameters
-
self(var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Returns
int
isEmpty inline ¶
proc isEmpty(self: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]): bool
Returns true if the queue is empty (atomic snapshot).
Checks if headSegment == tailSegment (and unconsumed slots).
Parameters
-
self(var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Returns
bool
pop ¶
proc pop(self: var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]): Option[T]
Spsc-absorbed pop — direct slot read + segment advance with
freeAligned(oldSeg). No pin (no retire-race; only one consumer
ever runs, only one producer ever writes). Lifted verbatim from
the unbounded SPSC pop path.
Parameters
-
self(var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads])
Returns
Option[T]
pop ¶
proc pop(self: var Queue[T, ccProd, ccSingle, ST, S, MaxThreads]; count: int): Option[seq[T]]
Batch pop for ccCons == ccSingle. Thin loop over single-item pop.
Parameters
-
self(var Queue[T, ccProd, ccSingle, ST, S, MaxThreads]) -
count(int)
Returns
Option[seq[T]]
pop error ¶
proc pop(self: var Queue[T, ccProd, ccMulti, ST, S, MaxThreads]): Option[T]
Parameters
-
self(var Queue[T, ccProd, ccMulti, ST, S, MaxThreads])
Returns
Option[T]
pop error ¶
proc pop(self: var Queue[T, ccProd, ccMulti, ST, S, MaxThreads]; count: int): Option[seq[T]]
Parameters
-
self(var Queue[T, ccProd, ccMulti, ST, S, MaxThreads]) -
count(int)
Returns
Option[seq[T]]
= error ¶
proc =(dst: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]; src: Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Compile-time copy ban. A Queue owns heap state (segment chain +
optionally the debra manager, recorded by ownsManager) that is
reclaimed exactly once in =destroy. A field-wise copy would
duplicate the owning ptrs and reclaim them twice. Move semantics
(the implicit =sink synthesized alongside =destroy) remain
available; only copies are rejected.
Parameters
-
dst(var Queue[T, ccProd, ccCons, ST, S, MaxThreads]) -
src(Queue[T, ccProd, ccCons, ST, S, MaxThreads])
= destructorTransition transitionError raises ¶
proc =(self: var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Destructor. Walks headSegment → next → ... freeing each
segment. For non-spsc cardinalities, additionally unbinds the
client refcount on the manager and (when ownsManager) runs the
manager's destructor.
Also drives the Lifecycle terminal transition
(QueueInit -> QueueDestroyed) via destructorTransition.
Precondition: all worker threads that attached to this queue
must be joined before =destroy runs. debra 0.8.0 has no
per-thread unregister; thread handles live until the manager is
destroyed here. Destroying the queue while an attached worker is
still pinning/popping is undefined. For shared-manager queues
(ownsManager == false) the destructor only unbinds the client
refcount; the manager is left intact for its owner to free.
Parameters
-
self(var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
newUnboundedSpscQueue inline ¶
proc newUnboundedSpscQueue(): Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
Unbounded SPSC (ccSingle × ccSingle) auto-create
smart-constructor. Skips manager allocation (SPSC has no debra integration).
Returns
Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]
newUnboundedMpscQueue inline ¶
proc newUnboundedMpscQueue(manager: ptr DebraManager[MaxThreads, nebr.ccSingle]): Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]
Unbounded MPSC (ccMulti × ccSingle) borrow
smart-constructor — manager-only form. The consumer thread calls
q.bindConsumer() on its own thread to register and obtain its
Bound endpoint.
The consumer's debra handle is owned by Bound (opaque storage).
bindConsumer wraps registration + binding in one call.
Parameters
-
manager(ptr DebraManager[MaxThreads, nebr.ccSingle])
Returns
Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]
newUnboundedMpscQueue inline ¶
proc newUnboundedMpscQueue(): Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]
Unbounded mpsc-equivalent (ccMulti × ccSingle) auto-create
smart-constructor. No thread is registered at construction: the
consumer thread calls attachConsumer() and producer threads call
getProducer().attach() on their own threads.
Returns
Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]
newUnboundedSpmcQueue inline ¶
proc newUnboundedSpmcQueue(manager: ptr DebraManager[MaxThreads, nebr.ccMulti]): Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]
Unbounded spmc-equivalent (ccSingle × ccMulti) borrow
smart-constructor.
Parameters
-
manager(ptr DebraManager[MaxThreads, nebr.ccMulti])
Returns
Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]
newUnboundedSpmcQueue inline ¶
proc newUnboundedSpmcQueue(): Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]
Unbounded spmc-equivalent (ccSingle × ccMulti) auto-create
smart-constructor.
Returns
Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]
newUnboundedMpmcQueue inline ¶
proc newUnboundedMpmcQueue(manager: ptr DebraManager[MaxThreads, nebr.ccMulti]): Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]
Unbounded mpmc-equivalent (ccMulti × ccMulti) borrow
smart-constructor.
Parameters
-
manager(ptr DebraManager[MaxThreads, nebr.ccMulti])
Returns
Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]
newUnboundedMpmcQueue inline ¶
proc newUnboundedMpmcQueue(): Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]
Unbounded mpmc-equivalent (ccMulti × ccMulti) auto-create
smart-constructor.
Returns
Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]
push tags raises notATransition ¶
proc push(self: Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; item: sink T)
Push a single item onto the unbounded queue (cardinality-dispatched).
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
item(sink T)
push tags raises notATransition ¶
proc push(self: Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; items: openArray[T])
Batch push (thin loop).
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
items(openArray[T])
pop tags raises notATransition ¶
proc pop(self: Bound[T, Tag, Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]): Option[T]
Pop for ccCons == ccSingle (SPSC + MPSC). Single consumer thread,
no pin required. Body consolidates the two direct-on-Queue pop
overloads with cardinality dispatch via the existing when arms.
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccSingle, ST, S, MaxThreads]])
Returns
Option[T]
pop tags raises notATransition ¶
proc pop(self: Bound[T, Tag, Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]; count: int): Option[seq[T]]
Batch pop for ccCons==ccSingle. Thin loop over single-item pop.
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccSingle, ST, S, MaxThreads]]) -
count(int)
Returns
Option[seq[T]]
pop tags raises notATransition ¶
proc pop(self: Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]]): Option[T]
SPMC pop — retire-bearing site. Pin claim via reconstructed
ThreadHandle from opaque Bound storage.
DEBRA Pin–Claim Ordering Invariant (both topologies):
- Pin opens BEFORE headSegment.load (this
pinScopescope). - Pin covers slot reservation through readItem/extract window
(
prevConsumerIdx.compareExchange->move(seg.data[mySlot])in SPMC;prevConsumerIdx.compareExchange->tryClaimDWCAS in MPMC). - Segment under pin == segment under claim (pointer linearity).
- headSegment advance (
retireOnCAS) retires oldSeg via DEBRA; oldSeg is not freed until all pins in its retire-epoch rotate. - Bulk variant acquires per-call pin satisfying (1)–(4).
- Pin coverage (topology-conditional):
a. Consumer (both SPMC and MPMC): pin opens BEFORE the consumer's slot
claim / close-CAS in pop, and remains open THROUGH the claim/close
completion AND the subsequent item extraction or segment transition.
Same pinScope as the headSegment.load that produced the segment
pointer.
b. MPMC producer: push publish-CAS occurs within the same pinScope
as the tailSegment.load that produced the segment pointer.
Required because peer producers can retire the segment under us
via headSegment-advance (peer-producer retire race).
Implementation: MPMC push wraps the retry loop in pinScope.
c. SPMC producer: push does NOT require pinScope. Single-producer
immunity: no peer producer can retire the segment under us.
Implementation: SPMC push explicitly opts out of pinScope.
DO NOT alter pinScope without re-establishing this invariant.
Parameters
-
self(Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]])
Returns
Option[T]
pop tags raises notATransition ¶
proc pop(self: Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]]): Option[T]
MPMC pop — retire-bearing site.
DEBRA Pin–Claim Ordering Invariant (both topologies):
- Pin opens BEFORE headSegment.load (this
pinScopescope). - Pin covers slot reservation through readItem/extract window
(
prevConsumerIdx.compareExchange->move(seg.data[mySlot])in SPMC;prevConsumerIdx.compareExchange->tryClaimDWCAS in MPMC). - Segment under pin == segment under claim (pointer linearity).
- headSegment advance (
retireOnCAS) retires oldSeg via DEBRA; oldSeg is not freed until all pins in its retire-epoch rotate. - Bulk variant acquires per-call pin satisfying (1)–(4).
- Pin coverage (topology-conditional):
a. Consumer (both SPMC and MPMC): pin opens BEFORE the consumer's slot
claim / close-CAS in pop, and remains open THROUGH the claim/close
completion AND the subsequent item extraction or segment transition.
Same pinScope as the headSegment.load that produced the segment
pointer.
b. MPMC producer: push publish-CAS occurs within the same pinScope
as the tailSegment.load that produced the segment pointer.
Required because peer producers can retire the segment under us
via headSegment-advance (peer-producer retire race).
Implementation: MPMC push wraps the retry loop in pinScope.
c. SPMC producer: push does NOT require pinScope. Single-producer
immunity: no peer producer can retire the segment under us.
Implementation: SPMC push explicitly opts out of pinScope.
DO NOT alter pinScope without re-establishing this invariant.
Parameters
-
self(Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]])
Returns
Option[T]
pop tags raises notATransition ¶
proc pop(self: Bound[T, Tag, Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]; count: int): Option[seq[T]]
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccMulti, ST, S, MaxThreads]]) -
count(int)
Returns
Option[seq[T]]
popBatch tags raises notATransition ¶
proc popBatch(self: var Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; dest: var openArray[T]; maxCount: int = -1): int
Batch pop for Queue into caller-supplied buffer dest.
For MPMC strict-LCRQ, amortizes slot reservations on prevConsumerIdx
via atomic slot advance and runs SMR epoch advancement once for the entire batch.
Parameters
-
self(var Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
dest(var openArray[T]) -
maxCount(int)
Returns
int
popBatch tags raises notATransition ¶
proc popBatch(self: Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; dest: var openArray[T]; maxCount: int = -1): int
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
dest(var openArray[T]) -
maxCount(int)
Returns
int
popChunk tags raises notATransition ¶
proc popChunk(self: var Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; chunkSize: int): seq[T]
Parameters
-
self(var Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
chunkSize(int)
Returns
seq[T]
popChunk tags raises notATransition ¶
proc popChunk(self: Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]; chunkSize: int): seq[T]
Parameters
-
self(Bound[T, Tag, Queue[T, ccProd, ccCons, ST, S, MaxThreads]]) -
chunkSize(int)
Returns
seq[T]
drain ¶
iterator drain(self: var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]): T
Drain SPSC-absorbed unbounded Queue. Bare-Queue pop is available
for SPSC only (the only arm where direct-on-Queue pop survived v5.0.0); other cardinalities route drain through Bound endpoints.
Parameters
-
self(var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads])
Returns
T
drain ¶
iterator drain(self: Bound[T, Tag, Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]]): T
Drain MPSC unbounded Queue via Bound consumer endpoint.
Parameters
-
self(Bound[T, Tag, Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]])
Returns
T
drain ¶
iterator drain(self: Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]]): T
Drain SPMC unbounded Queue via Bound consumer endpoint.
Parameters
-
self(Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]])
Returns
T
drain ¶
iterator drain(self: Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]]): T
Drain MPMC unbounded Queue via Bound consumer endpoint. The MPMC
consumer-CAS pop succeeds on the first try under the single- consumer drain contract (no contention).
Parameters
-
self(Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]])
Returns
T
items ¶
iterator items(self: var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]): T
Parameters
-
self(var Queue[T, ccSingle, ccSingle, ST, S, MaxThreads])
Returns
T
items ¶
iterator items(self: Bound[T, Tag, Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]]): T
Parameters
-
self(Bound[T, Tag, Queue[T, ccMulti, ccSingle, ST, S, MaxThreads]])
Returns
T
items ¶
iterator items(self: Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]]): T
Parameters
-
self(Bound[T, Tag, Queue[T, ccSingle, ccMulti, ST, S, MaxThreads]])
Returns
T
items ¶
iterator items(self: Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]]): T
Parameters
-
self(Bound[T, Tag, Queue[T, ccMulti, ccMulti, ST, S, MaxThreads]])
Returns
T
destroyAndDrain ¶
proc destroyAndDrain(self: sink Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]; cleanup: proc (item: T) {.gcsafe, raises: [].})
Drain SPSC unbounded Queue, run cleanup per item, then destroy.
Takes sink of the queue so the caller's binding is moved-from
(preventing scope-end double-destroy). Under mm:none this is the
ONLY safe teardown path for a non-empty queue with ref / string /
seq payloads.
Parameters
-
self(sink Queue[T, ccSingle, ccSingle, ST, S, MaxThreads]) -
cleanup(proc (item: T) {.gcsafe, raises: [].})
destroyAndDrain ¶
proc destroyAndDrain(self: sink Queue[T, ccSingle, ccSingle, ST, S, MaxThreads])
SPSC POD discard overload. Takes sink to prevent caller-side
scope-end double-destroy.
Parameters
-
self(sink Queue[T, ccSingle, ccSingle, ST, S, MaxThreads])
headSegmentForTest ¶
proc headSegmentForTest(self: var Queue[T, ccProd, ccCons, ST, S, MaxThreads]): pointer
Test-only accessor: returns the queue's current head-segment
pointer so the cache-line padding audit can verify base alignment.
Parameters
-
self(var Queue[T, ccProd, ccCons, ST, S, MaxThreads])
Returns
pointer
segmentTailOffsetForTest ¶
proc segmentTailOffsetForTest(_: typedesc[Segment[T, ccProd, ccCons, S]]): int
Test-only accessor: returns offset of the cache-line-padded tail
field within the unified Segment for any cardinality.
Parameters
-
_(typedesc[Segment[T, ccProd, ccCons, S]])
Returns
int
segmentHeadOffsetForTest ¶
proc segmentHeadOffsetForTest(_: typedesc[Segment[T, ccProd, ccCons, S]]): int
Test-only accessor: returns offset of head for cardinality
combos that carry it (ccCons == ccSingle). For shapes that
lack head Nim's offsetOf will compile-fail at the call site.
Parameters
-
_(typedesc[Segment[T, ccProd, ccCons, S]])
Returns
int
segmentCommittedOffsetForTest ¶
proc segmentCommittedOffsetForTest(_: typedesc[Segment[T, ccProd, ccCons, S]]): int
Test-only accessor: returns offset of committed for cardinality
combos that carry it (ccProd == ccMulti and ccCons == ccSingle,
i.e. MPSC only). MPMC carries cells; use
segmentCellsOffsetForTest on the MPMC arm. Calling this with an
MPMC cardinality fails at the offsetOf site (field absent).
Parameters
-
_(typedesc[Segment[T, ccProd, ccCons, S]])
Returns
int
segmentCellsOffsetForTest ¶
proc segmentCellsOffsetForTest(_: typedesc[Segment[T, ccProd, ccCons, S]]): int
Test-only accessor: returns offset of the strict-LCRQ cells array
for the MPMC arm (ccProd == ccMulti and ccCons == ccMulti). Other
cardinality combos lack the field; calling there compile-fails at
the offsetOf site.
Parameters
-
_(typedesc[Segment[T, ccProd, ccCons, S]])
Returns
int
segmentPrevConsumerIdxOffsetForTest ¶
proc segmentPrevConsumerIdxOffsetForTest(_: typedesc[Segment[T, ccProd, ccCons, S]]): int
Test-only accessor: returns offset of prevConsumerIdx for
cardinality combos that carry it (ccCons == ccMulti).
Parameters
-
_(typedesc[Segment[T, ccProd, ccCons, S]])
Returns
int