Examples¶
Practical examples demonstrating lock-free queue patterns.
Running Examples¶
Or run individual examples:
nim c --threads:on -r examples/audio_buffer.nim
nim c --threads:on -r examples/task_fanout.nim
nim c --threads:on -r examples/event_collector.nim
nim c --threads:on -r examples/job_scheduler.nim
Audio Buffer (Bounded SPSC)¶
Real-time audio processing with fixed latency. The producer captures audio samples, the consumer plays them.
Key properties:
- Fixed latency determined by buffer size
- No allocation during operation
- Wait-free operations for predictable timing
Use cases: Audio engines, DSP pipelines, real-time signal processing
## Audio Buffer Example
##
## Demonstrates using a bounded Spsc (SPSC) queue for real-time audio processing.
## The producer thread captures audio samples, the consumer thread plays them.
##
## Key properties:
## - Fixed latency: buffer size determines latency (64 samples @ 44.1kHz = 1.45ms)
## - No allocation: ring buffer operations never allocate
## - Wait-free: both push and pop complete in bounded time
##
## This pattern is essential for audio applications where:
## - Latency must be predictable and minimal
## - Glitches from GC pauses or allocation are unacceptable
## - Sample rate is fixed and known
import lockfree/atomics
import lockfree/atomics/dsl
import math
import os
import options
import lockfree
const
SampleRate = 44100
BufferSize = 64 # ~1.45ms latency at 44.1kHz
DurationMs = 100 # Simulate 100ms of audio
type AudioSample = object
left: float32
right: float32
timestamp: int64
var
q = newSpscQueue[AudioSample, BufferSize]()
running: Atomic[bool]
samplesProduced: Atomic[int]
samplesConsumed: Atomic[int]
underruns: Atomic[int] # Consumer needed data but queue was empty
overruns: Atomic[int] # Producer couldn't push because queue was full
proc captureThread() {.thread.} =
## Simulates audio capture hardware filling the buffer.
## In a real application, this would be driven by hardware interrupts.
var sampleIndex = 0
let totalSamples = (SampleRate * DurationMs) div 1000
while sampleIndex < totalSamples:
# Generate a simple sine wave (440Hz tone)
let t = float32(sampleIndex) / float32(SampleRate)
let value = sin(t * 440.0 * 2.0 * PI).float32 * 0.5
let sample = AudioSample(left: value, right: value, timestamp: sampleIndex)
if q.push(sample):
discard samplesProduced.fetchAdd(1, moRelaxed)
inc sampleIndex
else:
# Buffer full - would cause overrun in real audio
discard overruns.fetchAdd(1, moRelaxed)
# In real audio, we'd drop the sample or wait for hardware timing
sleep(0) # Yield to let consumer catch up
running.store(false, moRelease)
proc playbackThread() {.thread.} =
## Simulates audio playback hardware draining the buffer.
## In a real application, this would be driven by hardware interrupts.
while running.load(moAcquire) or q.pop().isSome:
let sample = q.pop()
if sample.isSome:
# In a real application, send to DAC
discard samplesConsumed.fetchAdd(1, moRelaxed)
else:
# Buffer empty - would cause underrun (glitch) in real audio
discard underruns.fetchAdd(1, moRelaxed)
# Simulate playback timing (~22.7us per sample at 44.1kHz)
# In reality, hardware timing drives this
sleep(0)
when isMainModule:
echo "Audio Buffer Example"
echo "===================="
echo "Buffer size: ", BufferSize, " samples"
echo "Latency: ", (BufferSize * 1000) div SampleRate, "ms"
echo "Duration: ", DurationMs, "ms"
echo ""
running.store(true, moRelease)
samplesProduced.store(0, moRelaxed)
samplesConsumed.store(0, moRelaxed)
underruns.store(0, moRelaxed)
overruns.store(0, moRelaxed)
var threads: array[2, Thread[void]]
threads[0].createThread(captureThread)
threads[1].createThread(playbackThread)
joinThreads(threads)
echo "Results:"
echo " Samples produced: ", samplesProduced.load(moRelaxed)
echo " Samples consumed: ", samplesConsumed.load(moRelaxed)
echo " Underruns: ", underruns.load(moRelaxed)
echo " Overruns: ", overruns.load(moRelaxed)
let produced = samplesProduced.load(moRelaxed)
let consumed = samplesConsumed.load(moRelaxed)
if produced == consumed and underruns.load(moRelaxed) == 0:
echo ""
echo "Clean audio stream - no glitches!"
else:
echo ""
echo "Note: Some underruns/overruns expected in simulation"
echo "(Real audio uses hardware timing, not sleep())"
Task Fan-Out (Bounded SPMC)¶
Work distribution from a single dispatcher to multiple workers.
Key properties:
- Work-stealing via CAS coordination
- Natural load balancing (faster workers get more tasks)
- Bounded memory provides backpressure
Use cases: HTTP request routing, image processing, game engine jobs
## Task Fan-Out Example
##
## Demonstrates using a bounded Spmc (SPMC) queue for distributing work
## from a single producer to multiple consumer workers.
##
## Pattern: Single dispatcher → Multiple workers
##
## Key properties:
## - Work-stealing: consumers compete for tasks via CAS
## - Load balancing: faster workers naturally get more tasks
## - Bounded memory: fixed queue size provides backpressure
## - Wait-free push: dispatcher never blocks
##
## Use cases:
## - HTTP request routing to worker pool
## - Image processing pipeline
## - Game engine job distribution
import lockfree/atomics
import lockfree/atomics/dsl
import os
import options
import std/monotimes
import times
import lockfree
import lockfree/endpoint
import lockfree/role_tags
const
QueueCapacity = 128
NumWorkers = 4
NumTasks = 1000
type
TaskKind = enum
tkCompute # CPU-bound work
tkIO # Simulated I/O
tkFast # Quick task
Task = object
id: int
kind: TaskKind
payload: int
var
q = newSpmcQueue[Task, QueueCapacity, NumWorkers]()
done: Atomic[bool]
tasksCompleted: array[NumWorkers, Atomic[int]]
totalLatency: array[NumWorkers, Atomic[int64]]
proc simulateWork(task: Task) =
## Simulate different types of work
case task.kind
of tkCompute:
# Simulate CPU work with busy loop
var sum = 0
for i in 0 ..< task.payload:
sum += i
discard sum
of tkIO:
# Simulate I/O wait
sleep(task.payload)
of tkFast:
# Minimal work
discard
proc workerThread(idx: int) {.thread.} =
## Worker consumes tasks from the queue until shutdown.
var consumer = q.getConsumerHere(idx)
var completed = 0
var latencySum: int64 = 0
while true:
let task = consumer.pop()
if task.isSome:
let start = getMonoTime()
simulateWork(task.get)
let elapsed = (getMonoTime() - start).inMicroseconds
latencySum += elapsed
inc completed
elif done.load(moAcquire):
break
else:
# No work available, brief pause before retry
sleep(0)
tasksCompleted[idx].store(completed, moRelease)
totalLatency[idx].store(latencySum, moRelease)
proc dispatcherThread() {.thread.} =
## Dispatcher generates and distributes tasks.
var taskId = 0
for i in 0 ..< NumTasks:
# Create varied task mix
let kind =
case i mod 10
of 0 .. 2: tkFast
of 3 .. 6: tkCompute
else: tkIO
let payload =
case kind
of tkFast:
0
of tkCompute:
1000 + (i mod 500)
of tkIO:
1 + (i mod 3)
let task = Task(id: taskId, kind: kind, payload: payload)
inc taskId
# Push with backpressure - wait if queue is full
while not q.push(task):
sleep(0)
done.store(true, moRelease)
when isMainModule:
echo "Task Fan-Out Example"
echo "===================="
echo "Queue capacity: ", QueueCapacity
echo "Workers: ", NumWorkers
echo "Tasks: ", NumTasks
echo ""
done.store(false, moRelaxed)
for i in 0 ..< NumWorkers:
tasksCompleted[i].store(0, moRelaxed)
totalLatency[i].store(0, moRelaxed)
let startTime = getMonoTime()
# Start workers
var workers: array[NumWorkers, Thread[int]]
for i in 0 ..< NumWorkers:
createThread(workers[i], workerThread, i)
# Start dispatcher
var dispatcher: Thread[void]
createThread(dispatcher, dispatcherThread)
# Wait for completion
joinThread(dispatcher)
sleep(50) # Let workers drain queue
for i in 0 ..< NumWorkers:
joinThread(workers[i])
let totalTime = (getMonoTime() - startTime).inMilliseconds
# Report results
echo "Results:"
var totalCompleted = 0
for i in 0 ..< NumWorkers:
let completed = tasksCompleted[i].load(moAcquire)
let latency = totalLatency[i].load(moAcquire)
let avgLatency =
if completed > 0:
latency div completed
else:
0
echo " Worker ", i, ": ", completed, " tasks, avg ", avgLatency, "us/task"
totalCompleted += completed
echo ""
echo "Total completed: ", totalCompleted, "/", NumTasks
echo "Total time: ", totalTime, "ms"
echo "Throughput: ", (NumTasks * 1000) div max(1, totalTime.int), " tasks/sec"
if totalCompleted == NumTasks:
echo ""
echo "All tasks processed successfully!"
Event Collector (Unbounded MPSC)¶
Collecting events from multiple sources into a single processing pipeline.
Key properties:
- Handles traffic bursts (queue grows during spikes)
- Never drops events
- Lock-free producers don't block each other
Use cases: Log aggregation, metrics collection, network packet capture
## Event Collector Example — v5.0.0 static-affinity endpoint API.
##
## Multi-producer event sources feed a single-consumer processor via an
## unbounded MPSC queue. Each producer thread calls `getProducerHere()`
## on its own thread (sugar for `getProducer().bindToThread()`,
## debra-registers per-thread). The single consumer uses
## `bindConsumer()` (the v5.0.0 replacement for `attachConsumer`).
import os
import options
import random
import std/monotimes
import lockfree/atomics
import times
import lockfree
import lockfree/endpoint
import lockfree/role_tags
from lockfree/smr/nebr import DebraManager, initDebraManager
const
SegmentSize = 64
NumSources = 4
DurationMs = 200
MaxThreads = NumSources + 4
type
EventKind = enum
ekClick
ekPageView
ekError
ekMetric
Event = object
sourceId: int
kind: EventKind
timestamp: int64
value: int
QueueT = Queue[Event, ccMulti, ccSingle, stEager, SegmentSize, MaxThreads]
SourceContext = object
queue: ptr QueueT
sourceId: int
startTime: MonoTime
ProcessorContext = object
queue: ptr QueueT
var
manager = initDebraManager[MaxThreads]()
queue = newUnboundedMpscQueue[Event, stEager, SegmentSize, MaxThreads](addr manager)
running: Atomic[bool]
eventsProduced: array[NumSources, Atomic[int]]
eventsConsumed: Atomic[int]
maxQueueDepth: Atomic[int]
proc eventSourceThread(ctx: ptr SourceContext) {.thread.} =
{.cast(gcsafe).}:
# v5.0.0: getProducerHere binds this thread's debra handle on entry.
var producer = ctx.queue[].getProducerHere()
var produced = 0
while running.load(moAcquire):
let burstMode = rand(100) < 10
let eventCount =
if burstMode:
rand(10 .. 20)
else:
rand(1 .. 3)
for _ in 0 ..< eventCount:
let event = Event(
sourceId: ctx.sourceId,
kind: EventKind(rand(ord(EventKind.high))),
timestamp: (getMonoTime() - ctx.startTime).inMicroseconds,
value: rand(1000),
)
producer.push(event)
inc produced
let delayMs =
if burstMode:
1
else:
rand(5 .. 15)
sleep(delayMs)
eventsProduced[ctx.sourceId].store(produced, moRelease)
proc processorThread(ctx: ptr ProcessorContext) {.thread.} =
{.cast(gcsafe).}:
# v5.0.0: bindConsumer is the one-shot wrapper that replaces v4.x
# attachConsumer. Single-consumer cardinality registers exactly
# one debra handle on the consuming thread.
var consumer = ctx.queue[].bindConsumer()
var consumed = 0
var maxDepth = 0
var eventCounts: array[EventKind, int]
while running.load(moAcquire) or ctx.queue[].len() > 0:
let event = consumer.pop()
if event.isSome:
inc eventCounts[event.get.kind]
inc consumed
let depth = ctx.queue[].len()
if depth > maxDepth:
maxDepth = depth
else:
sleep(1)
eventsConsumed.store(consumed, moRelease)
maxQueueDepth.store(maxDepth, moRelease)
echo ""
echo "Event breakdown:"
for kind in EventKind:
echo " ", kind, ": ", eventCounts[kind]
when isMainModule:
randomize()
echo "Event Collector Example (v5.0.0 static-affinity API)"
echo "===================================================="
echo "Segment size: ", SegmentSize
echo "Event sources: ", NumSources
echo "Duration: ", DurationMs, "ms"
echo ""
running.store(true, moRelease)
eventsConsumed.store(0, moRelaxed)
maxQueueDepth.store(0, moRelaxed)
for i in 0 ..< NumSources:
eventsProduced[i].store(0, moRelaxed)
let startTime = getMonoTime()
var procCtx = ProcessorContext(queue: addr queue)
var processor: Thread[ptr ProcessorContext]
createThread(processor, processorThread, addr procCtx)
var sourceContexts: array[NumSources, SourceContext]
for i in 0 ..< NumSources:
sourceContexts[i] =
SourceContext(queue: addr queue, sourceId: i, startTime: startTime)
var sources: array[NumSources, Thread[ptr SourceContext]]
for i in 0 ..< NumSources:
createThread(sources[i], eventSourceThread, addr sourceContexts[i])
sleep(DurationMs)
running.store(false, moRelease)
for i in 0 ..< NumSources:
joinThread(sources[i])
sleep(50)
joinThread(processor)
let totalTime = (getMonoTime() - startTime).inMilliseconds
echo ""
echo "Results:"
var totalProduced = 0
for i in 0 ..< NumSources:
let produced = eventsProduced[i].load(moAcquire)
echo " Source ", i, ": ", produced, " events"
totalProduced += produced
let consumed = eventsConsumed.load(moAcquire)
let maxDepth = maxQueueDepth.load(moAcquire)
echo ""
echo "Total produced: ", totalProduced
echo "Total consumed: ", consumed
echo "Max queue depth: ", maxDepth
echo "Final segments: ", queue.segmentCount()
echo "Total time: ", totalTime, "ms"
echo "Throughput: ", (consumed * 1000) div max(1, totalTime.int), " events/sec"
if totalProduced == consumed:
echo ""
echo "All events processed - no data loss!"
Job Scheduler (Unbounded MPMC)¶
Dynamic job scheduling with multiple submitters and workers.
Key properties:
- Elastic capacity grows with pending work
- Lock-free operations for high concurrency
- Dynamic scaling of producers and consumers
Use cases: Background job processing, build systems, database query scheduling
## Job Scheduler Example
##
## Demonstrates using an unbounded Mpmc (MPMC) queue for a dynamic job
## scheduling system with multiple producers and consumers.
##
## Pattern: Multiple submitters → Job queue → Multiple workers
##
## Key properties:
## - Elastic capacity: queue grows with pending work
## - Lock-free operations: submitters and workers don't block each other
## - Dynamic scaling: workers can be added/removed at runtime
## - Fair scheduling: committed flag ensures FIFO within segments
##
## Use cases:
## - Web server request handling
## - Background job processing (like Sidekiq/Celery)
## - Build system task execution
## - Database query scheduling
import os
import options
import random
import std/monotimes
import strutils
import times
import lockfree
import lockfree/endpoint
import lockfree/role_tags
import ./debra_cc_helpers
const
SegmentSize = 32
NumSubmitters = 3
NumWorkers = 4
JobsPerSubmitter = 50
MaxThreads = NumSubmitters + NumWorkers + 4 # producers + consumers + main + slack
type
Priority = enum
pLow
pNormal
pHigh
Job = object
id: int
submitterId: int
priority: Priority
workMs: int # Simulated work duration
# v5.0.0 strict-LCRQ MPMC requires `sizeof(T) <= 8` because the
# publish slot is `Atomic[Pair[uint, T]]` (128-bit DWCAS, seq + T).
# `Job` is 32 bytes, so we push `ptr Job` (8 bytes) instead and
# let producers heap-allocate / consumers free.
JobQueue = Queue[ptr Job, ccMulti, ccMulti, stEager, SegmentSize, MaxThreads]
SubmitterContext = object
queue: ptr JobQueue
submitterId: int
WorkerContext = object
queue: ptr JobQueue
workerId: int
var
# `initMultiConsumerManager` (from `./debra_cc_helpers`) walls off the
# `debra` import so the `ccMulti` token here resolves unambiguously
# to `lockfree/internal/pinscope_stub.ccMulti` when the Queue
# type is instantiated below.
manager = initMultiConsumerManager[MaxThreads]()
queue = newUnboundedMpmcQueue[ptr Job, stEager, SegmentSize, MaxThreads](addr manager)
running: Atomic[bool]
jobsSubmitted: array[NumSubmitters, Atomic[int]]
jobsCompleted: array[NumWorkers, Atomic[int]]
nextJobId: Atomic[int]
proc submitterThread(ctx: ptr SubmitterContext) {.thread.} =
## Job submitter - creates and enqueues jobs.
{.cast(gcsafe).}:
# Sugar: `getProducerHere()` combines `getProducer()` + `attach()`,
# registering THIS submitter thread with debra before any push.
var producer = ctx.queue[].getProducerHere()
var submitted = 0
for i in 0 ..< JobsPerSubmitter:
let jobId = nextJobId.fetchAdd(1, moRelaxed)
# Varied job characteristics
let priority =
case rand(10)
of 0: pHigh
of 1 .. 3: pNormal
else: pLow
let workMs =
case priority
of pHigh:
rand(1 .. 5)
of pNormal:
rand(5 .. 15)
of pLow:
rand(10 .. 30)
# Allocate Job on the heap; the consumer frees after work completes.
# `create` returns a non-nil `ptr Job`; `Option[ptr Job]` cannot
# transport nil through pop (design §11.2 guard), so we never push
# a nil pointer here.
let job = create(Job)
job[] =
Job(id: jobId, submitterId: ctx.submitterId, priority: priority, workMs: workMs)
producer.push(job)
inc submitted
# Variable submission rate
sleep(rand(1 .. 10))
jobsSubmitted[ctx.submitterId].store(submitted, moRelease)
proc workerThread(ctx: ptr WorkerContext) {.thread.} =
## Worker - fetches and executes jobs.
{.cast(gcsafe).}:
# Sugar: `getConsumerHere()` combines `getConsumer()` + `attach()`,
# registering THIS worker thread with debra before any pop.
var consumer = ctx.queue[].getConsumerHere()
var completed = 0
var workTime: int64 = 0
while running.load(moAcquire) or ctx.queue[].len() > 0:
let job = consumer.pop()
if job.isSome:
let jp = job.get
# Simulate work
let start = getMonoTime()
sleep(jp.workMs)
workTime += (getMonoTime() - start).inMilliseconds
# Free the producer-allocated Job now that work is done.
dealloc(jp)
inc completed
else:
sleep(1)
jobsCompleted[ctx.workerId].store(completed, moRelease)
echo "Worker ",
ctx.workerId, " completed ", completed, " jobs (", workTime, "ms work)"
when isMainModule:
randomize()
echo "Job Scheduler Example"
echo "====================="
echo "Segment size: ", SegmentSize
echo "Submitters: ", NumSubmitters
echo "Workers: ", NumWorkers
echo "Jobs per submitter: ", JobsPerSubmitter
echo "Total jobs: ", NumSubmitters * JobsPerSubmitter
echo ""
running.store(true, moRelease)
nextJobId.store(0, moRelaxed)
for i in 0 ..< NumSubmitters:
jobsSubmitted[i].store(0, moRelaxed)
for i in 0 ..< NumWorkers:
jobsCompleted[i].store(0, moRelaxed)
let startTime = getMonoTime()
# Start workers
var workerContexts: array[NumWorkers, WorkerContext]
for i in 0 ..< NumWorkers:
workerContexts[i] = WorkerContext(queue: addr queue, workerId: i)
var workers: array[NumWorkers, Thread[ptr WorkerContext]]
for i in 0 ..< NumWorkers:
createThread(workers[i], workerThread, addr workerContexts[i])
# Start submitters
var submitterContexts: array[NumSubmitters, SubmitterContext]
for i in 0 ..< NumSubmitters:
submitterContexts[i] = SubmitterContext(queue: addr queue, submitterId: i)
var submitters: array[NumSubmitters, Thread[ptr SubmitterContext]]
for i in 0 ..< NumSubmitters:
createThread(submitters[i], submitterThread, addr submitterContexts[i])
# Wait for submitters to finish
for i in 0 ..< NumSubmitters:
joinThread(submitters[i])
echo "All jobs submitted, waiting for workers..."
echo ""
# Signal shutdown and wait for workers
sleep(100) # Let workers drain
running.store(false, moRelease)
for i in 0 ..< NumWorkers:
joinThread(workers[i])
let totalTime = (getMonoTime() - startTime).inMilliseconds
# Report results
echo ""
echo "Submission summary:"
var totalSubmitted = 0
for i in 0 ..< NumSubmitters:
let submitted = jobsSubmitted[i].load(moAcquire)
echo " Submitter ", i, ": ", submitted, " jobs"
totalSubmitted += submitted
echo ""
echo "Completion summary:"
var totalCompleted = 0
for i in 0 ..< NumWorkers:
let completed = jobsCompleted[i].load(moAcquire)
totalCompleted += completed
echo ""
echo "Total submitted: ", totalSubmitted
echo "Total completed: ", totalCompleted
echo "Queue remaining: ", queue.len()
echo "Segments used: ", queue.segmentCount()
echo "Total time: ", totalTime, "ms"
if totalSubmitted == totalCompleted:
echo ""
echo "All jobs completed successfully!"
Pattern Guide¶
| Pattern | Shape | Constructor | Example |
|---|---|---|---|
| Real-time audio/video | Bounded SPSC | newSpscQueue |
Audio buffer |
| Work distribution | Bounded SPMC | newSpmcQueue |
Task fan-out |
| Event aggregation | Unbounded MPSC | newUnboundedMpscQueue |
Event collector |
| Job scheduling | Unbounded MPMC | newUnboundedMpmcQueue |
Job scheduler |
| Sensor data | Bounded SPSC/MPSC | newSpscQueue / newMpscQueue |
- |
| Request routing | Bounded SPMC | newSpmcQueue |
Task fan-out |
| Log collection | Unbounded MPSC | newUnboundedMpscQueue |
Event collector |
| Thread pool | Unbounded MPMC | newUnboundedMpmcQueue |
Job scheduler |