diff --git a/aimdb-mqtt-connector/src/connector.rs b/aimdb-mqtt-connector/src/connector.rs index ccf9ac59..c410350b 100644 --- a/aimdb-mqtt-connector/src/connector.rs +++ b/aimdb-mqtt-connector/src/connector.rs @@ -39,6 +39,11 @@ type BuildFuture<'a> = Pin>> + S pub struct Native; /// The `mountain-mqtt` backend over a caller-supplied transport. +/// +/// Inbound publishes are delivered into their records by the session task +/// itself. At QoS 1 the PUBACK is sent before delivery, so it means the +/// message reached AimDB, not that every record kept it: a record whose +/// buffer is full drops it, and the broker does not resend. #[cfg(feature = "embedded")] pub struct Embedded { pub(crate) dialer: D, diff --git a/aimdb-mqtt-connector/src/embedded/manager.rs b/aimdb-mqtt-connector/src/embedded/manager.rs index 206b7425..12fa5852 100644 --- a/aimdb-mqtt-connector/src/embedded/manager.rs +++ b/aimdb-mqtt-connector/src/embedded/manager.rs @@ -1,7 +1,7 @@ -//! Session cadence and the channels a session talks over. +//! Session cadence and the channel a session takes actions from. //! -//! Channels use `CriticalSectionRawMutex`, so they are `Sync` and the sink and -//! source need no force-`Send` wrapper. Time comes from core's +//! The channel uses `CriticalSectionRawMutex`, so it is `Sync` and the sink +//! needs no force-`Send` wrapper. Time comes from core's //! [`aimdb_core::session::Delay`], so nothing here names an executor. use core::time::Duration; @@ -9,11 +9,7 @@ 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, EventHandlerError}; -use mountain_mqtt::packets::publish::ApplicationMessage; - -/// The event channel: broker session to `pump_source`. -pub(crate) type EventChannel = Channel; +use mountain_mqtt::client::ClientError; /// The action channel: `pump_sink` to broker session. pub(crate) type ActionChannel = Channel; @@ -23,13 +19,6 @@ pub(crate) fn now_ms(runtime: &dyn RuntimeOps) -> u64 { runtime.now_nanos() / 1_000_000 } -/// Convert a received [`ApplicationMessage`] into an application event. -pub trait FromApplicationMessage: Sized { - /// Build the event, or reject the message. - fn from_application_message(message: &ApplicationMessage

) - -> Result; -} - /// Why a session ended. #[derive(Debug, PartialEq, Clone, Copy)] pub enum Error { diff --git a/aimdb-mqtt-connector/src/embedded/mod.rs b/aimdb-mqtt-connector/src/embedded/mod.rs index 7fe77910..5933aa9e 100644 --- a/aimdb-mqtt-connector/src/embedded/mod.rs +++ b/aimdb-mqtt-connector/src/embedded/mod.rs @@ -1,8 +1,9 @@ //! The `mountain-mqtt` backend: broker session plus the data-plane bridges. //! -//! Outbound publishes and inbound routing ride core's [`pump_sink`] / -//! [`pump_source`] directly — the session channels are `Sync`, so nothing -//! force-`Send` stands between them and the runner. +//! Outbound publishes ride core's [`pump_sink`] — the action channel is +//! `Sync`, so nothing force-`Send` stands between it and the runner. Inbound +//! publishes are dispatched into their records by the session loop itself, +//! through an [`InboundDispatch`](aimdb_core::InboundDispatch). //! //! See the crate docs for a usage example. @@ -24,7 +25,7 @@ pub mod tls; extern crate alloc; use aimdb_core::connector::ConnectorUrl; -use aimdb_core::session::{pump_sink, pump_source, Payload}; +use aimdb_core::session::pump_sink; use aimdb_core::transport::{ConnectorConfig, PublishError}; use alloc::boxed::Box; use alloc::format; @@ -47,7 +48,7 @@ pub use crate::embedded::tls::TlsOptions; #[cfg(feature = "embedded-tls")] use crate::embedded::tls::{host_ip_literal, READ_BUF_MIN, WRITE_BUF_MIN}; -/// Maximum number of pending MQTT actions and events +/// Maximum number of pending MQTT actions pub(crate) const CHANNEL_SIZE: usize = 32; /// Buffer size for MQTT packets (4KB) @@ -61,15 +62,13 @@ pub(crate) const MAX_PROPERTIES: usize = 32; /// The runner's collected future type. type EmbassyBoxFuture = Pin + Send + 'static>>; -/// What a transport's setup hands back: the two channel ends the pumps ride, -/// plus the tasks that serve them. -type ManagerSetup = (Arc, Arc, Vec); +/// What a transport's setup hands back: the action channel `pump_sink` rides, +/// plus the tasks that serve it. +type ManagerSetup = (Arc, Vec); /// Outbound publishes and subscriptions: pumps to broker session. pub(crate) type ActionChannel = crate::embedded::manager::ActionChannel; -/// Inbound messages: broker session to pumps. -pub(crate) type EventChannel = crate::embedded::manager::EventChannel; /// What the pumps ask the session to put on the wire. /// @@ -91,42 +90,10 @@ pub enum AimdbMqttAction { }, } -/// What the session hands back for `pump_source` to route. -#[derive(Clone)] -pub enum AimdbMqttEvent { - /// A message was received from a subscribed topic - MessageReceived { - /// The topic the message was received on - topic: String, - /// The message payload, built once from the wire bytes. - payload: Payload, - }, -} - -impl crate::embedded::manager::FromApplicationMessage for AimdbMqttEvent { - fn from_application_message( - message: &mountain_mqtt::packets::publish::ApplicationMessage, - ) -> Result { - #[cfg(feature = "defmt")] - defmt::debug!( - "Received message on topic '{}', {} bytes", - message.topic_name, - message.payload.len() - ); - - Ok(Self::MessageReceived { - topic: message.topic_name.to_string(), - // Straight to `Payload` — one allocation and one copy, where a - // `Vec` here would be converted again on the way out. - payload: Payload::from(message.payload), - }) - } -} - // =========================================================================== -// Data-plane bridges — core's pumps drive these directly. The channels are -// `Sync` (their mutex is `CriticalSectionRawMutex`), so no force-`Send` -// wrapper stands between them and the runner. +// Data-plane bridge — `pump_sink` drives it directly. The action channel is +// `Sync` (its mutex is `CriticalSectionRawMutex`), so no force-`Send` wrapper +// stands between it and the runner. // =========================================================================== /// Turns a `pump_sink` publish into an `AimdbMqttAction::Publish` on the @@ -167,20 +134,6 @@ impl aimdb_core::transport::Connector for MqttSink { } } -/// Drains the session's event channel as `(topic, payload)` for `pump_source`. -struct MqttSource { - events: Arc, -} - -impl aimdb_core::session::Source for MqttSource { - fn next(&mut self) -> aimdb_core::BoxFut<'_, Option<(String, Payload)>> { - Box::pin(async move { - let AimdbMqttEvent::MessageReceived { topic, payload } = self.events.receive().await; - Some((topic, payload)) - }) - } -} - /// Force-`Send + Sync` slot for the TLS materials: [`TlsOptions`] holds /// `&'static mut` exclusive resources, so it is neither `Sync` nor takeable /// through the `&self` that [`ConnectorBuilder::build`] receives. @@ -205,8 +158,8 @@ where + 'static, { Box::pin(async move { - let router = db.inbound_router("mqtt", &crate::MqttGrammar)?; - let topics = inbound_topics(&router); + let inbound = aimdb_core::InboundDispatch::new(db, "mqtt", &crate::MqttGrammar)?; + let topics = inbound_topics(&inbound); warn_unsupported_qos(db); let broker = parse_broker_url(broker_url)?; if broker.tls { @@ -215,15 +168,16 @@ where let connection_settings = static_connection_settings(client_id, credentials, broker.credentials.as_ref()); - let (actions, events, manager_tasks) = setup_manager( + let (actions, manager_tasks) = setup_manager( &broker, connection_settings, dialer.clone(), topics, + inbound, Settings::from_keep_alive_secs(keep_alive_secs), db.runtime_ops(), )?; - Ok(collect_pumps(db, router, actions, events, manager_tasks)) + Ok(collect_pumps(db, actions, manager_tasks)) }) } @@ -246,8 +200,8 @@ where + 'static, { Box::pin(async move { - let router = db.inbound_router("mqtt", &crate::MqttGrammar)?; - let topics = inbound_topics(&router); + let inbound = aimdb_core::InboundDispatch::new(db, "mqtt", &crate::MqttGrammar)?; + let topics = inbound_topics(&inbound); warn_unsupported_qos(db); let broker = parse_broker_url(broker_url)?; if !broker.tls { @@ -260,22 +214,23 @@ where let connection_settings = static_connection_settings(client_id, credentials, broker.credentials.as_ref()); - let (actions, events, manager_tasks) = setup_tls_manager( + let (actions, manager_tasks) = setup_tls_manager( &broker, options, connection_settings, backend.dialer.clone(), topics, + inbound, Settings::from_keep_alive_secs(keep_alive_secs), db.runtime_ops(), )?; - Ok(collect_pumps(db, router, actions, events, manager_tasks)) + Ok(collect_pumps(db, actions, manager_tasks)) }) } /// The inbound topics the session must subscribe on every connection. -fn inbound_topics(router: &aimdb_core::Router) -> Vec { - let topics: Vec = router +fn inbound_topics(inbound: &aimdb_core::InboundDispatch) -> Vec { + let topics: Vec = inbound .subscriptions() .iter() .map(|t| t.to_string()) @@ -287,17 +242,14 @@ fn inbound_topics(router: &aimdb_core::Router) -> Vec { topics } -/// Outbound publishes and inbound routing ride core's pumps; the session tasks -/// join them. +/// Outbound publishes ride core's `pump_sink`; the session tasks, which also +/// dispatch inbound publishes, join them. fn collect_pumps( db: &aimdb_core::builder::AimDb, - router: aimdb_core::Router, actions: Arc, - events: Arc, manager_tasks: Vec, ) -> Vec { let mut futures = pump_sink(db, "mqtt", Arc::new(MqttSink { actions })); - futures.extend(pump_source(db, router, MqttSource { events })); futures.extend(manager_tasks); futures } @@ -391,6 +343,7 @@ fn setup_manager( connection_settings: ConnectionSettings<'static>, dialer: D, topics: Vec, + inbound: aimdb_core::InboundDispatch, settings: Settings, runtime: Arc, ) -> Result @@ -403,7 +356,6 @@ where + 'static, { let actions: Arc = Arc::new(ActionChannel::new()); - let events: Arc = Arc::new(EventChannel::new()); // The dialer is both the transport and the clock the session runs on. let host = broker.host.clone(); @@ -415,7 +367,6 @@ where let manager_task: EmbassyBoxFuture = Box::pin(unsafe { crate::embedded::session::SendSession::new({ let actions = actions.clone(); - let events = events.clone(); async move { #[cfg(feature = "defmt")] defmt::info!("MQTT background task starting"); @@ -427,7 +378,7 @@ where topics, connection_settings, settings, - events, + inbound, actions, runtime, ) @@ -436,18 +387,20 @@ where }) }); - Ok((actions, events, alloc::vec![manager_task])) + Ok((actions, alloc::vec![manager_task])) } /// Set up the TLS broker manager ([`run_tls`]) plus the SNTP time-source task. /// Synchronous — no `.await` — so the caller's `build` future stays `Send`. #[cfg(feature = "embedded-tls")] +#[allow(clippy::too_many_arguments)] fn setup_tls_manager( broker: &BrokerUrl, options: TlsOptions, connection_settings: ConnectionSettings<'static>, dialer: D, topics: Vec, + inbound: aimdb_core::InboundDispatch, settings: Settings, runtime: Arc, ) -> Result @@ -487,7 +440,6 @@ where } let actions: Arc = Arc::new(ActionChannel::new()); - let events: Arc = Arc::new(EventChannel::new()); let host = broker.host.clone(); let port = broker.port; @@ -502,7 +454,6 @@ where let mut tasks: Vec = alloc::vec![Box::pin(unsafe { crate::embedded::session::SendSession::new({ let actions = actions.clone(); - let events = events.clone(); async move { #[cfg(feature = "defmt")] defmt::info!("MQTT-TLS background task starting"); @@ -517,7 +468,7 @@ where topics, connection_settings, settings, - events, + inbound, actions, delay, runtime, @@ -539,7 +490,7 @@ where })); } - Ok((actions, events, tasks)) + Ok((actions, tasks)) } /// Map a QoS level to mountain-mqtt's `QualityOfService`. diff --git a/aimdb-mqtt-connector/src/embedded/packet_reader.rs b/aimdb-mqtt-connector/src/embedded/packet_reader.rs index ce376f3d..d77273f6 100644 --- a/aimdb-mqtt-connector/src/embedded/packet_reader.rs +++ b/aimdb-mqtt-connector/src/embedded/packet_reader.rs @@ -89,6 +89,12 @@ impl PacketReader { Ok(Some(total)) } + /// Whether the packet at the head is a PUBLISH above QoS 0, which the + /// client answers with a PUBACK. Reads only the first header byte. + pub(crate) fn head_needs_ack(&self) -> bool { + self.len > 0 && self.buf[0] >> 4 == 3 && (self.buf[0] >> 1) & 0b11 != 0 + } + /// Parse the complete packet at the head of the buffer. /// /// `total` must come from [`framed_len`](Self::framed_len). Takes `&self`, @@ -162,6 +168,25 @@ mod tests { got } + #[test] + fn only_a_publish_above_qos_0_needs_an_ack() { + let mut reader = PacketReader::<64>::new(); + assert!(!reader.head_needs_ack(), "empty"); + + reader.feed(&publish_bytes("t", b"x")).unwrap(); + assert!(!reader.head_needs_ack(), "QoS 0 publish"); + reader.consume(reader.framed_len().unwrap().unwrap()); + + let mut qos1 = publish_bytes("t", b"x"); + qos1[0] |= 0b0010; + reader.feed(&qos1).unwrap(); + assert!(reader.head_needs_ack(), "QoS 1 publish"); + reader.consume(qos1.len()); + + reader.feed(CONNACK).unwrap(); + assert!(!reader.head_needs_ack(), "CONNACK"); + } + #[test] fn one_byte_at_a_time_both_packets_parse() { let mut wire = Vec::new(); diff --git a/aimdb-mqtt-connector/src/embedded/session.rs b/aimdb-mqtt-connector/src/embedded/session.rs index 1d2b1c8d..12c8c0b0 100644 --- a/aimdb-mqtt-connector/src/embedded/session.rs +++ b/aimdb-mqtt-connector/src/embedded/session.rs @@ -53,7 +53,7 @@ pub(crate) async fn run_sessions( topics: alloc::vec::Vec, connection_settings: mountain_mqtt::client::ConnectionSettings<'static>, settings: crate::embedded::manager::Settings, - events: alloc::sync::Arc, + inbound: aimdb_core::InboundDispatch, actions: alloc::sync::Arc, runtime: alloc::sync::Arc, ) -> ! @@ -94,7 +94,7 @@ where tx, &connection_settings, &subscribe_topics, - &events, + &inbound, &actions, &ring, &settings, diff --git a/aimdb-mqtt-connector/src/embedded/session_loop.rs b/aimdb-mqtt-connector/src/embedded/session_loop.rs index 0c737fb6..acf4c1a4 100644 --- a/aimdb-mqtt-connector/src/embedded/session_loop.rs +++ b/aimdb-mqtt-connector/src/embedded/session_loop.rs @@ -17,7 +17,7 @@ use core::task::Poll; use core::time::Duration; use aimdb_core::session::{ByteRead, ByteWrite, Delay}; -use aimdb_core::RuntimeOps; +use aimdb_core::{InboundDispatch, RuntimeOps}; use embassy_futures::select::{select3, Either3}; use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex; use embassy_sync::channel::Channel; @@ -30,12 +30,10 @@ use mountain_mqtt::error::{PacketReadError, PacketWriteError}; use mountain_mqtt::packets::connect::Connect; use mountain_mqtt::packets::packet_generic::PacketGeneric; -use crate::embedded::manager::{now_ms, Error, FromApplicationMessage, Settings}; +use crate::embedded::manager::{now_ms, Error, Settings}; use crate::embedded::packet_reader::PacketReader; use crate::embedded::write_ring::{encoded_len, WriteRing, CONTROL_RESERVE}; -use crate::embedded::{ - ActionChannel, AimdbMqttAction, AimdbMqttEvent, EventChannel, BUFFER_SIZE, MAX_PROPERTIES, -}; +use crate::embedded::{ActionChannel, AimdbMqttAction, BUFFER_SIZE, MAX_PROPERTIES}; /// Bytes lifted off the socket at a time, and the size of one `inbound` slot. const RX_CHUNK: usize = 256; @@ -46,9 +44,8 @@ const RX_CHUNK: usize = 256; /// packets go into the connector's write ring instead. const PACKET_BUFFER_SIZE: usize = BUFFER_SIZE - 2 * RX_CHUNK; -/// Ring space waited for before each received packet is parsed: enough for -/// the PUBACK a QoS 1 publish needs (6 bytes; mountain-mqtt adds no -/// properties). +/// Ring space waited for before a QoS 1 publish is parsed: enough for its +/// PUBACK (6 bytes; mountain-mqtt adds no properties). const PUBACK_ROOM: usize = 16; /// The largest packet the reader takes whatever arrives before it: the @@ -63,8 +60,9 @@ type Chunk = heapless::Vec; /// Drive one MQTT session over a split stream until an error ends it. /// -/// Connects, subscribes `subscribe_topics`, then dispatches actions and -/// forwards events. Returns only on failure — the caller reconnects. +/// Connects, subscribes `subscribe_topics`, then performs actions and +/// dispatches every inbound publish into its records through `dispatch`. +/// Returns only on failure — the caller reconnects. /// /// `ring` is the connector's, reused across sessions; whatever an old session /// left in it is discarded first, so nothing reaches the new socket ahead of @@ -79,7 +77,7 @@ pub(crate) async fn run_session( tx: W, connection_settings: &ConnectionSettings<'static>, subscribe_topics: &[(&str, QualityOfService)], - events: &EventChannel, + dispatch: &InboundDispatch, actions: &ActionChannel, ring: &WriteRing, settings: &Settings, @@ -99,7 +97,7 @@ where ring, connection_settings, subscribe_topics, - events, + dispatch, actions, settings, delay, @@ -148,7 +146,7 @@ async fn client_loop( ring: &WriteRing, connection_settings: &ConnectionSettings<'static>, subscribe_topics: &[(&str, QualityOfService)], - events: &EventChannel, + dispatch: &InboundDispatch, actions: &ActionChannel, settings: &Settings, delay: &D, @@ -285,7 +283,7 @@ async fn client_loop( &mut reader, &mut state, ring, - events, + dispatch, runtime, &mut last_ack_ms, &mut connected, @@ -307,13 +305,14 @@ async fn client_loop( } } -/// Parse and dispatch every whole packet the reader now holds. +/// Parse every whole packet the reader now holds, dispatching publishes into +/// their records. #[allow(clippy::too_many_arguments)] async fn drain_packets( reader: &mut PacketReader, state: &mut ClientStateNoQueue, ring: &WriteRing, - events: &EventChannel, + dispatch: &InboundDispatch, runtime: &dyn RuntimeOps, last_ack_ms: &mut u64, connected: &mut bool, @@ -322,12 +321,13 @@ async fn drain_packets( // A burst of QoS 1 publishes needs a PUBACK each, so room for one is // waited for here, before parsing: the parsed packet is too large to // hold across an await in the session future. - ring.wait_room(PUBACK_ROOM).await; + if reader.head_needs_ack() { + ring.wait_room(PUBACK_ROOM).await; + } - // The packet borrows the reader's buffer, so everything that outlives - // it — the application event — is made owned inside this scope. - // `consume` can then take `&mut`. - let received = { + // The packet borrows the reader's buffer, so it is handled entirely + // inside this scope; `consume` can then take `&mut`. + { let packet: PacketGeneric<'_, MAX_PROPERTIES, 0, 0> = reader.parse(total).map_err(client_error)?; @@ -346,8 +346,8 @@ async fn drain_packets( } let event = state.receive(packet).map_err(client_error)?; - Received::of(event)? - }; + deliver(event, dispatch)?; + } reader.consume(total); // Every packet the state accepted proves the broker is alive. @@ -358,58 +358,53 @@ async fn drain_packets( if !*connected && matches!(state, ClientStateNoQueue::Connected(_)) { *connected = true; } - - if let Received::Event(event) = received { - events.send(event).await; - } } Ok(()) } -/// What a received packet leaves for the loop to do, owned so the reader's -/// buffer can be compacted first. -enum Received { - /// An acknowledgement: liveness only, nothing to forward. - Ack, - /// A message for `pump_source` to route. - Event(AimdbMqttEvent), -} - -impl Received { - fn of(event: ClientStateReceiveEvent<'_, '_, MAX_PROPERTIES>) -> Result { - Ok(match event { - ClientStateReceiveEvent::Ack => Self::Ack, - - ClientStateReceiveEvent::Publish { publish } - | ClientStateReceiveEvent::PublishAndPuback { publish, .. } => { - if publish.topic_name().is_empty() { - return Err(Error::Client( - ClientError::EmptyTopicNameWithAliasesDisabled, - )); - } - let message = publish.into(); - let event = AimdbMqttEvent::from_application_message(&message) - .map_err(|e| Error::Client(ClientError::EventHandler(e)))?; - Self::Event(event) +/// Hand a received publish to its records; everything else only proves +/// liveness. +/// +/// The PUBACK for a QoS 1 publish is already queued, so a record buffer that +/// is full drops a message the broker considers delivered. +fn deliver( + event: ClientStateReceiveEvent<'_, '_, MAX_PROPERTIES>, + dispatch: &InboundDispatch, +) -> Result<(), Error> { + match event { + ClientStateReceiveEvent::Ack => {} + + ClientStateReceiveEvent::Publish { publish } + | ClientStateReceiveEvent::PublishAndPuback { publish, .. } => { + if publish.topic_name().is_empty() { + return Err(Error::Client( + ClientError::EmptyTopicNameWithAliasesDisabled, + )); } + #[cfg(feature = "defmt")] + defmt::debug!( + "Received message on topic '{}', {} bytes", + publish.topic_name(), + publish.payload().len() + ); + dispatch.dispatch(publish.topic_name(), publish.payload()); + } - // Liveness, and nothing else. The broker is telling us a - // subscription was granted below the QoS asked for, that a publish - // matched no subscriber, or that an unsubscribe named a - // subscription it did not hold. AimDB has nowhere to deliver any of - // that: `pump_source` owns the channel an application would have - // read it from, and a record has no connection-state callback. - ClientStateReceiveEvent::SubscriptionGrantedBelowMaximumQos { .. } - | ClientStateReceiveEvent::PublishedMessageHadNoMatchingSubscribers - | ClientStateReceiveEvent::NoSubscriptionExisted => Self::Ack, - - ClientStateReceiveEvent::Disconnect { disconnect } => { - return Err(Error::Client(ClientError::Disconnected( - *disconnect.reason_code(), - ))) - } - }) + // Liveness, and nothing else. The broker is telling us a subscription + // was granted below the QoS asked for, that a publish matched no + // subscriber, or that an unsubscribe named a subscription it did not + // hold. A record has no connection-state callback to deliver it to. + ClientStateReceiveEvent::SubscriptionGrantedBelowMaximumQos { .. } + | ClientStateReceiveEvent::PublishedMessageHadNoMatchingSubscribers + | ClientStateReceiveEvent::NoSubscriptionExisted => {} + + ClientStateReceiveEvent::Disconnect { disconnect } => { + return Err(Error::Client(ClientError::Disconnected( + *disconnect.reason_code(), + ))) + } } + Ok(()) } /// Turn one queued action into a packet on the wire. @@ -559,12 +554,18 @@ mod tests { /// absorb codegen drift, but not loose enough to fit another buffer. #[test] fn the_session_future_has_not_outgrown_the_loop_it_replaced() { - let events = EventChannel::new(); let actions = ActionChannel::new(); let settings = Settings::default(); let connection_settings = ConnectionSettings::unauthenticated("size-probe"); let runtime = aimdb_core::executor::test_support::NoopRuntimeOps; let ring = WriteRing::new(64); + let (db, _runner) = futures::executor::block_on( + aimdb_core::AimDbBuilder::new() + .runtime(alloc::sync::Arc::new(runtime)) + .build(), + ) + .expect("empty database"); + let dispatch = InboundDispatch::new(&db, "mqtt", &crate::MqttGrammar).expect("no links"); // Built, never polled: `size_of_val` on the future is the whole point. let session = run_session( @@ -572,7 +573,7 @@ mod tests { NullWrite, &connection_settings, &[], - &events, + &dispatch, &actions, &ring, &settings, diff --git a/aimdb-mqtt-connector/src/embedded/tls.rs b/aimdb-mqtt-connector/src/embedded/tls.rs index 65571de7..c98f77f3 100644 --- a/aimdb-mqtt-connector/src/embedded/tls.rs +++ b/aimdb-mqtt-connector/src/embedded/tls.rs @@ -321,7 +321,7 @@ pub(crate) async fn run_tls( topics: Vec, connection_settings: ConnectionSettings<'static>, settings: Settings, - events: Arc, + inbound: aimdb_core::InboundDispatch, actions: Arc, delay: D, runtime: Arc, @@ -412,7 +412,7 @@ where TlsWrite(tls_tx), &connection_settings, &subscribe_topics, - &events, + &inbound, &actions, &ring, &settings, diff --git a/aimdb-mqtt-connector/tests/tokio_broker.rs b/aimdb-mqtt-connector/tests/tokio_broker.rs index 3acbfccb..32e7c303 100644 --- a/aimdb-mqtt-connector/tests/tokio_broker.rs +++ b/aimdb-mqtt-connector/tests/tokio_broker.rs @@ -349,3 +349,31 @@ async fn a_retained_message_over_the_maximum_packet_size_is_withheld() { assert_eq!(seen.withheld, 1); assert_eq!(received, None); } + +/// Inbound publishes are dispatched by the session task itself: with inbound +/// links and no outbound ones, the connector contributes one future. +#[tokio::test] +async fn the_embedded_backend_dispatches_inbound_on_its_session_task() { + use aimdb_core::buffer::BufferCfg; + use aimdb_core::connector::ConnectorBuilder; + use aimdb_core::AimDbBuilder; + use aimdb_mqtt_connector::MqttConnector; + use aimdb_tokio_adapter::net::TokioNet; + use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; + + let connector = || MqttConnector::new("mqtt://127.0.0.1:1").transport(TokioNet::tcp()); + let mut builder = AimDbBuilder::new() + .runtime(Arc::new(TokioAdapter)) + .with_connector(connector()); + builder.configure::("temperature", |reg| { + reg.buffer(BufferCfg::SingleLatest) + .link_from("mqtt://sensors/temperature") + .with_deserializer(|_ctx, data: &[u8]| Ok::(data.len() as u64)) + .finish(); + }); + let (db, _runner) = builder.build().await.expect("build db"); + + // Built, never polled: nothing dials. + let futures = connector().build(&db).await.expect("build connector"); + assert_eq!(futures.len(), 1, "the session task, and no inbound pump"); +}