Skip to content
Draft
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
84 changes: 84 additions & 0 deletions internal/mountsync/state_path.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"os"
"path/filepath"
"strings"
"time"
)

const (
Expand All @@ -19,6 +20,7 @@ const (
MountKindInitialSync = "initial-sync"
maxStatePathLength = 1024
maxStatePathComponentSize = 255
staleMountStateTempAge = 24 * time.Hour
)

type MountStatePathOptions struct {
Expand Down Expand Up @@ -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
Expand Down
116 changes: 116 additions & 0 deletions internal/mountsync/state_path_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"path/filepath"
"strings"
"testing"
"time"
)

func TestResolveMountStatePathUsesStateDirAndStableMountID(t *testing.T) {
Expand Down Expand Up @@ -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")
}
}
9 changes: 9 additions & 0 deletions internal/mountsync/syncer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading