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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions cmd/ingestor/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,14 @@ type Config struct {
obsIATAWhitelistCached map[string]bool
obsIATAWhitelistOnce sync.Once

// IATAWarnIntervalSec is how often a region dropped by
// ObserverIATAWhitelist is re-logged while it keeps arriving (#110).
// 0 or less means the default, 6 hours. See iata_drop_warn.go.
IATAWarnIntervalSec int `json:"iataWarnIntervalSec,omitempty"`

// iataDropWarn is the bounded per-region throttle for that warning.
iataDropWarn iataDropThrottle

// ObserverBlacklist is a list of observer public keys to drop at ingest.
// Messages from blacklisted observers are silently discarded — no DB writes,
// no UpsertObserver, no observations, no metrics.
Expand Down
146 changes: 146 additions & 0 deletions cmd/ingestor/iata_drop_log_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
package main

import (
"bytes"
"fmt"
"log"
"strings"
"testing"
)

// #110: when observerIATAWhitelist rejects a region, the drop used to be
// silent, so a legitimate region missing from the list lost all its traffic
// with nothing in the log. These tests drive the real handleMessage and read
// the log it writes.

// captureLog runs fn with the standard logger writing to a buffer.
func captureLog(t *testing.T, fn func()) string {
t.Helper()
var buf bytes.Buffer
orig, flags := log.Writer(), log.Flags()
log.SetOutput(&buf)
log.SetFlags(0)
defer func() { log.SetOutput(orig); log.SetFlags(flags) }()
fn()
return buf.String()
}

func regionFilterLines(out string) []string {
var lines []string
for _, l := range strings.Split(out, "\n") {
if strings.Contains(l, "[region-filter]") {
lines = append(lines, l)
}
}
return lines
}

func statusMsg(region, observer string) *mockMessage {
return &mockMessage{
topic: "meshcore/" + region + "/" + observer + "/status",
payload: []byte(`{"origin":"n","noise_floor":-110}`),
}
}

func TestIATAWhitelistDropIsLoggedOncePerRegion(t *testing.T) {
store := newTestStore(t)
cfg := &Config{ObserverIATAWhitelist: []string{"ARN"}}
out := captureLog(t, func() {
for i := 0; i < 5; i++ {
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg("got", fmt.Sprintf("obs%d", i)), nil, nil, cfg)
}
})
lines := regionFilterLines(out)
if len(lines) != 1 {
t.Fatalf("want exactly 1 [region-filter] line for 5 drops of one region, got %d:\n%s", len(lines), out)
}
if !strings.Contains(lines[0], `"GOT"`) || !strings.Contains(lines[0], "observerIATAWhitelist") {
t.Errorf("warning must name the normalized region and the setting: %q", lines[0])
}
var n int
store.db.QueryRow("SELECT COUNT(*) FROM observers").Scan(&n)
if n != 0 {
t.Errorf("dropped messages must still be dropped, %d observers stored", n)
}
}

func TestIATAWhitelistDropWarnsPerDistinctRegion(t *testing.T) {
store := newTestStore(t)
cfg := &Config{ObserverIATAWhitelist: []string{"ARN"}}
out := captureLog(t, func() {
for _, r := range []string{"GOT", "MMX", " got ", "CPH", "mmx"} {
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg(r, "o"), nil, nil, cfg)
}
})
lines := regionFilterLines(out)
if len(lines) != 3 {
t.Fatalf("want one line each for GOT, MMX, CPH (codes are normalized), got %d:\n%s", len(lines), out)
}
}

func TestIATAWhitelistAllowedAndEmptyAreSilent(t *testing.T) {
store := newTestStore(t)
out := captureLog(t, func() {
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg("ARN", "a1"), nil, nil, &Config{ObserverIATAWhitelist: []string{"arn"}})
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg("GOT", "a2"), nil, nil, &Config{})
})
if lines := regionFilterLines(out); len(lines) != 0 {
t.Fatalf("allowed traffic and an empty whitelist must not warn:\n%s", out)
}
var n int
store.db.QueryRow("SELECT COUNT(*) FROM observers").Scan(&n)
if n != 2 {
t.Errorf("allowed traffic must still be stored, got %d observers", n)
}
}

// The region is a topic segment the publisher controls: it must not be able
// to forge log lines or make one line arbitrarily long.
func TestIATAWhitelistDropWarningIsOneEscapedBoundedLine(t *testing.T) {
store := newTestStore(t)
cfg := &Config{ObserverIATAWhitelist: []string{"ARN"}}
evil := "XX\nMQTT [src] fake line\r" + strings.Repeat("A", 10000)
out := captureLog(t, func() {
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg(evil, "o"), nil, nil, cfg)
})
lines := regionFilterLines(out)
if len(lines) != 1 {
t.Fatalf("want 1 warning line, got %d:\n%.500s", len(lines), out)
}
if strings.Count(out, "\n") != 1 {
t.Errorf("the warning spans %d lines; control characters must be escaped", strings.Count(out, "\n"))
}
if len(lines[0]) > 300 {
t.Errorf("warning line is %d bytes; the region must be truncated", len(lines[0]))
}
}

