Make the process-wide lock atomic with flock (closes #227)
check / check (push) Waiting to run
check / check (push) Waiting to run
Acquire read vaultik.pid, checked whether that PID was alive, then wrote its own, so two writers started together could both pass the check and both run. The lock is now an flock on vaultik.pid, held while the file stays open; the kernel drops it when the process exits, so the stale-PID check is gone. Release empties the file instead of deleting it. Deleting it would let a process that opened the old file a moment earlier lock it while another creates and locks a new one. The new concurrent test fails against the old code only when the race is hit, not on every run; against the fix it cannot admit two callers. Model: opus-5-5
This commit was merged in pull request #258.
This commit is contained in:
@@ -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
@@ -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
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
Reference in New Issue
Block a user