Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
0ac1fea
feat(mqtt): implement wildcard subscriptions and enhance QoS handling
lxsaah Sep 27, 2026
9026ffb
docs(design): 055 moves topic matching into connectors
lxsaah Sep 28, 2026
be5f70b
feat(core): topic pattern syntax and grammar traits (055)
lxsaah Sep 28, 2026
1f11b51
feat(core): per-record inbound key table (055)
lxsaah Sep 28, 2026
d08e434
feat(core): pattern routes, TopicMatch and subscriptions in Router (055)
lxsaah Sep 28, 2026
cca20b6
feat(core): with_match_deserializer and .key() on inbound links (055)
lxsaah Sep 28, 2026
2ae708f
feat(core): implement inbound key tables and enhance router functiona…
lxsaah Sep 28, 2026
e3ae290
feat(core): refactor client pump functions to use router for inbound …
lxsaah Sep 28, 2026
46361b5
refactor(core)!: pump_client takes the inbound router
lxsaah Sep 28, 2026
856ce0f
refactor(core)!: remove the grammar-less inbound route API (055)
lxsaah Sep 28, 2026
8fcec93
refactor(core)!: remove the grammar-less inbound route API (055)
lxsaah Sep 28, 2026
d564051
feat(mqtt): MqttGrammar (055)
lxsaah Sep 28, 2026
32aafef
feat(mqtt): route inbound links through MqttGrammar (055)
lxsaah Sep 28, 2026
9bd2546
chore: wrap up 055 wildcard inbound links
lxsaah Sep 28, 2026
12df4ce
fix(mqtt): ensure captures are whole levels in MqttGrammar
lxsaah Sep 29, 2026
703f421
docs: clarify dropped message behavior in metadata and design documents
lxsaah Sep 29, 2026
7586a2f
fix(docs): correct example in builder documentation for inbound router
lxsaah Sep 29, 2026
0ad07f7
fix(mqtt): update error messages for outbound topic pattern restrictions
lxsaah Sep 29, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,30 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added — Design 055: wildcard inbound links

- **One inbound link can feed many topics into one record.**
`link_from("mqtt://sensors/{device}/temp")` matches every device;
`with_match_deserializer(|ctx, m, bytes| …)` sees the topic and its captures
(`m.get("device")`), and `.key("device", 1024)` turns a capture into a small
`KeyId` from a bounded table per record, reported in record metadata
(`inbound_keys`). Routing a pattern, and a known key, allocates nothing.
([aimdb-core](aimdb-core/CHANGELOG.md))
- **MQTT understands topic filters.** `MqttGrammar` implements MQTT 3.1.1
§4.7; both backends subscribe only filters no other filter covers, so the
MQTT 3.1.1 and MQTT 5 backends receive an overlapping topic once each, and a
hand-written `+`/`#` topic now matches.
Comment on lines +41 to +44

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

(nit) This still says both backends "receive an overlapping topic once each", which now contradicts 055 §3.3's known difference: partly overlapping filters (r/+/c and r/b/+) are delivered once per filter on the embedded (MQTT 5) backend. I re-confirmed on Mosquitto 2 at 0ad07f7: native ingests r/b/c 2×, embedded 4×. Suggest "a covered topic once each; partly overlapping filters are a known difference (055)". The same sentence is in aimdb-mqtt-connector/CHANGELOG.md (the "Both backends route through" entry).


Generated by Claude Code

([aimdb-mqtt-connector](aimdb-mqtt-connector/CHANGELOG.md))

### Changed (breaking) — one inbound path

Every connector builds its router with `AimDb::inbound_router(scheme,
grammar)`; `collect_inbound_routes`, `RouterBuilder`, `Route` and the public
`Router::new` are gone, `IngestFn` receives the `TopicMatch`, and `pump_source`
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))

## [2.0.0] - 2026-09-18

### Added
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions aimdb-bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -154,3 +154,5 @@ critical-section = { version = "1.1", features = ["std"] }
defmt = { workspace = true }
embassy-time-driver = "0.2.2"
postcard = { version = "1.0", default-features = false, features = ["alloc"] }
# `MqttGrammar` for the pattern rows of b0_alloc_connector (no backend).
aimdb-mqtt-connector = { path = "../aimdb-mqtt-connector", default-features = false }
126 changes: 110 additions & 16 deletions aimdb-bench/benches/b0_alloc_connector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@
//! Baseline for design 054. Measures what AimDB's connector interfaces cost
//! per message, independent of any real transport:
//!
//! - **Inbound:** `Router::route` alone, and the real `pump_source` driven by
//! - **Inbound:** `Router::route` alone — 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).
Expand Down Expand Up @@ -33,10 +35,10 @@ use aimdb_core::connector::{
ConnectorBuilder, SerializeError, SerializedPayload, SerializedReader, SerializedValueInto,
TopicProvider,
};
use aimdb_core::router::RouterBuilder;
use aimdb_core::session::{pump_source, Payload, Source};
use aimdb_core::transport::{Connector, ConnectorConfig, PublishError};
use aimdb_core::{AimDb, AimDbBuilder, BoxFut, DbResult, RuntimeContext, StringKey};
use aimdb_core::{AimDb, AimDbBuilder, BoxFut, DbResult, ExactGrammar, RuntimeContext, StringKey};
use aimdb_mqtt_connector::MqttGrammar;
use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt};

