Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 27 additions & 1 deletion aimdb-core/src/outbound/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,8 +72,12 @@ impl OutboundPayload<'_> {
/// Values taken from one route's buffer, by outcome.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct RouteStats {
/// Staged for the connector.
/// Staged and handed to the connector, including any it then rejected.
pub sent: u64,
/// Handed to the connector, which could not send them (for example,
/// larger than its transport accepts). Reported with
/// [`OutboundRoutes::reject`].
pub rejected: u64,
/// Missed because the reader fell behind.
pub lagged: u64,
/// Skipped: the written topic did not fit.
Expand Down Expand Up @@ -358,6 +362,15 @@ impl OutboundRoutes {
})
}

/// Count a message from route `id` that the connector took but could not
/// send. The connector logs why; this keeps the count beside the route's
/// other outcomes.
pub fn reject(&mut self, id: RouteId) {
if let Some(stats) = self.stats.get_mut(id) {
stats.rejected += 1;
}
}

/// [`poll_stage`](Self::poll_stage), then [`take_staged`](Self::take_staged),
/// for hand-written `poll` code.
pub fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll<Option<OutboundMessage<'_>>> {
Expand Down Expand Up @@ -394,4 +407,17 @@ mod tests {
let moved = payload.into_vec();
assert_eq!(moved.as_ptr(), ptr, "moved, not copied");
}

#[tokio::test]
async fn reject_counts_beside_the_routes_other_outcomes() {
let (db, _runner) = crate::AimDbBuilder::new()
.runtime(Arc::new(crate::executor::test_support::NoopRuntimeOps))
.build()
.await
.expect("empty database");
let mut routes = OutboundRoutes::new(&db, "mqtt").unwrap();
// No routes: an unknown id is ignored rather than a panic.
routes.reject(0);
assert_eq!(routes.stats(0), None);
}
}
41 changes: 40 additions & 1 deletion aimdb-mqtt-connector/src/connector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ pub struct Native;
#[cfg(feature = "embedded")]
pub struct Embedded<D> {
pub(crate) dialer: D,
pub(crate) write_buffer: usize,
}

/// The `mountain-mqtt` backend over `embedded-tls`, on the same
Expand All @@ -55,6 +56,7 @@ pub struct Embedded<D> {
pub struct EmbeddedTls<D> {
pub(crate) dialer: D,
pub(crate) options: crate::embedded::TlsSlot,
pub(crate) write_buffer: usize,
}

/// An MQTT connector over the backend `B`.
Expand Down Expand Up @@ -90,7 +92,10 @@ impl MqttConnector<Native> {
client_id: self.client_id,
credentials: self.credentials,
keep_alive: self.keep_alive,
backend: Embedded { dialer },
backend: Embedded {
dialer,
write_buffer: crate::embedded::DEFAULT_WRITE_BUFFER,
},
}
}

Expand All @@ -110,6 +115,7 @@ impl MqttConnector<Native> {
backend: EmbeddedTls {
dialer,
options: crate::embedded::TlsSlot::new(options),
write_buffer: crate::embedded::DEFAULT_WRITE_BUFFER,
},
}
}
Expand Down Expand Up @@ -143,6 +149,38 @@ impl<B> MqttConnector<B> {
}
}

/// The write buffer's documentation, shared by both embedded backends.
#[cfg(feature = "embedded")]
macro_rules! write_buffer_doc {
() => {
"Size the session's write ring, in bytes (default 4,096). Allocated once \
and reused across reconnects.\n\n\
An outbound PUBLISH frame plus a 64-byte reserve must fit in half the \
ring (1,984 bytes of frame at the default). `build()` fails for a \
route whose largest frame does not fit, and for a CONNECT or \
SUBSCRIBE that does not; an owned payload over the limit at runtime \
is skipped and counted as rejected in the route's `RouteStats`."
};
}

#[cfg(feature = "embedded")]
impl<D> MqttConnector<Embedded<D>> {
#[doc = write_buffer_doc!()]
pub fn with_write_buffer(mut self, bytes: usize) -> Self {
self.backend.write_buffer = bytes;
self
}
}

#[cfg(feature = "embedded-tls")]
impl<D> MqttConnector<EmbeddedTls<D>> {
#[doc = write_buffer_doc!()]
pub fn with_write_buffer(mut self, bytes: usize) -> Self {
self.backend.write_buffer = bytes;
self
}
}

/// Whole seconds for the wire, or the reason this keep-alive cannot be used.
fn keep_alive_secs(keep_alive: Duration) -> DbResult<u16> {
let secs = keep_alive.as_secs();
Expand Down Expand Up @@ -228,6 +266,7 @@ where
credentials,
keep_alive_secs,
&self.dialer,
self.write_buffer,
)
}
}
Expand Down
10 changes: 2 additions & 8 deletions aimdb-mqtt-connector/src/embedded/manager.rs
Original file line number Diff line number Diff line change
@@ -1,19 +1,13 @@
//! Session cadence and the channel a session takes actions from.
//! Session cadence and the reasons a session ends.
//!
//! The channel uses `CriticalSectionRawMutex`, so it is `Sync` and the sink
//! needs no force-`Send` wrapper. Time comes from core's
//! Time comes from core's
//! [`aimdb_core::session::Delay`], so nothing here names an executor.

use core::time::Duration;

use aimdb_core::RuntimeOps;
use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex;
use embassy_sync::channel::Channel;
use mountain_mqtt::client::ClientError;

/// The action channel: `pump_sink` to broker session.
pub(crate) type ActionChannel<A, const Q: usize> = Channel<CriticalSectionRawMutex, A, Q>;

/// Monotonic milliseconds. Only differences are meaningful.
pub(crate) fn now_ms(runtime: &dyn RuntimeOps) -> u64 {
runtime.now_nanos() / 1_000_000
Expand Down
Loading