amqp

    AMQP 0-9-1 编解码、会话状态机与 Node TCP/TLS 消息客户端

    Download zip
    Version
    0.25.0
    License
    MIT AND BSD-3-Clause AND BSD-2-Clause
    Last updated
    6 hours ago
    Downloads
    3

    #可执行 API 示例

    按协商帧上限拆分消息体,完整帧长度包含 8 字节开销。这些例子调用公开 API,并随 moon test 执行。

    ///|
    test "outbound body frames respect negotiated full frame size" {
    let frames = @amqp.body_frames(1, b"hello world", 12)
    assert_eq(frames.length(), 3)
    let d = @amqp.Decoder::new()
    let out = []
    for f in frames {
    let wire = f.encode()
    assert_true(wire.length() <= 12)
    for x in d.feed(wire) {
    out.push(x.payload)
    }
    }
    assert_eq(out, [b"hell", b"o wo", b"rld"])
    d.finish()
    assert_eq(@amqp.body_frames(1, b"", 12), [])
    assert_true(
    try {
    ignore(@amqp.body_frames(0, b"x", 12))
    false
    } catch {
    _ => true
    },
    )
    }

    MoonBit 核心提供方法参数、字段表、Basic 属性和 Session 会话 API;Node 宿主提供 TCP/TLS、RPC、心跳、消息确认、自动恢复、多种认证及流式消息收发。恢复不自动重发旧消息或重放旧正文源,完整生产规模验证仍未完成。各功能的使用方式、实测范围和剩余限制见 README.md 与 FEATURES.md。

    FrameError

    pub suberror FrameError {
    Invalid(String)
    } derive(
    Debug
    )

    Argument

    pub(all) enum Argument {
    Bit(Bool)
    Octet(Int)
    Short(Int)
    Long(UInt)
    LongLong(UInt64)
    ShortString(String)
    LongString(Bytes)
    Table(Array[(String, FieldValue)])
    } derive(Eq,
    Debug
    )

    Argument::equal

    fn Argument::equal(Argument, Argument) -> Bool

    Argument::not_equal

    fn Argument::not_equal(x : Argument, y : Argument) -> Bool

    Argument::to_repr

    ArgumentKind

    pub(all) enum ArgumentKind {
    BitKind
    OctetKind
    ShortKind
    LongKind
    LongLongKind
    ShortStringKind
    LongStringKind
    TableKind
    } derive(Eq,
    Debug
    )

    ArgumentKind::equal

    ArgumentKind::not_equal

    fn ArgumentKind::not_equal(x : ArgumentKind, y : ArgumentKind) -> Bool

    Assembler

    pub struct Assembler {
    // private fields
    }

    Content assembly for AMQP Basic.Publish/Return/Deliver/Get-Ok. Non-content methods and heartbeats produce no Content; callers retain those frames.

    Assembler::finish

    fn Assembler::finish(self : Assembler) -> Unit raise FrameError

    Assembler::new

    fn Assembler::new(max_body? : Int, strict_methods? : Bool) -> Assembler raise FrameError

    Assembler::push

    fn Assembler::push(self : Assembler, frame : Frame) -> Content? raise FrameError

    Errors poison the assembler. Limits: 64 in-flight channels and max_body total buffered bytes.

    Authentication

    pub struct Authentication {
    // private fields
    }

    Initial-response SASL credential. Credentials deliberately have no Debug impl.

    Authentication::amqplain

    fn Authentication::amqplain(username : String, password : String) -> Authentication raise FrameError

    Authentication::custom

    fn Authentication::custom(mechanism : String, response : Bytes) -> Authentication raise FrameError

    An initial response, preserved byte-for-byte. No Unicode normalization.

    Authentication::deferred

    fn Authentication::deferred(mechanism : String) -> Authentication raise FrameError

    Ask the host for an initial response only if this candidate is selected.

    Authentication::external

    fn Authentication::external() -> Authentication

    RabbitMQ EXTERNAL response used by the pinned amqp091-go implementation.

    Authentication::mechanism

    fn Authentication::mechanism(self : Authentication) -> String

    Authentication::plain

    fn Authentication::plain(username : String, password : String) -> Authentication raise FrameError

    BasicHeader

    pub(all) struct BasicHeader {
    body_size : UInt64
    properties : Array[(String, Argument)]
    } derive(Eq,
    Debug
    )

    Property names use XML spelling; absent properties differ from present empty/zero values.

    BasicHeader::decode

    fn BasicHeader::decode(frame : Frame) -> BasicHeader raise FrameError

    BasicHeader::encode

    fn BasicHeader::encode(self : BasicHeader, channel : Int, max_size? : Int) -> Frame raise FrameError

    BasicHeader::equal

    fn BasicHeader::equal(BasicHeader, BasicHeader) -> Bool

    BasicHeader::not_equal

    fn BasicHeader::not_equal(x : BasicHeader, y : BasicHeader) -> Bool

    ConfirmLedger

    pub struct ConfirmLedger {
    // private fields
    }

    Pure bookkeeping after confirm.select activation. The transport owns mode negotiation, queueing and the exact point a validated publish enters output. Registered but unconfirmed publishes become Unknown on close, never Nacked.

    ConfirmLedger::buffered_notices

    fn ConfirmLedger::buffered_notices(self : ConfirmLedger) -> Int

    ConfirmLedger::close

    Idempotent epoch termination. Previously confirmed/nacked handles stay terminal; only pending publishes become Unknown. Do not replay them here. No ordered ack/nack notices are fabricated to fill a disconnect gap.

    ConfirmLedger::confirm

    fn ConfirmLedger::confirm(self : ConfirmLedger, epoch : String, tag : UInt64, multiple~ : Bool, ack~ : Bool) -> ConfirmUpdate raise FrameError

    Range confirmations settle ONLY unresolved slots. A later cumulative nack must not overwrite an earlier individual ack waiting for ordered delivery. Zero-tag multiple is an explicit compatibility policy, defaulting to reject; it is not inferred from the consumer-ack direction of the protocol.

    ConfirmLedger::issue

    Register once, only after local validation and immediately before output. Capacity includes settled entries blocked behind earlier missing confirms.

    ConfirmLedger::new

    fn ConfirmLedger::new(epoch : String, capacity? : Int, zero_multiple_all? : Bool) -> ConfirmLedger raise FrameError

    ConfirmLedger::next_sequence

    fn ConfirmLedger::next_sequence(self : ConfirmLedger) -> UInt64?

    ConfirmLedger::pending_count

    fn ConfirmLedger::pending_count(self : ConfirmLedger) -> Int

    ConfirmUpdate

    pub struct ConfirmUpdate {
    settled : Array[PublishDecision]
    ordered : Array[PublishDecision]
    } derive(Eq,
    Debug
    )

    ConfirmUpdate::equal

    ConfirmUpdate::not_equal

    fn ConfirmUpdate::not_equal(x : ConfirmUpdate, y : ConfirmUpdate) -> Bool

    Content

    pub struct Content {
    channel : Int
    method_payload : Bytes
    header : Bytes
    body : Bytes
    } derive(Eq,
    Debug
    )

    Content::basic_header

    fn Content::basic_header(self : Content) -> BasicHeader raise FrameError

    Content::equal

    fn Content::equal(Content, Content) -> Bool

    Content::not_equal

    fn Content::not_equal(x : Content, y : Content) -> Bool

    Content::to_repr

    Decoder

    pub struct Decoder {
    // private fields
    } derive(
    Debug
    )

    Decoder::feed

    fn Decoder::feed(self : Decoder, input : Bytes) -> Array[Frame] raise FrameError

    A malformed stream poisons the decoder. Frames from the failing feed call are discarded.

    Decoder::finish

    fn Decoder::finish(self : Decoder) -> Unit raise FrameError

    Decoder::new

    fn Decoder::new(max_size? : Int) -> Decoder raise FrameError

    Decoder::to_repr

    FieldValue

    pub(all) enum FieldValue {
    Boolean(Bool)
    Signed8(Int)
    Unsigned8(Int)
    Signed16(Int)
    Unsigned16(Int)
    Signed32(Int)
    Unsigned32(UInt)
    Signed64(Int64)
    Float32Bits(UInt)
    Float64Bits(UInt64)
    Decimal(Int, Int)
    LongString(Bytes)
    ByteArray(Bytes)
    Timestamp(UInt64)
    ArrayValue(Array[FieldValue])
    TableValue(Array[(String, FieldValue)])
    Void
    } derive(Eq,
    Debug
    )

    RabbitMQ field-table dialect. Floating point values preserve their raw IEEE bits. LongString retains arbitrary bytes; table order and duplicate keys are preserved.

    FieldValue::equal

    fn FieldValue::equal(FieldValue, FieldValue) -> Bool

    FieldValue::not_equal

    fn FieldValue::not_equal(x : FieldValue, y : FieldValue) -> Bool

    Frame

    pub(all) struct Frame {
    kind : Int
    channel : Int
    payload : Bytes
    } derive(Eq,
    Debug
    )

    Frame::encode

    fn Frame::encode(self : Frame, max_size? : Int) -> Bytes raise FrameError

    Frame::equal

    fn Frame::equal(Frame, Frame) -> Bool

    Frame::method_ids

    fn Frame::method_ids(self : Frame) -> (Int, Int)?

    Frame::not_equal

    fn Frame::not_equal(x : Frame, y : Frame) -> Bool

    Frame::to_repr

    Method

    pub(all) struct Method {
    class_id : Int
    method_id : Int
    arguments : Array[Argument]
    } derive(Eq,
    Debug
    )

    Method::decode

    fn Method::decode(frame : Frame) -> Method raise FrameError

    Decode exactly one known method, rejecting truncated and extra arguments. Unused packed bit positions are ignored on input and zeroed on encoding.

    Method::encode

    fn Method::encode(self : Method, channel : Int, max_size? : Int) -> Frame raise FrameError

    Method::equal

    fn Method::equal(Method, Method) -> Bool

    Method::new

    fn Method::new(name : String, arguments : Array[Argument]) -> Method raise FrameError

    Names use the specification spelling, for example "queue.declare" or "basic.get-ok".

    Method::not_equal

    fn Method::not_equal(x : Method, y : Method) -> Bool

    Method::to_repr

    MethodSpec

    pub struct MethodSpec {
    class_id : Int
    method_id : Int
    name : String
    fields : Array[(String, ArgumentKind)]
    carries_content : Bool
    } derive(Eq,
    Debug
    )

    MethodSpec::equal

    fn MethodSpec::equal(MethodSpec, MethodSpec) -> Bool

    MethodSpec::not_equal

    fn MethodSpec::not_equal(x : MethodSpec, y : MethodSpec) -> Bool

    PublishDecision

    pub struct PublishDecision {
    token : PublishToken
    outcome : PublishOutcome
    } derive(Eq,
    Debug
    )

    PublishDecision::equal

    PublishDecision::not_equal

    fn PublishDecision::not_equal(x : PublishDecision, y : PublishDecision) -> Bool

    PublishOutcome

    pub(all) enum PublishOutcome {
    Confirmed
    Nacked
    Unknown
    } derive(Eq,
    Debug
    )

    PublishOutcome::equal

    PublishOutcome::not_equal

    fn PublishOutcome::not_equal(x : PublishOutcome, y : PublishOutcome) -> Bool

    PublishToken

    pub struct PublishToken {
    epoch : String
    sequence : UInt64
    } derive(Eq,
    Debug
    )

    Publisher confirms are channel-local, not consumer acknowledgements. An epoch must identify one physical channel lifetime, never just its reusable id.

    PublishToken::equal

    PublishToken::not_equal

    fn PublishToken::not_equal(x : PublishToken, y : PublishToken) -> Bool

    Session

    pub struct Session {
    // private fields
    }

    Transport-independent client. The host owns timers, sockets and RPC scheduling.

    Session::authentication_mechanism

    fn Session::authentication_mechanism(self : Session) -> String

    Session::client_properties

    fn Session::client_properties(self : Session) -> Array[(String, FieldValue)] raise FrameError

    Independent deep copy; changing it cannot alter this connection or recovery.

    Session::feed

    fn Session::feed(self : Session, input : Bytes) -> Array[SessionEvent] raise FrameError

    Session::finish

    fn Session::finish(self : Session) -> Unit raise FrameError

    Session::heartbeat

    fn Session::heartbeat(self : Session) -> Bytes raise FrameError

    Session::limits

    fn Session::limits(self : Session) -> (Int, Int, Int)

    Session::locale

    fn Session::locale(self : Session) -> String

    Session::new

    fn Session::new(username : String, password : String, vhost? : String, channel_max? : Int, frame_max? : Int, heartbeat? : Int, stream_bodies? : Bool, properties? : Array[(String, FieldValue)]) -> Session raise FrameError

    Session::publish

    fn Session::publish(self : Session, channel : Int, exchange : String, routing_key : String, body : Bytes, properties? : Array[(String, Argument)], mandatory? : Bool, immediate? : Bool) -> Unit raise FrameError

    Session::publish_body

    fn Session::publish_body(self : Session, channel : Int, data : Bytes) -> Unit raise FrameError

    One bounded body frame. Invalid fragments do not consume the remaining size.

    Session::publish_start

    fn Session::publish_start(self : Session, channel : Int, exchange : String, routing_key : String, body_size : UInt64, properties? : Array[(String, Argument)], mandatory? : Bool, immediate? : Bool) -> Unit raise FrameError

    Start a publication without retaining its body. All metadata is validated before output is exposed. Hosts must serialize content on each channel.

    Session::respond_authentication

    fn Session::respond_authentication(self : Session, response : Bytes) -> Unit raise FrameError

    Supply a deferred initial response. This is not a connection.secure challenge.

    Session::send

    fn Session::send(self : Session, channel : Int, command : Method) -> Unit raise FrameError

    Session::server_locales

    fn Session::server_locales(self : Session) -> Array[String]

    Session::server_properties

    fn Session::server_properties(self : Session) -> Array[(String, FieldValue)] raise FrameError

    Independent deep copy, available after connection.start; empty before start.

    Session::server_version

    fn Session::server_version(self : Session) -> (Int, Int)?

    Session::status

    fn Session::status(self : Session) -> String

    Session::take_output

    fn Session::take_output(self : Session) -> Array[Bytes]

    Session::virtual_host

    fn Session::virtual_host(self : Session) -> String

    Session::with_authentication

    fn Session::with_authentication(authentication : Array[Authentication], vhost? : String, locale? : String, channel_max? : Int, frame_max? : Int, heartbeat? : Int, stream_bodies? : Bool, properties? : Array[(String, FieldValue)]) -> Session raise FrameError

    Candidates are considered in client order. No shared candidate fails closed.

    SessionEvent

    pub(all) enum SessionEvent {
    Ready
    AuthenticationRequested(Int, String)
    Received(Int, Method)
    Message(Content)
    MessageStart(Int, Method, BasicHeader)
    MessageData(Int, Bytes)
    MessageEnd(Int)
    ChannelClosed(Int, Int, String)
    Closed(Int, String)
    } derive(Eq,
    Debug
    )

    SessionEvent::equal

    SessionEvent::not_equal

    fn SessionEvent::not_equal(x : SessionEvent, y : SessionEvent) -> Bool

    URI

    pub struct URI {
    scheme : String
    host : String
    port : Int
    username : String
    password : String
    vhost : String
    cert_file : String
    ca_cert_file : String
    key_file : String
    server_name : String
    auth_mechanism : Array[String]
    heartbeat_seconds : Int64?
    connection_timeout : Int64
    channel_max : Int
    }

    Parsed AMQP address. No Debug/Show implementation: this contains credentials. Parsing preserves values; the transport separately validates usable ranges.

    URI::amqplain_auth

    fn URI::amqplain_auth(self : URI) -> Authentication raise FrameError

    URI::plain_auth

    fn URI::plain_auth(self : URI) -> Authentication raise FrameError

    URI::redacted

    fn URI::redacted(self : URI) -> String

    URI::to_string

    fn URI::to_string(self : URI) -> String

    Canonical Go-compatible URI string. Includes the password; use redacted for logs. Like upstream URI.String, serializes TLS parameters only, not tuning/SASL queries.

    basic_property_spec

    fn basic_property_spec() -> Array[(String, ArgumentKind)]

    body_frames

    fn body_frames(channel : Int, body : Bytes, max_frame_size : Int) -> Array[Frame] raise FrameError

    Split a body into frame payloads. max_frame_size includes the 8 byte envelope.

    content_frames

    fn content_frames(command : Method, channel : Int, properties : Array[(String, Argument)], body : Bytes, max_frame_size? : Int) -> Array[Frame] raise FrameError

    Build method/header/body frames together, with the negotiated frame envelope limit. Bodies use the existing bounded body_frames helper (up to 1 MiB).

    decode_table

    fn decode_table(bytes : Bytes) -> Array[(String, FieldValue)] raise FrameError

    encode_table

    fn encode_table(entries : Array[(String, FieldValue)]) -> Bytes raise FrameError

    Encode a length-prefixed AMQP table; maximum depth 32, 65536 visited nodes.

    method_names

    fn method_names() -> Array[String]

    method_spec

    fn method_spec(class_id : Int, method_id : Int) -> MethodSpec?

    method_spec_by_name

    fn method_spec_by_name(name : String) -> MethodSpec?

    new_connection_properties

    fn new_connection_properties() -> Array[(String, FieldValue)]

    Fresh client identity table. Values identify this implementation, not Go.

    normalize_connection_properties

    fn normalize_connection_properties(properties : Array[(String, FieldValue)]) -> Bytes raise FrameError

    Encode once to validate and snapshot every nested mutable container.

    parse_uri

    fn parse_uri(input : String) -> URI raise FrameError

    Pinned amqp091-go URI defaults and escaping, with strict UTF-8 and a 64 KiB limit. Invalid query pairs are ignored, as in net/url.Values; unknown parameters are ignored. Diagnostics deliberately never include the input URI or its credentials.

    protocol_header

    fn protocol_header() -> Bytes