Files
lingniu-vehicle-ingest/docs/superpowers/plans/2026-07-02-nats-kafka-ingest-refactor.md
2026-07-02 15:07:44 +08:00

14 KiB

NATS Kafka Ingest Refactor Implementation Plan

For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (- [ ]) syntax for tracking.

Goal: Move vehicle ingress publishing from direct Kafka writes to NATS JetStream, then bridge NATS to existing Kafka topics so downstream history, realtime, and stats services keep working.

Architecture: gateway publishes parsed RAW and UNIFIED envelopes to NATS JetStream subjects. A new nats-kafka-bridge durable pull consumer reads from JetStream, writes batches to Kafka, and only ACKs NATS after Kafka succeeds. Existing Kafka consumers remain unchanged in this phase.

Tech Stack: Go 1.26, github.com/nats-io/nats.go, NATS JetStream, Kafka via segmentio/kafka-go, existing eventbus.Sink, systemd single-binary deployment on ECS.


File Structure

  • Create go/vehicle-gateway/internal/eventbus/nats_sink.go: NATS JetStream implementation of eventbus.Sink.
  • Create go/vehicle-gateway/internal/eventbus/nats_sink_test.go: topic/subject routing and envelope serialization tests.
  • Modify go/vehicle-gateway/cmd/gateway/main.go: choose NATS sink when NATS_URL is configured; keep Kafka path as fallback.
  • Create go/vehicle-gateway/cmd/nats-kafka-bridge/main.go: bridge process that creates JetStream stream/consumer, pulls NATS messages, writes Kafka batches, then ACKs.
  • Create go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go: bridge batch ACK behavior and topic routing tests.
  • Modify go/vehicle-gateway/go.mod: add github.com/nats-io/nats.go.
  • Create deploy/nats/nats-server.conf: single-node JetStream config for Kafka ECS.
  • Modify ECS deploy process: include nats-kafka-bridge binary and systemd service.

Task 1: Add NATS JetStream Sink

Files:

  • Create: go/vehicle-gateway/internal/eventbus/nats_sink.go

  • Create: go/vehicle-gateway/internal/eventbus/nats_sink_test.go

  • Modify: go/vehicle-gateway/go.mod

  • Step 1: Write failing subject routing test

Create go/vehicle-gateway/internal/eventbus/nats_sink_test.go with a fake JetStream publisher:

package eventbus