#[global_allocator]
Expand All @@ -47,11 +49,16 @@ const WARMUP_ITERS: usize = 100;
const MEASURE_ITERS: usize = 2_000;
const DECOY_ROUTES: usize = 63;
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`.
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),
Expand Down Expand Up @@ -139,8 +146,20 @@ async fn build_db(configure: impl FnOnce(&mut AimDbBuilder)) -> AimDb {

// --- Inbound ----------------------------------------------------------------

/// One linked record plus `DECOY_ROUTES` others, so routing scans a realistic
/// table.
/// `DECOY_ROUTES` exact links, so routing scans a realistic table.
fn add_decoys(b: &mut AimDbBuilder) {
for i in 0..DECOY_ROUTES {
let topic = format!("bench://in/decoy/{i}");
b.configure::<Reading>(StringKey::intern(format!("in.decoy{i}")), |reg| {
reg.buffer(BufferCfg::SingleLatest)
.link_from(&topic)
.with_deserializer(|_ctx, _bytes| Ok(reading(0)))
.finish();
});
}
}

/// One linked record plus the decoys.
async fn inbound_db() -> AimDb {
build_db(|b| {
b.configure::<Reading>("in.target", |reg| {
Expand All @@ -149,23 +168,67 @@ async fn inbound_db() -> AimDb {
.with_deserializer(|_ctx, bytes| Ok(reading(bytes[0] as usize)))
.finish();
});
for i in 0..DECOY_ROUTES {
let topic = format!("bench://in/decoy/{i}");
b.configure::<Reading>(StringKey::intern(format!("in.decoy{i}")), |reg| {
reg.buffer(BufferCfg::SingleLatest)
.link_from(&topic)
.with_deserializer(|_ctx, _bytes| Ok(reading(0)))
.finish();
});
}
add_decoys(b);
})
.await
}

/// One record on `in/{device}/target`, keyed or not, plus the decoys.
async fn pattern_db(keyed: bool) -> AimDb {
build_db(|b| {
b.configure::<Reading>("in.pattern", move |reg| {
let link = reg
.buffer(BufferCfg::SpmcRing { capacity: 64 })
.link_from("bench://in/{device}/target");
let link = if keyed {
link.key("device", KEY_CAPACITY)
} else {
link
};
link.with_match_deserializer(|_ctx, m, bytes| {
Ok(reading(
bytes[0] as usize + m.key().map_or(0, |k| k.index()),
))
})
.finish();
});
add_decoys(b);
})
.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<String> {
(0..n)
.map(|i| {
let device = if distinct { first_device + i } else { 0 };
format!("in/dev{device}/target")
})
.collect()
}

async fn measure_route() -> (u64, u64) {
let db = inbound_db().await;
let ctx = db.runtime_ctx();
let router = RouterBuilder::from_routes(db.collect_inbound_routes(SCHEME)).build();
let router = db.inbound_router(SCHEME, &ExactGrammar).unwrap();
let payload = [1u8; 8];
for _ in 0..WARMUP_ITERS {
router.route("in/target", &payload, &ctx).unwrap();
Expand Down Expand Up @@ -207,7 +270,8 @@ async fn pump_run(db: &AimDb, messages: usize) -> (u64, u64) {
remaining: messages,
};
reset();
for fut in pump_source(db, SCHEME, source) {
let router = db.inbound_router(SCHEME, &ExactGrammar).unwrap();
for fut in pump_source(db, router, source) {
fut.await;
}
snapshot()
Expand Down Expand Up @@ -318,6 +382,36 @@ 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",
Expand Down
27 changes: 27 additions & 0 deletions aimdb-bench/data/baselines/b0_alloc_connector.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,33 @@
"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",
Expand Down
33 changes: 33 additions & 0 deletions aimdb-core/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,41 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added

- **Topic patterns on inbound links (design 055).** `{name}` captures one
level, `{name..}` the rest; the syntax is checked at `build()`, and the
connector's `TopicGrammar` compiles each pattern into a `TopicFilter` when it
builds. `ExactGrammar` is the grammar for connectors without wildcards.
- **`InboundConnectorBuilder::with_match_deserializer`** passes a `TopicMatch`
(`topic()`, `get(name)`, `key()`) borrowed from the router's stack, and a
borrowed `&RuntimeContext`, so no reference count changes per message.
- **`.key(name, capacity)`** interns a capture into a `KeyId` from one table
per record, shared by all its keyed links. The table grows as values arrive;
when full, the message is dropped and counted. `AimDb::inbound_key_name`
resolves a key; `RecordMetadata::inbound_keys` (`InboundKeysInfo`) reports
captures, capacity, assigned and dropped.
- **`AimDb::inbound_router(scheme, grammar)`** compiles a scheme's links,
including patterns a `TopicResolverFn` returns, and reports every link it
cannot compile at once. `Router::subscriptions()` lists the filters to
subscribe, without those another filter covers.

### Changed (breaking, API)

- **One inbound path.** Removed `AimDb::collect_inbound_routes`,
`RouterBuilder`, `Route` and the public `Router::new`; a `Router` comes from
`inbound_router`. `IngestFn` takes the `TopicMatch`, and
`InboundConnectorLink` has one `ingest_factory` of that type.
`pump_source(db, router, src)` and `pump_client(db, scheme, router, handle)`
take the router, so a connector subscribes and routes with the same one.
- **`InboundConnectorLink` gains `key`, `RecordMetadata` gains
`inbound_keys`; both are now `#[non_exhaustive]`.** The serde form of
`RecordMetadata` stays backward compatible.
- **Outbound links reject `{…}` topics** at `build()`: a filter cannot be
published to.
- **`{` and `}` in topics are pattern syntax** on every connector, with no
escape: a topic containing a literal brace can no longer be linked.

- **`ConnectorConfig` gains `record_index: Option<usize>`**, the id of the
record an outbound publish comes from — its registration index, the same
`record_id` `AimDb::list_records` reports. A topic alone cannot identify the
Expand Down
Loading
Loading