///|
async fn Client::dial(self : Client) -> Session {
@async.with_timeout(self.config.connect_timeout_ms, () => {
let tcp = @socket.Tcp::connect_to_host(
self.config.host[:],
port=self.config.port,
)
errdefer tcp.close()
let (reader, writer, close_transport) : (
&@io.Reader,
&@io.Writer,
() -> Unit,
) = match self.config.tls {
Plain => (tcp, tcp, () => tcp.close())
mode => {
let trust = match mode {
SystemRoots => @tls.SystemRoot
CustomCA(path) => @tls.CustomPemFile(path)
Plain => @tls.SystemRoot
}
let tls = @tls.Tls::client(tcp, host=self.config.host, trust~)
(
tls,
tls,
() => {
tls.close()
tcp.close()
},
)
}
}
errdefer close_transport()
writer.write(
encode(self.config.connect_packet(), self.config.max_packet_size),
)
match read_packet(reader, self.config.max_packet_size) {
@codec.ConnackPacket(ack) => {
if ack.return_code != @codec.ConnectionAccepted {
raise ConnectionRefused(@debug.to_string(ack.return_code))
}
if ack.session_present {
raise ProtocolError("CleanSession=true CONNACK set Session Present")
}
}
_ => raise ProtocolError("expected CONNACK")
}
self.generation += 1
{
generation: self.generation,
reader,
writer,
close_transport,
outgoing: @aqueue.Queue(kind=@aqueue.Blocking(self.config.send_capacity)),
pending: {},
next_id: 1,
alive: true,
last_write: @async.now(),
ping_sent: None,
early_messages: [],
failure_reason: None,
}
})
}
///|
async fn Client::supervise(self : Client) -> Unit {
let mut attempts = 0
let mut delay = self.config.reconnect_delay_ms
let mut ever_ready = false
defer {
self.ready = false
self.stopped = true
if self.session is Some(session) {
session.abort("worker ended")
}
self.events.close(error=Closed)
self.changed.broadcast()
}
while !self.stopped {
try {
let session = self.dial()
self.session = Some(session)
defer session.abort("connection ended")
@async.with_task_group(group => {
group.spawn_bg(() => session.write_loop(self.config.ack_timeout_ms))
group.spawn_bg(() => session.heartbeat(self.config))
group.spawn_bg(() => self.read_loop(session))
let restore = self.subscriptions
.iter()
.map(pair => { topic: pair.0, qos: pair.1, })
.to_array()
for subscription in restore {
let result = self.subscribe_on(session, [subscription])
if result.contains(Rejected) {
raise ProtocolError("subscription restore rejected")
}
}
self.ready = true
ever_ready = true
attempts = 0
delay = self.config.reconnect_delay_ms
self.emit(Connected(session.generation))
for message in session.early_messages {
self.emit(MessageReceived(message))
}
session.early_messages.clear()
self.changed.broadcast()
while session.alive && !self.stopped {
self.changed.wait()
}
}) catch {
error => {
session.abort(error.to_string())
raise error
}
}
} catch {
error => {
let reason = match self.session {
Some(session) => session.failure_reason.unwrap_or(error.to_string())
None => error.to_string()
}
self.ready = false
self.session = None
self.changed.broadcast()
if self.stopped {
return
}
match error {
Backpressure(_) => raise error
_ => ()
}
if !ever_ready {
raise error
}
self.emit(Disconnected(self.generation, reason))
attempts += 1
if attempts > self.config.reconnect_attempts {
raise ReconnectExhausted(reason)
}
let jitter = (self.generation * 17 + attempts * 31) % 97
@async.sleep(delay + jitter)
delay = (delay + delay / 2).min(self.config.max_reconnect_delay_ms)
}
}
}
}
///|
async fn Client::read_loop(self : Client, session : Session) -> Unit {
while session.alive {
match read_packet(session.reader, self.config.max_packet_size) {
@codec.PublishPacket(packet) => {
let qos = match packet.qos {
@codec.QoS0 => AtMostOnce
@codec.QoS1 => AtLeastOnce
@codec.QoS2 => raise ProtocolError("QoS 2 is not supported")
}
let message : Message = {
topic: packet.topic,
payload: packet.payload,
qos,
retain: packet.retain,
dup: packet.dup,
generation: session.generation,
}
if self.ready {
self.emit(MessageReceived(message))
} else {
if session.early_messages.length() >= self.config.receive_capacity {
raise Backpressure(
"messages arrived faster than subscription restoration",
)
}
session.early_messages.push(message)
}
if packet.message_id is Some(id) && qos == AtLeastOnce {
session.control(
@codec.PubackPacket({ message_id: id, }),
self.config.max_packet_size,
)
}
}
@codec.PubackPacket(ack) =>
match session.pending.get(ack.message_id) {
Some(request) if request.kind is Publish1 && request.started => {
session.pending.remove(ack.message_id)
request.finish(Ok(Done))
}
_ => raise ProtocolError("unexpected PUBACK identifier")
}
@codec.SubackPacket(ack) =>
match session.pending.get(ack.message_id) {
Some(request) =>
match request.kind {
Subscribe(topics) if request.started &&
topics.length() == ack.return_codes.length() => {
let results : Array[SubscriptionResult] = []
for i in 0.. Granted(AtMostOnce)
@codec.SuccessQoS1 => Granted(AtLeastOnce)
@codec.Failure => Rejected
@codec.SuccessQoS2 =>
raise ProtocolError("SUBACK granted unsupported QoS")
}
match result {
Granted(qos) =>
if qos == AtLeastOnce && topics[i].qos == AtMostOnce {
raise ProtocolError("SUBACK exceeds requested QoS")
}
Rejected => ()
}
results.push(result)
}
// Commit desired subscriptions only after the whole ACK validates.
for i in 0.. raise ProtocolError("SUBACK does not match pending request")
}
None => raise ProtocolError("unexpected SUBACK identifier")
}
@codec.UnsubackPacket(ack) =>
match session.pending.get(ack.message_id) {
Some(request) =>
match request.kind {
Unsubscribe(topics) if request.started => {
for topic in topics {
self.subscriptions.remove(topic)
}
session.pending.remove(ack.message_id)
request.finish(Ok(Done))
}
_ =>
raise ProtocolError("UNSUBACK does not match pending request")
}
None => raise ProtocolError("unexpected UNSUBACK identifier")
}
@codec.PingrespPacket => session.ping_sent = None
_ => raise ProtocolError("unexpected server packet")
}
@async.pause()
}
}