Make the process-wide lock atomic with flock #258

Merged
clawbot merged 1 commits from issue-227-atomic-pid-lock into next 2026-10-07 06:12:11 +02:00
3 changed files with 140 additions and 55 deletions
+10
View File
@@ -22,6 +22,16 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps # Completed Steps
- 2026-10-07: Made taking the process-wide lock atomic
([issue #227](https://git.eeqj.de/sneak/vaultik/issues/227)). The lock
read `vaultik.pid`, checked whether that PID was alive and then wrote
its own, so two writers started together could both pass the check and
both run. It is now an `flock` on `vaultik.pid`, held until the run
ends; the kernel drops it when the process exits, so a crash leaves no
lock behind. A clean exit now empties the file instead of deleting it,
because deleting it would let two later runs each lock a different
file.
- 2026-10-06: Made `snapshot remove --json` write only its document to - 2026-10-06: Made `snapshot remove --json` write only its document to
stdout when the destination store cannot be reached stdout when the destination store cannot be reached
([issue #251](https://git.eeqj.de/sneak/vaultik/issues/251)). Its ([issue #251](https://git.eeqj.de/sneak/vaultik/issues/251)). Its
+67 -52
View File
@@ -10,15 +10,18 @@ import (
"path/filepath" "path/filepath"
"strconv" "strconv"
"strings" "strings"
"syscall"
"golang.org/x/sys/unix"
) )
// ErrAlreadyRunning indicates another vaultik instance is running. // ErrAlreadyRunning indicates another vaultik instance is running.
var ErrAlreadyRunning = errors.New("another vaultik instance is already running") var ErrAlreadyRunning = errors.New("another vaultik instance is already running")
// Lock represents an acquired PID lock. // Lock represents an acquired PID lock: an flock(2) on the PID file,
// held while the file stays open. The kernel drops it when the process
// exits, however it exits, so a crashed run never leaves the lock held.
type Lock struct { type Lock struct {
path string file *os.File
} }
const ( const (
@@ -29,10 +32,9 @@ const (
) )
// Acquire attempts to acquire a PID lock in the specified directory. // Acquire attempts to acquire a PID lock in the specified directory.
// If the lock file exists and the process is still running, it returns // If another process holds the lock, it returns ErrAlreadyRunning with
// ErrAlreadyRunning with details about the existing process. // that process's PID. On success, it writes the current PID to the lock
// On success, it writes the current PID to the lock file and returns // file and returns a Lock that must be released with Release().
// a Lock that must be released with Release().
func Acquire(lockDir string) (*Lock, error) { func Acquire(lockDir string) (*Lock, error) {
// Ensure lock directory exists // Ensure lock directory exists
err := os.MkdirAll(lockDir, lockDirPerm) err := os.MkdirAll(lockDir, lockDirPerm)
@@ -42,56 +44,82 @@ func Acquire(lockDir string) (*Lock, error) {
lockPath := filepath.Join(lockDir, "vaultik.pid") lockPath := filepath.Join(lockDir, "vaultik.pid")
// Check for existing lock // No O_TRUNC: the file may hold the PID of the process that has the
existingPID, err := readPIDFile(lockPath) // lock, which the error below reports.
if err == nil { file, err := os.OpenFile( //nolint:gosec // G304: path is our own lock file
// Lock file exists, check if process is running lockPath, os.O_RDWR|os.O_CREATE, pidFilePerm)
if isProcessRunning(existingPID) {
return nil, fmt.Errorf("%w (PID %d)", ErrAlreadyRunning, existingPID)
}
// Process is not running, stale lock file - we can take over
}
// Write our PID
pid := os.Getpid()
err = os.WriteFile(lockPath, []byte(strconv.Itoa(pid)), pidFilePerm)
if err != nil { if err != nil {
return nil, fmt.Errorf("writing PID file: %w", err) return nil, fmt.Errorf("opening PID file: %w", err)
} }
return &Lock{path: lockPath}, nil err = unix.Flock(int(file.Fd()), unix.LOCK_EX|unix.LOCK_NB)
if err != nil {
_ = file.Close()
if errors.Is(err, unix.EWOULDBLOCK) {
return nil, alreadyRunningError(lockPath)
}
return nil, fmt.Errorf("locking PID file: %w", err)
}
err = writePID(file)
if err != nil {
_ = file.Close()
return nil, err
}
return &Lock{file: file}, nil
} }
// Release removes the PID lock file. // Release empties the PID file and closes it, which drops the lock.
// It is safe to call Release multiple times. // It is safe to call Release multiple times.
func (l *Lock) Release() error { func (l *Lock) Release() error {
if l == nil || l.path == "" { if l == nil || l.file == nil {
return nil return nil
} }
// Verify we still own the lock (our PID is in the file) file := l.file
existingPID, err := readPIDFile(l.path) l.file = nil
// Do not remove the file here. A process that opened it a moment
// earlier could then lock the removed file while another creates and
// locks a new one, and both would run.
truncateErr := file.Truncate(0)
closeErr := file.Close()
return errors.Join(truncateErr, closeErr)
}
// writePID replaces the contents of the locked PID file with the current
// PID.
func writePID(file *os.File) error {
err := file.Truncate(0)
if err != nil { if err != nil {
// File already gone or unreadable - that's fine return fmt.Errorf("truncating PID file: %w", err)
return nil //nolint:nilerr // unreadable lock file means nothing to release
} }
if existingPID != os.Getpid() { _, err = file.WriteAt([]byte(strconv.Itoa(os.Getpid())), 0)
// Someone else wrote to our lock file - don't remove it if err != nil {
return nil return fmt.Errorf("writing PID file: %w", err)
} }
err = os.Remove(l.path)
if err != nil && !os.IsNotExist(err) {
return fmt.Errorf("removing PID file: %w", err)
}
l.path = "" // Prevent double-release
return nil return nil
} }
// alreadyRunningError reports that another process holds the lock,
// naming its PID when the file holds one. The holder writes its PID just
// after it locks, so the file can briefly be empty.
func alreadyRunningError(lockPath string) error {
pid, err := readPIDFile(lockPath)
if err != nil {
return ErrAlreadyRunning
}
return fmt.Errorf("%w (PID %d)", ErrAlreadyRunning, pid)
}
// readPIDFile reads and parses the PID from a lock file. // readPIDFile reads and parses the PID from a lock file.
func readPIDFile(path string) (int, error) { func readPIDFile(path string) (int, error) {
data, err := os.ReadFile(path) //nolint:gosec // G304: path is our own lock file data, err := os.ReadFile(path) //nolint:gosec // G304: path is our own lock file
@@ -106,16 +134,3 @@ func readPIDFile(path string) (int, error) {
return pid, nil return pid, nil
} }
// isProcessRunning checks if a process with the given PID is running.
func isProcessRunning(pid int) bool {
process, err := os.FindProcess(pid)
if err != nil {
return false
}
// On Unix, FindProcess always succeeds. We need to send signal 0 to check.
err = process.Signal(syscall.Signal(0))
return err == nil
}
+63 -3
View File
@@ -4,6 +4,7 @@ import (
"os" "os"
"path/filepath" "path/filepath"
"strconv" "strconv"
"sync"
"testing" "testing"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
@@ -33,9 +34,10 @@ func TestAcquireAndRelease(t *testing.T) {
err = lock.Release() err = lock.Release()
require.NoError(t, err) require.NoError(t, err)
// Verify PID file is gone // Verify PID file is empty
_, err = os.Stat(pidPath) data, err = os.ReadFile(pidPath) //nolint:gosec // G304: test's own temp file
assert.True(t, os.IsNotExist(err)) require.NoError(t, err)
assert.Empty(t, data)
} }
func TestAcquireBlocksSecondInstance(t *testing.T) { func TestAcquireBlocksSecondInstance(t *testing.T) {
@@ -55,6 +57,64 @@ func TestAcquireBlocksSecondInstance(t *testing.T) {
lock2, err := pidlock.Acquire(tmpDir) lock2, err := pidlock.Acquire(tmpDir)
require.ErrorIs(t, err, pidlock.ErrAlreadyRunning) require.ErrorIs(t, err, pidlock.ErrAlreadyRunning)
assert.Nil(t, lock2) assert.Nil(t, lock2)
// Once the first lock is released, the next Acquire succeeds
require.NoError(t, lock1.Release())
lock3, err := pidlock.Acquire(tmpDir)
require.NoError(t, err)
require.NoError(t, lock3.Release())
}
// TestConcurrentAcquireAdmitsOne starts many Acquire calls at the same
// moment, as two cron entries firing together would, and checks that
// exactly one of them gets the lock.
func TestConcurrentAcquireAdmitsOne(t *testing.T) {
t.Parallel()
const callers = 50
tmpDir := t.TempDir()
start := make(chan struct{})
var (
mu sync.Mutex
acquired []*pidlock.Lock
failures []error
wg sync.WaitGroup
)
for range callers {
wg.Go(func() {
<-start
lock, err := pidlock.Acquire(tmpDir)
mu.Lock()
defer mu.Unlock()
if err != nil {
failures = append(failures, err)
return
}
acquired = append(acquired, lock)
})
}
close(start)
wg.Wait()
for _, lock := range acquired {
require.NoError(t, lock.Release())
}
assert.Len(t, acquired, 1, "exactly one caller should hold the lock")
for _, err := range failures {
require.ErrorIs(t, err, pidlock.ErrAlreadyRunning)
}
} }
func TestAcquireWithStaleLock(t *testing.T) { func TestAcquireWithStaleLock(t *testing.T) {