diff --git a/.changeset/restart-forecast-learning.md b/.changeset/restart-forecast-learning.md new file mode 100644 index 00000000..ca6483cd --- /dev/null +++ b/.changeset/restart-forecast-learning.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Restart solar-production or household-consumption learning after a lasting site change. Each action resets the selected primary model, its legacy fallback and error calibration, preserves measured history and the other model, and shows the new learning period. Saved reset intent survives a restart and prevents old observations from restoring the previous model. diff --git a/go/cmd/ftw/forecast_learning.go b/go/cmd/ftw/forecast_learning.go new file mode 100644 index 00000000..d0a3d37c --- /dev/null +++ b/go/cmd/ftw/forecast_learning.go @@ -0,0 +1,281 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" +) + +const forecastLearningKey = "forecast/learning_periods_v1" + +func (f *forecastTracker) setReplan(fn func(string)) { + f.learningMu.Lock() + defer f.learningMu.Unlock() + f.requestReplan = fn +} + +type forecastLearningPeriods struct { + ConfigRevision string `json:"config_revision"` + PVMS int64 `json:"pv_ms"` + LoadMS int64 `json:"load_ms"` +} + +func (f *forecastTracker) learningSiteLocked(site forecastSite) forecastSite { + if f.learningPeriods.ConfigRevision == rustConfigRevision(site) { + site.PVLearningStartedMS, site.LoadLearningStartedMS = f.learningPeriods.PVMS, f.learningPeriods.LoadMS + } + return site +} + +func (f *forecastTracker) restoreLearning(ctx context.Context) error { + if data, ok := f.store.LoadConfig(forecastLearningKey); ok { + if err := json.Unmarshal([]byte(data), &f.learningPeriods); err != nil { + return fmt.Errorf("read forecast learning periods: %w", err) + } + if f.learningPeriods.ConfigRevision == "" || f.learningPeriods.PVMS < 0 || f.learningPeriods.LoadMS < 0 { + return errors.New("invalid forecast learning periods") + } + } + if f.configMu != nil { + f.configMu.RLock() + defer f.configMu.RUnlock() + } + f.learningMu.Lock() + defer f.learningMu.Unlock() + for _, signal := range []string{"pv", "load"} { + if err := f.applyLearningLocked(ctx, signal); err != nil { + slog.Warn("forecast learning restart pending", "signal", signal, "err", err) + } + } + return nil +} + +func (f *forecastTracker) applyLearningLocked(ctx context.Context, signal string) error { + site := f.learningSiteLocked(f.site()) + if site.IdentityPending { + return nil + } + var errs []error + if signal == "pv" && site.PVLearningStartedMS > 0 && f.pv != nil { + errs = append(errs, f.pv.RestartLearning(time.UnixMilli(site.PVLearningStartedMS))) + } + if signal == "load" && site.LoadLearningStartedMS > 0 && f.load != nil { + errs = append(errs, f.load.RestartLearning(time.UnixMilli(site.LoadLearningStartedMS))) + } + if r, ok := f.candidate.(*rustForecast); ok { + errs = append(errs, r.RestartLearning(ctx, site, signal)) + } + err := errors.Join(errs...) + if f.learningErrors == nil { + f.learningErrors = make(map[string]error) + } + f.learningErrors[signal] = err + return err +} + +func (f *forecastTracker) learningPendingLocked(site forecastSite, signal string) bool { + if site.IdentityPending || f.learningPeriods.ConfigRevision != rustConfigRevision(site) { + return false + } + cutoff := site.PVLearningStartedMS + if signal == "load" { + cutoff = site.LoadLearningStartedMS + } + if cutoff == 0 { + return false + } + if f.learningErrors[signal] != nil { + return true + } + if signal == "pv" && f.pv != nil && f.pv.LearningStartedMS() < cutoff { + return true + } + if signal == "load" && f.load != nil && f.load.LearningStartedMS() < cutoff { + return true + } + if r, ok := f.candidate.(*rustForecast); ok { + return r.LearningStatus(site, signal).Status == "unavailable" + } + return false +} + +// Both samplers and forecast captures reconcile before using any legacy state. +// In particular, startup may first need to wait for the hardware identity. +func (f *forecastTracker) reconcileLearning(ctx context.Context) { + if f.configMu != nil { + f.configMu.RLock() + defer f.configMu.RUnlock() + } + f.learningMu.Lock() + defer f.learningMu.Unlock() + site := f.learningSiteLocked(f.site()) + if site.IdentityPending { + return + } + for _, signal := range []string{"pv", "load"} { + if !f.learningPendingLocked(site, signal) { + continue + } + workCtx, cancel := context.WithTimeout(ctx, 2*time.Second) + err := f.applyLearningLocked(workCtx, signal) + cancel() + if err == nil && f.requestReplan != nil { + f.requestReplan(signal + "_learning_recovered") + } + } +} + +// RestartLearning saves the intent before changing either pipeline. Startup +// replays it idempotently, including after a crash between the model saves. +func (f *forecastTracker) RestartLearning(ctx context.Context, signal string) error { + if f == nil || (signal != "pv" && signal != "load") { + return errors.New("forecast learning unavailable") + } + if f.refreshIdentity != nil { + f.refreshIdentity() + } + if f.configMu != nil { + f.configMu.RLock() + defer f.configMu.RUnlock() + } + f.learningMu.Lock() + defer f.learningMu.Unlock() + site := f.site() + if site.IdentityPending { + return errors.New("waiting for site identity before restarting learning") + } + if (signal == "pv" && !site.HasLocation && f.pv == nil) || (signal == "load" && f.load == nil) { + return errors.New("forecast model disabled") + } + if r, ok := f.candidate.(*rustForecast); ok { + if !r.resetSupported { + return errors.New("forecast worker does not support restarting learning") + } + if signal == "pv" && !site.HasLocation { + return errors.New("solar forecast needs a site location") + } + } + next := f.learningPeriods + if next.ConfigRevision != rustConfigRevision(site) { + next = forecastLearningPeriods{ConfigRevision: rustConfigRevision(site)} + } + cutoff := f.now().UnixMilli() + previous := next.PVMS + if signal == "load" { + previous = next.LoadMS + } + if cutoff <= 0 || cutoff < previous { + return errors.New("clock precedes the current learning period") + } + // A retry of an incomplete operation retains its original cutoff. + if f.learningErrors[signal] != nil && previous > 0 { + cutoff = previous + } + if signal == "pv" { + next.PVMS = cutoff + } else { + next.LoadMS = cutoff + } + data, err := json.Marshal(next) + if err != nil { + return err + } + // Cancel a solve that captured the old models. Its replacement waits for + // learningMu, so it cannot capture a partially reset pair of pipelines. + if f.requestReplan != nil { + f.requestReplan(signal + "_learning_restarted") + } + if err = f.store.SaveConfig(forecastLearningKey, string(data)); err != nil { + return fmt.Errorf("save forecast learning period: %w", err) + } + f.learningPeriods = next + if err = f.applyLearningLocked(ctx, signal); err != nil { + return &forecasting.LearningRestartPendingError{Err: err} + } + return nil +} + +func (f *forecastTracker) LearningStatus(signal string) forecasting.LearningStatus { + status := forecasting.LearningStatus{Engine: "legacy", Status: "unavailable"} + if f == nil { + return status + } + if f.configMu != nil { + f.configMu.RLock() + defer f.configMu.RUnlock() + } + f.learningMu.RLock() + defer f.learningMu.RUnlock() + site := f.learningSiteLocked(f.site()) + if r, ok := f.candidate.(*rustForecast); ok { + status = r.LearningStatus(site, signal) + } else { + status.Status, status.ResetAvailable = "cold_start", !site.IdentityPending + if signal == "pv" { + status.StartedMS = site.PVLearningStartedMS + if f.pv == nil { + status.Status, status.ResetAvailable = "unavailable", false + } else { + m := f.pv.Model() + status.LatestTrainingMS = m.LastMs + if m.Samples > 0 { + status.Status = "learning" + } + if m.Quality() >= 1 { + status.Status = "ready" + } + } + } else { + status.StartedMS = site.LoadLearningStartedMS + if f.load == nil { + status.Status, status.ResetAvailable = "unavailable", false + } else { + m := f.load.Model() + status.LatestTrainingMS = m.LastMs + if m.Samples > 0 { + status.Status = "learning" + } + if m.Quality() >= 1 { + status.Status = "ready" + } + } + } + } + if site.IdentityPending || (f.learningPeriods.ConfigRevision == rustConfigRevision(site) && f.learningErrors[signal] != nil) { + status.Status = "unavailable" + } + return status +} + +// Calibration is per signal: a PV reset also drops joint net errors, while +// retaining load evidence. Neither the archive nor its measured truth is erased. +func learningEvidence(history []forecasting.ErrorSample, observations []forecasting.Observation, pvMS, loadMS int64) ([]forecasting.ErrorSample, []forecasting.Observation) { + if pvMS == 0 && loadMS == 0 { + return history, observations + } + errors := append([]forecasting.ErrorSample(nil), history...) + for i := range errors { + e := &errors[i] + if e.OriginMS < pvMS || e.StartMS < pvMS { + e.PVKnown = false + } + if e.OriginMS < loadMS || e.StartMS < loadMS { + e.LoadKnown = false + } + } + truth := append([]forecasting.Observation(nil), observations...) + for i := range truth { + if truth[i].StartMS < pvMS { + truth[i].PVKnown = false + } + if truth[i].StartMS < loadMS { + truth[i].LoadKnown = false + } + } + return errors, truth +} diff --git a/go/cmd/ftw/forecast_learning_backup_test.go b/go/cmd/ftw/forecast_learning_backup_test.go new file mode 100644 index 00000000..bf0e46d9 --- /dev/null +++ b/go/cmd/ftw/forecast_learning_backup_test.go @@ -0,0 +1,292 @@ +package main + +import ( + "bytes" + "context" + "os" + "path/filepath" + "reflect" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/backup" + "github.com/srcfl/ftw/go/internal/config" + "github.com/srcfl/ftw/go/internal/loadmodel" + "github.com/srcfl/ftw/go/internal/modelstate" + "github.com/srcfl/ftw/go/internal/pvmodel" + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func TestForecastLearningNativeMigrationBackupRestore(t *testing.T) { + for _, signal := range []string{"pv", "load"} { + t.Run(signal, func(t *testing.T) { + root := t.TempDir() + dataDir := filepath.Join(root, "data") + if err := os.Mkdir(dataDir, 0o700); err != nil { + t.Fatal(err) + } + configPath := filepath.Join(dataDir, "config.yaml") + databasePath := filepath.Join(dataDir, "state.db") + cfg, err := config.Parse([]byte(` +site: + name: Forecast backup test +fuse: + max_amps: 16 +api: + port: 8080 +app_link: + enabled: false +weather: + provider: open_meteo + latitude: 59 + longitude: 18 + timezone: Europe/Stockholm + pv_rated_w: 8000 + heating_w_per_degc: 275 +planner: + pv_forecast_safety_k: 0 +`), dataDir) + if err != nil { + t.Fatal(err) + } + st, err := state.Open(databasePath) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = st.Close() }) + + siteConfig := newForecastSiteConfig(st) + siteConfig.Configure(cfg, nil) + site := siteConfig.Snapshot() + if site.IdentityPending || !site.HasLocation || site.LearningRevision == "" { + t.Fatalf("invalid forecast fixture site: %+v", site) + } + start := time.Now().UTC().Truncate(15 * time.Minute).Add(-30 * time.Minute) + seedForecastGoModels(t, st, site, start) + + native := learningNative(t, st) + if err := native.Update(context.Background(), site, learningObservation(site, start), nil, false); err != nil { + t.Fatal(err) + } + tracker := migrationLearningTracker(t, st, &site, native, cfg) + beforePV := tracker.pv.Model() + beforeLoad := tracker.load.Snapshot() + other := "load" + if signal == "load" { + other = "pv" + } + beforeNativeOther := append([]byte(nil), learningModel(t, native, other)...) + + cutoff := start.Add(22 * time.Minute) + tracker.clock = func() time.Time { return cutoff } + if err := tracker.RestartLearning(context.Background(), signal); err != nil { + t.Fatal(err) + } + assertMigrationLearningState(t, tracker, signal, cutoff.UnixMilli(), beforePV, beforeLoad, beforeNativeOther) + + cfg, err = config.InitializeStorage(configPath, databasePath, cfg, st) + if err != nil { + t.Fatal(err) + } + loaded, err := config.Load(configPath) + if err != nil { + t.Fatal(err) + } + migratedSiteConfig := newForecastSiteConfig(st) + migratedSiteConfig.Configure(loaded, nil) + migratedSite := migratedSiteConfig.Snapshot() + assertForecastSiteIdentity(t, site, migratedSite) + + postResetPV := tracker.pv.Model() + postResetLoad := tracker.load.Snapshot() + postResetNative := append([]byte(nil), native.Snapshot()...) + periods := tracker.learningPeriods + + info, err := backup.Create(context.Background(), backup.CreateOptions{ + ConfigPath: configPath, + State: st, + StatePath: databasePath, + DataDir: dataDir, + OutputDir: filepath.Join(root, "backups"), + Now: cutoff.Add(time.Hour), + }) + if err != nil { + t.Fatal(err) + } + if _, err := backup.Verify(info.Path); err != nil { + t.Fatal(err) + } + restoredDir := filepath.Join(root, "restored") + if _, err := backup.Restore(info.Path, restoredDir, cutoff.Add(2*time.Hour)); err != nil { + t.Fatal(err) + } + + restoredConfig, err := config.Load(filepath.Join(restoredDir, "config.yaml")) + if err != nil { + t.Fatal(err) + } + restoredStore, err := state.Open(restoredConfig.ConfigDatabase) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = restoredStore.Close() }) + restoredSiteConfig := newForecastSiteConfig(restoredStore) + restoredSiteConfig.Configure(restoredConfig, nil) + restoredSite := restoredSiteConfig.Snapshot() + assertForecastSiteIdentity(t, site, restoredSite) + + restoredNative := learningNative(t, restoredStore) + restored := migrationLearningTracker(t, restoredStore, &restoredSite, restoredNative, restoredConfig) + if err := restored.restoreLearning(context.Background()); err != nil { + t.Fatal(err) + } + if restored.learningPeriods != periods { + t.Fatalf("learning periods changed across backup restore: got %+v want %+v", restored.learningPeriods, periods) + } + if got := restored.pv.Model(); !reflect.DeepEqual(got, postResetPV) { + t.Fatalf("PV model changed across backup restore:\ngot %+v\nwant %+v", got, postResetPV) + } + if got := restored.load.Snapshot(); !reflect.DeepEqual(got, postResetLoad) { + t.Fatalf("load models changed across backup restore:\ngot %+v\nwant %+v", got, postResetLoad) + } + if got := restoredNative.Snapshot(); !bytes.Equal(got, postResetNative) { + t.Fatalf("native model changed across backup restore:\ngot %s\nwant %s", got, postResetNative) + } + assertMigrationLearningState(t, restored, signal, cutoff.UnixMilli(), beforePV, beforeLoad, beforeNativeOther) + + effectiveSite := restored.learningSiteLocked(restoredSite) + eligible := learningObservation(effectiveSite, start.Add(30*time.Minute)) + if err := restoredNative.Update(context.Background(), effectiveSite, eligible, nil, false); err != nil { + t.Fatal(err) + } + if status := restoredNative.LearningStatus(effectiveSite, signal); status.Status != "learning" || status.LatestTrainingMS != eligible.EndMS { + t.Fatalf("new post-reset observation did not train %s: %+v", signal, status) + } + beforeStale := append([]byte(nil), restoredNative.Snapshot()...) + if err := restoredNative.Update(context.Background(), effectiveSite, learningObservation(effectiveSite, start), nil, false); err == nil { + t.Fatal("pre-reset observation was accepted after restore") + } + if got := restoredNative.Snapshot(); !bytes.Equal(got, beforeStale) { + t.Fatal("rejected pre-reset observation changed native state") + } + }) + } +} + +func seedForecastGoModels(t *testing.T, st *state.Store, site forecastSite, at time.Time) { + t.Helper() + pv := pvmodel.NewModel(8000) + pv.ConfigRevision = site.LearningRevision + if !pv.Update(800, 20, at, 3000) { + t.Fatal("PV fixture did not learn") + } + encoded, err := modelstate.Wrap(pvmodel.FeatureHash(), pv) + if err != nil { + t.Fatal(err) + } + if err := st.SaveConfig("pvmodel/state_utc", encoded); err != nil { + t.Fatal(err) + } + + values := map[string]string{ + "loadmodel/profile": string(loadmodel.ProfileAway), + "loadmodel/timezone": site.Timezone, + } + for _, profile := range loadmodel.Profiles() { + model := loadmodel.NewModel(4000) + if profile == loadmodel.ProfileAway { + model.PriorScale = 0.25 + } + model.ConfigRevision = site.LearningRevision + model.Timezone = site.Timezone + model.HeatingW_per_degC = 900 + if !model.Update(at, 10_000, 10) { + t.Fatalf("%s load fixture did not learn", profile) + } + encoded, err := modelstate.Wrap(loadmodel.FeatureHash(), model) + if err != nil { + t.Fatal(err) + } + values["loadmodel/state_utc:"+string(profile)] = encoded + } + if err := st.SaveConfigValues(values); err != nil { + t.Fatal(err) + } +} + +func migrationLearningTracker(t *testing.T, st *state.Store, site *forecastSite, native *rustForecast, cfg *config.Config) *forecastTracker { + t.Helper() + tel := telemetry.NewStore() + pv := pvmodel.NewService(st, tel, func(time.Time) float64 { return 800 }, nil, 8000) + pv.Reconfigure(func(time.Time) float64 { return 800 }, site.LearningRevision) + load := loadmodel.NewService(st, tel, site.Meter, 4000, 0) + if err := load.Reconfigure(site.Meter, site.Options, site.Timezone, site.LearningRevision); err != nil { + t.Fatal(err) + } + load.SeedHeatingCoef(cfg.Weather.HeatingWPerDegC) + return &forecastTracker{ + store: st, + tele: tel, + pv: pv, + load: load, + site: func() forecastSite { return *site }, + candidate: native, + } +} + +func assertForecastSiteIdentity(t *testing.T, want, got forecastSite) { + t.Helper() + if got.IdentityPending || got.SiteID != want.SiteID || got.LearningRevision != want.LearningRevision || got.Revision != want.Revision || got.WeatherSinceMS != want.WeatherSinceMS { + t.Fatalf("forecast identity changed across config migration or restore:\ngot %+v\nwant %+v", got, want) + } +} + +func assertMigrationLearningState(t *testing.T, tracker *forecastTracker, signal string, cutoff int64, beforePV pvmodel.Model, beforeLoad loadmodel.Snapshot, beforeNativeOther []byte) { + t.Helper() + wantPeriods := forecastLearningPeriods{ConfigRevision: rustConfigRevision(tracker.site())} + if signal == "pv" { + wantPeriods.PVMS = cutoff + } else { + wantPeriods.LoadMS = cutoff + } + if tracker.learningPeriods != wantPeriods { + t.Fatalf("%s reset saved wrong learning periods: got %+v want %+v", signal, tracker.learningPeriods, wantPeriods) + } + other := "load" + if signal == "load" { + other = "pv" + } + if got := learningModel(t, tracker.candidate.(*rustForecast), other); !bytes.Equal(got, beforeNativeOther) { + t.Fatalf("%s reset changed native %s state", signal, other) + } + if status := tracker.LearningStatus(signal); status.Status != "cold_start" || status.StartedMS != cutoff || status.LatestTrainingMS != 0 { + t.Fatalf("%s reset status=%+v", signal, status) + } + if signal == "pv" { + got := tracker.pv.Model() + if got.Samples != 0 || got.LastMs != 0 || got.LearningStartedMS != cutoff { + t.Fatalf("PV model was not cleared at cutoff: %+v", got) + } + if load := tracker.load.Snapshot(); !reflect.DeepEqual(load, beforeLoad) { + t.Fatalf("PV reset changed load models:\ngot %+v\nwant %+v", load, beforeLoad) + } + return + } + if got := tracker.pv.Model(); !reflect.DeepEqual(got, beforePV) { + t.Fatalf("load reset changed PV model:\ngot %+v\nwant %+v", got, beforePV) + } + load := tracker.load.Snapshot() + if load.ActiveProfile != beforeLoad.ActiveProfile { + t.Fatalf("load reset changed active profile: got %s want %s", load.ActiveProfile, beforeLoad.ActiveProfile) + } + for _, profile := range loadmodel.Profiles() { + got := load.Profiles[profile] + if got.Samples != 0 || got.LastMs != 0 || got.MAE != 0 || got.HasTemperature || got.LearningStartedMS != cutoff { + t.Fatalf("load profile %s was not cleared at cutoff: %+v", profile, got) + } + if got.HeatingW_per_degC != 275 || got.ConfigRevision != beforeLoad.Profiles[profile].ConfigRevision || got.Timezone != beforeLoad.Profiles[profile].Timezone || got.PeakW != beforeLoad.Profiles[profile].PeakW { + t.Fatalf("load profile %s lost configured state: %+v", profile, got) + } + } +} diff --git a/go/cmd/ftw/forecast_learning_test.go b/go/cmd/ftw/forecast_learning_test.go new file mode 100644 index 00000000..a7de4b55 --- /dev/null +++ b/go/cmd/ftw/forecast_learning_test.go @@ -0,0 +1,221 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "reflect" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/energyforecast" + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/loadmodel" + "github.com/srcfl/ftw/go/internal/pvmodel" + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func learningTracker(st *state.Store, site *forecastSite, at time.Time, r *rustForecast) *forecastTracker { + tel := telemetry.NewStore() + f := &forecastTracker{store: st, tele: tel, site: func() forecastSite { return *site }, clock: func() time.Time { return at }, + pv: pvmodel.NewService(st, tel, func(time.Time) float64 { return 800 }, nil, 8000), load: loadmodel.NewService(st, tel, "site", 4000, 0)} + if r != nil { + f.candidate = r + } + return f +} + +func TestForecastLearningSerializesIdentityAcrossIntentAndApply(t *testing.T) { + st := hostForecastDB(t) + site := hostForecastSite() + at := time.Date(2026, 6, 15, 12, 22, 0, 0, time.UTC) + f := learningTracker(st, &site, at, nil) + f.configMu = &sync.RWMutex{} + + var siteReads int + var identityWriteCouldInterleave bool + f.site = func() forecastSite { + siteReads++ + if f.configMu.TryLock() { + identityWriteCouldInterleave = true + f.configMu.Unlock() + } + return site + } + f.requestReplan = func(string) { + if f.configMu.TryLock() { + identityWriteCouldInterleave = true + f.configMu.Unlock() + } + } + + if err := f.RestartLearning(context.Background(), "pv"); err != nil { + t.Fatal(err) + } + if identityWriteCouldInterleave { + t.Fatal("site identity write could interleave with learning restart") + } + if siteReads != 2 { + t.Fatalf("site reads=%d want2", siteReads) + } + if got := f.pv.LearningStartedMS(); got != at.UnixMilli() { + t.Fatalf("PV learning start=%d want%d", got, at.UnixMilli()) + } +} + +func TestForecastLearningNativePendingExchangeRecoversOnlySelectedSignal(t *testing.T) { + st := hostForecastDB(t) + r := learningNative(t, st) + site := hostForecastSite() + at := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + if err := r.Update(context.Background(), site, learningObservation(site, at), nil, false); err != nil { + t.Fatal(err) + } + f := learningTracker(st, &site, at.Add(22*time.Minute), r) + transport := r.transport.(energyforecast.RoundTripper) + var fail atomic.Bool + fail.Store(true) + r.client = energyforecast.NewClient(hostForecastExchange(func(ctx context.Context, p []byte) ([]byte, error) { + var header struct { + Action string `json:"action"` + } + _ = json.Unmarshal(p, &header) + if header.Action == "reset" && fail.Load() { + return nil, errors.New("injected worker exchange failure") + } + return transport.RoundTrip(ctx, p) + })) + if err := f.RestartLearning(context.Background(), "pv"); err == nil { + t.Fatal("failed worker reset reported success") + } + if _, ok := st.LoadConfig(forecastLearningKey); !ok { + t.Fatal("reset intent was not durable before worker call") + } + if s := f.LearningStatus("pv"); s.Status != "unavailable" { + t.Fatalf("pending PV=%+v", s) + } + if s := f.LearningStatus("load"); s.Status != "learning" { + t.Fatalf("pending PV hid load=%+v", s) + } + if f.pv.LearningStartedMS() != at.Add(22*time.Minute).UnixMilli() { + t.Fatal("legacy fallback retained old epoch") + } + fail.Store(false) + f.Snapshot(time.Now(), nil) + if s := f.LearningStatus("pv"); s.Status != "cold_start" { + t.Fatalf("automatic reconciliation left PV pending=%+v", s) + } + if s := f.LearningStatus("load"); s.Status != "learning" { + t.Fatalf("recovery changed load=%+v", s) + } + restored := learningNative(t, st) + if s := restored.LearningStatus(f.learningSiteLocked(site), "pv"); s.Status != "cold_start" { + t.Fatalf("recovered reset not persisted=%+v", s) + } +} + +func TestForecastLearningIdentityReadyAppliesSavedIntentBeforeCapture(t *testing.T) { + st := hostForecastDB(t) + site := hostForecastSite() + site.IdentityPending = true + at := time.Now().Add(-time.Minute) + cutoff := at.UnixMilli() + data, _ := json.Marshal(forecastLearningPeriods{ConfigRevision: rustConfigRevision(site), PVMS: cutoff}) + if err := st.SaveConfig(forecastLearningKey, string(data)); err != nil { + t.Fatal(err) + } + f := learningTracker(st, &site, time.Now(), nil) + if err := f.restoreLearning(context.Background()); err != nil { + t.Fatal(err) + } + if f.pv.LearningStartedMS() != 0 { + t.Fatal("identity-pending startup changed model") + } + site.IdentityPending = false + f.Snapshot(time.Now(), nil) + if f.pv.LearningStartedMS() != cutoff { + t.Fatal("first capture after identity resolution missed durable reset") + } + if f.load.LearningStartedMS() != 0 { + t.Fatal("PV reset changed load epoch") + } + restored := learningTracker(st, &site, time.Now(), nil) + if restored.pv.LearningStartedMS() != cutoff { + t.Fatal("legacy reset did not survive DB reload") + } +} + +func TestForecastLearningNativeDisabledPVDoesNotPersistIntentOrBlockLoad(t *testing.T) { + st := hostForecastDB(t) + r := learningNative(t, st) + site := hostForecastSite() + site.HasLocation = false + at := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + f := learningTracker(st, &site, at.Add(22*time.Minute), r) + if err := f.RestartLearning(context.Background(), "pv"); err == nil { + t.Fatal("disabled native PV reset accepted") + } + if _, ok := st.LoadConfig(forecastLearningKey); ok { + t.Fatal("disabled PV reset saved a blocking intent") + } + if s := f.LearningStatus("pv"); s.ResetAvailable { + t.Fatalf("disabled PV reset advertised: %+v", s) + } + // An older saved PV intent must not block updates in a load-only model. + site.PVLearningStartedMS = at.Add(10 * time.Minute).UnixMilli() + if err := r.Update(context.Background(), site, learningObservation(site, at), nil, false); err != nil { + t.Fatal(err) + } + if s := r.LearningStatus(site, "load"); s.Status != "learning" { + t.Fatalf("disabled PV blocked load=%+v", s) + } +} + +func TestForecastLearningCalibrationDropsPVAndNetRetainsLoadAndArchive(t *testing.T) { + base := time.Date(2026, 6, 1, 0, 0, 0, 0, time.UTC).UnixMilli() + var history []forecasting.ErrorSample + var truth []forecasting.Observation + for day := 0; day < 7; day++ { + for hour := 0; hour < 8; hour++ { + start := base + int64(day*24+hour)*3600000 + origin := start - 2*3600000 + p := forecasting.Point{StartMS: start, EndMS: start + 3600000, PVW: 2000, LoadW: 1000, PVKnown: true, LoadKnown: true, PVQuality: "ready", LoadQuality: "ready", PVBand: forecasting.Band{Method: forecasting.BandMethodColdStart}, LoadBand: forecasting.Band{Method: forecasting.BandMethodColdStart}, NetBand: forecasting.Band{Method: forecasting.BandMethodColdStart}} + e := forecasting.ErrorSample{Series: "energyplan", ConfigVersion: "cfg", IssueID: fmt.Sprintf("%d-%d", day, hour), OriginMS: origin, IssuedAtMS: origin, StartMS: start, EndMS: start + 3600000, AvailableAtMS: start + 3600000, Lead: forecasting.LeadBucket(origin, start), PVErrorW: float64(hour + day), LoadErrorW: float64(hour - day), PVKnown: true, LoadKnown: true, Prediction: p} + if err := e.Validate(); err != nil { + t.Fatal(err) + } + history = append(history, e) + truth = append(truth, forecasting.Observation{StartMS: start, EndMS: start + 3600000, PVKnown: true, LoadKnown: true}) + } + } + archived := append([]forecasting.ErrorSample(nil), history...) + rawTruth := append([]forecasting.Observation(nil), truth...) + cutoff := base + 8*24*3600000 + before := forecasting.NewCalibrator(history, "cfg", cutoff) + filtered, observations := learningEvidence(history, truth, cutoff, 0) + after := forecasting.NewCalibrator(filtered, "cfg", cutoff) + for _, signal := range []string{"pv", "net", "load"} { + previous := before.Band("energyplan", signal, cutoff+2*3600000, 1000) + got := after.Band("energyplan", signal, cutoff+2*3600000, 1000) + if previous.Method != forecasting.BandMethodEmpirical { + t.Fatalf("fixture not calibrated for %s: %+v", signal, previous) + } + if signal == "load" { + if got != previous { + t.Fatalf("PV reset changed load calibration: %+v %+v", previous, got) + } + } else if got.Samples != 0 { + t.Fatalf("%s retained pre-reset errors: %+v", signal, got) + } + } + if !reflect.DeepEqual(history, archived) || !reflect.DeepEqual(truth, rawTruth) { + t.Fatal("reset rewrote archived evidence") + } + if observations[0].PVKnown || !observations[0].LoadKnown { + t.Fatal("observation evidence reset crossed signals") + } +} diff --git a/go/cmd/ftw/forecast_rust.go b/go/cmd/ftw/forecast_rust.go index 3289c58c..1ab339a9 100644 --- a/go/cmd/ftw/forecast_rust.go +++ b/go/cmd/ftw/forecast_rust.go @@ -22,20 +22,27 @@ import ( const forecastRustStateKey = "forecast/energyplan_state_v1" type savedForecastState struct { - SiteID string `json:"site_id"` - ConfigRevision string `json:"config_revision"` - ModelRevision uint64 `json:"model_revision"` - LatestAvailableMS int64 `json:"latest_available_ms"` - State json.RawMessage `json:"state"` + SiteID string `json:"site_id"` + ConfigRevision string `json:"config_revision"` + ModelRevision uint64 `json:"model_revision"` + LatestAvailableMS int64 `json:"latest_available_ms"` + LatestTraining energyforecast.LatestInput `json:"latest_training_ms"` + PVLearningStartedMS int64 `json:"pv_learning_started_ms,omitempty"` + LoadLearningStartedMS int64 `json:"load_learning_started_ms,omitempty"` + State json.RawMessage `json:"state"` } type rustForecast struct { - version string - client *energyforecast.Client - transport interface{ Close() error } - store *state.Store - mu sync.RWMutex - saved savedForecastState + version string + client *energyforecast.Client + transport interface{ Close() error } + store *state.Store + mu sync.RWMutex + updateMu sync.Mutex // serialize complete update/reset exchanges and persistence + saved savedForecastState + resetSupported bool + pvQuality, loadQuality string + predictedTraining energyforecast.LatestInput } func newRustForecast(st *state.Store, binary string) (*rustForecast, error) { @@ -47,9 +54,10 @@ func newRustForecast(st *state.Store, binary string) (*rustForecast, error) { defer cancel() line, err := transport.RoundTrip(ctx, []byte(`{"type":"handshake","protocol_version":1}`)) var info struct { - Name string `json:"name"` - Version string `json:"version"` - ForecastVersion int `json:"forecast_protocol_version"` + Name string `json:"name"` + Version string `json:"version"` + ForecastVersion int `json:"forecast_protocol_version"` + Features []string `json:"features"` } if err == nil { err = json.Unmarshal(line, &info) @@ -59,6 +67,11 @@ func newRustForecast(st *state.Store, binary string) (*rustForecast, error) { return nil, fmt.Errorf("Energyplan forecast v1 unavailable: name=%q version=%d: %v", info.Name, info.ForecastVersion, err) } r := &rustForecast{version: forecastBinaryIdentity(binary), client: energyforecast.NewClient(transport), transport: transport, store: st} + for _, feature := range info.Features { + if feature == "forecast_reset" { + r.resetSupported = true + } + } if data, ok := st.LoadConfig(forecastRustStateKey); ok { if len(data) <= energyforecast.MaxStateBytes+4096 { var saved savedForecastState @@ -106,14 +119,19 @@ func rustFeatures(t time.Time, site forecastSite, home bool, row *state.Forecast } func (r *rustForecast) Update(ctx context.Context, site forecastSite, o forecasting.Observation, weather *state.ForecastPoint, away bool) error { + r.updateMu.Lock() + defer r.updateMu.Unlock() + if err := r.applyLearningPeriodsLocked(ctx, site, max(time.Now().UnixMilli(), o.AvailableAtMS)); err != nil { + return err + } r.mu.RLock() saved := r.saved r.mu.RUnlock() if saved.SiteID != site.SiteID || saved.ConfigRevision != rustConfigRevision(site) { saved = savedForecastState{SiteID: site.SiteID, ConfigRevision: rustConfigRevision(site)} } - origin := o.AvailableAtMS - if weather != nil && weather.FetchedAtMs > origin { + origin := max(o.AvailableAtMS, saved.LatestAvailableMS) + if weather != nil && weather.FetchedAtMs > o.AvailableAtMS { return errors.New("observation weather arrived after update origin") } input := energyforecast.Observation{Interval: energyforecast.Interval{ValidStartMs: o.StartMS, ValidEndMs: o.EndMS}, @@ -133,8 +151,9 @@ func (r *rustForecast) Update(ctx context.Context, site forecastSite, o forecast if err != nil { return err } - next := savedForecastState{SiteID: site.SiteID, ConfigRevision: rustConfigRevision(site), ModelRevision: reply.ModelRevision, - LatestAvailableMS: origin, State: reply.State} + next := saved + next.ModelRevision, next.LatestAvailableMS, next.State = reply.ModelRevision, origin, reply.State + next.LatestTraining = reply.LatestTrainingMs data, err := json.Marshal(next) if err != nil { return err @@ -144,6 +163,10 @@ func (r *rustForecast) Update(ctx context.Context, site forecastSite, o forecast return err } r.mu.Lock() + if r.saved.SiteID != next.SiteID || r.saved.ConfigRevision != next.ConfigRevision { + r.pvQuality, r.loadQuality = "", "" + r.predictedTraining = energyforecast.LatestInput{} + } r.saved = next r.mu.Unlock() return nil @@ -210,6 +233,17 @@ func (r *rustForecast) Predict(ctx context.Context, site forecastSite, issued fo if err != nil { return forecasting.Issue{}, err } + // A durable reset intent can outlive a failed worker exchange. Until that + // signal has restarted, only its freshly reset legacy fallback may serve it. + for i := range reply.Predictions { + if saved.PVLearningStartedMS < site.PVLearningStartedMS { + reply.Predictions[i].PV = &energyforecast.Estimate{Quality: "cold_start", Uncertainty: "unknown"} + } + if saved.LoadLearningStartedMS < site.LoadLearningStartedMS { + reply.Predictions[i].Load = &energyforecast.Estimate{Quality: "cold_start", Uncertainty: "unknown"} + } + } + r.recordLearningQuality(saved, reply) now := time.Now().UnixMilli() out := forecasting.Issue{Schema: forecasting.Schema, ID: uuid.NewString(), DecisionID: issued.DecisionID, OriginMS: issued.OriginMS, IssuedAtMS: now, ConfigVersion: issued.ConfigVersion, Site: forecastSiteContext(site), Weather: issued.Weather, diff --git a/go/cmd/ftw/forecast_rust_learning.go b/go/cmd/ftw/forecast_rust_learning.go new file mode 100644 index 00000000..4685f03d --- /dev/null +++ b/go/cmd/ftw/forecast_rust_learning.go @@ -0,0 +1,151 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "time" + + "github.com/google/uuid" + "github.com/srcfl/ftw/go/internal/energyforecast" + "github.com/srcfl/ftw/go/internal/forecasting" +) + +func (r *rustForecast) RestartLearning(ctx context.Context, site forecastSite, signal string) error { + if signal == "pv" { + site.LoadLearningStartedMS = 0 + } else { + site.PVLearningStartedMS = 0 + } + r.updateMu.Lock() + defer r.updateMu.Unlock() + return r.applyLearningPeriodsLocked(ctx, site, time.Now().UnixMilli()) +} + +// Caller holds updateMu across the request, durable save and memory update. +func (r *rustForecast) applyLearningPeriodsLocked(ctx context.Context, site forecastSite, origin int64) error { + for _, signal := range []string{"pv", "load"} { + if signal == "pv" && !site.HasLocation { + continue + } + cutoff := site.PVLearningStartedMS + if signal == "load" { + cutoff = site.LoadLearningStartedMS + } + if cutoff == 0 { + continue + } + r.mu.RLock() + saved := r.saved + r.mu.RUnlock() + if saved.SiteID != site.SiteID || saved.ConfigRevision != rustConfigRevision(site) { + saved = savedForecastState{SiteID: site.SiteID, ConfigRevision: rustConfigRevision(site)} + } + applied := saved.PVLearningStartedMS + if signal == "load" { + applied = saved.LoadLearningStartedMS + } + if applied >= cutoff { + continue + } + if !r.resetSupported { + return errors.New("forecast worker does not support restarting learning") + } + origin = max(origin, saved.LatestAvailableMS, cutoff) + reply, err := r.client.Reset(ctx, energyforecast.ResetRequest{RequestContext: energyforecast.RequestContext{ + RequestID: uuid.NewString(), SiteID: site.SiteID, ConfigRevision: rustConfigRevision(site), OriginMs: origin, + Config: rustForecastConfig(site), State: saved.State}, Signal: signal, LearningStartedMs: cutoff}) + if err != nil { + return err + } + next := saved + next.ModelRevision, next.LatestAvailableMS, next.State = reply.ModelRevision, origin, reply.State + next.LatestTraining = reply.LatestTrainingMs + if signal == "pv" { + next.PVLearningStartedMS = reply.LearningStartedMs + } else { + next.LoadLearningStartedMS = reply.LearningStartedMs + } + data, err := json.Marshal(next) + if err != nil { + return err + } + if err = r.store.SaveConfig(forecastRustStateKey, string(data)); err != nil { + return err + } + r.mu.Lock() + r.saved = next + if signal == "pv" { + r.pvQuality = "cold_start" + r.predictedTraining.PV = nil + } else { + r.loadQuality = "cold_start" + r.predictedTraining.Load = nil + } + r.mu.Unlock() + } + return nil +} + +func (r *rustForecast) LearningStatus(site forecastSite, signal string) forecasting.LearningStatus { + r.mu.RLock() + defer r.mu.RUnlock() + status := forecasting.LearningStatus{Engine: "energyplan", Status: "cold_start", ResetAvailable: r.resetSupported && !site.IdentityPending} + if signal == "pv" && !site.HasLocation { + status.Status, status.ResetAvailable = "unavailable", false + return status + } + cutoff, applied, latest, quality := site.PVLearningStartedMS, r.saved.PVLearningStartedMS, r.saved.LatestTraining.PV, r.pvQuality + if signal == "load" { + cutoff, applied, latest, quality = site.LoadLearningStartedMS, r.saved.LoadLearningStartedMS, r.saved.LatestTraining.Load, r.loadQuality + } + predictedLatest := r.predictedTraining.PV + if signal == "load" { + predictedLatest = r.predictedTraining.Load + } + if predictedLatest != nil && (latest == nil || *predictedLatest > *latest) { + latest = predictedLatest + } + status.StartedMS = cutoff + if site.IdentityPending { + status.Status = "unavailable" + return status + } + if r.saved.SiteID != site.SiteID || r.saved.ConfigRevision != rustConfigRevision(site) { + applied, latest, quality = 0, nil, "cold_start" + } + if applied < cutoff { + status.Status = "unavailable" + return status + } + if latest != nil && *latest > cutoff { + status.LatestTrainingMS, status.Status = *latest, "learning" + if quality == "ready" { + status.Status = "ready" + } + } + return status +} + +func (r *rustForecast) recordLearningQuality(saved savedForecastState, reply energyforecast.PredictReply) { + r.mu.Lock() + defer r.mu.Unlock() + if r.saved.SiteID != saved.SiteID || r.saved.ConfigRevision != saved.ConfigRevision || r.saved.ModelRevision != saved.ModelRevision { + return + } + pv, load := "ready", "ready" + for _, p := range reply.Predictions { + if p.PV == nil || !p.PV.Known { + pv = leastForecastQuality(pv, "cold_start") + } else { + pv = leastForecastQuality(pv, p.PV.Quality) + } + if p.Load == nil || !p.Load.Known { + load = leastForecastQuality(load, "cold_start") + } else { + load = leastForecastQuality(load, p.Load.Quality) + } + } + r.pvQuality, r.loadQuality = pv, load + r.predictedTraining = reply.LatestTrainingMs +} diff --git a/go/cmd/ftw/forecast_rust_learning_test.go b/go/cmd/ftw/forecast_rust_learning_test.go new file mode 100644 index 00000000..ff4085cb --- /dev/null +++ b/go/cmd/ftw/forecast_rust_learning_test.go @@ -0,0 +1,183 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "os" + "sync" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/energyforecast" + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/state" +) + +func learningNative(t *testing.T, st *state.Store) *rustForecast { + t.Helper() + binary := os.Getenv("FTW_FORECAST_WORKER") + if binary == "" { + t.Skip("set FTW_FORECAST_WORKER for native learning integration") + } + r, err := newRustForecast(st, binary) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = r.Close() }) + return r +} + +func learningObservation(site forecastSite, start time.Time) forecasting.Observation { + return forecasting.Observation{StartMS: start.UnixMilli(), EndMS: start.Add(15 * time.Minute).UnixMilli(), AvailableAtMS: start.Add(15 * time.Minute).UnixMilli(), PVW: 2350.125, PVKnown: true, LoadW: 677.375, LoadKnown: true, Quality: "complete", ConfigVersion: site.Revision} +} + +func learningModel(t *testing.T, r *rustForecast, signal string) json.RawMessage { + t.Helper() + var saved savedForecastState + if err := json.Unmarshal(r.Snapshot(), &saved); err != nil { + t.Fatal(err) + } + var models map[string]json.RawMessage + if err := json.Unmarshal(saved.State, &models); err != nil { + t.Fatal(err) + } + return models[signal] +} + +func TestForecastLearningNativeSelectiveResetAndDurableRecovery(t *testing.T) { + for _, signal := range []string{"pv", "load"} { + t.Run(signal, func(t *testing.T) { + st := hostForecastDB(t) + r := learningNative(t, st) + site := hostForecastSite() + start := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + if err := r.Update(context.Background(), site, learningObservation(site, start), nil, false); err != nil { + t.Fatal(err) + } + other := "load" + if signal == "load" { + other = "pv" + } + originalOther := learningModel(t, r, other) + cutoff := start.Add(22 * time.Minute).UnixMilli() + if signal == "pv" { + site.PVLearningStartedMS = cutoff + } else { + site.LoadLearningStartedMS = cutoff + } + if err := r.RestartLearning(context.Background(), site, signal); err != nil { + t.Fatal(err) + } + if !bytes.Equal(originalOther, learningModel(t, r, other)) { + t.Fatal("reset changed other model") + } + if s := r.LearningStatus(site, signal); s.Status != "cold_start" || s.LatestTrainingMS != 0 { + t.Fatalf("reset status=%+v", s) + } + cleared := learningModel(t, r, signal) + restored := learningNative(t, st) + if !bytes.Equal(r.Snapshot(), restored.Snapshot()) { + t.Fatal("DB restore changed state") + } + // Replaying an already-consumed pre-reset interval cannot repopulate + // the cleared signal or publish a partial update of the other model. + beforeOld := restored.Snapshot() + if err := restored.Update(context.Background(), site, learningObservation(site, start), nil, false); err == nil { + t.Fatal("duplicate pre-reset interval unexpectedly accepted") + } + if !bytes.Equal(beforeOld, restored.Snapshot()) { + t.Fatal("failed pre-reset replay changed durable state") + } + // A delayed quarter straddles the reset cutoff. Only the other signal learns. + if err := restored.Update(context.Background(), site, learningObservation(site, start.Add(15*time.Minute)), nil, false); err != nil { + t.Fatal(err) + } + if !bytes.Equal(cleared, learningModel(t, restored, signal)) { + t.Fatal("straddling interval repopulated reset model") + } + if bytes.Equal(originalOther, learningModel(t, restored, other)) { + t.Fatal("reset blocked other signal's next interval") + } + eligible := learningObservation(site, start.Add(30*time.Minute)) + futureWeather := &state.ForecastPoint{FetchedAtMs: eligible.AvailableAtMS + 1, SolarWm2: hostForecastPtr(500.0)} + beforeWeather := restored.Snapshot() + if err := restored.Update(context.Background(), site, eligible, futureWeather, false); err == nil { + t.Fatal("newer request origin admitted weather unavailable to the observation") + } + if !bytes.Equal(beforeWeather, restored.Snapshot()) { + t.Fatal("rejected weather changed state") + } + if err := restored.Update(context.Background(), site, learningObservation(site, start.Add(30*time.Minute)), nil, false); err != nil { + t.Fatal(err) + } + if s := restored.LearningStatus(site, signal); s.Status != "learning" || s.LatestTrainingMS != start.Add(45*time.Minute).UnixMilli() { + t.Fatalf("post-reset learning=%+v", s) + } + learned := restored.Snapshot() + if err := restored.RestartLearning(context.Background(), site, signal); err != nil { + t.Fatal(err) + } + if !bytes.Equal(learned, restored.Snapshot()) { + t.Fatal("durable reset replay erased new evidence") + } + }) + } +} + +func TestForecastLearningNativeInFlightUpdateCannotUndoReset(t *testing.T) { + st := hostForecastDB(t) + r := learningNative(t, st) + site := hostForecastSite() + start := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + transport := r.transport.(energyforecast.RoundTripper) + entered, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + r.client = energyforecast.NewClient(hostForecastExchange(func(ctx context.Context, p []byte) ([]byte, error) { + reply, err := transport.RoundTrip(ctx, p) + var header struct { + Action string `json:"action"` + } + _ = json.Unmarshal(p, &header) + if header.Action == "update" { + once.Do(func() { + close(entered) + select { + case <-release: + case <-ctx.Done(): + } + }) + } + return reply, err + })) + done := make(chan error, 1) + go func() { done <- r.Update(context.Background(), site, learningObservation(site, start), nil, false) }() + select { + case <-entered: + case <-time.After(5 * time.Second): + t.Fatal("update never reached persistence boundary") + } + site.PVLearningStartedMS = start.Add(20 * time.Minute).UnixMilli() + resetDone := make(chan error, 1) + go func() { resetDone <- r.RestartLearning(context.Background(), site, "pv") }() + select { + case err := <-resetDone: + close(release) + t.Fatalf("reset bypassed in-flight update lock: %v", err) + case <-time.After(30 * time.Millisecond): + } + close(release) + if err := <-done; err != nil { + t.Fatal(err) + } + if err := <-resetDone; err != nil { + t.Fatal(err) + } + reloaded := learningNative(t, st) + if s := reloaded.LearningStatus(site, "pv"); s.Status != "cold_start" || s.LatestTrainingMS != 0 { + t.Fatalf("old update restored PV: %+v", s) + } + if s := reloaded.LearningStatus(site, "load"); s.LatestTrainingMS != start.Add(15*time.Minute).UnixMilli() { + t.Fatalf("load history lost: %+v", s) + } +} diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index 47c1a3df..1e272089 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -22,17 +22,19 @@ import ( ) type forecastSite struct { - IdentityPending bool - WeatherSinceMS int64 - LearningRevision string - Revision string - SiteID string - Meter string - Latitude, Longitude float64 - HasLocation bool - HasPVScale bool - Timezone string - Options telemetry.ForecastOptions + IdentityPending bool + WeatherSinceMS int64 + LearningRevision string + Revision string + SiteID string + Meter string + Latitude, Longitude float64 + HasLocation bool + HasPVScale bool + Timezone string + Options telemetry.ForecastOptions + PVLearningStartedMS int64 + LoadLearningStartedMS int64 } // forecastCandidate is a host adapter around the compiled model. State is @@ -70,9 +72,16 @@ type forecastTracker struct { pvAccumulator telemetry.ForecastAccumulator stopped bool clock func() time.Time + learningMu sync.RWMutex + learningPeriods forecastLearningPeriods + learningErrors map[string]error + requestReplan func(string) } func (f *forecastTracker) Start(ctx context.Context) error { + if err := f.restoreLearning(ctx); err != nil { + return err + } initCtx, cancel := context.WithTimeout(ctx, 3*time.Second) err := f.store.InitForecastArchive(initCtx) cancel() @@ -167,6 +176,7 @@ func (f *forecastTracker) observe(ctx context.Context, now time.Time) { if f.refreshIdentity != nil { f.refreshIdentity() } + f.reconcileLearning(ctx) if f.configMu != nil { f.configMu.RLock() } @@ -205,7 +215,10 @@ func (f *forecastTracker) observe(ctx context.Context, now time.Time) { } away := f.away != nil && f.away(time.UnixMilli(o.StartMS)) workCtx, cancel := context.WithTimeout(ctx, 2*time.Second) + f.learningMu.RLock() + site = f.learningSiteLocked(site) err = f.candidate.Update(workCtx, site, o, weather, away) + f.learningMu.RUnlock() cancel() if err != nil { slog.Debug("forecast candidate update unavailable", "err", err) @@ -356,12 +369,15 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m if f.refreshIdentity != nil { f.refreshIdentity() } + f.reconcileLearning(context.Background()) if f.configMu != nil { f.configMu.RLock() defer f.configMu.RUnlock() } + f.learningMu.RLock() + defer f.learningMu.RUnlock() captureAt := f.now() - site := f.site() + site := f.learningSiteLocked(f.site()) if site.IdentityPending { // Keep persisted learning untouched until the running hardware proves its // binding. An explicit empty weather slice prevents fallback to old rows. @@ -374,6 +390,7 @@ func (f *forecastTracker) Snapshot(_ time.Time, weather []state.ForecastPoint) m history := f.errors observations := f.observations f.mu.RUnlock() + history, observations = learningEvidence(history, observations, site.PVLearningStartedMS, site.LoadLearningStartedMS) candidateState := json.RawMessage(nil) if f.candidate != nil { candidateState = append(json.RawMessage(nil), f.candidate.Snapshot()...) diff --git a/go/cmd/ftw/main.go b/go/cmd/ftw/main.go index d40981be..a5243f54 100644 --- a/go/cmd/ftw/main.go +++ b/go/cmd/ftw/main.go @@ -1698,6 +1698,7 @@ func main() { } if forecastTrackerSvc != nil { mpcSvc.ForecastSnapshot = forecastTrackerSvc.Snapshot + forecastTrackerSvc.setReplan(mpcSvc.RequestReplan) } mpcSvc.Start(ctx) defer mpcSvc.Stop() @@ -2375,6 +2376,7 @@ func main() { PlannerPrefs: plannerPrefs, PVModel: pvSvc, LoadModel: loadSvc, + ForecastLearning: forecastTrackerSvc, Loadpoints: lpMgr, LoadpointCtrl: lpController, OCPPChargers: ocppChargersFn, diff --git a/go/internal/api/api.go b/go/internal/api/api.go index b8e15425..7e60f888 100644 --- a/go/internal/api/api.go +++ b/go/internal/api/api.go @@ -149,7 +149,8 @@ type Deps struct { PVModel *pvmodel.Service // Optional: load digital-twin self-learner. - LoadModel *loadmodel.Service + LoadModel *loadmodel.Service + ForecastLearning ForecastLearning // Optional: EV loadpoint state consumed by the API and MPC. Loadpoints *loadpoint.Manager @@ -2935,16 +2936,12 @@ func (s *Server) handlePVModel(w http.ResponseWriter, r *http.Request) { "pv_residual_mean_w": rd.MeanW, "pv_residual_std_w": rd.StdW, "pv_residual_window_minutes": rd.WindowMinutes, + "learning": s.forecastLearningStatus("pv"), }) } func (s *Server) handlePVModelReset(w http.ResponseWriter, r *http.Request) { - if s.deps.PVModel == nil { - writeJSON(w, 400, map[string]string{"error": "pvmodel disabled"}) - return - } - s.deps.PVModel.Reset() - writeJSON(w, 200, map[string]string{"status": "reset"}) + s.handleForecastLearningReset(w, r, "pv") } // ---- static ---- diff --git a/go/internal/api/api_forecast_learning.go b/go/internal/api/api_forecast_learning.go new file mode 100644 index 00000000..fe6415e2 --- /dev/null +++ b/go/internal/api/api_forecast_learning.go @@ -0,0 +1,40 @@ +package api + +import ( + "context" + "errors" + "net/http" + + "github.com/srcfl/ftw/go/internal/forecasting" +) + +// ForecastLearning restarts the primary, fallback and calibration together. +// It owns the durable learning boundary; the API never edits model state. +type ForecastLearning interface { + LearningStatus(string) forecasting.LearningStatus + RestartLearning(context.Context, string) error +} + +func (s *Server) forecastLearningStatus(signal string) forecasting.LearningStatus { + if s.deps.ForecastLearning == nil { + return forecasting.LearningStatus{Engine: "legacy", Status: "unavailable"} + } + return s.deps.ForecastLearning.LearningStatus(signal) +} + +func (s *Server) handleForecastLearningReset(w http.ResponseWriter, r *http.Request, signal string) { + if s.deps.ForecastLearning == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "forecast learning service unavailable"}) + return + } + if err := s.deps.ForecastLearning.RestartLearning(r.Context(), signal); err != nil { + status := "error" + var pending *forecasting.LearningRestartPendingError + if errors.As(err, &pending) { + status = "pending" + } + writeJSON(w, http.StatusServiceUnavailable, map[string]any{"status": status, "error": err.Error(), "learning": s.forecastLearningStatus(signal)}) + return + } + writeJSON(w, http.StatusOK, map[string]any{"status": "reset", "learning": s.forecastLearningStatus(signal)}) +} diff --git a/go/internal/api/api_forecast_learning_test.go b/go/internal/api/api_forecast_learning_test.go new file mode 100644 index 00000000..a08ece3f --- /dev/null +++ b/go/internal/api/api_forecast_learning_test.go @@ -0,0 +1,183 @@ +package api + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/loadmodel" + "github.com/srcfl/ftw/go/internal/pvmodel" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +type forecastLearningAPIMock struct { + status map[string]forecasting.LearningStatus + err error + calls []string +} + +func (m *forecastLearningAPIMock) LearningStatus(signal string) forecasting.LearningStatus { + return m.status[signal] +} + +func (m *forecastLearningAPIMock) RestartLearning(_ context.Context, signal string) error { + m.calls = append(m.calls, signal) + if m.err != nil { + return m.err + } + status := m.status[signal] + status.Status = "learning" + status.StartedMS += 1000 + m.status[signal] = status + return nil +} + +type forecastLearningAPIResponse struct { + Status string `json:"status"` + Error string `json:"error"` + Learning forecasting.LearningStatus `json:"learning"` +} + +func decodeForecastLearningResponse(t *testing.T, rr *httptest.ResponseRecorder) forecastLearningAPIResponse { + t.Helper() + var got forecastLearningAPIResponse + if err := json.Unmarshal(rr.Body.Bytes(), &got); err != nil { + t.Fatalf("decode response: %v; body: %s", err, rr.Body.String()) + } + return got +} + +func TestForecastLearningEndpointsDelegateExactSignalAndReportCurrentStatus(t *testing.T) { + learning := &forecastLearningAPIMock{status: map[string]forecasting.LearningStatus{ + "pv": {Engine: "energyplan", Status: "ready", StartedMS: 100, LatestTrainingMS: 150, ResetAvailable: true}, + "load": {Engine: "energyplan", Status: "ready", StartedMS: 200, LatestTrainingMS: 250, ResetAvailable: true}, + }} + srv := New(&Deps{ + PVModel: pvmodel.NewService(nil, nil, nil, nil, 5000), + LoadModel: loadmodel.NewService(nil, telemetry.NewStore(), "site", 4000, 0), + ForecastLearning: learning, + }) + + for _, tc := range []struct { + signal string + path string + }{ + {signal: "pv", path: "/api/pvmodel"}, + {signal: "load", path: "/api/loadmodel"}, + } { + t.Run(tc.signal, func(t *testing.T) { + get := httptest.NewRecorder() + srv.Handler().ServeHTTP(get, httptest.NewRequest(http.MethodGet, tc.path, nil)) + if get.Code != http.StatusOK { + t.Fatalf("GET status = %d, body: %s", get.Code, get.Body.String()) + } + before := decodeForecastLearningResponse(t, get) + if before.Learning != learning.status[tc.signal] { + t.Fatalf("GET learning = %+v, want %+v", before.Learning, learning.status[tc.signal]) + } + + post := httptest.NewRecorder() + srv.Handler().ServeHTTP(post, httptest.NewRequest(http.MethodPost, tc.path+"/reset", nil)) + if post.Code != http.StatusOK { + t.Fatalf("POST status = %d, body: %s", post.Code, post.Body.String()) + } + after := decodeForecastLearningResponse(t, post) + if after.Status != "reset" || after.Learning != learning.status[tc.signal] { + t.Fatalf("POST response = %+v, current status = %+v", after, learning.status[tc.signal]) + } + }) + } + if got := strings.Join(learning.calls, ","); got != "pv,load" { + t.Fatalf("restart signals = %q, want pv,load", got) + } +} + +func TestForecastLearningResetRefusesWithoutCoordinatorAndKeepsProfileActionSeparate(t *testing.T) { + pv := pvmodel.NewService(nil, nil, nil, nil, 5000) + pv.Residuals.Add(timeForForecastLearningTest(), 1000, 1500) + load := loadmodel.NewService(nil, telemetry.NewStore(), "site", 4000, 0) + srv := New(&Deps{PVModel: pv, LoadModel: load}) + profile := httptest.NewRecorder() + profileRequest := httptest.NewRequest(http.MethodPost, "/api/loadmodel/profile", strings.NewReader(`{"profile":"away"}`)) + profileRequest.Header.Set("Content-Type", "application/json") + srv.Handler().ServeHTTP(profile, profileRequest) + if profile.Code != http.StatusOK { + t.Fatalf("profile action status = %d, body: %s", profile.Code, profile.Body.String()) + } + + for _, path := range []string{"/api/pvmodel/reset", "/api/loadmodel/reset"} { + rr := httptest.NewRecorder() + srv.Handler().ServeHTTP(rr, httptest.NewRequest(http.MethodPost, path, nil)) + if rr.Code != http.StatusServiceUnavailable { + t.Fatalf("%s status = %d, want 503; body: %s", path, rr.Code, rr.Body.String()) + } + got := decodeForecastLearningResponse(t, rr) + if got.Status == "reset" || got.Error == "" { + t.Fatalf("%s claimed success: %+v", path, got) + } + } + if pv.Residuals.Len() != 1 { + t.Fatal("PV legacy state was reset without the coordinator") + } + if load.Profile() != loadmodel.ProfileAway { + t.Fatal("load learning reset changed the separate profile action") + } +} + +func TestForecastLearningPersistenceFailureReturns503WithoutLegacyReset(t *testing.T) { + pv := pvmodel.NewService(nil, nil, nil, nil, 5000) + pv.Residuals.Add(timeForForecastLearningTest(), 1000, 1500) + wantErr := errors.New("persist learning boundary: disk full") + learning := &forecastLearningAPIMock{ + status: map[string]forecasting.LearningStatus{"pv": {Engine: "energyplan", Status: "error", ResetAvailable: true}}, + err: wantErr, + } + srv := New(&Deps{PVModel: pv, ForecastLearning: learning}) + + rr := httptest.NewRecorder() + srv.Handler().ServeHTTP(rr, httptest.NewRequest(http.MethodPost, "/api/pvmodel/reset", nil)) + + if rr.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d, want 503; body: %s", rr.Code, rr.Body.String()) + } + got := decodeForecastLearningResponse(t, rr) + if got.Status != "error" || !strings.Contains(got.Error, wantErr.Error()) { + t.Fatalf("failure response = %+v", got) + } + if pv.Residuals.Len() != 1 { + t.Fatal("handler ran a legacy-only reset after coordinator failure") + } + if len(learning.calls) != 1 || learning.calls[0] != "pv" { + t.Fatalf("restart calls = %v", learning.calls) + } +} + +func TestForecastLearningDurablePendingErrorReturns503Pending(t *testing.T) { + cause := errors.New("load model save failed") + learning := &forecastLearningAPIMock{ + status: map[string]forecasting.LearningStatus{"load": {Engine: "energyplan", Status: "pending", StartedMS: 1234, ResetAvailable: true}}, + err: &forecasting.LearningRestartPendingError{Err: cause}, + } + srv := New(&Deps{LoadModel: loadmodel.NewService(nil, telemetry.NewStore(), "site", 4000, 0), ForecastLearning: learning}) + + rr := httptest.NewRecorder() + srv.Handler().ServeHTTP(rr, httptest.NewRequest(http.MethodPost, "/api/loadmodel/reset", nil)) + + if rr.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d, want 503; body: %s", rr.Code, rr.Body.String()) + } + got := decodeForecastLearningResponse(t, rr) + if got.Status != "pending" || got.Learning != learning.status["load"] || !strings.Contains(got.Error, cause.Error()) { + t.Fatalf("pending response = %+v", got) + } +} + +func timeForForecastLearningTest() time.Time { + return time.Unix(1, 0) +} diff --git a/go/internal/api/api_loadmodel.go b/go/internal/api/api_loadmodel.go index 20766480..bb412a3f 100644 --- a/go/internal/api/api_loadmodel.go +++ b/go/internal/api/api_loadmodel.go @@ -44,6 +44,7 @@ func (s *Server) handleLoadModel(w http.ResponseWriter, r *http.Request) { "heating_w_per_degc": stats.HeatingWPerDegC, "buckets_warm": stats.BucketsWarm, "buckets_total": stats.BucketsTotal, + "learning": s.forecastLearningStatus("load"), }) } @@ -80,16 +81,7 @@ func (s *Server) handleLoadModelProfile(w http.ResponseWriter, r *http.Request) } func (s *Server) handleLoadModelReset(w http.ResponseWriter, r *http.Request) { - if s.deps.LoadModel == nil { - writeJSON(w, 400, map[string]string{"error": "loadmodel disabled"}) - return - } - profile := s.deps.LoadModel.Profile() - s.deps.LoadModel.Reset() - if s.deps.MPC != nil { - s.deps.MPC.ReplanWithReason(r.Context(), "load_profile_reset") - } - writeJSON(w, 200, map[string]any{"status": "reset", "profile": profile}) + s.handleForecastLearningReset(w, r, "load") } func loadModelStatsFrom(m loadmodel.Model) loadModelStats { diff --git a/go/internal/energyforecast/native_reset_test.go b/go/internal/energyforecast/native_reset_test.go new file mode 100644 index 00000000..36f9a1be --- /dev/null +++ b/go/internal/energyforecast/native_reset_test.go @@ -0,0 +1,121 @@ +package energyforecast_test + +import ( + "context" + "os" + "path/filepath" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/energyforecast" + "github.com/srcfl/ftw/go/internal/mpc" +) + +func TestNativeForecastResetIsSignalScopedAndDurable(t *testing.T) { + binary := os.Getenv("FTW_FORECAST_WORKER") + if binary == "" { + t.Skip("set FTW_FORECAST_WORKER to an Energyplan worker with forecast reset support") + } + newClient := func() (*energyforecast.Client, *mpc.ProcessTransport) { + transport, err := mpc.NewProcessTransport(mpc.ProcessTransportConfig{Command: []string{binary}, ModuleDir: filepath.Dir(binary)}) + if err != nil { + t.Fatal(err) + } + return energyforecast.NewClient(transport), transport + } + start := time.Date(2026, 6, 15, 12, 0, 0, 0, time.UTC) + quarter := 15 * time.Minute + end := start.Add(quarter) + features := func(at time.Time) energyforecast.Features { + ghi := 500.0 + available := start.UnixMilli() + return energyforecast.Features{ + LocalDay: at.Unix() / 86400, LocalWeekday: (int(at.Weekday()) + 6) % 7, + LocalMinute: at.Hour()*60 + at.Minute(), GHIWm2: &ghi, WeatherAvailableAtMs: &available, + } + } + config := energyforecast.Config{PV: &energyforecast.PVConfig{LatitudeDeg: 57, LongitudeDeg: 15}, Load: &energyforecast.LoadConfig{}} + observation := func(at time.Time, pv, load float64) energyforecast.Observation { + return energyforecast.Observation{ + Interval: energyforecast.Interval{ValidStartMs: at.UnixMilli(), ValidEndMs: at.Add(quarter).UnixMilli()}, + Features: features(at), AvailableAtMs: at.Add(quarter).UnixMilli(), + PVAvailableW: &pv, HouseholdLoadW: &load, PVQuality: energyforecast.QualityGood, LoadQuality: energyforecast.QualityGood, + } + } + for _, selected := range []string{"pv", "load"} { + t.Run(selected, func(t *testing.T) { + client, transport := newClient() + meta := energyforecast.RequestContext{RequestID: "native-reset-train-" + selected, SiteID: "native-reset-site", ConfigRevision: "v1", OriginMs: end.UnixMilli(), Config: config} + pv, load := 2300.0, 470.0 + trained, err := client.Update(context.Background(), energyforecast.UpdateRequest{RequestContext: meta, Observations: []energyforecast.Observation{observation(start, pv, load)}}) + if err != nil { + t.Fatal(err) + } + if trained.LatestTrainingMs.PV == nil || trained.LatestTrainingMs.Load == nil { + t.Fatalf("training clocks = %#v", trained.LatestTrainingMs) + } + cutoff := end.UnixMilli() + meta.RequestID = "native-reset-" + selected + meta.OriginMs = cutoff + meta.State = trained.State + reset, err := client.Reset(context.Background(), energyforecast.ResetRequest{RequestContext: meta, Signal: selected, LearningStartedMs: cutoff}) + if err != nil { + t.Fatal(err) + } + if reset.Signal != selected || reset.LearningStartedMs != cutoff || len(reset.State) == 0 { + t.Fatalf("reset response = %#v", reset) + } + if reset.ModelRevision != trained.ModelRevision+1 { + t.Fatalf("reset revision = %d, want %d", reset.ModelRevision, trained.ModelRevision+1) + } + if selected == "pv" { + if reset.LatestTrainingMs.PV != nil || reset.LatestTrainingMs.Load == nil { + t.Fatalf("PV reset clocks = %#v", reset.LatestTrainingMs) + } + } else if reset.LatestTrainingMs.Load != nil || reset.LatestTrainingMs.PV == nil { + t.Fatalf("load reset clocks = %#v", reset.LatestTrainingMs) + } + other := "pv" + if selected == "pv" { + other = "load" + } + meta.RequestID = "native-reset-other-" + other + meta.State = reset.State + both, err := client.Reset(context.Background(), energyforecast.ResetRequest{RequestContext: meta, Signal: other, LearningStartedMs: cutoff}) + if err != nil { + t.Fatal(err) + } + if both.LatestTrainingMs.PV != nil || both.LatestTrainingMs.Load != nil { + t.Fatalf("both reset clocks = %#v", both.LatestTrainingMs) + } + + // State is the saved worker snapshot. A fresh worker must retain the + // two independently applied cutoffs before it receives old history. + if err := transport.Close(); err != nil { + t.Fatal(err) + } + client, transport = newClient() + t.Cleanup(func() { _ = transport.Close() }) + meta.RequestID = "native-reset-old-input-" + selected + meta.OriginMs = cutoff + quarter.Milliseconds() + meta.State = both.State + oldPV, oldLoad := 9000.0, 9000.0 + afterOld, err := client.Update(context.Background(), energyforecast.UpdateRequest{RequestContext: meta, Observations: []energyforecast.Observation{observation(start, oldPV, oldLoad)}}) + if err != nil { + t.Fatal(err) + } + selectedCounts, otherCounts := afterOld.Updates.PV, afterOld.Updates.Load + selectedClock, otherClock := afterOld.LatestTrainingMs.PV, afterOld.LatestTrainingMs.Load + if selected == "load" { + selectedCounts, otherCounts = afterOld.Updates.Load, afterOld.Updates.PV + selectedClock, otherClock = afterOld.LatestTrainingMs.Load, afterOld.LatestTrainingMs.PV + } + if selectedCounts.Applied != 0 || selectedCounts.Skipped != 1 || selectedClock != nil { + t.Fatalf("old %s history resurrected: counts=%#v clock=%v", selected, selectedCounts, selectedClock) + } + if otherCounts.Applied != 0 || otherCounts.Skipped != 1 || otherClock != nil { + t.Fatalf("old %s history resurrected: counts=%#v clock=%v", other, otherCounts, otherClock) + } + }) + } +} diff --git a/go/internal/energyforecast/reset.go b/go/internal/energyforecast/reset.go new file mode 100644 index 00000000..42131837 --- /dev/null +++ b/go/internal/energyforecast/reset.go @@ -0,0 +1,37 @@ +package energyforecast + +import ( + "context" + "errors" +) + +func (c *Client) Reset(ctx context.Context, request ResetRequest) (ResetReply, error) { + var reply ResetReply + if err := validateContext(request.RequestContext); err != nil { + return reply, err + } + if (request.Signal != "pv" && request.Signal != "load") || + (request.Signal == "pv" && request.Config.PV == nil) || + (request.Signal == "load" && request.Config.Load == nil) { + return reply, errors.New("forecast reset needs an enabled pv or load signal") + } + if request.LearningStartedMs <= 0 || request.LearningStartedMs > request.OriginMs { + return reply, errors.New("forecast learning start must be positive and no later than origin") + } + payload := struct { + Op string `json:"op"` + Version int `json:"version"` + Action string `json:"action"` + ResetRequest + }{"forecast", ProtocolVersion, "reset", request} + if err := c.call(ctx, payload, request.RequestContext, "reset", &reply); err != nil { + return ResetReply{}, err + } + if err := validateState(reply.State, true); err != nil { + return ResetReply{}, err + } + if reply.Signal != request.Signal || reply.LearningStartedMs < request.LearningStartedMs || reply.LearningStartedMs > request.OriginMs { + return ResetReply{}, errors.New("forecast reset reply changed signal or learning start") + } + return reply, nil +} diff --git a/go/internal/energyforecast/reset_test.go b/go/internal/energyforecast/reset_test.go new file mode 100644 index 00000000..63341051 --- /dev/null +++ b/go/internal/energyforecast/reset_test.go @@ -0,0 +1,115 @@ +package energyforecast + +import ( + "context" + "encoding/json" + "testing" +) + +func resetRequest() ResetRequest { + return ResetRequest{RequestContext: requestContext(), Signal: "pv", LearningStartedMs: requestContext().OriginMs} +} + +func resetReply(payload []byte) map[string]any { + var request map[string]any + _ = json.Unmarshal(payload, &request) + reply := map[string]any{ + "ok": true, + "model_revision": 2, + "latest_input_ms": map[string]any{"pv": nil, "load": nil}, + "latest_training_ms": map[string]any{"pv": nil, "load": nil}, + "latest_available_at_ms": map[string]any{"pv": nil, "load": nil}, + "state": map[string]any{"opaque": "reset-state"}, + "signal": request["signal"], + "learning_started_ms": request["learning_started_ms"], + } + for _, key := range []string{"op", "version", "action", "request_id", "site_id", "config_revision", "origin_ms"} { + reply[key] = request[key] + } + return reply +} + +func resetReplying(mutate func(map[string]any)) (*Client, *map[string]any) { + var request map[string]any + c := NewClient(exchangeFunc(func(_ context.Context, payload []byte) ([]byte, error) { + if err := json.Unmarshal(payload, &request); err != nil { + return nil, err + } + reply := resetReply(payload) + if mutate != nil { + mutate(reply) + } + return json.Marshal(reply) + })) + return c, &request +} + +func TestClientResetSendsTargetAndAcceptsVersionedState(t *testing.T) { + c, sent := resetReplying(nil) + request := resetRequest() + reply, err := c.Reset(context.Background(), request) + if err != nil { + t.Fatal(err) + } + if (*sent)["op"] != "forecast" || (*sent)["action"] != "reset" || (*sent)["signal"] != "pv" { + t.Fatalf("reset target = %#v", *sent) + } + if got := int64((*sent)["learning_started_ms"].(float64)); got != request.LearningStartedMs { + t.Fatalf("reset cutoff = %d, want %d", got, request.LearningStartedMs) + } + if reply.Signal != request.Signal || reply.LearningStartedMs != request.LearningStartedMs || len(reply.State) == 0 { + t.Fatalf("reset reply = %#v", reply) + } + if reply.Action != "reset" || reply.RequestID != request.RequestID || reply.SiteID != request.SiteID || reply.ConfigRevision != request.ConfigRevision || reply.OriginMs != request.OriginMs { + t.Fatalf("reset metadata = %#v", reply.ReplyContext) + } +} + +func TestClientResetRejectsInvalidReply(t *testing.T) { + request := resetRequest() + for name, mutate := range map[string]func(map[string]any){ + "missing_state": func(r map[string]any) { delete(r, "state") }, + "array_state": func(r map[string]any) { r["state"] = []any{"not opaque"} }, + "wrong_signal": func(r map[string]any) { r["signal"] = "load" }, + "missing_signal": func(r map[string]any) { delete(r, "signal") }, + "earlier_cutoff": func(r map[string]any) { r["learning_started_ms"] = request.LearningStartedMs - 1 }, + "future_cutoff": func(r map[string]any) { r["learning_started_ms"] = request.OriginMs + 1 }, + "future_clock": func(r map[string]any) { + r["latest_training_ms"] = map[string]any{"pv": request.OriginMs + 1, "load": nil} + }, + "wrong_action": func(r map[string]any) { r["action"] = "update" }, + "predictions": func(r map[string]any) { r["predictions"] = []any{} }, + } { + t.Run(name, func(t *testing.T) { + c, _ := resetReplying(mutate) + if _, err := c.Reset(context.Background(), request); err == nil { + t.Fatal("invalid reset reply accepted") + } + }) + } +} + +func TestClientResetRejectsInvalidRequestBeforeIO(t *testing.T) { + calls := 0 + c := NewClient(exchangeFunc(func(context.Context, []byte) ([]byte, error) { + calls++ + return nil, nil + })) + for _, mutate := range []func(*ResetRequest){ + func(r *ResetRequest) { r.Signal = "battery" }, + func(r *ResetRequest) { r.Signal = "pv"; r.Config.PV = nil }, + func(r *ResetRequest) { r.Signal = "load"; r.Config.Load = nil }, + func(r *ResetRequest) { r.LearningStartedMs = 0 }, + func(r *ResetRequest) { r.LearningStartedMs = r.OriginMs + 1 }, + func(r *ResetRequest) { r.State = json.RawMessage("[]") }, + } { + request := resetRequest() + mutate(&request) + if _, err := c.Reset(context.Background(), request); err == nil { + t.Fatal("invalid reset request reached client call") + } + } + if calls != 0 { + t.Fatalf("invalid reset requests reached worker %d times", calls) + } +} diff --git a/go/internal/energyforecast/types.go b/go/internal/energyforecast/types.go index 870e9196..24e8b94a 100644 --- a/go/internal/energyforecast/types.go +++ b/go/internal/energyforecast/types.go @@ -100,6 +100,14 @@ type PredictRequest struct { Horizon []HorizonSlot `json:"horizon"` } +// ResetRequest starts a new learning period for one signal. OriginMs is the +// current request time; LearningStartedMs also permits replay of a saved intent. +type ResetRequest struct { + RequestContext + Signal string `json:"signal"` + LearningStartedMs int64 `json:"learning_started_ms"` +} + type ReplyContext struct { Op string `json:"op"` Version int `json:"version"` @@ -129,6 +137,12 @@ type UpdateReply struct { State json.RawMessage `json:"state"` Updates UpdateCounts `json:"updates"` } +type ResetReply struct { + ReplyContext + State json.RawMessage `json:"state"` + Signal string `json:"signal"` + LearningStartedMs int64 `json:"learning_started_ms"` +} type LatestInput struct { PV *int64 `json:"pv,omitempty"` Load *int64 `json:"load,omitempty"` diff --git a/go/internal/forecasting/learning.go b/go/internal/forecasting/learning.go new file mode 100644 index 00000000..402cc89f --- /dev/null +++ b/go/internal/forecasting/learning.go @@ -0,0 +1,22 @@ +package forecasting + +import "fmt" + +// LearningRestartPendingError means the boundary is durable, but one of the +// model saves still needs retry. Callers must not report the intent as rejected. +type LearningRestartPendingError struct{ Err error } + +func (e *LearningRestartPendingError) Error() string { + return fmt.Sprintf("Learning period saved; model restart pending: %v", e.Err) +} +func (e *LearningRestartPendingError) Unwrap() error { return e.Err } + +// LearningStatus describes the model serving one forecast signal. Historical +// measurements and issued forecasts remain available across learning periods. +type LearningStatus struct { + Engine string `json:"engine"` + Status string `json:"status"` + StartedMS int64 `json:"started_ms"` + LatestTrainingMS int64 `json:"latest_training_ms"` + ResetAvailable bool `json:"reset_available"` +} diff --git a/go/internal/loadmodel/model.go b/go/internal/loadmodel/model.go index b2fe3940..79c17b48 100644 --- a/go/internal/loadmodel/model.go +++ b/go/internal/loadmodel/model.go @@ -58,6 +58,7 @@ type Bucket struct { type Model struct { ConfigRevision string `json:"config_revision,omitempty"` + LearningStartedMS int64 `json:"learning_started_ms,omitempty"` Timezone string `json:"timezone,omitempty"` LastTemperatureC float64 `json:"last_temperature_c"` HasTemperature bool `json:"has_temperature"` diff --git a/go/internal/loadmodel/restart_learning_test.go b/go/internal/loadmodel/restart_learning_test.go new file mode 100644 index 00000000..4055e554 --- /dev/null +++ b/go/internal/loadmodel/restart_learning_test.go @@ -0,0 +1,200 @@ +package loadmodel + +import ( + "encoding/json" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/modelstate" + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func TestRestartLearningClearsEveryLoadProfileAndRestoresConfiguredPrior(t *testing.T) { + s := NewService(nil, telemetry.NewStore(), "site", 4000, 0) + if err := s.Reconfigure("site", telemetry.ForecastOptions{}, "Europe/Stockholm", "site-a"); err != nil { + t.Fatal(err) + } + s.SeedHeatingCoef(275) + for _, profile := range Profiles() { + m := s.models[profile] + m.Update(time.Now(), 2000, 5) + m.HeatingW_per_degC = 900 + } + if err := s.SetProfile(ProfileAway); err != nil { + t.Fatal(err) + } + cutoff := time.Date(2026, 9, 7, 12, 0, 0, 0, time.UTC) + + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + snap := s.Snapshot() + if snap.ActiveProfile != ProfileAway { + t.Fatalf("active profile changed to %q", snap.ActiveProfile) + } + for _, profile := range Profiles() { + got := snap.Profiles[profile] + if got.Samples != 0 || got.MAE != 0 || got.LastMs != 0 || got.HasTemperature { + t.Fatalf("learned %s state survived: %+v", profile, got) + } + if got.HeatingW_per_degC != 275 { + t.Fatalf("%s heating prior = %.0f, want configured 275", profile, got.HeatingW_per_degC) + } + if got.PeakW != 4000 || got.Timezone != "Europe/Stockholm" || got.ConfigRevision != "site-a" { + t.Fatalf("configured %s state changed: %+v", profile, got) + } + if got.LearningStartedMS != cutoff.UnixMilli() { + t.Fatalf("%s cutoff = %d", profile, got.LearningStartedMS) + } + } +} + +func TestRestartLearningRejectsLoadSampleCapturedBeforeCutoff(t *testing.T) { + tel := telemetry.NewStore() + tel.Update("site", telemetry.DerMeter, 1000, nil, nil) + tel.RecordDriverSuccess("site") + s := NewService(nil, tel, "site", 4000, 0) + cutoff := time.Now().Add(time.Millisecond) + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + s.sampleAt(cutoff.Add(10 * time.Millisecond)) + + if got := s.Model().Samples; got != 0 { + t.Fatalf("pre-cutoff telemetry trained load model: %d samples", got) + } +} + +func TestRestartLearningRejectsMixedLoadBalanceWithOldInput(t *testing.T) { + tel := telemetry.NewStore() + now := time.Now() + cutoff := now.Add(-time.Second) + powerData := func(watts float64, measured time.Time) []byte { + data, err := json.Marshal(map[string]telemetry.ForecastPowerSample{ + "forecast_power": { + Version: 1, + Known: true, + Watts: watts, + MeasuredAtMS: measured.UnixMilli(), + ReceivedAtMS: now.UnixMilli(), + }, + }) + if err != nil { + t.Fatal(err) + } + return data + } + tel.Update("site", telemetry.DerMeter, 1000, nil, powerData(1000, cutoff.Add(-time.Second))) + tel.RecordDriverSuccess("site") + tel.Update("battery", telemetry.DerBattery, 500, nil, powerData(500, cutoff.Add(time.Second))) + tel.RecordDriverSuccess("battery") + s := NewService(nil, tel, "site", 4000, 0) + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + s.sampleAt(now.Add(time.Second)) + + if got := s.Model().Samples; got != 0 { + t.Fatalf("balance with pre-cutoff meter trained load model: %d samples", got) + } +} + +func TestRestartLearningInvalidatesLoadSampleInFlight(t *testing.T) { + tel := telemetry.NewStore() + tel.Update("site", telemetry.DerMeter, 1000, nil, nil) + tel.RecordDriverSuccess("site") + s := NewService(nil, tel, "site", 4000, 0) + entered := make(chan struct{}) + release := make(chan struct{}) + s.Temp = func(time.Time) (float64, bool) { + close(entered) + <-release + return 5, true + } + done := make(chan struct{}) + go func() { + defer close(done) + s.sampleAt(time.Now()) + }() + <-entered + if err := s.RestartLearning(time.Now()); err != nil { + t.Fatal(err) + } + close(release) + <-done + + if got := s.Model().Samples; got != 0 { + t.Fatalf("in-flight sample survived restart: %d", got) + } +} + +func TestRestartLearningLoadRetryPersistsAllProfilesWithoutErasingNewSamples(t *testing.T) { + bad := openTestDB(t) + s := NewService(bad, telemetry.NewStore(), "site", 4000, 0) + s.SeedHeatingCoef(275) + if err := bad.Close(); err != nil { + t.Fatal(err) + } + cutoff := time.Date(2026, 9, 7, 13, 0, 0, 0, time.UTC) + if err := s.RestartLearning(cutoff); err == nil { + t.Fatal("closed store did not surface persistence failure") + } + for _, profile := range Profiles() { + if !s.models[profile].Update(cutoff.Add(time.Minute), 1000, 10) { + t.Fatalf("post-reset %s sample rejected", profile) + } + } + + good, err := state.Open(t.TempDir() + "/retry.db") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { good.Close() }) + s.Store = good + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + restored := NewService(good, telemetry.NewStore(), "site", 4000, 0) + if restored.LearningStartedMS() != cutoff.UnixMilli() { + t.Fatalf("restored cutoff = %d", restored.LearningStartedMS()) + } + for _, profile := range Profiles() { + if got := restored.models[profile]; got.Samples != 1 || got.LearningStartedMS != cutoff.UnixMilli() { + t.Fatalf("same-cutoff retry lost %s learning: %+v", profile, got) + } + } +} + +func TestLoadRestoreReconcilesPartiallyPersistedEpoch(t *testing.T) { + st := openTestDB(t) + s := NewService(st, telemetry.NewStore(), "site", 4000, 0) + at := time.Date(2026, 9, 7, 14, 0, 0, 0, time.UTC) + for _, profile := range Profiles() { + s.models[profile].Update(at, 1000, 10) + } + if err := s.persist(); err != nil { + t.Fatal(err) + } + + cutoff := at.Add(time.Hour).UnixMilli() + home := freshProfile(s.models[ProfileHome], ProfileHome, "UTC", cutoff, nil) + js, err := modelstate.Wrap(FeatureHash(), home) + if err != nil { + t.Fatal(err) + } + if err := st.SaveConfig(stateKey(ProfileHome), js); err != nil { + t.Fatal(err) + } + + restored := NewService(st, telemetry.NewStore(), "site", 4000, 0) + for _, profile := range Profiles() { + got := restored.models[profile] + if got.Samples != 0 || got.LearningStartedMS != cutoff { + t.Fatalf("partial epoch left stale %s state: %+v", profile, got) + } + } +} diff --git a/go/internal/loadmodel/service.go b/go/internal/loadmodel/service.go index 2e9a2dd5..8add3c56 100644 --- a/go/internal/loadmodel/service.go +++ b/go/internal/loadmodel/service.go @@ -74,6 +74,7 @@ type Service struct { forecastOptions telemetry.ForecastOptions timezone string lastForecastInput time.Time + configuredHeating *float64 stop chan struct{} done chan struct{} @@ -141,12 +142,25 @@ func NewService(st *state.Store, tel *telemetry.Store, siteMeter string, peakW, zone = "UTC" } if old.Samples > 0 && zone != s.timezone { - m := newProfileModel(old.PeakW, profile) - m.HeatingW_per_degC = old.HeatingW_per_degC - s.models[profile] = m + heating := old.HeatingW_per_degC + s.models[profile] = freshProfile(old, profile, s.timezone, old.LearningStartedMS, &heating) } s.models[profile].Timezone = s.timezone } + var latestEpoch int64 + for _, profile := range Profiles() { + if m := s.models[profile]; m != nil && m.LearningStartedMS > latestEpoch { + latestEpoch = m.LearningStartedMS + } + } + if latestEpoch > 0 { + for _, profile := range Profiles() { + old := s.models[profile] + if old != nil && old.LearningStartedMS < latestEpoch { + s.models[profile] = freshProfile(old, profile, s.timezone, latestEpoch, nil) + } + } + } } return s } @@ -387,9 +401,14 @@ func (s *Service) Reconfigure(siteMeter string, opts telemetry.ForecastOptions, s.lastForecastInput = time.Time{} if changed { peak := s.activeModelLocked().PeakW + epoch := s.learningStartedMSLocked() for _, p := range Profiles() { m := newProfileModel(peak, p) m.Timezone, m.ConfigRevision = zone, revision + m.LearningStartedMS = epoch + if s.configuredHeating != nil { + m.HeatingW_per_degC = *s.configuredHeating + } s.models[p] = m } } @@ -412,6 +431,9 @@ func (s *Service) SetTimezone(zone string) error { m := newProfileModel(old.PeakW, p) m.Timezone = zone m.HeatingW_per_degC = old.HeatingW_per_degC + m.LearningStartedMS = old.LearningStartedMS + m.ConfigRevision = old.ConfigRevision + m.MaxPlausibleW = old.MaxPlausibleW s.models[p] = m } } @@ -427,7 +449,7 @@ func (s *Service) sampleAt(now time.Time) { return } s.mu.RLock() - site, opts, profile, generation := s.SiteMeter, s.forecastOptions, s.active, s.generation + site, opts, profile, generation, learningStartedMS := s.SiteMeter, s.forecastOptions, s.active, s.generation, s.learningStartedMSLocked() s.mu.RUnlock() reading := s.Tele.ForecastMeasurement(now, site, opts) temp := math.NaN() @@ -443,7 +465,7 @@ func (s *Service) sampleAt(now time.Time) { } model := s.activeModelLocked() updated := false - if reading.Valid && reading.Latest.After(s.lastForecastInput) { + if reading.Valid && reading.Earliest.UnixMilli() >= learningStartedMS && reading.Latest.After(s.lastForecastInput) { s.lastForecastInput = reading.Latest updated = model.Update(now, reading.HouseholdW, temp) } @@ -466,6 +488,10 @@ func (s *Service) persistProfile(profile Profile) error { func (s *Service) persist() error { s.persistMu.Lock() defer s.persistMu.Unlock() + return s.persistLocked() +} + +func (s *Service) persistLocked() error { if s.Store == nil { return nil } @@ -510,6 +536,8 @@ func (s *Service) SetHeatingCoef(w float64) { return } s.mu.Lock() + prior := w + s.configuredHeating = &prior for _, model := range s.models { model.HeatingW_per_degC = w } @@ -533,6 +561,8 @@ func (s *Service) SeedHeatingCoef(w float64) { return } s.mu.Lock() + prior := w + s.configuredHeating = &prior for _, model := range s.models { if model.Samples > 0 { continue @@ -557,6 +587,8 @@ func (s *Service) Reset() { s.models[profile].HeatingW_per_degC = heating s.models[profile].Timezone = s.timezone s.models[profile].ConfigRevision = old.ConfigRevision + s.models[profile].LearningStartedMS = old.LearningStartedMS + s.models[profile].MaxPlausibleW = old.MaxPlausibleW s.generation++ s.lastForecastInput = time.Time{} s.mu.Unlock() @@ -564,3 +596,68 @@ func (s *Service) Reset() { slog.Warn("loadmodel persist", "err", err) } } + +// RestartLearning clears every profile and starts one shared learned-model +// epoch. Replaying the same or an older cutoff persists the current models but +// does not erase samples learned after the first reset. +func (s *Service) RestartLearning(at time.Time) error { + if s == nil { + return nil + } + startedMS := at.UnixMilli() + if at.IsZero() || startedMS <= 0 { + return fmt.Errorf("loadmodel learning cutoff must be after the Unix epoch") + } + s.persistMu.Lock() + defer s.persistMu.Unlock() + s.mu.Lock() + if startedMS > s.learningStartedMSLocked() { + for _, profile := range Profiles() { + s.models[profile] = freshProfile(s.models[profile], profile, s.timezone, startedMS, s.configuredHeating) + } + s.generation++ + s.lastForecastInput = time.Time{} + } + s.mu.Unlock() + return s.persistLocked() +} + +// LearningStartedMS returns the cutoff shared by all load profiles. +func (s *Service) LearningStartedMS() int64 { + if s == nil { + return 0 + } + s.mu.RLock() + defer s.mu.RUnlock() + return s.learningStartedMSLocked() +} + +func (s *Service) learningStartedMSLocked() int64 { + var startedMS int64 + for _, profile := range Profiles() { + if m := s.models[profile]; m != nil && m.LearningStartedMS > startedMS { + startedMS = m.LearningStartedMS + } + } + return startedMS +} + +func freshProfile(old *Model, profile Profile, zone string, startedMS int64, heating *float64) *Model { + peak := 0.0 + revision := "" + maxPlausibleW := 0.0 + if old != nil { + peak = old.PeakW + revision = old.ConfigRevision + maxPlausibleW = old.MaxPlausibleW + } + m := newProfileModel(peak, profile) + m.Timezone = zone + m.ConfigRevision = revision + m.MaxPlausibleW = maxPlausibleW + m.LearningStartedMS = startedMS + if heating != nil { + m.HeatingW_per_degC = *heating + } + return m +} diff --git a/go/internal/mpc/energyplan_fault_test.go b/go/internal/mpc/energyplan_fault_test.go new file mode 100644 index 00000000..82cf698a --- /dev/null +++ b/go/internal/mpc/energyplan_fault_test.go @@ -0,0 +1,154 @@ +package mpc + +import ( + "context" + "encoding/json" + "math" + "strings" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/telemetry" +) + +// Faults enter after the bundled worker's process boundary, so the wrapper, +// request binding, plan validation, DP fallback and publication gate all run. +type energyplanFaultTransport struct { + OptimizerTransport + fault string + requests []externalRequest +} + +func (f *energyplanFaultTransport) RoundTrip(ctx context.Context, payload []byte) ([]byte, error) { + var request externalRequest + if err := json.Unmarshal(payload, &request); err != nil { + return nil, err + } + f.requests = append(f.requests, request) + reply, err := f.OptimizerTransport.RoundTrip(ctx, payload) + if err != nil { + return nil, err + } + switch f.fault { + case "timeout": + <-ctx.Done() + return nil, ctx.Err() + case "invalid_plan": + var response externalResponse + if err := json.Unmarshal(reply, &response); err != nil { + return nil, err + } + response.OK, response.Error = true, nil + if len(response.Plan.Actions) == 0 { + response.Plan.Actions = []externalAction{{SlotStartMs: request.Slots[0].StartMs, SlotLenMin: request.Slots[0].LenMin}} + } + response.Plan.Actions[0].GridW = 1e12 + return json.Marshal(response) + default: + return reply, nil + } +} + +func TestNativeEnergyplanEVRecoveryFallbackFaults(t *testing.T) { + for _, fault := range []string{"timeout", "invalid_plan"} { + for _, scenario := range []string{"feasible", "infeasible_without_previous", "infeasible_keep_previous"} { + t.Run(fault+"/"+scenario, func(t *testing.T) { + external := nativeWorker(t, 500*time.Millisecond) + t.Cleanup(func() { _ = external.Close() }) + wrapper := &EnergyplanOptimizer{ExternalOptimizer: external} + health, err := wrapper.Health(context.Background()) + if err != nil || health.Version != "0.2.2" { + t.Fatalf("bundle health=%+v err=%v", health, err) + } + svc := shadowTestService(t) + svc.Optimizer = wrapper + svc.Defaults.Mode = ModeArbitrage + svc.Defaults.SoCLevels, svc.Defaults.ActionLevels = 21, 21 + svc.Tele = telemetry.NewStore() + measuredSoC := .025 + svc.Tele.Update("battery", telemetry.DerBattery, 0, &measuredSoC, nil) + _, fixture := nativeFixture() + ev := *fixture.Loadpoint + svc.Loadpoint = func(int) *LoadpointSpec { copy := ev; return © } + var previous *Plan + var previousBytes []byte + if scenario == "infeasible_keep_previous" { + previous = svc.Replan(context.Background()) + if previous == nil || previous.Solver == nil || previous.Solver.Fallback || previous.InitialSoC != measuredSoC { + t.Fatalf("native baseline did not publish from measured low SoC: %+v", previous) + } + if err := ValidatePlan(svc.lastSlots, svc.lastParams, previous); err != nil { + t.Fatal(err) + } + svc.shadowWG.Wait() + previous = svc.Latest() + previousBytes, _ = json.Marshal(previous) + } + if scenario != "feasible" { + svc.BaseLoad, svc.FuseMaxW = 20000, 1000 + } + injected := &energyplanFaultTransport{OptimizerTransport: external.transport, fault: fault} + external.transport = injected + // Keep the native solve budget unchanged; only bound the lost reply. + if fault == "timeout" { + external.cfg.Timeout = 150 * time.Millisecond + } + plan := svc.Replan(context.Background()) + if len(injected.requests) != 1 { + t.Fatalf("wrapper sent %d requests", len(injected.requests)) + } + request := injected.requests[0] + if len(request.Storages) != 1 || len(request.FlexLoads) != 1 { + t.Fatalf("fault did not exercise battery+EV: storages=%d flex=%d", len(request.Storages), len(request.FlexLoads)) + } + if math.Abs(request.Storages[0].InitialEnergyWh-measuredSoC*svc.Defaults.CapacityWh) > 1e-9 { + t.Fatalf("worker request clamped measured initial energy: %+v", request.Storages[0]) + } + if scenario != "feasible" { + if previous == nil { + if plan != nil || svc.Latest() != nil { + t.Fatalf("infeasible fallback published without previous plan: %+v", plan) + } + } else { + got, _ := json.Marshal(svc.Latest()) + returned, _ := json.Marshal(plan) + if string(got) != string(previousBytes) || string(returned) != string(previousBytes) { + t.Fatal("rejected fallback replaced or altered previous plan") + } + } + return + } + if plan == nil || plan.Solver == nil || !plan.Solver.Fallback || plan.Solver.Engine != "core" || plan.Solver.Status != "fallback" { + t.Fatalf("missing Core DP fallback: %+v", plan) + } + reason := "optimizer plan rejected:" + if fault == "timeout" { + reason = "optimizer timeout after" + } + if !strings.Contains(plan.Solver.FallbackReason, reason) { + t.Fatalf("wrong fallback reason: %+v", plan.Solver) + } + if plan.InitialSoC != measuredSoC || svc.lastParams.InitialSoC != measuredSoC || svc.lastParams.InitialSoC >= svc.lastParams.SoCMin { + t.Fatalf("fallback clamped recovery start: plan=%g params=%g floor=%g", plan.InitialSoC, svc.lastParams.InitialSoC, svc.lastParams.SoCMin) + } + if err := ValidatePlan(svc.lastSlots, svc.lastParams, plan); err != nil { + t.Fatalf("published fallback failed physical replay: %v", err) + } + if svc.lastParams.Loadpoint == nil || svc.lastParams.Loadpoint.ID != ev.ID { + t.Fatal("fallback dropped active EV") + } + charged := false + for _, a := range plan.Actions { + charged = charged || a.LoadpointW > 0 + } + if !charged || plan.Actions[ev.TargetSlotIdx].LoadpointSoC+1e-9 < ev.TargetSoC { + t.Fatalf("fallback did not meet feasible EV target: %+v", plan.Actions) + } + if svc.Latest() == nil || svc.Latest().DecisionID != plan.DecisionID { + t.Fatal("valid fallback was not published") + } + svc.shadowWG.Wait() + }) + } + } +} diff --git a/go/internal/mpc/energyplan_test.go b/go/internal/mpc/energyplan_test.go index e6844cfe..cd5aa599 100644 --- a/go/internal/mpc/energyplan_test.go +++ b/go/internal/mpc/energyplan_test.go @@ -60,7 +60,7 @@ func TestNativeEnergyplanDownsideAndAsyncShadow(t *testing.T) { svc := shadowTestService(t) svc.Optimizer = &EnergyplanOptimizer{ExternalOptimizer: o} info, err := svc.Optimizer.(*EnergyplanOptimizer).Health(context.Background()) - if err != nil || info.Name != "ftw-solver" || info.Version != "0.2.1" { + if err != nil || info.Name != "ftw-solver" || info.Version != "0.2.2" { t.Fatalf("bundled worker health: %+v %v", info, err) } start := time.Now().UTC().Truncate(time.Hour) diff --git a/go/internal/pvmodel/model.go b/go/internal/pvmodel/model.go index 3c89c08a..f2e4a856 100644 --- a/go/internal/pvmodel/model.go +++ b/go/internal/pvmodel/model.go @@ -51,13 +51,14 @@ func finite(v float64) bool { return !math.IsNaN(v) && !math.IsInf(v, 0) } // Model is the learned PV predictor. type Model struct { - ConfigRevision string `json:"config_revision,omitempty"` - Beta [NFeat]float64 `json:"beta"` - P [NFeat][NFeat]float64 `json:"p"` // covariance - Forgetting float64 `json:"forgetting"` - Samples int64 `json:"samples"` - LastMs int64 `json:"last_ms"` - MAE float64 `json:"mae"` // EMA of |err| (W) + ConfigRevision string `json:"config_revision,omitempty"` + LearningStartedMS int64 `json:"learning_started_ms,omitempty"` + Beta [NFeat]float64 `json:"beta"` + P [NFeat][NFeat]float64 `json:"p"` // covariance + Forgetting float64 `json:"forgetting"` + Samples int64 `json:"samples"` + LastMs int64 `json:"last_ms"` + MAE float64 `json:"mae"` // EMA of |err| (W) // RelMAE is MAE expressed as a share of the prediction it belongs to // (0..1), over the same EMA window. The planner sizes each slot's PV // downside against that slot's own expected generation, which a watt diff --git a/go/internal/pvmodel/restart_learning_test.go b/go/internal/pvmodel/restart_learning_test.go new file mode 100644 index 00000000..82946453 --- /dev/null +++ b/go/internal/pvmodel/restart_learning_test.go @@ -0,0 +1,171 @@ +package pvmodel + +import ( + "encoding/json" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func TestRestartLearningClearsLearnedPVState(t *testing.T) { + s := NewService(nil, nil, func(time.Time) float64 { return 800 }, nil, 5000) + s.SetACLimit(4200) + s.model.ConfigRevision = "site-a" + at := time.Date(2026, 9, 7, 10, 0, 0, 0, time.UTC) + for i := 0; i < 20; i++ { + s.model.Update(800, 20, at.Add(time.Duration(i)*time.Minute), 3000) + s.Residuals.Add(at.Add(time.Duration(i)*time.Minute), 2500, 3000) + } + + if err := s.RestartLearning(at.Add(time.Hour)); err != nil { + t.Fatal(err) + } + + got := s.Model() + if got.Samples != 0 || got.MAE != 0 || got.RelMAE != 0 || got.InferredScaleKnown || got.ScaleSamples != 0 { + t.Fatalf("learned state survived: %+v", got) + } + if s.Residuals.Len() != 0 { + t.Fatal("PV residuals survived") + } + if got.RatedW != 5000 || got.ACLimitW != 4200 || got.ConfigRevision != "site-a" { + t.Fatalf("configured PV state changed: %+v", got) + } + if got.LearningStartedMS != at.Add(time.Hour).UnixMilli() || s.LearningStartedMS() != got.LearningStartedMS { + t.Fatalf("learning cutoff = %d", got.LearningStartedMS) + } +} + +func TestRestartLearningRejectsSampleCapturedBeforeCutoff(t *testing.T) { + tel := telemetry.NewStore() + tel.Update("pv", telemetry.DerPV, -3000, nil, nil) + tel.RecordDriverSuccess("pv") + s := NewService(nil, tel, func(time.Time) float64 { return 800 }, func(time.Time) (float64, bool) { return 20, true }, 5000) + cutoff := time.Now().Add(time.Millisecond) + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + s.sampleAt(cutoff.Add(10 * time.Millisecond)) + + if got := s.Model().Samples; got != 0 { + t.Fatalf("pre-cutoff telemetry trained PV model: %d samples", got) + } +} + +func TestRestartLearningRejectsMixedPVWithOldInput(t *testing.T) { + tel := telemetry.NewStore() + now := time.Now() + cutoff := now.Add(-time.Second) + tel.Update("old-pv", telemetry.DerPV, -1000, nil, forecastPowerData(t, -1000, cutoff.Add(-time.Second), now)) + tel.RecordDriverSuccess("old-pv") + tel.Update("new-pv", telemetry.DerPV, -2000, nil, forecastPowerData(t, -2000, cutoff.Add(time.Second), now)) + tel.RecordDriverSuccess("new-pv") + s := NewService(nil, tel, func(time.Time) float64 { return 800 }, func(time.Time) (float64, bool) { return 20, true }, 5000) + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + s.sampleAt(now.Add(time.Second)) + + if got := s.Model().Samples; got != 0 { + t.Fatalf("PV sum with pre-cutoff source trained model: %d samples", got) + } +} + +func TestRestartLearningIgnoresOldNonPVInputForPVTraining(t *testing.T) { + tel := telemetry.NewStore() + now := time.Now() + cutoff := now.Add(-time.Second) + tel.Update("pv", telemetry.DerPV, -3000, nil, forecastPowerData(t, -3000, cutoff.Add(time.Second), now)) + tel.RecordDriverSuccess("pv") + tel.Update("battery", telemetry.DerBattery, 500, nil, forecastPowerData(t, 500, cutoff.Add(-time.Second), now)) + tel.RecordDriverSuccess("battery") + s := NewService(nil, tel, func(time.Time) float64 { return 800 }, func(time.Time) (float64, bool) { return 20, true }, 5000) + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + + s.sampleAt(now.Add(time.Second)) + + if got := s.Model().Samples; got != 1 { + t.Fatalf("old non-PV input blocked fresh PV training: %d samples", got) + } +} + +func forecastPowerData(t *testing.T, watts float64, measured, received time.Time) []byte { + t.Helper() + data, err := json.Marshal(map[string]telemetry.ForecastPowerSample{ + "forecast_power": { + Version: 1, + Known: true, + Watts: watts, + MeasuredAtMS: measured.UnixMilli(), + ReceivedAtMS: received.UnixMilli(), + }, + }) + if err != nil { + t.Fatal(err) + } + return data +} + +func TestRestartLearningInvalidatesPVSampleInFlight(t *testing.T) { + tel := telemetry.NewStore() + tel.Update("pv", telemetry.DerPV, -3000, nil, nil) + tel.RecordDriverSuccess("pv") + entered := make(chan struct{}) + release := make(chan struct{}) + s := NewService(nil, tel, func(time.Time) float64 { + close(entered) + <-release + return 800 + }, func(time.Time) (float64, bool) { return 20, true }, 5000) + done := make(chan struct{}) + go func() { + defer close(done) + s.sampleAt(time.Now()) + }() + <-entered + cutoff := time.Now() + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + close(release) + <-done + + if got := s.Model().Samples; got != 0 || s.Residuals.Len() != 0 { + t.Fatalf("in-flight sample survived restart: samples=%d residuals=%d", got, s.Residuals.Len()) + } +} + +func TestRestartLearningPVRetryPersistsWithoutErasingNewSamples(t *testing.T) { + bad := openTestDB(t) + s := NewService(bad, nil, func(time.Time) float64 { return 800 }, nil, 5000) + if err := bad.Close(); err != nil { + t.Fatal(err) + } + cutoff := time.Date(2026, 9, 7, 11, 0, 0, 0, time.UTC) + if err := s.RestartLearning(cutoff); err == nil { + t.Fatal("closed store did not surface persistence failure") + } + if !s.model.Update(800, 20, cutoff.Add(time.Minute), 3000) { + t.Fatal("post-reset fixture sample rejected") + } + + good, err := state.Open(t.TempDir() + "/retry.db") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { good.Close() }) + s.Store = good + if err := s.RestartLearning(cutoff); err != nil { + t.Fatal(err) + } + restored := NewService(good, nil, func(time.Time) float64 { return 800 }, nil, 5000) + if got := restored.Model(); got.Samples != 1 || got.LearningStartedMS != cutoff.UnixMilli() { + t.Fatalf("same-cutoff retry erased or failed to persist new learning: %+v", got) + } +} diff --git a/go/internal/pvmodel/service.go b/go/internal/pvmodel/service.go index 01b216cc..6a5ce584 100644 --- a/go/internal/pvmodel/service.go +++ b/go/internal/pvmodel/service.go @@ -2,6 +2,7 @@ package pvmodel import ( "context" + "fmt" "log/slog" "sync" "time" @@ -344,14 +345,19 @@ func (s *Service) liveActualPV() (float64, bool) { } func (s *Service) liveActualPVAt(now time.Time) (float64, bool) { + pvW, _, valid := s.livePVMeasurementAt(now) + return pvW, valid +} + +func (s *Service) livePVMeasurementAt(now time.Time) (float64, time.Time, bool) { if s.Tele == nil || (s.CurtailmentActive != nil && s.CurtailmentActive()) { - return 0, false + return 0, time.Time{}, false } s.mu.RLock() options := s.forecastOptions s.mu.RUnlock() m := s.Tele.ForecastMeasurement(now, "", options) - return -m.PVW, m.PVValid + return -m.PVW, m.PVEarliest, m.PVValid } // PredictNow returns the twin's prediction for right now using the @@ -416,7 +422,7 @@ func (s *Service) sample() { func (s *Service) sampleAt(now time.Time) { s.mu.RLock() - clearSky, generation := s.ClearSky, s.generation + clearSky, generation, learningStartedMS := s.ClearSky, s.generation, s.model.LearningStartedMS s.mu.RUnlock() if clearSky == nil { return @@ -434,10 +440,13 @@ func (s *Service) sampleAt(now time.Time) { } // Aggregate PV across all drivers. PV telemetry is stored as // site-sign (negative = generating), so flip to positive. - pvW, valid := s.liveActualPVAt(now) + pvW, earliest, valid := s.livePVMeasurementAt(now) if !valid { return } + if learningStartedMS > 0 && earliest.UnixMilli() < learningStartedMS { + return + } // Capture (predicted_at_now, actual_at_now) for the residual buffer // BEFORE running the RLS update — otherwise the model's prediction @@ -470,24 +479,31 @@ func (s *Service) sampleAt(now time.Time) { } } -func (s *Service) persist() { - if s.Store == nil { - return - } +func (s *Service) persist() error { // Serialise the entire marshal+save so a sample-loop persist that // started before a Reset cannot finish after Reset's persist and // clobber the clean state with stale coefficients. s.persistMu.Lock() defer s.persistMu.Unlock() + return s.persistLocked() +} + +// persistLocked writes one snapshot while persistMu is held. +func (s *Service) persistLocked() error { + if s.Store == nil { + return nil + } s.mu.RLock() js, err := modelstate.Wrap(FeatureHash(), s.model) s.mu.RUnlock() if err != nil { - return + return err } if err := s.Store.SaveConfig(stateKey, string(js)); err != nil { slog.Warn("pvmodel persist", "err", err) + return err } + return nil } // Reset clears the model to a fresh prior (useful after a system change @@ -501,20 +517,56 @@ func (s *Service) Reset() { s.mu.Lock() s.resetLocked() s.mu.Unlock() - s.persist() + if err := s.persist(); err != nil { + slog.Warn("pvmodel persist", "err", err) + } } func (s *Service) resetLocked() { rated := s.model.RatedW ac := s.model.ACLimitW revision := s.model.ConfigRevision + learningStartedMS := s.model.LearningStartedMS s.model = NewModel(rated) s.model.ACLimitW = ac s.model.ConfigRevision = revision + s.model.LearningStartedMS = learningStartedMS s.Residuals = NewResidualBuffer() s.generation++ } +// RestartLearning starts a new learned-model epoch. A replay of the same or an +// older cutoff only persists the current state, so a boot-time retry repairs a +// failed write without erasing samples learned after the reset. +func (s *Service) RestartLearning(at time.Time) error { + if s == nil { + return nil + } + startedMS := at.UnixMilli() + if at.IsZero() || startedMS <= 0 { + return fmt.Errorf("pvmodel learning cutoff must be after the Unix epoch") + } + s.persistMu.Lock() + defer s.persistMu.Unlock() + s.mu.Lock() + if startedMS > s.model.LearningStartedMS { + s.resetLocked() + s.model.LearningStartedMS = startedMS + } + s.mu.Unlock() + return s.persistLocked() +} + +// LearningStartedMS returns the current learned-model epoch cutoff. +func (s *Service) LearningStartedMS() int64 { + if s == nil { + return 0 + } + s.mu.RLock() + defer s.mu.RUnlock() + return s.model.LearningStartedMS +} + // ResidualCorrect is the integration point for the MPC. Returns the // additive correction (W) to apply to a base prediction targeting // tTarget, computed at wall time `now`. Returns 0 when residuals are diff --git a/go/internal/telemetry/forecast.go b/go/internal/telemetry/forecast.go index d2344cc3..8ed431ad 100644 --- a/go/internal/telemetry/forecast.go +++ b/go/internal/telemetry/forecast.go @@ -39,7 +39,7 @@ type ForecastOptions struct { // Valid describes the household balance. PVValid only describes the PV sources. // Neither flag claims that the instantaneous sample covers a time interval. type ForecastReading struct { - At, Earliest, Latest, PVLatest time.Time + At, Earliest, Latest, PVEarliest, PVLatest time.Time GridW, PVW, BatteryW, EVW, V2XW, HouseholdW float64 Valid, PVValid bool Reason, PVReason string @@ -208,7 +208,7 @@ func (s *Store) forecastMeasurementLocked(now time.Time, siteMeter string, opts if out.Latest.Sub(out.Earliest) > opts.MaxSkew { fail("time_skew", false) } - out.PVLatest = pvLast + out.PVEarliest, out.PVLatest = pvFirst, pvLast if pvLast.Sub(pvFirst) > opts.MaxSkew { fail("pv_time_skew", true) } diff --git a/optimizer/native/bundle/forecast-v1.response.schema.json b/optimizer/native/bundle/forecast-v1.response.schema.json index a8370357..31f920b7 100644 --- a/optimizer/native/bundle/forecast-v1.response.schema.json +++ b/optimizer/native/bundle/forecast-v1.response.schema.json @@ -175,6 +175,85 @@ "predictions" ] }, + { + "type": "object", + "additionalProperties": false, + "properties": { + "op": { + "const": "forecast" + }, + "version": { + "const": 1 + }, + "action": { + "const": "reset" + }, + "request_id": { + "type": "string" + }, + "site_id": { + "type": "string" + }, + "config_revision": { + "type": "string" + }, + "origin_ms": { + "type": "integer", + "minimum": 0 + }, + "ok": { + "const": true + }, + "model_revision": { + "type": "integer", + "minimum": 0, + "maximum": 18446744073709551615, + "description": "Increments when reset applies a newer cutoff. Replaying an already-applied cutoff keeps the revision." + }, + "latest_input_ms": { + "$ref": "#/$defs/watermarks" + }, + "latest_training_ms": { + "$ref": "#/$defs/watermarks" + }, + "latest_available_at_ms": { + "$ref": "#/$defs/watermarks" + }, + "state": { + "type": "object", + "description": "Opaque snapshot after the selected signal reset. Persist unchanged." + }, + "signal": { + "enum": [ + "pv", + "load" + ], + "description": "The selected model signal." + }, + "learning_started_ms": { + "type": "integer", + "minimum": 0, + "description": "Applied learning cutoff for signal. It is no later than origin_ms and can be later than a replayed request cutoff." + } + }, + "required": [ + "op", + "version", + "action", + "request_id", + "site_id", + "config_revision", + "origin_ms", + "ok", + "model_revision", + "latest_input_ms", + "latest_training_ms", + "latest_available_at_ms", + "state", + "signal", + "learning_started_ms" + ] + }, { "type": "object", "required": [ diff --git a/optimizer/native/bundle/forecast-v1.schema.json b/optimizer/native/bundle/forecast-v1.schema.json index 2487b9dc..e5df5e97 100644 --- a/optimizer/native/bundle/forecast-v1.schema.json +++ b/optimizer/native/bundle/forecast-v1.schema.json @@ -2,7 +2,7 @@ "$schema": "https://json-schema.org/draft/2020-12/schema", "$id": "urn:sourceful:energyplan:forecast:1", "title": "Energyplan forecasting request v1", - "description": "One JSON object per stdin line, at most 2 MiB including the newline. The caller must split update batches by encoded byte size, not just count; a larger line terminates the worker. Update consumes complete UTC-aligned 15-minute means available by origin_ms and returns an opaque model snapshot. Predict uses an immutable snapshot. Each requested future segment lasts 1..900000 ms and stays within one UTC quarter-hour; full aligned quarters remain valid. Returned bounds are the actual requested segment. For prediction, site-local clock, weather and occupancy fields refer to the containing quarter-hour start. The caller splits longer spans at each UTC quarter boundary and handles timezone or DST changes. Site/config identity, interval order, state bounds and causal availability are checked at runtime. Household load excludes EV, battery and V2X. PV denotes nonnegative available AC generation, distinct from signed planner pv_w. The caller records actual completion/issue time and persists the snapshot. See forecast-v1.response.schema.json for responses.", + "description": "One JSON object per stdin line, at most 2 MiB including the newline. The caller must split update batches by encoded byte size, not just count; a larger line terminates the worker. Update consumes complete UTC-aligned 15-minute means available by origin_ms and returns an opaque model snapshot. Predict uses an immutable snapshot. Reset starts a new learning period for exactly one enabled signal while retaining the other signal's state. Each requested future segment lasts 1..900000 ms and stays within one UTC quarter-hour; full aligned quarters remain valid. Returned bounds are the actual requested segment. For prediction, site-local clock, weather and occupancy fields refer to the containing quarter-hour start. The caller splits longer spans at each UTC quarter boundary and handles timezone or DST changes. Site/config identity, interval order, state bounds and causal availability are checked at runtime. Household load excludes EV, battery and V2X. PV denotes nonnegative available AC generation, distinct from signed planner pv_w. The caller records actual completion/issue time and persists the snapshot. See forecast-v1.response.schema.json for responses.", "type": "object", "additionalProperties": false, "properties": { @@ -15,7 +15,8 @@ "action": { "enum": [ "update", - "predict" + "predict", + "reset" ] }, "request_id": { @@ -175,6 +176,18 @@ "items": { "$ref": "#/$defs/horizon" } + }, + "signal": { + "enum": [ + "pv", + "load" + ], + "description": "The one enabled model whose learned state reset replaces. Only valid for action reset." + }, + "learning_started_ms": { + "type": "integer", + "minimum": 0, + "description": "Durable start cutoff for the selected signal. It must not be later than origin_ms; runtime validates that relation. A replay may return a later already-applied cutoff. Only valid for action reset." } }, "required": [ @@ -201,8 +214,22 @@ "observations" ], "not": { - "required": [ - "horizon" + "anyOf": [ + { + "required": [ + "horizon" + ] + }, + { + "required": [ + "signal" + ] + }, + { + "required": [ + "learning_started_ms" + ] + } ] } } @@ -220,8 +247,51 @@ "horizon" ], "not": { - "required": [ - "observations" + "anyOf": [ + { + "required": [ + "observations" + ] + }, + { + "required": [ + "signal" + ] + }, + { + "required": [ + "learning_started_ms" + ] + } + ] + } + } + }, + { + "if": { + "properties": { + "action": { + "const": "reset" + } + } + }, + "then": { + "required": [ + "signal", + "learning_started_ms" + ], + "not": { + "anyOf": [ + { + "required": [ + "observations" + ] + }, + { + "required": [ + "horizon" + ] + } ] } } diff --git a/optimizer/native/bundle/ftw-solver-darwin-arm64 b/optimizer/native/bundle/ftw-solver-darwin-arm64 index e657e916..46b556c6 100755 Binary files a/optimizer/native/bundle/ftw-solver-darwin-arm64 and b/optimizer/native/bundle/ftw-solver-darwin-arm64 differ diff --git a/optimizer/native/bundle/ftw-solver-linux-amd64 b/optimizer/native/bundle/ftw-solver-linux-amd64 index 5d5f8bf2..30cd7272 100755 Binary files a/optimizer/native/bundle/ftw-solver-linux-amd64 and b/optimizer/native/bundle/ftw-solver-linux-amd64 differ diff --git a/optimizer/native/bundle/ftw-solver-linux-arm64 b/optimizer/native/bundle/ftw-solver-linux-arm64 index 6ca0867c..41c510c4 100755 Binary files a/optimizer/native/bundle/ftw-solver-linux-arm64 and b/optimizer/native/bundle/ftw-solver-linux-arm64 differ diff --git a/optimizer/native/bundle/manifest.json b/optimizer/native/bundle/manifest.json index 24120a1d..a5ec4a84 100644 --- a/optimizer/native/bundle/manifest.json +++ b/optimizer/native/bundle/manifest.json @@ -1,9 +1,9 @@ { "schema_version": 1, "product": "energyplan", - "version": "0.2.1", + "version": "0.2.2", "source_repository": "srcfl/energyplan", - "source_commit": "820383e99ff21ae1a1428ebbf97bfaa9768e1cd5", + "source_commit": "b78cde268432671762e8c76c2c739f7a71f91087", "rustc": "rustc 1.95.0 (59807616e 2026-04-14)", "protocol_version": 1, "forecast_protocol_version": 1, @@ -31,24 +31,24 @@ "bytes": 14042 }, "forecast-v1.response.schema.json": { - "sha256": "b1eaac87ac7234d133e7bef6cc343321c78a0f9e3d1958e45770637770436999", - "bytes": 10989 + "sha256": "b7f06c00679b3798cbedd458f8b9b6bfd556e96ac442bb3820f9696c6c7bdd99", + "bytes": 13008 }, "forecast-v1.schema.json": { - "sha256": "131040c1ce6a7269f99d3114c84dec6d6b853c64bc8de182540928659b99233a", - "bytes": 13696 + "sha256": "6b374bd19d351e568761e1050016202d78aea2fc2e6783a2e9245aa23ffb2d3f", + "bytes": 15377 }, "ftw-solver-darwin-arm64": { - "sha256": "88d4f39d5d3432ecd1050d99c17d2d7f4d08c2b297523135a46be7e9d31cef53", - "bytes": 825040 + "sha256": "04e4e02dc7bdfe672f71d4f37a955fc39af6ee657d6315d304b28a5d7ccb2805", + "bytes": 858240 }, "ftw-solver-linux-amd64": { - "sha256": "5a0a02d46d495f885f53c99c791502f2aa115c9eed40c541b70b567b65a9072e", - "bytes": 1035304 + "sha256": "70a950fc6b51f227a1441183c8571030d9f45ea89774d7231f2416c773f8927f", + "bytes": 1079848 }, "ftw-solver-linux-arm64": { - "sha256": "482e6926be51483a3dda0aef574d3b3bf7fc40f34e3b81768b323f8ebc8ca4ef", - "bytes": 867432 + "sha256": "5bc5f3f98a10a0faf65b1a2630289d24bfe04402a20aa44c90f2a6058ea6af5a", + "bytes": 902840 }, "rust-runtime/COPYRIGHT-library.html": { "sha256": "90567e2718bf7fd65a71a3a43c5596488e80e5f51ed02bfea6fec54458b5f3d1", diff --git a/web/twins.js b/web/twins.js index 5249fcc3..7982731c 100644 --- a/web/twins.js +++ b/web/twins.js @@ -1,46 +1,56 @@ -// twins.js — advanced-mode ML diagnostics panel. -// Renders a small card per twin (PV, load, price forecaster) with -// sample count, MAE, quality bar, and last-updated time. Polls only -// while advanced mode is visible. +// twins.js — advanced-mode forecast-learning diagnostics. +// Polls only while advanced mode is visible. The box owns the model state; +// this view only asks it to begin a new learning period for one signal. (function () { 'use strict'; const REFRESH_MS = 10000; + const actions = { + '/api/pvmodel/reset': { signal: 'solar production', other: 'consumption', label: 'Relearn solar production', id: 'pv' }, + '/api/loadmodel/reset': { signal: 'consumption', other: 'solar production', label: 'Relearn consumption', id: 'load' }, + }; let refreshTimer = null; - - function apiFetch(path, opts) { - return fetch(path, opts); + let refreshRevision = 0; + let lastPV = null; + let lastLoad = null; + const pending = new Set(); + const availableActions = new Set(); + const actionMessages = new Map(); + const actionFocus = new Set(); + + function apiFetch(path, opts) { return fetch(path, opts); } + + async function fetchModel(path) { + const response = await apiFetch(path); + if (!response.ok) throw new Error('HTTP ' + response.status); + return response.json(); } async function fetchAll() { + const revision = ++refreshRevision; const [pv, load] = await Promise.all([ - apiFetch('/api/pvmodel').then(r => r.json()).catch(() => ({ enabled: false })), - apiFetch('/api/loadmodel').then(r => r.json()).catch(() => ({ enabled: false })), + fetchModel('/api/pvmodel').catch(() => ({ unavailable: true })), + fetchModel('/api/loadmodel').catch(() => ({ unavailable: true })), ]); + if (revision !== refreshRevision) return; + lastPV = pv; + lastLoad = load; render(pv, load); } - function advancedVisible() { - return !!(document.body && document.body.classList.contains('advanced')); - } - + function advancedVisible() { return !!(document.body && document.body.classList.contains('advanced')); } function startPolling() { if (refreshTimer) return; fetchAll(); refreshTimer = setInterval(fetchAll, REFRESH_MS); } - function stopPolling() { if (!refreshTimer) return; clearInterval(refreshTimer); refreshTimer = null; } - - function syncPolling() { - if (advancedVisible()) startPolling(); - else stopPolling(); - } + function syncPolling() { if (advancedVisible()) startPolling(); else stopPolling(); } function fmtAge(ms) { if (!ms) return '—'; @@ -49,18 +59,43 @@ if (s < 3600) return Math.round(s / 60) + 'm ago'; return Math.round(s / 3600) + 'h ago'; } + function fmtLocalTime(ms) { + if (!ms || !Number.isFinite(Number(ms))) return '—'; + const date = new Date(Number(ms)); + return Number.isNaN(date.getTime()) ? '—' : date.toLocaleString(); + } + function esc(value) { + return String(value == null ? '' : value).replace(/&/g, '&').replace(//g, '>').replace(/"/g, '"').replace(/'/g, '''); + } + function learningStatus(learning) { + if (!learning) return 'Model state unavailable'; + switch (learning.status) { + case 'cold_start': return 'Cold start'; + case 'learning': return 'Learning'; + case 'ready': return 'Ready'; + case 'unavailable': return 'Unavailable'; + default: return 'Model state unavailable'; + } + } + function engineLabel(learning) { + if (!learning) return 'Unavailable'; + return learning.engine === 'energyplan' ? 'Energyplan' : learning.engine === 'legacy' ? 'Legacy' : 'Unavailable'; + } + function actionAvailable(d) { + const learning = d && d.learning; + return !!(d && d.enabled && learning && learning.reset_available === true); + } - // Matches the battery-model reset button in models.js — same class - // `.btn-reset-model` so both paths share the theme-aware styling - // declared in app.css (ghost look per the shared design system: transparent bg, - // --line border, --fg text, theme-aware). The old inline style - // referenced the legacy --surface2 / --border / --text-dim hex - // tokens that don't flip with the light-mode switch, so this - // button read as a black-on-white blob on paper. - function resetButton(endpoint, label) { - return `'; + function resetButton(endpoint, d) { + const action = actions[endpoint]; + const isPending = pending.has(endpoint); + const enabled = actionAvailable(d); + const statusID = 'twin-status-' + action.id; + if (enabled) availableActions.add(endpoint); + const unavailable = !enabled; + const label = isPending ? 'Starting new learning period…' : action.label; + return ``; } function loadProfileControl(d) { @@ -71,84 +106,139 @@ return ``; } return '
profile' + - `
` + - btn('home', 'Home') + - btn('away', 'Away') + - '
'; + `
` + btn('home', 'Home') + btn('away', 'Away') + '
'; } - function twinCard(title, d, resetEndpoint, resetLabel, extraHtml) { - if (!d || !d.enabled) return `

${title}

disabled
`; - const q = Math.max(0, Math.min(1, d.quality || 0)); - const qPct = (q * 100).toFixed(0); - const qColor = q >= 0.7 ? '#22c55e' : q >= 0.3 ? '#fbbf24' : '#ef4444'; + function modelRows(d) { + const learning = d.learning; + const legacy = learning && learning.engine === 'energyplan'; + const prefix = legacy ? 'legacy ' : ''; const rows = []; - rows.push(`
samples${d.samples || 0}
`); - if (d.mae_w != null) rows.push(`
MAE${d.mae_w.toFixed(0)} W
`); - if (d.peak_w != null) rows.push(`
peak ref${(d.peak_w/1000).toFixed(1)} kW
`); - if (d.rated_w != null) rows.push(`
rated${(d.rated_w/1000).toFixed(1)} kW
`); - if (d.heating_w_per_degc != null && d.heating_w_per_degc > 0) { - rows.push(`
heating${d.heating_w_per_degc.toFixed(0)} W/°C
`); + if (legacy) rows.push('
legacy model statssecondary
'); + rows.push(`
${prefix}samples${d.samples || 0}
`); + if (d.mae_w != null) rows.push(`
${prefix}MAE${d.mae_w.toFixed(0)} W
`); + if (d.peak_w != null) rows.push(`
${prefix}peak ref${(d.peak_w / 1000).toFixed(1)} kW
`); + if (d.rated_w != null) rows.push(`
${prefix}rated${(d.rated_w / 1000).toFixed(1)} kW
`); + if (d.heating_w_per_degc != null && d.heating_w_per_degc > 0) rows.push(`
${prefix}heating${d.heating_w_per_degc.toFixed(0)} W/°C
`); + if (d.buckets_warm != null) rows.push(`
${prefix}buckets warm${d.buckets_warm}/${d.buckets_total}
`); + rows.push(`
${prefix}last update${fmtAge(d.last_ms)}
`); + if (d.quality != null) { + const quality = Math.max(0, Math.min(1, d.quality)); + const qualityPct = (quality * 100).toFixed(0); + const qualityColor = quality >= 0.7 ? '#22c55e' : quality >= 0.3 ? '#fbbf24' : '#ef4444'; + rows.push(`
${prefix}quality${qualityPct}%
`); + rows.push(`
`); + } + return rows.join(''); + } + + function twinCard(title, d, endpoint, extraHtml) { + const action = actions[endpoint]; + const statusID = 'twin-status-' + action.id; + if (!d || d.unavailable) return `

${title}

model stateUnavailable
The box did not provide this model state.
`; + if (!d.enabled) return `

${title}

model stateDisabled
This model is disabled on the box.
`; + const learning = d.learning; + const message = actionMessages.get(endpoint) || (!actionAvailable(d) ? 'Relearning is unavailable for this model.' : ''); + const learningRows = [ + `
engine${engineLabel(learning)}
`, + `
learning state${learningStatus(learning)}
`, + ]; + if (learning) { + learningRows.push(`
learning started${fmtLocalTime(learning.started_ms)}
`); + learningRows.push(`
latest training${fmtLocalTime(learning.latest_training_ms)}
`); } - if (d.buckets_warm != null) rows.push(`
buckets warm${d.buckets_warm}/${d.buckets_total}
`); - rows.push(`
last update${fmtAge(d.last_ms)}
`); - rows.push(`
quality${qPct}%
`); - rows.push(`
`); - const btn = resetEndpoint ? resetButton(resetEndpoint, resetLabel) : ''; - return `

${title}

${extraHtml || ''}${rows.join('')}${btn}
`; + return `

${title}

${extraHtml || ''}${learningRows.join('')}${modelRows(d)}${resetButton(endpoint, d)}
${esc(message)}
`; } - function render(pv, load) { + function focusedAction() { + const active = document.activeElement; + return active && active.dataset && actions[active.dataset.resetTwin] ? active.dataset.resetTwin : null; + } + function focusedControl() { + const active = document.activeElement; + if (!active || !active.dataset) return null; + if (actions[active.dataset.resetTwin]) return `[data-reset-twin="${active.dataset.resetTwin}"]`; + if (active.dataset.loadmodelProfile) return `[data-loadmodel-profile="${active.dataset.loadmodelProfile}"]`; + return null; + } + function focusMayReturn() { + const active = document.activeElement; + return !active || active === document.body || active === document.documentElement; + } + function render(pv, load, restoreFocus) { const grid = document.getElementById('twins-grid'); if (!grid) return; - grid.innerHTML = twinCard('PV twin', pv, '/api/pvmodel/reset', 'PV twin') - + twinCard('Load twin', load, '/api/loadmodel/reset', 'load twin', loadProfileControl(load)); + const selector = restoreFocus || focusedControl(); + availableActions.clear(); + grid.innerHTML = twinCard('Solar production', pv, '/api/pvmodel/reset') + twinCard('Consumption', load, '/api/loadmodel/reset', loadProfileControl(load)); + if (selector && grid.querySelector) { + const button = grid.querySelector(selector); + if (button && !button.disabled && typeof button.focus === 'function') button.focus(); + } const sub = document.getElementById('twins-subtitle'); - if (sub) sub.textContent = 'self-learning digital twins — feed MPC + UI forecasts'; + if (sub) sub.textContent = 'Forecast learning for solar production and consumption'; + } + + async function failedResponse(response) { + let body = null; + try { + body = await response.json(); + } catch (_) { /* A non-JSON error still has an HTTP status worth showing. */ } + const detail = body && (body.error || body.message) ? ': ' + (body.error || body.message) : ''; + return { body, error: new Error('The box did not accept this request (HTTP ' + response.status + ')' + detail) }; + } + async function startRelearn(endpoint) { + const action = actions[endpoint]; + if (!action || pending.has(endpoint) || !availableActions.has(endpoint)) return; + if (!confirm(`Start a new learning period for ${action.signal}?\n\nThis keeps measured history and the ${action.other} model. Forecast confidence will be lower while ${action.signal} learns.`)) return; + if (focusedAction() === endpoint) actionFocus.add(endpoint); + pending.add(endpoint); + actionMessages.set(endpoint, 'Requesting a new learning period…'); + render(lastPV, lastLoad); + try { + const response = await apiFetch(endpoint, { method: 'POST' }); + if (!response.ok) { + const failure = await failedResponse(response); + if (response.status === 503 && failure.body && failure.body.status === 'pending') { + actionMessages.set(endpoint, failure.body.error || 'Learning period saved; model restart pending.'); + await fetchAll(); + return; + } + throw failure.error; + } + actionMessages.set(endpoint, 'The box accepted the request. The learning state above reports its progress.'); + await fetchAll(); + } catch (err) { + const cancelled = err && err.name === 'AbortError'; + actionMessages.set(endpoint, cancelled ? 'The request was cancelled. The box did not confirm a new learning period.' : 'The box did not confirm a new learning period: ' + ((err && err.message) || 'request failed')); + } finally { + pending.delete(endpoint); + const restoreFocus = actionFocus.has(endpoint) && focusMayReturn() + ? `[data-reset-twin="${endpoint}"]` : null; + actionFocus.delete(endpoint); + render(lastPV, lastLoad, restoreFocus); + } } - // Wired once on init — delegation off the grid so dynamically-rendered - // buttons pick up the handler without rebinding each refresh. function onGridClick(e) { const profile = e.target && e.target.dataset && e.target.dataset.loadmodelProfile; if (profile) { if (e.target.classList.contains('active')) return; - apiFetch('/api/loadmodel/profile', { - method: 'POST', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ profile }) - }) + apiFetch('/api/loadmodel/profile', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ profile }) }) .then(r => { if (!r.ok) throw new Error('HTTP ' + r.status); return r.json(); }) .then(() => fetchAll()) - .catch(err => alert('Load profile switch failed: ' + err.message)); + .catch(err => { actionMessages.set('/api/loadmodel/reset', 'Load profile switch failed: ' + err.message); render(lastPV, lastLoad); }); return; } - const endpoint = e.target && e.target.dataset && e.target.dataset.resetTwin; - if (!endpoint) return; - const twinName = endpoint.indexOf('pv') >= 0 ? 'PV twin' : 'load twin'; - if (!confirm(`Reset ${twinName} to fresh defaults?\n\n` + - 'All learned samples will be wiped. The model re-trains from the ' + - 'physics / bucket prior; expect ~50 minutes of lower-quality ' + - 'predictions while it collects samples again.')) { - return; - } - apiFetch(endpoint, { method: 'POST' }) - .then(r => { if (!r.ok) throw new Error('HTTP ' + r.status); return r.json(); }) - .then(() => fetchAll()) - .catch(err => alert('Reset failed: ' + err.message)); + if (endpoint) startRelearn(endpoint); } - function init() { const grid = document.getElementById('twins-grid'); if (grid) grid.addEventListener('click', onGridClick); document.addEventListener('ftw-ui-mode-change', syncPolling); syncPolling(); } - - if (document.readyState === 'loading') { - document.addEventListener('DOMContentLoaded', init); - } else { - init(); - } + if (document.readyState === 'loading') document.addEventListener('DOMContentLoaded', init); + else init(); })(); diff --git a/web/twins.test.mjs b/web/twins.test.mjs new file mode 100644 index 00000000..616a998a --- /dev/null +++ b/web/twins.test.mjs @@ -0,0 +1,213 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; +import vm from "node:vm"; + +const source = readFileSync(new URL("./twins.js", import.meta.url), "utf8"); +const now = Date.UTC(2026, 8, 7, 10, 30, 0); + +function model(status = "learning") { + return { + enabled: true, + samples: 41, + mae_w: 182, + quality: 0.6, + last_ms: now - 60_000, + learning: { + engine: "energyplan", + status, + started_ms: now - 3_600_000, + latest_training_ms: now - 60_000, + reset_available: status !== "unavailable", + }, + }; +} + +async function settle() { + for (let i = 0; i < 5; i++) await new Promise(resolve => setImmediate(resolve)); +} + +function load({ pv = model(), loadModel = model("ready"), confirm = () => true, post } = {}) { + const listeners = new Map(); + let html = ""; + let focused = null; + let document; + const grid = { + addEventListener(type, handler) { listeners.set(type, handler); }, + querySelector(selector) { + const match = selector.match(/data-(reset-twin|loadmodel-profile)="([^"]+)"/); + if (!match || !html.includes(`data-${match[1]}="${match[2]}"`)) return null; + const tag = html.match(new RegExp(`]*data-${match[1]}="${match[2]}"[^>]*>`)); + return { + disabled: /\sdisabled(?:\s|=|>)/.test(tag?.[0] || ""), + focus() { + document.activeElement = { dataset: { [match[1]]: match[2] } }; + focused = match[1] === "reset-twin" ? match[2] : "profile:" + match[2]; + }, + }; + }, + }; + Object.defineProperty(grid, "innerHTML", { + get() { return html; }, + set(value) { + if (document?.activeElement?.dataset?.resetTwin || document?.activeElement?.dataset?.loadmodelProfile) document.activeElement = document.body; + html = String(value); + }, + }); + const subtitle = { textContent: "" }; + const requests = []; + document = { + readyState: "complete", + body: { classList: { contains: value => value === "advanced" } }, + activeElement: null, + addEventListener() {}, + getElementById(id) { + return id === "twins-grid" ? grid : id === "twins-subtitle" ? subtitle : null; + }, + }; + const response = (body, ok = true, status = ok ? 200 : 500) => ({ ok, status, json: async () => body }); + const fetch = (path, options = {}) => { + requests.push({ path, options }); + if (options.method === "POST") return post ? post(path, options) : Promise.resolve(response({})); + if (path === "/api/pvmodel") return Promise.resolve(response(pv)); + if (path === "/api/loadmodel") return Promise.resolve(response(loadModel)); + return Promise.resolve(response({})); + }; + vm.runInNewContext(source, { + document, fetch, confirm, setInterval: () => 1, clearInterval() {}, Date: class extends Date { + static now() { return now; } + }, Number, Math, String, Map, Set, Promise, Error, + }, { filename: "twins.js" }); + return { + grid, subtitle, requests, + click(endpoint) { + listeners.get("click")({ target: { dataset: { resetTwin: endpoint } } }); + }, + setFocused(endpoint) { document.activeElement = { dataset: { resetTwin: endpoint } }; }, + clickProfile(profile) { + listeners.get("click")({ target: { dataset: { loadmodelProfile: profile }, classList: { contains: () => false } } }); + }, + setFocusedProfile(profile) { document.activeElement = { dataset: { loadmodelProfile: profile } }; focused = "profile:" + profile; }, + focused: () => focused, + }; +} + +test("renders the Energyplan learning state, local training times, and legacy stats as secondary", async () => { + const ui = load(); + await settle(); + + assert.match(ui.grid.innerHTML, /

Solar production<\/h3>/); + assert.match(ui.grid.innerHTML, /engine<\/span>Energyplan<\/b>/); + assert.match(ui.grid.innerHTML, /learning state<\/span>Learning<\/b>/); + assert.match(ui.grid.innerHTML, /learning started<\/span>.*2026.*<\/b>/); + assert.match(ui.grid.innerHTML, /legacy model stats/); + assert.match(ui.grid.innerHTML, /legacy samples/); + assert.match(ui.grid.innerHTML, /role="status" aria-live="polite"/); + assert.doesNotMatch(ui.grid.innerHTML, /50 minutes/i); + assert.equal(ui.subtitle.textContent, "Forecast learning for solar production and consumption"); +}); + +test("confirmation cancellation sends no reset and names the protected history and other model", async () => { + const questions = []; + const ui = load({ confirm: question => { questions.push(question); return false; } }); + await settle(); + ui.click("/api/pvmodel/reset"); + await settle(); + + assert.equal(ui.requests.filter(request => request.options.method === "POST").length, 0); + assert.match(questions[0], /solar production/); + assert.match(questions[0], /measured history/i); + assert.match(questions[0], /consumption model/i); + assert.match(questions[0], /lower/i); +}); + +test("a pending relearn disables only its action, keeps focus through a repaint, and prevents a second POST", async () => { + let resolvePost; + const post = () => new Promise(resolve => { resolvePost = resolve; }); + const ui = load({ post }); + await settle(); + ui.setFocused("/api/pvmodel/reset"); + ui.click("/api/pvmodel/reset"); + ui.click("/api/pvmodel/reset"); + + assert.match(ui.grid.innerHTML, /Starting new learning period…/); + assert.match(ui.grid.innerHTML, /data-reset-twin="\/api\/pvmodel\/reset"[^>]*disabled/); + assert.equal(ui.focused(), null, "the pending disabled button cannot keep focus"); + assert.equal(ui.requests.filter(request => request.options.method === "POST").length, 1); + + resolvePost({ ok: true, status: 200, json: async () => ({}) }); + await settle(); + assert.match(ui.grid.innerHTML, /The box accepted the request/); + assert.equal(ui.focused(), "/api/pvmodel/reset", "completion returns keyboard focus to the enabled action"); +}); + +test("completion does not steal focus moved to another control, and profile focus survives its refresh", async () => { + let resolvePost; + const ui = load({ post: () => new Promise(resolve => { resolvePost = resolve; }) }); + await settle(); + ui.setFocused("/api/pvmodel/reset"); + ui.click("/api/pvmodel/reset"); + ui.setFocusedProfile("away"); + resolvePost({ ok: true, status: 200, json: async () => ({}) }); + await settle(); + assert.equal(ui.focused(), "profile:away", "a user focus move wins over automatic restore"); + + ui.setFocusedProfile("away"); + ui.clickProfile("away"); + await settle(); + assert.equal(ui.focused(), "profile:away", "profile changes retain keyboard focus after their GET repaint"); +}); + +test("an unavailable model has no action, and an aborted request does not claim it started learning", async () => { + const unavailable = load({ pv: { enabled: true, learning: { engine: "energyplan", status: "unavailable", started_ms: 0, latest_training_ms: 0, reset_available: false } } }); + await settle(); + assert.match(unavailable.grid.innerHTML, /Relearning is unavailable for this model/); + assert.match(unavailable.grid.innerHTML, /data-reset-twin="\/api\/pvmodel\/reset"[^>]*disabled/); + + const aborted = load({ post: () => Promise.reject(Object.assign(new Error("offline"), { name: "AbortError" })) }); + await settle(); + aborted.click("/api/loadmodel/reset"); + await settle(); + assert.match(aborted.grid.innerHTML, /request was cancelled/i); + assert.doesNotMatch(aborted.grid.innerHTML, /started learning/i); +}); + +test("a refused relearn reports the box response and never claims that it began", async () => { + const ui = load({ + post: () => Promise.resolve({ ok: false, status: 409, json: async () => ({ error: "solar meter is stale" }) }), + }); + await settle(); + ui.click("/api/pvmodel/reset"); + await settle(); + + assert.match(ui.grid.innerHTML, /did not confirm a new learning period/i); + assert.match(ui.grid.innerHTML, /HTTP 409.*solar meter is stale/i); + assert.doesNotMatch(ui.grid.innerHTML, /The box accepted the request/); +}); + +test("a saved reset with a pending worker restart is retried from the GET state", async () => { + const retryModel = model("unavailable"); + retryModel.learning.reset_available = true; + const ui = load({ + pv: retryModel, + post: () => Promise.resolve({ + ok: false, + status: 503, + json: async () => ({ + status: "pending", + error: "Learning period saved; model restart pending: worker is unavailable", + learning: retryModel.learning, + }), + }), + }); + await settle(); + assert.doesNotMatch(ui.grid.innerHTML, /data-reset-twin="\/api\/pvmodel\/reset"[^>]*disabled/); + + ui.click("/api/pvmodel/reset"); + await settle(); + + assert.equal(ui.requests.filter(request => request.path === "/api/pvmodel").length, 2, "pending response must refresh model state"); + assert.match(ui.grid.innerHTML, /Learning period saved; model restart pending: worker is unavailable/); + assert.doesNotMatch(ui.grid.innerHTML, /did not confirm a new learning period/i); + assert.doesNotMatch(ui.grid.innerHTML, /data-reset-twin="\/api\/pvmodel\/reset"[^>]*disabled/); +});