feat: spool gateway publishes to disk

This commit is contained in:
lingniu
2026-07-01 23:00:06 +08:00
parent 4a32e2181a
commit 243de5c2b2
4 changed files with 360 additions and 5 deletions

View File

@@ -27,7 +27,7 @@ func main() {
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
sink, err := buildSink(logger)
sink, err := buildSink(ctx, logger)
if err != nil {
logger.Error("build sink failed", "error", err)
os.Exit(1)
@@ -132,7 +132,7 @@ func buildIdentityResolver(ctx context.Context, logger *slog.Logger) (identity.R
return identity.NewMySQLResolver(db, table), func() { _ = db.Close() }, nil
}
func buildSink(logger *slog.Logger) (eventbus.Sink, error) {
func buildSink(ctx context.Context, 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")
@@ -150,11 +150,22 @@ func buildSink(logger *slog.Logger) (eventbus.Sink, error) {
if err != nil {
return nil, err
}
return eventbus.NewRetryingSink(sink, eventbus.RetryConfig{
var out eventbus.Sink = eventbus.NewRetryingSink(sink, eventbus.RetryConfig{
Attempts: envInt("KAFKA_PUBLISH_ATTEMPTS", 3),
Backoff: time.Duration(envInt("KAFKA_PUBLISH_BACKOFF_MS", 100)) * time.Millisecond,
AttemptTimeout: time.Duration(envInt("KAFKA_PUBLISH_TIMEOUT_MS", 3000)) * time.Millisecond,
}), nil
})
spoolDir := strings.TrimSpace(os.Getenv("KAFKA_SPOOL_DIR"))
if spoolDir != "" {
durable := eventbus.NewDurableSink(out, eventbus.DurableConfig{Directory: spoolDir})
interval := time.Duration(envInt("KAFKA_SPOOL_REPLAY_INTERVAL_MS", 1000)) * time.Millisecond
go durable.ReplayLoop(ctx, interval, func(err error) {
logger.Warn("kafka spool replay failed", "error", err)
})
logger.Info("kafka durable spool enabled", "dir", spoolDir, "replay_interval_ms", interval.Milliseconds())
out = durable
}
return out, nil
}
func env(key string, fallback string) string {

View File

@@ -0,0 +1,187 @@
package eventbus
import (
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
type DurableConfig struct {
Directory string
}
type DurableSink struct {
delegate Sink
dir string
mu sync.Mutex
seq uint64
rawPending map[string]struct{}
}
type durableRecord struct {
Kind string `json:"kind"`
Envelope envelope.FrameEnvelope `json:"envelope"`
}
func NewDurableSink(delegate Sink, cfg DurableConfig) *DurableSink {
if delegate == nil {
panic("durable delegate sink must not be nil")
}
return &DurableSink{
delegate: delegate,
dir: strings.TrimSpace(cfg.Directory),
rawPending: map[string]struct{}{},
}
}
func (s *DurableSink) PublishRaw(ctx context.Context, env envelope.FrameEnvelope) error {
if err := s.delegate.PublishRaw(ctx, env); err == nil {
return nil
}
if err := s.spool("raw", env); err != nil {
return err
}
s.markRawPending(env)
return nil
}
func (s *DurableSink) PublishUnified(ctx context.Context, env envelope.FrameEnvelope) error {
if s.isRawPending(env) {
return s.spool("unified", env)
}
if err := s.delegate.PublishUnified(ctx, env); err == nil {
return nil
}
return s.spool("unified", env)
}
func (s *DurableSink) ReplayOnce(ctx context.Context) error {
files, err := filepath.Glob(filepath.Join(s.dir, "*.json"))
if err != nil {
return err
}
sort.Strings(files)
for _, file := range files {
record, err := readDurableRecord(file)
if err != nil {
return err
}
if err := s.publishRecord(ctx, record); err != nil {
return err
}
if err := os.Remove(file); err != nil {
return err
}
if record.Kind == "raw" {
s.clearRawPending(record.Envelope)
}
}
return nil
}
func (s *DurableSink) ReplayLoop(ctx context.Context, interval time.Duration, onError func(error)) {
if interval <= 0 {
interval = time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := s.ReplayOnce(ctx); err != nil && onError != nil {
onError(err)
}
}
}
}
func (s *DurableSink) Close() error {
return s.delegate.Close()
}
func (s *DurableSink) publishRecord(ctx context.Context, record durableRecord) error {
switch record.Kind {
case "raw":
return s.delegate.PublishRaw(ctx, record.Envelope)
case "unified":
return s.delegate.PublishUnified(ctx, record.Envelope)
default:
return fmt.Errorf("unknown durable record kind %q", record.Kind)
}
}
func (s *DurableSink) spool(kind string, env envelope.FrameEnvelope) error {
if s.dir == "" {
return fmt.Errorf("durable spool directory is empty")
}
if env.EventID == "" {
env.EventID = env.StableEventID()
}
if env.ParseStatus == "" {
env.ParseStatus = envelope.ParseOK
}
if err := os.MkdirAll(s.dir, 0o750); err != nil {
return err
}
payload, err := json.Marshal(durableRecord{Kind: kind, Envelope: env})
if err != nil {
return err
}
name := s.nextFileName(env, kind)
path := filepath.Join(s.dir, name)
tmp := path + ".tmp"
if err := os.WriteFile(tmp, payload, 0o640); err != nil {
return err
}
return os.Rename(tmp, path)
}
func (s *DurableSink) nextFileName(env envelope.FrameEnvelope, kind string) string {
s.mu.Lock()
defer s.mu.Unlock()
s.seq++
eventID := env.StableEventID()
if len(eventID) > 12 {
eventID = eventID[:12]
}
return fmt.Sprintf("%020d-%06d-%s-%s.json", time.Now().UnixNano(), s.seq, kind, eventID)
}
func (s *DurableSink) markRawPending(env envelope.FrameEnvelope) {
s.mu.Lock()
defer s.mu.Unlock()
s.rawPending[env.StableEventID()] = struct{}{}
}
func (s *DurableSink) clearRawPending(env envelope.FrameEnvelope) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.rawPending, env.StableEventID())
}
func (s *DurableSink) isRawPending(env envelope.FrameEnvelope) bool {
s.mu.Lock()
defer s.mu.Unlock()
_, ok := s.rawPending[env.StableEventID()]
return ok
}
func readDurableRecord(path string) (durableRecord, error) {
var record durableRecord
payload, err := os.ReadFile(path)
if err != nil {
return record, err
}
return record, json.Unmarshal(payload, &record)
}

View File

@@ -0,0 +1,153 @@
package eventbus
import (
"context"
"errors"
"os"
"path/filepath"
"testing"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
func TestDurableSinkSpoolsUnifiedWhenRawWasSpooled(t *testing.T) {
dir := t.TempDir()
delegate := &scriptedSink{rawErrors: []error{errSpoolTest}}
sink := NewDurableSink(delegate, DurableConfig{Directory: dir})
env := durableTestEnvelope()
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 delegate.rawCalls != 1 {
t.Fatalf("raw calls = %d, want 1", delegate.rawCalls)
}
if delegate.unifiedCalls != 0 {
t.Fatalf("unified should not be delegated while raw is spooled, calls = %d", delegate.unifiedCalls)
}
files := spoolFiles(t, dir)
if len(files) != 2 {
t.Fatalf("spool files = %d, want 2: %#v", len(files), files)
}
assertSpoolKind(t, files[0], "raw")
assertSpoolKind(t, files[1], "unified")
}
func TestDurableSinkReplayPublishesInFileOrderAndDeletesFiles(t *testing.T) {
dir := t.TempDir()
delegate := &scriptedSink{rawErrors: []error{errSpoolTest}}
sink := NewDurableSink(delegate, DurableConfig{Directory: dir})
env := durableTestEnvelope()
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)
}
delegate.rawErrors = nil
if err := sink.ReplayOnce(context.Background()); err != nil {
t.Fatalf("ReplayOnce() error = %v", err)
}
got := delegate.calls
want := []string{"raw", "raw", "unified"}
if len(got) != len(want) {
t.Fatalf("calls = %#v, want %#v", got, want)
}
for i := range want {
if got[i] != want[i] {
t.Fatalf("calls[%d] = %q, want %q; all=%#v", i, got[i], want[i], got)
}
}
if files := spoolFiles(t, dir); len(files) != 0 {
t.Fatalf("spool files after replay = %#v, want none", files)
}
}
func durableTestEnvelope() envelope.FrameEnvelope {
return envelope.FrameEnvelope{
Protocol: envelope.ProtocolJT808,
MessageID: "0x0200",
Phone: "013079963379",
VIN: "LKLG7C4E3NA774736",
Sequence: 183,
EventTimeMS: 1782903940000,
ReceivedAtMS: 1782917751903,
RawHex: "0200002201307996337900B7",
}
}
func spoolFiles(t *testing.T, dir string) []string {
t.Helper()
matches, err := filepath.Glob(filepath.Join(dir, "*.json"))
if err != nil {
t.Fatalf("glob spool files: %v", err)
}
return matches
}
func assertSpoolKind(t *testing.T, path string, kind string) {
t.Helper()
payload, err := os.ReadFile(path)
if err != nil {
t.Fatalf("read spool file: %v", err)
}
if !containsString(string(payload), `"kind":"`+kind+`"`) {
t.Fatalf("spool file %s payload = %s, want kind %q", path, payload, kind)
}
}
func containsString(value string, needle string) bool {
return len(needle) == 0 || (len(value) >= len(needle) && indexString(value, needle) >= 0)
}
func indexString(value string, needle string) int {
for i := 0; i+len(needle) <= len(value); i++ {
if value[i:i+len(needle)] == needle {
return i
}
}
return -1
}
var errSpoolTest = errors.New("delegate unavailable")
type scriptedSink struct {
rawErrors []error
unifiedErrors []error
rawCalls int
unifiedCalls int
calls []string
}
func (s *scriptedSink) PublishRaw(context.Context, envelope.FrameEnvelope) error {
s.rawCalls++
s.calls = append(s.calls, "raw")
if len(s.rawErrors) > 0 {
err := s.rawErrors[0]
s.rawErrors = s.rawErrors[1:]
return err
}
return nil
}
func (s *scriptedSink) PublishUnified(context.Context, envelope.FrameEnvelope) error {
s.unifiedCalls++
s.calls = append(s.calls, "unified")
if len(s.unifiedErrors) > 0 {
err := s.unifiedErrors[0]
s.unifiedErrors = s.unifiedErrors[1:]
return err
}
return nil
}
func (s *scriptedSink) Close() error {
return nil
}