package stats import ( "context" "database/sql" "encoding/json" "errors" "math" "sort" "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 hydrogenActiveCurrentA = 1.0 hydrogenMolarMassKGPerMol = 0.00201588 universalGasConstant = 8.314472 hydrogenSegmentWindowSize = 5 hydrogenSegmentMinSamples = hydrogenSegmentWindowSize * 2 hydrogenSegmentMaxGap = 5 * time.Minute hydrogenRecoveryWindow = time.Minute hydrogenRecoveryPressure = 8.0 ) 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 hydrogenCurrentFieldKeys = []string{ "gb32960.fuel_cell.fuel_cell_current_a", "fuel_cell_current_a", } 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 '', last_fuel_cell_active TINYINT(1) NOT NULL DEFAULT 0, last_fuel_cell_state_known TINYINT(1) NOT NULL DEFAULT 0, 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` const HydrogenSegmentStreamStateTableSQL = `CREATE TABLE IF NOT EXISTS vehicle_open_hydrogen_segment_stream_state ( vin VARCHAR(64) NOT NULL, stat_date DATE NOT NULL, source_endpoint VARCHAR(128) NOT NULL DEFAULT '', finalized_consumption_kg DECIMAL(18,3) NOT NULL DEFAULT 0, projected_consumption_kg DECIMAL(18,3) NOT NULL DEFAULT 0, sample_count BIGINT UNSIGNED NOT NULL DEFAULT 0, refuel_count INT UNSIGNED NOT NULL DEFAULT 0, abnormal_drop_count INT UNSIGNED NOT NULL DEFAULT 0, eligible_interval_count BIGINT UNSIGNED NOT NULL DEFAULT 0, qualified_segment_count INT UNSIGNED NOT NULL DEFAULT 0, last_mass_kg DECIMAL(18,3) NOT NULL DEFAULT 0, last_event_time DATETIME(3) NOT NULL, last_event_id VARCHAR(128) NOT NULL DEFAULT '', state_json JSON NOT NULL, calculation_method VARCHAR(48) NOT NULL DEFAULT 'PRESSURE_NIST_SEGMENT_MEDIAN_5', quality_status VARCHAR(24) NOT NULL DEFAULT 'NO_DATA', quality_reason VARCHAR(255) NOT NULL DEFAULT '有效工作段关键帧不足10条', 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), KEY idx_hydrogen_segment_stream_date_quality (stat_date, quality_status, vin), CONSTRAINT chk_hydrogen_segment_stream_quality CHECK (quality_status IN ('OK', 'NO_DATA', 'SUSPECT')) ) 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 FuelCellActive bool FuelCellStateKnown bool ConsumptionEligible bool } type HydrogenStreamResult struct { Found int Written int Duplicate int Invalid int NotCurrent int } type hydrogenStreamEndpoint struct { MassKG float64 `json:"massKg"` NoiseKG float64 `json:"noiseKg"` } type hydrogenStreamSegmentWindow struct { Count int64 `json:"count"` First []hydrogenStreamEndpoint `json:"first"` Tail []hydrogenStreamEndpoint `json:"tail"` CycleMinimumMassKG float64 `json:"cycleMinimumMassKg"` } type hydrogenSegmentStreamState struct { SourceEndpoint string `json:"sourceEndpoint"` FinalizedConsumptionKG float64 `json:"finalizedConsumptionKg"` SampleCount int64 `json:"sampleCount"` RefuelCount int64 `json:"refuelCount"` AbnormalDropCount int64 `json:"abnormalDropCount"` EligibleIntervalCount int64 `json:"eligibleIntervalCount"` QualifiedSegmentCount int64 `json:"qualifiedSegmentCount"` FirstMassKG float64 `json:"firstMassKg"` RemainingMassKG float64 `json:"remainingMassKg"` MetadataCycleMinimumKG float64 `json:"metadataCycleMinimumKg"` LastObservedMassKG float64 `json:"lastObservedMassKg"` TankCapacityLiters float64 `json:"tankCapacityLiters"` FirstPressureMPa float64 `json:"firstPressureMpa"` LastPressureMPa float64 `json:"lastPressureMpa"` FirstTemperatureC float64 `json:"firstTemperatureC"` LastTemperatureC float64 `json:"lastTemperatureC"` FirstEventTime time.Time `json:"firstEventTime"` LastEventTime time.Time `json:"lastEventTime"` LastEventID string `json:"lastEventId"` LastFuelCellActive bool `json:"lastFuelCellActive"` LastFuelCellStateKnown bool `json:"lastFuelCellStateKnown"` Segment hydrogenStreamSegmentWindow `json:"segment"` } type hydrogenStreamBaseline struct { SourceEndpoint string ConsumptionKG float64 FirstMassKG float64 LastMassKG float64 SampleCount int64 RefuelCount int64 QualityStatus string } func newHydrogenSegmentStreamState(sample HydrogenStreamSample, baseline *hydrogenStreamBaseline) hydrogenSegmentStreamState { state := hydrogenSegmentStreamState{ SourceEndpoint: sample.SourceEndpoint, FirstMassKG: sample.MassKG, RemainingMassKG: sample.MassKG, MetadataCycleMinimumKG: sample.MassKG, LastObservedMassKG: sample.MassKG, TankCapacityLiters: sample.TankCapacityLiters, FirstPressureMPa: sample.PressureMPa, LastPressureMPa: sample.PressureMPa, FirstTemperatureC: sample.TemperatureC, LastTemperatureC: sample.TemperatureC, FirstEventTime: sample.EventTime, LastEventTime: sample.EventTime, LastEventID: sample.EventID, LastFuelCellActive: sample.FuelCellActive, LastFuelCellStateKnown: sample.FuelCellStateKnown, SampleCount: 1, } if baseline != nil { state.SourceEndpoint = firstNonEmptyHydrogen(sample.SourceEndpoint, baseline.SourceEndpoint) state.FinalizedConsumptionKG = baseline.ConsumptionKG state.SampleCount += baseline.SampleCount state.RefuelCount = baseline.RefuelCount if baseline.FirstMassKG >= 0 { state.FirstMassKG = baseline.FirstMassKG } if baseline.LastMassKG >= 0 { state.RemainingMassKG = baseline.LastMassKG } if baseline.QualityStatus == "OK" || baseline.QualityStatus == "SUSPECT" { state.EligibleIntervalCount = 1 state.QualifiedSegmentCount = 1 } if baseline.QualityStatus == "SUSPECT" { state.AbnormalDropCount = 1 } } if hydrogenStreamCanStartSegment(sample) { state.startSegment(sample) } return state } func (state *hydrogenSegmentStreamState) add(sample HydrogenStreamSample) bool { if !sample.EventTime.After(state.LastEventTime) { return false } previousTime := state.LastEventTime previousMass := state.LastObservedMassKG previousPressure := state.LastPressureMPa intervalEligible := hydrogenStreamIntervalEligible( state.LastFuelCellActive, state.LastFuelCellStateKnown, sample.FuelCellActive, sample.FuelCellStateKnown, ) if intervalEligible || (sample.FuelCellStateKnown && sample.FuelCellActive) || sample.MassKG-state.RemainingMassKG > sample.NoiseKG { state.RemainingMassKG = sample.MassKG } delta := state.MetadataCycleMinimumKG - sample.MassKG switch { case sample.MassKG-state.MetadataCycleMinimumKG > sample.RefuelThresholdKG: if !hydrogenStreamRapidRecovery(previousTime, previousPressure, sample) { state.RefuelCount++ } state.MetadataCycleMinimumKG = sample.MassKG case delta > sample.NoiseKG && delta <= defaultHydrogenMaxDropKG: state.MetadataCycleMinimumKG = sample.MassKG } if intervalEligible { state.EligibleIntervalCount++ } switch { case state.Segment.Count == 0 || !intervalEligible: state.flushSegment() if hydrogenStreamCanStartSegment(sample) { state.startSegment(sample) } case sample.EventTime.Sub(previousTime) > hydrogenSegmentMaxGap: state.flushSegment() state.startSegment(sample) case previousMass-sample.MassKG > defaultHydrogenMaxDropKG: state.AbnormalDropCount++ state.flushSegment() state.startSegment(sample) case sample.MassKG-state.Segment.CycleMinimumMassKG > sample.RefuelThresholdKG: state.flushSegment() state.startSegment(sample) default: state.appendSegmentSample(sample) if sample.MassKG < state.Segment.CycleMinimumMassKG { state.Segment.CycleMinimumMassKG = sample.MassKG } } state.SourceEndpoint = sample.SourceEndpoint state.SampleCount++ state.LastObservedMassKG = sample.MassKG state.TankCapacityLiters = sample.TankCapacityLiters state.LastPressureMPa = sample.PressureMPa state.LastTemperatureC = sample.TemperatureC state.LastEventTime = sample.EventTime state.LastEventID = sample.EventID state.LastFuelCellActive = sample.FuelCellActive state.LastFuelCellStateKnown = sample.FuelCellStateKnown return true } func (state *hydrogenSegmentStreamState) startSegment(sample HydrogenStreamSample) { state.Segment = hydrogenStreamSegmentWindow{CycleMinimumMassKG: sample.MassKG} state.appendSegmentSample(sample) } func (state *hydrogenSegmentStreamState) appendSegmentSample(sample HydrogenStreamSample) { endpoint := hydrogenStreamEndpoint{MassKG: sample.MassKG, NoiseKG: sample.NoiseKG} state.Segment.Count++ if len(state.Segment.First) < hydrogenSegmentWindowSize { state.Segment.First = append(state.Segment.First, endpoint) } if len(state.Segment.Tail) == hydrogenSegmentWindowSize { copy(state.Segment.Tail, state.Segment.Tail[1:]) state.Segment.Tail[len(state.Segment.Tail)-1] = endpoint return } state.Segment.Tail = append(state.Segment.Tail, endpoint) } func (state *hydrogenSegmentStreamState) flushSegment() { if state.Segment.Count >= hydrogenSegmentMinSamples { state.QualifiedSegmentCount++ state.FinalizedConsumptionKG += hydrogenStreamEndpointDrop(state.Segment.First, state.Segment.Tail) } state.Segment = hydrogenStreamSegmentWindow{} } func (state hydrogenSegmentStreamState) projectedConsumptionKG() float64 { consumption := state.FinalizedConsumptionKG if state.Segment.Count >= hydrogenSegmentMinSamples { consumption += hydrogenStreamEndpointDrop(state.Segment.First, state.Segment.Tail) } return roundHydrogenKG(consumption) } func (state hydrogenSegmentStreamState) quality() (string, string) { qualified := state.QualifiedSegmentCount if state.Segment.Count >= hydrogenSegmentMinSamples { qualified++ } switch { case state.SampleCount < 2: return "NO_DATA", "有效车载氢量样本不足2条" case state.EligibleIntervalCount == 0: return "NO_DATA", "无燃料电池工作状态下的有效压力区间" case qualified == 0: return "NO_DATA", "有效工作段关键帧不足10条" case state.AbnormalDropCount > 0: return "SUSPECT", "存在超过阈值的异常下降" default: return "OK", "" } } func hydrogenStreamCanStartSegment(sample HydrogenStreamSample) bool { return !sample.FuelCellStateKnown || sample.FuelCellActive } func hydrogenStreamIntervalEligible(previousActive, previousKnown, currentActive, currentKnown bool) bool { if previousKnown || currentKnown { return previousKnown && previousActive && currentKnown && currentActive } return true } func hydrogenStreamRapidRecovery(previousTime time.Time, previousPressure float64, sample HydrogenStreamSample) bool { elapsed := sample.EventTime.Sub(previousTime) return elapsed > 0 && elapsed <= hydrogenRecoveryWindow && sample.PressureMPa-previousPressure >= hydrogenRecoveryPressure } func hydrogenStreamEndpointDrop(first, tail []hydrogenStreamEndpoint) float64 { if len(first) < hydrogenSegmentWindowSize || len(tail) < hydrogenSegmentWindowSize { return 0 } startMass := hydrogenStreamMedianMass(first) endMass := hydrogenStreamMedianMass(tail) noise := defaultHydrogenNoiseKG for _, endpoint := range first { noise = math.Max(noise, endpoint.NoiseKG) } for _, endpoint := range tail { noise = math.Max(noise, endpoint.NoiseKG) } drop := startMass - endMass if drop <= noise { return 0 } return drop } func hydrogenStreamMedianMass(values []hydrogenStreamEndpoint) float64 { masses := make([]float64, len(values)) for index, value := range values { masses[index] = value.MassKG } sort.Float64s(masses) return masses[len(masses)/2] } func roundHydrogenKG(value float64) float64 { return math.Round(value*1000) / 1000 } func firstNonEmptyHydrogen(values ...string) string { for _, value := range values { if strings.TrimSpace(value) != "" { return strings.TrimSpace(value) } } return "" } 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 !validHydrogenPressureTemperature(pressure, temperature) { return HydrogenStreamSample{}, "invalid_pressure_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) fuelCellActive, fuelCellStateKnown := FuelCellActiveModeFromEnvelope(env) if !fuelCellStateKnown { if currentA, currentOK := firstHydrogenNumber(env.Fields, hydrogenCurrentFieldKeys); currentOK && currentA >= 0 { fuelCellActive = currentA > hydrogenActiveCurrentA fuelCellStateKnown = true } } 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, FuelCellActive: fuelCellActive, FuelCellStateKnown: fuelCellStateKnown, ConsumptionEligible: !fuelCellStateKnown || fuelCellActive, }, "", true } 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 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() state, found, err := selectHydrogenSegmentStreamState(ctx, tx, sample.VIN, sample.Date) if err != nil { return result, err } if found { if !state.add(sample) { result.Duplicate = 1 return result, tx.Commit() } if err := updateHydrogenSegmentStreamState(ctx, tx, sample.VIN, sample.Date, state); err != nil { return result, err } } else { baseline, baselineFound, err := selectHydrogenStreamBaseline(ctx, tx, sample.VIN, sample.Date) if err != nil { return result, err } if !baselineFound { baseline = nil } state = newHydrogenSegmentStreamState(sample, baseline) if err := insertHydrogenSegmentStreamState(ctx, tx, sample.VIN, sample.Date, state); err != nil { return result, err } } if err := projectHydrogenSegmentStreamDaily(ctx, tx, sample.VIN, sample.Date, state); err != nil { return result, err } if err := tx.Commit(); err != nil { return result, err } result.Written = 1 return result, nil } func selectHydrogenSegmentStreamState(ctx context.Context, tx *sql.Tx, vin, date string) (hydrogenSegmentStreamState, bool, error) { var encoded []byte err := tx.QueryRowContext(ctx, selectHydrogenSegmentStreamStateSQL, vin, date).Scan(&encoded) if errors.Is(err, sql.ErrNoRows) { return hydrogenSegmentStreamState{}, false, nil } if err != nil { return hydrogenSegmentStreamState{}, false, err } var state hydrogenSegmentStreamState if err := json.Unmarshal(encoded, &state); err != nil { return hydrogenSegmentStreamState{}, false, err } return state, true, nil } func selectHydrogenStreamBaseline(ctx context.Context, tx *sql.Tx, vin, date string) (*hydrogenStreamBaseline, bool, error) { var baseline hydrogenStreamBaseline var firstMass, lastMass sql.NullFloat64 err := tx.QueryRowContext(ctx, selectHydrogenStreamBaselineSQL, vin, date).Scan( &baseline.SourceEndpoint, &baseline.ConsumptionKG, &firstMass, &lastMass, &baseline.SampleCount, &baseline.RefuelCount, &baseline.QualityStatus, ) if errors.Is(err, sql.ErrNoRows) { return nil, false, nil } if err != nil { return nil, false, err } baseline.FirstMassKG = -1 baseline.LastMassKG = -1 if firstMass.Valid { baseline.FirstMassKG = firstMass.Float64 } if lastMass.Valid { baseline.LastMassKG = lastMass.Float64 } return &baseline, true, nil } func insertHydrogenSegmentStreamState(ctx context.Context, tx *sql.Tx, vin, date string, state hydrogenSegmentStreamState) error { encoded, err := json.Marshal(state) if err != nil { return err } qualityStatus, qualityReason := state.quality() _, err = tx.ExecContext(ctx, insertHydrogenSegmentStreamStateSQL, vin, date, state.SourceEndpoint, state.FinalizedConsumptionKG, state.projectedConsumptionKG(), state.SampleCount, state.RefuelCount, state.AbnormalDropCount, state.EligibleIntervalCount, state.QualifiedSegmentCount, state.RemainingMassKG, state.LastEventTime, state.LastEventID, encoded, qualityStatus, qualityReason, ) return err } func updateHydrogenSegmentStreamState(ctx context.Context, tx *sql.Tx, vin, date string, state hydrogenSegmentStreamState) error { encoded, err := json.Marshal(state) if err != nil { return err } qualityStatus, qualityReason := state.quality() _, err = tx.ExecContext(ctx, updateHydrogenSegmentStreamStateSQL, state.SourceEndpoint, state.FinalizedConsumptionKG, state.projectedConsumptionKG(), state.SampleCount, state.RefuelCount, state.AbnormalDropCount, state.EligibleIntervalCount, state.QualifiedSegmentCount, state.RemainingMassKG, state.LastEventTime, state.LastEventID, encoded, qualityStatus, qualityReason, vin, date, ) return err } func projectHydrogenSegmentStreamDaily(ctx context.Context, tx *sql.Tx, vin, date string, state hydrogenSegmentStreamState) error { qualityStatus, qualityReason := state.quality() _, err := tx.ExecContext(ctx, projectHydrogenSegmentStreamDailySQL, vin, date, state.SourceEndpoint, state.projectedConsumptionKG(), state.FirstMassKG, state.RemainingMassKG, state.SampleCount, state.RefuelCount, qualityStatus, qualityReason, ) return err } const selectHydrogenSegmentStreamStateSQL = `SELECT state_json FROM vehicle_open_hydrogen_segment_stream_state WHERE vin=? AND stat_date=? FOR UPDATE` const selectHydrogenStreamBaselineSQL = `SELECT source_endpoint,consumption_kg,first_mass_kg,last_mass_kg, sample_count,refuel_count,quality_status FROM vehicle_open_daily_energy WHERE vin=? AND stat_date=? AND energy_type='HYDROGEN' LIMIT 1` const insertHydrogenSegmentStreamStateSQL = `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_SEGMENT_MEDIAN_5',?,?)` const updateHydrogenSegmentStreamStateSQL = `UPDATE vehicle_open_hydrogen_segment_stream_state SET 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='PRESSURE_NIST_SEGMENT_MEDIAN_5',quality_status=?,quality_reason=? WHERE vin=? AND stat_date=?` const projectHydrogenSegmentStreamDailySQL = `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 ) VALUES(?,?,'HYDROGEN',?,?,'kg',?,?,?,?,?,?,NOW(3)) 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)`