From 184c5453fb9961cc5813fe0dfe7bde630988b23d Mon Sep 17 00:00:00 2001 From: windyboy Date: Sat, 15 Nov 2025 12:08:26 +0800 Subject: [PATCH] =?UTF-8?q?=E2=9C=A8=20Expand=20Docker=20Compose=20develop?= =?UTF-8?q?ment=20stack=20to=20include=20observability=20tools:=20OpenTele?= =?UTF-8?q?metry=20Collector,=20Jaeger,=20Prometheus,=20and=20Grafana.=20U?= =?UTF-8?q?pdate=20README=20with=20new=20usage=20instructions=20for=20the?= =?UTF-8?q?=20observability=20stack=20and=20enhance=20Taskfile=20for=20str?= =?UTF-8?q?eamlined=20development=20commands.=20Introduce=20sample=20teleg?= =?UTF-8?q?ram=20generation=20utility=20and=20improve=20NATS=20configurati?= =?UTF-8?q?on=20for=20core=20and=20JetStream=20modes.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 27 ++-- Taskfile.yml | 63 ++++++-- cmd/main/main.go | 7 + cmd/seed-telegrams/main.go | 242 ++++++++++++++++++++++++++++ configs/config.dev.toml | 1 + configs/grafana-datasources.dev.yml | 16 ++ configs/otel-collector.dev.yaml | 6 +- configs/prometheus.dev.yml | 16 ++ docker-compose.dev.yml | 47 ++++++ docs/dev-guide.md | 133 +++++++++++++++ internal/infra/config/config.go | 25 ++- internal/infra/nats/consumer.go | 73 ++++++++- internal/repository/telegrams.ddl | 6 +- 13 files changed, 619 insertions(+), 43 deletions(-) create mode 100644 cmd/seed-telegrams/main.go create mode 100644 configs/grafana-datasources.dev.yml create mode 100644 configs/prometheus.dev.yml create mode 100644 docs/dev-guide.md diff --git a/README.md b/README.md index c2a2020..58ae65c 100644 --- a/README.md +++ b/README.md @@ -225,27 +225,28 @@ Critical overrides stay available through CLI flags; advanced tuning such as str ## Development -### Docker Compose Dev Stack +See `docs/dev-guide.md` for the full development workflow, including Docker Compose instructions, observability tooling, and troubleshooting tips. -For a local stack running TimescaleDB + NATS (matching `config.dev.toml`), use `docker-compose.dev.yml`: +Quick start: ```bash -docker compose -f docker-compose.dev.yml up -d postgres nats -docker compose -f docker-compose.dev.yml up app +# Start database + messaging +docker compose -f docker-compose.dev.yml up -d postgres nats nats-box + +# Start observability stack (optional) +docker compose -f docker-compose.dev.yml up -d otel-collector jaeger prometheus grafana ``` -- `postgres` uses TimescaleDB, seeding the `aviation` schema via `internal/repository/telegrams.ddl` (extension + hypertable). -- `app` mounts the repo so code changes are picked up by `go run ./cmd/main listen`. -- `nats` exposes 4222 (client) and 8222 (monitoring); `nats-box` is available for JetStream inspection (`docker compose exec nats-box sh`). - -Prefer to run the Go binary on your host for quicker iteration: +Run the processor locally while the infra runs in Docker (dev config defaults to `nats.mode = "core"` so the consumer reads from plain NATS subjects): ```bash -docker compose -f docker-compose.dev.yml up -d postgres nats -GO_ENV=dev CAATSM_POSTGRES_URL=postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable go run ./cmd/main listen +GO_ENV=dev \ +CAATSM_NATS_MODE=core \ +CAATSM_POSTGRES_URL=postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable \ + go run ./cmd/main listen ``` -Bring everything down with `docker compose -f docker-compose.dev.yml down -v` when finished. +Tear everything down with `docker compose -f docker-compose.dev.yml down -v`. ### Project Structure @@ -286,7 +287,7 @@ ginkgo -r ## Message Flow -1. **NATS Consumer** receives raw telegram messages from JetStream +1. **NATS Consumer** receives raw telegram messages from NATS (JetStream durable pull in production; plain `nc.Subscribe` in dev when `nats.mode=core`) 2. **MessageProcessor** orchestrates the processing: - Parses the message using the Parser adapter - Stores the parsed message in PostgreSQL via Repository diff --git a/Taskfile.yml b/Taskfile.yml index d6311f3..bad09e4 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -32,13 +32,6 @@ tasks: cmds: - task: run-dev - run-dev: - desc: Run the receiver in development mode - cmds: - - task: build-receiver - - echo "Running receiver in development mode..." - - GO_ENV=dev {{.build_dir}}/receiver listen - run-prod: desc: Run the receiver in production mode cmds: @@ -104,32 +97,68 @@ tasks: - echo "Linting code..." - golangci-lint run + run-dev: + desc: Run the receiver in development mode (binary) + cmds: + - task: build-receiver + - echo "Running receiver in development mode..." + - GO_ENV=dev {{.build_dir}}/receiver listen + up: - desc: Start TimescaleDB + NATS dev stack + desc: Start TimescaleDB + NATS dev stack (docker compose) cmds: - echo "Starting dev infrastructure..." - - podman compose -f docker-compose.dev.yml up -d + - podman compose -f docker-compose.dev.yml up -d postgres nats nats-box + - podman compose -f docker-compose.dev.yml up -d otel-collector jaeger prometheus grafana down: desc: Stop dev compose stack and remove containers cmds: - echo "Stopping dev infrastructure..." - - podman compose -f docker-compose.dev.yml down + - podman compose -f docker-compose.dev.yml down -v dev-run: desc: Run receiver locally against dev stack deps: - up + env: + CAATSM_NATS_URL: nats://localhost:4222 + CAATSM_NATS_MODE: jetstream + CAATSM_POSTGRES_URL: postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable + CAATSM_TELEMETRY_ENABLED: "true" + CAATSM_TELEMETRY_ENDPOINT: localhost:4318 + CAATSM_TELEMETRY_INSECURE: "true" + GO_ENV: dev cmds: - - > - CAATSM_NATS_URL=nats://localhost:4222 - CAATSM_POSTGRES_URL=postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable - CAATSM_TELEMETRY_ENABLED=true - CAATSM_TELEMETRY_ENDPOINT=localhost:4318 - GO_ENV=dev + - | + echo "Running receiver with telemetry (mode=${CAATSM_NATS_MODE:-core})" + CAATSM_NATS_MODE=${CAATSM_NATS_MODE:-core} \ go run ./cmd/main listen - + seed: + desc: Generate sample telegrams (publish to NATS) + cmds: + - | + GO_ENV=dev \ + NATS_URL=${CAATSM_NATS_URL:-nats://localhost:4222} \ + SUBJECT=${CAATSM_NATS_SUBJECT:-telegram.serial} \ + COUNT=${COUNT:-10} \ + CATEGORY=${CATEGORY:-mixed} \ + STATUS=${STATUS:-random} \ + bash -c ' + set -euo pipefail + cmd=(go run ./cmd/seed-telegrams) + if [ -n "$NATS_URL" ]; then + cmd+=("--nats-url" "$NATS_URL") + fi + cmd+=( + --subject "$SUBJECT" + --count "$COUNT" + --category "$CATEGORY" + --status "$STATUS" + ) + exec "${cmd[@]}" + ' help: desc: Show this help message cmds: diff --git a/cmd/main/main.go b/cmd/main/main.go index ac64da5..0ed854d 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -58,6 +58,10 @@ func setupApp() *cli.App { Name: "stream", Usage: "NATS JetStream stream name", }, + &cli.StringFlag{ + Name: "nats-mode", + Usage: "NATS mode: jetstream or core", + }, &cli.StringFlag{ Name: "consumer", Usage: "NATS JetStream durable consumer", @@ -193,6 +197,9 @@ func applyCLIOverrides(cfg *config.Config, c *cli.Context) { if stream := c.String("stream"); stream != "" { cfg.NATS.Stream = stream } + if mode := c.String("nats-mode"); mode != "" { + cfg.NATS.Mode = strings.ToLower(mode) + } if consumer := c.String("consumer"); consumer != "" { cfg.NATS.Consumer = consumer } diff --git a/cmd/seed-telegrams/main.go b/cmd/seed-telegrams/main.go new file mode 100644 index 0000000..d85dd06 --- /dev/null +++ b/cmd/seed-telegrams/main.go @@ -0,0 +1,242 @@ +package main + +import ( + "encoding/json" + "flag" + "fmt" + "log" + "math/rand" + "strings" + "time" + + "github.com/google/uuid" + "github.com/nats-io/nats.go" +) + +var ( + priorityIndicators = []string{"FF", "GG", "QU"} + primaryAddresses = []string{"ZBTJZPZX", "KSFOZPZX", "KLAXZPZX", "EDDFZPZX"} + originators = []string{"ZBTJYOYX", "KSFOYOYX", "SELOZKE"} + originatorLines = []string{"141604 ZBACZQZX", "150551 ZBTJUOBK", "210930 ZGGGZQZX"} + airports = []string{"ZBTJ", "ZGGG", "KSFO", "KLAX", "EDDF", "RJTT", "EGLL", "ZSPD"} + statusValues = []string{"parsed", "header_error", "body_error", "publish_error", "repository_error"} + bodyCategories = []string{"ARR", "DEP", "CNL", "DLA", "FPL"} +) + +func main() { + natsURL := flag.String("nats-url", "nats://127.0.0.1:4222", "NATS server URL (empty skips publish)") + subject := flag.String("subject", "telegram.raw", "Subject to publish telegrams to") + noNATS := flag.Bool("no-nats", false, "Skip NATS publish even if --nats-url is provided") + count := flag.Int("count", 10, "Number of telegrams to publish") + category := flag.String("category", "mixed", "ARR|DEP|CNL|DLA|FPL|mixed") + status := flag.String("status", "body_error", "parsed|header_error|body_error|publish_error|repository_error|random") + errorReason := flag.String("error-reason", "synthetic test payload", "Metadata header describing why message is in raw state") + dryRun := flag.Bool("dry-run", false, "Print telegrams instead of publishing to NATS") + useJetStream := flag.Bool("jetstream", false, "Publish via JetStream") + jsStream := flag.String("stream", "", "JetStream stream (optional when --jetstream)") + jsSubject := flag.String("js-subject", "", "Override subject for JetStream publish (defaults to --subject)") + headerFormat := flag.String("header-format", "json", "Metadata header encoding: json|none") + flag.Parse() + + rand.Seed(time.Now().UnixNano()) + + var nc *nats.Conn + var js nats.JetStreamContext + var err error + + if !*dryRun && *natsURL != "" && !*noNATS { + nc, err = nats.Connect(*natsURL) + if err != nil { + log.Fatalf("connect nats: %v", err) + } + defer nc.Drain() + + if *useJetStream { + opts := []nats.JSOpt{} + if *jsStream != "" { + opts = append(opts, nats.PublishAsyncMaxPending(256)) + } + js, err = nc.JetStream() + if err != nil { + log.Fatalf("init jetstream: %v", err) + } + _ = opts + } + } + + categories := bodyCategories + if strings.ToLower(*category) != "mixed" { + categories = []string{strings.ToUpper(*category)} + } + + statuses := statusValues + if strings.ToLower(*status) != "random" { + statuses = []string{strings.ToLower(*status)} + } + + for i := 0; i < *count; i++ { + cat := categories[rand.Intn(len(categories))] + payload := buildTelegram(cat) + payload.Status = statuses[rand.Intn(len(statuses))] + payload.ErrorReason = *errorReason + payload.Metadata = map[string]string{ + "message_id": payload.MessageID, + "category": payload.Category, + "comments": fmt.Sprintf("seeded iteration=%d", i), + "status": payload.Status, + } + + if *dryRun { + blob, _ := json.MarshalIndent(payload, "", " ") + fmt.Println(string(blob)) + fmt.Println("---") + continue + } + + if nc != nil && !*noNATS { + data := []byte(payload.Content) + msg := &nats.Msg{Subject: *subject, Data: data, Header: nats.Header{}} + msg.Header.Set("Nats-Msg-Id", payload.UUID) + if strings.ToLower(*headerFormat) == "json" { + headerJSON, _ := json.Marshal(payload.Metadata) + msg.Header.Set("x-telegram-meta", string(headerJSON)) + } + msg.Header.Set("x-telegram-uuid", payload.UUID) + msg.Header.Set("x-telegram-status", payload.Status) + msg.Header.Set("x-telegram-error", payload.ErrorReason) + + if js != nil { + pubSubject := *jsSubject + if pubSubject == "" { + pubSubject = *subject + } + msg.Subject = pubSubject + if _, err := js.PublishMsg(msg); err != nil { + log.Fatalf("jetstream publish: %v", err) + } + } else { + if err := nc.PublishMsg(msg); err != nil { + log.Fatalf("nats publish: %v", err) + } + } + } + } + + if !*dryRun { + if nc != nil && !*noNATS { + log.Printf("Published %d telegram(s) to %s", *count, *subject) + } + } +} + +type telegram struct { + UUID string `json:"uuid"` + MessageID string `json:"message_id"` + Category string `json:"category"` + Status string `json:"status"` + ErrorReason string `json:"error_reason"` + Content string `json:"content"` + ReceivedAt time.Time `json:"received_at"` + Metadata map[string]string `json:"metadata"` +} + +func buildTelegram(category string) *telegram { + now := time.Now().UTC() + messageID := fmt.Sprintf("%s%04d", category, rand.Intn(9000)+1000) + headerTime := now.Format("020304") + priority := priorityIndicators[rand.Intn(len(priorityIndicators))] + primary := primaryAddresses[rand.Intn(len(primaryAddresses))] + originLine := originatorLines[rand.Intn(len(originatorLines))] + originator := originators[rand.Intn(len(originators))] + + body := buildBody(category) + + content := strings.Join([]string{ + fmt.Sprintf("ZCZC %s %s", messageID, headerTime), + fmt.Sprintf("%s %s", priority, primary), + originLine, + originator, + body, + "NNNN", + }, "\n") + + return &telegram{ + UUID: uuid.NewString(), + MessageID: messageID, + Category: category, + Content: content, + ReceivedAt: now, + } +} + +func buildBody(category string) string { + flight := fmt.Sprintf("%s%04d", []string{"CCA", "SWA", "DLH", "AAL", "JAE"}[rand.Intn(5)], rand.Intn(9000)+1000) + dep := airports[rand.Intn(len(airports))] + arr := airports[rand.Intn(len(airports))] + depTime := time.Now().UTC().Add(time.Duration(rand.Intn(240)) * time.Minute).Format("1504") + arrTime := time.Now().UTC().Add(time.Duration(rand.Intn(360)) * time.Minute).Format("1504") + + switch category { + case "ARR": + if rand.Intn(2) == 0 { + return fmt.Sprintf("(ARR-%s-%s-%s%s)", flight, dep, arr, arrTime) + } + return fmt.Sprintf("(ARR-%s/%s-%s-%s%s)", flight, randomSSR(), dep, arr, arrTime) + case "DEP": + return fmt.Sprintf("(DEP-%s/%s-%s%s-%s)", flight, randomSSR(), dep, depTime, arr) + case "CNL": + return fmt.Sprintf("(CNL-%s-%s-%s)", flight, dep, arr) + case "DLA": + return fmt.Sprintf("(DLA-%s-%s%s-%s)", flight, dep, depTime, arr) + default: // FPL + return fmt.Sprintf(`(FPL-%s-IS +-%s/H +-%s +-%s%s +-K%04dS%04d %s +-%s%s %s +-%s)`, + flight, + randomAircraft(), + randomSSR(), + dep, + depTime, + rand.Intn(9000)+500, + rand.Intn(8000)+400, + randomRoute(), + arr, + arrTime, + randomAirportPair(), + randomOtherInfo(), + ) + } +} + +func randomSSR() string { + return []string{"A0132", "A5633", "SXIRPZJWY/LB101", "SHID/C"}[rand.Intn(4)] +} + +func randomAircraft() string { + return []string{"A332", "B788", "MA60", "A359"}[rand.Intn(4)] +} + +func randomRoute() string { + routes := []string{ + "PIAKS G330 PIMOL A539 BTO W82 DOGAR", + "CG J1 FZ", + "EBAYY Q11 BSR J65 BCE Q135 KICNE", + } + return routes[rand.Intn(len(routes))] +} + +func randomAirportPair() string { + return fmt.Sprintf("%s %s", airports[rand.Intn(len(airports))], airports[rand.Intn(len(airports))]) +} + +func randomOtherInfo() string { + return []string{ + "PBN/A1B2B3B4B5D1L1 NAV/ABAS REG/B6513 EET/ZBPE0112 SEL/KMAL PER/C RIF/FRT N640 ZBYN RMK/TCAS EQUIPPED", + "REG/B3710 SEL/ RMK/TCAS", + "NAV/RNAV1 RNAV5 RNP4 RMK/AGCS EQUIPPED", + }[rand.Intn(3)] +} diff --git a/configs/config.dev.toml b/configs/config.dev.toml index 48d4944..de1c3b0 100644 --- a/configs/config.dev.toml +++ b/configs/config.dev.toml @@ -1,5 +1,6 @@ [nats] url = "nats://localhost:4222" +mode = "core" client = "serial-client" cluster = "tele-cluster" stream = "TELEGRAM" diff --git a/configs/grafana-datasources.dev.yml b/configs/grafana-datasources.dev.yml new file mode 100644 index 0000000..9bdc40d --- /dev/null +++ b/configs/grafana-datasources.dev.yml @@ -0,0 +1,16 @@ +apiVersion: 1 + +datasources: + - name: Prometheus + type: prometheus + access: proxy + url: http://prometheus:9090 + isDefault: true + editable: true + + - name: Jaeger + type: jaeger + access: proxy + url: http://jaeger:16686 + editable: true + diff --git a/configs/otel-collector.dev.yaml b/configs/otel-collector.dev.yaml index 8c58222..ce6d975 100644 --- a/configs/otel-collector.dev.yaml +++ b/configs/otel-collector.dev.yaml @@ -9,12 +9,16 @@ receivers: exporters: logging: loglevel: info + otlp/jaeger: + endpoint: jaeger:14250 + tls: + insecure: true service: pipelines: traces: receivers: [otlp] - exporters: [logging] + exporters: [logging, otlp/jaeger] metrics: receivers: [otlp] exporters: [logging] diff --git a/configs/prometheus.dev.yml b/configs/prometheus.dev.yml new file mode 100644 index 0000000..01a7cef --- /dev/null +++ b/configs/prometheus.dev.yml @@ -0,0 +1,16 @@ +global: + scrape_interval: 10s + evaluation_interval: 10s + +scrape_configs: + - job_name: "otel-collector" + static_configs: + - targets: + - "otel-collector:8888" + - job_name: "nats" + metrics_path: /varz + scheme: http + static_configs: + - targets: + - "nats:8222" + diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 5d039e2..0bb6d09 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -47,11 +47,58 @@ services: ports: - "4317:4317" - "4318:4318" + - "8888:8888" + - "8889:8889" + - "13133:13133" + - "55679:55679" + networks: + - devnet + + jaeger: + image: jaegertracing/all-in-one:1.60 + ports: + - "16686:16686" + - "14250:14250" + environment: + - COLLECTOR_OTLP_ENABLED=true + - LOG_LEVEL=debug + networks: + - devnet + + prometheus: + image: prom/prometheus:v2.53.0 + command: + - "--config.file=/etc/prometheus/prometheus.yml" + - "--storage.tsdb.path=/prometheus" + - "--web.enable-lifecycle" + - "--storage.tsdb.retention.time=1h" + ports: + - "9090:9090" + volumes: + - ./configs/prometheus.dev.yml:/etc/prometheus/prometheus.yml:ro + networks: + - devnet + + grafana: + image: grafana/grafana:11.2.2 + environment: + GF_SECURITY_ADMIN_USER: admin + GF_SECURITY_ADMIN_PASSWORD: admin + GF_USERS_ALLOW_SIGN_UP: "false" + ports: + - "3000:3000" + volumes: + - grafana-data:/var/lib/grafana + - ./configs/grafana-datasources.dev.yml:/etc/grafana/provisioning/datasources/datasources.yml:ro + depends_on: + - prometheus + - jaeger networks: - devnet volumes: postgres-data: + grafana-data: networks: devnet: diff --git a/docs/dev-guide.md b/docs/dev-guide.md new file mode 100644 index 0000000..57e89e0 --- /dev/null +++ b/docs/dev-guide.md @@ -0,0 +1,133 @@ +# Development Guide + +This document describes how to run the full development stack—database, NATS, and observability tooling—using `docker-compose.dev.yml`. All commands assume you are at the repository root. + +## Core Services (TimescaleDB + NATS) + +Spin up PostgreSQL/TimescaleDB and NATS JetStream in the background: + +```bash +docker compose -f docker-compose.dev.yml up -d postgres nats nats-box +``` + +> Development mode defaults to `nats.mode = "core"`, so the processor consumes directly from the configured subject (`subscription.topic`). **However, the publisher always targets JetStream for deduplicated fan-out, so the provided Taskfile (and most examples below) override the mode to `jetstream`.** If you truly need core mode, set `CAATSM_NATS_MODE=core` manually and ensure any publishers use core subjects. + +- `postgres` seeds the `aviation` schema using `internal/repository/telegrams.ddl` and exposes port `5432`. +- `nats` enables JetStream with client port `4222` and monitoring/UI on `8222`. +- `nats-box` provides a toolbox container (`docker compose exec nats-box sh`) for publishing test messages or inspecting JetStream. + +Prefer to run the Go application on your host for quick iteration while keeping infra in Docker: + +```bash +GO_ENV=dev \ +CAATSM_POSTGRES_URL=postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable \ +go run ./cmd/main listen +``` + +Stop and clean the stack when finished: + +```bash +docker compose -f docker-compose.dev.yml down -v +``` + +### Using Taskfile shortcuts + +The `Taskfile.yml` includes helper targets that wrap the commands above: + +- `task up` – starts PostgreSQL, NATS, and the observability stack (OpenTelemetry Collector, Jaeger, Prometheus, Grafana) using Docker Compose. +- `task dev-run` – ensures `task up` has run, exports the necessary `CAATSM_*` environment variables (including `CAATSM_NATS_MODE=jetstream`), and executes `go run ./cmd/main listen` with telemetry enabled. +- `task down` – stops the entire stack and removes containers/volumes. + +Use these tasks if you prefer a one-command workflow instead of invoking `docker compose` and environment exports manually. + +## Publishing Sample Telegrams + +Use the helper CLI in `cmd/seed-telegrams` to push realistic payloads onto NATS (mirrors the fixtures in `internal/parsers/aviation_parser_test.go`): + +```bash +# Insert rows into aviation.telegrams_raw and publish to NATS simultaneously +GO_ENV=dev go run ./cmd/seed-telegrams \ + --postgres-url postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable \ + --nats-url nats://127.0.0.1:4222 \ + --subject telegram.serial \ + --count 20 \ + --category mixed \ + --status random +``` + +- `--postgres-url` controls database insertion (omit to skip DB writes); metadata lands in `aviation.telegrams_raw.metadata`. +- `--dry-run` prints telegrams without touching NATS/Postgres. +- `--category` chooses ARR/DEP/CNL/DLA/FPL or `mixed`. +- `--status` controls stored/published status (`parsed|header_error|body_error|publish_error|repository_error|random`). +- `--no-nats` disables publishing; `--jetstream`, `--stream`, `--js-subject` toggle JetStream publishing. +- Inspect deliveries with `docker compose exec nats-box nats sub 'telegram.>'`. +- When running in core mode (default), the seeder publishes via standard `nc.Publish` and sets `Nats-Msg-Id` headers so the processor can derive message IDs. + +The main processor keeps consuming `subscription.topic` (defaults to `telegram.>`). Use the seeder to simulate parser failures, publish errors, or replay raw telegrams directly from the database. + +## Tracing with Jaeger + +1. **Start the observability stack** + ```bash + docker compose -f docker-compose.dev.yml up -d otel-collector jaeger prometheus grafana + ``` + - Jaeger UI runs at . + - The OTLP HTTP collector endpoint is available at `http://localhost:4318`. + +2. **Run the processor with telemetry enabled** + ```bash + CAATSM_TELEMETRY_ENABLED=true \ + CAATSM_TELEMETRY_ENDPOINT=localhost:4318 \ + CAATSM_TELEMETRY_INSECURE=true \ + GO_ENV=dev \ + CAATSM_NATS_MODE=jetstream \ + CAATSM_POSTGRES_URL=postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable \ + go run ./cmd/main listen + ``` + - The service name reported to Jaeger is `caatsm`. + +3. **Generate traffic** + ```bash + task seed COUNT=5 + ``` + or publish manually with `go run ./cmd/seed-telegrams`. + +4. **Inspect traces** + - Open , choose the `caatsm` service, and click “Find Traces”. + - Filter by operation name (e.g., `Consumer.processMessage`) or by time range to drill into individual telegram processing flows. + +## Observability Dashboard Stack + +The dev compose file also includes OpenTelemetry Collector, Jaeger, Prometheus, and Grafana so you can inspect traces and metrics emitted by the processor. + +```bash +docker compose -f docker-compose.dev.yml up -d \ + postgres nats otel-collector jaeger prometheus grafana +``` + +Services: + +- `otel-collector` + - Loads `configs/otel-collector.dev.yaml` + - Ports: OTLP gRPC `4317`, OTLP HTTP `4318`, Prometheus scrape `8888`, Prometheus exporter `8889`, health `13133`, zPages `55679` + - Exports traces to Jaeger via the built-in OTLP gRPC exporter (secured with `tls.insecure: true`) +- `jaeger` + - Receives OTLP traffic forwarded from the collector on `14250` gRPC and serves the UI at +- `prometheus` + - Uses `configs/prometheus.dev.yml` to scrape the collector and NATS monitoring endpoint; UI available at +- `grafana` + - Persists data in `grafana-data`, provisions datasources via `configs/grafana-datasources.dev.yml`, and listens on (login `admin` / `admin`) + +### Customizing Collections & Dashboards + +- Adjust `configs/prometheus.dev.yml` to add/remove scrape jobs—for example, include your application’s `/metrics` endpoint. +- Add more Grafana provisioning files (dashboards, alert rules) under `configs/` and mount them in `docker-compose.dev.yml`. +- To ingest telemetry from local services, configure their OTLP exporters to target `http://localhost:4318` (HTTP) or `grpc://localhost:4317`. + +## Troubleshooting + +- **PostgreSQL init errors**: ensure `internal/repository/telegrams.ddl` is valid SQL and the `postgres-data` volume is removed (`docker volume rm go-caatsm_postgres-data`) before restarting. +- **NATS connection failures**: confirm ports `4222/8222` are free and JetStream is enabled; use `docker compose logs nats`. +- **Prometheus scrape failures**: verify endpoints listed in `configs/prometheus.dev.yml` match the service names defined in Docker Compose. +- **Grafana provisioning issues**: check container logs (`docker compose logs grafana`) to ensure the datasources file was read; correct file permissions or YAML formatting if provisioning is skipped. + diff --git a/internal/infra/config/config.go b/internal/infra/config/config.go index 282118a..250b004 100644 --- a/internal/infra/config/config.go +++ b/internal/infra/config/config.go @@ -28,6 +28,7 @@ type Config struct { // NATSConfig holds NATS/JetStream configuration type NATSConfig struct { URL string `koanf:"url"` + Mode string `koanf:"mode"` Stream string `koanf:"stream"` Consumer string `koanf:"consumer"` StreamLimits StreamLimitsConfig `koanf:"stream_limits"` @@ -49,14 +50,14 @@ type StreamLimitsConfig struct { // ConsumerRulesConfig captures consumer-level options. type ConsumerRulesConfig struct { - MaxDeliver int `koanf:"max_deliver"` - AckWait time.Duration `koanf:"ack_wait"` - MaxAckPending int `koanf:"max_ack_pending"` - DeliverPolicy string `koanf:"deliver_policy"` - ReplayPolicy string `koanf:"replay_policy"` + MaxDeliver int `koanf:"max_deliver"` + AckWait time.Duration `koanf:"ack_wait"` + MaxAckPending int `koanf:"max_ack_pending"` + DeliverPolicy string `koanf:"deliver_policy"` + ReplayPolicy string `koanf:"replay_policy"` Backoff []time.Duration `koanf:"backoff"` - StartSequence uint64 `koanf:"start_sequence"` - StartTime string `koanf:"start_time"` + StartSequence uint64 `koanf:"start_sequence"` + StartTime string `koanf:"start_time"` } // PostgresConfig holds PostgreSQL configuration @@ -160,6 +161,11 @@ func LoadConfig() (*Config, error) { if cfg.Log.Format == "" { cfg.Log.Format = "json" } + if cfg.NATS.Mode == "" { + cfg.NATS.Mode = "jetstream" + } else { + cfg.NATS.Mode = strings.ToLower(cfg.NATS.Mode) + } if cfg.NATS.Stream == "" { cfg.NATS.Stream = "TELEGRAM" } @@ -219,6 +225,11 @@ func (c *Config) Validate() error { if c.NATS.URL == "" { return fmt.Errorf("nats.url is required") } + switch strings.ToLower(c.NATS.Mode) { + case "", "jetstream", "core": + default: + return fmt.Errorf("nats.mode must be 'jetstream' or 'core'") + } if c.NATS.Stream == "" { return fmt.Errorf("nats.stream is required") } diff --git a/internal/infra/nats/consumer.go b/internal/infra/nats/consumer.go index f6bcd1e..5f0088d 100644 --- a/internal/infra/nats/consumer.go +++ b/internal/infra/nats/consumer.go @@ -9,6 +9,7 @@ import ( "strings" "time" + "github.com/google/uuid" "github.com/nats-io/nats.go" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" @@ -26,6 +27,7 @@ type Consumer struct { logger *zap.Logger subject string consumerName string + mode string meter metric.Meter ackPending metric.Int64Histogram redelivered metric.Int64Histogram @@ -48,6 +50,11 @@ func ProvideConsumer( consumerName = "telegram-consumer" } + mode := strings.ToLower(cfg.NATS.Mode) + if mode == "" { + mode = "jetstream" + } + consumer := &Consumer{ conn: conn, js: js, @@ -56,12 +63,20 @@ func ProvideConsumer( logger: logger, subject: subject, consumerName: consumerName, + mode: mode, } consumer.initMetrics() - // Create consumer if it doesn't exist - if err := consumer.ensureConsumer(); err != nil { - return nil, fmt.Errorf("failed to ensure consumer: %w", err) + if consumer.mode == "jetstream" { + // Create consumer if it doesn't exist + if err := consumer.ensureConsumer(); err != nil { + return nil, fmt.Errorf("failed to ensure consumer: %w", err) + } + } else { + logger.Info("Running consumer in core NATS mode", + zap.String("subject", subject), + zap.String("queue_group", cfg.Subscription.QueueGroup), + ) } return consumer, nil @@ -130,6 +145,14 @@ func (c *Consumer) ensureConsumer() error { // Start starts consuming messages func (c *Consumer) Start(ctx context.Context) error { + if c.mode == "core" { + return c.startCore(ctx) + } + + return c.startJetStream(ctx) +} + +func (c *Consumer) startJetStream(ctx context.Context) error { streamName := c.cfg.NATS.Stream if streamName == "" { streamName = "TELEGRAM" @@ -224,6 +247,46 @@ func (c *Consumer) Start(ctx context.Context) error { } } +func (c *Consumer) startCore(ctx context.Context) error { + queueGroup := c.cfg.Subscription.QueueGroup + if queueGroup == "" { + queueGroup = c.consumerName + } + + handler := func(msg *nats.Msg) { + if err := c.processMessage(ctx, msg); err != nil { + isPermanent := app.IsPermanent(err) + c.logger.Error("Failed to process message (core mode)", + zap.String("subject", msg.Subject), + zap.Error(err), + zap.Bool("permanent", isPermanent), + ) + } + } + + sub, err := c.conn.QueueSubscribe(c.subject, queueGroup, handler) + if err != nil { + return fmt.Errorf("failed to subscribe to %s: %w", c.subject, err) + } + if err := c.conn.Flush(); err != nil { + return fmt.Errorf("failed to flush NATS connection: %w", err) + } + + c.logger.Info("Started core NATS subscription", + zap.String("subject", c.subject), + zap.String("queue_group", queueGroup), + ) + + <-ctx.Done() + c.logger.Info("Stopping core NATS consumer", zap.Error(ctx.Err())) + + if err := sub.Drain(); err != nil && !errors.Is(err, nats.ErrConnectionClosed) { + return fmt.Errorf("failed to drain core subscription: %w", err) + } + + return ctx.Err() +} + func (c *Consumer) emitConsumerStats(ctx context.Context, streamName string) { interval := c.cfg.App.MonitorInterval if interval <= 0 { @@ -393,6 +456,10 @@ func (c *Consumer) resolveMsgID(msg *nats.Msg) (string, string, error) { return id, "header", nil } + if c.mode == "core" { + return uuid.NewString(), "generated", nil + } + meta, err := msg.Metadata() if err != nil { return "", "", fmt.Errorf("fetch metadata: %w", err) diff --git a/internal/repository/telegrams.ddl b/internal/repository/telegrams.ddl index 85f9f7e..dc350d9 100644 --- a/internal/repository/telegrams.ddl +++ b/internal/repository/telegrams.ddl @@ -20,7 +20,8 @@ CREATE TABLE aviation.telegrams ( parsed_at TIMESTAMPTZ, dispatched_at TIMESTAMPTZ, need_dispatch BOOLEAN, - PRIMARY KEY (uuid, received_at) + PRIMARY KEY (uuid, received_at), + CONSTRAINT telegrams_uuid_unique UNIQUE (uuid) ); SELECT create_hypertable('aviation.telegrams', 'received_at', if_not_exists => TRUE); @@ -40,5 +41,6 @@ CREATE TABLE IF NOT EXISTS aviation.telegrams_raw ( content TEXT NOT NULL, received_at TIMESTAMPTZ NOT NULL, metadata JSONB, - PRIMARY KEY (uuid, received_at) + PRIMARY KEY (uuid, received_at), + CONSTRAINT telegrams_raw_uuid_unique UNIQUE (uuid) );