Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/workflows/deploy.yml
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ jobs:
set -e -o pipefail
(cd internal/channelregistry && go test ./...)
(cd internal/dbschema && go test ./...)
(cd internal/brokerurl && go test ./...)

# internal/anomaly is its own module (not wired into any binary), so the
# server/ingestor steps above do not reach it. No -fuzz and no -bench:
Expand Down
2 changes: 2 additions & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ COPY internal/dbconfig/ ../../internal/dbconfig/
COPY internal/dbschema/ ../../internal/dbschema/
COPY internal/prunequeue/ ../../internal/prunequeue/
COPY internal/perfio/ ../../internal/perfio/
COPY internal/brokerurl/ ../../internal/brokerurl/
COPY internal/mbcapqueue/ ../../internal/mbcapqueue/
COPY internal/lora/ ../../internal/lora/
COPY internal/regions/ ../../internal/regions/
Expand All @@ -41,6 +42,7 @@ COPY internal/dbconfig/ ../../internal/dbconfig/
COPY internal/dbschema/ ../../internal/dbschema/
COPY internal/prunequeue/ ../../internal/prunequeue/
COPY internal/perfio/ ../../internal/perfio/
COPY internal/brokerurl/ ../../internal/brokerurl/
COPY internal/mbcapqueue/ ../../internal/mbcapqueue/
COPY internal/regions/ ../../internal/regions/
COPY internal/channelregistry/ ../../internal/channelregistry/
Expand Down
15 changes: 12 additions & 3 deletions cmd/ingestor/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,16 @@ a `procIO` block sampled from `/proc/self/io` (read/write/cancelled bytes per
second + syscall counts). The server reads this file and surfaces the data on
the Perf page so operators can self-diagnose write-volume anomalies.

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

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

