feat: add go tcp gateway runtime

This commit is contained in:
lingniu
2026-07-01 21:45:28 +08:00
parent 1fdbc29e27
commit 967981c040
4 changed files with 507 additions and 2 deletions

View File

@@ -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
}

View File

@@ -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)
}

View File

@@ -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)
}