diff --git a/.changeset/clear-history-import-progress.md b/.changeset/clear-history-import-progress.md new file mode 100644 index 000000000..c6664d813 --- /dev/null +++ b/.changeset/clear-history-import-progress.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Show startup and background history import progress, keep incomplete history visible, and resume the update view when another client starts the work. diff --git a/.changeset/safe-background-history-import.md b/.changeset/safe-background-history-import.md new file mode 100644 index 000000000..f046e0bc9 --- /dev/null +++ b/.changeset/safe-background-history-import.md @@ -0,0 +1,9 @@ +--- +"ftw": patch +--- + +Start live collection after importing the catalog and energy accounting, then move large sample archives into DuckDB in the background. Resume verified progress after an interruption, preserve live writes, and report incomplete history until every source is checked. Release native database buffers only between active SQL connections. + +Give Core startup more time and leave its image and data in place if readiness fails. Preserve the previous image ID before replacement. Check the updater before opening data, and stop with a clear message when an older installation needs to update its updater first. + +Keep catalog ID allocation safe after an abrupt process exit, and defer raw-history retention until import finishes. diff --git a/docs/self-update.md b/docs/self-update.md index 67c14e6e3..a6634d67c 100644 --- a/docs/self-update.md +++ b/docs/self-update.md @@ -122,16 +122,83 @@ newer updater, update Core and updater together using the paired commands below. A normal Core update also asks the updater to replace itself with the same tag after Core passes its health check. -For manual Core + updater operation, first set `FTW_IMAGE_TAG` and -`FTW_UPDATER_IMAGE_TAG` in the project's `.env` to the same published immutable -tag. Run these commands with the Core service name from your Compose file: +For manual updates, install the updater first while the existing Core still +runs. Set `FTW_UPDATER_IMAGE_TAG` in the project's `.env` to the published +immutable release tag. Keep the current `FTW_IMAGE_TAG` until the updater is +installed. Use your installation's Compose project and override files: ```bash cd ~/ftw -docker compose pull ftw ftw-updater -docker compose up -d --no-deps ftw ftw-updater +docker compose pull ftw-updater +docker compose up -d --no-deps ftw-updater +docker compose images ftw-updater +docker compose logs --tail 20 ftw-updater ``` +Check that the running updater uses the intended release image and has started +its socket listener. Then use Update Center to update Core. A manual Core +replacement must first create and verify a full backup; pin `FTW_IMAGE_TAG` to +the same tag, then pull and recreate only the Core service. + +### First DuckDB upgrade + +The updater shipped before the fix for [#1164](https://github.com/srcfl/ftw/issues/1164) +waits only 30 minutes and can revert the Core image without its matching data. +It replaces itself only after Core becomes ready. **Installing a new Core does +not fix the old updater before that first upgrade.** Install an updater release +that contains the fix using the updater-only steps above before starting the +DuckDB upgrade. `v3.2.0-beta.1` does not contain this fix. + +When `FTW_SELFUPDATE_ENABLED=1`, the new Core checks `GET /capabilities` on the +updater Unix socket before loading config or opening state. It requires protocol +1 and `preserve_core_on_readiness_failure: true`. It waits up to 30 seconds for +the socket to start. An old, unknown or unavailable updater makes Core exit with +`update ftw-updater first`, before it changes any data. The old updater can then +revert that refused image safely. This guard cannot protect a Core already +running `v3.2.0-beta.1`; use the stopped-updater procedure below for that case. +Native installations with self-update disabled do not need an updater. + +Core first imports the catalog, site history and energy accounting, including +the latest counters. It then starts live collection and serves its API while +the large SQLite sample table and older Parquet files import in the background. +The startup page and history status show progress. Historical coverage stays +incomplete until all sources pass verification; a missing sample during this +phase does not mean zero consumption. Unknown time bounds apply to the whole +historical view. New writes go to DuckDB throughout the background import. + +An interrupted import resumes from committed progress on the next Core start. +An import error leaves live collection running and reports incomplete history. +Keep the original sources and the verified pre-update full backup. Core rejects +a new full-history export until import finishes. Fix the reported source or +storage problem before restarting Core to resume the import. + +Raw-history retention waits until the full import passes. Site-history rollups +remain available after the initial seed. Import scratch directories belong to +the importer; it removes leftovers on resume while keeping original sources. + +The fixed updater allows six hours for Core startup. This is a waiting budget, +not a promise that every migration will finish within it. Core still needs a +working `/api/status` before the updater reports success. A failed readiness +check leaves the new container and data in place and reports failure. To return +to an older version, stop Core and restore a verified full backup with its +matching image; changing only the image can lose history. + +### An updater stopped during migration + +Leave the running Core alone until its import and API readiness have been +checked. Save the updater's existing `state.json`, the pre-update full backup, +snapshot directory and verified previous Core image identity before starting +the updater again. Older updaters hold the current job's previous image ID only +in memory until success; a saved `previous_image_id` may belong to an earlier +update. Use the backup inventory and pre-update Docker evidence to resolve it. + +After readiness, start or replace **only** `ftw-updater` with `--no-deps` using +the fixed release. Its startup recovery marks an interrupted update as failed; +it does not restart Core or replay an image rollback. This records the lost +supervision honestly even if Core has since finished. Keep that record and +check Core's installed image and history status separately. Do not repeat the +Core update or overwrite the job as `done` merely to clear its failed status. + Manage Optimizer and Drivers independently in Update Center. A blanket `docker compose pull` is intentionally not the documented upgrade procedure. diff --git a/go/cmd/ftw-updater/main.go b/go/cmd/ftw-updater/main.go index 209973c22..90cb717f5 100644 --- a/go/cmd/ftw-updater/main.go +++ b/go/cmd/ftw-updater/main.go @@ -35,6 +35,8 @@ import ( "time" "gopkg.in/yaml.v3" + + "github.com/srcfl/ftw/go/internal/updateipc" ) const ( @@ -111,13 +113,9 @@ type server struct { // pins to the requested version. nil/empty means "inherit only". runner func(ctx context.Context, env []string, args ...string) error // imageID captures the image backing the running service before an update. - // imageRef captures the exact image reference used to create it so an - // automatic rollback can restore a beta's runtime identity as well as its - // bytes. - // healthCheck waits for the recreated service to become healthy. Both are - // injectable so the rollback path is testable without Docker. + // healthCheck waits for the recreated service to become ready. + // These hooks let tests exercise recovery without Docker. imageID func(ctx context.Context, service string) (string, error) - imageRef func(ctx context.Context, service string) (string, error) containerID func(ctx context.Context, service string) (string, error) healthCheck func(ctx context.Context, service string) error // selfReplace brings the updater sidecar to the release Core just moved to. @@ -293,7 +291,6 @@ func main() { } return } - srv.imageRef = srv.currentServiceImageRef srv.containerID = srv.serviceContainerID srv.healthCheck = srv.waitForServiceHealth srv.selfReplace = func(target string) error { @@ -317,6 +314,7 @@ func main() { mux := http.NewServeMux() mux.HandleFunc("POST /update", srv.handleUpdate) mux.HandleFunc("GET /status", srv.handleStatus) + mux.HandleFunc("GET "+updateipc.CapabilitiesPath, updateipc.ServeCapabilities) // Remove a stale socket — common pattern; the listener would EADDRINUSE otherwise. _ = os.Remove(*socket) @@ -546,26 +544,21 @@ func (s *server) runComponentJob(action, target, component string, startedAt tim } } - // Capture the current immutable image ID before pulling. Docker retains the - // old image object after the tag moves, which lets us retag and recreate it - // if the new container never becomes healthy. - var previousImageID, previousImageTag string + // Keep the previous image in durable job history before replacing Core. + // Starting it again also requires its matching full backup: the new Core + // may have changed the data before its API becomes ready. + var previousImageID string if action == "update" && s.imageID != nil { inspectCtx, cancelInspect := context.WithTimeout(context.Background(), 30*time.Second) var err error previousImageID, err = s.imageID(inspectCtx, spec.service) - if err == nil && s.imageRef != nil { - if previousRef, refErr := s.imageRef(inspectCtx, spec.service); refErr == nil { - previousImageTag, _ = imageTagFromReference(previousRef) - } else { - slog.Warn("cannot capture current image tag; rollback will use a synthetic tag", "service", spec.service, "err", refErr) - } - } cancelInspect() if err != nil { s.writeState(State{State: "failed", Action: action, Component: component, Target: target, StartedAt: now, UpdatedAt: time.Now(), Message: "cannot capture current image for rollback: " + err.Error()}) return } + pullState.PreviousImageID = previousImageID + s.writeState(pullState) } if !s.skipPull { @@ -633,15 +626,10 @@ func (s *server) runComponentJob(action, target, component string, startedAt tim }) cancelHealth() if healthErr != nil { - if action == "update" && previousImageID != "" { - if rollbackErr := s.restorePreviousComponentImageWithTag(previousImageID, previousImageTag, spec); rollbackErr == nil { - s.writeState(State{State: "failed", Action: action, Component: component, Target: target, StartedAt: now, UpdatedAt: time.Now(), Message: "new image failed health check; previous image restored: " + healthErr.Error()}) - return - } else { - healthErr = fmt.Errorf("%v; automatic image rollback failed: %w", healthErr, rollbackErr) - } - } - s.writeState(State{State: "failed", Action: action, Component: component, Target: target, StartedAt: now, UpdatedAt: time.Now(), Message: "health check failed: " + healthErr.Error()}) + checkState.State = "failed" + checkState.UpdatedAt = time.Now() + checkState.Message = "Core readiness failed: " + healthErr.Error() + ". Core and data were left in place; migration may still be running. Check Core status before retrying. To return to an older release, stop Core and restore a verified full backup with its matching image." + s.writeState(checkState) return } } @@ -678,56 +666,6 @@ func (s *server) componentSpec(component string) (componentSpec, error) { } } -func (s *server) restorePreviousImage(imageID string) error { - spec, err := s.componentSpec("core") - if err != nil { - return err - } - return s.restorePreviousComponentImage(imageID, spec) -} - -func (s *server) restorePreviousComponentImage(imageID string, spec componentSpec) error { - return s.restorePreviousComponentImageWithTag(imageID, "", spec) -} - -func (s *server) restorePreviousComponentImageWithTag(imageID, previousTag string, spec componentSpec) error { - image, ok, err := serviceImageFromComposeFiles(s.composeFiles(), spec.service) - if err != nil { - return err - } - if !ok { - return fmt.Errorf("service %q has no image", spec.service) - } - repository, ok := composeImageRepositoryForTag(image, spec.tagVariable) - if !ok { - return fmt.Errorf("service %q image %q does not reference %s", spec.service, image, spec.tagVariable) - } - rollbackTag := previousTag - if rollbackTag == "" { - rollbackTag = fmt.Sprintf("ftw-rollback-%d", time.Now().Unix()) - } - rollbackRef := repository + ":" + rollbackTag - timeout := 10 * time.Minute - if spec.name == "core" { - timeout = 35 * time.Minute - } - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() - if err := s.runner(ctx, nil, "image", "tag", imageID, rollbackRef); err != nil { - return fmt.Errorf("tag previous image: %w", err) - } - env := []string{spec.tagEnv + "=" + rollbackTag} - if err := s.runner(ctx, env, s.composeArgs("up", "-d", spec.service)...); err != nil { - return fmt.Errorf("recreate previous image: %w", err) - } - if s.healthCheck != nil { - if err := s.healthCheck(ctx, spec.service); err != nil { - return fmt.Errorf("previous image health check: %w", err) - } - } - return nil -} - // prepareUpdateImagePin makes old Compose layouts safe for immutable updates. // Older and developer installations may hard-code an image tag or omit the // container-side FTW_IMAGE_TAG mapping needed to report the selected release. @@ -1312,6 +1250,8 @@ func (s *server) recoverCrashedState() { return } st.Message = fmt.Sprintf("updater restarted during rollback; safety recovery unavailable: container=%v image=%v", containerErr, imageErr) + } else if st.Action == "update" { + st.Message = "Updater stopped before it confirmed Core readiness. Core and data were left in place; check Core status before retrying. Keep the pre-update full backup and matching image for recovery." } else if st.Message == "" { st.Message = "updater process restarted while in-flight" } @@ -1444,39 +1384,6 @@ func (s *server) currentServiceImageID(ctx context.Context, service string) (str return imageID, nil } -func (s *server) currentServiceImageRef(ctx context.Context, service string) (string, error) { - containerID, err := s.serviceContainerID(ctx, service) - if err != nil { - return "", err - } - out, err := dockerOutput(ctx, "inspect", "--format", "{{.Config.Image}}", containerID) - if err != nil { - return "", err - } - imageRef := strings.TrimSpace(out) - if imageRef == "" { - return "", errors.New("running container has no image reference") - } - return imageRef, nil -} - -func imageTagFromReference(imageRef string) (string, bool) { - imageRef = strings.TrimSpace(imageRef) - if imageRef == "" || strings.Contains(imageRef, "@") { - return "", false - } - lastSlash := strings.LastIndexByte(imageRef, '/') - lastColon := strings.LastIndexByte(imageRef, ':') - if lastColon <= lastSlash || lastColon == len(imageRef)-1 { - return "", false - } - tag := imageRef[lastColon+1:] - if !isImmutableImageTag(tag) { - return "", false - } - return tag, true -} - func (s *server) waitForServiceHealth(ctx context.Context, service string) error { containerID, err := s.serviceContainerID(ctx, service) if err != nil { @@ -1520,7 +1427,9 @@ func (s *server) waitForServiceHealth(ctx context.Context, service string) error func componentHealthTimeout(component string) time.Duration { if component == "core" { - return 30 * time.Minute + // First-start imports on microSD can take hours. The finite deadline + // reports failure without stopping Core or reverting its image. + return 6 * time.Hour } return 3 * time.Minute } diff --git a/go/cmd/ftw-updater/main_test.go b/go/cmd/ftw-updater/main_test.go index 286da3035..c3f68e878 100644 --- a/go/cmd/ftw-updater/main_test.go +++ b/go/cmd/ftw-updater/main_test.go @@ -647,7 +647,7 @@ func TestHandleUpdate_RestartKeepsEveryRunningImage(t *testing.T) { s, runner := newTestServer(t) writeCompose(t, s.composeFile, "services:\n ftw:\n image: "+tc.image+"\n") writeCompose(t, filepath.Join(filepath.Dir(s.composeFile), ".env"), "FTW_IMAGE_TAG=v2.0.0-beta.2\n") - s.imageRef = func(context.Context, string) (string, error) { + s.imageID = func(context.Context, string) (string, error) { t.Error("restart must not resolve an image") return "", errors.New("inspect unavailable") } @@ -991,25 +991,6 @@ func TestSelectMainServiceRejectsAmbiguousDataOwners(t *testing.T) { } } -func TestImageTagFromReferenceAcceptsOnlyImmutableReleaseTags(t *testing.T) { - for _, tc := range []struct { - ref string - want string - ok bool - }{ - {ref: "ghcr.io/srcfl/ftw:v2.0.0-beta.7", want: "v2.0.0-beta.7", ok: true}, - {ref: "ghcr.io/srcfl/ftw:v2.0.0", want: "v2.0.0", ok: true}, - {ref: "ghcr.io/srcfl/ftw:latest"}, - {ref: "ghcr.io/srcfl/ftw@sha256:deadbeef"}, - {ref: "ghcr.io/srcfl/ftw"}, - } { - got, ok := imageTagFromReference(tc.ref) - if got != tc.want || ok != tc.ok { - t.Errorf("imageTagFromReference(%q) = %q, %v; want %q, %v", tc.ref, got, ok, tc.want, tc.ok) - } - } -} - func TestRecoverCrashedState(t *testing.T) { s, _ := newTestServer(t) s.writeState(State{State: "pulling", UpdatedAt: time.Now()}) @@ -1062,3 +1043,44 @@ func TestContainerIDsFromDockerPS_IgnoresNoiseAndDedupes(t *testing.T) { t.Fatalf("clean id = %v", got) } } + +func TestUpdateReadinessFailureNeverRevertsImage(t *testing.T) { + for _, healthErr := range []error{context.DeadlineExceeded, errors.New("container status is unhealthy")} { + t.Run(healthErr.Error(), func(t *testing.T) { + s, runner := newTestServer(t) + s.healthCheck = func(ctx context.Context, _ string) error { + deadline, ok := ctx.Deadline() + if !ok || time.Until(deadline) < 5*time.Hour { + t.Error("migration got less than the Core startup budget") + } + if st := s.readState(); st.State != "checking" || st.PreviousImageID != "sha256:current" { + t.Errorf("missing in-flight image history: %+v", st) + } + return healthErr + } + s.selfReplace = func(string) error { t.Error("replaced updater before Core became ready"); return nil } + s.runJob("update", "v3.2.0-beta.2") + st := s.readState() + if st.State != "failed" || st.Target != "v3.2.0-beta.2" || st.PreviousImageID != "sha256:current" || !strings.Contains(st.Message, "verified full backup") { + t.Fatalf("state=%+v", st) + } + calls := runner.snapshot() + if len(calls) != 2 || !strings.Contains(strings.Join(calls[0], " "), "pull ftw") || !strings.Contains(strings.Join(calls[1], " "), "up -d ftw") { + t.Fatalf("readiness failure changed the running image: %v", calls) + } + }) + } +} + +func TestInterruptedUpdateRetainsImageHistoryWithoutTouchingCore(t *testing.T) { + s, runner := newTestServer(t) + s.writeState(State{State: "checking", Action: "update", Component: "core", Target: "v3.2.0-beta.2", PreviousImageID: "sha256:before", Message: "Waiting for the new service to become ready"}) + s.recoverCrashedState() + st := s.readState() + if st.State != "failed" || st.PreviousImageID != "sha256:before" || !strings.Contains(st.Message, "Core and data were left in place") { + t.Fatalf("state=%+v", st) + } + if calls := runner.snapshot(); len(calls) != 0 { + t.Fatalf("restart changed Core: %v", calls) + } +} diff --git a/go/cmd/ftw/boothealth.go b/go/cmd/ftw/boothealth.go index 946c60d39..cf98e2779 100644 --- a/go/cmd/ftw/boothealth.go +++ b/go/cmd/ftw/boothealth.go @@ -1,8 +1,14 @@ package main import ( + "encoding/json" "net/http" + "os" + "path/filepath" + "strings" "sync/atomic" + + "github.com/srcfl/ftw/go/internal/state" ) // swappableHandler lets the API port be bound before slow boot work (state @@ -28,21 +34,61 @@ func (s *swappableHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { (*s.h.Load()).ServeHTTP(w, r) } -// bootPhaseHandler answers health probes 200 while the process initializes, -// so a legitimately slow boot is distinguishable from a dead one. Everything -// else gets 503 + Retry-After so clients and the UI know to come back. +type bootHandler struct { + webDir string + migration atomic.Pointer[json.RawMessage] +} + func bootPhaseHandler() http.Handler { - mux := http.NewServeMux() - mux.HandleFunc("/api/health", func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte(`{"status":"starting","phase":"initializing state"}`)) - }) - mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - w.Header().Set("Retry-After", "10") + return newBootPhaseHandler("") +} + +func newBootPhaseHandler(webDir string) *bootHandler { + return &bootHandler{webDir: webDir} +} + +func (b *bootHandler) setMigration(progress state.HistoryMigrationStatus) { + data, err := json.Marshal(progress) + if err != nil { + return + } + snapshot := json.RawMessage(data) + b.migration.Store(&snapshot) +} + +// Health proves process liveness. Every other API stays unavailable until +// the fully wired handler replaces this one. Browser reloads get a progress +// page instead of a JSON error, without changing the updater's readiness test. +func (b *bootHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Cache-Control", "no-store") + if b.webDir != "" && (r.Method == http.MethodGet || r.Method == http.MethodHead) { + if r.URL.Path == "/history-migration.js" { + http.ServeFile(w, r, filepath.Join(b.webDir, "history-migration.js")) + return + } + if !strings.HasPrefix(r.URL.Path, "/api/") { + if page, err := os.ReadFile(filepath.Join(b.webDir, "boot.html")); err == nil { + w.Header().Set("Content-Type", "text/html; charset=utf-8") + w.Header().Set("Retry-After", "2") + w.WriteHeader(http.StatusServiceUnavailable) + if r.Method != http.MethodHead { + _, _ = w.Write(page) + } + return + } + } + } + w.Header().Set("Content-Type", "application/json") + payload := map[string]any{"phase": "initializing state"} + if migration := b.migration.Load(); migration != nil { + payload["migration"] = *migration + } + if r.URL.Path == "/api/health" { + payload["status"] = "starting" + } else { + payload["error"] = "starting" + w.Header().Set("Retry-After", "2") w.WriteHeader(http.StatusServiceUnavailable) - _, _ = w.Write([]byte(`{"error":"starting","phase":"initializing state"}`)) - }) - return mux + } + _ = json.NewEncoder(w).Encode(payload) } diff --git a/go/cmd/ftw/boothealth_test.go b/go/cmd/ftw/boothealth_test.go new file mode 100644 index 000000000..a015963da --- /dev/null +++ b/go/cmd/ftw/boothealth_test.go @@ -0,0 +1,82 @@ +package main + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + + "github.com/srcfl/ftw/go/internal/state" +) + +func TestBootProgressNeverClaimsAPIReadiness(t *testing.T) { + dir := t.TempDir() + for name, body := range map[string]string{"boot.html": "

Preparing FTW

", "history-migration.js": "export const fixture=true;"} { + if err := os.WriteFile(filepath.Join(dir, name), []byte(body), 0600); err != nil { + t.Fatal(err) + } + } + boot := newBootPhaseHandler(dir) + boot.setMigration(state.HistoryMigrationStatus{State: "starting", Phase: "seed", RowsDone: 2048}) + for _, path := range []string{"/api/health", "/api/status", "/api/version/update/status", "/", "/history-migration.js"} { + t.Run(path, func(t *testing.T) { + w := httptest.NewRecorder() + boot.ServeHTTP(w, httptest.NewRequest(http.MethodGet, path, nil)) + want := http.StatusServiceUnavailable + if path == "/api/health" || path == "/history-migration.js" { + want = http.StatusOK + } + if w.Code != want { + t.Fatalf("status=%d want=%d", w.Code, want) + } + if w.Header().Get("Cache-Control") != "no-store" { + t.Fatal("startup responses must not be cached") + } + if strings.HasPrefix(path, "/api/") { + var body struct { + Status string + Error string + Migration state.HistoryMigrationStatus + } + if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if body.Migration.RowsDone != 2048 || body.Migration.HistoryComplete { + t.Fatalf("wrong progress: %+v", body) + } + if path == "/api/health" && body.Status != "starting" { + t.Fatal("health must say starting") + } + if path != "/api/health" && body.Error != "starting" { + t.Fatal("full API must stay unavailable") + } + } else if path == "/" && !strings.Contains(w.Body.String(), "Preparing FTW") { + t.Fatal("browser must get startup page") + } + }) + } +} + +func TestBootProgressConcurrentReadersAndUpdates(t *testing.T) { + boot := newBootPhaseHandler("") + var wg sync.WaitGroup + for worker := 0; worker < 4; worker++ { + wg.Add(1) + go func() { + defer wg.Done() + for i := int64(0); i < 50; i++ { + boot.setMigration(state.HistoryMigrationStatus{State: "starting", Phase: "seed", RowsDone: i}) + w := httptest.NewRecorder() + boot.ServeHTTP(w, httptest.NewRequest(http.MethodGet, "/api/health", nil)) + if !json.Valid(w.Body.Bytes()) { + t.Error("partial progress response") + } + } + }() + } + wg.Wait() +} diff --git a/go/cmd/ftw/main.go b/go/cmd/ftw/main.go index 109526844..fbb819fbf 100644 --- a/go/cmd/ftw/main.go +++ b/go/cmd/ftw/main.go @@ -67,6 +67,7 @@ import ( "github.com/srcfl/ftw/go/internal/selfupdate" "github.com/srcfl/ftw/go/internal/state" "github.com/srcfl/ftw/go/internal/telemetry" + "github.com/srcfl/ftw/go/internal/updateipc" ) // Version gets injected at build time via -ldflags. Defaults to "dev" for @@ -359,6 +360,17 @@ func main() { slog.Warn("ignoring FTW_IMAGE_TAG that does not match a built release identity", "built_version", builtVersion, "built_candidate", CandidateTag, "image_tag", imageTag) } slog.Info("FTW starting", "version", Version, "config", *configPath) + // The previous updater may revert this image if startup fails. Confirm + // its failure behavior before config/bootstrap/state can write any data. + if envBool("FTW_SELFUPDATE_ENABLED") { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + err := updateipc.RequireSafeUpdater(ctx, envOr("FTW_UPDATER_SOCKET", "/run/ftw-update/sock")) + cancel() + if err != nil { + slog.Error("updater preflight", "err", err) + os.Exit(1) + } + } // Route "drivers/.lua" path resolution through the drivers dir // (from -drivers). Picked up by both the initial Load below and every @@ -432,7 +444,8 @@ func main() { bootPolicy := apiMutationPolicy() bootPolicy.LANAuthEnabled = lanAuth.Enabled bootPolicy.VerifyLANSecret = lanAuth.Verify - apiHandler := newSwappableHandler(bootPhaseHandler()) + boot := newBootPhaseHandler(*webDir) + apiHandler := newSwappableHandler(boot) httpSrv := &http.Server{ Addr: fmt.Sprintf(":%d", cfg.API.Port), Handler: api.WithSecurityHeaders(api.Authenticate(apiHandler, bootPolicy)), @@ -445,7 +458,7 @@ func main() { } }() - st, err := state.OpenWithLegacyHistory(statePath, coldDir) + st, err := state.OpenWithBackgroundHistory(statePath, coldDir, boot.setMigration) if err != nil { slog.Error("open state", "err", err) os.Exit(1) diff --git a/go/cmd/ftw/updater_preflight_test.go b/go/cmd/ftw/updater_preflight_test.go new file mode 100644 index 000000000..1325c0cab --- /dev/null +++ b/go/cmd/ftw/updater_preflight_test.go @@ -0,0 +1,81 @@ +package main + +import ( + "context" + "flag" + "net" + "net/http" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/updateipc" +) + +func TestUpdaterPreflightBeforePersistentState(t *testing.T) { + if path := os.Getenv("FTW_PREFLIGHT_TEST_CONFIG"); path != "" { + flag.CommandLine = flag.NewFlagSet("ftw", flag.ExitOnError) + os.Args = []string{"ftw", "-config", path} + main() + return + } + for _, mode := range []string{"old", "fixed", "native"} { + t.Run(mode, func(t *testing.T) { + dir, err := os.MkdirTemp("", "ftw-preflight-") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(dir) + socket := filepath.Join(dir, "sock") + ln, err := net.Listen("unix", socket) + if err != nil { + t.Fatal(err) + } + srv := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if mode == "fixed" { + updateipc.ServeCapabilities(w, r) + return + } + http.NotFound(w, r) + })} + go srv.Serve(ln) + defer srv.Close() + data := t.TempDir() + path := filepath.Join(data, "config.yaml") + original := []byte("invalid: [yaml\n") + if err := os.WriteFile(path, original, 0600); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + cmd := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestUpdaterPreflightBeforePersistentState$") + enabled := "1" + if mode == "native" { + enabled = "0" + } + cmd.Env = append(os.Environ(), "FTW_PREFLIGHT_TEST_CONFIG="+path, "FTW_UPDATER_SOCKET="+socket, "FTW_SELFUPDATE_ENABLED="+enabled) + out, err := cmd.CombinedOutput() + if err == nil { + t.Fatalf("unexpected startup success: %s", out) + } + guarded := strings.Contains(string(out), "update ftw-updater first") + if guarded != (mode == "old") { + t.Fatalf("wrong startup path: %s", out) + } + if mode != "old" && !strings.Contains(string(out), "load config") { + t.Fatalf("did not pass preflight: %s", out) + } + files, err := os.ReadDir(data) + if err != nil || len(files) != 1 { + t.Fatalf("startup changed data: %v %v", files, err) + } + got, err := os.ReadFile(path) + if err != nil || string(got) != string(original) { + t.Fatalf("startup changed config: %q %v", got, err) + } + }) + } +} diff --git a/go/internal/api/api.go b/go/internal/api/api.go index dc8f7b334..23c291777 100644 --- a/go/internal/api/api.go +++ b/go/internal/api/api.go @@ -705,7 +705,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 != "" || (writer.LastRejectMS > 0 && time.Now().UnixMilli()-writer.LastRejectMS < time.Minute.Milliseconds()) { resp["status"] = "degraded" } } diff --git a/go/internal/backup/archive_test.go b/go/internal/backup/archive_test.go index de4abd7e6..5920c39e3 100644 --- a/go/internal/backup/archive_test.go +++ b/go/internal/backup/archive_test.go @@ -146,6 +146,7 @@ func TestDuckDBBackupOmitsImportedSamplesAndLiveFiles(t *testing.T) { t.Fatal(err) } writeTestFile(t, filepath.Join(liveTmp, "spill.bin"), "transient history data") + writeTestFile(t, filepath.Join(state.HistoryDatabasePath(statePath)+".import-abandoned", "staging.duckdb"), "abandoned import staging") writeTestFile(t, filepath.Join(coldDir, "diagnostics", "2026", "01", "01.parquet"), "diagnostic archive") info, err := Create(context.Background(), CreateOptions{State: st, StatePath: statePath, DataDir: dataDir, OutputDir: filepath.Join(root, "backups")}) if err != nil { diff --git a/go/internal/state/history_connection.go b/go/internal/state/history_connection.go new file mode 100644 index 000000000..305bd0e6c --- /dev/null +++ b/go/internal/state/history_connection.go @@ -0,0 +1,175 @@ +package state + +import ( + "context" + "database/sql/driver" + "errors" + "sync" + "time" + + duckdb "github.com/duckdb/duckdb-go/v2" +) + +// historyConnector keeps database/sql stable while rotating a native instance +// between import files. A physical connection holds a read lease until SQL has +// closed its Rows, Tx, statements and Raw calls and finally closes the connection. +// The pool MUST have MaxIdleConns(0); an idle connection would retain its lease. +// Rotation uses TryLock so waiting for a reader never prevents new live writes. +// https://duckdb.org/docs/current/guides/performance/indexing#indexes-and-memory +// explains why connection closure alone cannot release DuckDB's index buffers. +type historyConnector struct { + mu sync.RWMutex + native *duckdb.Connector + dsn string + closed bool +} + +func newHistoryConnector(dsn string) (*historyConnector, error) { + native, err := duckdb.NewConnector(dsn, nil) + if err != nil { + return nil, err + } + return &historyConnector{native: native, dsn: dsn}, nil +} + +func (c *historyConnector) Driver() driver.Driver { return &duckdb.Driver{} } + +func (c *historyConnector) lock(ctx context.Context, exclusive bool) error { + for { + if err := ctx.Err(); err != nil { + return err + } + if exclusive { + if c.mu.TryLock() { + return nil + } + } else { + if c.mu.TryRLock() { + return nil + } + } + timer := time.NewTimer(time.Millisecond) + select { + case <-ctx.Done(): + timer.Stop() + return ctx.Err() + case <-timer.C: + } + } +} + +func (c *historyConnector) Connect(ctx context.Context) (driver.Conn, error) { + for { + if err := c.lock(ctx, false); err != nil { + return nil, err + } + if c.closed { + c.mu.RUnlock() + return nil, errors.New("history database is closed") + } + if c.native != nil { + conn, err := c.native.Connect(ctx) + if err != nil { + c.mu.RUnlock() + return nil, err + } + return &historyLeaseConn{Conn: conn.(*duckdb.Conn), release: c.mu.RUnlock}, nil + } + c.mu.RUnlock() + // A failed reopen remains a visible storage error. Future connections may + // retry opening the same durable file after the underlying problem clears. + if err := c.lock(ctx, true); err != nil { + return nil, err + } + var err error + if c.native == nil && !c.closed { + c.native, err = duckdb.NewConnector(c.dsn, nil) + } + c.mu.Unlock() + if err != nil { + return nil, err + } + } +} + +func (c *historyConnector) rotate(ctx context.Context) error { + if err := c.lock(ctx, true); err != nil { + return err + } + defer c.mu.Unlock() + if c.closed { + return errors.New("history database is closed") + } + if c.native != nil { + // Flush while there are no live connections, and retain this instance if + // checkpoint fails. A rotation never interrupts a reader or transaction. + conn, err := c.native.Connect(ctx) + if err != nil { + return err + } + _, err = conn.(driver.ExecerContext).ExecContext(ctx, "CHECKPOINT", nil) + err = errors.Join(err, conn.Close()) + if err != nil { + return err + } + if err := c.native.Close(); err != nil { + return err + } + c.native = nil + } + var err error + c.native, err = duckdb.NewConnector(c.dsn, nil) + return err +} + +func (c *historyConnector) Close() error { + c.mu.Lock() + defer c.mu.Unlock() + c.closed = true + if c.native == nil { + return nil + } + err := c.native.Close() + c.native = nil + return err +} + +// Embedding preserves the driver's optional context and value interfaces. +// database/sql serializes a physical connection; Close is called once. +type historyLeaseConn struct { + *duckdb.Conn + release func() +} + +func (c *historyLeaseConn) Close() error { + err := c.Conn.Close() + c.release() + return err +} + +// Call only within sql.Conn.Raw, while database/sql still holds the lease. +func nativeHistoryConn(raw any) driver.Conn { + if c, ok := raw.(*historyLeaseConn); ok { + return c.Conn + } + return raw.(driver.Conn) +} + +func (s *Store) rotateHistory(ctx context.Context) error { + if s.historyConnector == nil { + return nil + } + for { + attempt, cancel := context.WithTimeout(ctx, 5*time.Second) + err := s.historyConnector.rotate(attempt) + cancel() + if !errors.Is(err, context.DeadlineExceeded) || ctx.Err() != nil { + return err + } + // A long query postpones import. Its lease remains intact and live writer + // connections can still proceed while the importer waits for a quiet point. + if s.historyMigration != nil { + s.historyMigration.update(func(*HistoryMigrationStatus) {}) + } + } +} diff --git a/go/internal/state/history_connection_test.go b/go/internal/state/history_connection_test.go new file mode 100644 index 000000000..49d36cb77 --- /dev/null +++ b/go/internal/state/history_connection_test.go @@ -0,0 +1,240 @@ +package state + +import ( + "context" + "database/sql/driver" + "errors" + "testing" + "time" +) + +func TestHistoryRotationWaitsForSQLLifetimes(t *testing.T) { + for _, kind := range []string{"rows", "row", "transaction", "raw", "statement"} { + t.Run(kind, func(t *testing.T) { + s := freshStore(t) + ctx := context.Background() + var release func() + switch kind { + case "rows": + rows, err := s.history.Query(`SELECT * FROM range(3)`) + if err != nil { + t.Fatal(err) + } + release = func() { + var n int64 + if !rows.Next() { + t.Error("rows lost during attempted rotation") + } + if err := rows.Scan(&n); err != nil { + t.Error(err) + } + rows.Close() + } + case "row": + row := s.history.QueryRow(`SELECT 42`) + release = func() { + var n int + if err := row.Scan(&n); err != nil || n != 42 { + t.Errorf("row=%d %v", n, err) + } + } + case "transaction": + tx, err := s.history.BeginTx(ctx, nil) + if err != nil { + t.Fatal(err) + } + release = func() { + var n int + if err := tx.QueryRow(`SELECT 42`).Scan(&n); err != nil || n != 42 { + t.Errorf("tx=%d %v", n, err) + } + tx.Rollback() + } + case "statement": + conn, err := s.history.Conn(ctx) + if err != nil { + t.Fatal(err) + } + stmt, err := conn.PrepareContext(ctx, `SELECT 42`) + if err != nil { + t.Fatal(err) + } + release = func() { + var n int + if err := stmt.QueryRowContext(ctx).Scan(&n); err != nil || n != 42 { + t.Errorf("stmt=%d %v", n, err) + } + stmt.Close() + conn.Close() + } + case "raw": + conn, err := s.history.Conn(ctx) + if err != nil { + t.Fatal(err) + } + entered, finish, done := make(chan struct{}), make(chan struct{}), make(chan error, 1) + go func() { + done <- conn.Raw(func(raw any) error { + close(entered) + <-finish + _, err := nativeHistoryConn(raw).(driver.ExecerContext).ExecContext(ctx, `SELECT 42`, nil) + return err + }) + }() + <-entered + release = func() { + close(finish) + if err := <-done; err != nil { + t.Error(err) + } + conn.Close() + } + } + native := s.historyConnector.native + attempt, cancel := context.WithTimeout(ctx, 20*time.Millisecond) + err := s.historyConnector.rotate(attempt) + cancel() + if !errors.Is(err, context.DeadlineExceeded) { + release() + t.Fatalf("rotation ignored %s lease: %v", kind, err) + } + if s.historyConnector.native != native { + release() + t.Fatal("rotation replaced an active instance") + } + // Waiting for the old reader must not block a different live connection. + if err := s.RecordSamples([]Sample{{TsMs: 1, Driver: "live", Metric: "power", Value: 42}}); err != nil { + release() + t.Fatal(err) + } + release() + if err := s.historyConnector.rotate(ctx); err != nil { + t.Fatal(err) + } + if s.historyConnector.native == native { + t.Fatal("quiescent instance did not rotate") + } + if got, err := s.LatestSample("live", "power"); err != nil || got.Value != 42 { + t.Fatalf("live sample lost: %+v %v", got, err) + } + }) + } +} + +func TestHistoryConnectHonorsCancellationDuringRotation(t *testing.T) { + s := freshStore(t) + s.historyConnector.mu.Lock() + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + _, err := s.historyConnector.Connect(ctx) + cancel() + s.historyConnector.mu.Unlock() + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("Connect ignored cancellation: %v", err) + } + if err := s.RecordHistory(HistoryPoint{TsMs: 1}); err != nil { + t.Fatal(err) + } +} + +func TestHistoryPreparedStatementSurvivesNativeRotation(t *testing.T) { + s := freshStore(t) + stmt, err := s.history.Prepare(`SELECT 42`) + if err != nil { + t.Fatal(err) + } + defer stmt.Close() + for range 2 { + if err := s.historyConnector.rotate(context.Background()); err != nil { + t.Fatal(err) + } + var n int + if err := stmt.QueryRow().Scan(&n); err != nil || n != 42 { + t.Fatalf("prepared statement=%d %v", n, err) + } + } +} + +func TestLiveWriterRotatesAfterCommittedRows(t *testing.T) { + s := freshStore(t) + s.historyWriter.maintenanceRowsLimit = 2 + before := s.historyConnector.native + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + for i := range 3 { + if err := s.EnqueueTelemetryTick(nil, []Sample{{TsMs: int64(i), Driver: "live", Metric: "power", Value: float64(i)}}, nil); err != nil { + t.Fatal(err) + } + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + for s.HistoryWriterStatus().MaintenanceRuns == 0 { + if ctx.Err() != nil { + t.Fatal(ctx.Err()) + } + time.Sleep(time.Millisecond) + } + s.historyConnector.mu.RLock() + same := s.historyConnector.native == before + s.historyConnector.mu.RUnlock() + if same { + t.Fatal("live commits did not rotate the native instance") + } + if st := s.HistoryWriterStatus(); st.Committed != 3 || st.MaintenanceError != "" || st.LastMaintenanceMS == 0 { + t.Fatalf("writer=%+v", st) + } + got, err := s.LoadSeries("live", "power", 0, 3, 0) + if err != nil || len(got) != 3 { + t.Fatalf("committed rows=%v %v", got, err) + } +} + +func TestLiveWriterRetriesMaintenanceAfterLongReader(t *testing.T) { + s := freshStore(t) + s.historyWriter.maintenanceRowsLimit = 1 + s.historyWriter.maintenanceRetryDelay = 0 + reader, err := s.history.Begin() + if err != nil { + t.Fatal(err) + } + defer reader.Rollback() + ctx, cancel := context.WithTimeout(context.Background(), 6*time.Second) + defer cancel() + if err := s.EnqueueTelemetryTick(nil, []Sample{{TsMs: 1, Driver: "live", Metric: "power", Value: 42}}, nil); err != nil { + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + for s.HistoryWriterStatus().MaintenanceError == "" { + if ctx.Err() != nil { + t.Fatal(ctx.Err()) + } + time.Sleep(time.Millisecond) + } + if st := s.HistoryWriterStatus(); st.Committed != 1 || st.LastError != "" { + t.Fatalf("maintenance hid a durable commit: %+v", st) + } + var n int + if err := reader.QueryRow(`SELECT 42`).Scan(&n); err != nil || n != 42 { + t.Fatalf("reader interrupted: %d %v", n, err) + } + if err := reader.Rollback(); err != nil { + t.Fatal(err) + } + if err := s.EnqueueTelemetryTick(nil, []Sample{{TsMs: 2, Driver: "live", Metric: "power", Value: 43}}, nil); err != nil { + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + for s.HistoryWriterStatus().MaintenanceRuns == 0 { + if ctx.Err() != nil { + t.Fatal(ctx.Err()) + } + time.Sleep(time.Millisecond) + } + if st := s.HistoryWriterStatus(); st.MaintenanceError != "" || st.Committed != 2 || st.Rejected != 0 { + t.Fatalf("maintenance retry=%+v", st) + } +} diff --git a/go/internal/state/history_duckdb.go b/go/internal/state/history_duckdb.go index 7fc25cc53..b43b7951a 100644 --- a/go/internal/state/history_duckdb.go +++ b/go/internal/state/history_duckdb.go @@ -47,12 +47,14 @@ func (s *Store) openHistory() error { if _, err := os.Stat(s.historyPath); errors.Is(err, os.ErrNotExist) && active != "" && restore == "" { return errors.New("primary DuckDB history is missing; restore a full backup") } - db, err := sql.Open("duckdb", s.historyPath+"?threads=2&memory_limit=128MB&max_temp_directory_size=512MB&autoload_known_extensions=false&autoinstall_known_extensions=false") + connector, err := newHistoryConnector(s.historyPath + "?threads=2&memory_limit=128MB&max_temp_directory_size=512MB&autoload_known_extensions=false&autoinstall_known_extensions=false") if err != nil { return fmt.Errorf("open DuckDB history: %w", err) } + db := sql.OpenDB(connector) db.SetMaxOpenConns(4) - db.SetMaxIdleConns(4) + db.SetMaxIdleConns(0) + s.historyConnector = connector s.history = db ok := false defer func() { @@ -70,14 +72,21 @@ func (s *Store) openHistory() error { if err := db.QueryRow(`SELECT COUNT(*) FROM history_migrations WHERE name='sqlite-v1'`).Scan(&complete); err != nil { return err } + var seeded int + if err := db.QueryRow(`SELECT COUNT(*) FROM history_migrations WHERE name='sqlite-seed-v1'`).Scan(&seeded); err != nil { + return err + } + if seeded != 0 && complete == 0 && s.historyMigration == nil { + return errors.New("raw history migration is pending; start Core to resume it before using offline history tools") + } var generation string - if complete != 0 { + if complete != 0 || seeded != 0 { if err := db.QueryRow(`SELECT name FROM history_migrations WHERE name LIKE 'generation:%'`).Scan(&generation); err != nil { return err } generation = strings.TrimPrefix(generation, "generation:") } - if restore != "" && complete != 0 && generation != restore { + if restore != "" && (complete != 0 || seeded != 0) && generation != restore { // A restore explicitly selects its SQLite snapshot as the source. // Preserve the previous DuckDB files; never silently reuse old history. db.Close() @@ -91,10 +100,10 @@ func (s *Store) openHistory() error { ok = true return s.openHistory() } - if complete != 0 && active == "" && restore == "" && intent != generation { + if (complete != 0 || seeded != 0) && active == "" && restore == "" && intent != generation { return errors.New("unbound DuckDB history beside SQLite; restore a full backup before starting") } - if complete == 0 { + if complete == 0 && seeded == 0 { if active != "" && restore == "" { return errors.New("primary DuckDB history is incomplete; restore a full backup") } @@ -150,6 +159,15 @@ func (s *Store) migrateSQLiteHistory(ctx context.Context, generation string) err } defer conn.Close() for _, table := range historyTables { + if s.historyMigration != nil && table == "ts_samples" { + continue + } + if s.historyMigration != nil { + s.historyMigration.update(func(st *HistoryMigrationStatus) { + st.CurrentSource = table + st.CurrentSourceRowsDone, st.CurrentSourceRowsTotal = 0, 0 + }) + } // The destination stays inactive until every table passes verification. // Bounded commits keep migration memory independent of source row count. if _, err := conn.ExecContext(ctx, `DELETE FROM `+table); err != nil { @@ -173,7 +191,7 @@ func (s *Store) migrateSQLiteHistory(ctx context.Context, generation string) err return err } err = conn.Raw(func(raw any) error { - app, err := duckdb.NewAppenderFromConn(raw.(driver.Conn), "", table) + app, err := duckdb.NewAppenderFromConn(nativeHistoryConn(raw), "", table) if err != nil { return err } @@ -219,6 +237,9 @@ func (s *Store) migrateSQLiteHistory(ctx context.Context, generation string) err rows.Close() return fmt.Errorf("migrate history table %s: %w", table, err) } + if s.historyMigration != nil { + s.historyMigration.update(func(st *HistoryMigrationStatus) { st.CurrentSourceRowsDone = count }) + } } rows.Close() actualHash := sha256.New() @@ -236,8 +257,9 @@ func (s *Store) migrateSQLiteHistory(ctx context.Context, generation string) err return err } defer conn.ExecContext(context.Background(), `ROLLBACK`) - // Seed generated IDs above the imported IDs. Sequences are deliberately - // not relied on for rollback; gaps in IDs have no semantic meaning. + // Seed IDs above the imported IDs. Inserts select these sequences explicitly: + // ALTER COLUMN SET DEFAULT nextval cannot replay safely from DuckDB's WAL. + // Gaps in IDs have no semantic meaning. for _, table := range []string{"ts_drivers", "ts_metrics"} { var next int64 if err := conn.QueryRowContext(ctx, `SELECT COALESCE(MAX(id), 0)+1 FROM `+table).Scan(&next); err != nil { @@ -246,14 +268,15 @@ func (s *Store) migrateSQLiteHistory(ctx context.Context, generation string) err if _, err := conn.ExecContext(ctx, fmt.Sprintf(`CREATE OR REPLACE SEQUENCE %s_next_id START %d`, table, next)); err != nil { return err } - if _, err := conn.ExecContext(ctx, fmt.Sprintf(`ALTER TABLE %s ALTER COLUMN id SET DEFAULT nextval('%s_next_id')`, table, table)); err != nil { - return err - } } if _, err := conn.ExecContext(ctx, `INSERT INTO history_migrations(name) VALUES (?)`, "generation:"+generation); err != nil { return err } - if _, err := conn.ExecContext(ctx, `INSERT INTO history_migrations(name) VALUES ('sqlite-v1')`); err != nil { + marker := "sqlite-v1" + if s.historyMigration != nil { + marker = "sqlite-seed-v1" + } + if _, err := conn.ExecContext(ctx, `INSERT INTO history_migrations(name) VALUES (?)`, marker); err != nil { return err } if _, err := conn.ExecContext(ctx, `COMMIT`); err != nil { @@ -319,6 +342,7 @@ func historyFloatBits(value float64) uint64 { func (s *Store) HistoryBackend() map[string]any { info := map[string]any{"engine": "duckdb", "version": "1.5.5", "role": "primary", "file": filepath.Base(s.historyPath), "writer": s.HistoryWriterStatus()} + info["migration"] = s.HistoryMigrationStatus() for key, path := range map[string]string{"file_bytes": s.historyPath, "wal_bytes": s.historyPath + ".wal"} { if stat, err := os.Stat(path); err == nil { info[key] = stat.Size() @@ -334,6 +358,12 @@ func (s *Store) CheckpointHistory(ctx context.Context) error { if s.history == nil { return nil } + if s.historyConnector != nil { + // A checkpoint alone does not evict all native index/table buffers. + // Wait for a gap between active connections, then reopen the native + // instance while retaining the public SQL pool and durable primary. + return s.rotateHistory(ctx) + } s.historyWriteMu.Lock() defer s.historyWriteMu.Unlock() _, err := s.history.ExecContext(ctx, `CHECKPOINT`) @@ -434,8 +464,15 @@ func (s *Store) exportHistoryToSQLite(path string) error { return err } defer src.Rollback() + var sqliteComplete int + if err := src.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_migrations WHERE name='sqlite-v1'`).Scan(&sqliteComplete); err != nil { + return err + } + if sqliteComplete == 0 || !s.HistoryMigrationStatus().HistoryComplete { + return errors.New("finish historical import before exporting a full backup; keep the verified pre-update backup") + } var pending int - if err := src.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_imports`).Scan(&pending); err != nil { + if err := src.QueryRowContext(ctx, `SELECT (SELECT COUNT(*) FROM history_parquet_imports) + (SELECT COUNT(*) FROM history_migrations WHERE name='legacy-import-pending')`).Scan(&pending); err != nil { return err } if pending != 0 { @@ -503,7 +540,7 @@ func (s *Store) exportHistoryToSQLite(path string) error { func appendHistoryRows(conn *sql.Conn, points []HistoryPoint) error { return conn.Raw(func(raw any) error { - app, err := duckdb.NewTableAppender(raw.(driver.Conn), `INSERT OR REPLACE INTO history_hot SELECT * FROM appended_data`, "", "", "history_hot", nil) + app, err := duckdb.NewTableAppender(nativeHistoryConn(raw), `INSERT OR REPLACE INTO history_hot SELECT * FROM appended_data`, "", "", "history_hot", nil) if err != nil { return err } diff --git a/go/internal/state/history_migration.go b/go/internal/state/history_migration.go new file mode 100644 index 000000000..6e6bd83e0 --- /dev/null +++ b/go/internal/state/history_migration.go @@ -0,0 +1,263 @@ +package state + +import ( + "context" + "database/sql" + "database/sql/driver" + "errors" + "fmt" + "log/slog" + "sync" + "time" + + duckdb "github.com/duckdb/duckdb-go/v2" +) + +// HistoryMigrationStatus describes historical coverage, independently of Core +// readiness and live-writer durability. Unknown totals and bounds are omitted. +type HistoryMigrationStatus struct { + State string `json:"state"` + Phase string `json:"phase"` + HistoryComplete bool `json:"history_complete"` + FilesDone int `json:"files_done"` + FilesTotal int `json:"files_total"` + RowsDone int64 `json:"rows_done"` + RowsTotal int64 `json:"rows_total,omitempty"` + CurrentSource string `json:"current_source,omitempty"` + CurrentSourceRowsDone int64 `json:"current_source_rows_done"` + CurrentSourceRowsTotal int64 `json:"current_source_rows_total,omitempty"` + StartedAtMS int64 `json:"started_at_ms"` + UpdatedAtMS int64 `json:"updated_at_ms"` + LastError string `json:"last_error,omitempty"` + IncompleteFromMS *int64 `json:"incomplete_from_ms,omitempty"` + IncompleteUntilMS *int64 `json:"incomplete_until_ms,omitempty"` +} + +type historyMigration struct { + mu sync.Mutex + status HistoryMigrationStatus + progress func(HistoryMigrationStatus) + ctx context.Context + cancel context.CancelFunc + done chan struct{} +} + +func newHistoryMigration(progress func(HistoryMigrationStatus)) *historyMigration { + ctx, cancel := context.WithCancel(context.Background()) + m := &historyMigration{ctx: ctx, cancel: cancel, done: make(chan struct{}), progress: progress} + m.update(func(st *HistoryMigrationStatus) { + st.State, st.Phase = "starting", "seed" + st.StartedAtMS = time.Now().UnixMilli() + }) + return m +} + +func (m *historyMigration) update(change func(*HistoryMigrationStatus)) { + m.mu.Lock() + change(&m.status) + m.status.UpdatedAtMS = time.Now().UnixMilli() + st := m.status + m.mu.Unlock() + if m.progress != nil { + m.progress(st) + } +} + +func (s *Store) HistoryMigrationStatus() HistoryMigrationStatus { + if s.historyMigration == nil { + return HistoryMigrationStatus{State: "complete", Phase: "complete", HistoryComplete: true} + } + m := s.historyMigration + m.mu.Lock() + defer m.mu.Unlock() + return m.status +} + +func (s *Store) runHistoryMigration(coldDir string) { + m := s.historyMigration + defer close(m.done) + m.update(func(st *HistoryMigrationStatus) { st.State, st.Phase = "running", "sqlite" }) + err := s.importSQLiteSamples(m.ctx) + if err == nil { + err = s.ImportLegacyParquet(m.ctx, coldDir) + } + if err == nil { + err = s.CheckpointHistory(m.ctx) + } + if err == nil { + s.historyWriteMu.Lock() + _, err = s.history.ExecContext(m.ctx, `DELETE FROM history_migrations WHERE name='legacy-import-pending'`) + s.historyWriteMu.Unlock() + } + if err != nil { + slog.Error("historical import paused; live collection continues; original sources retained", "err", err) + m.update(func(st *HistoryMigrationStatus) { + st.State = "failed" + st.LastError = "Historical import paused. Keep the original history files and pre-update backup; check Core logs before retrying." + }) + return + } + m.update(func(st *HistoryMigrationStatus) { + st.State, st.Phase, st.HistoryComplete = "complete", "complete", true + st.CurrentSource = "" + st.CurrentSourceRowsDone, st.CurrentSourceRowsTotal = 0, 0 + st.IncompleteFromMS, st.IncompleteUntilMS = nil, nil + }) + slog.Info("historical import complete; all source rows verified") +} + +// SQLite's legacy sample table is frozen after Core selects DuckDB. Each +// verified merge commits its source cursor in the same primary transaction. +// A restart never clears primary rows written by the live collector. +func (s *Store) importSQLiteSamples(ctx context.Context) error { + var complete int + if err := s.history.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_migrations WHERE name='sqlite-v1'`).Scan(&complete); err != nil { + return err + } + if complete != 0 { + var count int64 + if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM ts_samples`).Scan(&count); err != nil { + return err + } + s.historyMigration.update(func(st *HistoryMigrationStatus) { st.RowsDone = count }) + return nil + } + var count int64 + if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM ts_samples`).Scan(&count); err != nil { + return err + } + var done, d, m, ts int64 + err := s.history.QueryRowContext(ctx, `SELECT rows_done,driver_id,metric_id,ts_ms FROM history_sqlite_progress WHERE source='ts_samples'`).Scan(&done, &d, &m, &ts) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + s.historyMigration.update(func(st *HistoryMigrationStatus) { + st.CurrentSource = "state.db/ts_samples" + st.RowsDone, st.RowsTotal = done, count + st.CurrentSourceRowsDone, st.CurrentSourceRowsTotal = done, count + }) + for { + if err := s.yieldHistoryImport(ctx); err != nil { + return err + } + q := `SELECT s.driver_id,s.metric_id,s.ts_ms,s.value,d.name,m.name FROM ts_samples s JOIN ts_drivers d ON d.id=s.driver_id JOIN ts_metrics m ON m.id=s.metric_id` + var args []any + if done > 0 { + q += ` WHERE (s.driver_id,s.metric_id,s.ts_ms) > (?,?,?)` + args = []any{d, m, ts} + } + q += ` ORDER BY s.driver_id,s.metric_id,s.ts_ms LIMIT 2048` + rows, err := s.db.QueryContext(ctx, q, args...) + if err != nil { + return err + } + batch := make([]Sample, 0, historyImportRows) + for rows.Next() { + var sm Sample + if err := rows.Scan(&d, &m, &sm.TsMs, &sm.Value, &sm.Driver, &sm.Metric); err != nil { + rows.Close() + return err + } + ts = sm.TsMs + batch = append(batch, sm) + } + err = errors.Join(rows.Err(), rows.Close()) + if err != nil { + return err + } + if len(batch) == 0 { + break + } + next := done + int64(len(batch)) + if err := s.mergeHistoricalSamples(ctx, batch, func(tx *sql.Tx) error { + _, err := tx.ExecContext(ctx, `INSERT INTO history_sqlite_progress VALUES ('ts_samples',?,?,?,?) ON CONFLICT(source) DO UPDATE SET rows_done=excluded.rows_done,driver_id=excluded.driver_id,metric_id=excluded.metric_id,ts_ms=excluded.ts_ms`, next, d, m, ts) + return err + }); err != nil { + return err + } + done = next + s.historyMigration.update(func(st *HistoryMigrationStatus) { st.RowsDone, st.CurrentSourceRowsDone = done, done }) + if done%(64*historyImportRows) == 0 { + if err := s.CheckpointHistory(ctx); err != nil { + return err + } + } + } + if done != count { + return fmt.Errorf("SQLite sample count changed: verified %d expected %d", done, count) + } + s.historyWriteMu.Lock() + defer s.historyWriteMu.Unlock() + _, err = s.history.ExecContext(ctx, `INSERT INTO history_migrations(name) VALUES ('sqlite-v1')`) + return err +} + +// Let queued live ticks commit before importing another bounded chunk. No +// database or catalog lock is held while the importer yields. +func (s *Store) yieldHistoryImport(ctx context.Context) error { + for { + if err := ctx.Err(); err != nil { + return err + } + if s.HistoryWriterStatus().Pending == 0 { + return nil + } + timer := time.NewTimer(10 * time.Millisecond) + select { + case <-ctx.Done(): + timer.Stop() + return ctx.Err() + case <-timer.C: + } + } +} + +// A private, chunk-sized staging table bounds the live instance's working set. +// Source-sized sorting and duplicate checks belong to a separate instance. +func (s *Store) mergeHistoricalSamples(ctx context.Context, samples []Sample, receipt func(*sql.Tx) error) error { + if err := validateHistorySamples(samples); err != nil { + return err + } + if err := s.hydrateIntern(); err != nil { + return err + } + values := make([][]driver.Value, 0, len(samples)) + for _, sm := range samples { + d, err := s.driverID(sm.Driver) + if err != nil { + return err + } + m, err := s.metricID(sm.Metric, "") + if err != nil { + return err + } + values = append(values, []driver.Value{sm.TsMs, d, m, canonicalHistoryFloat(sm.Value)}) + } + s.historyWriteMu.Lock() + defer s.historyWriteMu.Unlock() + conn, err := s.history.Conn(ctx) + if err != nil { + return err + } + defer conn.Close() + if _, err := conn.ExecContext(ctx, `CREATE TEMP TABLE history_import_source (ts_ms BIGINT NOT NULL,driver_id BIGINT NOT NULL,metric_id BIGINT NOT NULL,value DOUBLE NOT NULL)`); err != nil { + return err + } + defer conn.ExecContext(context.Background(), `DROP TABLE IF EXISTS history_import_source`) + if err := conn.Raw(func(raw any) error { + app, err := duckdb.NewAppender(nativeHistoryConn(raw), "temp", "main", "history_import_source") + if err != nil { + return err + } + var appendErr error + for _, row := range values { + if appendErr = app.AppendRow(row...); appendErr != nil { + break + } + } + return errors.Join(appendErr, app.Close()) + }); err != nil { + return err + } + return importHistoryChunkCommit(ctx, conn, 0, int64(len(samples)), receipt) +} diff --git a/go/internal/state/history_migration_test.go b/go/internal/state/history_migration_test.go new file mode 100644 index 000000000..3acfe5ce0 --- /dev/null +++ b/go/internal/state/history_migration_test.go @@ -0,0 +1,619 @@ +package state + +import ( + "context" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "sync" + "testing" + "time" +) + +// A frozen SQLite source beside a fresh DuckDB destination, as on first update. +func legacyMigrationFixture(t *testing.T, n int) (string, string) { + t.Helper() + dir := t.TempDir() + path := filepath.Join(dir, "state.db") + cold := filepath.Join(dir, "cold") + s, err := Open(path) + if err != nil { + t.Fatal(err) + } + if _, err := s.db.Exec(`INSERT INTO ts_drivers(id,name) VALUES(7,'meter'); INSERT INTO ts_metrics(id,name,unit) VALUES(9,'power','W'); DELETE FROM config WHERE key LIKE 'history_%'`); err != nil { + t.Fatal(err) + } + tx, err := s.db.Begin() + if err != nil { + t.Fatal(err) + } + stmt, err := tx.Prepare(`INSERT INTO ts_samples(driver_id,metric_id,ts_ms,value) VALUES (7,9,?,?)`) + if err != nil { + t.Fatal(err) + } + for i := 1; i <= n; i++ { + if _, err := stmt.Exec(i, float64(i)/7); err != nil { + t.Fatal(err) + } + } + stmt.Close() + if err := tx.Commit(); err != nil { + t.Fatal(err) + } + if _, err := s.db.Exec(`INSERT INTO energy_ledger_cursors VALUES ('site','import','counter',1234,10)`); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + if err := os.Remove(historyDatabasePath(path)); err != nil { + t.Fatal(err) + } + return path, cold +} + +func TestBackgroundHistoryKeepsLiveWritesAcrossInterruptedImport(t *testing.T) { + path, cold := legacyMigrationFixture(t, historyImportRows*3+11) + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + points := make([]parquetSampleRow, historyImportRows*2+17) + for i := range points { + points[i] = parquetSampleRow{TsMs: int64(5000 + i), Driver: "meter", Metric: "power", Value: float64(i) + 0.25} + } + if err := writeParquetDay(filepath.Join(day, "01.parquet"), points); err != nil { + t.Fatal(err) + } + reached, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + s, err := OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "sqlite" && st.RowsDone == historyImportRows { + once.Do(func() { close(reached); <-release }) + } + }) + if err != nil { + t.Fatal(err) + } + defer func() { s.Close() }() + select { + case <-reached: + case <-time.After(10 * time.Second): + close(release) + t.Fatal("import never reached first committed chunk") + } + primary := s.history + s.historyWriter.maintenanceRowsLimit = 2 + var cursor float64 + if err := s.history.QueryRow(`SELECT value FROM energy_ledger_cursors WHERE asset_id='site'`).Scan(&cursor); err != nil || cursor != 1234 { + close(release) + t.Fatalf("accounting was not seeded: %v %v", cursor, err) + } + if s.HistoryMigrationStatus().HistoryComplete { + close(release) + t.Fatal("partial history reported complete") + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + live := []Sample{{TsMs: 3000, Driver: "meter", Metric: "power", Value: 900}, {TsMs: 8000, Driver: "meter", Metric: "power", Value: 901}, {TsMs: 999999, Driver: "new meter", Metric: "new power", Value: 902}} + if err := s.EnqueueTelemetryTick(nil, live, nil); err != nil { + close(release) + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + close(release) + t.Fatal(err) + } + if err := s.BackupToCompressed(filepath.Join(filepath.Dir(path), "partial.gz")); err == nil { + close(release) + t.Fatal("published an incomplete history backup") + } + // Cancel after a durable chunk, then close the real database. The next open + // must keep live data and resume from the committed source cursor. + s.historyMigration.cancel() + close(release) + <-s.historyMigration.done + if s.history != primary { + t.Fatal("import replaced the live database") + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + s, err = OpenWithBackgroundHistory(path, cold, nil) + if err != nil { + t.Fatal(err) + } + primary = s.history + var reads int + for { + select { + case <-s.historyMigration.done: + goto finished + default: + } + if _, err := s.LoadSeries("meter", "power", 0, 10000, 48); err != nil { + t.Fatal(err) + } + if err := s.EnqueueTelemetryTick(nil, []Sample{{TsMs: 1000000 + int64(reads), Driver: "live", Metric: "power", Value: float64(reads)}}, nil); err != nil { + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + reads++ + if ctx.Err() != nil { + t.Fatal(ctx.Err()) + } + } +finished: + if s.history != primary { + t.Fatal("background import replaced the primary database") + } + st := s.HistoryMigrationStatus() + if !st.HistoryComplete || st.State != "complete" || st.FilesDone != 1 { + t.Fatalf("migration=%+v", st) + } + if reads == 0 { + t.Fatal("concurrent reader/writer did not run") + } + got, err := s.LoadSeries("meter", "power", 0, 10000, 0) + if err != nil { + t.Fatal(err) + } + want := map[int64]float64{} + for i := 1; i <= historyImportRows*3+11; i++ { + want[int64(i)] = float64(i) / 7 + } + for _, p := range points { + if _, ok := want[p.TsMs]; !ok { + want[p.TsMs] = p.Value + } + } + want[3000], want[8000] = 900, 901 + if len(got) != len(want) { + t.Fatalf("got %d samples, want %d", len(got), len(want)) + } + for _, p := range got { + if expected, ok := want[p.TsMs]; !ok || historyFloatBits(expected) != historyFloatBits(p.Value) { + t.Fatalf("sample %+v expected %.17g", p, expected) + } + } + if newest, err := s.LatestSample("new meter", "new power"); err != nil || newest.Value != 902 { + t.Fatalf("live catalog/sample lost: %+v %v", newest, err) + } + if _, err := os.Stat(filepath.Join(day, "01.parquet")); err != nil { + t.Fatal("original Parquet source removed", err) + } + if s.HistoryWriterStatus().Rejected != 0 { + t.Fatalf("writer=%+v", s.HistoryWriterStatus()) + } +} + +func TestBackgroundHistoryFailureKeepsCoreStoreUsable(t *testing.T) { + path, cold := legacyMigrationFixture(t, 3) + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + file := filepath.Join(day, "01.parquet") + if err := os.WriteFile(file, []byte("invalid parquet"), 0600); err != nil { + t.Fatal(err) + } + s, err := OpenWithBackgroundHistory(path, cold, nil) + if err != nil { + t.Fatal(err) + } + defer s.Close() + select { + case <-s.historyMigration.done: + case <-time.After(10 * time.Second): + t.Fatal("import did not stop") + } + st := s.HistoryMigrationStatus() + if st.State != "failed" || st.HistoryComplete || st.LastError == "" { + t.Fatalf("migration=%+v", st) + } + if err := s.RecordSamples([]Sample{{Driver: "live", Metric: "power", TsMs: 99, Value: 42}}); err != nil { + t.Fatal(err) + } + if got, err := s.LatestSample("live", "power"); err != nil || got.Value != 42 { + t.Fatalf("live store=%+v %v", got, err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + offline, err := OpenBackupSource(path) + if err != nil { + t.Fatal(err) + } + defer offline.Close() + if err := offline.BackupToCompressed(filepath.Join(filepath.Dir(path), "partial-offline.gz")); err == nil { + t.Fatal("offline export omitted unfinished history") + } + if _, err := os.Stat(file); err != nil { + t.Fatal(fmt.Errorf("source was removed: %w", err)) + } +} + +func TestBackgroundParquetResumesAfterCommittedChunk(t *testing.T) { + path, cold := legacyMigrationFixture(t, 0) + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + file := filepath.Join(day, "01.parquet") + points := make([]parquetSampleRow, historyImportRows*2+19) + for i := range points { + points[i] = parquetSampleRow{TsMs: int64(i + 1), Driver: "meter", Metric: "power", Value: float64(i)} + } + if err := writeParquetDay(file, points); err != nil { + t.Fatal(err) + } + reached, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + s, err := OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "parquet" && st.CurrentSourceRowsDone == historyImportRows { + once.Do(func() { close(reached); <-release }) + } + }) + if err != nil { + t.Fatal(err) + } + select { + case <-reached: + case <-time.After(10 * time.Second): + close(release) + s.Close() + t.Fatal("no committed Parquet chunk") + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := s.EnqueueTelemetryTick(nil, []Sample{{TsMs: 3000, Driver: "meter", Metric: "power", Value: 99999}}, nil); err != nil { + close(release) + s.Close() + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + close(release) + s.Close() + t.Fatal(err) + } + s.historyMigration.cancel() + close(release) + if err := s.Close(); err != nil { + t.Fatal(err) + } + resumed := int64(-1) + s, err = OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "parquet" && st.CurrentSourceRowsTotal > 0 && resumed < 0 { + resumed = st.CurrentSourceRowsDone + } + }) + if err != nil { + t.Fatal(err) + } + defer s.Close() + select { + case <-s.historyMigration.done: + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + if !s.HistoryMigrationStatus().HistoryComplete || resumed != historyImportRows { + t.Fatalf("resumed=%d status=%+v", resumed, s.HistoryMigrationStatus()) + } + got, err := s.LoadSeries("meter", "power", 0, 10000, 0) + if err != nil || len(got) != len(points) { + t.Fatalf("rows=%d %v", len(got), err) + } + for i, p := range got { + want := float64(i) + if p.TsMs == 3000 { + want = 99999 + } + if p.Value != want { + t.Fatalf("sample %+v expected %v", p, want) + } + } +} + +func TestBackgroundManifestDetectsSourceLostBeforeFirstChunk(t *testing.T) { + path, cold := legacyMigrationFixture(t, 1) + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + file := filepath.Join(day, "01.parquet") + if err := writeParquetDay(file, []parquetSampleRow{{TsMs: 99, Driver: "meter", Metric: "power", Value: 42}}); err != nil { + t.Fatal(err) + } + reached, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + s, err := OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "sqlite" { + once.Do(func() { close(reached); <-release }) + } + }) + if err != nil { + t.Fatal(err) + } + defer s.Close() + <-reached + if err := os.Remove(file); err != nil { + close(release) + t.Fatal(err) + } + close(release) + select { + case <-s.historyMigration.done: + case <-time.After(10 * time.Second): + t.Fatal("missing source did not stop import") + } + if st := s.HistoryMigrationStatus(); st.State != "failed" || st.HistoryComplete { + t.Fatalf("lost source reported complete: %+v", st) + } +} + +func TestHistoricalImportRemovesOnlyOwnedAbandonedStaging(t *testing.T) { + s := freshStore(t) + abandoned := s.historyPath + ".import-abandoned" + ordinary := s.historyPath + ".original-source" + for _, dir := range []string{abandoned, ordinary} { + if err := os.MkdirAll(dir, 0700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "marker"), []byte("keep source"), 0600); err != nil { + t.Fatal(err) + } + } + cold := t.TempDir() + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + source := filepath.Join(day, "01.parquet") + if err := writeParquetDay(source, []parquetSampleRow{{TsMs: 1, Driver: "meter", Metric: "power", Value: 42}}); err != nil { + t.Fatal(err) + } + if err := s.ImportLegacyParquet(context.Background(), cold); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(abandoned); !os.IsNotExist(err) { + t.Fatalf("abandoned staging still exists: %v", err) + } + for _, path := range []string{source, filepath.Join(ordinary, "marker"), s.mainDBPath, s.historyPath} { + if _, err := os.Stat(path); err != nil { + t.Fatalf("removed retained source %s: %v", path, err) + } + } + var receipts int + if err := s.history.QueryRow(`SELECT COUNT(*) FROM history_parquet_sources`).Scan(&receipts); err != nil || receipts != 1 { + t.Fatalf("receipt=%d %v", receipts, err) + } +} + +func TestHistoricalImportSurvivesAbruptProcessExit(t *testing.T) { + if phase := os.Getenv("FTW_MIGRATION_CRASH_PHASE"); phase != "" { + ready := make(chan struct{}) + s, err := OpenWithBackgroundHistory(os.Getenv("FTW_MIGRATION_CRASH_DB"), os.Getenv("FTW_MIGRATION_CRASH_COLD"), func(st HistoryMigrationStatus) { + if st.State == "running" { + <-ready + } + if st.Phase == phase && st.CurrentSource != "" && st.CurrentSourceRowsDone >= historyImportRows { + os.Exit(23) + } + }) + if err != nil { + t.Fatal(err) + } + if err := s.RecordSamples([]Sample{{TsMs: 999999, Driver: "live before crash", Metric: "power", Value: 123}}); err != nil { + t.Fatal(err) + } + close(ready) + <-s.historyMigration.done + t.Fatal("child did not exit after a committed chunk") + } + for _, phase := range []string{"sqlite", "parquet"} { + t.Run(phase, func(t *testing.T) { + path, cold := legacyMigrationFixture(t, historyImportRows*2+9) + if phase == "parquet" { + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + points := make([]parquetSampleRow, historyImportRows*2+9) + for i := range points { + points[i] = parquetSampleRow{TsMs: int64(10000 + i), Driver: "archive", Metric: "power", Value: float64(i)} + } + if err := writeParquetDay(filepath.Join(day, "01.parquet"), points); err != nil { + t.Fatal(err) + } + } + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + child := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestHistoricalImportSurvivesAbruptProcessExit$") + child.Env = append(os.Environ(), "FTW_MIGRATION_CRASH_PHASE="+phase, "FTW_MIGRATION_CRASH_DB="+path, "FTW_MIGRATION_CRASH_COLD="+cold) + out, err := child.CombinedOutput() + var exit *exec.ExitError + if !errors.As(err, &exit) || exit.ExitCode() != 23 { + t.Fatalf("child did not stop at durable progress: %v %s", err, out) + } + s, err := OpenWithBackgroundHistory(path, cold, nil) + if err != nil { + t.Fatal(err) + } + defer s.Close() + select { + case <-s.historyMigration.done: + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + if !s.HistoryMigrationStatus().HistoryComplete { + t.Fatalf("recovery=%+v", s.HistoryMigrationStatus()) + } + if p, err := s.LatestSample("live before crash", "power"); err != nil || p.Value != 123 { + t.Fatalf("committed live data lost: %+v %v", p, err) + } + got, err := s.LoadSeries("meter", "power", 0, 10000, 0) + if err != nil || len(got) != historyImportRows*2+9 { + t.Fatalf("SQLite rows=%d %v", len(got), err) + } + for _, p := range got { + if p.Value != float64(p.TsMs)/7 { + t.Fatalf("SQLite value changed: %+v", p) + } + } + if phase == "parquet" { + got, err := s.LoadSeries("archive", "power", 10000, 20000, 0) + if err != nil || len(got) != historyImportRows*2+9 { + t.Fatalf("archive rows=%d %v", len(got), err) + } + for _, p := range got { + if p.Value != float64(p.TsMs-10000) { + t.Fatalf("archive value changed: %+v", p) + } + } + } + }) + } +} + +func TestRawRetentionWaitsForHistoryImport(t *testing.T) { + path, cold := legacyMigrationFixture(t, historyImportRows+7) + reached, release := make(chan struct{}), make(chan struct{}) + var once sync.Once + s, err := OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "sqlite" && st.CurrentSourceRowsDone == historyImportRows { + once.Do(func() { close(reached); <-release }) + } + }) + if err != nil { + t.Fatal(err) + } + var released sync.Once + defer func() { released.Do(func() { close(release) }); s.Close() }() + select { + case <-reached: + case <-time.After(10 * time.Second): + t.Fatal("import did not reach committed chunk") + } + if err := s.PruneHistorySamples(context.Background(), 1, time.Now()); err != nil { + t.Fatal(err) + } + var n int + if err := s.history.QueryRow(`SELECT COUNT(*) FROM ts_samples`).Scan(&n); err != nil || n != historyImportRows { + t.Fatalf("retention deleted pending import rows: %d %v", n, err) + } + released.Do(func() { close(release) }) + select { + case <-s.historyMigration.done: + case <-time.After(10 * time.Second): + t.Fatal("import did not finish") + } + if !s.HistoryMigrationStatus().HistoryComplete { + t.Fatalf("import failed: %+v", s.HistoryMigrationStatus()) + } + if err := s.history.QueryRow(`SELECT COUNT(*) FROM ts_samples`).Scan(&n); err != nil || n != historyImportRows+7 { + t.Fatalf("import lost rows: %d %v", n, err) + } + if err := s.PruneHistorySamples(context.Background(), 1, time.Now()); err != nil { + t.Fatal(err) + } + if err := s.history.QueryRow(`SELECT COUNT(*) FROM ts_samples`).Scan(&n); err != nil || n != 0 { + t.Fatalf("retention did not resume: %d %v", n, err) + } +} + +func TestBackgroundHistoryContinuesBetaOneReceiptsAndSequences(t *testing.T) { + path, cold := legacyMigrationFixture(t, 17) + s, err := Open(path) // beta.1 completes SQLite before importing Parquet. + if err != nil { + t.Fatal(err) + } + defer func() { s.Close() }() + generation, err := s.historyConfig("history_duckdb_generation") + if err != nil || generation == "" { + t.Fatalf("generation=%q %v", generation, err) + } + day := filepath.Join(cold, "2026", "01") + if err := os.MkdirAll(day, 0700); err != nil { + t.Fatal(err) + } + first := filepath.Join(day, "01.parquet") + if err := writeParquetDay(first, []parquetSampleRow{{TsMs: 100, Driver: "meter", Metric: "power", Value: 100}}); err != nil { + t.Fatal(err) + } + if err := s.ImportLegacyParquet(context.Background(), cold); err != nil { + t.Fatal(err) + } + second := filepath.Join(day, "02.parquet") + points := []parquetSampleRow{{TsMs: 200, Driver: "meter", Metric: "power", Value: 200}, {TsMs: 201, Driver: "meter", Metric: "power", Value: 201}} + if err := writeParquetDay(second, points); err != nil { + t.Fatal(err) + } + digest, err := historyFileHash(second) + if err != nil { + t.Fatal(err) + } + if err := s.RecordSamples([]Sample{{TsMs: 200, Driver: "meter", Metric: "power", Value: 999}}); err != nil { + t.Fatal(err) + } + // beta.1 has a bound file hash, but no per-file resume cursor. Its table + // defaults depend on the seeded sequences; leave that catalog unchanged. + for _, stmt := range []string{ + `ALTER TABLE ts_drivers ALTER COLUMN id SET DEFAULT nextval('ts_drivers_next_id')`, + `ALTER TABLE ts_metrics ALTER COLUMN id SET DEFAULT nextval('ts_metrics_next_id')`, + `DROP TABLE history_parquet_progress`, + `DROP TABLE history_parquet_manifest`, + `DROP TABLE history_sqlite_progress`, + } { + if _, err := s.history.Exec(stmt); err != nil { + t.Fatal(err) + } + } + if _, err := s.history.Exec(`INSERT INTO history_parquet_imports VALUES (?,?)`, second, digest); err != nil { + t.Fatal(err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + seedRepeated := false + s, err = OpenWithBackgroundHistory(path, cold, func(st HistoryMigrationStatus) { + if st.Phase == "seed" && st.CurrentSource != "" { + seedRepeated = true + } + }) + if err != nil { + t.Fatal(err) + } + select { + case <-s.historyMigration.done: + case <-time.After(10 * time.Second): + t.Fatal("legacy import did not finish") + } + st := s.HistoryMigrationStatus() + if seedRepeated || !st.HistoryComplete || st.FilesDone != 2 || st.RowsDone != 20 { + t.Fatalf("seedRepeated=%v status=%+v", seedRepeated, st) + } + if got, err := s.historyConfig("history_duckdb_generation"); err != nil || got != generation { + t.Fatalf("generation changed: %q %v", got, err) + } + if p, err := s.LatestSample("meter", "power"); err != nil || p.Value != 201 { + t.Fatalf("missing remaining source: %+v %v", p, err) + } + got, err := s.LoadSeries("meter", "power", 200, 200, 0) + if err != nil || len(got) != 1 || got[0].Value != 999 { + t.Fatalf("overwrote prior primary: %+v %v", got, err) + } + if err := s.RecordSamples([]Sample{{TsMs: 300, Driver: "new meter", Metric: "new metric", Value: 123}}); err != nil { + t.Fatal(err) + } + var d, m int64 + if err := s.history.QueryRow(`SELECT driver_id,metric_id FROM ts_samples WHERE ts_ms=300`).Scan(&d, &m); err != nil || d <= 7 || m <= 9 { + t.Fatalf("reused seeded IDs: %d %d %v", d, m, err) + } + if got, err := historyFileHash(second); err != nil || got != digest { + t.Fatalf("changed original source: %s %v", got, err) + } +} diff --git a/go/internal/state/history_parquet_import.go b/go/internal/state/history_parquet_import.go index affad67ae..ac6568c3b 100644 --- a/go/internal/state/history_parquet_import.go +++ b/go/internal/state/history_parquet_import.go @@ -11,6 +11,7 @@ import ( "math" "os" "path/filepath" + "strings" duckdb "github.com/duckdb/duckdb-go/v2" "github.com/parquet-go/parquet-go" @@ -18,104 +19,124 @@ import ( const historyImportRows = 2048 -// ImportLegacyParquet imports frozen daily files once. Existing recent samples -// win overlap. The source files remain as evidence after verification. -// Call only during startup, before readers or telemetry producers can use the -// Store. Production uses OpenWithLegacyHistory to enforce that lifecycle. +// ImportLegacyParquet imports frozen files while the primary remains open. +// Native instances rotate only after all active connections have closed. +// Existing primary samples win overlap. func (s *Store) ImportLegacyParquet(ctx context.Context, coldDir string) error { - if s.HistoryWriterStatus().Accepted != 0 { - return errors.New("legacy history import must finish before telemetry starts") + s.historyImportMu.Lock() + defer s.historyImportMu.Unlock() + // Only this importer owns these disposable directories. A killed process + // may leave one behind; it contains no authoritative rows or receipts. + entries, err := os.ReadDir(filepath.Dir(s.historyPath)) + if err != nil { + return err } - if coldDir == "" { - var pending int - if err := s.history.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_imports`).Scan(&pending); err != nil { - return err - } - if pending != 0 { - return errors.New("an interrupted Parquet import requires its original cold directory") + for _, entry := range entries { + if entry.IsDir() && strings.HasPrefix(entry.Name(), filepath.Base(s.historyPath)+".import-") { + if err := os.RemoveAll(filepath.Join(filepath.Dir(s.historyPath), entry.Name())); err != nil { + return err + } } - return nil } - paths, err := filepath.Glob(filepath.Join(coldDir, "[0-9][0-9][0-9][0-9]", "[0-9][0-9]", "[0-9][0-9].parquet")) - if err != nil { + if err := s.bindLegacyParquetSources(coldDir); err != nil { return err } - s.ts.allocMu.Lock() - defer s.ts.allocMu.Unlock() - s.historyWriteMu.Lock() - defer s.historyWriteMu.Unlock() - // An interrupted import may have created new catalog entries too. - defer func() { s.ts.mu.Lock(); s.ts.loaded = false; s.ts.mu.Unlock() }() - conn, err := s.history.Conn(ctx) + rows, err := s.history.QueryContext(ctx, `SELECT m.path,COALESCE(s.rows,0) FROM history_parquet_manifest m LEFT JOIN history_parquet_sources s ON s.path=m.path ORDER BY m.path`) if err != nil { return err } - defer func() { - if conn != nil { - conn.Close() - } - }() - reopen := func() error { - if conn != nil { - if err := conn.Close(); err != nil { - return err - } - conn = nil - } - if err := s.history.Close(); err != nil { - return err - } - s.history = nil - if err := s.openHistory(); err != nil { + paths := []string{} + var completedRows int64 + for rows.Next() { + var path string + var count int64 + if err := rows.Scan(&path, &count); err != nil { + rows.Close() return err } - conn, err = s.history.Conn(ctx) + paths = append(paths, path) + completedRows += count + } + if err := errors.Join(rows.Err(), rows.Close()); err != nil { return err } - imported := false + if s.historyMigration != nil { + s.historyMigration.update(func(st *HistoryMigrationStatus) { + st.Phase = "parquet" + st.FilesTotal = len(paths) + st.CurrentSource = "" + st.RowsDone += completedRows + st.RowsTotal = 0 + }) + } for _, path := range paths { - abs, err := filepath.Abs(path) - if err != nil { + if err := s.yieldHistoryImport(ctx); err != nil { return err } - var complete int - if err := conn.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_sources WHERE path=?`, abs).Scan(&complete); err != nil { + abs, err := filepath.Abs(path) + if err != nil { return err } - if complete == 0 { - // CHECKPOINT releases dirty segments and ART buffers, but DuckDB - // can retain other table buffers for the native instance's life. - // Each new file starts a fresh session on the same durable primary. - if err := reopen(); err != nil { - return fmt.Errorf("reopen history before import: %w", err) - } - imported = true - } - if err := importHistoryFile(ctx, conn, abs); err != nil { + if err := s.importHistoryFile(ctx, abs, coldDir); err != nil { return fmt.Errorf("import cold history %s: %w", abs, err) } + if s.historyMigration != nil { + s.historyMigration.update(func(st *HistoryMigrationStatus) { st.FilesDone++ }) + } } var pending int - if err := conn.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_imports`).Scan(&pending); err != nil { + if err := s.history.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_imports`).Scan(&pending); err != nil { return err } if pending != 0 { - return errors.New("an interrupted Parquet source is missing; restore the original source before starting") - } - if imported { - return reopen() + return errors.New("an interrupted Parquet source is missing; restore its original source to resume historical import") } return nil } -func importHistoryFile(ctx context.Context, conn *sql.Conn, path string) error { +func (s *Store) bindLegacyParquetSources(coldDir string) error { + if coldDir == "" { + return nil + } + paths, err := filepath.Glob(filepath.Join(coldDir, "[0-9][0-9][0-9][0-9]", "[0-9][0-9]", "[0-9][0-9].parquet")) + if err != nil { + return err + } + s.historyWriteMu.Lock() + defer s.historyWriteMu.Unlock() + tx, err := s.history.Begin() + if err != nil { + return err + } + defer tx.Rollback() + for _, path := range paths { + abs, err := filepath.Abs(path) + if err != nil { + return err + } + if _, err := tx.Exec(`INSERT INTO history_parquet_manifest VALUES (?) ON CONFLICT DO NOTHING`, abs); err != nil { + return err + } + } + return tx.Commit() +} + +func (s *Store) importHistoryFile(ctx context.Context, path, coldDir string) error { digest, err := historyFileHash(path) if err != nil { + // A verified file may have been removed after a complete backup. Its + // receipt still proves coverage; unimported missing files stay errors. + if errors.Is(err, os.ErrNotExist) { + var complete int + if checkErr := s.history.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_parquet_sources WHERE path=?`, path).Scan(&complete); checkErr == nil && complete != 0 { + return nil + } + } return err } for _, table := range []string{"history_parquet_sources", "history_parquet_imports"} { var prior string - err := conn.QueryRowContext(ctx, `SELECT sha256 FROM `+table+` WHERE path=?`, path).Scan(&prior) + err := s.history.QueryRowContext(ctx, `SELECT sha256 FROM `+table+` WHERE path=?`, path).Scan(&prior) if err == nil { if prior != digest { return errors.New("previously imported or pending Parquet source changed") @@ -127,26 +148,56 @@ func importHistoryFile(ctx context.Context, conn *sql.Conn, path string) error { return err } } - // Release buffers from prior committed files before staging another day. - if _, err := conn.ExecContext(ctx, `CHECKPOINT`); err != nil { - return fmt.Errorf("checkpoint before staging: %w", err) + if s.historyMigration != nil { + name, _ := filepath.Rel(coldDir, path) + s.historyMigration.update(func(st *HistoryMigrationStatus) { + st.CurrentSource = filepath.ToSlash(name) + st.CurrentSourceRowsDone = 0 + st.CurrentSourceRowsTotal = 0 + st.RowsTotal = 0 + }) + } + if err := s.rotateHistory(ctx); err != nil { + return err + } + // This instance owns all file-sized buffers and may spill to disk. Closing + // it cannot close the live database or any reader's connection. + dir, err := os.MkdirTemp(filepath.Dir(s.historyPath), filepath.Base(s.historyPath)+".import-") + if err != nil { + return err + } + defer os.RemoveAll(dir) + db, err := sql.Open("duckdb", filepath.Join(dir, "staging.duckdb")+"?threads=1&memory_limit=64MB&max_temp_directory_size=512MB&autoload_known_extensions=false&autoinstall_known_extensions=false") + if err != nil { + return err } - // This table can spill to disk. A file-sized unique index cannot, so detect - // duplicate keys with a spillable ordered window before primary writes. - if _, err := conn.ExecContext(ctx, `CREATE TEMP TABLE history_import_source ( - ts_ms BIGINT NOT NULL, driver_id BIGINT NOT NULL, metric_id BIGINT NOT NULL, - value DOUBLE NOT NULL CHECK(isfinite(value)))`); err != nil { + defer db.Close() + db.SetMaxOpenConns(1) + conn, err := db.Conn(ctx) + if err != nil { return err } - defer conn.ExecContext(context.Background(), `DROP TABLE IF EXISTS history_import_source`) - count, err := stageHistoryParquet(ctx, conn, path) + defer conn.Close() + for _, q := range []string{ + `CREATE SEQUENCE ts_drivers_id START 1`, `CREATE SEQUENCE ts_metrics_id START 1`, + `CREATE TABLE ts_drivers(id BIGINT DEFAULT nextval('ts_drivers_id'),name VARCHAR UNIQUE)`, + `CREATE TABLE ts_metrics(id BIGINT DEFAULT nextval('ts_metrics_id'),name VARCHAR UNIQUE)`, + `CREATE TEMP TABLE history_import_source(ts_ms BIGINT NOT NULL,driver_id BIGINT NOT NULL,metric_id BIGINT NOT NULL,value DOUBLE NOT NULL CHECK(isfinite(value)))`, + } { + if _, err := conn.ExecContext(ctx, q); err != nil { + return err + } + } + count, err := stageHistoryParquet(ctx, conn, path, func(int64) { + if s.historyMigration != nil { + s.historyMigration.update(func(*HistoryMigrationStatus) {}) + } + }) if err != nil { return fmt.Errorf("stage source: %w", err) } var duplicates int64 - if err := conn.QueryRowContext(ctx, `SELECT COUNT(*) FROM ( - SELECT ts_ms,LAG(ts_ms) OVER(PARTITION BY driver_id,metric_id ORDER BY ts_ms) AS previous - FROM history_import_source) WHERE ts_ms=previous`).Scan(&duplicates); err != nil { + if err := conn.QueryRowContext(ctx, `SELECT COUNT(*) FROM (SELECT ts_ms,LAG(ts_ms) OVER(PARTITION BY driver_id,metric_id ORDER BY ts_ms) AS previous FROM history_import_source) WHERE ts_ms=previous`).Scan(&duplicates); err != nil { return fmt.Errorf("validate source keys: %w", err) } if duplicates != 0 { @@ -155,27 +206,73 @@ func importHistoryFile(ctx context.Context, conn *sql.Conn, path string) error { if after, err := historyFileHash(path); err != nil || after != digest { return errors.Join(err, errors.New("Parquet changed during staging")) } - // Bind every committed chunk to immutable source bytes. A retry rechecks - // all rows, retaining the values verified by an earlier completed chunk. - if _, err := conn.ExecContext(ctx, `INSERT INTO history_parquet_imports VALUES (?,?) ON CONFLICT DO NOTHING`, path, digest); err != nil { + var offset int64 + if err := s.history.QueryRowContext(ctx, `SELECT rows_done FROM history_parquet_progress WHERE path=?`, path).Scan(&offset); err != nil && !errors.Is(err, sql.ErrNoRows) { return err } - for offset := int64(0); offset < count; offset += historyImportRows { - if err := importHistoryChunk(ctx, conn, offset, count); err != nil { + if offset > count { + return errors.New("Parquet source is shorter than its committed progress") + } + s.historyWriteMu.Lock() + _, err = s.history.ExecContext(ctx, `INSERT INTO history_parquet_imports VALUES (?,?) ON CONFLICT DO NOTHING`, path, digest) + s.historyWriteMu.Unlock() + if err != nil { + return err + } + if s.historyMigration != nil { + s.historyMigration.update(func(st *HistoryMigrationStatus) { + st.CurrentSourceRowsDone = offset + st.CurrentSourceRowsTotal = count + st.RowsDone += offset + }) + } + for offset < count { + if err := s.yieldHistoryImport(ctx); err != nil { + return err + } + end := min(offset+historyImportRows, count) + rows, err := conn.QueryContext(ctx, `SELECT p.ts_ms,d.name,m.name,p.value FROM history_import_source p JOIN ts_drivers d ON d.id=p.driver_id JOIN ts_metrics m ON m.id=p.metric_id WHERE p.rowid>=? AND p.rowid + + + +FTW is starting + +
+
FTW
+

FTW is starting

Preparing your site. Control has not started yet.

Keep the box powered. This page checks progress automatically.

+

Connecting to the box…

+

Reload status · Reloading this page does not restart the box.

+
+ + diff --git a/web/components/ftw-update-check.js b/web/components/ftw-update-check.js index 43cb9064f..fe5683232 100644 --- a/web/components/ftw-update-check.js +++ b/web/components/ftw-update-check.js @@ -32,6 +32,7 @@ import { FtwElement } from "./ftw-element.js"; import { apiFetch } from "./api-fetch.js"; +import { migrationHTML } from "../history-migration.js"; import "./ftw-modal.js"; const STATUS_POLL_MS = 2000; @@ -193,6 +194,9 @@ class FtwUpdateCheck extends FtwElement { this._statusTimer = null; this._pollAbort = null; // AbortController for in-flight status fetches this._updateStartedAt = 0; + this._bootHealth = null; + this._bootConnected = true; + this._healthInFlight = false; this.classList.add("hidden"); } @@ -205,6 +209,15 @@ class FtwUpdateCheck extends FtwElement { this._stopPolling(); } + update() { + const action = this.shadowRoot.activeElement?.dataset?.action; + super.update(); + if (action) { + const button = [...this.shadowRoot.querySelectorAll("[data-action]")].find(el => el.dataset.action === action); + if (button && !button.disabled) button.focus({ preventScroll:true }); + } + } + // ---- data ---- _check() { apiFetch("/api/version/check") @@ -236,6 +249,12 @@ class FtwUpdateCheck extends FtwElement { .then((r) => r.json().then((b) => ({ ok: r.ok, body: b }))) .then((res) => { if (!res.ok) { + if (res.body?.error === "starting") { + this._bootHealth = { status:"starting", migration:res.body.migration }; + this._startPolling(); + this.update(); + return; + } this._fail((res.body && res.body.error) || "failed to start"); return; } @@ -274,6 +293,23 @@ class FtwUpdateCheck extends FtwElement { _tick() { const signal = this._pollAbort ? this._pollAbort.signal : undefined; + if (!this._healthInFlight) { + this._healthInFlight = true; + apiFetch("/api/health", { cache:"no-store", signal:AbortSignal.timeout(8000) }) + .then(r => r.ok ? r.json() : null) + .then(health => { + if (signal?.aborted || this._phase !== "updating") return; + this._bootConnected = !!health; + if (health) this._bootHealth = health.status === "starting" ? health : null; + this.update(); + }) + .catch(() => { + if (signal?.aborted || this._phase !== "updating") return; + this._bootConnected = false; + this.update(); + }) + .finally(() => { this._healthInFlight = false; }); + } apiFetch("/api/version/update/status", { signal }) .then((r) => (r.ok ? r.json() : null)) .then((st) => { @@ -309,7 +345,7 @@ class FtwUpdateCheck extends FtwElement { const timeout = this._status && this._status.state === "snapshotting" ? SNAPSHOT_SOFT_TIMEOUT_MS : UPDATE_SOFT_TIMEOUT_MS; - if (Date.now() - timeoutStartedAt > timeout && this._phase === "updating") { + if (!this._bootHealth && Date.now() - timeoutStartedAt > timeout && this._phase === "updating") { this._status = Object.assign({}, this._status, { timed_out: true }); this.update(); } @@ -411,6 +447,13 @@ class FtwUpdateCheck extends FtwElement { } _overlayHTML() { + if (this._bootHealth) { + return `Starting FTW
+ ${migrationHTML(this._bootHealth.migration, {boot:true, connected:this._bootConnected}) || '

Core is preparing to start. Control has not started yet. Keep the box powered.

'} +

${this._bootConnected ? "The box is responding." : "Cannot reach the box. The last report may be out of date."}

+
+ Reloading this page does not restart the box.
`; + } const st = this._status || { state: "starting" }; const busy = this._phase === "updating"; const failed = this._phase === "failed"; diff --git a/web/header-status-marks.test.mjs b/web/header-status-marks.test.mjs index 96a9be3e8..152b276e4 100644 --- a/web/header-status-marks.test.mjs +++ b/web/header-status-marks.test.mjs @@ -69,11 +69,11 @@ function renderBadge({ update = false, degraded = false, connected = true, lates } describe("header status marks", () => { - it("stays disabled when slower component requests finish after the 503 gate", async () => { + it("stays disabled when slower component requests finish after the explicit feature gate", async () => { const delayed = new Map(); const fetchImpl = (url) => { if (url === "/api/version/check") { - return Promise.resolve({ status: 503, ok: false }); + return Promise.resolve({ status: 503, ok: false, json:async () => ({error:"self-update disabled"}) }); } if (url === "/api/version/update/status") { return Promise.resolve({ status: 503, ok: false }); diff --git a/web/history-migration.js b/web/history-migration.js new file mode 100644 index 000000000..d5bb8050e --- /dev/null +++ b/web/history-migration.js @@ -0,0 +1,89 @@ +// Shared by the startup page, update dialog and live dashboard. +const count = value => Math.max(0, Number.isFinite(Number(value)) ? Number(value) : 0); +const number = value => count(value).toLocaleString("en-US"); +const escape = value => String(value ?? "").replace(/[&<>"']/g, c => ({ + "&": "&", "<": "<", ">": ">", '"': """, "'": "'", +})[c]); + +export function migrationView(migration, { boot = false, connected = true, now = Date.now() } = {}) { + if (!migration || migration.history_complete === true) return null; + const failed = migration.state === "failed"; + const archive = migration.phase === "parquet"; + const seed = migration.phase === "seed"; + const total = count(archive || seed ? migration.current_source_rows_total : migration.rows_total); + const done = count(archive || seed ? migration.current_source_rows_done : migration.rows_done); + const progress = total > 0 ? { value: Math.min(done, total), max: total, label: archive ? "Current archive file" : seed ? "Current preparation step" : "History import progress" } : null; + const details = []; + if (seed && done > 0) details.push(`${number(done)}${total > 0 ? ` of ${number(total)}` : ""} saved items copied in this step`); + else if (total > 0) details.push(`${archive ? "Current archive: " : ""}${number(done)} of ${number(total)} readings imported`); + if ((archive || !total) && count(migration.rows_done) > 0) details.push(`${number(migration.rows_done)} readings imported in total`); + if (count(migration.files_total) > 0) details.push(`${number(migration.files_done)} of ${number(migration.files_total)} archive files complete`); + const from = Number(migration.incomplete_from_ms), until = Number(migration.incomplete_until_ms); + const coverage = Number.isFinite(from) && Number.isFinite(until) && from > 0 && until >= from + ? `History may be incomplete from ${new Date(from).toLocaleDateString("en-GB")} to ${new Date(until).toLocaleDateString("en-GB")}.` + : "Older history is not complete yet."; + const updated = Number(migration.updated_at_ms); + const age = Number.isFinite(updated) && updated > 0 ? Math.max(0, Math.floor((now - updated) / 1000)) : null; + const elapsed = age < 60 ? age : Math.floor(age / 60); + const unit = age < 60 ? "second" : "minute"; + const activity = age === null ? "Waiting for the first import report." + : `Last import report ${elapsed} ${unit}${elapsed === 1 ? "" : "s"} ago.`; + return { + title: !connected ? "History import status unavailable" : failed ? "History import paused" : boot ? "Preparing FTW" : "Importing older history", + description: !connected ? "The box is not responding. Showing its last import report." + : failed ? "The import needs attention. Your original history files are still kept." + : boot ? "FTW is preparing the data it needs to start. Control has not started yet." + : "Core is running while FTW imports older readings in the background.", + details: details.join(" · "), coverage, activity, progress, + error: failed ? migration.last_error || "The box could not finish importing history." : "", + guidance: failed || !connected ? "Keep the original data and backup. Reload this page to check the current status." + : "Keep the box powered. You can close this page and return later. Import progress is saved so it can resume after a restart.", + }; +} + +export function migrationHTML(migration, options) { + const view = migrationView(migration, options); + if (!view) return ""; + const style = 'style="width:100%;accent-color:var(--accent-e,#fbbf24)"'; + const progress = view.progress + ? `` + : ``; + const heading = options?.boot ? "h1" : "h3"; + return `
+ <${heading} role="status">${escape(view.title)}

${escape(view.description)}

+ ${progress}

${escape(view.details)}

${escape(view.coverage)}

+

${escape(view.activity)}

+ ${view.error ? `

${escape(view.error)}

` : ""} +

${escape(view.guidance)}

`; +} + +let lastMigration = null; +let lastHealth = null; +export function updateMigrationBanner(health) { + const current = health?.history_storage?.migration || health?.migration; + if (current && (current.history_complete === true || !lastMigration || count(current.updated_at_ms) >= count(lastMigration.updated_at_ms))) lastMigration = current; + if (health) lastHealth = health; + const view = migrationView(lastMigration, { boot: lastHealth?.status === "starting", connected: !!health }); + let banner = document.getElementById("history-import-banner"); + if (!view) { banner?.remove(); return; } + if (!banner) { + const main = document.querySelector("main"); + if (!main) return; + banner = document.createElement("section"); + banner.id = "history-import-banner"; + banner.className = "storage-banner"; + banner.setAttribute("aria-label", "History import"); + banner.style.cssText = "display:block;text-align:left;padding:16px 24px"; + banner.innerHTML = '

'; + main.parentNode.insertBefore(banner, main); + } + for (const key of ["title", "description", "details", "coverage", "activity", "error", "guidance"]) { + const element = banner.querySelector(`[data-field="${key}"]`); + if (element.textContent !== view[key]) element.textContent = view[key]; + element.hidden = !view[key]; + } + const progress = banner.querySelector("progress"); + progress.setAttribute("aria-label", view.progress?.label || "History import progress"); + if (view.progress) { progress.max = view.progress.max; progress.value = view.progress.value; } + else progress.removeAttribute("value"); +} diff --git a/web/history-migration.test.mjs b/web/history-migration.test.mjs new file mode 100644 index 000000000..5227d7c88 --- /dev/null +++ b/web/history-migration.test.mjs @@ -0,0 +1,36 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { migrationHTML, migrationView } from "./history-migration.js"; + +const running = { state:"running", phase:"parquet", history_complete:false, files_done:3, files_total:12, rows_done:6000, updated_at_ms:100000 }; +test("unknown totals stay indeterminate and incomplete history stays explicit", () => { + const view = migrationView(running, {now:130000}); + assert.equal(view.progress, null); + assert.match(view.details, /6,000 readings.*3 of 12/); + assert.match(view.coverage, /not complete yet/); + assert.equal(view.activity, "Last import report 30 seconds ago."); + assert.match(migrationHTML(running), /]*\bvalue=)[^>]*><\/progress>/); +}); +test("completed history, not a row percentage, removes the notice", () => { + assert.ok(migrationView({...running, rows_total:6000})); + assert.equal(migrationView({...running, history_complete:true}), null); +}); +test("startup uses the current seed step count before overall totals exist", () => { + const view = migrationView({state:"starting", phase:"seed", history_complete:false, rows_done:0, current_source_rows_done:65536}, {boot:true}); + assert.equal(view.details, "65,536 saved items copied in this step"); + assert.equal(view.progress, null); +}); +test("startup, live background import and lost contact make distinct claims", () => { + assert.match(migrationView(running, {boot:true}).description, /Control has not started/); + assert.match(migrationView(running).description, /Core is running/); + assert.match(migrationView(running, {connected:false}).description, /last import report/); + assert.doesNotMatch(migrationView(running, {connected:false}).description, /Core is running/); +}); +test("failure keeps incomplete coverage and escapes diagnostic text", () => { + const failed = {...running, state:"failed", last_error:'bad '}; + assert.equal(migrationView(failed).title, "History import paused"); + const html = migrationHTML(failed); + assert.doesNotMatch(html, /