From 6bee0499ae115a04f6a80701a0fe9d6e903a61e5 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 23:39:10 +0800 Subject: [PATCH] feat: expose raw frame query api --- deploy/portainer/docker-compose-go.yml | 3 + go/vehicle-gateway/cmd/realtime-api/main.go | 28 ++ go/vehicle-gateway/internal/history/query.go | 327 ++++++++++++++++++ .../internal/history/query_test.go | 122 +++++++ 4 files changed, 480 insertions(+) create mode 100644 go/vehicle-gateway/internal/history/query.go create mode 100644 go/vehicle-gateway/internal/history/query_test.go diff --git a/deploy/portainer/docker-compose-go.yml b/deploy/portainer/docker-compose-go.yml index 4bbd5ec8..9abbef86 100644 --- a/deploy/portainer/docker-compose-go.yml +++ b/deploy/portainer/docker-compose-go.yml @@ -96,6 +96,9 @@ services: REDIS_DB: ${REDIS_DB:-50} ONLINE_TTL_SECONDS: ${ONLINE_TTL_SECONDS:-600} MYSQL_DSN: ${MYSQL_DSN:-} + TDENGINE_DRIVER: ${TDENGINE_DRIVER:-taosWS} + TDENGINE_DSN: ${TDENGINE_DSN:-} + TDENGINE_DATABASE: ${TDENGINE_DATABASE:-lingniu_vehicle_ts} ports: - "${GO_REALTIME_HTTP_PORT:-20210}:20210" diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 8d2dba67..a98411bc 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -15,8 +15,10 @@ import ( _ "github.com/go-sql-driver/mysql" "github.com/redis/go-redis/v9" "github.com/segmentio/kafka-go" + _ "github.com/taosdata/driver-go/v3/taosWS" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/history" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats" @@ -74,6 +76,32 @@ func main() { logger.Warn("MYSQL_DSN is empty; stats query api disabled") } defer closeStats() + closeHistory := func() {} + if dsn := strings.TrimSpace(os.Getenv("TDENGINE_DSN")); dsn != "" { + driver := env("TDENGINE_DRIVER", "taosWS") + db, err := sql.Open(driver, dsn) + if err != nil { + logger.Error("history tdengine open failed", "driver", driver, "error", err) + os.Exit(1) + } + if err := db.PingContext(ctx); err != nil { + _ = db.Close() + logger.Error("history tdengine ping failed", "driver", driver, "error", err) + os.Exit(1) + } + closeHistory = func() { _ = db.Close() } + database := env("TDENGINE_DATABASE", history.DefaultDatabase) + mux.Handle("/api/history/raw-frames", history.NewRawFrameHandler(history.NewRawFrameRepository(db, database))) + logger.Info("history raw frame query enabled", "driver", driver, "database", database) + } else { + mux.HandleFunc("/api/history/raw-frames", func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusServiceUnavailable) + _ = json.NewEncoder(w).Encode(map[string]any{"error": "TDENGINE_DSN is not configured"}) + }) + logger.Warn("TDENGINE_DSN is empty; raw frame query api disabled") + } + defer closeHistory() server := &http.Server{ Addr: env("HTTP_ADDR", ":20210"), diff --git a/go/vehicle-gateway/internal/history/query.go b/go/vehicle-gateway/internal/history/query.go new file mode 100644 index 00000000..93a28877 --- /dev/null +++ b/go/vehicle-gateway/internal/history/query.go @@ -0,0 +1,327 @@ +package history + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "net/http" + "strconv" + "strings" + "time" +) + +type Queryer interface { + QueryContext(context.Context, string, ...any) (*sql.Rows, error) +} + +type RawFrameQuery struct { + Protocol string + VIN string + Phone string + DeviceID string + MessageID string + DateFrom string + DateTo string + Limit int + Offset int +} + +type RawFrameRow struct { + TS string `json:"ts"` + FrameID string `json:"frame_id"` + EventID string `json:"event_id"` + MessageID int64 `json:"message_id"` + MessageIDHex string `json:"message_id_hex"` + EventTime string `json:"event_time"` + ReceivedAt string `json:"received_at"` + RawSizeBytes int64 `json:"raw_size_bytes"` + RawHex string `json:"raw_hex,omitempty"` + RawText string `json:"raw_text,omitempty"` + ParsedJSON string `json:"parsed_json,omitempty"` + FieldsJSON string `json:"fields_json,omitempty"` + ParseStatus string `json:"parse_status"` + ParseError string `json:"parse_error,omitempty"` + SourceEndpoint string `json:"source_endpoint"` + Protocol string `json:"protocol"` + VehicleKey string `json:"vehicle_key"` + VIN string `json:"vin"` + Phone string `json:"phone,omitempty"` + DeviceID string `json:"device_id,omitempty"` +} + +type RawFrameRepository struct { + db Queryer + database string +} + +func NewRawFrameRepository(db Queryer, database string) *RawFrameRepository { + if db == nil { + panic("raw frame query db must not be nil") + } + database = strings.TrimSpace(database) + if database != "" && !safeIdentifier(database) { + database = "" + } + return &RawFrameRepository{db: db, database: database} +} + +func (r *RawFrameRepository) Query(ctx context.Context, query RawFrameQuery) ([]RawFrameRow, error) { + query = normalizeRawFrameQuery(query) + sqlText, args := buildRawFrameSQL(r.tableName(), query) + rows, err := r.db.QueryContext(ctx, sqlText, args...) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []RawFrameRow + for rows.Next() { + var row RawFrameRow + var ts scanDateTime + var eventTime scanDateTime + var receivedAt scanDateTime + if err := rows.Scan( + &ts, + &row.FrameID, + &row.EventID, + &row.MessageID, + &eventTime, + &receivedAt, + &row.RawSizeBytes, + &row.RawHex, + &row.RawText, + &row.ParsedJSON, + &row.FieldsJSON, + &row.ParseStatus, + &row.ParseError, + &row.SourceEndpoint, + &row.Protocol, + &row.VehicleKey, + &row.VIN, + &row.Phone, + &row.DeviceID, + ); err != nil { + return nil, err + } + row.TS = ts.String + row.EventTime = eventTime.String + row.ReceivedAt = receivedAt.String + row.MessageIDHex = fmt.Sprintf("0x%04X", row.MessageID) + out = append(out, row) + } + return out, rows.Err() +} + +func (r *RawFrameRepository) tableName() string { + if r.database == "" { + return "raw_frames" + } + return r.database + ".raw_frames" +} + +func normalizeRawFrameQuery(query RawFrameQuery) RawFrameQuery { + query.Protocol = strings.ToUpper(strings.TrimSpace(query.Protocol)) + query.VIN = strings.TrimSpace(query.VIN) + query.Phone = strings.TrimSpace(query.Phone) + query.DeviceID = strings.TrimSpace(query.DeviceID) + query.MessageID = strings.TrimSpace(query.MessageID) + query.DateFrom = strings.TrimSpace(query.DateFrom) + query.DateTo = strings.TrimSpace(query.DateTo) + if query.Limit <= 0 { + query.Limit = 20 + } + return query +} + +func buildRawFrameSQL(table string, query RawFrameQuery) (string, []any) { + var where []string + var args []any + add := func(clause string, value any) { + where = append(where, clause) + args = append(args, value) + } + if query.Protocol != "" { + add("protocol = ?", query.Protocol) + } + if query.VIN != "" { + add("vin = ?", query.VIN) + } + if query.Phone != "" { + add("phone = ?", query.Phone) + } + if query.DeviceID != "" { + add("device_id = ?", query.DeviceID) + } + if query.MessageID != "" { + if parsed, ok := parseMessageID(query.MessageID); ok { + add("message_id = ?", parsed) + } + } + if query.DateFrom != "" { + add("ts >= ?", query.DateFrom) + } + if query.DateTo != "" { + add("ts <= ?", query.DateTo) + } + sqlText := `SELECT ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, raw_hex, raw_text, parsed_json, fields_json, parse_status, parse_error, source_endpoint, protocol, vehicle_key, vin, phone, device_id FROM ` + table + if len(where) > 0 { + sqlText += " WHERE " + strings.Join(where, " AND ") + } + sqlText += " ORDER BY ts DESC LIMIT ? OFFSET ?" + args = append(args, query.Limit, query.Offset) + return sqlText, args +} + +type RawFrameHandler struct { + repository *RawFrameRepository +} + +func NewRawFrameHandler(repository *RawFrameRepository) *RawFrameHandler { + if repository == nil { + panic("raw frame repository must not be nil") + } + return &RawFrameHandler{repository: repository} +} + +func (h *RawFrameHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + writeHistoryError(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + if strings.Trim(r.URL.Path, "/") != "api/history/raw-frames" { + writeHistoryError(w, http.StatusNotFound, "route not found") + return + } + query, err := parseRawFrameQuery(r) + if err != nil { + writeHistoryError(w, http.StatusBadRequest, err.Error()) + return + } + rows, err := h.repository.Query(r.Context(), query) + if err != nil { + writeHistoryError(w, http.StatusInternalServerError, err.Error()) + return + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{ + "items": rows, + "total": len(rows), + "limit": query.Limit, + "offset": query.Offset, + }) +} + +func parseRawFrameQuery(r *http.Request) (RawFrameQuery, error) { + values := r.URL.Query() + limit, err := parseBoundedInt(values.Get("limit"), 20, 1, 500, "limit") + if err != nil { + return RawFrameQuery{}, err + } + offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset") + if err != nil { + return RawFrameQuery{}, err + } + query := RawFrameQuery{ + Protocol: values.Get("protocol"), + VIN: values.Get("vin"), + Phone: values.Get("phone"), + DeviceID: values.Get("deviceId"), + MessageID: values.Get("messageId"), + DateFrom: values.Get("dateFrom"), + DateTo: values.Get("dateTo"), + Limit: limit, + Offset: offset, + } + if !validDateTime(query.DateFrom) || !validDateTime(query.DateTo) { + return RawFrameQuery{}, errors.New("dateFrom/dateTo must use YYYY-MM-DD or YYYY-MM-DD HH:mm:ss") + } + if strings.TrimSpace(query.MessageID) != "" { + if _, ok := parseMessageID(query.MessageID); !ok { + return RawFrameQuery{}, errors.New("messageId must be decimal or hex like 0x0200") + } + } + return normalizeRawFrameQuery(query), nil +} + +func parseBoundedInt(raw string, fallback int, min int, max int, name string) (int, error) { + raw = strings.TrimSpace(raw) + if raw == "" { + return fallback, nil + } + value, err := strconv.Atoi(raw) + if err != nil || value < min || value > max { + return 0, fmt.Errorf("%s must be between %d and %d", name, min, max) + } + return value, nil +} + +func parseMessageID(value string) (int64, bool) { + value = strings.TrimSpace(strings.ToLower(value)) + if value == "" { + return 0, false + } + if strings.HasPrefix(value, "0x") { + parsed, err := strconv.ParseInt(strings.TrimPrefix(value, "0x"), 16, 32) + return parsed, err == nil + } + parsed, err := strconv.ParseInt(value, 10, 32) + return parsed, err == nil +} + +func validDateTime(value string) bool { + value = strings.TrimSpace(value) + if value == "" { + return true + } + for _, layout := range []string{"2006-01-02", "2006-01-02 15:04:05", time.RFC3339} { + if _, err := time.Parse(layout, value); err == nil { + return true + } + } + return false +} + +func safeIdentifier(value string) bool { + if value == "" { + return false + } + for _, r := range value { + if (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') || r == '_' { + continue + } + return false + } + return true +} + +type scanDateTime struct { + String string +} + +func (s *scanDateTime) Scan(value any) error { + s.String = formatSQLTime(value, "2006-01-02 15:04:05") + return nil +} + +func formatSQLTime(value any, layout string) string { + switch typed := value.(type) { + case nil: + return "" + case time.Time: + return typed.Format(layout) + case []byte: + return strings.TrimSpace(string(typed)) + case string: + return strings.TrimSpace(typed) + default: + return strings.TrimSpace(fmt.Sprint(typed)) + } +} + +func writeHistoryError(w http.ResponseWriter, status int, message string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(map[string]any{"error": message}) +} diff --git a/go/vehicle-gateway/internal/history/query_test.go b/go/vehicle-gateway/internal/history/query_test.go new file mode 100644 index 00000000..15c6c994 --- /dev/null +++ b/go/vehicle-gateway/internal/history/query_test.go @@ -0,0 +1,122 @@ +package history + +import ( + "context" + "database/sql" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/DATA-DOG/go-sqlmock" +) + +func TestRawFrameRepositoryQueriesRawFramesWithFilters(t *testing.T) { + db, mock, err := sqlmock.New() + if err != nil { + t.Fatalf("sqlmock.New() error = %v", err) + } + defer db.Close() + mock.ExpectQuery("SELECT ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, raw_hex, raw_text, parsed_json, fields_json, parse_status, parse_error, source_endpoint, protocol, vehicle_key, vin, phone, device_id FROM lingniu_vehicle_ts.raw_frames"). + WithArgs("JT808", "LKLG7C4E3NA774736", "2026-07-01 00:00:00", "2026-07-01 23:59:59", 20, 0). + WillReturnRows(sqlmock.NewRows([]string{ + "ts", "frame_id", "event_id", "message_id", "event_time", "received_at", "raw_size_bytes", + "raw_hex", "raw_text", "parsed_json", "fields_json", "parse_status", "parse_error", "source_endpoint", + "protocol", "vehicle_key", "vin", "phone", "device_id", + }).AddRow( + time.Date(2026, 7, 1, 23, 25, 36, 0, time.FixedZone("Asia/Shanghai", 8*3600)), + "go_frame", "event-1", 0x0200, + time.Date(2026, 7, 1, 23, 25, 36, 0, time.FixedZone("Asia/Shanghai", 8*3600)), + time.Date(2026, 7, 1, 23, 25, 37, 0, time.FixedZone("Asia/Shanghai", 8*3600)), + 64, "7E0200", "", `{"header":{"message_id":"0x0200"}}`, `{"speed_kmh":0}`, + "OK", "", "222.66.200.68:29646", "JT808", "013079963379", "LKLG7C4E3NA774736", "013079963379", "9963379", + )) + + repository := NewRawFrameRepository(db, "lingniu_vehicle_ts") + rows, err := repository.Query(context.Background(), RawFrameQuery{ + Protocol: "JT808", + VIN: "LKLG7C4E3NA774736", + DateFrom: "2026-07-01 00:00:00", + DateTo: "2026-07-01 23:59:59", + Limit: 20, + }) + if err != nil { + t.Fatalf("Query() error = %v", err) + } + if len(rows) != 1 { + t.Fatalf("row count = %d", len(rows)) + } + if rows[0].MessageID != 512 || rows[0].MessageIDHex != "0x0200" || rows[0].ParsedJSON == "" { + t.Fatalf("unexpected row: %#v", rows[0]) + } + if rows[0].TS != "2026-07-01 23:25:36" { + t.Fatalf("timestamp = %q", rows[0].TS) + } + if err := mock.ExpectationsWereMet(); err != nil { + t.Fatalf("sql expectations: %v", err) + } +} + +func TestRawFrameHandlerReturnsRawFrames(t *testing.T) { + db, mock, err := sqlmock.New() + if err != nil { + t.Fatalf("sqlmock.New() error = %v", err) + } + defer db.Close() + mock.ExpectQuery("SELECT ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, raw_hex, raw_text, parsed_json, fields_json, parse_status, parse_error, source_endpoint, protocol, vehicle_key, vin, phone, device_id FROM lingniu_vehicle_ts.raw_frames"). + WithArgs("GB32960", "LB9A32A21R0LS1707", 5, 0). + WillReturnRows(sqlmock.NewRows([]string{ + "ts", "frame_id", "event_id", "message_id", "event_time", "received_at", "raw_size_bytes", + "raw_hex", "raw_text", "parsed_json", "fields_json", "parse_status", "parse_error", "source_endpoint", + "protocol", "vehicle_key", "vin", "phone", "device_id", + }).AddRow( + "2026-07-01 22:28:25", "go_frame", "event-2", 2, "2026-07-01 22:28:25", "2026-07-01 22:28:25", + 128, "2323", "", `{"command":"REALTIME"}`, `{"total_mileage_km":53490.9}`, + "OK", "", "8.134.95.166:53702", "GB32960", "LB9A32A21R0LS1707", "LB9A32A21R0LS1707", "", "", + )) + + handler := NewRawFrameHandler(NewRawFrameRepository(db, "lingniu_vehicle_ts")) + request := httptest.NewRequest(http.MethodGet, "/api/history/raw-frames?protocol=GB32960&vin=LB9A32A21R0LS1707&limit=5", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", response.Code, response.Body.String()) + } + body := response.Body.String() + for _, want := range []string{`"vin":"LB9A32A21R0LS1707"`, `"message_id_hex":"0x0002"`, `"total":1`} { + if !strings.Contains(body, want) { + t.Fatalf("response missing %s: %s", want, body) + } + } + if err := mock.ExpectationsWereMet(); err != nil { + t.Fatalf("sql expectations: %v", err) + } +} + +func TestRawFrameHandlerRejectsInvalidLimit(t *testing.T) { + handler := NewRawFrameHandler(NewRawFrameRepository(&sql.DB{}, "")) + request := httptest.NewRequest(http.MethodGet, "/api/history/raw-frames?limit=501", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusBadRequest { + t.Fatalf("status = %d body=%s", response.Code, response.Body.String()) + } +} + +func TestParseMessageIDSupportsDecimalAndHex(t *testing.T) { + for raw, want := range map[string]int64{ + "512": 512, + "0x0200": 512, + "0X0100": 256, + } { + got, ok := parseMessageID(raw) + if !ok || got != want { + t.Fatalf("parseMessageID(%q) = %d,%v want %d,true", raw, got, ok, want) + } + } +}