Skip to content

feat: wildcard inbound links (055) - #267

Open
lxsaah wants to merge 14 commits into
mainfrom
feat/mqtt-wildcards
Open

lxsaah wants to merge 14 commits into
mainfrom
feat/mqtt-wildcards

Conversation

@lxsaah

@lxsaah lxsaah commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Description

Implements design 055: one inbound link with a topic pattern feeds many topics into one record.

reg.buffer(BufferCfg::SpmcRing { capacity: 256 })
   .link_from("mqtt://sensors/{device}/temp")
   .key("device", 1024)
   .with_match_deserializer(|ctx, m, bytes| {
       let device = m.get("device").unwrap();   // &str borrowed from the topic
       let key = m.key();                       // Some(KeyId), 0-based via .index()
       Reading::decode(key, bytes)
   })
   .finish();

How it is split. Core owns the {name} / {name..} syntax, capture numbering, keys and routing. Each connector owns its wildcard rules via a TopicGrammar, which compiles a pattern into a matcher and decides which filters cover others. This revises the spike's shape, where core held one matcher for every protocol (055 §3.3).

aimdb-core

  • TopicPattern, and the TopicGrammar / TopicFilter traits. ExactGrammar is for connectors without wildcards.
  • with_match_deserializer receives a TopicMatch borrowed from the router's stack (topic(), get(name), key()) and a borrowed &RuntimeContext.
  • .key(name, capacity): one key table per record, shared by its keyed links. The table grows as values arrive. When it is full, the message is dropped and counted.
  • inbound_key_name resolves a key. RecordMetadata::inbound_keys reports captures, capacity, assigned and dropped.
  • AimDb::inbound_router(scheme, grammar) compiles every link on a scheme, including patterns a topic resolver returns. It reports every link it cannot compile at once.
  • Router::subscriptions() lists the filters to subscribe, leaving out those another filter covers.

aimdb-mqtt-connector

  • MqttGrammar implements MQTT 3.1.1 §4.7: +/{name} match one level, #/{name..} the rest (last only), and a leading wildcard does not match a $… topic.
  • Both backends subscribe only the covering filters. Mosquitto delivers overlapping subscriptions once to MQTT 3.1.1 clients and twice to MQTT 5 clients, so without covering the native and embedded backends would disagree (055 §3.1).
  • A hand-written +/# topic now matches. Before, it subscribed but never delivered.

Breaking: one inbound path

  • Removed collect_inbound_routes, RouterBuilder, Route and the public Router::new. A Router now only comes from inbound_router.
  • IngestFn receives the TopicMatch.
  • 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, and both become #[non_exhaustive].
  • The user-facing link API is unchanged.
  • KNX, WebSocket, TCP, UDS and serial use ExactGrammar: a {…} link on them now fails the build instead of being skipped.
  • Every in-tree connector is migrated. Crate version bumps follow separately.

Out of scope (055 §2)

  • Inbound subscribe QoS stays at 1 on both backends, as before. The with_qos doc no longer claims otherwise.
  • Outbound topic templates.

Allocations (b0_alloc_connector, 64 routes)

Path Allocs/msg
exact route 0 (unchanged)
pattern route 0
keyed, known key 0
keyed, new key 1 (the key name)

The existing rows are identical to the committed baseline.

Tests

  • MqttGrammar: the §4.7 examples, captures at the first, middle and last level, {rest..} matching zero levels, and covering, including $ topics.
  • Core: syntax errors, key validation, inbound_router errors (grammar rejection, invalid resolver pattern, keyed capture missing after resolution), shared keys, and key metadata.
  • Backend parity: both backends against one broker, with a keyed pattern link next to an exact link it covers. Each backend subscribes only parity/+/in; each record receives the message exactly once; the capture and key reach the deserializer.

Related Issue

Checklist

  • I have read the CONTRIBUTING.md document.
  • My code follows the project's coding standards.
  • I have added tests to cover my changes.
  • All new and existing tests passed (make check): targeted clippy and test runs per affected crate and feature leg passed locally; the full matrix runs in CI.
  • I have updated the documentation accordingly (055, CHANGELOGs).

@lxsaah lxsaah left a comment •

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.

Blocking:

  1. {x}{y} in one MQTT level drops capture x, and the shared spans buffer then leaks stale values from earlier routes (grammar.rs / router.rs).
  2. Backend parity still breaks for partially overlapping filters: embedded (MQTT 5) ingests twice what native (3.1.1) does, reproduced on Mosquitto (router.rs subscriptions).
  3. inbound_keys.dropped counts per route, not per message.

Design: the key table can be exhausted permanently by any publisher, since keys are assigned before deserialization and never freed.

