Don't use it in production!
///|
test {
let h = spawn(fn() { 40 + 2 })
inspect(h.join(), content="42")
}///|
test {
let (tx, rx) : (Sender[Int], Receiver[Int]) = channel(16)
let tx1 = tx.clone()
tx.destroy()
let h = spawn(fn() {
defer tx1.destroy()
for i in 0..<10 {
tx1.send(i) |> ignore
}
})
let mut sum = 0
let mut cnt = 0
while true {
match rx.recv() {
Some(v) => {
sum v
cnt 1
}
None => break
}
}
rx.destroy()
h.join()
inspect(cnt, content="10")
inspect(sum, content="45")
}///|
test {
let b : BroadcastSender[Int] = broadcast(4)
let r1 = b.subscribe()
let r2 = b.subscribe()
inspect(b.send(1), content="2")
inspect(b.send(2), content="2")
b.destroy()
inspect(r1.recv(), content="Some(1)")
inspect(r1.recv(), content="Some(2)")
inspect(r1.recv(), content="None")
inspect(r2.recv(), content="Some(1)")
inspect(r2.recv(), content="Some(2)")
inspect(r2.recv(), content="None")
r1.destroy()
r2.destroy()
}///|
test {
let pool = ThreadPool::new(4, 64)
defer pool.shutdown()
let rx = pool.submit_with_result(fn() { 40 + 2 })
defer rx.destroy()
inspect(rx.recv(), content="Some(42)")
}///|
test {
let pool = ThreadPool::new(4, 64)
defer pool.shutdown()
let xs : Array[Int] = []
for i in 0..<1000 {
xs.push(i)
}
let cfg = ParConfig::default(pool)
match par_map_collect_unordered(xs.iterator(), pool, cfg, fn(x) { x * 2 }) {
Some(ys) => {
inspect(ys.length(), content="1000")
let mut sum = 0
for y in ys {
sum y
}
inspect(sum, content="999000")
}
None => fail("par_map_collect_unordered failed")
}
}///|
test {
let pool = ThreadPool::new(4, 64)
defer pool.shutdown()
let xs : Array[Int] = []
for i in 0..<1000 {
xs.push(i)
}
let cfg = ParConfig::default(pool)
match
par_filter_collect_unordered(xs.iterator(), pool, cfg, fn(x) { x % 2 == 0 }) {
Some(ys) => {
inspect(ys.length(), content="500")
let mut sum = 0
for y in ys {
sum y
}
inspect(sum, content="249500")
}
None => fail("par_filter_collect_unordered failed")
}
}///|
test {
let pool = ThreadPool::new(4, 64)
defer pool.shutdown()
let xs : Array[Int] = []
for i in 0..<1000 {
xs.push(i)
}
let (tx, rx) : (Sender[Int], Receiver[Int]) = channel(128)
let consumer = spawn(fn() {
defer rx.destroy()
let mut sum = 0
while true {
match rx.recv() {
Some(v) => sum v
None => break
}
}
sum
})
let cfg = ParConfig::default(pool)
let ok = par_each(xs.iterator(), pool, cfg, fn(x) { tx.send(x) |> ignore })
tx.destroy()
let sum = consumer.join()
inspect(ok, content="true")
inspect(sum, content="499500")
}pub struct BroadcastReceiver[T] {
// private fields
}pub struct BroadcastSender[T] {
// private fields
}type Handle[_]pub struct ParConfig {
chunk_size : Int
max_in_flight : Int
}pub struct Receiver[T] {
// private fields
}pub struct Sender[T] {
// private fields
}pub struct ThreadPool {
// private fields
}fn[T, U] par_array_map_reduce(xs : ArrayView[T], pool : ThreadPool, cfg : ParConfig, map : (T) -> U, init : () -> U, reduce : (U, U) -> U) -> U?fn[T] par_filter_collect_unordered(iter : Iterator[T], pool : ThreadPool, cfg : ParConfig, pred : (T) -> Bool) -> Array[T]?fn[T, U] par_map_collect_unordered(iter : Iterator[T], pool : ThreadPool, cfg : ParConfig, f : (T) -> U) -> Array[U]?fn[T, U] par_map_reduce_unordered(iter : Iterator[T], pool : ThreadPool, cfg : ParConfig, map : (T) -> U, reduce : (U, U) -> U) -> U?Don't use it in production!