package main import ( "context" "encoding/json" "errors" "strings" "sync" "sync/atomic" "testing" "time" "github.com/nats-io/nats.go" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func TestBridgeBatchWritesKafkaThenAcks(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} payload, err := json.Marshal(env) if err != nil { t.Fatal(err) } writer := &recordingBridgeWriter{} acked := 0 err = bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.raw.jt808.v1", data: payload, ack: func() error { acked++ return nil }}, }, map[string]string{"vehicle.raw.jt808.v1": "vehicle.raw.jt808.v1"}) if err != nil { t.Fatalf("bridgeBatch() error = %v", err) } if len(writer.messages) != 1 { t.Fatalf("kafka writes = %d, want 1", len(writer.messages)) } if got, want := writer.messages[0].Topic, "vehicle.raw.jt808.v1"; got != want { t.Fatalf("topic = %q, want %q", got, want) } if got, want := string(writer.messages[0].Key), "JT808:13307795425"; got != want { t.Fatalf("key = %q, want %q", got, want) } if acked != 1 { t.Fatalf("acks = %d, want 1", acked) } } func TestBridgeBatchDoesNotAckWhenKafkaFails(t *testing.T) { writer := &recordingBridgeWriter{err: errors.New("kafka unavailable")} acked := 0 err := bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.event.unified.v1", data: []byte(`{"protocol":"JT808"}`), ack: func() error { acked++ return nil }}, }, map[string]string{"vehicle.event.unified.v1": "vehicle.event.unified.v1"}) if err == nil { t.Fatal("bridgeBatch() error is nil, want kafka error") } if acked != 0 { t.Fatalf("acks = %d, want 0", acked) } } func TestBridgeBatchProjectsFieldsFromCanonicalRawAndAcksOnce(t *testing.T) { raw := envelope.FrameEnvelope{ EventID: "raw-event-1", EventKind: envelope.EventKindRaw, Protocol: envelope.ProtocolJT808, MessageID: "0x0200", VIN: "LTEST000000000001", EventTimeMS: 1_700_000_000_000, ReceivedAtMS: 1_700_000_000_100, ParseStatus: envelope.ParseOK, ParsedFields: map[string]any{ "jt808.location.speed_kmh": 52.3, "jt808.location.total_mileage_km": 12345.6, }, } payload, err := raw.MarshalJSONBytes() if err != nil { t.Fatal(err) } acked := 0 writer := &recordingBridgeWriter{} registry := metrics.NewRegistry() err = bridgeBatchWithProjection( context.Background(), registry, writer, []bridgeMessage{{subject: topics.RawJT808, data: payload, ack: func() error { acked++; return nil }}}, map[string]string{topics.RawJT808: topics.RawJT808}, map[string]fieldsProjectionRoute{topics.RawJT808: {Protocol: envelope.ProtocolJT808, Topic: topics.FieldsJT808}}, true, ) if err != nil { t.Fatalf("bridgeBatchWithProjection() error = %v", err) } if acked != 1 { t.Fatalf("acks = %d, want exactly one ack for raw plus derived fields", acked) } if len(writer.messages) != 2 { t.Fatalf("kafka messages = %d, want raw and fields", len(writer.messages)) } byTopic := map[string]kafka.Message{} for _, message := range writer.messages { byTopic[message.Topic] = message } if len(byTopic[topics.RawJT808].Value) == 0 || len(byTopic[topics.FieldsJT808].Value) == 0 { t.Fatalf("projected topics = %#v", byTopic) } var fields envelope.FrameEnvelope if err := json.Unmarshal(byTopic[topics.FieldsJT808].Value, &fields); err != nil { t.Fatal(err) } if fields.EventKind != envelope.EventKindFields || fields.SourceEventID != raw.EventID { t.Fatalf("fields envelope = %#v", fields) } if got := fields.Fields["jt808.location.total_mileage_km"]; got != 12345.6 { t.Fatalf("projected mileage = %#v", got) } metricText := registry.Render() for _, want := range []string{ `vehicle_bridge_fields_projection_total{protocol="JT808",status="published"} 1`, `vehicle_bridge_fields_projection_count{protocol="JT808",status="published"} 2`, `vehicle_bridge_nats_acks_total{status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, } { if !strings.Contains(metricText, want) { t.Fatalf("projection metric missing %s:\n%s", want, metricText) } } } func TestBridgeBatchDoesNotAckRawWhenDerivedFieldsWriteFails(t *testing.T) { raw := envelope.FrameEnvelope{ EventKind: envelope.EventKindRaw, Protocol: envelope.ProtocolJT808, MessageID: "0x0200", VIN: "LTEST000000000001", EventTimeMS: 1_700_000_000_000, ReceivedAtMS: 1_700_000_000_100, ParseStatus: envelope.ParseOK, ParsedFields: map[string]any{"jt808.location.total_mileage_km": 12345.6}, } payload, err := raw.MarshalJSONBytes() if err != nil { t.Fatal(err) } acked := 0 writer := &recordingBridgeWriter{topicErr: map[string]error{topics.FieldsJT808: errors.New("fields unavailable")}} err = bridgeBatchWithProjection( context.Background(), nil, writer, []bridgeMessage{{subject: topics.RawJT808, data: payload, ack: func() error { acked++; return nil }}}, map[string]string{topics.RawJT808: topics.RawJT808}, map[string]fieldsProjectionRoute{topics.RawJT808: {Protocol: envelope.ProtocolJT808, Topic: topics.FieldsJT808}}, true, ) if err == nil { t.Fatal("bridgeBatchWithProjection() error = nil, want fields Kafka error") } if acked != 0 { t.Fatalf("acks = %d, want raw left pending until both Kafka outputs succeed", acked) } if len(writer.messages) != 1 || writer.messages[0].Topic != topics.RawJT808 { t.Fatalf("successful kafka messages = %#v, want only raw before replay", writer.messages) } } func TestBridgeBatchSkipsFieldsProjectionForNonRealtimeRaw(t *testing.T) { raw := envelope.FrameEnvelope{ EventKind: envelope.EventKindRaw, Protocol: envelope.ProtocolJT808, MessageID: "0x0100", Phone: "13307795425", EventTimeMS: 1_700_000_000_000, ReceivedAtMS: 1_700_000_000_100, ParseStatus: envelope.ParseOK, } payload, err := raw.MarshalJSONBytes() if err != nil { t.Fatal(err) } acked := 0 writer := &recordingBridgeWriter{} registry := metrics.NewRegistry() err = bridgeBatchWithProjection( context.Background(), registry, writer, []bridgeMessage{{subject: topics.RawJT808, data: payload, ack: func() error { acked++; return nil }}}, map[string]string{topics.RawJT808: topics.RawJT808}, map[string]fieldsProjectionRoute{topics.RawJT808: {Protocol: envelope.ProtocolJT808, Topic: topics.FieldsJT808}}, true, ) if err != nil { t.Fatalf("bridgeBatchWithProjection() error = %v", err) } if acked != 1 || len(writer.messages) != 1 || writer.messages[0].Topic != topics.RawJT808 { t.Fatalf("acks=%d messages=%#v", acked, writer.messages) } if text := registry.Render(); !strings.Contains(text, `vehicle_bridge_fields_projection_total{protocol="JT808",status="skipped_non_realtime"} 1`) { t.Fatalf("non-realtime projection metric missing:\n%s", text) } } func TestBridgeBatchAcksSuccessfulTopicWhenAnotherTopicFails(t *testing.T) { rawEnv := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} fieldsEnv := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} rawPayload, err := json.Marshal(rawEnv) if err != nil { t.Fatal(err) } fieldsPayload, err := json.Marshal(fieldsEnv) if err != nil { t.Fatal(err) } registry := metrics.NewRegistry() writer := &recordingBridgeWriter{ topicErr: map[string]error{ "vehicle.fields.go.jt808.v1": errors.New("fields topic unavailable"), }, } acks := map[string]int{} err = bridgeBatch(context.Background(), registry, writer, []bridgeMessage{ {subject: "vehicle.fields.go.jt808.v1", data: fieldsPayload, ack: func() error { acks["fields"]++ return nil }}, {subject: "vehicle.raw.go.jt808.v1", data: rawPayload, ack: func() error { acks["raw"]++ return nil }}, }, map[string]string{ "vehicle.fields.go.jt808.v1": "vehicle.fields.go.jt808.v1", "vehicle.raw.go.jt808.v1": "vehicle.raw.go.jt808.v1", }) if err == nil { t.Fatal("bridgeBatch() error = nil, want failed fields topic") } if len(writer.calls) != 2 { t.Fatalf("writer calls = %d, want one call per topic", len(writer.calls)) } if len(writer.messages) != 1 || writer.messages[0].Topic != "vehicle.raw.go.jt808.v1" { t.Fatalf("successful kafka messages = %#v, want only raw topic", writer.messages) } if acks["raw"] != 1 || acks["fields"] != 0 { t.Fatalf("acks = %#v, want raw acked and fields left unacked", acks) } text := registry.Render() for _, want := range []string{ `vehicle_bridge_kafka_writes_total{status="error",topic="vehicle.fields.go.jt808.v1"} 1`, `vehicle_bridge_kafka_writes_total{status="ok",topic="vehicle.raw.go.jt808.v1"} 1`, `vehicle_bridge_nats_acks_total{status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, } { if !strings.Contains(text, want) { t.Fatalf("partial bridge metric missing %s:\n%s", want, text) } } if strings.Contains(text, `vehicle_bridge_nats_acks_total{status="ok",subject="vehicle.fields.go.jt808.v1"}`) { t.Fatalf("failed topic should not be acked:\n%s", text) } } func TestBridgeBatchAcksAndDropsUnroutedSubject(t *testing.T) { writer := &recordingBridgeWriter{} acked := 0 registry := metrics.NewRegistry() err := bridgeBatch(context.Background(), registry, writer, []bridgeMessage{ {subject: "vehicle.unconfigured.v1", data: []byte(`{"protocol":"JT808"}`), ack: func() error { acked++ return nil }}, }, map[string]string{"vehicle.raw.go.jt808.v1": "vehicle.raw.go.jt808.v1"}) if err != nil { t.Fatalf("bridgeBatch() error = %v", err) } if acked != 1 { t.Fatalf("acks = %d, want unrouted message acked", acked) } if len(writer.messages) != 0 { t.Fatalf("kafka writes = %d, want 0", len(writer.messages)) } text := registry.Render() for _, want := range []string{ `vehicle_bridge_messages_total{status="route_error",subject="vehicle.unconfigured.v1"} 1`, `vehicle_bridge_nats_acks_total{status="dropped_route_error",subject="vehicle.unconfigured.v1"} 1`, `vehicle_bridge_batch_pending_kafka_messages 0`, } { if !strings.Contains(text, want) { t.Fatalf("metrics missing %s:\n%s", want, text) } } } func TestBridgeBatchKeepsValidMessagesWhenUnroutedSubjectIsPresent(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} payload, err := json.Marshal(env) if err != nil { t.Fatal(err) } writer := &recordingBridgeWriter{} acked := map[string]int{} err = bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.unconfigured.v1", data: []byte(`{"protocol":"JT808"}`), ack: func() error { acked["unknown"]++ return nil }}, {subject: "vehicle.raw.go.jt808.v1", data: payload, ack: func() error { acked["valid"]++ return nil }}, }, map[string]string{"vehicle.raw.go.jt808.v1": "vehicle.raw.go.jt808.v1"}) if err != nil { t.Fatalf("bridgeBatch() error = %v", err) } if len(writer.messages) != 1 { t.Fatalf("kafka writes = %d, want 1 valid message", len(writer.messages)) } if writer.messages[0].Topic != "vehicle.raw.go.jt808.v1" { t.Fatalf("topic = %q", writer.messages[0].Topic) } if acked["unknown"] != 1 || acked["valid"] != 1 { t.Fatalf("acks = %#v, want both messages acked", acked) } } func TestBridgeBatchReturnsErrorWhenUnroutedAckFails(t *testing.T) { writer := &recordingBridgeWriter{} err := bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.unconfigured.v1", data: []byte(`{"protocol":"JT808"}`), ack: func() error { return errors.New("nats ack unavailable") }}, }, map[string]string{"vehicle.raw.go.jt808.v1": "vehicle.raw.go.jt808.v1"}) if err == nil { t.Fatal("bridgeBatch() error = nil, want ack error") } if len(writer.messages) != 0 { t.Fatalf("kafka writes = %d, want 0", len(writer.messages)) } } func TestBridgeBatchRecordsMetrics(t *testing.T) { env := envelope.FrameEnvelope{ Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200", ReceivedAtMS: time.Now().Add(-20 * time.Millisecond).UnixMilli(), } payload, err := json.Marshal(env) if err != nil { t.Fatal(err) } registry := metrics.NewRegistry() writer := &recordingBridgeWriter{} err = bridgeBatch(context.Background(), registry, writer, []bridgeMessage{ {subject: "vehicle.raw.jt808.v1", data: payload, ack: func() error { return nil }}, }, map[string]string{"vehicle.raw.jt808.v1": "vehicle.raw.jt808.v1"}) if err != nil { t.Fatalf("bridgeBatch() error = %v", err) } text := registry.Render() for _, want := range []string{ `vehicle_bridge_messages_total{status="received",subject="vehicle.raw.jt808.v1"} 1`, `vehicle_bridge_last_message_unix_seconds{status="received",subject="vehicle.raw.jt808.v1"} `, `vehicle_bridge_kafka_writes_total{status="ok",topic="vehicle.raw.jt808.v1"} 1`, `vehicle_bridge_last_kafka_write_unix_seconds{status="ok",topic="vehicle.raw.jt808.v1"} `, `vehicle_bridge_kafka_write_duration_ms_histogram_count{status="ok",topic="vehicle.raw.jt808.v1"} 1`, `vehicle_bridge_kafka_e2e_duration_ms_histogram_count{topic="vehicle.raw.jt808.v1"} 1`, `vehicle_bridge_kafka_e2e_recent_p99_ms{topic="vehicle.raw.jt808.v1"} `, `vehicle_bridge_kafka_e2e_recent_samples{topic="vehicle.raw.jt808.v1"} `, `vehicle_bridge_last_kafka_e2e_unix_seconds{topic="vehicle.raw.jt808.v1"} `, `vehicle_bridge_nats_acks_total{status="ok",subject="vehicle.raw.jt808.v1"} 1`, `vehicle_bridge_last_ack_unix_seconds{status="ok",subject="vehicle.raw.jt808.v1"} `, } { if !strings.Contains(text, want) { t.Fatalf("metrics missing %s:\n%s", want, text) } } } func TestRecordBridgeKafkaE2EDurationSkipsMissingReceiveTime(t *testing.T) { registry := metrics.NewRegistry() recordBridgeKafkaE2EDuration(registry, "vehicle.raw.go.jt808.v1", 0) if text := registry.Render(); strings.Contains(text, "vehicle_bridge_kafka_e2e_duration_ms_histogram") || strings.Contains(text, "vehicle_bridge_kafka_e2e_recent") { t.Fatalf("e2e metric should be skipped when received_at_ms is missing:\n%s", text) } } func TestBridgeBatchExposesPendingAndDurationMetrics(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", MessageID: "0x02"} payload, err := json.Marshal(env) if err != nil { t.Fatal(err) } registry := metrics.NewRegistry() writer := &recordingBridgeWriter{ onWrite: func() { text := registry.Render() for _, want := range []string{ `vehicle_bridge_batch_pending_messages 2`, `vehicle_bridge_batch_pending_kafka_messages 2`, } { if !strings.Contains(text, want) { t.Fatalf("pending bridge metric missing %s during write:\n%s", want, text) } } }, } err = bridgeBatch(context.Background(), registry, writer, []bridgeMessage{ {subject: "vehicle.raw.go.gb32960.v1", data: payload, ack: func() error { return nil }}, {subject: "vehicle.raw.go.gb32960.v1", data: payload, ack: func() error { return nil }}, }, map[string]string{"vehicle.raw.go.gb32960.v1": "vehicle.raw.go.gb32960.v1"}) if err != nil { t.Fatalf("bridgeBatch() error = %v", err) } text := registry.Render() for _, want := range []string{ `vehicle_bridge_batch_pending_messages 0`, `vehicle_bridge_batch_pending_kafka_messages 0`, `vehicle_bridge_batch_duration_ms_histogram_bucket{le="+Inf",status="ok"} 1`, `vehicle_bridge_batch_duration_ms_histogram_count{status="ok"} 1`, `vehicle_bridge_batch_duration_ms_histogram_sum{status="ok"}`, } { if !strings.Contains(text, want) { t.Fatalf("bridge batch metric missing %s:\n%s", want, text) } } } func TestBridgeBatchPendingAggregatesConcurrentWorkers(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", MessageID: "0x02"} payload, err := json.Marshal(env) if err != nil { t.Fatal(err) } registry := metrics.NewRegistry() firstStarted := make(chan struct{}) secondStarted := make(chan struct{}) release := make(chan struct{}) writerFor := func(started chan struct{}) *recordingBridgeWriter { return &recordingBridgeWriter{ onWrite: func() { close(started) <-release }, } } messages := func(count int) []bridgeMessage { out := make([]bridgeMessage, 0, count) for i := 0; i < count; i++ { out = append(out, bridgeMessage{subject: "vehicle.raw.go.gb32960.v1", data: payload, ack: func() error { return nil }}) } return out } route := map[string]string{"vehicle.raw.go.gb32960.v1": "vehicle.raw.go.gb32960.v1"} var wg sync.WaitGroup wg.Add(2) var firstErr error var secondErr error go func() { defer wg.Done() firstErr = bridgeBatch(context.Background(), registry, writerFor(firstStarted), messages(2), route) }() go func() { defer wg.Done() secondErr = bridgeBatch(context.Background(), registry, writerFor(secondStarted), messages(3), route) }() <-firstStarted <-secondStarted text := registry.Render() for _, want := range []string{ `vehicle_bridge_batch_pending_messages 5`, `vehicle_bridge_batch_pending_kafka_messages 5`, } { if !strings.Contains(text, want) { t.Fatalf("aggregate pending bridge metric missing %s during concurrent writes:\n%s", want, text) } } close(release) wg.Wait() if firstErr != nil || secondErr != nil { t.Fatalf("bridgeBatch errors = %v / %v", firstErr, secondErr) } text = registry.Render() for _, want := range []string{ `vehicle_bridge_batch_pending_messages 0`, `vehicle_bridge_batch_pending_kafka_messages 0`, } { if !strings.Contains(text, want) { t.Fatalf("aggregate pending bridge metric should reset after concurrent writes, missing %s:\n%s", want, text) } } } func TestRecordNATSConsumerInfoMetrics(t *testing.T) { registry := metrics.NewRegistry() cfg := config{NATSStream: "VEHICLE_INGEST", NATSDurable: "vehicle-kafka-bridge"} recordNATSConsumerInfoMetrics(registry, cfg, &nats.ConsumerInfo{ NumPending: 17, NumAckPending: 3, NumWaiting: 2, }) text := registry.Render() for _, want := range []string{ `vehicle_bridge_nats_consumer_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 17`, `vehicle_bridge_nats_consumer_ack_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 3`, `vehicle_bridge_nats_consumer_waiting{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 2`, } { if !strings.Contains(text, want) { t.Fatalf("metrics missing %s:\n%s", want, text) } } } func TestIsBridgeShutdownFetchError(t *testing.T) { if isBridgeShutdownFetchError(context.Background(), errors.New("temporary nats failure")) { t.Fatal("temporary failure should not be treated as shutdown") } if !isBridgeShutdownFetchError(context.Background(), nats.ErrConnectionClosed) { t.Fatal("connection closed should stop worker without noisy error log") } ctx, cancel := context.WithCancel(context.Background()) cancel() if !isBridgeShutdownFetchError(ctx, errors.New("any fetch error after cancellation")) { t.Fatal("cancelled context should stop worker without noisy error log") } } func TestIsTransientBridgeFetchError(t *testing.T) { for _, err := range []error{ errors.New("nats: disconnected during fetch"), errors.New("nats: connection closed"), errors.New("read tcp: connection reset by peer"), errors.New("temporary network unavailable"), errors.New("i/o timeout"), } { if !isTransientBridgeFetchError(err) { t.Fatalf("isTransientBridgeFetchError(%v) = false, want true", err) } } if isTransientBridgeFetchError(errors.New("permission denied")) { t.Fatal("non-transient bridge fetch error should stay non-transient") } if isTransientBridgeFetchError(nil) { t.Fatal("nil should not be transient") } } func TestLoadConfigDefaultsToGoSubjectRoutes(t *testing.T) { cfg := loadConfig() if err := cfg.Validate(); err != nil { t.Fatalf("default config Validate() error = %v", err) } want := map[string]string{ "vehicle.raw.go.gb32960.v1": "vehicle.raw.go.gb32960.v1", "vehicle.raw.go.jt808.v1": "vehicle.raw.go.jt808.v1", "vehicle.raw.go.yutong-mqtt.v1": "vehicle.raw.go.yutong-mqtt.v1", "vehicle.fields.go.gb32960.v1": "vehicle.fields.go.gb32960.v1", "vehicle.fields.go.jt808.v1": "vehicle.fields.go.jt808.v1", "vehicle.fields.go.yutong-mqtt.v1": "vehicle.fields.go.yutong-mqtt.v1", } if len(cfg.Route) != len(want) { t.Fatalf("route len = %d, want %d: %#v", len(cfg.Route), len(want), cfg.Route) } for subject, topic := range want { if got := cfg.Route[subject]; got != topic { t.Fatalf("route[%q] = %q, want %q; route=%#v", subject, got, topic, cfg.Route) } } if !cfg.DeriveFieldsFromRaw { t.Fatal("DeriveFieldsFromRaw = false, want canonical raw projection enabled by default") } for subject, wantTopic := range map[string]string{ topics.RawGB32960: topics.FieldsGB32960, topics.RawJT808: topics.FieldsJT808, topics.RawYutongMQTT: topics.FieldsYutongMQTT, } { if got := cfg.RawFieldRoutes[subject].Topic; got != wantTopic { t.Fatalf("RawFieldRoutes[%q].Topic = %q, want %q", subject, got, wantTopic) } } } func TestConfigValidateRejectsRawFieldsSubjectOverlap(t *testing.T) { cfg := config{ RawRoutes: []subjectRoute{ {Subject: "vehicle.same.jt808", Topic: "vehicle.raw.go.jt808.v1"}, }, FieldsRoutes: []subjectRoute{ {Subject: "vehicle.same.jt808", Topic: "vehicle.fields.go.jt808.v1"}, }, } err := cfg.Validate() if err == nil { t.Fatal("Validate() error = nil, want subject overlap rejection") } if !strings.Contains(err.Error(), "nats subject") { t.Fatalf("Validate() error = %q, want nats subject hint", err) } } func TestConfigValidateRejectsKnownProtocolTopicMismatch(t *testing.T) { cfg := config{ RawRoutes: []subjectRoute{ {Protocol: envelope.ProtocolJT808, Subject: "vehicle.raw.go.jt808.v1", Topic: "vehicle.raw.go.gb32960.v1"}, }, } err := cfg.Validate() if err == nil { t.Fatal("Validate() error = nil, want known protocol topic mismatch") } if !strings.Contains(err.Error(), "must match protocol") { t.Fatalf("Validate() error = %q, want protocol mismatch hint", err) } } func TestConfigValidateRejectsFieldsRouteToRawKafkaTopic(t *testing.T) { cfg := config{ RawRoutes: []subjectRoute{ {Subject: "vehicle.raw.go.jt808.v1", Topic: "vehicle.raw.go.jt808.v1"}, }, FieldsRoutes: []subjectRoute{ {Subject: "vehicle.fields.go.jt808.v1", Topic: "vehicle.raw.go.jt808.v1"}, }, } err := cfg.Validate() if err == nil { t.Fatal("Validate() error = nil, want fields topic family rejection") } if !strings.Contains(err.Error(), "fields kafka topic") { t.Fatalf("Validate() error = %q, want fields kafka topic hint", err) } } func TestLoadConfigDefaultsStreamMaxBytes(t *testing.T) { cfg := loadConfig() if got, want := cfg.StreamMaxBytes, int64(20*1024*1024*1024); got != want { t.Fatalf("StreamMaxBytes = %d, want %d", got, want) } if got, want := cfg.StreamEnsureWait, 60*time.Second; got != want { t.Fatalf("StreamEnsureWait = %v, want %v", got, want) } if got, want := cfg.KafkaBatchTimeout, 20*time.Millisecond; got != want { t.Fatalf("KafkaBatchTimeout = %v, want %v", got, want) } if got, want := cfg.KafkaWriteConcurrency, 6; got != want { t.Fatalf("KafkaWriteConcurrency = %d, want %d", got, want) } if got, want := cfg.FetchWait, 20*time.Millisecond; got != want { t.Fatalf("FetchWait = %v, want %v", got, want) } if got, want := cfg.Workers, 4; got != want { t.Fatalf("Workers = %d, want %d", got, want) } } func TestLoadConfigReadsStreamMaxBytesOverride(t *testing.T) { t.Setenv("NATS_STREAM_MAX_BYTES", "1073741824") t.Setenv("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", "90") t.Setenv("BRIDGE_KAFKA_BATCH_TIMEOUT_MS", "35") t.Setenv("BRIDGE_KAFKA_WRITE_CONCURRENCY", "4") t.Setenv("BRIDGE_FETCH_WAIT_MS", "45") t.Setenv("BRIDGE_WORKERS", "6") cfg := loadConfig() if got, want := cfg.StreamMaxBytes, int64(1073741824); got != want { t.Fatalf("StreamMaxBytes = %d, want %d", got, want) } if got, want := cfg.StreamEnsureWait, 90*time.Second; got != want { t.Fatalf("StreamEnsureWait = %v, want %v", got, want) } if got, want := cfg.KafkaBatchTimeout, 35*time.Millisecond; got != want { t.Fatalf("KafkaBatchTimeout = %v, want %v", got, want) } if got, want := cfg.KafkaWriteConcurrency, 4; got != want { t.Fatalf("KafkaWriteConcurrency = %d, want %d", got, want) } if got, want := cfg.FetchWait, 45*time.Millisecond; got != want { t.Fatalf("FetchWait = %v, want %v", got, want) } if got, want := cfg.Workers, 6; got != want { t.Fatalf("Workers = %d, want %d", got, want) } } func TestRecordBridgeConfigMetrics(t *testing.T) { registry := metrics.NewRegistry() recordBridgeConfigMetrics(registry, config{ KafkaBatchTimeout: 35 * time.Millisecond, KafkaWriteConcurrency: 4, BatchSize: 600, FetchWait: 20 * time.Millisecond, Workers: 6, DeriveFieldsFromRaw: true, }) text := registry.Render() for _, want := range []string{ `vehicle_bridge_config{setting="batch_size"} 600`, `vehicle_bridge_config{setting="fetch_wait_ms"} 20`, `vehicle_bridge_config{setting="kafka_batch_timeout_ms"} 35`, `vehicle_bridge_config{setting="kafka_write_concurrency"} 4`, `vehicle_bridge_config{setting="workers"} 6`, `vehicle_bridge_config{setting="derive_fields_from_raw_enabled"} 1`, } { if !strings.Contains(text, want) { t.Fatalf("bridge config metric missing %s:\n%s", want, text) } } } func TestNewKafkaWriterUsesConfiguredBatchTimeout(t *testing.T) { writer := newKafkaWriter(config{ KafkaBrokers: []string{"127.0.0.1:9092"}, KafkaBatchTimeout: 17 * time.Millisecond, }) defer writer.Close() if got, want := writer.BatchTimeout, 17*time.Millisecond; got != want { t.Fatalf("BatchTimeout = %v, want %v", got, want) } } func TestLoadConfigIncludesUnifiedOnlyWhenExplicitlyConfigured(t *testing.T) { t.Setenv("NATS_SUBJECT_UNIFIED", "vehicle.event.go.unified.v1") t.Setenv("KAFKA_TOPIC_UNIFIED", "vehicle.event.go.unified.v1") cfg := loadConfig() if got := cfg.Route["vehicle.event.go.unified.v1"]; got != "vehicle.event.go.unified.v1" { t.Fatalf("unified route = %q, route=%#v", got, cfg.Route) } } type recordingBridgeWriter struct { mu sync.Mutex messages []kafka.Message calls [][]kafka.Message err error topicErr map[string]error onWrite func() } func (w *recordingBridgeWriter) WriteMessages(_ context.Context, messages ...kafka.Message) error { if w.onWrite != nil { w.onWrite() } w.mu.Lock() defer w.mu.Unlock() call := append([]kafka.Message(nil), messages...) w.calls = append(w.calls, call) if len(messages) > 0 && w.topicErr != nil { if err := w.topicErr[messages[0].Topic]; err != nil { return err } } if w.err != nil { return w.err } w.messages = append(w.messages, messages...) return nil } type concurrentBridgeWriter struct { started chan string release chan struct{} active atomic.Int32 max atomic.Int32 } func (w *concurrentBridgeWriter) WriteMessages(_ context.Context, messages ...kafka.Message) error { active := w.active.Add(1) defer w.active.Add(-1) for { current := w.max.Load() if active <= current || w.max.CompareAndSwap(current, active) { break } } w.started <- messages[0].Topic <-w.release return nil } func TestBridgeWritesIndependentKafkaTopicsConcurrently(t *testing.T) { writer := &concurrentBridgeWriter{ started: make(chan string, 2), release: make(chan struct{}), } messages := []bridgeMessage{ {subject: topics.RawJT808, data: []byte(`{"protocol":"JT808"}`), ack: func() error { return nil }}, {subject: topics.RawGB32960, data: []byte(`{"protocol":"GB32960"}`), ack: func() error { return nil }}, } route := map[string]string{ topics.RawJT808: topics.RawJT808, topics.RawGB32960: topics.RawGB32960, } done := make(chan error, 1) go func() { done <- bridgeBatchWithProjectionConcurrency(context.Background(), nil, writer, messages, route, nil, false, 2) }() for i := 0; i < 2; i++ { select { case <-writer.started: case <-time.After(time.Second): t.Fatal("independent Kafka topic writes did not start concurrently") } } if got := writer.max.Load(); got != 2 { t.Fatalf("max concurrent Kafka writes = %d, want 2", got) } close(writer.release) if err := <-done; err != nil { t.Fatalf("bridgeBatchWithProjectionConcurrency() error = %v", err) } }