From 1271102a5ef2607e922b36d2c365688c23730904 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=8B=92=E5=B8=83=E6=9C=97-=E8=A9=B9=E5=A7=86=E6=96=AF?= <318569545+waterbro-8@users.noreply.github.com> Date: Fri, 18 Sep 2026 09:39:33 +0000 Subject: [PATCH] fix(ingest): bound cursor lock waits so a wedged peer cannot hang ingest (#139) #217 already moved the OS lock into server/internal/ingest. Port the remaining #191 review items onto that package: non-blocking flock / LockFileEx, a 5s give-up, and a subprocess test that a contended lock returns instead of blocking forever. --- CHANGELOG.md | 4 + server/internal/ingest/cursor_lock.go | 26 +++++- server/internal/ingest/cursor_lock_aix.go | 13 ++- server/internal/ingest/cursor_lock_other.go | 6 +- server/internal/ingest/cursor_lock_test.go | 90 +++++++++++++++++++ server/internal/ingest/cursor_lock_unix.go | 9 +- server/internal/ingest/cursor_lock_windows.go | 9 +- 7 files changed, 144 insertions(+), 13 deletions(-) create mode 100644 server/internal/ingest/cursor_lock_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 30320bd..a4fdaf0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/server/internal/ingest/cursor_lock.go b/server/internal/ingest/cursor_lock.go index 3091d8d..6619d2b 100644 --- a/server/internal/ingest/cursor_lock.go +++ b/server/internal/ingest/cursor_lock.go @@ -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 @@ -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 { diff --git a/server/internal/ingest/cursor_lock_aix.go b/server/internal/ingest/cursor_lock_aix.go index eaccd7d..b5bd033 100644 --- a/server/internal/ingest/cursor_lock_aix.go +++ b/server/internal/ingest/cursor_lock_aix.go @@ -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, diff --git a/server/internal/ingest/cursor_lock_other.go b/server/internal/ingest/cursor_lock_other.go index ef18b73..c542b67 100644 --- a/server/internal/ingest/cursor_lock_other.go +++ b/server/internal/ingest/cursor_lock_other.go @@ -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 } diff --git a/server/internal/ingest/cursor_lock_test.go b/server/internal/ingest/cursor_lock_test.go new file mode 100644 index 0000000..356d30b --- /dev/null +++ b/server/internal/ingest/cursor_lock_test.go @@ -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) + } +} diff --git a/server/internal/ingest/cursor_lock_unix.go b/server/internal/ingest/cursor_lock_unix.go index 5d87009..71e2e5a 100644 --- a/server/internal/ingest/cursor_lock_unix.go +++ b/server/internal/ingest/cursor_lock_unix.go @@ -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 { diff --git a/server/internal/ingest/cursor_lock_windows.go b/server/internal/ingest/cursor_lock_windows.go index 533f241..2503e43 100644 --- a/server/internal/ingest/cursor_lock_windows.go +++ b/server/internal/ingest/cursor_lock_windows.go @@ -3,6 +3,7 @@ package ingest import ( + "errors" "os" "golang.org/x/sys/windows" @@ -10,10 +11,10 @@ import ( // 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, @@ -21,6 +22,10 @@ func lockCursorFile(file *os.File) error { ) } +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{}) }