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 ORDER BY p.rowid`, offset, end)
+ if err != nil {
+ return err
+ }
+ batch := make([]Sample, 0, historyImportRows)
+ for rows.Next() {
+ var sm Sample
+ if err := rows.Scan(&sm.TsMs, &sm.Driver, &sm.Metric, &sm.Value); err != nil {
+ rows.Close()
+ return err
+ }
+ batch = append(batch, sm)
+ }
+ err = errors.Join(rows.Err(), rows.Close())
+ if err != nil {
+ return err
+ }
+ if int64(len(batch)) != end-offset {
+ return errors.New("staging row count changed")
+ }
+ if err := s.mergeHistoricalSamples(ctx, batch, func(tx *sql.Tx) error {
+ _, err := tx.ExecContext(ctx, `INSERT INTO history_parquet_progress VALUES (?,?) ON CONFLICT(path) DO UPDATE SET rows_done=excluded.rows_done`, path, end)
+ return err
+ }); err != nil {
return fmt.Errorf("import rows at %d: %w", offset, err)
}
- // Committing alone does not move all new row segments out of memory.
- // Bound that work within large files as well as between daily files.
- if (offset+historyImportRows)%(64*historyImportRows) == 0 {
- if _, err := conn.ExecContext(ctx, `CHECKPOINT`); err != nil {
- return fmt.Errorf("checkpoint imported rows: %w", err)
+ if s.historyMigration != nil {
+ s.historyMigration.update(func(st *HistoryMigrationStatus) { st.CurrentSourceRowsDone = end; st.RowsDone += end - offset })
+ }
+ offset = end
+ if offset%(64*historyImportRows) == 0 {
+ if err := s.CheckpointHistory(ctx); err != nil {
+ return err
}
}
}
if after, err := historyFileHash(path); err != nil || after != digest {
return errors.Join(err, errors.New("Parquet changed during import; restore the original source"))
}
- tx, err := conn.BeginTx(ctx, nil)
+ s.historyWriteMu.Lock()
+ defer s.historyWriteMu.Unlock()
+ tx, err := s.history.BeginTx(ctx, nil)
if err != nil {
return err
}
@@ -186,6 +283,9 @@ func importHistoryFile(ctx context.Context, conn *sql.Conn, path string) error {
if _, err := tx.ExecContext(ctx, `DELETE FROM history_parquet_imports WHERE path=?`, path); err != nil {
return err
}
+ if _, err := tx.ExecContext(ctx, `DELETE FROM history_parquet_progress WHERE path=?`, path); err != nil {
+ return err
+ }
if err := tx.Commit(); err != nil {
return err
}
@@ -193,7 +293,7 @@ func importHistoryFile(ctx context.Context, conn *sql.Conn, path string) error {
return nil
}
-func stageHistoryParquet(ctx context.Context, conn *sql.Conn, path string) (int64, error) {
+func stageHistoryParquet(ctx context.Context, conn *sql.Conn, path string, progress ...func(int64)) (int64, error) {
f, err := os.Open(path)
if err != nil {
return 0, err
@@ -305,7 +405,7 @@ func stageHistoryParquet(ctx context.Context, conn *sql.Conn, path string) (int6
return count, writeErr
}
err := conn.Raw(func(raw any) error {
- app, err := duckdb.NewAppender(raw.(driver.Conn), "temp", "main", "history_import_source")
+ app, err := duckdb.NewAppender(nativeHistoryConn(raw), "temp", "main", "history_import_source")
if err != nil {
return err
}
@@ -321,6 +421,9 @@ func stageHistoryParquet(ctx context.Context, conn *sql.Conn, path string) (int6
return count, err
}
count += int64(n)
+ for _, notify := range progress {
+ notify(count)
+ }
}
if readErr == io.EOF {
return count, nil
@@ -329,6 +432,10 @@ func stageHistoryParquet(ctx context.Context, conn *sql.Conn, path string) (int6
}
func importHistoryChunk(ctx context.Context, conn *sql.Conn, offset, total int64) error {
+ return importHistoryChunkCommit(ctx, conn, offset, total, nil)
+}
+
+func importHistoryChunkCommit(ctx context.Context, conn *sql.Conn, offset, total int64, receipt func(*sql.Tx) error) error {
tx, err := conn.BeginTx(ctx, nil)
if err != nil {
return err
@@ -387,5 +494,10 @@ func importHistoryChunk(ctx context.Context, conn *sql.Conn, offset, total int64
return err
}
}
+ if receipt != nil {
+ if err := receipt(tx); err != nil {
+ return err
+ }
+ }
return tx.Commit()
}
diff --git a/go/internal/state/history_parquet_import_test.go b/go/internal/state/history_parquet_import_test.go
index bb1cd7486..a7bfccd96 100644
--- a/go/internal/state/history_parquet_import_test.go
+++ b/go/internal/state/history_parquet_import_test.go
@@ -45,8 +45,11 @@ func TestOpenImportsLegacyHistoryBeforeTelemetryStarts(t *testing.T) {
if err := s.FlushHistory(ctx); err != nil {
t.Fatal(err)
}
- if err := s.ImportLegacyParquet(ctx, cold); err == nil {
- t.Fatal("allowed native session replacement after telemetry had started")
+ if err := s.ImportLegacyParquet(ctx, cold); err != nil {
+ t.Fatal(err)
+ }
+ if s.history != primary {
+ t.Fatal("historical import replaced the live native database")
}
}
diff --git a/go/internal/state/history_schema.go b/go/internal/state/history_schema.go
index 4a53d74e7..c3ca365d2 100644
--- a/go/internal/state/history_schema.go
+++ b/go/internal/state/history_schema.go
@@ -2,6 +2,9 @@ package state
// HistorySchema is separate from SQLite configuration and model state.
var historySchema = []string{
+ `CREATE TABLE IF NOT EXISTS history_parquet_manifest(path VARCHAR PRIMARY KEY)`,
+ `CREATE TABLE IF NOT EXISTS history_parquet_progress(path VARCHAR PRIMARY KEY, rows_done BIGINT NOT NULL)`,
+ `CREATE TABLE IF NOT EXISTS history_sqlite_progress (source VARCHAR PRIMARY KEY, rows_done BIGINT NOT NULL, driver_id BIGINT NOT NULL, metric_id BIGINT NOT NULL, ts_ms BIGINT NOT NULL)`,
`CREATE TABLE IF NOT EXISTS history_parquet_imports (path VARCHAR PRIMARY KEY, sha256 VARCHAR NOT NULL)`,
`CREATE TABLE IF NOT EXISTS history_parquet_sources (path VARCHAR PRIMARY KEY, sha256 VARCHAR NOT NULL, rows BIGINT NOT NULL, imported_at TIMESTAMP DEFAULT current_timestamp)`,
`CREATE SEQUENCE IF NOT EXISTS history_commit_sequence START 1`,
diff --git a/go/internal/state/history_writer.go b/go/internal/state/history_writer.go
index 3027696bb..151e91db5 100644
--- a/go/internal/state/history_writer.go
+++ b/go/internal/state/history_writer.go
@@ -46,6 +46,9 @@ type HistoryWriterStatus struct {
LastRejectMS int64 `json:"last_reject_ms,omitempty"`
LastRejectError string `json:"last_reject_error,omitempty"`
Stopping bool `json:"stopping"`
+ MaintenanceError string `json:"maintenance_error,omitempty"`
+ LastMaintenanceMS int64 `json:"last_maintenance_ms,omitempty"`
+ MaintenanceRuns uint64 `json:"maintenance_runs"`
}
type historyWriter struct {
@@ -57,11 +60,19 @@ type historyWriter struct {
done chan struct{}
ctx context.Context
cancel context.CancelFunc
+ // Owned by run, outside the short status mutex. Tests set limits before
+ // sending the first tick; production uses bounded rows, time and retries.
+ maintenanceRows int
+ maintenanceRowsLimit int
+ maintenanceDue time.Time
+ maintenanceRetry time.Time
+ maintenanceRetryDelay time.Duration
}
func newHistoryWriter(s *Store) *historyWriter {
ctx, cancel := context.WithCancel(context.Background())
- w := &historyWriter{store: s, queue: make(chan historyBatch, historyQueueTicks), changed: make(chan struct{}), done: make(chan struct{}), ctx: ctx, cancel: cancel}
+ w := &historyWriter{store: s, queue: make(chan historyBatch, historyQueueTicks), changed: make(chan struct{}), done: make(chan struct{}), ctx: ctx, cancel: cancel,
+ maintenanceRowsLimit: 64 * historyImportRows, maintenanceDue: time.Now().Add(time.Hour), maintenanceRetryDelay: 30 * time.Second}
go w.run()
return w
}
@@ -172,6 +183,11 @@ func (w *historyWriter) run() {
w.mu.Unlock()
if err == nil {
acknowledgedSequence = seq
+ rows := len(b.payload.Samples) + len(b.payload.Observations)
+ if b.payload.Point != nil {
+ rows++
+ }
+ w.maintainHistory(rows)
break
}
timer := time.NewTimer(time.Second)
@@ -185,6 +201,37 @@ func (w *historyWriter) run() {
}
}
+// Maintenance follows a durable commit. It never holds a catalog/write/status
+// lock, and admission can continue into the bounded queue. A long read only
+// postpones maintenance; its transaction and the new committed data stay intact.
+func (w *historyWriter) maintainHistory(rows int) {
+ w.maintenanceRows += rows
+ now := time.Now()
+ if now.Before(w.maintenanceRetry) || (w.maintenanceRows < w.maintenanceRowsLimit && now.Before(w.maintenanceDue)) {
+ return
+ }
+ ctx, cancel := context.WithTimeout(w.ctx, 2*time.Second)
+ err := w.store.CheckpointHistory(ctx)
+ cancel()
+ if err == nil {
+ w.maintenanceRows = 0
+ w.maintenanceDue = time.Now().Add(time.Hour)
+ } else {
+ w.maintenanceRetry = time.Now().Add(w.maintenanceRetryDelay)
+ slog.Warn("history maintenance postponed; committed data retained", "err", err)
+ }
+ w.mu.Lock()
+ defer w.mu.Unlock()
+ if err == nil {
+ w.status.MaintenanceError = ""
+ w.status.LastMaintenanceMS = time.Now().UnixMilli()
+ w.status.MaintenanceRuns++
+ } else {
+ w.status.MaintenanceError = "History maintenance could not finish; committed data is retained and maintenance will retry."
+ }
+ w.signal()
+}
+
func (s *Store) HistoryWriterStatus() HistoryWriterStatus {
if s.historyWriter == nil {
return HistoryWriterStatus{}
diff --git a/go/internal/state/store.go b/go/internal/state/store.go
index 54363ce17..f45b0082e 100644
--- a/go/internal/state/store.go
+++ b/go/internal/state/store.go
@@ -41,10 +41,13 @@ const (
//
// See heal.go for the boot-time integrity gate that populates healEvents.
type Store struct {
- history *sql.DB
- historyPath string
- historyWriteMu sync.Mutex
- historyWriter *historyWriter
+ historyConnector *historyConnector
+ history *sql.DB
+ historyPath string
+ historyImportMu sync.Mutex
+ historyWriteMu sync.Mutex
+ historyWriter *historyWriter
+ historyMigration *historyMigration
db *sql.DB
cache *sql.DB
@@ -76,17 +79,30 @@ type Store struct {
// then runs all migrations. The connection pragmas (WAL, synchronous(NORMAL),
// foreign_keys, busy_timeout) and a small pool live in openRaw — see heal.go.
func Open(path string) (*Store, error) {
- return openStore(path, "", false)
+ return openStore(path, "", false, nil)
}
-// OpenWithLegacyHistory finishes the cold-history import before starting the
-// writer or returning a Store to readers. Native history sessions may reopen
-// during this one-time import to release buffers retained by DuckDB.
+// OpenWithLegacyHistory is the synchronous entry point for offline tools.
+// Core uses OpenWithBackgroundHistory so raw history does not block startup.
func OpenWithLegacyHistory(path, coldDir string) (*Store, error) {
- return openStore(path, coldDir, true)
+ return openStore(path, coldDir, true, nil)
}
-func openStore(path, coldDir string, importLegacy bool) (*Store, error) {
+// OpenWithBackgroundHistory seeds the catalog and energy accounting before
+// starting telemetry. Frozen raw samples and Parquet then import in bounded
+// transactions while the same primary database serves live readers and writers.
+func OpenWithBackgroundHistory(path, coldDir string, onProgress func(HistoryMigrationStatus)) (*Store, error) {
+ m := newHistoryMigration(onProgress)
+ s, err := openStore(path, coldDir, false, m)
+ if err != nil {
+ m.cancel()
+ return nil, err
+ }
+ go s.runHistoryMigration(coldDir)
+ return s, nil
+}
+
+func openStore(path, coldDir string, importLegacy bool, migration *historyMigration) (*Store, error) {
nowMs := time.Now().UnixMilli()
cachePath := filepath.Join(filepath.Dir(path), "cache.db")
@@ -113,7 +129,7 @@ func openStore(path, coldDir string, importLegacy bool) (*Store, error) {
slog.Info("state: integrity gate complete", "elapsed", time.Since(tGate).Round(time.Millisecond))
s := &Store{
- db: db, cache: cache, ts: newInternCache(), mainDBPath: absolutePath,
+ db: db, cache: cache, ts: newInternCache(), mainDBPath: absolutePath, historyMigration: migration,
}
for _, ev := range []*HealEvent{stEv, caEv} {
if ev != nil {
@@ -159,6 +175,22 @@ func openStore(path, coldDir string, importLegacy bool) (*Store, error) {
return nil, err
}
}
+ if migration != nil {
+ // Remember all source paths before live work starts. A file that goes
+ // missing before its first chunk must not disappear from coverage.
+ if err := s.bindLegacyParquetSources(coldDir); err != nil {
+ s.history.Close()
+ db.Close()
+ cache.Close()
+ return nil, err
+ }
+ if _, err := s.history.Exec(`INSERT INTO history_migrations(name) VALUES ('legacy-import-pending') ON CONFLICT DO NOTHING`); err != nil {
+ s.history.Close()
+ db.Close()
+ cache.Close()
+ return nil, err
+ }
+ }
s.historyWriter = newHistoryWriter(s)
writeCleanMarker(path)
return s, nil
@@ -235,6 +267,10 @@ func (s *Store) Close() error {
cancel()
}
s.verifyWG.Wait()
+ if s.historyMigration != nil {
+ s.historyMigration.cancel()
+ <-s.historyMigration.done
+ }
var err error
if s.historyWriter != nil {
diff --git a/go/internal/state/store_ts.go b/go/internal/state/store_ts.go
index 4ef657a6f..333da8146 100644
--- a/go/internal/state/store_ts.go
+++ b/go/internal/state/store_ts.go
@@ -163,7 +163,7 @@ func (s *Store) driverID(name string) (int64, error) {
// a caller that raced us before hydrate finished) resolves to the same
// id rather than failing the whole sample batch.
if _, err := s.history.Exec(
- `INSERT INTO ts_drivers (name) VALUES (?) ON CONFLICT(name) DO NOTHING`, name,
+ `INSERT INTO ts_drivers (id, name) VALUES (nextval('ts_drivers_next_id'), ?) ON CONFLICT(name) DO NOTHING`, name,
); err != nil {
return 0, err
}
@@ -207,7 +207,7 @@ func (s *Store) metricID(name, unit string) (int64, error) {
// One statement covers both jobs: allocate the row, or relabel an
// existing one once the driver supplies a unit. An empty unit never
// erases a label already stored.
- if _, err := s.history.Exec(`INSERT INTO ts_metrics (name, unit) VALUES (?, NULLIF(?, ''))
+ if _, err := s.history.Exec(`INSERT INTO ts_metrics (id, name, unit) VALUES (nextval('ts_metrics_next_id'), ?, NULLIF(?, ''))
ON CONFLICT(name) DO UPDATE SET unit = COALESCE(NULLIF(excluded.unit, ''), ts_metrics.unit)`,
name, unit,
); err != nil {
@@ -655,7 +655,9 @@ func (s *Store) DriverNames() ([]string, error) {
// A nonpositive retention keeps all samples. The oldest hour is removed per
// transaction, releasing the writer between batches.
func (s *Store) PruneHistorySamples(ctx context.Context, retentionDays int, now time.Time) error {
- if retentionDays <= 0 {
+ // Import receipts refer to committed source chunks. Do not delete their rows
+ // until every source has been verified, including after a failed import.
+ if retentionDays <= 0 || !s.HistoryMigrationStatus().HistoryComplete {
return nil
}
cutoff := now.UTC().AddDate(0, 0, -retentionDays)
diff --git a/go/internal/updateipc/capabilities.go b/go/internal/updateipc/capabilities.go
new file mode 100644
index 000000000..79535a2c3
--- /dev/null
+++ b/go/internal/updateipc/capabilities.go
@@ -0,0 +1,62 @@
+// Package updateipc defines the read-only startup contract with the updater.
+package updateipc
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "io"
+ "net"
+ "net/http"
+ "time"
+)
+
+const CapabilitiesPath = "/capabilities"
+
+type Capabilities struct {
+ Protocol int `json:"protocol"`
+ PreserveCoreOnReadinessFailure bool `json:"preserve_core_on_readiness_failure"`
+}
+
+// ServeCapabilities advertises the fixed updater's failure behavior. Reading
+// it does not create a job, alter state.json or run a container command.
+func ServeCapabilities(w http.ResponseWriter, r *http.Request) {
+ w.Header().Set("Content-Type", "application/json")
+ w.Header().Set("Cache-Control", "no-store")
+ json.NewEncoder(w).Encode(Capabilities{Protocol: 1, PreserveCoreOnReadinessFailure: true})
+}
+
+// RequireSafeUpdater runs before Core opens any persistent state. Only a
+// positive capability response permits startup; an older sidecar can then
+// revert a refused Core without any new data format having been written.
+func RequireSafeUpdater(ctx context.Context, socket string) error {
+ transport := &http.Transport{DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
+ return (&net.Dialer{}).DialContext(ctx, "unix", socket)
+ }}
+ defer transport.CloseIdleConnections()
+ client := &http.Client{Transport: transport, Timeout: 2 * time.Second, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
+ for {
+ req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://unix"+CapabilitiesPath, nil)
+ if err != nil {
+ return err
+ }
+ resp, err := client.Do(req)
+ if err == nil {
+ var capabilities Capabilities
+ decodeErr := json.NewDecoder(io.LimitReader(resp.Body, 4096)).Decode(&capabilities)
+ resp.Body.Close()
+ if resp.StatusCode == http.StatusOK && decodeErr == nil && capabilities.Protocol == 1 && capabilities.PreserveCoreOnReadinessFailure {
+ return nil
+ }
+ return fmt.Errorf("update ftw-updater first: the running updater does not confirm safe history upgrades (HTTP %d); Core has not opened its data", resp.StatusCode)
+ }
+ // Compose can start Core before the sidecar has bound its socket.
+ timer := time.NewTimer(250 * time.Millisecond)
+ select {
+ case <-ctx.Done():
+ timer.Stop()
+ return fmt.Errorf("update ftw-updater first: cannot verify the updater at %s; Core has not opened its data: %w", socket, ctx.Err())
+ case <-timer.C:
+ }
+ }
+}
diff --git a/go/internal/updateipc/capabilities_test.go b/go/internal/updateipc/capabilities_test.go
new file mode 100644
index 000000000..1b85d0d55
--- /dev/null
+++ b/go/internal/updateipc/capabilities_test.go
@@ -0,0 +1,70 @@
+package updateipc
+
+import (
+ "context"
+ "net"
+ "net/http"
+ "os"
+ "path/filepath"
+ "strings"
+ "testing"
+ "time"
+)
+
+func TestRequireSafeUpdater(t *testing.T) {
+ for _, tc := range []struct {
+ name, body string
+ status int
+ accepted bool
+ }{
+ {"fixed", "", 200, true},
+ {"old", `{}`, 404, false},
+ {"missing", `{"protocol":1}`, 200, false},
+ {"unsafe", `{"protocol":1,"preserve_core_on_readiness_failure":false}`, 200, false},
+ {"unknown", `{"protocol":2,"preserve_core_on_readiness_failure":true}`, 200, false},
+ {"malformed", `{"protocol":1,`, 200, false},
+ } {
+ t.Run(tc.name, func(t *testing.T) {
+ dir, err := os.MkdirTemp("", "ftw-ipc-")
+ 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 r.Method != http.MethodGet || r.URL.Path != CapabilitiesPath {
+ t.Errorf("unexpected write request: %s %s", r.Method, r.URL.Path)
+ }
+ if tc.accepted {
+ ServeCapabilities(w, r)
+ return
+ }
+ w.WriteHeader(tc.status)
+ w.Write([]byte(tc.body))
+ })}
+ go srv.Serve(ln)
+ defer srv.Close()
+ ctx, cancel := context.WithTimeout(context.Background(), time.Second)
+ defer cancel()
+ err = RequireSafeUpdater(ctx, socket)
+ if (err == nil) != tc.accepted {
+ t.Fatalf("accepted=%v: %v", tc.accepted, err)
+ }
+ if err != nil && !strings.Contains(err.Error(), "update ftw-updater first") {
+ t.Fatal(err)
+ }
+ })
+ }
+}
+
+func TestUnavailableUpdaterIsBounded(t *testing.T) {
+ ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
+ defer cancel()
+ if err := RequireSafeUpdater(ctx, filepath.Join(t.TempDir(), "missing")); err == nil || !strings.Contains(err.Error(), "Core has not opened its data") {
+ t.Fatalf("unavailable updater accepted: %v", err)
+ }
+}
diff --git a/web/app.js b/web/app.js
index 2a3ed3609..d94ab9698 100644
--- a/web/app.js
+++ b/web/app.js
@@ -4,6 +4,11 @@
"use strict";
const POLL_INTERVAL = 2000; // status poll cadence — snappier cards
+ const historyMigrationUI = import("/history-migration.js").catch(function () { return null; });
+ function updateHistoryMigration(health) {
+ historyMigrationUI.then(function (ui) { if (ui) ui.updateMigrationBanner(health); }).catch(function () {});
+ if (updateBadge && typeof updateBadge.setBootHealth === "function") updateBadge.setBootHealth(health);
+ }
// Prices arrive as minor units per kWh; what to call them depends on the
// configured currency. window.FTWUnits is set when
@@ -2200,13 +2205,21 @@
function fetchStatus() {
return Promise.all([
boundedApiRead("/api/status", function (r) {
- if (!r.ok) throw new Error("HTTP " + r.status);
+ if (!r.ok) return r.json().catch(function () { return {}; }).then(function (body) {
+ var error = new Error("HTTP " + r.status);
+ error.starting = body.error === "starting";
+ throw error;
+ });
return r.json();
}),
boundedApiRead("/api/loadpoints", function (r) { return r.ok ? r.json() : null; })
.catch(function () { return null; }),
boundedApiRead("/api/health", function (r) { return r.ok ? r.json() : null; })
- .catch(function () { return null; }),
+ .catch(function () { return null; })
+ .then(function (health) {
+ updateHistoryMigration(health);
+ return health;
+ }),
])
.then(function (results) {
var data = results[0];
@@ -2243,7 +2256,7 @@
console.warn("status fetch failed:", e);
updateChargingNotice(null);
setConnected(false);
- if (firstLoad) { showSetupBanner(); }
+ if (firstLoad && !e.starting) { showSetupBanner(); }
});
}
diff --git a/web/boot.html b/web/boot.html
new file mode 100644
index 000000000..4d4cd94bc
--- /dev/null
+++ b/web/boot.html
@@ -0,0 +1,49 @@
+
+
+
+
+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.