mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
fix(agent/x/agentmcp): watch MCP config files for late-appearing or rewritten config (#25172)
## Bug `agent/x/agentmcp/Manager` resolves its config paths once at boot, calls `Reload`, and then only re-stats them lazily when a `GET /api/v0/mcp/tools` request arrives (PR #24700). If any of the manager's MCP config files (`~/.mcp.json` by default, or whatever paths `agentcontextconfig.MCPConfigFiles()` resolves from `CODER_AGENT_EXP_MCP_CONFIG_FILES`) is created, atomically rewritten, or removed _after_ that initial `Reload` and _before_ the next tools HTTP request, the manager keeps serving the stale (often empty) snapshot. `parseAndDedup` silently swallows `fs.ErrNotExist`, so a late-appearing file looks indistinguishable from "no config" until something pokes the manager again. This affects any workspace where the file lands after `MarkStartupSettled` fires, including: - a startup script that writes `~/.mcp.json` after MCP init - a user creating or editing the file mid-session - an installer, dotfiles step, or sync tool writing the file later in startup - another agent process (Claude Code, etc.) producing the file out-of-band - editor rewrites that land as `Write + Chmod + Rename` bursts ### Concrete repro (dual-agent workspace) The race is easiest to reproduce on dual-agent workspaces (inner sandbox + outer host), where the inner agent's `scriptRunner` has `script_count=0` and `mcpManager.Reload` fires at ~t+0.3 s while the host agent writes `~/.mcp.json` ~21 s later. Timeline from `workspace-otto-aa16`: - agent up `20:23:38.918` - lifecycle Ready `20:23:39.200` - MCP config file Birth `20:24:00.460` (~21 s gap) - no MCP log lines for 8 minutes - first `GET /api/v0/mcp/tools` at `20:32:11.812` logs `[warn] mcp: mcp reload canceled by caller`, takes 4946 ms; subsequent turns are cached at 2 ms. The single-agent case has the same race; it's just usually narrow enough that the next HTTP request masks it, at the cost of a multi-second stall on the first call that has to do the lazy reload itself. PR #25034's `MarkStartupSettled` does not help: "settled" fires before the file is necessarily on disk. ## Fix Add an fsnotify-backed `configWatcher` to `agent/x/agentmcp/Manager`. The watcher consumes whatever paths the manager is told to reload, which is the same `[]string` returned by `agentcontextconfig.MCPConfigFiles()`. For each path the watcher: - Watches the **parent directory** of the path, not the file itself. This handles late creation, atomic rewrite (rename + create), and deletion uniformly because inotify watches on individual non-existent files return `ENOENT` and are lost across renames. The pattern matches `agent/agentcontainers/watcher`. - Walks up to the first existing **ancestor directory** when the parent does not yet exist, and re-arms deeper on `Create` events that promote an unrealized path. - Refcounts directory watches so multiple configured paths sharing a parent dir only register one inotify watch. - Resolves symlinks **once at arming time** via `filepath.EvalSymlinks`; never chases arbitrary symlink targets on events. - Debounces multi-event editor writes through a single `quartz.AfterFunc` timer so a `Write + Chmod + Rename` burst produces one reload. - Fires a debounced callback that calls `Manager.Reload`, which routes through the existing singleflight so concurrent triggers coalesce. - Re-syncs on every `Reload` call so a future path-list change is picked up. Lifecycle: the watcher is armed lazily on the first `Reload` (no goroutine cost for unit tests that never reload). `Manager.Close` marks the manager closed and closes its `closedCh` before tearing down the watcher, so any in-flight watcher-driven reload observes the close via `waitReload` and returns `ErrManagerClosed` instead of blocking `firesWG.Wait()` on a stuck connect. The watcher then waits for its goroutine and any in-flight debounced `fire` callback before returning. `parseAndDedup` behavior is unchanged: `fs.ErrNotExist` still records an empty snapshot. With the watcher armed before that snapshot is committed, any `Create` event that races `parseAndDedup` is still delivered. This is the agent-side complement to PR #25169, which fixes chatd's mid-turn workspace MCP discovery. `MarkStartupSettled` semantics are not changed. ## Tests (`agent/x/agentmcp/configwatcher_internal_test.go`) All new tests pass with the fix and fail without it. No `time.Sleep`; synchronization uses `testutil.Eventually` for fsnotify-driven assertions and the quartz mock clock for debounce assertions. - `TestWatcher_LateFileTriggersReload` - the late-file regression: empty dir, settle startup, `Reload` sees no file, write the file later, watcher reloads, tools appear. - `TestWatcher_RewriteTriggersReload` - existing file overwritten with a new server list, watcher reloads, cache reflects new server. - `TestWatcher_RemovalTransitionsToEmpty` - delete the file, watcher reloads, manager transitions to empty cleanly. - `TestWatcher_DebouncesBurst` - quartz mock clock; three back-to-back `scheduleFire` calls produce exactly one `onChange` after `AdvanceNext`. - `TestWatcher_CloseStopsGoroutine` - construct/Reload/Close five times to surface goroutine or fd leaks under `-race`. - `TestWatcher_DualAgentHTTPNoStall` - integration: write file after `Reload`, wait for watcher reload, then `GET /tools` returns the MCP tools in less than `testutil.WaitShort` instead of the multi-second "reload canceled" stall. - `TestWatcher_LateParentDirTriggersReload` - parent dir doesn't exist at `Reload` time; create the dir then the file; watcher re-arms deeper and reloads. - `TestWatcher_SharedParentRefcount` - two configured paths share a parent dir; only one inotify watch is registered and both reload on changes. - `TestWatcher_CloseDoesNotStallOnInFlightReload` - installs a `connectStartedHook` to block a watcher-driven reload mid-`connectAll`, then asserts `Close()` returns within `WaitMedium`. Regression-verified: reverting the close-ordering causes the test to time out. ## Acceptance checklist - [x] `go test ./agent/x/agentmcp/... -race -count=1` passes (also `-count=5`). - [x] All `./agent/...` tests pass under `-race -short`. - [x] No emdash, endash, or ` -- `; `scripts/check_emdash.sh` clean. - [x] No `time.Sleep` in tests. - [x] New tests fail without the fix and pass with it (verified by temporarily disabling `m.armWatcher(paths)` and by reverting `Close()` ordering). ## Out of scope No changes to chatd's mid-turn workspace MCP discovery (PR #25169) or `MarkStartupSettled` semantics. --- <sub>This pull request was prepared by a [Coder Agents](https://coder.com/docs/admin/ai-coder) run.</sub>
This commit is contained in:
@@ -0,0 +1,435 @@
|
||||
package agentmcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/fsnotify/fsnotify"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"cdr.dev/slog/v3"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
// defaultWatchDebounce coalesces editor-style multi-event writes
|
||||
// (truncate plus rename plus chmod) into a single reload. The
|
||||
// value is small enough to keep the late-file recovery latency
|
||||
// well under a second.
|
||||
const defaultWatchDebounce = 250 * time.Millisecond
|
||||
|
||||
// configWatcher watches the parent directories of one or more
|
||||
// .mcp.json paths and fires a single debounced callback when any
|
||||
// of those paths is created, modified, removed, or renamed.
|
||||
//
|
||||
// The watcher is deliberately tolerant of late-arriving config:
|
||||
// if the parent directory does not exist yet, it walks up to the
|
||||
// first existing ancestor and re-arms deeper as ancestors appear.
|
||||
// Symlinks are resolved once at arming time; the watcher does not
|
||||
// chase arbitrary symlink targets on every event.
|
||||
type configWatcher struct {
|
||||
logger slog.Logger
|
||||
clock quartz.Clock
|
||||
debounce time.Duration
|
||||
|
||||
// onChange is invoked once per debounce window when a watched
|
||||
// path is touched. It runs on a clock-managed timer goroutine
|
||||
// and must return promptly; callers should hand off to a
|
||||
// singleflight or background goroutine.
|
||||
onChange func()
|
||||
|
||||
mu sync.Mutex
|
||||
watcher *fsnotify.Watcher
|
||||
files map[string]string // resolved path -> watched ancestor dir.
|
||||
dirs map[string]int // ancestor dir -> refcount.
|
||||
timer *quartz.Timer
|
||||
closed bool
|
||||
closedCh chan struct{}
|
||||
closeOnce sync.Once
|
||||
runDoneCh chan struct{} // closed when run() exits.
|
||||
firesWG sync.WaitGroup // tracks in-flight fire callbacks.
|
||||
}
|
||||
|
||||
// newConfigWatcher creates a configWatcher and starts its event
|
||||
// loop. Sync registers the actual paths to watch. The watcher does
|
||||
// nothing until Sync is called.
|
||||
func newConfigWatcher(
|
||||
logger slog.Logger,
|
||||
clock quartz.Clock,
|
||||
debounce time.Duration,
|
||||
onChange func(),
|
||||
) (*configWatcher, error) {
|
||||
if onChange == nil {
|
||||
return nil, xerrors.New("onChange callback is required")
|
||||
}
|
||||
if debounce <= 0 {
|
||||
debounce = defaultWatchDebounce
|
||||
}
|
||||
|
||||
w, err := fsnotify.NewWatcher()
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("create fsnotify watcher: %w", err)
|
||||
}
|
||||
|
||||
cw := &configWatcher{
|
||||
logger: logger,
|
||||
clock: clock,
|
||||
debounce: debounce,
|
||||
onChange: onChange,
|
||||
watcher: w,
|
||||
files: make(map[string]string),
|
||||
dirs: make(map[string]int),
|
||||
closedCh: make(chan struct{}),
|
||||
runDoneCh: make(chan struct{}),
|
||||
}
|
||||
go cw.run()
|
||||
return cw, nil
|
||||
}
|
||||
|
||||
// Sync replaces the watched set with paths. Files no longer in the
|
||||
// list are removed; new files are added. Symlinks are resolved
|
||||
// once. Individual arm failures are logged and skipped; partial
|
||||
// arming is acceptable because parseAndDedup is the source of
|
||||
// truth and the watcher exists purely to trigger a fresh stat.
|
||||
//
|
||||
// Sync is idempotent and safe to call repeatedly.
|
||||
func (cw *configWatcher) Sync(paths []string) {
|
||||
if cw == nil {
|
||||
return
|
||||
}
|
||||
|
||||
resolved := make(map[string]struct{}, len(paths))
|
||||
for _, p := range paths {
|
||||
rp := resolvePath(p)
|
||||
if rp == "" {
|
||||
continue
|
||||
}
|
||||
resolved[rp] = struct{}{}
|
||||
}
|
||||
|
||||
cw.mu.Lock()
|
||||
if cw.closed {
|
||||
cw.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
// Remove paths that are no longer wanted.
|
||||
for rp, dir := range cw.files {
|
||||
if _, keep := resolved[rp]; keep {
|
||||
continue
|
||||
}
|
||||
delete(cw.files, rp)
|
||||
cw.releaseDirLocked(dir)
|
||||
}
|
||||
|
||||
// Add new paths.
|
||||
for rp := range resolved {
|
||||
if _, already := cw.files[rp]; already {
|
||||
continue
|
||||
}
|
||||
dir, err := cw.armAncestorLocked(rp)
|
||||
if err != nil {
|
||||
cw.logger.Warn(context.Background(),
|
||||
"failed to arm config file watch",
|
||||
slog.F("path", rp), slog.Error(err))
|
||||
continue
|
||||
}
|
||||
cw.files[rp] = dir
|
||||
}
|
||||
cw.mu.Unlock()
|
||||
}
|
||||
|
||||
// armAncestorLocked walks up the parent chain from rp until it
|
||||
// finds an existing directory, then watches that directory.
|
||||
// Returns the actual directory it ended up watching. The last
|
||||
// fsnotify Add error is preserved so callers can distinguish a
|
||||
// missing-ancestor failure from an inotify-limit (ENOSPC) failure.
|
||||
// Callers must hold cw.mu.
|
||||
func (cw *configWatcher) armAncestorLocked(rp string) (string, error) {
|
||||
dir := filepath.Dir(rp)
|
||||
var lastAddDir string
|
||||
var lastAddErr error
|
||||
for {
|
||||
// Bail out if we somehow reached the root without finding
|
||||
// an existing directory. filepath.Dir("/") == "/" on POSIX
|
||||
// and "C:\" == "C:\" on Windows, so guard against an
|
||||
// infinite loop.
|
||||
if dir == "" || dir == "." {
|
||||
return "", noAncestorErr(rp, lastAddDir, lastAddErr)
|
||||
}
|
||||
|
||||
if cw.dirs[dir] > 0 {
|
||||
cw.dirs[dir]++
|
||||
return dir, nil
|
||||
}
|
||||
|
||||
err := cw.watcher.Add(dir)
|
||||
if err == nil {
|
||||
cw.dirs[dir] = 1
|
||||
return dir, nil
|
||||
}
|
||||
lastAddDir = dir
|
||||
lastAddErr = err
|
||||
|
||||
parent := filepath.Dir(dir)
|
||||
if parent == dir {
|
||||
return "", noAncestorErr(rp, lastAddDir, lastAddErr)
|
||||
}
|
||||
dir = parent
|
||||
}
|
||||
}
|
||||
|
||||
// noAncestorErr formats the failure to register a watch on any
|
||||
// ancestor of path. If the loop tried at least one Add, the
|
||||
// underlying error (usually inotify ENOSPC) is wrapped so the
|
||||
// operator sees the actual kernel-level cause instead of a generic
|
||||
// "no existing ancestor" message.
|
||||
func noAncestorErr(path, lastDir string, lastErr error) error {
|
||||
if lastErr != nil {
|
||||
return xerrors.Errorf("cannot watch any ancestor of %q (last attempt on %q): %w", path, lastDir, lastErr)
|
||||
}
|
||||
return xerrors.Errorf("no existing ancestor for %q", path)
|
||||
}
|
||||
|
||||
// releaseDirLocked decrements the refcount for dir and removes the
|
||||
// watch when no remaining file points at it. Callers must hold
|
||||
// cw.mu.
|
||||
func (cw *configWatcher) releaseDirLocked(dir string) {
|
||||
cw.dirs[dir]--
|
||||
if cw.dirs[dir] > 0 {
|
||||
return
|
||||
}
|
||||
delete(cw.dirs, dir)
|
||||
if err := cw.watcher.Remove(dir); err != nil {
|
||||
// Removal can fail when the directory no longer exists;
|
||||
// fsnotify already dropped the watch, so this is benign.
|
||||
cw.logger.Debug(context.Background(),
|
||||
"failed to remove config dir watch",
|
||||
slog.F("dir", dir), slog.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
// run is the watcher loop. It exits when the underlying
|
||||
// fsnotify.Watcher closes its channels or Close is called.
|
||||
func (cw *configWatcher) run() {
|
||||
defer close(cw.runDoneCh)
|
||||
ctx := context.Background()
|
||||
for {
|
||||
select {
|
||||
case <-cw.closedCh:
|
||||
return
|
||||
case evt, ok := <-cw.watcher.Events:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
cw.handleEvent(ctx, evt)
|
||||
case err, ok := <-cw.watcher.Errors:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
cw.logger.Warn(ctx,
|
||||
"fsnotify watch error; config file changes may not be detected until the next HTTP request",
|
||||
slog.Error(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// handleEvent decides whether the event concerns one of the
|
||||
// watched files (or could promote an ancestor watch) and, if so,
|
||||
// schedules a debounced fire.
|
||||
func (cw *configWatcher) handleEvent(ctx context.Context, evt fsnotify.Event) {
|
||||
cw.mu.Lock()
|
||||
if cw.closed {
|
||||
cw.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
// Match against any watched file. fsnotify event names are
|
||||
// already absolute when the watched directory is absolute,
|
||||
// which it is because armAncestorLocked called filepath.Dir
|
||||
// on a path resolved to absolute. The filepath.Abs call below
|
||||
// is a defensive normalization.
|
||||
evtAbs, err := filepath.Abs(evt.Name)
|
||||
if err != nil {
|
||||
cw.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
matchedFile := ""
|
||||
for rp := range cw.files {
|
||||
if rp == evtAbs {
|
||||
matchedFile = rp
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// If a directory we are watching for an ancestor of an
|
||||
// unrealized path just gained a new child, try to re-arm
|
||||
// deeper. This handles `mkdir ~/.config; touch
|
||||
// ~/.config/.mcp.json` cases.
|
||||
if matchedFile == "" && evt.Has(fsnotify.Create) {
|
||||
for rp, dir := range cw.files {
|
||||
// Only re-arm files whose final parent is not yet
|
||||
// being watched directly.
|
||||
expected := filepath.Dir(rp)
|
||||
if dir == expected {
|
||||
continue
|
||||
}
|
||||
// If this event is a directory inside our currently
|
||||
// watched ancestor that lies on the way to rp,
|
||||
// re-arm.
|
||||
if isAncestorPathSegment(evtAbs, rp) {
|
||||
cw.releaseDirLocked(dir)
|
||||
newDir, armErr := cw.armAncestorLocked(rp)
|
||||
if armErr != nil {
|
||||
cw.logger.Debug(ctx,
|
||||
"failed to re-arm config file watch on ancestor create",
|
||||
slog.F("path", rp), slog.Error(armErr))
|
||||
// Leave the file unarmed for now;
|
||||
// next Sync will retry.
|
||||
delete(cw.files, rp)
|
||||
continue
|
||||
}
|
||||
cw.files[rp] = newDir
|
||||
// The new dir may already contain the
|
||||
// target file. Treat that as a match.
|
||||
matchedFile = rp
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
cw.mu.Unlock()
|
||||
|
||||
if matchedFile == "" {
|
||||
return
|
||||
}
|
||||
cw.scheduleFire()
|
||||
}
|
||||
|
||||
// isAncestorPathSegment reports whether candidate is on the path
|
||||
// from the currently watched ancestor toward target.
|
||||
func isAncestorPathSegment(candidate, target string) bool {
|
||||
// candidate must be a prefix of target's directory chain.
|
||||
tdir := filepath.Dir(target)
|
||||
for {
|
||||
if tdir == candidate {
|
||||
return true
|
||||
}
|
||||
parent := filepath.Dir(tdir)
|
||||
if parent == tdir {
|
||||
return false
|
||||
}
|
||||
tdir = parent
|
||||
}
|
||||
}
|
||||
|
||||
// scheduleFire arms or extends a single debounce timer.
|
||||
func (cw *configWatcher) scheduleFire() {
|
||||
cw.mu.Lock()
|
||||
defer cw.mu.Unlock()
|
||||
if cw.closed {
|
||||
return
|
||||
}
|
||||
if cw.timer != nil {
|
||||
// Reset existing timer to extend the debounce window.
|
||||
// Stop reports whether the call stopped the timer before
|
||||
// it fired; if so we owe a Done because Add was called
|
||||
// when the timer was created.
|
||||
if cw.timer.Stop() {
|
||||
cw.firesWG.Done()
|
||||
}
|
||||
}
|
||||
cw.firesWG.Add(1)
|
||||
cw.timer = cw.clock.AfterFunc(cw.debounce, cw.fire, "agentmcp", "watch_debounce")
|
||||
}
|
||||
|
||||
// fire is called once per debounce window. It invokes onChange
|
||||
// outside the lock so reload code can re-enter Sync safely.
|
||||
func (cw *configWatcher) fire() {
|
||||
defer cw.firesWG.Done()
|
||||
|
||||
cw.mu.Lock()
|
||||
if cw.closed {
|
||||
cw.mu.Unlock()
|
||||
return
|
||||
}
|
||||
cw.timer = nil
|
||||
cw.mu.Unlock()
|
||||
|
||||
cw.onChange()
|
||||
}
|
||||
|
||||
// Close stops the watcher and waits for the run goroutine and
|
||||
// any in-flight debounced fire callbacks to exit. Close is
|
||||
// idempotent.
|
||||
func (cw *configWatcher) Close() error {
|
||||
if cw == nil {
|
||||
return nil
|
||||
}
|
||||
var closeErr error
|
||||
cw.closeOnce.Do(func() {
|
||||
cw.mu.Lock()
|
||||
cw.closed = true
|
||||
if cw.timer != nil {
|
||||
// Stop returns true if the call prevented the timer
|
||||
// callback from running. Account for the Add() that
|
||||
// scheduleFire performed when arming this timer.
|
||||
if cw.timer.Stop() {
|
||||
cw.firesWG.Done()
|
||||
}
|
||||
cw.timer = nil
|
||||
}
|
||||
cw.mu.Unlock()
|
||||
|
||||
close(cw.closedCh)
|
||||
if err := cw.watcher.Close(); err != nil {
|
||||
closeErr = xerrors.Errorf("close fsnotify watcher: %w", err)
|
||||
}
|
||||
// Wait for run() to exit, then wait for any in-flight
|
||||
// fire callback to return. Callers should not observe a
|
||||
// stale onChange after Close returns; this is critical
|
||||
// for tests that use slogtest, which panics on log
|
||||
// calls made after the test has finished.
|
||||
<-cw.runDoneCh
|
||||
cw.firesWG.Wait()
|
||||
})
|
||||
return closeErr
|
||||
}
|
||||
|
||||
// resolvePath converts a path to an absolute, symlink-resolved
|
||||
// form. If the file does not exist, falls back to filepath.Abs so
|
||||
// the caller can still arm an ancestor directory.
|
||||
func resolvePath(p string) string {
|
||||
if p == "" {
|
||||
return ""
|
||||
}
|
||||
if abs, err := filepath.Abs(p); err == nil {
|
||||
// EvalSymlinks fails on non-existent paths. Resolve as
|
||||
// far as possible without erroring out: walk up until
|
||||
// we find an existing ancestor, eval its symlinks, and
|
||||
// re-join the trailing segments.
|
||||
if resolved, err := filepath.EvalSymlinks(abs); err == nil {
|
||||
return resolved
|
||||
}
|
||||
return resolvePathBestEffort(abs)
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func resolvePathBestEffort(abs string) string {
|
||||
dir := filepath.Dir(abs)
|
||||
base := filepath.Base(abs)
|
||||
for dir != "" && dir != "." {
|
||||
if resolved, err := filepath.EvalSymlinks(dir); err == nil {
|
||||
return filepath.Join(resolved, base)
|
||||
}
|
||||
parent := filepath.Dir(dir)
|
||||
base = filepath.Join(filepath.Base(dir), base)
|
||||
if parent == dir {
|
||||
break
|
||||
}
|
||||
dir = parent
|
||||
}
|
||||
return abs
|
||||
}
|
||||
@@ -0,0 +1,514 @@
|
||||
package agentmcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"cdr.dev/slog/v3"
|
||||
"cdr.dev/slog/v3/sloggers/slogtest"
|
||||
"github.com/coder/coder/v2/agent/agentexec"
|
||||
"github.com/coder/coder/v2/codersdk/workspacesdk"
|
||||
"github.com/coder/coder/v2/testutil"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
// These tests exercise the dual-agent late-file regression: the
|
||||
// inner sandbox agent settles startup quickly and calls Reload
|
||||
// while `~/.mcp.json` still does not exist on disk. The host
|
||||
// agent then writes the file ~20s later. Before this fix, the
|
||||
// manager cached an empty snapshot and stayed empty until a
|
||||
// subsequent HTTP call lazily re-statted the file. With the
|
||||
// fsnotify-backed watcher, the manager picks up the late file
|
||||
// without external prompting.
|
||||
|
||||
// awaitTools polls cachedTools until the predicate succeeds or
|
||||
// the context expires. It avoids time.Sleep loops in callers.
|
||||
func awaitTools(ctx context.Context, t *testing.T, m *Manager, pred func([]workspacesdk.MCPToolInfo) bool) []workspacesdk.MCPToolInfo {
|
||||
t.Helper()
|
||||
var final []workspacesdk.MCPToolInfo
|
||||
testutil.Eventually(ctx, t, func(context.Context) bool {
|
||||
final = m.cachedTools()
|
||||
return pred(final)
|
||||
}, testutil.IntervalFast)
|
||||
return final
|
||||
}
|
||||
|
||||
// useFastDebounce shortens the watcher's debounce window so
|
||||
// real-clock tests do not stall on the 250 ms default. Must be
|
||||
// called before any Reload arms the watcher.
|
||||
func useFastDebounce(t *testing.T, m *Manager) {
|
||||
t.Helper()
|
||||
m.mu.Lock()
|
||||
m.watchDebounce = 10 * time.Millisecond
|
||||
m.mu.Unlock()
|
||||
}
|
||||
|
||||
func TestWatcher_LateFileTriggersReload(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, ".mcp.json")
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
// First Reload arms the watcher but finds nothing on disk.
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
require.Empty(t, m.cachedTools(), "manager should start with no tools")
|
||||
|
||||
// Write the file after the manager has already settled. The
|
||||
// watcher must observe the Create event, debounce it, and
|
||||
// trigger a fresh Reload without any external HTTP call.
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
tools := awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 1
|
||||
})
|
||||
require.Len(t, tools, 1)
|
||||
assert.Contains(t, tools[0].Name, "echo")
|
||||
|
||||
// The snapshot must now reflect the on-disk file so the
|
||||
// next Reload short-circuits.
|
||||
assert.False(t, m.SnapshotChanged([]string{configPath}))
|
||||
}
|
||||
|
||||
func TestWatcher_RewriteTriggersReload(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
configPath := writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
tools := m.cachedTools()
|
||||
require.Len(t, tools, 1)
|
||||
assert.Contains(t, tools[0].Name, "srv")
|
||||
|
||||
// Overwrite the config with a different server name. The
|
||||
// watcher should fire and the cache should reflect the new
|
||||
// server without any caller-driven Reload.
|
||||
_, entry2 := fakeMCPServerConfig(t, "srv2")
|
||||
writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv2": entry2})
|
||||
|
||||
tools = awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 1 && len(tools[0].Name) > 0 &&
|
||||
(tools[0].ServerName == "srv2")
|
||||
})
|
||||
require.Len(t, tools, 1)
|
||||
assert.Equal(t, "srv2", tools[0].ServerName)
|
||||
}
|
||||
|
||||
func TestWatcher_RemovalTransitionsToEmpty(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
configPath := writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
require.Len(t, m.cachedTools(), 1)
|
||||
|
||||
require.NoError(t, os.Remove(configPath))
|
||||
|
||||
awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 0
|
||||
})
|
||||
assert.Empty(t, m.cachedTools())
|
||||
}
|
||||
|
||||
// TestWatcher_DebouncesBurst uses the quartz mock clock to
|
||||
// confirm that three writes inside a single debounce window
|
||||
// produce exactly one onChange invocation. This is the
|
||||
// guarantee that lets the watcher coalesce editor-style
|
||||
// multi-event writes (write + chmod + rename) into a single
|
||||
// Reload.
|
||||
func TestWatcher_DebouncesBurst(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
mClock := quartz.NewMock(t)
|
||||
|
||||
var fires atomic.Int64
|
||||
fired := make(chan struct{}, 4)
|
||||
cw, err := newConfigWatcher(logger, mClock, 100*time.Millisecond, func() {
|
||||
fires.Add(1)
|
||||
fired <- struct{}{}
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = cw.Close() })
|
||||
|
||||
dir := t.TempDir()
|
||||
target := filepath.Join(dir, ".mcp.json")
|
||||
cw.Sync([]string{target})
|
||||
|
||||
// First burst: simulate three fsnotify events landing within
|
||||
// the debounce window. We do this by directly calling
|
||||
// scheduleFire, which is exactly what handleEvent does for
|
||||
// each matching event.
|
||||
cw.scheduleFire()
|
||||
cw.scheduleFire()
|
||||
cw.scheduleFire()
|
||||
|
||||
// Before the timer fires, no callback should have run.
|
||||
require.Equal(t, int64(0), fires.Load())
|
||||
|
||||
// Advance past the debounce window. Only one fire is
|
||||
// expected because all three scheduleFire calls reused the
|
||||
// same timer.
|
||||
_, waiter := mClock.AdvanceNext()
|
||||
waiter.MustWait(testutil.Context(t, testutil.WaitShort))
|
||||
|
||||
select {
|
||||
case <-fired:
|
||||
case <-time.After(testutil.WaitShort):
|
||||
t.Fatal("expected one fire after debounce window")
|
||||
}
|
||||
|
||||
// Drain any spurious extra fire briefly.
|
||||
select {
|
||||
case <-fired:
|
||||
t.Fatal("unexpected additional fire within debounce window")
|
||||
default:
|
||||
}
|
||||
require.Equal(t, int64(1), fires.Load())
|
||||
|
||||
// A second burst after the first window settles must fire
|
||||
// again (debounce per-window, not global).
|
||||
cw.scheduleFire()
|
||||
cw.scheduleFire()
|
||||
_, waiter = mClock.AdvanceNext()
|
||||
waiter.MustWait(testutil.Context(t, testutil.WaitShort))
|
||||
|
||||
select {
|
||||
case <-fired:
|
||||
case <-time.After(testutil.WaitShort):
|
||||
t.Fatal("expected fire after second window")
|
||||
}
|
||||
require.Equal(t, int64(2), fires.Load())
|
||||
}
|
||||
|
||||
// TestWatcher_CloseStopsGoroutine asserts that Close releases the
|
||||
// fsnotify watcher fd and stops its goroutine. We rely on the
|
||||
// race detector and on creating a fresh manager on the same path
|
||||
// to surface fd or goroutine leaks.
|
||||
func TestWatcher_CloseStopsGoroutine(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, ".mcp.json")
|
||||
|
||||
for range 5 {
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
require.NoError(t, m.Close())
|
||||
|
||||
// After Close the watcher field is cleared and the
|
||||
// fsnotify watcher is shut down.
|
||||
m.mu.RLock()
|
||||
w := m.watcher
|
||||
m.mu.RUnlock()
|
||||
require.Nil(t, w, "watcher must be nil after Close")
|
||||
}
|
||||
}
|
||||
|
||||
// TestWatcher_DualAgentHTTPNoStall mimics the dual-agent
|
||||
// workspace scenario from workspace-otto-aa16: the inner sandbox
|
||||
// agent calls MarkStartupSettled and Reload while the host agent
|
||||
// has not yet written ~/.mcp.json. Once the file appears, an
|
||||
// HTTP request to /tools must return the MCP tools quickly
|
||||
// instead of triggering a multi-second "reload canceled" stall.
|
||||
func TestWatcher_DualAgentHTTPNoStall(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, ".mcp.json")
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
// First Reload races ahead of the host agent: empty config.
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
require.Empty(t, m.cachedTools())
|
||||
|
||||
api := NewAPI(logger, m, func() []string { return []string{configPath} })
|
||||
|
||||
// Host agent writes the file later.
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
// Wait for the watcher to pick up the file so we know the
|
||||
// cache is warm before issuing the HTTP request.
|
||||
awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 1
|
||||
})
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/tools", nil).WithContext(ctx)
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
start := time.Now()
|
||||
api.Routes().ServeHTTP(rec, req)
|
||||
elapsed := time.Since(start)
|
||||
|
||||
require.Equal(t, http.StatusOK, rec.Code)
|
||||
require.Less(t, elapsed, testutil.WaitShort,
|
||||
"warm HTTP request should not stall on watcher reload; took %s", elapsed)
|
||||
|
||||
var resp workspacesdk.ListMCPToolsResponse
|
||||
require.NoError(t, json.NewDecoder(rec.Body).Decode(&resp))
|
||||
require.Len(t, resp.Tools, 1)
|
||||
assert.Contains(t, resp.Tools[0].Name, "echo")
|
||||
}
|
||||
|
||||
// TestWatcher_LateParentDirTriggersReload exercises the
|
||||
// ancestor-walk-up branch (handleEvent re-arm path,
|
||||
// armAncestorLocked walk-up). The watcher is started with the
|
||||
// final parent directory missing; once that directory is
|
||||
// created, the watcher must promote its watch deeper and then
|
||||
// fire on the file write.
|
||||
func TestWatcher_LateParentDirTriggersReload(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
root := t.TempDir()
|
||||
// Parent directory does not exist yet: armAncestorLocked
|
||||
// will watch root instead.
|
||||
missing := filepath.Join(root, "config")
|
||||
configPath := filepath.Join(missing, ".mcp.json")
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
require.Empty(t, m.cachedTools())
|
||||
|
||||
// Create the missing parent directory. fsnotify will deliver
|
||||
// a Create event on root; handleEvent must release the root
|
||||
// watch, re-arm on the new parent, and schedule a reload.
|
||||
require.NoError(t, os.MkdirAll(missing, 0o755))
|
||||
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
writeMCPConfig(t, missing, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
tools := awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 1
|
||||
})
|
||||
require.Len(t, tools, 1)
|
||||
assert.Contains(t, tools[0].Name, "echo")
|
||||
}
|
||||
|
||||
// TestWatcher_SharedParentRefcount covers the multi-path
|
||||
// directory-watch refcount path: two configured paths in the
|
||||
// same parent dir should produce a single fsnotify watch, and
|
||||
// removing one path via a subsequent Sync must keep the
|
||||
// remaining path armed.
|
||||
func TestWatcher_SharedParentRefcount(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
pathA := filepath.Join(dir, "a.mcp.json")
|
||||
pathB := filepath.Join(dir, "b.mcp.json")
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
t.Cleanup(func() { _ = m.Close() })
|
||||
|
||||
// First Reload arms both paths, sharing the dir watch.
|
||||
require.NoError(t, m.Reload(ctx, []string{pathA, pathB}))
|
||||
|
||||
m.mu.RLock()
|
||||
w := m.watcher
|
||||
m.mu.RUnlock()
|
||||
require.NotNil(t, w, "watcher must be armed")
|
||||
|
||||
w.mu.Lock()
|
||||
require.Equal(t, 2, len(w.files), "two files tracked")
|
||||
require.Equal(t, 1, len(w.dirs), "shared parent dir")
|
||||
require.Equal(t, 2, w.dirs[dir], "refcount equals number of files")
|
||||
w.mu.Unlock()
|
||||
|
||||
// Second Reload removes pathB, so the dir refcount drops to
|
||||
// 1 but the watch must remain in place for pathA.
|
||||
require.NoError(t, m.Reload(ctx, []string{pathA}))
|
||||
|
||||
w.mu.Lock()
|
||||
require.Equal(t, 1, len(w.files), "one file tracked after removal")
|
||||
require.Equal(t, 1, w.dirs[dir], "refcount decremented but not zero")
|
||||
w.mu.Unlock()
|
||||
|
||||
// Writing pathA should still trigger a reload via the
|
||||
// surviving dir watch.
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
cfg := mcpConfigFile{MCPServers: make(map[string]json.RawMessage)}
|
||||
raw, err := json.Marshal(entry)
|
||||
require.NoError(t, err)
|
||||
cfg.MCPServers["srv"] = raw
|
||||
data, err := json.Marshal(cfg)
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, os.WriteFile(pathA, data, 0o600))
|
||||
|
||||
tools := awaitTools(ctx, t, m, func(tools []workspacesdk.MCPToolInfo) bool {
|
||||
return len(tools) == 1
|
||||
})
|
||||
require.Len(t, tools, 1)
|
||||
}
|
||||
|
||||
// TestWatcher_CloseDoesNotStallOnInFlightReload guards the
|
||||
// shutdown-ordering invariant: Close() must mark the manager
|
||||
// closed before w.Close() so an in-flight watcher-driven Reload
|
||||
// short-circuits instead of blocking firesWG.Wait() for the full
|
||||
// connect timeout. Without the ordering, this test would block
|
||||
// at Close() for ~30 s.
|
||||
//
|
||||
// The test installs a connectStartedHook that signals when a
|
||||
// watcher-driven reload has reached connectAll and then blocks
|
||||
// until released. While the hook is blocking the singleflight
|
||||
// reload goroutine, the test calls Close() and asserts it
|
||||
// returns quickly: the DEREM-5 ordering ensures m.closedCh is
|
||||
// closed before w.Close()'s firesWG.Wait(), so waitReload
|
||||
// observes the close, fire() returns, and firesWG drains. If
|
||||
// the ordering is reverted, w.Close() blocks on firesWG.Wait()
|
||||
// while fire() is stuck inside waitReload waiting for the
|
||||
// connect that will never finish.
|
||||
func TestWatcher_CloseDoesNotStallOnInFlightReload(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if os.Getenv("TEST_MCP_FAKE_SERVER") == "1" {
|
||||
runFakeMCPServer()
|
||||
return
|
||||
}
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
logger := slogtest.Make(t, nil).Leveled(slog.LevelDebug)
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, ".mcp.json")
|
||||
|
||||
m := NewManager(ctx, logger, agentexec.DefaultExecer, nil)
|
||||
useFastDebounce(t, m)
|
||||
m.MarkStartupSettled()
|
||||
|
||||
// Arm the watcher with an initial empty Reload. We install the
|
||||
// hook after this so the first connectAll (with empty
|
||||
// toConnect) is not blocked.
|
||||
require.NoError(t, m.Reload(ctx, []string{configPath}))
|
||||
|
||||
reached := make(chan struct{})
|
||||
release := make(chan struct{})
|
||||
var releaseOnce sync.Once
|
||||
releaseHook := func() { releaseOnce.Do(func() { close(release) }) }
|
||||
t.Cleanup(releaseHook)
|
||||
|
||||
m.mu.Lock()
|
||||
var hookOnce sync.Once
|
||||
m.connectStartedHook = func() {
|
||||
hookOnce.Do(func() { close(reached) })
|
||||
<-release
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
// Write the file. The watcher will fire a debounced reload
|
||||
// that hits the connectStartedHook and blocks there.
|
||||
_, entry := fakeMCPServerConfig(t, "srv")
|
||||
writeMCPConfig(t, dir, map[string]mcpServerEntry{"srv": entry})
|
||||
|
||||
select {
|
||||
case <-reached:
|
||||
case <-time.After(testutil.WaitLong):
|
||||
t.Fatal("watcher-driven reload never reached connectAll")
|
||||
}
|
||||
|
||||
// Reload is in-flight: connectAll is blocked inside the hook,
|
||||
// the singleflight body has not returned, and fire() is
|
||||
// blocked in waitReload. Now call Close. With the correct
|
||||
// ordering (m.closedCh closed before w.Close()), this returns
|
||||
// quickly even though the hook is still blocking.
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- m.Close() }()
|
||||
|
||||
select {
|
||||
case err := <-done:
|
||||
require.NoError(t, err)
|
||||
case <-time.After(testutil.WaitMedium):
|
||||
t.Fatal("Close stalled; ordering bug: w.Close before m.closed=true")
|
||||
}
|
||||
|
||||
// Release the hook so the leaked singleflight goroutine can
|
||||
// drain. The manager is already closed, so its work has no
|
||||
// observable effect.
|
||||
releaseHook()
|
||||
}
|
||||
+135
-3
@@ -106,6 +106,28 @@ type Manager struct {
|
||||
// caller and may outlive Close).
|
||||
closedCh chan struct{}
|
||||
closeOnce sync.Once
|
||||
|
||||
// lastPaths records the most recent config paths passed to
|
||||
// Reload/Tools. The fsnotify-backed watcher uses these to
|
||||
// drive its own reloads when ~/.mcp.json appears late on
|
||||
// dual-agent workspaces.
|
||||
lastPaths []string
|
||||
|
||||
// watcher fires a debounced Reload when any watched config
|
||||
// file is created, written, removed, or renamed. It is armed
|
||||
// lazily on the first Reload call so tests that never call
|
||||
// Reload do not pay for an extra goroutine and file
|
||||
// descriptor.
|
||||
watcher *configWatcher
|
||||
watcherOnce sync.Once
|
||||
watchDebounce time.Duration
|
||||
|
||||
// connectStartedHook is a test hook invoked at the start of
|
||||
// connectAll, before any client is dialed. Production code
|
||||
// leaves this nil; tests set it to coordinate with an
|
||||
// in-flight reload (for example, to verify Close()'s
|
||||
// shutdown ordering does not stall on a stuck connect).
|
||||
connectStartedHook func()
|
||||
}
|
||||
|
||||
// serverEntry pairs a server config with its connected client.
|
||||
@@ -136,6 +158,7 @@ func NewManager(
|
||||
snapshot: make(map[string]fileSnapshot),
|
||||
startupSettled: make(chan struct{}),
|
||||
closedCh: make(chan struct{}),
|
||||
watchDebounce: defaultWatchDebounce,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -258,6 +281,14 @@ func (m *Manager) startReloadIfNeeded(paths []string) (<-chan reloadResult, bool
|
||||
}
|
||||
return nil, false, err
|
||||
}
|
||||
// Arm the fsnotify watcher before deciding whether to short
|
||||
// circuit. The first call lazily creates it; subsequent calls
|
||||
// re-sync the watched path set if it changed. Arming before
|
||||
// the SnapshotChanged check ensures any Create event that
|
||||
// races with parseAndDedup is still delivered: the watcher
|
||||
// is running when parseAndDedup returns the empty snapshot.
|
||||
m.armWatcher(paths)
|
||||
|
||||
if firstSyncSettled && !m.SnapshotChanged(paths) {
|
||||
return nil, false, nil
|
||||
}
|
||||
@@ -270,6 +301,82 @@ func (m *Manager) startReloadIfNeeded(paths []string) (<-chan reloadResult, bool
|
||||
return ch, true, nil
|
||||
}
|
||||
|
||||
// armWatcher lazily initializes the fsnotify-backed configWatcher
|
||||
// and syncs it to the latest config paths. Lazy initialization
|
||||
// keeps unit tests that never call Reload free of extra goroutines
|
||||
// and file descriptors.
|
||||
//
|
||||
// If the underlying watcher cannot be created (e.g. inotify limit
|
||||
// reached), the error is logged once and the manager continues
|
||||
// without a watcher. The lazy stat-on-request path remains the
|
||||
// primary mechanism; the watcher is an optimization that closes
|
||||
// the dual-agent race window.
|
||||
func (m *Manager) armWatcher(paths []string) {
|
||||
m.watcherOnce.Do(func() {
|
||||
cw, err := newConfigWatcher(
|
||||
m.logger.Named("config_watcher"),
|
||||
m.clock,
|
||||
m.watchDebounce,
|
||||
m.handleWatchedConfigChange,
|
||||
)
|
||||
if err != nil {
|
||||
m.logger.Warn(m.ctx,
|
||||
"failed to start MCP config watcher; falling back to lazy stat",
|
||||
slog.Error(err))
|
||||
return
|
||||
}
|
||||
// Close the watcher if the manager was closed between
|
||||
// newConfigWatcher returning and us acquiring m.mu.
|
||||
// Otherwise its goroutine and inotify fd leak.
|
||||
m.mu.Lock()
|
||||
if m.closed {
|
||||
m.mu.Unlock()
|
||||
_ = cw.Close()
|
||||
return
|
||||
}
|
||||
m.watcher = cw
|
||||
m.mu.Unlock()
|
||||
})
|
||||
|
||||
m.mu.Lock()
|
||||
m.lastPaths = slices.Clone(paths)
|
||||
w := m.watcher
|
||||
closed := m.closed
|
||||
m.mu.Unlock()
|
||||
if w == nil || closed {
|
||||
return
|
||||
}
|
||||
w.Sync(paths)
|
||||
}
|
||||
|
||||
// handleWatchedConfigChange is invoked by the watcher on a
|
||||
// debounced fire. It triggers a singleflight Reload using the
|
||||
// most recently observed path set so the cached server map and
|
||||
// snapshot are refreshed without waiting for the next HTTP
|
||||
// request.
|
||||
func (m *Manager) handleWatchedConfigChange() {
|
||||
m.mu.RLock()
|
||||
paths := slices.Clone(m.lastPaths)
|
||||
closed := m.closed
|
||||
m.mu.RUnlock()
|
||||
if closed || len(paths) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
logger := m.logger.With(slog.F("trigger", "fsnotify"))
|
||||
logger.Debug(m.ctx, "reloading due to config change")
|
||||
if err := m.Reload(m.ctx, paths); err != nil {
|
||||
if errors.Is(err, ErrManagerClosed) ||
|
||||
errors.Is(err, context.Canceled) {
|
||||
logger.Debug(m.ctx,
|
||||
"watched reload short-circuited by shutdown",
|
||||
slog.Error(err))
|
||||
return
|
||||
}
|
||||
logger.Warn(m.ctx, "watched reload failed", slog.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Manager) waitReload(ctx context.Context, ch <-chan reloadResult, timeout time.Duration) error {
|
||||
// Prefer caller cancellation when it already happened before the
|
||||
// wait. Otherwise select may choose a ready reload result instead.
|
||||
@@ -515,6 +622,10 @@ func (m *Manager) classifyServers(wanted map[string]ServerConfig) (*serverDiff,
|
||||
func (m *Manager) connectAll(ctx context.Context, toConnect []ServerConfig) []connectedServer {
|
||||
logger := m.logger.With(agentchat.Fields(ctx)...)
|
||||
|
||||
if hook := m.connectStartedHook; hook != nil {
|
||||
hook()
|
||||
}
|
||||
|
||||
var (
|
||||
mu sync.Mutex
|
||||
connected []connectedServer
|
||||
@@ -756,13 +867,34 @@ func (m *Manager) RefreshTools(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// Close terminates all MCP server connections and child
|
||||
// processes.
|
||||
// processes, stops the config file watcher, and waits for any
|
||||
// in-flight watcher-driven reload to complete.
|
||||
func (m *Manager) Close() error {
|
||||
// Mark the manager closed and signal closedCh first, then
|
||||
// hand the watcher off and release the lock. Marking closed
|
||||
// before w.Close() ensures that any in-flight
|
||||
// handleWatchedConfigChange short-circuits and any Reload
|
||||
// blocked in waitReload observes m.closedCh, instead of
|
||||
// blocking firesWG.Wait() inside w.Close() until a 30 s
|
||||
// connectAll times out.
|
||||
m.mu.Lock()
|
||||
m.closed = true
|
||||
m.closeOnce.Do(func() { close(m.closedCh) })
|
||||
w := m.watcher
|
||||
m.watcher = nil
|
||||
m.mu.Unlock()
|
||||
|
||||
// Close the watcher outside the manager lock. Its goroutine
|
||||
// may call handleWatchedConfigChange, which takes m.mu, so
|
||||
// holding m.mu while waiting for the watcher to drain would
|
||||
// deadlock. Close on a nil watcher is a no-op.
|
||||
if w != nil {
|
||||
_ = w.Close()
|
||||
}
|
||||
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
m.closed = true
|
||||
m.closeOnce.Do(func() { close(m.closedCh) })
|
||||
var errs []error
|
||||
for _, entry := range m.servers {
|
||||
if err := entry.client.Close(); err != nil {
|
||||
|
||||
Reference in New Issue
Block a user