Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/forecast-source-from-errors.md
Original file line number Diff line number Diff line change
@@ -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.
7 changes: 5 additions & 2 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/energyplan-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
16 changes: 10 additions & 6 deletions go/cmd/ftw/forecast_load_risk_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion go/cmd/ftw/forecast_site.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
115 changes: 115 additions & 0 deletions go/cmd/ftw/forecast_source_choice.go
Original file line number Diff line number Diff line change
@@ -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))
}
}
152 changes: 152 additions & 0 deletions go/cmd/ftw/forecast_source_choice_test.go
Original file line number Diff line number Diff line change
@@ -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])
}
})
}
}
Loading
Loading