feat(operations): add automated reconciliation center

This commit is contained in:
lingniu
2026-07-16 19:14:53 +08:00
parent bbab018d55
commit d67f42b4f4
25 changed files with 1697 additions and 39 deletions

View File

@@ -232,6 +232,12 @@ func requiredMenu(r *http.Request) string {
}
func requiredRole(r *http.Request) string {
if strings.HasPrefix(r.URL.Path, "/api/v2/reconciliation/") {
if r.Method == http.MethodGet || r.Method == http.MethodHead || r.Method == http.MethodPost {
return "operator"
}
return "admin"
}
if strings.HasPrefix(r.URL.Path, "/api/v2/operations/") {
if r.Method == http.MethodGet || r.Method == http.MethodHead {
return "operator"

View File

@@ -106,6 +106,14 @@ func TestAPIAuthEnforcesTokensAndRoleBoundaries(t *testing.T) {
if operatorSourcePolicy.Code != http.StatusForbidden {
t.Fatalf("operator source policy mutation should be forbidden, status=%d", operatorSourcePolicy.Code)
}
viewerReconciliation := authRequest(t, cfg, http.MethodGet, "/api/v2/reconciliation/summary", viewerToken)
if viewerReconciliation.Code != http.StatusForbidden {
t.Fatalf("viewer reconciliation summary should be forbidden, status=%d", viewerReconciliation.Code)
}
operatorReconciliation := authRequest(t, cfg, http.MethodPost, "/api/v2/reconciliation/issues/reconciliation-1/actions", operatorToken)
if operatorReconciliation.Code != http.StatusNoContent {
t.Fatalf("operator reconciliation action status=%d body=%s", operatorReconciliation.Code, operatorReconciliation.Body.String())
}
adminProfile := authRequest(t, cfg, http.MethodPut, "/api/v2/vehicles/VIN001/profile", adminToken)
if adminProfile.Code != http.StatusNoContent || adminProfile.Header().Get("X-Principal") != "admin-a:admin" {
t.Fatalf("admin profile mutation status=%d principal=%s", adminProfile.Code, adminProfile.Header().Get("X-Principal"))

View File

@@ -63,6 +63,10 @@ func (h *Handler) routes() {
h.mux.HandleFunc("GET /api/v2/vehicles/{vin}/source-evidence", h.handleVehicleSourceEvidence)
h.mux.HandleFunc("GET /api/v2/operations/vehicles/{vin}/sources", h.handleVehicleSourceDiagnostic)
h.mux.HandleFunc("PUT /api/v2/operations/vehicles/{vin}/sources/{sourceRef}", h.handleUpdateVehicleSourcePolicy)
h.mux.HandleFunc("GET /api/v2/reconciliation/summary", h.handleReconciliationSummary)
h.mux.HandleFunc("POST /api/v2/reconciliation/issues", h.handleReconciliationIssues)
h.mux.HandleFunc("GET /api/v2/reconciliation/issues/{id}", h.handleReconciliationIssue)
h.mux.HandleFunc("POST /api/v2/reconciliation/issues/{id}/actions", h.handleReconciliationAction)
h.mux.HandleFunc("POST /api/v2/vehicle-profiles/sync", h.handleSyncVehicleProfiles)
h.mux.HandleFunc("GET /api/v2/tracks", h.handleTrackPlayback)
h.mux.HandleFunc("GET /api/v2/metrics", h.handleMetricCatalog)
@@ -89,6 +93,34 @@ func (h *Handler) routes() {
h.mux.HandleFunc("POST /api/v2/alerts/notifications/read", h.handleAlertNotificationsRead)
}
func (h *Handler) handleReconciliationSummary(w http.ResponseWriter, r *http.Request) {
data, err := h.service.ReconciliationSummary(r.Context(), parsePositive(r.URL.Query().Get("days"), 30))
h.write(w, r, data, err)
}
func (h *Handler) handleReconciliationIssues(w http.ResponseWriter, r *http.Request) {
var query ReconciliationQuery
if !decodeJSONBody(w, r, &query) {
return
}
data, err := h.service.ReconciliationIssues(r.Context(), query)
h.write(w, r, data, err)
}
func (h *Handler) handleReconciliationIssue(w http.ResponseWriter, r *http.Request) {
data, err := h.service.ReconciliationIssue(r.Context(), r.PathValue("id"))
h.write(w, r, data, err)
}
func (h *Handler) handleReconciliationAction(w http.ResponseWriter, r *http.Request) {
var request ReconciliationActionRequest
if !decodeJSONBody(w, r, &request) {
return
}
data, err := h.service.UpdateReconciliationIssue(r.Context(), r.PathValue("id"), request)
h.write(w, r, data, err)
}
func (h *Handler) handleAccessUnresolvedIdentities(w http.ResponseWriter, r *http.Request) {
var query AccessUnresolvedIdentityQuery
if !decodeJSONBody(w, r, &query) {

View File

@@ -29,6 +29,36 @@ func TestHandlerDashboardSummary(t *testing.T) {
}
}
func TestHandlerReconciliationQueueAndReview(t *testing.T) {
handler := NewHandler(NewService(NewMockStore()))
summary := httptest.NewRecorder()
handler.ServeHTTP(summary, httptest.NewRequest(http.MethodGet, "/api/v2/reconciliation/summary?days=30", nil))
if summary.Code != http.StatusOK || !strings.Contains(summary.Body.String(), `"active":1`) {
t.Fatalf("summary status=%d body=%s", summary.Code, summary.Body.String())
}
list := httptest.NewRecorder()
handler.ServeHTTP(list, httptest.NewRequest(http.MethodPost, "/api/v2/reconciliation/issues", strings.NewReader(`{"status":"active","limit":20}`)))
if list.Code != http.StatusOK || !strings.Contains(list.Body.String(), "POSITION_DRIFT") {
t.Fatalf("list status=%d body=%s", list.Code, list.Body.String())
}
detail := httptest.NewRecorder()
handler.ServeHTTP(detail, httptest.NewRequest(http.MethodGet, "/api/v2/reconciliation/issues/reconciliation-demo-position", nil))
if detail.Code != http.StatusOK || !strings.Contains(detail.Body.String(), "规则首次发现差异") {
t.Fatalf("detail status=%d body=%s", detail.Code, detail.Body.String())
}
review := httptest.NewRecorder()
request := httptest.NewRequest(http.MethodPost, "/api/v2/reconciliation/issues/reconciliation-demo-position/actions", strings.NewReader(`{"version":1,"status":"confirmed_source_a","note":"已核对原始报文,来源 A 可信"}`))
request = request.WithContext(WithPrincipal(request.Context(), Principal{Name: "运维甲", Role: "operator", UserType: "operator"}))
handler.ServeHTTP(review, request)
if review.Code != http.StatusOK || !strings.Contains(review.Body.String(), `"status":"confirmed_source_a"`) || !strings.Contains(review.Body.String(), `"actor":"运维甲"`) {
t.Fatalf("review status=%d body=%s", review.Code, review.Body.String())
}
}
func TestHandlerV2MonitorSummaryAndMap(t *testing.T) {
handler := NewHandler(NewService(NewMockStore()))
for _, test := range []struct {

View File

@@ -11,20 +11,22 @@ import (
)
type MockStore struct {
vehicles []VehicleRow
locations []RealtimeLocationRow
accessMu sync.RWMutex
accessThresholds AccessThresholdConfig
sourcePolicyMu sync.RWMutex
sourcePolicies map[string]VehicleSourcePolicyConfig
profileMu sync.RWMutex
profiles map[string]VehicleProfile
alertMu sync.RWMutex
alertRules []AlertRule
alertEvents []AlertEvent
alertNotifications []AlertNotification
nextAlertActionID int64
nextNotificationID int64
vehicles []VehicleRow
locations []RealtimeLocationRow
accessMu sync.RWMutex
accessThresholds AccessThresholdConfig
sourcePolicyMu sync.RWMutex
sourcePolicies map[string]VehicleSourcePolicyConfig
profileMu sync.RWMutex
profiles map[string]VehicleProfile
alertMu sync.RWMutex
alertRules []AlertRule
alertEvents []AlertEvent
alertNotifications []AlertNotification
nextAlertActionID int64
nextNotificationID int64
reconciliationMu sync.RWMutex
reconciliationIssues []ReconciliationIssue
}
func NewMockStore() *MockStore {
@@ -49,11 +51,169 @@ func NewMockStore() *MockStore {
},
}
store.seedAlertCenter()
store.seedReconciliationCenter()
return store
}
func int64Pointer(value int64) *int64 { return &value }
func (m *MockStore) seedReconciliationCenter() {
m.reconciliationIssues = []ReconciliationIssue{
{
ID: "reconciliation-demo-position", RuleCode: "POSITION_DRIFT", Category: "location", Severity: "major", Status: "pending",
VIN: "LB9A32A24R0LS1426", Plate: "粤AG18312", ProtocolA: "GB32960", ProtocolB: "JT808",
Title: "多来源实时位置漂移", Summary: "两个有效来源位置相差 1,286 米,系统不自动判定哪一来源正确",
Evidence: map[string]any{"distanceM": 1286, "sourceA": map[string]any{"protocol": "GB32960"}, "sourceB": map[string]any{"protocol": "JT808"}},
FirstSeenAt: "2026-07-16 08:00:00", LastSeenAt: "2026-07-16 10:00:00", OccurrenceCount: 3, Version: 1,
Actions: []ReconciliationAction{{ID: 1, Action: "detect", FromStatus: "", ToStatus: "pending", Actor: "reconciliation-evaluator", Note: "规则首次发现差异", CreatedAt: "2026-07-16 08:00:00"}},
},
{
ID: "reconciliation-demo-source", RuleCode: "SOURCE_MISSING", Category: "source", Severity: "minor", Status: "recovered",
VIN: "LB9A32A24P0LS1230", Plate: "粤AFF7936", Title: "主车辆暂无真实数据来源",
Summary: "车辆已建档,但当前没有真实来源证据", Evidence: map[string]any{"vin": "LB9A32A24P0LS1230"},
FirstSeenAt: "2026-07-14 08:00:00", LastSeenAt: "2026-07-15 08:00:00", OccurrenceCount: 2,
RecoveredAt: "2026-07-16 02:15:00", Version: 2,
},
}
}
func cloneReconciliationIssue(item ReconciliationIssue) ReconciliationIssue {
item.Actions = append([]ReconciliationAction(nil), item.Actions...)
if item.Evidence == nil {
item.Evidence = map[string]any{}
} else {
next := make(map[string]any, len(item.Evidence))
for key, value := range item.Evidence {
next[key] = value
}
item.Evidence = next
}
return item
}
func (m *MockStore) ReconciliationSummary(_ context.Context, _ int) (ReconciliationSummary, error) {
m.reconciliationMu.RLock()
defer m.reconciliationMu.RUnlock()
result := ReconciliationSummary{
ByRule: []ReconciliationBucket{}, BySeverity: []ReconciliationBucket{},
Trend: []ReconciliationTrendPoint{{Date: "2026-07-16", Detected: 1, New: 1, Active: 1, Recovered: 1}},
LastRunAt: "2026-07-16 02:15:00", AsOf: time.Now().Format(time.RFC3339),
}
rules := map[string]int{}
severities := map[string]int{}
for _, item := range m.reconciliationIssues {
switch item.Status {
case "pending":
result.Active++
result.Pending++
case "confirmed_source_a", "confirmed_source_b":
result.Active++
result.Confirmed++
case "recovered":
result.Recovered++
}
if item.Status == "pending" || strings.HasPrefix(item.Status, "confirmed_source_") {
rules[item.RuleCode]++
severities[item.Severity]++
}
}
for name, count := range rules {
result.ByRule = append(result.ByRule, ReconciliationBucket{Name: name, Count: count})
}
for name, count := range severities {
result.BySeverity = append(result.BySeverity, ReconciliationBucket{Name: name, Count: count})
}
sort.Slice(result.ByRule, func(i, j int) bool { return result.ByRule[i].Count > result.ByRule[j].Count })
sort.Slice(result.BySeverity, func(i, j int) bool { return result.BySeverity[i].Count > result.BySeverity[j].Count })
return result, nil
}
func (m *MockStore) ReconciliationIssues(_ context.Context, query ReconciliationQuery) (Page[ReconciliationIssue], error) {
m.reconciliationMu.RLock()
defer m.reconciliationMu.RUnlock()
items := make([]ReconciliationIssue, 0, len(m.reconciliationIssues))
for _, item := range m.reconciliationIssues {
if query.Keyword != "" && !strings.Contains(strings.ToLower(item.VIN+" "+item.Plate+" "+item.Title+" "+item.Summary), strings.ToLower(query.Keyword)) {
continue
}
if query.RuleCode != "" && query.RuleCode != "all" && item.RuleCode != query.RuleCode {
continue
}
if query.Category != "" && query.Category != "all" && item.Category != query.Category {
continue
}
if query.Severity != "" && query.Severity != "all" && item.Severity != query.Severity {
continue
}
if query.Status == "active" && item.Status != "pending" && !strings.HasPrefix(item.Status, "confirmed_source_") {
continue
}
if query.Status != "" && query.Status != "all" && query.Status != "active" && item.Status != query.Status {
continue
}
cloned := cloneReconciliationIssue(item)
cloned.Actions = nil
items = append(items, cloned)
}
total := len(items)
if query.Offset >= total {
items = []ReconciliationIssue{}
} else {
end := min(total, query.Offset+query.Limit)
items = items[query.Offset:end]
}
return Page[ReconciliationIssue]{Items: items, Total: total, Limit: query.Limit, Offset: query.Offset}, nil
}
func (m *MockStore) ReconciliationIssue(_ context.Context, id string) (ReconciliationIssue, error) {
m.reconciliationMu.RLock()
defer m.reconciliationMu.RUnlock()
for _, item := range m.reconciliationIssues {
if item.ID == id {
return cloneReconciliationIssue(item), nil
}
}
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOT_FOUND", Message: "差异记录不存在"}
}
func (m *MockStore) UpdateReconciliationIssue(_ context.Context, id string, request ReconciliationActionRequest) (ReconciliationIssue, error) {
m.reconciliationMu.Lock()
defer m.reconciliationMu.Unlock()
for index := range m.reconciliationIssues {
item := &m.reconciliationIssues[index]
if item.ID != id {
continue
}
if item.Version != request.Version {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_VERSION_CONFLICT", Message: "差异记录已更新,请刷新后重试"}
}
from := item.Status
item.Status = request.Status
item.ResolutionNote = request.Note
item.ResolvedBy = request.Actor
item.Version++
if request.Status == "fixed" || request.Status == "no_action" {
item.RecoveredAt = time.Now().Format("2006-01-02 15:04:05")
} else {
item.RecoveredAt = ""
}
item.Actions = append(item.Actions, ReconciliationAction{
ID: int64(len(item.Actions) + 1), Action: "review", FromStatus: from, ToStatus: request.Status,
Actor: request.Actor, Note: request.Note, CreatedAt: time.Now().Format("2006-01-02 15:04:05"),
})
return cloneReconciliationIssue(*item), nil
}
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOT_FOUND", Message: "差异记录不存在"}
}
func (m *MockStore) EvaluateReconciliation(_ context.Context) (ReconciliationEvaluationResult, error) {
summary, _ := m.ReconciliationSummary(context.Background(), 30)
return ReconciliationEvaluationResult{
RunID: "reconciliation-run-demo", Detected: summary.Active, Active: summary.Active,
RuleCounts: map[string]int{"POSITION_DRIFT": 1}, AsOf: time.Now().Format(time.RFC3339),
}, nil
}
func (m *MockStore) VehicleLocationSourceHistory(_ context.Context, vin string) ([]vehicleLocationSourceHistory, error) {
if vin != "LB9A32A24R0LS1426" {
return []vehicleLocationSourceHistory{}, nil

View File

@@ -736,6 +736,92 @@ type AlertNotificationReadRequest struct {
Actor string `json:"actor"`
}
type ReconciliationQuery struct {
Keyword string `json:"keyword"`
RuleCode string `json:"ruleCode"`
Category string `json:"category"`
Severity string `json:"severity"`
Status string `json:"status"`
Limit int `json:"limit"`
Offset int `json:"offset"`
}
type ReconciliationIssue struct {
ID string `json:"id"`
RuleCode string `json:"ruleCode"`
Category string `json:"category"`
Severity string `json:"severity"`
Status string `json:"status"`
VIN string `json:"vin"`
Plate string `json:"plate"`
ProtocolA string `json:"protocolA"`
ProtocolB string `json:"protocolB"`
Title string `json:"title"`
Summary string `json:"summary"`
Evidence map[string]any `json:"evidence"`
FirstSeenAt string `json:"firstSeenAt"`
LastSeenAt string `json:"lastSeenAt"`
OccurrenceCount int64 `json:"occurrenceCount"`
RecoveredAt string `json:"recoveredAt"`
ResolutionNote string `json:"resolutionNote"`
ResolvedBy string `json:"resolvedBy"`
Version int `json:"version"`
Actions []ReconciliationAction `json:"actions,omitempty"`
}
type ReconciliationAction struct {
ID int64 `json:"id"`
Action string `json:"action"`
FromStatus string `json:"fromStatus"`
ToStatus string `json:"toStatus"`
Actor string `json:"actor"`
Note string `json:"note"`
CreatedAt string `json:"createdAt"`
}
type ReconciliationActionRequest struct {
Version int `json:"version"`
Status string `json:"status"`
Note string `json:"note"`
Actor string `json:"actor"`
}
type ReconciliationBucket struct {
Name string `json:"name"`
Count int `json:"count"`
}
type ReconciliationTrendPoint struct {
Date string `json:"date"`
Detected int `json:"detected"`
New int `json:"new"`
Active int `json:"active"`
Recovered int `json:"recovered"`
}
type ReconciliationSummary struct {
Active int `json:"active"`
Pending int `json:"pending"`
Confirmed int `json:"confirmed"`
Recovered int `json:"recovered"`
OverSLA int `json:"overSla"`
ByRule []ReconciliationBucket `json:"byRule"`
BySeverity []ReconciliationBucket `json:"bySeverity"`
Trend []ReconciliationTrendPoint `json:"trend"`
LastRunAt string `json:"lastRunAt"`
AsOf string `json:"asOf"`
}
type ReconciliationEvaluationResult struct {
RunID string `json:"runId"`
Detected int `json:"detected"`
New int `json:"new"`
Active int `json:"active"`
Recovered int `json:"recovered"`
RuleCounts map[string]int `json:"ruleCounts"`
AsOf string `json:"asOf"`
}
type AlertEvaluationResult struct {
RulesEvaluated int `json:"rulesEvaluated"`
VehiclesScanned int `json:"vehiclesScanned"`

View File

@@ -12,19 +12,21 @@ import (
)
type ProductionStore struct {
db *sql.DB
tdengine *sql.DB
tdDatabase string
redisOnline redisOnlineKeyCounter
capacityCheck capacityChecker
alertStreamGroup string
alertStreamMode string
accessSchemaOnce sync.Once
accessSchemaErr error
alertSchemaOnce sync.Once
alertSchemaErr error
profileSchemaOnce sync.Once
profileSchemaErr error
db *sql.DB
tdengine *sql.DB
tdDatabase string
redisOnline redisOnlineKeyCounter
capacityCheck capacityChecker
alertStreamGroup string
alertStreamMode string
accessSchemaOnce sync.Once
accessSchemaErr error
alertSchemaOnce sync.Once
alertSchemaErr error
profileSchemaOnce sync.Once
profileSchemaErr error
reconciliationSchemaOnce sync.Once
reconciliationSchemaErr error
}
type redisOnlineKeyCounter interface {

View File

@@ -0,0 +1,78 @@
package platform
import (
"context"
"fmt"
"strings"
)
func (s *Service) reconciliationStore() (ReconciliationStore, error) {
store, ok := s.store.(ReconciliationStore)
if !ok {
return nil, fmt.Errorf("store does not provide reconciliation center")
}
return store, nil
}
func (s *Service) ReconciliationSummary(ctx context.Context, days int) (ReconciliationSummary, error) {
store, err := s.reconciliationStore()
if err != nil {
return ReconciliationSummary{}, err
}
if days <= 0 || days > 90 {
days = 30
}
return store.ReconciliationSummary(ctx, days)
}
func (s *Service) ReconciliationIssues(ctx context.Context, query ReconciliationQuery) (Page[ReconciliationIssue], error) {
store, err := s.reconciliationStore()
if err != nil {
return Page[ReconciliationIssue]{}, err
}
query.Keyword = strings.TrimSpace(query.Keyword)
query.RuleCode = strings.TrimSpace(query.RuleCode)
query.Category = strings.TrimSpace(query.Category)
query.Severity = strings.TrimSpace(query.Severity)
query.Status = strings.TrimSpace(query.Status)
if query.Limit <= 0 || query.Limit > 200 {
query.Limit = 50
}
if query.Offset < 0 {
query.Offset = 0
}
return store.ReconciliationIssues(ctx, query)
}
func (s *Service) ReconciliationIssue(ctx context.Context, id string) (ReconciliationIssue, error) {
id = strings.TrimSpace(id)
if id == "" || len(id) > 64 {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_ID_INVALID", Message: "差异记录编号无效"}
}
store, err := s.reconciliationStore()
if err != nil {
return ReconciliationIssue{}, err
}
return store.ReconciliationIssue(ctx, id)
}
func (s *Service) UpdateReconciliationIssue(ctx context.Context, id string, request ReconciliationActionRequest) (ReconciliationIssue, error) {
id = strings.TrimSpace(id)
if id == "" || len(id) > 64 {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_ID_INVALID", Message: "差异记录编号无效"}
}
request.Actor = ActorFromContext(ctx)
store, err := s.reconciliationStore()
if err != nil {
return ReconciliationIssue{}, err
}
return store.UpdateReconciliationIssue(ctx, id, request)
}
func (s *Service) EvaluateReconciliation(ctx context.Context) (ReconciliationEvaluationResult, error) {
store, err := s.reconciliationStore()
if err != nil {
return ReconciliationEvaluationResult{}, err
}
return store.EvaluateReconciliation(ctx)
}

View File

@@ -0,0 +1,619 @@
package platform
import (
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"fmt"
"strings"
"time"
)
const reconciliationFindingLimit = 50000
type ReconciliationStore interface {
ReconciliationSummary(context.Context, int) (ReconciliationSummary, error)
ReconciliationIssues(context.Context, ReconciliationQuery) (Page[ReconciliationIssue], error)
ReconciliationIssue(context.Context, string) (ReconciliationIssue, error)
UpdateReconciliationIssue(context.Context, string, ReconciliationActionRequest) (ReconciliationIssue, error)
EvaluateReconciliation(context.Context) (ReconciliationEvaluationResult, error)
}
type reconciliationFinding struct {
SubjectKey string
RuleCode string
Category string
Severity string
VIN string
Plate string
ProtocolA string
ProtocolB string
Title string
Summary string
Evidence map[string]any
}
type reconciliationExisting struct {
ID string
Status string
Version int
}
func (s *ProductionStore) ensureReconciliationSchema(ctx context.Context) error {
s.reconciliationSchemaOnce.Do(func() {
var count int
if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM information_schema.tables
WHERE table_schema=DATABASE() AND table_name IN ('vehicle_reconciliation_run','vehicle_reconciliation_issue','vehicle_reconciliation_action')`).Scan(&count); err != nil {
s.reconciliationSchemaErr = err
return
}
if count != 3 {
s.reconciliationSchemaErr = fmt.Errorf("reconciliation schema unavailable; apply deploy/migrations/017_reconciliation_center.sql")
}
})
return s.reconciliationSchemaErr
}
func (s *ProductionStore) EvaluateReconciliation(ctx context.Context) (ReconciliationEvaluationResult, error) {
if err := s.ensureReconciliationSchema(ctx); err != nil {
return ReconciliationEvaluationResult{}, err
}
startedAt := time.Now()
runID, err := newAlertID("reconciliation-run")
if err != nil {
return ReconciliationEvaluationResult{}, err
}
findings, err := s.collectReconciliationFindings(ctx)
if err != nil {
return ReconciliationEvaluationResult{}, err
}
if len(findings) > reconciliationFindingLimit {
return ReconciliationEvaluationResult{}, fmt.Errorf("reconciliation findings exceed limit: %d", len(findings))
}
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
if err != nil {
return ReconciliationEvaluationResult{}, err
}
defer tx.Rollback()
existing, err := loadReconciliationExisting(ctx, tx)
if err != nil {
return ReconciliationEvaluationResult{}, err
}
now := time.Now()
seen := make(map[string]bool, len(findings))
result := ReconciliationEvaluationResult{RunID: runID, RuleCounts: map[string]int{}, AsOf: now.Format(time.RFC3339)}
for _, finding := range findings {
fingerprint := reconciliationFingerprint(finding)
if seen[fingerprint] {
continue
}
seen[fingerprint] = true
result.Detected++
result.RuleCounts[finding.RuleCode]++
evidence, err := json.Marshal(finding.Evidence)
if err != nil {
return ReconciliationEvaluationResult{}, err
}
current, exists := existing[fingerprint]
if !exists {
id, err := newAlertID("reconciliation")
if err != nil {
return ReconciliationEvaluationResult{}, err
}
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_issue(
id,fingerprint,rule_code,category,severity,status,vin,plate,protocol_a,protocol_b,title,summary,evidence_json,
first_seen_at,last_seen_at,occurrence_count,version
) VALUES(?,?,?,?,?,'pending',?,?,?,?,?,?,?,?,?,1,1)`,
id, fingerprint, finding.RuleCode, finding.Category, finding.Severity, finding.VIN, finding.Plate,
finding.ProtocolA, finding.ProtocolB, finding.Title, finding.Summary, string(evidence), now, now); err != nil {
return ReconciliationEvaluationResult{}, err
}
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_action(
issue_id,action,from_status,to_status,actor,note
) VALUES(?,'detect','','pending','reconciliation-evaluator','规则首次发现差异')`, id); err != nil {
return ReconciliationEvaluationResult{}, err
}
result.New++
continue
}
nextStatus := current.Status
if current.Status == "recovered" || current.Status == "fixed" {
nextStatus = "pending"
}
if _, err := tx.ExecContext(ctx, `UPDATE vehicle_reconciliation_issue SET
rule_code=?,category=?,severity=?,status=?,vin=?,plate=?,protocol_a=?,protocol_b=?,title=?,summary=?,evidence_json=?,
last_seen_at=?,occurrence_count=occurrence_count+1,recovered_at=NULL,version=version+1
WHERE id=?`, finding.RuleCode, finding.Category, finding.Severity, nextStatus, finding.VIN, finding.Plate,
finding.ProtocolA, finding.ProtocolB, finding.Title, finding.Summary, string(evidence), now, current.ID); err != nil {
return ReconciliationEvaluationResult{}, err
}
if nextStatus != current.Status {
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_action(
issue_id,action,from_status,to_status,actor,note
) VALUES(?,'reopen',?,?,'reconciliation-evaluator','已恢复或已修复的差异再次出现')`, current.ID, current.Status, nextStatus); err != nil {
return ReconciliationEvaluationResult{}, err
}
}
}
for fingerprint, current := range existing {
if seen[fingerprint] || !reconciliationAutoRecoverable(current.Status) {
continue
}
if _, err := tx.ExecContext(ctx, `UPDATE vehicle_reconciliation_issue SET
status='recovered',recovered_at=?,version=version+1 WHERE id=?`, now, current.ID); err != nil {
return ReconciliationEvaluationResult{}, err
}
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_action(
issue_id,action,from_status,to_status,actor,note
) VALUES(?,'recover',?,'recovered','reconciliation-evaluator','本轮检测未再发现该差异')`, current.ID, current.Status); err != nil {
return ReconciliationEvaluationResult{}, err
}
result.Recovered++
}
if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM vehicle_reconciliation_issue
WHERE status IN ('pending','confirmed_source_a','confirmed_source_b')`).Scan(&result.Active); err != nil {
return ReconciliationEvaluationResult{}, err
}
ruleCountsJSON, _ := json.Marshal(result.RuleCounts)
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_run(
run_id,status,detected_count,new_count,active_count,recovered_count,rule_counts_json,started_at,finished_at
) VALUES(?,'completed',?,?,?,?,?,?,?)`, runID, result.Detected, result.New, result.Active, result.Recovered, string(ruleCountsJSON), startedAt, now); err != nil {
return ReconciliationEvaluationResult{}, err
}
if err := tx.Commit(); err != nil {
return ReconciliationEvaluationResult{}, err
}
return result, nil
}
func reconciliationAutoRecoverable(status string) bool {
return status == "pending" || status == "confirmed_source_a" || status == "confirmed_source_b"
}
func reconciliationFingerprint(finding reconciliationFinding) string {
key := strings.Join([]string{
strings.TrimSpace(finding.RuleCode), strings.TrimSpace(finding.SubjectKey),
strings.ToUpper(strings.TrimSpace(finding.VIN)), strings.TrimSpace(finding.ProtocolA), strings.TrimSpace(finding.ProtocolB),
}, "|")
sum := sha256.Sum256([]byte(key))
return hex.EncodeToString(sum[:])
}
func loadReconciliationExisting(ctx context.Context, tx *sql.Tx) (map[string]reconciliationExisting, error) {
rows, err := tx.QueryContext(ctx, `SELECT fingerprint,id,status,version FROM vehicle_reconciliation_issue FOR UPDATE`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]reconciliationExisting{}
for rows.Next() {
var fingerprint string
var item reconciliationExisting
if err := rows.Scan(&fingerprint, &item.ID, &item.Status, &item.Version); err != nil {
return nil, err
}
out[fingerprint] = item
}
return out, rows.Err()
}
func (s *ProductionStore) collectReconciliationFindings(ctx context.Context) ([]reconciliationFinding, error) {
queries := []struct {
name string
text string
}{
{name: "identity", text: reconciliationIdentitySQL},
{name: "coverage", text: reconciliationCoverageSQL},
{name: "position", text: reconciliationPositionSQL},
{name: "mileage", text: reconciliationMileageSQL},
{name: "fleet_count", text: reconciliationFleetCountSQL},
{name: "business_scope", text: reconciliationBusinessScopeSQL},
}
out := make([]reconciliationFinding, 0, 1024)
for _, query := range queries {
rows, err := s.db.QueryContext(ctx, query.text)
if err != nil {
return nil, fmt.Errorf("reconciliation query %s: %w", query.name, err)
}
for rows.Next() {
var item reconciliationFinding
var evidenceJSON string
if err := rows.Scan(&item.SubjectKey, &item.RuleCode, &item.Category, &item.Severity, &item.VIN, &item.Plate,
&item.ProtocolA, &item.ProtocolB, &item.Title, &item.Summary, &evidenceJSON); err != nil {
rows.Close()
return nil, fmt.Errorf("scan reconciliation query %s: %w", query.name, err)
}
item.Evidence = map[string]any{}
if err := json.Unmarshal([]byte(evidenceJSON), &item.Evidence); err != nil {
rows.Close()
return nil, fmt.Errorf("decode reconciliation evidence for %s: %w", item.RuleCode, err)
}
out = append(out, item)
if len(out) > reconciliationFindingLimit {
rows.Close()
return nil, fmt.Errorf("reconciliation findings exceed limit")
}
}
if err := rows.Close(); err != nil {
return nil, fmt.Errorf("close reconciliation query %s: %w", query.name, err)
}
}
return out, nil
}
const reconciliationIdentitySQL = `
SELECT CONCAT('plate:',duplicate.plate),'DUPLICATE_PLATE','identity','major','',duplicate.plate,'','','同一车牌关联多个 VIN',
CONCAT('车牌 ',duplicate.plate,' 当前关联 ',duplicate.vin_count,' 个 VIN需核对权威身份绑定'),
JSON_OBJECT('plate',duplicate.plate,'vins',duplicate.vins,'vinCount',duplicate.vin_count)
FROM (
SELECT TRIM(plate) plate,GROUP_CONCAT(DISTINCT vin ORDER BY vin) vins,COUNT(DISTINCT vin) vin_count
FROM vehicle_identity_binding
WHERE plate IS NOT NULL AND TRIM(plate)<>''
GROUP BY TRIM(plate) HAVING COUNT(DISTINCT vin)>1
) duplicate
UNION ALL
SELECT CONCAT('phone:',duplicate.identifier),'DUPLICATE_PHONE','identity','major','',duplicate.plate,'JT808','','同一终端关联多个 VIN',
CONCAT('JT808 终端标识当前关联 ',duplicate.vin_count,' 个 VIN需核对设备换绑或误挂'),
JSON_OBJECT('identifierHash',SHA2(duplicate.identifier,256),'vins',duplicate.vins,'vinCount',duplicate.vin_count)
FROM (
SELECT TRIM(phone) identifier,MAX(COALESCE(plate,'')) plate,
GROUP_CONCAT(DISTINCT vin ORDER BY vin) vins,COUNT(DISTINCT vin) vin_count
FROM vehicle_identity_binding
WHERE phone IS NOT NULL AND TRIM(phone)<>''
GROUP BY TRIM(phone) HAVING COUNT(DISTINCT vin)>1
) duplicate`
const reconciliationCoverageSQL = `
SELECT CONCAT('unbound:',s.vin),'UNBOUND_SOURCE','source','major',s.vin,MAX(COALESCE(s.plate,'')),
GROUP_CONCAT(DISTINCT s.protocol ORDER BY s.protocol),'','实时来源未绑定主车辆',
CONCAT('来源快照存在 VIN但 vehicle_identity_binding 无对应主车辆:',s.vin),
JSON_OBJECT('vin',s.vin,'protocols',GROUP_CONCAT(DISTINCT s.protocol ORDER BY s.protocol),'latestSeen',DATE_FORMAT(MAX(s.updated_at),'%Y-%m-%d %H:%i:%s'))
FROM vehicle_realtime_snapshot s
LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY s.vin
WHERE s.vin IS NOT NULL AND s.vin<>'' AND b.vin IS NULL
GROUP BY s.vin
UNION ALL
SELECT CONCAT('missing:',b.vin),'SOURCE_MISSING','source','minor',b.vin,COALESCE(b.plate,''),
'','','主车辆暂无真实数据来源','车辆已建档,但当前没有 GB32960、JT808 或 YUTONG_MQTT 来源证据',
JSON_OBJECT('vin',b.vin,'plate',COALESCE(b.plate,''),'bindingSource','vehicle_identity_binding')
FROM vehicle_identity_binding b
LEFT JOIN vehicle_realtime_snapshot s ON BINARY s.vin=BINARY b.vin
WHERE b.vin IS NOT NULL AND b.vin<>''
GROUP BY b.vin,b.plate HAVING COUNT(s.protocol)=0`
const reconciliationPositionSQL = `
SELECT CONCAT('position:',a.vin,':',LEAST(a.protocol,b.protocol),':',GREATEST(a.protocol,b.protocol),':',SHA2(CONCAT(LEAST(a.source_key,b.source_key),'|',GREATEST(a.source_key,b.source_key)),256)),
'POSITION_DRIFT','location',
CASE WHEN ST_Distance_Sphere(POINT(a.longitude,a.latitude),POINT(b.longitude,b.latitude))>=5000 THEN 'critical' ELSE 'major' END,
a.vin,COALESCE(binding.plate,''),a.protocol,b.protocol,'多来源实时位置漂移',
CONCAT('两个有效来源位置相差 ',ROUND(ST_Distance_Sphere(POINT(a.longitude,a.latitude),POINT(b.longitude,b.latitude))), ' 米,系统不自动判定哪一来源正确'),
JSON_OBJECT(
'sourceA',JSON_OBJECT('protocol',a.protocol,'sourceHash',SHA2(a.source_key,256),'longitude',a.longitude,'latitude',a.latitude,'receivedAt',DATE_FORMAT(a.received_at,'%Y-%m-%d %H:%i:%s')),
'sourceB',JSON_OBJECT('protocol',b.protocol,'sourceHash',SHA2(b.source_key,256),'longitude',b.longitude,'latitude',b.latitude,'receivedAt',DATE_FORMAT(b.received_at,'%Y-%m-%d %H:%i:%s')),
'distanceM',ROUND(ST_Distance_Sphere(POINT(a.longitude,a.latitude),POINT(b.longitude,b.latitude)))
)
FROM vehicle_realtime_location_source a
JOIN vehicle_realtime_location_source b ON BINARY b.vin=BINARY a.vin
AND CONCAT(b.protocol,'|',b.source_key)>CONCAT(a.protocol,'|',a.source_key)
LEFT JOIN vehicle_identity_binding binding ON BINARY binding.vin=BINARY a.vin
WHERE a.received_at>=DATE_SUB(NOW(),INTERVAL 5 MINUTE)
AND b.received_at>=DATE_SUB(NOW(),INTERVAL 5 MINUTE)
AND a.longitude BETWEEN -180 AND 180 AND b.longitude BETWEEN -180 AND 180
AND a.latitude BETWEEN -90 AND 90 AND b.latitude BETWEEN -90 AND 90
AND NOT (a.longitude=0 AND a.latitude=0) AND NOT (b.longitude=0 AND b.latitude=0)
AND ST_Distance_Sphere(POINT(a.longitude,a.latitude),POINT(b.longitude,b.latitude))>=1000`
const reconciliationMileageSQL = `
SELECT CONCAT('mileage:',m.vin,':',m.protocol,':',DATE_FORMAT(m.stat_date,'%Y-%m-%d')),
CASE WHEN m.daily_mileage_km<0 OR m.latest_total_mileage_km<m.daily_mileage_km THEN 'MILEAGE_REVERSE' ELSE 'MILEAGE_JUMP' END,
'mileage',
CASE WHEN m.daily_mileage_km<0 OR m.latest_total_mileage_km<m.daily_mileage_km THEN 'critical' ELSE 'major' END,
m.vin,COALESCE(b.plate,''),m.protocol,'',
CASE WHEN m.daily_mileage_km<0 OR m.latest_total_mileage_km<m.daily_mileage_km THEN '里程倒退或起始值无效' ELSE '单日里程异常跳变' END,
CONCAT(DATE_FORMAT(m.stat_date,'%Y-%m-%d'),' ',m.protocol,' 日里程 ',ROUND(m.daily_mileage_km,3),' km需核对原始总里程与单位'),
JSON_OBJECT('date',DATE_FORMAT(m.stat_date,'%Y-%m-%d'),'protocol',m.protocol,'dailyMileageKm',m.daily_mileage_km,
'latestTotalMileageKm',m.latest_total_mileage_km,'derivedStartMileageKm',m.latest_total_mileage_km-m.daily_mileage_km)
FROM vehicle_daily_mileage m
LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY m.vin
WHERE m.stat_date>=DATE_SUB(CURDATE(),INTERVAL 7 DAY)
AND (m.daily_mileage_km<0 OR m.daily_mileage_km>2000 OR m.latest_total_mileage_km<m.daily_mileage_km)
UNION ALL
SELECT CONCAT('mileage-spread:',m.vin,':',DATE_FORMAT(m.stat_date,'%Y-%m-%d')),
'MILEAGE_SOURCE_DIVERGENCE','mileage','major',m.vin,COALESCE(MAX(b.plate),''),
MIN(m.protocol),MAX(m.protocol),'多来源日里程差异偏大',
CONCAT(DATE_FORMAT(m.stat_date,'%Y-%m-%d'),' 多来源日里程最大差 ',ROUND(MAX(m.daily_mileage_km)-MIN(m.daily_mileage_km),3),' km'),
JSON_OBJECT('date',DATE_FORMAT(m.stat_date,'%Y-%m-%d'),'minimumKm',MIN(m.daily_mileage_km),'maximumKm',MAX(m.daily_mileage_km),
'differenceKm',MAX(m.daily_mileage_km)-MIN(m.daily_mileage_km),'protocols',GROUP_CONCAT(DISTINCT m.protocol ORDER BY m.protocol))
FROM vehicle_daily_mileage m
LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY m.vin
WHERE m.stat_date>=DATE_SUB(CURDATE(),INTERVAL 7 DAY)
GROUP BY m.vin,m.stat_date
HAVING COUNT(DISTINCT m.protocol)>1 AND MAX(m.daily_mileage_km)-MIN(m.daily_mileage_km)>20`
const reconciliationFleetCountSQL = `
SELECT 'fleet-count','FLEET_COUNT_MISMATCH','fleet','major','','','','','主车辆与来源车辆总数不一致',
CONCAT('主车辆 ',x.bound_count,' 辆,来源并集 ',x.source_union_count,' 辆,来源未绑定 ',x.unbound_count,' 辆'),
JSON_OBJECT('boundVehicles',x.bound_count,'sourceUnionVehicles',x.source_union_count,'unboundSourceVehicles',x.unbound_count)
FROM (
SELECT
(SELECT COUNT(DISTINCT vin) FROM vehicle_identity_binding WHERE vin IS NOT NULL AND vin<>'') bound_count,
(SELECT COUNT(*) FROM (
SELECT CAST(vin AS BINARY) vin FROM vehicle_identity_binding WHERE vin IS NOT NULL AND vin<>''
UNION SELECT CAST(vin AS BINARY) vin FROM vehicle_realtime_snapshot WHERE vin IS NOT NULL AND vin<>''
) u) source_union_count,
(SELECT COUNT(DISTINCT s.vin) FROM vehicle_realtime_snapshot s LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY s.vin
WHERE s.vin IS NOT NULL AND s.vin<>'' AND b.vin IS NULL) unbound_count
) x
WHERE x.bound_count<>x.source_union_count OR x.unbound_count>0`
const reconciliationBusinessScopeSQL = `
SELECT
CONVERT(CONCAT('business-unbound:',s.customer_id,':',s.vin) USING utf8mb4) COLLATE utf8mb4_unicode_ci,
'BUSINESS_SCOPE_UNBOUND','business','critical',
CONVERT(s.vin USING utf8mb4) COLLATE utf8mb4_unicode_ci,
CONVERT(s.plate_number USING utf8mb4) COLLATE utf8mb4_unicode_ci,
'ONEOS','','OneOS 业务范围车辆未绑定主车辆',
CONVERT(CONCAT('客户 ',s.customer_id,' 的业务范围包含 VIN ',s.vin,',但中台主车辆不存在') USING utf8mb4) COLLATE utf8mb4_unicode_ci,
JSON_OBJECT('scopeVersion',s.source_version,'customerId',CAST(s.customer_id AS CHAR),'customerName',s.customer_name,
'contractCode',s.contract_code,'projectName',s.project_name,'departmentName',s.department_name,'responsibleUserName',s.responsible_user_name)
FROM business_scope_state st
JOIN business_customer_vehicle_scope s ON BINARY s.source_version=BINARY st.active_version
LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY s.vin
WHERE st.id=1 AND st.active_version IS NOT NULL AND b.vin IS NULL
UNION ALL
SELECT
CONVERT(CONCAT('auth-business:',u.id,':',g.vin) USING utf8mb4) COLLATE utf8mb4_unicode_ci,
'AUTH_SCOPE_BUSINESS_MISMATCH','business','major',
CONVERT(g.vin USING utf8mb4) COLLATE utf8mb4_unicode_ci,
CONVERT(COALESCE(b.plate,'') USING utf8mb4) COLLATE utf8mb4_unicode_ci,
'PLATFORM_AUTH','ONEOS','人工授权与 OneOS 业务范围不一致',
CONVERT(CONCAT('客户账号 ',u.username,' 当前授权 VIN 不在同客户的 OneOS 活跃范围内') USING utf8mb4) COLLATE utf8mb4_unicode_ci,
JSON_OBJECT('scopeVersion',st.active_version,'userId',CAST(u.id AS CHAR),'username',u.username,'customerRef',u.customer_ref,
'grantValidFrom',DATE_FORMAT(COALESCE(g.valid_from,g.granted_at),'%Y-%m-%d %H:%i:%s'))
FROM business_scope_state st
JOIN platform_user u ON u.user_type='customer' AND u.status='enabled' AND u.customer_ref<>''
JOIN platform_user_vehicle g ON g.user_id=u.id AND (g.valid_to IS NULL OR g.valid_to>NOW())
LEFT JOIN business_customer_vehicle_scope s ON BINARY s.source_version=BINARY st.active_version
AND BINARY CAST(s.customer_id AS CHAR)=BINARY u.customer_ref AND BINARY s.vin=BINARY g.vin
LEFT JOIN vehicle_identity_binding b ON BINARY b.vin=BINARY g.vin
WHERE st.id=1 AND st.active_version IS NOT NULL AND s.vin IS NULL`
func buildReconciliationWhere(query ReconciliationQuery) (string, []any) {
where := []string{"1=1"}
args := []any{}
if value := strings.TrimSpace(query.Keyword); value != "" {
like := "%" + value + "%"
where = append(where, "(i.vin LIKE ? OR i.plate LIKE ? OR i.title LIKE ? OR i.summary LIKE ?)")
args = append(args, like, like, like, like)
}
for column, value := range map[string]string{
"i.rule_code": query.RuleCode, "i.category": query.Category, "i.severity": query.Severity,
} {
if value = strings.TrimSpace(value); value != "" && value != "all" {
where = append(where, column+"=?")
args = append(args, value)
}
}
if status := strings.TrimSpace(query.Status); status == "active" {
where = append(where, "i.status IN ('pending','confirmed_source_a','confirmed_source_b')")
} else if status != "" && status != "all" {
where = append(where, "i.status=?")
args = append(args, status)
}
return strings.Join(where, " AND "), args
}
const reconciliationSelect = `SELECT i.id,i.rule_code,i.category,i.severity,i.status,i.vin,i.plate,i.protocol_a,i.protocol_b,
i.title,i.summary,CAST(i.evidence_json AS CHAR),DATE_FORMAT(i.first_seen_at,'%Y-%m-%d %H:%i:%s'),
DATE_FORMAT(i.last_seen_at,'%Y-%m-%d %H:%i:%s'),i.occurrence_count,
COALESCE(DATE_FORMAT(i.recovered_at,'%Y-%m-%d %H:%i:%s'),''),i.resolution_note,i.resolved_by,i.version
FROM vehicle_reconciliation_issue i `
func scanReconciliationIssue(scanner interface{ Scan(...any) error }) (ReconciliationIssue, error) {
var item ReconciliationIssue
var evidence string
err := scanner.Scan(&item.ID, &item.RuleCode, &item.Category, &item.Severity, &item.Status, &item.VIN, &item.Plate,
&item.ProtocolA, &item.ProtocolB, &item.Title, &item.Summary, &evidence, &item.FirstSeenAt, &item.LastSeenAt,
&item.OccurrenceCount, &item.RecoveredAt, &item.ResolutionNote, &item.ResolvedBy, &item.Version)
if err == nil {
item.Evidence = map[string]any{}
err = json.Unmarshal([]byte(evidence), &item.Evidence)
}
return item, err
}
func (s *ProductionStore) ReconciliationIssues(ctx context.Context, query ReconciliationQuery) (Page[ReconciliationIssue], error) {
if err := s.ensureReconciliationSchema(ctx); err != nil {
return Page[ReconciliationIssue]{}, err
}
if query.Limit <= 0 || query.Limit > 200 {
query.Limit = 50
}
if query.Offset < 0 {
query.Offset = 0
}
where, args := buildReconciliationWhere(query)
var total int
if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM vehicle_reconciliation_issue i WHERE `+where, args...).Scan(&total); err != nil {
return Page[ReconciliationIssue]{}, err
}
listArgs := append(append([]any(nil), args...), query.Limit, query.Offset)
rows, err := s.db.QueryContext(ctx, reconciliationSelect+`WHERE `+where+`
ORDER BY FIELD(i.severity,'critical','major','minor'),i.last_seen_at DESC,i.id DESC LIMIT ? OFFSET ?`, listArgs...)
if err != nil {
return Page[ReconciliationIssue]{}, err
}
defer rows.Close()
items := make([]ReconciliationIssue, 0, query.Limit)
for rows.Next() {
item, err := scanReconciliationIssue(rows)
if err != nil {
return Page[ReconciliationIssue]{}, err
}
items = append(items, item)
}
return Page[ReconciliationIssue]{Items: items, Total: total, Limit: query.Limit, Offset: query.Offset}, rows.Err()
}
func (s *ProductionStore) ReconciliationIssue(ctx context.Context, id string) (ReconciliationIssue, error) {
if err := s.ensureReconciliationSchema(ctx); err != nil {
return ReconciliationIssue{}, err
}
item, err := scanReconciliationIssue(s.db.QueryRowContext(ctx, reconciliationSelect+`WHERE i.id=?`, id))
if err == sql.ErrNoRows {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOT_FOUND", Message: "差异记录不存在"}
}
if err != nil {
return ReconciliationIssue{}, err
}
rows, err := s.db.QueryContext(ctx, `SELECT id,action,from_status,to_status,actor,note,
DATE_FORMAT(created_at,'%Y-%m-%d %H:%i:%s')
FROM vehicle_reconciliation_action WHERE issue_id=? ORDER BY created_at,id`, id)
if err != nil {
return ReconciliationIssue{}, err
}
defer rows.Close()
item.Actions = []ReconciliationAction{}
for rows.Next() {
var action ReconciliationAction
if err := rows.Scan(&action.ID, &action.Action, &action.FromStatus, &action.ToStatus, &action.Actor, &action.Note, &action.CreatedAt); err != nil {
return ReconciliationIssue{}, err
}
item.Actions = append(item.Actions, action)
}
return item, rows.Err()
}
func (s *ProductionStore) UpdateReconciliationIssue(ctx context.Context, id string, request ReconciliationActionRequest) (ReconciliationIssue, error) {
if err := s.ensureReconciliationSchema(ctx); err != nil {
return ReconciliationIssue{}, err
}
allowed := map[string]bool{
"pending": true, "confirmed_source_a": true, "confirmed_source_b": true,
"no_action": true, "fixed": true,
}
request.Status = strings.TrimSpace(request.Status)
request.Note = strings.TrimSpace(request.Note)
if !allowed[request.Status] {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_STATUS_INVALID", Message: "差异处置状态无效"}
}
if request.Status != "pending" && request.Note == "" {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOTE_REQUIRED", Message: "确认结论必须填写说明"}
}
if len([]rune(request.Note)) > 500 {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOTE_TOO_LONG", Message: "处置说明不能超过 500 个字符"}
}
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{})
if err != nil {
return ReconciliationIssue{}, err
}
defer tx.Rollback()
var currentStatus string
var currentVersion int
if err := tx.QueryRowContext(ctx, `SELECT status,version FROM vehicle_reconciliation_issue WHERE id=? FOR UPDATE`, id).Scan(&currentStatus, &currentVersion); err != nil {
if err == sql.ErrNoRows {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_NOT_FOUND", Message: "差异记录不存在"}
}
return ReconciliationIssue{}, err
}
if currentVersion != request.Version {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_VERSION_CONFLICT", Message: "差异记录已更新,请刷新后重试"}
}
result, err := tx.ExecContext(ctx, `UPDATE vehicle_reconciliation_issue SET
status=?,resolution_note=?,resolved_by=?,
recovered_at=CASE WHEN ? IN ('fixed','no_action') THEN NOW(3) ELSE NULL END,
version=version+1
WHERE id=? AND version=?`, request.Status, request.Note, request.Actor, request.Status, id, request.Version)
if err != nil {
return ReconciliationIssue{}, err
}
affected, _ := result.RowsAffected()
if affected != 1 {
return ReconciliationIssue{}, clientError{Code: "RECONCILIATION_VERSION_CONFLICT", Message: "差异记录更新冲突"}
}
if _, err := tx.ExecContext(ctx, `INSERT INTO vehicle_reconciliation_action(
issue_id,action,from_status,to_status,actor,note
) VALUES(?,'review',?,?,?,?)`, id, currentStatus, request.Status, request.Actor, request.Note); err != nil {
return ReconciliationIssue{}, err
}
if err := tx.Commit(); err != nil {
return ReconciliationIssue{}, err
}
return s.ReconciliationIssue(ctx, id)
}
const reconciliationTrendSQL = `SELECT DATE_FORMAT(MIN(finished_at),'%Y-%m-%d'),
CAST(SUBSTRING_INDEX(GROUP_CONCAT(detected_count ORDER BY finished_at DESC),',',1) AS UNSIGNED),
SUM(new_count),
CAST(SUBSTRING_INDEX(GROUP_CONCAT(active_count ORDER BY finished_at DESC),',',1) AS UNSIGNED),
SUM(recovered_count)
FROM vehicle_reconciliation_run
WHERE status='completed' AND finished_at>=DATE_SUB(CURDATE(),INTERVAL ? DAY)
GROUP BY DATE(finished_at) ORDER BY DATE(finished_at)`
func (s *ProductionStore) ReconciliationSummary(ctx context.Context, days int) (ReconciliationSummary, error) {
if err := s.ensureReconciliationSchema(ctx); err != nil {
return ReconciliationSummary{}, err
}
if days <= 0 || days > 90 {
days = 30
}
var result ReconciliationSummary
if err := s.db.QueryRowContext(ctx, `SELECT
COALESCE(SUM(status IN ('pending','confirmed_source_a','confirmed_source_b')),0),
COALESCE(SUM(status='pending'),0),
COALESCE(SUM(status IN ('confirmed_source_a','confirmed_source_b')),0),
COALESCE(SUM(status='recovered'),0),
COALESCE(SUM(status IN ('pending','confirmed_source_a','confirmed_source_b') AND first_seen_at<DATE_SUB(NOW(),INTERVAL 24 HOUR)),0)
FROM vehicle_reconciliation_issue`).Scan(&result.Active, &result.Pending, &result.Confirmed, &result.Recovered, &result.OverSLA); err != nil {
return ReconciliationSummary{}, err
}
var err error
if result.ByRule, err = s.reconciliationBuckets(ctx, "rule_code"); err != nil {
return ReconciliationSummary{}, err
}
if result.BySeverity, err = s.reconciliationBuckets(ctx, "severity"); err != nil {
return ReconciliationSummary{}, err
}
rows, err := s.db.QueryContext(ctx, reconciliationTrendSQL, days-1)
if err != nil {
return ReconciliationSummary{}, err
}
defer rows.Close()
result.Trend = []ReconciliationTrendPoint{}
for rows.Next() {
var item ReconciliationTrendPoint
if err := rows.Scan(&item.Date, &item.Detected, &item.New, &item.Active, &item.Recovered); err != nil {
return ReconciliationSummary{}, err
}
result.Trend = append(result.Trend, item)
}
_ = s.db.QueryRowContext(ctx, `SELECT COALESCE(DATE_FORMAT(MAX(finished_at),'%Y-%m-%dT%H:%i:%s.%fZ'),'')
FROM vehicle_reconciliation_run WHERE status='completed'`).Scan(&result.LastRunAt)
result.AsOf = time.Now().Format(time.RFC3339)
return result, rows.Err()
}
func (s *ProductionStore) reconciliationBuckets(ctx context.Context, column string) ([]ReconciliationBucket, error) {
if column != "rule_code" && column != "severity" {
return nil, fmt.Errorf("unsupported reconciliation bucket")
}
rows, err := s.db.QueryContext(ctx, `SELECT `+column+`,COUNT(*) FROM vehicle_reconciliation_issue
WHERE status IN ('pending','confirmed_source_a','confirmed_source_b')
GROUP BY `+column+` ORDER BY COUNT(*) DESC,`+column)
if err != nil {
return nil, err
}
defer rows.Close()
out := []ReconciliationBucket{}
for rows.Next() {
var item ReconciliationBucket
if err := rows.Scan(&item.Name, &item.Count); err != nil {
return nil, err
}
out = append(out, item)
}
return out, rows.Err()
}

View File

@@ -0,0 +1,60 @@
package platform
import (
"strings"
"testing"
)
func TestReconciliationSQLUsesCollationSafeCrossTableKeys(t *testing.T) {
checks := map[string][]string{
"coverage": {
"BINARY b.vin=BINARY s.vin",
"BINARY s.vin=BINARY b.vin",
},
"position": {
"BINARY b.vin=BINARY a.vin",
"BINARY binding.vin=BINARY a.vin",
},
"mileage": {
"BINARY b.vin=BINARY m.vin",
},
"fleet count": {
"SELECT CAST(vin AS BINARY) vin FROM vehicle_identity_binding",
"UNION SELECT CAST(vin AS BINARY) vin FROM vehicle_realtime_snapshot",
"BINARY b.vin=BINARY s.vin",
},
"business scope": {
"BINARY s.source_version=BINARY st.active_version",
"BINARY b.vin=BINARY s.vin",
"BINARY CAST(s.customer_id AS CHAR)=BINARY u.customer_ref",
"BINARY s.vin=BINARY g.vin",
"BINARY b.vin=BINARY g.vin",
"CONVERT(s.vin USING utf8mb4) COLLATE utf8mb4_unicode_ci",
"CONVERT(g.vin USING utf8mb4) COLLATE utf8mb4_unicode_ci",
"CONVERT(COALESCE(b.plate,'') USING utf8mb4) COLLATE utf8mb4_unicode_ci",
},
}
sqlByName := map[string]string{
"coverage": reconciliationCoverageSQL,
"position": reconciliationPositionSQL,
"mileage": reconciliationMileageSQL,
"fleet count": reconciliationFleetCountSQL,
"business scope": reconciliationBusinessScopeSQL,
}
for name, fragments := range checks {
for _, fragment := range fragments {
if !strings.Contains(sqlByName[name], fragment) {
t.Errorf("%s reconciliation SQL is missing %q", name, fragment)
}
}
}
}
func TestReconciliationTrendSQLSupportsOnlyFullGroupBy(t *testing.T) {
if !strings.Contains(reconciliationTrendSQL, "DATE_FORMAT(MIN(finished_at),'%Y-%m-%d')") {
t.Fatalf("trend date must be derived from an aggregate: %s", reconciliationTrendSQL)
}
if strings.Contains(reconciliationTrendSQL, "SELECT DATE_FORMAT(DATE(finished_at)") {
t.Fatalf("trend query must not select a non-aggregated finished_at expression: %s", reconciliationTrendSQL)
}
}