package stats import ( "context" "database/sql" "math" "strings" "time" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/telemetry" ) const ( defaultHydrogenNoiseKG = 0.05 defaultHydrogenMaxDropKG = 20.0 hydrogenMolarMassKGPerMol = 0.00201588 universalGasConstant = 8.314472 ) var hydrogenPressureFieldKeys = []string{ "gb32960.fuel_cell.max_hydrogen_pressure_mpa", "fuel_cell_max_hydrogen_pressure_mpa", } var hydrogenTemperatureFieldKeys = []string{ "gb32960.fuel_cell.max_hydrogen_temperature_c", "fuel_cell_max_hydrogen_temperature_c", } var hydrogenDensityA = [...]float64{0.05888460, -0.06136111, -0.002650473, 0.002731125, 0.001802374, -0.001150707, 0.00009588528, -0.0000001109040, 0.0000000001264403} var hydrogenDensityB = [...]float64{1.325, 1.87, 2.5, 2.8, 2.938, 3.14, 3.37, 3.75, 4.0} var hydrogenDensityC = [...]float64{1, 1, 2, 2, 2.42, 2.63, 3, 4, 5} const HydrogenStreamStateTableSQL = `CREATE TABLE IF NOT EXISTS vehicle_open_hydrogen_stream_state ( vin VARCHAR(64) NOT NULL, stat_date DATE NOT NULL, source_endpoint VARCHAR(128) NOT NULL DEFAULT '', first_mass_kg DECIMAL(18,3) NOT NULL, last_mass_kg DECIMAL(18,3) NOT NULL, cycle_min_mass_kg DECIMAL(18,3) NOT NULL, consumption_kg DECIMAL(18,3) NOT NULL DEFAULT 0, sample_count BIGINT UNSIGNED NOT NULL DEFAULT 1, refuel_count INT UNSIGNED NOT NULL DEFAULT 0, abnormal_drop_count INT UNSIGNED NOT NULL DEFAULT 0, tank_capacity_l DECIMAL(12,2) NOT NULL, first_pressure_mpa DECIMAL(10,3) NOT NULL, last_pressure_mpa DECIMAL(10,3) NOT NULL, first_temperature_c DECIMAL(10,3) NOT NULL, last_temperature_c DECIMAL(10,3) NOT NULL, calculation_method VARCHAR(32) NOT NULL DEFAULT 'PRESSURE_NIST', first_event_time DATETIME(3) NOT NULL, last_event_time DATETIME(3) NOT NULL, last_event_id VARCHAR(128) NOT NULL DEFAULT '', quality_status VARCHAR(24) NOT NULL DEFAULT 'NO_DATA', quality_reason VARCHAR(255) NOT NULL DEFAULT '有效压力质量样本不足2条', created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), updated_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3), PRIMARY KEY (vin, stat_date, source_endpoint), KEY idx_hydrogen_stream_date_quality (stat_date, quality_status, vin) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci` type HydrogenStreamSample struct { VIN string Date string SourceEndpoint string EventID string EventTime time.Time MassKG float64 TankCapacityLiters float64 PressureMPa float64 TemperatureC float64 NoiseKG float64 RefuelThresholdKG float64 } type HydrogenStreamResult struct { Found int Written int Duplicate int Invalid int NotCurrent int } func HydrogenDensityKGPerM3(pressureMPa, temperatureC float64) (float64, bool) { temperatureK := temperatureC + 273.15 if pressureMPa < 0 || pressureMPa > 70 || temperatureK < 220 || temperatureK > 1000 { return 0, false } z := 1.0 for i := range hydrogenDensityA { z += hydrogenDensityA[i] * math.Pow(100/temperatureK, hydrogenDensityB[i]) * math.Pow(pressureMPa, hydrogenDensityC[i]) } if z <= 0 { return 0, false } molPerLiter := pressureMPa * 1000 / (universalGasConstant * temperatureK * z) return molPerLiter * hydrogenMolarMassKGPerMol * 1000, true } func PressureHydrogenMassKG(pressureMPa, temperatureC, tankCapacityLiters float64) (float64, bool) { if tankCapacityLiters <= 0 || tankCapacityLiters > 10000 { return 0, false } density, ok := HydrogenDensityKGPerM3(pressureMPa, temperatureC) if !ok { return 0, false } mass := density * tankCapacityLiters / 1000 return mass, !math.IsNaN(mass) && !math.IsInf(mass, 0) && mass >= 0 && mass <= 500 } func HydrogenStreamSampleFromEnvelope(env envelope.FrameEnvelope, loc *time.Location, now time.Time, tankCapacityLiters float64) (HydrogenStreamSample, string, bool) { if env.Protocol != envelope.ProtocolGB32960 { return HydrogenStreamSample{}, "unsupported_protocol", false } vin := strings.ToUpper(strings.TrimSpace(env.VIN)) if len(vin) != 17 { return HydrogenStreamSample{}, "invalid_vin", false } pressure, pressureOK := firstHydrogenNumber(env.Fields, hydrogenPressureFieldKeys) if !pressureOK { return HydrogenStreamSample{}, "missing_pressure", false } temperature, temperatureOK := firstHydrogenNumber(env.Fields, hydrogenTemperatureFieldKeys) if !temperatureOK { return HydrogenStreamSample{}, "missing_temperature", false } if tankCapacityLiters <= 0 || tankCapacityLiters > 10000 { return HydrogenStreamSample{}, "missing_tank_capacity", false } mass, ok := PressureHydrogenMassKG(pressure, temperature, tankCapacityLiters) if !ok { return HydrogenStreamSample{}, "invalid_pressure_mass", false } pressureStepMass, _ := PressureHydrogenMassKG(math.Max(0, pressure-0.2), temperature, tankCapacityLiters) noiseKG := math.Max(defaultHydrogenNoiseKG, mass-pressureStepMass) noiseKG = math.Min(noiseKG, 1.0) refuelThresholdKG := math.Max(1.0, mass*0.05) if loc == nil { loc = time.FixedZone("Asia/Shanghai", 8*3600) } eventMS, _, ok := envelope.NormalizedEventTimeMSWithReason(env) if !ok { return HydrogenStreamSample{}, "missing_event_time", false } eventTime := time.UnixMilli(eventMS).In(loc) if now.IsZero() { now = time.Now() } date := eventTime.Format("2006-01-02") if date != now.In(loc).Format("2006-01-02") { return HydrogenStreamSample{}, "not_current_date", false } return HydrogenStreamSample{ VIN: vin, Date: date, SourceEndpoint: strings.TrimSpace(env.SourceEndpoint), EventID: env.StableEventID(), EventTime: eventTime, MassKG: mass, TankCapacityLiters: tankCapacityLiters, PressureMPa: pressure, TemperatureC: temperature, NoiseKG: noiseKG, RefuelThresholdKG: refuelThresholdKG, }, "", true } func firstHydrogenNumber(fields map[string]any, keys []string) (float64, bool) { for _, key := range keys { if value, ok := telemetry.Number(fields, key); ok { return value, true } } return 0, false } func AppendHydrogenStream(ctx context.Context, exec Execer, env envelope.FrameEnvelope, loc *time.Location, now time.Time, tankCapacityLiters float64) (HydrogenStreamResult, error) { sample, reason, ok := HydrogenStreamSampleFromEnvelope(env, loc, now, tankCapacityLiters) if !ok { result := HydrogenStreamResult{} switch reason { case "unsupported_protocol", "missing_pressure", "missing_temperature": return result, nil case "not_current_date": result.NotCurrent = 1 default: result.Invalid = 1 } return result, nil } result := HydrogenStreamResult{Found: 1} beginner, ok := exec.(txBeginner) if !ok { return result, sql.ErrTxDone } tx, err := beginner.BeginTx(ctx, nil) if err != nil { return result, err } defer tx.Rollback() write, err := tx.ExecContext(ctx, upsertHydrogenStreamStateSQL, sample.VIN, sample.Date, sample.SourceEndpoint, sample.MassKG, sample.MassKG, sample.MassKG, sample.TankCapacityLiters, sample.PressureMPa, sample.PressureMPa, sample.TemperatureC, sample.TemperatureC, sample.EventTime, sample.EventTime, sample.EventID, sample.NoiseKG, defaultHydrogenMaxDropKG, sample.RefuelThresholdKG, defaultHydrogenMaxDropKG, defaultHydrogenMaxDropKG, defaultHydrogenMaxDropKG, sample.RefuelThresholdKG, sample.NoiseKG, defaultHydrogenMaxDropKG, ) if err != nil { return result, err } affected, _ := write.RowsAffected() if affected == 0 { result.Duplicate = 1 return result, tx.Commit() } if _, err := tx.ExecContext(ctx, projectHydrogenStreamDailySQL, sample.VIN, sample.Date); err != nil { return result, err } if err := tx.Commit(); err != nil { return result, err } result.Written = 1 return result, nil } const upsertHydrogenStreamStateSQL = ` INSERT INTO vehicle_open_hydrogen_stream_state( vin,stat_date,source_endpoint,first_mass_kg,last_mass_kg,cycle_min_mass_kg,consumption_kg,sample_count, refuel_count,abnormal_drop_count,tank_capacity_l,first_pressure_mpa,last_pressure_mpa, first_temperature_c,last_temperature_c,calculation_method, first_event_time,last_event_time,last_event_id,quality_status,quality_reason ) VALUES(?,?,?,?,?,?,0,1,0,0,?,?,?,?,?,'PRESSURE_NIST',?,?,?,'NO_DATA','有效压力质量样本不足2条') ON DUPLICATE KEY UPDATE consumption_kg = consumption_kg + IF(VALUES(last_event_time)>last_event_time AND cycle_min_mass_kg-VALUES(last_mass_kg)>? AND cycle_min_mass_kg-VALUES(last_mass_kg)<=?, cycle_min_mass_kg-VALUES(last_mass_kg),0), refuel_count = refuel_count + IF(VALUES(last_event_time)>last_event_time AND VALUES(last_mass_kg)-cycle_min_mass_kg>?,1,0), quality_status = IF(VALUES(last_event_time)<=last_event_time,quality_status, IF(abnormal_drop_count+IF(cycle_min_mass_kg-VALUES(last_mass_kg)>?,1,0)>0,'SUSPECT','OK')), quality_reason = IF(VALUES(last_event_time)<=last_event_time,quality_reason, IF(abnormal_drop_count+IF(cycle_min_mass_kg-VALUES(last_mass_kg)>?,1,0)>0,'存在超过阈值的异常下降','')), abnormal_drop_count = abnormal_drop_count + IF(VALUES(last_event_time)>last_event_time AND cycle_min_mass_kg-VALUES(last_mass_kg)>?,1,0), cycle_min_mass_kg = IF(VALUES(last_event_time)>last_event_time, IF(VALUES(last_mass_kg)-cycle_min_mass_kg>?,VALUES(last_mass_kg), IF(cycle_min_mass_kg-VALUES(last_mass_kg)>? AND cycle_min_mass_kg-VALUES(last_mass_kg)<=?,VALUES(last_mass_kg),cycle_min_mass_kg)), cycle_min_mass_kg), sample_count = sample_count + IF(VALUES(last_event_time)>last_event_time,1,0), tank_capacity_l = IF(VALUES(last_event_time)>last_event_time,VALUES(tank_capacity_l),tank_capacity_l), last_pressure_mpa = IF(VALUES(last_event_time)>last_event_time,VALUES(last_pressure_mpa),last_pressure_mpa), last_temperature_c = IF(VALUES(last_event_time)>last_event_time,VALUES(last_temperature_c),last_temperature_c), calculation_method = 'PRESSURE_NIST', last_mass_kg = IF(VALUES(last_event_time)>last_event_time,VALUES(last_mass_kg),last_mass_kg), last_event_id = IF(VALUES(last_event_time)>last_event_time,VALUES(last_event_id),last_event_id), last_event_time = GREATEST(last_event_time,VALUES(last_event_time))` const projectHydrogenStreamDailySQL = ` 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 ) SELECT vin,stat_date,'HYDROGEN',source_endpoint,consumption_kg,'kg',first_mass_kg,last_mass_kg, sample_count,refuel_count,quality_status,quality_reason,NOW(3) FROM vehicle_open_hydrogen_stream_state WHERE vin=? AND stat_date=? ORDER BY CASE quality_status WHEN 'OK' THEN 0 WHEN 'SUSPECT' THEN 1 ELSE 2 END, sample_count DESC,source_endpoint ASC LIMIT 1 ON DUPLICATE KEY UPDATE source_endpoint=VALUES(source_endpoint),consumption_kg=VALUES(consumption_kg),unit='kg', first_mass_kg=VALUES(first_mass_kg),last_mass_kg=VALUES(last_mass_kg), sample_count=VALUES(sample_count),refuel_count=VALUES(refuel_count), quality_status=VALUES(quality_status),quality_reason=VALUES(quality_reason),calculated_at=VALUES(calculated_at)`