~/bend-docscommunity

concurrency.bend checks

raw source on the hub · import bend-kit-concurrency@0.2.0.0/concurrency.bend as Concurrency

Parallel map and reduce over lists and arrays, a worker pool, select over channels, and timeouts.

1 import
import Base

Types

type Worker source · line 145 · raw

@-S:Type -> @-A:Type -> @-B:Type -> Type

A pool worker between channel steps. It owns its state s alone.

type Relay source · line 244 · raw

@-A:Type -> Type

A select relay between channel steps.

Definitions

def count.cons source · line 18 · raw

@-A:Type -> @h:A -> @r:Pair(Nat, List<&1, A>) -> Pair(Nat, List<&1, A>)

The length of xs, beside xs.

def count source · line 22 · raw

@-A:Type -> @xs:List<&1, A> -> Pair(Nat, List<&1, A>)

def split.cons source · line 30 · raw

@-A:Type -> @h:A -> @r:Pair(List<&1, A>, List<&1, A>) -> Pair(List<&1, A>, List<&1, A>)

The first n elements of xs, and the rest.

def split source · line 34 · raw

@-A:Type -> @xs:List<&1, A> -> @n:Nat -> Pair(List<&1, A>, List<&1, A>)

def pool.join source · line 205 · raw

@-A:Type -> @cs:List<&1, Chan(A)> -> IO(List<&1, A>)

def select.relay source · line 251 · raw

@-A:Type -> @fuel:Nat -> @+i:U32 -> @+src:Chan(A) -> @+out:Chan(Pair(U32, A)) -> @+done:Chan(Unit) -> @r:Relay<A> -> IO(Result<&1, &1, Unit, Unit>)

Forwards each value of src to out as (i, value), then reports on done. It stops when src closes, or when out is closed under it.

def select.start source · line 282 · raw

@-A:Type -> @chans:List<&1, Chan(A)> -> @+i:U32 -> @+out:Chan(Pair(U32, A)) -> @+done:Chan(Unit) -> IO(Nat)

ponytail: relay fuel is the largest Nat literal, and a value takes three steps, so a relay stops after about 1.4 billion values. Drop the fuel if Base gains an unbounded loop.

def select.close source · line 295 · raw

@-A:Type -> @n:Nat -> @+out:Chan(Pair(U32, A)) -> @+done:Chan(Unit) -> IO(Unit)

Closes out once all n relays have reported.

def select.go source · line 306 · raw

@-A:Type -> @chans:List<&1, Chan(A)> -> @+out:Chan(Pair(U32, A)) -> @+done:Chan(Unit) -> IO(Chan(Pair(U32, A)))

def select source · line 318 · raw

@-A:Type -> @chans:List<&1, Chan(A)> -> IO(Chan(Pair(U32, A)))

One channel that carries (i, value) for each value received on chans[i], in the order they arrive. Chan.recv on it waits for whichever source is ready first, and answers None once every source is closed. The selector owns the receiving side of its sources. Close it to stop early: each relay then drops the one value it holds, if any, and stops.

def timeout.won source · line 324 · raw

@-A:Type -> @+c:Chan(Maybe<&1, A>) -> @got:Maybe<&1, Maybe<&1, A>> -> IO(Maybe<&1, A>)

def timeout source · line 337 · raw

@-A:Type -> @ms:U32 -> @act:IO(A) -> IO(Maybe<&1, A>)

Some{result} if act finishes within ms milliseconds, else None at the deadline. act is not cancelled: it runs on, its late result is dropped, and the program waits for it (and for the timer) before it exits.

Templates

template par_map.list.go source · line 44 · raw

@-A:Type -> @-B:Type -> @-f:(@_:A -> B) -> @d:Nat -> @+nl:Nat -> @+nr:Nat -> @lr:Pair(List<&1, A>, List<&1, A>) -> List<&1, B>

lr holds two halves of nl and nr elements. Depth d runs 2^(d+1) parallel parts.

template par_map.list.top source · line 61 · raw

@-A:Type -> @-B:Type -> @-f:(@_:A -> B) -> @d:Nat -> @nx:Pair(Nat, List<&1, A>) -> List<&1, B>

