From 5145d156bc829b455f29bdf22925c833574577af Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 14:11:43 +0800 Subject: [PATCH] feat(go): decouple gateway publish with async sink --- go/vehicle-gateway/cmd/gateway/main.go | 14 +++ .../internal/eventbus/async_sink.go | 119 ++++++++++++++++++ .../internal/eventbus/async_sink_test.go | 105 ++++++++++++++++ 3 files changed, 238 insertions(+) create mode 100644 go/vehicle-gateway/internal/eventbus/async_sink.go create mode 100644 go/vehicle-gateway/internal/eventbus/async_sink_test.go diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 883161c1..3aed504c 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -178,6 +178,20 @@ func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) logger.Info("kafka durable spool enabled", "dir", spoolDir, "replay_interval_ms", interval.Milliseconds(), "replay_batch_size", replayBatchSize) out = durable } + if envBool("KAFKA_ASYNC_ENABLED", true) { + queueSize := envInt("KAFKA_ASYNC_QUEUE_SIZE", 10000) + workers := envInt("KAFKA_ASYNC_WORKERS", 1) + timeout := time.Duration(envInt("KAFKA_ASYNC_PUBLISH_TIMEOUT_MS", 30000)) * time.Millisecond + out = eventbus.NewAsyncSink(out, eventbus.AsyncConfig{ + QueueSize: queueSize, + Workers: workers, + OperationTimeout: timeout, + OnError: func(err error) { + logger.Warn("kafka async publish failed", "error", err) + }, + }) + logger.Info("kafka async publish enabled", "queue_size", queueSize, "workers", workers, "publish_timeout_ms", timeout.Milliseconds()) + } return out, nil } diff --git a/go/vehicle-gateway/internal/eventbus/async_sink.go b/go/vehicle-gateway/internal/eventbus/async_sink.go new file mode 100644 index 00000000..5b6d2591 --- /dev/null +++ b/go/vehicle-gateway/internal/eventbus/async_sink.go @@ -0,0 +1,119 @@ +package eventbus + +import ( + "context" + "errors" + "sync" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +type AsyncConfig struct { + QueueSize int + Workers int + OperationTimeout time.Duration + OnError func(error) +} + +type AsyncSink struct { + delegate Sink + jobs chan asyncJob + timeout time.Duration + onError func(error) + + closeOnce sync.Once + closed chan struct{} + done chan struct{} + wg sync.WaitGroup +} + +type asyncJob struct { + kind string + env envelope.FrameEnvelope +} + +var ErrAsyncSinkClosed = errors.New("async sink is closed") + +func NewAsyncSink(delegate Sink, cfg AsyncConfig) *AsyncSink { + if delegate == nil { + panic("async delegate sink must not be nil") + } + if cfg.QueueSize <= 0 { + cfg.QueueSize = 10_000 + } + if cfg.Workers <= 0 { + cfg.Workers = 1 + } + if cfg.OperationTimeout <= 0 { + cfg.OperationTimeout = 30 * time.Second + } + s := &AsyncSink{ + delegate: delegate, + jobs: make(chan asyncJob, cfg.QueueSize), + timeout: cfg.OperationTimeout, + onError: cfg.OnError, + closed: make(chan struct{}), + done: make(chan struct{}), + } + s.wg.Add(cfg.Workers) + for i := 0; i < cfg.Workers; i++ { + go s.worker() + } + go func() { + s.wg.Wait() + close(s.done) + }() + return s +} + +func (s *AsyncSink) PublishRaw(ctx context.Context, env envelope.FrameEnvelope) error { + return s.enqueue(ctx, asyncJob{kind: "raw", env: env}) +} + +func (s *AsyncSink) PublishUnified(ctx context.Context, env envelope.FrameEnvelope) error { + return s.enqueue(ctx, asyncJob{kind: "unified", env: env}) +} + +func (s *AsyncSink) Close() error { + s.closeOnce.Do(func() { + close(s.closed) + close(s.jobs) + }) + <-s.done + return s.delegate.Close() +} + +func (s *AsyncSink) enqueue(ctx context.Context, job asyncJob) error { + select { + case <-s.closed: + return ErrAsyncSinkClosed + default: + } + select { + case s.jobs <- job: + return nil + case <-s.closed: + return ErrAsyncSinkClosed + case <-ctx.Done(): + return ctx.Err() + } +} + +func (s *AsyncSink) worker() { + defer s.wg.Done() + for job := range s.jobs { + ctx, cancel := context.WithTimeout(context.Background(), s.timeout) + var err error + switch job.kind { + case "raw": + err = s.delegate.PublishRaw(ctx, job.env) + case "unified": + err = s.delegate.PublishUnified(ctx, job.env) + } + cancel() + if err != nil && s.onError != nil { + s.onError(err) + } + } +} diff --git a/go/vehicle-gateway/internal/eventbus/async_sink_test.go b/go/vehicle-gateway/internal/eventbus/async_sink_test.go new file mode 100644 index 00000000..fc58548c --- /dev/null +++ b/go/vehicle-gateway/internal/eventbus/async_sink_test.go @@ -0,0 +1,105 @@ +package eventbus + +import ( + "context" + "testing" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestAsyncSinkPublishRawReturnsAfterQueueing(t *testing.T) { + delegate := newBlockingSink() + sink := NewAsyncSink(delegate, AsyncConfig{QueueSize: 1, Workers: 1, OperationTimeout: time.Second}) + + start := time.Now() + if err := sink.PublishRaw(context.Background(), envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "LNBSCB3D4R1234567"}); err != nil { + t.Fatalf("PublishRaw() error = %v", err) + } + if elapsed := time.Since(start); elapsed > 100*time.Millisecond { + t.Fatalf("PublishRaw() blocked for %s, want quick queueing", elapsed) + } + + select { + case <-delegate.rawStarted: + case <-time.After(time.Second): + t.Fatal("delegate raw publish was not started") + } + delegate.release() + if err := sink.Close(); err != nil { + t.Fatalf("Close() error = %v", err) + } +} + +func TestAsyncSinkPublishesInFIFOOrderWithSingleWorker(t *testing.T) { + delegate := &orderedSink{} + sink := NewAsyncSink(delegate, AsyncConfig{QueueSize: 4, Workers: 1, OperationTimeout: time.Second}) + + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425"} + 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 err := sink.Close(); err != nil { + t.Fatalf("Close() error = %v", err) + } + + if got, want := delegate.calls, []string{"raw", "unified"}; len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("calls = %#v, want %#v", got, want) + } +} + +type blockingSink struct { + rawStarted chan struct{} + releaseRaw chan struct{} +} + +func newBlockingSink() *blockingSink { + return &blockingSink{ + rawStarted: make(chan struct{}), + releaseRaw: make(chan struct{}), + } +} + +func (s *blockingSink) PublishRaw(context.Context, envelope.FrameEnvelope) error { + close(s.rawStarted) + <-s.releaseRaw + return nil +} + +func (s *blockingSink) PublishUnified(context.Context, envelope.FrameEnvelope) error { + return nil +} + +func (s *blockingSink) Close() error { + s.release() + return nil +} + +func (s *blockingSink) release() { + select { + case <-s.releaseRaw: + default: + close(s.releaseRaw) + } +} + +type orderedSink struct { + calls []string +} + +func (s *orderedSink) PublishRaw(context.Context, envelope.FrameEnvelope) error { + s.calls = append(s.calls, "raw") + return nil +} + +func (s *orderedSink) PublishUnified(context.Context, envelope.FrameEnvelope) error { + s.calls = append(s.calls, "unified") + return nil +} + +func (s *orderedSink) Close() error { + return nil +}