From d3b0ee73f888bebed3c041a2a067717aa530ca5e Mon Sep 17 00:00:00 2001 From: Fredrik Ahlgren Date: Thu, 17 Sep 2026 09:06:45 +0200 Subject: [PATCH 1/2] fix(state): bound maintenance writes and retain forecast observations Signed-off-by: Fredrik Ahlgren --- .changeset/storage-maintenance-recovery.md | 5 + go/cmd/ftw/forecast_learning.go | 58 ++++- go/cmd/ftw/forecast_observation_queue.go | 123 +++++++++++ go/cmd/ftw/forecast_observation_queue_test.go | 206 ++++++++++++++++++ go/cmd/ftw/forecast_rust.go | 2 +- go/cmd/ftw/forecast_rust_learning.go | 2 +- go/cmd/ftw/forecast_tracking.go | 107 +++++---- go/cmd/ftw/main.go | 44 +--- go/cmd/ftw/rolloff_test.go | 6 +- go/internal/api/api.go | 2 +- go/internal/api/api_history_storage_test.go | 38 ++++ go/internal/forecasting/learning.go | 2 + go/internal/state/configuration.go | 35 ++- go/internal/state/energy_ledger.go | 41 ---- go/internal/state/forecast_issues.go | 104 +++------ go/internal/state/forecast_prune.go | 124 +++++++++++ go/internal/state/history_maintenance.go | 84 +++++++ go/internal/state/history_rollup.go | 158 ++++++++++++++ .../state/history_rollup_contention_test.go | 131 +++++++++++ go/internal/state/history_sqlite.go | 1 + go/internal/state/history_writer.go | 6 + .../state/history_writer_batch_test.go | 3 + go/internal/state/maintenance.go | 16 +- .../state/maintenance_contention_test.go | 144 ++++++++++++ go/internal/state/parquet_stream.go | 7 +- go/internal/state/store.go | 49 +---- 26 files changed, 1227 insertions(+), 271 deletions(-) create mode 100644 .changeset/storage-maintenance-recovery.md create mode 100644 go/cmd/ftw/forecast_observation_queue.go create mode 100644 go/cmd/ftw/forecast_observation_queue_test.go create mode 100644 go/internal/state/forecast_prune.go create mode 100644 go/internal/state/history_maintenance.go create mode 100644 go/internal/state/history_rollup.go create mode 100644 go/internal/state/history_rollup_contention_test.go create mode 100644 go/internal/state/maintenance_contention_test.go diff --git a/.changeset/storage-maintenance-recovery.md b/.changeset/storage-maintenance-recovery.md new file mode 100644 index 000000000..23de29654 --- /dev/null +++ b/.changeset/storage-maintenance-recovery.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Keep live storage responsive during archive and history maintenance. Bound write transactions, retain completed forecast observations for retry, and expose current storage and forecast health with failure history for diagnostics. diff --git a/go/cmd/ftw/forecast_learning.go b/go/cmd/ftw/forecast_learning.go index d0a3d37cb..d3c78ab44 100644 --- a/go/cmd/ftw/forecast_learning.go +++ b/go/cmd/ftw/forecast_learning.go @@ -201,7 +201,7 @@ func (f *forecastTracker) RestartLearning(ctx context.Context, signal string) er } func (f *forecastTracker) LearningStatus(signal string) forecasting.LearningStatus { - status := forecasting.LearningStatus{Engine: "legacy", Status: "unavailable"} + status := forecasting.LearningStatus{Engine: "legacy", Status: "unavailable", Health: "unknown"} if f == nil { return status } @@ -249,9 +249,65 @@ func (f *forecastTracker) LearningStatus(signal string) forecasting.LearningStat if site.IdentityPending || (f.learningPeriods.ConfigRevision == rustConfigRevision(site) && f.learningErrors[signal] != nil) { status.Status = "unavailable" } + status.Health, status.HealthReason = f.learningHealth(signal, status) return status } +// Health is current pipeline health, separate from the model's learning stage. +// A ready model alone cannot establish working collection or persistence. +func (f *forecastTracker) learningHealth(signal string, status forecasting.LearningStatus) (string, string) { + f.mu.RLock() + defer f.mu.RUnlock() + if f.stopped { + return "unknown", "stopped" + } + if f.observationOverflow { + return "degraded", "observation_queue_full" + } + if f.observationError != "" { + return "degraded", f.observationError + } + if f.issueArchiveError { + return "degraded", "forecast_archive_pending" + } + if f.store != nil { + writer := f.store.HistoryWriterStatus() + if writer.LastError != "" { + return "degraded", "history_write_pending" + } + if writer.MaintenanceError != "" || f.store.HistoryMaintenanceStatus().LastError != "" { + return "degraded", "history_maintenance_failed" + } + } + if status.Status == "unavailable" { + return "unknown", "model_unavailable" + } + if f.observationCheckedMS == 0 { + return "unknown", "not_checked" + } + age := f.now().UnixMilli() - f.observationCheckedMS + valid := f.observationLoadValid + trainingAgeLimit := 2 * time.Hour + if signal == "pv" { + valid = f.observationPVValid + trainingAgeLimit = 36 * time.Hour + } + if age < 0 || age > (30*time.Second).Milliseconds() || !valid { + return "waiting_for_data", "measurements_unavailable" + } + if status.Status != "ready" { + return "unknown", "learning" + } + trainingAge := f.now().UnixMilli() - status.LatestTrainingMS + if status.LatestTrainingMS <= 0 || trainingAge < 0 || trainingAge > trainingAgeLimit.Milliseconds() { + return "waiting_for_data", "training_outdated" + } + if status.Engine != "energyplan" { + return "unknown", "model_persistence_unchecked" + } + return "healthy", "" +} + // 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) { diff --git a/go/cmd/ftw/forecast_observation_queue.go b/go/cmd/ftw/forecast_observation_queue.go new file mode 100644 index 000000000..8eb1982e7 --- /dev/null +++ b/go/cmd/ftw/forecast_observation_queue.go @@ -0,0 +1,123 @@ +package main + +import ( + "context" + "errors" + "log/slog" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/state" +) + +// Keep a bounded backlog of completed intervals. The observer owns this queue; +// issuing forecasts and scoring them must not postpone the next measurement. +const maxPendingForecastObservations = 64 + +type forecastObservationJob struct { + observation forecasting.Observation + site forecastSite + weather *state.ForecastPoint + away bool + archived bool +} + +func (f *forecastTracker) runObservations(ctx context.Context) { + defer func() { + drain, cancel := context.WithTimeout(context.WithoutCancel(ctx), 3*time.Second) + defer cancel() + f.flushObservations(drain) + if len(f.pendingObservations) > 0 { + slog.Error("forecast archive: shutdown observations still pending", "intervals", len(f.pendingObservations)) + } + }() + tick := time.NewTicker(10 * time.Second) + defer tick.Stop() + for { + select { + case <-ctx.Done(): + return + case <-tick.C: + f.observe(ctx) + } + } +} + +func (f *forecastTracker) enqueueObservation(job forecastObservationJob) { + if len(f.pendingObservations) == maxPendingForecastObservations { + f.mu.Lock() + f.observationOverflow = true + f.mu.Unlock() + slog.Error("forecast archive: observation queue full", "start_ms", job.observation.StartMS) + return + } + f.pendingObservations = append(f.pendingObservations, job) +} + +func (f *forecastTracker) flushObservations(ctx context.Context) { + f.flushObservationsWith(ctx, f.store.SaveForecastObservation, f.updateObservation) +} + +// Save and update retries keep the original evidence and captured features. +// A successful save followed by a failed model update skips saving on retry. +func (f *forecastTracker) flushObservationsWith(ctx context.Context, + save func(context.Context, forecasting.Observation) error, + update func(context.Context, forecastObservationJob) error) { + for len(f.pendingObservations) > 0 { + if ctx.Err() != nil { + return + } + job := &f.pendingObservations[0] + writeCtx, cancel := context.WithTimeout(ctx, 2*time.Second) + phase := "observation_save_pending" + var err error + if !job.archived { + err = save(writeCtx, job.observation) + if err == nil { + job.archived = true + } + } + if err == nil { + phase = "model_save_pending" + err = update(writeCtx, *job) + } + if err != nil && writeCtx.Err() != nil { + err = errors.Join(err, writeCtx.Err()) + } + cancel() + f.mu.Lock() + if err != nil { + f.observationError = phase + } else { + f.observationError = "" + } + f.mu.Unlock() + if err != nil { + slog.Warn("forecast archive: observation retained for retry", "phase", phase, "start_ms", job.observation.StartMS, "err", err) + return + } + f.pendingObservations[0] = forecastObservationJob{} + f.pendingObservations = f.pendingObservations[1:] + f.requestScoring() + } +} + +func (f *forecastTracker) updateObservation(ctx context.Context, job forecastObservationJob) error { + if f.candidate == nil { + return nil + } + if f.configMu != nil { + f.configMu.RLock() + defer f.configMu.RUnlock() + } + f.learningMu.RLock() + defer f.learningMu.RUnlock() + current := f.site() + // Old configuration evidence remains in the archive, but must not switch + // the current model back to a previous identity when the queue recovers. + if current.IdentityPending || current.Revision != job.site.Revision { + return nil + } + site := f.learningSiteLocked(job.site) + return f.candidate.Update(ctx, site, job.observation, job.weather, job.away) +} diff --git a/go/cmd/ftw/forecast_observation_queue_test.go b/go/cmd/ftw/forecast_observation_queue_test.go new file mode 100644 index 000000000..e6422bd3a --- /dev/null +++ b/go/cmd/ftw/forecast_observation_queue_test.go @@ -0,0 +1,206 @@ +package main + +import ( + "context" + "database/sql" + "errors" + "path/filepath" + "reflect" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/forecasting" + "github.com/srcfl/ftw/go/internal/state" +) + +func observationJobFixture() forecastObservationJob { + at := time.Now().Add(-time.Hour).Truncate(time.Hour) + return forecastObservationJob{site: trackerSite(), observation: forecasting.Observation{ + StartMS: at.UnixMilli(), EndMS: at.Add(15 * time.Minute).UnixMilli(), AvailableAtMS: at.Add(15 * time.Minute).UnixMilli(), + LoadW: 500, PVW: 1000, LoadKnown: true, PVKnown: true, Quality: "complete_balance_v1", ConfigVersion: trackerSite().Revision, + }} +} + +func TestForecastObservationRetainedAcrossRealSQLiteDeadline(t *testing.T) { + path := filepath.Join(t.TempDir(), "state.db") + s, err := state.Open(path) + if err != nil { + t.Fatal(err) + } + defer s.Close() + if err := s.InitForecastArchive(context.Background()); err != nil { + t.Fatal(err) + } + blocker, err := sql.Open("sqlite", path) + if err != nil { + t.Fatal(err) + } + defer blocker.Close() + blocker.SetMaxOpenConns(1) + if _, err := blocker.Exec("BEGIN IMMEDIATE"); err != nil { + t.Fatal(err) + } + defer blocker.Exec("ROLLBACK") + f := trackerFixture(time.Now()) + f.store = s + job := observationJobFixture() + f.enqueueObservation(job) + f.flushObservations(context.Background()) + if len(f.pendingObservations) != 1 || f.observationError != "observation_save_pending" { + t.Fatalf("failed interval lost or hidden: queue=%d error=%q", len(f.pendingObservations), f.observationError) + } + if _, err := blocker.Exec("ROLLBACK"); err != nil { + t.Fatal(err) + } + f.flushObservations(context.Background()) + if len(f.pendingObservations) != 0 || f.observationError != "" { + t.Fatal("queue did not recover") + } + rows, err := s.LoadForecastObservations(context.Background(), job.observation.StartMS, job.observation.EndMS) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(rows, []forecasting.Observation{job.observation}) { + t.Fatalf("retry changed measured evidence: %+v", rows) + } +} + +func TestForecastObservationRetryPreservesArchiveAndModelStages(t *testing.T) { + for _, phase := range []string{"ambiguous archive", "model save"} { + t.Run(phase, func(t *testing.T) { + f := trackerFixture(time.Now()) + job := observationJobFixture() + f.enqueueObservation(job) + saves, updates := 0, 0 + save := func(_ context.Context, o forecasting.Observation) error { + saves++ + if o != job.observation { + t.Fatal("retry changed interval") + } + if phase == "ambiguous archive" && saves == 1 { + return context.DeadlineExceeded + } + return nil + } + update := func(_ context.Context, got forecastObservationJob) error { + updates++ + if got.observation != job.observation { + t.Fatal("retry changed training input") + } + if phase == "model save" && updates == 1 { + return errors.New("disk busy") + } + return nil + } + f.flushObservationsWith(context.Background(), save, update) + if len(f.pendingObservations) != 1 || f.observationError == "" { + t.Fatal("failure not retained") + } + f.flushObservationsWith(context.Background(), save, update) + if len(f.pendingObservations) != 0 || f.observationError != "" { + t.Fatal("recovery did not complete") + } + if phase == "model save" && saves != 1 { + t.Fatal("already committed archive saved again") + } + }) + } +} + +func TestForecastObservationQueueBoundAndCancellation(t *testing.T) { + f := trackerFixture(time.Now()) + job := observationJobFixture() + for i := 0; i <= maxPendingForecastObservations; i++ { + f.enqueueObservation(job) + } + if len(f.pendingObservations) != maxPendingForecastObservations || !f.observationOverflow { + t.Fatal("overflow unbounded or hidden") + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + f.flushObservationsWith(ctx, func(context.Context, forecasting.Observation) error { t.Fatal("save after cancellation"); return nil }, nil) + if len(f.pendingObservations) != maxPendingForecastObservations { + t.Fatal("cancellation dropped pending evidence") + } + status := forecasting.LearningStatus{Engine: "energyplan", Status: "ready", LatestTrainingMS: time.Now().UnixMilli()} + if health, _ := f.learningHealth("load", status); health != "degraded" { + t.Fatal("queue overflow reported healthy") + } +} + +func TestForecastLearningHealthRequiresCurrentPipelineEvidence(t *testing.T) { + now := time.Now() + f := trackerFixture(now) + status := forecasting.LearningStatus{Engine: "energyplan", Status: "ready", LatestTrainingMS: now.UnixMilli()} + check := func(want string) { + t.Helper() + if health, reason := f.learningHealth("load", status); health != want { + t.Fatalf("health=%s reason=%s want=%s", health, reason, want) + } + } + check("unknown") + f.observationCheckedMS, f.observationLoadValid = now.UnixMilli(), true + check("healthy") + f.observationError = "model_save_pending" + check("degraded") + f.observationError = "" + f.issueArchiveError = true + check("degraded") + f.issueArchiveError = false + f.observationLoadValid = false + check("waiting_for_data") + f.observationLoadValid = true + status.LatestTrainingMS = now.Add(-3 * time.Hour).UnixMilli() + check("waiting_for_data") +} + +func TestForecastObservationShutdownDrainsWithFreshContext(t *testing.T) { + s, err := state.Open(filepath.Join(t.TempDir(), "state.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + if err := s.InitForecastArchive(context.Background()); err != nil { + t.Fatal(err) + } + f := trackerFixture(time.Now()) + f.store = s + f.enqueueObservation(observationJobFixture()) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + f.runObservations(ctx) + if len(f.pendingObservations) != 0 { + t.Fatal("shutdown did not save pending observation") + } +} + +type observationCandidate struct { + trackingCandidate + calls int +} + +func (c *observationCandidate) Update(context.Context, forecastSite, forecasting.Observation, *state.ForecastPoint, bool) error { + c.calls++ + return nil +} + +func TestForecastObservationOldRevisionCannotReplaceCurrentModel(t *testing.T) { + f := trackerFixture(time.Now()) + candidate := &observationCandidate{} + f.candidate = candidate + job := observationJobFixture() + job.site.Revision = "previous" + if err := f.updateObservation(context.Background(), job); err != nil { + t.Fatal(err) + } + if candidate.calls != 0 { + t.Fatal("old configuration changed the current model") + } + job.site = trackerSite() + if err := f.updateObservation(context.Background(), job); err != nil { + t.Fatal(err) + } + if candidate.calls != 1 { + t.Fatal("current observation did not update model") + } +} diff --git a/go/cmd/ftw/forecast_rust.go b/go/cmd/ftw/forecast_rust.go index 1ab339a94..6b8e1b995 100644 --- a/go/cmd/ftw/forecast_rust.go +++ b/go/cmd/ftw/forecast_rust.go @@ -159,7 +159,7 @@ func (r *rustForecast) Update(ctx context.Context, site forecastSite, o forecast return err } // Persist the complete response atomically before exposing it to planning. - if err = r.store.SaveConfig(forecastRustStateKey, string(data)); err != nil { + if err = r.store.SaveConfigContext(ctx, forecastRustStateKey, string(data)); err != nil { return err } r.mu.Lock() diff --git a/go/cmd/ftw/forecast_rust_learning.go b/go/cmd/ftw/forecast_rust_learning.go index 4685f03d5..70ced24d6 100644 --- a/go/cmd/ftw/forecast_rust_learning.go +++ b/go/cmd/ftw/forecast_rust_learning.go @@ -70,7 +70,7 @@ func (r *rustForecast) applyLearningPeriodsLocked(ctx context.Context, site fore if err != nil { return err } - if err = r.store.SaveConfig(forecastRustStateKey, string(data)); err != nil { + if err = r.store.SaveConfigContext(ctx, forecastRustStateKey, string(data)); err != nil { return err } r.mu.Lock() diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index 45144f090..30cab6e0a 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -53,32 +53,39 @@ type forecastJob struct { } type forecastTracker struct { - refreshIdentity func() - configMu *sync.RWMutex - store *state.Store - tele *telemetry.Store - pv *pvmodel.Service - load *loadmodel.Service - site func() forecastSite - away func(time.Time) bool - curtailed func(time.Time) bool - candidate forecastCandidate - queue chan forecastJob - scoreWake chan struct{} - cancel context.CancelFunc - wg sync.WaitGroup - mu sync.RWMutex - errors []forecasting.ErrorSample - observations []forecasting.Observation - accumulator telemetry.ForecastAccumulator - lastRevision string - pvAccumulator telemetry.ForecastAccumulator - stopped bool - clock func() time.Time - learningMu sync.RWMutex - learningPeriods forecastLearningPeriods - learningErrors map[string]error - requestReplan func(string) + refreshIdentity func() + configMu *sync.RWMutex + store *state.Store + tele *telemetry.Store + pv *pvmodel.Service + load *loadmodel.Service + site func() forecastSite + away func(time.Time) bool + curtailed func(time.Time) bool + candidate forecastCandidate + queue chan forecastJob + scoreWake chan struct{} + cancel context.CancelFunc + wg sync.WaitGroup + mu sync.RWMutex + errors []forecasting.ErrorSample + observations []forecasting.Observation + accumulator telemetry.ForecastAccumulator + lastRevision string + pvAccumulator telemetry.ForecastAccumulator + stopped bool + clock func() time.Time + learningMu sync.RWMutex + learningPeriods forecastLearningPeriods + learningErrors map[string]error + requestReplan func(string) + pendingObservations []forecastObservationJob // owned by runObservations + observationError string // guarded by mu + observationOverflow bool + observationCheckedMS int64 + observationPVValid bool + observationLoadValid bool + issueArchiveError bool } func (f *forecastTracker) Start(ctx context.Context) error { @@ -95,9 +102,10 @@ func (f *forecastTracker) Start(ctx context.Context) error { f.scoreWake = make(chan struct{}, 1) workerCtx, stop := context.WithCancel(ctx) f.cancel = stop - f.wg.Add(2) + f.wg.Add(3) go func() { defer f.wg.Done(); f.run(workerCtx) }() go func() { defer f.wg.Done(); f.runScoring(workerCtx) }() + go func() { defer f.wg.Done(); f.runObservations(workerCtx) }() return nil } func (f *forecastTracker) Stop() { @@ -151,8 +159,6 @@ func (f *forecastTracker) run(ctx context.Context) { } }() f.refreshEvidence(ctx, f.now()) - tick := time.NewTicker(10 * time.Second) - defer tick.Stop() for { select { case <-ctx.Done(): @@ -160,6 +166,9 @@ func (f *forecastTracker) run(ctx context.Context) { case job := <-f.queue: pending = &job err := saveForecastIssueWithRetry(ctx, job.issue, f.store.SaveForecastIssue) + f.mu.Lock() + f.issueArchiveError = err != nil + f.mu.Unlock() if err != nil { if ctx.Err() != nil { return @@ -169,8 +178,6 @@ func (f *forecastTracker) run(ctx context.Context) { continue } pending = nil - case <-tick.C: - f.observe(ctx) } } } @@ -253,33 +260,23 @@ func (f *forecastTracker) observe(ctx context.Context) { r.PVValid = false r.PVReason = "curtailed" } + f.mu.Lock() + f.observationCheckedMS = now.UnixMilli() + f.observationPVValid, f.observationLoadValid = r.PVValid, r.Valid + f.mu.Unlock() for _, o := range f.observationIntervals(r, site, now) { - writeCtx, cancel := context.WithTimeout(ctx, 2*time.Second) - err := f.store.SaveForecastObservation(writeCtx, o) - cancel() - if err != nil { - slog.Warn("forecast archive: observation not saved", "err", err) - continue + job := forecastObservationJob{observation: o, site: site, away: f.away != nil && f.away(time.UnixMilli(o.StartMS))} + rows, _ := f.store.LoadForecasts(o.StartMS-time.Hour.Milliseconds(), o.EndMS) + if weather := forecastRow(usableTrackingWeather(rows, o.AvailableAtMS), o.StartMS, o.AvailableAtMS); weather != nil && weather.FetchedAtMs >= site.WeatherSinceMS { + frozen := *weather + job.weather = &frozen } - if f.candidate != nil { - rows, _ := f.store.LoadForecasts(o.StartMS-time.Hour.Milliseconds(), o.EndMS) - weather := forecastRow(usableTrackingWeather(rows, now.UnixMilli()), o.StartMS, now.UnixMilli()) - if weather != nil && weather.FetchedAtMs < site.WeatherSinceMS { - weather = nil - } - 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) - } - } - f.requestScoring() + f.enqueueObservation(job) } + // A slow disk must not make draining a backlog postpone measurement capture. + flushCtx, cancel := context.WithTimeout(ctx, 3*time.Second) + defer cancel() + f.flushObservations(flushCtx) } // observationIntervals keeps the two measurement claims independent. Unknown diff --git a/go/cmd/ftw/main.go b/go/cmd/ftw/main.go index 0067037e7..70db98a09 100644 --- a/go/cmd/ftw/main.go +++ b/go/cmd/ftw/main.go @@ -3310,21 +3310,11 @@ func rolloffLoop(ctx context.Context, st *state.Store, coldDir string, retention if retentionDays != nil { days = retentionDays() } - doRolloff(ctx, st, coldDir) - if err := st.PruneHistorySamples(ctx, days, time.Now()); err != nil { - slog.Warn("history retention failed", "err", err) + if err := st.MaintainHistory(ctx, coldDir, days, time.Now()); err != nil { + slog.Warn("history maintenance incomplete", "err", err) } - - // The bulk DELETEs above just generated a WAL burst; reclaim it now - // instead of letting the -wal file ratchet upward on the SD card. st.CheckpointWAL() - if removed, err := state.PruneDiagnosticsParquet(coldDir, days, time.Now()); err != nil { - slog.Warn("cold parquet retention prune failed", "err", err) - } else if len(removed) > 0 { - slog.Info("cold parquet retention", "removed_files", len(removed), "retention_days", days) - } - // Disk watch: an SD card that fills up takes SQLite down with it. // Warn loudly (log + event feed) at most once per day. if avail, err := state.DiskAvail(coldDir); err == nil { @@ -3352,36 +3342,6 @@ func rolloffLoop(ctx context.Context, st *state.Store, coldDir string, retention } } -func doRolloff(ctx context.Context, st *state.Store, coldDir string) { - // Age the fixed-column dashboard history on the same cadence as the - // long-format TS + diagnostics rolloff below. Prune is idempotent and pure - // SQL; without this call history_hot/history_warm grow forever even though - // ts_samples is correctly moved to Parquet. - if err := st.Prune(ctx); err != nil { - slog.Warn("history tier prune failed", "err", err) - } - if rolled, expired, err := st.PruneEnergyLedger(ctx, time.Now()); err != nil { - slog.Warn("energy ledger retention failed", "err", err) - } else if rolled > 0 || expired > 0 { - slog.Info("energy ledger retention", "detailed_rows_rolled_up", rolled, "hourly_rows_rolled_to_days", expired) - } - - // Planner diagnostics roll off on the same cadence but keep a - // longer hot tier (30 d vs. the 14 d of ts_samples) — they're - // sparse enough (~100/day) that the extra month in SQLite - // costs < 60 MB and makes the time-travel UI snappy for - // recent-incident debugging. - dRows, dFiles, err := st.RolloffDiagnosticsToParquet(ctx, coldDir) - if err != nil { - slog.Warn("diagnostics parquet rolloff failed", "err", err) - return - } - if dRows > 0 { - slog.Info("diagnostics parquet rolloff", - "rows", dRows, "files", len(dFiles)) - } -} - func flushHistoryOnStop(st *state.Store) { if err := st.StopHistory(); err != nil { slog.Warn("history flush on shutdown", "err", err) diff --git a/go/cmd/ftw/rolloff_test.go b/go/cmd/ftw/rolloff_test.go index 3e20e47f0..c93103ac9 100644 --- a/go/cmd/ftw/rolloff_test.go +++ b/go/cmd/ftw/rolloff_test.go @@ -9,7 +9,7 @@ import ( "github.com/srcfl/ftw/go/internal/state" ) -func TestDoRolloffPrunesFixedColumnHistory(t *testing.T) { +func TestHistoryMaintenancePrunesFixedColumnHistory(t *testing.T) { dir := t.TempDir() st, err := state.Open(filepath.Join(dir, "state.db")) if err != nil { @@ -30,7 +30,9 @@ func TestDoRolloffPrunesFixedColumnHistory(t *testing.T) { t.Fatal(err) } - doRolloff(context.Background(), st, filepath.Join(dir, "cold")) + if err := st.MaintainHistory(context.Background(), filepath.Join(dir, "cold"), 0, time.Now()); err != nil { + t.Fatal(err) + } hot, warm, _, err := st.HistoryCounts() if err != nil { diff --git a/go/internal/api/api.go b/go/internal/api/api.go index 0353dc027..000b3bbf1 100644 --- a/go/internal/api/api.go +++ b/go/internal/api/api.go @@ -707,7 +707,7 @@ func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) { if s.deps.State != nil { resp["history_storage"] = s.deps.State.HistoryBackend() writer := s.deps.State.HistoryWriterStatus() - if writer.LastError != "" || (writer.LastRejectMS > 0 && time.Now().UnixMilli()-writer.LastRejectMS < time.Minute.Milliseconds()) { + if writer.LastError != "" || writer.MaintenanceError != "" || s.deps.State.HistoryMaintenanceStatus().LastError != "" || (writer.LastRejectMS > 0 && time.Now().UnixMilli()-writer.LastRejectMS < time.Minute.Milliseconds()) { resp["status"] = "degraded" } } diff --git a/go/internal/api/api_history_storage_test.go b/go/internal/api/api_history_storage_test.go index f8d5376ed..cbbe73f00 100644 --- a/go/internal/api/api_history_storage_test.go +++ b/go/internal/api/api_history_storage_test.go @@ -1,11 +1,13 @@ package api import ( + "context" "encoding/json" "math" "net/http" "net/http/httptest" "testing" + "time" "github.com/srcfl/ftw/go/internal/state" "github.com/srcfl/ftw/go/internal/telemetry" @@ -34,3 +36,39 @@ func TestHealthShowsRejectedPrimaryHistory(t *testing.T) { t.Fatalf("health hid a collection error: %s", rr.Body.String()) } } + +func TestHealthShowsFailedArchiveAndRecovery(t *testing.T) { + srv, st, _ := newSeriesTestServer(t) + srv.deps.Tel = telemetry.NewStore() + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := st.MaintainHistory(ctx, t.TempDir(), 0, time.Now()); err == nil { + t.Fatal("cancelled maintenance succeeded") + } + check := func(want string) { + t.Helper() + rr := httptest.NewRecorder() + srv.Handler().ServeHTTP(rr, httptest.NewRequest(http.MethodGet, "/api/health", nil)) + var body struct { + Status string `json:"status"` + History struct { + Maintenance state.HistoryMaintenanceStatus `json:"maintenance"` + } `json:"history_storage"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if body.Status != want || body.History.Maintenance.Failures != 1 || body.History.Maintenance.LastFailureMS == 0 { + t.Fatalf("maintenance status missing: %s", rr.Body.String()) + } + } + check("degraded") + if err := st.MaintainHistory(context.Background(), t.TempDir(), 0, time.Now()); err != nil { + t.Fatal(err) + } + check("ok") + status := st.HistoryMaintenanceStatus() + if status.State != "complete" || status.LastError != "" || status.LastFailureError == "" || status.LastSuccessMS == 0 { + t.Fatalf("recovery status=%+v", status) + } +} diff --git a/go/internal/forecasting/learning.go b/go/internal/forecasting/learning.go index 402cc89f1..6b53d19f5 100644 --- a/go/internal/forecasting/learning.go +++ b/go/internal/forecasting/learning.go @@ -19,4 +19,6 @@ type LearningStatus struct { StartedMS int64 `json:"started_ms"` LatestTrainingMS int64 `json:"latest_training_ms"` ResetAvailable bool `json:"reset_available"` + Health string `json:"health"` + HealthReason string `json:"health_reason,omitempty"` } diff --git a/go/internal/state/configuration.go b/go/internal/state/configuration.go index 1a0e867bb..2e1b68889 100644 --- a/go/internal/state/configuration.go +++ b/go/internal/state/configuration.go @@ -89,30 +89,53 @@ func (s *Store) ConfigValue(key string) (string, bool, error) { // History keeps its existing policy. FULL syncs the WAL before acknowledging // settings, rather than waiting for a later checkpoint. func (s *Store) durableConfigWrite(write func(*sql.Tx) error) error { - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + return s.durableConfigWriteContext(context.Background(), write) +} + +func (s *Store) durableConfigWriteContext(parent context.Context, write func(*sql.Tx) error) error { + ctx, cancel := context.WithTimeout(parent, 5*time.Second) defer cancel() conn, err := s.db.Conn(ctx) if err != nil { return err } defer conn.Close() - if _, err := conn.ExecContext(ctx, `PRAGMA synchronous=FULL`); err != nil { + var synchronous int + if err := conn.QueryRowContext(ctx, `PRAGMA synchronous`).Scan(&synchronous); err != nil { return err } defer func() { - if _, err := conn.ExecContext(ctx, `PRAGMA synchronous=NORMAL`); err != nil { + restore, stop := context.WithTimeout(context.Background(), time.Second) + defer stop() + if _, err := conn.ExecContext(restore, fmt.Sprintf(`PRAGMA synchronous=%d`, synchronous)); err != nil { _ = conn.Raw(func(any) error { return driver.ErrBadConn }) } }() + if _, err := conn.ExecContext(ctx, `PRAGMA synchronous=FULL`); err != nil { + return err + } tx, err := conn.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback() - if err := write(tx); err != nil { - return err + err = write(tx) + if err == nil { + err = tx.Commit() + } + if err != nil && ctx.Err() != nil { + return errors.Join(err, ctx.Err()) } - return tx.Commit() + return err +} + +// SaveConfigContext preserves the durable config contract while allowing a +// background model update to cancel without spending another full five seconds. +func (s *Store) SaveConfigContext(ctx context.Context, key, value string) error { + return s.durableConfigWriteContext(ctx, func(tx *sql.Tx) error { + _, err := tx.ExecContext(ctx, `INSERT INTO config (key,value) VALUES (?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value`, key, value) + return err + }) } func saveConfigValues(tx *sql.Tx, values map[string]string) error { diff --git a/go/internal/state/energy_ledger.go b/go/internal/state/energy_ledger.go index c17287ed0..7aa1427ef 100644 --- a/go/internal/state/energy_ledger.go +++ b/go/internal/state/energy_ledger.go @@ -500,47 +500,6 @@ func (s *Store) rollupEnergyLedgerChunk(ctx context.Context, fromMS, toMS int64) return s.rollupEnergyLedgerWidth(ctx, fromMS, toMS, EnergyLedgerRollupBucketMS) } -func (s *Store) rollupEnergyLedgerWidth(ctx context.Context, fromMS, toMS, width int64) (int64, error) { - s.historyWriteMu.Lock() - defer s.historyWriteMu.Unlock() - tx, err := s.history.BeginTx(ctx, nil) - if err != nil { - return 0, err - } - defer tx.Rollback() - if _, err := tx.ExecContext(ctx, `INSERT INTO energy_ledger_entries( - schema_version, asset_id, flow, bucket_start_ms, bucket_len_ms, - energy_wh, source, quality, provenance, sample_count, observed_at_ms - ) - SELECT schema_version, asset_id, flow, - (bucket_start_ms / ?) * ?, ?, SUM(energy_wh), source, quality, provenance, - SUM(sample_count), MAX(observed_at_ms) - FROM energy_ledger_entries - WHERE bucket_len_ms < ? AND bucket_start_ms >= ? AND bucket_start_ms < ? - GROUP BY schema_version, asset_id, flow, - 4, source, quality, provenance - ON CONFLICT(schema_version, asset_id, flow, bucket_start_ms, bucket_len_ms, source, quality, provenance) - DO UPDATE SET - energy_wh = energy_ledger_entries.energy_wh + excluded.energy_wh, - sample_count = energy_ledger_entries.sample_count + excluded.sample_count, - observed_at_ms = MAX(energy_ledger_entries.observed_at_ms, excluded.observed_at_ms)`, - width, width, width, - width, fromMS, toMS); err != nil { - return 0, err - } - res, err := tx.ExecContext(ctx, `DELETE FROM energy_ledger_entries - WHERE bucket_len_ms < ? AND bucket_start_ms >= ? AND bucket_start_ms < ?`, - width, fromMS, toMS) - if err != nil { - return 0, err - } - n, _ := res.RowsAffected() - if err := tx.Commit(); err != nil { - return 0, err - } - return n, nil -} - func min64(a, b int64) int64 { if a < b { return a diff --git a/go/internal/state/forecast_issues.go b/go/internal/state/forecast_issues.go index 551c73253..cf5f4b4dc 100644 --- a/go/internal/state/forecast_issues.go +++ b/go/internal/state/forecast_issues.go @@ -222,38 +222,6 @@ func storePreparedForecastModelStates(ctx context.Context, tx *sql.Tx, prepared return nil } -func cleanForecastModelStateRefs(ctx context.Context, tx *sql.Tx) error { - if _, err := tx.ExecContext(ctx, `DELETE FROM forecast_issue_model_states - WHERE issue_id NOT IN (SELECT id FROM forecast_issues)`); err != nil { - return err - } - _, err := tx.ExecContext(ctx, `DELETE FROM forecast_model_states - WHERE id NOT IN (SELECT state_id FROM forecast_issue_model_states)`) - return err -} - -func enforceForecastModelStateBudget(ctx context.Context, tx *sql.Tx) error { - for { - var total int64 - if err := tx.QueryRowContext(ctx, "SELECT COALESCE(SUM(length(payload)),0) FROM forecast_model_states").Scan(&total); err != nil { - return err - } - if total <= MaxForecastModelStateBytes { - return nil - } - var oldest string - if err := tx.QueryRowContext(ctx, "SELECT id FROM forecast_issues ORDER BY issued_at_ms,id LIMIT 1").Scan(&oldest); err != nil { - return err - } - if _, err := tx.ExecContext(ctx, "DELETE FROM forecast_issues WHERE id=?", oldest); err != nil { - return err - } - if err := cleanForecastModelStateRefs(ctx, tx); err != nil { - return err - } - } -} - // SaveForecastIssue is append-only within a bounded retention window. // An identical retry is allowed; reusing an ID for a changed issue is an error. func (s *Store) SaveForecastIssue(ctx context.Context, issue forecasting.Issue) error { @@ -302,25 +270,12 @@ func (s *Store) SaveForecastIssue(ctx context.Context, issue forecasting.Issue) return fmt.Errorf("store forecast model state reference: %w", err) } } - if _, err = tx.ExecContext(ctx, "DELETE FROM forecast_issues WHERE issued_at_ms < ?", now-ForecastIssueRetention.Milliseconds()); err != nil { - return fmt.Errorf("prune expired forecast issues: %w", err) - } - if _, err = tx.ExecContext(ctx, `DELETE FROM forecast_issues WHERE id IN ( - SELECT id FROM (SELECT id,ROW_NUMBER() OVER (ORDER BY issued_at_ms DESC,id DESC) AS n, - SUM(length(payload)) OVER (ORDER BY issued_at_ms DESC,id DESC) AS bytes FROM forecast_issues) - WHERE n>? OR bytes>?)`, MaxForecastIssues, MaxForecastArchiveBytes); err != nil { - return fmt.Errorf("prune forecast issue budget: %w", err) - } - if err = cleanForecastModelStateRefs(ctx, tx); err != nil { - return fmt.Errorf("clean forecast model states: %w", err) - } - if err = enforceForecastModelStateBudget(ctx, tx); err != nil { - return fmt.Errorf("enforce forecast model state budget: %w", err) - } if err = tx.Commit(); err != nil { return fmt.Errorf("commit forecast issue: %w", err) } - return nil + // The immutable record is committed before retention starts. If maintenance + // fails, an identical retry completes it without changing the issued plan. + return s.pruneForecastIssues(ctx, now) } func (s *Store) SaveForecastObservation(ctx context.Context, observation forecasting.Observation) error { @@ -358,16 +313,10 @@ func (s *Store) SaveForecastObservation(ctx context.Context, observation forecas return errors.New("forecast observation is immutable") } } - if _, err = tx.ExecContext(ctx, "DELETE FROM forecast_observations WHERE available_at_ms < ?", now-ForecastIssueRetention.Milliseconds()); err != nil { - return err - } - if _, err = tx.ExecContext(ctx, `DELETE FROM forecast_observations WHERE rowid IN ( - SELECT rowid FROM (SELECT rowid,ROW_NUMBER() OVER (ORDER BY available_at_ms DESC,start_ms DESC,end_ms DESC,config_version DESC) AS n, - SUM(length(payload)) OVER (ORDER BY available_at_ms DESC,start_ms DESC,end_ms DESC,config_version DESC) AS bytes FROM forecast_observations) - WHERE n>? OR bytes>?)`, MaxForecastObservations, MaxForecastObservationBytes); err != nil { + if err := tx.Commit(); err != nil { return err } - return tx.Commit() + return s.pruneForecastRecords(ctx, "forecast_observations", "available_at_ms DESC,start_ms DESC,end_ms DESC,config_version DESC", "available_at_ms", now, MaxForecastObservations, MaxForecastObservationBytes) } func decodeForecastIssue(data []byte) (forecasting.Issue, int, error) { @@ -546,25 +495,36 @@ func (s *Store) SaveForecastErrors(ctx context.Context, samples []forecasting.Er if now <= 0 { return errors.New("invalid forecast score time") } - tx, err := s.db.BeginTx(ctx, nil) - if err != nil { - return err + // Validate and encode the full page before taking SQLite's writer lock. + if len(samples) > MaxForecastErrors { + return errors.New("forecast score page exceeds row limit") } - defer tx.Rollback() - for _, sample := range samples { - if err = sample.Validate(); err != nil { + payloads := make([]string, len(samples)) + totalBytes := 0 + for i, sample := range samples { + if err := sample.Validate(); err != nil { return err } if sample.AvailableAtMS > now || sample.EndMS > now { return errors.New("future forecast score") } - data, marshalErr := json.Marshal(sample) - if marshalErr != nil { - return marshalErr + data, err := json.Marshal(sample) + if err != nil { + return err } - if len(data) > maxForecastRecordBytes { - return errors.New("forecast error exceeds payload limit") + totalBytes += len(data) + if len(data) > maxForecastRecordBytes || totalBytes > maxForecastSliceBytes { + return errors.New("forecast score page exceeds payload limit") } + payloads[i] = string(data) + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + for i, sample := range samples { + data := payloads[i] result, execErr := tx.ExecContext(ctx, `INSERT INTO forecast_errors(series,config_version,lead,start_ms,end_ms,origin_ms,issued_at_ms,issue_id,available_at_ms,payload) VALUES(?,?,?,?,?,?,?,?,?,?) ON CONFLICT(series,config_version,lead,start_ms,end_ms) DO UPDATE SET origin_ms=excluded.origin_ms,issued_at_ms=excluded.issued_at_ms,issue_id=excluded.issue_id, @@ -591,16 +551,10 @@ func (s *Store) SaveForecastErrors(ctx context.Context, samples []forecasting.Er } } } - if _, err = tx.ExecContext(ctx, "DELETE FROM forecast_errors WHERE end_ms < ?", now-ForecastIssueRetention.Milliseconds()); err != nil { - return err - } - if _, err = tx.ExecContext(ctx, `DELETE FROM forecast_errors WHERE rowid IN ( - SELECT rowid FROM (SELECT rowid,ROW_NUMBER() OVER (ORDER BY end_ms DESC,start_ms DESC,series DESC,lead DESC) AS n, - SUM(length(payload)) OVER (ORDER BY end_ms DESC,start_ms DESC,series DESC,lead DESC) AS bytes FROM forecast_errors) - WHERE n>? OR bytes>?)`, MaxForecastErrors, MaxForecastErrorBytes); err != nil { + if err := tx.Commit(); err != nil { return err } - return tx.Commit() + return s.pruneForecastRecords(ctx, "forecast_errors", "end_ms DESC,start_ms DESC,series DESC,lead DESC,config_version DESC", "end_ms", now, MaxForecastErrors, MaxForecastErrorBytes) } // LoadForecastErrors reads a bounded set of valid residuals. With hourly set, diff --git a/go/internal/state/forecast_prune.go b/go/internal/state/forecast_prune.go new file mode 100644 index 000000000..5aafb4932 --- /dev/null +++ b/go/internal/state/forecast_prune.go @@ -0,0 +1,124 @@ +package state + +import ( + "context" + "errors" + "fmt" + "strings" + "time" +) + +const forecastPruneBatch = 64 + +// Budget scans run in a read snapshot. Only the selected, bounded deletes +// acquire SQLite's writer lock. A concurrent insert invalidates the snapshot; +// retry the selection so it cannot delete from an outdated budget calculation. +func (s *Store) pruneForecastRows(ctx context.Context, table, selection string, args ...any) (int64, error) { + var total int64 + for { + n, err := s.tryPruneForecastRows(ctx, table, selection, args...) + if err != nil && !historyWriteBusy(err) && !(ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded)) { + return total, err + } + if err == nil { + total += n + if n == 0 { + return total, nil + } + } + if err := pauseMaintenance(ctx); err != nil { + return total, err + } + } +} + +// Table and selection are static SQL from this package, never client input. +func (s *Store) tryPruneForecastRows(parent context.Context, table, selection string, args ...any) (int64, error) { + txCtx, cancelTx := context.WithCancel(parent) + defer cancelTx() + tx, err := s.db.BeginTx(txCtx, nil) + if err != nil { + return 0, err + } + defer tx.Rollback() + readCtx, cancelRead := context.WithTimeout(parent, 3*time.Second) + rows, err := tx.QueryContext(readCtx, selection, args...) + if err != nil { + cancelRead() + return 0, err + } + ids := make([]any, 0, forecastPruneBatch) + for rows.Next() { + var id int64 + if err = rows.Scan(&id); err != nil { + break + } + ids = append(ids, id) + if len(ids) > forecastPruneBatch { + err = errors.New("forecast prune exceeds batch limit") + break + } + } + err = errors.Join(err, rows.Err(), rows.Close()) + cancelRead() + if err != nil { + return 0, err + } + if len(ids) == 0 { + return 0, nil + } + writeCtx, cancelWrite := context.WithTimeout(parent, time.Second) + defer cancelWrite() + stop := context.AfterFunc(writeCtx, cancelTx) + defer stop() + res, err := tx.ExecContext(writeCtx, "DELETE FROM "+table+" WHERE rowid IN ("+strings.TrimSuffix(strings.Repeat("?,", len(ids)), ",")+")", ids...) + if err == nil { + err = tx.Commit() + } + if err != nil { + if writeCtx.Err() != nil { + err = errors.Join(err, writeCtx.Err()) + } + return 0, err + } + return res.RowsAffected() +} + +func (s *Store) pruneForecastRecords(ctx context.Context, table, order, expiry string, now int64, maxRows, maxBytes int) error { + _, err := s.pruneForecastRows(ctx, table, fmt.Sprintf(`SELECT rowid FROM ( + SELECT rowid,%s AS expires,ROW_NUMBER() OVER (ORDER BY %s) AS n, + SUM(length(payload)) OVER (ORDER BY %s) AS bytes FROM %s) + WHERE expires? OR bytes>? LIMIT %d`, expiry, order, order, table, forecastPruneBatch), + now-ForecastIssueRetention.Milliseconds(), maxRows, maxBytes) + return err +} + +func (s *Store) pruneForecastIssues(ctx context.Context, now int64) error { + if err := s.pruneForecastRecords(ctx, "forecast_issues", "issued_at_ms DESC,id DESC", "issued_at_ms", now, MaxForecastIssues, MaxForecastArchiveBytes); err != nil { + return err + } + for { + if _, err := s.pruneForecastRows(ctx, "forecast_issue_model_states", `SELECT rowid FROM forecast_issue_model_states + WHERE issue_id NOT IN (SELECT id FROM forecast_issues) LIMIT 64`); err != nil { + return err + } + if _, err := s.pruneForecastRows(ctx, "forecast_model_states", `SELECT rowid FROM forecast_model_states + WHERE id NOT IN (SELECT state_id FROM forecast_issue_model_states) LIMIT 64`); err != nil { + return err + } + // Evict one oldest issue at a time, then release its unshared states. + // The sum and reference scans never run after acquiring the writer lock. + n, err := s.tryPruneForecastRows(ctx, "forecast_issues", `SELECT rowid FROM forecast_issues + WHERE (SELECT COALESCE(SUM(length(payload)),0) FROM forecast_model_states)>? + ORDER BY issued_at_ms,id LIMIT 1`, MaxForecastModelStateBytes) + if err == nil && n == 0 { + return nil + } + if err != nil && !historyWriteBusy(err) && !(ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded)) { + return err + } + if err := pauseMaintenance(ctx); err != nil { + return err + } + } +} diff --git a/go/internal/state/history_maintenance.go b/go/internal/state/history_maintenance.go new file mode 100644 index 000000000..c13492ca2 --- /dev/null +++ b/go/internal/state/history_maintenance.go @@ -0,0 +1,84 @@ +package state + +import ( + "context" + "errors" + "fmt" + "time" +) + +type HistoryMaintenanceStatus struct { + State string `json:"state"` + Phase string `json:"phase,omitempty"` + StartedMS int64 `json:"started_ms,omitempty"` + LastSuccessMS int64 `json:"last_success_ms,omitempty"` + LastError string `json:"last_error,omitempty"` + LastFailureMS int64 `json:"last_failure_ms,omitempty"` + LastFailureError string `json:"last_failure_error,omitempty"` + Runs uint64 `json:"runs"` + Failures uint64 `json:"failures"` +} + +func (s *Store) HistoryMaintenanceStatus() HistoryMaintenanceStatus { + s.maintenanceStatusMu.Lock() + defer s.maintenanceStatusMu.Unlock() + status := s.maintenanceStatus + if status.State == "" { + status.State = "not_started" + } + return status +} + +// MaintainHistory runs the existing retention work and records its result. +// The caller serializes it with backup/restore. Each stage must yield to live +// writes; the overall deadline also bounds a permanently stalled archive. +func (s *Store) MaintainHistory(parent context.Context, coldDir string, days int, now time.Time) error { + ctx, cancel := context.WithTimeout(parent, 2*time.Hour) + defer cancel() + s.maintenanceStatusMu.Lock() + s.maintenanceStatus.State, s.maintenanceStatus.StartedMS = "running", time.Now().UnixMilli() + s.maintenanceStatus.Runs++ + s.maintenanceStatusMu.Unlock() + var failures []error + for _, stage := range []struct { + name string + run func() error + }{ + {"dashboard_rollup", func() error { return s.Prune(ctx) }}, + {"energy_rollup", func() error { _, _, err := s.PruneEnergyLedger(ctx, now); return err }}, + {"diagnostic_archive", func() error { _, _, err := s.RolloffDiagnosticsToParquet(ctx, coldDir); return err }}, + {"sample_archive", func() error { return s.PruneHistorySamples(ctx, days, now) }}, + {"diagnostic_retention", func() error { _, err := PruneDiagnosticsParquet(coldDir, days, now); return err }}, + } { + s.maintenanceStatusMu.Lock() + s.maintenanceStatus.Phase = stage.name + s.maintenanceStatusMu.Unlock() + if err := ctx.Err(); err != nil { + failures = append(failures, err) + break + } + if err := stage.run(); err != nil { + failures = append(failures, fmt.Errorf("%s: %w", stage.name, err)) + // A later archive stage may run for a long time. Report this + // failure now, rather than hiding it until the full cycle ends. + s.maintenanceStatusMu.Lock() + s.maintenanceStatus.LastError = errors.Join(failures...).Error() + s.maintenanceStatus.LastFailureMS = time.Now().UnixMilli() + s.maintenanceStatus.LastFailureError = s.maintenanceStatus.LastError + s.maintenanceStatusMu.Unlock() + } + } + err := errors.Join(failures...) + s.maintenanceStatusMu.Lock() + defer s.maintenanceStatusMu.Unlock() + s.maintenanceStatus.Phase = "" + if err == nil { + s.maintenanceStatus.State, s.maintenanceStatus.LastError = "complete", "" + s.maintenanceStatus.LastSuccessMS = time.Now().UnixMilli() + } else { + s.maintenanceStatus.State, s.maintenanceStatus.LastError = "failed", err.Error() + s.maintenanceStatus.Failures++ + s.maintenanceStatus.LastFailureMS, s.maintenanceStatus.LastFailureError = time.Now().UnixMilli(), err.Error() + } + return err +} diff --git a/go/internal/state/history_rollup.go b/go/internal/state/history_rollup.go new file mode 100644 index 000000000..6f2287913 --- /dev/null +++ b/go/internal/state/history_rollup.go @@ -0,0 +1,158 @@ +package state + +import ( + "context" + "database/sql" + "errors" + "fmt" + "strings" +) + +// Unlike prepared hourly archives, rollups return a short write deadline so +// their caller can reduce the batch. Busy snapshots still reread their input. +func (s *Store) rollupTransaction(ctx context.Context, prepare, write func(context.Context, *sql.Tx) error) error { + conflicts := 0 + for { + if err := ctx.Err(); err != nil { + return err + } + if s.HistoryWriterStatus().Pending < historyCommitMaxTicks/2 { + var err error + if conflicts < 3 { + err = s.tryArchiveBatch(ctx, prepare, write) + } else { + // Constant live commits can invalidate every read snapshot. One + // attempt may read under the lock, within the same one-second + // total budget. The caller shrinks the batch if it cannot fit. + err = s.tryArchiveBatch(ctx, nil, func(ctx context.Context, tx *sql.Tx) error { + // Reserve SQLite's writer as well: compatibility writers and + // another process do not share the Go mutex. No rows change. + if _, err := tx.ExecContext(ctx, `UPDATE history_receipts SET batch_id=batch_id WHERE 0`); err != nil { + return err + } + if err := prepare(ctx, tx); err != nil { + return err + } + return write(ctx, tx) + }) + } + conflicts++ + if !historyWriteBusy(err) { + return err + } + } + if err := pauseMaintenance(ctx); err != nil { + return err + } + } +} + +// Read the full source buckets before acquiring the live writer mutex. The +// same SQLite snapshot binds the averages and newest JSON to the deleted rows. +func (s *Store) pruneChunk(ctx context.Context, src, dst string, fromMS, toMS, bucketMS int64) (int64, error) { + var points []HistoryPoint + var deleted int64 + err := s.rollupTransaction(ctx, func(ctx context.Context, tx *sql.Tx) error { + points = nil + q := fmt.Sprintf(`SELECT (ts_ms / %d) * %d + %d, + AVG(grid_w),AVG(pv_w),AVG(bat_w),AVG(load_w),AVG(bat_soc),json,MAX(ts_ms) + FROM %s WHERE ts_ms>=? AND ts_ms=? AND ts_ms=? AND bucket_start_ms0 LIMIT 64`) + if err == nil && n != 130 { + err = fmt.Errorf("deleted %d rows, want 130", n) + } + done <- err + }() + select { + case <-reading: + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + goal := make(chan error, 1) + go func() { goal <- s.SaveConfigValues(map[string]string{"test-goal": "80"}) }() + select { + case err := <-goal: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + released.Do(func() { close(release) }) + <-done + <-goal + t.Fatal("forecast budget scan held the state writer lock") + } + // The config commit invalidates the pruning snapshot. The retry must + // reselect, then finish all three delete batches without losing counts. + released.Do(func() { close(release) }) + if err := <-done; err != nil { + t.Fatal(err) + } + if value, ok := s.LoadConfig("test-goal"); !ok || value != "80" { + t.Fatal("durable goal changed") + } +} + +func TestArchiveSingleRowDeadlineRetainsCurrentBatch(t *testing.T) { + s := freshStore(t) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := s.RecordSamples([]Sample{{Driver: "test", Metric: "power", TsMs: 1, Value: 42}}); err != nil { + t.Fatal(err) + } + d, err := s.driverID("test") + if err != nil { + t.Fatal(err) + } + m, err := s.metricID("power", "") + if err != nil { + t.Fatal(err) + } + // One long reader outlasts the archive's write budget. The verified row + // must stay in this job and retry after the reader releases the view. + s.archiveViewMu.RLock() + locked := true + defer func() { + if locked { + s.archiveViewMu.RUnlock() + } + }() + done := make(chan error, 1) + go func() { + n, err := s.pruneArchivedSamples(ctx, []resolvedSample{{dID: d, mID: m, ts: 1, v: 42}}) + if err == nil && n != 1 { + err = sql.ErrNoRows + } + done <- err + }() + select { + case err := <-done: + t.Fatalf("archive abandoned its verified row after one write deadline: %v", err) + case <-time.After(archiveWriteTimeout + 150*time.Millisecond): + } + s.archiveViewMu.RUnlock() + locked = false + if err := <-done; err != nil { + t.Fatal(err) + } +} diff --git a/go/internal/state/parquet_stream.go b/go/internal/state/parquet_stream.go index 2fffaafbf..d7075b83f 100644 --- a/go/internal/state/parquet_stream.go +++ b/go/internal/state/parquet_stream.go @@ -368,8 +368,13 @@ func (s *Store) pruneArchivedSamples(ctx context.Context, batch []resolvedSample return nil }) if err != nil { - if ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded) && n > 1 { + if ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded) { + // Even one row can meet a busy reader or a slow sync. Retain + // this verified prefix instead of staging the whole day again. limit = max(1, n/2) + if err := pauseMaintenance(ctx); err != nil { + return deleted, err + } continue } return deleted, err diff --git a/go/internal/state/store.go b/go/internal/state/store.go index 197398ff3..156fc45bb 100644 --- a/go/internal/state/store.go +++ b/go/internal/state/store.go @@ -83,6 +83,9 @@ type Store struct { seriesHourWG sync.WaitGroup seriesHourMu sync.Mutex seriesHourCancel context.CancelFunc + + maintenanceStatusMu sync.Mutex + maintenanceStatus HistoryMaintenanceStatus } // Open initializes (or creates) the precious state.db at path plus the @@ -1786,6 +1789,7 @@ func (s *Store) Prune(ctx context.Context) error { func (s *Store) pruneTier(ctx context.Context, src, dst string, cutoffMs, bucketMs int64) (aged int64, chunks int, err error) { // Only age complete buckets: align the cutoff down to a bucket boundary. cutoffMs = (cutoffMs / bucketMs) * bucketMs + spanMS := pruneChunkSpanMS for { if err := ctx.Err(); err != nil { @@ -1801,7 +1805,7 @@ func (s *Store) pruneTier(ctx context.Context, src, dst string, cutoffMs, bucket } // Chunk upper bound: at most pruneChunkSpanMS of rows, never past the // cutoff, always on a bucket boundary. - chunkEnd := minTs.Int64 + pruneChunkSpanMS + chunkEnd := minTs.Int64 + spanMS if chunkEnd > cutoffMs { chunkEnd = cutoffMs } @@ -1815,6 +1819,13 @@ func (s *Store) pruneTier(ctx context.Context, src, dst string, cutoffMs, bucket n, err := s.pruneChunk(ctx, src, dst, minTs.Int64, chunkEnd, bucketMs) if err != nil { + if ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded) { + spanMS = max(bucketMs, (chunkEnd-minTs.Int64)/2/bucketMs*bucketMs) + if err := pauseMaintenance(ctx); err != nil { + return aged, chunks, err + } + continue + } return aged, chunks, err } aged += n @@ -1838,42 +1849,6 @@ func (s *Store) pruneTier(ctx context.Context, src, dst string, cutoffMs, bucket // tests can shrink it. var pruneChunkPause = 250 * time.Millisecond -// pruneChunk aggregates+deletes src rows in [fromMs, toMs) in one short -// transaction. -func (s *Store) pruneChunk(ctx context.Context, src, dst string, fromMs, toMs, bucketMs int64) (int64, error) { - s.historyWriteMu.Lock() - defer s.historyWriteMu.Unlock() - tx, err := s.history.BeginTx(ctx, nil) - if err != nil { - return 0, err - } - defer tx.Rollback() - - // Bare-column rule: exactly one MAX() aggregate in the grouped inner - // query makes the un-aggregated json column come from that newest row. - q := fmt.Sprintf(` - INSERT OR REPLACE INTO %s (ts_ms, grid_w, pv_w, bat_w, load_w, bat_soc, json) - SELECT b_ts, a_grid, a_pv, a_bat, a_load, a_soc, json FROM ( - SELECT (ts_ms / %d) * %d + %d AS b_ts, - AVG(grid_w) AS a_grid, AVG(pv_w) AS a_pv, AVG(bat_w) AS a_bat, - AVG(load_w) AS a_load, AVG(bat_soc) AS a_soc, - json AS json, MAX(ts_ms) AS newest - FROM %s - WHERE ts_ms >= ? AND ts_ms < ? - GROUP BY ts_ms / %d - )`, dst, bucketMs, bucketMs, bucketMs/2, src, bucketMs) - if _, err := tx.ExecContext(ctx, q, fromMs, toMs); err != nil { - return 0, fmt.Errorf("aggregate: %w", err) - } - res, err := tx.ExecContext(ctx, - `DELETE FROM `+src+` WHERE ts_ms >= ? AND ts_ms < ?`, fromMs, toMs) - if err != nil { - return 0, fmt.Errorf("delete: %w", err) - } - n, _ := res.RowsAffected() - return n, tx.Commit() -} - // ---- Prices ---- // PricePoint is one time-slot's spot price row. Slot length varies by source: From b293a2ce5ce447bf73df59e8bb8d7c4384f642ff Mon Sep 17 00:00:00 2001 From: Fredrik Ahlgren Date: Thu, 17 Sep 2026 09:46:34 +0200 Subject: [PATCH 2/2] fix(forecast): recover observation capacity and health Signed-off-by: Fredrik Ahlgren --- go/cmd/ftw/forecast_observation_queue.go | 20 +++++++++++ go/cmd/ftw/forecast_observation_queue_test.go | 35 +++++++++++++++++++ go/cmd/ftw/forecast_tracking.go | 6 ++-- 3 files changed, 59 insertions(+), 2 deletions(-) diff --git a/go/cmd/ftw/forecast_observation_queue.go b/go/cmd/ftw/forecast_observation_queue.go index 8eb1982e7..11d1430b1 100644 --- a/go/cmd/ftw/forecast_observation_queue.go +++ b/go/cmd/ftw/forecast_observation_queue.go @@ -43,10 +43,24 @@ func (f *forecastTracker) runObservations(ctx context.Context) { } } +// Capture new intervals before disk work, but let a recovered backlog free +// capacity before rejecting the captured evidence. The same deadline bounds +// both drains; recovery never postpones the next observation indefinitely. +func (f *forecastTracker) acceptObservations(ctx context.Context, jobs []forecastObservationJob, flush func(context.Context)) { + if len(f.pendingObservations)+len(jobs) > maxPendingForecastObservations { + flush(ctx) + } + for _, job := range jobs { + f.enqueueObservation(job) + } + flush(ctx) +} + func (f *forecastTracker) enqueueObservation(job forecastObservationJob) { if len(f.pendingObservations) == maxPendingForecastObservations { f.mu.Lock() f.observationOverflow = true + f.observationDrops++ f.mu.Unlock() slog.Error("forecast archive: observation queue full", "start_ms", job.observation.StartMS) return @@ -100,6 +114,12 @@ func (f *forecastTracker) flushObservationsWith(ctx context.Context, f.pendingObservations = f.pendingObservations[1:] f.requestScoring() } + f.mu.Lock() + if f.observationOverflow { + slog.Info("forecast archive: observation queue recovered", "dropped_intervals", f.observationDrops) + } + f.observationOverflow = false + f.mu.Unlock() } func (f *forecastTracker) updateObservation(ctx context.Context, job forecastObservationJob) error { diff --git a/go/cmd/ftw/forecast_observation_queue_test.go b/go/cmd/ftw/forecast_observation_queue_test.go index e6422bd3a..e22551de4 100644 --- a/go/cmd/ftw/forecast_observation_queue_test.go +++ b/go/cmd/ftw/forecast_observation_queue_test.go @@ -204,3 +204,38 @@ func TestForecastObservationOldRevisionCannotReplaceCurrentModel(t *testing.T) { t.Fatal("current observation did not update model") } } + +func TestForecastObservationQueueRecoversBeforeAcceptingNewInterval(t *testing.T) { + f := trackerFixture(time.Now()) + old := observationJobFixture() + for i := 0; i < maxPendingForecastObservations; i++ { + f.enqueueObservation(old) + } + newJob := old + newJob.observation.StartMS++ + var saved []int64 + drain := func(ctx context.Context) { + f.flushObservationsWith(ctx, func(_ context.Context, o forecasting.Observation) error { + saved = append(saved, o.StartMS) + return nil + }, func(context.Context, forecastObservationJob) error { return nil }) + } + f.acceptObservations(context.Background(), []forecastObservationJob{newJob}, drain) + if len(saved) != maxPendingForecastObservations+1 || saved[len(saved)-1] != newJob.observation.StartMS || f.observationDrops != 0 { + t.Fatalf("recovery discarded new evidence: saved=%v drops=%d", saved, f.observationDrops) + } +} + +func TestForecastObservationQueueOverflowHealthRecovers(t *testing.T) { + f := trackerFixture(time.Now()) + for i := 0; i <= maxPendingForecastObservations; i++ { + f.enqueueObservation(observationJobFixture()) + } + if !f.observationOverflow || f.observationDrops != 1 { + t.Fatal("overflow was not recorded") + } + f.flushObservationsWith(context.Background(), func(context.Context, forecasting.Observation) error { return nil }, func(context.Context, forecastObservationJob) error { return nil }) + if f.observationOverflow || len(f.pendingObservations) != 0 || f.observationDrops != 1 { + t.Fatal("current health did not recover or historical loss was erased") + } +} diff --git a/go/cmd/ftw/forecast_tracking.go b/go/cmd/ftw/forecast_tracking.go index 30cab6e0a..b2e4bb208 100644 --- a/go/cmd/ftw/forecast_tracking.go +++ b/go/cmd/ftw/forecast_tracking.go @@ -82,6 +82,7 @@ type forecastTracker struct { pendingObservations []forecastObservationJob // owned by runObservations observationError string // guarded by mu observationOverflow bool + observationDrops uint64 observationCheckedMS int64 observationPVValid bool observationLoadValid bool @@ -264,6 +265,7 @@ func (f *forecastTracker) observe(ctx context.Context) { f.observationCheckedMS = now.UnixMilli() f.observationPVValid, f.observationLoadValid = r.PVValid, r.Valid f.mu.Unlock() + var jobs []forecastObservationJob for _, o := range f.observationIntervals(r, site, now) { job := forecastObservationJob{observation: o, site: site, away: f.away != nil && f.away(time.UnixMilli(o.StartMS))} rows, _ := f.store.LoadForecasts(o.StartMS-time.Hour.Milliseconds(), o.EndMS) @@ -271,12 +273,12 @@ func (f *forecastTracker) observe(ctx context.Context) { frozen := *weather job.weather = &frozen } - f.enqueueObservation(job) + jobs = append(jobs, job) } // A slow disk must not make draining a backlog postpone measurement capture. flushCtx, cancel := context.WithTimeout(ctx, 3*time.Second) defer cancel() - f.flushObservations(flushCtx) + f.acceptObservations(flushCtx, jobs, f.flushObservations) } // observationIntervals keeps the two measurement claims independent. Unknown