From 2ae9794ddacb6c816193057e8471d08e84980a04 Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 20:22:47 +0800 Subject: [PATCH] refactor(go): default to go vehicle topics --- go/vehicle-gateway/cmd/gateway/main.go | 17 +++++++++-------- go/vehicle-gateway/cmd/gateway/main_test.go | 19 +++++++++++++++++++ go/vehicle-gateway/cmd/history-writer/main.go | 3 ++- .../cmd/history-writer/main_test.go | 10 ++++++++++ .../cmd/nats-kafka-bridge/main.go | 9 +++++---- .../cmd/nats-kafka-bridge/main_test.go | 19 +++++++++++++++++++ go/vehicle-gateway/cmd/realtime-api/main.go | 10 ++++++++-- .../cmd/realtime-api/main_test.go | 8 ++++++++ go/vehicle-gateway/cmd/stat-writer/main.go | 3 ++- .../cmd/stat-writer/main_test.go | 10 ++++++++++ .../internal/eventbus/kafka_sink.go | 9 +++++---- .../internal/eventbus/kafka_sink_test.go | 6 +++--- .../internal/eventbus/nats_sink.go | 9 +++++---- .../internal/eventbus/nats_sink_test.go | 17 +++++++++++++++++ go/vehicle-gateway/internal/topics/topics.go | 8 ++++++++ 15 files changed, 130 insertions(+), 27 deletions(-) create mode 100644 go/vehicle-gateway/internal/topics/topics.go diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 07e9dfcf..2061d539 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -22,6 +22,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/gb32960" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/jt808" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { @@ -192,11 +193,11 @@ func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) sink, err := eventbus.NewKafkaSink(eventbus.KafkaConfig{ Brokers: brokers, RawTopics: map[envelope.Protocol]string{ - envelope.ProtocolGB32960: env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1"), - envelope.ProtocolJT808: env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1"), - envelope.ProtocolYutongMQTT: env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1"), + envelope.ProtocolGB32960: env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960), + envelope.ProtocolJT808: env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808), + envelope.ProtocolYutongMQTT: env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT), }, - UnifiedTopic: env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1"), + UnifiedTopic: env("KAFKA_TOPIC_UNIFIED", topics.Unified), }) if err != nil { return nil, err @@ -239,11 +240,11 @@ func natsSinkConfigFromEnv() eventbus.NATSConfig { URL: env("NATS_URL", ""), Name: env("NATS_CLIENT_NAME", "lingniu-vehicle-gateway"), RawSubjects: map[envelope.Protocol]string{ - envelope.ProtocolGB32960: env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1")), - envelope.ProtocolJT808: env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1")), - envelope.ProtocolYutongMQTT: env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1")), + envelope.ProtocolGB32960: env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)), + envelope.ProtocolJT808: env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)), + envelope.ProtocolYutongMQTT: env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)), }, - UnifiedSubject: env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1")), + UnifiedSubject: env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", topics.Unified)), } } diff --git a/go/vehicle-gateway/cmd/gateway/main_test.go b/go/vehicle-gateway/cmd/gateway/main_test.go index 87378480..ff1e476e 100644 --- a/go/vehicle-gateway/cmd/gateway/main_test.go +++ b/go/vehicle-gateway/cmd/gateway/main_test.go @@ -31,3 +31,22 @@ func TestNATSSinkConfigFromEnvUsesExplicitSubjects(t *testing.T) { t.Fatalf("unified subject = %q, want %q", got, want) } } + +func TestNATSSinkConfigFromEnvDefaultsToGoSubjects(t *testing.T) { + t.Setenv("NATS_URL", "nats://172.17.111.56:4222") + + cfg := natsSinkConfigFromEnv() + + if got, want := cfg.RawSubjects[envelope.ProtocolGB32960], "vehicle.raw.go.gb32960.v1"; got != want { + t.Fatalf("gb32960 subject = %q, want %q", got, want) + } + if got, want := cfg.RawSubjects[envelope.ProtocolJT808], "vehicle.raw.go.jt808.v1"; got != want { + t.Fatalf("jt808 subject = %q, want %q", got, want) + } + if got, want := cfg.RawSubjects[envelope.ProtocolYutongMQTT], "vehicle.raw.go.yutong-mqtt.v1"; got != want { + t.Fatalf("yutong subject = %q, want %q", got, want) + } + if got, want := cfg.UnifiedSubject, "vehicle.event.go.unified.v1"; got != want { + t.Fatalf("unified subject = %q, want %q", got, want) + } +} diff --git a/go/vehicle-gateway/cmd/history-writer/main.go b/go/vehicle-gateway/cmd/history-writer/main.go index 543b055a..4c5f11ba 100644 --- a/go/vehicle-gateway/cmd/history-writer/main.go +++ b/go/vehicle-gateway/cmd/history-writer/main.go @@ -18,6 +18,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/history" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { @@ -143,7 +144,7 @@ type config struct { func loadConfig() config { return config{ KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")), - KafkaTopics: splitCSV(env("KAFKA_TOPICS", "vehicle.raw.gb32960.v1,vehicle.raw.jt808.v1,vehicle.raw.yutong-mqtt.v1")), + KafkaTopics: splitCSV(env("KAFKA_TOPICS", strings.Join([]string{topics.RawGB32960, topics.RawJT808, topics.RawYutongMQTT}, ","))), KafkaGroup: env("KAFKA_GROUP", "go-history-writer"), TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: env("TDENGINE_DSN", ""), diff --git a/go/vehicle-gateway/cmd/history-writer/main_test.go b/go/vehicle-gateway/cmd/history-writer/main_test.go index 2e4844dc..f6460f20 100644 --- a/go/vehicle-gateway/cmd/history-writer/main_test.go +++ b/go/vehicle-gateway/cmd/history-writer/main_test.go @@ -74,6 +74,16 @@ func TestProcessHistoryMessageRecordsMetrics(t *testing.T) { } } +func TestLoadConfigDefaultsToGoRawTopics(t *testing.T) { + cfg := loadConfig() + + got := strings.Join(cfg.KafkaTopics, ",") + want := "vehicle.raw.go.gb32960.v1,vehicle.raw.go.jt808.v1,vehicle.raw.go.yutong-mqtt.v1" + if got != want { + t.Fatalf("KafkaTopics = %q, want %q", got, want) + } +} + type contextCheckingHistoryAppender struct { ctxErr error count int diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go index a260705c..379c5046 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go @@ -20,6 +20,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { @@ -101,10 +102,10 @@ type config struct { func loadConfig() config { route := map[string]string{ - env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1")): env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1"), - env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1")): env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1"), - env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1")): env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1"), - env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1")): env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1"), + env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)): env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960), + env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)): env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808), + env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)): env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT), + env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", topics.Unified)): env("KAFKA_TOPIC_UNIFIED", topics.Unified), } return config{ NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go index 15ba57d5..d28e9a1c 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go @@ -113,6 +113,25 @@ func TestRecordNATSConsumerInfoMetrics(t *testing.T) { } } +func TestLoadConfigDefaultsToGoSubjectRoutes(t *testing.T) { + cfg := loadConfig() + + 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.event.go.unified.v1": "vehicle.event.go.unified.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) + } + } +} + type recordingBridgeWriter struct { messages []kafka.Message err error diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 2b99554c..ebfd4cad 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -24,6 +24,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { @@ -163,15 +164,16 @@ func consumeKafka(ctx context.Context, logger interface { Error(string, ...any) Warn(string, ...any) }, registry *metrics.Registry, updater realtimeUpdater, brokers []string) { + kafkaTopics := kafkaTopicsFromEnv() reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: brokers, GroupID: env("KAFKA_GROUP", "go-realtime-api"), - GroupTopics: splitCSV(env("KAFKA_TOPICS", "vehicle.event.unified.v1")), + GroupTopics: kafkaTopics, MinBytes: 1, MaxBytes: 10e6, }) defer reader.Close() - logger.Info("realtime kafka consumer started", "topics", env("KAFKA_TOPICS", "vehicle.event.unified.v1")) + logger.Info("realtime kafka consumer started", "topics", strings.Join(kafkaTopics, ",")) for { message, err := reader.FetchMessage(ctx) if err != nil { @@ -185,6 +187,10 @@ func consumeKafka(ctx context.Context, logger interface { } } +func kafkaTopicsFromEnv() []string { + return splitCSV(env("KAFKA_TOPICS", topics.Unified)) +} + const kafkaMessageOperationTimeout = 30 * time.Second type realtimeUpdater interface { diff --git a/go/vehicle-gateway/cmd/realtime-api/main_test.go b/go/vehicle-gateway/cmd/realtime-api/main_test.go index 782c19c9..cd99efce 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main_test.go +++ b/go/vehicle-gateway/cmd/realtime-api/main_test.go @@ -95,6 +95,14 @@ func TestCompositeRealtimeUpdaterUpdatesBothStores(t *testing.T) { } } +func TestKafkaTopicsFromEnvDefaultsToGoUnifiedTopic(t *testing.T) { + got := strings.Join(kafkaTopicsFromEnv(), ",") + want := "vehicle.event.go.unified.v1" + if got != want { + t.Fatalf("topics = %q, want %q", got, want) + } +} + type contextCheckingRealtimeUpdater struct { ctxErr error count int diff --git a/go/vehicle-gateway/cmd/stat-writer/main.go b/go/vehicle-gateway/cmd/stat-writer/main.go index dcead6fa..cb5c7d18 100644 --- a/go/vehicle-gateway/cmd/stat-writer/main.go +++ b/go/vehicle-gateway/cmd/stat-writer/main.go @@ -18,6 +18,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { @@ -142,7 +143,7 @@ func loadConfig() config { } return config{ KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")), - KafkaTopics: splitCSV(env("KAFKA_TOPICS", "vehicle.raw.gb32960.v1,vehicle.raw.jt808.v1")), + KafkaTopics: splitCSV(env("KAFKA_TOPICS", strings.Join([]string{topics.RawGB32960, topics.RawJT808, topics.RawYutongMQTT}, ","))), KafkaGroup: env("KAFKA_GROUP", "go-stat-writer"), MySQLDSN: env("MYSQL_DSN", ""), EnsureSchema: env("MYSQL_ENSURE_SCHEMA", "true") != "false", diff --git a/go/vehicle-gateway/cmd/stat-writer/main_test.go b/go/vehicle-gateway/cmd/stat-writer/main_test.go index a1845c86..945c14ba 100644 --- a/go/vehicle-gateway/cmd/stat-writer/main_test.go +++ b/go/vehicle-gateway/cmd/stat-writer/main_test.go @@ -74,6 +74,16 @@ func TestProcessStatMessageRecordsMetrics(t *testing.T) { } } +func TestLoadConfigDefaultsToGoRawTopics(t *testing.T) { + cfg := loadConfig() + + got := strings.Join(cfg.KafkaTopics, ",") + want := "vehicle.raw.go.gb32960.v1,vehicle.raw.go.jt808.v1,vehicle.raw.go.yutong-mqtt.v1" + if got != want { + t.Fatalf("KafkaTopics = %q, want %q", got, want) + } +} + type contextCheckingStatAppender struct { ctxErr error count int diff --git a/go/vehicle-gateway/internal/eventbus/kafka_sink.go b/go/vehicle-gateway/internal/eventbus/kafka_sink.go index 84856f4f..4f32f4d2 100644 --- a/go/vehicle-gateway/internal/eventbus/kafka_sink.go +++ b/go/vehicle-gateway/internal/eventbus/kafka_sink.go @@ -8,6 +8,7 @@ import ( "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) type Sink interface { @@ -48,9 +49,9 @@ func NewKafkaSink(cfg KafkaConfig) (*KafkaSink, error) { func newKafkaSinkWithWriter(writer kafkaWriter, cfg KafkaConfig) *KafkaSink { rawTopics := map[envelope.Protocol]string{ - envelope.ProtocolGB32960: "vehicle.raw.gb32960.v1", - envelope.ProtocolJT808: "vehicle.raw.jt808.v1", - envelope.ProtocolYutongMQTT: "vehicle.raw.yutong-mqtt.v1", + envelope.ProtocolGB32960: topics.RawGB32960, + envelope.ProtocolJT808: topics.RawJT808, + envelope.ProtocolYutongMQTT: topics.RawYutongMQTT, } for protocol, topic := range cfg.RawTopics { if topic != "" { @@ -59,7 +60,7 @@ func newKafkaSinkWithWriter(writer kafkaWriter, cfg KafkaConfig) *KafkaSink { } unifiedTopic := cfg.UnifiedTopic if unifiedTopic == "" { - unifiedTopic = "vehicle.event.unified.v1" + unifiedTopic = topics.Unified } return &KafkaSink{writer: writer, rawTopics: rawTopics, unifiedTopic: unifiedTopic} } diff --git a/go/vehicle-gateway/internal/eventbus/kafka_sink_test.go b/go/vehicle-gateway/internal/eventbus/kafka_sink_test.go index 6483b781..77087962 100644 --- a/go/vehicle-gateway/internal/eventbus/kafka_sink_test.go +++ b/go/vehicle-gateway/internal/eventbus/kafka_sink_test.go @@ -23,8 +23,8 @@ func TestKafkaSinkRoutesRawTopics(t *testing.T) { t.Fatalf("messages = %d", len(writer.messages)) } msg := writer.messages[0] - if msg.Topic != "vehicle.raw.jt808.v1" { - t.Fatalf("topic = %q", msg.Topic) + if got, want := msg.Topic, "vehicle.raw.go.jt808.v1"; got != want { + t.Fatalf("topic = %q, want %q", got, want) } if string(msg.Key) != "JT808:13307795425" { t.Fatalf("key = %q", string(msg.Key)) @@ -82,7 +82,7 @@ func TestKafkaSinkPublishesDurableRecordsInSingleWriterCall(t *testing.T) { if len(writer.messages) != 2 { t.Fatalf("messages = %d, want 2", len(writer.messages)) } - if writer.messages[0].Topic != "vehicle.raw.jt808.v1" || writer.messages[1].Topic != "custom.unified" { + if writer.messages[0].Topic != "vehicle.raw.go.jt808.v1" || writer.messages[1].Topic != "custom.unified" { t.Fatalf("topics = %q, %q", writer.messages[0].Topic, writer.messages[1].Topic) } } diff --git a/go/vehicle-gateway/internal/eventbus/nats_sink.go b/go/vehicle-gateway/internal/eventbus/nats_sink.go index 21d7153b..c0629dfa 100644 --- a/go/vehicle-gateway/internal/eventbus/nats_sink.go +++ b/go/vehicle-gateway/internal/eventbus/nats_sink.go @@ -9,6 +9,7 @@ import ( "github.com/nats-io/nats.go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) type NATSConfig struct { @@ -55,9 +56,9 @@ func NewNATSSink(cfg NATSConfig) (*NATSSink, error) { func newNATSSinkWithPublisher(publisher natsPublisher, cfg NATSConfig) *NATSSink { rawSubjects := map[envelope.Protocol]string{ - envelope.ProtocolGB32960: "vehicle.raw.gb32960.v1", - envelope.ProtocolJT808: "vehicle.raw.jt808.v1", - envelope.ProtocolYutongMQTT: "vehicle.raw.yutong-mqtt.v1", + envelope.ProtocolGB32960: topics.RawGB32960, + envelope.ProtocolJT808: topics.RawJT808, + envelope.ProtocolYutongMQTT: topics.RawYutongMQTT, } for protocol, subject := range cfg.RawSubjects { if subject != "" { @@ -66,7 +67,7 @@ func newNATSSinkWithPublisher(publisher natsPublisher, cfg NATSConfig) *NATSSink } unifiedSubject := cfg.UnifiedSubject if unifiedSubject == "" { - unifiedSubject = "vehicle.event.unified.v1" + unifiedSubject = topics.Unified } return &NATSSink{ publisher: publisher, diff --git a/go/vehicle-gateway/internal/eventbus/nats_sink_test.go b/go/vehicle-gateway/internal/eventbus/nats_sink_test.go index 36938d41..1eb501bb 100644 --- a/go/vehicle-gateway/internal/eventbus/nats_sink_test.go +++ b/go/vehicle-gateway/internal/eventbus/nats_sink_test.go @@ -40,6 +40,23 @@ func TestNATSSinkRoutesRawAndUnifiedSubjects(t *testing.T) { } } +func TestNATSSinkDefaultsToGoRawSubjects(t *testing.T) { + publisher := &recordingNATSPublisher{} + sink := newNATSSinkWithPublisher(publisher, NATSConfig{}) + + err := sink.PublishRaw(context.Background(), envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + Phone: "13307795425", + ReceivedAtMS: 1782918600000, + }) + if err != nil { + t.Fatalf("PublishRaw() error = %v", err) + } + if got, want := publisher.messages[0].subject, "vehicle.raw.go.jt808.v1"; got != want { + t.Fatalf("subject = %q, want %q", got, want) + } +} + type recordingNATSPublisher struct { messages []recordedNATSMessage } diff --git a/go/vehicle-gateway/internal/topics/topics.go b/go/vehicle-gateway/internal/topics/topics.go new file mode 100644 index 00000000..8d42a7fd --- /dev/null +++ b/go/vehicle-gateway/internal/topics/topics.go @@ -0,0 +1,8 @@ +package topics + +const ( + RawGB32960 = "vehicle.raw.go.gb32960.v1" + RawJT808 = "vehicle.raw.go.jt808.v1" + RawYutongMQTT = "vehicle.raw.go.yutong-mqtt.v1" + Unified = "vehicle.event.go.unified.v1" +)