Skip to content

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 (S items 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 standalone UnboundedSpsc type) 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 a DebraManager (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 type
  • ccProd: static PinScopeCardinality — Producer cardinality (ccSingle or ccMulti)
  • ccCons: static PinScopeCardinality — Consumer cardinality (ccSingle or ccMulti)
  • ST: static DeallocationStrategy — Reclamation policy (stManual or stEager). 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 subsequently push() / pop() through it. attach() performs the debra registration and may raise DebraRegistrationError.

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

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.

assertQueueParams

template assertQueueParams()

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):

  1. Pin opens BEFORE headSegment.load (this pinScope scope).
  2. Pin covers slot reservation through readItem/extract window (prevConsumerIdx.compareExchange -> move(seg.data[mySlot]) in SPMC; prevConsumerIdx.compareExchange -> tryClaim DWCAS in MPMC).
  3. Segment under pin == segment under claim (pointer linearity).
  4. headSegment advance (retireOnCAS) retires oldSeg via DEBRA; oldSeg is not freed until all pins in its retire-epoch rotate.
  5. Bulk variant acquires per-call pin satisfying (1)–(4).
  6. 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):

  1. Pin opens BEFORE headSegment.load (this pinScope scope).
  2. Pin covers slot reservation through readItem/extract window (prevConsumerIdx.compareExchange -> move(seg.data[mySlot]) in SPMC; prevConsumerIdx.compareExchange -> tryClaim DWCAS in MPMC).
  3. Segment under pin == segment under claim (pointer linearity).
  4. headSegment advance (retireOnCAS) retires oldSeg via DEBRA; oldSeg is not freed until all pins in its retire-epoch rotate.
  5. Bulk variant acquires per-call pin satisfying (1)–(4).
  6. 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