import (
	"context"
	"encoding/json"
	"testing"

	"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)

func TestNATSSinkRoutesRawAndUnifiedSubjects(t *testing.T) {
	publisher := &recordingNATSPublisher{}
	sink := newNATSSinkWithPublisher(publisher, NATSConfig{
		RawSubjects: map[envelope.Protocol]string{
			envelope.ProtocolJT808: "vehicle.raw.jt808.v1",
		},
		UnifiedSubject: "vehicle.event.unified.v1",
	})
	env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"}

	if err := sink.PublishRaw(context.Background(), env); err != nil {
		t.Fatalf("PublishRaw() error = %v", err)
	}
	if err := sink.PublishUnified(context.Background(), env); err != nil {
		t.Fatalf("PublishUnified() error = %v", err)
	}

	if got, want := publisher.messages[0].subject, "vehicle.raw.jt808.v1"; got != want {
		t.Fatalf("raw subject = %q, want %q", got, want)
	}
	if got, want := publisher.messages[1].subject, "vehicle.event.unified.v1"; got != want {
		t.Fatalf("unified subject = %q, want %q", got, want)
	}
	var decoded envelope.FrameEnvelope
	if err := json.Unmarshal(publisher.messages[0].data, &decoded); err != nil {
		t.Fatalf("raw payload is not envelope json: %v", err)
	}
	if decoded.Phone != "13307795425" {
		t.Fatalf("decoded phone = %q", decoded.Phone)
	}
}

type recordingNATSPublisher struct {
	messages []recordedNATSMessage
}

type recordedNATSMessage struct {
	subject string
	data    []byte
}

func (p *recordingNATSPublisher) Publish(_ context.Context, subject string, data []byte, _ ...NATSPublishOption) error {
	p.messages = append(p.messages, recordedNATSMessage{subject: subject, data: append([]byte(nil), data...)})
	return nil
}

Run: go test ./internal/eventbus -run TestNATSSinkRoutesRawAndUnifiedSubjects -count=1

Expected: FAIL because NATSConfig, newNATSSinkWithPublisher, and NATSPublishOption do not exist.

  • Step 2: Implement NATS sink

Create go/vehicle-gateway/internal/eventbus/nats_sink.go:

package eventbus

import (
	"context"
	"errors"
	"fmt"

	"github.com/nats-io/nats.go"

	"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)

type NATSConfig struct {
	URL            string
	Stream         string
	RawSubjects    map[envelope.Protocol]string
	UnifiedSubject string
}

type NATSSink struct {
	conn           *nats.Conn
	publisher      natsPublisher
	rawSubjects    map[envelope.Protocol]string
	unifiedSubject string
}

type NATSPublishOption = nats.PubOpt

type natsPublisher interface {
	Publish(context.Context, string, []byte, ...NATSPublishOption) error
}

func NewNATSSink(cfg NATSConfig) (*NATSSink, error) {
	if cfg.URL == "" {
		return nil, errors.New("nats url is required")
	}
	conn, err := nats.Connect(cfg.URL, nats.Name("lingniu-vehicle-gateway"))
	if err != nil {
		return nil, err
	}
	js, err := conn.JetStream()
	if err != nil {
		conn.Close()
		return nil, err
	}
	sink := newNATSSinkWithPublisher(natsJetStreamPublisher{js: js}, cfg)
	sink.conn = conn
	return sink, nil
}

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",
	}
	for protocol, subject := range cfg.RawSubjects {
		if subject != "" {
			rawSubjects[protocol] = subject
		}
	}
	unifiedSubject := cfg.UnifiedSubject
	if unifiedSubject == "" {
		unifiedSubject = "vehicle.event.unified.v1"
	}
	return &NATSSink{publisher: publisher, rawSubjects: rawSubjects, unifiedSubject: unifiedSubject}
}

func (s *NATSSink) PublishRaw(ctx context.Context, env envelope.FrameEnvelope) error {
	subject, ok := s.rawSubjects[env.Protocol]
	if !ok || subject == "" {
		return fmt.Errorf("raw subject not configured for protocol %s", env.Protocol)
	}
	return s.publish(ctx, subject, env)
}

func (s *NATSSink) PublishUnified(ctx context.Context, env envelope.FrameEnvelope) error {
	return s.publish(ctx, s.unifiedSubject, env)
}

func (s *NATSSink) Close() error {
	if s != nil && s.conn != nil {
		s.conn.Drain()
		s.conn.Close()
	}
	return nil
}

func (s *NATSSink) publish(ctx context.Context, subject string, env envelope.FrameEnvelope) error {
	payload, err := env.MarshalJSONBytes()
	if err != nil {
		return err
	}
	return s.publisher.Publish(ctx, subject, payload, nats.MsgId(env.StableEventID()))
}

type natsJetStreamPublisher struct {
	js nats.JetStreamContext
}

func (p natsJetStreamPublisher) Publish(ctx context.Context, subject string, data []byte, opts ...NATSPublishOption) error {
	_, err := p.js.Publish(subject, data, append(opts, nats.Context(ctx))...)
	return err
}
  • Step 3: Run eventbus tests

Run: go get github.com/nats-io/nats.go@latest && go test ./internal/eventbus -count=1

Expected: PASS.

Task 2: Gateway Sink Selection

Files:

  • Modify: go/vehicle-gateway/cmd/gateway/main.go

  • Step 1: Write failing buildSink test

Create go/vehicle-gateway/cmd/gateway/main_test.go if it does not exist, or add:

func TestBuildSinkPrefersNATSWhenConfigured(t *testing.T) {
	t.Setenv("NATS_URL", "nats://127.0.0.1:4222")
	t.Setenv("KAFKA_BROKERS", "127.0.0.1:9092")
	if got := chooseSinkMode(); got != "nats" {
		t.Fatalf("sink mode = %q, want nats", got)
	}
}

Run: go test ./cmd/gateway -run TestBuildSinkPrefersNATSWhenConfigured -count=1

Expected: FAIL because chooseSinkMode does not exist.

  • Step 2: Implement mode choice and NATS sink branch

In cmd/gateway/main.go, add:

