diff --git a/cmd/ingestor/client_rx_sources.go b/cmd/ingestor/client_rx_sources.go new file mode 100644 index 000000000..95dc6933d --- /dev/null +++ b/cmd/ingestor/client_rx_sources.go @@ -0,0 +1,126 @@ +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) + } + } + 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, ", ")) + } + 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..37b756b21 --- /dev/null +++ b/cmd/ingestor/client_rx_sources_test.go @@ -0,0 +1,289 @@ +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) + } + 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 +// 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).