diff --git a/.changeset/fail-closed-mpc-physics.md b/.changeset/fail-closed-mpc-physics.md new file mode 100644 index 00000000..97cd15ff --- /dev/null +++ b/.changeset/fail-closed-mpc-physics.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Reject invalid MPC physics before either solver and keep the prior plan. Use one PV uncertainty snapshot per replan, and do not use the Go fallback while a battery is recovering into its operating SoC range. diff --git a/go/internal/mpc/params_validation.go b/go/internal/mpc/params_validation.go new file mode 100644 index 00000000..f424e50a --- /dev/null +++ b/go/internal/mpc/params_validation.go @@ -0,0 +1,412 @@ +package mpc + +import ( + "fmt" + "math" +) + +func validatePlanningSlots(slots []Slot) error { + if len(slots) == 0 { + return fmt.Errorf("slots must be non-empty") + } + for i, slot := range slots { + field := fmt.Sprintf("slots[%d]", i) + for _, value := range []struct { + name string + v float64 + }{ + {"price_ore", slot.PriceOre}, + {"spot_ore", slot.SpotOre}, + {"pv_w", slot.PVW}, + {"load_w", slot.LoadW}, + {"confidence", slot.Confidence}, + {"max_import_w", slot.Limits.MaxImportW}, + {"max_export_w", slot.Limits.MaxExportW}, + } { + if !finite(value.v) { + return fmt.Errorf("%s.%s must be finite", field, value.name) + } + } + if slot.PVW > 0 { + return fmt.Errorf("%s.pv_w must be non-positive in the site sign convention", field) + } + if slot.LoadW < 0 { + return fmt.Errorf("%s.load_w must be non-negative in the site sign convention", field) + } + if slot.Confidence <= 0 || slot.Confidence > 1 { + return fmt.Errorf("%s.confidence must be within (0, 1]", field) + } + if slot.Limits.MaxImportW < 0 || slot.Limits.MaxExportW < 0 { + return fmt.Errorf("%s grid limits must be non-negative", field) + } + } + return nil +} + +func validateBatteryFleetMembers(fleet []BatteryFleetMember) error { + seen := make(map[string]struct{}, len(fleet)) + for i, battery := range fleet { + field := fmt.Sprintf("battery_fleet[%d]", i) + if battery.Driver == "" { + return fmt.Errorf("%s.driver must be non-empty", field) + } + if _, exists := seen[battery.Driver]; exists { + return fmt.Errorf("%s.driver %q is duplicated", field, battery.Driver) + } + seen[battery.Driver] = struct{}{} + if err := requirePositivePlanningValue(field+".capacity_wh", battery.CapacityWh); err != nil { + return err + } + if err := requireNonNegativePlanningValue(field+".max_charge_w", battery.MaxChargeW); err != nil { + return err + } + if err := requireNonNegativePlanningValue(field+".max_discharge_w", battery.MaxDischargeW); err != nil { + return err + } + } + return nil +} + +// planningParamsRequireRecovery marks states the external model can replay but +// the discrete Go DP cannot yet represent. The DP snaps its initial SoC onto +// the operating grid, so using it here would plan from energy the site does not +// have (or discard energy it does have). +func planningParamsRequireRecovery(p Params) bool { + if p.InitialSoCPct < p.SoCMinPct || p.InitialSoCPct > p.SoCMaxPct { + return true + } + for _, storage := range p.Storages { + if storage.InitialEnergyWh < storage.MinEnergyWh || storage.InitialEnergyWh > storage.MaxEnergyWh { + return true + } + } + return false +} + +// validatePlanningParams checks the effective inputs that both planning +// engines receive. Service calls it after live battery and loadpoint state has +// been applied, so an invalid value cannot reach either the external optimizer +// or the Go fallback with different defaulting or failure semantics. +func validatePlanningParams(p Params) error { + switch p.Mode { + case ModeSelfConsumption, ModeCheapCharge, ModePassiveArbitrage, ModeArbitrage: + default: + return fmt.Errorf("unsupported mode %q", p.Mode) + } + + if p.SoCLevels < 3 { + return fmt.Errorf("soc_levels must be at least 3, got %d", p.SoCLevels) + } + if p.ActionLevels < 3 { + return fmt.Errorf("action_levels must be at least 3, got %d", p.ActionLevels) + } + if err := requirePositivePlanningValue("capacity_wh", p.CapacityWh); err != nil { + return err + } + if !finite(p.SoCMinPct) || !finite(p.SoCMaxPct) || + p.SoCMinPct < 0 || p.SoCMinPct >= p.SoCMaxPct || p.SoCMaxPct > 100 { + return fmt.Errorf("soc bounds must satisfy 0 <= min < max <= 100, got %.6g..%.6g", + p.SoCMinPct, p.SoCMaxPct) + } + // A live battery may start outside the configured operating band and + // recover toward it. Only the physical 0..100 percent range is hard here. + if !finite(p.InitialSoCPct) || p.InitialSoCPct < 0 || p.InitialSoCPct > 100 { + return fmt.Errorf("initial_soc_pct must be within 0..100, got %.6g", p.InitialSoCPct) + } + if err := requireNonNegativePlanningValue("max_charge_w", p.MaxChargeW); err != nil { + return err + } + if err := requireNonNegativePlanningValue("max_discharge_w", p.MaxDischargeW); err != nil { + return err + } + if err := requirePlanningEfficiency("charge_efficiency", p.ChargeEfficiency, false); err != nil { + return err + } + if err := requirePlanningEfficiency("discharge_efficiency", p.DischargeEfficiency, false); err != nil { + return err + } + + for _, value := range []struct { + name string + v float64 + }{ + {"terminal_soc_price", p.TerminalSoCPrice}, + {"export_ore_per_kwh", p.ExportOrePerKWh}, + {"export_bonus_ore_kwh", p.ExportBonusOreKwh}, + {"export_fee_ore_kwh", p.ExportFeeOreKwh}, + } { + if !finite(value.v) { + return fmt.Errorf("%s must be finite", value.name) + } + } + if p.ExportFloorOreKwh != nil && !finite(*p.ExportFloorOreKwh) { + return fmt.Errorf("export_floor_ore_kwh must be finite") + } + for _, value := range []struct { + name string + v float64 + }{ + {"pv_charge_bonus_ore_kwh", p.PVChargeBonusOreKwh}, + {"min_arbitrage_spread_ore_kwh", p.MinArbitrageSpreadOreKwh}, + {"pv_uncertainty_w", p.PVUncertaintyW}, + {"pv_forecast_safety_k", p.PVForecastSafetyK}, + } { + if err := requireNonNegativePlanningValue(value.name, value.v); err != nil { + return err + } + } + + assetIDs := make(map[string]string, len(p.Storages)+len(p.Loadpoints)+1) + if err := validateStorageSpecs(p, assetIDs); err != nil { + return err + } + storageIDs := make(map[string]string, len(assetIDs)) + for id, field := range assetIDs { + storageIDs[id] = field + } + if err := validateLoadpointSpecs(planningLoadpointSpecs(p), assetIDs); err != nil { + return err + } + activeList := activeLoadpointSpecs(p.Loadpoints) + if len(p.Loadpoints) > 0 && p.Loadpoint != nil && p.Loadpoint.active() { + // The external planner consumes Loadpoints while the Go fallback consumes + // Loadpoint. Validate the fallback on its own so a stale copy cannot hide + // behind a valid list and reach only the fallback engine. + if err := validateLoadpointSpecs([]*LoadpointSpec{p.Loadpoint}, storageIDs); err != nil { + return fmt.Errorf("loadpoint fallback: %w", err) + } + } + if len(activeList) > 0 && !planningLoadpointsEquivalent(p.Loadpoint, activeList[0]) { + return fmt.Errorf("loadpoint fallback must match first active loadpoint %q", activeList[0].ID) + } + return nil +} + +func validateStorageSpecs(p Params, assetIDs map[string]string) error { + var totalCapacityWh, totalInitialWh, totalMinWh, totalMaxWh float64 + var totalChargeW, totalDischargeW float64 + for i, storage := range p.Storages { + field := fmt.Sprintf("storages[%d]", i) + if storage.ID == "" { + return fmt.Errorf("%s.id must be non-empty", field) + } + if previous, exists := assetIDs[storage.ID]; exists { + return fmt.Errorf("%s.id %q duplicates %s", field, storage.ID, previous) + } + assetIDs[storage.ID] = field + if err := requirePositivePlanningValue(field+".capacity_wh", storage.CapacityWh); err != nil { + return err + } + if !finite(storage.MinEnergyWh) || !finite(storage.MaxEnergyWh) || + storage.MinEnergyWh < 0 || storage.MinEnergyWh > storage.MaxEnergyWh || + storage.MaxEnergyWh > storage.CapacityWh { + return fmt.Errorf("%s energy bounds must satisfy 0 <= min <= max <= capacity", field) + } + // Initial energy outside the operating band is recoverable, but energy + // outside the physical battery is not. + if !finite(storage.InitialEnergyWh) || storage.InitialEnergyWh < 0 || + storage.InitialEnergyWh > storage.CapacityWh { + return fmt.Errorf("%s.initial_energy_wh must be within 0..capacity", field) + } + if err := requireNonNegativePlanningValue(field+".max_charge_w", storage.MaxChargeW); err != nil { + return err + } + if err := requireNonNegativePlanningValue(field+".max_discharge_w", storage.MaxDischargeW); err != nil { + return err + } + if err := requirePlanningEfficiency(field+".charge_efficiency", storage.ChargeEfficiency, false); err != nil { + return err + } + if err := requirePlanningEfficiency(field+".discharge_efficiency", storage.DischargeEfficiency, false); err != nil { + return err + } + if !planningValuesEqual(storage.ChargeEfficiency, p.ChargeEfficiency) || + !planningValuesEqual(storage.DischargeEfficiency, p.DischargeEfficiency) { + return fmt.Errorf("%s efficiencies must match aggregate fallback efficiencies", field) + } + + totalCapacityWh += storage.CapacityWh + totalInitialWh += storage.InitialEnergyWh + totalMinWh += storage.MinEnergyWh + totalMaxWh += storage.MaxEnergyWh + totalChargeW += storage.MaxChargeW + totalDischargeW += storage.MaxDischargeW + } + if len(p.Storages) == 0 { + return nil + } + + energyToleranceWh := math.Max(1, p.CapacityWh*0.0002) + checks := []struct { + name string + got float64 + want float64 + tol float64 + }{ + {"capacity_wh", totalCapacityWh, p.CapacityWh, 1}, + {"initial_energy_wh", totalInitialWh, p.CapacityWh * p.InitialSoCPct / 100, energyToleranceWh}, + {"min_energy_wh", totalMinWh, p.CapacityWh * p.SoCMinPct / 100, energyToleranceWh}, + {"max_energy_wh", totalMaxWh, p.CapacityWh * p.SoCMaxPct / 100, energyToleranceWh}, + {"max_charge_w", totalChargeW, p.MaxChargeW, 2}, + {"max_discharge_w", totalDischargeW, p.MaxDischargeW, 2}, + } + for _, check := range checks { + if math.Abs(check.got-check.want) > check.tol { + return fmt.Errorf("storage aggregate %s %.6g does not match fallback %.6g", + check.name, check.got, check.want) + } + } + return nil +} + +func planningLoadpointSpecs(p Params) []*LoadpointSpec { + if len(p.Loadpoints) > 0 { + return p.Loadpoints + } + if p.Loadpoint != nil { + return []*LoadpointSpec{p.Loadpoint} + } + return nil +} + +func activeLoadpointSpecs(loadpoints []*LoadpointSpec) []*LoadpointSpec { + active := make([]*LoadpointSpec, 0, len(loadpoints)) + for _, loadpoint := range loadpoints { + if loadpoint.active() { + active = append(active, loadpoint) + } + } + return active +} + +func planningLoadpointsEquivalent(fallback, primary *LoadpointSpec) bool { + if fallback == nil || primary == nil || fallback.PluggedIn != primary.PluggedIn || + fallback.ID != primary.ID || fallback.Levels != primary.Levels || + fallback.TargetSlotIdx != primary.TargetSlotIdx || + fallback.SurplusOnly != primary.SurplusOnly || + fallback.NoBatteryToEV != primary.NoBatteryToEV { + return false + } + fallbackMin, fallbackMax := planningLoadpointBounds(fallback) + primaryMin, primaryMax := planningLoadpointBounds(primary) + for _, values := range [][2]float64{ + {fallback.CapacityWh, primary.CapacityWh}, + {fallbackMin, primaryMin}, + {fallbackMax, primaryMax}, + {fallback.InitialSoCPct, primary.InitialSoCPct}, + {fallback.TargetSoCPct, primary.TargetSoCPct}, + {fallback.MaxChargeW, primary.MaxChargeW}, + {planningLoadpointEfficiency(fallback), planningLoadpointEfficiency(primary)}, + } { + if !planningValuesEqual(values[0], values[1]) { + return false + } + } + fallbackSteps := fallback.normalizedSteps() + primarySteps := primary.normalizedSteps() + if len(fallbackSteps) != len(primarySteps) { + return false + } + for i := range fallbackSteps { + if !planningValuesEqual(fallbackSteps[i], primarySteps[i]) { + return false + } + } + return true +} + +func planningLoadpointBounds(loadpoint *LoadpointSpec) (float64, float64) { + minPct, maxPct := loadpoint.MinPct, loadpoint.MaxPct + if minPct == 0 && maxPct == 0 { + maxPct = 100 + } + return minPct, maxPct +} + +func planningLoadpointEfficiency(loadpoint *LoadpointSpec) float64 { + if loadpoint.ChargeEfficiency == 0 { + return 0.9 + } + return loadpoint.ChargeEfficiency +} + +func validateLoadpointSpecs(loadpoints []*LoadpointSpec, assetIDs map[string]string) error { + for i, loadpoint := range loadpoints { + if loadpoint == nil || !loadpoint.PluggedIn { + continue + } + field := fmt.Sprintf("loadpoints[%d]", i) + if loadpoint.ID == "" { + return fmt.Errorf("%s.id must be non-empty", field) + } + if previous, exists := assetIDs[loadpoint.ID]; exists { + return fmt.Errorf("%s.id %q duplicates %s", field, loadpoint.ID, previous) + } + assetIDs[loadpoint.ID] = field + if err := requirePositivePlanningValue(field+".capacity_wh", loadpoint.CapacityWh); err != nil { + return err + } + if loadpoint.Levels < 2 { + return fmt.Errorf("%s.levels must be at least 2, got %d", field, loadpoint.Levels) + } + minPct, maxPct := planningLoadpointBounds(loadpoint) + if !finite(minPct) || !finite(maxPct) || + minPct < 0 || minPct >= maxPct || maxPct > 100 { + return fmt.Errorf("%s SoC bounds must satisfy 0 <= min < max <= 100", field) + } + if !finite(loadpoint.InitialSoCPct) || loadpoint.InitialSoCPct < minPct || + loadpoint.InitialSoCPct > maxPct { + return fmt.Errorf("%s.initial_soc_pct must be within the loadpoint SoC bounds", field) + } + if !finite(loadpoint.TargetSoCPct) || loadpoint.TargetSoCPct < 0 || + loadpoint.TargetSoCPct > 100 || + (loadpoint.TargetSoCPct != 0 && + (loadpoint.TargetSoCPct < minPct || loadpoint.TargetSoCPct > maxPct)) { + return fmt.Errorf("%s.target_soc_pct must be zero or within the loadpoint SoC bounds", field) + } + if err := requireNonNegativePlanningValue(field+".max_charge_w", loadpoint.MaxChargeW); err != nil { + return err + } + if err := requirePlanningEfficiency(field+".charge_efficiency", loadpoint.ChargeEfficiency, true); err != nil { + return err + } + for stepIdx, stepW := range loadpoint.AllowedStepsW { + if !finite(stepW) || stepW < 0 { + return fmt.Errorf("%s.allowed_steps_w[%d] must be finite and non-negative", field, stepIdx) + } + if loadpoint.MaxChargeW > 0 && stepW > loadpoint.MaxChargeW { + return fmt.Errorf("%s.allowed_steps_w[%d] exceeds max_charge_w", field, stepIdx) + } + } + } + return nil +} + +func requirePositivePlanningValue(name string, value float64) error { + if !finite(value) || value <= 0 { + return fmt.Errorf("%s must be finite and greater than zero", name) + } + return nil +} + +func requireNonNegativePlanningValue(name string, value float64) error { + if !finite(value) || value < 0 { + return fmt.Errorf("%s must be finite and non-negative", name) + } + return nil +} + +func requirePlanningEfficiency(name string, value float64, allowDefaultZero bool) error { + if !finite(value) || value < 0 || value > 1 || (!allowDefaultZero && value == 0) { + if allowDefaultZero { + return fmt.Errorf("%s must be zero or within (0, 1]", name) + } + return fmt.Errorf("%s must be within (0, 1]", name) + } + return nil +} + +func planningValuesEqual(a, b float64) bool { + tolerance := math.Max(1e-6, math.Max(math.Abs(a), math.Abs(b))*1e-9) + return math.Abs(a-b) <= tolerance +} diff --git a/go/internal/mpc/params_validation_test.go b/go/internal/mpc/params_validation_test.go new file mode 100644 index 00000000..234cdaf2 --- /dev/null +++ b/go/internal/mpc/params_validation_test.go @@ -0,0 +1,707 @@ +package mpc + +import ( + "context" + "errors" + "math" + "path/filepath" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func validPlanningParams() Params { + p := baseParams(ModePassiveArbitrage) + p.PVChargeBonusOreKwh = 10 + p.ExportBonusOreKwh = 5 + p.ExportFeeOreKwh = 2 + p.MinArbitrageSpreadOreKwh = 20 + p.PVUncertaintyW = 300 + p.PVForecastSafetyK = 1 + return p +} + +func validPlanningStorageParams() Params { + p := validPlanningParams() + p.Storages = []StorageAssetSpec{{ + ID: "battery", CapacityWh: p.CapacityWh, + InitialEnergyWh: p.CapacityWh * p.InitialSoCPct / 100, + MinEnergyWh: p.CapacityWh * p.SoCMinPct / 100, + MaxEnergyWh: p.CapacityWh * p.SoCMaxPct / 100, + MaxChargeW: p.MaxChargeW, MaxDischargeW: p.MaxDischargeW, + ChargeEfficiency: p.ChargeEfficiency, DischargeEfficiency: p.DischargeEfficiency, + }} + return p +} + +func validPlanningLoadpoint() *LoadpointSpec { + return &LoadpointSpec{ + ID: "garage", CapacityWh: 60000, Levels: 11, + MinPct: 0, MaxPct: 100, InitialSoCPct: 20, PluggedIn: true, + TargetSoCPct: 80, TargetSlotIdx: 8, + MaxChargeW: 11000, AllowedStepsW: []float64{0, 1400, 4100, 11000}, + ChargeEfficiency: 0.9, + } +} + +func TestValidatePlanningParamsAcceptsSupportedPhysicalStates(t *testing.T) { + modes := []Mode{ModeSelfConsumption, ModeCheapCharge, ModePassiveArbitrage, ModeArbitrage} + for _, mode := range modes { + p := validPlanningParams() + p.Mode = mode + if err := validatePlanningParams(p); err != nil { + t.Fatalf("mode %q rejected: %v", mode, err) + } + } + + tests := []struct { + name string + mutate func(*Params) + }{ + {"physical minimum and maximum", func(p *Params) { p.SoCMinPct, p.SoCMaxPct = 0, 100 }}, + {"below operating minimum recovery", func(p *Params) { p.InitialSoCPct = 5 }}, + {"above operating maximum recovery", func(p *Params) { p.InitialSoCPct = 97 }}, + {"disabled power", func(p *Params) { p.MaxChargeW, p.MaxDischargeW = 0, 0 }}, + {"ideal efficiency", func(p *Params) { p.ChargeEfficiency, p.DischargeEfficiency = 1, 1 }}, + {"finite negative prices", func(p *Params) { + p.TerminalSoCPrice, p.ExportOrePerKWh, p.ExportBonusOreKwh, p.ExportFeeOreKwh = -10, -20, -30, -40 + }}, + {"storage below operating minimum recovery", func(p *Params) { + *p = validPlanningStorageParams() + p.InitialSoCPct = 5 + p.Storages[0].InitialEnergyWh = 500 + }}, + {"storage above operating maximum recovery", func(p *Params) { + *p = validPlanningStorageParams() + p.InitialSoCPct = 97 + p.Storages[0].InitialEnergyWh = 9700 + }}, + {"loadpoint documented default efficiency", func(p *Params) { + lp := validPlanningLoadpoint() + lp.MinPct, lp.MaxPct = 0, 0 + lp.ChargeEfficiency = 0 + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = lp + }}, + {"loadpoint zero max shorthand", func(p *Params) { + lp := validPlanningLoadpoint() + lp.MaxChargeW = 0 + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = lp + }}, + {"loadpoint zero target sentinel", func(p *Params) { + lp := validPlanningLoadpoint() + lp.MinPct, lp.TargetSoCPct = 10, 0 + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = lp + }}, + {"equivalent normalized loadpoint fallback", func(p *Params) { + lp := validPlanningLoadpoint() + lp.MinPct, lp.MaxPct = 0, 0 + lp.ChargeEfficiency = 0 + lp.AllowedStepsW = []float64{4100, 0, 1400, 4100} + fallback := *lp + fallback.MinPct, fallback.MaxPct = 0, 100 + fallback.ChargeEfficiency = 0.9 + fallback.AllowedStepsW = []float64{1400, 4100} + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = &fallback + }}, + {"unplugged invalid loadpoint ignored", func(p *Params) { + p.Loadpoints = []*LoadpointSpec{{PluggedIn: false, CapacityWh: math.NaN()}} + }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + p := validPlanningParams() + tc.mutate(&p) + if err := validatePlanningParams(p); err != nil { + t.Fatalf("valid state rejected: %v", err) + } + }) + } +} + +func TestValidatePlanningParamsRejectsInvalidAggregateValues(t *testing.T) { + floorNaN := math.NaN() + tests := []struct { + name string + want string + mutate func(*Params) + }{ + {"empty mode", "unsupported mode", func(p *Params) { p.Mode = "" }}, + {"unknown mode", "unsupported mode", func(p *Params) { p.Mode = "fast" }}, + {"soc levels", "soc_levels", func(p *Params) { p.SoCLevels = 2 }}, + {"action levels", "action_levels", func(p *Params) { p.ActionLevels = 2 }}, + {"zero capacity", "capacity_wh", func(p *Params) { p.CapacityWh = 0 }}, + {"nan capacity", "capacity_wh", func(p *Params) { p.CapacityWh = math.NaN() }}, + {"negative soc minimum", "soc bounds", func(p *Params) { p.SoCMinPct = -1 }}, + {"equal soc bounds", "soc bounds", func(p *Params) { p.SoCMaxPct = p.SoCMinPct }}, + {"reversed soc bounds", "soc bounds", func(p *Params) { p.SoCMinPct = 96 }}, + {"soc maximum above physical", "soc bounds", func(p *Params) { p.SoCMaxPct = 101 }}, + {"nan soc bound", "soc bounds", func(p *Params) { p.SoCMinPct = math.NaN() }}, + {"initial soc below physical", "initial_soc_pct", func(p *Params) { p.InitialSoCPct = -0.1 }}, + {"initial soc above physical", "initial_soc_pct", func(p *Params) { p.InitialSoCPct = 100.1 }}, + {"initial soc infinite", "initial_soc_pct", func(p *Params) { p.InitialSoCPct = math.Inf(1) }}, + {"negative charge power", "max_charge_w", func(p *Params) { p.MaxChargeW = -1 }}, + {"infinite discharge power", "max_discharge_w", func(p *Params) { p.MaxDischargeW = math.Inf(1) }}, + {"zero charge efficiency", "charge_efficiency", func(p *Params) { p.ChargeEfficiency = 0 }}, + {"high charge efficiency", "charge_efficiency", func(p *Params) { p.ChargeEfficiency = 1.01 }}, + {"nan discharge efficiency", "discharge_efficiency", func(p *Params) { p.DischargeEfficiency = math.NaN() }}, + {"nan terminal price", "terminal_soc_price", func(p *Params) { p.TerminalSoCPrice = math.NaN() }}, + {"infinite export flat", "export_ore_per_kwh", func(p *Params) { p.ExportOrePerKWh = math.Inf(1) }}, + {"nan export bonus", "export_bonus_ore_kwh", func(p *Params) { p.ExportBonusOreKwh = math.NaN() }}, + {"infinite export fee", "export_fee_ore_kwh", func(p *Params) { p.ExportFeeOreKwh = math.Inf(-1) }}, + {"nan export floor", "export_floor_ore_kwh", func(p *Params) { p.ExportFloorOreKwh = &floorNaN }}, + {"negative pv bonus", "pv_charge_bonus", func(p *Params) { p.PVChargeBonusOreKwh = -1 }}, + {"negative spread", "min_arbitrage_spread", func(p *Params) { p.MinArbitrageSpreadOreKwh = -1 }}, + {"negative uncertainty", "pv_uncertainty", func(p *Params) { p.PVUncertaintyW = -1 }}, + {"nan safety", "pv_forecast_safety", func(p *Params) { p.PVForecastSafetyK = math.NaN() }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + p := validPlanningParams() + tc.mutate(&p) + err := validatePlanningParams(p) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("error = %v, want path %q", err, tc.want) + } + }) + } +} + +func TestValidatePlanningParamsRejectsInvalidStoragePhysics(t *testing.T) { + tests := []struct { + name string + want string + mutate func(*Params) + }{ + {"empty id", ".id", func(p *Params) { p.Storages[0].ID = "" }}, + {"duplicate id", "duplicates", func(p *Params) { p.Storages = append(p.Storages, p.Storages[0]) }}, + {"zero capacity", ".capacity_wh", func(p *Params) { p.Storages[0].CapacityWh = 0 }}, + {"nan capacity", ".capacity_wh", func(p *Params) { p.Storages[0].CapacityWh = math.NaN() }}, + {"negative minimum", "energy bounds", func(p *Params) { p.Storages[0].MinEnergyWh = -1 }}, + {"nan minimum", "energy bounds", func(p *Params) { p.Storages[0].MinEnergyWh = math.NaN() }}, + {"infinite maximum", "energy bounds", func(p *Params) { p.Storages[0].MaxEnergyWh = math.Inf(1) }}, + {"maximum above capacity", "energy bounds", func(p *Params) { p.Storages[0].MaxEnergyWh = 10001 }}, + {"initial below physical", "initial_energy_wh", func(p *Params) { p.Storages[0].InitialEnergyWh = -1 }}, + {"nan initial", "initial_energy_wh", func(p *Params) { p.Storages[0].InitialEnergyWh = math.NaN() }}, + {"initial above physical", "initial_energy_wh", func(p *Params) { p.Storages[0].InitialEnergyWh = 10001 }}, + {"negative charge power", ".max_charge_w", func(p *Params) { p.Storages[0].MaxChargeW = -1 }}, + {"nan discharge power", ".max_discharge_w", func(p *Params) { p.Storages[0].MaxDischargeW = math.NaN() }}, + {"zero efficiency", ".charge_efficiency", func(p *Params) { p.Storages[0].ChargeEfficiency = 0 }}, + {"nan efficiency", ".charge_efficiency", func(p *Params) { p.Storages[0].ChargeEfficiency = math.NaN() }}, + {"high efficiency", ".discharge_efficiency", func(p *Params) { p.Storages[0].DischargeEfficiency = 1.01 }}, + {"different fallback efficiency", "fallback efficiencies", func(p *Params) { p.Storages[0].ChargeEfficiency = 0.9 }}, + {"capacity aggregate mismatch", "aggregate capacity", func(p *Params) { p.Storages[0].CapacityWh += 10 }}, + {"initial aggregate mismatch", "aggregate initial", func(p *Params) { p.Storages[0].InitialEnergyWh += 10 }}, + {"minimum aggregate mismatch", "aggregate min", func(p *Params) { p.Storages[0].MinEnergyWh += 10 }}, + {"maximum aggregate mismatch", "aggregate max", func(p *Params) { p.Storages[0].MaxEnergyWh -= 10 }}, + {"charge aggregate mismatch", "aggregate max_charge", func(p *Params) { p.Storages[0].MaxChargeW -= 3 }}, + {"discharge aggregate mismatch", "aggregate max_discharge", func(p *Params) { p.Storages[0].MaxDischargeW -= 3 }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + p := validPlanningStorageParams() + tc.mutate(&p) + err := validatePlanningParams(p) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("error = %v, want path %q", err, tc.want) + } + }) + } +} + +func TestValidatePlanningParamsRejectsInvalidLoadpointPhysics(t *testing.T) { + tests := []struct { + name string + want string + mutate func(*Params) + }{ + {"empty id", ".id", func(p *Params) { p.Loadpoints[0].ID = "" }}, + {"duplicate id", "duplicates", func(p *Params) { p.Loadpoints = append(p.Loadpoints, validPlanningLoadpoint()) }}, + {"zero capacity", ".capacity_wh", func(p *Params) { p.Loadpoints[0].CapacityWh = 0 }}, + {"nan capacity", ".capacity_wh", func(p *Params) { p.Loadpoints[0].CapacityWh = math.NaN() }}, + {"too few levels", ".levels", func(p *Params) { p.Loadpoints[0].Levels = 1 }}, + {"invalid bounds", "SoC bounds", func(p *Params) { p.Loadpoints[0].MinPct, p.Loadpoints[0].MaxPct = 60, 50 }}, + {"nan bound", "SoC bounds", func(p *Params) { p.Loadpoints[0].MinPct = math.NaN() }}, + {"initial outside bounds", "initial_soc_pct", func(p *Params) { p.Loadpoints[0].InitialSoCPct = 101 }}, + {"initial below operating minimum", "initial_soc_pct", func(p *Params) { + p.Loadpoints[0].MinPct, p.Loadpoints[0].InitialSoCPct = 30, 20 + }}, + {"initial above operating maximum", "initial_soc_pct", func(p *Params) { + p.Loadpoints[0].MaxPct, p.Loadpoints[0].InitialSoCPct = 10, 20 + }}, + {"nan initial", "initial_soc_pct", func(p *Params) { p.Loadpoints[0].InitialSoCPct = math.NaN() }}, + {"target outside bounds", "target_soc_pct", func(p *Params) { p.Loadpoints[0].TargetSoCPct = 101 }}, + {"target below operating minimum", "target_soc_pct", func(p *Params) { + p.Loadpoints[0].MinPct, p.Loadpoints[0].InitialSoCPct, p.Loadpoints[0].TargetSoCPct = 30, 40, 20 + }}, + {"target above operating maximum", "target_soc_pct", func(p *Params) { + p.Loadpoints[0].MaxPct, p.Loadpoints[0].TargetSoCPct = 70, 80 + }}, + {"infinite target", "target_soc_pct", func(p *Params) { p.Loadpoints[0].TargetSoCPct = math.Inf(1) }}, + {"negative max power", ".max_charge_w", func(p *Params) { p.Loadpoints[0].MaxChargeW = -1 }}, + {"infinite max power", ".max_charge_w", func(p *Params) { p.Loadpoints[0].MaxChargeW = math.Inf(1) }}, + {"negative efficiency", ".charge_efficiency", func(p *Params) { p.Loadpoints[0].ChargeEfficiency = -0.1 }}, + {"nan efficiency", ".charge_efficiency", func(p *Params) { p.Loadpoints[0].ChargeEfficiency = math.NaN() }}, + {"high efficiency", ".charge_efficiency", func(p *Params) { p.Loadpoints[0].ChargeEfficiency = 1.01 }}, + {"nan step", "allowed_steps", func(p *Params) { p.Loadpoints[0].AllowedStepsW = []float64{0, math.NaN()} }}, + {"negative step", "allowed_steps", func(p *Params) { p.Loadpoints[0].AllowedStepsW = []float64{0, -1} }}, + {"step above max", "exceeds max_charge_w", func(p *Params) { p.Loadpoints[0].AllowedStepsW = []float64{0, 12000} }}, + {"fallback different physics", "fallback must match", func(p *Params) { + fallback := *p.Loadpoint + fallback.MaxChargeW = 12000 + p.Loadpoint = &fallback + }}, + {"fallback invalid physics", "loadpoint fallback", func(p *Params) { + fallback := *p.Loadpoint + fallback.InitialSoCPct = math.NaN() + p.Loadpoint = &fallback + }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + p := validPlanningParams() + lp := validPlanningLoadpoint() + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = lp + tc.mutate(&p) + err := validatePlanningParams(p) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("error = %v, want path %q", err, tc.want) + } + }) + } + + p := validPlanningStorageParams() + lp := validPlanningLoadpoint() + lp.ID = p.Storages[0].ID + p.Loadpoints = []*LoadpointSpec{lp} + p.Loadpoint = lp + if err := validatePlanningParams(p); err == nil || !strings.Contains(err.Error(), "duplicates") { + t.Fatalf("cross-asset duplicate error = %v", err) + } +} + +func TestValidatePlanningSlotsAndFleet(t *testing.T) { + validSlot := Slot{StartMs: 1, LenMin: 15, PriceOre: -100, SpotOre: -200, PVW: -500, LoadW: 1000, Confidence: 1} + if err := validatePlanningSlots([]Slot{validSlot}); err != nil { + t.Fatalf("valid slot rejected: %v", err) + } + if err := validateBatteryFleetMembers([]BatteryFleetMember{{ + Driver: "battery", CapacityWh: 10000, MaxChargeW: 0, MaxDischargeW: 0, + }}); err != nil { + t.Fatalf("valid disabled fleet member rejected: %v", err) + } + + slotTests := []struct { + name string + want string + mutate func(*Slot) + }{ + {"nan price", "price_ore", func(s *Slot) { s.PriceOre = math.NaN() }}, + {"infinite spot", "spot_ore", func(s *Slot) { s.SpotOre = math.Inf(1) }}, + {"positive pv", "pv_w", func(s *Slot) { s.PVW = 1 }}, + {"negative load", "load_w", func(s *Slot) { s.LoadW = -1 }}, + {"zero confidence", "confidence", func(s *Slot) { s.Confidence = 0 }}, + {"high confidence", "confidence", func(s *Slot) { s.Confidence = 1.1 }}, + {"negative import limit", "grid limits", func(s *Slot) { s.Limits.MaxImportW = -1 }}, + {"nan export limit", "max_export_w", func(s *Slot) { s.Limits.MaxExportW = math.NaN() }}, + } + for _, tc := range slotTests { + t.Run("slot "+tc.name, func(t *testing.T) { + slot := validSlot + tc.mutate(&slot) + err := validatePlanningSlots([]Slot{slot}) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("error = %v, want path %q", err, tc.want) + } + }) + } + + fleetTests := []struct { + name string + fleet []BatteryFleetMember + want string + }{ + {"empty driver", []BatteryFleetMember{{CapacityWh: 1}}, ".driver"}, + {"duplicate driver", []BatteryFleetMember{{Driver: "a", CapacityWh: 1}, {Driver: "a", CapacityWh: 1}}, "duplicated"}, + {"zero capacity", []BatteryFleetMember{{Driver: "a", CapacityWh: 0}}, ".capacity_wh"}, + {"nan capacity", []BatteryFleetMember{{Driver: "a", CapacityWh: math.NaN()}}, ".capacity_wh"}, + {"negative charge", []BatteryFleetMember{{Driver: "a", CapacityWh: 1, MaxChargeW: -1}}, ".max_charge_w"}, + {"infinite discharge", []BatteryFleetMember{{Driver: "a", CapacityWh: 1, MaxDischargeW: math.Inf(1)}}, ".max_discharge_w"}, + } + for _, tc := range fleetTests { + t.Run("fleet "+tc.name, func(t *testing.T) { + err := validateBatteryFleetMembers(tc.fleet) + if err == nil || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("error = %v, want path %q", err, tc.want) + } + }) + } +} + +type physicsGateCountingOptimizer struct { + calls atomic.Int32 +} + +func (o *physicsGateCountingOptimizer) Optimize(_ context.Context, slots []Slot, p Params) (Plan, error) { + o.calls.Add(1) + plan := Optimize(slots, p) + plan.Solver = &SolverInfo{Engine: "test", Backend: "counting", Status: "optimal"} + return plan, nil +} + +func (*physicsGateCountingOptimizer) Close() error { return nil } + +type physicsGateRecoveryOptimizer struct { + calls atomic.Int32 +} + +func (o *physicsGateRecoveryOptimizer) Optimize(_ context.Context, slots []Slot, p Params) (Plan, error) { + o.calls.Add(1) + plan := Plan{ + GeneratedAtMs: time.Now().UnixMilli(), Mode: p.Mode, + HorizonSlots: len(slots), CapacityWh: p.CapacityWh, + InitialSoCPct: p.InitialSoCPct, + Actions: make([]Action, len(slots)), + Solver: &SolverInfo{Engine: "test", Backend: "recovery", Status: "optimal"}, + } + for i, slot := range slots { + gridW := slot.LoadW + slot.PVW + cost := SlotGridCostOre(slot, gridW*float64(slot.LenMin)/60/1000, p) + plan.Actions[i] = Action{ + SlotStartMs: slot.StartMs, SlotLenMin: slot.LenMin, + PriceOre: slot.PriceOre, SpotOre: slot.SpotOre, + PVW: slot.PVW, LoadW: slot.LoadW, Confidence: slot.Confidence, + GridW: gridW, SoCPct: p.InitialSoCPct, CostOre: cost, + } + plan.TotalCostOre += cost + } + return plan, nil +} + +func (*physicsGateRecoveryOptimizer) Close() error { return nil } + +type physicsGateFailingOptimizer struct { + calls atomic.Int32 +} + +func (o *physicsGateFailingOptimizer) Optimize(context.Context, []Slot, Params) (Plan, error) { + o.calls.Add(1) + return Plan{}, errors.New("primary unavailable") +} + +func (*physicsGateFailingOptimizer) Close() error { return nil } + +func configurePhysicsGateFleet(svc *Service, firstSoCPct, secondSoCPct float64) { + svc.Tele = telemetry.NewStore() + svc.BatteryFleet = []BatteryFleetMember{ + {Driver: "battery-a", CapacityWh: 5000, MaxChargeW: 1500, MaxDischargeW: 1500}, + {Driver: "battery-b", CapacityWh: 5000, MaxChargeW: 1500, MaxDischargeW: 1500}, + } + for i, socPct := range []float64{firstSoCPct, secondSoCPct} { + driver := svc.BatteryFleet[i].Driver + soc := socPct / 100 + svc.Tele.Update(driver, telemetry.DerBattery, 0, &soc, nil) + svc.Tele.DriverHealthMut(driver).RecordSuccess() + } +} + +func TestReplanRejectsInvalidPhysicsBeforeSolver(t *testing.T) { + tests := []struct { + name string + mutate func(*Service) + }{ + {"unknown mode", func(s *Service) { s.Defaults.Mode = "fast" }}, + {"nan capacity", func(s *Service) { s.Defaults.CapacityWh = math.NaN() }}, + {"efficiency above one", func(s *Service) { s.Defaults.ChargeEfficiency = 1.01 }}, + {"reversed soc bounds", func(s *Service) { s.Defaults.SoCMinPct = 96 }}, + {"negative power", func(s *Service) { s.Defaults.MaxDischargeW = -1 }}, + {"nan pv uncertainty", func(s *Service) { + s.PVUncertaintyW = func() float64 { return math.NaN() } + }}, + {"negative pv uncertainty", func(s *Service) { + s.PVUncertaintyW = func() float64 { return -1 } + }}, + {"negative load input", func(s *Service) { + s.Load = func(time.Time) float64 { return -1 } + }}, + {"nan load input", func(s *Service) { + s.Load = func(time.Time) float64 { return math.NaN() } + }}, + {"invalid plugged loadpoint", func(s *Service) { + s.Loadpoints = func(int) []*LoadpointSpec { + return []*LoadpointSpec{{ID: "garage", PluggedIn: true, CapacityWh: 0, Levels: 1}} + } + }}, + {"loadpoint initial outside operating bounds", func(s *Service) { + s.Loadpoints = func(int) []*LoadpointSpec { + lp := validPlanningLoadpoint() + lp.MinPct, lp.InitialSoCPct = 30, 20 + return []*LoadpointSpec{lp} + } + }}, + {"loadpoint target outside operating bounds", func(s *Service) { + s.Loadpoints = func(int) []*LoadpointSpec { + lp := validPlanningLoadpoint() + lp.MaxPct, lp.TargetSoCPct = 70, 80 + return []*LoadpointSpec{lp} + } + }}, + {"stale loadpoint fallback", func(s *Service) { + primary := validPlanningLoadpoint() + fallback := *primary + fallback.MaxChargeW = 12000 + s.Defaults.Loadpoints = []*LoadpointSpec{primary} + s.Defaults.Loadpoint = &fallback + }}, + {"duplicate fleet", func(s *Service) { + s.Tele = telemetry.NewStore() + soc := 0.5 + s.Tele.Update("battery", telemetry.DerBattery, 0, &soc, nil) + s.Tele.DriverHealthMut("battery").RecordSuccess() + s.BatteryFleet = []BatteryFleetMember{ + {Driver: "battery", CapacityWh: 10000, MaxChargeW: 3000, MaxDischargeW: 3000}, + {Driver: "battery", CapacityWh: 10000, MaxChargeW: 3000, MaxDischargeW: 3000}, + } + }}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + st, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + now := time.Now().UTC().Truncate(time.Minute) + if err := st.SavePrices([]state.PricePoint{{ + Zone: "SE3", SlotTsMs: now.Add(-5 * time.Minute).UnixMilli(), SlotLenMin: 60, + SpotOreKwh: 50, TotalOreKwh: 100, Source: "test", FetchedAtMs: now.UnixMilli(), + }}); err != nil { + t.Fatal(err) + } + + optimizer := &physicsGateCountingOptimizer{} + svc := New(st, nil, "SE3", Params{ + Mode: ModeSelfConsumption, SoCLevels: 11, ActionLevels: 5, + CapacityWh: 10000, SoCMinPct: 10, SoCMaxPct: 95, InitialSoCPct: 50, + MaxChargeW: 3000, MaxDischargeW: 3000, + ChargeEfficiency: 0.95, DischargeEfficiency: 0.95, + }) + svc.BaseLoad = 500 + svc.Optimizer = optimizer + var idCalls, saveCalls atomic.Int32 + svc.decisionIDFactory = func() string { + idCalls.Add(1) + return "00000000-0000-4000-8000-000000000099" + } + svc.SaveDiag = func(*Diagnostic, string) error { + saveCalls.Add(1) + return nil + } + + accepted := svc.Replan(context.Background()) + if accepted == nil || accepted.DecisionID == "" { + t.Fatalf("valid baseline plan = %+v", accepted) + } + tc.mutate(svc) + got := svc.Replan(context.Background()) + if got != accepted || svc.Latest() != accepted { + t.Fatalf("invalid inputs replaced prior plan: got=%p accepted=%p latest=%p", got, accepted, svc.Latest()) + } + if optimizer.calls.Load() != 1 || idCalls.Load() != 1 || saveCalls.Load() != 1 { + t.Fatalf("calls after rejection: optimizer=%d id=%d save=%d, want 1/1/1", + optimizer.calls.Load(), idCalls.Load(), saveCalls.Load()) + } + if d := svc.Diagnose(); d == nil || d.DecisionID != accepted.DecisionID { + t.Fatalf("diagnostic changed after rejected inputs: %+v", d) + } + }) + } +} + +func TestReplanRejectsInvalidPhysicsBeforeGoDP(t *testing.T) { + st, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + now := time.Now().UTC().Truncate(time.Minute) + if err := st.SavePrices([]state.PricePoint{{ + Zone: "SE3", SlotTsMs: now.Add(-5 * time.Minute).UnixMilli(), SlotLenMin: 60, + SpotOreKwh: 50, TotalOreKwh: 100, Source: "test", FetchedAtMs: now.UnixMilli(), + }}); err != nil { + t.Fatal(err) + } + + svc := New(st, nil, "SE3", validPlanningParams()) + svc.BaseLoad = 500 + var idCalls, saveCalls atomic.Int32 + svc.decisionIDFactory = func() string { + idCalls.Add(1) + return "00000000-0000-4000-8000-000000000097" + } + svc.SaveDiag = func(*Diagnostic, string) error { + saveCalls.Add(1) + return nil + } + accepted := svc.Replan(context.Background()) + if accepted == nil { + t.Fatal("valid Go-DP baseline returned nil") + } + svc.Defaults.DischargeEfficiency = 1.01 + if got := svc.Replan(context.Background()); got != accepted || svc.Latest() != accepted { + t.Fatalf("invalid Go-DP inputs replaced prior plan: got=%p accepted=%p latest=%p", got, accepted, svc.Latest()) + } + if idCalls.Load() != 1 || saveCalls.Load() != 1 { + t.Fatalf("calls after rejection: id=%d save=%d, want 1/1", idCalls.Load(), saveCalls.Load()) + } +} + +func TestReplanRecoveryUsesExternalPlannerWithoutDPDerivedResults(t *testing.T) { + tests := []struct { + name string + setup func(*Service) + }{ + {"aggregate recovery", func(svc *Service) { svc.Defaults.InitialSoCPct = 5 }}, + {"storage member recovery", func(svc *Service) { configurePhysicsGateFleet(svc, 5, 55) }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + st, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + now := time.Now().UTC().Truncate(time.Minute) + if err := st.SavePrices([]state.PricePoint{{ + Zone: "SE3", SlotTsMs: now.Add(-5 * time.Minute).UnixMilli(), SlotLenMin: 60, + SpotOreKwh: 50, TotalOreKwh: 100, Source: "test", FetchedAtMs: now.UnixMilli(), + }}); err != nil { + t.Fatal(err) + } + + p := validPlanningParams() + p.Mode = ModeArbitrage + optimizer := &physicsGateRecoveryOptimizer{} + svc := New(st, nil, "SE3", p) + svc.BaseLoad = 500 + svc.Optimizer = optimizer + tc.setup(svc) + plan := svc.Replan(context.Background()) + if plan == nil || plan.Solver == nil || plan.Solver.Backend != "recovery" { + t.Fatalf("external recovery plan = %+v", plan) + } + if optimizer.calls.Load() != 1 { + t.Fatalf("external optimizer calls = %d, want 1", optimizer.calls.Load()) + } + if plan.DPEvaluationShadow != nil || plan.DPShadow != nil || plan.Baselines != nil { + t.Fatalf("DP-derived results must be absent during recovery: evaluation=%+v downside=%+v baselines=%+v", + plan.DPEvaluationShadow, plan.DPShadow, plan.Baselines) + } + }) + } +} + +func TestReplanRecoveryKeepsPreviousPlanWhenDPWouldBeRequired(t *testing.T) { + tests := []struct { + name string + baseline PlanOptimizer + recovery PlanOptimizer + wantFailCalls int32 + setupRecovery func(*Service) + }{ + {name: "aggregate go only", setupRecovery: func(svc *Service) { svc.Defaults.InitialSoCPct = 5 }}, + {name: "aggregate failed primary", baseline: &physicsGateCountingOptimizer{}, recovery: &physicsGateFailingOptimizer{}, wantFailCalls: 1, + setupRecovery: func(svc *Service) { svc.Defaults.InitialSoCPct = 5 }}, + {name: "storage member go only", setupRecovery: func(svc *Service) { configurePhysicsGateFleet(svc, 5, 55) }}, + {name: "storage member failed primary", baseline: &physicsGateCountingOptimizer{}, recovery: &physicsGateFailingOptimizer{}, wantFailCalls: 1, + setupRecovery: func(svc *Service) { configurePhysicsGateFleet(svc, 5, 55) }}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + st, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + now := time.Now().UTC().Truncate(time.Minute) + if err := st.SavePrices([]state.PricePoint{{ + Zone: "SE3", SlotTsMs: now.Add(-5 * time.Minute).UnixMilli(), SlotLenMin: 60, + SpotOreKwh: 50, TotalOreKwh: 100, Source: "test", FetchedAtMs: now.UnixMilli(), + }}); err != nil { + t.Fatal(err) + } + + p := validPlanningParams() + svc := New(st, nil, "SE3", p) + svc.BaseLoad = 500 + svc.Optimizer = tc.baseline + var idCalls, saveCalls atomic.Int32 + svc.decisionIDFactory = func() string { + idCalls.Add(1) + return "00000000-0000-4000-8000-000000000098" + } + svc.SaveDiag = func(*Diagnostic, string) error { + saveCalls.Add(1) + return nil + } + accepted := svc.Replan(context.Background()) + if accepted == nil || accepted.DecisionID == "" { + t.Fatalf("baseline plan = %+v", accepted) + } + + tc.setupRecovery(svc) + svc.Optimizer = tc.recovery + got := svc.Replan(context.Background()) + if got != accepted || svc.Latest() != accepted { + t.Fatalf("recovery replaced prior plan through Go DP: got=%p accepted=%p latest=%p", + got, accepted, svc.Latest()) + } + if idCalls.Load() != 1 || saveCalls.Load() != 1 { + t.Fatalf("calls after recovery rejection: id=%d save=%d, want 1/1", idCalls.Load(), saveCalls.Load()) + } + if failing, ok := tc.recovery.(*physicsGateFailingOptimizer); ok && failing.calls.Load() != tc.wantFailCalls { + t.Fatalf("failing optimizer calls = %d, want %d", failing.calls.Load(), tc.wantFailCalls) + } + }) + } +} + +func TestReplanSnapshotsPVUncertaintyOnce(t *testing.T) { + st, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + now := time.Now().UTC().Truncate(time.Minute) + if err := st.SavePrices([]state.PricePoint{{ + Zone: "SE3", SlotTsMs: now.Add(-5 * time.Minute).UnixMilli(), SlotLenMin: 60, + SpotOreKwh: 50, TotalOreKwh: 100, Source: "test", FetchedAtMs: now.UnixMilli(), + }}); err != nil { + t.Fatal(err) + } + + svc := New(st, nil, "SE3", validPlanningParams()) + svc.BaseLoad = 500 + svc.PVForecastSafetyK = 1 + svc.Optimizer = &physicsGateCountingOptimizer{} + var calls atomic.Int32 + svc.PVUncertaintyW = func() float64 { + calls.Add(1) + return 300 + } + if plan := svc.Replan(context.Background()); plan == nil { + t.Fatal("Replan returned nil") + } + if calls.Load() != 1 { + t.Fatalf("PV uncertainty sampled %d times, want exactly 1", calls.Load()) + } +} diff --git a/go/internal/mpc/service.go b/go/internal/mpc/service.go index ddecf402..4a5e55ac 100644 --- a/go/internal/mpc/service.go +++ b/go/internal/mpc/service.go @@ -1112,7 +1112,14 @@ func (s *Service) runReplan(request replanRequest) *Plan { // a separate downside copy for the emergency Go-DP path, preserving the // previous safety behavior if the worker is unavailable. fallbackSlots := append([]Slot(nil), slots...) - s.applyPVDownsideToSlots(fallbackSlots) + var pvUncertaintyW float64 + pvUncertainty := s.PVUncertaintyW + if pvUncertainty != nil { + // One replan must use one uncertainty snapshot. Reading the live model + // twice could give the external scenarios and Go fallback different + // physics for the same request. + pvUncertaintyW = pvUncertainty() + } // Plumb the site fuse + export ceiling into per-slot limits so the DP // joint-plans battery + EV under the grid constraints instead of @@ -1128,13 +1135,20 @@ func (s *Service) runReplan(request replanRequest) *Plan { clampSlotGridLimits(fallbackSlots, s.FuseMaxW, s.MaxExportW) p := request.params + if p.Mode == "" { + p.Mode = ModeSelfConsumption + } fleet := request.fleet if len(fleet) > 0 { + if err := validateBatteryFleetMembers(fleet); err != nil { + slog.Error("mpc: invalid optimization parameters; keeping previous plan", "err", err) + return s.Latest() + } var ok bool p, ok = s.onlineFleetParams(p, fleet) if !ok { slog.Warn("mpc: no online battery capacity with SoC — keeping previous plan") - return nil + return s.Latest() } } else { p.InitialSoCPct = currentSoCPct(s.Tele, p.InitialSoCPct) @@ -1149,9 +1163,10 @@ func (s *Service) runReplan(request replanRequest) *Plan { p.MinArbitrageSpreadOreKwh = s.MinArbitrageSpreadOreKwh p.ExportFloorOreKwh = s.ExportFloorOreKwh p.PVForecastSafetyK = s.PVForecastSafetyK - if s.PVUncertaintyW != nil { - p.PVUncertaintyW = math.Max(0, s.PVUncertaintyW()) + if pvUncertainty != nil { + p.PVUncertaintyW = pvUncertaintyW } + applyPVDownside(fallbackSlots, p.PVForecastSafetyK, p.PVUncertaintyW) // Default terminal valuation. Mode-dependent because self-consumption // is a constrained game: the battery can only offset local load, not @@ -1205,6 +1220,13 @@ func (s *Service) runReplan(request replanRequest) *Plan { p.Loadpoints = []*LoadpointSpec{spec} } } + // Check the probe output before activeLoadpoints filters it. A plugged-in + // spec with broken physics must stop planning, not disappear and make the + // service silently solve a different battery-only problem. + if err := validateLoadpointSpecs(planningLoadpointSpecs(p), make(map[string]string)); err != nil { + slog.Error("mpc: invalid optimization parameters; keeping previous plan", "err", err) + return s.Latest() + } active := p.activeLoadpoints() if len(active) > 0 { p.Loadpoints = active @@ -1230,6 +1252,26 @@ func (s *Service) runReplan(request replanRequest) *Plan { p.TerminalSoCPrice = selfConsumptionTerminalPrice(prices, s.ExportBonusOreKwh, s.ExportFeeOreKwh) } + if err := validatePlanningSlots(slots); err != nil { + slog.Error("mpc: invalid optimization inputs; keeping previous plan", "basis", "primary", "err", err) + return s.Latest() + } + if err := validatePlanningSlots(fallbackSlots); err != nil { + slog.Error("mpc: invalid optimization inputs; keeping previous plan", "basis", "go-fallback", "err", err) + return s.Latest() + } + if err := validatePlanningParams(p); err != nil { + slog.Error("mpc: invalid optimization parameters; keeping previous plan", "err", err) + return s.Latest() + } + recoveryRequired := planningParamsRequireRecovery(p) + if recoveryRequired && s.Optimizer == nil { + slog.Error("mpc: battery state requires operating-bound recovery that Go DP cannot model; keeping previous plan", + "soc_start", p.InitialSoCPct, + "soc_min", p.SoCMinPct, + "soc_max", p.SoCMaxPct) + return s.Latest() + } slog.Info("mpc: optimize params", "mode", p.Mode, @@ -1263,34 +1305,56 @@ func (s *Service) runReplan(request replanRequest) *Plan { return s.canceledReplan(request, "primary-solve") } if err == nil { - dpEvaluation := Optimize(slots, p) - dpEvaluation.Solver = &SolverInfo{ - Engine: "go-dp", Backend: "bellman", Status: "optimal-grid", - Formulation: "discrete-dp", - } - candidate.DPEvaluationShadow = compareDPShadow(candidate, dpEvaluation) - candidate.DPEvaluationShadow.ForecastBasis = "same base forecast input" - candidate.DPEvaluationShadow.Solver = dpEvaluation.Solver - candidate.DPEvaluationShadow.TotalCostOre = dpEvaluation.TotalCostOre - candidate.DPEvaluationShadow.ActiveMinusShadowOre = candidate.TotalCostOre - dpEvaluation.TotalCostOre - if candidate.DPEvaluationShadow.FirstAction != nil { - mode, _, _ := actionToSlot(*candidate.DPEvaluationShadow.FirstAction, p.Mode) - candidate.DPEvaluationShadow.FirstAction.EMSMode = mode - } + if recoveryRequired { + candidate.DPEvaluationShadow = nil + candidate.DPShadow = nil + candidate.Baselines = nil + slog.Info("mpc: skipping Go DP shadows while battery state recovers into operating bounds", + "soc_start", p.InitialSoCPct, + "soc_min", p.SoCMinPct, + "soc_max", p.SoCMaxPct) + } else { + dpEvaluation := Optimize(slots, p) + dpEvaluation.Solver = &SolverInfo{ + Engine: "go-dp", Backend: "bellman", Status: "optimal-grid", + Formulation: "discrete-dp", + } + candidate.DPEvaluationShadow = compareDPShadow(candidate, dpEvaluation) + candidate.DPEvaluationShadow.ForecastBasis = "same base forecast input" + candidate.DPEvaluationShadow.Solver = dpEvaluation.Solver + candidate.DPEvaluationShadow.TotalCostOre = dpEvaluation.TotalCostOre + candidate.DPEvaluationShadow.ActiveMinusShadowOre = candidate.TotalCostOre - dpEvaluation.TotalCostOre + if candidate.DPEvaluationShadow.FirstAction != nil { + mode, _, _ := actionToSlot(*candidate.DPEvaluationShadow.FirstAction, p.Mode) + candidate.DPEvaluationShadow.FirstAction.EMSMode = mode + } - dpShadow := Optimize(fallbackSlots, p) - dpShadow.Solver = &SolverInfo{ - Engine: "go-dp", Backend: "bellman", Status: "optimal-grid", - Formulation: "discrete-dp", - } - candidate.DPShadow = compareDPShadow(candidate, dpShadow) - candidate.DPShadow.ForecastBasis = "downside-pv fallback input" - candidate.DPShadow.Solver = dpShadow.Solver - candidate.DPShadow.TotalCostOre = dpShadow.TotalCostOre - candidate.DPShadow.ActiveMinusShadowOre = candidate.TotalCostOre - dpShadow.TotalCostOre - if candidate.DPShadow.FirstAction != nil { - mode, _, _ := actionToSlot(*candidate.DPShadow.FirstAction, p.Mode) - candidate.DPShadow.FirstAction.EMSMode = mode + dpShadow := Optimize(fallbackSlots, p) + dpShadow.Solver = &SolverInfo{ + Engine: "go-dp", Backend: "bellman", Status: "optimal-grid", + Formulation: "discrete-dp", + } + candidate.DPShadow = compareDPShadow(candidate, dpShadow) + candidate.DPShadow.ForecastBasis = "downside-pv fallback input" + candidate.DPShadow.Solver = dpShadow.Solver + candidate.DPShadow.TotalCostOre = dpShadow.TotalCostOre + candidate.DPShadow.ActiveMinusShadowOre = candidate.TotalCostOre - dpShadow.TotalCostOre + if candidate.DPShadow.FirstAction != nil { + mode, _, _ := actionToSlot(*candidate.DPShadow.FirstAction, p.Mode) + candidate.DPShadow.FirstAction.EMSMode = mode + } + + optimizerSolveMs := 0.0 + if candidate.Solver != nil { + optimizerSolveMs = candidate.Solver.SolveMs + } + slog.Info("mpc: active optimizer vs DP shadow", + "optimizer_cost_ore", candidate.TotalCostOre, + "dp_evaluation_cost_ore", dpEvaluation.TotalCostOre, + "active_minus_evaluation_ore", candidate.DPEvaluationShadow.ActiveMinusShadowOre, + "dp_shadow_cost_ore", dpShadow.TotalCostOre, + "active_minus_shadow_ore", candidate.DPShadow.ActiveMinusShadowOre, + "optimizer_solve_ms", optimizerSolveMs) } if request.wasCanceledByService() { return s.canceledReplan(request, "dp-shadow") @@ -1343,17 +1407,6 @@ func (s *Service) runReplan(request replanRequest) *Plan { } } } - optimizerSolveMs := 0.0 - if candidate.Solver != nil { - optimizerSolveMs = candidate.Solver.SolveMs - } - slog.Info("mpc: active optimizer vs DP shadow", - "optimizer_cost_ore", candidate.TotalCostOre, - "dp_evaluation_cost_ore", dpEvaluation.TotalCostOre, - "active_minus_evaluation_ore", candidate.DPEvaluationShadow.ActiveMinusShadowOre, - "dp_shadow_cost_ore", dpShadow.TotalCostOre, - "active_minus_shadow_ore", candidate.DPShadow.ActiveMinusShadowOre, - "optimizer_solve_ms", optimizerSolveMs) if candidate.RecourseShadow != nil { slog.Info("mpc: champion vs stochastic shadow", "champion_cost_ore", candidate.TotalCostOre, @@ -1366,6 +1419,14 @@ func (s *Service) runReplan(request replanRequest) *Plan { if request.wasCanceledByService() { return s.canceledReplan(request, "primary-fallback") } + if recoveryRequired { + slog.Error("mpc: primary optimizer failed and Go DP cannot model operating-bound recovery; keeping previous plan", + "err", err, + "soc_start", p.InitialSoCPct, + "soc_min", p.SoCMinPct, + "soc_max", p.SoCMaxPct) + return s.Latest() + } slog.Error("mpc: primary optimizer failed; using Go DP fallback", "err", err) slots = fallbackSlots plan = Optimize(slots, p) @@ -1397,7 +1458,7 @@ func (s *Service) runReplan(request replanRequest) *Plan { // self-consumption mode: the SC baseline is the plan itself, which // makes the badge trivially zero and distracts from the price // signal. For SC runs the UI still has the plan cost on its own. - if p.Mode != ModeSelfConsumption { + if p.Mode != ModeSelfConsumption && !recoveryRequired { bl := ComputeBaselines(slots, p) plan.Baselines = &bl }