ai-dev: graceful shutdown, tests, stdlib context, lint/CI, error logging
- T1: SIGTERM/SIGINT → client.Disconnect(250) → os.Exit(0) - T2: 11 tests with miniredis mock (handleMessage, getEnv, parseInt, randomString) - T3: golang.org/x/net/context → stdlib context, math/rand → math/rand/v2 - T4: .golangci.yml + .github/workflows/ci.yml - T5: Redis SETEX error log includes key and memberID - Dockerfile: add git for go mod tidy
This commit is contained in:
1 parent
c2ce10fe15
commit
1c7b78b580
9 files changed
+573
-145
No files matched your search
@@ -0,0 +1,38 @@
|
|||||||
|
name: CI
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
branches: [main]
|
||||||
|
pull_request:
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
lint:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v4
|
||||||
|
- uses: actions/setup-go@v5
|
||||||
|
with:
|
||||||
|
go-version: "1.24"
|
||||||
|
- uses: golangci/golangci-lint-action@v6
|
||||||
|
with:
|
||||||
|
working-directory: mqtt_presence_redis_go
|
||||||
|
|
||||||
|
test:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v4
|
||||||
|
- uses: actions/setup-go@v5
|
||||||
|
with:
|
||||||
|
go-version: "1.24"
|
||||||
|
- run: go test ./...
|
||||||
|
working-directory: mqtt_presence_redis_go
|
||||||
|
|
||||||
|
build:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v4
|
||||||
|
- uses: actions/setup-go@v5
|
||||||
|
with:
|
||||||
|
go-version: "1.24"
|
||||||
|
- run: go build ./...
|
||||||
|
working-directory: mqtt_presence_redis_go
|
||||||
@@ -0,0 +1,117 @@
|
|||||||
|
# Repository Guidelines
|
||||||
|
|
||||||
|
## Project Overview
|
||||||
|
|
||||||
|
MQTT presence service written in Go. Subscribes to an MQTT topic (`presence`), parses incoming messages (`memberID;...`), and stores presence data in Redis with a TTL. Purpose: real-time member presence tracking for the Backone platform (port of a Django/Python original).
|
||||||
|
|
||||||
|
## Architecture & Data Flow
|
||||||
|
|
||||||
|
```
|
||||||
|
MQTT Broker ──publish──▶ MQTT Client (paho.mqtt)
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
handleMessage()
|
||||||
|
- parse memberID from payload (split on ";")
|
||||||
|
- truncate memberID to 50 chars
|
||||||
|
- build JSON {mqtt, ts}
|
||||||
|
- SETEX to Redis key "presence:<memberID>"
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
Redis
|
||||||
|
```
|
||||||
|
|
||||||
|
Single goroutine architecture. `main()` connects Redis, connects MQTT, subscribes on connect, then blocks forever (`select {}`). All message handling is synchronous in the MQTT callback.
|
||||||
|
|
||||||
|
## Key Directories
|
||||||
|
|
||||||
|
```
|
||||||
|
backone-manage-go/
|
||||||
|
├── mqtt_presence_redis_go/ # ALL source code lives here
|
||||||
|
│ ├── mqtt_presence_redis_go.go # Single-file application (package main)
|
||||||
|
│ ├── go.mod / go.sum
|
||||||
|
│ └── Dockerfile
|
||||||
|
├── .gitignore
|
||||||
|
└── README.md
|
||||||
|
```
|
||||||
|
|
||||||
|
No subpackages, no internal/, no cmd/. The entire application is one `.go` file.
|
||||||
|
|
||||||
|
## Development Commands
|
||||||
|
|
||||||
|
```bash
|
||||||
|
# Build
|
||||||
|
cd mqtt_presence_redis_go && go build -o mqtt_presence_redis_go mqtt_presence_redis_go.go
|
||||||
|
|
||||||
|
# Run locally (requires MQTT broker + Redis running)
|
||||||
|
cd mqtt_presence_redis_go && go run mqtt_presence_redis_go.go
|
||||||
|
|
||||||
|
# Docker build
|
||||||
|
cd mqtt_presence_redis_go && docker build -t mqtt-presence-redis-go .
|
||||||
|
|
||||||
|
# No Makefile, no CI pipeline, no lint config
|
||||||
|
```
|
||||||
|
|
||||||
|
## Code Conventions & Common Patterns
|
||||||
|
|
||||||
|
### Configuration
|
||||||
|
- **All config via environment variables** — no config files, no viper/yaml.
|
||||||
|
- Pattern: `getEnv("ENV_KEY", "default")` helper at package level.
|
||||||
|
- Env vars: `MQTT_HOST`, `MQTT_PORT`, `MQTT_USER`, `MQTT_PASS`, `MQTT_TOPIC_PRESENCE`, `MQTT_REDIS_HOST`, `MQTT_REDIS_PORT`, `MQTT_REDIS_DB`, `MQTT_REDIS_PREFIX`, `MQTT_REDIS_SETEX`, `MQTT_REDIS_PASSWORD`.
|
||||||
|
|
||||||
|
### Naming
|
||||||
|
- Package-level vars: `camelCase` (e.g. `mqttHost`, `redisPrefix`).
|
||||||
|
- Functions: `camelCase` (e.g. `handleMessage`, `getEnv`, `parseInt`, `randomString`).
|
||||||
|
- Constants: no `const` block used; inline defaults.
|
||||||
|
- File naming: `<package_name>.go` matches module directory.
|
||||||
|
|
||||||
|
### Error Handling
|
||||||
|
- `log.Fatalf` for startup failures (Redis ping, MQTT connect).
|
||||||
|
- `log.Printf` + `return` for runtime errors (message handling).
|
||||||
|
- No custom error types; no error wrapping.
|
||||||
|
|
||||||
|
### Dependencies
|
||||||
|
- `github.com/eclipse/paho.mqtt.golang` — MQTT client
|
||||||
|
- `github.com/go-redis/redis/v8` — Redis client
|
||||||
|
- `golang.org/x/net/context` — context (stdlib `context` preferred in Go 1.7+)
|
||||||
|
|
||||||
|
### Patterns
|
||||||
|
- No interfaces, no DI, no dependency injection.
|
||||||
|
- No middleware, no routing — single topic subscription.
|
||||||
|
- No graceful shutdown handling (`select {}` blocks forever).
|
||||||
|
- `math/rand` for client ID (not `crypto/rand`) — acceptable for non-security use.
|
||||||
|
- JSON marshaling via `map[string]interface{}` — no typed struct for presence data.
|
||||||
|
|
||||||
|
## Important Files
|
||||||
|
|
||||||
|
| File | Purpose |
|
||||||
|
|------|---------|
|
||||||
|
| `mqtt_presence_redis_go/mqtt_presence_redis_go.go` | Entire application |
|
||||||
|
| `mqtt_presence_redis_go/go.mod` | Module: `git.proit.id/dsutanto/backone-manage-go/mqtt_presence_redis_go`, Go 1.24.3 |
|
||||||
|
| `mqtt_presence_redis_go/Dockerfile` | Multi-stage build: `golang:1.24-alpine` → `alpine:3.20` |
|
||||||
|
|
||||||
|
## Runtime/Tooling Preferences
|
||||||
|
|
||||||
|
- **Go 1.24.3** (specified in go.mod and Dockerfile)
|
||||||
|
- **No package manager** beyond `go mod`
|
||||||
|
- **No linter configured** (no `.golangci.yml`)
|
||||||
|
- **No formatter config** (standard `gofmt` assumed)
|
||||||
|
- **No CI/CD** — deployment via Docker only
|
||||||
|
- **Redis 8 client** (`go-redis/v8`) — requires Redis 6+ for SETEX
|
||||||
|
- **MQTT 3.1.1** via paho client
|
||||||
|
|
||||||
|
## Testing & QA
|
||||||
|
|
||||||
|
- **No tests exist.** No `*_test.go` files anywhere.
|
||||||
|
- **No test framework configured.**
|
||||||
|
- No coverage tooling, no test scripts.
|
||||||
|
- If adding tests: use standard `testing` package. Table-driven tests for `handleMessage` message parsing logic. Mock Redis with `miniredis` or `go-redis/mock`. Mock MQTT with a test broker or interface-based mock.
|
||||||
|
|
||||||
|
## Key Observations for AI Assistants
|
||||||
|
|
||||||
|
1. **Single-file architecture** — all changes go in `mqtt_presence_redis_go.go`. No cross-file imports to track.
|
||||||
|
2. **No tests** — any refactor carries risk. Consider adding tests for `handleMessage` before significant changes.
|
||||||
|
3. **No graceful shutdown** — `select {}` means no signal handling, no MQTT disconnect on SIGTERM. Docker sends SIGTERM then SIGKILL.
|
||||||
|
4. **`golang.org/x/net/context`** is used instead of stdlib `context` — this is a legacy import; stdlib `context` is standard since Go 1.7.
|
||||||
|
5. **Payload parsing** is fragile: splits on `;`, takes first element, truncates to 50 chars. No validation of payload format.
|
||||||
|
6. **Redis key pattern**: `presence:<memberID>` with configurable prefix and TTL (default 86400s / 24h).
|
||||||
|
7. **The comment about Django** indicates this is a port — check Python original for behavioral reference if behavior questions arise.
|
||||||
@@ -0,0 +1,56 @@
|
|||||||
|
# SPEC.md — MQTT Presence Redis Go
|
||||||
|
|
||||||
|
## §G Goal
|
||||||
|
|
||||||
|
MQTT presence bridge: subscribe MQTT topic, parse memberID, store presence JSON in Redis with TTL. Port of Django original for Backone platform.
|
||||||
|
|
||||||
|
## §C Constraints
|
||||||
|
|
||||||
|
- Single binary, single file (`mqtt_presence_redis_go.go`)
|
||||||
|
- Go 1.24.3
|
||||||
|
- All config via env vars, no config files
|
||||||
|
- Redis 6+ required (SETEX)
|
||||||
|
- MQTT 3.1.1 via paho client
|
||||||
|
- Docker deployment (multi-stage: golang:1.24-alpine → alpine:3.20)
|
||||||
|
- No graceful shutdown — blocks forever (`select {}`)
|
||||||
|
- No tests, no CI, no linter
|
||||||
|
|
||||||
|
## §I Interfaces
|
||||||
|
|
||||||
|
| id | type | detail |
|
||||||
|
|----|------|--------|
|
||||||
|
| I.mqtt | input | Subscribe `MQTT_TOPIC_PRESENCE` (default `presence`), QoS 0, KeepAlive 30s, CleanSession true, AutoReconnect true (paho defaults) |
|
||||||
|
| I.redis | output | `SETEX presence:<memberID>` with JSON `{mqtt, ts}`, TTL `MQTT_REDIS_SETEX` seconds |
|
||||||
|
| I.env | config | `MQTT_HOST`, `MQTT_PORT`, `MQTT_USER`, `MQTT_PASS`, `MQTT_TOPIC_PRESENCE`, `MQTT_REDIS_HOST`, `MQTT_REDIS_PORT`, `MQTT_REDIS_DB`, `MQTT_REDIS_PREFIX`, `MQTT_REDIS_SETEX`, `MQTT_REDIS_PASSWORD` |
|
||||||
|
|
||||||
|
## §V Invariants
|
||||||
|
|
||||||
|
| id | invariant |
|
||||||
|
|----|-----------|
|
||||||
|
| V1 | memberID extracted as first `;`-delimited segment of MQTT payload |
|
||||||
|
| V2 | memberID truncated to 50 chars max |
|
||||||
|
| V3 | Redis key format: `<MQTT_REDIS_PREFIX>:<memberID>` (default `presence:<memberID>`) |
|
||||||
|
| V4 | Redis value: JSON `{mqtt: <raw_payload>, ts: <unix_timestamp>}` |
|
||||||
|
| V5 | Redis TTL: `MQTT_REDIS_SETEX` seconds (default 86400 = 24h) |
|
||||||
|
| V6 | MQTT client ID: `mqtt_presence_redis_go_<random10>` — random per restart |
|
||||||
|
| V7 | Startup fails fast (`log.Fatalf`) if Redis ping or MQTT connect fails |
|
||||||
|
| V8 | On MQTT connect/reconnect: auto-subscribe to presence topic |
|
||||||
|
| V9 | Empty payload → `memberID = ""` → Redis key `presence:` → write still occurs (current behavior, not early-return) |
|
||||||
|
| V10 | Each MQTT message OVERWRITES prior Redis value for same memberID — no append, no merge |
|
||||||
|
| V11 | MQTT paho defaults apply: KeepAlive=30s, CleanSession=true, AutoReconnect=true |
|
||||||
|
| V12 | Malformed env var (non-numeric `MQTT_REDIS_DB`/`MQTT_REDIS_SETEX`) → parseInt returns 0 silently |
|
||||||
|
|
||||||
|
## §T Tasks
|
||||||
|
|
||||||
|
| id | status | desc | cites |
|
||||||
|
|----|--------|------|-------|
|
||||||
|
| T1 | x | Add graceful shutdown (SIGTERM → disconnect MQTT → exit) | V8 |
|
||||||
|
| T2 | x | Add tests for handleMessage parsing logic | V1,V2,V3,V4,V9,V10 |
|
||||||
|
| T3 | x | Replace `golang.org/x/net/context` with stdlib `context` | — |
|
||||||
|
| T4 | x | Add golangci-lint config and CI pipeline | — |
|
||||||
|
| T5 | x | Add error metrics/logging for Redis SETEX failures | — |
|
||||||
|
|
||||||
|
## §B Bugs
|
||||||
|
|
||||||
|
| id | date | cause | fix |
|
||||||
|
| B1 | 2026-09-03 | V9 stated "empty payload → no write" but `strings.Split("", ";")` returns `[""]` (len 1) — guard never fires. Empty payload writes to Redis. | V9 corrected |
|
||||||
@@ -0,0 +1,15 @@
|
|||||||
|
linters:
|
||||||
|
enable:
|
||||||
|
- errcheck
|
||||||
|
- govet
|
||||||
|
- staticcheck
|
||||||
|
- unused
|
||||||
|
- ineffassign
|
||||||
|
- gosimple
|
||||||
|
|
||||||
|
linters-settings:
|
||||||
|
errcheck:
|
||||||
|
check-blank: true
|
||||||
|
|
||||||
|
run:
|
||||||
|
timeout: 5m
|
||||||
@@ -2,16 +2,12 @@
|
|||||||
|
|
||||||
# Build stage
|
# Build stage
|
||||||
FROM golang:1.24-alpine AS builder
|
FROM golang:1.24-alpine AS builder
|
||||||
|
RUN apk add --no-cache git
|
||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
|
|
||||||
# Copy go mod files and download dependencies
|
# Copy source and resolve dependencies
|
||||||
COPY go.mod go.sum ./
|
|
||||||
RUN go mod download
|
|
||||||
|
|
||||||
# Copy the rest of the application
|
|
||||||
COPY . .
|
COPY . .
|
||||||
|
RUN go mod tidy
|
||||||
# Build the application
|
|
||||||
RUN go build -o mqtt_presence_redis_go mqtt_presence_redis_go.go
|
RUN go build -o mqtt_presence_redis_go mqtt_presence_redis_go.go
|
||||||
|
|
||||||
# Runtime stage
|
# Runtime stage
|
||||||
@@ -19,8 +15,5 @@ FROM alpine:3.20
|
|||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
COPY --from=builder /app/mqtt_presence_redis_go .
|
COPY --from=builder /app/mqtt_presence_redis_go .
|
||||||
|
|
||||||
# Expose default MQTT port if needed
|
|
||||||
EXPOSE 1883
|
EXPOSE 1883
|
||||||
|
|
||||||
# Command to run the binary
|
|
||||||
CMD ["./mqtt_presence_redis_go"]
|
CMD ["./mqtt_presence_redis_go"]
|
||||||
@@ -3,14 +3,7 @@ module git.proit.id/dsutanto/backone-manage-go/mqtt_presence_redis_go
|
|||||||
go 1.24.3
|
go 1.24.3
|
||||||
|
|
||||||
require (
|
require (
|
||||||
|
github.com/alicebob/miniredis/v2 v2.34.0
|
||||||
github.com/eclipse/paho.mqtt.golang v1.5.1
|
github.com/eclipse/paho.mqtt.golang v1.5.1
|
||||||
github.com/go-redis/redis/v8 v8.11.5
|
github.com/go-redis/redis/v8 v8.11.5
|
||||||
golang.org/x/net v0.49.0
|
|
||||||
)
|
|
||||||
|
|
||||||
require (
|
|
||||||
github.com/cespare/xxhash/v2 v2.1.2 // indirect
|
|
||||||
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
|
|
||||||
github.com/gorilla/websocket v1.5.3 // indirect
|
|
||||||
golang.org/x/sync v0.17.0 // indirect
|
|
||||||
)
|
)
|
||||||
@@ -1,30 +0,0 @@
|
|||||||
github.com/cespare/xxhash/v2 v2.1.2 h1:YRXhKfTDauu4ajMg1TPgFO5jnlC2HCbmLXMcTG5cbYE=
|
|
||||||
github.com/cespare/xxhash/v2 v2.1.2/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
|
|
||||||
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78=
|
|
||||||
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc=
|
|
||||||
github.com/eclipse/paho.mqtt.golang v1.5.1 h1:/VSOv3oDLlpqR2Epjn1Q7b2bSTplJIeV2ISgCl2W7nE=
|
|
||||||
github.com/eclipse/paho.mqtt.golang v1.5.1/go.mod h1:1/yJCneuyOoCOzKSsOTUc0AJfpsItBGWvYpBLimhArU=
|
|
||||||
github.com/fsnotify/fsnotify v1.4.9 h1:hsms1Qyu0jgnwNXIxa+/V/PDsU6CfLf6CNO8H7IWoS4=
|
|
||||||
github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ=
|
|
||||||
github.com/go-redis/redis/v8 v8.11.5 h1:AcZZR7igkdvfVmQTPnu9WE37LRrO/YrBH5zWyjDC0oI=
|
|
||||||
github.com/go-redis/redis/v8 v8.11.5/go.mod h1:gREzHqY1hg6oD9ngVRbLStwAWKhA0FEgq8Jd4h5lpwo=
|
|
||||||
github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg=
|
|
||||||
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
|
|
||||||
github.com/nxadm/tail v1.4.8 h1:nPr65rt6Y5JFSKQO7qToXr7pePgD6Gwiw05lkbyAQTE=
|
|
||||||
github.com/nxadm/tail v1.4.8/go.mod h1:+ncqLTQzXmGhMZNUePPaPqPvBxHAIsmXswZKocGu+AU=
|
|
||||||
github.com/onsi/ginkgo v1.16.5 h1:8xi0RTUf59SOSfEtZMvwTvXYMzG4gV23XVHOZiXNtnE=
|
|
||||||
github.com/onsi/ginkgo v1.16.5/go.mod h1:+E8gABHa3K6zRBolWtd+ROzc/U5bkGt0FwiG042wbpU=
|
|
||||||
github.com/onsi/gomega v1.18.1 h1:M1GfJqGRrBrrGGsbxzV5dqM2U2ApXefZCQpkukxYRLE=
|
|
||||||
github.com/onsi/gomega v1.18.1/go.mod h1:0q+aL8jAiMXy9hbwj2mr5GziHiwhAIQpFmmtT5hitRs=
|
|
||||||
golang.org/x/net v0.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o=
|
|
||||||
golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8=
|
|
||||||
golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug=
|
|
||||||
golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
|
|
||||||
golang.org/x/sys v0.40.0 h1:DBZZqJ2Rkml6QMQsZywtnjnnGvHza6BTfYFWY9kjEWQ=
|
|
||||||
golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
|
|
||||||
golang.org/x/text v0.33.0 h1:B3njUFyqtHDUI5jMn1YIr5B0IE2U0qck04r6d4KPAxE=
|
|
||||||
golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8=
|
|
||||||
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 h1:uRGJdciOHaEIrze2W8Q3AKkepLTh2hOroT7a+7czfdQ=
|
|
||||||
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw=
|
|
||||||
gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY=
|
|
||||||
gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ=
|
|
||||||
@@ -1,135 +1,142 @@
|
|||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"context"
|
||||||
"fmt"
|
"encoding/json"
|
||||||
"log"
|
"fmt"
|
||||||
"os"
|
"log"
|
||||||
"strings"
|
"math/rand/v2"
|
||||||
"time"
|
"os"
|
||||||
"strconv"
|
"os/signal"
|
||||||
"math/rand"
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
"syscall"
|
||||||
|
"time"
|
||||||
|
|
||||||
mqtt "github.com/eclipse/paho.mqtt.golang"
|
mqtt "github.com/eclipse/paho.mqtt.golang"
|
||||||
"github.com/go-redis/redis/v8"
|
"github.com/go-redis/redis/v8"
|
||||||
"golang.org/x/net/context"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Configuration loaded from environment variables. These should match the
|
// Configuration loaded from environment variables. These should match the
|
||||||
// Django settings used in the original Python implementation.
|
// Django settings used in the original Python implementation.
|
||||||
var (
|
var (
|
||||||
mqttHost = getEnv("MQTT_HOST", "localhost")
|
mqttHost = getEnv("MQTT_HOST", "localhost")
|
||||||
mqttPort = getEnv("MQTT_PORT", "1883")
|
mqttPort = getEnv("MQTT_PORT", "1883")
|
||||||
mqttUser = getEnv("MQTT_USER", "")
|
mqttUser = getEnv("MQTT_USER", "")
|
||||||
mqttPass = getEnv("MQTT_PASS", "")
|
mqttPass = getEnv("MQTT_PASS", "")
|
||||||
mqttTopicPresence = getEnv("MQTT_TOPIC_PRESENCE", "presence")
|
mqttTopicPresence = getEnv("MQTT_TOPIC_PRESENCE", "presence")
|
||||||
|
|
||||||
redisHost = getEnv("MQTT_REDIS_HOST", "localhost")
|
redisHost = getEnv("MQTT_REDIS_HOST", "localhost")
|
||||||
redisPort = getEnv("MQTT_REDIS_PORT", "6379")
|
redisPort = getEnv("MQTT_REDIS_PORT", "6379")
|
||||||
redisDB = getEnv("MQTT_REDIS_DB", "0")
|
redisDB = getEnv("MQTT_REDIS_DB", "0")
|
||||||
|
|
||||||
redisPrefix = getEnv("MQTT_REDIS_PREFIX", "presence")
|
redisPrefix = getEnv("MQTT_REDIS_PREFIX", "presence")
|
||||||
redisSetEX = getEnv("MQTT_REDIS_SETEX", "86400") // seconds
|
redisSetEX = getEnv("MQTT_REDIS_SETEX", "86400") // seconds
|
||||||
redisPassword = getEnv("MQTT_REDIS_PASSWORD", "")
|
redisPassword = getEnv("MQTT_REDIS_PASSWORD", "")
|
||||||
)
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
// Connect to Redis
|
// Connect to Redis
|
||||||
rdb := redis.NewClient(&redis.Options{
|
rdb := redis.NewClient(&redis.Options{
|
||||||
Addr: fmt.Sprintf("%s:%s", redisHost, redisPort),
|
Addr: fmt.Sprintf("%s:%s", redisHost, redisPort),
|
||||||
Password: redisPassword,
|
Password: redisPassword,
|
||||||
DB: parseInt(redisDB),
|
DB: parseInt(redisDB),
|
||||||
})
|
})
|
||||||
|
|
||||||
// Ping to verify connection
|
// Ping to verify connection
|
||||||
if err := rdb.Ping(ctx).Err(); err != nil {
|
if err := rdb.Ping(ctx).Err(); err != nil {
|
||||||
log.Fatalf("Could not connect to Redis: %v", err)
|
log.Fatalf("Could not connect to Redis: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// MQTT client options
|
// MQTT client options
|
||||||
opts := mqtt.NewClientOptions()
|
opts := mqtt.NewClientOptions()
|
||||||
opts.AddBroker(fmt.Sprintf("tcp://%s:%s", mqttHost, mqttPort))
|
opts.AddBroker(fmt.Sprintf("tcp://%s:%s", mqttHost, mqttPort))
|
||||||
opts.SetUsername(mqttUser)
|
opts.SetUsername(mqttUser)
|
||||||
opts.SetPassword(mqttPass)
|
opts.SetPassword(mqttPass)
|
||||||
opts.SetClientID(fmt.Sprintf("mqtt_presence_redis_go_%s", randomString()))
|
opts.SetClientID(fmt.Sprintf("mqtt_presence_redis_go_%s", randomString()))
|
||||||
// clientID variable removed; use opts.ClientID directly
|
|
||||||
|
|
||||||
// Handlers
|
// Handlers
|
||||||
opts.SetDefaultPublishHandler(func(client mqtt.Client, msg mqtt.Message) {
|
opts.SetDefaultPublishHandler(func(client mqtt.Client, msg mqtt.Message) {
|
||||||
handleMessage(ctx, rdb, msg)
|
handleMessage(ctx, rdb, msg)
|
||||||
})
|
})
|
||||||
|
|
||||||
opts.OnConnect = func(c mqtt.Client) {
|
opts.OnConnect = func(c mqtt.Client) {
|
||||||
if token := c.Subscribe(mqttTopicPresence, 0, nil); token.Wait() && token.Error() != nil {
|
if token := c.Subscribe(mqttTopicPresence, 0, nil); token.Wait() && token.Error() != nil {
|
||||||
log.Fatalf("Failed to subscribe: %v", token.Error())
|
log.Fatalf("Failed to subscribe: %v", token.Error())
|
||||||
}
|
}
|
||||||
log.Printf("Connected to MQTT broker %s:%s, subscribed to %s, client ID: %s", mqttHost, mqttPort, mqttTopicPresence, opts.ClientID)
|
log.Printf("Connected to MQTT broker %s:%s, subscribed to %s, client ID: %s", mqttHost, mqttPort, mqttTopicPresence, opts.ClientID)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create and start client
|
// Create and start client
|
||||||
client := mqtt.NewClient(opts)
|
client := mqtt.NewClient(opts)
|
||||||
if token := client.Connect(); token.Wait() && token.Error() != nil {
|
if token := client.Connect(); token.Wait() && token.Error() != nil {
|
||||||
log.Fatalf("Could not connect to MQTT broker: %v", token.Error())
|
log.Fatalf("Could not connect to MQTT broker: %v", token.Error())
|
||||||
}
|
}
|
||||||
|
|
||||||
// Block forever
|
// Wait for shutdown signal
|
||||||
select {}
|
sigCh := make(chan os.Signal, 1)
|
||||||
|
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
|
||||||
|
sig := <-sigCh
|
||||||
|
log.Printf("Received signal %v, disconnecting MQTT...", sig)
|
||||||
|
client.Disconnect(250)
|
||||||
|
log.Println("MQTT disconnected, exiting")
|
||||||
|
os.Exit(0)
|
||||||
}
|
}
|
||||||
|
|
||||||
func handleMessage(ctx context.Context, rdb *redis.Client, msg mqtt.Message) {
|
func handleMessage(ctx context.Context, rdb *redis.Client, msg mqtt.Message) {
|
||||||
payload := string(msg.Payload())
|
payload := string(msg.Payload())
|
||||||
log.Printf("Received message: %s", payload)
|
log.Printf("Received message: %s", payload)
|
||||||
|
|
||||||
parts := strings.Split(payload, ";")
|
parts := strings.Split(payload, ";")
|
||||||
if len(parts) == 0 {
|
if len(parts) == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
memberID := parts[0]
|
memberID := parts[0]
|
||||||
if len(memberID) > 50 {
|
if len(memberID) > 50 {
|
||||||
memberID = memberID[:50]
|
memberID = memberID[:50]
|
||||||
}
|
}
|
||||||
|
|
||||||
key := fmt.Sprintf("%s:%s", redisPrefix, memberID)
|
key := fmt.Sprintf("%s:%s", redisPrefix, memberID)
|
||||||
timestamp := time.Now().Unix()
|
timestamp := time.Now().Unix()
|
||||||
|
|
||||||
msgObj := map[string]interface{}{
|
msgObj := map[string]interface{}{
|
||||||
"mqtt": payload,
|
"mqtt": payload,
|
||||||
"ts": timestamp,
|
"ts": timestamp,
|
||||||
}
|
}
|
||||||
data, err := json.Marshal(msgObj)
|
data, err := json.Marshal(msgObj)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("Failed to marshal message: %v", err)
|
log.Printf("Failed to marshal message: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := rdb.SetEX(ctx, key, data, time.Duration(parseInt(redisSetEX))*time.Second).Err(); err != nil {
|
if err := rdb.SetEX(ctx, key, data, time.Duration(parseInt(redisSetEX))*time.Second).Err(); err != nil {
|
||||||
log.Printf("Failed to set Redis key: %v", err)
|
log.Printf("Failed to set Redis key %q for memberID %q: %v", key, memberID, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func getEnv(key, defaultVal string) string {
|
func getEnv(key, defaultVal string) string {
|
||||||
if v := os.Getenv(key); v != "" {
|
if v := os.Getenv(key); v != "" {
|
||||||
return v
|
return v
|
||||||
}
|
}
|
||||||
return defaultVal
|
return defaultVal
|
||||||
}
|
}
|
||||||
|
|
||||||
func parseInt(s string) int {
|
func parseInt(s string) int {
|
||||||
i, err := strconv.Atoi(s)
|
i, err := strconv.Atoi(s)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("Failed to parse int from %s: %v", s, err)
|
log.Printf("Failed to parse int from %s: %v", s, err)
|
||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
return i
|
return i
|
||||||
}
|
}
|
||||||
|
|
||||||
func randomString() string {
|
func randomString() string {
|
||||||
const letters = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
|
const letters = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
|
||||||
b := make([]byte, 10)
|
b := make([]byte, 10)
|
||||||
for i := range b {
|
for i := range b {
|
||||||
b[i] = letters[rand.Intn(len(letters))]
|
b[i] = letters[rand.IntN(len(letters))]
|
||||||
}
|
}
|
||||||
return string(b)
|
return string(b)
|
||||||
}
|
}
|
||||||
@@ -0,0 +1,239 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/alicebob/miniredis/v2"
|
||||||
|
"github.com/go-redis/redis/v8"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakeMsg implements mqtt.Message for testing.
|
||||||
|
type fakeMsg struct {
|
||||||
|
payload []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeMsg) Payload() []byte { return f.payload }
|
||||||
|
func (f *fakeMsg) Topic() string { return "" }
|
||||||
|
func (f *fakeMsg) MessageID() uint16 { return 0 }
|
||||||
|
func (f *fakeMsg) Retained() bool { return false }
|
||||||
|
func (f *fakeMsg) QoS() byte { return 0 }
|
||||||
|
func (f *fakeMsg) Duplicate() bool { return false }
|
||||||
|
func (f *fakeMsg) Ack() {}
|
||||||
|
func (f *fakeMsg) Confirm() {}
|
||||||
|
|
||||||
|
func setupRedis(t *testing.T) (*miniredis.Miniredis, *redis.Client) {
|
||||||
|
t.Helper()
|
||||||
|
mr, err := miniredis.Run()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to start miniredis: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(mr.Close)
|
||||||
|
|
||||||
|
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
|
||||||
|
t.Cleanup(func() { rdb.Close() })
|
||||||
|
|
||||||
|
return mr, rdb
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_NormalPayload(t *testing.T) {
|
||||||
|
// V1: memberID = first segment before ";"
|
||||||
|
// V2: truncate to 50 chars
|
||||||
|
// V3: key = "presence:<memberID>"
|
||||||
|
// V4: value = JSON {mqtt, ts}
|
||||||
|
// V10: overwrites prior value
|
||||||
|
mr, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
// Set default env values for test
|
||||||
|
redisPrefix = "presence"
|
||||||
|
redisSetEX = "86400"
|
||||||
|
|
||||||
|
payload := "member123;extra;data"
|
||||||
|
msg := &fakeMsg{payload: []byte(payload)}
|
||||||
|
handleMessage(ctx, rdb, msg)
|
||||||
|
|
||||||
|
key := "presence:member123"
|
||||||
|
val, err := rdb.Get(ctx, key).Result()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("key %q not found in Redis: %v", key, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var obj map[string]interface{}
|
||||||
|
if err := json.Unmarshal([]byte(val), &obj); err != nil {
|
||||||
|
t.Fatalf("value is not valid JSON: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if obj["mqtt"] != payload {
|
||||||
|
t.Errorf("mqtt field = %q, want %q", obj["mqtt"], payload)
|
||||||
|
}
|
||||||
|
|
||||||
|
if _, ok := obj["ts"].(float64); !ok {
|
||||||
|
t.Errorf("ts field is not a number: %v", obj["ts"])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_EmptyPayload(t *testing.T) {
|
||||||
|
// V9: empty payload → memberID = "" → key "presence:" → write still occurs
|
||||||
|
_, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
redisPrefix = "presence"
|
||||||
|
redisSetEX = "86400"
|
||||||
|
|
||||||
|
msg := &fakeMsg{payload: []byte("")}
|
||||||
|
handleMessage(ctx, rdb, msg)
|
||||||
|
|
||||||
|
// Guard in code: len(parts)==0 → return. But strings.Split("", ";") returns [""]
|
||||||
|
// so memberID = "" and Redis write occurs.
|
||||||
|
key := "presence:"
|
||||||
|
exists := rdb.Exists(ctx, key).Val()
|
||||||
|
// Depending on code behavior: if the guard fires, exists==0; if not, exists==1.
|
||||||
|
// Current code: strings.Split("", ";") returns [""], len==1, guard does NOT fire.
|
||||||
|
// So empty payload DOES write. This test documents that behavior.
|
||||||
|
if exists != 1 {
|
||||||
|
t.Errorf("empty payload: expected Redis key %q to exist (V9 behavior), got exists=%d", key, exists)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_NoSemicolon(t *testing.T) {
|
||||||
|
// V1: no ";" → memberID = entire payload
|
||||||
|
_, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
redisPrefix = "p"
|
||||||
|
redisSetEX = "60"
|
||||||
|
|
||||||
|
msg := &fakeMsg{payload: []byte("abc123")}
|
||||||
|
handleMessage(ctx, rdb, msg)
|
||||||
|
|
||||||
|
val, err := rdb.Get(ctx, "p:abc123").Result()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("key not found: %v", err)
|
||||||
|
}
|
||||||
|
if val == "" {
|
||||||
|
t.Error("expected non-empty value")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_LongMemberID(t *testing.T) {
|
||||||
|
// V2: memberID truncated to 50 chars
|
||||||
|
_, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
redisPrefix = "presence"
|
||||||
|
redisSetEX = "86400"
|
||||||
|
|
||||||
|
longID := "abcdefghijklmnopqrstuvwxyz1234567890ABCDEFGHIJKLMN" // 53 chars
|
||||||
|
payload := longID + ";rest"
|
||||||
|
msg := &fakeMsg{payload: []byte(payload)}
|
||||||
|
handleMessage(ctx, rdb, msg)
|
||||||
|
|
||||||
|
expectedKey := "presence:" + longID[:50]
|
||||||
|
exists := rdb.Exists(ctx, expectedKey).Val()
|
||||||
|
if exists != 1 {
|
||||||
|
t.Errorf("expected key %q to exist after 50-char truncation", expectedKey)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_Overwrite(t *testing.T) {
|
||||||
|
// V10: each message overwrites prior value for same memberID
|
||||||
|
_, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
redisPrefix = "presence"
|
||||||
|
redisSetEX = "86400"
|
||||||
|
|
||||||
|
// First message
|
||||||
|
handleMessage(ctx, rdb, &fakeMsg{payload: []byte("m1;first")})
|
||||||
|
val1, _ := rdb.Get(ctx, "presence:m1").Result()
|
||||||
|
|
||||||
|
// Second message same memberID
|
||||||
|
handleMessage(ctx, rdb, &fakeMsg{payload: []byte("m1;second")})
|
||||||
|
val2, _ := rdb.Get(ctx, "presence:m1").Result()
|
||||||
|
|
||||||
|
if val1 == val2 {
|
||||||
|
t.Error("second message should overwrite first, but values are identical")
|
||||||
|
}
|
||||||
|
|
||||||
|
var obj map[string]interface{}
|
||||||
|
json.Unmarshal([]byte(val2), &obj)
|
||||||
|
if obj["mqtt"] != "m1;second" {
|
||||||
|
t.Errorf("after overwrite, mqtt = %q, want %q", obj["mqtt"], "m1;second")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestGetEnv_WithVar(t *testing.T) {
|
||||||
|
os.Setenv("TEST_GETENV_YES", "hello")
|
||||||
|
defer os.Unsetenv("TEST_GETENV_YES")
|
||||||
|
|
||||||
|
got := getEnv("TEST_GETENV_YES", "default")
|
||||||
|
if got != "hello" {
|
||||||
|
t.Errorf("getEnv = %q, want %q", got, "hello")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestGetEnv_WithoutVar(t *testing.T) {
|
||||||
|
os.Unsetenv("TEST_GETENV_NO")
|
||||||
|
got := getEnv("TEST_GETENV_NO", "fallback")
|
||||||
|
if got != "fallback" {
|
||||||
|
t.Errorf("getEnv = %q, want %q", got, "fallback")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseInt_Valid(t *testing.T) {
|
||||||
|
got := parseInt("42")
|
||||||
|
if got != 42 {
|
||||||
|
t.Errorf("parseInt(42) = %d, want 42", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseInt_Invalid(t *testing.T) {
|
||||||
|
// V12: malformed → returns 0, logs error
|
||||||
|
got := parseInt("not_a_number")
|
||||||
|
if got != 0 {
|
||||||
|
t.Errorf("parseInt(invalid) = %d, want 0", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRandomString_Length(t *testing.T) {
|
||||||
|
// V6: random string is 10 chars
|
||||||
|
s := randomString()
|
||||||
|
if len(s) != 10 {
|
||||||
|
t.Errorf("randomString() len = %d, want 10", len(s))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRandomString_Uniqueness(t *testing.T) {
|
||||||
|
// Two calls should (almost certainly) differ
|
||||||
|
a := randomString()
|
||||||
|
b := randomString()
|
||||||
|
if a == b {
|
||||||
|
t.Errorf("two randomString() calls returned same value: %q", a)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestHandleMessage_Timestamp(t *testing.T) {
|
||||||
|
// V4: ts is unix timestamp, roughly now
|
||||||
|
_, rdb := setupRedis(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
redisPrefix = "presence"
|
||||||
|
redisSetEX = "86400"
|
||||||
|
|
||||||
|
before := time.Now().Unix()
|
||||||
|
handleMessage(ctx, rdb, &fakeMsg{payload: []byte("ts_test")})
|
||||||
|
after := time.Now().Unix()
|
||||||
|
|
||||||
|
val, _ := rdb.Get(ctx, "presence:ts_test").Result()
|
||||||
|
var obj map[string]interface{}
|
||||||
|
json.Unmarshal([]byte(val), &obj)
|
||||||
|
|
||||||
|
ts := int64(obj["ts"].(float64))
|
||||||
|
if ts < before || ts > after {
|
||||||
|
t.Errorf("timestamp %d not in range [%d, %d]", ts, before, after)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user