toolsrv: add path-based lock table for cross-agent serialization
All foreground tool calls now acquire a lock based on the tool's declared scope and file path before execution: - scope "read": no lock (reads never conflict) - scope "write": exclusive lock on the file path - scope "global": exclusive global lock (serializes with everything) This ensures writes to the same path serialize regardless of which agent initiated the call, enabling safe parallel sub-agents within a session without explicit coordination. Cross-session serialization (shared toolsrv per host) is left as future work.
This commit is contained in:
parent
ab79f6edd1
commit
5dcde264e6
|
|
@ -0,0 +1,71 @@
|
|||
// pathlock.go — Path-based read/write lock coordination for tool execution.
|
||||
//
|
||||
// All foreground tool calls acquire a lock before execution based on the
|
||||
// tool's declared scope and the file path from its arguments:
|
||||
//
|
||||
// - scope "read": no lock (reads never conflict)
|
||||
// - scope "write": exclusive lock on the file path (serializes writes to same path)
|
||||
// - scope "global": exclusive global lock (serializes with everything)
|
||||
//
|
||||
// This ensures correctness across all agents in a session (and eventually
|
||||
// across sessions sharing the same toolsrv), without the agents needing
|
||||
// to coordinate explicitly.
|
||||
package server
|
||||
|
||||
import "sync"
|
||||
|
||||
// pathLockTable provides path-granularity serialization for tool execution.
|
||||
type pathLockTable struct {
|
||||
mu sync.Mutex // protects the paths map
|
||||
global sync.RWMutex // global barrier: write-locked by "global" scope tools
|
||||
paths map[string]*sync.Mutex // per-path exclusive locks for "write" scope tools
|
||||
}
|
||||
|
||||
// acquire locks the appropriate resources for a tool call and returns a
|
||||
// release function. The caller MUST call release when execution completes.
|
||||
//
|
||||
// Semantics:
|
||||
// - "read": returns immediately (no-op release)
|
||||
// - "write": holds global.RLock + exclusive path lock
|
||||
// - "global" (or empty): holds global.Lock (exclusive — blocks all writes and other globals)
|
||||
func (t *pathLockTable) acquire(scope, path string) func() {
|
||||
switch scope {
|
||||
case "read":
|
||||
return func() {} // reads never conflict
|
||||
|
||||
case "write":
|
||||
if path == "" {
|
||||
// No path extractable — fall through to global barrier.
|
||||
t.global.Lock()
|
||||
return func() { t.global.Unlock() }
|
||||
}
|
||||
// Shared global lock (allows concurrent writes to different paths,
|
||||
// blocks while a global-scope tool is running).
|
||||
t.global.RLock()
|
||||
pl := t.pathLock(path)
|
||||
pl.Lock()
|
||||
return func() {
|
||||
pl.Unlock()
|
||||
t.global.RUnlock()
|
||||
}
|
||||
|
||||
default: // "global", "", or unknown
|
||||
t.global.Lock()
|
||||
return func() { t.global.Unlock() }
|
||||
}
|
||||
}
|
||||
|
||||
// pathLock returns (or creates) the mutex for a given path.
|
||||
func (t *pathLockTable) pathLock(path string) *sync.Mutex {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
if t.paths == nil {
|
||||
t.paths = make(map[string]*sync.Mutex)
|
||||
}
|
||||
l, ok := t.paths[path]
|
||||
if !ok {
|
||||
l = &sync.Mutex{}
|
||||
t.paths[path] = l
|
||||
}
|
||||
return l
|
||||
}
|
||||
|
|
@ -0,0 +1,156 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestPathLock_ReadNeverBlocks(t *testing.T) {
|
||||
var tbl pathLockTable
|
||||
// Multiple concurrent reads should never block each other.
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 10; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
release := tbl.acquire("read", "/foo.go")
|
||||
time.Sleep(1 * time.Millisecond)
|
||||
release()
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func TestPathLock_WritesSerializeOnSamePath(t *testing.T) {
|
||||
var tbl pathLockTable
|
||||
var counter atomic.Int32
|
||||
var maxConcurrent atomic.Int32
|
||||
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 5; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
release := tbl.acquire("write", "/foo.go")
|
||||
cur := counter.Add(1)
|
||||
if cur > maxConcurrent.Load() {
|
||||
maxConcurrent.Store(cur)
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
counter.Add(-1)
|
||||
release()
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if maxConcurrent.Load() != 1 {
|
||||
t.Errorf("expected max concurrency 1 for same path writes, got %d", maxConcurrent.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPathLock_WritesParallelOnDifferentPaths(t *testing.T) {
|
||||
var tbl pathLockTable
|
||||
var counter atomic.Int32
|
||||
var maxConcurrent atomic.Int32
|
||||
|
||||
var wg sync.WaitGroup
|
||||
paths := []string{"/a.go", "/b.go", "/c.go", "/d.go", "/e.go"}
|
||||
for _, p := range paths {
|
||||
wg.Add(1)
|
||||
go func(path string) {
|
||||
defer wg.Done()
|
||||
release := tbl.acquire("write", path)
|
||||
cur := counter.Add(1)
|
||||
for {
|
||||
mc := maxConcurrent.Load()
|
||||
if cur <= mc || maxConcurrent.CompareAndSwap(mc, cur) {
|
||||
break
|
||||
}
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
counter.Add(-1)
|
||||
release()
|
||||
}(p)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if maxConcurrent.Load() < 2 {
|
||||
t.Errorf("expected parallel writes on different paths, got max concurrency %d", maxConcurrent.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPathLock_GlobalBlocksWrites(t *testing.T) {
|
||||
var tbl pathLockTable
|
||||
var sequence []string
|
||||
var mu sync.Mutex
|
||||
|
||||
record := func(s string) {
|
||||
mu.Lock()
|
||||
sequence = append(sequence, s)
|
||||
mu.Unlock()
|
||||
}
|
||||
|
||||
// Start a global lock holder.
|
||||
globalReady := make(chan struct{})
|
||||
globalDone := make(chan struct{})
|
||||
go func() {
|
||||
release := tbl.acquire("global", "")
|
||||
close(globalReady)
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
record("global-done")
|
||||
release()
|
||||
close(globalDone)
|
||||
}()
|
||||
|
||||
<-globalReady
|
||||
// Give the global lock time to be fully acquired.
|
||||
time.Sleep(2 * time.Millisecond)
|
||||
|
||||
// Start a write — should block until global is done.
|
||||
writeDone := make(chan struct{})
|
||||
go func() {
|
||||
release := tbl.acquire("write", "/foo.go")
|
||||
record("write-done")
|
||||
release()
|
||||
close(writeDone)
|
||||
}()
|
||||
|
||||
<-globalDone
|
||||
<-writeDone
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if len(sequence) < 2 || sequence[0] != "global-done" || sequence[1] != "write-done" {
|
||||
t.Errorf("expected global to complete before write, got %v", sequence)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPathLock_WriteNoPathFallsToGlobal(t *testing.T) {
|
||||
var tbl pathLockTable
|
||||
var counter atomic.Int32
|
||||
var maxConcurrent atomic.Int32
|
||||
|
||||
var wg sync.WaitGroup
|
||||
// Multiple writes with empty path should serialize (fall through to global).
|
||||
for i := 0; i < 3; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
release := tbl.acquire("write", "")
|
||||
cur := counter.Add(1)
|
||||
if cur > maxConcurrent.Load() {
|
||||
maxConcurrent.Store(cur)
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
counter.Add(-1)
|
||||
release()
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if maxConcurrent.Load() != 1 {
|
||||
t.Errorf("expected max concurrency 1 for write with no path, got %d", maxConcurrent.Load())
|
||||
}
|
||||
}
|
||||
|
|
@ -31,6 +31,12 @@ type State struct {
|
|||
registry *registry.Registry
|
||||
yolo bool
|
||||
|
||||
// Path-based lock coordination for cross-agent serialization.
|
||||
// All foreground tool calls acquire a lock based on scope+path before
|
||||
// execution, ensuring writes to the same path serialize regardless of
|
||||
// which agent initiated the call.
|
||||
pathLocks pathLockTable
|
||||
|
||||
// Process management
|
||||
procMu sync.Mutex
|
||||
procs map[int]*Proc
|
||||
|
|
@ -301,6 +307,14 @@ func (st *State) NewProc(ctx context.Context, payload string, background bool) (
|
|||
defer close(proc.done)
|
||||
defer cancel()
|
||||
|
||||
// Acquire path-based lock for foreground procs to serialize
|
||||
// conflicting writes across agents. Background procs are
|
||||
// long-running and don't hold locks.
|
||||
if !background {
|
||||
release := st.pathLocks.acquire(info.Scope, args["path"])
|
||||
defer release()
|
||||
}
|
||||
|
||||
out, exitCode := st.executeTool(procCtx, info, args, cwd, envCopy, yolo, timeout, outputWriter, startedCh)
|
||||
|
||||
proc.mu.Lock()
|
||||
|
|
|
|||
Loading…
Reference in New Issue