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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ The project publishes 0.x prerelease versions; a stable release line is not yet

## [Unreleased]

### Changed

- Ingest cursor locks try non-blocking exclusive locks and give up after 5s so a wedged peer becomes a warning instead of a silent hang. Refs #139.

### Added

- Additive `durable-memory.v1` envelope for derived RoleWeave/mem records
Expand Down
26 changes: 22 additions & 4 deletions server/internal/ingest/cursor_lock.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,12 @@ package ingest
import (
"fmt"
"os"
"time"
)

const (
cursorLockWait = 5 * time.Second
cursorLockRetry = 10 * time.Millisecond
)

// cursorLock holds an advisory lock on one cursor sidecar. The sidecar
Expand All @@ -14,15 +20,27 @@ type cursorLock struct {
}

func acquireCursorLock(cursorPath string) (*cursorLock, error) {
return acquireCursorLockWithTimeout(cursorPath, cursorLockWait)
}

func acquireCursorLockWithTimeout(cursorPath string, timeout time.Duration) (*cursorLock, error) {
file, err := os.OpenFile(cursorPath+".lock", os.O_CREATE|os.O_RDWR, 0o600)
if err != nil {
return nil, fmt.Errorf("open lock file: %w", err)
}
if err := lockCursorFile(file); err != nil {
_ = file.Close()
return nil, fmt.Errorf("acquire OS lock: %w", err)
deadline := time.Now().Add(timeout)
for {
if err := tryLockCursorFile(file); err == nil {
return &cursorLock{file: file}, nil
} else if !isCursorLockBusy(err) {
_ = file.Close()
return nil, fmt.Errorf("acquire OS lock: %w", err)
} else if !time.Now().Before(deadline) {
_ = file.Close()
return nil, fmt.Errorf("acquire OS lock: timed out after %s: %w", timeout, err)
}
time.Sleep(cursorLockRetry)
}
return &cursorLock{file: file}, nil
}

func (l *cursorLock) release() error {
Expand Down
13 changes: 9 additions & 4 deletions server/internal/ingest/cursor_lock_aix.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,20 +3,25 @@
package ingest

import (
"errors"
"os"

"golang.org/x/sys/unix"
)

// AIX does not expose flock(2) through x/sys, so use the blocking fcntl record
// lock equivalent for the first byte of the persistent sidecar inode.
func lockCursorFile(file *os.File) error {
return unix.FcntlFlock(file.Fd(), unix.F_SETLKW, &unix.Flock_t{
// AIX does not expose flock(2) through x/sys, so use the non-blocking fcntl
// record lock equivalent for the first byte of the persistent sidecar inode.
func tryLockCursorFile(file *os.File) error {
return unix.FcntlFlock(file.Fd(), unix.F_SETLK, &unix.Flock_t{
Type: unix.F_WRLCK,
Len: 1,
})
}

func isCursorLockBusy(err error) bool {
return errors.Is(err, unix.EAGAIN) || errors.Is(err, unix.EACCES)
}

func unlockCursorFile(file *os.File) error {
return unix.FcntlFlock(file.Fd(), unix.F_SETLK, &unix.Flock_t{
Type: unix.F_UNLCK,
Expand Down
6 changes: 5 additions & 1 deletion server/internal/ingest/cursor_lock_other.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,14 @@ import (
"os"
)

func lockCursorFile(_ *os.File) error {
func tryLockCursorFile(_ *os.File) error {
return fmt.Errorf("cursor locks are not supported on this operating system")
}

func isCursorLockBusy(_ error) bool {
return false
}

func unlockCursorFile(_ *os.File) error {
return nil
}
90 changes: 90 additions & 0 deletions server/internal/ingest/cursor_lock_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
package ingest

import (
"context"
"os"
"os/exec"
"path/filepath"
"runtime"
"testing"
"time"
)

const cursorLockHelperEnv = "MEM_CURSOR_LOCK_HELPER"

func TestCursorLockHelperProcess(t *testing.T) {
if os.Getenv(cursorLockHelperEnv) != "1" {
return
}
path := os.Getenv("MEM_CURSOR_LOCK_PATH")
ready := os.Getenv("MEM_CURSOR_LOCK_READY")
release := os.Getenv("MEM_CURSOR_LOCK_RELEASE")
lock, err := acquireCursorLock(path)
if err != nil {
os.Stderr.WriteString("cursor lock helper: " + err.Error() + "\n")
os.Exit(2)
}
defer runtime.KeepAlive(lock)
if err := os.WriteFile(ready, []byte("ready\n"), 0o600); err != nil {
os.Stderr.WriteString("cursor lock helper: write ready: " + err.Error() + "\n")
os.Exit(2)
}
for {
if _, err := os.Stat(release); err == nil {
os.Exit(0)
} else if !os.IsNotExist(err) {
os.Stderr.WriteString("cursor lock helper: inspect release: " + err.Error() + "\n")
os.Exit(2)
}
time.Sleep(5 * time.Millisecond)
}
}

func TestAcquireCursorLockTimesOut(t *testing.T) {
if runtime.GOOS != "darwin" && runtime.GOOS != "linux" && runtime.GOOS != "windows" {
t.Skip("cursor locks are unsupported on this operating system")
}
dir := t.TempDir()
path := filepath.Join(dir, "cursor.json")
ready := filepath.Join(dir, "ready")
release := filepath.Join(dir, "release")

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
t.Cleanup(cancel)
cmd := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestCursorLockHelperProcess$")
cmd.Stdout = os.Stderr
cmd.Stderr = os.Stderr
cmd.Env = append(os.Environ(),
cursorLockHelperEnv+"=1",
"MEM_CURSOR_LOCK_PATH="+path,
"MEM_CURSOR_LOCK_READY="+ready,
"MEM_CURSOR_LOCK_RELEASE="+release,
)
if err := cmd.Start(); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
_ = os.WriteFile(release, []byte("release\n"), 0o600)
_ = cmd.Wait()
})

deadline := time.Now().Add(3 * time.Second)
for {
if _, err := os.Stat(ready); err == nil {
break
} else if !os.IsNotExist(err) {
t.Fatal(err)
}
if time.Now().After(deadline) {
t.Fatal("timed out waiting for helper to hold the cursor lock")
}
time.Sleep(5 * time.Millisecond)
}

started := time.Now()
if _, err := acquireCursorLockWithTimeout(path, 25*time.Millisecond); err == nil {
t.Fatal("second cursor lock unexpectedly acquired")
} else if elapsed := time.Since(started); elapsed > time.Second {
t.Fatalf("lock timeout took too long: %s", elapsed)
}
}
9 changes: 7 additions & 2 deletions server/internal/ingest/cursor_lock_unix.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,18 @@
package ingest

import (
"errors"
"os"

"golang.org/x/sys/unix"
)

func lockCursorFile(file *os.File) error {
return unix.Flock(int(file.Fd()), unix.LOCK_EX)
func tryLockCursorFile(file *os.File) error {
return unix.Flock(int(file.Fd()), unix.LOCK_EX|unix.LOCK_NB)
}

func isCursorLockBusy(err error) bool {
return errors.Is(err, unix.EAGAIN) || errors.Is(err, unix.EWOULDBLOCK)
}

func unlockCursorFile(file *os.File) error {
Expand Down
9 changes: 7 additions & 2 deletions server/internal/ingest/cursor_lock_windows.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,24 +3,29 @@
package ingest

import (
"errors"
"os"

"golang.org/x/sys/windows"
)

// Lock a one-byte range. Windows releases a LockFileEx lock when the owning
// process or file handle exits, matching the Unix advisory-lock lifecycle.
func lockCursorFile(file *os.File) error {
func tryLockCursorFile(file *os.File) error {
return windows.LockFileEx(
windows.Handle(file.Fd()),
windows.LOCKFILE_EXCLUSIVE_LOCK,
windows.LOCKFILE_EXCLUSIVE_LOCK|windows.LOCKFILE_FAIL_IMMEDIATELY,
0,
1,
0,
&windows.Overlapped{},
)
}

func isCursorLockBusy(err error) bool {
return errors.Is(err, windows.ERROR_LOCK_VIOLATION)
}

func unlockCursorFile(file *os.File) error {
return windows.UnlockFileEx(windows.Handle(file.Fd()), 0, 1, 0, &windows.Overlapped{})
}
Loading