Skip to content

Examples

Practical examples demonstrating lock-free queue patterns.

Running Examples

nimble 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