diff --git a/go/vehicle-gateway/cmd/history-writer/main.go b/go/vehicle-gateway/cmd/history-writer/main.go index 7a99a4a2..3017381f 100644 --- a/go/vehicle-gateway/cmd/history-writer/main.go +++ b/go/vehicle-gateway/cmd/history-writer/main.go @@ -8,6 +8,7 @@ import ( "os/signal" "strings" "syscall" + "time" "github.com/segmentio/kafka-go" _ "github.com/taosdata/driver-go/v3/taosWS" @@ -65,19 +66,39 @@ func main() { logger.Error("kafka fetch failed", "error", err) continue } - var env envelope.FrameEnvelope - if err := json.Unmarshal(message.Value, &env); err != nil { - logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - _ = reader.CommitMessages(ctx, message) - continue - } - if err := writer.AppendAll(ctx, env); err != nil { - logger.Error("tdengine append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) - continue - } - if err := reader.CommitMessages(ctx, message); err != nil { - logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - } + processHistoryMessage(ctx, logger, writer, reader, message) + } +} + +const kafkaMessageOperationTimeout = 30 * time.Second + +type historyAppender interface { + AppendAll(context.Context, envelope.FrameEnvelope) error +} + +type kafkaMessageCommitter interface { + CommitMessages(context.Context, ...kafka.Message) error +} + +func processHistoryMessage(ctx context.Context, logger interface { + Error(string, ...any) + Warn(string, ...any) +}, appender historyAppender, committer kafkaMessageCommitter, message kafka.Message) { + messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) + defer cancel() + + var env envelope.FrameEnvelope + if err := json.Unmarshal(message.Value, &env); err != nil { + logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + _ = committer.CommitMessages(messageCtx, message) + return + } + if err := appender.AppendAll(messageCtx, env); err != nil { + logger.Error("tdengine append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) + return + } + if err := committer.CommitMessages(messageCtx, message); err != nil { + logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) } } diff --git a/go/vehicle-gateway/cmd/history-writer/main_test.go b/go/vehicle-gateway/cmd/history-writer/main_test.go new file mode 100644 index 00000000..e6479c4b --- /dev/null +++ b/go/vehicle-gateway/cmd/history-writer/main_test.go @@ -0,0 +1,70 @@ +package main + +import ( + "context" + "encoding/json" + "testing" + + "github.com/segmentio/kafka-go" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestProcessHistoryMessageUsesUncancelledContextForAppendAndCommit(t *testing.T) { + parent, cancel := context.WithCancel(context.Background()) + cancel() + + env := envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "13307795425", + EventTimeMS: 1782918600000, + ReceivedAtMS: 1782918600000, + ParseStatus: envelope.ParseOK, + } + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + appender := &contextCheckingHistoryAppender{} + committer := &contextCheckingHistoryCommitter{} + + processHistoryMessage(parent, discardHistoryLogger{}, appender, committer, kafka.Message{Value: payload}) + + if appender.ctxErr != nil { + t.Fatalf("appender saw cancelled context: %v", appender.ctxErr) + } + if committer.ctxErr != nil { + t.Fatalf("committer saw cancelled context: %v", committer.ctxErr) + } + if appender.count != 1 || committer.count != 1 { + t.Fatalf("appends=%d commits=%d", appender.count, committer.count) + } +} + +type contextCheckingHistoryAppender struct { + ctxErr error + count int +} + +func (a *contextCheckingHistoryAppender) AppendAll(ctx context.Context, _ envelope.FrameEnvelope) error { + a.ctxErr = ctx.Err() + a.count++ + return a.ctxErr +} + +type contextCheckingHistoryCommitter struct { + ctxErr error + count int +} + +func (c *contextCheckingHistoryCommitter) CommitMessages(ctx context.Context, _ ...kafka.Message) error { + c.ctxErr = ctx.Err() + c.count++ + return c.ctxErr +} + +type discardHistoryLogger struct{} + +func (discardHistoryLogger) Error(string, ...any) {} +func (discardHistoryLogger) Warn(string, ...any) {} diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 273aca8f..e55f56d3 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -150,19 +150,39 @@ func consumeKafka(ctx context.Context, logger interface { logger.Error("kafka fetch failed", "error", err) continue } - var env envelope.FrameEnvelope - if err := json.Unmarshal(message.Value, &env); err != nil { - logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - _ = reader.CommitMessages(ctx, message) - continue - } - if err := repository.Update(ctx, env); err != nil { - logger.Error("redis realtime update failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) - continue - } - if err := reader.CommitMessages(ctx, message); err != nil { - logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - } + processRealtimeMessage(ctx, logger, repository, reader, message) + } +} + +const kafkaMessageOperationTimeout = 30 * time.Second + +type realtimeUpdater interface { + Update(context.Context, envelope.FrameEnvelope) error +} + +type kafkaMessageCommitter interface { + CommitMessages(context.Context, ...kafka.Message) error +} + +func processRealtimeMessage(ctx context.Context, logger interface { + Error(string, ...any) + Warn(string, ...any) +}, updater realtimeUpdater, committer kafkaMessageCommitter, message kafka.Message) { + messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) + defer cancel() + + var env envelope.FrameEnvelope + if err := json.Unmarshal(message.Value, &env); err != nil { + logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + _ = committer.CommitMessages(messageCtx, message) + return + } + if err := updater.Update(messageCtx, env); err != nil { + logger.Error("redis realtime update failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) + return + } + if err := committer.CommitMessages(messageCtx, message); err != nil { + logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) } } diff --git a/go/vehicle-gateway/cmd/realtime-api/main_test.go b/go/vehicle-gateway/cmd/realtime-api/main_test.go new file mode 100644 index 00000000..d7102f7a --- /dev/null +++ b/go/vehicle-gateway/cmd/realtime-api/main_test.go @@ -0,0 +1,71 @@ +package main + +import ( + "context" + "encoding/json" + "testing" + + "github.com/segmentio/kafka-go" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestProcessRealtimeMessageUsesUncancelledContextForUpdateAndCommit(t *testing.T) { + parent, cancel := context.WithCancel(context.Background()) + cancel() + + env := envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "13307795425", + EventTimeMS: 1782918600000, + ReceivedAtMS: 1782918600000, + ParseStatus: envelope.ParseOK, + } + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + updater := &contextCheckingRealtimeUpdater{} + committer := &contextCheckingMessageCommitter{} + + processRealtimeMessage(parent, discardRealtimeLogger{}, updater, committer, kafka.Message{Value: payload}) + + if updater.ctxErr != nil { + t.Fatalf("updater saw cancelled context: %v", updater.ctxErr) + } + if committer.ctxErr != nil { + t.Fatalf("committer saw cancelled context: %v", committer.ctxErr) + } + if updater.count != 1 || committer.count != 1 { + t.Fatalf("updates=%d commits=%d", updater.count, committer.count) + } +} + +type contextCheckingRealtimeUpdater struct { + ctxErr error + count int +} + +func (u *contextCheckingRealtimeUpdater) Update(ctx context.Context, _ envelope.FrameEnvelope) error { + u.ctxErr = ctx.Err() + u.count++ + return u.ctxErr +} + +type contextCheckingMessageCommitter struct { + ctxErr error + count int +} + +func (c *contextCheckingMessageCommitter) CommitMessages(ctx context.Context, _ ...kafka.Message) error { + c.ctxErr = ctx.Err() + c.count++ + return c.ctxErr +} + +type discardRealtimeLogger struct{} + +func (discardRealtimeLogger) Info(string, ...any) {} +func (discardRealtimeLogger) Error(string, ...any) {} +func (discardRealtimeLogger) Warn(string, ...any) {} diff --git a/go/vehicle-gateway/cmd/stat-writer/main.go b/go/vehicle-gateway/cmd/stat-writer/main.go index fb6716d4..b11de4c2 100644 --- a/go/vehicle-gateway/cmd/stat-writer/main.go +++ b/go/vehicle-gateway/cmd/stat-writer/main.go @@ -62,19 +62,39 @@ func main() { logger.Error("kafka fetch failed", "error", err) continue } - var env envelope.FrameEnvelope - if err := json.Unmarshal(message.Value, &env); err != nil { - logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - _ = reader.CommitMessages(ctx, message) - continue - } - if err := writer.Append(ctx, env); err != nil { - logger.Error("mysql append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) - continue - } - if err := reader.CommitMessages(ctx, message); err != nil { - logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) - } + processStatMessage(ctx, logger, writer, reader, message) + } +} + +const kafkaMessageOperationTimeout = 30 * time.Second + +type statAppender interface { + Append(context.Context, envelope.FrameEnvelope) error +} + +type kafkaMessageCommitter interface { + CommitMessages(context.Context, ...kafka.Message) error +} + +func processStatMessage(ctx context.Context, logger interface { + Error(string, ...any) + Warn(string, ...any) +}, appender statAppender, committer kafkaMessageCommitter, message kafka.Message) { + messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) + defer cancel() + + var env envelope.FrameEnvelope + if err := json.Unmarshal(message.Value, &env); err != nil { + logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + _ = committer.CommitMessages(messageCtx, message) + return + } + if err := appender.Append(messageCtx, env); err != nil { + logger.Error("mysql append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) + return + } + if err := committer.CommitMessages(messageCtx, message); err != nil { + logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) } } diff --git a/go/vehicle-gateway/cmd/stat-writer/main_test.go b/go/vehicle-gateway/cmd/stat-writer/main_test.go new file mode 100644 index 00000000..d595f844 --- /dev/null +++ b/go/vehicle-gateway/cmd/stat-writer/main_test.go @@ -0,0 +1,70 @@ +package main + +import ( + "context" + "encoding/json" + "testing" + + "github.com/segmentio/kafka-go" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestProcessStatMessageUsesUncancelledContextForAppendAndCommit(t *testing.T) { + parent, cancel := context.WithCancel(context.Background()) + cancel() + + env := envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "13307795425", + EventTimeMS: 1782918600000, + ReceivedAtMS: 1782918600000, + ParseStatus: envelope.ParseOK, + } + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + appender := &contextCheckingStatAppender{} + committer := &contextCheckingStatCommitter{} + + processStatMessage(parent, discardStatLogger{}, appender, committer, kafka.Message{Value: payload}) + + if appender.ctxErr != nil { + t.Fatalf("appender saw cancelled context: %v", appender.ctxErr) + } + if committer.ctxErr != nil { + t.Fatalf("committer saw cancelled context: %v", committer.ctxErr) + } + if appender.count != 1 || committer.count != 1 { + t.Fatalf("appends=%d commits=%d", appender.count, committer.count) + } +} + +type contextCheckingStatAppender struct { + ctxErr error + count int +} + +func (a *contextCheckingStatAppender) Append(ctx context.Context, _ envelope.FrameEnvelope) error { + a.ctxErr = ctx.Err() + a.count++ + return a.ctxErr +} + +type contextCheckingStatCommitter struct { + ctxErr error + count int +} + +func (c *contextCheckingStatCommitter) CommitMessages(ctx context.Context, _ ...kafka.Message) error { + c.ctxErr = ctx.Err() + c.count++ + return c.ctxErr +} + +type discardStatLogger struct{} + +func (discardStatLogger) Error(string, ...any) {} +func (discardStatLogger) Warn(string, ...any) {}