1734 lines
60 KiB
Go
1734 lines
60 KiB
Go
package stats
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"net/http"
|
||
"net/url"
|
||
"regexp"
|
||
"strconv"
|
||
"strings"
|
||
"time"
|
||
)
|
||
|
||
type DataSourceQuery struct {
|
||
Protocol string
|
||
SourceIP string
|
||
SourceCode string
|
||
SourceKind string
|
||
SourceCodeMissing *bool
|
||
Enabled *bool
|
||
IncludeTotal bool
|
||
Limit int
|
||
Offset int
|
||
}
|
||
|
||
type DataSourceDiagnosticsQuery struct {
|
||
Protocol string
|
||
SourceIP string
|
||
SourceKind string
|
||
SourceCodeMissing *bool
|
||
MappingIssueOnly bool
|
||
Enabled *bool
|
||
IncludeTotal bool
|
||
Limit int
|
||
Offset int
|
||
}
|
||
|
||
type JT808IdentityGapQuery struct {
|
||
SourceIP string
|
||
SourceCode string
|
||
RecentSeconds int64
|
||
IncludeTotal bool
|
||
Limit int
|
||
Offset int
|
||
}
|
||
|
||
type JT808MappingGapQuery struct {
|
||
SourceIP string
|
||
SourceCode string
|
||
RecentSeconds int64
|
||
IncludeTotal bool
|
||
Limit int
|
||
Offset int
|
||
}
|
||
|
||
type DataSourceRow struct {
|
||
ID int64 `json:"id"`
|
||
Protocol string `json:"protocol"`
|
||
SourceIP string `json:"source_ip"`
|
||
LatestSourceEndpoint string `json:"latest_source_endpoint,omitempty"`
|
||
PlatformName string `json:"platform_name,omitempty"`
|
||
SourceCode string `json:"source_code,omitempty"`
|
||
SourceKind string `json:"source_kind"`
|
||
TrustPriority int `json:"trust_priority"`
|
||
Enabled bool `json:"enabled"`
|
||
FirstSeenAt string `json:"first_seen_at,omitempty"`
|
||
LatestSeenAt string `json:"latest_seen_at,omitempty"`
|
||
Remark string `json:"remark,omitempty"`
|
||
UpdatedAt string `json:"updated_at"`
|
||
}
|
||
|
||
type DataSourceUpdate struct {
|
||
PlatformName *string `json:"platform_name"`
|
||
SourceCode *string `json:"source_code"`
|
||
SourceKind *string `json:"source_kind"`
|
||
TrustPriority *int `json:"trust_priority"`
|
||
Enabled *bool `json:"enabled"`
|
||
Remark *string `json:"remark"`
|
||
}
|
||
|
||
type DataSourceDiagnosticRow struct {
|
||
ID int64 `json:"id"`
|
||
Protocol string `json:"protocol"`
|
||
SourceIP string `json:"source_ip"`
|
||
LatestSourceEndpoint string `json:"latest_source_endpoint,omitempty"`
|
||
PlatformName string `json:"platform_name,omitempty"`
|
||
SourceCode string `json:"source_code,omitempty"`
|
||
SourceKind string `json:"source_kind"`
|
||
FirstSeenAt string `json:"first_seen_at,omitempty"`
|
||
LatestSeenAt string `json:"latest_seen_at,omitempty"`
|
||
ActiveSpanSeconds int64 `json:"active_span_seconds"`
|
||
LatestSeenAgeSeconds int64 `json:"latest_seen_age_seconds"`
|
||
RegistrationRows int64 `json:"registration_rows"`
|
||
PhoneCount int64 `json:"phone_count"`
|
||
UnknownVINRows int64 `json:"unknown_vin_rows"`
|
||
IdentifierMatchedPhones int64 `json:"identifier_matched_phones"`
|
||
UnmappedPhoneCount int64 `json:"unmapped_phone_count"`
|
||
IdentifierMatchRatio float64 `json:"identifier_match_ratio"`
|
||
ConfiguredSourceCodeMatchedPhones int64 `json:"configured_source_code_matched_phones"`
|
||
ConfiguredSourceCodePlatformName string `json:"configured_source_code_platform_name,omitempty"`
|
||
ConfiguredSourceCodeConflict bool `json:"configured_source_code_conflict"`
|
||
SourcePlatformNameMismatch bool `json:"source_platform_name_mismatch"`
|
||
MatchedSourceCodeCount int64 `json:"matched_source_code_count"`
|
||
CandidateSourceCode string `json:"candidate_source_code,omitempty"`
|
||
CandidatePlatformName string `json:"candidate_platform_name,omitempty"`
|
||
MatchedSourceCodes []string `json:"matched_source_codes,omitempty"`
|
||
MatchedPlatformNames []string `json:"matched_platform_names,omitempty"`
|
||
SamplePhones []string `json:"sample_phones,omitempty"`
|
||
Reason string `json:"reason"`
|
||
RecommendedOperatorAction string `json:"recommended_operator_action"`
|
||
SuggestedSourceKind string `json:"suggested_source_kind"`
|
||
SuggestionConfidence string `json:"suggestion_confidence"`
|
||
SuggestionReason string `json:"suggestion_reason"`
|
||
}
|
||
|
||
type JT808IdentityGapRow struct {
|
||
Phone string `json:"phone"`
|
||
DeviceID string `json:"device_id,omitempty"`
|
||
Plate string `json:"plate,omitempty"`
|
||
VIN string `json:"vin,omitempty"`
|
||
SourceIP string `json:"source_ip,omitempty"`
|
||
SourceEndpoint string `json:"source_endpoint,omitempty"`
|
||
SourceCode string `json:"source_code,omitempty"`
|
||
PlatformName string `json:"platform_name,omitempty"`
|
||
SourceKind string `json:"source_kind,omitempty"`
|
||
FirstRegisteredAt string `json:"first_registered_at,omitempty"`
|
||
LatestRegisteredAt string `json:"latest_registered_at,omitempty"`
|
||
LatestAuthenticatedAt string `json:"latest_authenticated_at,omitempty"`
|
||
LatestSeenAt string `json:"latest_seen_at,omitempty"`
|
||
LatestSeenAgeSeconds int64 `json:"latest_seen_age_seconds"`
|
||
Reason string `json:"reason"`
|
||
RecommendedOperatorAction string `json:"recommended_operator_action"`
|
||
RawFrameQueryPath string `json:"raw_frame_query_path,omitempty"`
|
||
DataSourceQueryPath string `json:"data_source_query_path,omitempty"`
|
||
}
|
||
|
||
type JT808MappingGapRow struct {
|
||
Phone string `json:"phone"`
|
||
DeviceID string `json:"device_id,omitempty"`
|
||
Plate string `json:"plate,omitempty"`
|
||
VIN string `json:"vin,omitempty"`
|
||
SourceIP string `json:"source_ip,omitempty"`
|
||
SourceEndpoint string `json:"source_endpoint,omitempty"`
|
||
SourceCode string `json:"source_code,omitempty"`
|
||
PlatformName string `json:"platform_name,omitempty"`
|
||
SourceKind string `json:"source_kind,omitempty"`
|
||
IdentifierVIN string `json:"identifier_vin,omitempty"`
|
||
IdentifierPlate string `json:"identifier_plate,omitempty"`
|
||
MatchedSourceCodes []string `json:"matched_source_codes,omitempty"`
|
||
MatchedPlatformNames []string `json:"matched_platform_names,omitempty"`
|
||
LatestSeenAt string `json:"latest_seen_at,omitempty"`
|
||
LatestSeenAgeSeconds int64 `json:"latest_seen_age_seconds"`
|
||
Reason string `json:"reason"`
|
||
RecommendedAction string `json:"recommended_action"`
|
||
SuggestedSourceCode string `json:"suggested_source_code,omitempty"`
|
||
SuggestedPlatformName string `json:"suggested_platform_name,omitempty"`
|
||
SuggestedIdentifierType string `json:"suggested_identifier_type,omitempty"`
|
||
SuggestedIdentifierValue string `json:"suggested_identifier_value,omitempty"`
|
||
SuggestedVIN string `json:"suggested_vin,omitempty"`
|
||
SuggestedPlate string `json:"suggested_plate,omitempty"`
|
||
RawFrameQueryPath string `json:"raw_frame_query_path,omitempty"`
|
||
DataSourceQueryPath string `json:"data_source_query_path,omitempty"`
|
||
VehicleIdentifierExample string `json:"vehicle_identifier_example,omitempty"`
|
||
}
|
||
|
||
type dataSourceDB interface {
|
||
QueryContext(context.Context, string, ...any) (*sql.Rows, error)
|
||
ExecContext(context.Context, string, ...any) (sql.Result, error)
|
||
}
|
||
|
||
type DataSourceRepository struct {
|
||
db dataSourceDB
|
||
}
|
||
|
||
func NewDataSourceRepository(db dataSourceDB) *DataSourceRepository {
|
||
if db == nil {
|
||
panic("data source db must not be nil")
|
||
}
|
||
return &DataSourceRepository{db: db}
|
||
}
|
||
|
||
func (r *DataSourceRepository) Query(ctx context.Context, query DataSourceQuery) ([]DataSourceRow, error) {
|
||
query = normalizeDataSourceQuery(query)
|
||
sqlText, args := buildDataSourceSQL(query)
|
||
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
out := make([]DataSourceRow, 0)
|
||
for rows.Next() {
|
||
var row DataSourceRow
|
||
var endpoint, platform, sourceCode, sourceKind, remark sql.NullString
|
||
var firstSeen, latestSeen, updatedAt scanDateTime
|
||
var enabled int
|
||
if err := rows.Scan(
|
||
&row.ID,
|
||
&row.Protocol,
|
||
&row.SourceIP,
|
||
&endpoint,
|
||
&platform,
|
||
&sourceCode,
|
||
&sourceKind,
|
||
&row.TrustPriority,
|
||
&enabled,
|
||
&firstSeen,
|
||
&latestSeen,
|
||
&remark,
|
||
&updatedAt,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
row.LatestSourceEndpoint = endpoint.String
|
||
row.PlatformName = platform.String
|
||
row.SourceCode = sourceCode.String
|
||
row.SourceKind = normalizeSourceKindForRead(sourceKind.String)
|
||
row.Enabled = enabled != 0
|
||
row.FirstSeenAt = firstSeen.String
|
||
row.LatestSeenAt = latestSeen.String
|
||
row.Remark = remark.String
|
||
row.UpdatedAt = updatedAt.String
|
||
out = append(out, row)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
func (r *DataSourceRepository) Count(ctx context.Context, query DataSourceQuery) (int64, error) {
|
||
query = normalizeDataSourceQuery(query)
|
||
sqlText, args := buildDataSourceCountSQL(query)
|
||
rows, err := r.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 *DataSourceRepository) QueryDiagnostics(ctx context.Context, query DataSourceDiagnosticsQuery) ([]DataSourceDiagnosticRow, error) {
|
||
query = normalizeDataSourceDiagnosticsQuery(query)
|
||
sqlText, args := buildDataSourceDiagnosticsSQL(query)
|
||
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
out := make([]DataSourceDiagnosticRow, 0)
|
||
for rows.Next() {
|
||
var row DataSourceDiagnosticRow
|
||
var endpoint, platform, sourceCode, sourceKind sql.NullString
|
||
var firstSeen, latestSeen scanDateTime
|
||
var configuredSourceCodePlatformName, candidateSourceCode, candidatePlatformName, matchedSourceCodes, matchedPlatformNames, samplePhones sql.NullString
|
||
if err := rows.Scan(
|
||
&row.ID,
|
||
&row.Protocol,
|
||
&row.SourceIP,
|
||
&endpoint,
|
||
&platform,
|
||
&sourceCode,
|
||
&sourceKind,
|
||
&firstSeen,
|
||
&latestSeen,
|
||
&row.ActiveSpanSeconds,
|
||
&row.LatestSeenAgeSeconds,
|
||
&row.RegistrationRows,
|
||
&row.PhoneCount,
|
||
&row.UnknownVINRows,
|
||
&row.IdentifierMatchedPhones,
|
||
&row.UnmappedPhoneCount,
|
||
&row.IdentifierMatchRatio,
|
||
&row.ConfiguredSourceCodeMatchedPhones,
|
||
&configuredSourceCodePlatformName,
|
||
&row.MatchedSourceCodeCount,
|
||
&candidateSourceCode,
|
||
&candidatePlatformName,
|
||
&matchedSourceCodes,
|
||
&matchedPlatformNames,
|
||
&samplePhones,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
row.LatestSourceEndpoint = endpoint.String
|
||
row.PlatformName = platform.String
|
||
row.SourceCode = sourceCode.String
|
||
row.SourceKind = normalizeSourceKindForRead(sourceKind.String)
|
||
row.FirstSeenAt = firstSeen.String
|
||
row.LatestSeenAt = latestSeen.String
|
||
row.CandidateSourceCode = candidateSourceCode.String
|
||
row.CandidatePlatformName = candidatePlatformName.String
|
||
row.MatchedSourceCodes = splitCommaList(matchedSourceCodes.String)
|
||
row.MatchedPlatformNames = splitCommaList(matchedPlatformNames.String)
|
||
row.SamplePhones = splitCommaList(samplePhones.String)
|
||
row.ConfiguredSourceCodeConflict = configuredSourceCodeConflict(row)
|
||
row.ConfiguredSourceCodePlatformName = configuredSourceCodePlatformName.String
|
||
row.SourcePlatformNameMismatch = sourcePlatformNameMismatch(row)
|
||
row.Reason, row.RecommendedOperatorAction = diagnoseDataSource(row)
|
||
row.SuggestedSourceKind, row.SuggestionConfidence, row.SuggestionReason = suggestSourceKind(row)
|
||
out = append(out, row)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
func (r *DataSourceRepository) CountDiagnostics(ctx context.Context, query DataSourceDiagnosticsQuery) (int64, error) {
|
||
query = normalizeDataSourceDiagnosticsQuery(query)
|
||
sqlText, args := buildDataSourceDiagnosticsCountSQL(query)
|
||
rows, err := r.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 *DataSourceRepository) QueryJT808IdentityGaps(ctx context.Context, query JT808IdentityGapQuery) ([]JT808IdentityGapRow, error) {
|
||
query = normalizeJT808IdentityGapQuery(query)
|
||
sqlText, args := buildJT808IdentityGapSQL(query)
|
||
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
out := make([]JT808IdentityGapRow, 0)
|
||
for rows.Next() {
|
||
var row JT808IdentityGapRow
|
||
var deviceID, plate, vin, sourceIP, endpoint, sourceCode, platform, sourceKind sql.NullString
|
||
var firstRegistered, latestRegistered, latestAuthenticated, latestSeen scanDateTime
|
||
if err := rows.Scan(
|
||
&row.Phone,
|
||
&deviceID,
|
||
&plate,
|
||
&vin,
|
||
&sourceIP,
|
||
&endpoint,
|
||
&sourceCode,
|
||
&platform,
|
||
&sourceKind,
|
||
&firstRegistered,
|
||
&latestRegistered,
|
||
&latestAuthenticated,
|
||
&latestSeen,
|
||
&row.LatestSeenAgeSeconds,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
row.DeviceID = deviceID.String
|
||
row.Plate = plate.String
|
||
row.VIN = vin.String
|
||
row.SourceIP = sourceIP.String
|
||
row.SourceEndpoint = endpoint.String
|
||
row.SourceCode = sourceCode.String
|
||
row.PlatformName = platform.String
|
||
row.SourceKind = normalizeSourceKindForRead(sourceKind.String)
|
||
row.FirstRegisteredAt = firstRegistered.String
|
||
row.LatestRegisteredAt = latestRegistered.String
|
||
row.LatestAuthenticatedAt = latestAuthenticated.String
|
||
row.LatestSeenAt = latestSeen.String
|
||
row.Reason, row.RecommendedOperatorAction = diagnoseJT808IdentityGap(row)
|
||
row.RawFrameQueryPath = jt808IdentityGapRawFrameQueryPath(row)
|
||
row.DataSourceQueryPath = jt808IdentityGapDataSourceQueryPath(row)
|
||
out = append(out, row)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
func (r *DataSourceRepository) CountJT808IdentityGaps(ctx context.Context, query JT808IdentityGapQuery) (int64, error) {
|
||
query = normalizeJT808IdentityGapQuery(query)
|
||
sqlText, args := buildJT808IdentityGapCountSQL(query)
|
||
rows, err := r.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 *DataSourceRepository) QueryJT808MappingGaps(ctx context.Context, query JT808MappingGapQuery) ([]JT808MappingGapRow, error) {
|
||
query = normalizeJT808MappingGapQuery(query)
|
||
sqlText, args := buildJT808MappingGapSQL(query)
|
||
rows, err := r.db.QueryContext(ctx, sqlText, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
out := make([]JT808MappingGapRow, 0)
|
||
for rows.Next() {
|
||
var row JT808MappingGapRow
|
||
var deviceID, plate, vin, sourceIP, endpoint, sourceCode, platform, sourceKind sql.NullString
|
||
var identifierVIN, identifierPlate, matchedSourceCodes, matchedPlatformNames sql.NullString
|
||
var latestSeen scanDateTime
|
||
if err := rows.Scan(
|
||
&row.Phone,
|
||
&deviceID,
|
||
&plate,
|
||
&vin,
|
||
&sourceIP,
|
||
&endpoint,
|
||
&sourceCode,
|
||
&platform,
|
||
&sourceKind,
|
||
&identifierVIN,
|
||
&identifierPlate,
|
||
&matchedSourceCodes,
|
||
&matchedPlatformNames,
|
||
&latestSeen,
|
||
&row.LatestSeenAgeSeconds,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
row.DeviceID = deviceID.String
|
||
row.Plate = plate.String
|
||
row.VIN = vin.String
|
||
row.SourceIP = sourceIP.String
|
||
row.SourceEndpoint = endpoint.String
|
||
row.SourceCode = sourceCode.String
|
||
row.PlatformName = platform.String
|
||
row.SourceKind = normalizeSourceKindForRead(sourceKind.String)
|
||
row.IdentifierVIN = identifierVIN.String
|
||
row.IdentifierPlate = identifierPlate.String
|
||
row.MatchedSourceCodes = splitCommaList(matchedSourceCodes.String)
|
||
row.MatchedPlatformNames = splitCommaList(matchedPlatformNames.String)
|
||
row.LatestSeenAt = latestSeen.String
|
||
annotateJT808MappingGap(&row)
|
||
out = append(out, row)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
func (r *DataSourceRepository) CountJT808MappingGaps(ctx context.Context, query JT808MappingGapQuery) (int64, error) {
|
||
query = normalizeJT808MappingGapQuery(query)
|
||
sqlText, args := buildJT808MappingGapCountSQL(query)
|
||
rows, err := r.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 *DataSourceRepository) Update(ctx context.Context, id int64, update DataSourceUpdate) (bool, error) {
|
||
if id <= 0 {
|
||
return false, errors.New("id must be positive")
|
||
}
|
||
if err := validateDataSourceUpdate(update); err != nil {
|
||
return false, err
|
||
}
|
||
var sets []string
|
||
var args []any
|
||
if update.PlatformName != nil {
|
||
sets = append(sets, "platform_name = ?")
|
||
args = append(args, nullableTrimmedString(*update.PlatformName))
|
||
}
|
||
if update.SourceCode != nil {
|
||
sets = append(sets, "source_code = ?")
|
||
args = append(args, nullableTrimmedString(*update.SourceCode))
|
||
}
|
||
if update.SourceKind != nil {
|
||
sets = append(sets, "source_kind = ?")
|
||
args = append(args, normalizeSourceKindForWrite(*update.SourceKind))
|
||
}
|
||
if update.TrustPriority != nil {
|
||
sets = append(sets, "trust_priority = ?")
|
||
args = append(args, *update.TrustPriority)
|
||
}
|
||
if update.Enabled != nil {
|
||
sets = append(sets, "enabled = ?")
|
||
if *update.Enabled {
|
||
args = append(args, 1)
|
||
} else {
|
||
args = append(args, 0)
|
||
}
|
||
}
|
||
if update.Remark != nil {
|
||
sets = append(sets, "remark = ?")
|
||
args = append(args, nullableTrimmedString(*update.Remark))
|
||
}
|
||
if len(sets) == 0 {
|
||
return false, errors.New("no mutable fields provided")
|
||
}
|
||
exists, err := r.exists(ctx, id)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
if !exists {
|
||
return false, nil
|
||
}
|
||
args = append(args, id)
|
||
_, err = r.db.ExecContext(ctx, "UPDATE vehicle_data_source SET "+strings.Join(sets, ", ")+", updated_at = CURRENT_TIMESTAMP WHERE id = ?", args...)
|
||
return err == nil, err
|
||
}
|
||
|
||
func (r *DataSourceRepository) exists(ctx context.Context, id int64) (bool, error) {
|
||
rows, err := r.db.QueryContext(ctx, "SELECT 1 FROM vehicle_data_source WHERE id = ? LIMIT 1", id)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
defer rows.Close()
|
||
return rows.Next(), rows.Err()
|
||
}
|
||
|
||
func normalizeDataSourceQuery(query DataSourceQuery) DataSourceQuery {
|
||
query.Protocol = strings.ToUpper(strings.TrimSpace(query.Protocol))
|
||
query.SourceIP = strings.TrimSpace(query.SourceIP)
|
||
query.SourceCode = strings.TrimSpace(query.SourceCode)
|
||
query.SourceKind = normalizeSourceKindForQuery(query.SourceKind)
|
||
if query.Limit <= 0 {
|
||
query.Limit = 50
|
||
}
|
||
return query
|
||
}
|
||
|
||
func normalizeDataSourceDiagnosticsQuery(query DataSourceDiagnosticsQuery) DataSourceDiagnosticsQuery {
|
||
query.Protocol = strings.ToUpper(strings.TrimSpace(query.Protocol))
|
||
if query.Protocol == "" {
|
||
query.Protocol = "JT808"
|
||
}
|
||
query.SourceIP = strings.TrimSpace(query.SourceIP)
|
||
query.SourceKind = normalizeSourceKindForQuery(query.SourceKind)
|
||
if query.Enabled == nil {
|
||
enabled := true
|
||
query.Enabled = &enabled
|
||
}
|
||
if query.Limit <= 0 {
|
||
query.Limit = 50
|
||
}
|
||
return query
|
||
}
|
||
|
||
func normalizeJT808IdentityGapQuery(query JT808IdentityGapQuery) JT808IdentityGapQuery {
|
||
query.SourceIP = strings.TrimSpace(query.SourceIP)
|
||
query.SourceCode = strings.TrimSpace(query.SourceCode)
|
||
if query.RecentSeconds < 0 {
|
||
query.RecentSeconds = 0
|
||
}
|
||
if query.Limit <= 0 {
|
||
query.Limit = 50
|
||
}
|
||
return query
|
||
}
|
||
|
||
func normalizeJT808MappingGapQuery(query JT808MappingGapQuery) JT808MappingGapQuery {
|
||
query.SourceIP = strings.TrimSpace(query.SourceIP)
|
||
query.SourceCode = strings.TrimSpace(query.SourceCode)
|
||
if query.RecentSeconds < 0 {
|
||
query.RecentSeconds = 0
|
||
}
|
||
if query.Limit <= 0 {
|
||
query.Limit = 50
|
||
}
|
||
return query
|
||
}
|
||
|
||
func buildDataSourceSQL(query DataSourceQuery) (string, []any) {
|
||
where, args := buildDataSourceWhere(query)
|
||
sqlText := `SELECT id, protocol, source_ip, latest_source_endpoint, platform_name, source_code, source_kind,
|
||
trust_priority, enabled, first_seen_at, latest_seen_at, remark, updated_at
|
||
FROM vehicle_data_source`
|
||
if len(where) > 0 {
|
||
sqlText += " WHERE " + strings.Join(where, " AND ")
|
||
}
|
||
sqlText += " ORDER BY latest_seen_at DESC, id ASC LIMIT ? OFFSET ?"
|
||
args = append(args, query.Limit, query.Offset)
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildDataSourceCountSQL(query DataSourceQuery) (string, []any) {
|
||
where, args := buildDataSourceWhere(query)
|
||
sqlText := `SELECT COUNT(*) FROM vehicle_data_source`
|
||
if len(where) > 0 {
|
||
sqlText += " WHERE " + strings.Join(where, " AND ")
|
||
}
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildDataSourceDiagnosticsSQL(query DataSourceDiagnosticsQuery) (string, []any) {
|
||
where, args := buildDataSourceDiagnosticsWhere(query)
|
||
sqlText := `SELECT
|
||
ds.id,
|
||
ds.protocol,
|
||
ds.source_ip,
|
||
ds.latest_source_endpoint,
|
||
ds.platform_name,
|
||
ds.source_code,
|
||
ds.source_kind,
|
||
ds.first_seen_at,
|
||
ds.latest_seen_at,
|
||
COALESCE(TIMESTAMPDIFF(SECOND, ds.first_seen_at, ds.latest_seen_at), 0) AS active_span_seconds,
|
||
COALESCE(TIMESTAMPDIFF(SECOND, ds.latest_seen_at, CURRENT_TIMESTAMP), 0) AS latest_seen_age_seconds,
|
||
COUNT(r.phone) AS registration_rows,
|
||
COUNT(DISTINCT r.phone) AS phone_count,
|
||
SUM(CASE WHEN r.phone IS NOT NULL AND (r.vin IS NULL OR TRIM(r.vin) = '' OR LOWER(TRIM(r.vin)) = 'unknown') THEN 1 ELSE 0 END) AS unknown_vin_rows,
|
||
COUNT(DISTINCT CASE WHEN vi.identifier_value IS NOT NULL THEN r.phone END) AS identifier_matched_phones,
|
||
COUNT(DISTINCT CASE WHEN r.phone IS NOT NULL AND vi.identifier_value IS NULL THEN r.phone END) AS unmapped_phone_count,
|
||
COALESCE(COUNT(DISTINCT CASE WHEN vi.identifier_value IS NOT NULL THEN r.phone END) / NULLIF(COUNT(DISTINCT r.phone), 0), 0) AS identifier_match_ratio,
|
||
COUNT(DISTINCT CASE
|
||
WHEN ds.source_code IS NOT NULL AND TRIM(ds.source_code) <> ''
|
||
AND vi.source_code IS NOT NULL AND TRIM(vi.source_code) <> ''
|
||
AND TRIM(vi.source_code) = TRIM(ds.source_code)
|
||
THEN r.phone
|
||
END) AS configured_source_code_matched_phones,
|
||
MAX(configured_catalog.platform_name) AS configured_source_code_platform_name,
|
||
COUNT(DISTINCT NULLIF(TRIM(vi.source_code), '')) AS matched_source_code_count,
|
||
MIN(NULLIF(TRIM(vi.source_code), '')) AS candidate_source_code,
|
||
MIN(COALESCE(NULLIF(TRIM(vi.oem), ''), NULLIF(TRIM(vi.source_code), ''))) AS candidate_platform_name,
|
||
GROUP_CONCAT(DISTINCT NULLIF(TRIM(vi.source_code), '') ORDER BY NULLIF(TRIM(vi.source_code), '') SEPARATOR ',') AS matched_source_codes,
|
||
GROUP_CONCAT(DISTINCT COALESCE(NULLIF(TRIM(vi.oem), ''), NULLIF(TRIM(vi.source_code), '')) ORDER BY COALESCE(NULLIF(TRIM(vi.oem), ''), NULLIF(TRIM(vi.source_code), '')) SEPARATOR ',') AS matched_platform_names,
|
||
SUBSTRING_INDEX(GROUP_CONCAT(DISTINCT r.phone ORDER BY r.latest_seen_at DESC SEPARATOR ','), ',', 5) AS sample_phones
|
||
FROM vehicle_data_source ds
|
||
LEFT JOIN jt808_registration r
|
||
ON ds.protocol = 'JT808'
|
||
AND r.source_ip = ds.source_ip
|
||
LEFT JOIN vehicle_identifier vi
|
||
ON vi.protocol = 'JT808'
|
||
AND vi.identifier_type = 'JT808_PHONE'
|
||
AND vi.identifier_value = r.phone
|
||
AND vi.enabled = 1
|
||
LEFT JOIN (
|
||
SELECT
|
||
source_code,
|
||
MIN(COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), ''))) AS platform_name
|
||
FROM vehicle_identifier
|
||
WHERE protocol = 'JT808'
|
||
AND enabled = 1
|
||
AND source_code IS NOT NULL
|
||
AND TRIM(source_code) <> ''
|
||
GROUP BY source_code
|
||
HAVING COUNT(DISTINCT COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), ''))) = 1
|
||
) configured_catalog
|
||
ON configured_catalog.source_code = ds.source_code`
|
||
if len(where) > 0 {
|
||
sqlText += " WHERE " + strings.Join(where, " AND ")
|
||
}
|
||
sqlText += ` GROUP BY ds.id, ds.protocol, ds.source_ip, ds.latest_source_endpoint, ds.platform_name, ds.source_code, ds.source_kind, ds.first_seen_at, ds.latest_seen_at
|
||
`
|
||
if having := buildDataSourceDiagnosticsHaving(query); len(having) > 0 {
|
||
sqlText += "HAVING " + strings.Join(having, " AND ") + "\n"
|
||
}
|
||
sqlText += `
|
||
ORDER BY ds.latest_seen_at DESC, ds.id ASC LIMIT ? OFFSET ?`
|
||
args = append(args, query.Limit, query.Offset)
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildDataSourceDiagnosticsCountSQL(query DataSourceDiagnosticsQuery) (string, []any) {
|
||
where, args := buildDataSourceDiagnosticsWhere(query)
|
||
having := buildDataSourceDiagnosticsHaving(query)
|
||
if len(having) > 0 {
|
||
sqlText := `SELECT COUNT(*)
|
||
FROM (
|
||
SELECT ds.id,
|
||
MAX(ds.platform_name) AS platform_name,
|
||
ds.source_code AS source_code,
|
||
COUNT(DISTINCT r.phone) AS phone_count,
|
||
COUNT(DISTINCT CASE WHEN vi.identifier_value IS NOT NULL THEN r.phone END) AS identifier_matched_phones,
|
||
COUNT(DISTINCT CASE WHEN r.phone IS NOT NULL AND vi.identifier_value IS NULL THEN r.phone END) AS unmapped_phone_count,
|
||
COALESCE(COUNT(DISTINCT CASE WHEN vi.identifier_value IS NOT NULL THEN r.phone END) / NULLIF(COUNT(DISTINCT r.phone), 0), 0) AS identifier_match_ratio,
|
||
COUNT(DISTINCT NULLIF(TRIM(vi.source_code), '')) AS matched_source_code_count,
|
||
COUNT(DISTINCT CASE
|
||
WHEN ds.source_code IS NOT NULL AND TRIM(ds.source_code) <> ''
|
||
AND vi.source_code IS NOT NULL AND TRIM(vi.source_code) <> ''
|
||
AND TRIM(vi.source_code) = TRIM(ds.source_code)
|
||
THEN r.phone
|
||
END) AS configured_source_code_matched_phones,
|
||
MAX(configured_catalog.platform_name) AS configured_source_code_platform_name
|
||
FROM vehicle_data_source ds
|
||
LEFT JOIN jt808_registration r
|
||
ON ds.protocol = 'JT808'
|
||
AND r.source_ip = ds.source_ip
|
||
LEFT JOIN vehicle_identifier vi
|
||
ON vi.protocol = 'JT808'
|
||
AND vi.identifier_type = 'JT808_PHONE'
|
||
AND vi.identifier_value = r.phone
|
||
AND vi.enabled = 1
|
||
LEFT JOIN (
|
||
SELECT
|
||
source_code,
|
||
MIN(COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), ''))) AS platform_name
|
||
FROM vehicle_identifier
|
||
WHERE protocol = 'JT808'
|
||
AND enabled = 1
|
||
AND source_code IS NOT NULL
|
||
AND TRIM(source_code) <> ''
|
||
GROUP BY source_code
|
||
HAVING COUNT(DISTINCT COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), ''))) = 1
|
||
) configured_catalog
|
||
ON configured_catalog.source_code = ds.source_code`
|
||
if len(where) > 0 {
|
||
sqlText += " WHERE " + strings.Join(where, " AND ")
|
||
}
|
||
sqlText += ` GROUP BY ds.id, ds.source_code
|
||
HAVING ` + strings.Join(having, " AND ") + `
|
||
) diagnostics`
|
||
return sqlText, args
|
||
}
|
||
sqlText := `SELECT COUNT(*) FROM vehicle_data_source ds`
|
||
if len(where) > 0 {
|
||
sqlText += " WHERE " + strings.Join(where, " AND ")
|
||
}
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildJT808IdentityGapSQL(query JT808IdentityGapQuery) (string, []any) {
|
||
fromSQL, args := buildJT808IdentityGapFromSQL(query)
|
||
sqlText := `SELECT
|
||
r.phone,
|
||
r.device_id,
|
||
r.plate,
|
||
r.vin,
|
||
r.source_ip,
|
||
r.source_endpoint,
|
||
ds.source_code,
|
||
ds.platform_name,
|
||
ds.source_kind,
|
||
r.first_registered_at,
|
||
r.latest_registered_at,
|
||
r.latest_authenticated_at,
|
||
r.latest_seen_at,
|
||
COALESCE(TIMESTAMPDIFF(SECOND, r.latest_seen_at, CURRENT_TIMESTAMP), 0) AS latest_seen_age_seconds
|
||
` + fromSQL + `
|
||
ORDER BY r.latest_seen_at DESC, r.phone ASC
|
||
LIMIT ? OFFSET ?`
|
||
args = append(args, query.Limit, query.Offset)
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildJT808IdentityGapCountSQL(query JT808IdentityGapQuery) (string, []any) {
|
||
fromSQL, args := buildJT808IdentityGapFromSQL(query)
|
||
return "SELECT COUNT(*) " + fromSQL, args
|
||
}
|
||
|
||
func buildJT808MappingGapSQL(query JT808MappingGapQuery) (string, []any) {
|
||
fromSQL, args := buildJT808MappingGapFromSQL(query)
|
||
sqlText := `SELECT
|
||
r.phone,
|
||
r.device_id,
|
||
r.plate,
|
||
r.vin,
|
||
r.source_ip,
|
||
r.source_endpoint,
|
||
ds.source_code,
|
||
ds.platform_name,
|
||
ds.source_kind,
|
||
source_vi.vin AS identifier_vin,
|
||
source_vi.plate AS identifier_plate,
|
||
all_vi.matched_source_codes,
|
||
all_vi.matched_platform_names,
|
||
r.latest_seen_at,
|
||
COALESCE(TIMESTAMPDIFF(SECOND, r.latest_seen_at, CURRENT_TIMESTAMP), 0) AS latest_seen_age_seconds
|
||
` + fromSQL + `
|
||
ORDER BY r.latest_seen_at DESC, r.phone ASC
|
||
LIMIT ? OFFSET ?`
|
||
args = append(args, query.Limit, query.Offset)
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildJT808MappingGapCountSQL(query JT808MappingGapQuery) (string, []any) {
|
||
fromSQL, args := buildJT808MappingGapFromSQL(query)
|
||
return "SELECT COUNT(*) " + fromSQL, args
|
||
}
|
||
|
||
func buildJT808IdentityGapFromSQL(query JT808IdentityGapQuery) (string, []any) {
|
||
where, args := buildJT808IdentityGapWhere(query)
|
||
sqlText := `FROM jt808_registration r
|
||
LEFT JOIN vehicle_data_source ds
|
||
ON ds.protocol = 'JT808' AND ds.source_ip = r.source_ip
|
||
LEFT JOIN (
|
||
SELECT identifier_value, COUNT(*) AS match_count
|
||
FROM vehicle_identifier
|
||
WHERE protocol = 'JT808'
|
||
AND identifier_type = 'JT808_PHONE'
|
||
AND enabled = 1
|
||
AND vin IS NOT NULL AND TRIM(vin) <> ''
|
||
GROUP BY identifier_value
|
||
) phone_vi ON phone_vi.identifier_value = r.phone
|
||
LEFT JOIN (
|
||
SELECT identifier_value, COUNT(*) AS match_count
|
||
FROM vehicle_identifier
|
||
WHERE protocol = 'JT808'
|
||
AND identifier_type = 'PLATE'
|
||
AND enabled = 1
|
||
AND vin IS NOT NULL AND TRIM(vin) <> ''
|
||
GROUP BY identifier_value
|
||
) plate_vi ON plate_vi.identifier_value = r.plate
|
||
LEFT JOIN vehicle_identity_binding phone_binding
|
||
ON phone_binding.phone = r.phone
|
||
AND phone_binding.vin IS NOT NULL AND TRIM(phone_binding.vin) <> ''
|
||
LEFT JOIN vehicle_identity_binding plate_binding
|
||
ON plate_binding.plate = r.plate
|
||
AND plate_binding.vin IS NOT NULL AND TRIM(plate_binding.vin) <> ''`
|
||
if len(where) > 0 {
|
||
sqlText += "\nWHERE " + strings.Join(where, " AND ")
|
||
}
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildJT808MappingGapFromSQL(query JT808MappingGapQuery) (string, []any) {
|
||
where, args := buildJT808MappingGapWhere(query)
|
||
sqlText := `FROM jt808_registration r
|
||
JOIN vehicle_data_source ds
|
||
ON ds.protocol = 'JT808'
|
||
AND ds.source_ip = r.source_ip
|
||
AND ds.enabled = 1
|
||
AND ds.source_code IS NOT NULL
|
||
AND TRIM(ds.source_code) <> ''
|
||
LEFT JOIN vehicle_identifier source_vi
|
||
ON source_vi.protocol = 'JT808'
|
||
AND source_vi.source_code = ds.source_code
|
||
AND source_vi.identifier_type = 'JT808_PHONE'
|
||
AND source_vi.identifier_value = r.phone
|
||
AND source_vi.enabled = 1
|
||
LEFT JOIN (
|
||
SELECT
|
||
identifier_value,
|
||
GROUP_CONCAT(DISTINCT NULLIF(TRIM(source_code), '') ORDER BY NULLIF(TRIM(source_code), '') SEPARATOR ',') AS matched_source_codes,
|
||
GROUP_CONCAT(DISTINCT COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), '')) ORDER BY COALESCE(NULLIF(TRIM(oem), ''), NULLIF(TRIM(source_code), '')) SEPARATOR ',') AS matched_platform_names
|
||
FROM vehicle_identifier
|
||
WHERE protocol = 'JT808'
|
||
AND identifier_type = 'JT808_PHONE'
|
||
AND enabled = 1
|
||
GROUP BY identifier_value
|
||
) all_vi
|
||
ON all_vi.identifier_value = r.phone`
|
||
if len(where) > 0 {
|
||
sqlText += "\nWHERE " + strings.Join(where, " AND ")
|
||
}
|
||
return sqlText, args
|
||
}
|
||
|
||
func buildJT808IdentityGapWhere(query JT808IdentityGapQuery) ([]string, []any) {
|
||
where := []string{
|
||
"r.phone IS NOT NULL AND TRIM(r.phone) <> ''",
|
||
"(r.vin IS NULL OR TRIM(r.vin) = '' OR LOWER(TRIM(r.vin)) = 'unknown')",
|
||
"COALESCE(phone_vi.match_count, 0) = 0",
|
||
"COALESCE(plate_vi.match_count, 0) = 0",
|
||
"phone_binding.vin IS NULL",
|
||
"plate_binding.vin IS NULL",
|
||
}
|
||
var args []any
|
||
if query.SourceIP != "" {
|
||
where = append(where, "r.source_ip = ?")
|
||
args = append(args, query.SourceIP)
|
||
}
|
||
if query.SourceCode != "" {
|
||
where = append(where, "ds.source_code = ?")
|
||
args = append(args, query.SourceCode)
|
||
}
|
||
if query.RecentSeconds > 0 {
|
||
where = append(where, "r.latest_seen_at >= DATE_SUB(CURRENT_TIMESTAMP, INTERVAL ? SECOND)")
|
||
args = append(args, query.RecentSeconds)
|
||
}
|
||
return where, args
|
||
}
|
||
|
||
func buildJT808MappingGapWhere(query JT808MappingGapQuery) ([]string, []any) {
|
||
where := []string{
|
||
"r.phone IS NOT NULL AND TRIM(r.phone) <> ''",
|
||
"(source_vi.identifier_value IS NULL OR source_vi.vin IS NULL OR TRIM(source_vi.vin) = '')",
|
||
}
|
||
var args []any
|
||
if query.SourceIP != "" {
|
||
where = append(where, "r.source_ip = ?")
|
||
args = append(args, query.SourceIP)
|
||
}
|
||
if query.SourceCode != "" {
|
||
where = append(where, "ds.source_code = ?")
|
||
args = append(args, query.SourceCode)
|
||
}
|
||
if query.RecentSeconds > 0 {
|
||
where = append(where, "r.latest_seen_at >= DATE_SUB(CURRENT_TIMESTAMP, INTERVAL ? SECOND)")
|
||
args = append(args, query.RecentSeconds)
|
||
}
|
||
return where, args
|
||
}
|
||
|
||
func buildDataSourceWhere(query DataSourceQuery) ([]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.SourceIP != "" {
|
||
add("source_ip = ?", query.SourceIP)
|
||
}
|
||
if query.SourceCode != "" {
|
||
add("source_code = ?", query.SourceCode)
|
||
}
|
||
if query.SourceKind != "" {
|
||
add("source_kind = ?", query.SourceKind)
|
||
}
|
||
if query.SourceCodeMissing != nil {
|
||
if *query.SourceCodeMissing {
|
||
where = append(where, "(source_code IS NULL OR TRIM(source_code) = '')")
|
||
} else {
|
||
where = append(where, "(source_code IS NOT NULL AND TRIM(source_code) <> '')")
|
||
}
|
||
}
|
||
if query.Enabled != nil {
|
||
if *query.Enabled {
|
||
add("enabled = ?", 1)
|
||
} else {
|
||
add("enabled = ?", 0)
|
||
}
|
||
}
|
||
return where, args
|
||
}
|
||
|
||
func buildDataSourceDiagnosticsWhere(query DataSourceDiagnosticsQuery) ([]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("ds.protocol = ?", query.Protocol)
|
||
}
|
||
if query.SourceIP != "" {
|
||
add("ds.source_ip = ?", query.SourceIP)
|
||
}
|
||
if query.SourceKind != "" {
|
||
add("ds.source_kind = ?", query.SourceKind)
|
||
}
|
||
if query.SourceCodeMissing != nil {
|
||
if *query.SourceCodeMissing {
|
||
where = append(where, "(ds.source_code IS NULL OR TRIM(ds.source_code) = '')")
|
||
} else {
|
||
where = append(where, "(ds.source_code IS NOT NULL AND TRIM(ds.source_code) <> '')")
|
||
}
|
||
}
|
||
if query.Enabled != nil {
|
||
if *query.Enabled {
|
||
add("ds.enabled = ?", 1)
|
||
} else {
|
||
add("ds.enabled = ?", 0)
|
||
}
|
||
}
|
||
return where, args
|
||
}
|
||
|
||
func buildDataSourceDiagnosticsHaving(query DataSourceDiagnosticsQuery) []string {
|
||
if !query.MappingIssueOnly {
|
||
return nil
|
||
}
|
||
return []string{`(
|
||
(
|
||
platform_name IS NOT NULL AND TRIM(platform_name) <> ''
|
||
AND configured_source_code_platform_name IS NOT NULL AND TRIM(configured_source_code_platform_name) <> ''
|
||
AND TRIM(platform_name) <> TRIM(configured_source_code_platform_name)
|
||
)
|
||
OR (
|
||
source_code IS NOT NULL AND TRIM(source_code) <> ''
|
||
AND matched_source_code_count > 0
|
||
AND configured_source_code_matched_phones = 0
|
||
)
|
||
OR (
|
||
source_code IS NOT NULL AND TRIM(source_code) <> ''
|
||
AND phone_count >= 10
|
||
AND unmapped_phone_count > 0
|
||
AND identifier_match_ratio < 0.8
|
||
)
|
||
)`}
|
||
}
|
||
|
||
type DataSourceHandler struct {
|
||
repository *DataSourceRepository
|
||
}
|
||
|
||
func NewDataSourceHandler(repository *DataSourceRepository) *DataSourceHandler {
|
||
if repository == nil {
|
||
panic("data source repository must not be nil")
|
||
}
|
||
return &DataSourceHandler{repository: repository}
|
||
}
|
||
|
||
func (h *DataSourceHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||
path := strings.Trim(r.URL.Path, "/")
|
||
switch {
|
||
case r.Method == http.MethodGet && path == "api/stats/data-sources/diagnostics":
|
||
h.handleDiagnostics(w, r)
|
||
case r.Method == http.MethodGet && path == "api/stats/data-sources/kind-suggestions":
|
||
h.handleKindSuggestions(w, r)
|
||
case r.Method == http.MethodGet && path == "api/stats/data-sources/jt808-identity-gaps":
|
||
h.handleJT808IdentityGaps(w, r)
|
||
case r.Method == http.MethodGet && path == "api/stats/data-sources/jt808-mapping-gaps":
|
||
h.handleJT808MappingGaps(w, r)
|
||
case r.Method == http.MethodGet && path == "api/stats/data-sources":
|
||
h.handleQuery(w, r)
|
||
case r.Method == http.MethodPatch && strings.HasPrefix(path, "api/stats/data-sources/"):
|
||
h.handlePatch(w, r)
|
||
default:
|
||
writeMetricError(w, http.StatusNotFound, "route not found")
|
||
}
|
||
}
|
||
|
||
func (h *DataSourceHandler) handleKindSuggestions(w http.ResponseWriter, r *http.Request) {
|
||
query, err := parseDataSourceKindSuggestionsQuery(r)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
var total int64
|
||
if query.IncludeTotal {
|
||
total, err = h.repository.CountDiagnostics(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
}
|
||
rows, err := h.repository.QueryDiagnostics(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
if !query.IncludeTotal {
|
||
total = int64(len(rows))
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||
"items": rows,
|
||
"total": total,
|
||
"limit": query.Limit,
|
||
"offset": query.Offset,
|
||
})
|
||
}
|
||
|
||
func (h *DataSourceHandler) handleDiagnostics(w http.ResponseWriter, r *http.Request) {
|
||
query, err := parseDataSourceDiagnosticsQuery(r)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
var total int64
|
||
if query.IncludeTotal {
|
||
total, err = h.repository.CountDiagnostics(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
}
|
||
rows, err := h.repository.QueryDiagnostics(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
if !query.IncludeTotal {
|
||
total = int64(len(rows))
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||
"items": rows,
|
||
"total": total,
|
||
"limit": query.Limit,
|
||
"offset": query.Offset,
|
||
})
|
||
}
|
||
|
||
func (h *DataSourceHandler) handleJT808IdentityGaps(w http.ResponseWriter, r *http.Request) {
|
||
query, err := parseJT808IdentityGapQuery(r)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
var total int64
|
||
if query.IncludeTotal {
|
||
total, err = h.repository.CountJT808IdentityGaps(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
}
|
||
rows, err := h.repository.QueryJT808IdentityGaps(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
if !query.IncludeTotal {
|
||
total = int64(len(rows))
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||
"items": rows,
|
||
"total": total,
|
||
"limit": query.Limit,
|
||
"offset": query.Offset,
|
||
"recentSeconds": query.RecentSeconds,
|
||
})
|
||
}
|
||
|
||
func (h *DataSourceHandler) handleJT808MappingGaps(w http.ResponseWriter, r *http.Request) {
|
||
query, err := parseJT808MappingGapQuery(r)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
var total int64
|
||
if query.IncludeTotal {
|
||
total, err = h.repository.CountJT808MappingGaps(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
}
|
||
rows, err := h.repository.QueryJT808MappingGaps(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
if !query.IncludeTotal {
|
||
total = int64(len(rows))
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||
"items": rows,
|
||
"total": total,
|
||
"limit": query.Limit,
|
||
"offset": query.Offset,
|
||
"recentSeconds": query.RecentSeconds,
|
||
})
|
||
}
|
||
|
||
func (h *DataSourceHandler) handleQuery(w http.ResponseWriter, r *http.Request) {
|
||
query, err := parseDataSourceQuery(r)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
var total int64
|
||
if query.IncludeTotal {
|
||
total, err = h.repository.Count(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
}
|
||
rows, err := h.repository.Query(r.Context(), query)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusInternalServerError, err.Error())
|
||
return
|
||
}
|
||
if !query.IncludeTotal {
|
||
total = int64(len(rows))
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||
"items": rows,
|
||
"total": total,
|
||
"limit": query.Limit,
|
||
"offset": query.Offset,
|
||
})
|
||
}
|
||
|
||
func (h *DataSourceHandler) handlePatch(w http.ResponseWriter, r *http.Request) {
|
||
id, err := parseDataSourceID(strings.TrimPrefix(strings.Trim(r.URL.Path, "/"), "api/stats/data-sources/"))
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
defer r.Body.Close()
|
||
var update DataSourceUpdate
|
||
if err := json.NewDecoder(r.Body).Decode(&update); err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, "invalid json body")
|
||
return
|
||
}
|
||
updated, err := h.repository.Update(r.Context(), id, update)
|
||
if err != nil {
|
||
writeMetricError(w, http.StatusBadRequest, err.Error())
|
||
return
|
||
}
|
||
if !updated {
|
||
writeMetricError(w, http.StatusNotFound, "data source not found")
|
||
return
|
||
}
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(map[string]any{"updated": true, "id": id})
|
||
}
|
||
|
||
func parseDataSourceQuery(r *http.Request) (DataSourceQuery, error) {
|
||
values := r.URL.Query()
|
||
limit, err := parseBoundedInt(values.Get("limit"), 50, 1, 1000, "limit")
|
||
if err != nil {
|
||
return DataSourceQuery{}, err
|
||
}
|
||
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
||
if err != nil {
|
||
return DataSourceQuery{}, err
|
||
}
|
||
enabled, err := parseOptionalBool(values.Get("enabled"))
|
||
if err != nil {
|
||
return DataSourceQuery{}, err
|
||
}
|
||
sourceCodeMissing, err := parseOptionalBool(values.Get("sourceCodeMissing"))
|
||
if err != nil {
|
||
return DataSourceQuery{}, errors.New("sourceCodeMissing must be true or false")
|
||
}
|
||
if strings.TrimSpace(values.Get("sourceCode")) != "" && sourceCodeMissing != nil {
|
||
return DataSourceQuery{}, errors.New("sourceCode and sourceCodeMissing cannot be used together")
|
||
}
|
||
if strings.TrimSpace(values.Get("sourceKind")) != "" {
|
||
if _, err := parseSourceKind(values.Get("sourceKind")); err != nil {
|
||
return DataSourceQuery{}, err
|
||
}
|
||
}
|
||
return normalizeDataSourceQuery(DataSourceQuery{
|
||
Protocol: values.Get("protocol"),
|
||
SourceIP: values.Get("sourceIP"),
|
||
SourceCode: values.Get("sourceCode"),
|
||
SourceKind: values.Get("sourceKind"),
|
||
SourceCodeMissing: sourceCodeMissing,
|
||
Enabled: enabled,
|
||
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
||
Limit: limit,
|
||
Offset: offset,
|
||
}), nil
|
||
}
|
||
|
||
func parseDataSourceDiagnosticsQuery(r *http.Request) (DataSourceDiagnosticsQuery, error) {
|
||
values := r.URL.Query()
|
||
limit, err := parseBoundedInt(values.Get("limit"), 50, 1, 1000, "limit")
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
sourceCodeMissing, err := parseOptionalBool(values.Get("sourceCodeMissing"))
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, errors.New("sourceCodeMissing must be true or false")
|
||
}
|
||
mappingIssueOnly, err := parseOptionalBool(values.Get("mappingIssueOnly"))
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, errors.New("mappingIssueOnly must be true or false")
|
||
}
|
||
enabled, err := parseOptionalBool(values.Get("enabled"))
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
if sourceCodeMissing == nil {
|
||
missing := true
|
||
sourceCodeMissing = &missing
|
||
}
|
||
query := normalizeDataSourceDiagnosticsQuery(DataSourceDiagnosticsQuery{
|
||
Protocol: values.Get("protocol"),
|
||
SourceIP: values.Get("sourceIP"),
|
||
SourceCodeMissing: sourceCodeMissing,
|
||
MappingIssueOnly: mappingIssueOnly != nil && *mappingIssueOnly,
|
||
Enabled: enabled,
|
||
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
||
Limit: limit,
|
||
Offset: offset,
|
||
})
|
||
return query, nil
|
||
}
|
||
|
||
func parseDataSourceKindSuggestionsQuery(r *http.Request) (DataSourceDiagnosticsQuery, error) {
|
||
values := r.URL.Query()
|
||
limit, err := parseBoundedInt(values.Get("limit"), 50, 1, 1000, "limit")
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
sourceKind := strings.TrimSpace(values.Get("sourceKind"))
|
||
if sourceKind == "" {
|
||
sourceKind = "UNKNOWN"
|
||
}
|
||
if _, err := parseSourceKind(sourceKind); err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
sourceCodeMissing, err := parseOptionalBool(values.Get("sourceCodeMissing"))
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, errors.New("sourceCodeMissing must be true or false")
|
||
}
|
||
enabled, err := parseOptionalBool(values.Get("enabled"))
|
||
if err != nil {
|
||
return DataSourceDiagnosticsQuery{}, err
|
||
}
|
||
query := normalizeDataSourceDiagnosticsQuery(DataSourceDiagnosticsQuery{
|
||
Protocol: values.Get("protocol"),
|
||
SourceIP: values.Get("sourceIP"),
|
||
SourceKind: sourceKind,
|
||
SourceCodeMissing: sourceCodeMissing,
|
||
Enabled: enabled,
|
||
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
||
Limit: limit,
|
||
Offset: offset,
|
||
})
|
||
return query, nil
|
||
}
|
||
|
||
func parseJT808IdentityGapQuery(r *http.Request) (JT808IdentityGapQuery, error) {
|
||
values := r.URL.Query()
|
||
limit, err := parseBoundedInt(values.Get("limit"), 50, 1, 1000, "limit")
|
||
if err != nil {
|
||
return JT808IdentityGapQuery{}, err
|
||
}
|
||
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
||
if err != nil {
|
||
return JT808IdentityGapQuery{}, err
|
||
}
|
||
recentSeconds, err := parseBoundedInt(values.Get("recentSeconds"), 86400, 0, 30*86400, "recentSeconds")
|
||
if err != nil {
|
||
return JT808IdentityGapQuery{}, err
|
||
}
|
||
sourceCode := strings.TrimSpace(values.Get("sourceCode"))
|
||
if len(sourceCode) > 64 || !sourceCodePattern.MatchString(sourceCode) {
|
||
return JT808IdentityGapQuery{}, errors.New("sourceCode may contain only letters, digits, underscore, dot and dash")
|
||
}
|
||
return normalizeJT808IdentityGapQuery(JT808IdentityGapQuery{
|
||
SourceIP: values.Get("sourceIP"),
|
||
SourceCode: sourceCode,
|
||
RecentSeconds: int64(recentSeconds),
|
||
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
||
Limit: limit,
|
||
Offset: offset,
|
||
}), nil
|
||
}
|
||
|
||
func parseJT808MappingGapQuery(r *http.Request) (JT808MappingGapQuery, error) {
|
||
values := r.URL.Query()
|
||
limit, err := parseBoundedInt(values.Get("limit"), 50, 1, 1000, "limit")
|
||
if err != nil {
|
||
return JT808MappingGapQuery{}, err
|
||
}
|
||
offset, err := parseBoundedInt(values.Get("offset"), 0, 0, 1_000_000, "offset")
|
||
if err != nil {
|
||
return JT808MappingGapQuery{}, err
|
||
}
|
||
recentSeconds, err := parseBoundedInt(values.Get("recentSeconds"), 86400, 0, 30*86400, "recentSeconds")
|
||
if err != nil {
|
||
return JT808MappingGapQuery{}, err
|
||
}
|
||
sourceCode := strings.TrimSpace(values.Get("sourceCode"))
|
||
if len(sourceCode) > 64 || !sourceCodePattern.MatchString(sourceCode) {
|
||
return JT808MappingGapQuery{}, errors.New("sourceCode may contain only letters, digits, underscore, dot and dash")
|
||
}
|
||
return normalizeJT808MappingGapQuery(JT808MappingGapQuery{
|
||
SourceIP: values.Get("sourceIP"),
|
||
SourceCode: sourceCode,
|
||
RecentSeconds: int64(recentSeconds),
|
||
IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"),
|
||
Limit: limit,
|
||
Offset: offset,
|
||
}), nil
|
||
}
|
||
|
||
func parseOptionalBool(value string) (*bool, error) {
|
||
value = strings.TrimSpace(value)
|
||
if value == "" {
|
||
return nil, nil
|
||
}
|
||
switch strings.ToLower(value) {
|
||
case "true", "1", "yes":
|
||
enabled := true
|
||
return &enabled, nil
|
||
case "false", "0", "no":
|
||
enabled := false
|
||
return &enabled, nil
|
||
default:
|
||
return nil, errors.New("enabled must be true or false")
|
||
}
|
||
}
|
||
|
||
func parseDataSourceID(value string) (int64, error) {
|
||
value = strings.Trim(strings.TrimSpace(value), "/")
|
||
if value == "" || strings.Contains(value, "/") {
|
||
return 0, errors.New("data source id is required")
|
||
}
|
||
id, err := strconv.ParseInt(value, 10, 64)
|
||
if err != nil || id <= 0 {
|
||
return 0, errors.New("data source id must be positive")
|
||
}
|
||
return id, nil
|
||
}
|
||
|
||
var sourceCodePattern = regexp.MustCompile(`^[A-Za-z0-9_.-]*$`)
|
||
|
||
func validateDataSourceUpdate(update DataSourceUpdate) error {
|
||
if update.PlatformName != nil && len([]rune(strings.TrimSpace(*update.PlatformName))) > 128 {
|
||
return errors.New("platform_name length exceeds 128")
|
||
}
|
||
if update.SourceCode != nil {
|
||
sourceCode := strings.TrimSpace(*update.SourceCode)
|
||
if len(sourceCode) > 64 {
|
||
return errors.New("source_code length exceeds 64")
|
||
}
|
||
if !sourceCodePattern.MatchString(sourceCode) {
|
||
return errors.New("source_code may only contain letters, numbers, dot, underscore or dash")
|
||
}
|
||
}
|
||
if update.SourceKind != nil {
|
||
if _, err := parseSourceKind(*update.SourceKind); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
if update.TrustPriority != nil && (*update.TrustPriority < 0 || *update.TrustPriority > 100000) {
|
||
return errors.New("trust_priority must be between 0 and 100000")
|
||
}
|
||
if update.Remark != nil && len([]rune(strings.TrimSpace(*update.Remark))) > 512 {
|
||
return errors.New("remark length exceeds 512")
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func parseSourceKind(value string) (string, error) {
|
||
kind := normalizeSourceKindForRead(value)
|
||
switch kind {
|
||
case "UNKNOWN", "PLATFORM", "DIRECT":
|
||
return kind, nil
|
||
default:
|
||
return "", errors.New("source_kind must be UNKNOWN, PLATFORM or DIRECT")
|
||
}
|
||
}
|
||
|
||
func normalizeSourceKindForRead(value string) string {
|
||
value = strings.ToUpper(strings.TrimSpace(value))
|
||
if value == "" {
|
||
return "UNKNOWN"
|
||
}
|
||
return value
|
||
}
|
||
|
||
func normalizeSourceKindForQuery(value string) string {
|
||
return strings.ToUpper(strings.TrimSpace(value))
|
||
}
|
||
|
||
func normalizeSourceKindForWrite(value string) string {
|
||
kind, err := parseSourceKind(value)
|
||
if err != nil {
|
||
return "UNKNOWN"
|
||
}
|
||
return kind
|
||
}
|
||
|
||
func nullableTrimmedString(value string) any {
|
||
value = strings.TrimSpace(value)
|
||
if value == "" {
|
||
return nil
|
||
}
|
||
return value
|
||
}
|
||
|
||
func splitCommaList(value string) []string {
|
||
value = strings.TrimSpace(value)
|
||
if value == "" {
|
||
return nil
|
||
}
|
||
parts := strings.Split(value, ",")
|
||
out := make([]string, 0, len(parts))
|
||
for _, part := range parts {
|
||
part = strings.TrimSpace(part)
|
||
if part != "" {
|
||
out = append(out, part)
|
||
}
|
||
}
|
||
return out
|
||
}
|
||
|
||
func diagnoseDataSource(row DataSourceDiagnosticRow) (string, string) {
|
||
if strings.ToUpper(strings.TrimSpace(row.Protocol)) != "JT808" {
|
||
if strings.TrimSpace(row.SourceCode) == "" {
|
||
return "source_code_missing", "补充来源平台名称和 source_code,用于多源统计选举、断链告警和大屏展示"
|
||
}
|
||
return "source_configured", "该来源已具备基础平台标识;如参与统计选举,请确认 source_kind、trust_priority 和 enabled"
|
||
}
|
||
switch {
|
||
case row.RegistrationRows == 0:
|
||
return "no_registration", "等待该来源产生 JT808 注册/鉴权/位置记录,或核对 source_endpoint 是否被正确写入 jt808_registration"
|
||
case lowIdentifierCoverage(row):
|
||
return "low_identifier_coverage", "该来源已配置平台,但大量注册手机号未维护到 vehicle_identifier,会影响 VIN 解析、在线统计和里程来源选举"
|
||
case row.SourcePlatformNameMismatch:
|
||
return "source_platform_name_mismatch", "当前 source_code 在 vehicle_identifier 中对应的平台名与 vehicle_data_source.platform_name 不一致;先确认并修正来源配置,再导入手机号映射"
|
||
case row.IdentifierMatchedPhones == 0:
|
||
return "no_identifier_match", "补充该来源手机号到 vehicle_identifier,或先导入对应平台的 808 车牌/手机号映射文件"
|
||
case row.ConfiguredSourceCodeConflict:
|
||
return "source_code_conflict", "已配置的 source_code 与该来源手机号绑定推断出的 source_code 不一致,需确认来源平台或修正 vehicle_identifier"
|
||
case strings.TrimSpace(row.SourceCode) != "":
|
||
return "source_configured", "该来源已具备基础平台标识;继续维护缺失手机号可提升 VIN 解析和统计覆盖率"
|
||
case row.MatchedSourceCodeCount > 1:
|
||
return "ambiguous_source_code", "同一来源 IP 匹配到多个 source_code,需人工确认平台并在 vehicle_data_source 固定 source_code"
|
||
case row.CandidateSourceCode != "":
|
||
return "candidate_available", "可用候选 source_code,执行 identity-import -sync-data-sources -apply 或在来源管理中确认"
|
||
default:
|
||
return "no_source_code", "已匹配到标识但没有 source_code,补齐 vehicle_identifier.source_code 后重新同步"
|
||
}
|
||
}
|
||
|
||
func diagnoseJT808IdentityGap(row JT808IdentityGapRow) (string, string) {
|
||
if strings.TrimSpace(row.Plate) != "" {
|
||
return "missing_phone_and_plate_binding", "将该 phone 或 plate 维护到 vehicle_identifier,并确认 VIN、source_code、oem;若仍使用旧表,也同步 vehicle_identity_binding"
|
||
}
|
||
return "missing_phone_binding", "将该 phone 维护到 vehicle_identifier,并确认 VIN、source_code、oem;如终端会补发注册帧,可等待 plate 后再复核"
|
||
}
|
||
|
||
func jt808IdentityGapRawFrameQueryPath(row JT808IdentityGapRow) string {
|
||
if strings.TrimSpace(row.Phone) == "" {
|
||
return ""
|
||
}
|
||
dateFrom, dateTo := jt808IdentityGapDateWindow(row.LatestSeenAt)
|
||
values := url.Values{}
|
||
values.Set("protocol", "JT808")
|
||
values.Set("phone", row.Phone)
|
||
values.Set("dateFrom", dateFrom)
|
||
values.Set("dateTo", dateTo)
|
||
values.Set("orderBy", "eventTime")
|
||
values.Set("includeFields", "true")
|
||
values.Set("limit", "20")
|
||
return "/api/history/raw-frames?" + values.Encode()
|
||
}
|
||
|
||
func jt808IdentityGapDataSourceQueryPath(row JT808IdentityGapRow) string {
|
||
if strings.TrimSpace(row.SourceIP) == "" {
|
||
return ""
|
||
}
|
||
values := url.Values{}
|
||
values.Set("protocol", "JT808")
|
||
values.Set("sourceIP", row.SourceIP)
|
||
values.Set("includeTotal", "true")
|
||
values.Set("limit", "20")
|
||
return "/api/stats/data-sources?" + values.Encode()
|
||
}
|
||
|
||
func annotateJT808MappingGap(row *JT808MappingGapRow) {
|
||
if row == nil {
|
||
return
|
||
}
|
||
row.SuggestedSourceCode = strings.TrimSpace(row.SourceCode)
|
||
row.SuggestedPlatformName = strings.TrimSpace(firstNonEmpty(row.PlatformName, row.SourceCode))
|
||
row.SuggestedIdentifierType = "JT808_PHONE"
|
||
row.SuggestedIdentifierValue = strings.TrimSpace(row.Phone)
|
||
row.SuggestedVIN = suggestedMappingVIN(*row)
|
||
row.SuggestedPlate = strings.TrimSpace(firstNonEmpty(row.Plate, row.IdentifierPlate))
|
||
row.Reason, row.RecommendedAction = diagnoseJT808MappingGap(*row)
|
||
row.RawFrameQueryPath = jt808MappingGapRawFrameQueryPath(*row)
|
||
row.DataSourceQueryPath = jt808MappingGapDataSourceQueryPath(*row)
|
||
row.VehicleIdentifierExample = vehicleIdentifierExample(*row)
|
||
}
|
||
|
||
func diagnoseJT808MappingGap(row JT808MappingGapRow) (string, string) {
|
||
if strings.TrimSpace(row.SourceCode) == "" {
|
||
return "source_code_missing", "先维护 vehicle_data_source.source_code,再按来源平台导入 phone 到 VIN 的 vehicle_identifier 映射"
|
||
}
|
||
if strings.TrimSpace(row.IdentifierVIN) == "" && strings.TrimSpace(row.VIN) != "" && !strings.EqualFold(strings.TrimSpace(row.VIN), "unknown") {
|
||
return "missing_source_phone_identifier", "按当前来源 source_code 维护 JT808_PHONE 到 VIN 的映射;如平台名与 source_code 不一致,先修正来源配置"
|
||
}
|
||
if len(row.MatchedSourceCodes) > 0 {
|
||
return "source_identifier_mismatch", "该 phone 存在其他 source_code 映射,但当前来源缺少映射;确认平台归属后补当前 source_code 或修正来源配置"
|
||
}
|
||
return "missing_source_phone_identifier", "补充该 phone 在当前来源 source_code 下的 vehicle_identifier 映射"
|
||
}
|
||
|
||
func suggestedMappingVIN(row JT808MappingGapRow) string {
|
||
vin := strings.TrimSpace(row.VIN)
|
||
if vin != "" && !strings.EqualFold(vin, "unknown") {
|
||
return vin
|
||
}
|
||
return strings.TrimSpace(row.IdentifierVIN)
|
||
}
|
||
|
||
func jt808MappingGapRawFrameQueryPath(row JT808MappingGapRow) string {
|
||
if strings.TrimSpace(row.Phone) == "" {
|
||
return ""
|
||
}
|
||
dateFrom, dateTo := jt808IdentityGapDateWindow(row.LatestSeenAt)
|
||
values := url.Values{}
|
||
values.Set("protocol", "JT808")
|
||
values.Set("phone", row.Phone)
|
||
values.Set("dateFrom", dateFrom)
|
||
values.Set("dateTo", dateTo)
|
||
values.Set("orderBy", "eventTime")
|
||
values.Set("includeFields", "true")
|
||
values.Set("limit", "20")
|
||
return "/api/history/raw-frames?" + values.Encode()
|
||
}
|
||
|
||
func jt808MappingGapDataSourceQueryPath(row JT808MappingGapRow) string {
|
||
if strings.TrimSpace(row.SourceIP) == "" {
|
||
return ""
|
||
}
|
||
values := url.Values{}
|
||
values.Set("protocol", "JT808")
|
||
values.Set("sourceIP", row.SourceIP)
|
||
values.Set("includeTotal", "true")
|
||
values.Set("limit", "20")
|
||
return "/api/stats/data-sources?" + values.Encode()
|
||
}
|
||
|
||
func vehicleIdentifierExample(row JT808MappingGapRow) string {
|
||
sourceCode := strings.TrimSpace(row.SuggestedSourceCode)
|
||
phone := strings.TrimSpace(row.SuggestedIdentifierValue)
|
||
vin := strings.TrimSpace(row.SuggestedVIN)
|
||
if sourceCode == "" || phone == "" || vin == "" {
|
||
return ""
|
||
}
|
||
plate := strings.ReplaceAll(strings.TrimSpace(row.SuggestedPlate), "'", "''")
|
||
oem := strings.ReplaceAll(strings.TrimSpace(row.SuggestedPlatformName), "'", "''")
|
||
sourceCode = strings.ReplaceAll(sourceCode, "'", "''")
|
||
phone = strings.ReplaceAll(phone, "'", "''")
|
||
vin = strings.ReplaceAll(vin, "'", "''")
|
||
return fmt.Sprintf("protocol=JT808, source_code=%s, identifier_type=JT808_PHONE, identifier_value=%s, vin=%s, plate=%s, oem=%s", sourceCode, phone, vin, plate, oem)
|
||
}
|
||
|
||
func jt808IdentityGapDateWindow(value string) (string, string) {
|
||
value = strings.TrimSpace(value)
|
||
for _, layout := range []string{"2006-01-02 15:04:05", time.RFC3339} {
|
||
seenAt, err := time.ParseInLocation(layout, value, time.FixedZone("Asia/Shanghai", 8*3600))
|
||
if err == nil {
|
||
start := time.Date(seenAt.Year(), seenAt.Month(), seenAt.Day(), 0, 0, 0, 0, seenAt.Location())
|
||
return start.Format("2006-01-02 15:04:05"), start.AddDate(0, 0, 1).Format("2006-01-02 15:04:05")
|
||
}
|
||
}
|
||
now := time.Now().In(time.FixedZone("Asia/Shanghai", 8*3600))
|
||
start := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location())
|
||
return start.Format("2006-01-02 15:04:05"), start.AddDate(0, 0, 1).Format("2006-01-02 15:04:05")
|
||
}
|
||
|
||
func suggestSourceKind(row DataSourceDiagnosticRow) (kind string, confidence string, reason string) {
|
||
current := normalizeSourceKindForRead(row.SourceKind)
|
||
if current != "UNKNOWN" {
|
||
return current, "HIGH", "already_classified"
|
||
}
|
||
if strings.ToUpper(strings.TrimSpace(row.Protocol)) != "JT808" &&
|
||
(strings.TrimSpace(row.SourceCode) != "" || strings.TrimSpace(row.PlatformName) != "") {
|
||
return "PLATFORM", "MEDIUM", "non_jt808_configured_source"
|
||
}
|
||
if row.Reason == "ambiguous_source_code" || row.MatchedSourceCodeCount > 1 {
|
||
return "UNKNOWN", "HIGH", "ambiguous_source_code_requires_manual_review"
|
||
}
|
||
if row.Reason == "no_identifier_match" {
|
||
return "UNKNOWN", "MEDIUM", "missing_vehicle_identifier_mapping"
|
||
}
|
||
if row.CandidateSourceCode != "" && row.MatchedSourceCodeCount == 1 {
|
||
if row.PhoneCount >= 5 || row.ActiveSpanSeconds >= 6*3600 {
|
||
return "PLATFORM", "HIGH", "single_source_code_with_many_phones_or_long_activity"
|
||
}
|
||
if row.PhoneCount <= 1 && row.ActiveSpanSeconds < 2*3600 && row.LatestSeenAgeSeconds >= 3600 {
|
||
return "DIRECT", "MEDIUM", "single_phone_short_lived_stale_source"
|
||
}
|
||
return "PLATFORM", "MEDIUM", "single_source_code_candidate"
|
||
}
|
||
if row.Reason == "no_registration" {
|
||
if strings.TrimSpace(row.SourceCode) != "" || strings.TrimSpace(row.PlatformName) != "" {
|
||
return "UNKNOWN", "MEDIUM", "manual_source_without_registration_evidence"
|
||
}
|
||
if row.LatestSeenAgeSeconds >= 24*3600 {
|
||
return "DIRECT", "LOW", "stale_unclassified_source_without_registration"
|
||
}
|
||
}
|
||
return "UNKNOWN", "LOW", "insufficient_evidence"
|
||
}
|
||
|
||
func configuredSourceCodeConflict(row DataSourceDiagnosticRow) bool {
|
||
sourceCode := strings.TrimSpace(row.SourceCode)
|
||
if sourceCode == "" || row.MatchedSourceCodeCount == 0 {
|
||
return false
|
||
}
|
||
for _, matched := range row.MatchedSourceCodes {
|
||
if strings.EqualFold(strings.TrimSpace(matched), sourceCode) {
|
||
return false
|
||
}
|
||
}
|
||
return true
|
||
}
|
||
|
||
func sourcePlatformNameMismatch(row DataSourceDiagnosticRow) bool {
|
||
platformName := strings.TrimSpace(row.PlatformName)
|
||
configuredPlatformName := strings.TrimSpace(row.ConfiguredSourceCodePlatformName)
|
||
if platformName == "" || configuredPlatformName == "" {
|
||
return false
|
||
}
|
||
return platformName != configuredPlatformName
|
||
}
|
||
|
||
func lowIdentifierCoverage(row DataSourceDiagnosticRow) bool {
|
||
if strings.ToUpper(strings.TrimSpace(row.Protocol)) != "JT808" {
|
||
return false
|
||
}
|
||
if strings.TrimSpace(row.SourceCode) == "" || row.PhoneCount < 10 {
|
||
return false
|
||
}
|
||
return row.UnmappedPhoneCount > 0 && row.IdentifierMatchRatio < 0.8
|
||
}
|
||
|
||
func (q DataSourceQuery) String() string {
|
||
enabled := "any"
|
||
if q.Enabled != nil {
|
||
enabled = fmt.Sprint(*q.Enabled)
|
||
}
|
||
sourceCodeMissing := "any"
|
||
if q.SourceCodeMissing != nil {
|
||
sourceCodeMissing = fmt.Sprint(*q.SourceCodeMissing)
|
||
}
|
||
return fmt.Sprintf("protocol=%s sourceIP=%s sourceCode=%s sourceKind=%s sourceCodeMissing=%s enabled=%s limit=%d offset=%d", q.Protocol, q.SourceIP, q.SourceCode, q.SourceKind, sourceCodeMissing, enabled, q.Limit, q.Offset)
|
||
}
|