Skip to content
Closed
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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@ The project publishes 0.x prerelease versions; a stable release line is not yet
`Content-Security-Policy: default-src 'none'` to every API response, ordered
outside the CORS handler so a preflight reply carries them too.

### Fixed

- Concurrent `mem ingest qoder` processes now serialize each transcript's
checkpoint through an OS-backed sidecar lock, use unique private staging
files, and retain the highest successfully committed line cursor. A process
crash releases its advisory lock automatically, so a later ingest can resume
rather than being blocked by an orphaned lock.

## [0.1.1] - 2026-08-31

### Changed
Expand Down
50 changes: 45 additions & 5 deletions server/cmd/mem/qoder_checkpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,22 +58,62 @@ func loadQoderCheckpoint(stateDir, abs string) qoderCheckpoint {
return cp
}

// saveQoderCheckpoint atomically persists a transcript cursor. Errors are
// saveQoderCheckpoint atomically persists a transcript cursor. The per-cursor
// OS-backed lock covers the read/merge/write sequence so independent ingest
// processes cannot move LastLine backwards or share a staging path. Errors are
// returned (callers may warn without failing the whole ingest).
func saveQoderCheckpoint(stateDir string, cp qoderCheckpoint) error {
func saveQoderCheckpoint(stateDir string, cp qoderCheckpoint) (err error) {
p := qoderCheckpointPath(stateDir, cp.Abs)
if err := os.MkdirAll(filepath.Dir(p), 0o700); err != nil {
return fmt.Errorf("create checkpoint dir: %w", err)
}
lock, err := acquireQoderCheckpointLock(p)
if err != nil {
return fmt.Errorf("lock checkpoint: %w", err)
}
defer func() {
if releaseErr := lock.release(); err == nil && releaseErr != nil {
err = fmt.Errorf("release checkpoint lock: %w", releaseErr)
}
}()

return saveQoderCheckpointLocked(stateDir, p, cp)
}

// saveQoderCheckpointLocked commits cp while the caller owns p's checkpoint
// lock. Keeping this small inner operation separate lets the lock span the
// current-cursor read as well as the atomic replacement.
func saveQoderCheckpointLocked(stateDir, p string, cp qoderCheckpoint) error {
current := loadQoderCheckpoint(stateDir, cp.Abs)
if current.LastLine > cp.LastLine {
cp = current
}
b, err := json.Marshal(cp)
if err != nil {
return fmt.Errorf("encode checkpoint: %w", err)
}
tmp := p + ".tmp"
if err := os.WriteFile(tmp, b, 0o600); err != nil {
tmp, err := os.CreateTemp(filepath.Dir(p), "."+filepath.Base(p)+".tmp-*")
if err != nil {
return fmt.Errorf("create checkpoint staging file: %w", err)
}
tmpName := tmp.Name()
defer func() {
_ = tmp.Close()
_ = os.Remove(tmpName)
}()
if err := tmp.Chmod(0o600); err != nil {
return fmt.Errorf("secure checkpoint staging file: %w", err)
}
if _, err := tmp.Write(b); err != nil {
return fmt.Errorf("write checkpoint: %w", err)
}
if err := os.Rename(tmp, p); err != nil {
if err := tmp.Sync(); err != nil {
return fmt.Errorf("sync checkpoint staging file: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("close checkpoint staging file: %w", err)
}
if err := os.Rename(tmpName, p); err != nil {
return fmt.Errorf("commit checkpoint: %w", err)
}
return nil
Expand Down
43 changes: 43 additions & 0 deletions server/cmd/mem/qoder_checkpoint_lock.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
package main

import (
"fmt"
"os"
)

// qoderCheckpointLock holds an advisory lock on one cursor sidecar. The
// sidecar deliberately remains on disk after release: unlinking a locked file
// can create a second inode that another process locks independently. The OS
// releases the advisory lock when this descriptor, or its owning process,
// exits.
type qoderCheckpointLock struct {
file *os.File
}

func acquireQoderCheckpointLock(checkpointPath string) (*qoderCheckpointLock, error) {
file, err := os.OpenFile(checkpointPath+".lock", os.O_CREATE|os.O_RDWR, 0o600)
if err != nil {
return nil, fmt.Errorf("open lock file: %w", err)
}
if err := lockQoderCheckpointFile(file); err != nil {
_ = file.Close()
return nil, fmt.Errorf("acquire OS lock: %w", err)
}
return &qoderCheckpointLock{file: file}, nil
}

func (l *qoderCheckpointLock) release() error {
if l == nil || l.file == nil {
return nil
}
unlockErr := unlockQoderCheckpointFile(l.file)
closeErr := l.file.Close()
l.file = nil
if unlockErr != nil {
return fmt.Errorf("unlock OS lock: %w", unlockErr)
}
if closeErr != nil {
return fmt.Errorf("close lock file: %w", closeErr)
}
return nil
}
25 changes: 25 additions & 0 deletions server/cmd/mem/qoder_checkpoint_lock_aix.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
//go:build aix

package main

import (
"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 lockQoderCheckpointFile(file *os.File) error {
return unix.FcntlFlock(file.Fd(), unix.F_SETLKW, &unix.Flock_t{
Type: unix.F_WRLCK,
Len: 1,
})
}

func unlockQoderCheckpointFile(file *os.File) error {
return unix.FcntlFlock(file.Fd(), unix.F_SETLK, &unix.Flock_t{
Type: unix.F_UNLCK,
Len: 1,
})
}
16 changes: 16 additions & 0 deletions server/cmd/mem/qoder_checkpoint_lock_other.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
//go:build !(aix || darwin || dragonfly || freebsd || linux || netbsd || openbsd || solaris || windows)

package main

import (
"fmt"
"os"
)

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

func unlockQoderCheckpointFile(_ *os.File) error {
return nil
}
17 changes: 17 additions & 0 deletions server/cmd/mem/qoder_checkpoint_lock_unix.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
//go:build darwin || dragonfly || freebsd || linux || netbsd || openbsd || solaris

package main

import (
"os"

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

func lockQoderCheckpointFile(file *os.File) error {
return unix.Flock(int(file.Fd()), unix.LOCK_EX)
}

func unlockQoderCheckpointFile(file *os.File) error {
return unix.Flock(int(file.Fd()), unix.LOCK_UN)
}
26 changes: 26 additions & 0 deletions server/cmd/mem/qoder_checkpoint_lock_windows.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
//go:build windows

package main

import (
"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 lockQoderCheckpointFile(file *os.File) error {
return windows.LockFileEx(
windows.Handle(file.Fd()),
windows.LOCKFILE_EXCLUSIVE_LOCK,
0,
1,
0,
&windows.Overlapped{},
)
}

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