func chooseSinkMode() string {
	if strings.TrimSpace(os.Getenv("NATS_URL")) != "" {
		return "nats"
	}
	return "kafka"
}

At the top of buildSink, add:

if chooseSinkMode() == "nats" {
	sink, err := eventbus.NewNATSSink(eventbus.NATSConfig{
		URL:    env("NATS_URL", ""),
		Stream: env("NATS_STREAM", "vehicle-ingest"),
		RawSubjects: map[envelope.Protocol]string{
			envelope.ProtocolGB32960:    env("NATS_SUBJECT_GB32960_RAW", "vehicle.raw.gb32960.v1"),
			envelope.ProtocolJT808:      env("NATS_SUBJECT_JT808_RAW", "vehicle.raw.jt808.v1"),
			envelope.ProtocolYutongMQTT: env("NATS_SUBJECT_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1"),
		},
		UnifiedSubject: env("NATS_SUBJECT_UNIFIED", "vehicle.event.unified.v1"),
	})
	if err != nil {
		return nil, err
	}
	logger.Info("nats jetstream sink enabled", "url", env("NATS_URL", ""), "stream", env("NATS_STREAM", "vehicle-ingest"))
	return eventbus.NewAsyncSink(sink, eventbus.AsyncConfig{
		QueueSize:        envInt("NATS_ASYNC_QUEUE_SIZE", 100000),
		Workers:          envInt("NATS_ASYNC_WORKERS", 8),
		OperationTimeout: time.Duration(envInt("NATS_PUBLISH_TIMEOUT_MS", 30000)) * time.Millisecond,
		OnError: func(err error) {
			logger.Warn("nats async publish failed", "error", err)
		},
	}), nil
}
  • Step 3: Run gateway tests

Run: go test ./cmd/gateway -count=1

Expected: PASS.

Task 3: NATS to Kafka Bridge

Files:

  • Create: go/vehicle-gateway/cmd/nats-kafka-bridge/main.go

  • Create: go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go

  • Step 1: Write failing bridge test

Create a test with fake NATS messages and fake Kafka writer. It must prove that the bridge ACKs only after Kafka write succeeds:

func TestBridgeAcksAfterKafkaWrite(t *testing.T) {
	nats := &fakeNATSBatchReader{messages: []bridgeMessage{
		{subject: "vehicle.raw.jt808.v1", data: []byte(`{"protocol":"JT808","phone":"13307795425"}`)},
	}}
	kafka := &fakeKafkaBatchWriter{}
	bridge := bridgeProcessor{nats: nats, kafka: kafka, topics: defaultBridgeTopics()}

	if err := bridge.processBatch(context.Background(), 10); err != nil {
		t.Fatalf("processBatch() error = %v", err)
	}
	if kafka.writeCalls != 1 {
		t.Fatalf("kafka write calls = %d, want 1", kafka.writeCalls)
	}
	if !nats.messages[0].acked {
		t.Fatal("nats message was not acked")
	}
}

Run: go test ./cmd/nats-kafka-bridge -run TestBridgeAcksAfterKafkaWrite -count=1

Expected: FAIL because bridge types do not exist.

  • Step 2: Implement minimal bridge processor

Implement small interfaces in main.go:

type bridgeMessage struct {
	subject string
	data    []byte
	ack     func() error
}

type bridgeNATSReader interface {
	Fetch(context.Context, int) ([]bridgeMessage, error)
}

type bridgeKafkaWriter interface {
	Write(context.Context, []kafka.Message) error
}

type bridgeProcessor struct {
	nats   bridgeNATSReader
	kafka  bridgeKafkaWriter
	topics map[string]string
}

func (p bridgeProcessor) processBatch(ctx context.Context, batchSize int) error {
	messages, err := p.nats.Fetch(ctx, batchSize)
	if err != nil {
		return err
	}
	kafkaMessages := make([]kafka.Message, 0, len(messages))
	for _, msg := range messages {
		topic, ok := p.topics[msg.subject]
		if !ok {
			return fmt.Errorf("no kafka topic for nats subject %s", msg.subject)
		}
		key := kafkaKeyFromEnvelope(msg.data)
		kafkaMessages = append(kafkaMessages, kafka.Message{Topic: topic, Key: key, Value: msg.data})
	}
	if len(kafkaMessages) == 0 {
		return nil
	}
	if err := p.kafka.Write(ctx, kafkaMessages); err != nil {
		return err
	}
	for _, msg := range messages {
		if msg.ack != nil {
			if err := msg.ack(); err != nil {
				return err
			}
		}
	}
	return nil
}
  • Step 3: Implement real NATS and Kafka adapters

