diff --git a/Taskfile.yml b/Taskfile.yml index bad09e4..46a7bf8 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -108,19 +108,17 @@ tasks: desc: Start TimescaleDB + NATS dev stack (docker compose) cmds: - echo "Starting dev infrastructure..." - - 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 + - docker compose -f docker-compose.dev.yml up -d postgres nats nats-box + - docker 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 -v + - docker 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 diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 0bb6d09..6b229eb 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -19,6 +19,23 @@ services: networks: - devnet + postgres-init: + image: timescale/timescaledb:2.15.2-pg16 + depends_on: + - postgres + environment: + PGPASSWORD: caatsm + command: > + bash -c " + until pg_isready -h postgres -U caatsm; do sleep 1; done; + psql --host=postgres --username=caatsm --dbname=aviation --file=/tmp/telegrams.ddl + " + volumes: + - ./internal/repository/telegrams.ddl:/tmp/telegrams.ddl:ro + restart: "no" + networks: + - devnet + nats: image: nats:2.10-alpine command: ["-js", "--http_port=8222", "-DV"] @@ -30,9 +47,35 @@ services: nats-box: image: synadia/nats-box:latest - entrypoint: ["/bin/sh", "-c", "sleep infinity"] + entrypoint: ["sleep", "infinity"] depends_on: - nats + restart: unless-stopped + networks: + - devnet + + nats-init: + image: synadia/nats-box:latest + depends_on: + - nats + command: > + /bin/sh -c " + echo 'waiting for NATS...' ; + until nats --server nats://nats:4222 server ping >/dev/null 2>&1; do sleep 1; done; + echo 'initializing JetStream...' ; + nats --server nats://nats:4222 stream add TELEGRAM \ + --subjects 'telegram.serial,telegram.json' \ + --storage file \ + --retention limits \ + --replicas 1 \ + --max-msgs -1 \ + --max-bytes -1 \ + --if-not-exists \ + --defaults + echo 'JetStream init done.'; + sleep 2; + " + restart: "no" networks: - devnet environment: @@ -60,8 +103,8 @@ services: - "16686:16686" - "14250:14250" environment: - - COLLECTOR_OTLP_ENABLED=true - - LOG_LEVEL=debug + COLLECTOR_OTLP_ENABLED: "true" + LOG_LEVEL: debug networks: - devnet @@ -103,4 +146,3 @@ volumes: networks: devnet: driver: bridge - diff --git a/internal/infra/nats/jetstream.go b/internal/infra/nats/jetstream.go index 930d3dc..b03d2a3 100644 --- a/internal/infra/nats/jetstream.go +++ b/internal/infra/nats/jetstream.go @@ -45,7 +45,14 @@ func ProvideJetStream(nc *nats.Conn, cfg *config.Config, logger *zap.Logger) (na // Create stream if it doesn't exist streamName := cfg.NATS.Stream - subject := cfg.EffectiveSubscriptionTopic() + consumerSubject := cfg.EffectiveSubscriptionTopic() + publisherSubject := strings.TrimSpace(cfg.Publisher.Topic) + + streamSubjects := dedupeSubjects([]string{consumerSubject, publisherSubject}) + if len(streamSubjects) == 0 { + nc.Close() + return nil, fmt.Errorf("no subjects configured for JetStream stream %s", streamName) + } streamLimits := cfg.NATS.StreamLimits storage := nats.FileStorage @@ -63,7 +70,7 @@ func ProvideJetStream(nc *nats.Conn, cfg *config.Config, logger *zap.Logger) (na streamConfig := &nats.StreamConfig{ Name: streamName, - Subjects: []string{subject}, + Subjects: streamSubjects, Retention: nats.LimitsPolicy, MaxMsgs: streamLimits.MaxMsgs, MaxBytes: streamLimits.MaxBytes, @@ -81,7 +88,10 @@ func ProvideJetStream(nc *nats.Conn, cfg *config.Config, logger *zap.Logger) (na nc.Close() return nil, fmt.Errorf("failed to create stream: %w", err) } - logger.Info("Created JetStream", zap.String("stream", streamName), zap.String("subject", subject)) + logger.Info("Created JetStream", + zap.String("stream", streamName), + zap.Strings("subjects", streamSubjects), + ) } else { nc.Close() return nil, fmt.Errorf("stream %s not found and auto-creation disabled", streamName) @@ -91,7 +101,7 @@ func ProvideJetStream(nc *nats.Conn, cfg *config.Config, logger *zap.Logger) (na return nil, fmt.Errorf("failed to fetch stream info: %w", err) } } else { - validateStreamConfig(info, subject, logger) + validateStreamConfig(info, streamSubjects, logger) } return js, nil @@ -106,21 +116,35 @@ func shouldBootstrapStream() bool { } } -func validateStreamConfig(info *nats.StreamInfo, expectedSubject string, logger *zap.Logger) { +func validateStreamConfig(info *nats.StreamInfo, expectedSubjects []string, logger *zap.Logger) { if info == nil { return } + defer func() { + if len(expectedSubjects) == 0 { + expectedSubjects = []string{""} + } + }() - if !subjectListContains(info.Config.Subjects, expectedSubject) { - logger.Warn("JetStream stream subjects do not match config", + missing := make([]string, 0) + for _, subj := range expectedSubjects { + if subj == "" { + continue + } + if !containsSubject(info.Config.Subjects, subj) { + missing = append(missing, subj) + } + } + if len(missing) > 0 { + logger.Warn("JetStream stream subjects missing expected entries", zap.String("stream", info.Config.Name), zap.Strings("stream_subjects", info.Config.Subjects), - zap.String("configured_subject", expectedSubject), + zap.Strings("missing_subjects", missing), ) } } -func subjectListContains(subjects []string, target string) bool { +func containsSubject(subjects []string, target string) bool { for _, s := range subjects { if s == target { return true @@ -128,3 +152,20 @@ func subjectListContains(subjects []string, target string) bool { } return false } + +func dedupeSubjects(subjects []string) []string { + seen := make(map[string]struct{}) + result := make([]string, 0, len(subjects)) + for _, subj := range subjects { + subj = strings.TrimSpace(subj) + if subj == "" { + continue + } + if _, ok := seen[subj]; ok { + continue + } + seen[subj] = struct{}{} + result = append(result, subj) + } + return result +} diff --git a/internal/infra/postgres/repository.go b/internal/infra/postgres/repository.go index 9be97d4..55563a3 100644 --- a/internal/infra/postgres/repository.go +++ b/internal/infra/postgres/repository.go @@ -56,7 +56,7 @@ func (r *Repository) InsertOne(ctx context.Context, msg *domain.ParsedMessage) e ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) - ON CONFLICT (uuid) DO NOTHING + ON CONFLICT (uuid, received_at) DO NOTHING ` tag, err := r.pool.Exec(ctx, query, @@ -167,11 +167,10 @@ func (r *Repository) InsertRaw(ctx context.Context, msg *domain.ParsedMessage) e ) VALUES ( $1, $2, $3, $4, $5, $6 ) - ON CONFLICT (uuid) DO UPDATE + ON CONFLICT (uuid, received_at) DO UPDATE SET status = EXCLUDED.status, error_reason = EXCLUDED.error_reason, content = EXCLUDED.content, - received_at = EXCLUDED.received_at, metadata = EXCLUDED.metadata ` diff --git a/internal/repository/telegrams.ddl b/internal/repository/telegrams.ddl index dc350d9..85f9f7e 100644 --- a/internal/repository/telegrams.ddl +++ b/internal/repository/telegrams.ddl @@ -20,8 +20,7 @@ CREATE TABLE aviation.telegrams ( parsed_at TIMESTAMPTZ, dispatched_at TIMESTAMPTZ, need_dispatch BOOLEAN, - PRIMARY KEY (uuid, received_at), - CONSTRAINT telegrams_uuid_unique UNIQUE (uuid) + PRIMARY KEY (uuid, received_at) ); SELECT create_hypertable('aviation.telegrams', 'received_at', if_not_exists => TRUE); @@ -41,6 +40,5 @@ CREATE TABLE IF NOT EXISTS aviation.telegrams_raw ( content TEXT NOT NULL, received_at TIMESTAMPTZ NOT NULL, metadata JSONB, - PRIMARY KEY (uuid, received_at), - CONSTRAINT telegrams_raw_uuid_unique UNIQUE (uuid) + PRIMARY KEY (uuid, received_at) );