fix(stats): clear stale backfill candidates before projection

This commit is contained in:
lingniu
2026-07-08 15:12:27 +08:00
parent a5eebfb32d
commit 6d633aa292
3 changed files with 234 additions and 17 deletions

View File

@@ -266,8 +266,15 @@ func writeAggregates(ctx context.Context, db *sql.DB, aggregates map[string]*met
return 0, nil
}
var written int64
targets := map[string]struct{}{}
clearedTargets := map[string]struct{}{}
for _, agg := range aggregates {
target := agg.VIN + "|" + agg.Date + "|" + string(agg.Protocol)
if _, ok := clearedTargets[target]; !ok {
if err := clearBackfillTargetMileage(ctx, db, agg.VIN, agg.Date, agg.Protocol); err != nil {
return written, err
}
clearedTargets[target] = struct{}{}
}
identity := stats.SourceIdentity{
Protocol: agg.Protocol,
SourceIP: stats.NormalizeSourceIP(agg.SourceEndpoint),
@@ -308,13 +315,6 @@ func writeAggregates(ctx context.Context, db *sql.DB, aggregates map[string]*met
if err := stats.UpsertSourceMileage(ctx, db, candidate); err != nil {
return written, err
}
target := agg.VIN + "|" + agg.Date + "|" + string(agg.Protocol)
if _, ok := targets[target]; !ok {
if err := clearBackfillFinalMileage(ctx, db, agg.VIN, agg.Date, agg.Protocol); err != nil {
return written, err
}
targets[target] = struct{}{}
}
if err := stats.ProjectDailyMileage(ctx, db, agg.VIN, agg.Date, agg.Protocol); err != nil {
return written, err
}
@@ -323,11 +323,15 @@ func writeAggregates(ctx context.Context, db *sql.DB, aggregates map[string]*met
return written, nil
}
func clearBackfillFinalMileage(ctx context.Context, db *sql.DB, vin string, statDate string, protocol envelope.Protocol) error {
func clearBackfillTargetMileage(ctx context.Context, db *sql.DB, vin string, statDate string, protocol envelope.Protocol) error {
if db == nil || strings.TrimSpace(vin) == "" || strings.TrimSpace(statDate) == "" || strings.TrimSpace(string(protocol)) == "" {
return nil
}
_, err := db.ExecContext(ctx, "DELETE FROM vehicle_daily_mileage WHERE vin = ? AND stat_date = ? AND protocol = ?", vin, statDate, string(protocol))
_, err := db.ExecContext(ctx, "DELETE FROM vehicle_daily_mileage_source WHERE vin = ? AND stat_date = ? AND protocol = ?", vin, statDate, string(protocol))
if err != nil {
return err
}
_, err = db.ExecContext(ctx, "DELETE FROM vehicle_daily_mileage WHERE vin = ? AND stat_date = ? AND protocol = ?", vin, statDate, string(protocol))
return err
}

View File

@@ -67,24 +67,27 @@ func TestDailySourceLastBuildsCandidateKeysBySourceIP(t *testing.T) {
}
}
func TestClearBackfillFinalMileageClearsExactKey(t *testing.T) {
func TestClearBackfillTargetMileageClearsExactKey(t *testing.T) {
db, mock, err := sqlmock.New()
if err != nil {
t.Fatalf("sqlmock.New() error = %v", err)
}
defer db.Close()
mock.ExpectExec(`DELETE FROM vehicle_daily_mileage_source WHERE vin = \? AND stat_date = \? AND protocol = \?`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 1))
mock.ExpectExec(`DELETE FROM vehicle_daily_mileage WHERE vin = \? AND stat_date = \? AND protocol = \?`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 1))
if err := clearBackfillFinalMileage(context.Background(), db, "LA9GG64L7PBAF4001", "2026-07-08", envelope.ProtocolJT808); err != nil {
t.Fatalf("clearBackfillFinalMileage() error = %v", err)
if err := clearBackfillTargetMileage(context.Background(), db, "LA9GG64L7PBAF4001", "2026-07-08", envelope.ProtocolJT808); err != nil {
t.Fatalf("clearBackfillTargetMileage() error = %v", err)
}
if err := mock.ExpectationsWereMet(); err != nil {
t.Fatalf("sql expectations: %v", err)
}
}
func TestWriteAggregatesClearsFinalBeforeProjection(t *testing.T) {
func TestWriteAggregatesClearsTargetRowsBeforeStaleCandidateUpsert(t *testing.T) {
db, mock, err := sqlmock.New()
if err != nil {
t.Fatalf("sqlmock.New() error = %v", err)
@@ -111,6 +114,12 @@ func TestWriteAggregatesClearsFinalBeforeProjection(t *testing.T) {
},
}
mock.ExpectExec(`DELETE FROM vehicle_daily_mileage_source WHERE vin = \? AND stat_date = \? AND protocol = \?`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 1))
mock.ExpectExec(`DELETE FROM vehicle_daily_mileage WHERE vin = \? AND stat_date = \? AND protocol = \?`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 1))
mock.ExpectExec(`INSERT INTO vehicle_data_source`).
WithArgs("JT808", "115.231.168.135", "115.231.168.135:20215", sqlmock.AnyArg(), sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 1))
@@ -135,9 +144,6 @@ func TestWriteAggregatesClearsFinalBeforeProjection(t *testing.T) {
"outside_daily_range",
).
WillReturnResult(sqlmock.NewResult(0, 1))
mock.ExpectExec(`DELETE FROM vehicle_daily_mileage WHERE vin = \? AND stat_date = \? AND protocol = \?`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 1))
mock.ExpectExec(`UPDATE vehicle_daily_mileage_source`).
WithArgs("LA9GG64L7PBAF4001", "2026-07-08", "JT808").
WillReturnResult(sqlmock.NewResult(0, 0))