Use nats.PullSubscribe, Fetch, Ack, and kafka.Writer.WriteMessages. The bridge must create or reuse a durable consumer named by NATS_CONSUMER, default nats-kafka-bridge.

  • Step 4: Run bridge tests

Run: go test ./cmd/nats-kafka-bridge -count=1

Expected: PASS.

Task 4: NATS Server Deployment

Files:

  • Create: deploy/nats/nats-server.conf

  • Step 1: Create NATS config

Create:

server_name: lingniu-nats-01
port: 4222
http_port: 8222

jetstream {
  store_dir: "/data/nats/jetstream"
  max_file_store: 20GB
  max_mem_store: 512MB
}
  • Step 2: Deploy on Kafka ECS

On Kafka ECS, install or run nats:2 with:

mkdir -p /opt/lingniu-nats /data/nats/jetstream
docker run -d --name lingniu-nats --restart unless-stopped \
  -p 4222:4222 -p 8222:8222 \
  -v /opt/lingniu-nats/nats-server.conf:/etc/nats/nats-server.conf:ro \
  -v /data/nats:/data/nats \
  nats:2 -c /etc/nats/nats-server.conf
  • Step 3: Verify NATS

Run from gateway ECS:

timeout 3 bash -c '</dev/tcp/172.17.111.56/4222'

Expected: exit 0.

Task 5: ECS Cutover and Verification

Files:

  • Modify ECS env files only.

  • Step 1: Deploy binaries

Build and upload: gateway, nats-kafka-bridge, history-writer, stat-writer, realtime-api.

  • Step 2: Configure gateway

Add to /opt/lingniu-go-native/env/gateway.env:

NATS_URL=nats://172.17.111.56:4222
NATS_STREAM=vehicle-ingest
NATS_ASYNC_QUEUE_SIZE=100000
NATS_ASYNC_WORKERS=8
NATS_PUBLISH_TIMEOUT_MS=30000
  • Step 3: Configure bridge service

Create /opt/lingniu-go-native/env/nats-kafka-bridge.env with:

NATS_URL=nats://172.17.111.56:4222
NATS_STREAM=vehicle-ingest
NATS_CONSUMER=nats-kafka-bridge
NATS_BATCH_SIZE=500
KAFKA_BROKERS=172.17.111.56:9092
KAFKA_TOPIC_GB32960_RAW=vehicle.raw.go.gb32960.v1
KAFKA_TOPIC_JT808_RAW=vehicle.raw.go.jt808.v1
KAFKA_TOPIC_YUTONG_MQTT_RAW=vehicle.raw.go.yutong-mqtt.v1
KAFKA_TOPIC_UNIFIED=vehicle.event.go.unified.v1
  • Step 4: Verify production

Run:

systemctl is-active lingniu-go-gateway.service lingniu-go-nats-kafka-bridge.service lingniu-go-history-writer.service lingniu-go-realtime-api.service
ss -Htan state established "( sport = :32960 )"
ss -Htan state established "( sport = :808 )" | wc -l
journalctl -u lingniu-go-gateway.service --since '2 minutes ago' --no-pager | grep -E 'ERROR|deadline|nats async publish failed'
journalctl -u lingniu-go-nats-kafka-bridge.service --since '2 minutes ago' --no-pager | grep -E 'ERROR|failed'

Expected:

  • All services active.
  • 32960 Recv-Q remains 0 or near 0.
  • 808 connections remain stable.
  • No gateway publish deadline errors.
  • Kafka downstream consumers continue updating TDengine and Redis.

Self-Review

  • Spec coverage: gateway NATS publish, bridge to Kafka, ECS NATS deployment, and production verification are covered.
  • Scope intentionally excludes rewriting history/stat/realtime to consume NATS directly; they remain Kafka consumers in this phase.
  • No placeholders remain; later downlink notice/reply subjects are intentionally deferred because they are not needed for first reliable ingest cutover.