- `mqttSources[]` — array of MQTT broker connections
- `name` — display name for logging
- `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
- `broker` — MQTT URL (`mqtt://`, `mqtts://`)
- `username` / `password` — auth credentials
- `topics` — array of topic patterns to subscribe
- `iataFilter` — optional regional filter
- `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.
- `mqtt` — legacy single-broker config (auto-converted to `mqttSources`)
- `dbPath` — SQLite DB path (default: `data/meshcore.db`)

Expand Down
4 changes: 4 additions & 0 deletions cmd/ingestor/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ type MQTTSource struct {
IATAFilter []string `json:"iataFilter,omitempty"`
ConnectTimeoutSec int `json:"connectTimeoutSec,omitempty"`
Region string `json:"region,omitempty"`
// ClientID is the MQTT client ID, used verbatim; it must be unique among
// the broker's concurrent clients. Empty: generated per client (#118,
// see mqtt_client_id.go).
ClientID string `json:"clientId,omitempty"`
}

// ConnectTimeoutOrDefault returns the per-source connect timeout in seconds,
Expand Down
4 changes: 4 additions & 0 deletions cmd/ingestor/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -59,3 +59,7 @@ replace github.com/meshcore-analyzer/regions => ../../internal/regions
require github.com/meshcore-analyzer/channelregistry v0.0.0

replace github.com/meshcore-analyzer/channelregistry => ../../internal/channelregistry

require github.com/meshcore-analyzer/brokerurl v0.0.0

replace github.com/meshcore-analyzer/brokerurl => ../../internal/brokerurl
70 changes: 11 additions & 59 deletions cmd/ingestor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -135,61 +135,13 @@ func main() {
// Connect to each MQTT source
var clients []mqtt.Client
connectedCount := 0
for _, source := range sources {
tag := source.Name
if tag == "" {
tag = source.Broker
}

opts := buildMQTTOpts(source)
tags := mqttSourceTags(sources)
for i, source := range sources {
tag := tags[i]
opts, status, liveness := prepareMQTTSource(source, tag)
connectTimeout := source.ConnectTimeoutOrDefault()
log.Printf("MQTT [%s] connect timeout: %ds", tag, connectTimeout)

// Pre-allocate the liveness pointer so OnConnect can reset its
// stale-message clock on reconnect (PR #1216 r1 item 2). IsConnectedFn
// is wired below once the client exists.
liveness := &SourceLivenessState{
Tag: tag,
Broker: source.Broker,
}

// #1043: per-source status registry. Idempotent — repeated
// registration across reconnects returns the same state so
// counters accumulate across the process lifetime.
status := RegisterSourceStatus(tag, source.Broker)

opts.SetOnConnectHandler(func(c mqtt.Client) {
log.Printf("MQTT [%s] connected to %s", tag, source.Broker)
status.MarkConnect(time.Now())
// PR #1216 r1 item 2: clear the stale LastMessageUnix from
// before the outage so the watchdog doesn't immediately scream
// "stalled for 2h". Also restarts the cold-start grace window
// and clears the alert cooldown so a fresh stall edge can fire.
liveness.MarkReconnected(time.Now())
topics := source.Topics
if len(topics) == 0 {
topics = []string{"meshcore/#"}
}
for _, t := range topics {
token := c.Subscribe(t, 0, nil)
token.Wait()
if token.Error() != nil {
log.Printf("MQTT [%s] subscribe error for %s: %v", tag, t, token.Error())
} else {
log.Printf("MQTT [%s] subscribed to %s", tag, t)
}
}
})

opts.SetConnectionLostHandler(func(c mqtt.Client, err error) {
log.Printf("MQTT [%s] disconnected from %s: %v", tag, source.Broker, err)
status.MarkDisconnect(time.Now(), err)
})

opts.SetReconnectingHandler(func(c mqtt.Client, options *mqtt.ClientOptions) {
log.Printf("MQTT [%s] reconnecting to %s", tag, source.Broker)
})

// Capture source for closure
src := source
opts.SetDefaultPublishHandler(func(c mqtt.Client, m mqtt.Message) {
Expand Down Expand Up @@ -236,7 +188,7 @@ func main() {
continue
}
if token.Error() != nil {
log.Printf("MQTT [%s] connection failed (non-fatal): %v", tag, token.Error())
log.Printf("MQTT [%s] connection failed (non-fatal): %s", tag, errForLog(token.Error(), mqttSourceSecrets(source)...))
// BL1 fix: Disconnect to stop Paho's internal retry goroutines.
// With ConnectRetry=true, Connect() spawns background goroutines
// that leak if the client is simply discarded.
Expand Down Expand Up @@ -527,10 +479,7 @@ func main() {
// #1212 (prod outage on 2026-05-15 where the disconnect was logged but no
// reconnect activity was ever visible).
func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions {
tag := source.Name
if tag == "" {
tag = source.Broker
}
tag := mqttSourceTag(source)
opts := mqtt.NewClientOptions().
AddBroker(source.Broker).
SetAutoReconnect(true).
Expand All @@ -547,7 +496,10 @@ func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions {
// (paho default 30s actually — making this explicit so it can't
// drift, and so operators reading the code know it's intentional
// per the #1335 RCA).
SetKeepAlive(30 * time.Second)
SetKeepAlive(30 * time.Second).
// #118: explicit ID, generated once per client when not configured;
// paho reuses it on every reconnect. See mqtt_client_id.go.
SetClientID(mqttClientID(source))

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

Expand Down
197 changes: 197 additions & 0 deletions cmd/ingestor/mqtt_client_id.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,197 @@
package main

import (
"crypto/rand"
"encoding/binary"
"encoding/hex"
"io"
"net/url"
"sort"
"strconv"
"strings"
"sync/atomic"
"time"

"github.com/meshcore-analyzer/brokerurl"
)

// MQTT client IDs (#118).
//
// Without SetClientID paho connects with a zero-length client ID and
// CleanSession=true; whether the broker then assigns one, rejects the client
// or lets two ingestors take over each other's session is broker-dependent.
// Every source therefore gets an explicit ID:
//
// - mqttSources[].clientId, used verbatim. It must be unique among all
// clients connected to that broker at the same time.
// - otherwise corescope-<name>-<8 hex>, generated once per client
// construction (buildMQTTOpts). <name> is the sanitized source name, else
// the broker host, else omitted; the suffix is 32 random bits from
// crypto/rand. paho reuses the options for its own reconnects and the
// watchdog's force-reconnect reuses the same client, so the ID is stable
// for the process lifetime but differs between sources and processes.
//
// MQTT 3.1.1 (what paho tries first) only guarantees IDs of 1–23 characters
// from [0-9A-Za-z], although brokers such as Mosquitto and EMQX accept longer
// ones. paho falls back to MQTT 3.1 after a failed first handshake, and a
// strict 3.1 broker rejects IDs over 23 characters. The default is
// 19 characters plus the name part, so for such a broker configure a short
// clientId. (The previous empty ID was invalid under 3.1 as well.)

const mqttClientIDMaxBase = 32

// clientIDRandom is swapped in tests to simulate an entropy failure.
var clientIDRandom io.Reader = rand.Reader

var clientIDFallbackSeq atomic.Uint64

// clientIDNow feeds the fallback suffix; swapped in tests to model a coarse
// clock that returns the same reading twice.
var clientIDNow = time.Now

// mqttClientID returns the client ID for a new client of this source.
func mqttClientID(source MQTTSource) string {
if strings.TrimSpace(source.ClientID) != "" {
return source.ClientID
}
base := sanitizeClientIDPart(source.Name)
if base == "" {
// the masked broker: url.Parse alone may take a user name for
// the host (see brokerForLog)
if u, err := url.Parse(brokerForLog(source.Broker)); err == nil {
base = sanitizeClientIDPart(u.Hostname())
}
}
id := "corescope-"
if base != "" {
id += base + "-"
}
return id + clientIDSuffix()
}

// sanitizeClientIDPart lowercases s, keeps [a-z0-9], turns every other run
// into one '-', trims '-' and caps the result at mqttClientIDMaxBase.
func sanitizeClientIDPart(s string) string {
var b strings.Builder
dash := false
for _, r := range strings.ToLower(s) {
if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') {
if dash && b.Len() > 0 {
b.WriteByte('-')
}
dash = false
b.WriteRune(r)
continue
}
dash = true
}
out := b.String()
if len(out) > mqttClientIDMaxBase {
out = strings.TrimRight(out[:mqttClientIDMaxBase], "-")
}
return out
}

// clientIDSuffix is 8 hex chars from crypto/rand. Should that ever fail it
// mixes the clock and a process-wide counter instead, so two constructions
// still never share a suffix within a process.
func clientIDSuffix() string {
var b [4]byte
if _, err := io.ReadFull(clientIDRandom, b[:]); err != nil {
v := uint64(clientIDNow().UnixNano()) ^ (clientIDFallbackSeq.Add(1) * 0x9E3779B97F4A7C15)
binary.BigEndian.PutUint32(b[:], uint32(v>>32)^uint32(v))
}
return hex.EncodeToString(b[:])
}

// mqttConnectedLogLine is the "connected" log line. The broker URL is logged
// with its user-info masked and without query or fragment (brokerForLog), so
// credentials or device tokens embedded in it never reach the log.
func mqttConnectedLogLine(tag, broker, clientID string) string {
return "MQTT [" + tag + "] connected to " + brokerForLog(broker) + " as client " + clientID
}

// mqttSourceTag is the base of a source's tag in logs and the
// liveness/status registries: its name or, for an unnamed source, the
// broker with its credentials masked (brokerForLog). mqttSourceTags makes the tags
// of a whole configuration unique.
func mqttSourceTag(source MQTTSource) string {
if source.Name != "" {
return source.Name
}
return brokerForLog(source.Broker)
}

// mqttSourceTags returns the tag of every source. Unnamed sources on the
// same broker (say, with different credentials) would share one tag: the
// second would lose watchdog tracking and the two would share status
// counters. So an unnamed source whose tag is taken, by any named source or
// an earlier unnamed one, gets " (2)", " (3)", … A duplicate Name is left as
// it is: that is a configuration error, reported by registerLivenessOrSkip.
func mqttSourceTags(sources []MQTTSource) []string {
used := make(map[string]bool, len(sources))
for _, s := range sources {
if s.Name != "" {
used[s.Name] = true
}
}
tags := make([]string, len(sources))
for i, s := range sources {
tag := mqttSourceTag(s)
if s.Name == "" {
base := tag
for n := 2; used[tag]; n++ {
tag = base + " (" + strconv.Itoa(n) + ")"
}
used[tag] = true
}
tags[i] = tag
}
return tags
}

// brokerForLog returns broker with its user-info replaced by "****" and
// without query or fragment (brokerurl.Mask), so credentials or tokens
// embedded in it never reach a log, the stats file or a client ID, while the
// "****@" shows that the URL carries credentials and that the host may be
// cut short. A broker without a scheme is read as tcp://, as paho's
// AddBroker does.
func brokerForLog(broker string) string {
if !strings.Contains(broker, "://") {
broker = "tcp://" + broker
}
return brokerurl.Mask(broker)
}

// mqttSourceSecrets returns the non-empty secrets of a source: its
// password and user name, and the user-info, query and fragment of its
// broker URL as configured (brokerurl.Secrets).
func mqttSourceSecrets(source MQTTSource) []string {
var out []string
for _, v := range append([]string{source.Password, source.Username}, brokerurl.Secrets(source.Broker)...) {
if v != "" {
out = append(out, v)
}
}
return out
}

// errForLog is err's text with each of secrets (mqttSourceSecrets) replaced
// by "****", longest first, and then any broker URL or user-info it quotes
// masked (brokerurl.MaskText): MaskText only spots URL-shaped text, so a
// query token or a password quoted on its own would pass it. paho's errors
// normally quote none, but they reach the log and the stats file.
func errForLog(err error, secrets ...string) string {
if err == nil {
return "<nil>"
}
s := err.Error()
sorted := append([]string(nil), secrets...)
sort.Slice(sorted, func(i, j int) bool { return len(sorted[i]) > len(sorted[j]) })
for _, v := range sorted {
if v != "" {
s = strings.ReplaceAll(s, v, brokerurl.Marker)
}
}
return brokerurl.MaskText(s)
}
Loading
Loading