Skip to content

Commit ecd3bbe

Browse files
authored
Merge pull request #141 from dborup/codex/issue-118-mqtt-client-ids
fix(ingestor): assign explicit, collision-resistant MQTT client IDs
2 parents 751f8fc + 41675d9 commit ecd3bbe

28 files changed

Lines changed: 2032 additions & 168 deletions

‎.github/workflows/deploy.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -101,6 +101,7 @@ jobs:
101101
set -e -o pipefail
102102
(cd internal/channelregistry && go test ./...)
103103
(cd internal/dbschema && go test ./...)
104+
(cd internal/brokerurl && go test ./...)
104105
105106
# internal/anomaly is its own module (not wired into any binary), so the
106107
# server/ingestor steps above do not reach it. No -fuzz and no -bench:

‎Dockerfile‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ COPY internal/dbconfig/ ../../internal/dbconfig/
2222
COPY internal/dbschema/ ../../internal/dbschema/
2323
COPY internal/prunequeue/ ../../internal/prunequeue/
2424
COPY internal/perfio/ ../../internal/perfio/
25+
COPY internal/brokerurl/ ../../internal/brokerurl/
2526
COPY internal/mbcapqueue/ ../../internal/mbcapqueue/
2627
COPY internal/lora/ ../../internal/lora/
2728
COPY internal/regions/ ../../internal/regions/
@@ -41,6 +42,7 @@ COPY internal/dbconfig/ ../../internal/dbconfig/
4142
COPY internal/dbschema/ ../../internal/dbschema/
4243
COPY internal/prunequeue/ ../../internal/prunequeue/
4344
COPY internal/perfio/ ../../internal/perfio/
45+
COPY internal/brokerurl/ ../../internal/brokerurl/
4446
COPY internal/mbcapqueue/ ../../internal/mbcapqueue/
4547
COPY internal/regions/ ../../internal/regions/
4648
COPY internal/channelregistry/ ../../internal/channelregistry/

‎cmd/ingestor/README.md‎

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -57,8 +57,16 @@ a `procIO` block sampled from `/proc/self/io` (read/write/cancelled bytes per
5757
second + syscall counts). The server reads this file and surfaces the data on
5858
the Perf page so operators can self-diagnose write-volume anomalies.
5959

60-
The writer uses `O_NOFOLLOW | O_CREAT | O_TRUNC` mode `0o600`, so a
61-
pre-planted symlink at the path cannot be used to clobber an arbitrary file.
60+
The writer uses `O_NOFOLLOW | O_CREAT` mode `0o600`, so a pre-planted
61+
symlink at the path cannot be used to clobber an arbitrary file. A stale tmp
62+
file that belongs to another user is refused before anything is changed, also
63+
when the ingestor runs as root (as in Docker), where `chmod` alone would
64+
succeed; one of its own is forced to `0o600` and truncated. Broker URLs in the
65+
file (per-source `broker`, `name` and `lastError`, and the source tags) have
66+
their user-info replaced by `****` and carry no query or fragment, as in the
67+
log; a `lastError` also has the source's configured user name, password and
68+
URL user-info and query masked where it quotes them. The server masks once
69+
more before serving `/api/mqtt/status` and `/api/healthz`.
6270

6371
**Security note:** the default lives in `/tmp`, which is world-writable on
6472
most hosts (sticky bit only protects deletion, not creation). On
@@ -86,11 +94,12 @@ the corescope user can write to.
8694
The ingestor reads these fields from the existing `config.json`:
8795

