diff --git a/CHANGELOG.md b/CHANGELOG.md index 583287f..744ebe3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/server/cmd/mem/qoder_checkpoint.go b/server/cmd/mem/qoder_checkpoint.go index a7fd069..a572985 100644 --- a/server/cmd/mem/qoder_checkpoint.go +++ b/server/cmd/mem/qoder_checkpoint.go @@ -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 diff --git a/server/cmd/mem/qoder_checkpoint_lock.go b/server/cmd/mem/qoder_checkpoint_lock.go new file mode 100644 index 0000000..f9cb299 --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_lock.go @@ -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 +} diff --git a/server/cmd/mem/qoder_checkpoint_lock_aix.go b/server/cmd/mem/qoder_checkpoint_lock_aix.go new file mode 100644 index 0000000..d757ba3 --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_lock_aix.go @@ -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, + }) +} diff --git a/server/cmd/mem/qoder_checkpoint_lock_other.go b/server/cmd/mem/qoder_checkpoint_lock_other.go new file mode 100644 index 0000000..86a50ff --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_lock_other.go @@ -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 +} diff --git a/server/cmd/mem/qoder_checkpoint_lock_unix.go b/server/cmd/mem/qoder_checkpoint_lock_unix.go new file mode 100644 index 0000000..fba5622 --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_lock_unix.go @@ -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) +} diff --git a/server/cmd/mem/qoder_checkpoint_lock_windows.go b/server/cmd/mem/qoder_checkpoint_lock_windows.go new file mode 100644 index 0000000..2f07483 --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_lock_windows.go @@ -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{}) +} diff --git a/server/cmd/mem/qoder_checkpoint_test.go b/server/cmd/mem/qoder_checkpoint_test.go new file mode 100644 index 0000000..4481552 --- /dev/null +++ b/server/cmd/mem/qoder_checkpoint_test.go @@ -0,0 +1,282 @@ +package main + +import ( + "context" + "fmt" + "os" + "os/exec" + "path/filepath" + "runtime" + "strconv" + "testing" + "time" +) + +const qoderCheckpointHelperEnv = "MEM_QODER_CHECKPOINT_HELPER" + +// TestQoderCheckpointHelperProcess is re-executed as an independent process by +// the tests below. The hold-lock action exits through os.Exit rather than +// release(), deliberately modelling a process that dies while it owns the +// advisory lock. +func TestQoderCheckpointHelperProcess(t *testing.T) { + if os.Getenv(qoderCheckpointHelperEnv) != "1" { + return + } + + action := os.Getenv("MEM_QODER_CHECKPOINT_ACTION") + stateDir := os.Getenv("MEM_QODER_CHECKPOINT_STATE_DIR") + abs := os.Getenv("MEM_QODER_CHECKPOINT_ABS") + lastLine, err := strconv.Atoi(os.Getenv("MEM_QODER_CHECKPOINT_LAST_LINE")) + if err != nil { + checkpointHelperExit("parse last line: %v", err) + } + size, err := strconv.ParseInt(os.Getenv("MEM_QODER_CHECKPOINT_SIZE"), 10, 64) + if err != nil { + checkpointHelperExit("parse checkpoint size: %v", err) + } + started := os.Getenv("MEM_QODER_CHECKPOINT_STARTED") + ready := os.Getenv("MEM_QODER_CHECKPOINT_READY") + release := os.Getenv("MEM_QODER_CHECKPOINT_RELEASE") + + if started != "" { + checkpointHelperWriteSignal(started) + } + cp := qoderCheckpoint{ + Abs: abs, + Size: size, + ModTime: os.Getenv("MEM_QODER_CHECKPOINT_MOD_TIME"), + LastLine: lastLine, + } + var heldLock *qoderCheckpointLock + switch action { + case "save": + if err := saveQoderCheckpoint(stateDir, cp); err != nil { + checkpointHelperExit("save checkpoint: %v", err) + } + case "save-hold-lock": + path := qoderCheckpointPath(stateDir, abs) + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + checkpointHelperExit("create checkpoint dir: %v", err) + } + lock, err := acquireQoderCheckpointLock(path) + if err != nil { + checkpointHelperExit("acquire checkpoint lock: %v", err) + } + heldLock = lock + if err := saveQoderCheckpointLocked(stateDir, path, cp); err != nil { + checkpointHelperExit("save locked checkpoint: %v", err) + } + default: + checkpointHelperExit("unknown action %q", action) + } + + checkpointHelperWriteSignal(ready) + if release != "" { + for { + if _, err := os.Stat(release); err == nil { + os.Exit(0) + } else if !os.IsNotExist(err) { + checkpointHelperExit("inspect release signal: %v", err) + } + time.Sleep(5 * time.Millisecond) + } + } + runtime.KeepAlive(heldLock) + os.Exit(0) +} + +func checkpointHelperWriteSignal(path string) { + if path == "" { + return + } + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + checkpointHelperExit("create signal directory: %v", err) + } + if err := os.WriteFile(path, []byte("ready\n"), 0o600); err != nil { + checkpointHelperExit("write signal: %v", err) + } +} + +func checkpointHelperExit(format string, args ...any) { + fmt.Fprintf(os.Stderr, "qoder checkpoint helper: "+format+"\n", args...) + os.Exit(2) +} + +type qoderCheckpointHelper struct { + cmd *exec.Cmd + cancel context.CancelFunc +} + +func startQoderCheckpointHelper(t *testing.T, values map[string]string) *qoderCheckpointHelper { + t.Helper() + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + cmd := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestQoderCheckpointHelperProcess$") + cmd.Stdout = os.Stderr + cmd.Stderr = os.Stderr + cmd.Env = append(os.Environ(), qoderCheckpointHelperEnv+"=1") + for key, value := range values { + cmd.Env = append(cmd.Env, key+"="+value) + } + if err := cmd.Start(); err != nil { + cancel() + t.Fatalf("start qoder checkpoint helper: %v", err) + } + helper := &qoderCheckpointHelper{cmd: cmd, cancel: cancel} + t.Cleanup(func() { + cancel() + if cmd.ProcessState == nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + } + }) + return helper +} + +func (h *qoderCheckpointHelper) wait(t *testing.T) { + t.Helper() + if err := h.cmd.Wait(); err != nil { + t.Fatalf("qoder checkpoint helper failed: %v", err) + } +} + +func waitForQoderCheckpointSignal(t *testing.T, path string) { + t.Helper() + deadline := time.Now().Add(3 * time.Second) + for { + if _, err := os.Stat(path); err == nil { + return + } else if !os.IsNotExist(err) { + t.Fatalf("inspect checkpoint signal %s: %v", path, err) + } + if time.Now().After(deadline) { + t.Fatalf("timed out waiting for checkpoint signal %s", path) + } + time.Sleep(5 * time.Millisecond) + } +} + +func requireNoQoderCheckpointSignal(t *testing.T, path string, duration time.Duration) { + t.Helper() + deadline := time.Now().Add(duration) + for time.Now().Before(deadline) { + if _, err := os.Stat(path); err == nil { + t.Fatalf("checkpoint writer escaped its live process lock before release: %s", path) + } else if !os.IsNotExist(err) { + t.Fatalf("inspect checkpoint signal %s: %v", path, err) + } + time.Sleep(5 * time.Millisecond) + } +} + +func writeQoderCheckpointRelease(t *testing.T, path string) { + t.Helper() + if err := os.WriteFile(path, []byte("release\n"), 0o600); err != nil { + t.Fatalf("release checkpoint helper: %v", err) + } +} + +func TestSaveQoderCheckpointKeepsHighestLastLineAcrossIndependentProcesses(t *testing.T) { + dir := t.TempDir() + stateDir := filepath.Join(dir, "state") + abs := filepath.Join(dir, "sessions", "s.jsonl") + if err := os.MkdirAll(filepath.Dir(abs), 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(abs, []byte("0123456789abcdef0123456789abcdef\n"), 0o600); err != nil { + t.Fatal(err) + } + + highReady := filepath.Join(dir, "high-ready") + highRelease := filepath.Join(dir, "high-release") + high := startQoderCheckpointHelper(t, map[string]string{ + "MEM_QODER_CHECKPOINT_ACTION": "save", + "MEM_QODER_CHECKPOINT_STATE_DIR": stateDir, + "MEM_QODER_CHECKPOINT_ABS": abs, + "MEM_QODER_CHECKPOINT_LAST_LINE": "12", + "MEM_QODER_CHECKPOINT_SIZE": "33", + "MEM_QODER_CHECKPOINT_MOD_TIME": "2026-09-01T00:00:12Z", + "MEM_QODER_CHECKPOINT_READY": highReady, + "MEM_QODER_CHECKPOINT_RELEASE": highRelease, + }) + waitForQoderCheckpointSignal(t, highReady) + + lowReady := filepath.Join(dir, "low-ready") + low := startQoderCheckpointHelper(t, map[string]string{ + "MEM_QODER_CHECKPOINT_ACTION": "save", + "MEM_QODER_CHECKPOINT_STATE_DIR": stateDir, + "MEM_QODER_CHECKPOINT_ABS": abs, + "MEM_QODER_CHECKPOINT_LAST_LINE": "4", + "MEM_QODER_CHECKPOINT_SIZE": "4", + "MEM_QODER_CHECKPOINT_MOD_TIME": "2026-09-01T00:00:04Z", + "MEM_QODER_CHECKPOINT_READY": lowReady, + }) + waitForQoderCheckpointSignal(t, lowReady) + low.wait(t) + + got := loadQoderCheckpoint(stateDir, abs) + if got.LastLine != 12 { + t.Fatalf("LastLine after high then stale low process = %d, want 12", got.LastLine) + } + if got.Size != 33 || got.ModTime != "2026-09-01T00:00:12Z" || got.Abs != abs { + t.Fatalf("stale process regressed checkpoint diagnostics: %+v", got) + } + + writeQoderCheckpointRelease(t, highRelease) + high.wait(t) +} + +func TestSaveQoderCheckpointSerializesAndRecoversAfterOwnerExit(t *testing.T) { + dir := t.TempDir() + stateDir := filepath.Join(dir, "state") + abs := filepath.Join(dir, "sessions", "s.jsonl") + if err := os.MkdirAll(filepath.Dir(abs), 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(abs, []byte("0123456789abcdef0123456789abcdef\n"), 0o600); err != nil { + t.Fatal(err) + } + + highReady := filepath.Join(dir, "high-ready") + highRelease := filepath.Join(dir, "high-release") + high := startQoderCheckpointHelper(t, map[string]string{ + "MEM_QODER_CHECKPOINT_ACTION": "save-hold-lock", + "MEM_QODER_CHECKPOINT_STATE_DIR": stateDir, + "MEM_QODER_CHECKPOINT_ABS": abs, + "MEM_QODER_CHECKPOINT_LAST_LINE": "12", + "MEM_QODER_CHECKPOINT_SIZE": "33", + "MEM_QODER_CHECKPOINT_MOD_TIME": "2026-09-01T00:00:12Z", + "MEM_QODER_CHECKPOINT_READY": highReady, + "MEM_QODER_CHECKPOINT_RELEASE": highRelease, + }) + waitForQoderCheckpointSignal(t, highReady) + + lowStarted := filepath.Join(dir, "low-started") + lowReady := filepath.Join(dir, "low-ready") + low := startQoderCheckpointHelper(t, map[string]string{ + "MEM_QODER_CHECKPOINT_ACTION": "save", + "MEM_QODER_CHECKPOINT_STATE_DIR": stateDir, + "MEM_QODER_CHECKPOINT_ABS": abs, + "MEM_QODER_CHECKPOINT_LAST_LINE": "4", + "MEM_QODER_CHECKPOINT_SIZE": "4", + "MEM_QODER_CHECKPOINT_MOD_TIME": "2026-09-01T00:00:04Z", + "MEM_QODER_CHECKPOINT_STARTED": lowStarted, + "MEM_QODER_CHECKPOINT_READY": lowReady, + }) + waitForQoderCheckpointSignal(t, lowStarted) + requireNoQoderCheckpointSignal(t, lowReady, 100*time.Millisecond) + + writeQoderCheckpointRelease(t, highRelease) + high.wait(t) + waitForQoderCheckpointSignal(t, lowReady) + low.wait(t) + + if got := loadQoderCheckpoint(stateDir, abs).LastLine; got != 12 { + t.Fatalf("LastLine after abandoned lock and lower writer = %d, want 12", got) + } + if err := saveQoderCheckpoint(stateDir, qoderCheckpoint{Abs: abs, LastLine: 13}); err != nil { + t.Fatalf("checkpoint remained locked after owner exit: %v", err) + } + if got := loadQoderCheckpoint(stateDir, abs).LastLine; got != 13 { + t.Fatalf("LastLine after recovery writer = %d, want 13", got) + } +}