package openplatform import ( "bytes" "context" "database/sql" "encoding/json" "errors" "fmt" "math" "sort" "strconv" "strings" "time" ) var hydrogenMassFields = []string{ "gb32960.gd_fc_vehicle_info.hydrogen_mass_kg", "gb32960.gd_fc_vehicle.hydrogen_mass_kg", "gb32960.gd_fc_vehicle_info.gd_fc_vehicle_hydrogen_mass_kg", "gd_fc_vehicle_hydrogen_mass_kg", } const ( hydrogenFuelCellActiveCurrentA = 1.0 hydrogenRapidRecoveryWindow = time.Minute hydrogenRapidRecoveryPressureMPa = 8.0 hydrogenSegmentEndpointWindow = 5 hydrogenSegmentMaxGap = 5 * time.Minute ) func BuildHydrogenDailyStats(observations []HydrogenObservation, date string, noiseKg, maxDropKg float64) []HydrogenDailyStat { return buildHydrogenDailyStatsWithParameters(observations, date, noiseKg, maxDropKg, false, nil) } // BuildHydrogenDailyStatsOrdered avoids sorting the same high-volume day again // when LoadHydrogenObservations has already returned VIN/event-time order. func BuildHydrogenDailyStatsOrdered(observations []HydrogenObservation, date string, noiseKg, maxDropKg float64) []HydrogenDailyStat { return buildHydrogenDailyStatsWithParameters(observations, date, noiseKg, maxDropKg, true, nil) } func buildHydrogenDailyStatsWithParameters(observations []HydrogenObservation, date string, noiseKg, maxDropKg float64, alreadyOrdered bool, parameters map[string]HydrogenCalculationParameters) []HydrogenDailyStat { if noiseKg <= 0 { noiseKg = 0.05 } if maxDropKg <= noiseKg { maxDropKg = 20 } grouped := map[string][]HydrogenObservation{} for _, observation := range observations { observation.VIN = strings.ToUpper(strings.TrimSpace(observation.VIN)) observation.Source = strings.TrimSpace(observation.Source) if len(observation.VIN) != 17 || math.IsNaN(observation.MassKg) || math.IsInf(observation.MassKg, 0) || observation.MassKg < 0 || observation.MassKg > 200 { continue } grouped[observation.VIN] = append(grouped[observation.VIN], observation) } vins := make([]string, 0, len(grouped)) for vin := range grouped { vins = append(vins, vin) } sort.Strings(vins) stats := make([]HydrogenDailyStat, 0, len(vins)) for _, vin := range vins { values, primarySource := mergeHydrogenObservations(grouped[vin], alreadyOrdered) stats = append(stats, buildTrustedHydrogenDailyStat(vin, primarySource, date, values, noiseKg, maxDropKg, parameters[vin])) } return stats } func mergeHydrogenObservations(values []HydrogenObservation, alreadyOrdered bool) ([]HydrogenObservation, string) { sourceCounts := make(map[string]int) for _, value := range values { sourceCounts[value.Source]++ } primarySource := "" primaryCount := -1 for source, count := range sourceCounts { if count > primaryCount || (count == primaryCount && source < primarySource) { primarySource = source primaryCount = count } } if !alreadyOrdered { sort.SliceStable(values, func(i, j int) bool { if !values[i].ObservedAt.Equal(values[j].ObservedAt) { return values[i].ObservedAt.Before(values[j].ObservedAt) } return betterHydrogenDuplicate(values[i], values[j], sourceCounts) }) } merged := make([]HydrogenObservation, 0, len(values)) for first := 0; first < len(values); { last := first + 1 best := values[first] for last < len(values) && values[last].ObservedAt.Equal(values[first].ObservedAt) { if betterHydrogenDuplicate(values[last], best, sourceCounts) { best = values[last] } last++ } merged = append(merged, best) first = last } return merged, primarySource } func betterHydrogenDuplicate(candidate, current HydrogenObservation, sourceCounts map[string]int) bool { if candidate.FuelCellStateKnown != current.FuelCellStateKnown { return candidate.FuelCellStateKnown } if sourceCounts[candidate.Source] != sourceCounts[current.Source] { return sourceCounts[candidate.Source] > sourceCounts[current.Source] } if candidate.Source != current.Source { return candidate.Source < current.Source } return candidate.MassKg < current.MassKg } func hydrogenDailyMetadata(values []HydrogenObservation, noiseKg, maxDropKg float64) (float64, float64, int) { cycleMinimum := values[0] remainingMass := values[0].MassKg refuelCount := 0 for index := 1; index < len(values); index++ { value := values[index] previous := values[index-1] intervalEligible := hydrogenConsumptionIntervalEligible(previous, value) sampleNoise := hydrogenObservationNoise(value, noiseKg) refuelThreshold := hydrogenObservationRefuelThreshold(value, cycleMinimum) if intervalEligible || (value.FuelCellStateKnown && value.FuelCellActive) || value.MassKg-remainingMass > sampleNoise { remainingMass = value.MassKg } delta := cycleMinimum.MassKg - value.MassKg switch { case value.MassKg-cycleMinimum.MassKg > refuelThreshold: if !rapidHydrogenPressureRecovery(previous, value) { refuelCount++ } cycleMinimum = value case delta > sampleNoise && delta <= maxDropKg: cycleMinimum = value } } return remainingMass, cycleMinimum.MassKg, refuelCount } func hydrogenSegmentConsumption(values []HydrogenObservation, noiseKg, maxDropKg float64) (float64, int, int, int) { accumulator := newHydrogenSegmentAccumulator(noiseKg, maxDropKg) for _, value := range values { accumulator.Add(value) } return accumulator.Finalize() } type hydrogenSegmentAccumulator struct { noiseKg float64 maxDropKg float64 previous HydrogenObservation hasPrevious bool cycleMinimum HydrogenObservation segmentCount int segmentFirst []HydrogenObservation segmentTail []HydrogenObservation consumption float64 qualifiedSegments int eligibleIntervals int abnormalDrops int } func newHydrogenSegmentAccumulator(noiseKg, maxDropKg float64) *hydrogenSegmentAccumulator { return &hydrogenSegmentAccumulator{noiseKg: noiseKg, maxDropKg: maxDropKg} } func (accumulator *hydrogenSegmentAccumulator) Add(value HydrogenObservation) { if !accumulator.hasPrevious { accumulator.hasPrevious = true accumulator.previous = value if hydrogenObservationCanStartSegment(value) { accumulator.start(value) } return } previous := accumulator.previous accumulator.previous = value intervalEligible := hydrogenConsumptionIntervalEligible(previous, value) if intervalEligible { accumulator.eligibleIntervals++ } if accumulator.segmentCount == 0 || !intervalEligible { accumulator.flush() if hydrogenObservationCanStartSegment(value) { accumulator.start(value) } return } if value.ObservedAt.Sub(previous.ObservedAt) > hydrogenSegmentMaxGap { accumulator.flush() accumulator.start(value) return } if previous.MassKg-value.MassKg > accumulator.maxDropKg { accumulator.abnormalDrops++ accumulator.flush() accumulator.start(value) return } if value.MassKg-accumulator.cycleMinimum.MassKg > hydrogenObservationRefuelThreshold(value, accumulator.cycleMinimum) { accumulator.flush() accumulator.start(value) return } accumulator.append(value) if value.MassKg < accumulator.cycleMinimum.MassKg { accumulator.cycleMinimum = value } } func (accumulator *hydrogenSegmentAccumulator) Finalize() (float64, int, int, int) { accumulator.flush() return accumulator.consumption, accumulator.qualifiedSegments, accumulator.eligibleIntervals, accumulator.abnormalDrops } func (accumulator *hydrogenSegmentAccumulator) start(value HydrogenObservation) { accumulator.segmentCount = 0 accumulator.segmentFirst = accumulator.segmentFirst[:0] accumulator.segmentTail = accumulator.segmentTail[:0] accumulator.cycleMinimum = value accumulator.append(value) } func (accumulator *hydrogenSegmentAccumulator) append(value HydrogenObservation) { accumulator.segmentCount++ if len(accumulator.segmentFirst) < hydrogenSegmentEndpointWindow { accumulator.segmentFirst = append(accumulator.segmentFirst, value) } if len(accumulator.segmentTail) == hydrogenSegmentEndpointWindow { copy(accumulator.segmentTail, accumulator.segmentTail[1:]) accumulator.segmentTail[len(accumulator.segmentTail)-1] = value return } accumulator.segmentTail = append(accumulator.segmentTail, value) } func (accumulator *hydrogenSegmentAccumulator) flush() { if accumulator.segmentCount >= hydrogenSegmentEndpointWindow*2 { accumulator.qualifiedSegments++ accumulator.consumption += hydrogenEndpointDrop(accumulator.segmentFirst, accumulator.segmentTail, accumulator.noiseKg) } accumulator.segmentCount = 0 accumulator.segmentFirst = accumulator.segmentFirst[:0] accumulator.segmentTail = accumulator.segmentTail[:0] } func hydrogenObservationCanStartSegment(value HydrogenObservation) bool { return !value.FuelCellStateKnown || value.FuelCellActive } func hydrogenEndpointDrop(start, end []HydrogenObservation, noiseKg float64) float64 { startMass := hydrogenMedianMass(start) endMass := hydrogenMedianMass(end) segmentNoise := noiseKg for _, value := range start { segmentNoise = math.Max(segmentNoise, value.NoiseKg) } for _, value := range end { segmentNoise = math.Max(segmentNoise, value.NoiseKg) } drop := startMass - endMass if drop <= segmentNoise { return 0 } return drop } func hydrogenMedianMass(values []HydrogenObservation) float64 { masses := make([]float64, len(values)) for index, value := range values { masses[index] = value.MassKg } sort.Float64s(masses) return masses[len(masses)/2] } func hydrogenObservationNoise(value HydrogenObservation, fallback float64) float64 { return math.Max(fallback, value.NoiseKg) } func hydrogenObservationRefuelThreshold(value, cycleMinimum HydrogenObservation) float64 { threshold := math.Max(1, cycleMinimum.MassKg*0.05) if value.RefuelThresholdKg > 0 { threshold = value.RefuelThresholdKg } return threshold } func ExtractHydrogenMass(parsedJSON string) (float64, bool) { var fields map[string]any if err := json.Unmarshal([]byte(parsedJSON), &fields); err != nil { return 0, false } for _, key := range hydrogenMassFields { if value, ok := numericValue(fields[key]); ok && value >= 0 && value <= 200 { return value, true } } return 0, false } func ExtractHydrogenRateAndMileage(parsedJSON string) (float64, float64, bool) { var fields map[string]any if err := json.Unmarshal([]byte(parsedJSON), &fields); err != nil { return 0, 0, false } rate, rateOK := numericValue(fields["gb32960.fuel_cell.hydrogen_consumption_kg_per_100km"]) mileage, mileageOK := numericValue(fields["gb32960.vehicle.total_mileage_km"]) return rate, mileage, rateOK && mileageOK && rate >= 0 && rate <= 50 && mileage > 0 } func BuildHydrogenRateDailyStats(observations []HydrogenRateObservation, date string, maxDeltaKm float64) []HydrogenRateDailyStat { if maxDeltaKm <= 0 { maxDeltaKm = 10 } grouped := map[string]map[string][]HydrogenRateObservation{} for _, value := range observations { value.VIN = strings.ToUpper(strings.TrimSpace(value.VIN)) if len(value.VIN) != 17 || value.Rate < 0 || value.Rate > 50 || value.MileageKm <= 0 { continue } if grouped[value.VIN] == nil { grouped[value.VIN] = map[string][]HydrogenRateObservation{} } grouped[value.VIN][strings.TrimSpace(value.Source)] = append(grouped[value.VIN][strings.TrimSpace(value.Source)], value) } result := make([]HydrogenRateDailyStat, 0, len(grouped)) for vin, sources := range grouped { var best *HydrogenRateDailyStat for source, values := range sources { sort.SliceStable(values, func(i, j int) bool { return values[i].ObservedAt.Before(values[j].ObservedAt) }) stat := HydrogenRateDailyStat{VIN: vin, Source: source, Date: date, SampleCount: len(values), QualityStatus: "NO_DATA", QualityReason: "尚无有效行驶里程区间"} movement, abnormal := 0, 0 for i := 1; i < len(values); i++ { delta := values[i].MileageKm - values[i-1].MileageKm if delta > 0 && delta <= maxDeltaKm { stat.ConsumptionKg += delta * (values[i-1].Rate + values[i].Rate) / 200 movement++ } else if delta < 0 || delta > maxDeltaKm { abnormal++ } } stat.ConsumptionKg = round3(stat.ConsumptionKg) if abnormal > 0 { stat.QualityStatus, stat.QualityReason = "SUSPECT", "存在异常里程跳变" } else if movement > 0 { stat.QualityStatus, stat.QualityReason = "OK", "按里程区间积分百公里氢耗" } if best == nil || (stat.QualityStatus == "OK" && best.QualityStatus != "OK") || (stat.QualityStatus == best.QualityStatus && stat.SampleCount > best.SampleCount) { copy := stat best = © } } if best != nil { result = append(result, *best) } } sort.Slice(result, func(i, j int) bool { return result[i].VIN < result[j].VIN }) return result } func LoadHydrogenCapacities(ctx context.Context, db *sql.DB) (map[string]float64, error) { rows, err := db.QueryContext(ctx, `SELECT UPPER(TRIM(vin)),tank_capacity_l FROM vehicle_hydrogen_tank_capacity WHERE active=1 AND tank_capacity_l>0`) if err != nil { return nil, err } defer rows.Close() capacities := map[string]float64{} for rows.Next() { var vin string var capacity float64 if err := rows.Scan(&vin, &capacity); err != nil { return nil, err } if len(vin) == 17 && capacity > 0 && capacity <= 10000 { capacities[vin] = capacity } } return capacities, rows.Err() } func LoadHydrogenBrandVINs(ctx context.Context, db *sql.DB, brand string) ([]string, error) { brand = strings.TrimSpace(brand) if brand == "" { return nil, errors.New("hydrogen vehicle brand is required") } rows, err := db.QueryContext(ctx, `SELECT UPPER(TRIM(vin)) FROM vehicle_profile WHERE TRIM(brand_name)=? ORDER BY vin`, brand) if err != nil { return nil, err } defer rows.Close() vins := make([]string, 0) for rows.Next() { var vin string if err := rows.Scan(&vin); err != nil { return nil, err } vin = strings.ToUpper(strings.TrimSpace(vin)) if validHydrogenVIN(vin) { vins = append(vins, vin) } } return vins, rows.Err() } func LoadHydrogenCalculationParameters(ctx context.Context, db *sql.DB, date time.Time) (map[string]HydrogenCalculationParameters, error) { rows, err := db.QueryContext(ctx, `SELECT UPPER(TRIM(vin)),battery_capacity_kwh,hydrogen_energy_kwh_per_kg FROM vehicle_hydrogen_energy_parameter WHERE active=1 AND effective_from<=? AND (effective_to IS NULL OR effective_to>=?) ORDER BY vin,effective_from DESC`, date.Format("2006-01-02"), date.Format("2006-01-02")) if err != nil { return nil, err } defer rows.Close() result := map[string]HydrogenCalculationParameters{} for rows.Next() { var vin string var batteryCapacity, conversion float64 if err := rows.Scan(&vin, &batteryCapacity, &conversion); err != nil { return nil, err } if _, exists := result[vin]; exists { continue } result[vin] = HydrogenCalculationParameters{ BatteryCapacityKWh: batteryCapacity, HydrogenEnergyKWhKg: conversion, } } return result, rows.Err() } func LoadHydrogenObservations(ctx context.Context, tdengine *sql.DB, database string, start, end time.Time, capacities map[string]float64) ([]HydrogenObservation, error) { return loadHydrogenObservations(ctx, tdengine, database, start, end, capacities, "") } func LoadHydrogenObservationVINs(ctx context.Context, tdengine *sql.DB, database string, start, end time.Time, capacities map[string]float64) ([]string, error) { query, err := buildHydrogenObservationVINQuery(database, start, end) if err != nil { return nil, err } rows, err := tdengine.QueryContext(ctx, query) if err != nil { return nil, err } defer rows.Close() vins := make([]string, 0) seen := make(map[string]struct{}) for rows.Next() { var vin string if err := rows.Scan(&vin); err != nil { return nil, err } vin = strings.ToUpper(strings.TrimSpace(vin)) if !validHydrogenVIN(vin) { continue } if _, ok := capacities[vin]; !ok { continue } if _, ok := seen[vin]; ok { continue } seen[vin] = struct{}{} vins = append(vins, vin) } if err := rows.Err(); err != nil { return nil, err } sort.Strings(vins) return vins, nil } func LoadHydrogenObservationsForVIN(ctx context.Context, tdengine *sql.DB, database string, start, end time.Time, capacities map[string]float64, vin string) ([]HydrogenObservation, error) { vin = strings.ToUpper(strings.TrimSpace(vin)) if !validHydrogenVIN(vin) { return nil, fmt.Errorf("invalid VIN %q", vin) } return loadHydrogenObservations(ctx, tdengine, database, start, end, capacities, vin) } func loadHydrogenObservations(ctx context.Context, tdengine *sql.DB, database string, start, end time.Time, capacities map[string]float64, vin string) ([]HydrogenObservation, error) { database = strings.TrimSpace(database) if database == "" { database = "lingniu_vehicle_ts" } query, err := buildHydrogenObservationQuery(database, start, end, vin) if err != nil { return nil, err } rows, err := tdengine.QueryContext(ctx, query) if err != nil { return nil, err } defer rows.Close() observations := make([]HydrogenObservation, 0) for rows.Next() { var vin, parsed string var source, eventID sql.NullString var unixMS int64 if err := rows.Scan(&vin, &source, &eventID, &unixMS, &parsed); err != nil { return nil, err } vin = strings.ToUpper(strings.TrimSpace(vin)) capacity, capacityOK := capacities[vin] pressure, temperature, fuelCellActive, fuelCellStateKnown, ok := extractHydrogenTelemetryFast(parsed) if !capacityOK || !ok { continue } mass, ok := PressureHydrogenMassKg(pressure, temperature, capacity) if !ok { continue } stepMass, _ := PressureHydrogenMassKg(math.Max(0, pressure-0.2), temperature, capacity) noise := math.Min(1, math.Max(0.05, mass-stepMass)) observation := HydrogenObservation{ VIN: vin, Source: source.String, EventID: eventID.String, ObservedAt: time.UnixMilli(unixMS), MassKg: mass, TankCapacityLiter: capacity, PressureMPa: pressure, TemperatureC: temperature, NoiseKg: noise, RefuelThresholdKg: math.Max(1, mass*0.05), FuelCellActive: fuelCellActive, FuelCellStateKnown: fuelCellStateKnown, } populateHydrogenObservationStateFast([]byte(parsed), &observation) observations = append(observations, observation) } return observations, rows.Err() } func buildHydrogenObservationQuery(database string, start, end time.Time, vin string) (string, error) { if !validTDIdentifier(database) { return "", fmt.Errorf("invalid TDengine database %q", database) } if !end.After(start) { return "", errors.New("hydrogen observation end must be after start") } vinFilter := "" if vin != "" { if !validHydrogenVIN(vin) { return "", fmt.Errorf("invalid VIN %q", vin) } vinFilter = "\n AND vin='" + vin + "'" } return `SELECT vin,source_endpoint,event_id,CAST(event_time AS BIGINT),parsed_json FROM ` + database + `.raw_frames WHERE protocol='GB32960' AND ts>='` + quoteTDTime(start) + `' AND ts<'` + quoteTDTime(end) + `' AND event_time>='` + quoteTDTime(start) + `' AND event_time<'` + quoteTDTime(end) + `'` + vinFilter + ` AND parse_status='OK' ORDER BY vin,event_time,source_endpoint,ts`, nil } func populateHydrogenObservationStateFast(data []byte, observation *HydrogenObservation) { if value, ok := jsonNumericField(data, "gb32960.vehicle.soc_percent"); ok && value >= 0 && value <= 100 { observation.SOCPercent, observation.SOCKnown = value, true } if value, ok := jsonNumericField(data, "gb32960.vehicle.total_mileage_km"); ok && value >= 0 { observation.MileageKm, observation.MileageKnown = value, true } if value, ok := jsonNumericField(data, "gb32960.vehicle.vehicle_status"); ok { observation.VehicleState, observation.VehicleStateKnown = int(value), value >= 0 && value <= 255 } if value, ok := jsonNumericField(data, "gb32960.vehicle.charge_status"); ok { observation.ChargeState, observation.ChargeStateKnown = int(value), value >= 0 && value <= 255 } if value, ok := jsonNumericField(data, "gb32960.vehicle.running_mode"); ok { observation.RunningMode, observation.RunningModeKnown = int(value), value >= 0 && value <= 255 } voltage, voltageOK := jsonNumericField(data, "gb32960.fuel_cell.fuel_cell_voltage_v") current, currentOK := jsonNumericField(data, "gb32960.fuel_cell.fuel_cell_current_a") if voltageOK && currentOK && voltage > 0 && voltage <= 1000 && current >= 0 && current <= 2000 { observation.FuelCellVoltageV = voltage observation.FuelCellCurrentA = current observation.FuelCellPowerKnown = true } } func buildHydrogenObservationVINQuery(database string, start, end time.Time) (string, error) { database = strings.TrimSpace(database) if database == "" { database = "lingniu_vehicle_ts" } if !validTDIdentifier(database) { return "", fmt.Errorf("invalid TDengine database %q", database) } if !end.After(start) { return "", errors.New("hydrogen observation end must be after start") } return `SELECT DISTINCT vin FROM ` + database + `.raw_frames WHERE protocol='GB32960' AND ts>='` + quoteTDTime(start) + `' AND ts<'` + quoteTDTime(end) + `' AND event_time>='` + quoteTDTime(start) + `' AND event_time<'` + quoteTDTime(end) + `' AND parse_status='OK' ORDER BY vin`, nil } func ExtractHydrogenPressureTemperature(parsedJSON string) (float64, float64, bool) { var fields map[string]any if err := json.Unmarshal([]byte(parsedJSON), &fields); err != nil { return 0, 0, false } return extractHydrogenPressureTemperatureFields(fields) } func ExtractHydrogenTelemetry(parsedJSON string) (pressure, temperature float64, active, known, ok bool) { var fields map[string]any if err := json.Unmarshal([]byte(parsedJSON), &fields); err != nil { return 0, 0, false, false, false } pressure, temperature, ok = extractHydrogenPressureTemperatureFields(fields) if !ok { return 0, 0, false, false, false } active, known = extractHydrogenFuelCellStateFields(fields) return pressure, temperature, active, known, true } func extractHydrogenTelemetryFast(parsedJSON string) (pressure, temperature float64, active, known, ok bool) { data := []byte(parsedJSON) pressure, pressureOK := jsonNumericField(data, "gb32960.fuel_cell.max_hydrogen_pressure_mpa") temperature, temperatureOK := jsonNumericField(data, "gb32960.fuel_cell.max_hydrogen_temperature_c") if !pressureOK || !temperatureOK || !validHydrogenPressureTemperature(pressure, temperature) { return 0, 0, false, false, false } if current, exists := jsonNumericField(data, "gb32960.fuel_cell.fuel_cell_current_a"); exists && current >= 0 { return pressure, temperature, current > hydrogenFuelCellActiveCurrentA, true, true } if state, exists := jsonNumericField(data, "gb32960.gd_fc_stack.engine_work_state"); exists { switch state { case 2: return pressure, temperature, true, true, true case 0: return pressure, temperature, false, true, true } } return pressure, temperature, false, false, true } func jsonNumericField(data []byte, key string) (float64, bool) { needle := []byte(`"` + key + `"`) for offset := 0; offset < len(data); { match := bytes.Index(data[offset:], needle) if match < 0 { return 0, false } cursor := offset + match + len(needle) for cursor < len(data) && (data[cursor] == ' ' || data[cursor] == '\t' || data[cursor] == '\r' || data[cursor] == '\n') { cursor++ } if cursor >= len(data) || data[cursor] != ':' { offset = cursor continue } cursor++ for cursor < len(data) && (data[cursor] == ' ' || data[cursor] == '\t' || data[cursor] == '\r' || data[cursor] == '\n') { cursor++ } quoted := cursor < len(data) && data[cursor] == '"' if quoted { cursor++ } start := cursor for cursor < len(data) && jsonNumberByte(data[cursor]) { cursor++ } if start == cursor || (quoted && (cursor >= len(data) || data[cursor] != '"')) { return 0, false } value, err := strconv.ParseFloat(string(data[start:cursor]), 64) return value, err == nil && !math.IsNaN(value) && !math.IsInf(value, 0) } return 0, false } func jsonNumberByte(value byte) bool { return value == '+' || value == '-' || value == '.' || value == 'e' || value == 'E' || (value >= '0' && value <= '9') } func extractHydrogenPressureTemperatureFields(fields map[string]any) (float64, float64, bool) { pressure, pressureOK := numericValue(fields["gb32960.fuel_cell.max_hydrogen_pressure_mpa"]) temperature, temperatureOK := numericValue(fields["gb32960.fuel_cell.max_hydrogen_temperature_c"]) return pressure, temperature, pressureOK && temperatureOK && validHydrogenPressureTemperature(pressure, temperature) } func ExtractHydrogenFuelCellState(parsedJSON string) (active, known bool) { var fields map[string]any if err := json.Unmarshal([]byte(parsedJSON), &fields); err != nil { return false, false } return extractHydrogenFuelCellStateFields(fields) } func extractHydrogenFuelCellStateFields(fields map[string]any) (active, known bool) { if current, ok := numericValue(fields["gb32960.fuel_cell.fuel_cell_current_a"]); ok && current >= 0 { return current > hydrogenFuelCellActiveCurrentA, true } if state, ok := numericValue(fields["gb32960.gd_fc_stack.engine_work_state"]); ok { switch state { case 2: return true, true case 0: return false, true } } return false, false } func hydrogenConsumptionIntervalEligible(previous, current HydrogenObservation) bool { if previous.FuelCellStateKnown || current.FuelCellStateKnown { return previous.FuelCellStateKnown && previous.FuelCellActive && current.FuelCellStateKnown && current.FuelCellActive } // Preserve compatibility for protocols or historical rows without a known // fuel-cell state. Those intervals still rely on pressure hysteresis alone. return true } func rapidHydrogenPressureRecovery(previous, current HydrogenObservation) bool { elapsed := current.ObservedAt.Sub(previous.ObservedAt) return elapsed > 0 && elapsed <= hydrogenRapidRecoveryWindow && current.PressureMPa-previous.PressureMPa >= hydrogenRapidRecoveryPressureMPa } func validHydrogenPressureTemperature(pressureMPa, temperatureC float64) bool { return pressureMPa > 0 && pressureMPa <= 70 && temperatureC > -40 && temperatureC <= 726.85 && !math.IsNaN(pressureMPa) && !math.IsInf(pressureMPa, 0) && !math.IsNaN(temperatureC) && !math.IsInf(temperatureC, 0) } func validHydrogenVIN(vin string) bool { if len(vin) != 17 { return false } for _, character := range vin { if character < '0' || character > '9' { if character < 'A' || character > 'Z' { return false } } } return true } func validTDIdentifier(value string) bool { if value == "" { return false } for _, character := range value { if character != '_' && (character < '0' || character > '9') && (character < 'A' || character > 'Z') && (character < 'a' || character > 'z') { return false } } return true } func HydrogenObservationDateRange(ctx context.Context, tdengine *sql.DB, database, vin string) (time.Time, time.Time, bool, error) { vin = strings.ToUpper(strings.TrimSpace(vin)) if !validHydrogenVIN(vin) { return time.Time{}, time.Time{}, false, errors.New("invalid VIN") } return hydrogenObservationDateRange(ctx, tdengine, database, vin) } func HydrogenObservationDateRangeForAll(ctx context.Context, tdengine *sql.DB, database string) (time.Time, time.Time, bool, error) { return hydrogenObservationDateRange(ctx, tdengine, database, "") } func hydrogenObservationDateRange(ctx context.Context, tdengine *sql.DB, database, vin string) (time.Time, time.Time, bool, error) { database = strings.TrimSpace(database) if database == "" { database = "lingniu_vehicle_ts" } if !validTDIdentifier(database) { return time.Time{}, time.Time{}, false, errors.New("invalid database") } vinFilter := "" if vin != "" { vinFilter = " AND vin='" + vin + "'" } loadBoundary := func(direction string) (time.Time, bool, error) { query := `SELECT event_time FROM ` + database + `.raw_frames WHERE protocol='GB32960'` + vinFilter + ` AND parse_status='OK' AND event_time IS NOT NULL ORDER BY event_time ` + direction + ` LIMIT 1` row := tdengine.QueryRowContext(ctx, query) var observed time.Time if err := row.Scan(&observed); err != nil { if errors.Is(err, sql.ErrNoRows) { return time.Time{}, false, nil } return time.Time{}, false, err } return observed, true, nil } first, found, err := loadBoundary("ASC") if err != nil || !found { return time.Time{}, time.Time{}, false, err } last, found, err := loadBoundary("DESC") if err != nil || !found { return time.Time{}, time.Time{}, false, err } return first, last, true, nil } func PressureHydrogenMassKg(pressureMPa, temperatureC, capacityLiter float64) (float64, bool) { temperatureK := temperatureC + 273.15 if pressureMPa < 0 || pressureMPa > 70 || temperatureK < 220 || temperatureK > 1000 || capacityLiter <= 0 || capacityLiter > 10000 { return 0, false } a := [...]float64{0.05888460, -0.06136111, -0.002650473, 0.002731125, 0.001802374, -0.001150707, 0.00009588528, -0.0000001109040, 0.0000000001264403} b := [...]float64{1.325, 1.87, 2.5, 2.8, 2.938, 3.14, 3.37, 3.75, 4.0} c := [...]float64{1, 1, 2, 2, 2.42, 2.63, 3, 4, 5} z := 1.0 for index := range a { z += a[index] * math.Pow(100/temperatureK, b[index]) * math.Pow(pressureMPa, c[index]) } if z <= 0 { return 0, false } density := pressureMPa * 1000 / (8.314472 * temperatureK * z) * 0.00201588 * 1000 mass := density * capacityLiter / 1000 return mass, !math.IsNaN(mass) && !math.IsInf(mass, 0) && mass >= 0 && mass <= 500 } func ReplaceHydrogenDailyStats(ctx context.Context, db *sql.DB, date string, stats []HydrogenDailyStat) error { return replaceHydrogenDailyStats(ctx, db, date, nil, stats, false) } // ReplaceHydrogenDailyStatsAndSeedStream performs the current-day deployment // handoff atomically. The stream resumes after the exact final observation // included in the rebuild, so Kafka backlog already covered by the rebuild is // ignored instead of counted twice. func ReplaceHydrogenDailyStatsAndSeedStream(ctx context.Context, db *sql.DB, date string, stats []HydrogenDailyStat) error { return replaceHydrogenDailyStats(ctx, db, date, nil, stats, true) } func ReplaceHydrogenDailyStatsForVIN(ctx context.Context, db *sql.DB, date, vin string, stats []HydrogenDailyStat) error { vin = strings.ToUpper(strings.TrimSpace(vin)) if !validHydrogenVIN(vin) { return fmt.Errorf("invalid VIN %q", vin) } return replaceHydrogenDailyStats(ctx, db, date, []string{vin}, stats, false) } func ReplaceHydrogenDailyStatsForVINs(ctx context.Context, db *sql.DB, date string, vins []string, stats []HydrogenDailyStat) error { normalized := make([]string, 0, len(vins)) seen := make(map[string]struct{}, len(vins)) for _, vin := range vins { vin = strings.ToUpper(strings.TrimSpace(vin)) if !validHydrogenVIN(vin) { return fmt.Errorf("invalid VIN %q", vin) } if _, exists := seen[vin]; exists { continue } seen[vin] = struct{}{} normalized = append(normalized, vin) } if len(normalized) == 0 { return errors.New("at least one scoped VIN is required") } return replaceHydrogenDailyStats(ctx, db, date, normalized, stats, false) } func replaceHydrogenDailyStats(ctx context.Context, db *sql.DB, date string, scopedVINs []string, stats []HydrogenDailyStat, seedStream bool) error { tx, err := db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback() // A pressure-based rebuild is authoritative for the whole day. Delete every // previous hydrogen row first so legacy rate/direct-mass results cannot remain // for vehicles without valid pressure observations in this run. allowedVINs := make(map[string]struct{}, len(scopedVINs)) for _, vin := range scopedVINs { allowedVINs[vin] = struct{}{} } if len(scopedVINs) == 0 { if _, err := tx.ExecContext(ctx, `DELETE FROM vehicle_open_daily_energy WHERE stat_date=? AND energy_type='HYDROGEN'`, date); err != nil { return err } } else if len(scopedVINs) == 1 { if _, err := tx.ExecContext(ctx, `DELETE FROM vehicle_open_daily_energy WHERE stat_date=? AND energy_type='HYDROGEN' AND vin=?`, date, scopedVINs[0]); err != nil { return err } } else { placeholders := strings.TrimSuffix(strings.Repeat("?,", len(scopedVINs)), ",") args := make([]any, 0, len(scopedVINs)+1) args = append(args, date) for _, vin := range scopedVINs { args = append(args, vin) } if _, err := tx.ExecContext(ctx, `DELETE FROM vehicle_open_daily_energy WHERE stat_date=? AND energy_type='HYDROGEN' AND vin IN (`+placeholders+`)`, args...); err != nil { return err } } if seedStream { if _, err := tx.ExecContext(ctx, `DELETE FROM vehicle_open_hydrogen_segment_stream_state WHERE stat_date=?`, date); err != nil { return err } } for _, stat := range stats { statVIN := strings.ToUpper(strings.TrimSpace(stat.VIN)) if len(allowedVINs) > 0 { if _, allowed := allowedVINs[statVIN]; !allowed { return fmt.Errorf("refusing to write VIN %q outside scoped rebuild", stat.VIN) } } parameterJSON, err := json.Marshal(stat.CalculationParameters) if err != nil { return fmt.Errorf("encode hydrogen calculation parameters for %s: %w", stat.VIN, err) } evidenceJSON, err := json.Marshal(stat.Intervals) if err != nil { return fmt.Errorf("encode hydrogen interval evidence for %s: %w", stat.VIN, err) } if _, err := tx.ExecContext(ctx, ` INSERT INTO vehicle_open_daily_energy( vin,stat_date,energy_type,source_endpoint,consumption_kg,unit,first_mass_kg,last_mass_kg, sample_count,refuel_count,quality_status,quality_reason,calculated_at, refuel_amount_kg, raw_consumption_kg,battery_soc_delta_pct,battery_discharge_kwh,battery_equivalent_kg, soc_balanced_consumption_kg,mixed_mileage_km,pure_electric_mileage_km, consumption_kg_per_100km,soc_balanced_kg_per_100km,charge_count,charge_energy_kwh,valid_segment_count, invalid_segment_count,algorithm_version,calculation_phase,parameter_json,evidence_json ) VALUES(?,?,'HYDROGEN',?,?,'kg',?,?,?,?,?,?,NOW(3),?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,'FINAL',?,?)`, stat.VIN, stat.Date, stat.Source, stat.ConsumptionKg, stat.FirstMassKg, stat.LastMassKg, stat.SampleCount, stat.RefuelCount, stat.QualityStatus, stat.QualityReason, stat.RefuelAmountKg, stat.ConsumptionKg, stat.BatterySOCDeltaPct, stat.BatteryDischargeKWh, stat.BatteryEquivalentKg, stat.SOCBalancedConsumptionKg, stat.MixedMileageKm, stat.PureElectricMileageKm, stat.ConsumptionKgPer100Km, stat.SOCBalancedKgPer100Km, stat.ChargeCount, stat.ChargeEnergyKWh, stat.QualifiedSegmentCount, stat.InvalidSegmentCount, stat.CalculationParameters.AlgorithmVersion, parameterJSON, evidenceJSON, ); err != nil { return err } // A NO_DATA row can legitimately have no effective boundary observation. // Keep its daily result, but do not write a zero timestamp into the stream // watermark table; the next valid realtime frame will create the state. if seedStream && !stat.LastObservation.ObservedAt.IsZero() { stateJSON, err := hydrogenSegmentStreamSeedJSON(stat) if err != nil { return err } if _, err := tx.ExecContext(ctx, ` INSERT INTO vehicle_open_hydrogen_segment_stream_state( vin,stat_date,source_endpoint,finalized_consumption_kg,projected_consumption_kg, sample_count,refuel_count,abnormal_drop_count,eligible_interval_count,qualified_segment_count, last_mass_kg,last_event_time,last_event_id,state_json,calculation_method,quality_status,quality_reason ) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,'PRESSURE_NIST_VALID_BOUNDARY_CHARGE_CYCLE_V3_5',?,?)`, stat.VIN, stat.Date, stat.Source, stat.ConsumptionKg, stat.ConsumptionKg, stat.SampleCount, stat.RefuelCount, stat.AbnormalDropCount, stat.EligibleIntervalCount, stat.QualifiedSegmentCount, stat.LastMassKg, stat.LastObservation.ObservedAt, stat.LastObservation.EventID, stateJSON, stat.QualityStatus, stat.QualityReason, ); err != nil { return err } } } return tx.Commit() } func hydrogenSegmentStreamSeedJSON(stat HydrogenDailyStat) ([]byte, error) { type point struct { EventID string `json:"eventId"` ObservedAt time.Time `json:"observedAt"` MassKg float64 `json:"massKg"` PressureMpa float64 `json:"pressureMpa"` TemperatureC float64 `json:"temperatureC"` NoiseKg float64 `json:"noiseKg"` MileageKm float64 `json:"mileageKm"` MileageKnown bool `json:"mileageKnown"` SocPercent float64 `json:"socPercent"` SocKnown bool `json:"socKnown"` RunningMode int `json:"runningMode"` FuelCellActive bool `json:"fuelCellActive"` FuelCellStateKnown bool `json:"fuelCellStateKnown"` FuelCellPowerKw float64 `json:"fuelCellPowerKw"` FuelCellPowerKnown bool `json:"fuelCellPowerKnown"` } toPoint := func(value HydrogenObservation) point { return point{ EventID: value.EventID, ObservedAt: value.ObservedAt, MassKg: value.MassKg, PressureMpa: value.PressureMPa, TemperatureC: value.TemperatureC, NoiseKg: value.NoiseKg, MileageKm: value.MileageKm, MileageKnown: value.MileageKnown, SocPercent: value.SOCPercent, SocKnown: value.SOCKnown, RunningMode: value.RunningMode, FuelCellActive: value.FuelCellActive, FuelCellStateKnown: value.FuelCellStateKnown, FuelCellPowerKw: value.FuelCellVoltageV * value.FuelCellCurrentA / 1000, FuelCellPowerKnown: value.FuelCellPowerKnown, } } first, last := toPoint(stat.FirstObservation), toPoint(stat.LastObservation) lastIntervalType := "" mixedCycles := 0 for _, interval := range stat.Intervals { if interval.Type == "MIXED" { mixedCycles++ } lastIntervalType = interval.Type } cycle := map[string]any{ "first": last, "last": last, "hasRun": !last.ObservedAt.IsZero(), "startsPure": lastIntervalType == "PURE_ELECTRIC", "mixedLocked": lastIntervalType == "MIXED", "mixedConfirmed": lastIntervalType == "MIXED", "mixedStart": last, } state := struct { AlgorithmVersion string `json:"algorithmVersion"` SourceEndpoint string `json:"sourceEndpoint"` SampleCount int `json:"sampleCount"` LastEventID string `json:"lastEventId"` LastEventTime time.Time `json:"lastEventTime"` Parameters map[string]float64 `json:"parameters"` FirstMassKg float64 `json:"firstMassKg"` LastMassKg float64 `json:"lastMassKg"` FirstEffective point `json:"firstEffective"` LastEffective point `json:"lastEffective"` HasEffective bool `json:"hasEffective"` EffectiveRunCount int `json:"effectiveRunCount"` Cycle map[string]any `json:"cycle"` InitialStateUsed bool `json:"initialStateUsed"` ChargeCount int `json:"chargeCount"` ChargeEnergyKWh float64 `json:"chargeEnergyKWh"` CompletedPureKm float64 `json:"completedPureKm"` CompletedMixedKm float64 `json:"completedMixedKm"` CompletedSOCDelta float64 `json:"completedSocDelta"` CompletedMixedCycles int `json:"completedMixedCycles"` HydrogenSegmentStart point `json:"hydrogenSegmentStart"` LastHydrogenRun point `json:"lastHydrogenRun"` HasHydrogenSegment bool `json:"hasHydrogenSegment"` HydrogenSegmentRuns int `json:"hydrogenSegmentRuns"` FinalizedHydrogenKg float64 `json:"finalizedHydrogenKg"` HydrogenSegments int `json:"hydrogenSegments"` RefuelCount int `json:"refuelCount"` RefuelAmountKg float64 `json:"refuelAmountKg"` MinimumPressureMpa float64 `json:"minimumPressureMpa"` MinimumMassKg float64 `json:"minimumMassKg"` MinimumObservedAt time.Time `json:"minimumObservedAt"` PreviousSample point `json:"previousSample"` HasPreviousSample bool `json:"hasPreviousSample"` AbnormalDropCount int `json:"abnormalDropCount"` InvalidSampleCount int `json:"invalidSampleCount"` SuspectedLeakCount int `json:"suspectedLeakCount"` SuspectedLeakMaxMpa float64 `json:"suspectedLeakMaxMpa"` }{ AlgorithmVersion: trustedHydrogenAlgorithmVersion, SourceEndpoint: stat.Source, SampleCount: stat.SampleCount, LastEventID: stat.LastObservation.EventID, LastEventTime: stat.LastObservation.ObservedAt, Parameters: map[string]float64{"batteryCapacityKWh": stat.CalculationParameters.BatteryCapacityKWh, "hydrogenEnergyKWhPerKg": stat.CalculationParameters.HydrogenEnergyKWhKg}, FirstMassKg: stat.FirstMassKg, LastMassKg: stat.LastMassKg, FirstEffective: first, LastEffective: last, HasEffective: !last.ObservedAt.IsZero(), EffectiveRunCount: stat.EligibleIntervalCount, Cycle: cycle, InitialStateUsed: true, ChargeCount: stat.ChargeCount, ChargeEnergyKWh: stat.ChargeEnergyKWh, CompletedPureKm: stat.PureElectricMileageKm, CompletedMixedKm: stat.MixedMileageKm, CompletedSOCDelta: valueOrZero(stat.BatterySOCDeltaPct), CompletedMixedCycles: mixedCycles, HydrogenSegmentStart: last, LastHydrogenRun: last, HasHydrogenSegment: !last.ObservedAt.IsZero(), HydrogenSegmentRuns: 1, FinalizedHydrogenKg: stat.ConsumptionKg, HydrogenSegments: stat.QualifiedSegmentCount, RefuelCount: stat.RefuelCount, RefuelAmountKg: stat.RefuelAmountKg, MinimumPressureMpa: last.PressureMpa, MinimumMassKg: last.MassKg, MinimumObservedAt: last.ObservedAt, PreviousSample: last, HasPreviousSample: !last.ObservedAt.IsZero(), AbnormalDropCount: stat.AbnormalDropCount, InvalidSampleCount: stat.InvalidSegmentCount, SuspectedLeakCount: stat.SuspectedLeakCount, SuspectedLeakMaxMpa: stat.SuspectedLeakMaxPressureDropMPa, } return json.Marshal(state) } func valueOrZero(value *float64) float64 { if value == nil { return 0 } return *value } func numericValue(value any) (float64, bool) { switch typed := value.(type) { case float64: return typed, true case json.Number: parsed, err := typed.Float64() return parsed, err == nil case string: parsed, err := strconv.ParseFloat(strings.TrimSpace(typed), 64) return parsed, err == nil default: return 0, false } } func quoteTDTime(value time.Time) string { return strings.ReplaceAll(value.Format(time.RFC3339Nano), "'", "''") }