import Base import ./types.bend as Types def AMQP.connect(host: String, port: U32, user: String, pass: String, vhost: String) -> IO(Result<&1, &1, U32 & String, Nat>): import "./effs/amqp_connect.c" import "./effs/amqp_connect.js" def AMQP.channel(conn: Nat) -> IO(Nat & Result<&1, &1, U32 & String, Nat>): import "./effs/amqp_channel.c" import "./effs/amqp_channel.js" def AMQP.qos(chan: Nat, prefetch: U32) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): import "./effs/amqp_qos.c" import "./effs/amqp_qos.js" def AMQP.consume(chan: Nat, +queue: String) -> IO(Nat & Result<&1, &1, U32 & String, String & Nat>): import "./effs/amqp_consume.c" import "./effs/amqp_consume.js" def AMQP.ack(chan: Nat, tag: Nat) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): import "./effs/amqp_ack.c" import "./effs/amqp_ack.js" def AMQP.publish(chan: Nat, exchange: String, routing_key: String, body: String) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): import "./effs/amqp_publish.c" import "./effs/amqp_publish.js" def AMQP.close(conn: Nat) -> IO(Unit): import "./effs/amqp_close.c" import "./effs/amqp_close.js" def AMQP.connect_cfg(+cfg: Types.Config) -> IO(Result<&1, &1, U32 & String, Nat>): do IO>: AMQP.connect( Types.Config.host(cfg), Types.Config.port(cfg), Types.Config.user(cfg), Types.Config.pass(cfg), Types.Config.vhost(cfg)) def AMQP.apply_qos(chan: Nat, +cfg: Types.Config) -> IO(Nat & Result<&1, &1, U32 & String, Unit>): do IO>: AMQP.qos(chan, Types.Config.prefetch(cfg)) def AMQP.recv_map( got: Nat & Result<&1, &1, U32 & String, String & Nat>, +queue: String ) -> Nat & Result<&1, &1, U32 & String, Types.Delivery>: match got: case (next, Done{(body, tag)}): (next, Done{Types.Msg{body, tag, queue}}) case (next, Fail{e}): (next, Fail{e}) def AMQP.recv(chan: Nat, +queue: String) -> IO(Nat & Result<&1, &1, U32 & String, Types.Delivery>): do IO>: got : Nat & Result<&1, &1, U32 & String, String & Nat> <- AMQP.consume(chan, queue) return AMQP.recv_map(got, queue)