diff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml index a79c2112e..8788d76d9 100644 --- a/.github/workflows/deploy.yml +++ b/.github/workflows/deploy.yml @@ -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: diff --git a/Dockerfile b/Dockerfile index daaaf73de..30ad45353 100644 --- a/Dockerfile +++ b/Dockerfile @@ -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/ @@ -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/ diff --git a/cmd/ingestor/README.md b/cmd/ingestor/README.md index 5e7bce73a..0837e2c17 100644 --- a/cmd/ingestor/README.md +++ b/cmd/ingestor/README.md @@ -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 @@ -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--<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`) diff --git a/cmd/ingestor/config.go b/cmd/ingestor/config.go index 9b0c08581..fb0e47338 100644 --- a/cmd/ingestor/config.go +++ b/cmd/ingestor/config.go @@ -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, diff --git a/cmd/ingestor/go.mod b/cmd/ingestor/go.mod index 988ba3b61..16cdd556f 100644 --- a/cmd/ingestor/go.mod +++ b/cmd/ingestor/go.mod @@ -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 diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index 3ab0626b1..d0aabf65b 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -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) { @@ -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. @@ -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). @@ -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 @@ -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 }) diff --git a/cmd/ingestor/mqtt_client_id.go b/cmd/ingestor/mqtt_client_id.go new file mode 100644 index 000000000..887ed14b6 --- /dev/null +++ b/cmd/ingestor/mqtt_client_id.go @@ -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--<8 hex>, generated once per client +// construction (buildMQTTOpts). 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 "" + } + 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) +} diff --git a/cmd/ingestor/mqtt_client_id_118_test.go b/cmd/ingestor/mqtt_client_id_118_test.go new file mode 100644 index 000000000..77870915e --- /dev/null +++ b/cmd/ingestor/mqtt_client_id_118_test.go @@ -0,0 +1,326 @@ +package main + +import ( + "encoding/json" + "errors" + "net" + "os" + "regexp" + "strings" + "sync" + "testing" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" + "github.com/eclipse/paho.mqtt.golang/packets" +) + +// Issue #118: every source connects with an explicit client ID: the +// configured mqttSources[].clientId verbatim, else corescope-- +// generated once per client, so paho's auto-reconnect and the watchdog's +// force-reconnect reuse it, while unconfigured sources and processes differ. + +func TestMQTTClientIDConfiguredIsVerbatim_118(t *testing.T) { + for _, id := range []string{"my-ingestor-01", "Edge Box/1", "a"} { + opts := buildMQTTOpts(MQTTSource{Name: "x", Broker: "tcp://h:1883", ClientID: id}) + if opts.ClientID != id { + t.Errorf("configured clientId %q became %q", id, opts.ClientID) + } + } +} + +func TestMQTTClientIDDefaultShape_118(t *testing.T) { + cases := []struct { + name, broker string + want string + }{ + {"local", "mqtt://localhost:1883", `^corescope-local-[0-9a-f]{8}$`}, + {"Local Feed #1", "tcp://h:1883", `^corescope-local-feed-1-[0-9a-f]{8}$`}, + {" ÆØÅ lincomatic__SJC ", "tcp://h:1883", `^corescope-lincomatic-sjc-[0-9a-f]{8}$`}, + // name empty or nothing left after sanitizing: broker host, no port + {"", "mqtts://user:secret@MQTT.Example.com:8883", `^corescope-mqtt-example-com-[0-9a-f]{8}$`}, + {"###", "wss://ws.example.net/mqtt", `^corescope-ws-example-net-[0-9a-f]{8}$`}, + // neither: plain corescope + {"", "", `^corescope-[0-9a-f]{8}$`}, + {"%%%", "tcp://:1883", `^corescope-[0-9a-f]{8}$`}, + {"", "not a url", `^corescope-[0-9a-f]{8}$`}, + } + for _, c := range cases { + id := buildMQTTOpts(MQTTSource{Name: c.name, Broker: c.broker}).ClientID + if !regexp.MustCompile(c.want).MatchString(id) { + t.Errorf("name %q broker %q: client id %q, want %s", c.name, c.broker, id, c.want) + } + if strings.Contains(id, "secret") || strings.Contains(id, "user") { + t.Errorf("client id %q leaks broker credentials", id) + } + } +} + +// A whitespace-only clientId counts as unset (a blank template value must +// not become a shared, colliding ID). +func TestMQTTClientIDBlankConfiguredIsGenerated_118(t *testing.T) { + id := buildMQTTOpts(MQTTSource{Name: "local", Broker: "tcp://h:1883", ClientID: " "}).ClientID + if !regexp.MustCompile(`^corescope-local-[0-9a-f]{8}$`).MatchString(id) { + t.Fatalf("blank clientId gave %q", id) + } +} + +func TestMQTTClientIDLongNameIsCapped_118(t *testing.T) { + id := buildMQTTOpts(MQTTSource{Name: strings.Repeat("abcdefghij", 20), Broker: "tcp://h:1883"}).ClientID + if len(id) > len("corescope-")+mqttClientIDMaxBase+1+8 { + t.Fatalf("client id is %d chars: %q", len(id), id) + } + base := strings.TrimSuffix(strings.TrimPrefix(id, "corescope-"), id[len(id)-9:]) + if base == "" || strings.HasSuffix(base, "-") || strings.Contains(id, "--") { + t.Fatalf("capped id %q has an empty base or a dangling separator", id) + } +} + +// Same source, same process or not: every construction is a new client and +// gets its own ID. The suffix is 32 random bits, so n constructions collide +// with probability about n²/2³³: 200 keeps that near 5e-6 (5000, used +// before, failed about one run in 350). +func TestMQTTClientIDUniqueAcrossConstructions_118(t *testing.T) { + seen := map[string]bool{} + src := MQTTSource{Name: "local", Broker: "mqtt://localhost:1883"} + for i := 0; i < 200; i++ { + id := buildMQTTOpts(src).ClientID + if seen[id] { + t.Fatalf("duplicate client id %q after %d constructions", id, i) + } + seen[id] = true + } +} + +type failingReader struct{} + +func (failingReader) Read([]byte) (int, error) { return 0, errors.New("entropy unavailable") } + +// Should crypto/rand ever fail, the ID is still non-empty and still differs +// between constructions instead of collapsing to a shared value. +func TestMQTTClientIDRandomFailureStillUnique_118(t *testing.T) { + old, oldNow := clientIDRandom, clientIDNow + clientIDRandom = failingReader{} + frozen := time.Unix(1790000000, 0) + clientIDNow = func() time.Time { return frozen } // a coarse clock: one reading + defer func() { clientIDRandom, clientIDNow = old, oldNow }() + a := buildMQTTOpts(MQTTSource{Name: "local", Broker: "tcp://h:1883"}).ClientID + if !regexp.MustCompile(`^corescope-local-[0-9a-f]{8}$`).MatchString(a) { + t.Fatalf("with crypto/rand failing: %q", a) + } + // back-to-back constructions (same clock reading on coarse clocks) differ + seen := map[string]bool{a: true} + for i := 0; i < 2000; i++ { + id := mqttClientID(MQTTSource{Name: "local"}) + if seen[id] { + t.Fatalf("with crypto/rand failing, construction %d repeated %q", i, id) + } + seen[id] = true + } +} + +// The rest of the options are unchanged by the ID. +func TestMQTTClientIDKeepsOtherOptions_118(t *testing.T) { + f := false + opts := buildMQTTOpts(MQTTSource{Name: "n", Broker: "mqtts://h:8883", Username: "u", Password: "p", RejectUnauthorized: &f}) + if !opts.CleanSession || opts.Username != "u" || opts.Password != "p" || opts.TLSConfig == nil || !opts.TLSConfig.InsecureSkipVerify { + t.Fatalf("options changed: clean=%v user=%q tls=%v", opts.CleanSession, opts.Username, opts.TLSConfig) + } +} + +func TestMQTTConnectedLogLine_118(t *testing.T) { + line := mqttConnectedLogLine("feed", "mqtts://dev-user:tok3n-secret@broker.example:8883", "corescope-feed-0a1b2c3d") + for _, bad := range []string{"tok3n-secret", "dev-user"} { + if strings.Contains(line, bad) { + t.Fatalf("log line leaks %q: %s", bad, line) + } + } + for _, want := range []string{"MQTT [feed] connected to mqtts://****@broker.example:8883 as client corescope-feed-0a1b2c3d"} { + if !strings.Contains(line, want) { + t.Fatalf("log line %q misses %q", line, want) + } + } +} + +// config.example.json is copied as a live config by some deployments, so it +// must document clientId without setting one. +func TestConfigExampleHasNoLiteralClientID_118(t *testing.T) { + raw, err := os.ReadFile("../../config.example.json") + if err != nil { + t.Fatal(err) + } + var cfg struct { + Sources []map[string]json.RawMessage `json:"mqttSources"` + Doc string `json:"_comment_mqttSources"` + } + if err := json.Unmarshal(raw, &cfg); err != nil { + t.Fatal(err) + } + for i, s := range cfg.Sources { + if _, ok := s["clientId"]; ok { + t.Errorf("mqttSources[%d] sets a literal clientId", i) + } + } + if !strings.Contains(cfg.Doc, "clientId") { + t.Error("_comment_mqttSources does not document clientId") + } +} + +// The legacy single-broker block becomes source "default" and gets a +// generated ID like any other unconfigured source. +func TestLegacyMQTTConfigGetsGeneratedClientID_118(t *testing.T) { + dir := t.TempDir() + p := dir + "/config.json" + if err := os.WriteFile(p, []byte(`{"mqtt":{"broker":"mqtt://old-mosquitto:1883","topic":"meshcore/#"}}`), 0o600); err != nil { + t.Fatal(err) + } + cfg, err := LoadConfig(p) + if err != nil { + t.Fatal(err) + } + id := buildMQTTOpts(cfg.MQTTSources[0]).ClientID + if !regexp.MustCompile(`^corescope-default-[0-9a-f]{8}$`).MatchString(id) { + t.Fatalf("legacy source client id %q", id) + } +} + +// ── loopback broker: the ID actually sent, across reconnects ── + +type idBroker struct { + ln net.Listener + mu sync.Mutex + ids []string + conns []net.Conn + wg sync.WaitGroup +} + +func newIDBroker(t *testing.T) *idBroker { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + b := &idBroker{ln: ln} + b.wg.Add(1) + go func() { + defer b.wg.Done() + for { + c, err := ln.Accept() + if err != nil { + return + } + b.wg.Add(1) + go b.serve(c) + } + }() + t.Cleanup(func() { + ln.Close() + b.dropAll() + b.wg.Wait() + }) + return b +} + +func (b *idBroker) serve(c net.Conn) { + defer b.wg.Done() + defer c.Close() + pkt, err := packets.ReadPacket(c) + if err != nil { + return + } + cp, ok := pkt.(*packets.ConnectPacket) + if !ok { + return + } + b.mu.Lock() + b.ids = append(b.ids, cp.ClientIdentifier) + b.conns = append(b.conns, c) + b.mu.Unlock() + ack := packets.NewControlPacket(packets.Connack).(*packets.ConnackPacket) + ack.ReturnCode = packets.Accepted + if ack.Write(c) != nil { + return + } + for { + p, err := packets.ReadPacket(c) + if err != nil { + return + } + switch p := p.(type) { + case *packets.PingreqPacket: + packets.NewControlPacket(packets.Pingresp).Write(c) + case *packets.SubscribePacket: + ack := packets.NewControlPacket(packets.Suback).(*packets.SubackPacket) + ack.MessageID = p.MessageID + ack.ReturnCodes = make([]byte, len(p.Topics)) + ack.Write(c) + case *packets.DisconnectPacket: + return + } + } +} + +func (b *idBroker) dropAll() { + b.mu.Lock() + defer b.mu.Unlock() + for _, c := range b.conns { + c.Close() + } +} + +func (b *idBroker) seen() []string { + b.mu.Lock() + defer b.mu.Unlock() + return append([]string(nil), b.ids...) +} + +func waitForID(t *testing.T, what string, cond func() bool) { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + for !cond() { + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for %s", what) + } + time.Sleep(5 * time.Millisecond) + } +} + +// First connect, paho auto-reconnect after the broker drops the socket, and +// the watchdog's force-reconnect all present the same non-empty ID; a second +// client for the same source presents a different one. +func TestMQTTClientIDStableAcrossReconnects_118(t *testing.T) { + b := newIDBroker(t) + src := MQTTSource{Name: "loop", Broker: "tcp://" + b.ln.Addr().String()} + opts := buildMQTTOpts(src). + SetConnectTimeout(time.Second). + SetMaxReconnectInterval(100 * time.Millisecond). + SetConnectRetryInterval(50 * time.Millisecond) + client := mqtt.NewClient(opts) + client.Connect() + defer client.Disconnect(0) + waitForID(t, "first connect", func() bool { return len(b.seen()) == 1 && client.IsConnectionOpen() }) + + b.dropAll() // paho auto-reconnect + waitForID(t, "auto-reconnect", func() bool { return len(b.seen()) == 2 && client.IsConnectionOpen() }) + + buildForceReconnectFn(client, "loop")() // watchdog force-reconnect + waitForID(t, "force-reconnect", func() bool { return len(b.seen()) >= 3 && client.IsConnectionOpen() }) + + ids := b.seen() + if ids[0] == "" { + t.Fatal("the client connected with an empty client id") + } + for i, id := range ids { + if id != ids[0] { + t.Fatalf("connect %d used client id %q, first used %q", i+1, id, ids[0]) + } + } + + other := mqtt.NewClient(buildMQTTOpts(src).SetConnectTimeout(time.Second)) + other.Connect() + defer other.Disconnect(0) + waitForID(t, "second client", func() bool { return len(b.seen()) > len(ids) }) + if last := b.seen()[len(b.seen())-1]; last == ids[0] || last == "" { + t.Fatalf("a second client for the same source reused id %q", last) + } +} diff --git a/cmd/ingestor/mqtt_credentials_r3_118_test.go b/cmd/ingestor/mqtt_credentials_r3_118_test.go new file mode 100644 index 000000000..b47ffa9f7 --- /dev/null +++ b/cmd/ingestor/mqtt_credentials_r3_118_test.go @@ -0,0 +1,168 @@ +package main + +import ( + "bytes" + "errors" + "os" + "path/filepath" + "runtime" + "strings" + "testing" + "time" +) + +// #118 review round 3. + +// What the ingestor logs, registers and publishes keeps "****@" where it +// cut user-info, so an operator can see that the URL carries credentials, +// and a host that is only what follows an '@' reads as cut short. +func TestBrokerForLogMarksRemovedUserinfo_118(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {"mqtt://u:p@broker:1883", "mqtt://****@broker:1883"}, + {"wss://host/mqtt?u=me@x.org", "wss://****@x.org"}, + {"wss://host/mqtt?token=abc", "wss://host/mqtt"}, + {"broker:1883", "tcp://broker:1883"}, + } { + if got := brokerForLog(c.in); got != c.want { + t.Errorf("brokerForLog(%q) = %q, want %q", c.in, got, c.want) + } + } + resetSourceStatusRegistry() + t.Cleanup(resetSourceStatusRegistry) + if got := RegisterSourceStatus("t", "mqtt://u:p@broker:1883").snapshot(time.Now()).Broker; got != "mqtt://****@broker:1883" { + t.Errorf("status broker = %q, want mqtt://****@broker:1883", got) + } + if got := mqttSourceTag(MQTTSource{Broker: "wss://host/mqtt?u=me@x.org"}); got != "wss://****@x.org" { + t.Errorf("tag = %q, want wss://****@x.org", got) + } +} + +// The secrets the ingestor knows are replaced literally before MaskText, +// which only spots URL-shaped text: a query token quoted without its +// scheme, or the configured user name and password, would pass it. +func TestErrForLogMasksKnownSecrets_118(t *testing.T) { + src := MQTTSource{ + Broker: "wss://" + credUser + ":" + credPass + "@host/mqtt?token=abc", + Username: "cfg-user", + Password: "cfg-pass", + } + secrets := mqttSourceSecrets(src) + for _, c := range []struct{ in, want string }{ + {"dial host/mqtt?token=abc: refused", "dial host/mqtt?****: refused"}, + {"not authorized: cfg-user/cfg-pass", "not authorized: ****/****"}, + {"auth " + credUser + ":" + credPass + " rejected", "auth **** rejected"}, + {"connect tcp://x:y@host failed", "connect tcp://****@host failed"}, // MaskText still runs + {"EOF", "EOF"}, + } { + got := errForLog(errors.New(c.in), secrets...) + assertNoSecretParts118(t, "errForLog("+c.in+")", got) + if got != c.want { + t.Errorf("errForLog(%q) = %q, want %q", c.in, got, c.want) + } + } + // longest first: a user name that is also the start of the URL's + // user-info must not leave the password behind + both := mqttSourceSecrets(MQTTSource{Broker: "tcp://" + credUser + ":" + credPass + "@host", Username: credUser}) + if got := errForLog(errors.New("auth "+credUser+":"+credPass), both...); got != "auth ****" { + t.Errorf("user name before user-info: %q, want %q", got, "auth ****") + } + // only non-empty values: an empty one would mask between every rune + if got := errForLog(errors.New("EOF"), mqttSourceSecrets(MQTTSource{Broker: "tcp://host?"})...); got != "EOF" { + t.Errorf("no secrets: %q", got) + } +} + +// The disconnect handler main() installs logs and stores the error with +// the source's secrets masked. +func TestDisconnectErrorMasksKnownSecrets_118(t *testing.T) { + resetSourceStatusRegistry() + t.Cleanup(resetSourceStatusRegistry) + buf := captureLog118(t) + src := MQTTSource{Name: "feed", Broker: "wss://host/mqtt?token=abc", Password: "cfg-pass"} + opts, status, _ := prepareMQTTSource(src, "feed") + opts.OnConnectionLost(nil, errors.New("dial host/mqtt?token=abc: cfg-pass rejected")) + const want = "dial host/mqtt?****: **** rejected" + if got := status.snapshot(time.Now()).LastError; got != want { + t.Errorf("lastError = %q, want %q", got, want) + } + if out := buf.String(); !strings.Contains(out, want) || strings.Contains(out, "abc") || strings.Contains(out, "cfg-pass") { + t.Errorf("log:\n%s", out) + } +} + +// All four random bytes reach the suffix, in order. +func TestMQTTClientIDSuffixUsesAllRandomBytes_118(t *testing.T) { + old := clientIDRandom + t.Cleanup(func() { clientIDRandom = old }) + clientIDRandom = bytes.NewReader([]byte{0xde, 0xad, 0xbe, 0xef}) + if got := mqttClientID(MQTTSource{Name: "Feed"}); got != "corescope-feed-deadbeef" { + t.Fatalf("client id %q, want corescope-feed-deadbeef", got) + } +} + +// A stale tmp file owned by another user is refused before anything is +// truncated, chmod-ed or written. Root's chmod succeeds on any file, and +// the owner may still hold it open, so failing chmod is no guard there. +func TestWriteStatsAtomicRefusesForeignTmp_118(t *testing.T) { + if runtime.GOOS == "windows" { + t.Skip("no Unix file owners") + } + old := statsFileEUID + t.Cleanup(func() { statsFileEUID = old }) + statsFileEUID = func() int { return os.Geteuid() + 1 } // the tmp is someone else's + assertForeignTmpRefused118(t, func(string) {}) +} + +// A stale tmp file of our own loses its old content: without O_TRUNC the +// file is truncated after the owner check, and a longer stale file must +// not leave a tail behind the new JSON. +func TestWriteStatsAtomicTruncatesOwnStaleTmp_118(t *testing.T) { + path := filepath.Join(t.TempDir(), "stats.json") + if err := os.WriteFile(path+".tmp", []byte(strings.Repeat("stale ", 100)), 0o644); err != nil { + t.Fatal(err) + } + if err := writeStatsAtomic(path, []byte(`{}`)); err != nil { + t.Fatal(err) + } + if b, err := os.ReadFile(path); err != nil || string(b) != `{}` { + t.Fatalf("stats file %q, %v; want {}", b, err) + } +} + +// The same with a real foreign owner, which needs root to set up. +func TestWriteStatsAtomicRefusesForeignTmpAsRoot_118(t *testing.T) { + if runtime.GOOS == "windows" || os.Geteuid() != 0 { + t.Skip("needs root to plant a file owned by another user") + } + assertForeignTmpRefused118(t, func(tmp string) { + if err := os.Chown(tmp, 65534, 65534); err != nil { + t.Fatal(err) + } + }) +} + +func assertForeignTmpRefused118(t *testing.T, plant func(tmp string)) { + t.Helper() + path := filepath.Join(t.TempDir(), "stats.json") + tmp := path + ".tmp" + if err := os.WriteFile(tmp, []byte("planted"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chmod(tmp, 0o644); err != nil { + t.Fatal(err) + } + plant(tmp) + if err := writeStatsAtomic(path, []byte(`{"source_statuses":[]}`)); err == nil { + t.Error("writeStatsAtomic accepted a tmp file owned by another user") + } + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Errorf("stats file published: %v", err) + } + fi, err := os.Stat(tmp) + if err != nil { + t.Fatal(err) + } + if b, _ := os.ReadFile(tmp); string(b) != "planted" || fi.Mode().Perm() != 0o644 { + t.Errorf("foreign tmp changed: %q mode %o", b, fi.Mode().Perm()) + } +} diff --git a/cmd/ingestor/mqtt_log_credentials_118_test.go b/cmd/ingestor/mqtt_log_credentials_118_test.go new file mode 100644 index 000000000..15c784a41 --- /dev/null +++ b/cmd/ingestor/mqtt_log_credentials_118_test.go @@ -0,0 +1,186 @@ +package main + +import ( + "bytes" + "go/ast" + "go/parser" + "go/token" + "log" + "net/url" + "strings" + "sync" + "testing" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +// #118 follow-up: the log tag of a source without a name used to be the raw +// broker URL, so user-info in it reached every MQTT log line. The tag and +// every logged broker must be free of credentials, also for a broker given +// without a scheme (paho then dials tcp://). + +const ( + credUser = "dev-user" + credPass = "s3cret-pass" +) + +func assertNoCredentials(t *testing.T, what, s string) { + t.Helper() + for _, bad := range []string{credUser, credPass} { + if strings.Contains(s, bad) { + t.Errorf("%s leaks %q: %s", what, bad, s) + } + } +} + +func TestMQTTSourceTagHasNoCredentials_118(t *testing.T) { + for _, tc := range []struct{ name, broker, want string }{ + {"", "tcp://" + credUser + ":" + credPass + "@broker.example:1883", "tcp://****@broker.example:1883"}, + {"", credUser + ":" + credPass + "@broker.example:1883", "tcp://****@broker.example:1883"}, + {"", "wss://" + credUser + ":" + credPass + "@broker.example/mqtt?password=" + credPass, "wss://****@broker.example/mqtt"}, + {"", "tcp://broker.example:1883", "tcp://broker.example:1883"}, + {"", "broker.example:1883", "tcp://broker.example:1883"}, + {"feed", "tcp://" + credUser + ":" + credPass + "@broker.example:1883", "feed"}, + } { + got := mqttSourceTag(MQTTSource{Name: tc.name, Broker: tc.broker}) + assertNoCredentials(t, "tag for "+tc.broker, got) + if got != tc.want { + t.Errorf("tag for name %q broker %q = %q, want %q", tc.name, tc.broker, got, tc.want) + } + } + // unparseable input: everything up to the last '@' is dropped + got := mqttSourceTag(MQTTSource{Broker: "tcp://" + credUser + ":" + credPass + "@bro ker%zz:1883"}) + assertNoCredentials(t, "tag for an unparseable broker", got) +} + +// logBuf118 is a log sink safe to read while paho's goroutines log. +type logBuf118 struct { + mu sync.Mutex + buf bytes.Buffer +} + +func (b *logBuf118) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.Write(p) +} + +func (b *logBuf118) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.String() +} + +// captureLog118 collects the standard logger's output for the test. +func captureLog118(t *testing.T) *logBuf118 { + t.Helper() + buf := &logBuf118{} + prevOut, prevFlags := log.Writer(), log.Flags() + log.SetOutput(buf) + t.Cleanup(func() { log.SetOutput(prevOut); log.SetFlags(prevFlags) }) + return buf +} + +// A real connect through buildMQTTOpts to a loopback broker, for an unnamed +// source whose broker URL carries credentials, with and without a scheme. +// Every line logged (connection attempt, connected, disconnected, +// reconnecting, watchdog) must be free of them. +func TestMQTTLogLinesHaveNoCredentials_118(t *testing.T) { + b := newIDBroker(t) + addr := b.ln.Addr().String() + for _, broker := range []string{ + "tcp://" + credUser + ":" + credPass + "@" + addr, + credUser + ":" + credPass + "@" + addr, + } { + t.Run(broker[:3], func(t *testing.T) { + buf := captureLog118(t) + src := MQTTSource{Broker: broker} + tag, logBroker := mqttSourceTag(src), brokerForLog(src.Broker) + opts := buildMQTTOpts(src).SetConnectTimeout(time.Second) + client := mqtt.NewClient(opts) + seen := len(b.seen()) + client.Connect() + defer client.Disconnect(0) + waitForID(t, "connect", func() bool { return len(b.seen()) > seen && client.IsConnectionOpen() }) + log.Print(mqttConnectedLogLine(tag, src.Broker, opts.ClientID)) + s := &SourceLivenessState{Tag: tag, Broker: logBroker, FirstConnectedAt: time.Now().Add(-time.Hour).Unix()} + if msg, _ := checkSourceLiveness(s, time.Minute, time.Now()); msg != "" { + log.Print(msg) + } + out := buf.String() + if !strings.Contains(out, "connection attempt #") || !strings.Contains(out, "connected to") { + t.Fatalf("expected the attempt and connected lines, got:\n%s", out) + } + for _, line := range strings.Split(strings.TrimSpace(out), "\n") { + assertNoCredentials(t, "log line", line) + } + }) + } +} + +// main() and prepareMQTTSource wire the tag and the logged broker. A quick +// static check of both: no log call or RegisterSourceStatus may take a raw +// X.Broker argument, the tag must not be a raw broker, and the liveness +// state (whose Broker the watchdog logs) must get the log-safe broker. It only sees direct uses; +// TestMQTTSourceWiringLeaksNoCredentials_118 runs the handlers. +func TestMainLogsNoRawBroker_118(t *testing.T) { + for _, file := range []string{"main.go", "mqtt_source.go"} { + checkNoRawBrokerLogged118(t, file) + } +} + +func checkNoRawBrokerLogged118(t *testing.T, file string) { + t.Helper() + fset := token.NewFileSet() + f, err := parser.ParseFile(fset, file, nil, 0) + if err != nil { + t.Fatal(err) + } + isBroker := func(e ast.Expr) bool { + sel, ok := e.(*ast.SelectorExpr) + return ok && sel.Sel.Name == "Broker" + } + ast.Inspect(f, func(n ast.Node) bool { + switch n := n.(type) { + case *ast.CallExpr: + if id, ok := n.Fun.(*ast.Ident); ok && id.Name == "RegisterSourceStatus" { + for _, a := range n.Args { + if isBroker(a) { + t.Errorf("%s: status registry given a raw broker URL; the stats file publishes it", fset.Position(a.Pos())) + } + } + } + if sel, ok := n.Fun.(*ast.SelectorExpr); ok { + if pkg, ok := sel.X.(*ast.Ident); ok && pkg.Name == "log" { + for _, a := range n.Args { + if isBroker(a) { + t.Errorf("%s: log call with a raw broker URL", fset.Position(a.Pos())) + } + } + } + } + case *ast.AssignStmt: + for i, l := range n.Lhs { + if id, ok := l.(*ast.Ident); ok && id.Name == "tag" && i < len(n.Rhs) && isBroker(n.Rhs[i]) { + t.Errorf("%s: tag assigned a raw broker URL", fset.Position(n.Pos())) + } + } + case *ast.CompositeLit: + if id, ok := n.Type.(*ast.Ident); ok && id.Name == "SourceLivenessState" { + for _, el := range n.Elts { + if kv, ok := el.(*ast.KeyValueExpr); ok && isBroker(kv.Value) { + t.Errorf("%s: SourceLivenessState.Broker set to a raw broker URL; the watchdog logs it", fset.Position(kv.Pos())) + } + } + } + } + return true + }) +} + +func TestBrokerForLogHasNoCredentials_118(t *testing.T) { + u, _ := url.Parse("tcp://" + credUser + ":" + credPass + "@broker.example:1883") + assertNoCredentials(t, "brokerForLog of paho's URL", brokerForLog(u.String())) + assertNoCredentials(t, "brokerForLog without a scheme", brokerForLog(credUser+":"+credPass+"@broker.example:1883")) +} diff --git a/cmd/ingestor/mqtt_source.go b/cmd/ingestor/mqtt_source.go new file mode 100644 index 000000000..d0ddbfd27 --- /dev/null +++ b/cmd/ingestor/mqtt_source.go @@ -0,0 +1,68 @@ +package main + +import ( + "log" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +// prepareMQTTSource is main()'s per-source setup up to the client itself: +// the paho options with the connect, connection-lost and reconnecting +// handlers, the source's status-registry entry and its liveness state +// (IsConnectedFn and ForceReconnectFn are wired by the caller once the +// client exists; registration is left to the caller too). The message +// handler needs the store and ingest buffer, so main() sets it. +func prepareMQTTSource(source MQTTSource, tag string) (*mqtt.ClientOptions, *sourceStatusState, *SourceLivenessState) { + logBroker := brokerForLog(source.Broker) + secrets := mqttSourceSecrets(source) + opts := buildMQTTOpts(source) + clientID := opts.ClientID + + // Pre-allocate the liveness pointer so OnConnect can reset its + // stale-message clock on reconnect (PR #1216 r1 item 2). + liveness := &SourceLivenessState{ + Tag: tag, + Broker: logBroker, // the watchdog logs it + } + + // #1043: per-source status registry. Idempotent — repeated + // registration across reconnects returns the same state so + // counters accumulate across the process lifetime. It is published + // in the stats file, so it gets the masked broker (#118). + status := RegisterSourceStatus(tag, logBroker) + + opts.SetOnConnectHandler(func(c mqtt.Client) { + log.Print(mqttConnectedLogLine(tag, source.Broker, clientID)) + 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: %s", tag, logBroker, errForLog(err, secrets...)) + status.MarkDisconnect(time.Now(), err, secrets...) + }) + + opts.SetReconnectingHandler(func(c mqtt.Client, options *mqtt.ClientOptions) { + log.Printf("MQTT [%s] reconnecting to %s", tag, logBroker) + }) + + return opts, status, liveness +} diff --git a/cmd/ingestor/mqtt_status_credentials_118_test.go b/cmd/ingestor/mqtt_status_credentials_118_test.go new file mode 100644 index 000000000..27ed4e253 --- /dev/null +++ b/cmd/ingestor/mqtt_status_credentials_118_test.go @@ -0,0 +1,233 @@ +package main + +import ( + "encoding/json" + "errors" + "os" + "path/filepath" + "regexp" + "strings" + "testing" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +// #118 re-review: credentials in a broker URL must not reach the log, the +// status registry, the stats file (which feeds the public /api/mqtt/status +// and /api/healthz) or the client ID, whatever form the URL takes. + +// secretParts118 are the credential parts used in the brokers below. +var secretParts118 = []string{credUser, credPass, "tok3n", "2024", "1234", "abc", "p%zz"} + +func assertNoSecretParts118(t *testing.T, what, s string) { + t.Helper() + for _, bad := range secretParts118 { + if strings.Contains(s, bad) { + t.Errorf("%s leaks %q: %s", what, bad, s) + } + } +} + +// An unescaped '/', '?' or '#' in the password ends the authority for +// url.Parse, so the password (or its first part) used to be logged. +func TestBrokerForLogAmbiguousPassword_118(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {"tcp://" + credUser + ":2024/" + credPass + "@host:1883", "tcp://****@host:1883"}, + {"tcp://" + credUser + ":1234?abc@host", "tcp://****@host"}, + {"tcp://" + credUser + ":1234#abc@host", "tcp://****@host"}, + {credUser + ":" + credPass + "@host:1883", "tcp://****@host:1883"}, + {"tcp://tok3n@host", "tcp://****@host"}, + {"wss://host/mqtt?token=abc", "wss://host/mqtt"}, + {"tcp://" + credUser + ":p%zz@host", "tcp://****@host"}, + } { + got := brokerForLog(c.in) + assertNoSecretParts118(t, "brokerForLog("+c.in+")", got) + if got != c.want { + t.Errorf("brokerForLog(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +// The generated client ID is logged ("as client …"), so its host part must +// come from the stripped broker, not from url.Parse's idea of the host. +func TestMQTTClientIDBaseHasNoCredentials_118(t *testing.T) { + id := mqttClientID(MQTTSource{Broker: "tcp://" + credUser + ":2024/" + credPass + "@host:1883"}) + if !regexp.MustCompile(`^corescope-host-[0-9a-f]{8}$`).MatchString(id) { + t.Fatalf("client id %q", id) + } +} + +// Two unnamed sources on the same host with different credentials used to +// share one tag: the second lost watchdog tracking and the status counters +// were merged. Tags are unique and credential-free; named sources keep +// their name (a duplicate name stays the operator's config error). +func TestMQTTSourceTagsAreUnique_118(t *testing.T) { + sources := []MQTTSource{ + {Broker: "tcp://" + credUser + ":" + credPass + "@host:1883"}, + {Broker: "tcp://tok3n@host:1883"}, + {Name: "feed", Broker: "tcp://host:1883"}, + {Broker: "host:1883"}, + // a later source named like a suffixed tag keeps its name + {Name: "tcp://host:1883 (3)", Broker: "tcp://other:1883"}, + {Broker: "tcp://other:1883"}, + {Name: "feed", Broker: "tcp://x:1883"}, + } + got := mqttSourceTags(sources) + want := []string{"tcp://****@host:1883", "tcp://****@host:1883 (2)", "feed", "tcp://host:1883", "tcp://host:1883 (3)", "tcp://other:1883", "feed"} + if strings.Join(got, "|") != strings.Join(want, "|") { + t.Fatalf("tags\n got %q\nwant %q", got, want) + } + for _, tag := range got { + assertNoSecretParts118(t, "tag", tag) + } + // the two unnamed sources on host:1883 now both get tracked and counted + saved := livenessRegistry + livenessRegistry = map[string]*SourceLivenessState{} + resetSourceStatusRegistry() + t.Cleanup(func() { livenessRegistry = saved; resetSourceStatusRegistry() }) + for i := 0; i < 2; i++ { + if !registerLivenessOrSkip(&SourceLivenessState{Tag: got[i]}) { + t.Fatalf("source %d not tracked by the watchdog", i) + } + if RegisterSourceStatus(got[i], sources[i].Broker) == lookupSourceStatus(got[1-i]) { + t.Fatalf("source %d shares its status counters", i) + } + } +} + +// The status registry is written to the stats file and served by the +// public /api/mqtt/status, so it holds the stripped broker, whoever calls +// RegisterSourceStatus, and a disconnect error that quotes a broker URL +// is stripped too. +func TestSourceStatusHoldsNoCredentials_118(t *testing.T) { + resetSourceStatusRegistry() + t.Cleanup(resetSourceStatusRegistry) + s := RegisterSourceStatus("t", "tcp://"+credUser+":1234?abc@host") + s.MarkDisconnect(time.Now(), errors.New(`dial "wss://`+credUser+`:`+credPass+`@host/mqtt?token=abc": refused`)) + snap := s.snapshot(time.Now()) + b, _ := json.Marshal(snap) + assertNoSecretParts118(t, "status snapshot", string(b)) + if snap.Broker != "tcp://****@host" { + t.Errorf("Broker = %q, want tcp://****@host", snap.Broker) + } +} + +// ── runtime wiring: main()'s per-source setup, against a loopback broker ── + +// prepareMQTTSource is what main() runs for every source. Driven here +// through connect, a broker-side drop, paho's reconnect and the stats +// file: every log line, the status registry and the file must be free +// of the credentials in the broker URL, whichever path they take. +func TestMQTTSourceWiringLeaksNoCredentials_118(t *testing.T) { + b := newIDBroker(t) + addr := b.ln.Addr().String() + for _, broker := range []string{ + "tcp://" + credUser + ":" + credPass + "@" + addr, + credUser + ":" + credPass + "@" + addr, + } { + t.Run("", func(t *testing.T) { + resetSourceStatusRegistry() + saved := livenessRegistry + livenessRegistry = map[string]*SourceLivenessState{} + t.Cleanup(func() { livenessRegistry = saved; resetSourceStatusRegistry() }) + + buf := captureLog118(t) + src := MQTTSource{Broker: broker, Topics: []string{"meshcore/#"}} + tag := mqttSourceTags([]MQTTSource{src})[0] + opts, _, liveness := prepareMQTTSource(src, tag) + opts.SetConnectTimeout(time.Second). + SetMaxReconnectInterval(100 * time.Millisecond). + SetConnectRetryInterval(50 * time.Millisecond) + client := mqtt.NewClient(opts) + liveness.IsConnectedFn = client.IsConnected + if !registerLivenessOrSkip(liveness) { + t.Fatal("liveness not registered") + } + seen := len(b.seen()) + client.Connect() + defer client.Disconnect(0) + waitForID(t, "connect", func() bool { return len(b.seen()) > seen && strings.Contains(buf.String(), "subscribed to") }) + + b.dropAll() // ConnectionLost, then paho's Reconnecting and OnConnect + waitForID(t, "reconnect", func() bool { + return len(b.seen()) > seen+1 && strings.Count(buf.String(), "subscribed to") >= 2 + }) + + out := buf.String() + for _, want := range []string{"connected to tcp://****@" + addr, "disconnected from tcp://****@" + addr, "reconnecting to tcp://****@" + addr, "subscribed to"} { + if !strings.Contains(out, want) { + t.Fatalf("no %q line in:\n%s", want, out) + } + } + for _, line := range strings.Split(strings.TrimSpace(out), "\n") { + assertNoSecretParts118(t, "log line", line) + } + + // what the stats file carries: statuses and liveness keys + stats := struct { + S []SourceStatusSnapshot `json:"source_statuses"` + L map[string]SourceLivenessSnapshot `json:"source_liveness"` + }{SnapshotSourceStatuses(time.Now()), SnapshotLivenessClocks()} + j, _ := json.Marshal(stats) + assertNoSecretParts118(t, "stats data", string(j)) + if len(stats.S) != 1 || stats.S[0].DisconnectCount < 1 || stats.S[0].ConnectCount < 2 { + t.Fatalf("status not wired: %s", j) + } + }) + } +} + +// The stats file itself: stripped content, owner-only permissions even +// when a stale tmp file with wider permissions is lying around. +func TestStatsFileHasNoCredentials_118(t *testing.T) { + dir := t.TempDir() + statsPath := filepath.Join(dir, "ingestor-stats.json") + t.Setenv("CORESCOPE_INGESTOR_STATS", statsPath) + if err := os.WriteFile(statsPath+".tmp", []byte("stale"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chmod(statsPath+".tmp", 0o644); err != nil { + t.Fatal(err) + } + + resetSourceStatusRegistry() + saved := livenessRegistry + livenessRegistry = map[string]*SourceLivenessState{} + t.Cleanup(func() { livenessRegistry = saved; resetSourceStatusRegistry() }) + sources := []MQTTSource{ + {Broker: "tcp://" + credUser + ":" + credPass + "@host:1883"}, + {Broker: "tcp://tok3n@host:1883"}, + } + for i, tag := range mqttSourceTags(sources) { + _, status, liveness := prepareMQTTSource(sources[i], tag) + status.MarkDisconnect(time.Now(), errors.New("connect "+sources[i].Broker+" refused")) + registerLivenessOrSkip(liveness) + } + + store, err := OpenStore(filepath.Join(dir, "test.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { store.Close() }) + t.Cleanup(StartStatsFileWriter(store, 20*time.Millisecond)) + + var raw []byte + waitForID(t, "stats file", func() bool { + raw, err = os.ReadFile(statsPath) + return err == nil && strings.Contains(string(raw), "source_statuses") + }) + assertNoSecretParts118(t, "stats file", string(raw)) + for _, want := range []string{`"tcp://****@host:1883"`, `"tcp://****@host:1883 (2)"`} { + if !strings.Contains(string(raw), want) { + t.Errorf("stats file misses %s: %s", want, raw) + } + } + fi, err := os.Stat(statsPath) + if err != nil { + t.Fatal(err) + } + if perm := fi.Mode().Perm(); perm != 0o600 { + t.Errorf("stats file mode %o, want 600", perm) + } +} diff --git a/cmd/ingestor/route_mask_backfill_test.go b/cmd/ingestor/route_mask_backfill_test.go index 7b27eddb8..dbd83503c 100644 --- a/cmd/ingestor/route_mask_backfill_test.go +++ b/cmd/ingestor/route_mask_backfill_test.go @@ -304,7 +304,16 @@ func TestMain_StartsRouteMaskBackfillAfterBufferReady(t *testing.T) { t.Fatal(err) } s := string(src) - subscribe := strings.Index(s, "c.Subscribe(") + // #118: the OnConnect handler that subscribes is installed by + // prepareMQTTSource, called from main's connect loop. + wiring, err := os.ReadFile("mqtt_source.go") + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(wiring), "c.Subscribe(") { + t.Fatal("prepareMQTTSource no longer subscribes") + } + subscribe := strings.Index(s, "prepareMQTTSource(") ready := strings.Index(s, "ingestBuffer.Ready()") start := strings.Index(s, "store.StartRouteMaskBackfill(") shutdown := strings.Index(s, `log.Println("Shutting down...")`) diff --git a/cmd/ingestor/source_status.go b/cmd/ingestor/source_status.go index d18637ded..c9815496d 100644 --- a/cmd/ingestor/source_status.go +++ b/cmd/ingestor/source_status.go @@ -37,7 +37,7 @@ type SourceStatusSnapshot struct { // slots = 2.4KB/source — fine. type sourceStatusState struct { name string - broker string // raw broker URL — server-side handler masks the password + broker string // credentials masked (brokerForLog); published in the stats file connected atomic.Bool lastConnectUnix atomic.Int64 @@ -72,14 +72,15 @@ func (s *sourceStatusState) MarkConnect(now time.Time) { s.errMu.Unlock() } -// MarkDisconnect records the broker dropping the connection. -func (s *sourceStatusState) MarkDisconnect(now time.Time, err error) { +// MarkDisconnect records the broker dropping the connection. The error is +// stored with the source's secrets masked (errForLog). +func (s *sourceStatusState) MarkDisconnect(now time.Time, err error, secrets ...string) { s.connected.Store(false) s.lastDisconnectUnix.Store(now.Unix()) s.disconnectCount.Add(1) if err != nil { s.errMu.Lock() - s.lastError = err.Error() + s.lastError = errForLog(err, secrets...) s.errMu.Unlock() } } @@ -137,8 +138,8 @@ func (s *sourceStatusState) snapshot(now time.Time) SourceStatusSnapshot { } // sourceStatusRegistry holds one sourceStatusState per source. Keyed by -// tag (which is the source Name, or the Broker URL if the operator left -// the name blank). +// tag (mqttSourceTags: the source Name, or the credential-free broker if +// the operator left the name blank). var ( sourceStatusRegistryMu sync.RWMutex sourceStatusRegistry = map[string]*sourceStatusState{} @@ -147,14 +148,16 @@ var ( // RegisterSourceStatus creates (or returns the existing) state for the // given source. Safe for cold-start use; idempotent — re-registering the // same tag returns the existing state so counters aren't reset across -// reconnects. +// reconnects. The broker is stored with credentials masked (brokerForLog), +// whatever the caller passes: the stats file that carries it is served by +// the public /api/mqtt/status (#118). func RegisterSourceStatus(tag, broker string) *sourceStatusState { sourceStatusRegistryMu.Lock() defer sourceStatusRegistryMu.Unlock() if s, ok := sourceStatusRegistry[tag]; ok { return s } - s := &sourceStatusState{name: tag, broker: broker} + s := &sourceStatusState{name: tag, broker: brokerForLog(broker)} sourceStatusRegistry[tag] = s return s } diff --git a/cmd/ingestor/source_status_test.go b/cmd/ingestor/source_status_test.go index 06b94d42b..6f7eb7a11 100644 --- a/cmd/ingestor/source_status_test.go +++ b/cmd/ingestor/source_status_test.go @@ -12,7 +12,7 @@ func TestSourceStatus_BasicLifecycle(t *testing.T) { resetSourceStatusRegistry() defer resetSourceStatusRegistry() - s := RegisterSourceStatus("local", "mqtt://broker.example.com:1883") + s := RegisterSourceStatus("local", "mqtt://obsuser:hunter2@broker.example.com:1883/?token=t0k") if s == nil { t.Fatal("RegisterSourceStatus returned nil") } @@ -43,8 +43,14 @@ func TestSourceStatus_BasicLifecycle(t *testing.T) { if snap.LastConnectUnix != now.Unix() { t.Errorf("LastConnectUnix = %d, want %d", snap.LastConnectUnix, now.Unix()) } - if snap.Broker != "mqtt://broker.example.com:1883" { - t.Errorf("Broker = %q, want raw URL passthrough (server masks)", snap.Broker) + // #118: this used to pin a raw passthrough ("server masks"). The + // stats file carrying it is served by the public /api/mqtt/status + // and is readable on disk, and the server's masking was incomplete, + // so the ingestor now stores the broker with user-info masked and + // without query or fragment (brokerForLog); the server masks again as + // a second layer. + if snap.Broker != "mqtt://****@broker.example.com:1883/" { + t.Errorf("Broker = %q, want it without credentials", snap.Broker) } // After 5 minutes idle, sliding window must be empty. diff --git a/cmd/ingestor/stats_file.go b/cmd/ingestor/stats_file.go index da6f6cde0..06a7878de 100644 --- a/cmd/ingestor/stats_file.go +++ b/cmd/ingestor/stats_file.go @@ -4,6 +4,7 @@ import ( "bufio" "bytes" "encoding/json" + "fmt" "log" "os" "sync" @@ -137,12 +138,29 @@ func statsFilePath() string { func writeStatsAtomic(path string, b []byte) error { tmp := path + ".tmp" // O_NOFOLLOW: if tmp is a pre-existing symlink, openat fails with ELOOP - // instead of clobbering the symlink target. O_TRUNC zeroes existing - // regular-file content. 0o600 — no need for world-readable. - f, err := os.OpenFile(tmp, os.O_CREATE|os.O_WRONLY|os.O_TRUNC|oNoFollow, 0o600) + // instead of clobbering the symlink target. 0o600 — no need for + // world-readable. + f, err := os.OpenFile(tmp, os.O_CREATE|os.O_WRONLY|oNoFollow, 0o600) if err != nil { return err } + // The mode above applies only to a new file. A stale tmp keeps its + // owner and mode, and the rename would publish both. One that belongs + // to another user is refused before anything is changed: root's chmod + // would succeed on it, and its owner may still hold it open. Ours gets + // 0o600 and loses its old content (#118). + if err := checkStatsTmpOwner(f); err != nil { + f.Close() + return err + } + if err := f.Chmod(0o600); err != nil { + f.Close() + return err + } + if err := f.Truncate(0); err != nil { + f.Close() + return err + } if _, err := f.Write(b); err != nil { f.Close() os.Remove(tmp) @@ -159,6 +177,23 @@ func writeStatsAtomic(path string, b []byte) error { return nil } +// statsFileEUID is the user the stats tmp file must belong to; swapped in +// tests to model a file owned by someone else. +var statsFileEUID = os.Geteuid + +// checkStatsTmpOwner fails when f belongs to a user other than +// statsFileEUID. Where files have no Unix owner (Windows) it passes. +func checkStatsTmpOwner(f *os.File) error { + fi, err := f.Stat() + if err != nil { + return err + } + if uid, ok := fileOwnerUID(fi); ok && uid != statsFileEUID() { + return fmt.Errorf("%s belongs to uid %d, not %d", f.Name(), uid, statsFileEUID()) + } + return nil +} + // procIOSnapshot is the raw counter snapshot used to compute per-second rates // across two consecutive ticks of the stats-file writer. type procIOSnapshot struct { diff --git a/cmd/ingestor/stats_file_owner_unix.go b/cmd/ingestor/stats_file_owner_unix.go new file mode 100644 index 000000000..7d42787c2 --- /dev/null +++ b/cmd/ingestor/stats_file_owner_unix.go @@ -0,0 +1,17 @@ +//go:build !windows + +package main + +import ( + "os" + "syscall" +) + +// fileOwnerUID returns the user ID that owns fi. +func fileOwnerUID(fi os.FileInfo) (int, bool) { + st, ok := fi.Sys().(*syscall.Stat_t) + if !ok { + return 0, false + } + return int(st.Uid), true +} diff --git a/cmd/ingestor/stats_file_owner_windows.go b/cmd/ingestor/stats_file_owner_windows.go new file mode 100644 index 000000000..b0739a357 --- /dev/null +++ b/cmd/ingestor/stats_file_owner_windows.go @@ -0,0 +1,9 @@ +//go:build windows + +package main + +import "os" + +// fileOwnerUID reports no owner on Windows, which has no Unix user IDs; the +// ingestor is only deployed on Linux. +func fileOwnerUID(os.FileInfo) (int, bool) { return 0, false } diff --git a/cmd/server/go.mod b/cmd/server/go.mod index 8e7cc2464..cf2705eb0 100644 --- a/cmd/server/go.mod +++ b/cmd/server/go.mod @@ -64,3 +64,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 diff --git a/cmd/server/mqtt_status.go b/cmd/server/mqtt_status.go index 6dab0ee35..3e03a4a02 100644 --- a/cmd/server/mqtt_status.go +++ b/cmd/server/mqtt_status.go @@ -3,92 +3,107 @@ package main import ( "encoding/json" "net/http" - "net/url" "os" - "regexp" + "sort" + "strconv" "strings" + + "github.com/meshcore-analyzer/brokerurl" ) -// mqttBrokerSchemes is the set of broker URL schemes whose embedded -// `user:pass@host` credentials we want to redact. We URL-parse for these -// (defense vs. passwords containing `@`); other strings fall through to -// the legacy regex pass for embedded user:pass occurrences in free-form -// error strings. -var mqttBrokerSchemes = map[string]bool{ - "mqtt": true, "mqtts": true, "tcp": true, "ssl": true, "ws": true, "wss": true, -} +// maskBrokerURL returns a broker URL fit for the public /api/mqtt/status: +// all user-info (a lone user name or token too) becomes "****", and the +// query and fragment are dropped, also without a scheme and for URLs that +// url.Parse rejects (brokerurl.Mask). The ingestor already masks its +// brokers (#118); this is the second layer, for any ingestor version. +// `mqtt://user:secret@host:1883` -> `mqtt://****@host:1883`. +func maskBrokerURL(s string) string { return brokerurl.Mask(s) } -// mqttBrokerURLRe locates a broker URL (with credentials) embedded inside -// a larger free-form string — e.g. an error message that quotes the -// failing broker. Each match is fed through url.Parse + redaction. We -// match greedily up through the LAST `@` followed by a host-shaped token -// so passwords containing `@` are not truncated (#1682 adversarial r1). -// -// Go's RE2 has no lookahead; we capture the host tail and emit it -// unchanged in the replacement. -var mqttBrokerURLRe = regexp.MustCompile(`(?i)(?:mqtt|mqtts|tcp|ssl|ws|wss)://[^\s]*`) +// maskBrokerText masks every broker URL or user-info in free text, such as +// an error message (brokerurl.MaskText). +func maskBrokerText(s string) string { return brokerurl.MaskText(s) } -// maskBrokerURL returns the broker URL with any inline password redacted. -// `mqtt://user:secret@host:1883` -> `mqtt://user:****@host:1883`. -// `mqtt://user:p@ss@host` -> `mqtt://user:****@host` (password with `@`). -// URLs without inline credentials are returned unchanged. -// -// Primary strategy: url.Parse — handles passwords with `@`, `:`, etc. -// Fallback: regex sweep for free-form strings (e.g. error messages that -// quote a URL fragment but aren't standalone-parseable). -func maskBrokerURL(s string) string { - if s == "" { - return s +// maskSourceName masks a source name or tag that is a raw broker URL: an +// older ingestor tagged an unnamed source with its raw broker. That is a +// name among rawBrokers (the brokers in the stats file's source_statuses) +// or one holding "://"; it is masked as one broker URL (Mask, which also +// covers whitespace in a password). Names the operator chose, such as +// "obs@north" or "Feed @ CPH", are returned as they are (#118). +func maskSourceName(s string, rawBrokers map[string]bool) string { + if rawBrokers[s] || strings.Contains(s, "://") { + return brokerurl.Mask(s) } - // Fast path: the whole string is the broker URL. - if masked, ok := redactBrokerURL(s); ok { - return masked - } - // Fallback: free-form string (e.g. error message) containing a URL. - // Find embedded broker URLs and redact each in-place. - return mqttBrokerURLRe.ReplaceAllStringFunc(s, func(m string) string { - if out, ok := redactBrokerURL(m); ok { - return out + return s +} + +// rawBrokerSet returns the non-empty brokers of statuses as the stats file +// holds them, for maskSourceName. +func rawBrokerSet(statuses []MqttSourceStatus) map[string]bool { + set := make(map[string]bool, len(statuses)) + for _, s := range statuses { + if s.Broker != "" { + set[s.Broker] = true } - return m - }) + } + return set } -// redactBrokerURL parses s as a URL and, if it has an mqtt-family scheme -// with userinfo containing a password, returns the URL with the password -// replaced by `****`. Returns ok=false when s is not such a URL. -func redactBrokerURL(s string) (string, bool) { - u, err := url.Parse(s) - if err != nil || u.Scheme == "" || u.User == nil { - return s, false +// maskLivenessKeys masks the keys (source tags) of the ingestor's liveness +// map (maskSourceName) for the public /api/healthz: an older ingestor +// tagged an unnamed source with its raw broker (#118). Keys that masking +// leaves unchanged keep their name; a masked key that coincides with +// another gets " (2)", " (3)", … so no entry is lost. Called on a cache +// refresh only. +// +// haveStatuses reports whether the stats file had source_statuses. An +// ingestor built 2026-06-07..06-12 wrote source_liveness without them, so +// rawBrokers is empty and a raw broker without a scheme +// ("user:pass@host:1883") would pass maskSourceName. Without statuses, +// every key holding '@' is therefore masked as a broker URL; a name the +// operator chose with an '@' loses what precedes it then. +func maskLivenessKeys(m map[string]SourceLivenessSnapshot, rawBrokers map[string]bool, haveStatuses bool) map[string]SourceLivenessSnapshot { + if len(m) == 0 { + return m } - if !mqttBrokerSchemes[strings.ToLower(u.Scheme)] { - return s, false + mask := func(k string) string { + if !haveStatuses && strings.Contains(k, "@") { + return brokerurl.Mask(k) + } + return maskSourceName(k, rawBrokers) } - if _, hasPass := u.User.Password(); !hasPass { - return s, false + keys := make([]string, 0, len(m)) + for k := range m { + keys = append(keys, k) } - // Re-assemble manually rather than via url.UserPassword + u.String() - // because the latter percent-encodes the `*` mask token into `%2A`, - // defeating the user-visible redaction marker. We only need to swap - // the userinfo segment of the original string. - hostAndAfter := s - if idx := strings.LastIndex(s, "@"); idx >= 0 { - hostAndAfter = s[idx+1:] + sort.Strings(keys) + out := make(map[string]SourceLivenessSnapshot, len(m)) + var masked []string + for _, k := range keys { + if mask(k) == k { + out[k] = m[k] + } else { + masked = append(masked, k) + } } - // Preserve original scheme casing (url.Parse lowercases u.Scheme). - schemeEnd := strings.Index(s, "://") - if schemeEnd < 0 { - return s, false + for _, k := range masked { + base := mask(k) + key := base + for n := 2; ; n++ { + if _, taken := out[key]; !taken { + break + } + key = base + " (" + strconv.Itoa(n) + ")" + } + out[key] = m[k] } - return s[:schemeEnd] + "://" + u.User.Username() + ":****@" + hostAndAfter, true + return out } // MqttSourceStatus is the per-MQTT-source status row surfaced via // /api/mqtt/status. Mirrors the on-disk shape the ingestor publishes // (cmd/ingestor SourceStatusSnapshot) but with the broker URL credentials -// redacted before serving — operators must not see the broker password -// in the API response (#1043 acceptance criterion). +// redacted before serving — the endpoint needs no API key, so nobody may +// see a broker password or token in the response (#1043, #118). type MqttSourceStatus struct { Name string `json:"name"` Broker string `json:"broker"` @@ -141,7 +156,7 @@ type ingestorMqttStatusEnvelope struct { } // handleMqttStatus serves GET /api/mqtt/status. Reads the ingestor stats -// file, masks broker-URL passwords, and returns the per-source status +// file, masks broker-URL credentials, and returns the per-source status // list. Returns an empty list (200 OK) when the stats file is missing // or unparseable — the UI panel renders a "no data yet" state. func (s *Server) handleMqttStatus(w http.ResponseWriter, r *http.Request) { @@ -160,11 +175,14 @@ func (s *Server) handleMqttStatus(w http.ResponseWriter, r *http.Request) { resp.WatchdogLastTickUnix = env.WatchdogLastTickUnix resp.WatchdogPanicCount = env.WatchdogPanicCount resp.WatchdogLogDropCount = env.WatchdogLogDropCount + rawBrokers := rawBrokerSet(env.SourceStatuses) for _, src := range env.SourceStatuses { + // An older ingestor tagged an unnamed source with its raw + // broker, and broker libraries occasionally quote the failing + // URL in the error string — mask both (#118). + src.Name = maskSourceName(src.Name, rawBrokers) src.Broker = maskBrokerURL(src.Broker) - // Broker libraries occasionally quote the failing URL in the - // error string — redact there too as defense-in-depth. - src.LastError = maskBrokerURL(src.LastError) + src.LastError = maskBrokerText(src.LastError) resp.Sources = append(resp.Sources, src) } writeJSON(w, resp) diff --git a/cmd/server/mqtt_status_118_test.go b/cmd/server/mqtt_status_118_test.go new file mode 100644 index 000000000..edbf5e5ba --- /dev/null +++ b/cmd/server/mqtt_status_118_test.go @@ -0,0 +1,348 @@ +package main + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" +) + +// #118 re-review: /api/mqtt/status is public (no API key). The broker, +// name and lastError it serves come from the ingestor stats file, and an +// ingestor of any version may have written them raw, so the server masks +// defensively: all user-info (a lone user name or token too), the query +// and the fragment, also without a scheme and for URLs url.Parse rejects. + +// secrets118 are the credential parts used in the inputs below; none may +// reach a response. +var secrets118 = []string{"secret", "tok3n", "abc", "dev-user", "p%zz", "2024", "1234", "frag", "sec ret"} + +func assertNoSecrets118(t *testing.T, what, s string) { + t.Helper() + for _, bad := range secrets118 { + if strings.Contains(s, bad) { + t.Errorf("%s leaks %q: %s", what, bad, s) + } + } +} + +func TestMaskBrokerURLStripsAllCredentials_118(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {"dev-user:secret@host:1883", "****@host:1883"}, // no scheme + {"tcp://tok3n@host", "tcp://****@host"}, // user name / token only + {"wss://host/mqtt?token=abc", "wss://host/mqtt"}, // query + {"wss://host/mqtt#frag", "wss://host/mqtt"}, // fragment + {"tcp://dev-user:p%zz@host", "tcp://****@host"}, // url.Parse fails + {"mqtt://dev-user:secret@host:1883", "mqtt://****@host:1883"}, // user name too + {"tcp://dev-user:2024/secret@host:1883", "tcp://****@host:1883"}, + {"tcp://dev-user:1234?abc@host", "tcp://****@host"}, + {"dev-user:sec://ret@host", "****@host"}, // "://" inside the password + {"MQTTS://dev-user:secret@host:8883/x?abc", "MQTTS://****@host:8883/x"}, + {"mqtt://u:p@broker:1883", "mqtt://****@broker:1883"}, + {"wss://host/mqtt?u=me@x.org", "wss://****@x.org"}, // the '@' in the query: cut, and marked as cut + // unchanged + {"mqtt://broker.example.com:1883", "mqtt://broker.example.com:1883"}, + {"wss://broker.example.com/mqtt", "wss://broker.example.com/mqtt"}, + {"host:1883", "host:1883"}, + {"", ""}, + } { + got := maskBrokerURL(c.in) + assertNoSecrets118(t, "maskBrokerURL("+c.in+")", got) + if got != c.want { + t.Errorf("maskBrokerURL(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +// Free-form text (lastError, and name when an older ingestor used the raw +// broker as tag): every broker-shaped token is masked, the rest is kept. +func TestMaskBrokerTextStripsCredentials_118(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {`dial "wss://host/mqtt?token=abc": bad handshake`, `dial "wss://host/mqtt bad handshake`}, + {"network error tcp://dev-user:secret@host:1883: refused", "network error tcp://****@host:1883: refused"}, + {"auth dev-user:secret@host failed", "auth ****@host failed"}, + {"EOF", "EOF"}, + {"read tcp 10.0.0.1:5000->10.0.0.2:1883: connection reset by peer", "read tcp 10.0.0.1:5000->10.0.0.2:1883: connection reset by peer"}, + {"Local Feed #1", "Local Feed #1"}, + } { + got := maskBrokerText(c.in) + assertNoSecrets118(t, "maskBrokerText("+c.in+")", got) + if got != c.want { + t.Errorf("maskBrokerText(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +func writeStats118(t *testing.T, v any) { + t.Helper() + p := filepath.Join(t.TempDir(), "ingestor-stats.json") + t.Setenv("CORESCOPE_INGESTOR_STATS", p) + b, err := json.Marshal(v) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(p, b, 0o600); err != nil { + t.Fatal(err) + } +} + +// End to end through the handler, with a stats file as an older ingestor +// writes it (raw broker, raw broker as the tag of an unnamed source). +func TestMqttStatusServesNoCredentials_118(t *testing.T) { + var sources []map[string]any + for _, raw := range []string{ + "dev-user:secret@host:1883", + "tcp://tok3n@host", + "wss://host/mqtt?token=abc", + "tcp://dev-user:p%zz@host", + "tcp://dev-user:2024/secret@host:1883", + "tcp://dev-user:sec ret@host:1883", // a space splits MaskText tokens + } { + sources = append(sources, map[string]any{ + "name": raw, + "broker": raw, + "lastError": "connect " + raw + " failed", + }) + } + writeStats118(t, map[string]any{"sampledAt": "2026-09-30T12:00:00Z", "source_statuses": sources}) + + rec := httptest.NewRecorder() + (&Server{}).handleMqttStatus(rec, httptest.NewRequest(http.MethodGet, "/api/mqtt/status", nil)) + if rec.Code != http.StatusOK { + t.Fatalf("status %d", rec.Code) + } + assertNoSecrets118(t, "/api/mqtt/status", rec.Body.String()) + var resp MqttStatusResponse + if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil { + t.Fatal(err) + } + if len(resp.Sources) != len(sources) { + t.Fatalf("got %d sources, want %d", len(resp.Sources), len(sources)) + } + for _, s := range resp.Sources { + if !strings.Contains(s.Broker, "host") || !strings.Contains(s.Name, "host") { + t.Errorf("masked row lost its host: name %q broker %q", s.Name, s.Broker) + } + // only the broker in the error text is masked, not the message + if !strings.HasPrefix(s.LastError, "connect ") || !strings.HasSuffix(s.LastError, " failed") || !strings.Contains(s.LastError, "host") { + t.Errorf("lastError masked beyond its broker: %q", s.LastError) + } + } +} + +// /api/healthz (public) keys ingest_liveness by source tag, which an older +// ingestor set to the raw broker of an unnamed source. Masked keys that +// coincide stay separate entries. +func TestIngestLivenessKeysHaveNoCredentials_118(t *testing.T) { + resetSourceLivenessCache() + t.Cleanup(resetSourceLivenessCache) + writeStats118(t, map[string]any{"source_liveness": map[string]any{ + "tcp://dev-user:secret@host:1883": map[string]int64{"lastReceiptUnix": 1}, + "tcp://tok3n@host:1883": map[string]int64{"lastReceiptUnix": 2}, + "tcp://dev-user:sec ret@host:1883": map[string]int64{"lastReceiptUnix": 4}, + "feed": map[string]int64{"lastReceiptUnix": 3}, + }}) + got := readIngestorSourceLiveness() + if len(got) != 4 { + t.Fatalf("got %d entries, want 4: %v", len(got), got) + } + if _, ok := got["feed"]; !ok { + t.Errorf("plain tag renamed: %v", got) + } + receipts := map[int64]bool{} + for k, v := range got { + assertNoSecrets118(t, "ingest_liveness key", k) + receipts[v.LastReceiptUnix] = true + } + if len(receipts) != 4 { + t.Errorf("entries merged or lost: %v", got) + } +} + +// A name is masked only when it is a raw broker URL: one of the stats +// file's raw brokers (an older ingestor tagged an unnamed source with it), +// or one holding "://". Whitespace in a password is covered then, also +// without a scheme; names the operator chose are left alone, '@' or not. +func TestMaskSourceName_118(t *testing.T) { + raw := map[string]bool{"dev-user:sec ret@host:1883": true} + for _, c := range []struct{ in, want string }{ + {"dev-user:sec ret@host:1883", "****@host:1883"}, + {"tcp://dev-user:sec ret@host:1883", "tcp://****@host:1883"}, + {"tcp://host:1883 (2)", "tcp://host:1883 (2)"}, + {"tcp://****@host:1883 (2)", "tcp://****@host:1883 (2)"}, + {"obs@north", "obs@north"}, + {"Feed @ CPH", "Feed @ CPH"}, + {"Local Feed #1", "Local Feed #1"}, + {"feed", "feed"}, + } { + got := maskSourceName(c.in, raw) + assertNoSecrets118(t, "maskSourceName("+c.in+")", got) + if got != c.want { + t.Errorf("maskSourceName(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +// #118 round 3: names the operator chose are served as they are, and no +// longer collide into " (2)"; a raw broker as name is still masked. +func TestMqttStatusKeepsChosenNames_118(t *testing.T) { + writeStats118(t, map[string]any{"source_statuses": []map[string]string{ + {"name": "obs@north", "broker": "tcp://north:1883"}, + {"name": "obs@south", "broker": "tcp://south:1883"}, + {"name": "Feed @ CPH", "broker": "tcp://cph:1883"}, + {"name": "dev-user:secret@host:1883", "broker": "dev-user:secret@host:1883"}, + {"name": "tcp://tok3n@host", "broker": "tcp://tok3n@host"}, + }}) + rec := httptest.NewRecorder() + (&Server{}).handleMqttStatus(rec, httptest.NewRequest(http.MethodGet, "/api/mqtt/status", nil)) + assertNoSecrets118(t, "/api/mqtt/status", rec.Body.String()) + var resp MqttStatusResponse + if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil { + t.Fatal(err) + } + var names []string + for _, s := range resp.Sources { + names = append(names, s.Name) + } + want := []string{"obs@north", "obs@south", "Feed @ CPH", "****@host:1883", "tcp://****@host"} + if strings.Join(names, "|") != strings.Join(want, "|") { + t.Errorf("names\n got %q\nwant %q", names, want) + } +} + +func TestIngestLivenessKeepsChosenNames_118(t *testing.T) { + resetSourceLivenessCache() + t.Cleanup(resetSourceLivenessCache) + writeStats118(t, map[string]any{ + "source_statuses": []map[string]string{ + {"name": "obs@north", "broker": "tcp://north:1883"}, + {"name": "dev-user:secret@host:1883", "broker": "dev-user:secret@host:1883"}, + }, + "source_liveness": map[string]any{ + "obs@north": map[string]int64{"lastReceiptUnix": 1}, + "obs@south": map[string]int64{"lastReceiptUnix": 2}, + "Feed @ CPH": map[string]int64{"lastReceiptUnix": 3}, + "dev-user:secret@host:1883": map[string]int64{"lastReceiptUnix": 4}, + "tcp://tok3n@host:1883": map[string]int64{"lastReceiptUnix": 5}, + }, + }) + got := readIngestorSourceLiveness() + want := map[string]int64{"obs@north": 1, "obs@south": 2, "Feed @ CPH": 3, "****@host:1883": 4, "tcp://****@host:1883": 5} + if len(got) != len(want) { + t.Fatalf("got %v, want %v", got, want) + } + for k, v := range want { + if got[k].LastReceiptUnix != v { + t.Errorf("key %q: got %v, want lastReceiptUnix %d (all: %v)", k, got[k], v, got) + } + } +} + +// healthzLiveness118 serves /api/healthz and returns the raw body and its +// ingest_liveness as key -> lastReceiptUnix. +func healthzLiveness118(t *testing.T) (string, map[string]int64) { + t.Helper() + resetSourceLivenessCache() + t.Cleanup(resetSourceLivenessCache) + readiness.Store(1) + t.Cleanup(func() { readiness.Store(0) }) + rec := httptest.NewRecorder() + (&Server{store: &PacketStore{}}).handleHealthz(rec, httptest.NewRequest(http.MethodGet, "/api/healthz", nil)) + if rec.Code != http.StatusOK { + t.Fatalf("/api/healthz: status %d", rec.Code) + } + var resp struct { + IngestLiveness map[string]struct { + LastReceiptUnix int64 `json:"lastReceiptUnix"` + } `json:"ingest_liveness"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil { + t.Fatal(err) + } + got := make(map[string]int64, len(resp.IngestLiveness)) + for k, v := range resp.IngestLiveness { + got[k] = v.LastReceiptUnix + } + return rec.Body.String(), got +} + +func assertLiveness118(t *testing.T, got, want map[string]int64) { + t.Helper() + if len(got) != len(want) { + t.Fatalf("ingest_liveness = %v, want %v", got, want) + } + for k, v := range want { + if r, ok := got[k]; !ok || r != v { + t.Errorf("ingest_liveness[%q] = %d (present %v), want %d (all: %v)", k, r, ok, v, got) + } + } +} + +// #118 R1: an ingestor built 2026-06-07..06-12 wrote source_liveness +// without source_statuses. There is then no raw-broker list, and a raw +// broker without a scheme as tag must still not reach /api/healthz. +func TestHealthzMasksBrokerTagsWithoutStatuses_118(t *testing.T) { + liveness := map[string]any{ + "user:secret@host:1883": map[string]int64{"lastReceiptUnix": 1}, + "mqtt://u:p@h:1883": map[string]int64{"lastReceiptUnix": 2}, + "feed": map[string]int64{"lastReceiptUnix": 3}, + } + for name, stats := range map[string]map[string]any{ + "no source_statuses": {"source_liveness": liveness}, + "empty source_statuses": {"source_liveness": liveness, "source_statuses": []any{}}, + } { + t.Run(name, func(t *testing.T) { + writeStats118(t, stats) + body, got := healthzLiveness118(t) + for _, bad := range []string{"secret", "u:p", "user:"} { + if strings.Contains(body, bad) { + t.Errorf("/api/healthz leaks %q: %s", bad, body) + } + } + assertLiveness118(t, got, map[string]int64{"****@host:1883": 1, "mqtt://****@h:1883": 2, "feed": 3}) + }) + } +} + +// With source_statuses present the raw-broker list decides, as before: a +// name the operator chose keeps its '@'. +func TestHealthzKeepsChosenNamesWithStatuses_118(t *testing.T) { + writeStats118(t, map[string]any{ + "source_statuses": []map[string]string{ + {"name": "obs@north", "broker": "tcp://north:1883"}, + {"name": "", "broker": "user:secret@host:1883"}, + }, + "source_liveness": map[string]any{ + "user:secret@host:1883": map[string]int64{"lastReceiptUnix": 1}, + "mqtt://u:p@h:1883": map[string]int64{"lastReceiptUnix": 2}, + "obs@north": map[string]int64{"lastReceiptUnix": 3}, + }, + }) + body, got := healthzLiveness118(t) + for _, bad := range []string{"secret", "u:p", "user:"} { + if strings.Contains(body, bad) { + t.Errorf("/api/healthz leaks %q: %s", bad, body) + } + } + assertLiveness118(t, got, map[string]int64{"****@host:1883": 1, "mqtt://****@h:1883": 2, "obs@north": 3}) +} + +// Without source_statuses, tags that mask to the same value each keep an +// entry, with " (2)", " (3)". +func TestHealthzMaskedTagCollisionsWithoutStatuses_118(t *testing.T) { + writeStats118(t, map[string]any{"source_liveness": map[string]any{ + "user:secret@host:1883": map[string]int64{"lastReceiptUnix": 1}, + "other:secret@host:1883": map[string]int64{"lastReceiptUnix": 2}, + "tok3n@host:1883": map[string]int64{"lastReceiptUnix": 3}, + }}) + body, got := healthzLiveness118(t) + assertNoSecrets118(t, "/api/healthz", body) + if strings.Contains(body, "other:") || strings.Contains(body, "user:") { + t.Errorf("/api/healthz leaks a user name: %s", body) + } + // Sorted order of the raw tags decides the suffixes. + assertLiveness118(t, got, map[string]int64{"****@host:1883": 2, "****@host:1883 (2)": 3, "****@host:1883 (3)": 1}) +} diff --git a/cmd/server/mqtt_status_test.go b/cmd/server/mqtt_status_test.go index db3415e12..9327698a8 100644 --- a/cmd/server/mqtt_status_test.go +++ b/cmd/server/mqtt_status_test.go @@ -11,14 +11,15 @@ import ( ) // TestMqttStatus_MasksBrokerPassword (#1043) asserts the /api/mqtt/status -// handler never leaks the broker password embedded in a mqtt:// URL. +// handler never leaks the broker password (nor, since #118, the user +// name) embedded in a mqtt:// URL. // Operators viewing the API response (or the Observers panel that // consumes it) must see `****` in place of the inline credential. // // Test shape: write a stub ingestor stats file with one source whose // broker URL contains a plaintext password, invoke the handler, assert -// the JSON response (a) contains the username + host, (b) does NOT -// contain the password substring. +// the JSON response (a) contains the host, (b) does NOT contain the +// password or user name substring. func TestMqttStatus_MasksBrokerPassword(t *testing.T) { const password = "hunter2supersecret" const rawBroker = "mqtt://obsuser:" + password + "@broker.example.com:1883" @@ -67,8 +68,10 @@ func TestMqttStatus_MasksBrokerPassword(t *testing.T) { if !strings.Contains(body, "broker.example.com") { t.Errorf("response missing broker host: %s", body) } - if !strings.Contains(body, "obsuser") { - t.Errorf("response missing broker username: %s", body) + // #118: the user name is masked too. The endpoint is public, and a + // user name alone is often the credential (a token). + if strings.Contains(body, "obsuser") { + t.Errorf("response leaks broker username: %s", body) } // Mask token must be present so operators can tell credentials were // redacted vs the broker URL never having a password to begin with. @@ -103,29 +106,31 @@ func TestMqttStatus_EmptyWhenNoStatsFile(t *testing.T) { // TestMaskBrokerURL_Patterns is a unit table-driven test for the masking // helper. Kept separate from the handler test so a regression in the -// regex localizes immediately. +// masking localizes immediately. Since #118 all user-info is masked, the +// user name included (the endpoint is public and the user name may be the +// credential), so the expectations are `scheme://****@host`; more forms +// in mqtt_status_118_test.go. func TestMaskBrokerURL_Patterns(t *testing.T) { cases := []struct { name, in, want string }{ {"plain mqtt no creds", "mqtt://broker.example.com:1883", "mqtt://broker.example.com:1883"}, - {"mqtt with creds", "mqtt://u:secret@broker.example.com:1883", "mqtt://u:****@broker.example.com:1883"}, - {"mqtts with creds", "mqtts://u:secret@broker.example.com:8883", "mqtts://u:****@broker.example.com:8883"}, - {"tcp with creds", "tcp://u:p@host:1883", "tcp://u:****@host:1883"}, - {"ssl with creds", "ssl://u:p@host:8883", "ssl://u:****@host:8883"}, - {"ws with creds", "ws://u:p@host:8080/mqtt", "ws://u:****@host:8080/mqtt"}, - {"wss with creds", "wss://u:p@host:443/mqtt", "wss://u:****@host:443/mqtt"}, - {"uppercase scheme", "MQTT://u:p@host:1883", "MQTT://u:****@host:1883"}, + {"mqtt with creds", "mqtt://u:secret@broker.example.com:1883", "mqtt://****@broker.example.com:1883"}, + {"mqtts with creds", "mqtts://u:secret@broker.example.com:8883", "mqtts://****@broker.example.com:8883"}, + {"tcp with creds", "tcp://u:p@host:1883", "tcp://****@host:1883"}, + {"ssl with creds", "ssl://u:p@host:8883", "ssl://****@host:8883"}, + {"ws with creds", "ws://u:p@host:8080/mqtt", "ws://****@host:8080/mqtt"}, + {"wss with creds", "wss://u:p@host:443/mqtt", "wss://****@host:443/mqtt"}, + {"uppercase scheme", "MQTT://u:p@host:1883", "MQTT://****@host:1883"}, {"empty", "", ""}, - {"long password", "mqtt://obsuser:hunter2supersecretXYZ123@host:1883", "mqtt://obsuser:****@host:1883"}, + {"long password", "mqtt://obsuser:hunter2supersecretXYZ123@host:1883", "mqtt://****@host:1883"}, {"no scheme bare host", "host:1883", "host:1883"}, // Adversarial r1 review (#1682): password contains @. The previous // regex-only impl matched only up to the FIRST @, exposing "ss" as - // part of the path: "mqtt://user:****@ss@host". url.Parse handles - // this correctly because Go interprets the LAST @ as the userinfo - // boundary. - {"password with single @", "mqtt://user:p@ss@host:1883", "mqtt://user:****@host:1883"}, - {"password with multiple @", "mqtt://user:p@ss@wo@host:1883", "mqtt://user:****@host:1883"}, + // part of the path: "mqtt://user:****@ss@host". The LAST @ is + // the user-info boundary. + {"password with single @", "mqtt://user:p@ss@host:1883", "mqtt://****@host:1883"}, + {"password with multiple @", "mqtt://user:p@ss@wo@host:1883", "mqtt://****@host:1883"}, } for _, c := range cases { t.Run(c.name, func(t *testing.T) { diff --git a/cmd/server/openapi.go b/cmd/server/openapi.go index 33f46f12e..0afe3f571 100644 --- a/cmd/server/openapi.go +++ b/cmd/server/openapi.go @@ -47,7 +47,7 @@ func routeDescriptions() map[string]routeMeta { "GET /api/health": {Summary: "Health check", Description: "Returns server health, uptime, and memory stats.", Tag: "admin"}, "GET /api/stats": {Summary: "Network statistics", Description: "Returns aggregate stats (node counts, packet counts, observer counts). Cached for 10s.", Tag: "admin"}, "GET /api/perf": {Summary: "Performance statistics", Description: "Returns per-endpoint request timing and slow query log.", Tag: "admin"}, - "GET /api/mqtt/status": {Summary: "MQTT source status", Description: "Returns per-MQTT-source connection state and counters (lastConnectUnix, lastPacketUnix, packetsTotal, etc.). Broker URL passwords are masked. Sourced from the ingestor stats file; empty list when unavailable. (#1043)", Tag: "admin"}, + "GET /api/mqtt/status": {Summary: "MQTT source status", Description: "Returns per-MQTT-source connection state and counters (lastConnectUnix, lastPacketUnix, packetsTotal, etc.). Broker URL credentials (user-info, query, fragment) are masked in broker and lastError, and in a name that is a raw broker URL. Sourced from the ingestor stats file; empty list when unavailable. (#1043)", Tag: "admin"}, "POST /api/perf/reset": {Summary: "Reset performance stats", Tag: "admin", Auth: true}, // "POST /api/admin/prune" removed in #1283 (ingestor owns prune). "GET /api/debug/affinity": {Summary: "Debug neighbor affinity scores", Tag: "admin", Auth: true}, diff --git a/cmd/server/perf_io.go b/cmd/server/perf_io.go index c20cdc2cd..d67b31533 100644 --- a/cmd/server/perf_io.go +++ b/cmd/server/perf_io.go @@ -405,6 +405,12 @@ func readIngestorSourceLiveness() map[string]SourceLivenessSnapshot { sourceLivenessCache.mtime = time.Time{} return nil } + // public via /api/healthz: keys that are a raw broker get masked (#118) + var statuses ingestorMqttStatusEnvelope + // data parsed into st, so only a type mismatch can fail here, and + // json still fills the other fields (Broker among them) then. + _ = json.Unmarshal(data, &statuses) + st.SourceLiveness = maskLivenessKeys(st.SourceLiveness, rawBrokerSet(statuses.SourceStatuses), len(statuses.SourceStatuses) > 0) sourceLivenessCache.path = path sourceLivenessCache.value = st.SourceLiveness sourceLivenessCache.cachedAt = now diff --git a/config.example.json b/config.example.json index b0282f606..5ed9f6f17 100644 --- a/config.example.json +++ b/config.example.json @@ -370,7 +370,7 @@ "criticalMv": 3000, "_comment": "Voltage cutoffs (millivolts) for the per-node battery trend chart on /node-analytics. Latest sample below lowMv shows the node as ⚠️ Low; below criticalMv shows 🪫 Critical. Both default to 3300 / 3000 if omitted. Source data: observer_metrics.battery_mv populated from observer status messages; only nodes that are themselves observers (matching pubkey ↔ observer id) yield a series. Issue #663." }, - "_comment_mqttSources": "Each source connects to an MQTT broker. Supported schemes: mqtt:// (plain TCP), mqtts:// (TLS), ws:// (WebSocket), wss:// (WebSocket TLS). topics: what to subscribe to. iataFilter: only ingest packets from these regions (optional). region: default IATA region for this source — used when packet/topic doesn't specify one (optional, priority: payload > topic > this field).", + "_comment_mqttSources": "Each source connects to an MQTT broker. Supported schemes: mqtt:// (plain TCP), mqtts:// (TLS), ws:// (WebSocket), wss:// (WebSocket TLS). topics: what to subscribe to. iataFilter: only ingest packets from these regions (optional). region: default IATA region for this source — used when packet/topic doesn't specify one (optional, priority: payload > topic > this field). clientId: MQTT client ID, used verbatim (optional; must be unique among the broker's concurrent clients, and at most 23 characters of [0-9A-Za-z] for strict MQTT 3.1 brokers). When omitted, each ingestor start generates corescope--<8 random hex>, so copies of this file never share an ID — do not add a fixed clientId here.", "compression": { "gzip": false, "websocket": false, diff --git a/internal/brokerurl/brokerurl.go b/internal/brokerurl/brokerurl.go new file mode 100644 index 000000000..29dea72af --- /dev/null +++ b/internal/brokerurl/brokerurl.go @@ -0,0 +1,141 @@ +// Package brokerurl removes credentials from MQTT broker URLs. The ingestor +// uses it for everything it logs or publishes in its stats file, and the +// server, which serves that file on public endpoints, applies it again to +// whatever an ingestor of any version wrote there. +// +// It deliberately does not use url.Parse. A password with an unescaped '/', +// '?' or '#' ends the authority for url.Parse, which then reports the user +// name as the host and the password as port, path or query; one with a bad +// %-escape makes it fail. Here everything up to the last '@' is user-info, +// and the query and fragment (which may carry a token) are dropped. A +// broker with an '@' in its path or query therefore shows only what +// follows that '@': showing too little is the safe side. +package brokerurl + +import ( + "regexp" + "strings" +) + +// Marker replaces removed user-info in Mask and MaskText, so a reader can +// tell credentials were present, and that a host shown after it may be +// only what followed an '@' in the path or query. +const Marker = "****" + +// Mask returns s with its user-info replaced by Marker and without query +// and fragment. A scheme ("tcp://") is kept when it is a valid URL scheme. +func Mask(s string) string { + p := split(s) + rest := p.rest + if p.hasUserinfo { + rest = Marker + "@" + rest + } + return join(p.scheme, rest) +} + +// Secrets returns what Mask removes from s, so that a caller can mask the +// same values where they appear without a URL around them (say, a query +// token quoted in an error): the user-info, the query and the fragment, +// each only when non-empty. +func Secrets(s string) []string { + p := split(s) + var out []string + for _, v := range []string{p.userinfo, p.query, p.fragment} { + if v != "" { + out = append(out, v) + } + } + return out +} + +// urlUserinfoRe matches a URL with user-info in free text: from the start of +// the token holding the scheme to the last '@' on the line and the rest of +// that token, so a password with whitespace is covered whole (and any text +// between the scheme and a later '@' is masked with it). +var urlUserinfoRe = regexp.MustCompile(`\S*[A-Za-z][A-Za-z0-9+.\-]*://[^\n]*@\S*`) + +// maskTokenRe matches the whitespace-separated tokens that may hold a +// broker URL or user-info. +var maskTokenRe = regexp.MustCompile(`\S*(?:://|@)\S*`) + +// schemeRe finds a URL scheme inside a token, e.g. after a quote. +var schemeRe = regexp.MustCompile(`[A-Za-z][A-Za-z0-9+.\-]*://`) + +// MaskText masks broker URLs and user-info in free text (an error message, +// a source name) and leaves the rest alone: first every URL with user-info +// (urlUserinfoRe), then every remaining whitespace-separated token that +// contains "://" or "@". Without a scheme, user-info with whitespace in it +// cannot be told from text, and only its last part is masked; a string +// that is a broker URL goes through Mask instead. +func MaskText(s string) string { + s = urlUserinfoRe.ReplaceAllStringFunc(s, maskSpan) + return maskTokenRe.ReplaceAllStringFunc(s, maskSpan) +} + +// maskSpan masks one match of MaskText. Punctuation before a scheme, such +// as the quote in `"wss://…"`, is kept; anything else before it may be +// user-info. +func maskSpan(m string) string { + if loc := schemeRe.FindStringIndex(m); loc != nil && isPunct(m[:loc[0]]) { + return m[:loc[0]] + Mask(m[loc[0]:]) + } + return Mask(m) +} + +// parts is a broker URL as Mask reads it. +type parts struct { + scheme string + userinfo string // everything before the last '@' + hasUserinfo bool // an '@' was present, even with empty user-info + rest string // host, port and path + query string // without '?' + fragment string // without '#' +} + +func split(s string) parts { + var p parts + rest := s + if i := strings.Index(rest, "://"); i > 0 && validScheme(rest[:i]) { + p.scheme, rest = rest[:i], rest[i+len("://"):] + } + if i := strings.LastIndex(rest, "@"); i >= 0 { + p.userinfo, p.hasUserinfo, rest = rest[:i], true, rest[i+1:] + } + if i := strings.IndexByte(rest, '#'); i >= 0 { + rest, p.fragment = rest[:i], rest[i+1:] + } + if i := strings.IndexByte(rest, '?'); i >= 0 { + rest, p.query = rest[:i], rest[i+1:] + } + p.rest = rest + return p +} + +func join(scheme, rest string) string { + if scheme == "" { + return rest + } + return scheme + "://" + rest +} + +// validScheme reports whether s is an RFC 3986 scheme. +func validScheme(s string) bool { + for i, r := range s { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z': + case i > 0 && (r >= '0' && r <= '9' || r == '+' || r == '-' || r == '.'): + default: + return false + } + } + return s != "" +} + +func isPunct(s string) bool { + for _, r := range s { + if r == '@' || r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r > 0x7f { + return false + } + } + return true +} diff --git a/internal/brokerurl/brokerurl_test.go b/internal/brokerurl/brokerurl_test.go new file mode 100644 index 000000000..3656f3f55 --- /dev/null +++ b/internal/brokerurl/brokerurl_test.go @@ -0,0 +1,110 @@ +package brokerurl + +import ( + "strings" + "testing" +) + +// secrets are the credential parts used in the inputs below. +var secrets = []string{"secret", "tok3n", "abc", "dev-user", "p%zz", "2024", "1234", "frag", "sec", "ret"} + +func assertClean(t *testing.T, what, s string) { + t.Helper() + for _, bad := range secrets { + if strings.Contains(s, bad) { + t.Errorf("%s leaks %q: %s", what, bad, s) + } + } +} + +func TestMask(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {"tcp://dev-user:secret@host:1883", "tcp://****@host:1883"}, + {"mqtt://dev-user:secret@host:1883", "mqtt://****@host:1883"}, + {"mqtt://u:p@broker:1883", "mqtt://****@broker:1883"}, + {"dev-user:secret@host:1883", "****@host:1883"}, + {"tcp://tok3n@host", "tcp://****@host"}, + {"wss://host/mqtt?token=abc", "wss://host/mqtt"}, + {"wss://host/mqtt#frag", "wss://host/mqtt"}, + {"wss://host/mqtt?#", "wss://host/mqtt"}, + {"tcp://dev-user:p%zz@host", "tcp://****@host"}, + // unescaped '/', '?' or '#' in the password: url.Parse ends the + // authority there, so only "everything after the last '@'" is safe + {"tcp://dev-user:2024/secret@host:1883", "tcp://****@host:1883"}, + {"tcp://dev-user:1234?abc@host", "tcp://****@host"}, + {"tcp://dev-user:1234#abc@host", "tcp://****@host"}, + {"tcp://dev-user:p@ss@secret@host:1883", "tcp://****@host:1883"}, + // an '@' in the query: what follows it is shown, marked as cut + {"wss://host/mqtt?u=me@x.org", "wss://****@x.org"}, + // "://" inside the password is not a scheme + {"dev-user:sec://ret@host", "****@host"}, + {"MQTTS://dev-user:secret@[::1]:8883/x?abc", "MQTTS://****@[::1]:8883/x"}, + // nothing to mask + {"mqtt://broker.example.com:1883", "mqtt://broker.example.com:1883"}, + {"wss://broker.example.com/mqtt", "wss://broker.example.com/mqtt"}, + {"host:1883", "host:1883"}, + {"", ""}, + } { + got := Mask(c.in) + assertClean(t, "Mask("+c.in+")", got) + if got != c.want { + t.Errorf("Mask(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +// Secrets lists exactly what Mask removes, for literal masking elsewhere. +func TestSecrets(t *testing.T) { + for _, c := range []struct { + in string + want []string + }{ + {"tcp://dev-user:secret@host:1883", []string{"dev-user:secret"}}, + {"wss://host/mqtt?token=abc#frag", []string{"token=abc", "frag"}}, + {"dev-user:1234?abc@host?x=1", []string{"dev-user:1234?abc", "x=1"}}, + {"tcp://@host", nil}, + {"wss://host/mqtt?#", nil}, + {"mqtt://host:1883", nil}, + {"", nil}, + } { + got := Secrets(c.in) + if strings.Join(got, "|") != strings.Join(c.want, "|") || len(got) != len(c.want) { + t.Errorf("Secrets(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +func TestMaskText(t *testing.T) { + for _, c := range []struct{ in, want string }{ + {`dial "wss://host/mqtt?token=abc": bad handshake`, `dial "wss://host/mqtt bad handshake`}, + {`dial "tcp://dev-user:secret@host": refused`, `dial "tcp://****@host": refused`}, + {"network error tcp://dev-user:secret@host:1883: refused", "network error tcp://****@host:1883: refused"}, + {"auth dev-user:secret@host failed", "auth ****@host failed"}, + {"dev-user:sec://ret@host", "****@host"}, + // whitespace in the password: from the scheme to the last '@' + {"connect tcp://dev-user:sec ret@host:1883 failed", "connect tcp://****@host:1883 failed"}, + {`dial "wss://dev-user:a b c@host/mqtt?token=abc": refused`, `dial "wss://****@host/mqtt refused`}, + {"x tcp://dev-user:sec://ret@host y", "x tcp://****@host y"}, + {"tcp://h:1883 (2)", "tcp://h:1883 (2)"}, + {"EOF", "EOF"}, + {"read tcp 10.0.0.1:5000->10.0.0.2:1883: connection reset by peer", "read tcp 10.0.0.1:5000->10.0.0.2:1883: connection reset by peer"}, + {"Local Feed #1", "Local Feed #1"}, + {"", ""}, + } { + got := MaskText(c.in) + assertClean(t, "MaskText("+c.in+")", got) + if got != c.want { + t.Errorf("MaskText(%q) = %q, want %q", c.in, got, c.want) + } + } +} + +// Idempotent: the server masks what the ingestor already masked. +func TestMaskIsIdempotent(t *testing.T) { + for _, in := range []string{"tcp://dev-user:secret@host:1883", "tcp://dev-user:1234?abc@host", "wss://host/mqtt?token=abc", "wss://host/mqtt?u=me@x.org"} { + m := Mask(in) + if Mask(m) != m || MaskText(m) != m { + t.Errorf("Mask not idempotent for %q: %q then %q / %q", in, m, Mask(m), MaskText(m)) + } + } +} diff --git a/internal/brokerurl/go.mod b/internal/brokerurl/go.mod new file mode 100644 index 000000000..eec6dd37c --- /dev/null +++ b/internal/brokerurl/go.mod @@ -0,0 +1,3 @@ +module github.com/meshcore-analyzer/brokerurl + +go 1.22