~/bend-docscommunity

http/admission.bend source

http/admission.bend on the hub · documented module

# Private bounded transport accounting. Replies contain copied data, never resources.import Baseimport ../resources/resources.bend as Restype 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)}  Quiesce{reply: Chan(Bool)}  Running{reply: Chan(Bool)}  Stop{}# One transport resource: its pool, and every lease taken from it joined in one grant.type Lease is Type:  Lease{pool: Res.Pool, held: Res.Grant}type Marks is Data:  Marks{connection_high: U32, active_high: U32, buffered_high: U32,    connection_rejected: Nat, active_rejected: Nat, buffered_rejected: Nat}type Ledger is Type:  Ledger{connections: Lease, active: Lease, buffered: Lease, marks: Marks}# Broken means a lease count is wrong, such as a give with no matching take.type Booked is Type:  Booked{ledger: Ledger, admitted: Bool}  Broken{}type Taken is Type:  Admitted{lease: Lease}  Full{lease: Lease}  Lost{}def lease.open(+owner: Nat, capacity: U32) -> Lease:  Lease{Res.new(owner, U32.to_nat(capacity)), Res.Grant{owner, 0n}}def ledger.open(caps: Caps) -> Ledger:  Caps{cc, ac, bc} = caps  Ledger{lease.open(0n, cc), lease.open(1n, ac), lease.open(2n, bc), Marks{0, 0, 0, 0n, 0n, 0n}}def lease.joined(combined: Res.Combine, pool: Res.Pool) -> Taken:  match combined:    case Res.Joined{grant}:      Admitted{Lease{pool, grant}}    case Res.Apart{first, second}:      Lost{}def lease.took(reserved: Res.Reserve, held: Res.Grant) -> Taken:  match reserved:    case Res.Took{pool, grant}:      lease.joined(Res.combine(held, grant), pool)    case Res.Refused{pool}:      Full{Lease{pool, held}}def lease.take(lease: Lease) -> Taken:  Lease{pool, held} = lease  lease.took(Res.reserve(pool, 1n), held)def lease.released(released: Res.Release, rest: Res.Grant) -> Maybe<&1, Lease>:  match released:    case Res.Released{pool}:      Some{Lease{pool, rest}}    case Res.Kept{pool, grant}:      None{}def lease.split(parts: Res.Split, pool: Res.Pool) -> Maybe<&1, Lease>:  match parts:    case Res.Parts{one, rest}:      lease.released(Res.release(pool, one), rest)    case Res.Whole{grant}:      None{}def lease.give(lease: Lease) -> Maybe<&1, Lease>:  Lease{pool, held} = lease  lease.split(Res.split(held, 1n), pool)def count(+pool: Res.Pool) -> U32:  U32.from_nat(Res.reserved(pool))def counts(ledger: Ledger) -> Counts & Ledger:  Ledger{Lease{+cp, cg}, Lease{+ap, ag}, Lease{+bp, bg}, +marks} = ledger  Marks{ch, ah, bh, cr, ar, br} = marks  (Counts{count(cp), count(ap), count(bp), ch, ah, bh, cr, ar, br}, Ledger{Lease{cp, cg}, Lease{ap, ag}, Lease{bp, bg}, marks})def final(ledger: Ledger) -> Counts:  Pair.fst(Counts, Ledger, counts(ledger))# Raises each high-water mark to its current count.def marked(ledger: Ledger) -> Ledger:  Ledger{Lease{+cp, cg}, Lease{+ap, ag}, Lease{+bp, bg}, marks} = ledger  Marks{ch, ah, bh, cr, ar, br} = marks  Ledger{Lease{cp, cg}, Lease{ap, ag}, Lease{bp, bg},    Marks{U32.max(ch, count(cp)), U32.max(ah, count(ap)), U32.max(bh, count(bp)), cr, ar, br}}def refused(kind: U32, marks: Marks) -> Marks:  match kind marks:    case 0 Marks{ch, ah, bh, cr, ar, br}:      Marks{ch, ah, bh, 1n+cr, ar, br}    case 1 Marks{ch, ah, bh, cr, ar, br}:      Marks{ch, ah, bh, cr, 1n+ar, br}    case _ Marks{ch, ah, bh, cr, ar, br}:      Marks{ch, ah, bh, cr, ar, 1n+br}# A refused buffer returns the connection lease it already took.def acquire.rollback(connections: Maybe<&1, Lease>, active: Lease, buffered: Lease, marks: Marks) -> Booked:  match connections:    case None{}:      Broken{}    case Some{lease}:      Booked{Ledger{lease, active, buffered, refused(2, marks)}, False{}}def acquire.buffered(taken: Taken, connections: Lease, active: Lease, marks: Marks) -> Booked:  match taken:    case Admitted{buffered}:      Booked{marked(Ledger{connections, active, buffered, marks}), True{}}    case Full{buffered}:      acquire.rollback(lease.give(connections), active, buffered, marks)    case Lost{}:      Broken{}# A connection holds one connection unit and one buffer unit, taken together.def acquire.connection(taken: Taken, active: Lease, buffered: Lease, marks: Marks) -> Booked:  match taken:    case Admitted{connections}:      acquire.buffered(lease.take(buffered), connections, active, marks)    case Full{connections}:      Booked{Ledger{connections, active, buffered, refused(0, marks)}, False{}}    case Lost{}:      Broken{}def acquire.active(taken: Taken, connections: Lease, buffered: Lease, marks: Marks) -> Booked:  match taken:    case Admitted{active}:      Booked{marked(Ledger{connections, active, buffered, marks}), True{}}    case Full{active}:      Booked{Ledger{connections, active, buffered, refused(1, marks)}, False{}}    case Lost{}:      Broken{}def acquire(kind: U32, ledger: Ledger) -> Booked:  match kind:    case 0:      Ledger{connections, active, buffered, marks} = ledger      acquire.connection(lease.take(connections), active, buffered, marks)    case _:      Ledger{connections, active, buffered, marks} = ledger      acquire.active(lease.take(active), connections, buffered, marks)def acquire.running(running: Bool, kind: U32, ledger: Ledger) -> Booked:  match running:    case False{}:      Booked{ledger, False{}}    case True{}:      acquire(kind, ledger)def release.connection(connections: Maybe<&1, Lease>, buffered: Maybe<&1, Lease>, active: Lease, marks: Marks) -> Booked:  match connections buffered:    case Some{c} Some{b}:      Booked{Ledger{c, active, b, marks}, True{}}    case _ _:      Broken{}def release.active(active: Maybe<&1, Lease>, connections: Lease, buffered: Lease, marks: Marks) -> Booked:  match active:    case Some{a}:      Booked{Ledger{connections, a, buffered, marks}, True{}}    case None{}:      Broken{}def release(kind: U32, ledger: Ledger) -> Booked:  match kind:    case 0:      Ledger{connections, active, buffered, marks} = ledger      release.connection(lease.give(connections), lease.give(buffered), active, marks)    case _:      Ledger{connections, active, buffered, marks} = ledger      release.active(lease.give(active), connections, buffered, marks)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)def idle(closing: Bool, counted: Counts & Ledger) -> Bool & Ledger:  (counts, ledger) = counted  (stopped(closing, counts), ledger)type Step is Type:  Waiting{ledger: Ledger}  Received{ledger: Ledger, command: Maybe<&1, Command>}  Acquired{reply: Chan(Bool), booked: Booked}  Released{reply: Chan(Unit), booked: Booked}  Counted{reply: Chan(Counts), counted: Counts & Ledger}  Finishing{checked: Bool & Ledger}# The private inbox is reusable: IO.join would close it after one receive.@unsafedef loop(+inbox: Chan(Command), +running: Bool, +closing: Bool, step: Step) -> IO(Counts):  match step:    case Waiting{ledger}:      do IO<Counts>:        got : Maybe<&1, Command> <- Chan.recv(Command, inbox)        loop(inbox, running, closing, Received{ledger, got})    case Acquired{reply, booked}:      match booked:        case Broken{}:          IO.die(Counts, 1, "Admission: a lease count broke")        case Booked{ledger, admitted}:          do IO<Counts>:            sent : Bool <- Chan.send(Bool, reply, admitted)            loop(inbox, running, closing, Waiting{ledger})    case Released{reply, booked}:      match booked:        case Broken{}:          IO.die(Counts, 1, "Admission.give: no lease of this kind is held")        case Booked{ledger, admitted}:          do IO<Counts>:            sent : Bool <- Chan.send(Unit, reply, Unit{})            loop(inbox, running, closing, Finishing{idle(closing, counts(ledger))})    case Counted{reply, counted}:      (snapshot, ledger) = counted      do IO<Counts>:        sent : Bool <- Chan.send(Counts, reply, snapshot)        loop(inbox, running, closing, Waiting{ledger})    case Finishing{checked}:      (done, ledger) = checked      match done:        case True{}:          do IO<Counts>:            Chan.close(Command, inbox)            return final(ledger)        case False{}:          loop(inbox, running, closing, Waiting{ledger})    case Received{ledger, got}:      match got:        case None{}:          IO.pure(Counts, final(ledger))        case Some{command}:          match command:            case Acquire{kind, reply}:              loop(inbox, running, closing, Acquired{reply, acquire.running(running, kind, ledger)})            case Release{kind, reply}:              loop(inbox, running, closing, Released{reply, release(kind, ledger)})            case Snapshot{reply}:              loop(inbox, running, closing, Counted{reply, counts(ledger)})            case Quiesce{reply}:              do IO<Counts>:                sent : Bool <- Chan.send(Bool, reply, running)                loop(inbox, False{}, closing, Waiting{ledger})            case Running{reply}:              do IO<Counts>:                sent : Bool <- Chan.send(Bool, reply, running)                loop(inbox, running, closing, Waiting{ledger})            case Stop{}:              loop(inbox, running, True{}, Finishing{idle(True{}, counts(ledger))})def open(caps: Caps) -> IO(Chan(Command)):  do IO<Chan(Command)>:    +inbox : Chan(Command) <- Chan.new(Command, 0)    IO.spawn(Counts, loop(inbox, True{}, False{}, Waiting{ledger.open(caps)}))    return inbox# The affine server owner joins this actual actor result after closing its listener.def open.owner(caps: Caps) -> IO(Chan(Command) & Chan(Counts)):  do IO<Chan(Command) & Chan(Counts)>:    +inbox : Chan(Command) <- Chan.new(Command, 0)    finished : Chan(Counts) <- IO.fork(Counts, loop(inbox, True{}, False{}, Waiting{ledger.open(caps)}))    return (inbox, finished)def boolean.sent(sent: Bool, +reply: Chan(Bool)) -> IO(Bool):  match sent:    case False{}:      do IO<Bool>:        Chan.close(Bool, reply)        return False{}    case True{}:      IO.join(Bool, reply)# False also covers copied controls used after actual controller closure.def quiesce(+inbox: Chan(Command)) -> IO(Bool):  do IO<Bool>:    +reply : Chan(Bool) <- Chan.new(Bool, 1)    sent : Bool <- Chan.send(Command, inbox, Quiesce{reply})    boolean.sent(sent, reply)def running(+inbox: Chan(Command)) -> IO(Bool):  do IO<Bool>:    +reply : Chan(Bool) <- Chan.new(Bool, 1)    sent : Bool <- Chan.send(Command, inbox, Running{reply})    boolean.sent(sent, reply)def 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})    boolean.sent(sent, reply)# Correct owners cannot retire the actor while a lease remains counted. If an# internal channel is externally closed, report the violation instead of joining# an unanswered reply or pretending a retained operation drained.def give.sent(sent: Bool, +reply: Chan(Unit)) -> IO(Unit):  match sent:    case False{}:      do IO<Unit>:        Chan.close(Unit, reply)        IO.die(Unit, 1, "Admission.give: controller closed while lease retained")    case True{}:      IO.join(Unit, 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})    give.sent(sent, 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)# Retire accounting after listener closure; quiesce is the separate admission# cutover. Finite internal fixtures may finish their already-accepted protocol.def close(inbox: Chan(Command)) -> IO(Unit):  do IO<Unit>:    sent : Bool <- Chan.send(Command, inbox, Stop{})    return Unit{}