Make the process-wide lock atomic with flock #258
@@ -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