package main import ( "context" "database/sql" "flag" "fmt" "log" "os" "sort" "strconv" "strings" "sync" "time" _ "github.com/go-sql-driver/mysql" _ "github.com/taosdata/driver-go/v3/taosWS" "lingniu/vehicle-data-platform/apps/api/internal/config" "lingniu/vehicle-data-platform/apps/api/internal/openplatform" ) func main() { date := flag.String("date", "", "last date to calculate in yyyy-MM-dd; defaults to yesterday in Asia/Shanghai") lookback := flag.Int("lookback-days", envInt("OPEN_STAT_LOOKBACK_DAYS", 2), "number of dates ending at -date") noise := flag.Float64("hydrogen-noise-kg", envFloat("OPEN_STAT_HYDROGEN_NOISE_KG", 0.05), "ignored mass jitter") maxDrop := flag.Float64("hydrogen-max-drop-kg", envFloat("OPEN_STAT_HYDROGEN_MAX_DROP_KG", 20), "maximum accepted drop between samples") vin := flag.String("vin", "", "optional 17-character VIN for a scoped rebuild") vinWorkers := flag.Int("vin-workers", envInt("OPEN_STAT_VIN_WORKERS", 4), "parallel per-VIN queries for an all-vehicle day") allDates := flag.Bool("all-dates", false, "rebuild every available completed event date; optionally scoped by -vin") dryRun := flag.Bool("dry-run", false, "calculate and print results without changing MySQL") seedStream := flag.Bool("seed-stream-state", false, "atomically seed current-day segment stream watermarks during a controlled writer handoff") flag.Parse() cfg := config.Load() if cfg.MySQLDSN == "" || cfg.TDengineDSN == "" { log.Fatal("MYSQL_DSN and TDENGINE_DSN are required") } mysqlDB, err := sql.Open("mysql", cfg.MySQLDSN) if err != nil { log.Fatal(err) } defer mysqlDB.Close() tdengineDB, err := sql.Open(cfg.TDengineDriver, cfg.TDengineDSN) if err != nil { log.Fatal(err) } defer tdengineDB.Close() location := time.FixedZone("Asia/Shanghai", 8*60*60) endDate := time.Now().In(location).AddDate(0, 0, -1) if *date != "" { endDate, err = time.ParseInLocation("2006-01-02", *date, location) if err != nil { log.Fatal("-date must use yyyy-MM-dd") } } if *lookback < 1 || *lookback > 31 { log.Fatal("-lookback-days must be between 1 and 31") } if *vinWorkers < 1 || *vinWorkers > 16 { log.Fatal("-vin-workers must be between 1 and 16") } normalizedVIN := strings.ToUpper(strings.TrimSpace(*vin)) if *seedStream && normalizedVIN != "" { log.Fatal("-seed-stream-state requires an all-vehicle rebuild") } if *seedStream && *dryRun { log.Fatal("-seed-stream-state cannot be combined with -dry-run") } tdengineDB.SetMaxOpenConns(*vinWorkers) timeout := 30 * time.Minute if *allDates { timeout = 6 * time.Hour } ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() capacities, err := openplatform.LoadHydrogenCapacities(ctx, mysqlDB) if err != nil { log.Fatalf("load hydrogen tank capacities: %v", err) } if normalizedVIN != "" { if _, ok := capacities[normalizedVIN]; !ok { log.Fatalf("no active hydrogen tank capacity for VIN %s", normalizedVIN) } } startDate := time.Date(endDate.Year(), endDate.Month(), endDate.Day(), 0, 0, 0, 0, location).AddDate(0, 0, -(*lookback - 1)) if *allDates { var first, last time.Time var found bool var rangeErr error if normalizedVIN == "" { first, last, found, rangeErr = openplatform.HydrogenObservationDateRangeForAll(ctx, tdengineDB, cfg.TDengineDatabase) } else { first, last, found, rangeErr = openplatform.HydrogenObservationDateRange(ctx, tdengineDB, cfg.TDengineDatabase, normalizedVIN) } if rangeErr != nil { log.Fatalf("load hydrogen date range: %v", rangeErr) } if !found { log.Fatal("no GB32960 history found") } first = first.In(location) last = last.In(location) startDate = time.Date(first.Year(), first.Month(), first.Day(), 0, 0, 0, 0, location) endDate = time.Date(last.Year(), last.Month(), last.Day(), 0, 0, 0, 0, location) now := time.Now().In(location) lastCompletedDate := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, location).AddDate(0, 0, -1) if endDate.After(lastCompletedDate) { endDate = lastCompletedDate } scope := "all-vehicles" if normalizedVIN != "" { scope = normalizedVIN } fmt.Printf("scope=%s date_from=%s date_to=%s mode=all-dates\n", scope, startDate.Format("2006-01-02"), endDate.Format("2006-01-02")) } for start := startDate; !start.After(endDate); start = start.AddDate(0, 0, 1) { date := start.Format("2006-01-02") energyParameters, parameterErr := openplatform.LoadHydrogenCalculationParameters(ctx, mysqlDB, start) if parameterErr != nil { log.Fatalf("load hydrogen energy parameters for %s: %v", date, parameterErr) } var stats []openplatform.HydrogenDailyStat observationCount := 0 if normalizedVIN == "" { var vins []string vins, err = openplatform.LoadHydrogenObservationVINs(ctx, tdengineDB, cfg.TDengineDatabase, start, start.AddDate(0, 0, 1), capacities) if err == nil { stats, observationCount, err = buildHydrogenDailyStatsByVIN(ctx, tdengineDB, cfg.TDengineDatabase, capacities, energyParameters, vins, start, *noise, *maxDrop, *vinWorkers) } } else { var observations []openplatform.HydrogenObservation observations, err = openplatform.LoadHydrogenObservationsForVIN(ctx, tdengineDB, cfg.TDengineDatabase, start, start.AddDate(0, 0, 1), capacities, normalizedVIN) observationCount = len(observations) stats = openplatform.BuildHydrogenDailyStatsOrderedWithParameters(observations, date, *noise, *maxDrop, energyParameters) } if err != nil { log.Fatalf("load hydrogen observations for %s: %v", date, err) } if !*dryRun { if *seedStream { err = openplatform.ReplaceHydrogenDailyStatsAndSeedStream(ctx, mysqlDB, date, stats) } else if normalizedVIN == "" { err = openplatform.ReplaceHydrogenDailyStats(ctx, mysqlDB, date, stats) } else { err = openplatform.ReplaceHydrogenDailyStatsForVIN(ctx, mysqlDB, date, normalizedVIN, stats) } if err != nil { log.Fatalf("persist hydrogen statistics for %s: %v", date, err) } } mode := "write" if *dryRun { mode = "dry-run" } if normalizedVIN != "" && len(stats) == 1 { fmt.Printf("date=%s vin=%s observations=%d consumption_kg=%.3f refuels=%d quality=%s mode=%s\n", date, normalizedVIN, observationCount, stats[0].ConsumptionKg, stats[0].RefuelCount, stats[0].QualityStatus, mode) } else { fmt.Printf("date=%s observations=%d vehicles=%d mode=%s\n", date, observationCount, len(stats), mode) } } } type hydrogenVINResult struct { VIN string ObservationCount int Stats []openplatform.HydrogenDailyStat Err error } func buildHydrogenDailyStatsByVIN( ctx context.Context, tdengineDB *sql.DB, database string, capacities map[string]float64, energyParameters map[string]openplatform.HydrogenCalculationParameters, vins []string, start time.Time, noiseKg float64, maxDropKg float64, workerCount int, ) ([]openplatform.HydrogenDailyStat, int, error) { if len(vins) == 0 { return nil, 0, nil } if workerCount > len(vins) { workerCount = len(vins) } workerCtx, cancel := context.WithCancel(ctx) defer cancel() jobs := make(chan string) results := make(chan hydrogenVINResult, len(vins)) var workers sync.WaitGroup for worker := 0; worker < workerCount; worker++ { workers.Add(1) go func() { defer workers.Done() for vin := range jobs { observations, err := openplatform.LoadHydrogenObservationsForVIN(workerCtx, tdengineDB, database, start, start.AddDate(0, 0, 1), capacities, vin) if err != nil { results <- hydrogenVINResult{VIN: vin, Err: err} continue } results <- hydrogenVINResult{ VIN: vin, ObservationCount: len(observations), Stats: openplatform.BuildHydrogenDailyStatsOrderedWithParameters(observations, start.Format("2006-01-02"), noiseKg, maxDropKg, energyParameters), } } }() } go func() { defer close(jobs) for _, vin := range vins { select { case jobs <- vin: case <-workerCtx.Done(): return } } }() go func() { workers.Wait() close(results) }() stats := make([]openplatform.HydrogenDailyStat, 0, len(vins)) observationCount := 0 var firstErr error for result := range results { if result.Err != nil { if firstErr == nil { firstErr = fmt.Errorf("VIN %s: %w", result.VIN, result.Err) cancel() } continue } observationCount += result.ObservationCount stats = append(stats, result.Stats...) } if firstErr != nil { return nil, 0, firstErr } sort.Slice(stats, func(i, j int) bool { return stats[i].VIN < stats[j].VIN }) return stats, observationCount, nil } func envInt(name string, fallback int) int { value, err := strconv.Atoi(os.Getenv(name)) if err != nil { return fallback } return value } func envFloat(name string, fallback float64) float64 { value, err := strconv.ParseFloat(os.Getenv(name), 64) if err != nil { return fallback } return value }