///|
struct Client {
config : Config
events : @aqueue.Queue[Event]
changed : @cond_var.Cond
subscriptions : Map[String, QoS]
mut session : Session?
mut ready : Bool
mut stopped : Bool
mut generation : Int
mut worker : @async.Task[Unit]?
}
///|
/// Own all client tasks within the callback's lifetime.
pub async fn[T] with_client(config : Config, action : async (Client) -> T) -> T {
config.validate()
ignore(encode(config.connect_packet(), config.max_packet_size))
let client : Client = {
config,
events: @aqueue.Queue(kind=@aqueue.Blocking(config.receive_capacity)),
changed: @cond_var.Cond(),
subscriptions: {},
session: None,
ready: false,
stopped: false,
generation: 0,
worker: None,
}
@async.with_task_group(group => {
let worker = group.spawn(() => client.supervise())
client.worker = Some(worker)
defer {
client.stopped = true
client.ready = false
if client.session is Some(session) {
session.abort("client scope ended")
}
worker.cancel()
client.events.close(error=Closed)
client.changed.broadcast()
}
client.wait_connected()
let result = action(client)
client.disconnect()
result
})
}
///|
pub async fn Client::wait_connected(self : Client) -> Unit {
while !self.ready || !(self.session is Some(session) && session.alive) {
if self.stopped {
raise Closed
}
self.changed.wait()
}
}
///|
pub async fn Client::next_event(self : Client) -> Event {
self.events.get()
}
///|
fn Client::emit(self : Client, event : Event) -> Unit raise {
if !self.events.try_put(event) {
raise Backpressure(
"event queue full; client terminates instead of dropping messages",
)
}
}
///|
fn Client::active(self : Client) -> Session raise ClientError {
if self.stopped {
raise Closed
}
match self.session {
Some(session) if self.ready && session.alive => session
_ => raise NotConnected
}
}
///|
async fn Session::request(
self : Session,
config : Config,
kind : RequestKind,
packet : (Int) -> @codec.Packet,
) -> Reply {
if !self.alive {
raise NotConnected
}
// DISCONNECT has no wire identifier. Slot zero reserves one bounded close
// completion even when every ordinary inflight slot is occupied.
let id = if kind is Disconnect {
0
} else {
self.allocate(config.max_inflight)
}
let bytes = encode(packet(id), config.max_packet_size)
let request : Pending = {
id,
kind,
changed: @cond_var.Cond(),
started: false,
result: None,
}
self.pending[id] = request
defer (if request.result is None { self.abort("request cancelled") })
if !(kind is Disconnect) &&
!self.outgoing.try_put({ bytes, pending: Some(request), }) {
self.pending.remove(id)
request.finish(Err(Backpressure("send queue full")))
raise Backpressure("send queue full")
}
@async.with_timeout(config.ack_timeout_ms, () => {
if kind is Disconnect {
self.outgoing.put({ bytes, pending: Some(request), })
}
while request.result is None {
request.changed.wait()
}
}) catch {
error => {
self.abort(error.to_string())
if @async.is_cancellation_error(error) {
raise error
}
}
}
match request.result {
Some(Ok(reply)) => reply
Some(Err(error)) => raise error
None => raise Closed
}
}
///|
pub async fn Client::publish(
self : Client,
topic : String,
payload : Bytes,
qos? : QoS = AtMostOnce,
retain? : Bool = false,
) -> Unit {
let session = self.active()
if payload.length() > self.config.max_packet_size - 4 {
raise ProtocolError("outgoing payload exceeds packet limit")
}
ignore(
session.request(
self.config,
if qos == AtMostOnce {
Publish0
} else {
Publish1
},
id => {
@codec.PublishPacket({
topic,
payload,
qos: qos.wire(),
message_id: if qos == AtMostOnce {
None
} else {
Some(id)
},
retain,
dup: false,
})
},
),
)
}
///|
pub async fn Client::subscribe(
self : Client,
topics : Array[Subscription],
) -> Array[SubscriptionResult] {
self.subscribe_on(self.active(), topics.copy())
}
///|
async fn Client::subscribe_on(
self : Client,
session : Session,
topics : Array[Subscription],
) -> Array[SubscriptionResult] {
match
session.request(self.config, Subscribe(topics), id => {
@codec.SubscribePacket({
message_id: id,
topics: topics.map(t => { topic: t.topic, qos: t.qos.wire(), }),
})
}) {
Subscriptions(results) => results
_ => raise ProtocolError("wrong subscribe completion")
}
}
///|
pub async fn Client::unsubscribe(self : Client, topics : Array[String]) -> Unit {
let session = self.active()
let topics = topics.copy()
ignore(
session.request(self.config, Unsubscribe(topics), id => {
@codec.UnsubscribePacket({ message_id: id, topics, })
}),
)
}
///|
/// Stop reconnect, send DISCONNECT when possible, and await worker cleanup.
pub async fn Client::disconnect(self : Client) -> Unit {
if self.stopped {
return
}
self.stopped = true
self.ready = false
defer {
if self.session is Some(session) {
session.abort("client disconnected")
}
if self.worker is Some(worker) {
worker.cancel()
}
self.events.close(error=Closed)
self.changed.broadcast()
}
if self.session is Some(session) && session.alive {
ignore(
session.request(self.config, Disconnect, _ => @codec.DisconnectPacket),
)
}
if self.session is Some(session) {
session.abort("client disconnected")
}
if self.worker is Some(worker) {
worker.cancel()
// Explicit shutdown already settles requests and closes the transport;
// the worker may finish with cancellation or a concurrent transport error.
worker.wait() catch {
_ => ()
}
}
}