8896
- `mqttSources[]` — array of MQTT broker connections
89-
- `name` — display name for logging
97+
- `name` — display name for logging. Without one the source is tagged with its broker, credentials masked as `****`, plus ` (2)`, ` (3)`, … when another source already has that tag
9098
- `broker` — MQTT URL (`mqtt://`, `mqtts://`)
9199
- `username` / `password` — auth credentials
92100
- `topics` — array of topic patterns to subscribe
93101
- `iataFilter` — optional regional filter
102+
- `clientId` — optional MQTT client ID, used verbatim; must be unique among the broker's concurrent clients. When omitted, each start generates `corescope-<name>-<8 random hex>` (broker host if the name is empty), kept across reconnects. Strict MQTT 3.1 brokers accept at most 23 characters: set a short `clientId` for those.
94103
- `mqtt` — legacy single-broker config (auto-converted to `mqttSources`)
95104
- `dbPath` — SQLite DB path (default: `data/meshcore.db`)
96105

‎cmd/ingestor/config.go‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,10 @@ type MQTTSource struct {
2525
IATAFilter []string `json:"iataFilter,omitempty"`
2626
ConnectTimeoutSec int `json:"connectTimeoutSec,omitempty"`
2727
Region string `json:"region,omitempty"`
28+
// ClientID is the MQTT client ID, used verbatim; it must be unique among
29+
// the broker's concurrent clients. Empty: generated per client (#118,
30+
// see mqtt_client_id.go).
31+
ClientID string `json:"clientId,omitempty"`
2832
}
2933

3034
// ConnectTimeoutOrDefault returns the per-source connect timeout in seconds,

‎cmd/ingestor/go.mod‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,3 +59,7 @@ replace github.com/meshcore-analyzer/regions => ../../internal/regions
5959
require github.com/meshcore-analyzer/channelregistry v0.0.0
6060

6161
replace github.com/meshcore-analyzer/channelregistry => ../../internal/channelregistry
62+
63+
require github.com/meshcore-analyzer/brokerurl v0.0.0
64+
65+
replace github.com/meshcore-analyzer/brokerurl => ../../internal/brokerurl

‎cmd/ingestor/main.go‎

