diff --git a/Makefile b/Makefile index eba1f63..4a8fd46 100644 --- a/Makefile +++ b/Makefile @@ -18,9 +18,14 @@ build: ## Build the receiver binary run: run-dev ## Alias for run-dev .PHONY: run-dev -run-dev: build ## Run the receiver in development mode (uses core NATS mode by default) +run-dev: build ## Run the receiver in development mode (uses core NATS mode by default). Use MONITORING_ADDR=ip:port to set monitoring address @echo "Running receiver in development mode (NATS mode: core by default)..." - @GO_ENV=dev $(BINARY) listen + @if [ -n "$(MONITORING_ADDR)" ]; then \ + echo "Monitoring address: $(MONITORING_ADDR)"; \ + GO_ENV=dev $(BINARY) listen --monitoring-addr $(MONITORING_ADDR); \ + else \ + GO_ENV=dev $(BINARY) listen; \ + fi .PHONY: run-prod run-prod: build ## Run the receiver in production mode (requires config.prod.toml and JetStream mode) @@ -38,9 +43,14 @@ run-test: build ## Run the receiver in test mode @GO_ENV=test $(BINARY) listen .PHONY: run-local -run-local: ## Run receiver directly via go run (uses core NATS mode by default in dev) +run-local: ## Run receiver directly via go run (uses core NATS mode by default in dev). Use MONITORING_ADDR=ip:port to set monitoring address @echo "Running receiver via go run (GO_ENV=$(GO_ENV), NATS mode: core by default in dev)..." - @GO_ENV=$(GO_ENV) go run $(CMD) listen + @if [ -n "$(MONITORING_ADDR)" ]; then \ + echo "Monitoring address: $(MONITORING_ADDR)"; \ + GO_ENV=$(GO_ENV) go run $(CMD) listen --monitoring-addr $(MONITORING_ADDR); \ + else \ + GO_ENV=$(GO_ENV) go run $(CMD) listen; \ + fi .PHONY: test test: ## Run unit tests (Ginkgo, verbose) diff --git a/Taskfile.yml b/Taskfile.yml index c1fe6fa..4d3f0b4 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -24,13 +24,22 @@ tasks: - task: run-dev run-dev: - desc: Run the receiver in development mode (uses core NATS mode by default) + desc: 'Run the receiver in development mode (uses core NATS mode by default). Use MONITORING_ADDR=ip:port to set monitoring address' + vars: + MONITORING_ADDR: + sh: echo "${MONITORING_ADDR:-}" deps: - build cmds: - | echo "Running receiver in development mode (NATS mode: core by default)..." - - GO_ENV=dev {{.binary}} listen + - | + cmd="GO_ENV=dev {{.binary}} listen" + if [ -n "${MONITORING_ADDR:-}" ]; then + echo "Monitoring address: ${MONITORING_ADDR}" + cmd="$cmd --monitoring-addr ${MONITORING_ADDR}" + fi + eval "$cmd" run-prod: desc: Run the receiver in production mode (requires config.prod.toml and JetStream mode) @@ -55,12 +64,20 @@ tasks: - GO_ENV=test {{.binary}} listen run-local: - desc: Run receiver via go run (default GO_ENV=dev, uses core NATS mode by default) + desc: 'Run receiver via go run (default GO_ENV=dev, uses core NATS mode by default). Use MONITORING_ADDR=ip:port to set monitoring address' + vars: + MONITORING_ADDR: + sh: echo "${MONITORING_ADDR:-}" cmds: - | echo "Running receiver via go run (GO_ENV=${GO_ENV:-dev}, NATS mode: core by default in dev)..." - | - GO_ENV=${GO_ENV:-dev} go run {{.cmd}} listen + cmd="GO_ENV=${GO_ENV:-dev} go run {{.cmd}} listen" + if [ -n "${MONITORING_ADDR:-}" ]; then + echo "Monitoring address: ${MONITORING_ADDR}" + cmd="$cmd --monitoring-addr ${MONITORING_ADDR}" + fi + eval "$cmd" test: desc: Run Ginkgo unit test suites (verbose) @@ -160,11 +177,21 @@ tasks: - rm -rf {{.build_dir}} coverage up: - desc: Start TimescaleDB + NATS dev stack (docker compose) + desc: 'Start TimescaleDB + NATS dev stack (docker compose). Use MONITORING_ADDR=ip:port to set Prometheus scrape target' + vars: + MONITORING_ADDR: + sh: echo "${MONITORING_ADDR:-host.docker.internal:2112}" cmds: - echo "Starting dev infrastructure..." - - 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 + - | + monitoring_addr="${MONITORING_ADDR:-host.docker.internal:2112}" + echo "Prometheus will scrape caatsm-receiver at: ${monitoring_addr}" + CAATSM_MONITORING_ADDR="${monitoring_addr}" \ + docker compose -f docker-compose.dev.yml up -d postgres nats nats-box nats-exporter + - | + monitoring_addr="${MONITORING_ADDR:-host.docker.internal:2112}" + CAATSM_MONITORING_ADDR="${monitoring_addr}" \ + docker compose -f docker-compose.dev.yml up -d otel-collector jaeger prometheus grafana down: desc: Stop dev compose stack and remove containers @@ -173,14 +200,27 @@ tasks: - docker compose -f docker-compose.dev.yml down -v dev-run-compose: - desc: Run receiver inside docker-compose dev stack + desc: 'Run receiver inside docker-compose dev stack. Use MONITORING_ADDR=ip:port to set Prometheus scrape target' + vars: + MONITORING_ADDR: + sh: echo "${MONITORING_ADDR:-host.docker.internal:2112}" 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 + - | + monitoring_addr="${MONITORING_ADDR:-host.docker.internal:2112}" + echo "Prometheus will scrape caatsm-receiver at: ${monitoring_addr}" + CAATSM_MONITORING_ADDR="${monitoring_addr}" \ + docker compose -f docker-compose.dev.yml up -d postgres nats nats-box nats-exporter + - | + monitoring_addr="${MONITORING_ADDR:-host.docker.internal:2112}" + CAATSM_MONITORING_ADDR="${monitoring_addr}" \ + 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 + desc: 'Run receiver locally against dev stack. Use MONITORING_ADDR=ip:port to set monitoring address (default: 0.0.0.0:2112)' + vars: + MONITORING_ADDR: + sh: echo "${MONITORING_ADDR:-0.0.0.0:2112}" deps: - up env: @@ -190,7 +230,6 @@ 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 @@ -198,9 +237,12 @@ tasks: - | echo "Ensuring observability stack is healthy (Jaeger @ http://localhost:16686)" docker compose -f docker-compose.dev.yml ps jaeger otel-collector >/dev/null - echo "Running receiver with telemetry (mode=${CAATSM_NATS_MODE:-core})" + monitoring_addr="${MONITORING_ADDR:-0.0.0.0:2112}" + echo "Running receiver with telemetry (mode=${CAATSM_NATS_MODE:-core}, monitoring=${monitoring_addr})" + echo "Note: Ensure Prometheus is configured to scrape ${monitoring_addr}" + echo " If using a different IP, restart Prometheus with: MONITORING_ADDR=${monitoring_addr} task up" CAATSM_NATS_MODE=${CAATSM_NATS_MODE:-core} \ - go run ./cmd/main listen + go run ./cmd/main listen --monitoring-addr "${monitoring_addr}" seed: desc: Generate sample telegrams (publishes to Core NATS by default, compatible with dev mode) diff --git a/cmd/main/main.go b/cmd/main/main.go index 5305807..50d87e5 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -110,6 +110,16 @@ func setupApp() *cli.App { Name: "telemetry-insecure", Usage: "Send OTLP data without TLS", }, + &cli.StringFlag{ + Name: "monitoring-addr", + Usage: "Monitoring server address (e.g., 192.168.1.100:2112 or :2112)", + EnvVars: []string{"CAATSM_MONITORING_ADDR"}, + }, + &cli.BoolFlag{ + Name: "monitoring-disabled", + Usage: "Disable monitoring server", + EnvVars: []string{"CAATSM_MONITORING_DISABLED"}, + }, }, Action: executeListen, }, @@ -263,6 +273,13 @@ func applyCLIOverrides(cfg *config.Config, c *cli.Context) { if c.IsSet("telemetry-insecure") { cfg.Telemetry.Insecure = c.Bool("telemetry-insecure") } + if addr := c.String("monitoring-addr"); addr != "" { + cfg.Monitoring.Addr = addr + cfg.Monitoring.Disabled = false + } + if c.IsSet("monitoring-disabled") { + cfg.Monitoring.Disabled = c.Bool("monitoring-disabled") + } } func applyReplayOverride(cfg *config.Config, value string) { diff --git a/configs/config.dev.toml b/configs/config.dev.toml index f31b8af..c709867 100644 --- a/configs/config.dev.toml +++ b/configs/config.dev.toml @@ -77,7 +77,7 @@ close = "10s" ack_wait = "5s" [postgres] -url = "postgres://user:password@localhost:5432/aviation?sslmode=disable" +url = "postgres://caatsm:caatsm@localhost:5432/aviation?sslmode=disable" max_conns = 10 min_conns = 2 diff --git a/configs/otel-collector.dev.yaml b/configs/otel-collector.dev.yaml index 387f108..360e050 100644 --- a/configs/otel-collector.dev.yaml +++ b/configs/otel-collector.dev.yaml @@ -9,12 +9,12 @@ receivers: exporters: logging: loglevel: info - otlp/jaeger: - endpoint: jaeger:4317 + otlphttp/jaeger: + endpoint: http://jaeger:4318 tls: insecure: true prometheus: - endpoint: "0.0.0.0:8888" + endpoint: "0.0.0.0:8889" const_labels: source: "otel-collector" @@ -22,7 +22,7 @@ service: pipelines: traces: receivers: [otlp] - exporters: [logging, otlp/jaeger] + exporters: [logging, otlphttp/jaeger] metrics: receivers: [otlp] exporters: [logging, prometheus] diff --git a/configs/prometheus.dev.yml b/configs/prometheus.dev.yml index 858947f..0422107 100644 --- a/configs/prometheus.dev.yml +++ b/configs/prometheus.dev.yml @@ -6,14 +6,19 @@ scrape_configs: - job_name: "otel-collector" static_configs: - targets: - - "otel-collector:8888" + - "otel-collector:8889" - job_name: "nats-exporter" static_configs: - targets: - "nats-exporter:7777" - job_name: "caatsm-receiver" - static_configs: - - targets: - - 172.23.189.12:2112 + file_sd_configs: + # CAATSM_MONITORING_ADDR is set via environment variable in docker-compose + # Default: host.docker.internal:2112 (for Docker Desktop on Mac/Windows) + # Can be overridden: e.g., 10.16.66.236:2112 or 172.17.0.1:2112 + # Target file is dynamically generated in docker-compose entrypoint + - files: + - "/etc/prometheus/targets/caatsm-receiver.json" + refresh_interval: 5s diff --git a/configs/prometheus/targets/caatsm-receiver.json b/configs/prometheus/targets/caatsm-receiver.json new file mode 100644 index 0000000..c60b698 --- /dev/null +++ b/configs/prometheus/targets/caatsm-receiver.json @@ -0,0 +1 @@ +[{"labels":{"job":"caatsm-receiver","instance":"10.16.66.236:2112"},"targets":["10.16.66.236:2112"]}] diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 81f5585..ac9227c 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -133,15 +133,26 @@ services: 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" + entrypoint: + - /bin/sh + - -c + - | + MONITORING_ADDR=$${CAATSM_MONITORING_ADDR:-host.docker.internal:2112} + mkdir -p /etc/prometheus/targets + echo "[{\"labels\":{\"job\":\"caatsm-receiver\",\"instance\":\"$$MONITORING_ADDR\"},\"targets\":[\"$$MONITORING_ADDR\"]}]" > /etc/prometheus/targets/caatsm-receiver.json + exec /bin/prometheus \ + --config.file=/etc/prometheus/prometheus.yml \ + --storage.tsdb.path=/prometheus \ + --web.enable-lifecycle \ + --storage.tsdb.retention.time=1h + environment: + # Default monitoring address, can be overridden via .env file or docker-compose override + CAATSM_MONITORING_ADDR: ${CAATSM_MONITORING_ADDR:-host.docker.internal:2112} ports: - "9090:9090" volumes: - ./configs/prometheus.dev.yml:/etc/prometheus/prometheus.yml:ro + - ./configs/prometheus/targets:/etc/prometheus/targets networks: - devnet diff --git a/internal/adapter/parser/aviation.go b/internal/adapter/parser/aviation.go index 7cb8a83..70eb58f 100644 --- a/internal/adapter/parser/aviation.go +++ b/internal/adapter/parser/aviation.go @@ -155,7 +155,7 @@ func (parser *BodyParser) createBodyData(data map[string]string) (string, interf }, nil case CategoryCancellation: return category, &domain.CNL{ - Category: data[category], + Category: data[Category], AircraftID: data[FlightNumber], DepartureAirport: data[DepartureCode], DestinationAirport: data[ArrivalCode], @@ -358,12 +358,27 @@ func parseRemainingLines(lines []string) (string, string, string, string) { switch { case line == EndHeaderMarker: case strings.HasPrefix(line, "."): + // Validate if dot-prefixed line matches originator format: .ORIGINATOR_CODE YYMMDD + // Originator code should be uppercase letters, date/time should be digits originatorInfo := strings.Fields(line[1:]) if len(originatorInfo) >= 2 { - originator = originatorInfo[0] - originatorDateTime = originatorInfo[1] + // Check if first token is all uppercase letters and second is all digits + firstToken := originatorInfo[0] + secondToken := originatorInfo[1] + if isAllUppercaseLetters(firstToken) && isAllDigits(secondToken) { + originator = firstToken + originatorDateTime = secondToken + headerEnded = true + } else { + // Doesn't match originator format, treat as body content + headerEnded = true + bodyAndFooter.WriteString(line + "\n") + } + } else { + // Not enough tokens for originator format, treat as body content + headerEnded = true + bodyAndFooter.WriteString(line + "\n") } - headerEnded = true case strings.HasPrefix(line, BeginPartMarker) || strings.HasPrefix(line, "("): headerEnded = true if strings.Index(line, "NNNN") > 0 { @@ -393,6 +408,32 @@ func getOriginator(line string) (string, string) { return "", "" } +// isAllUppercaseLetters checks if a string contains only uppercase letters +func isAllUppercaseLetters(s string) bool { + if len(s) == 0 { + return false + } + for _, r := range s { + if r < 'A' || r > 'Z' { + return false + } + } + return true +} + +// isAllDigits checks if a string contains only digits +func isAllDigits(s string) bool { + if len(s) == 0 { + return false + } + for _, r := range s { + if r < '0' || r > '9' { + return false + } + } + return true +} + func parseOther(text string) map[string]string { data := make(map[string]string) for _, re := range otherPatterns { diff --git a/internal/adapter/parser/aviation_parser_test.go b/internal/adapter/parser/aviation_parser_test.go index 5d5e517..d785b9c 100644 --- a/internal/adapter/parser/aviation_parser_test.go +++ b/internal/adapter/parser/aviation_parser_test.go @@ -243,6 +243,7 @@ NNNN` Expect(category).To(Equal("CNL")) Expect(parsedBody).To(BeAssignableToTypeOf(&domain.CNL{})) cnlMessage := parsedBody.(*domain.CNL) + Expect(cnlMessage.Category).To(Equal("CNL")) Expect(cnlMessage.AircraftID).To(Equal("YZR7979")) }) }) diff --git a/internal/adapter/parser/constants.go b/internal/adapter/parser/constants.go index 6102520..16fc4cc 100644 --- a/internal/adapter/parser/constants.go +++ b/internal/adapter/parser/constants.go @@ -69,4 +69,5 @@ var ( eetPattern = regexp.MustCompile(`(?s)(-?EET\/(?P(?:[A-Z]{4}\d{4}\s*)+))`) performancePattern = regexp.MustCompile(`(?s)-?PER\/(?P\w)`) reroutePattern = regexp.MustCompile(`(?m)RIF\/(?P.*)[A-Z]{3}\/`) + cancelledPattern = regexp.MustCompile(`\bCNL\b`) ) diff --git a/internal/adapter/parser/schedule.go b/internal/adapter/parser/schedule.go index 53a8d30..e4a2984 100644 --- a/internal/adapter/parser/schedule.go +++ b/internal/adapter/parser/schedule.go @@ -53,7 +53,7 @@ func ParseWithDef(line string, parserDef *LineParser) *domain.ScheduleLine { var flightSchedule = &domain.ScheduleLine{ Reference: line, } - if strings.Contains(line, CANCELLED) { + if cancelledPattern.MatchString(line) { flightSchedule.Comments = "Cancelled" return flightSchedule } @@ -88,7 +88,16 @@ func ParseWithDef(line string, parserDef *LineParser) *domain.ScheduleLine { } } if len(words) > parserDef.WaypointStart { - flightSchedule.Waypoints, _ = parseWaypoints(words[parserDef.WaypointStart:]) + waypoints, err := parseWaypoints(words[parserDef.WaypointStart:]) + if err != nil { + log.Warnf("Failed to parse waypoints: %v", err) + if flightSchedule.Comments != "" { + flightSchedule.Comments += "; " + } + flightSchedule.Comments += "Waypoint parsing error: " + err.Error() + } else { + flightSchedule.Waypoints = waypoints + } } else { log.Warn("No waypoints found") flightSchedule.Comments = "No waypoints found" @@ -127,6 +136,11 @@ func ParseLine(line string) (*domain.ScheduleLine, error) { flightSchedule.Comments = "No waypoints found" } + // Validate the parsed schedule line before returning + if err := flightSchedule.Validate(); err != nil { + return nil, err + } + return flightSchedule, nil }