From 02772cd6f32e4e6201c186b31a9edb7fd75f47c0 Mon Sep 17 00:00:00 2001 From: Fredrik Ahlgren Date: Thu, 1 Oct 2026 19:18:58 +0200 Subject: [PATCH 1/3] fix(forecast): choose each forecast source from measured errors Core picked PV and load sources from the worker's quality label alone. On the home box that put cold Energyplan solar and a legacy load that overshot by 2 kW at night into the plan, while the archive already showed the other source was better for each signal. Compare Energyplan and legacy_shadow errors from the same issues over the last week. With 48 scored hours over three days and a 10% gap, use the better source for that signal; otherwise keep the quality rule. Log each change of choice. Bump the pipeline policy so bands do not mix errors from the old rule. Fixes srcfl/ftw#1490 Co-Authored-By: Claude Opus 5.5 Signed-off-by: Fredrik Ahlgren --- .changeset/forecast-source-from-errors.md | 5 ++ docs/architecture.md | 4 +- docs/energyplan-contract.md | 2 + go/cmd/ftw/forecast_site.go | 2 +- go/cmd/ftw/forecast_source_choice.go | 92 ++++++++++++++++++++ go/cmd/ftw/forecast_source_choice_test.go | 101 ++++++++++++++++++++++ go/cmd/ftw/forecast_tracking.go | 12 ++- 7 files changed, 212 insertions(+), 6 deletions(-) create mode 100644 .changeset/forecast-source-from-errors.md create mode 100644 go/cmd/ftw/forecast_source_choice.go create mode 100644 go/cmd/ftw/forecast_source_choice_test.go 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..c3fae1fba 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -243,7 +243,9 @@ 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 capture for a fair later comparison. 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_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..1a5057381 --- /dev/null +++ b/go/cmd/ftw/forecast_source_choice.go @@ -0,0 +1,92 @@ +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) + } + } + metrics := forecasting.CompareFrozenSeries(recent, "energyplan", "legacy_shadow") + return forecastSourceChoice{PV: pickForecastSource(metrics, "pv_daylight"), Load: pickForecastSource(metrics, "load")} +} + +func pickForecastSource(metrics []forecasting.PairMetric, signal string) forecastSourcePick { + var pick forecastSourcePick + for _, m := range metrics { + if m.Signal != signal { + continue + } + pick.Samples += m.Samples + pick.Hours = max(pick.Hours, m.Samples) // one sample per scored hour in each lead bucket + pick.Days = max(pick.Days, m.Days) + pick.EnergyplanMAEW += m.ChampionMAEW * float64(m.Samples) + pick.LegacyMAEW += m.CandidateMAEW * float64(m.Samples) + } + if pick.Samples == 0 { + return pick + } + pick.EnergyplanMAEW /= float64(pick.Samples) + pick.LegacyMAEW /= float64(pick.Samples) + 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 +} + +// 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..f7ebec194 --- /dev/null +++ b/go/cmd/ftw/forecast_source_choice_test.go @@ -0,0 +1,101 @@ +package main + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" +) + +type forecastPair struct{ energyplanPV, legacyPV, energyplanLoad, legacyLoad float64 } + +// pairedForecastErrors scores both sources of each issue against one actual: +// 2,000 W PV and 1,000 W load in every hour before end. +func pairedForecastErrors(t *testing.T, cohort string, end time.Time, hours int, p forecastPair) []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) + for _, s := range []struct { + series string + pv, load float64 + }{{"energyplan", p.energyplanPV, p.energyplanLoad}, {"legacy_shadow", p.legacyPV, p.legacyLoad}} { + e := forecasting.ErrorSample{Series: s.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 - s.pv, LoadErrorW: 1000 - s.load, PVKnown: true, Daylight: true, LoadKnown: true, + Prediction: forecasting.Point{StartMS: at.UnixMilli(), EndMS: at.Add(15 * time.Minute).UnixMilli(), PVW: s.pv, LoadW: s.load, + PVKnown: true, LoadKnown: true, PVQuality: "test", LoadQuality: "test", PVBand: band, LoadBand: band, NetBand: band}} + if err := e.Validate(); err != nil { + t.Fatal(err) + } + out = append(out, e) + } + } + return out +} + +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) + }) + f.errors = pairedForecastErrors(t, f.site().Revision, at, 72, + forecastPair{energyplanPV: 1250, legacyPV: 1850, energyplanLoad: 1100, legacyLoad: 3000}) + in := f.Snapshot(at, trackerWeather(at, at)) + legacy := trackerSlots(at, 2) + 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) + } + in.Record(got, got, "decision", at.UnixMilli()) + point := primarySeries(t, (<-f.queue).issue, "champion").Points[0] + if point.PVSource != "legacy" || point.LoadSource != "energyplan" || point.LoadQuality != "cold_start" || point.ModelLoad == nil { + t.Fatalf("chosen sources not recorded: %+v", point) + } + + // 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) + } +} diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index a17c3d043..a68247577 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,8 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m frozen = fresh } calibrator := forecasting.NewCalibrator(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 +600,15 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m continue } selected[i].PredictionStartMS = point.PredictionStartMS - if usablePrimaryForecast(point.PVKnown, point.PVQuality, point.PVW) { + if choice.PV.Source != "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 From f03eb771acb99db6ddd6c7a4eb0b2ea783dfb0f9 Mon Sep 17 00:00:00 2001 From: Fredrik Ahlgren Date: Thu, 1 Oct 2026 19:35:42 +0200 Subject: [PATCH 2/3] fix(forecast): calibrate bands on the sources chosen now After the choice moves a signal to another source, champion errors from the old source stayed in the calibration window for up to 256 hours, so the margin described a source the plan no longer used. Keep only champion errors whose recorded sources match the current choice. Until the new choice has its own errors, bands start cold. Co-Authored-By: Claude Opus 5.5 Signed-off-by: Fredrik Ahlgren --- go/cmd/ftw/forecast_source_choice.go | 18 +++++++ go/cmd/ftw/forecast_source_choice_test.go | 62 ++++++++++++++++------- go/cmd/ftw/forecast_tracking.go | 2 +- 3 files changed, 64 insertions(+), 18 deletions(-) diff --git a/go/cmd/ftw/forecast_source_choice.go b/go/cmd/ftw/forecast_source_choice.go index 1a5057381..788c53dc2 100644 --- a/go/cmd/ftw/forecast_source_choice.go +++ b/go/cmd/ftw/forecast_source_choice.go @@ -69,6 +69,24 @@ func pickForecastSource(metrics []forecasting.PairMetric, signal string) forecas return pick } +// calibrationEvidence keeps only champion errors made with the sources chosen +// now, so a band never describes a source the plan stopped using. Until a new +// choice has its own errors, the bands start cold. +func calibrationEvidence(history []forecasting.ErrorSample, c forecastSourceChoice) []forecasting.ErrorSample { + if c.PV.Source == "" && c.Load.Source == "" { + return history + } + out := make([]forecasting.ErrorSample, 0, len(history)) + for _, e := range history { + if e.Series == "champion" && ((c.PV.Source != "" && e.Prediction.PVSource != c.PV.Source) || + (c.Load.Source != "" && e.Prediction.LoadSource != c.Load.Source)) { + continue + } + out = append(out, e) + } + return out +} + // noteSourceChoice logs when measured errors move a signal to another source. func (f *forecastTracker) noteSourceChoice(c forecastSourceChoice) { f.mu.Lock() diff --git a/go/cmd/ftw/forecast_source_choice_test.go b/go/cmd/ftw/forecast_source_choice_test.go index f7ebec194..0cb2c322e 100644 --- a/go/cmd/ftw/forecast_source_choice_test.go +++ b/go/cmd/ftw/forecast_source_choice_test.go @@ -7,38 +7,42 @@ import ( "time" "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/mpc" ) type forecastPair struct{ energyplanPV, legacyPV, energyplanLoad, legacyLoad float64 } -// pairedForecastErrors scores both sources of each issue against one actual: -// 2,000 W PV and 1,000 W load in every hour before end. -func pairedForecastErrors(t *testing.T, cohort string, end time.Time, hours int, p forecastPair) []forecasting.ErrorSample { +// 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) - for _, s := range []struct { - series string - pv, load float64 - }{{"energyplan", p.energyplanPV, p.energyplanLoad}, {"legacy_shadow", p.legacyPV, p.legacyLoad}} { - e := forecasting.ErrorSample{Series: s.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 - s.pv, LoadErrorW: 1000 - s.load, PVKnown: true, Daylight: true, LoadKnown: true, - Prediction: forecasting.Point{StartMS: at.UnixMilli(), EndMS: at.Add(15 * time.Minute).UnixMilli(), PVW: s.pv, LoadW: s.load, - PVKnown: true, LoadKnown: true, PVQuality: "test", LoadQuality: "test", PVBand: band, LoadBand: band, NetBand: band}} - if err := e.Validate(); err != nil { - t.Fatal(err) - } - out = append(out, e) + 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. @@ -99,3 +103,27 @@ func TestForecastSourceChoiceSteersResolve(t *testing.T) { t.Fatalf("legacy load or default PV rule lost: %+v", got) } } + +func TestForecastSourceChoiceCalibratesOnlyTheChosenSources(t *testing.T) { + at := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + f := primaryFixture(at, primaryReply) + cohort := f.site().Revision + // Legacy load measured far better, so it now feeds the plan. The champion + // errors come from before, when Energyplan load overshot by 2 kW. + f.errors = append(pairedForecastErrors(t, cohort, at, 8*24, forecastPair{energyplanPV: 2000, legacyPV: 2000, energyplanLoad: 3000, legacyLoad: 1100}), + scoredForecastErrors(t, "champion", cohort, at, 8*24, 2000, 3000, "energyplan", "energyplan")...) + 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) + if base[0].LoadW != 1000 { + t.Fatalf("legacy load not chosen: %+v", base[0]) + } + planning := append([]mpc.Slot(nil), base...) + in.Risk(base, planning, 1) + // A band from the old errors would sit 2 kW below the forecast and add + // nothing. The cold band for a 1 kW load at lead 1 adds 400 W. + if planning[0].LoadW != 1400 { + t.Fatalf("band used errors of a source the plan no longer uses: %+v", planning[0]) + } +} diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index a68247577..21aebe7dd 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -450,9 +450,9 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m } frozen = fresh } - calibrator := forecasting.NewCalibrator(history, site.Revision, origin.UnixMilli()) choice := chooseForecastSources(history, site.Revision, origin.UnixMilli()) f.noteSourceChoice(choice) + calibrator := forecasting.NewCalibrator(calibrationEvidence(history, choice), site.Revision, origin.UnixMilli()) pvFn := mpc.PVPredictor(nil) if f.pv != nil && site.HasLocation { pvFn = func(t time.Time, cloud float64) float64 { From c9dc72b8d3eb883e6ab17cf7e05adf63512ca0a2 Mon Sep 17 00:00:00 2001 From: Fredrik Ahlgren Date: Thu, 1 Oct 2026 20:03:50 +0200 Subject: [PATCH 3/3] fix(forecast): margin from each slot's own sources; pool hours Three findings from a local Codex review: - Calibration filtered champion errors only while a measured choice held. When the choice fell back to the quality rule, old errors of another source set the margin again. Each slot's margin now comes from the errors its own sources made on earlier issues: Energyplan, legacy, or a mix composed from both series on the same issues. A switch keeps its history instead of starting cold. - A measured legacy PV choice also blocked Energyplan PV where legacy had no weather. Keep legacy PV only where it is known. - The evidence gate counted hours in the largest lead bucket. Count distinct scored hours and days across all buckets. Co-Authored-By: Claude Opus 5.5 Signed-off-by: Fredrik Ahlgren --- docs/architecture.md | 3 +- go/cmd/ftw/forecast_load_risk_test.go | 16 ++-- go/cmd/ftw/forecast_source_choice.go | 65 ++++++++-------- go/cmd/ftw/forecast_source_choice_test.go | 77 ++++++++++++------- go/cmd/ftw/forecast_tracking.go | 24 ++++-- go/internal/forecasting/evaluate.go | 94 +++++++++++++++++++---- go/internal/forecasting/pool_test.go | 74 ++++++++++++++++++ 7 files changed, 270 insertions(+), 83 deletions(-) create mode 100644 go/internal/forecasting/pool_test.go diff --git a/docs/architecture.md b/docs/architecture.md index c3fae1fba..f0a681e46 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -247,7 +247,8 @@ 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/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_source_choice.go b/go/cmd/ftw/forecast_source_choice.go index 788c53dc2..16cd0b5ae 100644 --- a/go/cmd/ftw/forecast_source_choice.go +++ b/go/cmd/ftw/forecast_source_choice.go @@ -36,27 +36,15 @@ func chooseForecastSources(history []forecasting.ErrorSample, cohort string, ori recent = append(recent, e) } } - metrics := forecasting.CompareFrozenSeries(recent, "energyplan", "legacy_shadow") - return forecastSourceChoice{PV: pickForecastSource(metrics, "pv_daylight"), Load: pickForecastSource(metrics, "load")} + return forecastSourceChoice{ + PV: pickForecastSource(forecasting.PoolFrozenSeries(recent, "energyplan", "legacy_shadow", "pv_daylight")), + Load: pickForecastSource(forecasting.PoolFrozenSeries(recent, "energyplan", "legacy_shadow", "load")), + } } -func pickForecastSource(metrics []forecasting.PairMetric, signal string) forecastSourcePick { - var pick forecastSourcePick - for _, m := range metrics { - if m.Signal != signal { - continue - } - pick.Samples += m.Samples - pick.Hours = max(pick.Hours, m.Samples) // one sample per scored hour in each lead bucket - pick.Days = max(pick.Days, m.Days) - pick.EnergyplanMAEW += m.ChampionMAEW * float64(m.Samples) - pick.LegacyMAEW += m.CandidateMAEW * float64(m.Samples) - } - if pick.Samples == 0 { - return pick - } - pick.EnergyplanMAEW /= float64(pick.Samples) - pick.LegacyMAEW /= float64(pick.Samples) +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 } @@ -69,20 +57,37 @@ func pickForecastSource(metrics []forecasting.PairMetric, signal string) forecas return pick } -// calibrationEvidence keeps only champion errors made with the sources chosen -// now, so a band never describes a source the plan stopped using. Until a new -// choice has its own errors, the bands start cold. -func calibrationEvidence(history []forecasting.ErrorSample, c forecastSourceChoice) []forecasting.ErrorSample { - if c.PV.Source == "" && c.Load.Source == "" { - return history +// forecastSeriesOf names the archived series that holds a source's own +// forecast. +func forecastSeriesOf(source string) string { + if source == "energyplan" { + return "energyplan" } - out := make([]forecasting.ErrorSample, 0, len(history)) + 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 == "champion" && ((c.PV.Source != "" && e.Prediction.PVSource != c.PV.Source) || - (c.Load.Source != "" && e.Prediction.LoadSource != c.Load.Source)) { - continue + if e.Series == "energyplan" || e.Series == "legacy_shadow" { + own = append(own, e) } - out = append(out, 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 } diff --git a/go/cmd/ftw/forecast_source_choice_test.go b/go/cmd/ftw/forecast_source_choice_test.go index 0cb2c322e..f0396393b 100644 --- a/go/cmd/ftw/forecast_source_choice_test.go +++ b/go/cmd/ftw/forecast_source_choice_test.go @@ -79,18 +79,27 @@ func TestForecastSourceChoiceSteersResolve(t *testing.T) { _, reply := hostForecastReply(p) return json.Marshal(reply) }) - f.errors = pairedForecastErrors(t, f.site().Revision, at, 72, + 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, 2) + 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) + 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()) - point := primarySeries(t, (<-f.queue).issue, "champion").Points[0] - if point.PVSource != "legacy" || point.LoadSource != "energyplan" || point.LoadQuality != "cold_start" || point.ModelLoad == nil { - t.Fatalf("chosen sources not recorded: %+v", point) + 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. @@ -100,30 +109,44 @@ func TestForecastSourceChoiceSteersResolve(t *testing.T) { 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) + t.Fatalf("legacy load or default PV rule lost: %+v", got[0]) } } -func TestForecastSourceChoiceCalibratesOnlyTheChosenSources(t *testing.T) { +// 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) - f := primaryFixture(at, primaryReply) - cohort := f.site().Revision - // Legacy load measured far better, so it now feeds the plan. The champion - // errors come from before, when Energyplan load overshot by 2 kW. - f.errors = append(pairedForecastErrors(t, cohort, at, 8*24, forecastPair{energyplanPV: 2000, legacyPV: 2000, energyplanLoad: 3000, legacyLoad: 1100}), - scoredForecastErrors(t, "champion", cohort, at, 8*24, 2000, 3000, "energyplan", "energyplan")...) - 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) - if base[0].LoadW != 1000 { - t.Fatalf("legacy load not chosen: %+v", base[0]) - } - planning := append([]mpc.Slot(nil), base...) - in.Risk(base, planning, 1) - // A band from the old errors would sit 2 kW below the forecast and add - // nothing. The cold band for a 1 kW load at lead 1 adds 400 W. - if planning[0].LoadW != 1400 { - t.Fatalf("band used errors of a source the plan no longer uses: %+v", planning[0]) + 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 21aebe7dd..086089af2 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -450,9 +450,10 @@ 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) - calibrator := forecasting.NewCalibrator(calibrationEvidence(history, choice), site.Revision, origin.UnixMilli()) pvFn := mpc.PVPredictor(nil) if f.pv != nil && site.HasLocation { pvFn = func(t time.Time, cloud float64) float64 { @@ -600,7 +601,13 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m continue } selected[i].PredictionStartMS = point.PredictionStartMS - if choice.PV.Source != "legacy" && 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 @@ -616,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) @@ -649,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) + } +}