~/bend-docscommunity

admission.bend source

admission.bend on the hub · documented module

# Private bounded transport accounting. Replies contain copied data, never resources.import Basetype Counts is Data:  Counts{connections: U32, active: U32, buffered: U32,    connection_high: U32, active_high: U32, buffered_high: U32,    connection_rejected: Nat, active_rejected: Nat, buffered_rejected: Nat}type Caps is Data:  Caps{connections: U32, active: U32, buffered: U32}type Command is Data:  Acquire{kind: U32, reply: Chan(Bool)}  Release{kind: U32, reply: Chan(Unit)}  Snapshot{reply: Chan(Counts)}  Stop{}def empty() -> Counts:  Counts{0, 0, 0, 0, 0, 0, 0n, 0n, 0n}def acquire.connection(full: Bool, buffer_full: Bool, counts: Counts) -> Counts & Bool:  match full buffer_full counts:    case True{} _ Counts{c, a, b, ch, ah, bh, cr, ar, br}:      (Counts{c, a, b, ch, ah, bh, 1n+cr, ar, br}, False{})    case False{} True{} Counts{c, a, b, ch, ah, bh, cr, ar, br}:      (Counts{c, a, b, ch, ah, bh, cr, ar, 1n+br}, False{})    case False{} False{} Counts{c, a, b, ch, ah, bh, cr, ar, br}:      +nc = (c + 1 : U32)      +nb = (b + 1 : U32)      (Counts{nc, a, nb, U32.max(ch, nc), ah, U32.max(bh, nb), cr, ar, br}, True{})def acquire.active(full: Bool, counts: Counts) -> Counts & Bool:  match full counts:    case True{} Counts{c, a, b, ch, ah, bh, cr, ar, br}:      (Counts{c, a, b, ch, ah, bh, cr, 1n+ar, br}, False{})    case False{} Counts{c, a, b, ch, ah, bh, cr, ar, br}:      +na = (a + 1 : U32)      (Counts{c, na, b, ch, U32.max(ah, na), bh, cr, ar, br}, True{})def acquire(kind: U32, caps: Caps, +counts: Counts) -> Counts & Bool:  match kind:    case 0:      Caps{cc, ac, bc} = caps      Counts{c, a, b, ch, ah, bh, cr, ar, br} = counts      acquire.connection(U32.is_ge(c, cc), U32.is_ge(b, bc), counts)    case _:      Caps{cc, ac, bc} = caps      Counts{c, a, b, ch, ah, bh, cr, ar, br} = counts      acquire.active(U32.is_ge(a, ac), counts)def release(kind: U32, counts: Counts) -> Counts:  match kind counts:    case 0 Counts{c, a, b, ch, ah, bh, cr, ar, br}:      Counts{(c - 1 : U32), a, (b - 1 : U32), ch, ah, bh, cr, ar, br}    case _ Counts{c, a, b, ch, ah, bh, cr, ar, br}:      Counts{c, (a - 1 : U32), b, ch, ah, bh, cr, ar, br}def stopped(closing: Bool, counts: Counts) -> Bool:  Counts{c, a, b, ch, ah, bh, cr, ar, br} = counts  closing && U32.is_zero(c) && U32.is_zero(a) && U32.is_zero(b)type Step is Type:  Waiting{}  Received{command: Maybe<&1, Command>}  Acquired{reply: Chan(Bool), result: Counts & Bool}# The private inbox is reusable: IO.join would close it after one receive.@unsafedef loop(+inbox: Chan(Command), +caps: Caps, +counts: Counts, +closing: Bool, step: Step) -> IO(Unit):  match step:    case Waiting{}:      do IO<Unit>:        got : Maybe<&1, Command> <- Chan.recv(Command, inbox)        loop(inbox, caps, counts, closing, Received{got})    case Acquired{reply, result}:      (next, admitted) = result      do IO<Unit>:        sent : Bool <- Chan.send(Bool, reply, admitted)        loop(inbox, caps, next, closing, Waiting{})    case Received{got}:      match got:        case None{}:          IO.pure(Unit, Unit{})        case Some{command}:          match command:            case Acquire{kind, reply}:              loop(inbox, caps, counts, closing, Acquired{reply, acquire(kind, caps, counts)})            case Release{kind, reply}:              +next = release(kind, counts)              do IO<Unit>:                sent : Bool <- Chan.send(Unit, reply, Unit{})                finish(stopped(closing, next), inbox, caps, next, closing)            case Snapshot{reply}:              do IO<Unit>:                sent : Bool <- Chan.send(Counts, reply, counts)                loop(inbox, caps, counts, closing, Waiting{})            case Stop{}:              finish(stopped(True{}, counts), inbox, caps, counts, True{})@unsafedef finish(done: Bool, +inbox: Chan(Command), caps: Caps, counts: Counts, closing: Bool) -> IO(Unit):  match done:    case True{}:      Chan.close(Command, inbox)    case False{}:      loop(inbox, caps, counts, closing, Waiting{})def open(caps: Caps) -> IO(Chan(Command)):  do IO<Chan(Command)>:    +inbox : Chan(Command) <- Chan.new(Command, 0)    IO.spawn(Unit, loop(inbox, caps, empty(), False{}, Waiting{}))    return inboxdef take(+inbox: Chan(Command), kind: U32) -> IO(Bool):  do IO<Bool>:    +reply : Chan(Bool) <- Chan.new(Bool, 1)    sent : Bool <- Chan.send(Command, inbox, Acquire{kind, reply})    IO.join(Bool, reply)def give(+inbox: Chan(Command), kind: U32) -> IO(Unit):  do IO<Unit>:    +reply : Chan(Unit) <- Chan.new(Unit, 1)    sent : Bool <- Chan.send(Command, inbox, Release{kind, reply})    IO.join(Unit, reply)def snapshot.sent(sent: Bool, +reply: Chan(Counts)) -> IO(Maybe<&2, Counts>):  match sent:    case False{}:      do IO<Maybe<&2, Counts>>:        Chan.close(Counts, reply)        return None{}    case True{}:      do IO<Maybe<&2, Counts>>:        counts : Counts <- IO.join(Counts, reply)        return Some{counts}# A successful rendezvous send was received: Stop cannot erase a queued request.# After controller closure a failed send returns None without waiting on its reply.def snapshot(+inbox: Chan(Command)) -> IO(Maybe<&2, Counts>):  do IO<Maybe<&2, Counts>>:    +reply : Chan(Counts) <- Chan.new(Counts, 1)    sent : Bool <- Chan.send(Command, inbox, Snapshot{reply})    snapshot.sent(sent, reply)def close(inbox: Chan(Command)) -> IO(Unit):  do IO<Unit>:    sent : Bool <- Chan.send(Command, inbox, Stop{})    return Unit{}