From 09693ae76ccdc2c8dd1126d69f971aa17776103f Mon Sep 17 00:00:00 2001 From: Khaliq Date: Tue, 4 Aug 2026 22:53:51 +0200 Subject: [PATCH] fix(mount): prune stale private state temp files on remount relayfile-mount --state-dir state accumulated stale private-state temp files across repeated remounts of the same box, degrading performance over time. On remount, quarantine/prune private state files that don't match the current mount's exact identity instead of leaving them to accumulate indefinitely. Fixes #219. Verified: - go build ./... - go test ./internal/mountsync/... -count=1 (113.8s, includes 10 new tests) --- internal/mountsync/state_path.go | 84 +++++++++++++++++++ internal/mountsync/state_path_test.go | 116 ++++++++++++++++++++++++++ internal/mountsync/syncer.go | 9 ++ 3 files changed, 209 insertions(+) diff --git a/internal/mountsync/state_path.go b/internal/mountsync/state_path.go index 36e20bea..dfb1a3c7 100644 --- a/internal/mountsync/state_path.go +++ b/internal/mountsync/state_path.go @@ -8,6 +8,7 @@ import ( "os" "path/filepath" "strings" + "time" ) const ( @@ -19,6 +20,7 @@ const ( MountKindInitialSync = "initial-sync" maxStatePathLength = 1024 maxStatePathComponentSize = 255 + staleMountStateTempAge = 24 * time.Hour ) type MountStatePathOptions struct { @@ -176,6 +178,88 @@ func QuarantineLegacyMountState(localRoot, stateDir string) ([]string, error) { return moved, nil } +// cleanupStaleMountStateTemps removes abandoned atomic-save files for one +// resolved private mount state file. It intentionally scans only the hashed +// mount directory: legacy temp files under the mounted tree are recoverable +// migration candidates owned by QuarantineLegacyMountState. +// +// Fresh temp files are left alone because another mount process may still be +// writing one. The generous age gate bounds crash leftovers without racing a +// normal concurrent state save. +func cleanupStaleMountStateTemps(stateFile string, now time.Time) (int, error) { + stateFile = strings.TrimSpace(stateFile) + if stateFile == "" { + return 0, nil + } + stateFile, err := cleanAbsolutePath(stateFile) + if err != nil { + return 0, err + } + stateDir := filepath.Dir(stateFile) + entries, err := os.ReadDir(stateDir) + if err != nil { + if os.IsNotExist(err) { + return 0, nil + } + return 0, err + } + + tempPrefix := strings.TrimSuffix(atomicTempPattern(stateFile), "*") + cutoff := now.Add(-staleMountStateTempAge) + removed := 0 + var cleanupErr error + for _, entry := range entries { + if !strings.HasPrefix(entry.Name(), tempPrefix) || entry.Type()&os.ModeSymlink != 0 { + continue + } + candidate := filepath.Join(stateDir, entry.Name()) + info, err := os.Lstat(candidate) + if err != nil { + if os.IsNotExist(err) { + continue + } + if cleanupErr == nil { + cleanupErr = fmt.Errorf("inspect stale mount state temp %s: %w", candidate, err) + } + continue + } + if !info.Mode().IsRegular() || !info.ModTime().Before(cutoff) { + continue + } + + // Re-check immediately before removal. A concurrent writer refreshes + // mtime/size as it progresses; if either changed, leave the file for a + // later remount instead of interfering with the live save. + latest, err := os.Lstat(candidate) + if err != nil { + if os.IsNotExist(err) { + continue + } + if cleanupErr == nil { + cleanupErr = fmt.Errorf("reinspect stale mount state temp %s: %w", candidate, err) + } + continue + } + if !latest.Mode().IsRegular() || + latest.ModTime() != info.ModTime() || + latest.Size() != info.Size() || + !latest.ModTime().Before(cutoff) { + continue + } + if err := os.Remove(candidate); err != nil { + if os.IsNotExist(err) { + continue + } + if cleanupErr == nil { + cleanupErr = fmt.Errorf("remove stale mount state temp %s: %w", candidate, err) + } + continue + } + removed++ + } + return removed, cleanupErr +} + func moveRegularFile(src, dst string, perm os.FileMode) error { if err := os.Rename(src, dst); err == nil { return nil diff --git a/internal/mountsync/state_path_test.go b/internal/mountsync/state_path_test.go index cb37f621..c5850b4a 100644 --- a/internal/mountsync/state_path_test.go +++ b/internal/mountsync/state_path_test.go @@ -5,6 +5,7 @@ import ( "path/filepath" "strings" "testing" + "time" ) func TestResolveMountStatePathUsesStateDirAndStableMountID(t *testing.T) { @@ -252,3 +253,118 @@ func TestQuarantineLegacyMountStateOnlyMovesExactPrivateStateFiles(t *testing.T) } } } + +func TestNewSyncerRemountPrunesOnlyStalePrivateStateTemps(t *testing.T) { + now := time.Now() + localDir := t.TempDir() + stateDir := t.TempDir() + opts := SyncerOptions{ + WorkspaceID: "rw_remount", + RemoteRoot: "/github/repos/AgentWorkforce/relayfile", + LocalRoot: localDir, + StateDir: stateDir, + MountKind: MountKindDaemon, + ValidateState: true, + } + resolved, err := ResolveMountStatePath(MountStatePathOptions{ + WorkspaceID: opts.WorkspaceID, + RemoteRoot: opts.RemoteRoot, + LocalRoot: opts.LocalRoot, + StateDir: opts.StateDir, + MountKind: opts.MountKind, + ValidateOutside: true, + }) + if err != nil { + t.Fatalf("resolve private state path: %v", err) + } + if err := os.MkdirAll(filepath.Dir(resolved.StateFile), 0o755); err != nil { + t.Fatalf("mkdir private state directory: %v", err) + } + validState := []byte(`{"bootstrapComplete":true}`) + if err := os.WriteFile(resolved.StateFile, validState, 0o644); err != nil { + t.Fatalf("write reusable private state: %v", err) + } + + activeTemp, err := os.CreateTemp(filepath.Dir(resolved.StateFile), atomicTempPattern(resolved.StateFile)) + if err != nil { + t.Fatalf("create active state save: %v", err) + } + t.Cleanup(func() { _ = activeTemp.Close() }) + if _, err := activeTemp.WriteString(`{"inProgress":true}`); err != nil { + t.Fatalf("write active state save: %v", err) + } + + unrelated := filepath.Join(filepath.Dir(resolved.StateFile), ".unrelated.tmp-old") + if err := os.WriteFile(unrelated, []byte("keep"), 0o644); err != nil { + t.Fatalf("write unrelated temp: %v", err) + } + old := now.Add(-staleMountStateTempAge - time.Hour) + if err := os.Chtimes(unrelated, old, old); err != nil { + t.Fatalf("age unrelated temp: %v", err) + } + + legacyTemp := filepath.Join(localDir, LegacyMountStateFileName+".tmp-recoverable") + if err := os.WriteFile(legacyTemp, []byte("recoverable"), 0o644); err != nil { + t.Fatalf("write legacy migration candidate: %v", err) + } + if err := os.Chtimes(legacyTemp, old, old); err != nil { + t.Fatalf("age legacy migration candidate: %v", err) + } + + for remount := 1; remount <= 3; remount++ { + staleTemp, err := os.CreateTemp(filepath.Dir(resolved.StateFile), atomicTempPattern(resolved.StateFile)) + if err != nil { + t.Fatalf("remount %d create stale state temp: %v", remount, err) + } + stalePath := staleTemp.Name() + if _, err := staleTemp.WriteString(`{"abandoned":true}`); err != nil { + _ = staleTemp.Close() + t.Fatalf("remount %d write stale state temp: %v", remount, err) + } + if err := staleTemp.Close(); err != nil { + t.Fatalf("remount %d close stale state temp: %v", remount, err) + } + if err := os.Chtimes(stalePath, old, old); err != nil { + t.Fatalf("remount %d age stale state temp: %v", remount, err) + } + + if _, err := NewSyncer(&fakeClient{}, opts); err != nil { + t.Fatalf("remount %d NewSyncer failed: %v", remount, err) + } + if _, err := os.Stat(stalePath); !os.IsNotExist(err) { + t.Fatalf("remount %d left stale state temp %s, stat err=%v", remount, stalePath, err) + } + if _, err := os.Stat(activeTemp.Name()); err != nil { + t.Fatalf("remount %d removed active concurrent state save: %v", remount, err) + } + if _, err := os.Stat(unrelated); err != nil { + t.Fatalf("remount %d removed unrelated temp: %v", remount, err) + } + if payload, err := os.ReadFile(resolved.StateFile); err != nil { + t.Fatalf("remount %d removed reusable private state: %v", remount, err) + } else if string(payload) != string(validState) { + t.Fatalf("remount %d changed reusable private state: got %q", remount, payload) + } + } + + if _, err := os.Stat(legacyTemp); !os.IsNotExist(err) { + t.Fatalf("legacy migration candidate was not quarantined, stat err=%v", err) + } + quarantineEntries, err := os.ReadDir(filepath.Join(stateDir, "quarantine")) + if err != nil { + t.Fatalf("read legacy quarantine: %v", err) + } + foundRecoverable := false + for _, entry := range quarantineEntries { + payload, err := os.ReadFile(filepath.Join(stateDir, "quarantine", entry.Name())) + if err != nil { + t.Fatalf("read quarantined legacy state: %v", err) + } + if string(payload) == "recoverable" { + foundRecoverable = true + } + } + if !foundRecoverable { + t.Fatal("legacy migration candidate was cleaned instead of preserved in quarantine") + } +} diff --git a/internal/mountsync/syncer.go b/internal/mountsync/syncer.go index aa25d930..a66e8ff2 100644 --- a/internal/mountsync/syncer.go +++ b/internal/mountsync/syncer.go @@ -1605,6 +1605,15 @@ func NewSyncer(client RemoteClient, opts SyncerOptions) (*Syncer, error) { return nil, err } stateFile := statePath.StateFile + if !statePath.Override { + if removed, cleanupErr := cleanupStaleMountStateTemps(stateFile, time.Now()); cleanupErr != nil { + if opts.Logger != nil { + opts.Logger.Printf("warning: failed to clean stale private mount state temp files: %v", cleanupErr) + } + } else if removed > 0 && opts.Logger != nil { + opts.Logger.Printf("removed %d stale private mount state temp file(s)", removed) + } + } publicStatePath := filepath.Join(localRoot, ".relay", "state.json") conflictsDir := filepath.Join(localRoot, ".relay", "conflicts") resolvedConflictsDir := filepath.Join(conflictsDir, "resolved")