912 lines
25 KiB
Go
912 lines
25 KiB
Go
package history
|
|
|
|
import (
|
|
"compress/gzip"
|
|
"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
|
|
VehicleKey string
|
|
VIN string
|
|
Phone string
|
|
DeviceID string
|
|
MessageID string
|
|
OrderBy string
|
|
IncludeFields bool
|
|
IncludePayload bool
|
|
IncludeTotal bool
|
|
ParsedFields []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"`
|
|
ParsedFields json.RawMessage `json:"parsed_fields,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 LocationQuery struct {
|
|
Protocol string
|
|
VIN string
|
|
IncludeTotal bool
|
|
DateFrom string
|
|
DateTo string
|
|
Limit int
|
|
Offset int
|
|
}
|
|
|
|
type LocationRow struct {
|
|
TS string `json:"ts"`
|
|
EventID string `json:"event_id"`
|
|
ReceivedAt string `json:"received_at"`
|
|
Longitude float64 `json:"longitude"`
|
|
Latitude float64 `json:"latitude"`
|
|
AltitudeM *float64 `json:"altitude_m,omitempty"`
|
|
SpeedKMH *float64 `json:"speed_kmh,omitempty"`
|
|
DirectionDeg *int64 `json:"direction_deg,omitempty"`
|
|
AlarmFlag *int64 `json:"alarm_flag,omitempty"`
|
|
StatusFlag *int64 `json:"status_flag,omitempty"`
|
|
TotalMileageKM *float64 `json:"total_mileage_km,omitempty"`
|
|
Protocol string `json:"protocol"`
|
|
VIN string `json:"vin"`
|
|
}
|
|
|
|
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}
|
|
}
|
|
|
|
type LocationRepository struct {
|
|
db Queryer
|
|
database string
|
|
}
|
|
|
|
func NewLocationRepository(db Queryer, database string) *LocationRepository {
|
|
if db == nil {
|
|
panic("location query db must not be nil")
|
|
}
|
|
database = strings.TrimSpace(database)
|
|
if database != "" && !safeIdentifier(database) {
|
|
database = ""
|
|
}
|
|
return &LocationRepository{db: db, database: database}
|
|
}
|
|
|
|
func (r *RawFrameRepository) Query(ctx context.Context, query RawFrameQuery) ([]RawFrameRow, error) {
|
|
query = normalizeRawFrameQuery(query)
|
|
sqlText, args := buildRawFrameSQL(r.tableName(query), query)
|
|
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := make([]RawFrameRow, 0)
|
|
for rows.Next() {
|
|
var row RawFrameRow
|
|
var ts scanDateTime
|
|
var eventTime scanDateTime
|
|
var receivedAt scanDateTime
|
|
var parsedFields string
|
|
if err := rows.Scan(
|
|
&ts,
|
|
&row.FrameID,
|
|
&row.EventID,
|
|
&row.MessageID,
|
|
&eventTime,
|
|
&receivedAt,
|
|
&row.RawSizeBytes,
|
|
&row.RawHex,
|
|
&row.RawText,
|
|
&parsedFields,
|
|
&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)
|
|
row.ParsedFields = rawJSONMessage(parsedFields)
|
|
out = append(out, row)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := r.hydratePayloadChunks(ctx, out); err != nil {
|
|
return nil, err
|
|
}
|
|
filterParsedFields(out, query.ParsedFields)
|
|
return out, nil
|
|
}
|
|
|
|
func (r *RawFrameRepository) Count(ctx context.Context, query RawFrameQuery) (int64, error) {
|
|
query = normalizeRawFrameQuery(query)
|
|
sqlText, args := buildRawFrameCountSQL(r.tableName(query), query)
|
|
return countRows(ctx, r.db, sqlText, args...)
|
|
}
|
|
|
|
func (r *LocationRepository) Query(ctx context.Context, query LocationQuery) ([]LocationRow, error) {
|
|
query = normalizeLocationQuery(query)
|
|
sqlText, args := buildLocationSQL(r.tableName(), query)
|
|
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
out := make([]LocationRow, 0)
|
|
for rows.Next() {
|
|
var row LocationRow
|
|
var ts scanDateTime
|
|
var receivedAt scanDateTime
|
|
var altitude sql.NullFloat64
|
|
var speed sql.NullFloat64
|
|
var direction sql.NullInt64
|
|
var alarm sql.NullInt64
|
|
var status sql.NullInt64
|
|
var mileage sql.NullFloat64
|
|
if err := rows.Scan(
|
|
&ts,
|
|
&row.EventID,
|
|
&receivedAt,
|
|
&row.Longitude,
|
|
&row.Latitude,
|
|
&altitude,
|
|
&speed,
|
|
&direction,
|
|
&alarm,
|
|
&status,
|
|
&mileage,
|
|
&row.Protocol,
|
|
&row.VIN,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
row.TS = ts.String
|
|
row.ReceivedAt = receivedAt.String
|
|
row.AltitudeM = nullableFloat(altitude)
|
|
row.SpeedKMH = nullableFloat(speed)
|
|
row.DirectionDeg = nullableInt(direction)
|
|
row.AlarmFlag = nullableInt(alarm)
|
|
row.StatusFlag = nullableInt(status)
|
|
row.TotalMileageKM = nullableFloat(mileage)
|
|
out = append(out, row)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
func (r *LocationRepository) Count(ctx context.Context, query LocationQuery) (int64, error) {
|
|
query = normalizeLocationQuery(query)
|
|
sqlText, args := buildLocationCountSQL(r.tableName(), query)
|
|
return countRows(ctx, r.db, sqlText, args...)
|
|
}
|
|
|
|
func countRows(ctx context.Context, db Queryer, sqlText string, args ...any) (int64, error) {
|
|
rows, err := db.QueryContext(ctx, sqlText, args...)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
defer rows.Close()
|
|
var total int64
|
|
if rows.Next() {
|
|
if err := rows.Scan(&total); err != nil {
|
|
return 0, err
|
|
}
|
|
}
|
|
return total, rows.Err()
|
|
}
|
|
|
|
func (r *RawFrameRepository) tableName(query ...RawFrameQuery) string {
|
|
if len(query) > 0 {
|
|
if child := r.rawChildTableName(query[0]); child != "" {
|
|
return child
|
|
}
|
|
}
|
|
if r.database == "" {
|
|
return "raw_frames"
|
|
}
|
|
return r.database + ".raw_frames"
|
|
}
|
|
|
|
func (r *RawFrameRepository) rawChildTableName(query RawFrameQuery) string {
|
|
protocol := strings.ToUpper(strings.TrimSpace(query.Protocol))
|
|
if protocol == "" {
|
|
return ""
|
|
}
|
|
vehicleKey := strings.TrimSpace(query.VehicleKey)
|
|
if vehicleKey == "" {
|
|
vehicleKey = strings.TrimSpace(query.VIN)
|
|
}
|
|
if vehicleKey == "" {
|
|
return ""
|
|
}
|
|
table := "raw_" + strings.ToLower(protocol) + "_" + hash16(vehicleKey)
|
|
if r.database != "" {
|
|
return r.database + "." + table
|
|
}
|
|
return table
|
|
}
|
|
|
|
func (r *RawFrameRepository) chunkTableName() string {
|
|
if r.database == "" {
|
|
return "raw_frame_payload_chunks"
|
|
}
|
|
return r.database + ".raw_frame_payload_chunks"
|
|
}
|
|
|
|
func (r *LocationRepository) tableName() string {
|
|
if r.database == "" {
|
|
return "vehicle_locations"
|
|
}
|
|
return r.database + ".vehicle_locations"
|
|
}
|
|
|
|
func normalizeRawFrameQuery(query RawFrameQuery) RawFrameQuery {
|
|
query.Protocol = strings.ToUpper(strings.TrimSpace(query.Protocol))
|
|
query.VehicleKey = strings.TrimSpace(query.VehicleKey)
|
|
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.ParsedFields = normalizeParsedFieldNames(query.ParsedFields)
|
|
if len(query.ParsedFields) > 0 {
|
|
query.IncludeFields = true
|
|
}
|
|
query.OrderBy = normalizeRawFrameOrderBy(query.OrderBy)
|
|
query.DateFrom = normalizeDateTimeLiteral(query.DateFrom)
|
|
query.DateTo = normalizeDateTimeLiteral(query.DateTo)
|
|
if query.Limit <= 0 {
|
|
query.Limit = 20
|
|
}
|
|
return query
|
|
}
|
|
|
|
func normalizeLocationQuery(query LocationQuery) LocationQuery {
|
|
query.Protocol = strings.ToUpper(strings.TrimSpace(query.Protocol))
|
|
query.VIN = strings.TrimSpace(query.VIN)
|
|
query.DateFrom = normalizeDateTimeLiteral(query.DateFrom)
|
|
query.DateTo = normalizeDateTimeLiteral(query.DateTo)
|
|
if query.Limit <= 0 {
|
|
query.Limit = 20
|
|
}
|
|
return query
|
|
}
|
|
|
|
func buildRawFrameSQL(table string, query RawFrameQuery) (string, []any) {
|
|
where := rawFrameWhere(query)
|
|
rawHexSelect := "'' AS raw_hex"
|
|
rawTextSelect := "'' AS raw_text"
|
|
parsedFieldsSelect := "'' AS parsed_fields"
|
|
if query.IncludePayload {
|
|
rawHexSelect = "raw_hex"
|
|
rawTextSelect = "raw_text"
|
|
}
|
|
if query.IncludeFields {
|
|
parsedFieldsSelect = "parsed_json AS parsed_fields"
|
|
}
|
|
sqlText := `SELECT ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, ` + rawHexSelect + `, ` + rawTextSelect + `, ` + parsedFieldsSelect + `, 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 " + rawFrameOrderColumn(query.OrderBy) + " DESC LIMIT " + strconv.Itoa(query.Limit) + " OFFSET " + strconv.Itoa(query.Offset)
|
|
return sqlText, nil
|
|
}
|
|
|
|
func buildRawFrameCountSQL(table string, query RawFrameQuery) (string, []any) {
|
|
sqlText := `SELECT COUNT(*) FROM ` + table
|
|
if where := rawFrameWhere(query); len(where) > 0 {
|
|
sqlText += " WHERE " + strings.Join(where, " AND ")
|
|
}
|
|
return sqlText, nil
|
|
}
|
|
|
|
type payloadChunkManifest struct {
|
|
Chunked bool `json:"chunked"`
|
|
PayloadKind string `json:"payload_kind"`
|
|
EventID string `json:"event_id"`
|
|
ChunkCount int `json:"chunk_count"`
|
|
}
|
|
|
|
func (r *RawFrameRepository) hydratePayloadChunks(ctx context.Context, rows []RawFrameRow) error {
|
|
type target struct {
|
|
rowIndex int
|
|
kind string
|
|
eventID string
|
|
count int
|
|
}
|
|
targets := make([]target, 0)
|
|
eventSeen := map[string]struct{}{}
|
|
for index := range rows {
|
|
for _, candidate := range []struct {
|
|
kind string
|
|
value string
|
|
}{
|
|
{kind: "raw_hex", value: rows[index].RawHex},
|
|
{kind: "raw_text", value: rows[index].RawText},
|
|
{kind: "parsed_fields", value: string(rows[index].ParsedFields)},
|
|
} {
|
|
manifest, ok := parsePayloadChunkManifest(candidate.value)
|
|
if !ok || !payloadKindMatches(candidate.kind, manifest.PayloadKind) || manifest.ChunkCount <= 0 {
|
|
continue
|
|
}
|
|
eventID := strings.TrimSpace(manifest.EventID)
|
|
if eventID == "" {
|
|
eventID = rows[index].EventID
|
|
}
|
|
targets = append(targets, target{
|
|
rowIndex: index,
|
|
kind: candidate.kind,
|
|
eventID: eventID,
|
|
count: manifest.ChunkCount,
|
|
})
|
|
eventSeen[eventID] = struct{}{}
|
|
}
|
|
}
|
|
if len(targets) == 0 {
|
|
return nil
|
|
}
|
|
|
|
eventIDs := make([]string, 0, len(eventSeen))
|
|
for eventID := range eventSeen {
|
|
eventIDs = append(eventIDs, eventID)
|
|
}
|
|
sqlText := `SELECT event_id, payload_kind, chunk_index, chunk_text FROM ` + r.chunkTableName() +
|
|
" WHERE event_id IN (" + quotedList(eventIDs) + ") ORDER BY event_id, payload_kind, chunk_index"
|
|
chunkRows, err := r.db.QueryContext(ctx, sqlText)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer chunkRows.Close()
|
|
|
|
chunks := map[string]map[string][]string{}
|
|
for chunkRows.Next() {
|
|
var eventID string
|
|
var kind string
|
|
var index int
|
|
var text string
|
|
if err := chunkRows.Scan(&eventID, &kind, &index, &text); err != nil {
|
|
return err
|
|
}
|
|
if chunks[eventID] == nil {
|
|
chunks[eventID] = map[string][]string{}
|
|
}
|
|
current := chunks[eventID][kind]
|
|
if index >= len(current) {
|
|
next := make([]string, index+1)
|
|
copy(next, current)
|
|
current = next
|
|
}
|
|
current[index] = text
|
|
chunks[eventID][kind] = current
|
|
}
|
|
if err := chunkRows.Err(); err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, target := range targets {
|
|
parts := chunks[target.eventID][target.kind]
|
|
if len(parts) != target.count {
|
|
continue
|
|
}
|
|
hydrated := strings.Join(parts, "")
|
|
switch target.kind {
|
|
case "raw_hex":
|
|
rows[target.rowIndex].RawHex = hydrated
|
|
case "raw_text":
|
|
rows[target.rowIndex].RawText = hydrated
|
|
case "parsed_fields":
|
|
rows[target.rowIndex].ParsedFields = rawJSONMessage(hydrated)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func rawJSONMessage(value string) json.RawMessage {
|
|
value = strings.TrimSpace(value)
|
|
if value == "" {
|
|
return nil
|
|
}
|
|
if json.Valid([]byte(value)) {
|
|
return json.RawMessage(value)
|
|
}
|
|
return json.RawMessage(jsonString(value))
|
|
}
|
|
|
|
func filterParsedFields(rows []RawFrameRow, fieldNames []string) {
|
|
fieldNames = normalizeParsedFieldNames(fieldNames)
|
|
if len(fieldNames) == 0 {
|
|
return
|
|
}
|
|
for index := range rows {
|
|
if len(rows[index].ParsedFields) == 0 {
|
|
continue
|
|
}
|
|
var fields map[string]any
|
|
if err := json.Unmarshal(rows[index].ParsedFields, &fields); err != nil {
|
|
continue
|
|
}
|
|
filtered := make(map[string]any, len(fieldNames))
|
|
for _, name := range fieldNames {
|
|
if value, ok := fields[name]; ok {
|
|
filtered[name] = value
|
|
}
|
|
}
|
|
if len(filtered) == 0 {
|
|
rows[index].ParsedFields = json.RawMessage(`{}`)
|
|
continue
|
|
}
|
|
rows[index].ParsedFields = rawJSONMessage(jsonString(filtered))
|
|
}
|
|
}
|
|
|
|
func normalizeParsedFieldNames(values []string) []string {
|
|
out := make([]string, 0, len(values))
|
|
seen := map[string]struct{}{}
|
|
for _, value := range values {
|
|
for _, part := range strings.Split(value, ",") {
|
|
name := strings.TrimSpace(part)
|
|
if name == "" {
|
|
continue
|
|
}
|
|
if _, ok := seen[name]; ok {
|
|
continue
|
|
}
|
|
seen[name] = struct{}{}
|
|
out = append(out, name)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func payloadKindMatches(candidate string, manifest string) bool {
|
|
if candidate == manifest {
|
|
return true
|
|
}
|
|
return candidate == "parsed_fields" && manifest == "parsed_json"
|
|
}
|
|
|
|
func parsePayloadChunkManifest(value string) (payloadChunkManifest, bool) {
|
|
var manifest payloadChunkManifest
|
|
if !strings.Contains(value, `"chunked"`) {
|
|
return manifest, false
|
|
}
|
|
if err := json.Unmarshal([]byte(value), &manifest); err != nil {
|
|
return manifest, false
|
|
}
|
|
return manifest, manifest.Chunked
|
|
}
|
|
|
|
func quotedList(values []string) string {
|
|
out := make([]string, 0, len(values))
|
|
for _, value := range values {
|
|
out = append(out, "'"+quote(value)+"'")
|
|
}
|
|
return strings.Join(out, ", ")
|
|
}
|
|
|
|
func buildLocationSQL(table string, query LocationQuery) (string, []any) {
|
|
where := locationWhere(query)
|
|
sqlText := `SELECT ts, event_id, received_at, longitude, latitude, altitude_m, speed_kmh, direction_deg, alarm_flag, status_flag, total_mileage_km, protocol, vin FROM ` + table
|
|
if len(where) > 0 {
|
|
sqlText += " WHERE " + strings.Join(where, " AND ")
|
|
}
|
|
sqlText += " ORDER BY ts DESC LIMIT " + strconv.Itoa(query.Limit) + " OFFSET " + strconv.Itoa(query.Offset)
|
|
return sqlText, nil
|
|
}
|
|
|
|
func buildLocationCountSQL(table string, query LocationQuery) (string, []any) {
|
|
sqlText := `SELECT COUNT(*) FROM ` + table
|
|
if where := locationWhere(query); len(where) > 0 {
|
|
sqlText += " WHERE " + strings.Join(where, " AND ")
|
|
}
|
|
return sqlText, nil
|
|
}
|
|
|
|
func rawFrameWhere(query RawFrameQuery) []string {
|
|
var where []string
|
|
add := func(clause string) {
|
|
where = append(where, clause)
|
|
}
|
|
if query.Protocol != "" {
|
|
add("protocol = '" + quote(query.Protocol) + "'")
|
|
}
|
|
if query.VehicleKey != "" {
|
|
add("vehicle_key = '" + quote(query.VehicleKey) + "'")
|
|
}
|
|
if query.VIN != "" {
|
|
add("vin = '" + quote(query.VIN) + "'")
|
|
}
|
|
if query.Phone != "" {
|
|
add("phone = '" + quote(query.Phone) + "'")
|
|
}
|
|
if query.DeviceID != "" {
|
|
add("device_id = '" + quote(query.DeviceID) + "'")
|
|
}
|
|
if query.MessageID != "" {
|
|
if parsed, ok := parseMessageID(query.MessageID); ok {
|
|
add("message_id = " + strconv.FormatInt(parsed, 10))
|
|
}
|
|
}
|
|
if query.DateFrom != "" {
|
|
add("ts >= '" + quote(normalizeDateTimeLiteral(query.DateFrom)) + "'")
|
|
}
|
|
if query.DateTo != "" {
|
|
add("ts <= '" + quote(normalizeDateTimeLiteral(query.DateTo)) + "'")
|
|
}
|
|
return where
|
|
}
|
|
|
|
func normalizeRawFrameOrderBy(value string) string {
|
|
value = strings.ToLower(strings.TrimSpace(value))
|
|
switch value {
|
|
case "receivedat", "received_at":
|
|
return "receivedAt"
|
|
default:
|
|
return "ts"
|
|
}
|
|
}
|
|
|
|
func rawFrameOrderColumn(value string) string {
|
|
if normalizeRawFrameOrderBy(value) == "receivedAt" {
|
|
return "received_at"
|
|
}
|
|
return "ts"
|
|
}
|
|
|
|
func locationWhere(query LocationQuery) []string {
|
|
var where []string
|
|
add := func(clause string) {
|
|
where = append(where, clause)
|
|
}
|
|
if query.Protocol != "" {
|
|
add("protocol = '" + quote(query.Protocol) + "'")
|
|
}
|
|
if query.VIN != "" {
|
|
add("vin = '" + quote(query.VIN) + "'")
|
|
}
|
|
if query.DateFrom != "" {
|
|
add("ts >= '" + quote(normalizeDateTimeLiteral(query.DateFrom)) + "'")
|
|
}
|
|
if query.DateTo != "" {
|
|
add("ts <= '" + quote(normalizeDateTimeLiteral(query.DateTo)) + "'")
|
|
}
|
|
return where
|
|
}
|
|
|
|
type RawFrameHandler struct {
|
|
repository *RawFrameRepository
|
|
}
|
|
|
|
type LocationHandler struct {
|
|
repository *LocationRepository
|
|
}
|
|
|
|
func NewRawFrameHandler(repository *RawFrameRepository) *RawFrameHandler {
|
|
if repository == nil {
|
|
panic("raw frame repository must not be nil")
|
|
}
|
|
return &RawFrameHandler{repository: repository}
|
|
}
|
|
|
|
func NewLocationHandler(repository *LocationRepository) *LocationHandler {
|
|
if repository == nil {
|
|
panic("location repository must not be nil")
|
|
}
|
|
return &LocationHandler{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
|
|
}
|
|
var total int64
|
|
if query.IncludeTotal {
|
|
total, err = h.repository.Count(r.Context(), query)
|
|
if err != nil {
|
|
writeHistoryError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
}
|
|
rows, err := h.repository.Query(r.Context(), query)
|
|
if err != nil {
|
|
writeHistoryError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
if !query.IncludeTotal {
|
|
total = int64(len(rows))
|
|
}
|
|
writeHistoryJSON(w, r, map[string]any{
|
|
"items": rows,
|
|
"total": total,
|
|
"limit": query.Limit,
|
|
"offset": query.Offset,
|
|
})
|
|
}
|
|
|
|
func (h *LocationHandler) 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/locations" {
|
|
writeHistoryError(w, http.StatusNotFound, "route not found")
|
|
return
|
|
}
|
|
query, err := parseLocationQuery(r)
|
|
if err != nil {
|
|
writeHistoryError(w, http.StatusBadRequest, err.Error())
|
|
return
|
|
}
|
|
var total int64
|
|
if query.IncludeTotal {
|
|
total, err = h.repository.Count(r.Context(), query)
|
|
if err != nil {
|
|
writeHistoryError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
}
|
|
rows, err := h.repository.Query(r.Context(), query)
|
|
if err != nil {
|
|
writeHistoryError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
if !query.IncludeTotal {
|
|
total = int64(len(rows))
|
|
}
|
|
writeHistoryJSON(w, r, map[string]any{
|
|
"items": rows,
|
|
"total": total,
|
|
"limit": query.Limit,
|
|
"offset": query.Offset,
|
|
})
|
|
}
|
|
|
|
func writeHistoryJSON(w http.ResponseWriter, r *http.Request, value any) {
|
|
w.Header().Set("Content-Type", "application/json")
|
|
w.Header().Add("Vary", "Accept-Encoding")
|
|
if strings.Contains(strings.ToLower(r.Header.Get("Accept-Encoding")), "gzip") {
|
|
w.Header().Set("Content-Encoding", "gzip")
|
|
gzipWriter := gzip.NewWriter(w)
|
|
defer gzipWriter.Close()
|
|
_ = json.NewEncoder(gzipWriter).Encode(value)
|
|
return
|
|
}
|
|
_ = json.NewEncoder(w).Encode(value)
|
|
}
|
|
|
|
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"),
|
|
VehicleKey: values.Get("vehicleKey"),
|
|
VIN: values.Get("vin"),
|
|
Phone: values.Get("phone"),
|
|
DeviceID: values.Get("deviceId"),
|
|
MessageID: values.Get("messageId"),
|
|
OrderBy: values.Get("orderBy"),
|
|
IncludeFields: strings.EqualFold(strings.TrimSpace(values.Get("includeFields")), "true"),
|
|
IncludePayload: strings.EqualFold(strings.TrimSpace(values.Get("includePayload")), "true"),
|
|
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
|
ParsedFields: append(values["fields"], values["parsedFields"]...),
|
|
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 parseLocationQuery(r *http.Request) (LocationQuery, error) {
|
|
values := r.URL.Query()
|
|
limit, err := parseBoundedInt(values.Get("limit"), 20, 1, 500, "limit")
|
|
if err != nil {
|
|
return LocationQuery{}, err
|
|
}
|
|
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
|
if err != nil {
|
|
return LocationQuery{}, err
|
|
}
|
|
query := LocationQuery{
|
|
Protocol: values.Get("protocol"),
|
|
VIN: values.Get("vin"),
|
|
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
|
DateFrom: values.Get("dateFrom"),
|
|
DateTo: values.Get("dateTo"),
|
|
Limit: limit,
|
|
Offset: offset,
|
|
}
|
|
if !validDateTime(query.DateFrom) || !validDateTime(query.DateTo) {
|
|
return LocationQuery{}, errors.New("dateFrom/dateTo must use YYYY-MM-DD or YYYY-MM-DD HH:mm:ss")
|
|
}
|
|
return normalizeLocationQuery(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 = normalizeDateTimeLiteral(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 normalizeDateTimeLiteral(value string) string {
|
|
value = strings.TrimSpace(value)
|
|
if value == "" {
|
|
return ""
|
|
}
|
|
shanghai := time.FixedZone("Asia/Shanghai", 8*3600)
|
|
for _, layout := range []string{"2006-01-02T15:04:05", "2006-01-02 15:04:05"} {
|
|
if parsed, err := time.ParseInLocation(layout, value, shanghai); err == nil {
|
|
return parsed.UTC().Format("2006-01-02 15:04:05")
|
|
}
|
|
}
|
|
if parsed, err := time.Parse(time.RFC3339, value); err == nil {
|
|
return parsed.UTC().Format("2006-01-02 15:04:05")
|
|
}
|
|
return value
|
|
}
|
|
|
|
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 nullableFloat(value sql.NullFloat64) *float64 {
|
|
if !value.Valid {
|
|
return nil
|
|
}
|
|
return &value.Float64
|
|
}
|
|
|
|
func nullableInt(value sql.NullInt64) *int64 {
|
|
if !value.Valid {
|
|
return nil
|
|
}
|
|
return &value.Int64
|
|
}
|
|
|
|
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})
|
|
}
|