From b6c8ec30604e175f95baf49adae8781f80022199 Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 5 Oct 2026 19:55:24 +0200 Subject: [PATCH 1/2] feat(ingestor): accept client RX coverage only from configured sources MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Trust in meshcore/client/{PUBLIC_KEY}/packets rests entirely on the broker binding the topic pubkey to the publisher holding that key. An instance that reads several brokers can mix one that enforces that binding with a legacy username/password broker where many accounts may write meshcore/#: enabling coverage there trusts the weakest source, and any account on it could inject coverage under any companion pubkey with any GPS position. clientRxCoverage.sources is an optional allowlist of mqttSources[].name. When set and non-empty, the client namespace is handled only for messages that arrived on a listed source; a message from any other source is dropped before any write (no client_receptions, no client_observers, and no observer row — the namespace still always returns), logged per source with a throttle and a line cap. Absent or empty keeps every source accepted, so the default is unchanged. A name matching no configured source is reported once at startup, since an allowlist that can never match would otherwise drop all coverage silently. The blacklist check still runs first, so it holds whatever the source. Relates to #265 --- cmd/ingestor/client_rx_sources.go | 120 +++++++++++ cmd/ingestor/client_rx_sources_test.go | 274 +++++++++++++++++++++++++ cmd/ingestor/config.go | 53 +++++ cmd/ingestor/main.go | 15 ++ config.example.json | 1 + docs/client-rx-coverage.md | 47 ++++- 6 files changed, 508 insertions(+), 2 deletions(-) create mode 100644 cmd/ingestor/client_rx_sources.go create mode 100644 cmd/ingestor/client_rx_sources_test.go diff --git a/cmd/ingestor/client_rx_sources.go b/cmd/ingestor/client_rx_sources.go new file mode 100644 index 000000000..cb94de118 --- /dev/null +++ b/cmd/ingestor/client_rx_sources.go @@ -0,0 +1,120 @@ +package main + +import ( + "log" + "strings" + "sync" + "time" +) + +// Client-RX source allowlist (#265): startup validation and the throttled, +// bounded warning for coverage dropped because it arrived on a source that is +// not in clientRxCoverage.sources. +// +// Dropping in silence fails in the dangerous direction — an operator who +// mistypes a source name would lose all coverage with nothing in the log — but +// a line per message is not an option either: an unlisted broker publishing +// meshcore/client/# produces mesh-rate traffic. So each source is logged on its +// first drop, re-logged at most every clientRxSourceWarnInterval while drops +// continue, and never more than clientRxSourceWarnMax times for the lifetime of +// the process (the last line says so). Mirrors the observerIATAWhitelist +// warning in iata_drop_warn.go; the throttle table needs no size cap here +// because the key is a configured source name, not a publisher-controlled +// topic segment. + +const ( + // clientRxSourceWarnInterval is how long a source stays quiet after a + // logged drop. + clientRxSourceWarnInterval = 10 * time.Minute + // clientRxSourceWarnMax caps how many lines one source may ever emit, so a + // permanently misconfigured deployment cannot fill the log over weeks. + clientRxSourceWarnMax = 10 +) + +// clientRxSourceDropThrottle is the per-source throttle state. The zero value +// is ready for use. +type clientRxSourceDropThrottle struct { + mu sync.Mutex + state map[string]*clientRxSourceDropEntry +} + +type clientRxSourceDropEntry struct { + last time.Time // when this source was last logged + lines int // lines emitted for this source so far +} + +// shouldWarn reports whether a drop from source name should be logged at now, +// and whether this is the final line allowed for that source. +func (t *clientRxSourceDropThrottle) shouldWarn(name string, now time.Time) (warn, final bool) { + key := strings.ToLower(strings.TrimSpace(name)) + t.mu.Lock() + defer t.mu.Unlock() + if t.state == nil { + t.state = make(map[string]*clientRxSourceDropEntry) + } + e := t.state[key] + if e == nil { + e = &clientRxSourceDropEntry{} + t.state[key] = e + } + if e.lines >= clientRxSourceWarnMax { + return false, false + } + if e.lines > 0 && now.Sub(e.last) < clientRxSourceWarnInterval { + return false, false + } + e.lines++ + e.last = now + return true, e.lines == clientRxSourceWarnMax +} + +// warnClientRxSourceDrop logs a throttled warning for a client-namespace +// message dropped because its MQTT source is not in clientRxCoverage.sources. +// The source name is quoted (%q) so an odd config value cannot forge log lines. +func (c *Config) warnClientRxSourceDrop(tag, name string, now time.Time) { + if c == nil { + return + } + warn, final := c.clientRxSrcDrop.shouldWarn(name, now) + if !warn { + return + } + if final { + log.Printf("MQTT [%s] [client-rx] dropping client coverage: source %q not in clientRxCoverage.sources; line cap reached (%d), further drops from this source are silent", + tag, name, clientRxSourceWarnMax) + return + } + log.Printf("MQTT [%s] [client-rx] dropping client coverage: source %q not in clientRxCoverage.sources; further drops from this source suppressed for %s", + tag, name, clientRxSourceWarnInterval) +} + +// checkClientRxSources reports, once at startup, every clientRxCoverage.sources +// entry that matches no configured mqttSources[].name — an allowlist of names +// that can never match would silently drop all coverage. It also logs the +// effective restriction so the gate is visible in the boot log. Returns the +// unknown names (nil when there is no allowlist or every name matches). +func checkClientRxSources(cfg *Config, sources []MQTTSource) []string { + allow := cfg.ClientRxCoverageSources() + if len(allow) == 0 { + return nil + } + var unknown []string + for _, want := range allow { + found := false + for _, src := range sources { + if strings.EqualFold(strings.TrimSpace(src.Name), want) { + found = true + break + } + } + if !found { + unknown = append(unknown, want) + } + } + log.Printf("[client-rx] coverage restricted to %d MQTT source(s): %s", len(allow), strings.Join(allow, ", ")) + if len(unknown) > 0 { + log.Printf("[client-rx] WARNING: %d clientRxCoverage.sources name(s) match no configured mqttSources[].name: %s — coverage will never be accepted for those names; check the spelling", + len(unknown), strings.Join(unknown, ", ")) + } + return unknown +} diff --git a/cmd/ingestor/client_rx_sources_test.go b/cmd/ingestor/client_rx_sources_test.go new file mode 100644 index 000000000..9d6a94889 --- /dev/null +++ b/cmd/ingestor/client_rx_sources_test.go @@ -0,0 +1,274 @@ +package main + +import ( + "strings" + "testing" + "time" +) + +// Tests for the #265 client-RX source allowlist: when clientRxCoverage.sources +// is set and non-empty, meshcore/client/... is only handled for messages that +// arrived on a listed mqttSources[].name. All ingest assertions drive the real +// handleMessage path with a named source. + +// clientObserverCount counts the mobile-client name rows the coverage path +// writes (handleClientPacket upserts one per message carrying an "origin"). +func clientObserverCount(t *testing.T, s *Store) int { + t.Helper() + var n int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM client_observers`).Scan(&n); err != nil { + t.Fatal(err) + } + return n +} + +// dispatchFrom drives handleMessage with a named MQTT source and returns the +// row-count snapshots either side of it, so a drop can be asserted across EVERY +// table: no client_receptions, no client_observers, and no observer row from a +// fall-through to the observer path. +func dispatchFrom(t *testing.T, s *Store, cfg *Config, sourceName string, m *mockMessage) (before, after map[string]int) { + t.Helper() + src := MQTTSource{Name: sourceName, Broker: "tcp://broker.invalid:1883"} + tag := sourceName + if tag == "" { + tag = src.Broker + } + before = dataTableCounts(t, s) + handleMessage(s, tag, src, m, nil, nil, cfg) + after = dataTableCounts(t, s) + return before, after +} + +// coverageCfgWithSources builds an enabled coverage config with the given +// source allowlist (no arguments ⇒ no allowlist). +func coverageCfgWithSources(sources ...string) *Config { + return &Config{ClientRxCoverage: &ClientRxCoverageConfig{Enabled: true, Sources: sources}} +} + +// TestClientRxSourceAllowedWithoutList pins the upstream-compatible default: +// an absent, empty or blank-only list accepts every source. +func TestClientRxSourceAllowedWithoutList(t *testing.T) { + cases := []struct { + name string + cfg *Config + }{ + {"nil clientRxCoverage", &Config{}}, + {"absent sources", coverageCfgWithSources()}, + {"blank-only sources", coverageCfgWithSources(" ", "")}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + for _, src := range []string{"device-auth", "legacy", ""} { + if !tc.cfg.ClientRxSourceAllowed(src) { + t.Errorf("source %q must be allowed when no allowlist is configured", src) + } + } + }) + } +} + +// TestClientRxSourceAllowedWithList covers the matching rule: listed names +// (including case and whitespace variants) pass, everything else is rejected — +// an unnamed source included, since it can never appear on the list. +func TestClientRxSourceAllowedWithList(t *testing.T) { + cfg := coverageCfgWithSources("device-auth", " Secondary ") + for _, name := range []string{"device-auth", "DEVICE-AUTH", " device-auth ", "secondary", "Secondary"} { + if !cfg.ClientRxSourceAllowed(name) { + t.Errorf("source %q must be allowed", name) + } + } + for _, name := range []string{"legacy", "device-auth2", "device", "", " "} { + if cfg.ClientRxSourceAllowed(name) { + t.Errorf("source %q must be rejected", name) + } + } +} + +// TestClientRxCoverageListedSourceIngests: the packet that is dropped from an +// unlisted source (next test) is ingested when the source is on the allowlist. +func TestClientRxCoverageListedSourceIngests(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources("device-auth", "other") + + dispatchFrom(t, store, cfg, "device-auth", clientCoverageMsg()) + + if n := clientReceptionCount(t, store); n != 1 { + t.Fatalf("listed source: expected 1 client_receptions row, got %d", n) + } + if n := clientObserverCount(t, store); n != 1 { + t.Fatalf("listed source: expected 1 client_observers row, got %d", n) + } +} + +// TestClientRxCoverageUnlistedSourceDropsEverything is the mutant killer for the +// source check: remove the check in handleMessage and these counts change. +// Nothing may be written anywhere — not client_receptions, not client_observers, +// and no observer row (the client namespace still returns). +func TestClientRxCoverageUnlistedSourceDropsEverything(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources("device-auth") + + before, after := dispatchFrom(t, store, cfg, "legacy", clientCoverageMsg()) + + assertOnlyDeltas(t, before, after, nil) + if n := clientReceptionCount(t, store); n != 0 { + t.Fatalf("unlisted source: expected 0 client_receptions rows, got %d", n) + } + if n := clientObserverCount(t, store); n != 0 { + t.Fatalf("unlisted source: expected 0 client_observers rows, got %d", n) + } +} + +// TestClientRxCoverageUnnamedSourceDropped: a source with no name cannot be on +// the allowlist, so with one configured it must not contribute coverage. +func TestClientRxCoverageUnnamedSourceDropped(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources("device-auth") + + before, after := dispatchFrom(t, store, cfg, "", clientCoverageMsg()) + + assertOnlyDeltas(t, before, after, nil) +} + +// TestClientRxCoverageNoAllowlistUnchanged: without sources, coverage is +// accepted from an arbitrary source name, exactly as before #265. +func TestClientRxCoverageNoAllowlistUnchanged(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources() // enabled, no allowlist + + dispatchFrom(t, store, cfg, "any-broker", clientCoverageMsg()) + + if n := clientReceptionCount(t, store); n != 1 { + t.Fatalf("no allowlist: expected 1 client_receptions row, got %d", n) + } +} + +// TestClientRxCoverageBlacklistBeatsAllowlist: the blacklist rule holds whatever +// the source — a blacklisted companion on a LISTED source is still dropped. +func TestClientRxCoverageBlacklistBeatsAllowlist(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources("device-auth") + cfg.ObserverBlacklist = []string{testCompanionPK} + + before, after := dispatchFrom(t, store, cfg, "device-auth", clientCoverageMsg()) + + assertOnlyDeltas(t, before, after, nil) +} + +// TestClientRxCoverageUnlistedSourceOtherSubtopicDropped: a non-"packets" client +// sub-topic from an unlisted source must also write nothing — the "client +// namespace always returns" rule applies whatever the source. +func TestClientRxCoverageUnlistedSourceOtherSubtopicDropped(t *testing.T) { + store := newTestStore(t) + cfg := coverageCfgWithSources("device-auth") + msg := &mockMessage{ + topic: "meshcore/client/" + testCompanionPK + "/status", + payload: []byte(`{"origin":"MyMob","noise_floor":-100}`), + } + + before, after := dispatchFrom(t, store, cfg, "legacy", msg) + + assertOnlyDeltas(t, before, after, nil) +} + +// clientRxDropLines returns the per-source drop warnings in captured log output. +func clientRxDropLines(out string) []string { + var lines []string + for _, l := range strings.Split(out, "\n") { + if strings.Contains(l, "not in clientRxCoverage.sources") { + lines = append(lines, l) + } + } + return lines +} + +// TestClientRxSourceDropLogThrottled: a burst of drops inside the throttle +// window produces one line, not one per message; the window reopens afterwards. +func TestClientRxSourceDropLogThrottled(t *testing.T) { + cfg := coverageCfgWithSources("device-auth") + t0 := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + + out := captureLog(t, func() { + for i := 0; i < 500; i++ { + cfg.warnClientRxSourceDrop("legacy", "legacy", t0.Add(time.Duration(i)*time.Millisecond)) + } + }) + if n := len(clientRxDropLines(out)); n != 1 { + t.Fatalf("burst inside the window: got %d lines, want 1:\n%s", n, out) + } + if !strings.Contains(out, `"legacy"`) { + t.Fatalf("the warning must name the rejected source:\n%s", out) + } + + // A second source gets its own first line; the window is per source. + out = captureLog(t, func() { + cfg.warnClientRxSourceDrop("other", "other", t0.Add(time.Second)) + }) + if n := len(clientRxDropLines(out)); n != 1 { + t.Fatalf("second source: got %d lines, want 1:\n%s", n, out) + } + + // Past the interval the same source may warn again. + out = captureLog(t, func() { + cfg.warnClientRxSourceDrop("legacy", "legacy", t0.Add(clientRxSourceWarnInterval)) + }) + if n := len(clientRxDropLines(out)); n != 1 { + t.Fatalf("after the interval: got %d lines, want 1:\n%s", n, out) + } +} + +// TestClientRxSourceDropLogCapped: the number of lines one source may ever emit +// is bounded, so a permanently misconfigured deployment cannot fill the log. +func TestClientRxSourceDropLogCapped(t *testing.T) { + cfg := coverageCfgWithSources("device-auth") + t0 := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + + out := captureLog(t, func() { + // Well past the cap, every call a full interval apart. + for i := 0; i < clientRxSourceWarnMax*5; i++ { + cfg.warnClientRxSourceDrop("legacy", "legacy", t0.Add(time.Duration(i)*clientRxSourceWarnInterval)) + } + }) + lines := clientRxDropLines(out) + if len(lines) != clientRxSourceWarnMax { + t.Fatalf("cap: got %d lines, want %d:\n%s", len(lines), clientRxSourceWarnMax, out) + } + if !strings.Contains(lines[len(lines)-1], "line cap reached") { + t.Fatalf("the last allowed line must say the cap was reached:\n%s", lines[len(lines)-1]) + } +} + +// TestCheckClientRxSourcesUnknownName: a name matching no configured source is +// reported once at startup; a configured name (any case) is not. +func TestCheckClientRxSourcesUnknownName(t *testing.T) { + sources := []MQTTSource{{Name: "device-auth"}, {Name: "Legacy"}} + + cfg := coverageCfgWithSources("device-auth", "legacy", "typo-broker") + var unknown []string + out := captureLog(t, func() { unknown = checkClientRxSources(cfg, sources) }) + if len(unknown) != 1 || unknown[0] != "typo-broker" { + t.Fatalf("unknown names = %v, want [typo-broker]", unknown) + } + if n := strings.Count(out, "WARNING"); n != 1 { + t.Fatalf("expected exactly one warning line, got %d:\n%s", n, out) + } + if !strings.Contains(out, "typo-broker") { + t.Fatalf("the warning must name the unknown entry:\n%s", out) + } +} + +// TestCheckClientRxSourcesSilentWithoutAllowlist: no allowlist ⇒ nothing to +// validate and nothing logged at startup. +func TestCheckClientRxSourcesSilentWithoutAllowlist(t *testing.T) { + sources := []MQTTSource{{Name: "device-auth"}} + for _, cfg := range []*Config{{}, coverageCfgWithSources(), coverageCfgWithSources(" ")} { + var unknown []string + out := captureLog(t, func() { unknown = checkClientRxSources(cfg, sources) }) + if len(unknown) != 0 { + t.Fatalf("no allowlist: unknown = %v, want none", unknown) + } + if strings.TrimSpace(out) != "" { + t.Fatalf("no allowlist: expected no log output, got:\n%s", out) + } + } +} diff --git a/cmd/ingestor/config.go b/cmd/ingestor/config.go index a54fdf778..cfc991092 100644 --- a/cmd/ingestor/config.go +++ b/cmd/ingestor/config.go @@ -84,6 +84,10 @@ type Config struct { // iataDropWarn is the bounded per-region throttle for that warning. iataDropWarn iataDropThrottle + // clientRxSrcDrop is the per-source throttle for the client-RX source + // allowlist drop warning (#265). See client_rx_sources.go. + clientRxSrcDrop clientRxSourceDropThrottle + // ObserverBlacklist is a list of observer public keys to drop at ingest. // Messages from blacklisted observers are silently discarded — no DB writes, // no UpsertObserver, no observations, no metrics. @@ -155,6 +159,16 @@ func (f *ForeignAdvertConfig) IsDropMode() bool { // ClientRxCoverageConfig controls the opt-in mobile client-RX coverage feature. type ClientRxCoverageConfig struct { Enabled bool `json:"enabled"` + + // Sources optionally restricts which MQTT sources may contribute client RX + // coverage, by mqttSources[].name (#265). Trust in meshcore/client/... rests + // entirely on the broker binding the topic pubkey to the publisher, so on an + // instance that also reads a broker without that binding the namespace must + // be accepted from the trusted source only. + // + // Absent or empty (or only blank entries) means every source is accepted, + // which keeps the upstream-compatible default. See ClientRxSourceAllowed. + Sources []string `json:"sources,omitempty"` } // ClientRxCoverageEnabled reports whether the opt-in mobile client-RX coverage @@ -163,6 +177,45 @@ func (c *Config) ClientRxCoverageEnabled() bool { return c.ClientRxCoverage != nil && c.ClientRxCoverage.Enabled } +// ClientRxCoverageSources returns the configured client-RX source allowlist +// with blank entries dropped. An empty result means "no restriction". +func (c *Config) ClientRxCoverageSources() []string { + if c == nil || c.ClientRxCoverage == nil || len(c.ClientRxCoverage.Sources) == 0 { + return nil + } + out := make([]string, 0, len(c.ClientRxCoverage.Sources)) + for _, name := range c.ClientRxCoverage.Sources { + if name = strings.TrimSpace(name); name != "" { + out = append(out, name) + } + } + return out +} + +// ClientRxSourceAllowed reports whether client RX coverage arriving on the MQTT +// source called name may be handled (#265). Matching is against +// mqttSources[].name, trimmed and case-insensitive (as elsewhere in this config). +// +// With no allowlist configured every source is allowed — unchanged behaviour. +// With one configured a source that is not on it is rejected, including a source +// with no name (which cannot be listed, so it cannot be trusted). +// +// The list is scanned linearly: it is bounded by the number of configured MQTT +// sources (a handful), so a set would cost more than it saves. +func (c *Config) ClientRxSourceAllowed(name string) bool { + allow := c.ClientRxCoverageSources() + if len(allow) == 0 { + return true + } + name = strings.TrimSpace(name) + for _, a := range allow { + if strings.EqualFold(a, name) { + return true + } + } + return false +} + // RetentionConfig controls how long stale nodes are kept before being moved to inactive_nodes. type RetentionConfig struct { NodeDays int `json:"nodeDays"` diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index 4de249f02..a44c5d1f4 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -74,6 +74,10 @@ func main() { sources := cfg.ResolvedSources() + // #265: surface a clientRxCoverage.sources entry that matches no configured + // source once, at boot, instead of silently dropping all coverage. + checkClientRxSources(cfg, sources) + store, err := OpenStoreWithInterval(cfg.DBPath, cfg.MetricsSampleInterval()) if err != nil { log.Fatalf("db: %v", err) @@ -673,6 +677,17 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, return } if cfg.ClientRxCoverageEnabled() && len(parts) >= 4 && parts[3] == "packets" { + // Optional per-source allowlist (#265). Trust in this topic is the + // broker's binding of the topic pubkey to the publisher, which not + // every configured source necessarily provides; when the operator + // names the sources that do, coverage from any other source is + // dropped here — before any client_receptions / client_observers + // write — with a throttled, bounded warning. No allowlist ⇒ every + // source is accepted, as before. + if !cfg.ClientRxSourceAllowed(source.Name) { + cfg.warnClientRxSourceDrop(tag, source.Name, time.Now()) + return + } handleClientPacket(store, tag, parts[2], msg, channelKeys) } return diff --git a/config.example.json b/config.example.json index 8cfd4e5f9..6e3a7b4bc 100644 --- a/config.example.json +++ b/config.example.json @@ -398,6 +398,7 @@ "_comment_compression": "Opt-in HTTP gzip middleware + WebSocket permessage-deflate. Both default to false — enable ONLY when your upstream reverse proxy is NOT already compressing. gzip: enables the gzipMiddleware wrapper around the HTTP handler. websocket: sets gorilla websocket Upgrader.EnableCompression. level: gzip compression level 1..9 (1=BestSpeed, 9=BestCompression, default 6). minSizeBytes: advisory minimum response size below which compression would not pay off. contentTypes: MIME allow-list — only responses with these Content-Type values are compressed. Already-compressed types (image/*, video/*, audio/*, application/zip, application/x-gzip, application/pdf, application/octet-stream) are always skipped, as are responses whose handler already set Content-Encoding. Omit contentTypes to use the built-in default allow-list.", "clientRxCoverage": { "enabled": false }, "_comment_clientRxCoverage": "Opt-in mobile client-RX coverage (corescope-rx companions publishing GPS-tagged receptions to meshcore/client//packets). Default OFF: when disabled the ingestor writes no client_receptions, the /api/rx-coverage|rx-leaderboard|nodes/{pubkey}/rx-coverage endpoints 404, and the UI hides the Coverage dashboard + Reach overlay. Set enabled=true to turn it on. SINGLE FLAG, BOTH PROCESSES: the ingestor and server each parse this same config.json, so this one clientRxCoverage.enabled entry gates both the ingest write path and the read endpoints — set it once, not per-process. TRUST: the feature requires an ACL-capable broker binding meshcore/client/{pubkey}/packets to that publisher; without ACLs the companion GPS is spoofable (see docs/client-rx-coverage.md). Retention: see retention.clientRxDays. Companion app + setup: https://github.com/efiten/corescope-rx.", + "_comment_clientRxCoverage_sources": "OPTIONAL allow-list of mqttSources[].name entries that may contribute client RX coverage. Omitted or empty (the default) = every source is accepted, unchanged behaviour. Set it when this instance reads MORE THAN ONE broker and they do not all bind meshcore/client//packets to the publisher holding that key: coverage is only as trustworthy as the weakest broker, so on a mixed deployment (one broker with per-device identity + one legacy username/password broker where many accounts may write meshcore/#) any account on the weak broker could otherwise inject coverage under any companion pubkey with any GPS position. Listing the trusted source names keeps the client namespace on those sources — a client message from any other source is dropped before any write (no client_receptions, no client_observers, no observer row), logged per source with throttling and a line cap. Matching is on the source name, trimmed and case-insensitive; a source with no name can never be listed. A name matching no configured source is logged once at startup. Example: \"clientRxCoverage\": { \"enabled\": true, \"sources\": [\"device-auth\"] }. See docs/client-rx-coverage.md.", "_comment_channelKeys": "Hex keys for decrypting channel messages. Key name = channel display name. public channel key is well-known.", "_comment_hashChannels": "Channel names whose keys are derived via SHA256. Key = SHA256(name)[:16]. Listed here so the ingestor can auto-derive keys.", "hashRegions": [ diff --git a/docs/client-rx-coverage.md b/docs/client-rx-coverage.md index 484688963..edb2da83b 100644 --- a/docs/client-rx-coverage.md +++ b/docs/client-rx-coverage.md @@ -25,13 +25,53 @@ Coverage is **off by default**. To turn it on: publish **only** under its own pubkey (e.g. an EMQX ACL keyed on the connected client's identity). This is the trust boundary, not an optimization — see [Trust](#trust). The ingestor already subscribes under `meshcore/#`. -3. Optionally set `retention.clientRxDays` to bound the coverage tables (see +3. **If you read more than one broker: name the ones that enforce that ACL** in + `clientRxCoverage.sources` — see [Restricting which MQTT sources may contribute](#restricting-which-mqtt-sources-may-contribute). +4. Optionally set `retention.clientRxDays` to bound the coverage tables (see [Storage](#storage--client_receptions-ingestor-owned)). -4. Point your users at [corescope-rx](https://github.com/efiten/corescope-rx) and they start +5. Point your users at [corescope-rx](https://github.com/efiten/corescope-rx) and they start contributing. Results show on each node's Reach page (coverage toggle) and the `#/rx-coverage` dashboard. **Warn them first that their contribution is world-readable and a per-observer view can reconstruct their movements — see [Privacy](#privacy--contributor-location-is-public).** +## Restricting which MQTT sources may contribute + +`clientRxCoverage.sources` is an **optional** allowlist of `mqttSources[].name` values. When it is set +and non-empty, the ingestor handles `meshcore/client/...` **only** for messages that arrived on a +listed source; a message from any other source is dropped before any write — no `client_receptions` +row, no `client_observers` row, and (as always for this namespace) no observer row either. When the +option is absent or empty, every source is accepted, which is the default and unchanged behaviour. + +```json +"mqttSources": [ + { "name": "device-auth", "broker": "mqtts://…" }, + { "name": "legacy", "broker": "mqtts://…" } +], +"clientRxCoverage": { "enabled": true, "sources": ["device-auth"] } +``` + +**Why this matters on a multi-broker deployment.** Coverage is only as trustworthy as the broker ACL +that binds `meshcore/client/{PUBLIC_KEY}/packets` to the publisher holding that key (see +[Trust](#trust)). An instance often reads several brokers with *different* authentication: one that +binds the topic to a per-device identity, and a legacy username/password broker where many accounts +may write `meshcore/#`. Without this option, enabling coverage trusts the weakest of them — any +account on the legacy broker could inject coverage under any companion pubkey, with any GPS position. +The allowlist keeps the namespace on the sources that actually enforce the binding, while the legacy +broker keeps contributing ordinary observer traffic as before. + +Notes: + +- Matching is on the configured source **name**, trimmed and case-insensitive. A source with no `name` + can never be listed, so it is rejected whenever an allowlist is set. +- A name that matches no configured source is logged once at startup — the allowlist would otherwise + silently drop all coverage from that (mistyped) name. +- Drops are logged per source, throttled and capped at a fixed number of lines, so a busy unlisted + broker cannot flood the log. +- The option gates only the client namespace. The observer blacklist and every other ingest rule stay + in force whatever the source. +- Unlike `clientRxCoverage.enabled`, which both processes read, `sources` concerns the ingest write + path only — the read endpoints are gated by `enabled` alone. Restart the ingestor after changing it. + The rest of this document is the MQTT payload contract the companion app implements. ## Companion BLE source (verified against firmware) @@ -172,6 +212,9 @@ Server/ingestor-side defense-in-depth (these reduce blast radius but do **not** - The ingestor rejects any topic pubkey that is not lowercase hex before writing, and never falls back to a payload-supplied id (`cmd/ingestor/client_reception.go`, #2/#10). +- `clientRxCoverage.sources` narrows the namespace to the brokers that enforce the ACL, so one weak + source among several cannot inject coverage (see + [Restricting which MQTT sources may contribute](#restricting-which-mqtt-sources-may-contribute)). - A blacklisted operator cannot contribute via the client topic (the blacklist is enforced before the coverage write, #1). - The frontend HTML-escapes the pubkey it renders, so a junk pubkey can't inject markup (#14). From e75c19017b2a6f219ebf41855d746704dbee0922 Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 5 Oct 2026 19:58:37 +0200 Subject: [PATCH 2/2] fix(ingestor): say so when a client-RX allowlist is set but coverage is off The boot line reported "coverage restricted to N MQTT source(s)" even with clientRxCoverage.enabled false, where the allowlist is inert and no coverage is ingested from any source. Name that state instead of implying the listed sources are contributing. Relates to #265 --- cmd/ingestor/client_rx_sources.go | 8 +++++++- cmd/ingestor/client_rx_sources_test.go | 15 +++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/cmd/ingestor/client_rx_sources.go b/cmd/ingestor/client_rx_sources.go index cb94de118..95dc6933d 100644 --- a/cmd/ingestor/client_rx_sources.go +++ b/cmd/ingestor/client_rx_sources.go @@ -111,7 +111,13 @@ func checkClientRxSources(cfg *Config, sources []MQTTSource) []string { unknown = append(unknown, want) } } - log.Printf("[client-rx] coverage restricted to %d MQTT source(s): %s", len(allow), strings.Join(allow, ", ")) + state := "" + if !cfg.ClientRxCoverageEnabled() { + // The allowlist is inert while the feature is off; say so rather than + // implying coverage is being ingested from the listed sources. + state = " (clientRxCoverage.enabled is false, so no coverage is ingested at all)" + } + log.Printf("[client-rx] coverage restricted to %d MQTT source(s): %s%s", len(allow), strings.Join(allow, ", "), state) if len(unknown) > 0 { log.Printf("[client-rx] WARNING: %d clientRxCoverage.sources name(s) match no configured mqttSources[].name: %s — coverage will never be accepted for those names; check the spelling", len(unknown), strings.Join(unknown, ", ")) diff --git a/cmd/ingestor/client_rx_sources_test.go b/cmd/ingestor/client_rx_sources_test.go index 9d6a94889..37b756b21 100644 --- a/cmd/ingestor/client_rx_sources_test.go +++ b/cmd/ingestor/client_rx_sources_test.go @@ -255,6 +255,21 @@ func TestCheckClientRxSourcesUnknownName(t *testing.T) { if !strings.Contains(out, "typo-broker") { t.Fatalf("the warning must name the unknown entry:\n%s", out) } + if strings.Contains(out, "enabled is false") { + t.Fatalf("coverage is enabled here, the boot line must not say otherwise:\n%s", out) + } +} + +// TestCheckClientRxSourcesDisabledSaysSo: an allowlist set while the feature is +// off is inert, and the boot line must not imply coverage is being ingested. +func TestCheckClientRxSourcesDisabledSaysSo(t *testing.T) { + cfg := &Config{ClientRxCoverage: &ClientRxCoverageConfig{Enabled: false, Sources: []string{"device-auth"}}} + + out := captureLog(t, func() { checkClientRxSources(cfg, []MQTTSource{{Name: "device-auth"}}) }) + + if !strings.Contains(out, "enabled is false") { + t.Fatalf("expected the boot line to flag the disabled feature:\n%s", out) + } } // TestCheckClientRxSourcesSilentWithoutAllowlist: no allowlist ⇒ nothing to