Compare commits
2 Commits
4a32e2181a
...
c803a19bf5
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c803a19bf5 | ||
|
|
243de5c2b2 |
@@ -308,3 +308,76 @@ total_mileage_km=0
|
|||||||
```
|
```
|
||||||
|
|
||||||
说明:当前重试层只能覆盖 Kafka 短暂抖动;如果 Kafka 长时间不可用,仍需后续增加本地磁盘 spool/WAL 和恢复补发。
|
说明:当前重试层只能覆盖 Kafka 短暂抖动;如果 Kafka 长时间不可用,仍需后续增加本地磁盘 spool/WAL 和恢复补发。
|
||||||
|
|
||||||
|
## 2026-07-01 23:02 切换 Go 到生产 32960/808 端口
|
||||||
|
|
||||||
|
本轮操作:
|
||||||
|
|
||||||
|
- 停止 ECS 上原有 Java Docker 容器:
|
||||||
|
- `gb32960-ingest-app`
|
||||||
|
- `jt808-ingest-app`
|
||||||
|
- `yutong-mqtt-app`
|
||||||
|
- `vehicle-history-app`
|
||||||
|
- `vehicle-analytics-app`
|
||||||
|
- `vehicle-state-app`
|
||||||
|
- `telemetry-field-parser-app`
|
||||||
|
- 保留 Go sidecar 和 Portainer agent。
|
||||||
|
- 修改 `/opt/lingniu-go/.env`:
|
||||||
|
- `GO_GB32960_TCP_PORT=32960`
|
||||||
|
- `GO_JT808_TCP_PORT=808`
|
||||||
|
- Go gateway 使用镜像:
|
||||||
|
- `crpi-85r4m0ackrm3qpje.cn-shanghai.personal.cr.aliyuncs.com/oneos/vehicle-gateway-go:go-243de5c-20260701230028`
|
||||||
|
- 开启 Kafka 本地 spool:
|
||||||
|
- host:`/opt/lingniu-go/spool/gateway`
|
||||||
|
- container:`/data/spool/gateway`
|
||||||
|
|
||||||
|
端口验证:
|
||||||
|
|
||||||
|
```text
|
||||||
|
0.0.0.0:32960 -> go-vehicle-gateway:32960
|
||||||
|
0.0.0.0:808 -> go-vehicle-gateway:808
|
||||||
|
18080/23296 -> 已释放
|
||||||
|
```
|
||||||
|
|
||||||
|
启动日志确认:
|
||||||
|
|
||||||
|
```text
|
||||||
|
kafka durable spool enabled dir=/data/spool/gateway replay_interval_ms=1000
|
||||||
|
tcp listener started protocol=GB32960 addr=[::]:32960
|
||||||
|
tcp listener started protocol=JT808 addr=[::]:808
|
||||||
|
```
|
||||||
|
|
||||||
|
真实连接确认:
|
||||||
|
|
||||||
|
```text
|
||||||
|
GB32960 external remote: 8.134.95.166:50084
|
||||||
|
JT808 external remote: 115.231.168.135:13801
|
||||||
|
```
|
||||||
|
|
||||||
|
生产 `808` 回放复验:
|
||||||
|
|
||||||
|
```text
|
||||||
|
source_endpoint=172.20.0.1:54680
|
||||||
|
vin=LKLG7C4E3NA774736
|
||||||
|
phone=013079963379
|
||||||
|
longitude=119.557997
|
||||||
|
latitude=29.049092
|
||||||
|
total_mileage_km=0
|
||||||
|
```
|
||||||
|
|
||||||
|
TDengine 最新 RAW 复验显示生产 808 已有真实外部数据持续落地:
|
||||||
|
|
||||||
|
```text
|
||||||
|
2026-07-01 15:03:02.960 | JT808 | LKLG7C4E1NA774802 | 013079963291 | 222.66.200.68:28394
|
||||||
|
2026-07-01 15:03:02.114 | JT808 | | 14894135583 | 122.152.221.156:33871
|
||||||
|
2026-07-01 15:03:01.352 | JT808 | | 013307795519 | 115.231.168.135:13801
|
||||||
|
2026-07-01 15:03:01.330 | JT808 | LKLG7C4E3NA774736 | 013079963379 | 222.66.200.68:28393
|
||||||
|
```
|
||||||
|
|
||||||
|
Kafka spool 目录确认:
|
||||||
|
|
||||||
|
```text
|
||||||
|
/opt/lingniu-go/spool/gateway file count = 0
|
||||||
|
```
|
||||||
|
|
||||||
|
说明:生产 808 中仍有部分 phone 未解析出 VIN,后续需要继续维护 `vehicle_identity_binding` 的 phone/device_id/plate 到 VIN 映射,或者增强 808 注册/鉴权帧中的 identity 回写。
|
||||||
|
|||||||
@@ -186,7 +186,9 @@ go/vehicle-gateway/
|
|||||||
|
|
||||||
- Kafka 生产端必须开启同步写入,RAW 写成功后才允许写 unified event。
|
- Kafka 生产端必须开启同步写入,RAW 写成功后才允许写 unified event。
|
||||||
- Kafka 写入必须支持可配置重试、单次写超时和短退避,默认值为 `KAFKA_PUBLISH_ATTEMPTS=3`、`KAFKA_PUBLISH_TIMEOUT_MS=3000`、`KAFKA_PUBLISH_BACKOFF_MS=100`。
|
- Kafka 写入必须支持可配置重试、单次写超时和短退避,默认值为 `KAFKA_PUBLISH_ATTEMPTS=3`、`KAFKA_PUBLISH_TIMEOUT_MS=3000`、`KAFKA_PUBLISH_BACKOFF_MS=100`。
|
||||||
- 当前阶段的重试只解决短暂网络抖动;后续阶段需要增加本地磁盘 spool/WAL,使 Kafka 长时间不可用时 RAW 不丢失,并支持恢复后补发。
|
- gateway 支持本地磁盘 spool/WAL,配置 `KAFKA_SPOOL_DIR` 后启用;Kafka 长时间不可用时,发布失败的 envelope 先原子写入本地 JSON 文件。
|
||||||
|
- spool 补发按文件名顺序执行,补发成功后删除文件;默认补发间隔由 `KAFKA_SPOOL_REPLAY_INTERVAL_MS=1000` 控制。
|
||||||
|
- 如果 RAW 写 Kafka 失败并已落盘,同一个 event 的 unified event 不允许抢先写 Kafka,也必须进入 spool,确保恢复后按 RAW -> unified 顺序补发。
|
||||||
- Kafka topic 不允许自动创建,topic 和分区数由部署脚本或运维初始化,避免生产拼写错误造成隐性分流。
|
- Kafka topic 不允许自动创建,topic 和分区数由部署脚本或运维初始化,避免生产拼写错误造成隐性分流。
|
||||||
|
|
||||||
## TDengine 数据库设计
|
## TDengine 数据库设计
|
||||||
@@ -407,6 +409,7 @@ flowchart LR
|
|||||||
- 协议坏帧仍写 RAW topic,`parse_status=BAD_FRAME`。
|
- 协议坏帧仍写 RAW topic,`parse_status=BAD_FRAME`。
|
||||||
- 可部分解析的帧写 `parse_status=PARTIAL`,保留 parse error。
|
- 可部分解析的帧写 `parse_status=PARTIAL`,保留 parse error。
|
||||||
- Kafka 写失败时接入层先按配置重试;重试耗尽后不写 unified event,必须打错误日志和 metrics。
|
- Kafka 写失败时接入层先按配置重试;重试耗尽后不写 unified event,必须打错误日志和 metrics。
|
||||||
|
- 启用 `KAFKA_SPOOL_DIR` 后,Kafka 重试耗尽会落本地 spool;如果 spool 写失败,才视为接入层最终失败。
|
||||||
- TDengine 写失败不提交 Kafka offset。
|
- TDengine 写失败不提交 Kafka offset。
|
||||||
- MySQL 统计写失败不提交 Kafka offset。
|
- MySQL 统计写失败不提交 Kafka offset。
|
||||||
- Redis 写失败不影响历史和统计,但要通过 metrics 暴露。
|
- Redis 写失败不影响历史和统计,但要通过 metrics 暴露。
|
||||||
@@ -473,6 +476,7 @@ TDengine 连接:
|
|||||||
- Go gateway 能在 ECS 上接收真实 32960、808、宇通 MQTT 数据。
|
- Go gateway 能在 ECS 上接收真实 32960、808、宇通 MQTT 数据。
|
||||||
- 三种协议 RAW 都进入 Kafka。
|
- 三种协议 RAW 都进入 Kafka。
|
||||||
- Kafka 短暂写失败时,gateway 按配置重试;单元测试覆盖首次失败后成功和重试耗尽返回错误。
|
- Kafka 短暂写失败时,gateway 按配置重试;单元测试覆盖首次失败后成功和重试耗尽返回错误。
|
||||||
|
- Kafka 长时间不可用时,gateway 可把 RAW/unified 写入本地 spool;单元测试覆盖 RAW 已落盘时 unified 也落盘、恢复后按顺序补发并删除文件。
|
||||||
- 三种协议 RAW 和 parsed JSON 都进入 TDengine。
|
- 三种协议 RAW 和 parsed JSON 都进入 TDengine。
|
||||||
- 32960 和 808 的位置进入 `vehicle_locations`。
|
- 32960 和 808 的位置进入 `vehicle_locations`。
|
||||||
- 32960 和 808 的总里程采样进入 `vehicle_mileage_points`。
|
- 32960 和 808 的总里程采样进入 `vehicle_mileage_points`。
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ func main() {
|
|||||||
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||||
defer stop()
|
defer stop()
|
||||||
|
|
||||||
sink, err := buildSink(logger)
|
sink, err := buildSink(ctx, logger)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("build sink failed", "error", err)
|
logger.Error("build sink failed", "error", err)
|
||||||
os.Exit(1)
|
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
|
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"))
|
brokers := splitCSV(os.Getenv("KAFKA_BROKERS"))
|
||||||
if len(brokers) == 0 {
|
if len(brokers) == 0 {
|
||||||
logger.Warn("KAFKA_BROKERS is empty; using log sink")
|
logger.Warn("KAFKA_BROKERS is empty; using log sink")
|
||||||
@@ -150,11 +150,22 @@ func buildSink(logger *slog.Logger) (eventbus.Sink, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return eventbus.NewRetryingSink(sink, eventbus.RetryConfig{
|
var out eventbus.Sink = eventbus.NewRetryingSink(sink, eventbus.RetryConfig{
|
||||||
Attempts: envInt("KAFKA_PUBLISH_ATTEMPTS", 3),
|
Attempts: envInt("KAFKA_PUBLISH_ATTEMPTS", 3),
|
||||||
Backoff: time.Duration(envInt("KAFKA_PUBLISH_BACKOFF_MS", 100)) * time.Millisecond,
|
Backoff: time.Duration(envInt("KAFKA_PUBLISH_BACKOFF_MS", 100)) * time.Millisecond,
|
||||||
AttemptTimeout: time.Duration(envInt("KAFKA_PUBLISH_TIMEOUT_MS", 3000)) * 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 {
|
func env(key string, fallback string) string {
|
||||||
|
|||||||
187
go/vehicle-gateway/internal/eventbus/durable_sink.go
Normal file
187
go/vehicle-gateway/internal/eventbus/durable_sink.go
Normal 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)
|
||||||
|
}
|
||||||
153
go/vehicle-gateway/internal/eventbus/durable_sink_test.go
Normal file
153
go/vehicle-gateway/internal/eventbus/durable_sink_test.go
Normal 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
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user