import Base import ./types.bend as Types import ./lapine.bend as Lapine import ./retry.bend as Retry def Consumer.acks_only_on_done(+d: Types.Delivery) -> Bool: True{} def Consumer.body_of(+d: Types.Delivery) -> String: match d: case Types.Msg{body, tag, queue}: body def Consumer.queue_of(+sub: Types.Subscription) -> String: match sub: case Types.Sub{queue, dlq_exchange, dlq_routing_key}: queue def Consumer.dlq_exchange_of(+sub: Types.Subscription) -> String: match sub: case Types.Sub{queue, dlq_exchange, dlq_routing_key}: dlq_exchange def Consumer.dlq_key_of(+sub: Types.Subscription) -> String: match sub: case Types.Sub{queue, dlq_exchange, dlq_routing_key}: dlq_routing_key def Consumer.when_empty(flag: Bool) -> Bool: match flag: case True{}: False{} case False{}: True{} def Consumer.should_route_dlq(+sub: Types.Subscription) -> Bool: Consumer.when_empty(String.is_empty(Consumer.dlq_exchange_of(sub))) def Consumer.on_message(d: Types.Delivery) -> IO(Result<&1, &1, U32 & String, Unit>): do IO>: IO.print(Consumer.body_of(d)) return Done{Unit{}} def Consumer.publish_dlq(chan: Nat, +sub: Types.Subscription, +body: String) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): Lapine.AMQP.publish( chan, Consumer.dlq_exchange_of(sub), Consumer.dlq_key_of(sub), body) def Consumer.next_chan(sent: Nat & Result<&1, &1, U32 & String, Unit>) -> IO(Nat): match sent: case (next, Done{Unit{}}): IO.pure(Nat, next) case (next, Fail{(code, message)}): IO.die(Nat, code, message) def Consumer.when_dlq(flag: Bool, chan: Nat, +sub: Types.Subscription, +body: String) -> IO(Nat): match flag: case True{}: do IO: sent : Nat & Result<&1, &1, U32 & String, Unit> <- Consumer.publish_dlq(chan, sub, body) Consumer.next_chan(sent) case False{}: IO.pure(Nat, chan) def Consumer.ack_delivery(chan: Nat, +d: Types.Delivery) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): match d: case Types.Msg{body, tag, queue}: Lapine.AMQP.ack(chan, tag) def Consumer.ack_and_return(chan: Nat, +d: Types.Delivery) -> IO(Nat): do IO: acked : Nat & Result<&1, &1, U32 & String, Unit> <- Consumer.ack_delivery(chan, d) Consumer.next_chan(acked) def Consumer.finish_outcome( chan: Nat, +sub: Types.Subscription, +d: Types.Delivery, outcome: Result<&1, &1, U32 & String, Unit> ) -> IO(Nat): match outcome: case Done{Unit{}}: Consumer.ack_and_return(chan, d) case Fail{_}: do IO: routed : Nat <- Consumer.when_dlq(Consumer.should_route_dlq(sub), chan, sub, Consumer.body_of(d)) Consumer.ack_and_return(routed, d) @unsafe def Consumer.process_attempt( chan: Nat, +sub: Types.Subscription, +d: Types.Delivery ) -> IO(Nat): do IO: outcome : Result<&1, &1, U32 & String, Unit> <- Consumer.on_message(d) Consumer.finish_outcome(chan, sub, d, outcome) @unsafe def Consumer.recv_next( recv: Nat & Result<&1, &1, U32 & String, Types.Delivery>, +sub: Types.Subscription ) -> IO(Nat): match recv: case (next, Done{d}): Consumer.process_attempt(next, sub, d) case (next, Fail{_}): IO.pure(Nat, next) @unsafe def Consumer.handle_once(chan: Nat, +sub: Types.Subscription) -> IO(Nat): do IO: queue : String = Consumer.queue_of(sub) recv : Nat & Result<&1, &1, U32 & String, Types.Delivery> <- Lapine.AMQP.recv(chan, queue) Consumer.recv_next(recv, sub) @unsafe def Consumer.loop(fuel: Nat, chan: Nat, +sub: Types.Subscription) -> IO(Unit): match fuel: case 0n: IO.pure(Unit, Unit{}) case 1n+rest: do IO: next : Nat <- Consumer.handle_once(chan, sub) Consumer.loop(rest, next, sub) def Consumer.channel_from_opened(opened: Nat & Result<&1, &1, U32 & String, Nat>) -> IO(Nat): match opened: case (conn_next, Done{chan}): IO.pure(Nat, chan) case (conn_next, Fail{(code, message)}): IO.die(Nat, code, message) def Consumer.open_channel(conn: Nat, +cfg: Types.Config) -> IO(Nat): do IO: opened : Nat & Result<&1, &1, U32 & String, Nat> <- Lapine.AMQP.channel(conn) chan : Nat <- Consumer.channel_from_opened(opened) q : Nat & Result<&1, &1, U32 & String, Unit> <- Lapine.AMQP.apply_qos(chan, cfg) Consumer.next_chan(q) def Consumer.first_subscription(subs: List) -> Types.Subscription: match subs: case Nil{}: Types.Sub{"", "", ""} case Con{sub, _}: sub @unsafe def Consumer.run( fuel: Nat, +cfg: Types.Config, subs: List, policy: Retry.RetryPolicy ) -> IO(Unit): do IO: conn : Nat <- IO.try(Nat, Lapine.AMQP.connect_cfg(cfg)) chan : Nat <- Consumer.open_channel(conn, cfg) sub : Types.Subscription = Consumer.first_subscription(subs) Consumer.loop(fuel, chan, sub)