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
75 changes: 9 additions & 66 deletions internal/cron/lock.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,31 +5,23 @@ import (
"fmt"
"os"
"path/filepath"
"sync/atomic"
"time"

"github.com/Gitlawb/zero/internal/lockutil"
)

// Cross-process lock tuning. The lock is held only for a single metadata
// read-modify-write (milliseconds), never across a job's exec, so the timeout is
// generous and the stale threshold sits far above any real hold.
// read-modify-write (milliseconds), never across a job's exec, so the timeout
// is generous relative to any real hold.
const (
cronLockTimeout = 10 * time.Second
cronLockStaleAfter = 60 * time.Second
cronLockRetryDelay = 20 * time.Millisecond
)

var cronLockSeq atomic.Uint64

// lockJob takes a cross-process exclusive lock for one job by O_EXCL-creating a
// sibling "<id>.lock" file next to the job directory (so Remove's RemoveAll of
// the job dir never deletes a live lock). It serializes the read-modify-write of
// a job's metadata across concurrent schedulers and commands. The lock file is
// removed on release; a stale lock from a crashed holder (older than
// cronLockStaleAfter) is reclaimed. Wall-clock time.Now is used deliberately
// (not the injectable Store.now) so lock timing never depends on a frozen test
// clock and the stale check compares against real file mtimes.
// lockJob takes a cross-process advisory lock for one job on a stable sibling
// "<id>.lock" file next to the job directory. It serializes the read-modify-write
// of a job's metadata across concurrent schedulers and commands. The kernel
// releases the lock when a holder exits; the file is never renamed or removed.
func (s *Store) lockJob(id string) (func(), error) {
if !validID(id) {
return nil, fmt.Errorf("invalid cron job id %q", id)
Expand All @@ -38,64 +30,15 @@ func (s *Store) lockJob(id string) (func(), error) {
if err := os.MkdirAll(filepath.Dir(lockPath), 0o700); err != nil {
return nil, err
}
token := fmt.Sprintf("%d-%d-%d", os.Getpid(), time.Now().UnixNano(), cronLockSeq.Add(1))
deadline := time.Now().Add(cronLockTimeout)
for {
f, err := os.OpenFile(lockPath, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
lock, err := lockutil.TryAcquireFileLock(lockPath)
if err == nil {
// A partial write would leave a lock file without our token, so the
// releaser could never delete it — stranding the lock. Fail closed.
if _, werr := f.WriteString(token); werr != nil {
_ = f.Close()
_ = lockutil.RemoveLockFile(lockPath)
return nil, fmt.Errorf("cron: write job lock: %w", werr)
}
if cerr := f.Close(); cerr != nil {
_ = lockutil.RemoveLockFile(lockPath)
return nil, fmt.Errorf("cron: close job lock: %w", cerr)
}
var released bool
return func() {
if released {
return
}
released = true
if data, rerr := os.ReadFile(lockPath); rerr == nil && string(data) == token {
_ = lockutil.RemoveLockFile(lockPath)
}
}, nil
return func() { _ = lock.Release() }, nil
}
// On Windows a concurrent holder's os.Remove leaves the lock file in a
// "delete pending" state, so an O_EXCL create races it with
// ERROR_ACCESS_DENIED (os.ErrPermission) rather than ErrExist. Treat that
// as contention and retry, exactly like ErrExist — otherwise the lock
// spuriously fails under concurrency on Windows.
if !errors.Is(err, os.ErrExist) && !errors.Is(err, os.ErrPermission) {
if !errors.Is(err, lockutil.ErrLockHeld) {
return nil, fmt.Errorf("cron: acquire job lock: %w", err)
}
// Reclaim a stale lock left by a crashed holder — atomically (H3). A blind
// Remove lets two racers both "reclaim" and recreate, so both hold the lock;
// reclaimStaleLock renames the file aside (only one rename of a given file
// wins) and restores it if it turns out fresh, so a live lock is never deleted
// out from under its holder.
if info, statErr := os.Stat(lockPath); statErr == nil && time.Since(info.ModTime()) > cronLockStaleAfter {
cleared, rerr := lockutil.ReclaimStaleLock(lockPath, token, func(reclaimedPath string) bool {
info, err := os.Stat(reclaimedPath)
return err == nil && time.Since(info.ModTime()) <= cronLockStaleAfter
})
if rerr != nil {
// Reclaim hit a hard failure: the rename aside failed outright, or a
// live holder's lock could not be put back (the lock path may be
// missing, so re-acquiring would break mutual exclusion). Fail closed
// instead of spinning to the deadline.
return nil, fmt.Errorf("cron: reclaim stale job lock: %w", rerr)
}
if cleared {
continue // cleared a genuinely stale lock; retry the O_EXCL create now
}
// Lost the reclaim race (or it was actually fresh) — fall through to the
// bounded wait instead of hot-spinning on a reclaim that never wins (L13).
}
if !time.Now().Before(deadline) {
return nil, fmt.Errorf("cron: timed out acquiring job lock for %q", id)
}
Expand Down
2 changes: 1 addition & 1 deletion internal/cron/mutate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ func TestMutateRemovedJobReturnsNotFound(t *testing.T) {
if !errors.Is(err, ErrJobNotFound) {
t.Fatalf("Mutate on removed job = %v, want ErrJobNotFound", err)
}
// And the lock file left no litter / did not resurrect the job dir.
// The stable advisory-lock file does not resurrect the removed job dir.
if _, gerr := store.Get(job.ID); !errors.Is(gerr, ErrJobNotFound) {
t.Fatalf("job should remain removed, Get = %v", gerr)
}
Expand Down
102 changes: 35 additions & 67 deletions internal/daemon/lock.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,107 +3,60 @@ package daemon
import (
"errors"
"fmt"
"io/fs"
"os"
"strconv"
"strings"
"sync/atomic"

"github.com/Gitlawb/zero/internal/lockutil"
)

// Single-instance lock. Mirrors reference-daemon-code-agent-js/supervisor.js's
// lock file: a PID file created with O_EXCL. A second start fails; a STALE lock
// left by a dead daemon (the recorded PID is no longer alive) is reclaimed so the
// daemon recovers from an unclean shutdown without manual cleanup.
// Single-instance lock. The stable PID file carries a kernel advisory lock for
// the daemon's full lifetime. Process exit releases the lock automatically, so
// stale metadata needs no rename/remove recovery protocol.

// ErrAlreadyRunning is returned when a live daemon already holds the lock.
var ErrAlreadyRunning = errors.New("daemon: another instance is already running")

// fileLock is an acquired single-instance lock.
type fileLock struct {
path string
lock *lockutil.FileLock
}

// processAlive reports whether pid is a live process. Implemented per-platform
// (lock_posix.go / lock_windows.go). It is a package var so tests can stub it.
var processAlive = osProcessAlive

// acquireLock takes the single-instance lock at path, reclaiming a stale lock
// whose recorded PID is dead. isAlive overrides the liveness check (tests pass a
// stub); nil uses the real processAlive.
// acquireLock takes the single-instance advisory lock at path. isAlive is used
// only to enrich a contention error with the current holder's PID; the kernel
// lock, not PID metadata, is authoritative.
func acquireLock(path string, isAlive func(pid int) bool) (*fileLock, error) {
if isAlive == nil {
isAlive = processAlive
}
// At most two passes: create, or detect-stale-then-retry once.
for attempt := 0; attempt < 2; attempt++ {
f, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
if err == nil {
// A failed PID write would leave a malformed lock file that another
// process reads as stale (unparsable PID) and wrongly reclaims, breaking
// the single-instance guarantee — so on write failure, remove it and fail.
if _, werr := fmt.Fprintf(f, "%d\n", os.Getpid()); werr != nil {
_ = f.Close()
_ = lockutil.RemoveLockFile(path)
return nil, werr
}
if cerr := f.Close(); cerr != nil {
_ = lockutil.RemoveLockFile(path)
return nil, cerr
}
return &fileLock{path: path}, nil
}
if !errors.Is(err, fs.ErrExist) && !errors.Is(err, os.ErrPermission) {
lock, err := lockutil.TryAcquireFileLock(path)
if err != nil {
if !errors.Is(err, lockutil.ErrLockHeld) {
return nil, err
}
pid, perr := readPidFile(path)
if perr == nil && pid > 0 && isAlive(pid) {
return nil, fmt.Errorf("%w (pid %d)", ErrAlreadyRunning, pid)
}
// Stale lock (dead PID or unreadable) — reclaim it atomically, then retry the
// O_EXCL create. A blind Remove here races: two daemons starting at once could
// both read the stale PID, both Remove, and then one Removes the OTHER's
// freshly-created lock — leaving both "holding" the single-instance lock.
// reclaimStaleLock renames the file aside so only one racer wins the rename,
// and restores it if a live holder reacquired in the gap (D6).
if _, rerr := reclaimStaleLock(path, isAlive); rerr != nil {
// Reclaim hit a hard failure: the rename aside failed outright, or a
// live holder's lock could not be put back (the lock path may be
// missing, so re-acquiring would break the single-instance guarantee).
// Fail closed instead of spinning to the deadline.
return nil, fmt.Errorf("daemon: reclaim stale lock: %w", rerr)
}
return nil, ErrAlreadyRunning
}
return nil, ErrAlreadyRunning
}

// daemonLockSeq makes each reclaim attempt's sidelined filename unique per process.
var daemonLockSeq atomic.Uint64

// reclaimStaleLock atomically reclaims a single-instance lock whose recorded
// PID is dead, via lockutil.ReclaimStaleLock: the lock file is renamed aside
// (only one racer can win the rename) and stolen only if the moved file's PID
// is still dead; a holder that reacquired the lock in the gap between the
// stale check and the rename carries a LIVE pid and is restored rather than
// stolen. Returns true only when a genuinely stale lock was removed. A
// non-nil error means the caller must fail closed instead of re-acquiring
// (see lockutil.ReclaimStaleLock).
func reclaimStaleLock(path string, isAlive func(pid int) bool) (bool, error) {
suffix := fmt.Sprintf("%d-%d", os.Getpid(), daemonLockSeq.Add(1))
return lockutil.ReclaimStaleLock(path, suffix, func(reclaimedPath string) bool {
pid, err := readPidFile(reclaimedPath)
return err == nil && pid > 0 && isAlive(pid)
})
if err := lock.WriteMetadata([]byte(fmt.Sprintf("%d\n", os.Getpid()))); err != nil {
return nil, errors.Join(err, lock.Release())
}
return &fileLock{lock: lock}, nil
}

// release removes the lock file. Safe to call once. An already-missing lock
// file is not an error (lockutil.RemoveLockFile swallows it on every platform).
// release drops the kernel lock and closes its handle. The stable PID file is
// deliberately retained so contenders always address the same inode.
func (l *fileLock) release() error {
if l == nil || l.path == "" {
if l == nil || l.lock == nil {
return nil
}
return lockutil.RemoveLockFile(l.path)
return l.lock.Release()
}

// readPidFile reads and parses the PID recorded in a lock file.
Expand All @@ -112,5 +65,20 @@ func readPidFile(path string) (int, error) {
if err != nil {
return 0, err
}
return strconv.Atoi(strings.TrimSpace(string(data)))
metadata := strings.TrimSpace(string(data))
if pid, err := strconv.Atoi(metadata); err == nil {
return pid, nil
}
pidText, sequenceText, ok := strings.Cut(metadata, "-")
if !ok || pidText == "" || sequenceText == "" || strings.Contains(sequenceText, "-") {
return 0, fmt.Errorf("daemon: parse lock PID metadata %q", metadata)
}
if _, err := strconv.ParseUint(sequenceText, 10, 64); err != nil {
return 0, fmt.Errorf("daemon: parse lock holder sequence: %w", err)
}
pid, err := strconv.Atoi(pidText)
if err != nil {
return 0, fmt.Errorf("daemon: parse lock PID: %w", err)
}
return pid, nil
}
96 changes: 54 additions & 42 deletions internal/daemon/lock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,27 +32,48 @@ func TestLockSingleInstance(t *testing.T) {
_ = l2.release()
}

func TestLockStaleRecovery(t *testing.T) {
// A PID liveness probe is diagnostic only. Even if it reports the recorded PID
// dead, a contended kernel lock proves that the holder is active and must win.
// The old rename-aside protocol trusted this stale observation and admitted a
// second daemon while the first one was still inside its critical section.
func TestLockDoesNotOverrideKernelHolderWithStalePIDProbe(t *testing.T) {
path := filepath.Join(t.TempDir(), "daemon.lock")
// Simulate a stale lock from a crashed daemon: a PID file whose process is
// dead.
first, err := acquireLock(path, func(int) bool { return true })
if err != nil {
t.Fatalf("first acquireLock: %v", err)
}
defer first.release()

second, err := acquireLock(path, func(int) bool { return false })
if second != nil {
_ = second.release()
}
if !errors.Is(err, ErrAlreadyRunning) {
t.Fatalf("contended acquire with stale PID probe = %v, want ErrAlreadyRunning", err)
}
}

func TestLockIgnoresStaleMetadataWhenKernelLockIsFree(t *testing.T) {
path := filepath.Join(t.TempDir(), "daemon.lock")
// A crashed daemon may leave PID metadata behind, but its kernel lock is
// released automatically, so the next daemon can acquire without moving it.
if err := os.WriteFile(path, []byte("4242\n"), 0o600); err != nil {
t.Fatalf("write stale lock: %v", err)
}
dead := func(int) bool { return false }
l, err := acquireLock(path, dead)
if err != nil {
t.Fatalf("stale-lock recovery failed: %v", err)
t.Fatalf("acquire over stale metadata: %v", err)
}
// The lock now records OUR pid, not the stale one.
data, _ := os.ReadFile(path)
if strings.TrimSpace(string(data)) != strconv.Itoa(os.Getpid()) {
t.Fatalf("reclaimed lock pid = %q, want %d", strings.TrimSpace(string(data)), os.Getpid())
t.Fatalf("lock pid = %q, want %d", strings.TrimSpace(string(data)), os.Getpid())
}
_ = l.release()
}

func TestLockReleaseRemovesFile(t *testing.T) {
func TestLockReleaseKeepsStableFile(t *testing.T) {
path := filepath.Join(t.TempDir(), "daemon.lock")
l, err := acquireLock(path, func(int) bool { return true })
if err != nil {
Expand All @@ -61,46 +82,37 @@ func TestLockReleaseRemovesFile(t *testing.T) {
if err := l.release(); err != nil {
t.Fatalf("release: %v", err)
}
if _, err := os.Stat(path); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("lock file still present after release: %v", err)
if _, err := os.Stat(path); err != nil {
t.Fatalf("release removed stable lock file: %v", err)
}
}

func TestReclaimStaleLockRemovesDeadHolder(t *testing.T) {
func TestReadPidFileMetadataFormats(t *testing.T) {
path := filepath.Join(t.TempDir(), "daemon.lock")
if err := os.WriteFile(path, []byte("4242\n"), 0o600); err != nil {
t.Fatalf("write stale lock: %v", err)
}
if ok, err := reclaimStaleLock(path, func(int) bool { return false }); err != nil || !ok {
t.Fatalf("reclaimStaleLock must report a genuinely stale (dead-PID) lock reclaimed (ok=%v err=%v)", ok, err)
}
if _, err := os.Stat(path); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("a reclaimed stale lock must be removed, stat = %v", err)
}
}

func TestReclaimStaleLockRestoresLiveHolder(t *testing.T) {
// If a holder reacquires the lock in the gap between the stale check and the
// rename, the moved file carries a LIVE pid; reclaim must restore it untouched
// rather than steal it — otherwise two daemons both "hold" the lock (D6).
path := filepath.Join(t.TempDir(), "daemon.lock")
if err := os.WriteFile(path, []byte("4242\n"), 0o600); err != nil {
t.Fatalf("write lock: %v", err)
}
if ok, err := reclaimStaleLock(path, func(int) bool { return true }); err != nil || ok {
t.Fatalf("reclaimStaleLock must NOT report a live-PID lock reclaimed (ok=%v err=%v)", ok, err)
}
data, err := os.ReadFile(path)
if err != nil {
t.Fatalf("live lock must be restored in place, read = %v", err)
}
if strings.TrimSpace(string(data)) != "4242" {
t.Fatalf("restored lock content = %q, want %q (unchanged)", strings.TrimSpace(string(data)), "4242")
}
// No sidelined ".stale" leftovers.
matches, _ := filepath.Glob(path + ".stale.*")
if len(matches) != 0 {
t.Fatalf("reclaim left sidelined files: %v", matches)
for _, test := range []struct {
name string
metadata string
wantPID int
wantErr bool
}{
{name: "legacy PID", metadata: "4242\n", wantPID: 4242},
{name: "PID and holder sequence", metadata: "4242-7\n", wantPID: 4242},
{name: "missing sequence", metadata: "4242-", wantErr: true},
{name: "non-numeric sequence", metadata: "4242-next", wantErr: true},
{name: "extra component", metadata: "4242-7-8", wantErr: true},
} {
t.Run(test.name, func(t *testing.T) {
if err := os.WriteFile(path, []byte(test.metadata), 0o600); err != nil {
t.Fatal(err)
}
got, err := readPidFile(path)
if (err != nil) != test.wantErr {
t.Fatalf("readPidFile() error = %v, wantErr %v", err, test.wantErr)
}
if !test.wantErr && got != test.wantPID {
t.Fatalf("readPidFile() = %d, want %d", got, test.wantPID)
}
})
}
}

Expand Down
Loading
Loading