package main import ( "context" "errors" "strings" "testing" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestProcessFastMessageWritesTDengineAndRedisBeforeAck(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", EventID: "evt-1"} payload, err := env.MarshalJSONBytes() if err != nil { t.Fatal(err) } appender := &recordingFastAppender{} updater := &recordingFastUpdater{} ackCount := 0 msg := &fastMessage{data: payload, ack: func() error { ackCount++ return nil }} if err := processFastMessage(context.Background(), nil, appender, updater, msg); err != nil { t.Fatalf("processFastMessage() error = %v", err) } if appender.count != 1 || updater.count != 1 || ackCount != 1 { t.Fatalf("appends=%d updates=%d acks=%d", appender.count, updater.count, ackCount) } } func TestProcessFastMessageDoesNotAckWhenRedisFails(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN001", EventID: "evt-2"} payload, err := env.MarshalJSONBytes() if err != nil { t.Fatal(err) } appender := &recordingFastAppender{} updater := &recordingFastUpdater{err: errors.New("redis down")} ackCount := 0 msg := &fastMessage{data: payload, ack: func() error { ackCount++ return nil }} if err := processFastMessage(context.Background(), nil, appender, updater, msg); err == nil { t.Fatal("processFastMessage() error = nil, want redis failure") } if ackCount != 0 { t.Fatalf("ack count = %d, want 0", ackCount) } } func TestProcessFastMessageRecordsStageDurationMetrics(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN001", EventID: "evt-3"} payload, err := env.MarshalJSONBytes() if err != nil { t.Fatal(err) } registry := metrics.NewRegistry() msg := &fastMessage{subject: "vehicle.raw.go.jt808.v1", data: payload, ack: func() error { return nil }} if err := processFastMessage(context.Background(), registry, &recordingFastAppender{}, &recordingFastUpdater{}, msg); err != nil { t.Fatalf("processFastMessage() error = %v", err) } text := registry.Render() for _, want := range []string{ `vehicle_fast_writer_stage_duration_ms_histogram_bucket{le="+Inf",stage="tdengine",status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, `vehicle_fast_writer_stage_duration_ms_histogram_count{stage="tdengine",status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, `vehicle_fast_writer_stage_duration_ms_histogram_bucket{le="+Inf",stage="redis",status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, `vehicle_fast_writer_stage_duration_ms_histogram_bucket{le="+Inf",stage="ack",status="ok",subject="vehicle.raw.go.jt808.v1"} 1`, } { if !strings.Contains(text, want) { t.Fatalf("fast writer stage metric missing %s:\n%s", want, text) } } } type recordingFastAppender struct { count int err error } func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelope) error { a.count++ return a.err } type recordingFastUpdater struct { count int err error } func (u *recordingFastUpdater) FastUpdate(context.Context, envelope.FrameEnvelope) error { u.count++ return u.err }