From 6616c7e10db27530ac28e8e77066e517f72af6bb Mon Sep 17 00:00:00 2001 From: windyboy Date: Sun, 16 Nov 2025 11:38:49 +0800 Subject: [PATCH] =?UTF-8?q?=E2=9C=A8=20Add=20new=20`dev-run-compose`=20tas?= =?UTF-8?q?k=20to=20Taskfile=20for=20running=20the=20CAATSM=20receiver=20i?= =?UTF-8?q?n=20a=20Docker=20Compose=20development=20stack.=20Update=20moni?= =?UTF-8?q?toring=20configuration=20in=20`config.dev.toml`=20to=20use=20co?= =?UTF-8?q?nsole=20format.=20Modify=20Prometheus=20scrape=20configuration?= =?UTF-8?q?=20to=20target=20a=20specific=20IP=20address.=20Enhance=20NATS?= =?UTF-8?q?=20consumer=20error=20handling=20with=20retry=20logic=20for=20J?= =?UTF-8?q?etStream=20availability=20and=20adjust=20logging=20level=20for?= =?UTF-8?q?=20consumer=20metrics.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Taskfile.yml | 16 +++++++++++++--- configs/config.dev.toml | 2 +- configs/prometheus.dev.yml | 4 +--- internal/infra/nats/consumer.go | 13 ++++++++++++- internal/infra/nats/jetstream.go | 2 ++ 5 files changed, 29 insertions(+), 8 deletions(-) diff --git a/Taskfile.yml b/Taskfile.yml index eb74482..f5e6342 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -147,6 +147,13 @@ tasks: - echo "Stopping dev infrastructure..." - docker compose -f docker-compose.dev.yml down -v + dev-run-compose: + desc: Run receiver inside docker-compose dev stack + cmds: + - echo "Starting dev stack including caatsm-receiver..." + - docker compose -f docker-compose.dev.yml up -d postgres nats nats-box nats-exporter + - docker compose -f docker-compose.dev.yml up -d otel-collector jaeger prometheus grafana caatsm-receiver + dev-run: desc: Run receiver locally against dev stack deps: @@ -158,6 +165,9 @@ tasks: CAATSM_TELEMETRY_ENABLED: "true" CAATSM_TELEMETRY_ENDPOINT: localhost:4318 CAATSM_TELEMETRY_INSECURE: "true" + CAATSM_MONITORING_ADDR: 0.0.0.0:2112 + CAATSM_MONITORING_ENABLE_METRICS: "true" + CAATSM_MONITORING_ENABLE_HEALTH: "true" GO_ENV: dev cmds: - | @@ -174,9 +184,9 @@ tasks: 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} \ + COUNT={{.COUNT | default "10"}} \ + CATEGORY={{.CATEGORY | default "mixed"}} \ + STATUS={{.STATUS | default "random"}} \ bash -c ' set -euo pipefail cmd=(go run ./cmd/seed-telegrams) diff --git a/configs/config.dev.toml b/configs/config.dev.toml index 1432001..1985196 100644 --- a/configs/config.dev.toml +++ b/configs/config.dev.toml @@ -49,7 +49,7 @@ monitor_interval = "30s" [log] level = "info" -format = "json" +format = "console" [telemetry] enabled = true diff --git a/configs/prometheus.dev.yml b/configs/prometheus.dev.yml index 705dbe5..858947f 100644 --- a/configs/prometheus.dev.yml +++ b/configs/prometheus.dev.yml @@ -15,7 +15,5 @@ scrape_configs: - job_name: "caatsm-receiver" static_configs: - targets: - # When go-caatsm runs on the host/WSL and Prometheus runs in Docker, - # use host.docker.internal to reach the host monitoring server. - - "host.docker.internal:2112" + - 172.23.189.12:2112 diff --git a/internal/infra/nats/consumer.go b/internal/infra/nats/consumer.go index 21a7189..e68ed0b 100644 --- a/internal/infra/nats/consumer.go +++ b/internal/infra/nats/consumer.go @@ -241,6 +241,17 @@ func (c *Consumer) startJetStream(ctx context.Context) error { // Timeout is expected when no messages are available continue } + if errors.Is(err, nats.ErrNoResponders) { + // JetStream API is currently unavailable (e.g., NATS just restarted or JetStream not ready). + // Back off a bit to avoid log spam while allowing the system to recover. + c.logger.Warn("JetStream not available, will retry", + zap.Error(err), + zap.String("stream", c.streamName), + zap.String("consumer", c.consumerName), + ) + time.Sleep(5 * time.Second) + continue + } c.logger.Error("Failed to fetch messages", zap.Error(err)) time.Sleep(time.Second) continue @@ -377,7 +388,7 @@ func (c *Consumer) emitConsumerStats(ctx context.Context) { continue } - c.logger.Info("JetStream consumer metrics", + c.logger.Debug("JetStream consumer metrics", zap.String("stream", c.streamName), zap.String("consumer", c.consumerName), zap.Uint64("num_ack_pending", uint64(info.NumAckPending)), diff --git a/internal/infra/nats/jetstream.go b/internal/infra/nats/jetstream.go index b03d2a3..0e732a2 100644 --- a/internal/infra/nats/jetstream.go +++ b/internal/infra/nats/jetstream.go @@ -18,6 +18,8 @@ func ProvideNATSConn(cfg *config.Config, logger *zap.Logger) (*nats.Conn, error) nats.RetryOnFailedConnect(true), nats.Timeout(cfg.Timeouts.Server), nats.ReconnectWait(cfg.Timeouts.ReconnectWait), + // Use infinite reconnects so the app survives long NATS outages (e.g. docker compose down/up). + nats.MaxReconnects(-1), nats.DisconnectErrHandler(func(nc *nats.Conn, err error) { if err != nil { logger.Warn("NATS disconnected", zap.Error(err))