AMQP 0-9-1 编解码、会话状态机与 Node TCP/TLS 消息客户端
///|
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
},
)
}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)pub struct Assembler {
// private fields
}pub struct Authentication {
// private fields
}fn Authentication::amqplain(username : String, password : String) -> Authentication raise FrameErrorfn BasicHeader::encode(self : BasicHeader, channel : Int, max_size? : Int) -> Frame raise FrameErrorpub struct ConfirmLedger {
// private fields
}fn ConfirmLedger::confirm(self : ConfirmLedger, epoch : String, tag : UInt64, multiple~ : Bool, ack~ : Bool) -> ConfirmUpdate raise FrameErrorfn ConfirmLedger::new(epoch : String, capacity? : Int, zero_multiple_all? : Bool) -> ConfirmLedger raise FrameErrorpub struct ConfirmUpdate {
settled : Array[PublishDecision]
ordered : Array[PublishDecision]
} derive(Eq, Debug)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)pub struct MethodSpec {
class_id : Int
method_id : Int
name : String
fields : Array[(String, ArgumentKind)]
carries_content : Bool
} derive(Eq, Debug)pub struct Session {
// private fields
}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 FrameErrorfn Session::publish(self : Session, channel : Int, exchange : String, routing_key : String, body : Bytes, properties? : Array[(String, Argument)], mandatory? : Bool, immediate? : Bool) -> Unit raise FrameErrorfn 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 FrameErrorfn 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 FrameErrorpub 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
}fn normalize_connection_properties(properties : Array[(String, FieldValue)]) -> Bytes raise FrameErrorInstall
Download zipAMQP 0-9-1 编解码、会话状态机与 Node TCP/TLS 消息客户端