From f087471194555109be1011c9ad90849e5badceed Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 15:38:26 +0000 Subject: [PATCH 1/2] docs(design): revise 055 after review Match reaches user code only via with_match_deserializer as a borrowed TopicMatch<'a> (no RuntimeContext change, no per-message allocation). Split validation between build() and connector build; add a single AimDb::inbound_router used for both subscription and routing; MQTT subscribes the covering set of overlapping filters; key tables are owned by the link; key() returns Option; multi-level captures may be keys; resolver-returned patterns are allowed. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01KD5fpkpzuuiEjpGZMAxCxC --- docs/design/055-wildcard-inbound-links.md | 227 ++++++++++++++-------- 1 file changed, 143 insertions(+), 84 deletions(-) diff --git a/docs/design/055-wildcard-inbound-links.md b/docs/design/055-wildcard-inbound-links.md index 1b9f057d..b0df841d 100644 --- a/docs/design/055-wildcard-inbound-links.md +++ b/docs/design/055-wildcard-inbound-links.md @@ -1,11 +1,12 @@ # 055 — Wildcard inbound links -**Status:** 📝 Proposed +**Status:** 📝 Proposed (revised after review, 2026-09-27) **Scope:** inbound links whose topic is a pattern: matching in -`aimdb-core`'s router, the matched topic and captures available to the -deserializer, optional per-link key interning, and the MQTT grammar in -`aimdb-mqtt-connector`. **Additive:** no existing public signature changes. +`aimdb-core`'s router, the matched topic and captures passed to a +match-aware deserializer, optional per-link key interning, and the MQTT +grammar in `aimdb-mqtt-connector`. **Additive:** no existing public signature +changes; `RuntimeContext` and `IngestFn` are untouched. **Independent of** [054](./054-zero-alloc-connector-boundary.md). This design ships on today's interfaces; §6 describes what changes when 054 lands. @@ -46,6 +47,7 @@ Two facts make a small change sufficient: first time a value is seen, with a bounded table per link. - G4. No change to existing links, routers or connectors that do not use patterns. +- G5. No per-message allocation on pattern routes, and none for a known key. **Non-goals** @@ -55,6 +57,7 @@ Two facts make a small change sufficient: - The AimX/WebSocket wildcard grammar over record keys (`session/topic_match.rs`). It matches dot-separated record keys, not broker topics. +- Making the match visible to plain `with_deserializer` closures (§4.3). ## 3. User API @@ -62,10 +65,10 @@ Two facts make a small change sufficient: builder.configure::("sensors.readings", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 256 }) .link_from("mqtt://sensors/{device}/temp") - .key("device", 1024) // optional (§4.4) + .key("device", 1024) // optional (§4.4) .with_match_deserializer(|ctx, m, bytes| { - let device = m.get("device").unwrap(); // borrowed &str - let key = m.key(); // KeyId, when .key(..) is set + let device = m.get("device").unwrap(); // &str borrowed from the topic + let key = m.key(); // Option; Some when .key(..) is set Reading::decode(key, bytes) }) .finish(); @@ -76,16 +79,17 @@ builder.configure::("sensors.readings", |reg| { (zero or more) and must be last. - The grammar's own wildcards (`+`, `#` for MQTT) are also accepted as unnamed levels, for users who write filters by hand. -- `with_deserializer(|ctx, bytes|)` keeps working on pattern links; the - match is available through `ctx.inbound_match()` (§4.3). +- `with_match_deserializer` is the only way to see the match. A pattern link + with a plain `with_deserializer(|ctx, bytes|)` still works (every matching + message is ingested) but cannot tell publishers apart. ## 4. Design ### 4.1 Grammar is chosen by the connector Wildcard syntax is protocol-specific (MQTT `/` `+` `#`; Zenoh `*` `**`; KNX -none), and the router is protocol-agnostic. Core defines the grammar -interface; connectors supply an implementation: +none), and the router is protocol-agnostic. Core defines the grammar as a +plain data struct; connectors supply a value: ```rust // aimdb-core @@ -101,36 +105,62 @@ pub struct TopicGrammar { impl TopicGrammar { /// No wildcards: every resource id is literal (today's behaviour). + /// A pattern link on an EXACT connector is a configuration error. pub const EXACT: Self = /* … */; } -impl RouterBuilder { - pub fn with_grammar(self, grammar: TopicGrammar) -> Self; +impl AimDb { + /// The one inbound router for `scheme`: pattern links compiled against + /// `grammar`, key tables attached (§4.4). Both subscription and routing + /// use this router, so they cannot disagree. + pub fn inbound_router(&self, scheme: &str, grammar: TopicGrammar) + -> DbResult; } pub fn pump_source_with(db: &AimDb, scheme: &str, src: impl Source + 'static, - grammar: TopicGrammar) -> Vec; -// pump_source(..) == pump_source_with(.., TopicGrammar::EXACT) + grammar: TopicGrammar) -> DbResult>; +// pump_source(..) keeps its signature and uses TopicGrammar::EXACT. // aimdb-mqtt-connector pub const MQTT_GRAMMAR: TopicGrammar = /* '/', "+", "#", leading '$' hidden */; ``` -A plain data struct keeps `Router` non-generic, which is what makes this -additive. The pattern walk itself is generic code in core; the struct only -supplies tokens. +The struct keeps `Router` non-generic, which is what makes this additive. +The pattern walk is one function in core, parameterised at runtime by the +struct's tokens. -### 4.2 Compilation and matching +`collect_inbound_routes` keeps its signature and returns exact links only; +it logs a warning for each pattern link it skips. Out-of-tree connectors that +use it keep working for exact links, and in-tree connectors move to +`inbound_router`. -At `build()` (route collection), each pattern route is parsed once into -levels: `Literal(Arc)`, `Capture { name, multi }` or `Wildcard { multi }`. -Invalid patterns fail `build()` with the record key and URL, like other link +### 4.2 Compilation, validation and matching + +Validation happens in two places, because the grammar is only known when the +connector builds. + +**At `AimDbBuilder::build()`** — the `{…}` syntax, which is the same for +every grammar. Errors name the record key and URL, like other link configuration errors: - a capture sharing a level with text (`sensors/dev-{id}`); -- a multi-level capture or wildcard that is not last; +- a multi-level capture that is not last; - two captures with the same name; -- more than 8 captures. +- more than 8 captures; +- `.key(name, ..)` naming a capture that the pattern does not have, or a + capacity outside `1..=65535`. + +**At connector build** (`inbound_router` / `pump_source_with`, which return +`DbResult`) — everything that needs the grammar or the connector: + +- a hand-written multi-level wildcard (`#`) that is not last; +- any pattern (captures or grammar wildcards) on a connector using + `TopicGrammar::EXACT`, e.g. KNX or the WebSocket server; +- patterns returned by a `TopicResolverFn` (018): the resolved string is + compiled and checked exactly like a URL pattern. + +Each pattern route is parsed once into levels: `Literal(Arc)`, +`Capture { name, multi }` or `Wildcard { multi }`. `Router::route` checks exact routes as today, then pattern routes by walking `topic.split(separator)` against the compiled levels. Capture positions are @@ -139,44 +169,42 @@ matching. A topic matching several routes (exact or pattern) is delivered to each, as today. `Router::resource_ids()` returns each pattern rendered in the grammar -(`sensors/{device}/temp` → `sensors/+/temp`), so connectors subscribe -unchanged. +(`sensors/{device}/temp` → `sensors/+/temp`). -### 4.3 The match reaches the deserializer through `RuntimeContext` +### 4.3 The match reaches the deserializer as a borrow -`IngestFn` is `Fn(&RuntimeContext, &[u8])`; changing it would break a public -type. Instead, `RuntimeContext` carries an optional match, set by the router -for the duration of one ingest call: +`IngestFn` is `Fn(&RuntimeContext, &[u8])` and stays unchanged; exact routes +keep using it. Pattern routes use a second, internal ingest type, and +`Route` holds one or the other: ```rust -impl RuntimeContext { - /// The inbound match being ingested, if any. - pub fn inbound_match(&self) -> Option<&TopicMatch>; -} +// aimdb-core (internal) +type MatchIngestFn = + Arc, &[u8]) -> Result<(), String> + Send + Sync>; -pub struct TopicMatch { - topic: Arc, +// public +pub struct TopicMatch<'a> { + topic: &'a str, spans: [(u16, u16); 8], - names: Arc<[Arc]>, // from the compiled route + names: &'a [Arc], // from the compiled route key: Option, } -impl TopicMatch { - pub fn topic(&self) -> &str; - pub fn get(&self, name: &str) -> Option<&str>; - pub fn key(&self) -> KeyId; // panics if the link has no key +impl<'a> TopicMatch<'a> { + pub fn topic(&self) -> &'a str; + pub fn get(&self, name: &str) -> Option<&'a str>; + pub fn key(&self) -> Option; // Some iff the link has .key(..) } ``` -- `RuntimeContext` gains one private field; `RuntimeContext::new` is - unchanged. -- **Cost:** for pattern routes, the topic is copied into an `Arc` once - per message, because the context cannot borrow it. That is one allocation - per message on pattern routes only (the `b0_alloc_connector` row added by - this design records it). Exact routes are unaffected. 054 removes this - cost (§6). -- `with_match_deserializer(|ctx, m, bytes|)` is convenience over - `with_deserializer` that reads `ctx.inbound_match()`. +- `Router::route(&self, topic: &str, ..)` already borrows the topic for the + whole call, so the router builds `TopicMatch` on its stack and passes a + reference. **No per-message allocation**, and no change to + `RuntimeContext`. +- The match cannot outlive its message: `TopicMatch` borrows the topic, so + a closure cannot keep it. +- `with_match_deserializer(|ctx, m, bytes|)` builds a `MatchIngestFn`; + `with_deserializer` on a pattern link builds one that ignores `m`. ### 4.4 Keys @@ -184,18 +212,24 @@ impl TopicMatch { #[derive(Copy, Clone, Eq, PartialEq, Hash, Debug)] pub struct KeyId(NonZeroU16); // Option is 2 bytes -.key("device", 1024) // capture name, capacity (required) +.key("device", 1024) // capture name, capacity 1..=65535 (required) ``` -- Each keyed link owns a table: `HashMap, KeyId>` (hashbrown, - created with full capacity at `build()`) and `Vec>` for the - reverse direction, under one `spin::Mutex`. Both dependencies are already - in `aimdb-core`. +- **Ownership.** The table belongs to the link and is created at `build()` + as an `Arc`. Every router built from the link holds a clone, so + `KeyId`s are the same everywhere, and `db.inbound_key_name` looks there. +- The table is `HashMap, KeyId>` (hashbrown, created with full + capacity) plus `Vec>` for the reverse direction, under one + `spin::Mutex`. Both dependencies are already in `aimdb-core`. - A value seen for the first time gets the next `KeyId`; storing its name is - one allocation, once. Later messages from the same value only look it up. -- When the table is full, the message is dropped and a per-link counter + one allocation, once. Later messages look it up by `&str` without + allocating. +- A `{name..}` capture can be a key; its value is the whole remainder + (`site/{loc..}` on `site/a/b/c` → key for `a/b/c`). +- **When the table is full**, the message is dropped and a per-link counter increases (logged; visible in record metadata with `observability`). - Capacity is therefore also an admission limit. + Capacity is therefore also an admission limit. A message that reaches the + deserializer of a keyed link always has `m.key() == Some(_)`. - Keys are never reused while the process runs. - `db.inbound_key_name("sensors.readings", key) -> Option>` resolves a key for display, AimX and logs. @@ -210,8 +244,20 @@ measurements show the lock. - `MQTT_GRAMMAR` follows MQTT 3.1.1 §4.7: `+` one level, `#` the rest including the parent level (`a/#` matches `a`), topics starting with `$` are not matched by a leading wildcard. -- Both backends switch from `pump_source` to `pump_source_with(.., - MQTT_GRAMMAR)`. +- Both backends build their subscription list from + `db.inbound_router("mqtt", MQTT_GRAMMAR)`, replacing the separate + `RouterBuilder::from_routes(..)` calls in `native.rs` and + `embedded/mod.rs::inbound_topics`, and switch from `pump_source` to + `pump_source_with(.., MQTT_GRAMMAR)`. +- **Overlapping filters.** Some brokers deliver one copy per matching + subscription, so `sensors/{d}/temp` next to `sensors/kitchen/temp` could + ingest a message twice. The connector subscribes only the **covering set**: + a filter is dropped from the subscription list when another filter matches + every topic it matches. The router still fans out locally, so each record + receives exactly one copy. Coverage is a level-by-level comparison over the + compiled patterns. The subscription QoS of a covering filter is the highest + QoS of the filters it covers (both backends subscribe at a fixed QoS 1 + today, so this only matters once per-link subscribe QoS is honoured). - Outbound links reject patterns at `build()`: you cannot publish to a filter. @@ -232,14 +278,14 @@ everything else from the payload. ## 6. Relation to 054 When 054 lands, `InboundDispatch::dispatch(topic, payload)` borrows the topic -for the whole ingest call. The router then passes a borrowed match instead of -copying the topic into `RuntimeContext`: - -- `with_match_deserializer`'s closure signature is unchanged; `m` becomes a - borrow of the dispatcher's stack, and the per-message allocation on - pattern routes disappears. -- `ctx.inbound_match()` stays available for `with_deserializer` users, set - from the same borrow. +for the whole ingest call, exactly as `Router::route` does today. Because +`TopicMatch<'a>` already borrows the topic: + +- `with_match_deserializer`'s closure signature and `TopicMatch<'a>` are + unchanged. +- `InboundDispatch::new` gains a grammar argument (or a `with_grammar` + variant) and replaces `inbound_router` + `pump_source_with` for migrated + connectors. - The `TopicGrammar` struct may become a trait for inlining if the 054 bench shows the pattern walk. @@ -257,36 +303,49 @@ copying the topic into `RuntimeContext`: larger than this. 5. **Positional captures only (`+` → index 0).** Indices shift when a pattern changes; names cost nothing at runtime. +6. **Match carried in `RuntimeContext` (`ctx.inbound_match()`).** Reaches + plain `with_deserializer` users, but the context cannot borrow the topic, + so it costs one allocation per message, the match can outlive the message + through a cloned context, and it breaks when 054 turns the topic into a + borrow. Rejected in review. +7. **Declaring grammars on the builder** so every check runs at `build()`. + Adds a registration step for every connector; connector-build errors are + early enough. +8. **Allowing overlapping subscriptions** and documenting duplicates, or + **rejecting overlaps** at build. The first gives records duplicate + messages on some brokers; the second rules out a legitimate layout. +9. **Delivering unkeyed messages when the key table is full.** Keeps + messages, but every consumer then has to handle `key() == None` on a + keyed link. Capacity as an admission limit is the simpler contract. ## 8. Open questions -- **Overlapping subscriptions.** With `sensors/{d}/temp` and - `sensors/kitchen/temp` on different records, some brokers deliver one - message per matching subscription, so the router would ingest it twice. - Check Mosquitto 2.x and EMQX; if needed, the MQTT connector subscribes only - the broadest filter of an overlapping set, since the router fans out - locally. - **`mountain-mqtt` subscriptions.** Confirm `subscribe_packet` accepts `+` and `#` unchanged. -- **`TopicResolverFn`.** A resolver (018) may return a pattern; it should be - compiled the same way. Confirm no existing resolver returns strings with - `{`. +- **Broker duplicate behaviour.** The covering set makes it irrelevant for + correctness, but record Mosquitto 2.x and EMQX behaviour in the + integration test notes so the rationale is checked. ## 9. Acceptance criteria 1. Router unit tests: the MQTT §4.7 cases above, captures at first, middle and last level, `{name..}` matching zero levels, exact and pattern routes on one topic both delivering, `TopicGrammar::EXACT` routers unchanged. -2. `build()` rejects each invalid pattern in §4.2 with the record key. -3. Tokio integration test against a local Mosquitto: two clients publish to +2. `build()` rejects each `{…}` error in §4.2 with the record key; + `inbound_router` rejects each connector-build error in §4.2, including a + pattern on an `EXACT` connector and an invalid resolver-returned pattern. +3. Covering-set unit tests: `sensors/+/temp` covers `sensors/kitchen/temp`; + `a/#` covers `a` and `a/+/b`; unrelated filters are all kept. +4. Tokio integration test against a local Mosquitto: two clients publish to `sensors/a/temp` and `sensors/b/temp`; one record receives both; the deserializer sees `device = a` and `b`; distinct `KeyId`s; `db.inbound_key_name` returns the names; a table of capacity 1 drops the - second device and counts it. -4. `b0_alloc_connector` gains `inbound_route_pattern` (1 alloc/msg, the - topic copy) and `inbound_route_keyed_known` (1 alloc/msg for a known - key). Existing rows unchanged. -5. `weather-station-gamma` and the embedded MQTT demo build for + second device and counts it. An exact link on `sensors/a/temp` next to the + pattern link receives each message once, and the pattern record once. +5. `b0_alloc_connector` gains `inbound_route_pattern` (0 allocs/msg), + `inbound_route_keyed_known` (0 allocs/msg) and `inbound_route_keyed_new` + (1 alloc, the key name). Existing rows unchanged. +6. `weather-station-gamma` and the embedded MQTT demo build for `thumbv7em-none-eabihf` with no behaviour change. ## 10. References From 997b5a9562af563604b4d46b1cefe1dfe372f882 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 16:09:26 +0000 Subject: [PATCH 2/2] docs(design): 055 validated by spike; grammar trait, per-record keys Lead with an evaluation of a spike implementation: allocation and latency measurements, both MQTT backends against one broker, Mosquitto delivering overlapping subscriptions once to MQTT 3.1.1 and twice to MQTT 5 clients, and key-table memory. Design changes from that evaluation: TopicGrammar becomes a trait (Zenoh's mid-pattern ** needs it); one key table per record, grown lazily, reported in RecordMetadata::inbound_keys; grammar-dependent checks move to connector build; all in-tree connectors move to inbound_router; compatibility caveat for two public structs gaining fields. Key reuse and a pattern index are non-goals. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01KD5fpkpzuuiEjpGZMAxCxC --- docs/design/055-wildcard-inbound-links.md | 508 +++++++++++++--------- 1 file changed, 311 insertions(+), 197 deletions(-) diff --git a/docs/design/055-wildcard-inbound-links.md b/docs/design/055-wildcard-inbound-links.md index b0df841d..1daf340c 100644 --- a/docs/design/055-wildcard-inbound-links.md +++ b/docs/design/055-wildcard-inbound-links.md @@ -1,15 +1,16 @@ # 055 — Wildcard inbound links -**Status:** 📝 Proposed (revised after review, 2026-09-27) +**Status:** 📝 Proposed — validated by a spike (§3), 2026-09-27 **Scope:** inbound links whose topic is a pattern: matching in `aimdb-core`'s router, the matched topic and captures passed to a -match-aware deserializer, optional per-link key interning, and the MQTT -grammar in `aimdb-mqtt-connector`. **Additive:** no existing public signature -changes; `RuntimeContext` and `IngestFn` are untouched. +match-aware deserializer, optional per-record key interning surfaced in +record metadata, a grammar trait implemented by connectors, and the MQTT +grammar in `aimdb-mqtt-connector`. `RuntimeContext` and `IngestFn` are +untouched; §5.8 lists the two public structs that gain fields. **Independent of** [054](./054-zero-alloc-connector-boundary.md). This -design ships on today's interfaces; §6 describes what changes when 054 lands. +design ships on today's interfaces; §7 describes what changes when 054 lands. --- @@ -44,10 +45,12 @@ Two facts make a small change sufficient: - G2. Named captures (`{device}`) are available to the deserializer as `&str`, together with the full topic. - G3. Optionally, a capture becomes a **key**: a small integer assigned the - first time a value is seen, with a bounded table per link. + first time a value is seen, from a bounded table per record. - G4. No change to existing links, routers or connectors that do not use patterns. - G5. No per-message allocation on pattern routes, and none for a known key. +- G6. Protocols with different wildcard rules (MQTT, Zenoh) plug in without + changing core. **Non-goals** @@ -57,15 +60,67 @@ Two facts make a small change sufficient: - The AimX/WebSocket wildcard grammar over record keys (`session/topic_match.rs`). It matches dot-separated record keys, not broker topics. -- Making the match visible to plain `with_deserializer` closures (§4.3). - -## 3. User API +- Making the match visible to plain `with_deserializer` closures (§5.5). +- Releasing or reusing keys. A key lives as long as the process. +- An index (trie) over pattern routes. Routes are scanned linearly (§5.4). + +## 3. Evaluation + +Before this revision, the design was implemented as a spike on the +in-tree code: core (`topic_pattern.rs`, router, builder, link builder, +metadata), both MQTT backends, the `b0_alloc_connector` bench rows and +end-to-end tests. About 775 changed lines plus 1,250 new ones, most of them +tests. Everything below was measured on that spike, not estimated. + +### 3.1 Results + +| Question | Result | +|---|---| +| Allocations per message: exact / pattern / known key / new key | **0 / 0 / 0 / 1** (the key name, 24 bytes) | +| Routing time, 64 other routes, release build | exact 76–83 ns, pattern with key 132–143 ns, no match 41–47 ns | +| Grammar as a trait (`&'static dyn`) instead of a data struct | same allocations; latency within run-to-run noise | +| Zenoh-style `a/**/b` (multi-level wildcard mid-pattern) | works through the trait, with a backtracking matcher; not expressible with the data struct | +| Both MQTT backends against one broker, pattern next to a covered exact link | each subscribes only `parity/+/in`; each record gets the message once; capture and key reach the deserializer | +| `mountain-mqtt` with `+` in SUBSCRIBE | accepted unchanged | +| Mosquitto 2.0.18, overlapping subscriptions, one publish | **MQTT 3.1.1 client: 1 copy. MQTT 5 client: 2 copies.** | +| Key table memory | about 66 bytes per key plus the name; 1,024 keys ≈ 67 KB, 65,535 keys ≈ 4.3 MB | +| Key table growing lazily instead of reserving capacity | still 1 allocation per new key on average (log₂ n extra in total) | +| thumbv7em-none-eabihf: core (`alloc`, `connector-session`, `remote`) and the embedded MQTT backend | builds | + +### 3.2 What the evaluation changed + +1. **The covering set is required, not an optimisation.** The native backend + speaks MQTT 3.1.1 and the embedded one MQTT 5, and Mosquitto sends + overlapping subscriptions one copy and two copies respectively. Without + the covering set (§5.7) the backends would disagree. +2. **The grammar is a trait.** Zenoh allows `**` anywhere and hides + verbatim `@…` chunks from wildcards at any level. A data struct with + tokens and a "must be last" rule cannot say that (§5.1). +3. **Fewer checks run at `build()`.** "Capture shares a level with text" + needs the separator, and "multi-level capture must be last" depends on + the grammar. Both run when the connector builds (§5.3). +4. **Hand-written `+`/`#` stay broken on the old path.** Whether `+` is a + wildcard depends on the grammar, so `collect_inbound_routes` cannot spot + such links. Every in-tree connector moves to `inbound_router` (§5.2). +5. **Keys are per record.** With per-link tables, two keyed links on one + record handed out overlapping `KeyId`s. One table per record fixes that + and matches the "array indexed by `KeyId`" guidance (§5.6). +6. **Key tables grow lazily.** Reserving full capacity cost 67 KB for 1,024 + keys before any device appeared; lazy growth costs nothing measurable on + the per-message path (§5.6). +7. **Two public structs gain fields** (§5.8), so "additive" needs a caveat. +8. **Rules the first draft left implicit:** setting both deserializers is an + error; a literal topic with a match-aware deserializer takes the pattern + path; the dropped counter is `AtomicU32` (thumbv7em has no 64-bit + atomics). + +## 4. User API ```rust builder.configure::("sensors.readings", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 256 }) .link_from("mqtt://sensors/{device}/temp") - .key("device", 1024) // optional (§4.4) + .key("device", 1024) // optional (§5.6) .with_match_deserializer(|ctx, m, bytes| { let device = m.get("device").unwrap(); // &str borrowed from the topic let key = m.key(); // Option; Some when .key(..) is set @@ -75,120 +130,139 @@ builder.configure::("sensors.readings", |reg| { }); ``` -- `{name}` matches exactly one level. `{name..}` matches the remaining levels - (zero or more) and must be last. +- `{name}` matches exactly one level. `{name..}` matches zero or more + levels; where it may stand is the grammar's rule (MQTT: last only). - The grammar's own wildcards (`+`, `#` for MQTT) are also accepted as unnamed levels, for users who write filters by hand. - `with_match_deserializer` is the only way to see the match. A pattern link with a plain `with_deserializer(|ctx, bytes|)` still works (every matching - message is ingested) but cannot tell publishers apart. + message is ingested) but cannot tell publishers apart. Setting both on one + link is a configuration error. -## 4. Design +## 5. Design -### 4.1 Grammar is chosen by the connector +### 5.1 The grammar is a connector-supplied trait -Wildcard syntax is protocol-specific (MQTT `/` `+` `#`; Zenoh `*` `**`; KNX -none), and the router is protocol-agnostic. Core defines the grammar as a -plain data struct; connectors supply a value: +Wildcard syntax is protocol-specific and the router is protocol-agnostic. +Core owns one matcher; a grammar only answers questions about single levels: ```rust // aimdb-core -#[derive(Clone, Copy)] -pub struct TopicGrammar { - pub separator: char, - /// The grammar's single-level and multi-level wildcard tokens. - pub single: &'static str, - pub multi: &'static str, - /// Topics that wildcards must not match (MQTT: leading `$`). - pub hidden: fn(&str) -> bool, -} - -impl TopicGrammar { - /// No wildcards: every resource id is literal (today's behaviour). - /// A pattern link on an EXACT connector is a configuration error. - pub const EXACT: Self = /* … */; -} - -impl AimDb { - /// The one inbound router for `scheme`: pattern links compiled against - /// `grammar`, key tables attached (§4.4). Both subscription and routing - /// use this router, so they cannot disagree. - pub fn inbound_router(&self, scheme: &str, grammar: TopicGrammar) - -> DbResult; +pub enum LevelKind { Literal, Single, Multi, Invalid(&'static str) } + +pub trait TopicGrammar: Send + Sync { + /// `false`: every `{…}` link on this connector is an error. + fn supports_patterns(&self) -> bool { true } + fn separator(&self) -> char; + /// What a hand-written level is (`+` → Single, `a+` → Invalid, …). + fn classify(&self, level: &str) -> LevelKind; + /// Tokens used to render `{name}` / `{name..}` for the subscription. + fn single_token(&self) -> &str; + fn multi_token(&self) -> &str; + /// Zenoh `a/**/b`: true. MQTT `#`: false (last only). + fn multi_anywhere(&self) -> bool { false } + /// Whether a wildcard at level `index` may match `level`. + /// MQTT: not a leading `$…`. Zenoh: not a verbatim `@…` chunk. + fn wildcard_matches(&self, index: usize, level: &str) -> bool { true } } -pub fn pump_source_with(db: &AimDb, scheme: &str, src: impl Source + 'static, - grammar: TopicGrammar) -> DbResult>; -// pump_source(..) keeps its signature and uses TopicGrammar::EXACT. +/// KNX, WebSocket, AimX session connectors: no wildcards. +pub struct ExactGrammar; // aimdb-mqtt-connector -pub const MQTT_GRAMMAR: TopicGrammar = /* '/', "+", "#", leading '$' hidden */; +pub struct MqttGrammar; // '/', "+", "#", `$` hidden at level 0 ``` -The struct keeps `Router` non-generic, which is what makes this additive. -The pattern walk is one function in core, parameterised at runtime by the -struct's tokens. +- Grammars are passed as `&'static dyn TopicGrammar` (unit structs: + `&MqttGrammar`). `Router` stays non-generic and the per-match dynamic + calls are one `separator()` at compile time and one + `wildcard_matches()` per wildcard level; §3.1 shows no measurable cost. +- Features outside the trait are rejected by `classify` with a reason, e.g. + Zenoh's `$*` sub-chunk wildcards. Supporting them later is a trait + addition, not a change to the router. -`collect_inbound_routes` keeps its signature and returns exact links only; -it logs a warning for each pattern link it skips. Out-of-tree connectors that -use it keep working for exact links, and in-tree connectors move to -`inbound_router`. +### 5.2 One inbound router per connector -### 4.2 Compilation, validation and matching +```rust +impl AimDb { + /// Exact links as `collect_inbound_routes`, plus pattern links compiled + /// against `grammar`, with key tables attached. + pub fn inbound_router(&self, scheme: &str, grammar: &'static dyn TopicGrammar) + -> DbResult; +} -Validation happens in two places, because the grammar is only known when the -connector builds. +impl Router { + /// Subscription filters with any filter covered by another removed. + pub fn covering_resource_ids(&self) -> Vec>; +} -**At `AimDbBuilder::build()`** — the `{…}` syntax, which is the same for -every grammar. Errors name the record key and URL, like other link -configuration errors: +pub fn pump_source_with(db: &AimDb, scheme: &str, src: impl Source + 'static, + grammar: &'static dyn TopicGrammar) + -> DbResult>; +// pump_source(..) keeps its signature and behaviour. +``` -- a capture sharing a level with text (`sensors/dev-{id}`); -- a multi-level capture that is not last; +- A connector subscribes and routes with the **same** router, so the two + cannot disagree. The spike replaced the separate + `RouterBuilder::from_routes(..)` calls in `native.rs` and + `embedded/mod.rs::inbound_topics`. +- **Every in-tree connector moves to `inbound_router`** in the same change: + MQTT with `&MqttGrammar`; KNX, the WebSocket server and client, and + core's AimX session client (TCP, UDS, serial) with `&ExactGrammar`. Only then is a `{…}` link on those + connectors an error instead of a warning, and a hand-written `+` on MQTT + starts working. +- `collect_inbound_routes` keeps its signature for out-of-tree connectors. + It returns exact links only and logs a warning for each `{…}` link it + skips. A hand-written `+`/`#` link without braces still goes through it + as an exact route: it subscribes and never matches, which is today's + behaviour. + +### 5.3 Validation + +**At `AimDbBuilder::build()`**: the grammar-independent `{…}` syntax. +Errors name the record key and URL: + +- unbalanced braces, empty names, names outside `[A-Za-z0-9_]`; - two captures with the same name; - more than 8 captures; -- `.key(name, ..)` naming a capture that the pattern does not have, or a - capacity outside `1..=65535`. +- `.key(name, ..)` naming a capture that the pattern does not have, a + capacity of 0, or a capacity that differs from another keyed link on the + same record (§5.6). -**At connector build** (`inbound_router` / `pump_source_with`, which return -`DbResult`) — everything that needs the grammar or the connector: +**At connector build** (`inbound_router` / `pump_source_with`, returning +`DbResult`): everything that needs the grammar: + +- a capture sharing a level with text (`sensors/dev-{id}`); +- a multi-level capture or wildcard where the grammar forbids it; +- a level `classify` rejects (`a+` in MQTT, `$*` in Zenoh); +- any `{…}` on a connector whose grammar has `supports_patterns() == false`; +- patterns returned by a `TopicResolverFn` (018), compiled and checked like + URL patterns. -- a hand-written multi-level wildcard (`#`) that is not last; -- any pattern (captures or grammar wildcards) on a connector using - `TopicGrammar::EXACT`, e.g. KNX or the WebSocket server; -- patterns returned by a `TopicResolverFn` (018): the resolved string is - compiled and checked exactly like a URL pattern. +### 5.4 Matching -Each pattern route is parsed once into levels: `Literal(Arc)`, -`Capture { name, multi }` or `Wildcard { multi }`. +Each pattern route is compiled once into levels: `Literal`, `Single` or +`Multi`, each wildcard optionally carrying a capture slot. -`Router::route` checks exact routes as today, then pattern routes by walking -`topic.split(separator)` against the compiled levels. Capture positions are -recorded as byte ranges into the topic in a fixed array; no allocation for -matching. A topic matching several routes (exact or pattern) is delivered to -each, as today. +`Router::route` checks exact routes as today, then pattern routes in +registration order. The matcher walks the topic in place, recording capture +positions as byte ranges in a fixed `[(u16, u16); 8]`. A `Multi` level that +is last takes the rest in one step; one followed by more levels (Zenoh) +backtracks over split points. A topic matching several routes is delivered +to each, as today. -`Router::resource_ids()` returns each pattern rendered in the grammar -(`sensors/{device}/temp` → `sensors/+/temp`). +Routes are scanned linearly. At 64 routes a pattern route costs about 60 ns +over an exact one (§3.1). An index can come later if a benchmark asks for it. -### 4.3 The match reaches the deserializer as a borrow +### 5.5 The match reaches the deserializer as a borrow -`IngestFn` is `Fn(&RuntimeContext, &[u8])` and stays unchanged; exact routes -keep using it. Pattern routes use a second, internal ingest type, and -`Route` holds one or the other: +Exact routes keep `IngestFn`. Pattern routes use a second ingest type: ```rust -// aimdb-core (internal) -type MatchIngestFn = +pub type MatchIngestFn = Arc, &[u8]) -> Result<(), String> + Send + Sync>; -// public -pub struct TopicMatch<'a> { - topic: &'a str, - spans: [(u16, u16); 8], - names: &'a [Arc], // from the compiled route - key: Option, -} +pub struct TopicMatch<'a> { /* topic: &'a str, spans, names, key */ } impl<'a> TopicMatch<'a> { pub fn topic(&self) -> &'a str; @@ -198,76 +272,104 @@ impl<'a> TopicMatch<'a> { ``` - `Router::route(&self, topic: &str, ..)` already borrows the topic for the - whole call, so the router builds `TopicMatch` on its stack and passes a - reference. **No per-message allocation**, and no change to - `RuntimeContext`. -- The match cannot outlive its message: `TopicMatch` borrows the topic, so - a closure cannot keep it. -- `with_match_deserializer(|ctx, m, bytes|)` builds a `MatchIngestFn`; - `with_deserializer` on a pattern link builds one that ignores `m`. - -### 4.4 Keys + whole call, so the router builds `TopicMatch` on its stack. No allocation, + no change to `RuntimeContext`, and the match cannot outlive its message. +- `with_match_deserializer` builds a `MatchIngestFn`. A `{…}` link with a + plain `with_deserializer` gets one that ignores the match. A **literal** + topic with `with_match_deserializer` becomes an all-literal pattern route + so the closure still receives the topic. +- The closure takes `RuntimeContext` by value, like `with_deserializer`. + That clones an `Arc` per message (an atomic increment, no allocation). + Passing `&RuntimeContext` would be cheaper but is a change for both + builders; it belongs in 054's breaking window. + +### 5.6 Keys ```rust -#[derive(Copy, Clone, Eq, PartialEq, Hash, Debug)] -pub struct KeyId(NonZeroU16); // Option is 2 bytes +pub struct KeyId(NonZeroU16); // Option is 2 bytes; .index() is 0-based -.key("device", 1024) // capture name, capacity 1..=65535 (required) +.key("device", 1024) // capture name, capacity 1..=65535 ``` -- **Ownership.** The table belongs to the link and is created at `build()` - as an `Arc`. Every router built from the link holds a clone, so - `KeyId`s are the same everywhere, and `db.inbound_key_name` looks there. -- The table is `HashMap, KeyId>` (hashbrown, created with full - capacity) plus `Vec>` for the reverse direction, under one - `spin::Mutex`. Both dependencies are already in `aimdb-core`. -- A value seen for the first time gets the next `KeyId`; storing its name is - one allocation, once. Later messages look it up by `&str` without - allocating. -- A `{name..}` capture can be a key; its value is the whole remainder - (`site/{loc..}` on `site/a/b/c` → key for `a/b/c`). -- **When the table is full**, the message is dropped and a per-link counter - increases (logged; visible in record metadata with `observability`). - Capacity is therefore also an admission limit. A message that reaches the +- **One table per record.** All keyed links of a record share it, so a + `KeyId` names one value per record, and the same device seen through + `temp/{dev}` and `hum/{id}` gets the same key. Keyed links on one record + must give the same capacity. +- The table is `HashMap, KeyId>` (hashbrown) plus + `Vec>` for the reverse direction, under one `spin::Mutex`. Both + dependencies are already in `aimdb-core`. +- **It grows as keys arrive**; capacity is a limit, not a reservation. + Memory is about 66 bytes per key plus the name (§3.1). +- A new value costs one allocation (its name); a known value costs none. +- A `{name..}` capture can be a key; its value is the whole remainder. +- **When the table is full**, the message is dropped and the table's + `dropped` counter (`AtomicU32`) increases. A message that reaches the deserializer of a keyed link always has `m.key() == Some(_)`. - Keys are never reused while the process runs. -- `db.inbound_key_name("sensors.readings", key) -> Option>` resolves - a key for display, AimX and logs. +- `db.inbound_key_name("sensors.readings", key) -> Option>` + resolves a key for display, AimX and logs. -The lock is taken once per message by the single task that runs routing for -the connector, so it is uncontended. 054's `InboundDispatch` does not remove -it by itself; moving the table to `&mut` ownership is a follow-up if -measurements show the lock. +**Record metadata.** `RecordMetadata` gains -### 4.5 MQTT specifics +```rust +#[serde(default, skip_serializing_if = "Option::is_none")] +pub inbound_keys: Option, + +pub struct InboundKeysInfo { + pub captures: Vec, // one per keyed link + pub capacity: u16, + pub assigned: usize, + pub dropped: u32, +} +``` -- `MQTT_GRAMMAR` follows MQTT 3.1.1 §4.7: `+` one level, `#` the rest - including the parent level (`a/#` matches `a`), topics starting with `$` - are not matched by a leading wildcard. -- Both backends build their subscription list from - `db.inbound_router("mqtt", MQTT_GRAMMAR)`, replacing the separate - `RouterBuilder::from_routes(..)` calls in `native.rs` and - `embedded/mod.rs::inbound_topics`, and switch from `pump_source` to - `pump_source_with(.., MQTT_GRAMMAR)`. -- **Overlapping filters.** Some brokers deliver one copy per matching - subscription, so `sensors/{d}/temp` next to `sensors/kitchen/temp` could - ingest a message twice. The connector subscribes only the **covering set**: - a filter is dropped from the subscription list when another filter matches - every topic it matches. The router still fans out locally, so each record - receives exactly one copy. Coverage is a level-by-level comparison over the - compiled patterns. The subscription QoS of a covering filter is the highest - QoS of the filters it covers (both backends subscribe at a fixed QoS 1 - today, so this only matters once per-link subscribe QoS is honoured). +It is filled on every build, not only with `observability`: a full table +silently turns away new publishers, so the numbers matter everywhere. The +field is optional in serde, so older AimX clients keep working. Key names +are not listed (there can be 65,535); `inbound_key_name` resolves one. + +The lock is taken once per keyed message by the single task that routes for +the connector, so it is uncontended. + +### 5.7 MQTT specifics + +- `MqttGrammar` follows MQTT 3.1.1 §4.7: `+` one level, `#` the rest + including the parent level (`a/#` matches `a`) and only last, a leading + wildcard does not match a `$…` topic, and a wildcard must be a whole + level. +- Both backends subscribe `covering_resource_ids()` of + `db.inbound_router("mqtt", &MqttGrammar)` and route with + `pump_source_with(.., &MqttGrammar)`. +- **Covering set.** A filter is left out when another matches every topic + it matches (`sensors/kitchen/temp` under `sensors/+/temp`; `a/+/b` under + `a/#`). Required for backend parity (§3.1). The router still fans each + message out to every route. Both backends subscribe at a fixed QoS 1 + today; once per-link subscribe QoS is honoured, a covering filter takes + the highest QoS of the filters it covers. - Outbound links reject patterns at `build()`: you cannot publish to a filter. -## 5. Guidance for pattern records +### 5.8 Compatibility + +No existing function or type signature changes. Two public structs with +public fields gain fields, which breaks code that builds them as struct +literals: + +- `InboundConnectorLink` gains `match_ingest_factory` and `key`. It has + `InboundConnectorLink::new`, and no in-tree code uses a literal. +- `RecordMetadata` gains `inbound_keys`. It has `RecordMetadata::new`; the + serde form is backward compatible. + +Mark both `#[non_exhaustive]` in the same change, so later fields are not +breaking. + +## 6. Guidance for pattern records A pattern record is one interleaved stream: its latest value is whichever publisher sent last. Applications that need per-publisher state keep it -downstream, for example a transform holding an array indexed by `KeyId`. -Use `SpmcRing`; `SingleLatest` would drop other publishers' values between -reads. +downstream, for example a transform holding an array indexed by +`KeyId::index()`. Use `SpmcRing`; `SingleLatest` would drop other +publishers' values between reads. When the broker binds topics to credentials (e.g. an ACL allowing a client to publish only under `sensors//…`), the topic is the @@ -275,7 +377,10 @@ authenticated part of a message and the payload is the publisher's claim. Keep them separate in the record: identity from `m.get(..)`/`m.key()`, everything else from the payload. -## 6. Relation to 054 +Watch `inbound_keys.dropped` in record metadata: a non-zero value means +publishers are being turned away because the key table is full. + +## 7. Relation to 054 When 054 lands, `InboundDispatch::dispatch(topic, payload)` borrows the topic for the whole ingest call, exactly as `Router::route` does today. Because @@ -283,72 +388,80 @@ for the whole ingest call, exactly as `Router::route` does today. Because - `with_match_deserializer`'s closure signature and `TopicMatch<'a>` are unchanged. -- `InboundDispatch::new` gains a grammar argument (or a `with_grammar` - variant) and replaces `inbound_router` + `pump_source_with` for migrated - connectors. -- The `TopicGrammar` struct may become a trait for inlining if the - 054 bench shows the pattern walk. +- `InboundDispatch::new` takes a `&'static dyn TopicGrammar` and replaces + `inbound_router` + `pump_source_with` for migrated connectors. +- 054's breaking window is where `&RuntimeContext` can replace the by-value + context in both deserializer builders (§5.5). -## 7. Alternatives considered +## 8. Alternatives considered 1. **Reuse `topic_matches` from `session/topic_match.rs`.** Wrong grammar: dot-separated with `*` for one level. Translating MQTT filters breaks on topics that contain dots. 2. **MQTT grammar hard-coded in core.** Makes the router protocol-aware; Zenoh (053) would add a second hard-coded grammar. -3. **Change `IngestFn` to take the topic.** Clean, but breaking; 054 is the +3. **Grammar as a data struct** (separator, tokens, a `hidden` function). + Fits MQTT but cannot express Zenoh's mid-pattern `**` or per-level hidden + chunks (§3.2). The trait costs nothing measurable. +4. **Change `IngestFn` to take the topic.** Clean, but breaking; 054 is the breaking window that removes the need. -4. **Records created per new topic at runtime.** Per-publisher buffers and +5. **Match carried in `RuntimeContext` (`ctx.inbound_match()`).** Reaches + plain `with_deserializer` users, but costs one allocation per message, + lets the match outlive the message through a cloned context, and breaks + when 054 turns the topic into a borrow. +6. **Records created per new topic at runtime.** Per-publisher buffers and AimX addresses, but it is the post-`run()` registration problem, far larger than this. -5. **Positional captures only (`+` → index 0).** Indices shift when a pattern +7. **Positional captures only (`+` → index 0).** Indices shift when a pattern changes; names cost nothing at runtime. -6. **Match carried in `RuntimeContext` (`ctx.inbound_match()`).** Reaches - plain `with_deserializer` users, but the context cannot borrow the topic, - so it costs one allocation per message, the match can outlive the message - through a cloned context, and it breaks when 054 turns the topic into a - borrow. Rejected in review. -7. **Declaring grammars on the builder** so every check runs at `build()`. - Adds a registration step for every connector; connector-build errors are - early enough. -8. **Allowing overlapping subscriptions** and documenting duplicates, or - **rejecting overlaps** at build. The first gives records duplicate - messages on some brokers; the second rules out a legitimate layout. -9. **Delivering unkeyed messages when the key table is full.** Keeps - messages, but every consumer then has to handle `key() == None` on a - keyed link. Capacity as an admission limit is the simpler contract. - -## 8. Open questions - -- **`mountain-mqtt` subscriptions.** Confirm `subscribe_packet` accepts `+` - and `#` unchanged. -- **Broker duplicate behaviour.** The covering set makes it irrelevant for - correctness, but record Mosquitto 2.x and EMQX behaviour in the - integration test notes so the rationale is checked. - -## 9. Acceptance criteria - -1. Router unit tests: the MQTT §4.7 cases above, captures at first, middle - and last level, `{name..}` matching zero levels, exact and pattern routes - on one topic both delivering, `TopicGrammar::EXACT` routers unchanged. -2. `build()` rejects each `{…}` error in §4.2 with the record key; - `inbound_router` rejects each connector-build error in §4.2, including a - pattern on an `EXACT` connector and an invalid resolver-returned pattern. +8. **All checks at `build()`** by declaring grammars on the builder. Adds a + registration step for every connector; connector-build errors are early + enough. +9. **Subscribing every filter as written.** Duplicates on MQTT 5 but not + 3.1.1 (§3.1), so the backends would disagree. **Rejecting overlaps** + instead rules out a legitimate layout. +10. **Delivering unkeyed messages when the key table is full.** Every + consumer of a keyed link would have to handle `key() == None`. +11. **One key table per link.** Overlapping `KeyId`s on a record with two + keyed links (§3.2). +12. **Reserving the key table's full capacity.** 67 KB for 1,024 keys up + front; lazy growth measured the same per message. +13. **An index over pattern routes.** Not needed at the measured cost; + revisit with a benchmark. + +## 9. Open questions + +- **Per-link subscribe QoS.** Both backends subscribe at QoS 1. When + per-link QoS is honoured, confirm the covering filter's QoS rule (§5.7) in + the parity test. + +## 10. Acceptance criteria + +1. Matcher unit tests: the MQTT §4.7 cases; captures at first, middle and + last level; `{name..}` matching zero levels; a Zenoh-style test grammar + with `a/**/b`, a mid-pattern `{path..}` and verbatim `@` chunks; + `ExactGrammar` routers unchanged. +2. `build()` rejects each §5.3 build-time error with the record key; + `inbound_router` rejects each connector-build error, including a pattern + on an `ExactGrammar` connector and an invalid resolver-returned pattern. 3. Covering-set unit tests: `sensors/+/temp` covers `sensors/kitchen/temp`; - `a/#` covers `a` and `a/+/b`; unrelated filters are all kept. -4. Tokio integration test against a local Mosquitto: two clients publish to - `sensors/a/temp` and `sensors/b/temp`; one record receives both; the - deserializer sees `device = a` and `b`; distinct `KeyId`s; - `db.inbound_key_name` returns the names; a table of capacity 1 drops the - second device and counts it. An exact link on `sensors/a/temp` next to the - pattern link receives each message once, and the pattern record once. -5. `b0_alloc_connector` gains `inbound_route_pattern` (0 allocs/msg), - `inbound_route_keyed_known` (0 allocs/msg) and `inbound_route_keyed_new` - (1 alloc, the key name). Existing rows unchanged. -6. `weather-station-gamma` and the embedded MQTT demo build for + `a/#` covers `a` and `a/+/b`; `a/**/b` covers `a/*/b`; unrelated filters + are all kept. +4. Parity test, both backends against one broker: a pattern link beside a + covered exact link subscribes only the covering filter; each record + receives the message once; capture and key reach the deserializer. +5. End-to-end test: two devices on one pattern record get distinct + `KeyId`s; `inbound_key_name` resolves them; a table of capacity 2 drops + the third device and counts it; two keyed links on one record share keys; + `RecordMetadata::inbound_keys` reports captures, capacity, assigned and + dropped. +6. `b0_alloc_connector` gains `inbound_route_pattern` (0 allocs/msg), + `inbound_route_keyed_known` (0) and `inbound_route_keyed_new` (1). + Existing rows unchanged. +7. `weather-station-gamma` and the embedded MQTT demo build for `thumbv7em-none-eabihf` with no behaviour change. -## 10. References +## 11. References - [018 — Dynamic MQTT topics](./018-M7-dynamic-mqtt-topics.md) - [052 — Runtime-neutral connectors](./052-runtime-neutral-connectors.md) @@ -356,3 +469,4 @@ for the whole ingest call, exactly as `Router::route` does today. Because - [053 — Zenoh connector](./053-zenoh-connector.md) (a second grammar) - [054 — Zero-allocation connector boundary](./054-zero-alloc-connector-boundary.md) - [MQTT 3.1.1 §4.7 — Topic names and topic filters](https://docs.oasis-open.org/mqtt/mqtt/v3.1.1/os/mqtt-v3.1.1-os.html#_Toc398718106) +- [Zenoh key expressions](https://github.com/eclipse-zenoh/roadmap/blob/main/rfcs/ALL/Key%20Expressions.md)