Skip to content

Threads & Channels (std.thread)

Flame cleanly separates I/O Concurrency (async/await) from Computational Multi-Core Concurrency (thread { ... }).


Architecture: Execution Model vs Concurrency Model

Section titled “Architecture: Execution Model vs Concurrency Model”

Rather than cloning interpreter instances per thread (which causes memory bloat and state divergence), Flame uses Lexical Snapshot Isolation:

Flame Multithreading Architecture

  1. Multi-Core Workers: Incoming network sockets and background workers execute in parallel across CPU cores.
  2. Atomic Memory Channels: Results and event payloads are passed safely over lock-free message channels.
  3. Deterministic State: Flame’s engine evaluates tasks deterministically without mutex deadlock risks.

Spawning Dedicated Compute Threads (thread)

Section titled “Spawning Dedicated Compute Threads (thread)”

To run CPU-intensive calculations in parallel on an independent physical CPU core:

import std.thread
print($"Main Thread ID: {thread.id()}")
let handle = thread {
print($"Worker Thread ID: {thread.id()}")
let mut sum = 0
let mut i = 0
while i < 50000000 {
sum = sum + i
i = i + 1
}
return sum
}
// Main thread continues executing without blocking...
print("Main thread is working...")
// Synchronize and receive the result
let result = await handle
print($"Worker computed result: {result}")

Manage thread execution, sleep durations, and pass messages cleanly across concurrent threads using channels:

import std.thread
// Display the current OS thread identifier
println($"Active thread ID: {thread.id()}")
// Pause execution of the current thread for 200 milliseconds
thread.sleep(200)
// Yield execution time back to the operating system scheduler
thread.yield()
// Create an asynchronous communication channel
let (tx, rx) = thread.channel()
// Send and receive messages across thread channels
tx.send("Message across channel!")
let received = rx.recv()
println($"Received: {received}")

  • Panic Isolation: If a compute thread encounters a runtime error or division-by-zero, the panic is trapped inside the thread boundary. Awaiting the handle resolves cleanly without crashing your main server.
  • Automatic Resource Cleanup: RAII ensures all memory allocated within a worker thread is reclaimed immediately upon completion.
  • Graceful Shutdown: The runtime monitors all active worker threads and waits for clean resolution before process exit.

You can spawn threads using the keyword thread { ... } or through an imported alias:

import std.thread as th
// Using the aliased block syntax
th {
println("Running inside aliased thread block!")
th.yield()
}

Thread channels provide safe, multi-producer, single-consumer (MPSC) lock-free communication between concurrent threads.

import std.thread as th
let (tx, rx) = th.channel()
let tx2 = tx.clone()
// Spawn Worker 1
th {
th.yield()
tx.send("Result from Worker 1")
}
// Spawn Worker 2
th {
th.yieldNow()
tx2.send("Result from Worker 2")
}
// Main thread collects both results
let msg1 = rx.recv()
let msg2 = rx.recv()
println($"Collected: {msg1}")
println($"Collected: {msg2}")

When running real-time services, game loops, or event-driven servers, blocking indefinitely with rx.recv() is undesirable. Use rx.tryRecv() or rx.isEmpty() to inspect the channel without pausing execution:

import std.thread as th
let (tx, rx) = th.channel()
// Attempt non-blocking receive
let maybe_msg = rx.tryRecv()
if maybe_msg != nil {
println($"Processed immediate message: {maybe_msg}")
} else {
println("No messages pending; proceeding with frame update...")
}

Distribute compute chunks across multiple worker threads and aggregate the results back on the main thread:

import std.thread as th
fn process_chunk(chunk_id: Int, tx: Any) {
th {
// Simulate heavy work
th.sleep(50)
let processed_value = chunk_id * 100
tx.send(processed_value)
}
}
let (tx, rx) = th.channel()
let total_workers = 4
// Launch workers
for id in 0..total_workers {
let worker_tx = tx.clone()
process_chunk(id, worker_tx)
}
// Aggregate responses from all workers
let mut total_sum = 0
for _ in 0..total_workers {
let result = rx.recv()
total_sum = total_sum + result
}
println($"All workers completed! Total sum: {total_sum}")

2. Producer-Consumer Pipeline with Non-Blocking Drain

Section titled “2. Producer-Consumer Pipeline with Non-Blocking Drain”

A producer thread continuously pushes telemetry or events, while the main consumer drains and batches items periodically:

import std.thread as th
let (event_tx, event_rx) = th.channel()
// Producer: Background metric collector
th {
for i in 1..=5 {
th.sleep(20)
event_tx.send($"Telemetry event #{i}")
}
}
// Consumer: Non-blocking drain loop
th.sleep(150) // Allow events to arrive
while !event_rx.isEmpty() {
let event = event_rx.recv()
println($"Flushing to analytics: {event}")
}

Method Arguments Returns Description
thread.sleep ms: Int Nil Suspends thread execution for the specified number of milliseconds.
thread.id None String Returns a formatted string representation of the active thread ID.
thread.yield None Nil Yields the processor, allowing other threads on the OS scheduler to run.
thread.yieldNow None Nil Immediately yields execution time slice to allow other threads to progress.
thread.channel None (Sender, Receiver) Creates a multi-producer, single-consumer communication channel.
thread.spawn callback: Function ThreadHandler Spawns a background concurrent execution thread.
Method Arguments Returns Description
tx.send(value) value: Any Nil Sends a message value across the channel to the connected Receiver. Thread-safe.
tx.clone() None Sender Creates a new Sender handle connected to the same Receiver for multi-thread dispatching.
Method Arguments Returns Description
rx.recv() None Any Blocks execution until a message is received from the channel.
rx.tryRecv() None Any | Nil Returns the next message if available, otherwise returns nil immediately without blocking.
rx.isEmpty() None Bool Returns true if no messages are currently queued in the channel.