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:

- Multi-Core Workers: Incoming network sockets and background workers execute in parallel across CPU cores.
- Atomic Memory Channels: Results and event payloads are passed safely over lock-free message channels.
- 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 resultlet result = await handleprint($"Worker computed result: {result}")Concurrency & Channels (std.thread)
Section titled “Concurrency & Channels (std.thread)”Manage thread execution, sleep durations, and pass messages cleanly across concurrent threads using channels:
import std.thread
// Display the current OS thread identifierprintln($"Active thread ID: {thread.id()}")
// Pause execution of the current thread for 200 millisecondsthread.sleep(200)
// Yield execution time back to the operating system schedulerthread.yield()
// Create an asynchronous communication channellet (tx, rx) = thread.channel()
// Send and receive messages across thread channelstx.send("Message across channel!")let received = rx.recv()println($"Received: {received}")Thread Safety & Panic Isolation
Section titled “Thread Safety & Panic Isolation”- 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.
Thread Block Syntax & Aliases
Section titled “Thread Block Syntax & Aliases”You can spawn threads using the keyword thread { ... } or through an imported alias:
import std.thread as th
// Using the aliased block syntaxth { println("Running inside aliased thread block!") th.yield()}Channels & Message Passing
Section titled “Channels & Message Passing”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 1th { th.yield() tx.send("Result from Worker 1")}
// Spawn Worker 2th { th.yieldNow() tx2.send("Result from Worker 2")}
// Main thread collects both resultslet msg1 = rx.recv()let msg2 = rx.recv()
println($"Collected: {msg1}")println($"Collected: {msg2}")Non-blocking Polling (tryRecv & isEmpty)
Section titled “Non-blocking Polling (tryRecv & isEmpty)”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 receivelet maybe_msg = rx.tryRecv()if maybe_msg != nil { println($"Processed immediate message: {maybe_msg}")} else { println("No messages pending; proceeding with frame update...")}Real-World Applications
Section titled “Real-World Applications”1. Parallel Map-Reduce Worker Pool
Section titled “1. Parallel Map-Reduce Worker Pool”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 workersfor id in 0..total_workers { let worker_tx = tx.clone() process_chunk(id, worker_tx)}
// Aggregate responses from all workerslet mut total_sum = 0for _ 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 collectorth { for i in 1..=5 { th.sleep(20) event_tx.send($"Telemetry event #{i}") }}
// Consumer: Non-blocking drain loopth.sleep(150) // Allow events to arrive
while !event_rx.isEmpty() { let event = event_rx.recv() println($"Flushing to analytics: {event}")}API Reference
Section titled “API Reference”Threading (std.thread)
Section titled “Threading (std.thread)”| 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. |
Sender Methods (Sender)
Section titled “Sender Methods (Sender)”| 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. |
Receiver Methods (Receiver)
Section titled “Receiver Methods (Receiver)”| 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. |
