From 967981c0406a9666d460ed2b6b026a216c3d2007 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 21:45:28 +0800 Subject: [PATCH] feat: add go tcp gateway runtime --- go/vehicle-gateway/cmd/gateway/main.go | 130 ++++++++++- .../internal/eventbus/log_sink.go | 33 +++ .../internal/gateway/tcp_server.go | 201 ++++++++++++++++++ .../internal/gateway/tcp_server_test.go | 145 +++++++++++++ 4 files changed, 507 insertions(+), 2 deletions(-) create mode 100644 go/vehicle-gateway/internal/eventbus/log_sink.go create mode 100644 go/vehicle-gateway/internal/gateway/tcp_server.go create mode 100644 go/vehicle-gateway/internal/gateway/tcp_server_test.go diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 1a4e4620..cb653141 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -1,8 +1,134 @@ package main -import "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" +import ( + "context" + "log/slog" + "os" + "os/signal" + "strings" + "sync" + "syscall" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/gateway" + "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" +) func main() { logger := observability.NewLogger("vehicle-gateway") - logger.Info("vehicle gateway scaffold started") + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + sink, err := buildSink(logger) + if err != nil { + logger.Error("build sink failed", "error", err) + os.Exit(1) + } + defer sink.Close() + + protocols := []gateway.TCPProtocol{ + { + Protocol: envelope.ProtocolGB32960, + Addr: env("GB32960_TCP_ADDR", ":32960"), + Extract: gb32960.ExtractFrames, + Parse: gb32960.ParseFrame, + }, + { + Protocol: envelope.ProtocolJT808, + Addr: env("JT808_TCP_ADDR", ":808"), + Extract: jt808.ExtractFrames, + Parse: jt808.ParseFrame, + }, + } + + var wg sync.WaitGroup + errs := make(chan error, len(protocols)) + for _, protocol := range protocols { + server, err := gateway.NewTCPServer(gateway.TCPServerConfig{ + Protocol: protocol, + Sink: sink, + Logger: logger, + ReadBufferSize: envInt("TCP_READ_BUFFER_BYTES", 64*1024), + IdleTimeout: time.Duration(envInt("TCP_IDLE_TIMEOUT_SECONDS", 180)) * time.Second, + MaxConnections: envInt("TCP_MAX_CONNECTIONS", 20_000), + }) + if err != nil { + logger.Error("build tcp server failed", "protocol", protocol.Protocol, "error", err) + os.Exit(1) + } + wg.Add(1) + go func() { + defer wg.Done() + if err := server.ListenAndServe(ctx); err != nil { + errs <- err + } + }() + } + + logger.Info("vehicle gateway started") + select { + case <-ctx.Done(): + case err := <-errs: + logger.Error("vehicle gateway listener failed", "error", err) + stop() + } + wg.Wait() +} + +func buildSink(logger *slog.Logger) (eventbus.Sink, error) { + brokers := splitCSV(os.Getenv("KAFKA_BROKERS")) + if len(brokers) == 0 { + logger.Warn("KAFKA_BROKERS is empty; using log sink") + return eventbus.NewLogSink(logger), nil + } + return 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"), + }, + UnifiedTopic: env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1"), + }) +} + +func env(key string, fallback string) string { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + return value +} + +func envInt(key string, fallback int) int { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + var out int + for _, r := range value { + if r < '0' || r > '9' { + return fallback + } + out = out*10 + int(r-'0') + } + if out <= 0 { + return fallback + } + return out +} + +func splitCSV(value string) []string { + var out []string + for _, item := range strings.Split(value, ",") { + item = strings.TrimSpace(item) + if item != "" { + out = append(out, item) + } + } + return out } diff --git a/go/vehicle-gateway/internal/eventbus/log_sink.go b/go/vehicle-gateway/internal/eventbus/log_sink.go new file mode 100644 index 00000000..68eed2e5 --- /dev/null +++ b/go/vehicle-gateway/internal/eventbus/log_sink.go @@ -0,0 +1,33 @@ +package eventbus + +import ( + "context" + "log/slog" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +type LogSink struct { + logger *slog.Logger +} + +func NewLogSink(logger *slog.Logger) *LogSink { + if logger == nil { + logger = slog.Default() + } + return &LogSink{logger: logger} +} + +func (s *LogSink) PublishRaw(_ context.Context, env envelope.FrameEnvelope) error { + s.logger.Info("raw envelope", "protocol", env.Protocol, "event_id", env.StableEventID(), "vehicle_key", env.VehicleKey(), "message_id", env.MessageID, "status", env.ParseStatus) + return nil +} + +func (s *LogSink) PublishUnified(_ context.Context, env envelope.FrameEnvelope) error { + s.logger.Info("unified envelope", "protocol", env.Protocol, "event_id", env.StableEventID(), "vehicle_key", env.VehicleKey(), "message_id", env.MessageID) + return nil +} + +func (s *LogSink) Close() error { + return nil +} diff --git a/go/vehicle-gateway/internal/gateway/tcp_server.go b/go/vehicle-gateway/internal/gateway/tcp_server.go new file mode 100644 index 00000000..05a3a426 --- /dev/null +++ b/go/vehicle-gateway/internal/gateway/tcp_server.go @@ -0,0 +1,201 @@ +package gateway + +import ( + "context" + "encoding/hex" + "errors" + "fmt" + "io" + "log/slog" + "net" + "strings" + "sync" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" +) + +type FrameExtractor func([]byte) (frames [][]byte, remainder []byte, err error) + +type FrameParser func(raw []byte, receivedAtMS int64, sourceEndpoint string) (envelope.FrameEnvelope, error) + +type TCPProtocol struct { + Protocol envelope.Protocol + Addr string + Extract FrameExtractor + Parse FrameParser +} + +type TCPServer struct { + protocol TCPProtocol + sink eventbus.Sink + logger *slog.Logger + readBufferSize int + idleTimeout time.Duration + maxConnections int +} + +type TCPServerConfig struct { + Protocol TCPProtocol + Sink eventbus.Sink + Logger *slog.Logger + ReadBufferSize int + IdleTimeout time.Duration + MaxConnections int +} + +func NewTCPServer(cfg TCPServerConfig) (*TCPServer, error) { + if cfg.Protocol.Protocol == "" { + return nil, errors.New("protocol is required") + } + if strings.TrimSpace(cfg.Protocol.Addr) == "" { + return nil, errors.New("listen addr is required") + } + if cfg.Protocol.Extract == nil { + return nil, errors.New("frame extractor is required") + } + if cfg.Protocol.Parse == nil { + return nil, errors.New("frame parser is required") + } + if cfg.Sink == nil { + return nil, errors.New("sink is required") + } + if cfg.Logger == nil { + cfg.Logger = slog.Default() + } + if cfg.ReadBufferSize <= 0 { + cfg.ReadBufferSize = 32 * 1024 + } + if cfg.IdleTimeout <= 0 { + cfg.IdleTimeout = 2 * time.Minute + } + if cfg.MaxConnections <= 0 { + cfg.MaxConnections = 10_000 + } + return &TCPServer{ + protocol: cfg.Protocol, + sink: cfg.Sink, + logger: cfg.Logger, + readBufferSize: cfg.ReadBufferSize, + idleTimeout: cfg.IdleTimeout, + maxConnections: cfg.MaxConnections, + }, nil +} + +func (s *TCPServer) ListenAndServe(ctx context.Context) error { + var lc net.ListenConfig + listener, err := lc.Listen(ctx, "tcp", s.protocol.Addr) + if err != nil { + return err + } + defer listener.Close() + + go func() { + <-ctx.Done() + _ = listener.Close() + }() + + s.logger.Info("tcp listener started", "protocol", s.protocol.Protocol, "addr", listener.Addr().String()) + sem := make(chan struct{}, s.maxConnections) + var wg sync.WaitGroup + defer wg.Wait() + + for { + conn, err := listener.Accept() + if err != nil { + if ctx.Err() != nil { + return nil + } + s.logger.Warn("tcp accept failed", "protocol", s.protocol.Protocol, "error", err) + continue + } + select { + case sem <- struct{}{}: + wg.Add(1) + go func() { + defer wg.Done() + defer func() { <-sem }() + s.handleConnection(ctx, conn) + }() + default: + s.logger.Warn("tcp connection rejected: max connections reached", "protocol", s.protocol.Protocol, "remote", conn.RemoteAddr().String()) + _ = conn.Close() + } + } +} + +func (s *TCPServer) handleConnection(ctx context.Context, conn net.Conn) { + defer conn.Close() + source := conn.RemoteAddr().String() + log := s.logger.With("protocol", s.protocol.Protocol, "remote", source) + log.Info("tcp connection opened") + defer log.Info("tcp connection closed") + + readBuffer := make([]byte, s.readBufferSize) + var pending []byte + for { + _ = conn.SetReadDeadline(time.Now().Add(s.idleTimeout)) + n, err := conn.Read(readBuffer) + if n > 0 { + pending = append(pending, readBuffer[:n]...) + frames, remainder, extractErr := s.protocol.Extract(pending) + if extractErr != nil { + log.Warn("frame extraction failed", "error", extractErr) + return + } + pending = remainder + for _, frame := range frames { + s.handleFrame(ctx, frame, source) + } + } + if err != nil { + if errors.Is(err, io.EOF) { + return + } + var netErr net.Error + if errors.As(err, &netErr) && netErr.Timeout() { + log.Warn("tcp connection idle timeout") + return + } + log.Warn("tcp read failed", "error", err) + return + } + if ctx.Err() != nil { + return + } + } +} + +func (s *TCPServer) handleFrame(ctx context.Context, raw []byte, source string) { + receivedAtMS := time.Now().UnixMilli() + env, err := s.protocol.Parse(raw, receivedAtMS, source) + if err != nil { + env = envelope.FrameEnvelope{ + Protocol: s.protocol.Protocol, + SourceEndpoint: source, + ReceivedAtMS: receivedAtMS, + EventTimeMS: receivedAtMS, + RawHex: strings.ToUpper(hex.EncodeToString(raw)), + ParseStatus: envelope.ParseBadFrame, + ParseError: err.Error(), + } + env.EventID = env.StableEventID() + } + + if err := s.sink.PublishRaw(ctx, env); err != nil { + s.logger.Error("publish raw failed", "protocol", s.protocol.Protocol, "event_id", env.StableEventID(), "error", err) + return + } + if env.ParseStatus == envelope.ParseBadFrame { + return + } + if err := s.sink.PublishUnified(ctx, env); err != nil { + s.logger.Error("publish unified failed", "protocol", s.protocol.Protocol, "event_id", env.StableEventID(), "error", err) + return + } +} + +func (p TCPProtocol) String() string { + return fmt.Sprintf("%s@%s", p.Protocol, p.Addr) +} diff --git a/go/vehicle-gateway/internal/gateway/tcp_server_test.go b/go/vehicle-gateway/internal/gateway/tcp_server_test.go new file mode 100644 index 00000000..291e72b6 --- /dev/null +++ b/go/vehicle-gateway/internal/gateway/tcp_server_test.go @@ -0,0 +1,145 @@ +package gateway + +import ( + "context" + "encoding/hex" + "log/slog" + "net" + "testing" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/gb32960" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/jt808" +) + +func TestTCPServerPublishesGoodFrameToRawAndUnified(t *testing.T) { + frame, err := hex.DecodeString("7E020000320133077954250001000000000048000301D2C4C707376139000A00E6004F26063016235701040001900C2504000000000202000030011F31010F867E") + if err != nil { + t.Fatal(err) + } + sink := &recordingSink{} + server := newTestServer(t, TCPProtocol{ + Protocol: envelope.ProtocolJT808, + Addr: ":0", + Extract: jt808.ExtractFrames, + Parse: jt808.ParseFrame, + }, sink) + + client, done := runPipe(t, server) + if _, err := client.Write(frame); err != nil { + t.Fatalf("client.Write() error = %v", err) + } + _ = client.Close() + <-done + + if len(sink.raw) != 1 || len(sink.unified) != 1 { + t.Fatalf("raw=%d unified=%d", len(sink.raw), len(sink.unified)) + } + if sink.raw[0].Phone != "013307795425" { + t.Fatalf("phone = %q", sink.raw[0].Phone) + } + if sink.raw[0].Fields[envelope.FieldTotalMileageKM] != 10241.2 { + t.Fatalf("total mileage = %#v", sink.raw[0].Fields[envelope.FieldTotalMileageKM]) + } +} + +func TestTCPServerPublishesBadFrameOnlyToRaw(t *testing.T) { + good := buildGBFrame(0x02, 0xfe, "LNBSCB3D4R1234567", nil) + good[len(good)-1] ^= 0xff + sink := &recordingSink{} + server := newTestServer(t, TCPProtocol{ + Protocol: envelope.ProtocolGB32960, + Addr: ":0", + Extract: gb32960.ExtractFrames, + Parse: gb32960.ParseFrame, + }, sink) + + client, done := runPipe(t, server) + if _, err := client.Write(good); err != nil { + t.Fatalf("client.Write() error = %v", err) + } + _ = client.Close() + <-done + + if len(sink.raw) != 1 || len(sink.unified) != 0 { + t.Fatalf("raw=%d unified=%d", len(sink.raw), len(sink.unified)) + } + if sink.raw[0].ParseStatus != envelope.ParseBadFrame { + t.Fatalf("parse status = %q", sink.raw[0].ParseStatus) + } + if sink.raw[0].ParseError == "" { + t.Fatal("parse error should be recorded") + } +} + +func newTestServer(t *testing.T, protocol TCPProtocol, sink *recordingSink) *TCPServer { + t.Helper() + server, err := NewTCPServer(TCPServerConfig{ + Protocol: protocol, + Sink: sink, + Logger: slog.New(slog.NewTextHandler(testWriter{t: t}, nil)), + ReadBufferSize: 1024, + IdleTimeout: time.Second, + MaxConnections: 1, + }) + if err != nil { + t.Fatalf("NewTCPServer() error = %v", err) + } + return server +} + +func runPipe(t *testing.T, server *TCPServer) (net.Conn, <-chan struct{}) { + t.Helper() + client, srv := net.Pipe() + done := make(chan struct{}) + go func() { + defer close(done) + server.handleConnection(context.Background(), srv) + }() + return client, done +} + +type recordingSink struct { + raw []envelope.FrameEnvelope + unified []envelope.FrameEnvelope +} + +func (s *recordingSink) PublishRaw(_ context.Context, env envelope.FrameEnvelope) error { + s.raw = append(s.raw, env) + return nil +} + +func (s *recordingSink) PublishUnified(_ context.Context, env envelope.FrameEnvelope) error { + s.unified = append(s.unified, env) + return nil +} + +func (s *recordingSink) Close() error { + return nil +} + +type testWriter struct { + t *testing.T +} + +func (w testWriter) Write(p []byte) (int, error) { + w.t.Log(string(p)) + return len(p), nil +} + +func buildGBFrame(command byte, response byte, vin string, body []byte) []byte { + frame := []byte{'#', '#', command, response} + vinBytes := []byte(vin) + if len(vinBytes) < 17 { + vinBytes = append(vinBytes, make([]byte, 17-len(vinBytes))...) + } + frame = append(frame, vinBytes[:17]...) + frame = append(frame, 0x01, byte(len(body)>>8), byte(len(body))) + frame = append(frame, body...) + var bcc byte + for _, value := range frame[2:] { + bcc ^= value + } + return append(frame, bcc) +}