Nits: literal {/} in topics is an undocumented breaking change and inbound/outbound are asymmetric; KeyId has no record identity; stale doc reference aimdb-websocket-connector/src/client/builder.rs:12 still mentions db.collect_inbound_routes.

}
PatternPart::Capture { multi, .. } => {
if let Some(level) = raw.last_mut() {
level.capture = Some((captures, *multi));

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.

Bug: two captures in one level silently lose the first one.

a/{x}{y} (and {x}{y..}) compiles: the second capture overwrites level.capture, the level text stays empty, so it becomes Single(Some(1)) and the filter is a/+. Capture 0 (x) is never written by matches.

Combined with the router reusing one Spans array across routes (see the comment in router.rs), m.get("x") then returns whatever an earlier, non-matching route left there. Repro through AimDb + inbound_router("mqtt", &MqttGrammar):

reg.link_from("mqtt://s/{a}/nomatch").with_match_deserializer(..).finish(); // tried first
reg.link_from("mqtt://s/{x}{y}").with_match_deserializer(/* reads m.get("x") */).finish();
// route("s/ab")  -> x == "ab"  (route 1's span for {a}, which then failed on level 2)

If x is the keyed capture, the wrong value gets keyed.

Suggested fix: reject when a level already has a capture, e.g.

PatternPart::Capture { multi, .. } => {
    if let Some(level) = raw.last_mut() {
        if level.capture.is_some() {
            return Err(format!("'{topic}': a capture must be a whole level, not part of one"));
        }
        level.capture = Some((captures, *multi));
    }
    captures += 1;
}

Generated by Claude Code

Comment thread aimdb-core/src/router.rs

// Linear search through all routes
// Note: Multiple routes may match the same resource_id (different types)
let mut spans: Spans = [(0, 0); MAX_CAPTURES];

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.

One spans buffer is shared across every route in the loop, and TopicFilter::matches writes spans as it goes, so a route that fails part-way leaves partial spans behind. Any grammar that doesn't write every capture on a successful match (the MqttGrammar {x}{y} case above, or a third-party TopicGrammar) then leaks stale values into TopicMatch::get and into key assignment.

Resetting before each attempt is cheap and makes this robust against any grammar (still zero-alloc):

for route in &self.routes {
    spans = [(0, 0); MAX_CAPTURES];
    if route.matches(resource_id, &mut spans) { … }
}

Alternatively, document on TopicFilter::matches that a successful match must write every capture.


Generated by Claude Code

Comment thread aimdb-core/src/router.rs
Router::new(self.routes)
/// Filters to subscribe: [`resource_ids`](Self::resource_ids) without
/// the filters another one covers.
pub fn subscriptions(&self) -> Vec<Arc<str>> {

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.

Backend parity still breaks on partially overlapping filters (reproduced against real Mosquitto 2).

Covering only removes filters that are fully covered. For r/{x}/c and r/b/{y} on one record, neither covers the other, so both are subscribed. One PUBLISH to r/b/c:

backend {x} link calls {y} link calls values in record
native (rumqttc, 3.1.1) 1 1 2
embedded (mountain-mqtt, MQTT 5) 2 2 4

(Test: both backends against a local mosquitto, publishes via mosquitto_pub -q 1. The fully-covered exact-link case from the PR's parity test does hold on the real broker: 1 call each on both backends.)

So the 055 §3.1 parity problem remains for partial overlaps. Options: MQTT 5 subscription identifiers on the embedded backend (drop the PUBLISH unless it carries the lowest matching id), or document it as a known difference. Separately, even with one delivery the record receives the message once per matching link. That may be intended, but worth a line in the docs.


Generated by Claude Code

Comment thread aimdb-core/src/router.rs
matched = true;
match (route.ingest)(ctx, payload) {
let Some(key) = route.key(resource_id, &spans) else {
log_debug!("Key table full, dropped message on '{}'", resource_id);

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.

dropped counts once per matching keyed route, not per message. With two keyed links on one record whose patterns both match (t/{d} and t/{dev}, capacity 1), one message on a new value reports inbound_keys.dropped == 2 in list_records(). The field doc says "Messages turned away because the table was full", so it either needs counting once per route() call per table, or the doc should say "per link".


Generated by Claude Code


/// The key for `name`, assigned if new. `None` when the table is full;
/// the caller drops the message and it is counted.
pub(crate) fn key(&self, name: &str) -> Option<KeyId> {

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.

Design: any publisher can permanently exhaust the key table. The key is assigned before the deserializer runs (Router::route → route.key() → ingest), and keys are never freed. With .key("device", 2), two junk publishes to sensors/junk1/temp and sensors/junk2/temp with unparseable payloads fill the table, and the real sensors/kitchen/temp is then dropped for the rest of the process's life.

On a shared broker that's an easy DoS for anyone with publish rights under the prefix. Some options:

  • assign provisionally and only commit the key if ingest returns Ok (the deserializer would see a tentative key);
  • an allowlist or validator hook on .key();
  • at minimum, document it next to .key() and in 055, and surface dropped prominently.

Generated by Claude Code

// Mutual exclusion with local producers (.source()/.transform()) is
// validated once, in build(), where the record key is known.
// The `{…}` syntax; the connector checks the rest.
let pattern = crate::TopicPattern::parse(url.resource_id()).map_err(|e| e.to_string())?;

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) Behavior change for topics containing literal braces. { and } are legal MQTT topic characters (and legal for KNX, WS and the other connectors). Before this PR, link_from("mqtt://a/{b}") subscribed the literal topic a/{b}. Now it's a capture subscribed as a/+, and a/x} or a/{{b}} fail the build, with no escape syntax.

Outbound is also asymmetric: link_to (line ~908) only rejects when parsing succeeds with captures, so link_to("mqtt://a/x}") is accepted as a literal, while the same topic inbound is rejected.

Maybe worth a CHANGELOG line, and either an escape (e.g. {{ → {) or consistent handling in both directions.


Generated by Claude Code


/// A capture value's key, assigned the first time the value is seen.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct KeyId(NonZeroU16);

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) A KeyId carries no record identity, so a key from record A passed to inbound_key_name("B", key) quietly resolves to B's value with the same index. The same happens across databases. That's fine for a dense u16, but a sentence on KeyId ("only meaningful for the record whose link produced it") would help.


Generated by Claude Code

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant