From 51ca2d1ea2584cf70043f27ce0d8c17cb1e84e72 Mon Sep 17 00:00:00 2001 From: "sounds.like.lx" <147444674+lxsaah@users.noreply.github.com> Date: Sun, 4 Oct 2026 21:08:42 +0200 Subject: [PATCH 1/2] fix(mqtt): advertise Maximum Packet Size in the embedded CONNECT (#284) --- CHANGELOG.md | 7 ++ aimdb-mqtt-connector/CHANGELOG.md | 10 +++ .../src/embedded/session_loop.rs | 69 +++++++++++++++++- aimdb-mqtt-connector/tests/common/mod.rs | 48 +++++++++++-- aimdb-mqtt-connector/tests/tokio_broker.rs | 72 +++++++++++++++++++ 5 files changed, 200 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 620b1899..cd7465e3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -53,6 +53,13 @@ and `pump_client` take the router. KNX, WebSocket, TCP, UDS and serial use `ExactGrammar`: a `{…}` link on them fails the build. The user-facing link API is unchanged. ([aimdb-core](aimdb-core/CHANGELOG.md)) +### Fixed + +- **The embedded MQTT backend advertises the largest packet it receives** + (`Maximum Packet Size` 3,328 in CONNECT), so an oversized retained message is + withheld by the broker instead of reconnecting the client forever. + ([aimdb-mqtt-connector](aimdb-mqtt-connector/CHANGELOG.md)) + ## [2.0.0] - 2026-09-18 ### Added diff --git a/aimdb-mqtt-connector/CHANGELOG.md b/aimdb-mqtt-connector/CHANGELOG.md index 4f0f62e0..5d499e8d 100644 --- a/aimdb-mqtt-connector/CHANGELOG.md +++ b/aimdb-mqtt-connector/CHANGELOG.md @@ -23,6 +23,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **The `with_qos` doc no longer claims an inbound subscribe QoS.** Inbound subscriptions stay at QoS 1, as before. +### Fixed + +- **An oversized retained message no longer reconnects the `Embedded` backend + forever.** The session receives packets of up to 3,328 bytes (its 3,584-byte + buffer minus one read), but its MQTT 5 CONNECT did not say so. A broker + could send a larger packet, which ends the session, and a retained one is + replayed after every SUBSCRIBE, so the client reconnected once per + reconnection delay. CONNECT now advertises `Maximum Packet Size` 3,328, and + the broker withholds anything larger instead of sending it. + ## [0.7.0] - 2026-09-18 ### Changed (breaking) diff --git a/aimdb-mqtt-connector/src/embedded/session_loop.rs b/aimdb-mqtt-connector/src/embedded/session_loop.rs index ec8975aa..2533670d 100644 --- a/aimdb-mqtt-connector/src/embedded/session_loop.rs +++ b/aimdb-mqtt-connector/src/embedded/session_loop.rs @@ -46,6 +46,13 @@ const RX_CHUNK: usize = 256; /// rather than a fixed buffer. const PACKET_BUFFER_SIZE: usize = BUFFER_SIZE - 2 * RX_CHUNK; +/// The largest packet the reader takes whatever arrives before it: the +/// buffer minus one feed chunk (see [`PacketReader`]). Advertised as the +/// CONNECT's Maximum Packet Size, so a broker never sends a larger packet — +/// one that would end the session, and a retained one would end every +/// session it is replayed into. +const MAX_INBOUND_PACKET: usize = PACKET_BUFFER_SIZE - RX_CHUNK; + /// One chunk of freshly read bytes, in flight from the read half to the loop. type Chunk = heapless::Vec; @@ -176,10 +183,13 @@ async fn client_loop( // Topic aliases are declined: honouring them would mean storing the // server's topic names for the life of the connection. let _ = properties.push(ConnectProperty::TopicAliasMaximum(0.into())); + let _ = properties.push(ConnectProperty::MaximumPacketSize( + (MAX_INBOUND_PACKET as u32).into(), + )); // Ours, not `connection_settings.keep_alive()`: that field has no // setter, so it is always mountain-mqtt's own 60 s constant. The // cadence below is derived from the value we actually send. - let connect: Connect<'_, 1, 0> = Connect::new( + let connect: Connect<'_, 2, 0> = Connect::new( settings.keep_alive_secs, *connection_settings.username(), *connection_settings.password(), @@ -664,6 +674,63 @@ mod tests { ); } + /// A QoS 0 PUBLISH to `t` that is exactly `total` bytes on the wire. + fn publish_of(total: usize) -> Vec { + let varint_len = if total - 2 < 128 { 1 } else { 2 }; + let remaining = total - 1 - varint_len; + let mut bytes = alloc::vec![0x30u8]; + if varint_len == 1 { + bytes.push(remaining as u8); + } else { + bytes.push((remaining % 128) as u8 | 0x80); + bytes.push((remaining / 128) as u8); + } + bytes.extend_from_slice(&[0x00, 0x01, b't', 0x00]); + bytes.resize(total, b'x'); + bytes + } + + /// Feeds a `first`-byte packet, a `second`-byte one and a trailing one in + /// `RX_CHUNK` reads, consuming packets as they complete, as the session + /// does. The trailing packet makes the read that completes `second` a full + /// one that also carries the head of the next packet: the worst case. + fn receive_after(first: usize, second: usize) -> Result<(), PacketReadError> { + let mut stream = publish_of(first); + stream.extend_from_slice(&publish_of(second)); + stream.extend_from_slice(&publish_of(RX_CHUNK)); + let mut reader = PacketReader::::new(); + let mut received = 0; + for chunk in stream.chunks(RX_CHUNK) { + reader.feed(chunk)?; + while let Some(total) = reader.framed_len()? { + reader.consume(total); + received += 1; + } + } + assert!(received >= 2); + Ok(()) + } + + /// The advertised Maximum Packet Size is one the reader takes wherever the + /// packet starts inside a read. The reader's stated limit (buffer minus + /// one read) is conservative by one byte; two bytes more fail at some + /// offset. + #[test] + fn the_advertised_maximum_packet_size_is_always_received() { + assert_eq!(MAX_INBOUND_PACKET, 3328); + let offsets = 8..8 + RX_CHUNK; + for first in offsets.clone() { + assert_eq!( + receive_after(first, MAX_INBOUND_PACKET), + Ok(()), + "after {first} bytes" + ); + } + assert!(offsets + .map(|first| receive_after(first, MAX_INBOUND_PACKET + 2)) + .any(|r| r == Err(PacketReadError::PacketTooLargeForBuffer))); + } + #[test] fn a_deadline_in_the_past_still_sleeps_a_tick() { // Never zero: a zero-length sleep would spin the loop. diff --git a/aimdb-mqtt-connector/tests/common/mod.rs b/aimdb-mqtt-connector/tests/common/mod.rs index e2b35ce3..b6a28d5f 100644 --- a/aimdb-mqtt-connector/tests/common/mod.rs +++ b/aimdb-mqtt-connector/tests/common/mod.rs @@ -26,6 +26,10 @@ pub struct Seen { pub keep_alives: Vec, pub subscribes: Vec>, pub published: Vec<(String, Vec)>, + /// The Maximum Packet Size each MQTT 5 CONNECT advertised, when it did. + pub max_packet_sizes: Vec, + /// Pushes withheld because they exceeded the client's Maximum Packet Size. + pub withheld: usize, } impl Seen { @@ -127,6 +131,30 @@ fn connect_keep_alive(body: &[u8]) -> Option { Some(u16::from_be_bytes([*body.get(8)?, *body.get(9)?])) } +/// The Maximum Packet Size property (0x27) of an MQTT 5 CONNECT, if present. +/// Knows the fixed-size properties a client sends; anything else ends the +/// scan with `None`. +fn connect_max_packet_size(body: &[u8]) -> Option { + let mut i = 10; + let len = take_varint(body, &mut i)?; + let end = i + len; + while i < end { + let id = *body.get(i)?; + i += 1; + match id { + 0x27 => { + let b = body.get(i..i + 4)?; + return Some(u32::from_be_bytes([b[0], b[1], b[2], b[3]])); + } + 0x11 => i += 4, // session expiry interval + 0x21 | 0x22 => i += 2, // receive maximum, topic alias maximum + 0x17 | 0x19 => i += 1, // request problem / response information + _ => return None, + } + } + None +} + /// The identity a CONNECT carries: client id, then the credentials its flags /// advertise. Nothing here sets a will, so the payload fields are contiguous. fn connect_identity(body: &[u8], v5: bool) -> Option<(String, Option<(String, String)>)> { @@ -261,6 +289,9 @@ where { let mut buf = Vec::new(); let mut v5 = true; + // The client's Maximum Packet Size: a broker must not send it anything + // larger, so a push over it is withheld. + let mut client_max: Option = None; loop { let Some((first, body)) = read_packet(socket, &mut buf).await else { @@ -280,6 +311,14 @@ where if let Some(keep_alive) = connect_keep_alive(&body) { seen.keep_alives.push(keep_alive); } + client_max = if v5 { + connect_max_packet_size(&body) + } else { + None + }; + if let Some(max) = client_max { + seen.max_packet_sizes.push(max); + } } let ack: &[u8] = if v5 { &[0x20, 0x03, 0x00, 0x00, 0x00] @@ -299,11 +338,10 @@ where return; } if let Some((topic, payload)) = after.push { - if socket - .write_all(&publish(topic, payload, v5)) - .await - .is_err() - { + let packet = publish(topic, payload, v5); + if client_max.is_some_and(|max| packet.len() > max as usize) { + seen.lock().unwrap().withheld += 1; + } else if socket.write_all(&packet).await.is_err() { return; } } diff --git a/aimdb-mqtt-connector/tests/tokio_broker.rs b/aimdb-mqtt-connector/tests/tokio_broker.rs index 0a452a7d..3acbfccb 100644 --- a/aimdb-mqtt-connector/tests/tokio_broker.rs +++ b/aimdb-mqtt-connector/tests/tokio_broker.rs @@ -277,3 +277,75 @@ async fn two_connectors_in_one_process_keep_their_own_client_ids() { ids.sort(); assert_eq!(ids, vec!["first-node", "second-node"]); } + +/// Connects a client subscribed to `sensors/temperature` against a broker that +/// pushes `payload_len` bytes after every SUBACK, as it would a retained +/// message, and returns what the broker saw after `wait` plus the length the +/// record last received. +async fn with_retained_push(payload_len: usize, wait: Duration) -> (Seen, Option) { + use aimdb_core::buffer::BufferCfg; + use aimdb_core::AimDbBuilder; + use aimdb_mqtt_connector::MqttConnector; + use aimdb_tokio_adapter::net::TokioNet; + use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let seen = Arc::new(Mutex::new(Seen::default())); + + let connector = MqttConnector::new(format!("mqtt://127.0.0.1:{port}")) + .transport(TokioNet::tcp()) + .with_client_id("max-packet-size"); + 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"); + let mut reader = db.subscribe::("temperature").expect("subscribe"); + + let payload = vec![b'x'; payload_len]; + let broker = fake_broker( + listener, + seen.clone(), + 0, + Some(("sensors/temperature", payload.as_slice())), + ); + let mut received = None; + let observe = async { + loop { + received = Some(reader.recv().await.expect("record open")); + } + }; + tokio::select! { + _ = runner.run() => panic!("the session loop returned"), + _ = broker => panic!("the broker returned"), + _ = observe => unreachable!(), + _ = tokio::time::sleep(wait) => {} + } + let seen = std::mem::take(&mut *seen.lock().unwrap()); + (seen, received) +} + +/// The CONNECT advertises the largest packet the session always receives, so +/// a broker withholds a larger retained message instead of sending one that +/// would end every session it is replayed into. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_retained_message_over_the_maximum_packet_size_is_withheld() { + let wait = Duration::from_secs(3); + + let (seen, received) = with_retained_push(3000, wait).await; + assert_eq!(seen.max_packet_sizes, [3328]); + assert_eq!((seen.connects, seen.withheld), (1, 0)); + assert_eq!(received, Some(3000), "a message within the limit arrives"); + + let (seen, received) = with_retained_push(4000, wait).await; + assert_eq!(seen.max_packet_sizes, [3328]); + assert_eq!(seen.connects, 1, "no reconnect loop"); + assert_eq!(seen.withheld, 1); + assert_eq!(received, None); +} From 2792eb70c4a61446fe796db51cd2a10281e051f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Schn=C3=B6rch?= Date: Mon, 5 Oct 2026 19:44:41 +0000 Subject: [PATCH 2/2] refactor(core)!: remove Source, pump_sink, Connector and TopicProvider (design 054) Delete everything `InboundDispatch` and `OutboundRoutes` replaced, so any caller still on the old path is a compile error. Removed from aimdb-core: `Source`, `pump_source`, `pump_sink` (`session/pump.rs`), the `Connector` trait (`ConnectorConfig` and `PublishError` stay), `SerializedReader`, `SerializedSource`, `SerializedValue`, `SerializedValueInto`, `SerializedPayload`, `RecvSerializedFuture`, `RecvSerializedIntoFuture`, `SourceFactoryFn`, `TopicProvider`, `with_topic_provider`, `OutboundRoute` and `AimDb::collect_outbound_routes`. `Router` and `AimDb::inbound_router` are crate-private. `ConnectorLink` carries only the route factory. Embassy adapter: `EmbassySinkRaw`, `EmbassySink`, `EmbassySourceRaw` and `EmbassySource` removed; `connectors.rs` keeps `into_box_future` and `NetStack`. KNX: `embassy-sync` and the optional `critical-section` dependency are dropped; `critical-section-std-impl` stays as a deprecated no-op and is removed from the README, the Tokio demo and codegen output. Tests: the fused-reader cases run against `OutboundRoutes`, the router tests are folded into the `InboundDispatch` ones, and `link_codec.rs`, MQTT `link_ext_tests.rs` and the Tokio adapter tests use the new objects. Bench: the old rows are dropped and the baseline replaced; every row matches the design's targets. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 2 - aimdb-bench/README.md | 11 +- aimdb-bench/benches/b0_alloc_connector.rs | 288 +---- aimdb-bench/benches/b0_alloc_linkable.rs | 5 +- aimdb-bench/benches/b1_outbound_wakeup.rs | 2 +- .../data/baselines/b0_alloc_connector.json | 88 +- aimdb-codegen/src/rust.rs | 14 +- aimdb-core/Cargo.toml | 8 +- aimdb-core/src/builder.rs | 72 +- aimdb-core/src/codec.rs | 2 +- aimdb-core/src/connector.rs | 297 +----- aimdb-core/src/inbound_dispatch.rs | 3 +- aimdb-core/src/lib.rs | 19 +- aimdb-core/src/outbound/routes.rs | 34 +- aimdb-core/src/router.rs | 4 +- aimdb-core/src/session/mod.rs | 30 +- aimdb-core/src/session/pump.rs | 164 --- aimdb-core/src/topic_pattern.rs | 2 +- aimdb-core/src/transport.rs | 111 +- aimdb-core/src/typed_api.rs | 984 +++++------------- aimdb-data-contracts/src/link_codec.rs | 133 +-- aimdb-embassy-adapter/src/connectors.rs | 97 +- aimdb-embassy-adapter/src/lib.rs | 4 +- aimdb-embassy-adapter/src/send_wrapper.rs | 9 +- aimdb-knx-connector/Cargo.toml | 37 +- aimdb-knx-connector/README.md | 8 +- aimdb-knx-connector/src/lib.rs | 6 +- .../tests/topic_writer_tests.rs | 34 - aimdb-mqtt-connector/Cargo.toml | 5 +- aimdb-mqtt-connector/tests/link_ext_tests.rs | 34 +- .../tests/topic_writer_tests.rs | 2 +- aimdb-tokio-adapter/tests/outbound_routes.rs | 58 +- .../tests/decouple_record_keys_topics.rs | 6 +- aimdb-websocket-connector/tests/e2e.rs | 4 +- examples/tokio-knx-connector-demo/Cargo.toml | 3 - 35 files changed, 472 insertions(+), 2108 deletions(-) delete mode 100644 aimdb-core/src/session/pump.rs diff --git a/Cargo.lock b/Cargo.lock index 46f2d412..b775a6e2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -207,9 +207,7 @@ dependencies = [ "aimdb-core", "aimdb-knx-pico", "aimdb-tokio-adapter", - "critical-section", "embassy-futures 0.1.2 (registry+https://github.com/rust-lang/crates.io-index)", - "embassy-sync", "heapless 0.8.0", "tokio", "tokio-test", diff --git a/aimdb-bench/README.md b/aimdb-bench/README.md index 7b79866d..5bdc3a09 100644 --- a/aimdb-bench/README.md +++ b/aimdb-bench/README.md @@ -26,10 +26,9 @@ Plus two informational benches that exercise the full runner-driven pipeline. `b1_b2_remote_json` (host). These compare issue #196's direct JSON bytes with the compatibility `serde_json::Value` tree through the real typed record, buffer, `Payload` and AimX envelope. Socket I/O and scheduling are excluded. -- **connector boundary** — `b0_alloc_connector` (host). Baseline for design - 054: `Router::route`, `pump_source` with a minimal `Source`, and the - per-message `recv_into` + `Connector::publish` calls of `pump_sink`, on a - no-op connector. +- **connector boundary** — `b0_alloc_connector` (host). The connector + interfaces from design 054: `InboundDispatch::dispatch` and + `OutboundRoutes::next`, on a no-op connector. - **Embassy** — `b0_alloc_embassy`, `b1_b2_embassy` (host). These drive the real [`EmbassyBuffer`] backend via `futures::executor::block_on` over embassy-sync's poll methods — no @@ -127,11 +126,11 @@ The committed baseline lives in `data/baselines/b0_alloc_tokio.json`. When a cha `b0_alloc_embassy` mirrors this against the Embassy buffer backend and writes `data/baselines/b0_alloc_embassy.json` — also **0 allocs/msg** across all three profiles, confirming the Embassy `poll_recv` path is allocation-free on the host. The on-target B3 bench (`examples/embassy-bench-stm32h5`) re-checks the same 0-alloc claim against the real embedded allocator. `b0_alloc_linkable` warms up for 1,000 iterations, then measures 10,000 generated-shape postcard `Linkable::encode_into` calls into one stack buffer. -The required result is **0 allocation calls and 0 allocated bytes**. It isolates the codec seam: `SerializedReader` still returns a boxed future, dynamic topics may allocate and connector implementations may copy payload ownership after the core pump lends them the scratch slice. +The required result is **0 allocation calls and 0 allocated bytes**. It isolates the codec seam: connector implementations may still copy the payload after `OutboundRoutes` lends them the scratch slice. `b0_alloc_remote_json` warms the production in-memory `record.get` and subscription-event paths, then compares 5000 tree/direct operations. Its gate is relative: direct JSON must reduce both allocation calls and allocated bytes. It does not require zero allocations because the owned JSON `Vec`, `Arc<[u8]>` payload and AimX envelope serialization still own storage. -`b0_alloc_connector` measures what the connector interfaces cost per message, with no transport. It asserts today's values exactly (0 for `route`, 2 for `pump_source`, 2–3 outbound), so a regression *or* an improvement fails it until `EXPECTED` in the bench and `data/baselines/b0_alloc_connector.json` are updated together. The inbound `pump_source` row is the difference of two runs, so pump setup cancels out. See design 054 for where each allocation comes from. +`b0_alloc_connector` measures what the connector interfaces cost per message, with no transport. It asserts its values exactly (0 everywhere except a new key, 1, and the owned serializer, 1), so a regression *or* an improvement fails it until `EXPECTED` in the bench and `data/baselines/b0_alloc_connector.json` are updated together. `make bench-gate` runs it. See design 054 for where each allocation comes from. > **Embassy eager registration (design 039 F8/F9).** An Embassy `SpmcRing` reader registers its embassy `Subscriber` eagerly, at `subscribe()` time — matching Tokio's `broadcast` — so no separate priming step is needed before the first `push`. diff --git a/aimdb-bench/benches/b0_alloc_connector.rs b/aimdb-bench/benches/b0_alloc_connector.rs index b79df06a..92888351 100644 --- a/aimdb-bench/benches/b0_alloc_connector.rs +++ b/aimdb-bench/benches/b0_alloc_connector.rs @@ -1,26 +1,19 @@ //! B0-Connector — per-message allocations at the connector boundary. //! -//! Baseline for design 054. Measures what AimDB's connector interfaces cost -//! per message, independent of any real transport: +//! Measures what AimDB's connector interfaces cost per message, independent of +//! any real transport: //! -//! - **Inbound:** `Router::route` alone — for an exact topic, a pattern +//! - **Inbound:** `InboundDispatch::dispatch` for an exact topic, a pattern //! (`{device}`, MQTT grammar), and a keyed pattern with a known and a new -//! key — and the real `pump_source` driven by -//! the smallest possible `Source` (it clones a pre-built topic `String` and -//! payload `Arc` — the least any `Source` can do, since the trait returns -//! owned values). -//! - **Outbound:** the per-message calls `pump_sink` makes — -//! `SerializedReader::recv_into` followed by `Connector::publish` on a no-op -//! connector — for the scratch and owned serializers, with a static and a -//! dynamic (`TopicProvider`) topic. -//! - **`InboundDispatch` and `OutboundRoutes`:** the same inbound cases through -//! `InboundDispatch::dispatch`, and `OutboundRoutes::next` with a static -//! topic, a written topic and the owned serializer, plus eight routes that -//! are all ready (round-robin) and one pull that parks before every value -//! (the waker path). +//! key. +//! - **Outbound:** `OutboundRoutes::next` with a static topic, a written topic +//! and the owned serializer, plus eight routes that are all ready +//! (round-robin) and one pull that parks before every value (the waker +//! path). //! -//! Buffers, ingest and routing allocate nothing (design 037, and the `route` -//! row here); every non-zero row is a cost of the connector interface. The +//! Buffers, ingest and routing allocate nothing (design 037, and the +//! `inbound_dispatch` row here); every non-zero row is a cost of the connector +//! interface or of the serializer the link chose. The //! expected values are asserted, so a regression *or* an improvement fails the //! bench until `EXPECTED` and the committed baseline are updated together. //! @@ -36,15 +29,9 @@ use std::sync::Arc; use aimdb_bench::alloc::{reset, snapshot}; use aimdb_bench::reports::{write_reports, AllocReport}; use aimdb_core::buffer::BufferCfg; -use aimdb_core::connector::{ - ConnectorBuilder, SerializeError, SerializedPayload, SerializedReader, SerializedValueInto, - TopicProvider, -}; -use aimdb_core::session::{pump_source, Payload, Source}; -use aimdb_core::transport::{Connector, ConnectorConfig, PublishError}; +use aimdb_core::connector::{ConnectorBuilder, SerializeError}; use aimdb_core::{ - AimDb, AimDbBuilder, BoxFut, DbResult, ExactGrammar, InboundDispatch, OutboundRoutes, - RuntimeContext, StringKey, + AimDb, AimDbBuilder, DbResult, ExactGrammar, InboundDispatch, OutboundRoutes, StringKey, }; use aimdb_mqtt_connector::MqttGrammar; use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; @@ -60,17 +47,9 @@ const SCRATCH_CAPACITY: usize = 64; /// Room for a new key on every warm-up and measured message. const KEY_CAPACITY: u16 = 4096; -/// Allocations per message on `main` when this bench was added. Update -/// together with `data/baselines/b0_alloc_connector.json`. +/// Allocations per message. Update together with +/// `data/baselines/b0_alloc_connector.json`. const EXPECTED: &[(&str, u64)] = &[ - ("inbound_route", 0), - ("inbound_route_pattern", 0), - ("inbound_route_keyed_known", 0), - ("inbound_route_keyed_new", 1), - ("inbound_pump_source_minimal", 2), - ("outbound_scratch_static_topic", 2), - ("outbound_scratch_dynamic_topic", 3), - ("outbound_owned_static_topic", 3), ("inbound_dispatch", 0), ("inbound_dispatch_pattern", 0), ("inbound_dispatch_keyed_known", 0), @@ -131,29 +110,6 @@ impl ConnectorBuilder for NoopConnectorBuilder { } } -/// Returns a ready future, boxed as the `Connector` trait requires. -struct NoopSink; - -impl Connector for NoopSink { - fn publish( - &self, - destination: &str, - _config: &ConnectorConfig, - payload: &[u8], - ) -> Pin> + Send + '_>> { - black_box((destination, payload)); - Box::pin(async { Ok(()) }) - } -} - -struct IdTopic; - -impl TopicProvider for IdTopic { - fn topic(&self, value: &Reading) -> Option { - Some(format!("out/{}", value.id)) - } -} - async fn build_db(configure: impl FnOnce(&mut AimDbBuilder)) -> AimDb { let runtime = Arc::new(TokioAdapter::new().expect("tokio adapter")); let mut builder = AimDbBuilder::new() @@ -217,24 +173,6 @@ async fn pattern_db(keyed: bool) -> AimDb { .await } -/// Routes `warmup` then `measured`, counting only the second. -async fn measure_pattern_route(keyed: bool, warmup: &[String], measured: &[String]) -> (u64, u64) { - let db = pattern_db(keyed).await; - let ctx = db.runtime_ctx(); - let router = db.inbound_router(SCHEME, &MqttGrammar).unwrap(); - let payload = [1u8; 8]; - for topic in warmup { - router.route(topic, &payload, &ctx).unwrap(); - } - reset(); - for topic in measured { - router - .route(black_box(topic), black_box(&payload), &ctx) - .unwrap(); - } - snapshot() -} - /// `n` copies of one topic, or `n` topics each naming a new device. fn pattern_topics(n: usize, first_device: usize, distinct: bool) -> Vec { (0..n) @@ -245,24 +183,7 @@ fn pattern_topics(n: usize, first_device: usize, distinct: bool) -> Vec .collect() } -async fn measure_route() -> (u64, u64) { - let db = inbound_db().await; - let ctx = db.runtime_ctx(); - let router = db.inbound_router(SCHEME, &ExactGrammar).unwrap(); - let payload = [1u8; 8]; - for _ in 0..WARMUP_ITERS { - router.route("in/target", &payload, &ctx).unwrap(); - } - reset(); - for _ in 0..MEASURE_ITERS { - router - .route("in/target", black_box(&payload), &ctx) - .unwrap(); - } - snapshot() -} - -/// `measure_route` through `InboundDispatch`. +/// One exact topic among the decoys. async fn measure_dispatch() -> (u64, u64) { let db = inbound_db().await; let inbound = InboundDispatch::new(&db, SCHEME, &ExactGrammar).unwrap(); @@ -277,7 +198,7 @@ async fn measure_dispatch() -> (u64, u64) { snapshot() } -/// `measure_pattern_route` through `InboundDispatch`. +/// Dispatches `warmup` then `measured`, counting only the second. async fn measure_pattern_dispatch( keyed: bool, warmup: &[String], @@ -296,138 +217,14 @@ async fn measure_pattern_dispatch( snapshot() } -/// Yields `remaining` copies of one message, then ends. -struct MinimalSource { - topic: String, - payload: Payload, - remaining: usize, -} - -impl Source for MinimalSource { - fn next(&mut self) -> BoxFut<'_, Option<(String, Payload)>> { - Box::pin(async move { - if self.remaining == 0 { - return None; - } - self.remaining -= 1; - Some((self.topic.clone(), self.payload.clone())) - }) - } -} - -/// Allocations of one complete `pump_source` run over `messages` messages, -/// including its one-off setup. -async fn pump_run(db: &AimDb, messages: usize) -> (u64, u64) { - let source = MinimalSource { - topic: "in/target".to_string(), - payload: Arc::from(&[1u8; 8][..]), - remaining: messages, - }; - reset(); - let router = db.inbound_router(SCHEME, &ExactGrammar).unwrap(); - for fut in pump_source(db, router, source) { - fut.await; - } - snapshot() -} - -/// Per-message cost as the difference of two runs, so pump setup (router -/// build, start-up logging) cancels out. -async fn measure_pump_source() -> (u64, u64) { - let db = inbound_db().await; - pump_run(&db, WARMUP_ITERS).await; - let short = pump_run(&db, WARMUP_ITERS).await; - let long = pump_run(&db, WARMUP_ITERS + MEASURE_ITERS).await; - (long.0 - short.0, long.1 - short.1) -} - // --- Outbound --------------------------------------------------------------- -/// The per-route state `pump_sink` keeps, and one iteration of its loop. -struct OutboundIo { - reader: Box, - scratch: Vec, - default_topic: String, - config: ConnectorConfig, -} - -impl OutboundIo { - async fn publish_one(&mut self, ctx: &RuntimeContext) { - let SerializedValueInto { dest, payload } = self - .reader - .recv_into(ctx, &mut self.scratch) - .await - .expect("recv_into"); - let dest = dest.as_deref().unwrap_or(&self.default_topic); - let bytes: &[u8] = match &payload { - SerializedPayload::Scratch { len } => &self.scratch[..*len], - SerializedPayload::Owned(v) => v, - }; - NoopSink - .publish(dest, &self.config, bytes) - .await - .expect("publish"); - } -} - #[derive(Clone, Copy)] enum Serializer { Scratch, Owned, } -async fn measure_outbound(serializer: Serializer, dynamic_topic: bool) -> (u64, u64) { - let db = build_db(|b| { - b.configure::("out.record", |reg| { - reg.buffer(BufferCfg::SpmcRing { capacity: 64 }); - let mut link = reg - .link_to("bench://out/default") - .with_serializer(|_ctx, r: &Reading| Ok(encode(r).to_vec())); - if let Serializer::Scratch = serializer { - link = link.with_serializer_into(SCRATCH_CAPACITY, |_ctx, r: &Reading, buf| { - let bytes = encode(r); - let dst = buf - .get_mut(..bytes.len()) - .ok_or(SerializeError::BufferTooSmall)?; - dst.copy_from_slice(&bytes); - Ok(bytes.len()) - }); - } - if dynamic_topic { - link = link.with_topic_provider(IdTopic); - } - link.finish(); - }); - }) - .await; - - let ctx = db.runtime_ctx(); - let route = db - .collect_outbound_routes(SCHEME) - .pop() - .expect("one outbound route"); - let mut io = OutboundIo { - reader: route.source.subscribe(), - scratch: vec![0u8; route.source.serializer_scratch_capacity().unwrap_or(0)], - default_topic: route.topic.clone(), - config: ConnectorConfig::from_query(&route.config), - }; - let producer = db.producer::("out.record").expect("producer"); - - for i in 0..WARMUP_ITERS { - producer.produce(reading(i)); - io.publish_one(&ctx).await; - } - reset(); - for i in 0..MEASURE_ITERS { - producer.produce(reading(i)); - io.publish_one(&ctx).await; - } - snapshot() -} - -// --- OutboundRoutes --------------------------------------------------------- - #[derive(Clone, Copy)] enum Topic { Static, @@ -587,57 +384,6 @@ fn main() { let measured: Vec<(&str, &str, (u64, u64))> = runtime.block_on(async { vec![ - ("inbound_route", "SpmcRing", measure_route().await), - ( - "inbound_route_pattern", - "SpmcRing", - measure_pattern_route( - false, - &pattern_topics(WARMUP_ITERS, 0, false), - &pattern_topics(MEASURE_ITERS, 0, false), - ) - .await, - ), - ( - "inbound_route_keyed_known", - "SpmcRing", - measure_pattern_route( - true, - &pattern_topics(WARMUP_ITERS, 0, false), - &pattern_topics(MEASURE_ITERS, 0, false), - ) - .await, - ), - ( - "inbound_route_keyed_new", - "SpmcRing", - measure_pattern_route( - true, - &pattern_topics(WARMUP_ITERS, 0, true), - &pattern_topics(MEASURE_ITERS, WARMUP_ITERS, true), - ) - .await, - ), - ( - "inbound_pump_source_minimal", - "SpmcRing", - measure_pump_source().await, - ), - ( - "outbound_scratch_static_topic", - "SpmcRing", - measure_outbound(Serializer::Scratch, false).await, - ), - ( - "outbound_scratch_dynamic_topic", - "SpmcRing", - measure_outbound(Serializer::Scratch, true).await, - ), - ( - "outbound_owned_static_topic", - "SpmcRing", - measure_outbound(Serializer::Owned, false).await, - ), ("inbound_dispatch", "SpmcRing", measure_dispatch().await), ( "inbound_dispatch_pattern", diff --git a/aimdb-bench/benches/b0_alloc_linkable.rs b/aimdb-bench/benches/b0_alloc_linkable.rs index 785ad28a..6a80c4dc 100644 --- a/aimdb-bench/benches/b0_alloc_linkable.rs +++ b/aimdb-bench/benches/b0_alloc_linkable.rs @@ -1,8 +1,7 @@ //! B0-Linkable — allocation gate for direct and per-link Postcard encoding. //! -//! This deliberately measures the codec seam, not a complete connector. The -//! core pump still uses a boxed `SerializedReader` future and connector adapters -//! may copy payload ownership; issue #177 only claims that a generated-shape +//! This deliberately measures the codec seam, not a complete connector, which +//! may still copy the payload; issue #177 only claims that a generated-shape //! Postcard codec writes into caller-owned storage with zero heap allocations. use std::hint::black_box; diff --git a/aimdb-bench/benches/b1_outbound_wakeup.rs b/aimdb-bench/benches/b1_outbound_wakeup.rs index bc769a36..84498f00 100644 --- a/aimdb-bench/benches/b1_outbound_wakeup.rs +++ b/aimdb-bench/benches/b1_outbound_wakeup.rs @@ -8,7 +8,7 @@ //! one wake-up of the transport. Two columns: //! //! - **task per route:** one task per route awaiting `Reader::recv`, the -//! shape of the per-route pumps. +//! shape of the per-route publishers `OutboundRoutes` replaced. //! - **OutboundRoutes:** one task pulling with `OutboundRoutes::next`, which //! polls only the routes that woke. //! diff --git a/aimdb-bench/data/baselines/b0_alloc_connector.json b/aimdb-bench/data/baselines/b0_alloc_connector.json index ae351427..21cb1034 100644 --- a/aimdb-bench/data/baselines/b0_alloc_connector.json +++ b/aimdb-bench/data/baselines/b0_alloc_connector.json @@ -1,76 +1,4 @@ [ - { - "profile": "inbound_route", - "buffer_type": "SpmcRing", - "total_allocs": 0, - "total_bytes": 0, - "batch_size": 2000, - "allocs_per_msg": 0.0, - "bytes_per_msg": 0.0 - }, - { - "profile": "inbound_route_pattern", - "buffer_type": "SpmcRing", - "total_allocs": 0, - "total_bytes": 0, - "batch_size": 2000, - "allocs_per_msg": 0.0, - "bytes_per_msg": 0.0 - }, - { - "profile": "inbound_route_keyed_known", - "buffer_type": "SpmcRing", - "total_allocs": 0, - "total_bytes": 0, - "batch_size": 2000, - "allocs_per_msg": 0.0, - "bytes_per_msg": 0.0 - }, - { - "profile": "inbound_route_keyed_new", - "buffer_type": "SpmcRing", - "total_allocs": 2010, - "total_bytes": 373456, - "batch_size": 2000, - "allocs_per_msg": 1.005, - "bytes_per_msg": 186.728 - }, - { - "profile": "inbound_pump_source_minimal", - "buffer_type": "SpmcRing", - "total_allocs": 4000, - "total_bytes": 50000, - "batch_size": 2000, - "allocs_per_msg": 2.0, - "bytes_per_msg": 25.0 - }, - { - "profile": "outbound_scratch_static_topic", - "buffer_type": "SpmcRing", - "total_allocs": 4001, - "total_bytes": 146064, - "batch_size": 2000, - "allocs_per_msg": 2.0005, - "bytes_per_msg": 73.032 - }, - { - "profile": "outbound_scratch_dynamic_topic", - "buffer_type": "SpmcRing", - "total_allocs": 6000, - "total_bytes": 162000, - "batch_size": 2000, - "allocs_per_msg": 3.0, - "bytes_per_msg": 81.0 - }, - { - "profile": "outbound_owned_static_topic", - "buffer_type": "SpmcRing", - "total_allocs": 6000, - "total_bytes": 162000, - "batch_size": 2000, - "allocs_per_msg": 3.0, - "bytes_per_msg": 81.0 - }, { "profile": "inbound_dispatch", "buffer_type": "SpmcRing", @@ -110,11 +38,11 @@ { "profile": "outbound_next_static_topic", "buffer_type": "SpmcRing", - "total_allocs": 0, - "total_bytes": 0, + "total_allocs": 1, + "total_bytes": 64, "batch_size": 2000, - "allocs_per_msg": 0.0, - "bytes_per_msg": 0.0 + "allocs_per_msg": 0.0005, + "bytes_per_msg": 0.032 }, { "profile": "outbound_next_written_topic", @@ -137,11 +65,11 @@ { "profile": "outbound_next_round_robin", "buffer_type": "SpmcRing", - "total_allocs": 0, - "total_bytes": 0, + "total_allocs": 1, + "total_bytes": 128, "batch_size": 2000, - "allocs_per_msg": 0.0, - "bytes_per_msg": 0.0 + "allocs_per_msg": 0.0005, + "bytes_per_msg": 0.064 }, { "profile": "outbound_next_parked", diff --git a/aimdb-codegen/src/rust.rs b/aimdb-codegen/src/rust.rs index 6b9940ec..9c6b1683 100644 --- a/aimdb-codegen/src/rust.rs +++ b/aimdb-codegen/src/rust.rs @@ -497,11 +497,8 @@ pub fn generate_binary_cargo_toml(state: &ArchitectureState, binary_name: &str) ); } if has_knx { - optional_connector_deps.push_str( - "# critical-section-std-impl: the KNX channels need an impl, and only \ -the binary may pick one.\n\ -aimdb-knx-connector = { version = \"0.5\", features = [\"std\", \"critical-section-std-impl\"] }\n", - ); + optional_connector_deps + .push_str("aimdb-knx-connector = { version = \"0.5\", features = [\"std\"] }\n"); } if has_ws { optional_connector_deps.push_str( @@ -1334,11 +1331,8 @@ pub fn generate_hub_cargo_toml(state: &ArchitectureState) -> String { ); } if has_knx { - connector_deps.push_str( - "# critical-section-std-impl: the KNX channels need an impl, and only \ -the binary may pick one.\n\ -aimdb-knx-connector = { version = \"0.5\", features = [\"std\", \"critical-section-std-impl\"] }\n", - ); + connector_deps + .push_str("aimdb-knx-connector = { version = \"0.5\", features = [\"std\"] }\n"); } if has_ws { connector_deps.push_str( diff --git a/aimdb-core/Cargo.toml b/aimdb-core/Cargo.toml index ff2da3fd..6b2da2c6 100644 --- a/aimdb-core/Cargo.toml +++ b/aimdb-core/Cargo.toml @@ -42,10 +42,10 @@ alloc = ["serde"] # Enable heap in no_std remote = ["alloc", "serde_json"] # The connector-session substrate (`crate::session`): the dyn-safe trait set -# (Connection/Listener/Dialer, Dispatch/EnvelopeCodec, Source + shared types) and -# the runtime-neutral engines built on it — the reactive server (`serve`/ -# `run_session`), the proactive client (`run_client`/`pump_client`), the -# `pump_sink`/`pump_source` data plane, and the generic session connectors. Engine +# (Connection/Listener/Dialer, Dispatch/EnvelopeCodec + shared types), the byte +# and datagram transports, and the runtime-neutral engines built on them — the +# reactive server (`serve`/`run_session`), the proactive client +# (`run_client`/`pump_client`), and the generic session connectors. Engine # logic included; all compiles on `no_std + alloc`. The AimX protocol port # (`session::aimx`) additionally needs `remote-access`. connector-session = ["alloc"] diff --git a/aimdb-core/src/builder.rs b/aimdb-core/src/builder.rs index 89c83283..dcd2d7a7 100644 --- a/aimdb-core/src/builder.rs +++ b/aimdb-core/src/builder.rs @@ -34,19 +34,6 @@ use crate::typed_api::RecordRegistrar; use crate::typed_record::{AnyRecord, AnyRecordExt, RecordFutureCollector, TypedRecord}; use crate::{DbError, DbResult}; -/// One outbound route returned by [`AimDb::collect_outbound_routes`] -pub struct OutboundRoute { - /// Default topic/destination from the URL path; used when the source - /// yields no per-value destination. - pub topic: String, - /// Fused wire-level source: its readers yield destination + serialized - /// payload directly (subscribe → recv → resolve topic → serialize, all - /// typed inside — no `Box` per message). - pub source: Box, - /// Configuration options from the URL query - pub config: Vec<(String, String)>, -} - /// One registered record: its key, concrete type, and type-erased storage. struct RecordEntry { key: StringKey, @@ -1182,14 +1169,13 @@ impl AimDb { /// The inbound router for `scheme`: every link compiled against the /// connector's `grammar`, keyed links sharing their record's key table. /// - /// A connector subscribes [`Router::subscriptions`](crate::Router::subscriptions) - /// and routes with this same router. Rejects every link the grammar or - /// its key cannot compile, naming the record and the resolved topic. - pub fn inbound_router( + /// Rejects every link the grammar or its key cannot compile, naming the + /// record and the resolved topic. + pub(crate) fn inbound_router( &self, scheme: &str, grammar: &'static dyn crate::TopicGrammar, - ) -> DbResult { + ) -> DbResult { let mut routes = Vec::new(); let mut errors = Vec::new(); @@ -1212,7 +1198,7 @@ impl AimDb { if !errors.is_empty() { return Err(DbError::InvalidConfiguration { errors }); } - Ok(crate::Router::new(grammar, routes)) + Ok(crate::router::Router::new(grammar, routes)) } fn inbound_route( @@ -1256,11 +1242,11 @@ impl AimDb { .name(key) } - /// Every outbound link of `scheme`, with its record's index and key. + /// Every outbound link of `scheme`, with its record's index. pub(crate) fn outbound_links<'a>( &'a self, scheme: &'a str, - ) -> impl Iterator + 'a { + ) -> impl Iterator + 'a { self.inner .storages .iter() @@ -1278,49 +1264,7 @@ impl AimDb { .outbound_connectors() .iter() .filter(move |link| link.url.scheme() == scheme) - .map(move |link| (i, entry.key.as_str(), link)) + .map(move |link| (i, link)) }) } - - /// Collects outbound routes for a specific protocol scheme - /// - /// Mirrors [`inbound_router`](Self::inbound_router). Iterates all records, - /// filters their outbound_connectors by scheme, and returns - /// [`OutboundRoute`]s carrying fused serialized sources (subscribe → - /// recv → resolve topic → serialize, all typed inside — no - /// `Box` per message). - /// - /// This method is called by connectors during their `build()` phase to - /// collect all configured outbound routes and spawn publisher tasks - /// (usually via `pump_sink`). - /// - /// # Arguments - /// * `scheme` - URL scheme to filter by (e.g., "mqtt", "kafka") - pub fn collect_outbound_routes(&self, scheme: &str) -> Vec { - let routes: Vec = self - .outbound_links(scheme) - .map(|(i, _, link)| { - // config must carry the record index - let mut config = link.config.clone(); - config.push(("record_index".to_string(), i.to_string())); - - // Create the fused source using the stored factory - OutboundRoute { - topic: link.url.resource_id().to_string(), - source: link.create_source(self), - config, - } - }) - .collect(); - - if !routes.is_empty() { - log_debug!( - "Collected {} outbound routes for scheme '{}'", - routes.len(), - scheme - ); - } - - routes - } } diff --git a/aimdb-core/src/codec.rs b/aimdb-core/src/codec.rs index 7ec2839c..3ba945d8 100644 --- a/aimdb-core/src/codec.rs +++ b/aimdb-core/src/codec.rs @@ -21,7 +21,7 @@ //! zero-sized [`SerdeJsonCodec`] implementation. A record stores //! `Option>>`; the AimX read/write/subscribe paths and //! `RecordValue::as_json` route through it. This mirrors the connector -//! layer's fused `SerializedSource` / `IngestFn` callbacks. +//! layer's typed route and `IngestFn` callbacks. use alloc::vec::Vec; use serde::{de::DeserializeOwned, Serialize}; diff --git a/aimdb-core/src/connector.rs b/aimdb-core/src/connector.rs index 4202a514..ba0ffecf 100644 --- a/aimdb-core/src/connector.rs +++ b/aimdb-core/src/connector.rs @@ -82,171 +82,6 @@ impl std::fmt::Display for SerializeError { #[cfg(feature = "std")] impl std::error::Error for SerializeError {} -/// One serialized record update, produced by a fused [`SerializedReader`] -/// -/// Carries the wire payload plus the destination resolved by the link's -/// [`TopicProvider`] while the typed value was still in hand — the last -/// erasure crossing the old `topic_any(&dyn Any)` path required. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct SerializedValue { - /// Dynamic destination resolved by the link's `TopicProvider`; - /// `None` means use the route's default topic (from the URL). - pub dest: Option, - /// Wire payload from the link's serializer. - /// - /// `Vec` requires heap allocation; works on `std` and - /// `no_std + alloc` (not bare-metal without an allocator). - pub payload: Vec, -} - -/// Location of a payload produced by [`SerializedReader::recv_into`]. -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum SerializedPayload { - /// The initialized prefix of the caller-provided scratch buffer. - /// - /// A custom reader returning this variant must write the prefix during the - /// same `recv_into` call and its source must advertise a sufficient - /// [`SerializedSource::serializer_scratch_capacity`]. The pump validates - /// the length before publishing. - Scratch { - /// Number of initialized payload bytes in the scratch buffer. - len: usize, - }, - /// Compatibility fallback from the existing `Vec` serializer. - Owned(Vec), -} - -/// One serialized record produced into caller-owned scratch storage or an owned fallback. -/// -/// No reference escapes the async reader call. The pump validates `len`, then -/// borrows its own scratch buffer only for the subsequent `publish().await`. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct SerializedValueInto { - /// Dynamic destination resolved while the typed record was available. - pub dest: Option, - /// Scratch prefix metadata or an owned compatibility payload. - pub payload: SerializedPayload, -} - -/// Type alias for the future returned by [`SerializedReader::recv`] -/// -/// Manual boxed future for object safety — same pattern as the rest of this -/// module (`#[async_trait]` would drag in `std`). -pub type RecvSerializedFuture<'a> = - Pin> + Send + 'a>>; - -/// Future returned by [`SerializedReader::recv_into`]. -pub type RecvSerializedIntoFuture<'a> = - Pin> + Send + 'a>>; - -/// A subscription to one record, fused with destination resolution and -/// serialization at registration time — no `dyn Any` crosses this boundary -///. -pub trait SerializedReader: Send { - /// Yield the next successfully serialized value. - /// - /// `ctx` is threaded per call (not captured) for context-aware - /// serializers. Buffer errors propagate unchanged: - /// `DbError::BufferLagged` means values were skipped but the reader - /// recovered; any other error means the buffer is gone. Serialization - /// failures are logged and skipped inside the reader. - fn recv<'a>(&'a mut self, ctx: &'a crate::RuntimeContext) -> RecvSerializedFuture<'a>; - - /// Yield the next serialized value into caller-owned scratch storage. - /// - /// Third-party readers remain source-compatible: the default adapter calls - /// [`recv`](Self::recv) and returns its `Vec` as - /// [`SerializedPayload::Owned`]. AimDB's fused reader overrides this method - /// when an into-slice serializer was registered. - fn recv_into<'a>( - &'a mut self, - ctx: &'a crate::RuntimeContext, - _scratch: &'a mut [u8], - ) -> RecvSerializedIntoFuture<'a> { - Box::pin(async move { - let value = self.recv(ctx).await?; - Ok(SerializedValueInto { - dest: value.dest, - payload: SerializedPayload::Owned(value.payload), - }) - }) - } -} - -/// A record's outbound wire interface, built where the record type `T` is -/// known (`OutboundConnectorBuilder::finish`) and consumed by the pumps as -/// bytes. Replaces the erased `ConsumerTrait` + serializer + topic-provider -/// triple. -pub trait SerializedSource: Send + Sync { - /// Scratch capacity requested by this source's into-slice serializer. - /// - /// `None` means the source only supports the existing owned serializer. A - /// source whose reader can return [`SerializedPayload::Scratch`] must - /// return `Some(capacity)` here so the route pump provides that storage. - fn serializer_scratch_capacity(&self) -> Option { - None - } - - /// Subscribe to the record's updates. - /// - /// Synchronous and infallible — the buffer handle is pre-resolved at - /// construction. - fn subscribe(&self) -> Box; -} - -/// Type alias for source factory callback (alloc feature) -/// -/// Takes the live [`AimDb`] and returns the fused [`SerializedSource`]. -/// This allows capturing the record type T at link_to() time while storing -/// the factory in a type-erased ConnectorLink. The factory runs once at -/// route-collection time, not per message. -/// -/// Available in both `std` and `no_std + alloc` environments. -pub type SourceFactoryFn = Arc Box + Send + Sync>; - -// ============================================================================ -// TopicProvider - Dynamic topic/destination routing -// ============================================================================ - -/// Trait for dynamic topic providers (outbound only) -/// -/// Implement this trait to dynamically determine MQTT topics (or KNX group addresses) -/// based on the data being published. This enables reusable routing logic that -/// can be shared across multiple record types. -/// -/// # Type Safety -/// -/// The trait is generic over `T`, providing compile-time type safety -/// at the implementation site. The provider stays typed end-to-end: it is -/// fused into the link's [`SerializedSource`] at registration time and -/// called with `&T` while the value is in hand. -/// -/// # no_std Compatibility -/// -/// Works in both `std` and `no_std + alloc` environments. -/// -/// # Example -/// -/// ```rust -/// use aimdb_core::connector::TopicProvider; -/// # #[derive(Clone, Debug)] struct Temperature { sensor_id: u32 } -/// -/// struct SensorTopicProvider; -/// -/// impl TopicProvider for SensorTopicProvider { -/// fn topic(&self, value: &Temperature) -> Option { -/// Some(format!("sensors/temp/{}", value.sensor_id)) -/// } -/// } -/// ``` -pub trait TopicProvider: Send + Sync { - /// Determine the topic/destination for a given value - /// - /// Returns `Some(topic)` to use a dynamic topic, or `None` to fall back - /// to the static topic from the `link_to()` URL. - fn topic(&self, value: &T) -> Option; -} - // ============================================================================ // TopicWriter - Destinations written into bounded storage // ============================================================================ @@ -571,11 +406,10 @@ impl fmt::Display for ConnectorUrl { } } -/// Configuration for a connector link +/// Configuration for an outbound connector link /// -/// Stores the parsed URL, configuration, and the fused source factory until -/// the record is built. The actual client creation and handler spawning -/// happens during the build phase. +/// Stores the parsed URL, configuration, and the route factory until the +/// database is built. `OutboundRoutes` runs the factory once per link. #[derive(Clone)] pub struct ConnectorLink { /// Parsed link address (`scheme://resource`) @@ -584,21 +418,8 @@ pub struct ConnectorLink { /// Additional configuration options (protocol-specific) pub config: Vec<(String, String)>, - /// Fused source factory (alloc feature) - /// - /// Takes the live [`AimDb`] and returns the [`SerializedSource`] whose - /// readers yield destination + payload directly (subscribe → recv → - /// resolve topic → serialize, all typed inside). Captures the record - /// type T at link_to() configuration time — `finish()` validates the - /// serializer is present before registering the link, so the factory is - /// always set. - /// - /// Available in both `std` and `no_std + alloc` environments. - pub source_factory: SourceFactoryFn, - - /// Builds the link's route for `OutboundRoutes`; set by - /// `OutboundConnectorBuilder::finish`. - pub(crate) route_factory: Option, + /// Builds the link's route for `OutboundRoutes`. + pub(crate) route_factory: crate::outbound::RouteFactoryFn, } impl Debug for ConnectorLink { @@ -606,31 +427,19 @@ impl Debug for ConnectorLink { f.debug_struct("ConnectorLink") .field("url", &self.url) .field("config", &self.config) - .field("source_factory", &"") - .finish() + .finish_non_exhaustive() } } impl ConnectorLink { - /// Creates a new connector link from a link address and source factory - pub fn new(url: LinkAddress, source_factory: SourceFactoryFn) -> Self { + /// Creates a new connector link from a link address and route factory + pub(crate) fn new(url: LinkAddress, route_factory: crate::outbound::RouteFactoryFn) -> Self { Self { url, config: Vec::new(), - source_factory, - route_factory: None, + route_factory, } } - - /// Creates the fused serialized source using the stored factory. - /// - /// Runs once at route-collection time; the readers it hands out are the - /// per-message path (no `Box`). - /// - /// Available in both `std` and `no_std + alloc` environments. - pub fn create_source(&self, db: &AimDb) -> Box { - (self.source_factory)(db) - } } /// Fused inbound ingest callback: deserialize + produce in one typed closure @@ -881,8 +690,9 @@ fn parse_connector_url(url: &str) -> DbResult { /// ) -> Pin>> + Send + 'a>> { /// Box::pin(async move { /// // Wildcard rules are the connector's; `&ExactGrammar` if it has none. -/// let router = db.inbound_router(self.scheme(), &MqttGrammar)?; -/// let connector = MqttConnector::new(&self.broker_url, router).await?; +/// let inbound = InboundDispatch::new(db, self.scheme(), &MqttGrammar)?; +/// let outbound = OutboundRoutes::new(db, self.scheme())?; +/// let connector = MqttConnector::new(&self.broker_url, inbound, outbound).await?; /// Ok(connector.futures()) /// }) /// } @@ -896,8 +706,8 @@ pub trait ConnectorBuilder: Send + Sync { /// Build the connector and return its driving futures. /// /// Called during `AimDbBuilder::build()` after the database has been - /// constructed. The returned futures (infrastructure loops + per-route - /// publishers) are appended to the builder's accumulator and driven by + /// constructed. The returned futures (typically the transport task) + /// are appended to the builder's accumulator and driven by /// `AimDbRunner::run()`. /// /// # Arguments @@ -928,12 +738,12 @@ pub trait ConnectorBuilder: Send + Sync { /// Whether registering a second connector under this scheme is an error. /// /// Say `true` when [`build`](Self::build) claims every route for its - /// scheme — [`inbound_router`](crate::AimDb::inbound_router), - /// [`collect_outbound_routes`](crate::AimDb::collect_outbound_routes), - /// and `crate::session`'s `pump_source`, `pump_sink` and `pump_client` - /// (left unlinked: that module is behind `connector-session`, and this - /// trait is not) all filter by scheme alone, so two such connectors each - /// collect *all* of it: every `link_to` gets two publishers, and the routes + /// scheme — [`InboundDispatch`](crate::InboundDispatch), + /// [`OutboundRoutes`](crate::OutboundRoutes) and `crate::session`'s + /// `pump_client` (left unlinked: that module is behind + /// `connector-session`, and this trait is not) all filter by scheme + /// alone, so two such connectors each collect *all* of it: every + /// `link_to` gets two publishers, and the routes /// cannot be divided between the two endpoints because nothing in a route /// names which connector it belongs to. That misconfiguration is otherwise /// silent, and it fails as duplicated or misdirected traffic at runtime @@ -1060,73 +870,6 @@ mod tests { assert!(result.is_err()); } - // ======================================================================== - // TopicProvider Tests - // ======================================================================== - - #[allow(dead_code)] - #[derive(Debug, Clone)] - struct TestTemperature { - sensor_id: String, - celsius: f32, - } - - struct TestTopicProvider; - - impl super::TopicProvider for TestTopicProvider { - fn topic(&self, value: &TestTemperature) -> Option { - Some(format!("sensors/temp/{}", value.sensor_id)) - } - } - - #[test] - fn test_topic_provider_as_trait_object() { - // Providers are stored as Arc> — typed, no - // erasure. - let provider: Arc> = Arc::new(TestTopicProvider); - let temp = TestTemperature { - sensor_id: "kitchen-001".into(), - celsius: 22.5, - }; - - assert_eq!( - provider.topic(&temp), - Some("sensors/temp/kitchen-001".into()) - ); - } - - #[test] - fn test_topic_provider_returns_none() { - struct OptionalTopicProvider; - - impl super::TopicProvider for OptionalTopicProvider { - fn topic(&self, temp: &TestTemperature) -> Option { - if temp.sensor_id.is_empty() { - None // Fall back to default topic - } else { - Some(format!("sensors/{}", temp.sensor_id)) - } - } - } - - let provider: Arc> = - Arc::new(OptionalTopicProvider); - - // Non-empty sensor_id returns dynamic topic - let temp_with_id = TestTemperature { - sensor_id: "abc".into(), - celsius: 20.0, - }; - assert_eq!(provider.topic(&temp_with_id), Some("sensors/abc".into())); - - // Empty sensor_id returns None (fallback) - let temp_without_id = TestTemperature { - sensor_id: String::new(), - celsius: 20.0, - }; - assert_eq!(provider.topic(&temp_without_id), None); - } - // ======================================================================== // TopicBuf Tests // ======================================================================== diff --git a/aimdb-core/src/inbound_dispatch.rs b/aimdb-core/src/inbound_dispatch.rs index 5e264c99..6fa6519a 100644 --- a/aimdb-core/src/inbound_dispatch.rs +++ b/aimdb-core/src/inbound_dispatch.rs @@ -7,7 +7,8 @@ use alloc::sync::Arc; use alloc::vec::Vec; -use crate::{AimDb, DbResult, Router, RuntimeContext, TopicGrammar}; +use crate::router::Router; +use crate::{AimDb, DbResult, RuntimeContext, TopicGrammar}; /// Every inbound link of one scheme, compiled against the connector's grammar. /// diff --git a/aimdb-core/src/lib.rs b/aimdb-core/src/lib.rs index 7be4a774..b5ab18a2 100644 --- a/aimdb-core/src/lib.rs +++ b/aimdb-core/src/lib.rs @@ -93,7 +93,7 @@ pub mod profiling; pub mod record_id; #[cfg(feature = "remote")] pub mod remote; -pub mod router; +mod router; #[cfg(feature = "connector-session")] pub mod session; pub mod signal; @@ -117,10 +117,9 @@ pub use executor::{BoxFuture, ExecutorError, ExecutorResult, LogLevel, RuntimeOp pub use buffer::JsonReader; pub use buffer::Reader; pub use buffer::TryProduceError; -pub use builder::OutboundRoute; pub use builder::{AimDb, AimDbBuilder}; pub use connector::ConnectorBuilder; -pub use transport::{Connector, ConnectorConfig, PublishError}; +pub use transport::{ConnectorConfig, PublishError}; pub use typed_api::{ Consumer, InboundConnectorBuilder, OutboundConnectorBuilder, Producer, RecordRegistrar, StageKind, @@ -142,10 +141,9 @@ pub use remote::topic_leaf; // compatible). See docs/design/remote-access-via-connectors.md. #[cfg(feature = "connector-session")] pub use session::{ - is_wildcard, pattern_contains, pump_sink, pump_source, topic_matches, AuthError, BoxFut, - BoxStream, CodecError, Connection, Dialer, Dispatch, EnvelopeCodec, Inbound, Listener, - Outbound, Payload, PeerInfo, RpcError, SessionCtx, SessionLimits, Source, SubUpdate, - TransportError, TransportResult, + is_wildcard, pattern_contains, topic_matches, AuthError, BoxFut, BoxStream, CodecError, + Connection, Dialer, Dispatch, EnvelopeCodec, Inbound, Listener, Outbound, Payload, PeerInfo, + RpcError, SessionCtx, SessionLimits, SubUpdate, TransportError, TransportResult, }; // Signal gauge handle (always available; inert without `observability`) @@ -162,18 +160,13 @@ pub use profiling::{ pub use connector::TopicResolverFn; pub use connector::{ConnectorLink, ConnectorUrl, LinkAddress, SerializeError}; pub use connector::{IngestFactoryFn, IngestFn}; -pub use connector::{ - SerializedPayload, SerializedReader, SerializedSource, SerializedValue, SerializedValueInto, - SourceFactoryFn, -}; -pub use connector::{TopicBuf, TopicOverflow, TopicProvider, TopicWriter}; +pub use connector::{TopicBuf, TopicOverflow, TopicWriter}; // Router exports for connector implementations pub use inbound_dispatch::InboundDispatch; pub use outbound::{ OutboundMessage, OutboundPayload, OutboundRoutes, RouteId, RouteInfo, RouteStats, }; -pub use router::Router; // Topic grammar for connectors with wildcard subscriptions pub use topic_pattern::{ diff --git a/aimdb-core/src/outbound/routes.rs b/aimdb-core/src/outbound/routes.rs index 871e7cef..2d19a355 100644 --- a/aimdb-core/src/outbound/routes.rs +++ b/aimdb-core/src/outbound/routes.rs @@ -2,7 +2,6 @@ //! connector's transport task. use alloc::boxed::Box; -use alloc::string::ToString; use alloc::sync::Arc; use alloc::vec::Vec; use core::future::poll_fn; @@ -12,7 +11,7 @@ use super::ready::{Polled, ReadyRoutes}; use super::RouteId; use crate::connector::SerializeError; use crate::transport::ConnectorConfig; -use crate::{AimDb, ConfigError, DbError, DbResult, RuntimeContext}; +use crate::{AimDb, DbError, DbResult, RuntimeContext}; /// One outbound route, for parsing per-route configuration once at build. #[derive(Debug, Clone)] @@ -161,8 +160,6 @@ pub(crate) struct RouteParts { pub(crate) route: Box, pub(crate) topic_capacity: usize, pub(crate) payload_capacity: usize, - /// The link uses `with_topic_provider`, which this path does not support. - pub(crate) topic_provider: bool, } /// Builds a link's [`RouteParts`], subscribing to its record. @@ -200,31 +197,12 @@ const _: fn() = || { impl OutboundRoutes { /// Subscribes every outbound link of `scheme` and allocates the scratch /// once: the largest topic capacity plus the largest payload capacity. - /// - /// Rejects every link that uses `with_topic_provider`. pub fn new(db: &AimDb, scheme: &str) -> DbResult { let mut routes = Vec::new(); let mut states = Vec::new(); - let mut errors = Vec::new(); - - for (record_index, record_key, link) in db.outbound_links(scheme) { - let Some(factory) = &link.route_factory else { - errors.push(ConfigError::new( - record_key, - Some(link.url.to_string()), - "link was not registered through `link_to`", - )); - continue; - }; - let parts = factory(db); - if parts.topic_provider { - errors.push(ConfigError::new( - record_key, - Some(link.url.to_string()), - "`with_topic_provider` is not supported here; use `with_topic_writer` or `with_topic_fn`", - )); - continue; - } + + for (record_index, link) in db.outbound_links(scheme) { + let parts = (link.route_factory)(db); let mut config = ConnectorConfig::from_query(&link.config); config.record_index = Some(record_index); routes.push(RouteInfo { @@ -237,10 +215,6 @@ impl OutboundRoutes { states.push(parts.route); } - if !errors.is_empty() { - return Err(DbError::InvalidConfiguration { errors }); - } - let topic_region = routes.iter().map(|r| r.topic_capacity).max().unwrap_or(0); let payload_region = routes.iter().map(|r| r.payload_capacity).max().unwrap_or(0); Ok(Self { diff --git a/aimdb-core/src/router.rs b/aimdb-core/src/router.rs index 88eee71b..38b34c2e 100644 --- a/aimdb-core/src/router.rs +++ b/aimdb-core/src/router.rs @@ -19,8 +19,8 @@ use crate::topic_pattern::{Spans, TopicFilter, TopicGrammar, TopicMatch, MAX_CAP /// Generic message router for connector dispatch /// -/// Built by [`AimDb::inbound_router`](crate::AimDb::inbound_router). Routes -/// incoming messages to the matching records' ingest callbacks. Uses linear +/// Built by `AimDb::inbound_router` and held by +/// [`InboundDispatch`](crate::InboundDispatch). Routes incoming messages to the matching records' ingest callbacks. Uses linear /// search which is efficient for <100 routes. /// /// # Performance diff --git a/aimdb-core/src/session/mod.rs b/aimdb-core/src/session/mod.rs index 491436ad..98b9cba4 100644 --- a/aimdb-core/src/session/mod.rs +++ b/aimdb-core/src/session/mod.rs @@ -5,9 +5,7 @@ //! ([`EnvelopeCodec`]), and dispatch ([`Dispatch`]/[`Session`]), over a //! role-neutral [`Inbound`]/[`Outbound`] message set shared by the reactive //! server engine (`serve`/`run_session`) and the proactive client engine -//! (`run_client`/`pump_client`). Data-plane connectors use `pump_sink`/ -//! `pump_source` over the [`Source`] / [`Connector`](crate::transport::Connector) -//! capabilities. +//! (`run_client`/`pump_client`). //! //! All contracts are `dyn`-safe and compile on `std` and `no_std + alloc`. @@ -28,8 +26,6 @@ mod endpoint; #[cfg(feature = "connector-session")] mod io; #[cfg(feature = "connector-session")] -mod pump; -#[cfg(feature = "connector-session")] mod server; // Concrete AimX protocol substrate. The transport lives in a separate connector @@ -54,8 +50,6 @@ pub use io::{ OneShotListener, StreamDialer, StreamListener, }; #[cfg(feature = "connector-session")] -pub use pump::{pump_sink, pump_source}; -#[cfg(feature = "connector-session")] pub use server::{run_session, serve, SessionConfig}; // =========================================================================== @@ -517,19 +511,6 @@ pub trait EnvelopeCodec: Send + Sync { fn decode_outbound<'a>(&self, frame: &'a [u8]) -> Result, CodecError>; } -// =========================================================================== -// Data-plane capabilities — connectionless (an external library owns any -// session). The outbound `Sink` is the canonical -// [`Connector`](crate::transport::Connector); the inbound `Source` is below. -// =========================================================================== - -/// External → AimDB data-plane: a stream of inbound frames, drained by -/// `pump_source`. -pub trait Source: Send { - /// Yield the next `(topic, payload)`, or `None` when the source is done. - fn next(&mut self) -> BoxFut<'_, Option<(String, Payload)>>; -} - // =========================================================================== // Taking each trait as `&dyn Trait` forces the dyn-compatibility check on all // targets, not just under `cargo test`. @@ -543,7 +524,6 @@ fn _assert_object_safe( _dispatch: &dyn Dispatch, _session: &dyn Session, _codec: &dyn EnvelopeCodec, - _source: &dyn Source, ) { } @@ -632,13 +612,6 @@ mod tests { } } - struct MockSource; - impl Source for MockSource { - fn next(&mut self) -> BoxFut<'_, Option<(String, Payload)>> { - unimplemented!() - } - } - /// Every trait is `dyn`-usable. #[test] fn traits_are_object_safe() { @@ -648,7 +621,6 @@ mod tests { let _dispatch: Box = Box::new(MockDispatch); let _session: Box = Box::new(MockSession); let _codec: Box = Box::new(MockCodec); - let _source: Box = Box::new(MockSource); } /// `Box` satisfies the `Dialer` bound, so a runtime-selected diff --git a/aimdb-core/src/session/pump.rs b/aimdb-core/src/session/pump.rs deleted file mode 100644 index 7030acbf..00000000 --- a/aimdb-core/src/session/pump.rs +++ /dev/null @@ -1,164 +0,0 @@ -//! Data-plane pump helpers. -//! -//! Two free functions that own the boilerplate a data-plane connector used to -//! hand-roll. The author writes only the pure I/O adapter — a -//! [`Connector`](crate::transport::Connector) (outbound) and a [`Source`] -//! (inbound) — and composes the helpers in `build()` -//! (illustrative — `sink()`/`subscription()` are the author's own constructors): -//! -//! ```rust,ignore -//! let mut f = pump_sink(db, "redis", self.sink().await?); // outbound -//! let router = db.inbound_router("redis", &ExactGrammar)?; -//! f.extend(pump_source(db, router, self.subscription().await?)); // inbound -//! Ok(f) -//! ``` -//! -//! Both are `no_std + alloc`-native (boxed futures, no `tokio`). - -use alloc::boxed::Box; -use alloc::sync::Arc; -use alloc::vec; -use alloc::vec::Vec; - -use super::Source; -use crate::builder::{AimDb, BoxFuture}; -use crate::router::Router; -use crate::transport::{Connector, ConnectorConfig}; - -/// Outbound pump: one publisher future per outbound route on `scheme`. -/// -/// Extracts the consume-and-publish loop a data-plane connector used to write by -/// hand. For each route from [`collect_outbound_routes`](AimDb::collect_outbound_routes), -/// the returned future subscribes to the route's fused -/// [`SerializedSource`](crate::connector::SerializedSource) — whose readers -/// yield destination + serialized payload directly (no `Box` per -/// message) — and publishes through `sink`. Per-route -/// configuration (`qos`/`retain`/…) is built once from the route's URL query -/// via [`ConnectorConfig::from_query`]. -/// -/// The publisher future terminates when its subscription yields an error (e.g. the -/// record buffer closed), matching the legacy hand-rolled loop. -pub fn pump_sink(db: &AimDb, scheme: &str, sink: Arc) -> Vec { - let routes = db.collect_outbound_routes(scheme); - let mut futures: Vec = Vec::with_capacity(routes.len()); - - for crate::OutboundRoute { - topic: default_topic, - source, - config, - } in routes - { - let sink = sink.clone(); - let runtime_ctx = db.runtime_ctx(); - let cfg = ConnectorConfig::from_query(&config); - let scratch_capacity = source.serializer_scratch_capacity().unwrap_or(0); - - futures.push(Box::pin(async move { - // Subscribe inside the pump future (not at collect time), so the - // ring-buffer cursor starts when the publisher actually runs. - let mut reader = source.subscribe(); - // One bounded allocation per route, reused for every successful - // into-slice serialization. Owned-only sources request zero bytes. - let mut scratch = vec![0_u8; scratch_capacity]; - - log_info!( - "pump_sink: publisher started for destination: {}", - default_topic - ); - - loop { - let msg = match reader - .recv_into(&runtime_ctx, scratch.as_mut_slice()) - .await - { - Ok(m) => m, - // SPMC-ring overflow: messages were missed, but the reader - // recovers (cursor resets to the oldest live value). Skip the - // gap and keep pumping — a transient lag must not permanently - // kill the publisher. - Err(crate::DbError::BufferLagged { .. }) => { - log_warn!("pump_sink: consumer lagged for '{}'", default_topic); - continue; - } - // Buffer closed / fatal — the record is gone; end the publisher. - Err(_e) => { - log_info!( - "pump_sink: publisher stopping for '{}': {:?}", - default_topic, - _e - ); - break; - } - }; - let crate::connector::SerializedValueInto { dest, payload } = msg; - // Destination: dynamic (resolved by the source) or default (from URL). - let dest = dest.as_deref().unwrap_or(&default_topic); - - let payload = match &payload { - crate::connector::SerializedPayload::Scratch { len } => { - let Some(payload) = scratch.get(..*len) else { - log_error!( - "pump_sink: serializer returned invalid length {} for {}-byte scratch buffer", - len, - scratch.len() - ); - continue; - }; - payload - } - crate::connector::SerializedPayload::Owned(payload) => payload.as_slice(), - }; - - // Publish through the connector's pure I/O adapter. - if let Err(_e) = sink.publish(dest, &cfg, payload).await { - log_error!("pump_sink: failed to publish to '{}': {:?}", dest, _e); - } else { - log_debug!("pump_sink: published to: {}", dest); - } - } - - log_info!( - "pump_sink: publisher stopped for destination: {}", - default_topic - ); - })); - } - - futures -} - -/// Inbound pump: a single multiplexed reader future. -/// -/// Drives one [`Source`] (never one task per topic), fanning each -/// `(topic, payload)` out to the matching producers via `router`, from -/// [`AimDb::inbound_router`], whose subscriptions the connector made. -/// -/// Backpressure: [`Router::route`] drops + logs on a full producer buffer rather -/// than blocking, so one slow record never stalls the shared source. Route errors -/// are non-fatal. -pub fn pump_source(db: &AimDb, router: Router, mut src: impl Source + 'static) -> Vec { - let router = Arc::new(router); - let ctx = db.runtime_ctx(); - - vec![Box::pin(async move { - log_info!( - "pump_source: reader started ({} topics)", - router.resource_ids().len() - ); - - while let Some((topic, payload)) = src.next().await { - // `route` deserializes and fans out to producers (synchronously — - // the fused ingest path never awaits); it drops + logs on a full - // producer buffer and never returns a fatal error. - if let Err(_e) = router.route(&topic, &payload, &ctx) { - log_error!( - "pump_source: failed to route message on '{}': {}", - topic, - _e - ); - } - } - - log_info!("pump_source: reader stopped"); - })] -} diff --git a/aimdb-core/src/topic_pattern.rs b/aimdb-core/src/topic_pattern.rs index 7f7610b6..8a1eb946 100644 --- a/aimdb-core/src/topic_pattern.rs +++ b/aimdb-core/src/topic_pattern.rs @@ -216,7 +216,7 @@ impl TopicFilter for ExactFilter { } } -/// A stub grammar for tests of the router and of `inbound_router`. +/// A stub grammar for tests of the router and of `InboundDispatch`. #[cfg(test)] pub(crate) mod test_support { use super::*; diff --git a/aimdb-core/src/transport.rs b/aimdb-core/src/transport.rs index da5cf6a5..ef9ebf56 100644 --- a/aimdb-core/src/transport.rs +++ b/aimdb-core/src/transport.rs @@ -1,23 +1,13 @@ -//! Transport connector traits for protocol-agnostic publishing +//! Protocol-agnostic per-route configuration and publish errors. //! -//! Provides a generic `Connector` trait that enables scheme-based routing -//! to different transport protocols. Each connector manages a single connection -//! to a specific endpoint (e.g., one MQTT broker). -//! -//! # Design Philosophy -//! -//! - **Scheme-based routing**: the URL scheme (e.g. `mqtt://`, `knx://`) determines which connector handles requests -//! - **Single endpoint per connector**: Each connector connects to ONE broker/resource -//! - **Multi-transport publishing**: Same data can be published to multiple protocols -//! - **Protocol-agnostic core**: Core knows schemes and key/value options, never protocol semantics +//! Core knows schemes and key/value options, never protocol semantics. -use alloc::{boxed::Box, string::String, vec::Vec}; -use core::future::Future; -use core::pin::Pin; +use alloc::{string::String, vec::Vec}; /// Protocol-agnostic connector configuration /// -/// Carries the route's key/value options to [`Connector::publish`]. Only the +/// Carries the route's key/value options to the connector +/// ([`RouteInfo::config`](crate::RouteInfo::config)). Only the /// genuinely protocol-agnostic `timeout_ms` is a typed field; every /// protocol-specific knob (e.g. MQTT's `qos`/`retain`) travels in /// [`protocol_options`](ConnectorConfig::protocol_options) and is interpreted @@ -53,12 +43,8 @@ impl Default for ConnectorConfig { impl ConnectorConfig { /// Build a config from a route's URL-query key/value pairs. /// - /// This is the shared seam the data-plane `pump_sink` helper uses to thread - /// per-route configuration through to [`Connector::publish`] without changing - /// the `publish` signature. - /// - /// Only the protocol-agnostic `timeout_ms` and `record_index` (stamped by - /// `AimDb::collect_outbound_routes`; the last occurrence wins) are lifted + /// Only the protocol-agnostic `timeout_ms` and `record_index` (the last + /// occurrence wins) are lifted /// into typed fields; every other key is passed through verbatim in /// [`protocol_options`](ConnectorConfig::protocol_options) for the /// connector to interpret with its own defaults. @@ -134,92 +120,9 @@ impl std::fmt::Display for PublishError { #[cfg(feature = "std")] impl std::error::Error for PublishError {} -/// Generic transport connector trait for protocol-agnostic publishing -/// -/// This trait enables multi-protocol publishing via scheme-based routing -/// (e.g. `mqtt://topic` → MQTT broker, `knx://1/0/6` → KNX group address). -/// -/// Each connector manages ONE connection/endpoint. For multiple brokers/endpoints, -/// create multiple connectors and register them with different schemes. -/// -/// # Example Implementation -/// -/// Illustrative sketch (not compiled: the MQTT client types are fictional — -/// see `aimdb-mqtt-connector` for a real implementation): -/// -/// ```rust,ignore -/// impl Connector for MqttConnector { -/// fn publish( -/// &self, -/// destination: &str, // "sensors/temperature" -/// config: &ConnectorConfig, -/// payload: &[u8], -/// ) -> Pin> + Send + '_>> { -/// // Protocol knobs come from the route's key/value options, -/// // with connector-chosen defaults. -/// let qos = config -/// .protocol_options -/// .iter() -/// .find(|(k, _)| k == "qos") -/// .and_then(|(_, v)| v.parse::().ok()) -/// .unwrap_or(1); -/// Box::pin(async move { -/// self.client.publish(destination, qos, payload).await -/// .map_err(|_| PublishError::ConnectionFailed) -/// }) -/// } -/// } -/// ``` -/// -/// # Thread Safety -/// -/// Requires Send + Sync for Tokio compatibility. For Embassy (single-threaded), -/// use `unsafe impl Send + Sync` with safety documentation. -pub trait Connector: Send + Sync { - /// Publish data to a protocol-specific destination - /// - /// # Arguments - /// * `destination` - Protocol-specific path, no broker/host info - /// (e.g. an MQTT topic like "sensors/temperature") - /// * `config` - Publishing configuration (timeout + protocol options) - /// * `payload` - Message payload as byte slice - /// - /// # Returns - /// `Ok(())` on success, `PublishError` on failure - fn publish( - &self, - destination: &str, - config: &ConnectorConfig, - payload: &[u8], - ) -> Pin> + Send + '_>>; -} - #[cfg(test)] mod tests { use super::*; - use alloc::sync::Arc; - - // Mock connector for testing - struct MockConnector; - - impl Connector for MockConnector { - fn publish( - &self, - _destination: &str, - _config: &ConnectorConfig, - _payload: &[u8], - ) -> Pin> + Send + '_>> { - Box::pin(async move { Ok(()) }) - } - } - - #[test] - fn test_connector_trait() { - let connector = Arc::new(MockConnector); - - // Verify the connector can be used as a trait object - let _trait_obj: Arc = connector; - } #[test] fn test_connector_config_default() { diff --git a/aimdb-core/src/typed_api.rs b/aimdb-core/src/typed_api.rs index 3a0e784d..8490cb3a 100644 --- a/aimdb-core/src/typed_api.rs +++ b/aimdb-core/src/typed_api.rs @@ -284,10 +284,10 @@ impl Clone for Consumer { } // ============================================================================ -// Fused outbound source +// Outbound routes // ============================================================================ -/// Type alias for the unified typed serializer captured by [`FusedSource`] +/// Type alias for the unified typed serializer captured by [`TypedRoute`] /// /// Raw and context-aware serializers collapse into this shape at `finish()`; /// the raw variant simply ignores the threaded context. @@ -302,7 +302,7 @@ type ConsumerFactoryFn = Arc Consumer + Send + Sync>; /// Optional allocation-free serializer captured beside [`FusedSerializeFn`]. /// -/// The callback writes into one pump-owned bounded scratch buffer. Returning +/// The callback writes into the `OutboundRoutes` scratch buffer. Returning /// `BufferTooSmall` selects the owned serializer for that value; other failures /// retain the existing skip-and-log behavior. type FusedSerializeIntoFn = Arc< @@ -311,88 +311,12 @@ type FusedSerializeIntoFn = Arc< + Sync, >; -/// How an outbound link picks each value's destination. Setting one replaces -/// the other. -enum TopicSelector { - None, - Provider(Arc>), - Writer { - capacity: usize, - writer: Arc>, - }, -} - -impl Clone for TopicSelector { - fn clone(&self) -> Self { - match self { - Self::None => Self::None, - Self::Provider(p) => Self::Provider(p.clone()), - Self::Writer { capacity, writer } => Self::Writer { - capacity: *capacity, - writer: writer.clone(), - }, - } - } -} - -/// The [`SerializedSource`](crate::connector::SerializedSource) built by -/// `OutboundConnectorBuilder::finish()` — holds the typed consumer, -/// serializer, and optional topic provider, so every per-message step stays -/// typed (no `Box`). -struct FusedSource { - consumer: Consumer, - serialize: FusedSerializeFn, - serialize_into: Option<(usize, FusedSerializeIntoFn)>, - topic: TopicSelector, -} - -impl crate::connector::SerializedSource for FusedSource -where - T: Send + Sync + 'static + Debug + Clone, -{ - fn serializer_scratch_capacity(&self) -> Option { - self.serialize_into.as_ref().map(|(capacity, _)| *capacity) - } - - fn subscribe(&self) -> Box { - let topic_capacity = match &self.topic { - TopicSelector::Writer { capacity, .. } => *capacity, - _ => 0, - }; - Box::new(FusedReader { - inner: self.consumer.subscribe(), - serialize: self.serialize.clone(), - serialize_into: self - .serialize_into - .as_ref() - .map(|(_, serialize_into)| serialize_into.clone()), - topic: self.topic.clone(), - topic_buf: alloc::vec![0; topic_capacity].into_boxed_slice(), - }) - } -} - -/// One subscription of a [`FusedSource`]: recv → resolve destination → -/// serialize, all on the typed value. -/// -/// The connector SPI keeps its boxed `RecvSerializedFuture` (BYOC stays -/// stable); only the *inner* per-message box is eliminated by reading through -/// the allocation-free [`Reader`](crate::buffer::Reader). -struct FusedReader { - inner: crate::buffer::Reader, - serialize: FusedSerializeFn, - serialize_into: Option>, - topic: TopicSelector, - /// Storage a [`TopicWriter`](crate::connector::TopicWriter) writes into; - /// empty without one. - topic_buf: Box<[u8]>, -} +/// An outbound link's topic writer and the capacity it writes into. +type TopicWriterCfg = (usize, Arc>); /// Runs `writer` for `value` into `out`. `Ok(true)`: publish to what `out` /// holds; `Ok(false)`: to the link's default topic; `Err`: the topic did not /// fit and the value is skipped, whatever the writer returned. -/// -/// Shared by [`FusedReader`] and [`TypedRoute`]. fn write_topic + ?Sized>( writer: &W, value: &T, @@ -407,8 +331,6 @@ fn write_topic + ?Sized>( /// Serializes one outbound value: into `scratch` through /// `with_serializer_into` when the link has one, else (or when the value does /// not fit) through the owned serializer. -/// -/// Shared by [`FusedReader`] and [`TypedRoute`]. fn serialize_outbound( serialize: &FusedSerializeFn, serialize_into: Option<&FusedSerializeIntoFn>, @@ -437,103 +359,6 @@ fn serialize_outbound( } } -impl FusedReader { - /// The value's destination (`None`: the route's default), or `Err` when - /// its written topic overflowed and the value must be skipped. - fn resolve_dest(&mut self, value: &T) -> Result, ()> { - match &self.topic { - TopicSelector::None => Ok(None), - TopicSelector::Provider(p) => Ok(p.topic(value)), - TopicSelector::Writer { writer, .. } => { - let mut out = crate::connector::TopicBuf::new(&mut self.topic_buf); - match write_topic(&**writer, value, &mut out) { - Ok(true) => Ok(Some(out.as_str().to_string())), - Ok(false) => Ok(None), - Err(_) => { - log_warn!( - "outbound link: topic for {} does not fit in {} bytes, value skipped", - core::any::type_name::(), - out.capacity() - ); - Err(()) - } - } - } - } - } -} - -impl crate::connector::SerializedReader for FusedReader { - fn recv<'a>( - &'a mut self, - ctx: &'a crate::RuntimeContext, - ) -> crate::connector::RecvSerializedFuture<'a> { - Box::pin(async move { - loop { - // Buffer errors propagate unchanged: `BufferLagged` lets the - // pump skip the gap and keep going; anything else ends it. - let value = self.inner.recv().await?; - // Resolve the destination while the typed value is in hand. - let Ok(dest) = self.resolve_dest(&value) else { - continue; - }; - match (self.serialize)(ctx, &value) { - Ok(payload) => return Ok(crate::connector::SerializedValue { dest, payload }), - Err(_e) => { - // Same skip-and-log the pumps used to do around the - // erased serializer. - log_error!( - "outbound link: failed to serialize {} (dest {:?}): {:?}", - core::any::type_name::(), - dest, - _e - ); - continue; - } - } - } - }) - } - - fn recv_into<'a>( - &'a mut self, - ctx: &'a crate::RuntimeContext, - scratch: &'a mut [u8], - ) -> crate::connector::RecvSerializedIntoFuture<'a> { - use crate::connector::SerializedPayload; - use crate::outbound::StagedPayload; - - Box::pin(async move { - loop { - let value = self.inner.recv().await?; - let Ok(dest) = self.resolve_dest(&value) else { - continue; - }; - let payload = match serialize_outbound( - &self.serialize, - self.serialize_into.as_ref(), - ctx, - &value, - scratch, - ) { - Ok(StagedPayload::Scratch(len)) => SerializedPayload::Scratch { len }, - Ok(StagedPayload::Owned(bytes)) => SerializedPayload::Owned(bytes), - Err(_failure) => { - log_error!( - "outbound link: {} (dest {:?}): {}, value skipped", - core::any::type_name::(), - dest, - _failure - ); - continue; - } - }; - return Ok(crate::connector::SerializedValueInto { dest, payload }); - } - }) - } -} - /// One outbound link's per-route state inside /// [`OutboundRoutes`](crate::OutboundRoutes): reader, topic writer and /// serializers, all typed. @@ -868,7 +693,7 @@ where config: Vec::new(), context_serializer: None, context_serializer_into: None, - topic: TopicSelector::None, + topic_writer: None, } } @@ -902,7 +727,7 @@ pub struct OutboundConnectorBuilder<'r, 'a, T: Send + Sync + 'static + Debug + C config: Vec<(String, String)>, context_serializer: Option>, context_serializer_into: Option<(usize, TypedContextSerializerIntoFn)>, - topic: TopicSelector, + topic_writer: Option>, } impl<'r, 'a, T> OutboundConnectorBuilder<'r, 'a, T> @@ -982,26 +807,6 @@ 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 - /// - /// 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 - where - P: crate::connector::TopicProvider + 'static, - { - // Stays typed: fused into the link's SerializedSource at finish(). - self.topic = TopicSelector::Provider(Arc::new(provider)); - self - } - /// Sets a [`TopicWriter`](crate::connector::TopicWriter) that writes each /// value's destination. /// @@ -1012,10 +817,7 @@ where where W: crate::connector::TopicWriter + 'static, { - self.topic = TopicSelector::Writer { - capacity, - writer: Arc::new(writer), - }; + self.topic_writer = Some((capacity, Arc::new(writer))); self } @@ -1172,33 +974,13 @@ where }) }; - // Fused source for the pumps: the serializer and topic selector ride - // along typed, so its readers yield destination + payload with no - // erasure crossing. - let source_factory: crate::connector::SourceFactoryFn = { - let make_consumer = make_consumer.clone(); - let (serialize, serialize_into, topic) = ( - serialize.clone(), - serialize_into.clone(), - self.topic.clone(), - ); - Arc::new(move |db: &AimDb| { - Box::new(FusedSource { - consumer: make_consumer(db), - serialize: serialize.clone(), - serialize_into: serialize_into.clone(), - topic: topic.clone(), - }) as Box - }) - }; - - // The same parts for `OutboundRoutes`, subscribed when it is built. + // Subscribed when `OutboundRoutes` is built. let route_factory: crate::outbound::RouteFactoryFn = { - let topic = self.topic; + let topic_writer = self.topic_writer; Arc::new(move |db: &AimDb| { - let (writer, topic_capacity) = match &topic { - TopicSelector::Writer { capacity, writer } => (Some(writer.clone()), *capacity), - _ => (None, 0), + let (topic_capacity, writer) = match &topic_writer { + Some((capacity, writer)) => (*capacity, Some(writer.clone())), + None => (0, None), }; crate::outbound::RouteParts { route: Box::new(TypedRoute { @@ -1209,17 +991,14 @@ where }), topic_capacity, payload_capacity: serialize_into.as_ref().map_or(0, |(capacity, _)| *capacity), - topic_provider: matches!(topic, TopicSelector::Provider(_)), } }) }; - let mut link = ConnectorLink::new(url, source_factory); + let mut link = ConnectorLink::new(url, route_factory); link.config = self.config; - link.route_factory = Some(route_factory); - // 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 } @@ -1489,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"))] @@ -2189,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, @@ -2239,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; @@ -2276,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| { @@ -2298,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]; @@ -2332,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 })) @@ -2358,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'")); @@ -2366,80 +2143,17 @@ 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")); } - // ==================================================================== - // InboundDispatch: the same cases through the connector entry point - // ==================================================================== - - #[tokio::test] - async fn inbound_dispatch_routes_patterns_with_shared_keys() { - let (db, last, count) = inbound_db(|reg| { - reg.link_from("mqtt://temp/{d}") - .key("d", 2) - .with_match_deserializer(key_deser) - .finish(); - reg.link_from("mqtt://hum/{id}") - .key("id", 2) - .with_match_deserializer(key_deser) - .finish(); - reg.link_from("mqtt://cmd/in") - .with_deserializer(|_ctx, _bytes: &[u8]| Ok(TestRecord { value: 100 })) - .finish(); - }) - .await; - - 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 dispatch = |topic: &str| { - inbound.dispatch(topic, b""); - last.load(Ordering::SeqCst) - }; - 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); - dispatch("temp/c"); - assert_eq!(count.load(Ordering::SeqCst), produced, "full table drops"); - } - - #[tokio::test] - async fn inbound_dispatch_rejects_links_it_cannot_compile() { - let (db, _, _) = inbound_db(|reg| { - reg.link_from("mqtt://r/one") - .with_topic_resolver(|| Some("r/{d".into())) - .with_deserializer(|_ctx, _bytes: &[u8]| Ok(TestRecord { value: 0 })) - .finish(); - reg.link_from("mqtt://k/{d}") - .key("d", 4) - .with_topic_resolver(|| Some("k/{x}".into())) - .with_match_deserializer(key_deser) - .finish(); - }) - .await; - - 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'")); - assert!(errors[1] - .message - .contains("key 'd' is not a capture of 'k/{x}'")); - } - #[tokio::test] async fn inbound_dispatch_clones_share_records() { fn assert_send_sync() {} @@ -2464,34 +2178,8 @@ mod tests { assert_eq!(last.load(Ordering::SeqCst), 4); } - #[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 }) - } - } - - let (db, last, _) = inbound_db(|reg| { - reg.link_from("mqtt://temp/{d}") - .key("d", 2) - .with_match_deserializer(key_deser) - .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); - } - // ==================================================================== - // Fused outbound reader tests + // OutboundRoutes over typed links // ==================================================================== /// Buffer reader that replays a fixed script, then reports the buffer @@ -2500,430 +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: TopicSelector, - ) -> FusedReader { - let topic_capacity = match &topic { - TopicSelector::Writer { capacity, .. } => *capacity, - _ => 0, - }; - FusedReader { - inner: crate::buffer::Reader::new(Box::new(ScriptedReader { script })), - serialize, - serialize_into: None, - topic, - topic_buf: alloc::vec![0; topic_capacity].into_boxed_slice(), - } + 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: TopicSelector::None, - topic_buf: Box::default(), - } + 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())), - TopicSelector::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()) - } - }), - TopicSelector::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"); - - assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 }); - assert_eq!(&scratch[..4], 2i32.to_le_bytes().as_slice()); - } - - #[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"); - - assert_eq!(msg.payload, SerializedPayload::Scratch { len: 4 }); - assert_eq!(&scratch[..4], 2i32.to_le_bytes().as_slice()); - } - - /// The destination is resolved from the typed value while it is in hand. - #[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())), - TopicSelector::Provider(Arc::new(PositiveTopic)), - ); - let ctx = test_ctx(); - - let first = reader.recv(&ctx).await.expect("value"); - assert_eq!(first.dest.as_deref(), Some("dyn/5")); - - let second = reader.recv(&ctx).await.expect("value"); - assert_eq!(second.dest, None); // falls back to the route default - } + 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; - fn writer( - capacity: usize, - f: impl Fn(&TestRecord, &mut crate::TopicBuf<'_>) -> Result - + Send - + Sync - + 'static, - ) -> TopicSelector { - TopicSelector::Writer { - capacity, - writer: Arc::new(f), - } + assert_eq!(pull(&mut o).await, Some(("tele/out".into(), 3, false))); + assert_eq!(o.stats(0).unwrap().serialize_failed, 2); } - /// The same case with a written topic; `Ok(false)` uses the default. + /// `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_resolves_written_topic() { + async fn with_topic_fn_writes_the_topic_or_the_default() { use core::fmt::Write as _; - 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())), - writer(8, |v, out| { - if v.value <= 0 { - return Ok(false); - } - write!(out, "dyn/{}", v.value)?; - Ok(true) - }), - ); - let ctx = test_ctx(); - - let first = reader.recv(&ctx).await.expect("value"); - assert_eq!(first.dest.as_deref(), Some("dyn/5")); + 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; - let second = reader.recv(&ctx).await.expect("value"); - assert_eq!(second.dest, None); + 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())); } /// A topic that overflows skips its value, even when the writer ignores /// the error and returns `Ok(true)`. #[tokio::test] - async fn fused_reader_skips_overflowing_topics() { + async fn overflowing_topics_are_skipped() { use core::fmt::Write as _; - let mut reader = fused_reader( - vec![ - Ok(TestRecord { value: 123_456 }), - Ok(TestRecord { value: 1_234_567 }), - Ok(TestRecord { value: 7 }), - ], - Arc::new(|_ctx, r| Ok(r.value.to_le_bytes().to_vec())), - writer(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) - }), - ); - let ctx = test_ctx(); + 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; // "t/123456" (8 bytes) and "t/1234567" (9 bytes) do not fit in 6. - let msg = reader.recv(&ctx).await.expect("value"); - assert_eq!(msg.dest.as_deref(), Some("t/7")); - assert_eq!(msg.payload, 7i32.to_le_bytes().to_vec()); - } - - /// `with_topic_fn` infers an unannotated closure's argument types, which - /// a generic `W: TopicWriter` bound cannot. - #[tokio::test] - async fn with_topic_fn_infers_and_writes_the_topic() { - use core::fmt::Write as _; - 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 - } - } - - let mut builder = crate::AimDbBuilder::new() - .runtime(Arc::new(MockRuntime)) - .with_connector(NoopConnectorBuilder); - builder.configure::("rec.out", |reg| { - reg.buffer_raw(Box::new(CannedBuffer)); - reg.link_to("mqtt://tele/out") - .with_topic_fn(16, |v, out| { - write!(out, "dyn/{}", v.value)?; - Ok(true) - }) - .with_serializer(|_ctx, r: &TestRecord| Ok(r.value.to_le_bytes().to_vec())) - .finish(); - }); - let (db, _runner) = builder.build().await.expect("build must succeed"); - - let routes = db.collect_outbound_routes("mqtt"); - let mut reader = routes[0].source.subscribe(); - let msg = reader.recv(&db.runtime_ctx()).await.expect("value"); - assert_eq!(msg.dest.as_deref(), Some("dyn/5")); + 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/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>; -} - -/// Force-`Send` bridge turning an [`EmbassySinkRaw`] into a [`Connector`], so an -/// Embassy outbound sink rides core's [`pump_sink`](aimdb_core::session::pump_sink) -/// unchanged. -pub struct EmbassySink(pub C); - -// SAFETY: single-core cooperative Embassy executor — see the module-level invariant. -unsafe impl Send for EmbassySink {} -// SAFETY: same invariant; `Connector` is shared behind `Arc`. -unsafe impl Sync for EmbassySink {} - -impl Connector for EmbassySink { - fn publish( - &self, - destination: &str, - config: &ConnectorConfig, - payload: &[u8], - ) -> Pin> + Send + '_>> { - // Own the args so the inner future borrows only `&self` (see trait doc). - Box::pin(SendFutureWrapper(self.0.publish( - destination.to_string(), - config.clone(), - payload.to_vec(), - ))) - } -} - -/// The pure inbound I/O a data-plane connector contributes: yield the next -/// `(topic, payload)`. The `!Send` dual of [`Source`]; [`EmbassySource`] -/// force-`Send`s it so it can drive core's -/// [`pump_source`](aimdb_core::session::pump_source). -pub trait EmbassySourceRaw { - /// Yield the next `(topic, payload)`, or `None` when the source is done. - fn next(&mut self) -> impl Future>; -} - -/// Force-`Send` bridge turning an [`EmbassySourceRaw`] into a [`Source`], so an -/// Embassy inbound stream rides core's -/// [`pump_source`](aimdb_core::session::pump_source) unchanged. -pub struct EmbassySource(pub S); - -// SAFETY: single-core cooperative Embassy executor — see the module-level invariant. -unsafe impl Send for EmbassySource {} - -impl Source for EmbassySource { - fn next(&mut self) -> BoxFut<'_, Option<(String, Payload)>> { - Box::pin(SendFutureWrapper(self.0.next())) - } -} - /// Force-`Send` + box a connector's long-lived **protocol task** (an MQTT broker /// manager, a KNX tunnelling state machine, …) so it can join the runner's /// `Send` future set without the connector touching [`SendFutureWrapper`]. diff --git a/aimdb-embassy-adapter/src/lib.rs b/aimdb-embassy-adapter/src/lib.rs index 63cdc955..263de655 100644 --- a/aimdb-embassy-adapter/src/lib.rs +++ b/aimdb-embassy-adapter/src/lib.rs @@ -28,11 +28,11 @@ pub mod buffer; #[cfg(not(feature = "std"))] mod runtime; -// Force-`Send` helper for Embassy data-plane connectors (see module docs). +// Force-`Send` helper for Embassy connectors (see module docs). #[cfg(not(feature = "std"))] pub mod send_wrapper; -// Centralized Embassy connector spines (session + data-plane) — the one audited +// Centralized Embassy connector helpers — the one audited // home for the single-core `unsafe` + `SendFutureWrapper`. #[cfg(all(not(feature = "std"), feature = "connectors"))] pub mod connectors; diff --git a/aimdb-embassy-adapter/src/send_wrapper.rs b/aimdb-embassy-adapter/src/send_wrapper.rs index 6a634a03..b2db959c 100644 --- a/aimdb-embassy-adapter/src/send_wrapper.rs +++ b/aimdb-embassy-adapter/src/send_wrapper.rs @@ -4,13 +4,8 @@ //! returns `Vec>>` and `AimDbRunner` drives them on a //! `Send` `BoxFuture`. Embassy's primitives (channels over `NoopRawMutex`, …) are //! `!Send` *by design* — single-core, cooperative, no preemption or thread -//! migration — so an Embassy connector's data-plane futures must be force-`Send`ed -//! to satisfy that bound. -//! -//! This is also *why* Embassy data-plane connectors hand-roll their outbound / -//! inbound loops instead of riding core's `pump_sink` / `pump_source`: those need a -//! `Send + Sync` `Connector` / `Send` `Source`, which `!Send` Embassy channels -//! cannot be without force-`Send`ing every primitive. +//! migration — so an Embassy connector's futures must be force-`Send`ed to +//! satisfy that bound. use core::future::Future; use core::pin::Pin; diff --git a/aimdb-knx-connector/Cargo.toml b/aimdb-knx-connector/Cargo.toml index d50e4329..198c074f 100644 --- a/aimdb-knx-connector/Cargo.toml +++ b/aimdb-knx-connector/Cargo.toml @@ -20,14 +20,9 @@ default = ["aimdb-core/alloc"] # by passing an adapter's binder and clock — a dependency the caller adds, not a # feature here — which is what lets a third runtime (FreeRTOS/lwIP) use this # crate with no edit to it. -# -# `embassy-sync` no longer backs a public type here (the connection task pulls -# from the record buffers directly); it is dropped together with the old -# connector SPI. connector = [ "aimdb-core/alloc", "aimdb-core/connector-session", # `DatagramBinder`/`Datagram`/`Delay` - "embassy-sync", ] # Orthogonal to `connector`: the code is `alloc`-only either way, so this only @@ -41,20 +36,10 @@ std = ["connector", "aimdb-core/std", "knx-pico/std"] tokio-runtime = ["std"] embassy-runtime = ["connector"] -# Selects `critical-section`'s std implementation. -# -# `embassy-sync`'s only `Sync` raw mutex is `CriticalSectionRawMutex`, which -# needs a `critical-section` impl to link. That impl is registered by symbol -# name and is global to the binary, so per `critical-section`'s own docs only -# the **final binary** may choose one: a library enabling it would hand every -# downstream binary a duplicate-symbol link error with no way to opt out. -# -# So this stays off by default and out of `connector`/`std`. A std binary that -# instantiates a `CriticalSectionRawMutex` channel and has no other impl in its -# graph can either depend on `critical-section` with `features = ["std"]` directly (the -# documented way) or enable this feature. On Embassy the HAL (cortex-m / -# embassy-rp / ...) already provides one. -critical-section-std-impl = ["critical-section/std"] +# Deprecated no-op. It selected `critical-section`'s std impl for the +# channels the connector used to take; nothing here needs one any more. Kept +# so existing consumers keep building; remove after a release. +critical-section-std-impl = [] # Design 050 §10.4/§10.5: the facade reaches both destinations through # `aimdb_core::__private`, so neither dependency is declared here any more. @@ -72,19 +57,11 @@ aimdb-core = { version = "2.0.0", path = "../aimdb-core", default-features = fal # Keyed `knx-pico` to keep imports and feature references unchanged. knx-pico = { package = "aimdb-knx-pico", version = "0.3.1", default-features = false } -# Executor-independent, despite the names: neither pulls an executor and both -# build on std, so one channel type and one select serve every runtime. -# `embassy-sync` is optional because it is only reachable through `connector`; -# `embassy-futures` is unconditional (its `[dependencies]` is empty bar optional -# defmt/log, so the std graph is unaffected). -embassy-sync = { version = "0.8.0", path = "../_external/embassy/embassy-sync", optional = true } +# Executor-independent, despite the name: it pulls no executor and builds on +# std, so one select serves every runtime. Its `[dependencies]` is empty bar +# optional defmt/log, so the std graph is unaffected. embassy-futures = { version = "0.1.2" } -# A `critical-section` impl must be linked wherever `CriticalSectionRawMutex` -# is used. Choosing one is the final binary's call, so it is reachable only -# through the opt-in `critical-section-std-impl` feature above. -critical-section = { version = "1.1", optional = true } - # Embedded utilities (heapless is unconditional: the shared sans-io tunnel # engine uses stack-allocated frames on both runtimes) heapless = { workspace = true } diff --git a/aimdb-knx-connector/README.md b/aimdb-knx-connector/README.md index 9587c430..4fff84d5 100644 --- a/aimdb-knx-connector/README.md +++ b/aimdb-knx-connector/README.md @@ -10,12 +10,8 @@ Add to your `Cargo.toml`: ```toml [dependencies] -# `std` is the host leg. `critical-section-std-impl` selects the impl the -# connector's channels need to link — only a final binary may pick one. -aimdb-knx-connector = { version = "0.5", features = [ - "std", - "critical-section-std-impl", -] } +# `std` is the host leg. +aimdb-knx-connector = { version = "0.5", features = ["std"] } # The host also needs the adapter's UDP socket and clock. aimdb-tokio-adapter = { version = "0.6", features = ["tokio-runtime", "net"] } diff --git a/aimdb-knx-connector/src/lib.rs b/aimdb-knx-connector/src/lib.rs index bd7f2312..6e5cdf03 100644 --- a/aimdb-knx-connector/src/lib.rs +++ b/aimdb-knx-connector/src/lib.rs @@ -16,14 +16,12 @@ //! enables. //! - `std`: `connector` plus core's `std`, knx-pico's std error impls, and the //! back-compat DPT re-exports. Lifts `no_std`; adds no runtime. -//! - `critical-section-std-impl`: **final binaries only** — selects -//! `critical-section`'s std impl for a host binary that needs one. An -//! Embassy HAL already provides one. //! - `tracing`: Debug logging support (std) //! - `defmt`: Debug logging support (no_std) //! //! `tokio-runtime` and `embassy-runtime` are deprecated aliases for `std` and -//! `connector` respectively, kept for one release. +//! `connector` respectively, and `critical-section-std-impl` a deprecated +//! no-op, kept for one release. //! //! ## Production Status //! diff --git a/aimdb-knx-connector/tests/topic_writer_tests.rs b/aimdb-knx-connector/tests/topic_writer_tests.rs index 42055a0b..c976ed45 100644 --- a/aimdb-knx-connector/tests/topic_writer_tests.rs +++ b/aimdb-knx-connector/tests/topic_writer_tests.rs @@ -342,40 +342,6 @@ async fn test_knx_topic_writer_with_connector_registration() { assert!(builder.build().await.is_ok()); } -/// The connector pulls through `OutboundRoutes`, which rejects -/// `with_topic_provider` links at build. -#[tokio::test] -async fn test_knx_rejects_a_topic_provider() { - use aimdb_core::connector::TopicProvider; - - struct Fixed; - impl TopicProvider for Fixed { - fn topic(&self, _value: &DimmerValue) -> Option { - Some("1/0/1".into()) - } - } - - let runtime = Arc::new(TokioAdapter::new().unwrap()); - let mut builder = AimDbBuilder::new() - .runtime(runtime) - .with_connector(connector()); - builder.configure::("knx.dimmer.living", |reg| { - reg.buffer(BufferCfg::SingleLatest) - .link_to("knx://1/0/0") - .with_topic_provider(Fixed) - .with_serializer(|_ctx, dimmer: &DimmerValue| Ok(dimmer.to_knx_bytes())) - .finish(); - }); - - let Err(err) = builder.build().await else { - panic!("a topic provider must be rejected"); - }; - assert!( - format!("{err}").contains("with_topic_provider"), - "unexpected error: {err}" - ); -} - #[tokio::test] async fn test_knx_topic_resolver_with_connector_registration() { let runtime = Arc::new(TokioAdapter::new().unwrap()); diff --git a/aimdb-mqtt-connector/Cargo.toml b/aimdb-mqtt-connector/Cargo.toml index 8680be7f..185a36b5 100644 --- a/aimdb-mqtt-connector/Cargo.toml +++ b/aimdb-mqtt-connector/Cargo.toml @@ -13,9 +13,8 @@ categories = ["network-programming", "embedded", "asynchronous"] [features] default = ["aimdb-core/alloc"] -# `aimdb-core/connector-session` provides the data-plane `pump_sink`/`pump_source` -# helpers the tokio client builds on (re-exported there; `std` implies it too). -# The `rumqttc` backend, which owns its socket, TLS and reconnect. +# `aimdb-core/connector-session` provides the `StreamDialer`/`Delay` traits +# `MqttConnector`'s transport bounds name. The `rumqttc` backend, which owns its socket, TLS and reconnect. std = [ "aimdb-core/std", "aimdb-core/alloc", diff --git a/aimdb-mqtt-connector/tests/link_ext_tests.rs b/aimdb-mqtt-connector/tests/link_ext_tests.rs index 1e4076f1..6952ebd0 100644 --- a/aimdb-mqtt-connector/tests/link_ext_tests.rs +++ b/aimdb-mqtt-connector/tests/link_ext_tests.rs @@ -7,7 +7,7 @@ #![cfg(feature = "std")] use aimdb_core::buffer::BufferCfg; -use aimdb_core::AimDbBuilder; +use aimdb_core::{AimDbBuilder, InboundDispatch, OutboundRoutes}; use aimdb_data_contracts::{link_codecs, LinkCodec, LinkCodecBuilderExt}; use aimdb_mqtt_connector::{MqttConnector, MqttLinkExt, MqttOutboundLinkExt}; use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; @@ -16,7 +16,6 @@ use std::sync::Arc; #[derive(Clone, Debug, Deserialize, Serialize)] struct Reading { - #[allow(dead_code)] value: f32, } @@ -103,30 +102,29 @@ async fn per_link_codec_preserves_mqtt_extensions_and_wiring() { let (db, _runner) = builder.build().await.expect("build must succeed"); - let outbound = db.collect_outbound_routes("mqtt"); - assert_eq!(outbound.len(), 1); - assert_eq!(outbound[0].topic, "sensors/codec"); - assert_eq!(outbound[0].source.serializer_scratch_capacity(), Some(64)); - assert!(outbound[0] - .config - .contains(&("qos".to_string(), "2".to_string()))); - assert!(outbound[0] - .config - .contains(&("retain".to_string(), "true".to_string()))); + let outbound = OutboundRoutes::new(&db, "mqtt").expect("outbound routes"); + let routes = outbound.routes(); + assert_eq!(routes.len(), 1); + assert_eq!(&*routes[0].default_topic, "sensors/codec"); + assert_eq!(routes[0].payload_capacity, 64); + let options = &routes[0].config.protocol_options; + assert!(options.contains(&("qos".to_string(), "2".to_string()))); + assert!(options.contains(&("retain".to_string(), "true".to_string()))); let id = db.inner().resolve_str("test.reading.codec").unwrap(); let record = db.inner().storage(id).unwrap(); let inbound_config = &record.inbound_connectors()[0].config; assert!(inbound_config.contains(&("qos".to_string(), "0".to_string()))); - let inbound = db - .inbound_router("mqtt", &aimdb_core::ExactGrammar) - .expect("inbound routes"); + let inbound = + InboundDispatch::new(&db, "mqtt", &aimdb_core::ExactGrammar).expect("inbound routes"); assert_eq!(inbound.subscriptions(), [Arc::from("commands/codec")]); let encoded = link_codecs::Postcard::<64> .encode(&Reading { value: 17.5 }) .expect("Postcard encode must succeed"); - inbound - .route("commands/codec", &encoded, &db.runtime_ctx()) - .expect("Postcard ingest must succeed"); + let mut reader = db + .subscribe::("test.reading.codec") + .expect("subscribe"); + inbound.dispatch("commands/codec", &encoded); + assert_eq!(reader.try_recv().expect("Postcard ingest").value, 17.5); } diff --git a/aimdb-mqtt-connector/tests/topic_writer_tests.rs b/aimdb-mqtt-connector/tests/topic_writer_tests.rs index 13e38657..3f5e8b38 100644 --- a/aimdb-mqtt-connector/tests/topic_writer_tests.rs +++ b/aimdb-mqtt-connector/tests/topic_writer_tests.rs @@ -504,7 +504,7 @@ fn test_connector_topic_resolution_with_threshold_writer() { /// Test that verifies the inbound topic resolver is called correctly #[test] fn test_inbound_topic_resolver_simulation() { - // Simulate what inbound_router does for TopicResolver + // Simulate how inbound routes resolve a TopicResolver fn resolve_inbound_topic( default_topic: &str, resolver: Option<&dyn Fn() -> Option>, diff --git a/aimdb-tokio-adapter/tests/outbound_routes.rs b/aimdb-tokio-adapter/tests/outbound_routes.rs index a713225d..ae4b5906 100644 --- a/aimdb-tokio-adapter/tests/outbound_routes.rs +++ b/aimdb-tokio-adapter/tests/outbound_routes.rs @@ -9,7 +9,7 @@ use std::task::{Context, Poll, Waker}; use std::time::Duration; use aimdb_core::buffer::BufferCfg; -use aimdb_core::connector::{ConnectorBuilder, SerializeError, TopicProvider}; +use aimdb_core::connector::{ConnectorBuilder, SerializeError}; use aimdb_core::{AimDb, AimDbBuilder, DbResult, OutboundPayload, OutboundRoutes, RecordRegistrar}; use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; @@ -422,37 +422,7 @@ async fn wakes_from_other_threads_lose_nothing() { } #[tokio::test] -async fn a_topic_provider_link_is_rejected() { - struct Fixed; - impl TopicProvider for Fixed { - fn topic(&self, _: &V) -> Option { - Some("x".into()) - } - } - let db = db(vec![Box::new(|reg| { - reg.buffer(BufferCfg::SingleLatest) - .link_to("test://r0") - .with_topic_provider(Fixed) - .with_serializer(le) - .finish(); - })]) - .await; - let Err(aimdb_core::DbError::InvalidConfiguration { errors }) = - OutboundRoutes::new(&db, "test") - else { - panic!("expected a configuration error"); - }; - assert_eq!(errors.len(), 1); - assert_eq!(errors[0].record_key, "r0"); - assert!( - errors[0].message.contains("with_topic_writer"), - "{}", - errors[0].message - ); -} - -#[tokio::test] -async fn route_info_matches_collect_outbound_routes() { +async fn route_info_carries_topic_config_and_record_index() { let db = db(vec![ Box::new(|reg| { reg.buffer(BufferCfg::SingleLatest); @@ -460,24 +430,24 @@ async fn route_info_matches_collect_outbound_routes() { spmc(1, 16), Box::new(|reg| { reg.buffer(BufferCfg::Mailbox) - .link_to("test://two?qos=1") + .link_to("test://two") + .with_config("qos", "1") .with_serializer(le) .finish(); }), ]) .await; let o = OutboundRoutes::new(&db, "test").unwrap(); - let old = db.collect_outbound_routes("test"); - assert_eq!(o.routes().len(), old.len()); - for (info, route) in o.routes().iter().zip(&old) { - let config = aimdb_core::transport::ConnectorConfig::from_query(&route.config); - assert!(info.config.record_index.is_some()); - assert_eq!(info.config.record_index, config.record_index); - assert_eq!(&*info.default_topic, route.topic); - assert_eq!(info.config.protocol_options, config.protocol_options); - } - assert_eq!(o.routes()[0].config.record_index, Some(1)); - assert_eq!(o.routes()[1].config.record_index, Some(2)); + let routes = o.routes(); + assert_eq!(routes.len(), 2); + assert_eq!(&*routes[0].default_topic, "r1"); + assert_eq!(routes[0].config.record_index, Some(1)); + assert_eq!(&*routes[1].default_topic, "two"); + assert_eq!(routes[1].config.record_index, Some(2)); + assert_eq!( + routes[1].config.protocol_options, + [("qos".to_string(), "1".to_string())] + ); } #[tokio::test] diff --git a/aimdb-websocket-connector/tests/decouple_record_keys_topics.rs b/aimdb-websocket-connector/tests/decouple_record_keys_topics.rs index 4d16b13e..2727200a 100644 --- a/aimdb-websocket-connector/tests/decouple_record_keys_topics.rs +++ b/aimdb-websocket-connector/tests/decouple_record_keys_topics.rs @@ -382,7 +382,7 @@ async fn test_fixture( ) } GrantType::Injected => { - // Grant for TopicProvider tests + // Grant for topic writer tests let perms = Permissions { read_patterns: vec!["injected.granted".to_string()], write_patterns: vec![], @@ -468,7 +468,7 @@ async fn test_fixture( }); }); - // Extra keys for TopicProvider test + // Extra keys for the topic writer test for key in ["injected.granted", "injected.denied"] { sb.configure::(key, |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 64 }) // don't coalesce successive values @@ -754,7 +754,7 @@ async fn client_wildcard_subscription_receives_public_only() { } #[tokio::test] -async fn topic_provider_injects_for_unsubscribed_client() { +async fn topic_writer_injects_for_unsubscribed_client() { let addr = free_addr(); let server_db = test_fixture(addr, GrantType::Injected, None).await; diff --git a/aimdb-websocket-connector/tests/e2e.rs b/aimdb-websocket-connector/tests/e2e.rs index c6c8e0b8..8f712c81 100644 --- a/aimdb-websocket-connector/tests/e2e.rs +++ b/aimdb-websocket-connector/tests/e2e.rs @@ -6,7 +6,7 @@ //! `run_client` + [`WsDialer`] engine). Server→client data is pushed by //! *producing a record* — an "injector" record whose dynamic topic + raw //! serializer let a test broadcast an arbitrary `(topic, payload)` through the -//! real `pump_sink` → bus → session path. +//! real `OutboundRoutes` → bus → session path. //! //! The parity block at the bottom locks the AimX WS wire to the semantics the //! retired ws-protocol offered (subscribe ack, wildcard fan-out, late-join @@ -879,7 +879,7 @@ async fn stalled_client_does_not_block_a_healthy_one() { tokio::time::sleep(Duration::from_millis(100)).await; // let the stalled sub register // Flood well past the bounded funnel (256). This also overruns the injector - // ring, so the outbound `pump_sink` consumer lags — it must skip the gap and + // ring, so the outbound route lags — it must skip the gap and // keep publishing (not die), while the stalled client's pump drops on overflow // and the healthy client keeps up. for i in 0..2000u32 { diff --git a/examples/tokio-knx-connector-demo/Cargo.toml b/examples/tokio-knx-connector-demo/Cargo.toml index 59f39adc..3d96256c 100644 --- a/examples/tokio-knx-connector-demo/Cargo.toml +++ b/examples/tokio-knx-connector-demo/Cargo.toml @@ -30,9 +30,6 @@ knx-connector-demo-common = { path = "../knx-connector-demo-common", features = aimdb-knx-connector = { path = "../../aimdb-knx-connector", features = [ # The host leg: the runtime-neutral connector plus core's `std`. "std", - # The connector's channels are `CriticalSectionRawMutex`; only the final - # binary may pick the impl they need to link. - "critical-section-std-impl", "tracing", ] }