diff --git a/.changeset/forecast-source-from-errors.md b/.changeset/forecast-source-from-errors.md new file mode 100644 index 000000000..ca24a969e --- /dev/null +++ b/.changeset/forecast-source-from-errors.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +The planner now takes each forecast signal from the source that measured better over the last week, once it has three days of scored hours. Before, it used Energyplan solar even while that model was new, and kept the older load model even when that model overshot by 2 kW at night. diff --git a/docs/architecture.md b/docs/architecture.md index 4f316eb7d..f0a681e46 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -243,9 +243,12 @@ a separate versioned contract. At the start of each replan, Core freezes the legacy forecast, weather, occupancy and saved model state. It calls the forecast worker once under a deadline, outside control and dispatch locks. Core accepts PV and load independently for each covered interval. If either signal is -missing, late, partial or invalid, Core retains the matching legacy value. The +missing, late, partial or invalid, Core retains the matching legacy value. +When a week of scored errors shows one source clearly better for a signal, +Core uses that source; otherwise the worker's quality label decides. The resulting `champion` can therefore contain Energyplan PV with legacy load, or -the reverse. `legacy_shadow` keeps both legacy signals from the same frozen +the reverse. Each slot's planning margin comes from the errors its own sources +made on earlier issues. `legacy_shadow` keeps both legacy signals from the same frozen capture for a fair later comparison. Complete qualified 15-minute observations update the local models outside diff --git a/docs/energyplan-contract.md b/docs/energyplan-contract.md index 152b041bd..69b80497b 100644 --- a/docs/energyplan-contract.md +++ b/docs/energyplan-contract.md @@ -204,6 +204,8 @@ Predictions state whether a signal is known, its quality, coverage and uncertainty. Provisional model bounds are not calibrated quantiles. Missing, late, incomplete or rejected predictions trigger per-signal Core fallback; current Core also keeps legacy load during Energyplan load cold start. +When a week of paired, scored errors shows one source clearly better for a +signal, Core uses that source instead of the quality label. Core records which signal supplied each planner input. Energyplan returns complete opaque model state. Core validates and stores an diff --git a/go/cmd/ftw/forecast_load_risk_test.go b/go/cmd/ftw/forecast_load_risk_test.go index 969aa3108..0dbd67eb1 100644 --- a/go/cmd/ftw/forecast_load_risk_test.go +++ b/go/cmd/ftw/forecast_load_risk_test.go @@ -72,12 +72,16 @@ func TestForecastRiskUsesCalibratedLoadWithoutPVAndJointErrorsOnce(t *testing.T) band := forecasting.Band{LowW: 0, HighW: 5000, Method: forecasting.BandMethodColdStart} p := forecasting.Point{StartMS: at.UnixMilli(), EndMS: end.UnixMilli(), LoadW: 1000, PVKnown: pvKnown, LoadKnown: true, PVQuality: "test", LoadQuality: "test", PVBand: band, LoadBand: band, NetBand: band} - f.errors = append(f.errors, forecasting.ErrorSample{Series: "champion", ConfigVersion: f.site().Revision, IssueID: "history", - OriginMS: issued.UnixMilli(), IssuedAtMS: issued.UnixMilli(), StartMS: p.StartMS, EndMS: p.EndMS, - AvailableAtMS: end.UnixMilli(), Lead: forecasting.LeadBucket(issued.UnixMilli(), p.StartMS), - LoadErrorW: 200, PVKnown: pvKnown, LoadKnown: true, Prediction: p}) - if err := f.errors[len(f.errors)-1].Validate(); err != nil { - t.Fatal(err) + // The slot plans with Energyplan load and legacy PV; both + // series scored the same issue. + for _, series := range []string{"energyplan", "legacy_shadow"} { + f.errors = append(f.errors, forecasting.ErrorSample{Series: series, ConfigVersion: f.site().Revision, IssueID: "history", + OriginMS: issued.UnixMilli(), IssuedAtMS: issued.UnixMilli(), StartMS: p.StartMS, EndMS: p.EndMS, + AvailableAtMS: end.UnixMilli(), Lead: forecasting.LeadBucket(issued.UnixMilli(), p.StartMS), + LoadErrorW: 200, PVKnown: pvKnown, LoadKnown: true, Prediction: p}) + if err := f.errors[len(f.errors)-1].Validate(); err != nil { + t.Fatal(err) + } } } } diff --git a/go/cmd/ftw/forecast_site.go b/go/cmd/ftw/forecast_site.go index cde0343aa..47aab852e 100644 --- a/go/cmd/ftw/forecast_site.go +++ b/go/cmd/ftw/forecast_site.go @@ -51,7 +51,7 @@ const forecastIdentityReceiptKey = "forecast/live_identity_v1" // The evaluation cohort keys on it rather than on the Core version, so error // bands and baselines survive updates that leave forecasting alone. Bump it // whenever Core changes what reaches the planner. -const forecastPipelinePolicy = "energyplan-primary-v2" +const forecastPipelinePolicy = "energyplan-primary-v3" func newForecastSiteConfig(st *state.Store) *forecastSiteConfig { id, _ := st.LoadConfig("forecast/site_id") diff --git a/go/cmd/ftw/forecast_source_choice.go b/go/cmd/ftw/forecast_source_choice.go new file mode 100644 index 000000000..16cd0b5ae --- /dev/null +++ b/go/cmd/ftw/forecast_source_choice.go @@ -0,0 +1,115 @@ +package main + +import ( + "log/slog" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" +) + +// forecastSourcePick is one signal's source, chosen from paired errors of the +// Energyplan and legacy forecasts in the same issue. An empty Source leaves +// the choice to the worker's quality label. +type forecastSourcePick struct { + Source string + Samples, Hours, Days int + EnergyplanMAEW, LegacyMAEW float64 +} + +type forecastSourceChoice struct{ PV, Load forecastSourcePick } + +// One sunny or unusual day must not pick the source: require scored hours +// spread over several days and a clear gap. Evidence older than a week is +// dropped because both models keep learning. +const ( + forecastChoiceWindow = 7 * 24 * time.Hour + forecastChoiceMinHours = 48 + forecastChoiceMinDays = 3 + forecastChoiceMargin = 0.1 +) + +func chooseForecastSources(history []forecasting.ErrorSample, cohort string, origin int64) forecastSourceChoice { + since := origin - forecastChoiceWindow.Milliseconds() + recent := make([]forecasting.ErrorSample, 0, len(history)) + for _, e := range history { + if e.ConfigVersion == cohort && e.StartMS >= since && e.AvailableAtMS <= origin { + recent = append(recent, e) + } + } + return forecastSourceChoice{ + PV: pickForecastSource(forecasting.PoolFrozenSeries(recent, "energyplan", "legacy_shadow", "pv_daylight")), + Load: pickForecastSource(forecasting.PoolFrozenSeries(recent, "energyplan", "legacy_shadow", "load")), + } +} + +func pickForecastSource(m forecasting.PooledPairMetric) forecastSourcePick { + pick := forecastSourcePick{Samples: m.Samples, Hours: m.Hours, Days: m.Days, + EnergyplanMAEW: m.ChampionMAEW, LegacyMAEW: m.CandidateMAEW} + if pick.Hours < forecastChoiceMinHours || pick.Days < forecastChoiceMinDays { + return pick + } + switch { + case pick.EnergyplanMAEW < pick.LegacyMAEW*(1-forecastChoiceMargin): + pick.Source = "energyplan" + case pick.LegacyMAEW < pick.EnergyplanMAEW*(1-forecastChoiceMargin): + pick.Source = "legacy" + } + return pick +} + +// forecastSeriesOf names the archived series that holds a source's own +// forecast. +func forecastSeriesOf(source string) string { + if source == "energyplan" { + return "energyplan" + } + return "legacy_shadow" +} + +// forecastMixSeries names the forecast that takes PV from one series and load +// from another. +func forecastMixSeries(pv, load string) string { + if pv == load { + return pv + } + return "mix:" + pv + "+" + load +} + +// riskEvidence holds each source's own errors and both mixes of them, scored +// on the same issues. A slot's margin then describes the sources it plans +// with, whichever sources earlier plans used. +func riskEvidence(history []forecasting.ErrorSample) []forecasting.ErrorSample { + own := make([]forecasting.ErrorSample, 0, len(history)) + for _, e := range history { + if e.Series == "energyplan" || e.Series == "legacy_shadow" { + own = append(own, e) + } + } + out := append([]forecasting.ErrorSample(nil), own...) + for _, mix := range [][2]string{{"energyplan", "legacy_shadow"}, {"legacy_shadow", "energyplan"}} { + out = append(out, forecasting.ComposeFrozenSeries(own, mix[0], mix[1], forecastMixSeries(mix[0], mix[1]))...) + } + return out +} + +// noteSourceChoice logs when measured errors move a signal to another source. +func (f *forecastTracker) noteSourceChoice(c forecastSourceChoice) { + f.mu.Lock() + old := f.sourceChoice + f.sourceChoice = c + f.mu.Unlock() + if old.PV.Source == c.PV.Source && old.Load.Source == c.Load.Source { + return + } + for _, s := range []struct { + signal string + pick forecastSourcePick + }{{"pv", c.PV}, {"load", c.Load}} { + source := s.pick.Source + if source == "" { + source = "quality_rule" + } + slog.Info("forecast source chosen", "signal", s.signal, "source", source, "hours", s.pick.Hours, + "days", s.pick.Days, "energyplan_mae_w", int(s.pick.EnergyplanMAEW), "legacy_mae_w", int(s.pick.LegacyMAEW)) + } +} diff --git a/go/cmd/ftw/forecast_source_choice_test.go b/go/cmd/ftw/forecast_source_choice_test.go new file mode 100644 index 000000000..f0396393b --- /dev/null +++ b/go/cmd/ftw/forecast_source_choice_test.go @@ -0,0 +1,152 @@ +package main + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/mpc" +) + +type forecastPair struct{ energyplanPV, legacyPV, energyplanLoad, legacyLoad float64 } + +// scoredForecastErrors scores one series in every hour before end against +// 2,000 W PV and 1,000 W load, from issues made two hours ahead. +func scoredForecastErrors(t *testing.T, series, cohort string, end time.Time, hours int, pvW, loadW float64, pvSource, loadSource string) []forecasting.ErrorSample { + t.Helper() + band := forecasting.Band{LowW: 0, HighW: 5000, Method: forecasting.BandMethodColdStart} + var out []forecasting.ErrorSample + for h := 1; h <= hours; h++ { + at := end.Add(-time.Duration(h) * time.Hour) + issued := at.Add(-2 * time.Hour) + e := forecasting.ErrorSample{Series: series, ConfigVersion: cohort, IssueID: "issue-" + issued.Format(time.RFC3339), + OriginMS: issued.UnixMilli(), IssuedAtMS: issued.UnixMilli(), StartMS: at.UnixMilli(), EndMS: at.Add(15 * time.Minute).UnixMilli(), + AvailableAtMS: at.Add(15 * time.Minute).UnixMilli(), Lead: forecasting.LeadBucket(issued.UnixMilli(), at.UnixMilli()), + PVErrorW: 2000 - pvW, LoadErrorW: 1000 - loadW, PVKnown: true, Daylight: true, LoadKnown: true, + Prediction: forecasting.Point{StartMS: at.UnixMilli(), EndMS: at.Add(15 * time.Minute).UnixMilli(), PVW: pvW, LoadW: loadW, + PVKnown: true, LoadKnown: true, PVQuality: "test", LoadQuality: "test", PVSource: pvSource, LoadSource: loadSource, + PVBand: band, LoadBand: band, NetBand: band}} + if err := e.Validate(); err != nil { + t.Fatal(err) + } + out = append(out, e) + } + return out +} + +// pairedForecastErrors scores both sources of the same issues. +func pairedForecastErrors(t *testing.T, cohort string, end time.Time, hours int, p forecastPair) []forecasting.ErrorSample { + t.Helper() + return append(scoredForecastErrors(t, "energyplan", cohort, end, hours, p.energyplanPV, p.energyplanLoad, "energyplan", "energyplan"), + scoredForecastErrors(t, "legacy_shadow", cohort, end, hours, p.legacyPV, p.legacyLoad, "legacy", "legacy")...) +} + +func TestForecastSourceChoiceFollowsMeasuredErrors(t *testing.T) { + now := time.Date(2026, 10, 1, 12, 0, 0, 0, time.UTC) + // Home box, 22 Sep–1 Oct: Energyplan load and legacy PV were clearly better. + measured := forecastPair{energyplanPV: 1250, legacyPV: 1850, energyplanLoad: 1100, legacyLoad: 3000} + for _, tc := range []struct { + name string + history []forecasting.ErrorSample + pv, load string + }{ + {"enough evidence", pairedForecastErrors(t, "cfg", now, 72, measured), "legacy", "energyplan"}, + {"one day is not enough", pairedForecastErrors(t, "cfg", now, 24, measured), "", ""}, + {"small gap keeps the quality rule", pairedForecastErrors(t, "cfg", now, 72, + forecastPair{energyplanPV: 1950, legacyPV: 1952, energyplanLoad: 1100, legacyLoad: 1105}), "", ""}, + {"another cohort", pairedForecastErrors(t, "old", now, 72, measured), "", ""}, + {"older than a week", pairedForecastErrors(t, "cfg", now.Add(-8*24*time.Hour), 72, measured), "", ""}, + } { + t.Run(tc.name, func(t *testing.T) { + got := chooseForecastSources(tc.history, "cfg", now.UnixMilli()) + if got.PV.Source != tc.pv || got.Load.Source != tc.load { + t.Fatalf("got pv=%q load=%q, want pv=%q load=%q (%+v)", got.PV.Source, got.Load.Source, tc.pv, tc.load, got) + } + }) + } + got := chooseForecastSources(pairedForecastErrors(t, "cfg", now, 72, measured), "cfg", now.UnixMilli()) + if got.Load.Hours != 72 || got.Load.EnergyplanMAEW != 100 || got.Load.LegacyMAEW != 2000 || got.PV.LegacyMAEW != 150 { + t.Fatalf("evidence not reported: %+v", got) + } +} + +func TestForecastSourceChoiceSteersResolve(t *testing.T) { + at := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + // The worker says load is cold and PV ready; measured errors say otherwise. + f := primaryFixture(at, func(_ context.Context, p []byte) ([]byte, error) { + _, reply := hostForecastReply(p) + return json.Marshal(reply) + }) + site := hostForecastSite() + site.HasPVScale = true // legacy PV is known where it has weather + f.site = func() forecastSite { return site } + f.errors = pairedForecastErrors(t, site.Revision, at, 72, + forecastPair{energyplanPV: 1250, legacyPV: 1850, energyplanLoad: 1100, legacyLoad: 3000}) + in := f.Snapshot(at, trackerWeather(at, at)) + legacy := trackerSlots(at, 5) // weather covers the first hour only + got := in.Resolve(context.Background(), legacy) + if got[0].PVW != legacy[0].PVW || got[0].LoadW != 1100 { + t.Fatalf("measured errors did not choose the sources: %+v", got[0]) + } + if got[4].PVW != -100 { + t.Fatalf("legacy PV without weather replaced Energyplan PV: %+v", got[4]) + } + in.Record(got, got, "decision", at.UnixMilli()) + points := primarySeries(t, (<-f.queue).issue, "champion").Points + if p := points[0]; p.PVSource != "legacy" || p.LoadSource != "energyplan" || p.LoadQuality != "cold_start" || p.ModelLoad == nil { + t.Fatalf("chosen sources not recorded: %+v", p) + } + if points[4].PVSource != "energyplan" { + t.Fatalf("fallback to Energyplan PV not recorded: %+v", points[4]) + } + + // Learned Energyplan load still yields to a clearly better legacy load. + f = primaryFixture(at, primaryReply) + f.errors = pairedForecastErrors(t, f.site().Revision, at, 72, + forecastPair{energyplanPV: 2000, legacyPV: 2000, energyplanLoad: 3000, legacyLoad: 1100}) + in = f.Snapshot(at, trackerWeather(at, at)) + got = in.Resolve(context.Background(), legacy) + if got[0].LoadW != legacy[0].LoadW || got[0].PVW != -100 { + t.Fatalf("legacy load or default PV rule lost: %+v", got[0]) + } +} + +// Each slot's margin comes from the errors its own sources made, not from +// champion errors of sources the plan used before. +func TestForecastRiskUsesTheSlotsOwnSources(t *testing.T) { + at := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + for _, tc := range []struct { + name string + energyplanLoad float64 // predicted against a 1,000 W actual + legacyLoad float64 + championSource string + championLoad float64 + wantMargin float64 + }{ + // Both loads miss by 1 kW, so the quality rule picks learned + // Energyplan load. Its errors set the margin, not legacy's. + {"quality rule after legacy", 0, 2000, "legacy", 2000, 1000}, + // Legacy load wins by measurement. The margin uses its 300 W + // shortfall with Energyplan PV, not Energyplan's old overshoot. + {"measured legacy load", 3000, 700, "energyplan", 3000, 300}, + } { + t.Run(tc.name, func(t *testing.T) { + f := primaryFixture(at, primaryReply) + cohort := f.site().Revision + f.errors = append(pairedForecastErrors(t, cohort, at, 8*24, + forecastPair{energyplanPV: 2000, legacyPV: 2000, energyplanLoad: tc.energyplanLoad, legacyLoad: tc.legacyLoad}), + scoredForecastErrors(t, "champion", cohort, at, 8*24, 2000, tc.championLoad, tc.championSource, tc.championSource)...) + in := f.Snapshot(at, nil) + base := trackerSlots(at.Add(2*time.Hour), 1) + base[0].PVW, base[0].LoadW = 0, 1000 + base = in.Resolve(context.Background(), base) + planning := append([]mpc.Slot(nil), base...) + in.Risk(base, planning, 1) + if got := (planning[0].LoadW + planning[0].PVW) - (base[0].LoadW + base[0].PVW); got != tc.wantMargin { + t.Fatalf("margin %v W, want %v W: base %+v, planning %+v", got, tc.wantMargin, base[0], planning[0]) + } + }) + } +} diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index a17c3d043..086089af2 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -88,6 +88,7 @@ type forecastTracker struct { observationPVValid bool observationLoadValid bool issueArchiveError bool + sourceChoice forecastSourceChoice // guarded by mu } func (f *forecastTracker) Start(ctx context.Context) error { @@ -450,6 +451,9 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m frozen = fresh } calibrator := forecasting.NewCalibrator(history, site.Revision, origin.UnixMilli()) + riskCalibrator := forecasting.NewCalibrator(riskEvidence(history), site.Revision, origin.UnixMilli()) + choice := chooseForecastSources(history, site.Revision, origin.UnixMilli()) + f.noteSourceChoice(choice) pvFn := mpc.PVPredictor(nil) if f.pv != nil && site.HasLocation { pvFn = func(t time.Time, cloud float64) float64 { @@ -597,14 +601,21 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m continue } selected[i].PredictionStartMS = point.PredictionStartMS - if usablePrimaryForecast(point.PVKnown, point.PVQuality, point.PVW) { + // Measured errors keep legacy PV only where legacy has weather for + // this interval; elsewhere the quality rule decides. + pvChoice := choice.PV.Source + if pvChoice == "legacy" && !selected[i].PVKnown { + pvChoice = "" + } + if pvChoice != "legacy" && usablePrimaryForecast(point.PVKnown, point.PVQuality, point.PVW) { resolved[i].PVW = -point.PVW selected[i].PVW, selected[i].PVKnown, selected[i].PVQuality = point.PVW, true, point.PVQuality selected[i].PVSource, selected[i].ModelPV = "energyplan", point.ModelPV } - // A cold load can still be the worker's generic prior. Keep the frozen - // profile and heating prior until this interval has learned support. - if point.LoadQuality != "cold_start" && usablePrimaryForecast(point.LoadKnown, point.LoadQuality, point.LoadW) { + // Without measured evidence a cold load can still be the worker's + // generic prior, so the frozen profile and heating prior stay. + loadReady := point.LoadQuality != "cold_start" || choice.Load.Source == "energyplan" + if choice.Load.Source != "legacy" && loadReady && usablePrimaryForecast(point.LoadKnown, point.LoadQuality, point.LoadW) { resolved[i].LoadW = point.LoadW selected[i].LoadW, selected[i].LoadKnown, selected[i].LoadQuality = point.LoadW, true, point.LoadQuality selected[i].LoadSource, selected[i].ModelLoad = "energyplan", point.ModelLoad @@ -612,16 +623,21 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m } return resolved } - // Calibrate the actual composed primary. A Rust prediction does not inherit - // the old model's residual spread or its training trust threshold. + // Calibrate each slot on the errors its own sources made. A Rust + // prediction does not inherit the old model's residual spread or its + // training trust threshold, and a source switch keeps its history. in.Risk = func(base, planning []mpc.Slot, k float64) { if k <= 0 || math.IsNaN(k) || math.IsInf(k, 0) { return } for i, s := range base { + pvSeries, loadSeries := "legacy_shadow", "legacy_shadow" + if i < len(selected) { + pvSeries, loadSeries = forecastSeriesOf(selected[i].PVSource), forecastSeriesOf(selected[i].LoadSource) + } net := s.LoadW + s.PVW end := s.StartMs + int64(s.LenMin)*time.Minute.Milliseconds() - band := calibrator.BandForInterval("champion", "net", max(s.StartMs, origin.UnixMilli()), end, net) + band := riskCalibrator.BandForInterval(forecastMixSeries(pvSeries, loadSeries), "net", max(s.StartMs, origin.UnixMilli()), end, net) if band.Method == forecasting.BandMethodEmpirical { extra := math.Max(0, k*(band.HighW-net)) pvLoss := math.Min(-s.PVW, extra) @@ -645,7 +661,7 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m if band.Method != forecasting.BandMethodEmpirical { // Until joint errors can set the margin, cover load upside as // well as PV downside. This also works at night and without PV. - loadBand := calibrator.BandForInterval("champion", "load", max(s.StartMs, origin.UnixMilli()), end, s.LoadW) + loadBand := riskCalibrator.BandForInterval(loadSeries, "load", max(s.StartMs, origin.UnixMilli()), end, s.LoadW) loadLoss := math.Max(0, loadBand.HighW-s.LoadW) if loadBand.Method != forecasting.BandMethodEmpirical && i < len(selected) && selected[i].ModelLoad != nil { loadLoss = math.Max(loadLoss, selected[i].ModelLoad.UpperW-s.LoadW) diff --git a/go/internal/forecasting/evaluate.go b/go/internal/forecasting/evaluate.go index 13eeef59b..c4d8af172 100644 --- a/go/internal/forecasting/evaluate.go +++ b/go/internal/forecasting/evaluate.go @@ -705,7 +705,16 @@ func CompareFrozenSeries(samples []ErrorSample, primary, shadow string) []PairMe return compareSeries(samples, primary, shadow, true) } -func compareSeries(samples []ErrorSample, champion, candidate string, sameIssue bool) []PairMetric { +type matchedPair struct { + config string + start, end int64 + lead int + champion, candidate ErrorSample +} + +// matchedPairs keeps the latest issue per series and target, then pairs the +// targets both series scored. sameIssue also requires both from one issue. +func matchedPairs(samples []ErrorSample, champion, candidate string, sameIssue bool) []matchedPair { type targetKey struct { config string start, end int64 @@ -726,12 +735,7 @@ func compareSeries(samples []ErrorSample, champion, candidate string, sameIssue target[key] = e } } - type accumulator struct { - metric PairMetric - days map[int64]bool - wins float64 - } - all := make(map[string]*accumulator) + out := make([]matchedPair, 0, len(championByTarget)) for key, championSample := range championByTarget { candidateSample, ok := candidateByTarget[key] if !ok { @@ -741,24 +745,37 @@ func compareSeries(samples []ErrorSample, champion, candidate string, sameIssue championSample.OriginMS != candidateSample.OriginMS || championSample.IssuedAtMS != candidateSample.IssuedAtMS) { continue } + out = append(out, matchedPair{key.config, key.start, key.end, key.lead, championSample, candidateSample}) + } + return out +} + +func compareSeries(samples []ErrorSample, champion, candidate string, sameIssue bool) []PairMetric { + type accumulator struct { + metric PairMetric + days map[int64]bool + wins float64 + } + all := make(map[string]*accumulator) + for _, pair := range matchedPairs(samples, champion, candidate, sameIssue) { for _, signal := range []string{"pv", "pv_daylight", "load", "net"} { - if !sameActual(championSample, candidateSample, signal) { + if !sameActual(pair.champion, pair.candidate, signal) { continue } - championError, _, _, _ := signalError(championSample, signal) - candidateError, _, _, _ := signalError(candidateSample, signal) - accKey := key.config + "/" + signal + "/" + string(rune('0'+key.lead)) + championError, _, _, _ := signalError(pair.champion, signal) + candidateError, _, _, _ := signalError(pair.candidate, signal) + accKey := pair.config + "/" + signal + "/" + string(rune('0'+pair.lead)) a := all[accKey] if a == nil { a = &accumulator{metric: PairMetric{ - Champion: champion, Candidate: candidate, ConfigVersion: key.config, Signal: signal, Lead: key.lead, + Champion: champion, Candidate: candidate, ConfigVersion: pair.config, Signal: signal, Lead: pair.lead, }, days: map[int64]bool{}} all[accKey] = a } a.metric.Samples++ a.metric.ChampionMAEW += math.Abs(championError) a.metric.CandidateMAEW += math.Abs(candidateError) - a.days[key.start/(24*hourMS)] = true + a.days[pair.start/(24*hourMS)] = true switch { case math.Abs(candidateError) < math.Abs(championError): a.wins++ @@ -785,3 +802,54 @@ func compareSeries(samples []ErrorSample, champion, candidate string, sameIssue } return out } + +// PooledPairMetric is one signal's comparison over every lead bucket. Hours +// and Days count distinct targets, so a target scored at several leads +// counts once. +type PooledPairMetric struct { + Samples, Hours, Days int + ChampionMAEW, CandidateMAEW float64 +} + +// PoolFrozenSeries compares primary and shadow from the same issues for one +// signal across all lead buckets. Pass samples from one config. +func PoolFrozenSeries(samples []ErrorSample, primary, shadow, signal string) PooledPairMetric { + var m PooledPairMetric + hours, days := map[int64]bool{}, map[int64]bool{} + for _, pair := range matchedPairs(samples, primary, shadow, true) { + if !sameActual(pair.champion, pair.candidate, signal) { + continue + } + primaryError, _, _, _ := signalError(pair.champion, signal) + shadowError, _, _, _ := signalError(pair.candidate, signal) + m.Samples++ + m.ChampionMAEW += math.Abs(primaryError) + m.CandidateMAEW += math.Abs(shadowError) + hours[pair.start/hourMS] = true + days[pair.start/(24*hourMS)] = true + } + if m.Samples > 0 { + m.ChampionMAEW /= float64(m.Samples) + m.CandidateMAEW /= float64(m.Samples) + } + m.Hours, m.Days = len(hours), len(days) + return m +} + +// ComposeFrozenSeries scores the forecast that takes PV from pvSeries and +// load from loadSeries, on targets both series scored in the same issue. The +// result is named name. +func ComposeFrozenSeries(samples []ErrorSample, pvSeries, loadSeries, name string) []ErrorSample { + pairs := matchedPairs(samples, pvSeries, loadSeries, true) + out := make([]ErrorSample, 0, len(pairs)) + for _, pair := range pairs { + pv, e := pair.champion, pair.candidate + e.Series = name + e.AvailableAtMS = max(pv.AvailableAtMS, e.AvailableAtMS) + e.PVErrorW, e.PVKnown, e.Daylight = pv.PVErrorW, pv.PVKnown, pv.Daylight + e.Prediction.PVW, e.Prediction.PVKnown, e.Prediction.PVQuality = pv.Prediction.PVW, pv.Prediction.PVKnown, pv.Prediction.PVQuality + e.Prediction.PVSource, e.Prediction.ModelPV, e.Prediction.PVBand = pv.Prediction.PVSource, pv.Prediction.ModelPV, pv.Prediction.PVBand + out = append(out, e) + } + return out +} diff --git a/go/internal/forecasting/pool_test.go b/go/internal/forecasting/pool_test.go new file mode 100644 index 000000000..4709bafe4 --- /dev/null +++ b/go/internal/forecasting/pool_test.go @@ -0,0 +1,74 @@ +package forecasting + +import ( + "fmt" + "testing" + "time" +) + +// scoredAt scores one series' forecast against 1,000 W PV and 1,000 W load. +func scoredAt(series, issue string, start, origin int64, pvW, loadW float64) ErrorSample { + e := testError(series, "cfg", issue, float64(start), float64(origin), 1000-pvW, 1000-loadW) + e.Prediction.PVW, e.Prediction.LoadW = pvW, loadW + e.Daylight = true + return e +} + +func TestPoolFrozenSeriesCountsEachTargetHourOnce(t *testing.T) { + start := time.Date(2026, 9, 7, 0, 0, 0, 0, time.UTC).UnixMilli() + var samples []ErrorSample + for h := int64(0); h < 60; h++ { + at := start + h*hourMS + // Alternate lead buckets, so each bucket alone holds only 30 hours. + ahead := 2 * hourMS + if h%2 == 1 { + ahead = 4 * hourMS + } + issue := fmt.Sprintf("issue-%d", h) + samples = append(samples, scoredAt("energyplan", issue, at, at-ahead, 900, 1100), + scoredAt("legacy_shadow", issue, at, at-ahead, 700, 1600)) + } + // The first target again, from an earlier issue in a third bucket. + samples = append(samples, scoredAt("energyplan", "early", start, start-8*hourMS, 900, 1100), + scoredAt("legacy_shadow", "early", start, start-8*hourMS, 700, 1600)) + for _, m := range CompareFrozenSeries(samples, "energyplan", "legacy_shadow") { + if m.Signal == "load" && m.Samples > 30 { + t.Fatalf("fixture does not split the hours across buckets: %+v", m) + } + } + m := PoolFrozenSeries(samples, "energyplan", "legacy_shadow", "load") + if m.Samples != 61 || m.Hours != 60 || m.Days != 3 || m.ChampionMAEW != 100 || m.CandidateMAEW != 600 { + t.Fatalf("pooled load comparison = %+v", m) + } + if pv := PoolFrozenSeries(samples, "energyplan", "legacy_shadow", "pv_daylight"); pv.ChampionMAEW != 100 || pv.CandidateMAEW != 300 { + t.Fatalf("pooled PV comparison = %+v", pv) + } +} + +func TestComposeFrozenSeriesTakesEachSignalFromItsSeries(t *testing.T) { + start := time.Date(2026, 9, 7, 12, 0, 0, 0, time.UTC).UnixMilli() + pv := scoredAt("energyplan", "issue", start, start-2*hourMS, 900, 1100) + pv.Prediction.PVSource, pv.Prediction.LoadSource = "energyplan", "energyplan" + load := scoredAt("legacy_shadow", "issue", start, start-2*hourMS, 700, 1600) + load.Prediction.PVSource, load.Prediction.LoadSource = "legacy", "legacy" + alone := scoredAt("legacy_shadow", "other", start+hourMS, start-hourMS, 700, 1600) + got := ComposeFrozenSeries([]ErrorSample{pv, load, alone}, "energyplan", "legacy_shadow", "mix") + if len(got) != 1 { + t.Fatalf("composed %d samples, want only the target both series scored", len(got)) + } + e := got[0] + if err := e.Validate(); err != nil { + t.Fatal(err) + } + if e.Series != "mix" || e.PVErrorW != 100 || e.LoadErrorW != -600 || e.Prediction.PVW != 900 || e.Prediction.LoadW != 1600 || + e.Prediction.PVSource != "energyplan" || e.Prediction.LoadSource != "legacy" { + t.Fatalf("mix did not take PV and load from their series: %+v", e) + } + if net, known, _, _ := signalError(e, "net"); !known || net != -700 { + t.Fatalf("mixed net error = %v (known %v), want -700", net, known) + } + load.IssueID = "newer" + if got := ComposeFrozenSeries([]ErrorSample{pv, load}, "energyplan", "legacy_shadow", "mix"); len(got) != 0 { + t.Fatalf("errors from different issues were mixed: %+v", got) + } +}