// Many distinct attacker-chosen regions: the warning keeps coming (never
// silent), but through a shared overflow path, so the number of lines stays
// bounded well below the number of regions.
func TestIATAWhitelistDropManyDistinctRegionsStaysBoundedAndVisible(t *testing.T) {
store := newTestStore(t)
cfg := &Config{ObserverIATAWhitelist: []string{"ARN"}}
const distinct = 3000
out := captureLog(t, func() {
for i := 0; i < distinct; i++ {
handleMessage(store, "src", MQTTSource{Name: "src"}, statusMsg(fmt.Sprintf("Z%05d", i), "o"), nil, nil, cfg)
}
})
lines := regionFilterLines(out)
if len(lines) == 0 {
t.Fatal("drops beyond the throttle bound must not be silent")
}
if len(lines) >= distinct/2 {
t.Fatalf("%d warning lines for %d regions: the per-region throttle must be bounded", len(lines), distinct)
}
overflow := 0
for _, l := range lines {
if strings.Contains(l, "throttle table full") {
overflow++
}
}
if overflow != 1 {
t.Errorf("want exactly one shared overflow warning within the interval, got %d", overflow)
}
}
133 changes: 133 additions & 0 deletions cmd/ingestor/iata_drop_warn.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
package main

import (
"log"
"strings"
"sync"
"time"
"unicode/utf8"
)

// Throttled warning for observerIATAWhitelist drops (#110).
//
// Dropping in silence fails in the dangerous direction: a legitimate region
// missing from the allow-list loses all its traffic with nothing in the log.
// Logging every drop is not an option either (a foreign feed is thousands of
// messages a day), so each region is logged once and re-logged at most every
// IATAWarnInterval while it keeps arriving; a strict log-once would scroll
// out of any log window and leave an active drop looking healthy.
//
// The region is a topic segment the publisher controls, so the per-region
// state is strictly bounded: at most iataWarnMaxTracked keys of at most
// iataWarnMaxKeyLen bytes. Past the cap the drop is still logged, on one
// shared overflow throttle. Entries older than the interval are reclaimed
// when the table is full, so the first codes ever seen cannot own it.

const (
// iataWarnMaxTracked bounds the per-region throttle table. Far above any
// real deployment's region count, small enough to be harmless when a
// hostile publisher sends a new region per message.
iataWarnMaxTracked = 512
// iataWarnMaxKeyLen caps the stored and logged region code. IATA codes
// are three letters; longer segments are truncated (and may share a key).
iataWarnMaxKeyLen = 32
// defaultIATAWarnIntervalSec is the default re-log interval (6h).
defaultIATAWarnIntervalSec = 6 * 60 * 60
)

// iataDropThrottle is the per-Config throttle state. The zero value is ready.
type iataDropThrottle struct {
mu sync.Mutex
last map[string]time.Time
overflowLast time.Time
// oldest is a lower bound on the oldest time in last: a sweep can free
// nothing until it has expired, so a full table of fresh entries costs
// O(1) per drop, not a full scan.
oldest time.Time
sweeps int // full-table sweeps, for tests
}

// IATAWarnInterval returns how often a dropped region is re-logged.
func (c *Config) IATAWarnInterval() time.Duration {
if c == nil || c.IATAWarnIntervalSec <= 0 {
return defaultIATAWarnIntervalSec * time.Second
}
return time.Duration(c.IATAWarnIntervalSec) * time.Second
}

// normalizeIATAForWarn is the key and the logged form of a region segment:
// trimmed, upper-cased and cut to iataWarnMaxKeyLen bytes on a rune boundary.
func normalizeIATAForWarn(iata string) string {
code := strings.ToUpper(strings.TrimSpace(iata))
if len(code) <= iataWarnMaxKeyLen {
return code
}
cut := iataWarnMaxKeyLen
for cut > 0 && !utf8.RuneStart(code[cut]) {
cut--
}
return code[:cut]
}

// shouldWarn reports whether a drop of code (already normalized) should be
// logged at now, and whether it goes through the shared overflow throttle.
func (t *iataDropThrottle) shouldWarn(code string, now time.Time, interval time.Duration) (warn, overflow bool) {
t.mu.Lock()
defer t.mu.Unlock()
if last, ok := t.last[code]; ok {
if now.Sub(last) < interval {
return false, false
}
t.last[code] = now
return true, false
}
if t.last == nil {
t.last = make(map[string]time.Time)
}
if len(t.last) >= iataWarnMaxTracked && now.Sub(t.oldest) >= interval {
t.sweeps++
oldest := now
for k, ts := range t.last {
if now.Sub(ts) >= interval {
delete(t.last, k)
} else if ts.Before(oldest) {
oldest = ts
}
}
t.oldest = oldest
}
if len(t.last) >= iataWarnMaxTracked {
if !t.overflowLast.IsZero() && now.Sub(t.overflowLast) < interval {
return false, true
}
t.overflowLast = now
return true, true
}
if len(t.last) == 0 || now.Before(t.oldest) {
t.oldest = now
}
t.last[code] = now
return true, false
}

// warnIATADrop logs a throttled warning for a message dropped by
// observerIATAWhitelist. The region is quoted (%q), so control characters in
// the publisher-controlled segment cannot break or forge log lines.
func (c *Config) warnIATADrop(tag, iata string, now time.Time) {
if c == nil {
return
}
code := normalizeIATAForWarn(iata)
interval := c.IATAWarnInterval()
warn, overflow := c.iataDropWarn.shouldWarn(code, now, interval)
if !warn {
return
}
if overflow {
log.Printf("MQTT [%s] [region-filter] dropping region %q: not in observerIATAWhitelist; throttle table full (%d regions), further untracked regions suppressed for %s",
tag, code, iataWarnMaxTracked, interval)
return
}
log.Printf("MQTT [%s] [region-filter] dropping region %q: not in observerIATAWhitelist; further messages from this region suppressed for %s",
tag, code, interval)
}
Loading
Loading