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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/clear-history-import-progress.md
Original file line number Diff line number Diff line change
@@ -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.
9 changes: 9 additions & 0 deletions .changeset/safe-background-history-import.md
Original file line number Diff line number Diff line change
@@ -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.
77 changes: 72 additions & 5 deletions docs/self-update.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
131 changes: 20 additions & 111 deletions go/cmd/ftw-updater/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@ import (
"time"

"gopkg.in/yaml.v3"

"github.com/srcfl/ftw/go/internal/updateipc"
)

const (
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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 {
Expand All @@ -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)
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}
}
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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"
}
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}
Expand Down
62 changes: 42 additions & 20 deletions go/cmd/ftw-updater/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Expand Down Expand Up @@ -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()})
Expand Down Expand Up @@ -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)
}
}
Loading