Lines changed: 11 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -135,61 +135,13 @@ func main() {
135135
// Connect to each MQTT source
136136
var clients []mqtt.Client
137137
connectedCount := 0
138-
for _, source := range sources {
139-
tag := source.Name
140-
if tag == "" {
141-
tag = source.Broker
142-
}
143-
144-
opts := buildMQTTOpts(source)
138+
tags := mqttSourceTags(sources)
139+
for i, source := range sources {
140+
tag := tags[i]
141+
opts, status, liveness := prepareMQTTSource(source, tag)
145142
connectTimeout := source.ConnectTimeoutOrDefault()
146143
log.Printf("MQTT [%s] connect timeout: %ds", tag, connectTimeout)
147144

148-
// Pre-allocate the liveness pointer so OnConnect can reset its
149-
// stale-message clock on reconnect (PR #1216 r1 item 2). IsConnectedFn
150-
// is wired below once the client exists.
151-
liveness := &SourceLivenessState{
152-
Tag: tag,
153-
Broker: source.Broker,
154-
}
155-
156-
// #1043: per-source status registry. Idempotent — repeated
157-
// registration across reconnects returns the same state so
158-
// counters accumulate across the process lifetime.
159-
status := RegisterSourceStatus(tag, source.Broker)
160-
161-
opts.SetOnConnectHandler(func(c mqtt.Client) {
162-
log.Printf("MQTT [%s] connected to %s", tag, source.Broker)
163-
status.MarkConnect(time.Now())
164-
// PR #1216 r1 item 2: clear the stale LastMessageUnix from
165-
// before the outage so the watchdog doesn't immediately scream
166-
// "stalled for 2h". Also restarts the cold-start grace window
167-
// and clears the alert cooldown so a fresh stall edge can fire.
168-
liveness.MarkReconnected(time.Now())
169-
topics := source.Topics
170-
if len(topics) == 0 {
171-
topics = []string{"meshcore/#"}
172-
}
173-
for _, t := range topics {
174-
token := c.Subscribe(t, 0, nil)
175-
token.Wait()
176-
if token.Error() != nil {
177-
log.Printf("MQTT [%s] subscribe error for %s: %v", tag, t, token.Error())
178-
} else {
179-
log.Printf("MQTT [%s] subscribed to %s", tag, t)
180-
}
181-
}
182-
})
183-
184-
opts.SetConnectionLostHandler(func(c mqtt.Client, err error) {
185-
log.Printf("MQTT [%s] disconnected from %s: %v", tag, source.Broker, err)
186-
status.MarkDisconnect(time.Now(), err)
187-
})
188-
189-
opts.SetReconnectingHandler(func(c mqtt.Client, options *mqtt.ClientOptions) {
190-
log.Printf("MQTT [%s] reconnecting to %s", tag, source.Broker)
191-
})
192-
193145
// Capture source for closure
194146
src := source
195147
opts.SetDefaultPublishHandler(func(c mqtt.Client, m mqtt.Message) {
@@ -236,7 +188,7 @@ func main() {
236188
continue
237189
}
238190
if token.Error() != nil {
239-
log.Printf("MQTT [%s] connection failed (non-fatal): %v", tag, token.Error())
191+
log.Printf("MQTT [%s] connection failed (non-fatal): %s", tag, errForLog(token.Error(), mqttSourceSecrets(source)...))
240192
// BL1 fix: Disconnect to stop Paho's internal retry goroutines.
241193
// With ConnectRetry=true, Connect() spawns background goroutines
242194
// that leak if the client is simply discarded.
@@ -527,10 +479,7 @@ func main() {
527479
// #1212 (prod outage on 2026-05-15 where the disconnect was logged but no
528480
// reconnect activity was ever visible).
529481
func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions {
530-
tag := source.Name
531-
if tag == "" {
532-
tag = source.Broker
533-
}
482+
tag := mqttSourceTag(source)
534483
opts := mqtt.NewClientOptions().
535484
AddBroker(source.Broker).
536485
SetAutoReconnect(true).
@@ -547,7 +496,10 @@ func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions {
547496
// (paho default 30s actually — making this explicit so it can't
548497
// drift, and so operators reading the code know it's intentional
549498
// per the #1335 RCA).
550-
SetKeepAlive(30 * time.Second)
499+
SetKeepAlive(30 * time.Second).
500+
// #118: explicit ID, generated once per client when not configured;
501+
// paho reuses it on every reconnect. See mqtt_client_id.go.
502+
SetClientID(mqttClientID(source))
551503

552504
opts.SetConnectionAttemptHandler(func(broker *url.URL, tlsCfg *tls.Config) *tls.Config {
553505
// Look up the per-source liveness state (registered in main) so we
@@ -560,7 +512,7 @@ func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions {
560512
if s != nil {
561513
attempt = atomic.AddInt64(&s.AttemptCount, 1)
562514
}
563-
log.Printf("MQTT [%s] connection attempt #%d to %s", tag, attempt, broker.String())
515+
log.Printf("MQTT [%s] connection attempt #%d to %s", tag, attempt, brokerForLog(broker.String()))
564516
return tlsCfg
565517
})
566518

‎cmd/ingestor/mqtt_client_id.go‎

Lines changed: 197 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,197 @@
1+
package main
2+
3+
import (
4+
"crypto/rand"
5+
"encoding/binary"
6+
"encoding/hex"
7+
"io"
8+
"net/url"
9+
"sort"
10+
"strconv"
11+
"strings"
12+
"sync/atomic"
13+
"time"
14+
15+
"github.com/meshcore-analyzer/brokerurl"
16+
)
17+
18+
// MQTT client IDs (#118).
19+
//
20+
// Without SetClientID paho connects with a zero-length client ID and
21+
// CleanSession=true; whether the broker then assigns one, rejects the client
22+
// or lets two ingestors take over each other's session is broker-dependent.
23+
// Every source therefore gets an explicit ID:
24+
//
25+
// - mqttSources[].clientId, used verbatim. It must be unique among all
26+
// clients connected to that broker at the same time.
27+
// - otherwise corescope-<name>-<8 hex>, generated once per client
28+
// construction (buildMQTTOpts). <name> is the sanitized source name, else
29+
// the broker host, else omitted; the suffix is 32 random bits from
30+
// crypto/rand. paho reuses the options for its own reconnects and the
31+
// watchdog's force-reconnect reuses the same client, so the ID is stable
32+
// for the process lifetime but differs between sources and processes.
33+
//
34+
// MQTT 3.1.1 (what paho tries first) only guarantees IDs of 1–23 characters
35+
// from [0-9A-Za-z], although brokers such as Mosquitto and EMQX accept longer
36+
// ones. paho falls back to MQTT 3.1 after a failed first handshake, and a
37+
// strict 3.1 broker rejects IDs over 23 characters. The default is
38+
// 19 characters plus the name part, so for such a broker configure a short
39+
// clientId. (The previous empty ID was invalid under 3.1 as well.)
40+
41+
const mqttClientIDMaxBase = 32
42+
43+
// clientIDRandom is swapped in tests to simulate an entropy failure.
44+
var clientIDRandom io.Reader = rand.Reader
45+
46+
var clientIDFallbackSeq atomic.Uint64
47+
48+
// clientIDNow feeds the fallback suffix; swapped in tests to model a coarse
49+
// clock that returns the same reading twice.
50+
var clientIDNow = time.Now
51+
52+
// mqttClientID returns the client ID for a new client of this source.
53+
func mqttClientID(source MQTTSource) string {
54+
if strings.TrimSpace(source.ClientID) != "" {
55+
return source.ClientID
56+
}
57+
base := sanitizeClientIDPart(source.Name)
58+
if base == "" {
59+
// the masked broker: url.Parse alone may take a user name for
60+
// the host (see brokerForLog)
61+
if u, err := url.Parse(brokerForLog(source.Broker)); err == nil {
62+
base = sanitizeClientIDPart(u.Hostname())
63+
}
64+
}
65+
id := "corescope-"
66+
if base != "" {
67+
id += base + "-"
68+
}
69+
return id + clientIDSuffix()
70+
}
71+
72+
// sanitizeClientIDPart lowercases s, keeps [a-z0-9], turns every other run
73+
// into one '-', trims '-' and caps the result at mqttClientIDMaxBase.
74+
func sanitizeClientIDPart(s string) string {
75+
var b strings.Builder
76+
dash := false
77+
for _, r := range strings.ToLower(s) {
78+
if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') {
79+
if dash && b.Len() > 0 {
80+
b.WriteByte('-')
81+
}
82+
dash = false
83+
b.WriteRune(r)
84+
continue
85+
}
86+
dash = true
87+
}
88+
out := b.String()
89+
if len(out) > mqttClientIDMaxBase {
90+
out = strings.TrimRight(out[:mqttClientIDMaxBase], "-")
91+
}
92+
return out
93+
}
94+
95+
// clientIDSuffix is 8 hex chars from crypto/rand. Should that ever fail it
96+
// mixes the clock and a process-wide counter instead, so two constructions
97+
// still never share a suffix within a process.
98+
func clientIDSuffix() string {
99+
var b [4]byte
100+
if _, err := io.ReadFull(clientIDRandom, b[:]); err != nil {
101+
v := uint64(clientIDNow().UnixNano()) ^ (clientIDFallbackSeq.Add(1) * 0x9E3779B97F4A7C15)
102+
binary.BigEndian.PutUint32(b[:], uint32(v>>32)^uint32(v))
103+
}
104+
return hex.EncodeToString(b[:])
105+
}
106+
107+
// mqttConnectedLogLine is the "connected" log line. The broker URL is logged
108+
// with its user-info masked and without query or fragment (brokerForLog), so
109+
// credentials or device tokens embedded in it never reach the log.
110+
func mqttConnectedLogLine(tag, broker, clientID string) string {
111+
return "MQTT [" + tag + "] connected to " + brokerForLog(broker) + " as client " + clientID
112+
}
113+
114+
// mqttSourceTag is the base of a source's tag in logs and the
115+
// liveness/status registries: its name or, for an unnamed source, the
116+
// broker with its credentials masked (brokerForLog). mqttSourceTags makes the tags
117+
// of a whole configuration unique.
118+
func mqttSourceTag(source MQTTSource) string {
119+
if source.Name != "" {
120+
return source.Name
121+
}
122+
return brokerForLog(source.Broker)
123+
}
124+
125+
// mqttSourceTags returns the tag of every source. Unnamed sources on the
126+
// same broker (say, with different credentials) would share one tag: the
127+
// second would lose watchdog tracking and the two would share status
128+
// counters. So an unnamed source whose tag is taken, by any named source or
129+
// an earlier unnamed one, gets " (2)", " (3)", … A duplicate Name is left as
130+
// it is: that is a configuration error, reported by registerLivenessOrSkip.
131+
func mqttSourceTags(sources []MQTTSource) []string {
132+
used := make(map[string]bool, len(sources))
133+
for _, s := range sources {
134+
if s.Name != "" {
135+
used[s.Name] = true
136+
}
137+
}
138+
tags := make([]string, len(sources))
139+
for i, s := range sources {
140+
tag := mqttSourceTag(s)
141+
if s.Name == "" {
142+
base := tag
143+
for n := 2; used[tag]; n++ {
144+
tag = base + " (" + strconv.Itoa(n) + ")"
145+
}
146+
used[tag] = true
147+
}
148+
tags[i] = tag
149+
}
150+
return tags
151+
}
152+
153+
// brokerForLog returns broker with its user-info replaced by "****" and
154+
// without query or fragment (brokerurl.Mask), so credentials or tokens
155+
// embedded in it never reach a log, the stats file or a client ID, while the
156+
// "****@" shows that the URL carries credentials and that the host may be
157+
// cut short. A broker without a scheme is read as tcp://, as paho's
158+
// AddBroker does.
159+
func brokerForLog(broker string) string {
160+
if !strings.Contains(broker, "://") {
161+
broker = "tcp://" + broker
162+
}
163+
return brokerurl.Mask(broker)
164+
}
165+
166+
// mqttSourceSecrets returns the non-empty secrets of a source: its
167+
// password and user name, and the user-info, query and fragment of its
168+
// broker URL as configured (brokerurl.Secrets).
169+
func mqttSourceSecrets(source MQTTSource) []string {
170+
var out []string
171+
for _, v := range append([]string{source.Password, source.Username}, brokerurl.Secrets(source.Broker)...) {
172+
if v != "" {
173+
out = append(out, v)
174+
}
175+
}
176+
return out
177+
}
178+
179+
// errForLog is err's text with each of secrets (mqttSourceSecrets) replaced
180+
// by "****", longest first, and then any broker URL or user-info it quotes
181+
// masked (brokerurl.MaskText): MaskText only spots URL-shaped text, so a
182+
// query token or a password quoted on its own would pass it. paho's errors
183+
// normally quote none, but they reach the log and the stats file.
184+
func errForLog(err error, secrets ...string) string {
185+
if err == nil {
186+
return "<nil>"
187+
}
188+
s := err.Error()
189+
sorted := append([]string(nil), secrets...)
190+
sort.Slice(sorted, func(i, j int) bool { return len(sorted[i]) > len(sorted[j]) })
191+
for _, v := range sorted {
192+
if v != "" {
193+
s = strings.ReplaceAll(s, v, brokerurl.Marker)
194+
}
195+
}
196+
return brokerurl.MaskText(s)
197+
}

0 commit comments

Comments
 (0)