Files
lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus/durable_sink.go
2026-07-02 13:29:49 +08:00

245 lines
5.6 KiB
Go

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
ReplayBatchSize int
}
type DurableSink struct {
delegate Sink
dir string
mu sync.Mutex
seq uint64
rawPending map[string]struct{}
replayBatchSize int
}
type durableRecord struct {
Kind string `json:"kind"`
Envelope envelope.FrameEnvelope `json:"envelope"`
}
type durableRecordFile struct {
path string
record durableRecord
}
const durableReplayOperationTimeout = 30 * time.Second
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{}{},
replayBatchSize: cfg.ReplayBatchSize,
}
}
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 {
return s.replay(ctx, 0)
}
func (s *DurableSink) replay(ctx context.Context, limit int) error {
files, err := durableFiles(s.dir, limit)
if err != nil {
return err
}
records := make([]durableRecordFile, 0, len(files))
for _, file := range files {
record, err := readDurableRecord(file)
if err != nil {
return err
}
records = append(records, durableRecordFile{path: file, record: record})
}
sortDurableRecords(records)
for _, item := range records {
if err := s.publishRecord(ctx, item.record); err != nil {
return err
}
if err := os.Remove(item.path); err != nil {
return err
}
if item.record.Kind == "raw" {
s.clearRawPending(item.record.Envelope)
}
}
return nil
}
func sortDurableRecords(records []durableRecordFile) {
sort.SliceStable(records, func(i, j int) bool {
leftEvent := records[i].record.Envelope.StableEventID()
rightEvent := records[j].record.Envelope.StableEventID()
if leftEvent == rightEvent {
return durableKindOrder(records[i].record.Kind) < durableKindOrder(records[j].record.Kind)
}
return records[i].path < records[j].path
})
}
func durableKindOrder(kind string) int {
if kind == "raw" {
return 0
}
if kind == "unified" {
return 1
}
return 2
}
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.replayOnceFromLoop(ctx); err != nil && onError != nil {
onError(err)
}
}
}
}
func (s *DurableSink) replayOnceFromLoop(ctx context.Context) error {
replayCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), durableReplayOperationTimeout)
defer cancel()
return s.replay(replayCtx, s.replayBatchSize)
}
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)
}
func durableFiles(dir string, limit int) ([]string, error) {
files, err := filepath.Glob(filepath.Join(dir, "*.json"))
if err != nil {
return nil, err
}
sort.Strings(files)
if limit > 0 && len(files) > limit {
files = files[:limit]
}
return files, nil
}