Concurrency
Beans-এর thread হলো OS thread। thread.spawn একটা closure-কে নতুন একটা
OS thread-এ চালায়, আর এই model-এ কোনো green thread বা coroutine নেই।
সহযোগিতামূলক, একক-thread-এর concurrency চাইলে বরং দেখুন async আর
await; thread.spawn হলো CPU-ভারী বা blocking কাজের জন্য
tool। closure আর std.thread মিলেই পুরো কাজটা সেরে দেয়।
import std.thread
// spawn: run a closure on another threadlet t: Thread<int> = thread.spawn(fn() -> int { return heavy_work()})let n: int = t.join() // wait, then take the value
// a mutex wraps the data itself, so you cannot touch it without the locklet ledger: Mutex<Ledger> = new Mutex(new Ledger())ledger.with_lock(fn(l: Ledger) { l.post(entry) // locked for exactly this block, auto-unlock})
// channels move work between threadslet ch: Channel<string> = new Channel(64) // bufferedch.send("job")let job: Option<string> = ch.receive() // none when closed and empty
// atomics for plain counterslet hits: AtomicInt = new AtomicInt(0)hits.add_and_get(1)spawn আর join
“spawn আর join” সেকশনthread.spawn(fn() -> T) একটা closure-কে নতুন OS thread-এ চালায় আর একটা
Thread<T> ফেরত দেয়। join() thread-টা শেষ হওয়া পর্যন্ত অপেক্ষা করে আর তার
value ফেরত দেয়। ফেরত আসা value T-কে Send হতে হবে, আর closure যা যা
capture করে সেগুলোকেও।
Send আর Sync
“Send আর Sync” সেকশনtype system data race ঘটার আগেই সেটা থামিয়ে দেয়। কোনো thread.spawn
closure শুধু Send value capture করতে পারে আর অবশ্যই একটা Send value
return করতে হয়।
Sendনা: সাধারণ class reference,List,Map,Box,Arena,Bytes,File,MMap। এগুলো local reference value।- পার হতে পারে: scalar, immutable string,
AtomicInt,Mutex,Sendvalue-এর একটাChannel, আরShared<T>/Weak<T>যেখানেTহলোSend & Sync।
এতেই একটা class default-এ local reference হয়ে যায়। spawn করা কোনো
closure-এ non-Send value capture করলে সেটা compile error — যে value আর তার
type, দুটোই বলে দেয়:
import std.thread
fn main() { var xs: List<int> = [1, 2, 3] let t: Thread<int> = thread.spawn(fn() -> int { return xs.len() // error: cannot capture non-Send List<int> }) t.join()}thread-জুড়ে mutable data ভাগ করতে চাইলে সেটাকে একটা Mutex-এ মুড়ে নিতে হয় —
Mutex Send। দেখুন Memory আর ownership।
Mutex
“Mutex” সেকশনMutex<T> value-টাকে নিজের ভেতরেই ধরে রাখে। with_lock lock করে,
closure-টা চালায়, আর যেকোনো পথে বের হলেই unlock করে দেয় — তাই unlock করতে ভুলে
যাওয়ার কোনো সুযোগ নেই। value-টা শুধু closure-এর ভেতরেই নাগালে, আর এই কারণেই
lock এড়ানো অসম্ভব।
Channel
“Channel” সেকশনChannel<T> thread-এর মাঝে value চালাচালি করে। new Channel(capacity)
একটা buffered channel বানায়; send(x) একটা value ভেতরে দেয়, receive()
some(v) ফেরত দেয়, অথবা channel বন্ধ ও খালি হলে none, আর close() সেটা
বন্ধ করে। channel-এর সাথে defer ch.close() জুড়ে দিলে যেকোনো পথে বের হলেই
সেটা বন্ধ হয়ে যায়।
Atomic
“Atomic” সেকশনসাধারণ counter আর flag-এর জন্য atomic ব্যবহার করা হয়:
AtomicIntহলো একটা সরল sequentially-consistent integer:new AtomicInt(0),load(),store(v),add_and_get(v)।Atomic<T>হলো যেকোনো integer বাbool-এর উপর একটা typed atomic — সাথে স্পষ্ট memory order আরcompare_exchange,fetch_add, wait/notify আরAtomic.fence। দেখুন Atomics আর MemoryOrder।
একটা পুরো program
“একটা পুরো program” সেকশনএকটা worker spawn করা হয়, shared state-কে Mutex দিয়ে পাহারা দেওয়া হয়, আর একটা
value একটা Channel-এর ভেতর দিয়ে পাঠানো হয়:
import std.ioimport std.thread
class Ledger { total: int = 0 fn add(n: int) { self.total += n }}
fn main() { let t: Thread<int> = thread.spawn(fn() -> int { return 21 * 2 }) io.println("worker said {t.join()}")
let ledger: Mutex<Ledger> = new Mutex(new Ledger()) ledger.with_lock(fn(l: Ledger) { l.add(10) })
let ch: Channel<int> = new Channel(4) ch.send(7) ch.close() match ch.receive() { some(v) => io.println("got {v}"), none => io.println("empty"), }
let hits: AtomicInt = new AtomicInt(0) hits.add_and_get(1) io.println("hits {hits.load()}")}Readiness wait
“Readiness wait” সেকশনকোনো thread block না করে একটা socket-এর জন্য অপেক্ষা করতে চাইলে
async model-টা std.net-এর readable/writable-এর সাথে
ব্যবহার করুন; আর একসাথে অনেক descriptor-এর জন্য
std.poll poller ব্যবহার করুন।