>,
}
impl<'r, 'a, T> OutboundConnectorBuilder<'r, 'a, T>
@@ -856,26 +807,42 @@ where
self
}
- /// Sets a dynamic topic provider
- ///
- /// The provider receives the value being published and returns
- /// the topic/destination to publish to. Return `None` to use the default
- /// static topic from the URL.
- ///
- /// # Type Safety
+ /// Sets a [`TopicWriter`](crate::connector::TopicWriter) that writes each
+ /// value's destination.
///
- /// The provider is type-checked at compile time against `T` and stays
- /// typed end-to-end: it is fused into the link's serialized source and
- /// called with `&T` per value.
- pub fn with_topic_provider(mut self, provider: P) -> Self
+ /// `capacity` is the longest topic the writer produces, in bytes. A value
+ /// whose topic does not fit is skipped and logged. For a closure, use
+ /// [`with_topic_fn`](Self::with_topic_fn).
+ pub fn with_topic_writer(mut self, capacity: usize, writer: W) -> Self
where
- P: crate::connector::TopicProvider + 'static,
+ W: crate::connector::TopicWriter + 'static,
{
- // Stays typed: fused into the link's SerializedSource at finish().
- self.topic_provider = Some(Arc::new(provider));
+ self.topic_writer = Some((capacity, Arc::new(writer)));
self
}
+ /// Sets a closure that writes each value's destination; see
+ /// [`with_topic_writer`](Self::with_topic_writer).
+ ///
+ /// ```rust,ignore
+ /// .with_topic_fn(32, |v, out| {
+ /// write!(out, "sensors/{}/{}", v.site, v.id)?;
+ /// Ok(true)
+ /// })
+ /// ```
+ pub fn with_topic_fn(self, capacity: usize, f: F) -> Self
+ where
+ F: Fn(
+ &T,
+ &mut crate::connector::TopicBuf<'_>,
+ ) -> Result
+ + Send
+ + Sync
+ + 'static,
+ {
+ self.with_topic_writer(capacity, f)
+ }
+
/// Finalizes the connector registration
///
/// Configuration mistakes — an invalid URL, a missing serializer, or an
@@ -969,38 +936,32 @@ where
self.registrar.last_stage = Some((StageKind::Link, 0));
}
- // Fused source factory that captures type T and record key.
- //
- // Resolves the record at route-collection time (not per-message) and
- // constructs a `Consumer` bound to a pre-resolved buffer handle —
- // same pattern as the build-time path in
- // `TypedRecord::collect_consumer_futures`. The serializer
- // and topic provider ride along typed, so the readers handed to the
- // pumps yield destination + payload with no erasure crossing.
+ // Resolves the record and builds a `Consumer` bound to its buffer
+ // handle, once per route (not per message) — same pattern as the
+ // build-time path in `TypedRecord::collect_consumer_futures`.
//
- // The factory runs during build() after every record is registered and
- // validated (including the linked-records-need-a-buffer check), so
- // failures here are aimdb bugs, not user mistakes.
+ // The factories run during build() after every record is registered
+ // and validated (including the linked-records-need-a-buffer check),
+ // so failures here are aimdb bugs, not user mistakes.
#[allow(
clippy::panic,
reason = "the factory returns no Result and these lookups were validated at build() time"
)]
- let source_factory: crate::connector::SourceFactoryFn = {
+ let make_consumer: ConsumerFactoryFn = {
let record_key = self.registrar.record_key.clone();
- let topic_provider = self.topic_provider;
Arc::new(move |db: &AimDb| {
let typed_rec = db
.inner()
.get_typed_record_by_key::(&record_key)
.unwrap_or_else(|e| {
panic!(
- "source factory: record '{record_key}' lookup failed ({e:?}) — \
+ "outbound link: record '{record_key}' lookup failed ({e:?}) — \
this is a bug in aimdb-core"
)
});
let buffer = typed_rec.buffer_handle().unwrap_or_else(|| {
panic!(
- "source factory: record '{record_key}' has no buffer despite \
+ "outbound link: record '{record_key}' has no buffer despite \
build()-time validation — this is a bug in aimdb-core"
)
});
@@ -1009,20 +970,35 @@ where
let mut consumer = Consumer::::new(buffer);
#[cfg(feature = "observability")]
consumer.set_profiling(link_metrics.clone(), db.profiling_clock().clone());
- Box::new(FusedSource {
- consumer,
- serialize: serialize.clone(),
- serialize_into: serialize_into.clone(),
- topic: topic_provider.clone(),
- }) as Box
+ consumer
+ })
+ };
+
+ // Subscribed when `OutboundRoutes` is built.
+ let route_factory: crate::outbound::RouteFactoryFn = {
+ let topic_writer = self.topic_writer;
+ Arc::new(move |db: &AimDb| {
+ let (topic_capacity, writer) = match &topic_writer {
+ Some((capacity, writer)) => (*capacity, Some(writer.clone())),
+ None => (0, None),
+ };
+ crate::outbound::RouteParts {
+ route: Box::new(TypedRoute {
+ reader: make_consumer(db).subscribe(),
+ writer,
+ serialize: serialize.clone(),
+ serialize_into: serialize_into.as_ref().map(|(_, f)| f.clone()),
+ }),
+ topic_capacity,
+ payload_capacity: serialize_into.as_ref().map_or(0, |(capacity, _)| *capacity),
+ }
})
};
- let mut link = ConnectorLink::new(url, source_factory);
+ let mut link = ConnectorLink::new(url, route_factory);
link.config = self.config;
- // Store the connector link - sources will be created later in build()
- // after connectors are actually built
+ // Routes are built later, when a connector builds its `OutboundRoutes`.
self.registrar.rec.add_outbound_connector(link);
self.registrar
}
@@ -1292,10 +1268,7 @@ where
#[cfg(test)]
mod tests {
use super::*;
- use crate::{
- connector::{SerializeError, SerializedPayload, SerializedReader as _, TopicProvider},
- DbResult,
- };
+ use crate::{connector::SerializeError, DbResult};
use core::pin::Pin;
#[cfg(not(feature = "std"))]
@@ -1992,10 +1965,11 @@ mod tests {
}
}
- /// Routes `payload` on `topic` through the `mqtt` inbound router.
+ /// Dispatches `payload` on `topic` through the `mqtt` inbound routes.
fn route(db: &crate::AimDb, topic: &str, payload: &[u8]) {
- let router = db.inbound_router("mqtt", &Plus).expect("routes compile");
- router.route(topic, payload, &db.runtime_ctx()).unwrap();
+ crate::InboundDispatch::new(db, "mqtt", &Plus)
+ .expect("routes compile")
+ .dispatch(topic, payload);
}
/// End-to-end inbound path: bytes → fused ingest → typed buffer push,
@@ -2042,7 +2016,7 @@ mod tests {
}
// ====================================================================
- // inbound_router: patterns, keys and connector-build errors
+ // InboundDispatch: patterns, keys and connector-build errors
// ====================================================================
use crate::topic_pattern::test_support::Plus;
@@ -2079,8 +2053,16 @@ mod tests {
})
}
+ fn config_errors(result: crate::DbResult) -> Vec {
+ match result {
+ Err(crate::DbError::InvalidConfiguration { errors }) => errors,
+ Err(e) => panic!("unexpected error {e:?}"),
+ Ok(_) => panic!("expected configuration errors"),
+ }
+ }
+
#[tokio::test]
- async fn inbound_router_routes_patterns_with_shared_keys() {
+ async fn inbound_dispatch_routes_patterns_with_shared_keys() {
let keys: Arc>> = Default::default();
let seen = keys.clone();
let (db, last, count) = inbound_db(move |reg| {
@@ -2101,25 +2083,25 @@ mod tests {
})
.await;
- let router = db.inbound_router("mqtt", &Plus).expect("routes compile");
- let subscriptions: Vec = router
+ let inbound = crate::InboundDispatch::new(&db, "mqtt", &Plus).expect("routes compile");
+ let subscriptions: Vec = inbound
.subscriptions()
.iter()
.map(|s| s.to_string())
.collect();
assert_eq!(subscriptions, ["cmd/in", "hum/+", "temp/+"]);
+ assert_eq!(inbound.route_count(), 3);
- let ctx = db.runtime_ctx();
- let route = |topic: &str| {
- router.route(topic, b"", &ctx).unwrap();
+ let dispatch = |topic: &str| {
+ inbound.dispatch(topic, b"");
last.load(Ordering::SeqCst)
};
- assert_eq!(route("temp/a"), 0);
- assert_eq!(route("hum/b"), 1);
- assert_eq!(route("hum/a"), 0, "one key table per record");
- assert_eq!(route("cmd/in"), 100);
+ assert_eq!(dispatch("temp/a"), 0);
+ assert_eq!(dispatch("hum/b"), 1);
+ assert_eq!(dispatch("hum/a"), 0, "one key table per record");
+ assert_eq!(dispatch("cmd/in"), 100);
let produced = count.load(Ordering::SeqCst);
- route("temp/c");
+ dispatch("temp/c");
assert_eq!(count.load(Ordering::SeqCst), produced, "full table drops");
let a = keys.lock()[0];
@@ -2135,16 +2117,8 @@ mod tests {
}
}
- fn config_errors(result: crate::DbResult) -> Vec {
- match result {
- Err(crate::DbError::InvalidConfiguration { errors }) => errors,
- Err(e) => panic!("unexpected error {e:?}"),
- Ok(_) => panic!("expected configuration errors"),
- }
- }
-
#[tokio::test]
- async fn inbound_router_rejects_links_it_cannot_compile() {
+ async fn inbound_dispatch_rejects_links_it_cannot_compile() {
let (db, _, _) = inbound_db(|reg| {
reg.link_from("mqtt://s/{d}/t")
.with_deserializer(|_ctx, _bytes: &[u8]| Ok(TestRecord { value: 0 }))
@@ -2161,7 +2135,7 @@ mod tests {
})
.await;
- let errors = config_errors(db.inbound_router("mqtt", &Plus));
+ let errors = config_errors(crate::InboundDispatch::new(&db, "mqtt", &Plus));
assert_eq!(errors.len(), 2, "{errors:?}");
assert!(errors.iter().all(|e| e.record_key == "rec.in"));
assert!(errors[0].message.contains("unbalanced '{' in 'r/{d'"));
@@ -2169,41 +2143,43 @@ mod tests {
.message
.contains("key 'd' is not a capture of 'k/{x}'"));
- let errors = config_errors(db.inbound_router("mqtt", &crate::ExactGrammar));
+ let errors = config_errors(crate::InboundDispatch::new(
+ &db,
+ "mqtt",
+ &crate::ExactGrammar,
+ ));
assert_eq!(errors.len(), 3, "{errors:?}");
assert!(errors[0]
.message
.contains("does not support topic patterns"));
}
- #[cfg(feature = "connector-session")]
#[tokio::test]
- async fn pump_source_routes_through_the_given_router() {
- struct Once(Option<(String, crate::Payload)>);
- impl crate::Source for Once {
- fn next(&mut self) -> crate::BoxFut<'_, Option<(String, crate::Payload)>> {
- let next = self.0.take();
- Box::pin(async move { next })
- }
- }
+ async fn inbound_dispatch_clones_share_records() {
+ fn assert_send_sync() {}
+ assert_send_sync::();
- let (db, last, _) = inbound_db(|reg| {
- reg.link_from("mqtt://temp/{d}")
- .key("d", 2)
- .with_match_deserializer(key_deser)
+ let (db, last, count) = inbound_db(|reg| {
+ reg.link_from("mqtt://cmd/in")
+ .with_deserializer(|_ctx, bytes: &[u8]| {
+ Ok(TestRecord {
+ value: bytes.len() as i32,
+ })
+ })
.finish();
})
.await;
- let router = db.inbound_router("mqtt", &Plus).unwrap();
- let source = Once(Some(("temp/a".into(), Arc::from(&b"x"[..]))));
- for pump in crate::pump_source(&db, router, source) {
- pump.await;
- }
- assert_eq!(last.load(Ordering::SeqCst), 0);
+
+ let inbound = crate::InboundDispatch::new(&db, "mqtt", &Plus).unwrap();
+ let clone = inbound.clone();
+ drop(inbound);
+ clone.dispatch("cmd/in", b"abcd");
+ assert_eq!(count.load(Ordering::SeqCst), 1);
+ assert_eq!(last.load(Ordering::SeqCst), 4);
}
// ====================================================================
- // Fused outbound reader tests
+ // OutboundRoutes over typed links
// ====================================================================
/// Buffer reader that replays a fixed script, then reports the buffer
@@ -2212,319 +2188,294 @@ mod tests {
script: Vec>,
}
- impl ScriptedReader {
- fn closed() -> crate::DbError {
- crate::DbError::BufferClosed {
- buffer_name: "scripted".to_string(),
- }
- }
- }
-
impl crate::buffer::BufferReader for ScriptedReader {
fn poll_recv(
&mut self,
_cx: &mut core::task::Context<'_>,
) -> core::task::Poll> {
let next = if self.script.is_empty() {
- Err(Self::closed())
+ Err(crate::DbError::BufferClosed {
+ buffer_name: "scripted".to_string(),
+ })
} else {
self.script.remove(0)
};
core::task::Poll::Ready(next)
}
fn try_recv(&mut self) -> Result {
- unimplemented!("not needed for fused reader tests")
+ unimplemented!("not needed for outbound route tests")
}
}
- fn lagged() -> crate::DbError {
- crate::DbError::BufferLagged {
- lag_count: 1,
- buffer_name: "scripted".to_string(),
+ type Script = fn() -> Vec>;
+
+ /// Buffer whose readers replay `script`.
+ struct ScriptedBuffer(Script);
+
+ impl crate::buffer::DynBuffer for ScriptedBuffer {
+ fn push(&self, _value: TestRecord) {}
+ fn subscribe_boxed(&self) -> Box + Send> {
+ Box::new(ScriptedReader { script: (self.0)() })
+ }
+ fn as_any(&self) -> &dyn core::any::Any {
+ self
}
}
- fn fused_reader(
- script: Vec>,
- serialize: FusedSerializeFn,
- topic: Option>>,
- ) -> FusedReader {
- FusedReader {
- inner: crate::buffer::Reader::new(Box::new(ScriptedReader { script })),
- serialize,
- serialize_into: None,
- topic,
- }
+ fn values(values: &[i32]) -> Vec> {
+ values
+ .iter()
+ .map(|&value| Ok(TestRecord { value }))
+ .collect()
}
- fn fused_reader_into(
- script: Vec>,
- serialize: FusedSerializeFn,
- serialize_into: FusedSerializeIntoFn,
- ) -> FusedReader {
- FusedReader {
- inner: crate::buffer::Reader::new(Box::new(ScriptedReader { script })),
- serialize,
- serialize_into: Some(serialize_into),
- topic: None,
- }
+ fn le(_ctx: crate::RuntimeContext, r: &TestRecord) -> Result, SerializeError> {
+ Ok(r.value.to_le_bytes().to_vec())
}
- fn test_ctx() -> crate::RuntimeContext {
- crate::RuntimeContext::new(Arc::new(MockRuntime))
+ /// Writes `r` into `out`.
+ fn le_into(
+ _ctx: crate::RuntimeContext,
+ r: &TestRecord,
+ out: &mut [u8],
+ ) -> Result {
+ let bytes = r.value.to_le_bytes();
+ out.get_mut(..bytes.len())
+ .ok_or(SerializeError::BufferTooSmall)?
+ .copy_from_slice(&bytes);
+ Ok(bytes.len())
}
- /// Buffer errors propagate through the fused reader unchanged, so the
- /// pumps keep their `BufferLagged => continue / Err => break` shape.
- #[tokio::test]
- async fn fused_reader_propagates_buffer_errors() {
- let mut reader = fused_reader(
- vec![
- Ok(TestRecord { value: 1 }),
- Err(lagged()),
- Ok(TestRecord { value: 2 }),
- ],
- Arc::new(|_ctx, r| Ok(r.value.to_le_bytes().to_vec())),
- None,
- );
- let ctx = test_ctx();
+ /// `OutboundRoutes` for the `mqtt` links `links` registers on a record
+ /// whose buffer replays `script`.
+ async fn outbound(
+ script: Script,
+ links: impl FnOnce(&mut RecordRegistrar<'_, TestRecord>) + Send + 'static,
+ ) -> (crate::AimDb, crate::OutboundRoutes) {
+ let mut builder = crate::AimDbBuilder::new()
+ .runtime(Arc::new(MockRuntime))
+ .with_connector(NoopConnectorBuilder);
+ builder.configure::("rec.out", move |reg| {
+ reg.buffer_raw(Box::new(ScriptedBuffer(script)));
+ links(reg);
+ });
+ let (db, _runner) = builder.build().await.expect("build must succeed");
+ let outbound = crate::OutboundRoutes::new(&db, "mqtt").expect("outbound routes");
+ (db, outbound)
+ }
- let first = reader.recv(&ctx).await.expect("first value");
- assert_eq!(first.payload, 1i32.to_le_bytes().to_vec());
- assert_eq!(first.dest, None);
+ /// The next message as `(topic, value, owned)`.
+ async fn pull(o: &mut crate::OutboundRoutes) -> Option<(String, i32, bool)> {
+ let msg = o.next().await?;
+ let owned = matches!(msg.payload, crate::OutboundPayload::Owned(_));
+ let bytes: [u8; 4] = msg.payload.as_slice().try_into().expect("4 bytes");
+ Some((msg.topic.to_string(), i32::from_le_bytes(bytes), owned))
+ }
- let err = reader.recv(&ctx).await.expect_err("lag must propagate");
- assert!(matches!(err, crate::DbError::BufferLagged { .. }));
+ /// A lag is counted and skipped; a closed buffer ends the route.
+ #[tokio::test]
+ async fn outbound_routes_count_lag_and_end_on_close() {
+ let (_db, mut o) = outbound(
+ || {
+ vec![
+ Ok(TestRecord { value: 1 }),
+ Err(lagged()),
+ Ok(TestRecord { value: 2 }),
+ ]
+ },
+ |reg| {
+ reg.link_to("mqtt://tele/out").with_serializer(le).finish();
+ },
+ )
+ .await;
- let second = reader.recv(&ctx).await.expect("second value");
- assert_eq!(second.payload, 2i32.to_le_bytes().to_vec());
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 1, true)));
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 2, true)));
+ assert_eq!(o.stats(0).unwrap().lagged, 1);
+ assert_eq!(pull(&mut o).await, None);
+ }
- let closed = reader.recv(&ctx).await.expect_err("closed must propagate");
- assert!(matches!(closed, crate::DbError::BufferClosed { .. }));
+ fn lagged() -> crate::DbError {
+ crate::DbError::BufferLagged {
+ lag_count: 1,
+ buffer_name: "scripted".to_string(),
+ }
}
- /// Serialization failures are skipped inside the reader (logged), exactly
- /// like the old pump-side `continue`.
+ /// A value whose serializer fails is counted and skipped.
#[tokio::test]
- async fn fused_reader_skips_serialize_failures() {
- let mut reader = fused_reader(
- vec![Ok(TestRecord { value: 13 }), Ok(TestRecord { value: 42 })],
- Arc::new(|_ctx, r| {
- if r.value == 13 {
- Err(SerializeError::InvalidData)
- } else {
- Ok(r.value.to_le_bytes().to_vec())
- }
- }),
- None,
- );
+ async fn outbound_routes_skip_serialize_failures() {
+ let (_db, mut o) = outbound(
+ || values(&[13, 42]),
+ |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_serializer(|ctx, r: &TestRecord| {
+ if r.value == 13 {
+ return Err(SerializeError::InvalidData);
+ }
+ le(ctx, r)
+ })
+ .finish();
+ },
+ )
+ .await;
- // One recv: the failing value is skipped, the next good one returned.
- let msg = reader.recv(&test_ctx()).await.expect("value");
- assert_eq!(msg.payload, 42i32.to_le_bytes().to_vec());
+ assert_eq!(pull(&mut o).await.map(|m| m.1), Some(42));
+ assert_eq!(o.stats(0).unwrap().serialize_failed, 1);
}
#[tokio::test]
- async fn fused_reader_into_uses_scratch_without_owned_fallback() {
+ async fn outbound_routes_use_scratch_without_owned_fallback() {
let owned_calls = Arc::new(AtomicUsize::new(0));
let into_calls = Arc::new(AtomicUsize::new(0));
- let owned_counter = owned_calls.clone();
- let into_counter = into_calls.clone();
- let mut reader = fused_reader_into(
- vec![Ok(TestRecord { value: 7 })],
- Arc::new(move |_ctx, r| {
- owned_counter.fetch_add(1, Ordering::SeqCst);
- Ok(r.value.to_le_bytes().to_vec())
- }),
- Arc::new(move |_ctx, r, out| {
- into_counter.fetch_add(1, Ordering::SeqCst);
- let bytes = r.value.to_le_bytes();
- out.get_mut(..bytes.len())
- .ok_or(SerializeError::BufferTooSmall)?
- .copy_from_slice(&bytes);
- Ok(bytes.len())
- }),
- );
- let mut scratch = [0_u8; 8];
-
- let msg = reader
- .recv_into(&test_ctx(), &mut scratch)
- .await
- .expect("value");
+ let (owned_counter, into_counter) = (owned_calls.clone(), into_calls.clone());
+ let (_db, mut o) = outbound(
+ || values(&[7]),
+ move |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_serializer(move |ctx, r: &TestRecord| {
+ owned_counter.fetch_add(1, Ordering::SeqCst);
+ le(ctx, r)
+ })
+ .with_serializer_into(8, move |ctx, r: &TestRecord, out| {
+ into_counter.fetch_add(1, Ordering::SeqCst);
+ le_into(ctx, r, out)
+ })
+ .finish();
+ },
+ )
+ .await;
- assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 });
- assert_eq!(&scratch[..4], 7i32.to_le_bytes().as_slice());
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 7, false)));
assert_eq!(into_calls.load(Ordering::SeqCst), 1);
assert_eq!(owned_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
- async fn fused_reader_into_falls_back_once_when_scratch_is_small() {
+ async fn outbound_routes_fall_back_once_when_scratch_is_small() {
let owned_calls = Arc::new(AtomicUsize::new(0));
let owned_counter = owned_calls.clone();
- let mut reader = fused_reader_into(
- vec![Ok(TestRecord { value: 9 })],
- Arc::new(move |_ctx, r| {
- owned_counter.fetch_add(1, Ordering::SeqCst);
- Ok(r.value.to_le_bytes().to_vec())
- }),
- Arc::new(|_ctx, _r, _out| Err(SerializeError::BufferTooSmall)),
- );
- let mut scratch = [0_u8; 2];
-
- let msg = reader
- .recv_into(&test_ctx(), &mut scratch)
- .await
- .expect("fallback value");
+ let (_db, mut o) = outbound(
+ || values(&[9]),
+ move |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_serializer(move |ctx, r: &TestRecord| {
+ owned_counter.fetch_add(1, Ordering::SeqCst);
+ le(ctx, r)
+ })
+ .with_serializer_into(2, |_ctx, _r: &TestRecord, _out| {
+ Err(SerializeError::BufferTooSmall)
+ })
+ .finish();
+ },
+ )
+ .await;
- assert_eq!(
- msg.payload,
- SerializedPayload::Owned(9i32.to_le_bytes().to_vec())
- );
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 9, true)));
assert_eq!(owned_calls.load(Ordering::SeqCst), 1);
}
+ /// An into-slice serializer that reports more bytes than the scratch
+ /// holds, or fails with anything but `BufferTooSmall`, skips the value.
#[tokio::test]
- async fn fused_reader_into_rejects_invalid_length_and_skips_value() {
- let mut reader = fused_reader_into(
- vec![Ok(TestRecord { value: 1 }), Ok(TestRecord { value: 2 })],
- Arc::new(|_ctx, r| Ok(r.value.to_le_bytes().to_vec())),
- Arc::new(|_ctx, r, out| {
- if r.value == 1 {
- return Ok(out.len() + 1);
- }
- let bytes = r.value.to_le_bytes();
- out[..bytes.len()].copy_from_slice(&bytes);
- Ok(bytes.len())
- }),
- );
- let mut scratch = [0_u8; 8];
-
- let msg = reader
- .recv_into(&test_ctx(), &mut scratch)
- .await
- .expect("second value");
+ async fn outbound_routes_skip_invalid_lengths_and_data() {
+ let (_db, mut o) = outbound(
+ || values(&[1, 2, 3]),
+ |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_serializer(le)
+ .with_serializer_into(8, |ctx, r: &TestRecord, out| match r.value {
+ 1 => Ok(out.len() + 1),
+ 2 => Err(SerializeError::InvalidData),
+ _ => le_into(ctx, r, out),
+ })
+ .finish();
+ },
+ )
+ .await;
- assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 });
- assert_eq!(&scratch[..4], 2i32.to_le_bytes().as_slice());
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 3, false)));
+ assert_eq!(o.stats(0).unwrap().serialize_failed, 2);
}
+ /// `with_topic_fn` infers an unannotated closure's argument types;
+ /// `Ok(false)` publishes to the link's default topic.
#[tokio::test]
- async fn fused_reader_into_skips_invalid_data() {
- let mut reader = fused_reader_into(
- vec![Ok(TestRecord { value: 1 }), Ok(TestRecord { value: 2 })],
- Arc::new(|_ctx, r| Ok(r.value.to_le_bytes().to_vec())),
- Arc::new(|_ctx, r, out| {
- if r.value == 1 {
- return Err(SerializeError::InvalidData);
- }
- let bytes = r.value.to_le_bytes();
- out[..bytes.len()].copy_from_slice(&bytes);
- Ok(bytes.len())
- }),
- );
- let mut scratch = [0_u8; 8];
-
- let msg = reader
- .recv_into(&test_ctx(), &mut scratch)
- .await
- .expect("second value");
+ async fn with_topic_fn_writes_the_topic_or_the_default() {
+ use core::fmt::Write as _;
+ let (_db, mut o) = outbound(
+ || values(&[5, 0]),
+ |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_topic_fn(8, |v, out| {
+ if v.value <= 0 {
+ return Ok(false);
+ }
+ write!(out, "dyn/{}", v.value)?;
+ Ok(true)
+ })
+ .with_serializer(le)
+ .finish();
+ },
+ )
+ .await;
- assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 });
- assert_eq!(&scratch[..4], 2i32.to_le_bytes().as_slice());
+ assert_eq!(o.routes()[0].topic_capacity, 8);
+ assert_eq!(pull(&mut o).await.map(|m| m.0), Some("dyn/5".into()));
+ assert_eq!(pull(&mut o).await.map(|m| m.0), Some("tele/out".into()));
}
- /// The destination is resolved from the typed value while it is in hand.
+ /// A topic that overflows skips its value, even when the writer ignores
+ /// the error and returns `Ok(true)`.
#[tokio::test]
- async fn fused_reader_resolves_dynamic_topic() {
- struct PositiveTopic;
- impl TopicProvider for PositiveTopic {
- fn topic(&self, value: &TestRecord) -> Option {
- (value.value > 0).then(|| alloc::format!("dyn/{}", value.value))
- }
- }
-
- let mut reader = fused_reader(
- vec![Ok(TestRecord { value: 5 }), Ok(TestRecord { value: 0 })],
- Arc::new(|_ctx, r| Ok(r.value.to_le_bytes().to_vec())),
- Some(Arc::new(PositiveTopic)),
- );
- let ctx = test_ctx();
-
- let first = reader.recv(&ctx).await.expect("value");
- assert_eq!(first.dest.as_deref(), Some("dyn/5"));
+ async fn overflowing_topics_are_skipped() {
+ use core::fmt::Write as _;
+ let (_db, mut o) = outbound(
+ || values(&[123_456, 1_234_567, 7]),
+ |reg| {
+ reg.link_to("mqtt://tele/out")
+ .with_topic_fn(6, |v, out| {
+ if v.value == 1_234_567 {
+ let _ = write!(out, "t/{}", v.value);
+ return Ok(true);
+ }
+ write!(out, "t/{}", v.value)?;
+ Ok(true)
+ })
+ .with_serializer(le)
+ .finish();
+ },
+ )
+ .await;
- let second = reader.recv(&ctx).await.expect("value");
- assert_eq!(second.dest, None); // falls back to the route default
+ // "t/123456" (8 bytes) and "t/1234567" (9 bytes) do not fit in 6.
+ assert_eq!(pull(&mut o).await, Some(("t/7".into(), 7, true)));
+ assert_eq!(o.stats(0).unwrap().topic_overflow, 2);
}
- /// End-to-end outbound path: registrar → build → collect → subscribe →
- /// recv, pinning the factory wiring (raw and context serializers).
+ /// Registrar → build → `OutboundRoutes`, pinning the factory wiring: the
+ /// serializer set last wins, whichever way its context is typed.
#[tokio::test]
async fn outbound_roundtrip_yields_serialized_values() {
- /// Buffer whose readers replay one canned value, then close.
- struct CannedBuffer;
- impl crate::buffer::DynBuffer for CannedBuffer {
- fn push(&self, _value: TestRecord) {}
- fn subscribe_boxed(&self) -> Box + Send> {
- Box::new(ScriptedReader {
- script: vec![Ok(TestRecord { value: 5 })],
- })
- }
- fn as_any(&self) -> &dyn core::any::Any {
- self
- }
- }
-
- struct FixedTopic;
- impl TopicProvider for FixedTopic {
- fn topic(&self, value: &TestRecord) -> Option {
- Some(alloc::format!("dyn/{}", value.value))
- }
- }
-
- let mut builder = crate::AimDbBuilder::new()
- .runtime(Arc::new(MockRuntime))
- .with_connector(NoopConnectorBuilder);
- builder.configure::("rec.out", |reg| {
- reg.buffer_raw(Box::new(CannedBuffer));
- // Raw set first, context set last — context must win (the kind
- // enum is gone; mutual exclusion is behavior now).
- reg.link_to("mqtt://tele/out")
- .with_topic_provider(FixedTopic)
- .with_serializer(|_ctx, _r: &TestRecord| Ok(vec![0]))
- .with_serializer(|_ctx: crate::RuntimeContext, r: &TestRecord| {
- Ok(r.value.to_le_bytes().to_vec())
- })
- .with_serializer_into(4, |_ctx, r: &TestRecord, out| {
- let encoded = r.value.to_le_bytes();
- let dest = out
- .get_mut(..encoded.len())
- .ok_or(SerializeError::BufferTooSmall)?;
- dest.copy_from_slice(&encoded);
- Ok(encoded.len())
- })
- .finish();
- });
- let (db, _runner) = builder.build().await.expect("build must succeed");
+ let (_db, mut o) = outbound(
+ || values(&[5]),
+ |reg| {
+ reg.link_to("mqtt://tele/out?qos=1")
+ .with_serializer(|_ctx, _r: &TestRecord| Ok(vec![0]))
+ .with_serializer(|_ctx: crate::RuntimeContext, r: &TestRecord| {
+ Ok(r.value.to_le_bytes().to_vec())
+ })
+ .with_serializer_into(4, |ctx, r: &TestRecord, out| le_into(ctx, r, out))
+ .finish();
+ },
+ )
+ .await;
- let routes = db.collect_outbound_routes("mqtt");
- assert_eq!(routes.len(), 1);
- assert_eq!(routes[0].topic, "tele/out");
- assert_eq!(routes[0].source.serializer_scratch_capacity(), Some(4));
-
- let mut reader = routes[0].source.subscribe();
- let ctx = db.runtime_ctx();
- let mut scratch = [0u8; 4];
- let msg = reader.recv_into(&ctx, &mut scratch).await.expect("value");
- assert_eq!(msg.dest.as_deref(), Some("dyn/5"));
- assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 });
- assert_eq!(scratch, 5i32.to_le_bytes());
-
- let closed = reader
- .recv_into(&ctx, &mut scratch)
- .await
- .expect_err("buffer closed");
- assert!(matches!(closed, crate::DbError::BufferClosed { .. }));
+ let route = &o.routes()[0];
+ assert_eq!(&*route.default_topic, "tele/out");
+ assert_eq!(route.payload_capacity, 4);
+ assert_eq!(route.config.record_index, Some(0));
+ assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 5, false)));
+ assert_eq!(pull(&mut o).await, None);
}
}
diff --git a/aimdb-data-contracts/src/link_codec.rs b/aimdb-data-contracts/src/link_codec.rs
index aa5a6b52..d49cbb54 100644
--- a/aimdb-data-contracts/src/link_codec.rs
+++ b/aimdb-data-contracts/src/link_codec.rs
@@ -226,13 +226,16 @@ mod tests {
#[cfg(all(feature = "linkable-json", feature = "linkable-postcard"))]
use aimdb_core::buffer::{BufferReader, DynBuffer};
- use aimdb_core::connector::SerializeError;
#[cfg(all(feature = "linkable-json", feature = "linkable-postcard"))]
- use aimdb_core::connector::{ConnectorBuilder, SerializedPayload};
+ use aimdb_core::connector::ConnectorBuilder;
+ use aimdb_core::connector::SerializeError;
#[cfg(all(feature = "linkable-json", feature = "linkable-postcard"))]
use aimdb_core::executor::test_support::NoopRuntimeOps;
#[cfg(all(feature = "linkable-json", feature = "linkable-postcard"))]
- use aimdb_core::{AimDb, AimDbBuilder, BoxFuture, DbError, DbResult};
+ use aimdb_core::{
+ AimDb, AimDbBuilder, BoxFuture, DbError, DbResult, InboundDispatch, OutboundPayload,
+ OutboundRoutes,
+ };
use serde::{Deserialize, Serialize};
#[cfg(all(feature = "linkable-json", feature = "linkable-postcard"))]
@@ -523,106 +526,72 @@ mod tests {
});
let (db, _runner) = builder.build().await.expect("codec routes build");
- let routes = db.collect_outbound_routes("test");
- assert_eq!(routes.len(), 4);
-
- let json_route = routes
- .iter()
- .find(|route| route.topic == "json")
- .expect("JSON route");
- assert_eq!(json_route.source.serializer_scratch_capacity(), None);
- let mut json_reader = json_route.source.subscribe();
- let mut unused_scratch = [];
- let json_message = json_reader
- .recv_into(&db.runtime_ctx(), &mut unused_scratch)
- .await
- .expect("JSON payload");
- let SerializedPayload::Owned(json_bytes) = json_message.payload else {
- panic!("JSON must use owned serialization");
+ let mut outbound = OutboundRoutes::new(&db, "test").expect("outbound routes");
+ let route = |topic: &str| {
+ outbound
+ .routes()
+ .iter()
+ .find(|route| &*route.default_topic == topic)
+ .expect("route")
+ .clone()
};
+ assert_eq!(outbound.routes().len(), 4);
+ // A payload capacity of 0 means owned serialization only.
+ assert_eq!(route("json").payload_capacity, 0);
+ assert_eq!(route("postcard").payload_capacity, 64);
+ assert!(route("postcard")
+ .config
+ .protocol_options
+ .contains(&("wire".to_string(), "binary".to_string())));
+ assert_eq!(route("postcard-owned-fallback").payload_capacity, 1);
+ assert_eq!(
+ route("postcard-replaced-by-json").payload_capacity,
+ 0,
+ "owned-only replacement must clear the previous scratch codec"
+ );
+
+ // Each canned reader yields one value, then closes its route.
+ let mut messages = std::collections::BTreeMap::new();
+ while let Some(msg) = outbound.next().await {
+ let owned = matches!(msg.payload, OutboundPayload::Owned(_));
+ messages.insert(msg.topic.to_string(), (msg.payload.into_vec(), owned));
+ }
+ let message = |topic: &str| messages.get(topic).expect("message").clone();
+
+ let (json_bytes, owned) = message("json");
+ assert!(owned, "JSON must use owned serialization");
let decoded_json: Reading = link_codecs::Json.decode(&json_bytes).expect("JSON decode");
assert_eq!(decoded_json, reading);
- let postcard_route = routes
- .iter()
- .find(|route| route.topic == "postcard")
- .expect("Postcard route");
- assert_eq!(
- postcard_route.source.serializer_scratch_capacity(),
- Some(64)
- );
- assert!(postcard_route
- .config
- .contains(&("wire".to_string(), "binary".to_string())));
- let mut postcard_reader = postcard_route.source.subscribe();
- let mut scratch = [0_u8; 64];
- let postcard_message = postcard_reader
- .recv_into(&db.runtime_ctx(), &mut scratch)
- .await
- .expect("Postcard payload");
- let SerializedPayload::Scratch { len } = postcard_message.payload else {
- panic!("Postcard must use route scratch storage");
- };
+ let (postcard_bytes, owned) = message("postcard");
+ assert!(!owned, "Postcard must use route scratch storage");
let decoded_postcard: Reading = link_codecs::Postcard::<64>
- .decode(&scratch[..len])
+ .decode(&postcard_bytes)
.expect("Postcard decode");
assert_eq!(decoded_postcard, reading);
- assert_ne!(json_bytes, scratch[..len]);
-
- let fallback_route = routes
- .iter()
- .find(|route| route.topic == "postcard-owned-fallback")
- .expect("undersized Postcard route");
- assert_eq!(fallback_route.source.serializer_scratch_capacity(), Some(1));
- let mut fallback_reader = fallback_route.source.subscribe();
- let mut undersized_scratch = [0_u8; 1];
- let fallback_message = fallback_reader
- .recv_into(&db.runtime_ctx(), &mut undersized_scratch)
- .await
- .expect("owned Postcard fallback");
- let SerializedPayload::Owned(fallback_bytes) = fallback_message.payload else {
- panic!("undersized Postcard scratch must use owned fallback");
- };
+ assert_ne!(json_bytes, postcard_bytes);
+
+ let (fallback_bytes, owned) = message("postcard-owned-fallback");
+ assert!(owned, "undersized Postcard scratch must use owned fallback");
let decoded_fallback: Reading = link_codecs::Postcard::<1>
.decode(&fallback_bytes)
.expect("fallback Postcard decode");
assert_eq!(decoded_fallback, reading);
- let replacement_route = routes
- .iter()
- .find(|route| route.topic == "postcard-replaced-by-json")
- .expect("replacement route");
- assert_eq!(
- replacement_route.source.serializer_scratch_capacity(),
- None,
- "owned-only replacement must clear the previous scratch codec"
- );
- let mut replacement_reader = replacement_route.source.subscribe();
- let replacement_message = replacement_reader
- .recv_into(&db.runtime_ctx(), &mut [])
- .await
- .expect("replacement JSON payload");
- let SerializedPayload::Owned(replacement_bytes) = replacement_message.payload else {
- panic!("replacement JSON codec must use owned serialization");
- };
+ let (replacement_bytes, owned) = message("postcard-replaced-by-json");
+ assert!(owned, "replacement JSON codec must use owned serialization");
assert_eq!(replacement_bytes, json_bytes);
- let inbound = db
- .inbound_router("test", &aimdb_core::ExactGrammar)
+ let inbound = InboundDispatch::new(&db, "test", &aimdb_core::ExactGrammar)
.expect("inbound routes");
assert_eq!(inbound.route_count(), 2);
- let ctx = db.runtime_ctx();
- inbound
- .route("json-in", &json_bytes, &ctx)
- .expect("JSON ingest");
+ inbound.dispatch("json-in", &json_bytes);
assert_eq!(
json_in.lock().expect("JSON capture lock").as_ref(),
Some(&reading)
);
- inbound
- .route("postcard-in", &scratch[..len], &ctx)
- .expect("Postcard ingest");
+ inbound.dispatch("postcard-in", &postcard_bytes);
assert_eq!(
postcard_in.lock().expect("Postcard capture lock").as_ref(),
Some(&reading)
diff --git a/aimdb-embassy-adapter/CHANGELOG.md b/aimdb-embassy-adapter/CHANGELOG.md
index bd8fdbe0..8f931a96 100644
--- a/aimdb-embassy-adapter/CHANGELOG.md
+++ b/aimdb-embassy-adapter/CHANGELOG.md
@@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
+### Removed (breaking)
+
+- **`EmbassySinkRaw`, `EmbassySink`, `EmbassySourceRaw` and `EmbassySource`**.
+ They force-`Send`ed a `!Send` sink or source so it could ride
+ core's `pump_sink` / `pump_source`, which no longer exist: a connector now
+ runs its own task and pulls from `OutboundRoutes`. The `into_box_future`
+ helper and `NetStack` remain.
+
## [0.7.0] - 2026-09-18
### Added
diff --git a/aimdb-embassy-adapter/src/connectors.rs b/aimdb-embassy-adapter/src/connectors.rs
index 899cfd9f..770b23b8 100644
--- a/aimdb-embassy-adapter/src/connectors.rs
+++ b/aimdb-embassy-adapter/src/connectors.rs
@@ -1,28 +1,20 @@
-//! The data-plane bridge — the one audited home for the single-core `unsafe` +
-//! [`SendFutureWrapper`] that every Embassy data-plane connector used to
-//! hand-roll.
+//! The one audited home for the single-core `unsafe` that Embassy connectors
+//! need.
//!
//! AimDB's connector contract is `Send`-everywhere (so a Tokio app can
//! `tokio::spawn(runner.run())`). Embassy's primitives (channels over
//! `NoopRawMutex`, a borrowed `embassy_net::Stack`, …) are `!Send` *by design* —
//! single-core, cooperative, no preemption or thread migration. Bridging the two
//! requires force-`Send`ing the Embassy futures; this module does that **once**,
-//! so a connector crate carries **no `unsafe` and no wrapper**.
+//! so a connector crate carries **no `unsafe` and no wrapper**: it boxes its
+//! protocol task with [`into_box_future`] and holds the network stack as a
+//! `NetStack`.
//!
-//! Data-plane transports (MQTT, KNX) contribute an [`EmbassySinkRaw`] (outbound)
-//! and/or [`EmbassySourceRaw`] (inbound) and ride core's
-//! [`pump_sink`](aimdb_core::session::pump_sink) /
-//! [`pump_source`](aimdb_core::session::pump_source) via the force-`Send`
-//! bridges [`EmbassySink`] / [`EmbassySource`].
-//!
-//! Session transports (serial, TCP, …) no longer come through here. They ride
+//! Session transports (serial, TCP, …) do not come through here. They ride
//! core's runtime-neutral spine directly — `SessionClientConnector` /
//! `SessionServerConnector` over `FramedConnection`, with the byte source from
//! this crate's `io` or `net` module (unlinked: neither exists in a
-//! `connectors`-only build) — so the Embassy duals this module used to
-//! carry (`EmbassySessionClient`/`Server`, `EmbassyConnection`, `OneShotCell`
-//! and the one-shot dialer/listener) are gone. Their one-shot semantics live in
-//! core as `OneShot`, `OneShotDialer` and `OneShotListener`.
+//! `connectors`-only build).
//!
//! # Safety invariant (shared by every `unsafe impl` below)
//!
@@ -34,87 +26,12 @@ use core::future::Future;
use core::pin::Pin;
use alloc::boxed::Box;
-use alloc::string::{String, ToString};
-use alloc::vec::Vec;
-
-use aimdb_core::session::{BoxFut, Payload, Source};
-use aimdb_core::transport::{Connector, ConnectorConfig, PublishError};
use crate::SendFutureWrapper;
/// The runner's collected future type (`Send`, as the std contract requires).
type BoxFuture = Pin + Send + 'static>>;
-// ===========================================================================
-// Data-plane bridges — let a `!Send` sink/source ride core's pumps.
-// ===========================================================================
-
-/// The pure outbound I/O a data-plane connector contributes: publish one payload.
-/// The `!Send` dual of [`Connector`]; [`EmbassySink`] force-`Send`s it so it can
-/// drive core's [`pump_sink`](aimdb_core::session::pump_sink).
-///
-/// Args are owned (a data-plane sink enqueues owned data onto its channel anyway),
-/// so the returned future borrows only `&self` — matching [`Connector::publish`]'s
-/// `'_` return shape.
-pub trait EmbassySinkRaw {
- /// Publish `payload` to `destination` (e.g. enqueue onto an Embassy channel).
- fn publish(
- &self,
- destination: String,
- config: ConnectorConfig,
- payload: Vec,
- ) -> impl Future