From f54182cfce37df265d0f12371135c12a81364351 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:59:14 +0000 Subject: [PATCH] docs(design): 054 zero-alloc connector boundary, 055 wildcard inbound links, 056 aimdb join CLI - bench: add b0_alloc_connector, measuring per-message allocations at the connector boundary (Router::route 0, pump_source 2, outbound 2-3) and asserting them against a committed baseline. - 054: zero-allocation connector boundary, grounded in that baseline; extends 037's poll-based SPI to inbound dispatch, outbound readers, per-route publishers and written topics. - 055: wildcard inbound links on today's interfaces (additive): connector-supplied grammar, named captures, optional key interning. - 056: aimdb join CLI, porting the planning-branch implementation and fixing its five review defects. - 043 rev 2: deployment-neutral join format with discovery, optional app values and credential rotation on re-join. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01E5ySyp4thSEktXg4Awwpn7 --- aimdb-bench/Cargo.toml | 5 + aimdb-bench/README.md | 9 + aimdb-bench/benches/b0_alloc_connector.rs | 373 ++++++++++++++++++ .../data/baselines/b0_alloc_connector.json | 47 +++ docs/design/043-join-endpoint-v1.md | 248 ++++++------ .../054-zero-alloc-connector-boundary.md | 296 ++++++++++++++ docs/design/055-wildcard-inbound-links.md | 299 ++++++++++++++ docs/design/056-aimdb-join-cli.md | 138 +++++++ 8 files changed, 1296 insertions(+), 119 deletions(-) create mode 100644 aimdb-bench/benches/b0_alloc_connector.rs create mode 100644 aimdb-bench/data/baselines/b0_alloc_connector.json create mode 100644 docs/design/054-zero-alloc-connector-boundary.md create mode 100644 docs/design/055-wildcard-inbound-links.md create mode 100644 docs/design/056-aimdb-join-cli.md diff --git a/aimdb-bench/Cargo.toml b/aimdb-bench/Cargo.toml index af4a5d3f..fcb2645b 100644 --- a/aimdb-bench/Cargo.toml +++ b/aimdb-bench/Cargo.toml @@ -62,6 +62,11 @@ harness = false name = "b1_b2_remote_json" harness = false +# Connector-boundary allocations per message (design 054 baseline). +[[bench]] +name = "b0_alloc_connector" +harness = false + [features] default = ["std"] # Gates `profiles`/`reports`/`harness` and their criterion/serde_json/ diff --git a/aimdb-bench/README.md b/aimdb-bench/README.md index 0d403b41..7b79866d 100644 --- a/aimdb-bench/README.md +++ b/aimdb-bench/README.md @@ -26,6 +26,10 @@ 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. - **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 @@ -65,6 +69,9 @@ cargo bench -p aimdb-bench --bench b0_alloc_linkable # B0 — direct-vs-tree JSON allocations through the AimX envelope boundary cargo bench -p aimdb-bench --bench b0_alloc_remote_json +# B0 — allocations at the connector boundary (design 054 baseline) +cargo bench -p aimdb-bench --bench b0_alloc_connector + # B1 + B2 — latency (time/iter) and throughput (msgs/sec), one Criterion suite cargo bench -p aimdb-bench --bench b1_b2_tokio @@ -124,6 +131,8 @@ The required result is **0 allocation calls and 0 allocated bytes**. It isolates `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. + > **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`. **Noise reduction:** a `new_current_thread()` Tokio executor is used so there are no work-stealing threads and Tokio's scheduler does not allocate per-poll in the hot path. diff --git a/aimdb-bench/benches/b0_alloc_connector.rs b/aimdb-bench/benches/b0_alloc_connector.rs new file mode 100644 index 00000000..9c8ab489 --- /dev/null +++ b/aimdb-bench/benches/b0_alloc_connector.rs @@ -0,0 +1,373 @@ +//! 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: +//! +//! - **Inbound:** `Router::route` alone, 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. +//! +//! 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 +//! expected values are asserted, so a regression *or* an improvement fails the +//! bench until `EXPECTED` and the committed baseline are updated together. +//! +//! Run `cargo bench -p aimdb-bench --bench b0_alloc_connector`; results are +//! written to `aimdb-bench/target/bench-results/b0_alloc_connector.json` and +//! compared by hand with `aimdb-bench/data/baselines/b0_alloc_connector.json`. + +use std::future::Future; +use std::hint::black_box; +use std::pin::Pin; +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::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_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; + +#[global_allocator] +static GLOBAL: aimdb_bench::alloc::CountingAllocator = + aimdb_bench::alloc::CountingAllocator(std::alloc::System); + +const WARMUP_ITERS: usize = 100; +const MEASURE_ITERS: usize = 2_000; +const DECOY_ROUTES: usize = 63; +const SCRATCH_CAPACITY: usize = 64; + +/// 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_pump_source_minimal", 2), + ("outbound_scratch_static_topic", 2), + ("outbound_scratch_dynamic_topic", 3), + ("outbound_owned_static_topic", 3), +]; + +#[derive(Clone, Copy, Debug)] +struct Reading { + id: u32, + value: f32, +} + +fn reading(i: usize) -> Reading { + Reading { + id: i as u32, + value: 21.5, + } +} + +fn encode(r: &Reading) -> [u8; 8] { + let mut out = [0u8; 8]; + out[..4].copy_from_slice(&r.id.to_le_bytes()); + out[4..].copy_from_slice(&r.value.to_le_bytes()); + out +} + +// --- A connector that owns no transport: it only registers the scheme. ---- + +const SCHEME: &str = "bench"; + +struct NoopConnectorBuilder; + +impl ConnectorBuilder for NoopConnectorBuilder { + #[allow(clippy::type_complexity)] + fn build<'a>( + &'a self, + _db: &'a AimDb, + ) -> Pin< + Box< + dyn Future + Send + 'static>>>>> + + Send + + 'a, + >, + > { + Box::pin(async { Ok(Vec::new()) }) + } + + fn scheme(&self) -> &str { + SCHEME + } +} + +/// 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() + .runtime(runtime) + .with_connector(NoopConnectorBuilder); + configure(&mut builder); + let (db, _runner) = builder.build().await.expect("build"); + db +} + +// --- Inbound ---------------------------------------------------------------- + +/// One linked record plus `DECOY_ROUTES` others, so routing scans a realistic +/// table. +async fn inbound_db() -> AimDb { + build_db(|b| { + b.configure::("in.target", |reg| { + reg.buffer(BufferCfg::SpmcRing { capacity: 64 }) + .link_from("bench://in/target") + .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::(StringKey::intern(format!("in.decoy{i}")), |reg| { + reg.buffer(BufferCfg::SingleLatest) + .link_from(&topic) + .with_deserializer(|_ctx, _bytes| Ok(reading(0))) + .finish(); + }); + } + }) + .await +} + +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 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() +} + +/// 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(); + for fut in pump_source(db, SCHEME, 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() +} + +// --- Driver ----------------------------------------------------------------- + +fn main() { + let runtime = tokio::runtime::Builder::new_current_thread() + .build() + .expect("bench runtime"); + + let measured: Vec<(&str, &str, (u64, u64))> = runtime.block_on(async { + vec![ + ("inbound_route", "SpmcRing", measure_route().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, + ), + ] + }); + + println!("=== B0 connector boundary: allocations per message ==="); + let mut reports = Vec::with_capacity(measured.len()); + let mut mismatches = Vec::new(); + for &(case, buffer, (allocs, bytes)) in &measured { + let report = AllocReport::new(case, buffer, MEASURE_ITERS, allocs, bytes); + println!( + "{case:<34} {:>5.2} allocs/msg {:>7.1} bytes/msg", + report.allocs_per_msg, report.bytes_per_msg + ); + let expected = EXPECTED + .iter() + .find(|(name, _)| *name == case) + .map(|(_, n)| *n) + .expect("every case has an expected value"); + if report.allocs_per_msg.round() as u64 != expected { + mismatches.push(format!( + "{case}: expected {expected}, measured {:.2}", + report.allocs_per_msg + )); + } + reports.push(report); + } + write_reports("b0_alloc_connector", &reports); + + assert!( + mismatches.is_empty(), + "connector-boundary allocations changed — update EXPECTED and the baseline:\n {}", + mismatches.join("\n ") + ); +} diff --git a/aimdb-bench/data/baselines/b0_alloc_connector.json b/aimdb-bench/data/baselines/b0_alloc_connector.json new file mode 100644 index 00000000..3afe97e6 --- /dev/null +++ b/aimdb-bench/data/baselines/b0_alloc_connector.json @@ -0,0 +1,47 @@ +[ + { + "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_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 + } +] \ No newline at end of file diff --git a/docs/design/043-join-endpoint-v1.md b/docs/design/043-join-endpoint-v1.md index e78a7bb4..4aa4d63c 100644 --- a/docs/design/043-join-endpoint-v1.md +++ b/docs/design/043-join-endpoint-v1.md @@ -1,179 +1,189 @@ # 043 — Join endpoint format v1 -**Status:** 📝 Proposed - -**Scope:** the open wire format of the provisioning endpoint `aimdb join` -speaks (`POST /v1/join`), and the station profile file the CLI writes. No -code in this repo implements the *server* side — the flagship's provisioning -service lives in the private ops repo ([042](./042-public-weather-mesh-flagship.md) -D6) — but the format is public so any AimDB deployment can implement it and -the CLI stays deployment-agnostic. - -**Consumers:** `aimdb join` ([`tools/aimdb-cli`](../../tools/aimdb-cli), -042 §7) is built against this document; `weather-station` -([`examples/weather-mesh-demo`](../../examples/weather-mesh-demo)) consumes -the profile file defined in §4. +**Status:** 📝 Proposed (rev 2) + +- **rev 1** — first draft, written for the public weather mesh (042): the + client collected station metadata by prompts and the response carried a + slot-based `station_id`. +- **rev 2** — made deployment-neutral. The client no longer collects + metadata; `app` becomes optional in both directions. Adds a discovery + document so the CLI needs no compiled-in OAuth client ID. Re-joining is + allowed to rotate credentials. No server implemented rev 1, so v1 is + revised in place. + +**Scope:** the open wire format a *provisioning endpoint* speaks, and the +profile file `aimdb join` writes. Implementing the server side is out of +scope for this repository; the format is public so any AimDB deployment can +implement it. The client is designed in +[056](./056-aimdb-join-cli.md). --- ## 1. Overview -A *provisioning endpoint* admits a new edge node to a deployment: it verifies -an identity, applies the deployment's admission rules, allocates whatever -resources a member needs (for the weather mesh flagship: a slot number and a -scoped MQTT credential), and returns a **station profile** the node runs -with. The CLI is generic over the endpoint URL: +A provisioning endpoint admits a new node to a deployment: it verifies an +identity, applies the deployment's admission rules, creates whatever the node +needs (typically a broker credential scoped to the node's own topics), and +returns a **profile** the node runs with. ``` -aimdb join [--token ] [--out station.toml] +aimdb join + GET /v1/join → discovery document (§2) + POST /v1/join → profile (§3, §4) ``` -`aimdb join https://mesh.aimdb.dev` sends `POST https://mesh.aimdb.dev/v1/join`. -The path is fixed; the base URL is the deployment. TLS is required — the -request carries a bearer-equivalent token and the response carries a broker -credential. +The path is fixed; the base URL is the deployment. TLS is required: the +request carries an identity token and the response carries a credential. -## 2. Request - -``` -POST /v1/join -Content-Type: application/json +## 2. Discovery — `GET /v1/join` +```json { - "auth": { "kind": "github", "token": "gho_…" }, - "app": { "name": "graz-balcony", "location": "Graz" } + "auth": [ + { "kind": "github", "client_id": "Iv1.0123456789abcdef" }, + { "kind": "claim-token" } + ], + "app_fields": [ + { "name": "label", "required": false, "help": "Shown next to your node" } + ] } ``` -### 2.1 `auth` — identity envelope +| Field | Contract | +|---|---| +| `auth` | The identity kinds this deployment accepts, in order of preference. `github` carries the public OAuth client ID for the device flow. | +| `app_fields` | Optional. Names of `app` values the deployment accepts in the join request, with a help text. Empty or absent: the deployment wants none. | -`auth.kind` keeps the format deployment-agnostic. v1 defines two kinds: +Unknown fields are ignored. A deployment that does not serve discovery (404) +is treated as accepting `claim-token` only. -| `kind` | `token` | Used by | -|---|---|---| -| `github` | A GitHub **user access token**, obtained by the CLI via the OAuth device flow (no scopes — public identity only). The service verifies it against the GitHub API and applies identity-based admission rules (one active slot per account, minimum account age, …). | The public flagship (self-serve). | -| `claim-token` | A pre-issued, one-time opaque token, distributed out of band by the deployment operator. | Private deployments without public self-serve. `aimdb join --token ` selects this kind. | +## 3. Request — `POST /v1/join` -Deployments accept the kinds they support and reject others with `403`. +```json +{ + "auth": { "kind": "github", "token": "gho_…" }, + "app": { "label": "balcony" } +} +``` -### 2.2 `app` — deployment-specific metadata +### 3.1 `auth` -`app` is an **opaque string-keyed map**: the CLI collects its fields via -interactive prompts (driven by the deployment's docs) and passes them through -untouched — no deployment-specific types exist in the CLI. The service -validates, normalizes, and echoes the final values back in the profile (§3). +| `kind` | `token` | +|---|---| +| `github` | A GitHub **user access token** obtained through the OAuth device flow with the client ID from discovery, requesting no scopes. The service verifies that the token belongs to that OAuth app, applies identity-based rules, and should revoke the token after use. | +| `claim-token` | A pre-issued, single-use opaque token distributed by the operator. | -For the weather mesh flagship, v1 `app` fields are: +Deployments reject kinds they do not accept with `403`. -| Field | Meaning | -|---|---| -| `name` | Station name — becomes the public identity page (`aimdb.dev/mesh/`). | -| `location` | City-level location. The service geocodes and **coarsens it to 2 decimal places (~1 km)** before it appears anywhere (042 D8); precise location is never collected. | +### 3.2 `app` + +Optional string-keyed map of strings, passed through by the client without +interpretation. Only names announced in `app_fields` are meaningful; the +service validates and may ignore or reject others (`422`). -## 3. Response +## 4. Response -### 3.1 Success — `200 OK` +### 4.1 Success — `200 OK` ```json { "profile_version": 1, - "station_id": "slot-17", + "node_id": "alice", "broker": { - "url": "mqtts://xxxx.eu-central-1.emqx.cloud:8883", - "username": "station-17", + "url": "mqtts://mqtt.example.com:8883", + "username": "alice", "password": "…" }, "app": { - "name": "graz-balcony", - "lat": 47.07, - "lon": 15.44 + "publish_prefix": "sensors/alice/" } } ``` | Field | Contract | |---|---| -| `profile_version` | Format version of this response, `1` for this document. See §5. | -| `station_id` | The node's identity within the deployment — opaque to the CLI. The flagship uses `slot-`, where `` is the assigned slot in the hub's pool (042 §4); the station binary derives its topic/record namespace (`station//…`, `station..…`) from it. | -| `broker.url` | Connector URL of the data-plane broker, without credentials. `mqtts://` for the flagship (broker is TLS-only). | -| `broker.username`, `broker.password` | The per-station credential, scoped server-side to the station's own topics (flagship: publish-only on `station//#`). Individually revocable (042 D9). | -| `app` | The **final** deployment-specific values (post validation/coarsening) — clients must use these, not what they sent. The flagship echoes `name` and the coarsened `lat`/`lon` the station feeds to its weather source. | +| `profile_version` | Format version, `1` for this document (§6). | +| `node_id` | The node's identity in the deployment. Opaque to the client. | +| `broker.url` | Connector URL of the data-plane broker, without credentials. | +| `broker.username`, `broker.password` | The node's credential. The service scopes it; the client never sees how. | +| `app` | Optional deployment-specific values for the node, e.g. the topic prefix its credential may publish under. Opaque to the client. | + +### 4.2 Re-joining + +A join from an identity that already holds a membership is the +deployment's decision: -### 3.2 Errors +- **rotate** (recommended): issue a new password for the existing + membership and return `200` with the full profile. A lost profile is then + recovered by joining again. +- **refuse**: return `409` with a `message` naming the existing membership. -Error responses are JSON with a human-readable `message` the CLI prints -**as-is** and exits non-zero — the service owns the failure UX text: +### 4.3 Errors + +Error bodies are JSON with a human-readable `message`, which the client +prints as-is: ```json -{ "message": "The mesh is full (64/64 slots) — try again later." } +{ "message": "Admissions are paused — try again later." } ``` -| Status | Meaning | Flagship examples | -|---|---|---| -| `403` | Identity rejected. | Account too new; admissions paused (kill switch); unsupported `auth.kind`. | -| `409` | The identity already holds a membership. The `message` points at the existing profile (station name/id) rather than erroring opaquely. | "Your account already runs `graz-balcony` (slot-17). Revoke it first or reuse its profile." | -| `503` | Deployment full or temporarily not admitting. | Mesh at cap N. | +| Status | Meaning | +|---|---| +| `403` | Identity rejected: unsupported kind, token invalid, rules not met, admissions paused. | +| `409` | Identity already holds a membership and the deployment does not rotate (§4.2). | +| `422` | Invalid `app` values. | +| `429` | Too many requests. | +| `503` | Deployment full or temporarily not admitting. | -Any other non-2xx status is unexpected; the CLI reports it with the body if -one is present. +Any other non-2xx status is unexpected; the client reports it with the body +if present. -## 4. The station profile file (`station.toml`) +## 5. Profile file -`aimdb join` writes the §3.1 response as TOML (mode `0600` — it contains a -credential), default path `station.toml`, `--out` to override: +The client writes the §4.1 response as TOML, readable by the owner only +(mode `0600` on Unix): ```toml profile_version = 1 -station_id = "slot-17" +node_id = "alice" [broker] -url = "mqtts://xxxx.eu-central-1.emqx.cloud:8883" -username = "station-17" +url = "mqtts://mqtt.example.com:8883" +username = "alice" password = "…" [app] -name = "graz-balcony" -lat = 47.07 -lon = 15.44 +publish_prefix = "sensors/alice/" ``` -The mapping is mechanical (JSON object → TOML table, `app` passed through); -consumers (`weather-station --config station.toml`, gamma's `MESH_CONFIG` -build embedding — 042 §8.2) read this file, never the HTTP response. - -## 5. Versioning and compatibility - -The CLI ships on the workspace release cadence and must not need a release -in lockstep with mesh-side changes (042 §6.3). Therefore: - -- **Unknown fields are ignored**, both in the response (CLI) and in the - profile (station binaries). Services may add fields without a version bump. -- `profile_version` is bumped only for **breaking** changes to the meaning - of existing fields. A client that receives a `profile_version` newer than - it knows writes the profile anyway and warns; a station binary that cannot - interpret a profile says which version it expected. -- The request envelope is versioned by the URL path (`/v1/join`). New auth - kinds may be added to v1 (a service rejects unknown kinds with `403`, which - the CLI reports verbatim). - -## 6. What this format deliberately is not - -- **Not a client of the broker's management API.** The response carries a - ready-made credential; how it was minted (EMQX Cloud API for the flagship) - is invisible to the client, so no admin key can end up in a public binary - (042 D6/D7). -- **Not a schema registry.** Payload validity is enforced by contract - deserialization at the hub (042 §9); join-time carries no topic lists or - schemas. -- **Not a session.** One request, one response, no state on the client - beyond the profile file. Re-joining with an identity that already holds a - membership is a `409`, not an idempotent re-issue — credential re-issue is - a service-side policy decision. - -## 7. References - -- [042 — Public weather mesh flagship](./042-public-weather-mesh-flagship.md) - §6 (admission flow), §7 (`aimdb join`), D5–D8 -- [Implementation plan — 042](../plan/042-public-weather-mesh-flagship.md) - WP3/WP4 +The mapping is mechanical (JSON object → TOML table). Consumers read this +file, never the HTTP response. + +## 6. Versioning + +The CLI ships on the workspace release cadence and must not need a release in +lockstep with deployments: + +- Unknown fields are ignored everywhere (discovery, response, profile). + Services may add fields without a version bump. +- `profile_version` changes only when the meaning of an existing field + changes. A client receiving a newer version writes the profile and warns; a + consumer that cannot interpret a profile names the version it expected. +- The request and discovery formats are versioned by the path (`/v1/join`). + New auth kinds may be added to v1; services reject unknown kinds with `403`. + +## 7. What this format is not + +- **Not a client of the broker's management API.** Credentials are minted + server-side; no admin secret can end up in a client binary. +- **Not a schema registry.** Join carries no topic lists or payload schemas; + payload validity is enforced by the consuming records' deserializers. +- **Not a session.** One request, one response; no client state beyond the + profile file. + +## 8. References + +- [056 — `aimdb join` CLI](./056-aimdb-join-cli.md) - [GitHub OAuth device flow](https://docs.github.com/en/apps/oauth-apps/building-oauth-apps/authorizing-oauth-apps#device-flow) +- [GitHub: check a token](https://docs.github.com/en/rest/apps/oauth-applications#check-a-token) + (how a service verifies a token belongs to its OAuth app) diff --git a/docs/design/054-zero-alloc-connector-boundary.md b/docs/design/054-zero-alloc-connector-boundary.md new file mode 100644 index 00000000..4911c78c --- /dev/null +++ b/docs/design/054-zero-alloc-connector-boundary.md @@ -0,0 +1,296 @@ +# 054 — Zero-allocation connector boundary + +**Status:** 📝 Proposed + +**Scope:** the per-message path between `aimdb-core` and connectors, in both +directions: `Source`/`pump_source`, `SerializedReader`, `Connector::publish`, +`TopicProvider`, and the MQTT connector as the first adopter. SPI break in +`aimdb-core` (next major). User-facing APIs (`link_from`, `link_to`, +`with_deserializer`, `with_serializer`, `Reader::recv`) do not change. + +**Builds on:** [037 — Zero-allocation consume path](./037-zero-alloc-consume-path.md), +which made buffers and the consume path allocation-free and left the +connector boundary open (037 §2, "Outbound connector"; §3.3). + +--- + +## 1. Measured state + +Baseline: `cargo bench -p aimdb-bench --bench b0_alloc_connector` +(committed in `aimdb-bench/data/baselines/b0_alloc_connector.json`), Tokio +current-thread runtime, 2,000 messages after warm-up, no-op connector. + +| Path | Allocs/msg | Bytes/msg | +|---|---|---| +| `Router::route` (64 routes, deserialize + produce) | **0** | 0 | +| `pump_source` with a minimal `Source` | **2** | 25 | +| Outbound `recv_into` + `publish`, scratch serializer, static topic | **2** | 73 | +| same, dynamic topic (`TopicProvider`) | **3** | 81 | +| same, owned serializer | **3** | 81 | + +Buffers, the consume path (037) and routing/ingest allocate nothing. Every +non-zero row is caused by the connector interface: + +| # | Direction | Cause | Where | +|---|---|---|---| +| A1 | in | `Source::next` returns a boxed future | `session/mod.rs` (`BoxFut`) | +| A2 | in | `Source::next` returns an owned topic `String` | `session/mod.rs` | +| A3 | out | `SerializedReader::recv_into` returns a boxed future | `connector.rs` (`RecvSerializedIntoFuture`) | +| A4 | out | `Connector::publish` returns a boxed future | `transport.rs` | +| A5 | out | `TopicProvider::topic` returns `Option` | `connector.rs` | +| A6 | out | owned serializer returns `Vec` (only without `with_serializer_into`) | `typed_api.rs` | + +**Connector-side copies** (read from code, same traits): the MQTT connector +adds one payload copy inbound (`Arc::from(payload)`, both backends) and two +copies outbound (`destination.to_string()`, `payload.to_vec()`, both +backends). With `rumqttc` the outbound copies are required by its owned-value +`publish` API; on the embedded backend they exist to pass an owned action +through a channel to the session task. + +**Values with heap data.** Each reader receives its own clone of `T` +(`T: Clone` delivery). A value with one `String` field costs one allocation +per reader per message (measured with 1 and 3 readers). This is a property of +the buffer contract, not of the connector boundary; §4.5 covers it as +guidance. + +## 2. Goals and non-goals + +**Goals** + +- G1. Zero AimDB-added allocations per message on both connector directions, + for scratch serializers and static or written topics. +- G2. The MQTT connector adds no copies beyond what its client library's API + requires, documented per backend. +- G3. `b0_alloc_connector` gates the result in CI. + +**Non-goals** + +- Topic wildcards and templates. See [055](./055-wildcard-inbound-links.md). +- Making `T` delivery cheaper for heap-containing types (§4.5 is guidance). +- Allocations inside third-party client libraries. +- Remote-access JSON paths (tracked by `b0_alloc_remote_json`). + +## 3. Approach + +The same move as 037: object safety and `async fn` conflict, object safety +and `poll` do not. Per-message boxed futures exist only to satisfy trait +signatures, so the signatures change to poll form. Owned values that exist +only to cross the interface become borrows. + +## 4. Design + +### 4.1 Inbound: connectors push borrowed messages + +Connectors hold each received message by reference: inside `rumqttc`'s +`Event::Incoming(Publish)`, inside `mountain-mqtt`'s application-message +handler. `Router::route(&str, &[u8], ctx)` already takes borrows and does not +allocate. The interface in between is what forces owned copies (A1, A2). + +Core exposes the router directly to connectors: + +```rust +/// Built once per connector from the inbound routes of its scheme. +pub struct InboundDispatch { /* Router + RuntimeContext */ } + +impl InboundDispatch { + pub fn new(db: &AimDb, scheme: &str) -> Self; + /// Deserialize and produce into every matching record. Synchronous, + /// allocation-free, never blocks (full buffers drop and log, as today). + pub fn dispatch(&self, topic: &str, payload: &[u8]); + /// Topics to subscribe at the transport. + pub fn resource_ids(&self) -> &[Arc]; +} +``` + +- Connectors call `dispatch` wherever they hold the borrow. No future, no + copy, no channel. +- `Source` stays for transports that naturally produce owned frames, but in + poll form (037 pattern), which removes A1: + `fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll>`. + A2 remains for those transports by construction. +- Running ingest on the connector's own task is already today's behaviour + (`pump_source` runs it inline on the reader future); it stays synchronous + and non-blocking. + +### 4.2 Outbound: poll-based reader + +`SerializedReader::recv_into` becomes a poll method over the buffer reader's +existing `poll_recv` (037). Topic and payload are both written into storage +the pump owns: + +```rust +pub struct OutboundScratch { + pub payload: Vec, // allocated once per route (exists today) + pub topic: Vec, // new; allocated once per route +} + +pub enum Dest { Default, Written(usize), Owned(String) } // Owned: TopicProvider adapter + +pub struct OutboundFrame { pub dest: Dest, pub payload: SerializedPayload } + +pub trait SerializedReader: Send { + fn poll_recv_into( + &mut self, + cx: &mut Context<'_>, + ctx: &RuntimeContext, + scratch: &mut OutboundScratch, + ) -> Poll>; +} +``` + +Removes A3. `SerializedPayload::Owned` stays as the fallback for serializers +without an into-slice path (A6 remains by choice of the user). + +### 4.3 Outbound: topics are written, not returned + +```rust +pub trait TopicWriter: Send + Sync { + /// Write the destination into `out`. `Ok(false)` = use the link's default. + fn write_topic(&self, value: &T, out: &mut TopicBuf<'_>) -> Result; +} + +// TopicBuf implements core::fmt::Write over the route's topic scratch. +// usage: write!(out, "sensors/{}/{}", v.site, v.id)?; Ok(true) +``` + +- `.with_topic_writer(capacity, writer)` on the outbound builder; capacity is + the topic scratch size for that route. +- An overflow skips the message and increments a per-link counter; the topic + is never truncated. +- `TopicProvider` keeps working through an adapter that yields + `Dest::Owned` (A5 remains for code that has not migrated). + +### 4.4 Outbound: per-route publishers in poll form + +`Connector::publish` is shared by all routes and returns a boxed future +(A4). It is replaced by a publisher created once per route: + +```rust +pub trait Connector: Send + Sync { + /// Called once per outbound route when the pump starts. + fn route_publisher(&self, config: &ConnectorConfig) -> Box; +} + +pub trait RoutePublisher: Send { + /// Ready to accept one message (e.g. queue space). + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll>; + /// Hand one message over. Synchronous; the publisher copies what it + /// needs into storage it owns. + fn start_send(&mut self, dest: &str, payload: &[u8]) -> Result<(), PublishError>; +} +``` + +- This is the `futures::Sink` shape: backpressure through `poll_ready`, a + synchronous hand-over, no future per message. `&mut self` per route means + no locks. +- `start_send` borrows `dest` and `payload` from the pump's scratch; the + publisher decides how to keep them (a pre-allocated ring, a client queue). +- Configuration (`qos`, `retain`, …) is parsed once in `route_publisher`, + not per message as the MQTT sinks do today. +- A provided `BoxedPublisher` adapter wraps an old-style async `publish` + for third-party connectors during migration, keeping today's cost. + +With 4.2–4.4 the pump's per-message loop is: + +```rust +let frame = poll_fn(|cx| reader.poll_recv_into(cx, &ctx, &mut scratch)).await?; +poll_fn(|cx| publisher.poll_ready(cx)).await?; +publisher.start_send(dest_of(&frame, &scratch, &default_topic), payload_of(&frame, &scratch))?; +``` + +### 4.5 Values with heap data (guidance) + +Delivery clones `T` once per reader. For records with more than one reader +at message rate: + +- prefer `Copy` values (numbers, fixed-size arrays, `heapless::String`); +- or wrap the value in `Arc`, so each reader's clone is a reference-count + increment; +- keep rarely changing, heap-heavy data in a separate record. + +Document this next to the buffer types; no API change. + +### 4.6 MQTT connector + +| Backend | Inbound | Outbound | +|---|---|---| +| Embedded (`mountain-mqtt`) | `dispatch` from the application-message handler inside the session task; the event channel for inbound messages is removed | `start_send` copies topic and payload into a byte ring allocated at build; the session task reads frames from it. `poll_ready` = ring space | +| Native (`rumqttc`) | `dispatch(&publish.topic, &publish.payload)` in the event-loop task | `rumqttc::AsyncClient::publish` takes owned `String` and payload, so the two copies stay. Its future is kept in a per-route `tokio_util::sync::ReusableBoxFuture` (037's Tokio technique), which removes the box but not the copies | + +`rumqttc`'s own allocations (building `Publish`, its request channel) are +outside G1. The embedded backend also runs on Tokio through +`TokioTcpDialer` (052), which gives std users a fully allocation-free MQTT +path when they need it. + +## 5. Measurement and gate + +`b0_alloc_connector` gains rows for the new interfaces. Targets after this +design: + +| Row | Before | After | +|---|---|---| +| `inbound_route` | 0 | 0 | +| `inbound_dispatch` (new) | — | 0 | +| `inbound_pump_source_minimal` (poll `Source`) | 2 | 1 (the owned topic) | +| `outbound_scratch_static_topic` | 2 | 0 | +| `outbound_scratch_written_topic` (new) | — | 0 | +| `outbound_scratch_dynamic_topic` (`TopicProvider` adapter) | 3 | 1 | +| `outbound_owned_static_topic` | 3 | 1 | + +The bench asserts exact values, so every step updates `EXPECTED` and the +baseline together. A connector-level row per MQTT backend (loopback broker +for native, the existing loopback harness for embedded) records G2. The gate +becomes required on PRs touching `aimdb-core/src/{connector,transport,router}.rs`, +`session/`, or a connector's data path. + +## 6. Migration + +| Who | Change | +|---|---| +| Application code | None | +| Connector authors | `impl Source` → `InboundDispatch` or poll `Source`; `Connector::publish` → `route_publisher` + `RoutePublisher` (or the `BoxedPublisher` adapter) | +| `TopicProvider` users | None; optional move to `TopicWriter` | + +In-tree connectors to migrate and measure: MQTT (both backends), KNX, +WebSocket (client and server), the AimX session client (used by TCP, UDS and +serial), and the Embassy adapter's connector glue. + +**Order:** (1) core interfaces with adapters for the old forms, bench rows +added; (2) MQTT embedded, then native; (3) remaining connectors, one PR each; +(4) remove the adapters' use in-tree and make the gate required. + +## 7. Alternatives considered + +1. **`async fn` in traits with generic pumps.** Removes boxes by + monomorphization, but the pumps are type-erased per route (`dyn + SerializedSource`), and `Send` bounds on the returned futures need extra + machinery. Poll form keeps `dyn` and is already the codebase's pattern + (037). +2. **`ReusableBoxFuture` everywhere.** One allocation per route instead of + per message, but it needs `unsafe` or `tokio-util` on `no_std` (037 §3.2 + rejected the hand-rolled version) and keeps an indirection per poll. +3. **Lending async `Source` (`async fn next(&mut self) -> Option>`).** + The borrow must outlive `.await`, which conflicts with how both MQTT + clients expose messages. Push (`dispatch`) fits them directly. +4. **Keep `Connector::publish` and pool its futures.** Futures of different + connectors have different sizes; pooling adds unsafe layout handling for + no gain over `poll_ready`/`start_send`. + +## 8. Open questions + +- **Session-task latency (embedded).** Ingest moves into the MQTT session + task. It is synchronous and bounded, but a slow user deserializer now + delays keep-alive handling. Measure on the STM32H5 rig (037's B3). +- **Ring sizing (embedded outbound).** One shared ring or one per route; + default size; what `poll_ready` reports when one frame exceeds the ring. +- **Other connectors' own copies.** KNX, WebSocket and the AimX session were + not measured at connector level; §5's per-connector rows will show them. + +## 9. References + +- [037 — Zero-allocation consume path](./037-zero-alloc-consume-path.md) +- [045 — Per-link codec selection](./045-per-link-codec-selection.md) + (`with_serializer_into` scratch path) +- [052 — Runtime-neutral connectors](./052-runtime-neutral-connectors.md) + (`pump_source`, `pump_sink`, `StreamDialer`) +- `aimdb-bench/benches/b0_alloc_connector.rs` diff --git a/docs/design/055-wildcard-inbound-links.md b/docs/design/055-wildcard-inbound-links.md new file mode 100644 index 00000000..1b9f057d --- /dev/null +++ b/docs/design/055-wildcard-inbound-links.md @@ -0,0 +1,299 @@ +# 055 — Wildcard inbound links + +**Status:** 📝 Proposed + +**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. + +**Independent of** [054](./054-zero-alloc-connector-boundary.md). This +design ships on today's interfaces; §6 describes what changes when 054 lands. + +--- + +## 1. Problem + +An inbound link maps one exact topic to one record: + +```rust +reg.link_from("mqtt://sensors/kitchen/temp").with_deserializer(parse).finish(); +``` + +`Router::route` compares an incoming topic with each route's `resource_id` by +string equality. Records and their routes are fixed at `build()`, so a +publisher that was not known at startup has no route and its messages are +dropped. Applications with a changing set of devices must pre-declare a pool +of records and assign pool slots out of band. + +Two facts make a small change sufficient: + +- Both MQTT backends subscribe to `Router::resource_ids()` verbatim, so a + filter such as `sensors/+/temp` already subscribes correctly at the broker. + Only the router drops the messages. +- The deserializer receives `(ctx, bytes)` but not the topic, so it cannot + tell which publisher sent a matched message. + +## 2. Goals and non-goals + +**Goals** + +- G1. One inbound link with a topic pattern feeds many topics into one + record. +- 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. +- G4. No change to existing links, routers or connectors that do not use + patterns. + +**Non-goals** + +- Creating records at runtime. One record per pattern is the unit. +- Outbound topic templates (`link_to("mqtt://out/{device}")`). Possible + follow-up on 054's `TopicWriter`. +- The AimX/WebSocket wildcard grammar over record keys + (`session/topic_match.rs`). It matches dot-separated record keys, not broker + topics. + +## 3. 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) + .with_match_deserializer(|ctx, m, bytes| { + let device = m.get("device").unwrap(); // borrowed &str + let key = m.key(); // KeyId, when .key(..) is set + Reading::decode(key, bytes) + }) + .finish(); +}); +``` + +- `{name}` matches exactly one level. `{name..}` matches the remaining levels + (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). + +## 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: + +```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). + pub const EXACT: Self = /* … */; +} + +impl RouterBuilder { + pub fn with_grammar(self, grammar: TopicGrammar) -> Self; +} + +pub fn pump_source_with(db: &AimDb, scheme: &str, src: impl Source + 'static, + grammar: TopicGrammar) -> Vec; +// pump_source(..) == pump_source_with(.., 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. + +### 4.2 Compilation and matching + +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 +configuration errors: + +- a capture sharing a level with text (`sensors/dev-{id}`); +- a multi-level capture or wildcard that is not last; +- two captures with the same name; +- more than 8 captures. + +`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::resource_ids()` returns each pattern rendered in the grammar +(`sensors/{device}/temp` → `sensors/+/temp`), so connectors subscribe +unchanged. + +### 4.3 The match reaches the deserializer through `RuntimeContext` + +`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: + +```rust +impl RuntimeContext { + /// The inbound match being ingested, if any. + pub fn inbound_match(&self) -> Option<&TopicMatch>; +} + +pub struct TopicMatch { + topic: Arc, + spans: [(u16, u16); 8], + names: Arc<[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 +} +``` + +- `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()`. + +### 4.4 Keys + +```rust +#[derive(Copy, Clone, Eq, PartialEq, Hash, Debug)] +pub struct KeyId(NonZeroU16); // Option is 2 bytes + +.key("device", 1024) // capture name, capacity (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`. +- 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 + increases (logged; visible in record metadata with `observability`). + Capacity is therefore also an admission limit. +- Keys are never reused while the process runs. +- `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. + +### 4.5 MQTT specifics + +- `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)`. +- Outbound links reject patterns at `build()`: you cannot publish to a + filter. + +## 5. 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. + +When the broker binds topics to credentials (e.g. an ACL allowing a client +to publish only under `sensors//…`), the topic is the +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 + +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. +- The `TopicGrammar` struct may become a trait for inlining if the + 054 bench shows the pattern walk. + +## 7. 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 + breaking window that removes the need. +4. **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 + changes; names cost nothing at runtime. + +## 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 + `{`. + +## 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 + `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 + `thumbv7em-none-eabihf` with no behaviour change. + +## 10. References + +- [018 — Dynamic MQTT topics](./018-M7-dynamic-mqtt-topics.md) +- [052 — Runtime-neutral connectors](./052-runtime-neutral-connectors.md) + (`pump_source`, shared by both MQTT backends) +- [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) diff --git a/docs/design/056-aimdb-join-cli.md b/docs/design/056-aimdb-join-cli.md new file mode 100644 index 00000000..5c3b63a7 --- /dev/null +++ b/docs/design/056-aimdb-join-cli.md @@ -0,0 +1,138 @@ +# 056 — `aimdb join` CLI + +**Status:** 📝 Proposed + +**Scope:** the `join` subcommand of `tools/aimdb-cli`, speaking the +provisioning format of [043](./043-join-endpoint-v1.md) (rev 2). Client only; +provisioning services are out of scope. + +--- + +## 1. Starting point + +An implementation exists on the `planning` branch (`e0dcf9b`, +`tools/aimdb-cli/src/commands/join.rs`, ~480 lines with 9 tests), written +against 043 rev 1. Ported onto `main` it compiles unchanged and its tests +pass. A review found five defects to fix before release: + +| # | Defect | Effect | +|---|---|---| +| D1 | The GitHub OAuth client ID is read with `option_env!` at build time | Binaries built without the variable (including `cargo install`) cannot run the device flow | +| D2 | `prompt()` treats end of input as an empty answer and asks again | `aimdb join … [--token ] [--app =]... + [--out ] [--force] +``` + +| Option | Behaviour | +|---|---| +| `` | Deployment base URL. Must be `https://`, except loopback hosts (`localhost`, `127.0.0.0/8`, `::1`) for local testing (D5). | +| `--token` | Use the `claim-token` kind; skips the device flow. | +| `--app name=value` | Repeatable. Values for the deployment's `app_fields` (043 §2). Never prompted (D2, D3). | +| `--out` | Profile path. Default `aimdb-profile.toml`. | +| `--force` | Replace an existing profile. Without it, an existing file is an error (D4). | + +## 3. Flow + +1. **Discover.** `GET /v1/join`. Read the accepted auth kinds and + `app_fields`. A 404 means `claim-token` only. +2. **Check inputs.** `--app` names not listed in `app_fields` produce a + warning; required fields that are missing are an error before any + network call to GitHub. +3. **Authenticate.** + - `--token` given: `claim-token`. + - Otherwise, `github` from discovery: run the OAuth device flow with the + **discovered** `client_id` (D1), no scopes. Print the user code and + verification URL, poll respecting `interval` and `slow_down`, stop at + `expires_in`. Print the GitHub login for confirmation. + - Neither available: error naming the kinds the deployment accepts. +4. **Join.** `POST /v1/join` with `auth` and `app`. +5. **Write.** On `200`, write the profile (§4) and print `node_id` and the + path. On an error status, print the body's `message` as-is and exit with + code 1 (§5). + +The token is held in memory only: never printed, logged or written. + +## 4. Writing the profile + +- Serialize the response as TOML (043 §5). +- Write to a temporary file in the target directory, created with mode + `0600` on Unix (`OpenOptions::mode` + `create_new`), then `fsync` and + rename over the target. The mode is correct from the first byte whether or + not a file existed (D4), and a crash never leaves a half-written profile. +- Refuse an existing target without `--force`. +- Warn when `profile_version` is newer than the CLI knows; write anyway + (043 §6). + +## 5. Exit codes + +| Code | Meaning | +|---|---| +| 0 | Profile written | +| 1 | Rejected by the deployment (403, 409, 422, 429, 503) — its `message` is printed | +| 2 | Usage error: bad URL, `http://` to a non-loopback host, existing profile without `--force`, missing required `--app` | +| 3 | Network, TLS or unexpected response | +| 4 | Device flow failed: denied, expired | + +Distinct codes let scripts retry network failures and not rejections. + +## 6. Build and dependencies + +- Behind a `join` cargo feature, on by default. +- `reqwest` with `rustls` and the platform trust store + (`rustls-tls-native-roots`), `json` feature; no OpenSSL. The bundled + `webpki-roots` are avoided because their license is not on the + `deny.toml` allowlist (as on `planning`, `5e728a8`). +- `toml` for the profile. +- No OAuth client ID or other deployment value is compiled in. + +## 7. Tests + +Kept from `planning` (adapted to rev 2): response parsing, unknown fields +ignored, error `message` surfaced as-is, malformed success body, mock-server +round trip, profile mode `0600`. + +Added: + +1. Discovery: `github` + `client_id` used for the device flow; 404 → + `claim-token` only; neither available → exit code 2 with the accepted + kinds. +2. `--app` parsing (`=` inside values, repeated names rejected), unknown + names warned, missing required names rejected before any request. +3. `http://example.com` rejected; `http://127.0.0.1:…` and + `http://localhost:…` accepted. +4. Existing profile: refused without `--force`; with `--force`, a + pre-existing `0644` file ends up `0600`. +5. Closed stdin (`