template par_map source · line 72 · raw

@-A:Type -> @-B:Type -> @-f:(@_:A -> B) -> @workers:U32 -> @xs:List<&1, A> -> List<&1, B>

List.map(f, xs), on at most workers parallel parts.

template fold source · line 75 · raw

@-A:Data -> @-f:(@_:A -> @_:A -> A) -> @xs:List<&1, A> -> @acc:A -> A

template par_reduce.list.go source · line 82 · raw

@-A:Data -> @-f:(@_:A -> @_:A -> A) -> @d:Nat -> @+nl:Nat -> @+nr:Nat -> @+z:A -> @lr:Pair(List<&1, A>, List<&1, A>) -> A

template par_reduce.list.top source · line 99 · raw

@-A:Data -> @-f:(@_:A -> @_:A -> A) -> @d:Nat -> @+z:A -> @nx:Pair(Nat, List<&1, A>) -> A

template par_reduce source · line 111 · raw

@-A:Data -> @-f:(@_:A -> @_:A -> A) -> @workers:U32 -> @+z:A -> @xs:List<&1, A> -> A

The left fold of xs from z, on at most workers parallel parts. f must be associative and z its identity.

template par_map.array.go source · line 115 · raw

@-A:Type -> @-B:Type -> @-f:(@_:A -> B) -> @a:Array<A> -> @d:Nat -> Array<B>

An Array is a balanced tree, so each ANode above depth d is a parallel call.

template par_map.array source · line 126 · raw

@-A:Type -> @-B:Type -> @-f:(@_:A -> B) -> @workers:U32 -> @a:Array<A> -> Array<B>

f on every slot of a, on at most workers parallel parts. Slots keep their index.

template par_reduce.array.go source · line 129 · raw

@-A:Type -> @-f:(@_:A -> @_:A -> A) -> @a:Array<A> -> @d:Nat -> A

template par_reduce.array source · line 141 · raw

@-A:Type -> @-f:(@_:A -> @_:A -> A) -> @workers:U32 -> @a:Array<A> -> A

The slots of a combined in index order with an associative f, on at most workers parallel parts. An Array is never empty, so no identity is needed.

template pool.worker source · line 152 · raw

@-S:Type -> @-A:Type -> @-B:Type -> @-work:(@_:S -> @_:A -> IO(Pair(S, B))) -> @fuel:Nat -> @+q:Chan(Pair(A, Chan(B))) -> @w:Worker<S, A, B> -> IO(S)

Takes jobs until the queue is closed and empty, then answers its state. Each job takes three steps, so 3 * (jobs + 1) fuel is enough.

template pool.put source · line 180 · raw

@-A:Type -> @-B:Type -> @xs:List<&1, A> -> @+q:Chan(Pair(A, Chan(B))) -> IO(List<&1, Chan(B)>)

Queues each job with its own reply channel, and answers the replies in order.

template pool.start source · line 192 · raw

@-S:Type -> @-A:Type -> @-B:Type -> @-work:(@_:S -> @_:A -> IO(Pair(S, B))) -> @states:List<&1, S> -> @+fuel:Nat -> @+q:Chan(Pair(A, Chan(B))) -> IO(List<&1, Chan(S)>)

template pool.run source · line 215 · raw

@-S:Type -> @-A:Type -> @-B:Type -> @-work:(@_:S -> @_:A -> IO(Pair(S, B))) -> @states:List<&1, S> -> @nx:Pair(Nat, List<&1, A>) -> IO(Pair(List<&1, S>, List<&1, B>))

template pool source · line 234 · raw

@-S:Type -> @-A:Type -> @-B:Type -> @-work:(@_:S -> @_:A -> IO(Pair(S, B))) -> @states:List<&1, S> -> @xs:List<&1, A> -> IO(Pair(List<&1, S>, List<&1, B>))

Runs work(s, x) for every x in xs, on one worker per state in states. A worker owns its state and threads it through its jobs, so an affine handle such as a connection pool is never shared. Answers the final states in worker order and the results in the order of xs.