From ca2e2cc268da7ba8e1b62ae12771e01d03103641 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 19 Sep 2026 02:39:47 -0700 Subject: [PATCH] improve(skills): reduce repeated scans for missing folders (#152563) * improve(skills): reduce repeated scans for missing folders * test(skills): normalize Windows watcher fixture paths * fix(skills): preserve watcher snapshots and portable lifecycle proof * test(skills): retain POSIX ancestor rename coverage * fix(skills): retain ancestor coverage after root promotion * test(skills): avoid returning timer from promise executor --- src/skills/runtime/refresh-ancestor-watch.ts | 163 ++++++ src/skills/runtime/refresh-file-stability.ts | 4 +- src/skills/runtime/refresh-watch-path.ts | 27 + src/skills/runtime/refresh-watch-targets.ts | 156 ++++++ src/skills/runtime/refresh.churn.test.ts | 5 + .../refresh.missing-root.integration.test.ts | 238 ++++++++- .../runtime/refresh.subscriptions.test.ts | 257 ++++++++++ src/skills/runtime/refresh.test.ts | 477 +++++++++--------- src/skills/runtime/refresh.ts | 368 ++++++-------- .../runtime/refresh.watcher.test-support.ts | 27 +- src/skills/runtime/refresh.windows.test.ts | 2 +- 11 files changed, 1266 insertions(+), 458 deletions(-) create mode 100644 src/skills/runtime/refresh-ancestor-watch.ts create mode 100644 src/skills/runtime/refresh-watch-targets.ts create mode 100644 src/skills/runtime/refresh.subscriptions.test.ts diff --git a/src/skills/runtime/refresh-ancestor-watch.ts b/src/skills/runtime/refresh-ancestor-watch.ts new file mode 100644 index 000000000000..de88827f7ef8 --- /dev/null +++ b/src/skills/runtime/refresh-ancestor-watch.ts @@ -0,0 +1,163 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { toErrorObject } from "@openclaw/normalization-core/error-coercion"; +import chokidar, { type FSWatcher } from "chokidar"; +import { isPathInside } from "../../infra/path-guards.js"; +import { teardownSkillsPathWatcher } from "./refresh-watch-close.js"; +import type { createSkillsWatchPathFilter } from "./refresh-watch-path.js"; + +type AncestorSubscription = { + path: string; + ignored: ReturnType["ignored"]; + ready: () => void; + changed: (event: string, path: string) => void; + raw: (event: string, path: unknown, details: unknown) => void; + error: (error: Error) => void; +}; +type AncestorWatcher = { + watcher: FSWatcher; + ready: boolean; + error?: Error; + subscriptions: Set; +}; + +// Imported with the refresh owner at Gateway startup, outside turn contexts. +const runInWatcherContext = AsyncLocalStorage.snapshot(); +const ancestorWatchers = new Map(); + +function createAncestorWatcher( + watchRoot: string, + usePolling: boolean, + subscriptions: Set, +): FSWatcher { + return runInWatcherContext(() => + chokidar.watch(watchRoot, { + ignoreInitial: true, + followSymlinks: false, + usePolling, + // Only observe the next directory in each missing path. Appearance moves + // subscriptions to the closest existing ancestor or their recursive root. + depth: 0, + ignored: (candidate, stats) => { + let ignored = true; + // Each logical filter records directory-symlink identity for unlink + // events. Evaluate all of them even when another target admits entry. + for (const current of subscriptions) { + if (!current.ignored(candidate, stats)) { + ignored = false; + } + } + return ignored; + }, + }), + ); +} + +function observeAncestorWatcher(current: AncestorWatcher): void { + const { watcher, subscriptions } = current; + const isCurrent = () => current.watcher === watcher && !watcher.closed; + watcher.on("ready", () => { + // Chokidar can emit ready after failing to install its native watch. Keep + // that generation failed so a later acquisition retries the physical watch. + if (!isCurrent() || current.error) { + return; + } + current.ready = true; + for (const target of Array.from(subscriptions)) { + if (isCurrent() && subscriptions.has(target)) { + target.ready(); + } + } + }); + watcher.on("all", (event, changedPath) => { + if (!isCurrent()) { + return; + } + for (const target of Array.from(subscriptions)) { + if ( + isCurrent() && + subscriptions.has(target) && + (isPathInside(changedPath, target.path) || isPathInside(target.path, changedPath)) + ) { + target.changed(event, changedPath); + } + } + }); + watcher.on("raw", (event, rawPath, details) => { + if (!isCurrent()) { + return; + } + for (const target of Array.from(subscriptions)) { + if (isCurrent() && subscriptions.has(target)) { + target.raw(event, rawPath, details); + } + } + }); + watcher.on("error", (error) => { + if (!isCurrent()) { + return; + } + current.ready = false; + const watchError = toErrorObject(error, "Skills ancestor watcher failed"); + current.error = watchError; + for (const target of Array.from(subscriptions)) { + if (isCurrent() && subscriptions.has(target)) { + target.error(watchError); + } + } + }); +} + +export function acquireSkillsAncestorWatcher( + watchRoot: string, + usePolling: boolean, + subscription: AncestorSubscription, +): { release: () => void } { + let group = ancestorWatchers.get(watchRoot); + if (!group) { + const subscriptions = new Set([subscription]); + group = { + watcher: createAncestorWatcher(watchRoot, usePolling, subscriptions), + ready: false, + subscriptions, + }; + ancestorWatchers.set(watchRoot, group); + observeAncestorWatcher(group); + } else { + group.subscriptions.add(subscription); + if (group.error) { + const retired = group.watcher; + group.watcher = createAncestorWatcher(watchRoot, usePolling, group.subscriptions); + group.error = undefined; + observeAncestorWatcher(group); + // Releases retain this group and its subscriptions across native retries. + void teardownSkillsPathWatcher({ watcher: retired }); + } + } + const current = group; + const watcher = current.watcher; + // A late subscriber cannot wait for another ready event. Recheck its current + // path before publishing readiness, including creation before registration. + if (current.ready || current.error) { + queueMicrotask(() => { + if ( + current.watcher === watcher && + !watcher.closed && + current.subscriptions.has(subscription) + ) { + if (current.ready) { + subscription.ready(); + } else if (current.error) { + subscription.error(current.error); + } + } + }); + } + return { + release: () => { + if (current.subscriptions.delete(subscription) && current.subscriptions.size === 0) { + ancestorWatchers.delete(watchRoot); + void teardownSkillsPathWatcher(current); + } + }, + }; +} diff --git a/src/skills/runtime/refresh-file-stability.ts b/src/skills/runtime/refresh-file-stability.ts index e0afdfed1bc4..e0173ad4b451 100644 --- a/src/skills/runtime/refresh-file-stability.ts +++ b/src/skills/runtime/refresh-file-stability.ts @@ -20,7 +20,7 @@ function readFileStabilitySnapshot(filePath: string): FileStabilitySnapshot | un async function waitForStableSkillFile( filePath: string, stabilityMs: number, - watcher: FSWatcher, + watcher: Pick, readRevision: () => number, ): Promise { if (watcher.closed || stabilityMs <= 0) { @@ -66,7 +66,7 @@ export function createRawSkillFileScheduler({ schedule, onError, }: { - watcher: FSWatcher; + watcher: Pick; stabilityMs: number; schedule: (filePath: string) => void; onError: (filePath: string, error: unknown) => void; diff --git a/src/skills/runtime/refresh-watch-path.ts b/src/skills/runtime/refresh-watch-path.ts index 680f470cf231..1a7940842947 100644 --- a/src/skills/runtime/refresh-watch-path.ts +++ b/src/skills/runtime/refresh-watch-path.ts @@ -188,6 +188,7 @@ export function resolveSkillsWatcherUsePolling(): boolean { export function makeSkillsWatchTarget( raw: string, depth: number, + previousWatchRoot?: string, ): { path: string; watchRoot: string; depth: number } { const watchPath = toWatchRoot(resolveSkillsWatchPath(raw)); let watchRoot = watchPath; @@ -198,6 +199,32 @@ export function makeSkillsWatchTarget( } watchRoot = parent; } + if ( + previousWatchRoot && + previousWatchRoot !== watchRoot && + isPathInside(previousWatchRoot, watchRoot) + ) { + // Explicit descendant watches bypass Chokidar's followSymlinks:false for + // intermediate components. Promotion must stop where recursive observation + // would stop; trusted realpath targets are registered by target preparation. + let cursor = previousWatchRoot; + const parts = path.relative(previousWatchRoot, watchRoot).split(path.sep).filter(Boolean); + for (const part of ["", ...parts]) { + cursor = path.join(cursor, part); + try { + if (!fs.lstatSync(cursor).isSymbolicLink()) { + continue; + } + } catch { + // An ancestor disappeared between discovery and promotion. + } + watchRoot = path.dirname(cursor); + while (!fs.existsSync(watchRoot) && path.dirname(watchRoot) !== watchRoot) { + watchRoot = path.dirname(watchRoot); + } + break; + } + } return { path: watchPath, watchRoot: toWatchRoot(watchRoot), depth }; } diff --git a/src/skills/runtime/refresh-watch-targets.ts b/src/skills/runtime/refresh-watch-targets.ts new file mode 100644 index 000000000000..d90ec4c73a68 --- /dev/null +++ b/src/skills/runtime/refresh-watch-targets.ts @@ -0,0 +1,156 @@ +import fs from "node:fs"; +import path from "node:path"; +import { resolveRealpathOrAbsolute } from "../../infra/boundary-path.js"; +import { tryRealpath } from "../loading/symlink-targets.js"; +import { + DEFAULT_SKILLS_WATCH_IGNORED, + isTrustedSymlinkSkillTarget, + makeSkillsWatchTarget, + readBudgetedDirEntries, +} from "./refresh-watch-path.js"; + +export type WatchTarget = { + path: string; + watchRoot: string; + depth: number; + executionOnly?: true; +}; + +export const GROUPED_SKILLS_WATCH_DEPTH = 6; +const CONFIGURED_ROOT_WATCH_DEPTH = 2; +const MAX_SYMLINK_WATCH_TARGETS_PER_ROOT = 100; +const MAX_SYMLINK_WATCH_DIRECTORY_SCANS_PER_ROOT = 200; +const MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT = 2_000; + +function addWatchTarget(targets: Map, raw: string, depth: number): void { + const target = makeSkillsWatchTarget(raw, depth); + target.depth = Math.max(target.depth, targets.get(target.path)?.depth ?? 0); + targets.set(target.path, target); +} + +function addSkillRootWatchTargets( + targets: Map, + root: string, + rootDepth: number, +): string { + addWatchTarget(targets, root, rootDepth); + const companionSkillsRoot = path.join(root, "skills"); + addWatchTarget(targets, companionSkillsRoot, GROUPED_SKILLS_WATCH_DEPTH); + return companionSkillsRoot; +} + +export function addSkillSourceWatchTargets( + targets: Map, + root: string, + source: string, + allowedSymlinkTargetRealPaths: readonly string[], + rootDepth = path.basename(root) === "skills" + ? GROUPED_SKILLS_WATCH_DEPTH + : CONFIGURED_ROOT_WATCH_DEPTH, +): void { + const companionSkillsRoot = addSkillRootWatchTargets(targets, root, rootDepth); + // Both bounded scans share the source's containment identity for this preparation. + // Trusted symlink leaves below remain registration-only, never recursive scans. + const rootRealPath = resolveRealpathOrAbsolute(root); + addTrustedSymlinkSkillWatchTargets( + targets, + root, + source, + allowedSymlinkTargetRealPaths, + rootDepth, + rootRealPath, + rootRealPath, + ); + addTrustedSymlinkSkillWatchTargets( + targets, + companionSkillsRoot, + source, + allowedSymlinkTargetRealPaths, + GROUPED_SKILLS_WATCH_DEPTH, + rootRealPath, + resolveRealpathOrAbsolute(companionSkillsRoot), + ); +} + +function addTrustedSymlinkSkillWatchTargets( + targets: Map, + root: string, + source: string, + allowedSymlinkTargetRealPaths: readonly string[], + maxDepth: number, + containmentRootRealPath: string, + rootRealPath: string, +): void { + try { + if ( + fs.lstatSync(root).isSymbolicLink() && + isTrustedSymlinkSkillTarget( + source, + containmentRootRealPath, + rootRealPath, + allowedSymlinkTargetRealPaths, + ) + ) { + addSkillRootWatchTargets(targets, rootRealPath, maxDepth); + } + } catch { + return; + } + const queue: Array<{ dir: string; depth: number }> = [{ dir: root, depth: 0 }]; + let watched = 0; + let directoryScans = 0; + let rawEntries = 0; + for (const queued of queue) { + if ( + watched >= MAX_SYMLINK_WATCH_TARGETS_PER_ROOT || + directoryScans >= MAX_SYMLINK_WATCH_DIRECTORY_SCANS_PER_ROOT || + rawEntries >= MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT + ) { + break; + } + const current = queued; + if (!current) { + continue; + } + const scan = readBudgetedDirEntries( + current.dir, + MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT - rawEntries, + ); + directoryScans += 1; + rawEntries += scan.scannedEntryCount; + if (!scan.ok) { + continue; + } + for (const entry of scan.entries.toSorted((a, b) => a.name.localeCompare(b.name))) { + if (watched >= MAX_SYMLINK_WATCH_TARGETS_PER_ROOT) { + break; + } + if (entry.name.startsWith(".") || entry.name === "node_modules") { + continue; + } + const childPath = path.join(current.dir, entry.name); + if (DEFAULT_SKILLS_WATCH_IGNORED.some((re) => re.test(childPath))) { + continue; + } + if (entry.isSymbolicLink()) { + const targetRealPath = tryRealpath(childPath); + if ( + targetRealPath && + isTrustedSymlinkSkillTarget( + source, + containmentRootRealPath, + targetRealPath, + allowedSymlinkTargetRealPaths, + ) + ) { + addSkillRootWatchTargets(targets, targetRealPath, GROUPED_SKILLS_WATCH_DEPTH); + watched += 1; + } + continue; + } + if (entry.isDirectory() && current.depth < maxDepth) { + queue.push({ dir: childPath, depth: current.depth + 1 }); + } + } + } +} diff --git a/src/skills/runtime/refresh.churn.test.ts b/src/skills/runtime/refresh.churn.test.ts index b4d5d5a2b49b..65c006f14b56 100644 --- a/src/skills/runtime/refresh.churn.test.ts +++ b/src/skills/runtime/refresh.churn.test.ts @@ -51,6 +51,11 @@ describe("skills watcher churn", () => { vi.useFakeTimers(); const workspaceDir = fixtureWorkspaceDir; const executionWorkspaceDir = await createFixtureDirectory("busy-worktree"); + for (const root of [workspaceDir, executionWorkspaceDir]) { + for (const directory of ["skills", ".agents/skills"]) { + await fs.mkdir(path.join(root, directory), { recursive: true }); + } + } refreshModule.ensureSkillsWatcher({ workspaceDir, executionWorkspaceDir }); const version = getSkillsSnapshotVersion(workspaceDir); const seen: SkillsChangeEvent[] = []; diff --git a/src/skills/runtime/refresh.missing-root.integration.test.ts b/src/skills/runtime/refresh.missing-root.integration.test.ts index 185a1a4c0b18..1c922f20f102 100644 --- a/src/skills/runtime/refresh.missing-root.integration.test.ts +++ b/src/skills/runtime/refresh.missing-root.integration.test.ts @@ -4,7 +4,10 @@ import fs from "node:fs/promises"; import { syncBuiltinESMExports } from "node:module"; import os from "node:os"; import path from "node:path"; -import { expect, it, vi } from "vitest"; +import chokidar from "chokidar"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; +import { createDeferredCore } from "../../shared/deferred.js"; vi.mock("../loading/plugin-skills.js", () => ({ resolvePluginSkillRoots: () => [], @@ -197,3 +200,236 @@ it("refreshes skills created beneath an initially missing project skills root", await fs.rm(root, { recursive: true, force: true }); } }); + +describe("shared missing skill ancestors", () => { + const roots = useAutoCleanupTempDirTracker((cleanup) => + afterEach(async () => { + const { closeSkillsWatchers } = await import("./refresh.js"); + await closeSkillsWatchers(true); + vi.restoreAllMocks(); + cleanup(); + }), + ); + + it.each(["higher", "intermediate"] as const)( + "preserves settled root discovery after moving its %s ancestor and retains sibling subscriptions", + async (ancestor) => { + const root = await fs.realpath(roots.make("skills-shared-ancestor-")); + const source = (name: string) => { + const sourceRoot = path.join(root, name, "nested", "skills"); + return { + workspaceDir: path.join(root, `workspace-${name}`), + sourceRoot, + config: { skills: { load: { extraDirs: [sourceRoot] } } }, + }; + }; + const first = source("left"); + const second = source("right"); + for (const current of [first, second]) { + await fs.mkdir(path.join(current.workspaceDir, "skills"), { recursive: true }); + } + const { ensureSkillsWatcher, registerSkillsChangeListener } = await import("./refresh.js"); + const { loadWorkspaceSkills } = await import("../loading/workspace-skill-loader.js"); + const originalWatch = chokidar.watch; + const observed: Array<{ watcher: ReturnType; ready: boolean }> = []; + const watcherErrors: unknown[] = []; + const watch = vi.spyOn(chokidar, "watch").mockImplementation((...args) => { + const watcher = originalWatch(...args); + const observation = { watcher, ready: false }; + observed.push(observation); + // Attach before returning: promotion can create more watchers during ready. + watcher.once("ready", () => { + observation.ready = true; + }); + watcher.on("error", (error) => watcherErrors.push(error)); + return watcher; + }); + const originalSetTimeout = globalThis.setTimeout; + const originalClearTimeout = globalThis.clearTimeout; + const pendingTimers = new Map< + Parameters[0], + { settled: Promise; finish: () => void } + >(); + vi.spyOn(globalThis, "setTimeout").mockImplementation((callback, delay, ...args) => { + const { promise: settled, resolve: finish } = createDeferredCore(); + const timer = originalSetTimeout(() => { + pendingTimers.delete(timer); + try { + callback.apply(timer, args); + } finally { + finish(); + } + }, delay); + pendingTimers.set(timer, { settled, finish }); + return timer; + }); + vi.spyOn(globalThis, "clearTimeout").mockImplementation((timer) => { + originalClearTimeout(timer); + pendingTimers.get(timer)?.finish(); + pendingTimers.delete(timer); + }); + const settleWatchers = async () => { + for (;;) { + await vi.waitFor(() => { + expect(watcherErrors).toEqual([]); + expect(observed.every(({ watcher, ready }) => ready || watcher.closed)).toBe(true); + }); + const generationCount = observed.length; + // Drain actual debounce/stability work, including timers chained by its + // continuations. Keep native time and watchers; do not sleep past a guess. + await Promise.all(Array.from(pendingTimers.values(), ({ settled }) => settled)); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(watcherErrors).toEqual([]); + if ( + pendingTimers.size === 0 && + observed.length === generationCount && + observed.every(({ watcher, ready }) => ready || watcher.closed) + ) { + return; + } + } + }; + for (const current of [first, second]) { + ensureSkillsWatcher(current); + } + await settleWatchers(); + expect( + watch.mock.calls.filter(([watched]) => watched === root.replaceAll("\\", "/")), + ).toHaveLength(1); + const changes: string[] = []; + const unregister = registerSkillsChangeListener((event) => { + if (event.workspaceDir) { + changes.push(event.workspaceDir); + } + }); + const read = (current: typeof first) => + loadWorkspaceSkills(current.workspaceDir, { + config: current.config, + bundledSkillsDir: "", + managedSkillsDir: path.join(root, "unused"), + }).map((entry) => entry.skill.name); + const writeSkill = async (current: typeof first, name: string) => { + const directory = path.join(current.sourceRoot, name); + await fs.mkdir(directory, { recursive: true }); + await fs.writeFile( + path.join(directory, "SKILL.md"), + `---\nname: ${name}\ndescription: Shared ancestor proof\n---\n`, + ); + }; + try { + expect(read(first)).toEqual([]); + expect(read(second)).toEqual([]); + await fs.writeFile(path.join(root, "unrelated.sqlite-wal"), "unrelated"); + await writeSkill(first, "first-proof"); + await expect.poll(() => read(first), { timeout: 3_000 }).toContain("first-proof"); + expect(changes).not.toContain(second.workspaceDir); + await vi.waitFor(() => { + expect( + watch.mock.calls.some( + ([watched], index) => + watched === first.sourceRoot.replaceAll("\\", "/") && + observed[index]?.ready && + !observed[index]?.watcher.closed, + ), + ).toBe(true); + }); + await settleWatchers(); + // Prime after promoted root/companion scans and queued refreshes settle: + // late initial reconciliation must not mask a missed ancestor move. + expect(read(first)).toContain("first-proof"); + const movedAncestor = + ancestor === "higher" ? path.join(root, "left") : path.join(root, "left", "nested"); + if (process.platform === "win32") { + // Windows cannot rename an ancestor with live descendant directory watches. + await fs.rm(movedAncestor, { recursive: true }); + } else { + await fs.rename(movedAncestor, `${movedAncestor}-away`); + } + await expect.poll(() => read(first), { timeout: 3_000 }).toEqual([]); + await writeSkill(first, "returned-proof"); + await expect.poll(() => read(first), { timeout: 3_000 }).toEqual(["returned-proof"]); + await settleWatchers(); + expect(read(first)).toEqual(["returned-proof"]); + // Retiring one logical workspace must not retire the shared missing-root observer. + ensureSkillsWatcher({ + workspaceDir: first.workspaceDir, + config: { skills: { load: { watch: false } } }, + }); + await writeSkill(second, "remaining-proof"); + await expect.poll(() => read(second), { timeout: 3_000 }).toContain("remaining-proof"); + const skillFile = path.join(second.sourceRoot, "remaining-proof", "SKILL.md"); + const renamedSkillFile = path.join(second.sourceRoot, "remaining-proof", "SKILL.saved"); + await fs.rename(skillFile, renamedSkillFile); + await expect.poll(() => read(second), { timeout: 3_000 }).toEqual([]); + await fs.rename(renamedSkillFile, skillFile); + await expect.poll(() => read(second), { timeout: 3_000 }).toContain("remaining-proof"); + await fs.rm(path.join(root, "right"), { recursive: true }); + await expect.poll(() => read(second), { timeout: 3_000 }).toEqual([]); + await writeSkill(second, "recreated-proof"); + await expect.poll(() => read(second), { timeout: 3_000 }).toContain("recreated-proof"); + } finally { + unregister(); + } + }, + ); + + it.runIf(process.platform !== "win32")( + "does not promote missing roots through newly created ancestor symlinks", + async () => { + const root = await fs.realpath(roots.make("skills-ancestor-symlink-")); + const outside = await fs.realpath(roots.make("skills-ancestor-outside-")); + const workspaceDir = path.join(root, "workspace"); + await fs.mkdir(path.join(workspaceDir, "skills"), { recursive: true }); + const link = path.join(root, "missing"); + const sourceRoot = path.join(link, "nested", "skills"); + await fs.mkdir(path.join(outside, "nested", "skills", "outside-proof"), { recursive: true }); + const config = { skills: { load: { extraDirs: [sourceRoot] } } }; + const watch = vi.spyOn(chokidar, "watch"); + const { ensureSkillsWatcher } = await import("./refresh.js"); + const { getSkillsSourceVersion } = await import("./refresh-state.js"); + const { loadWorkspaceSkills } = await import("../loading/workspace-skill-loader.js"); + const read = () => + loadWorkspaceSkills(workspaceDir, { + config, + bundledSkillsDir: "", + managedSkillsDir: path.join(root, "unused"), + }).map((entry) => entry.skill.name); + ensureSkillsWatcher({ workspaceDir, config }); + expect(read()).toEqual([]); + await Promise.all( + watch.mock.results.map((result) => { + if (result.type !== "return") { + throw new Error("Watcher acquisition failed"); + } + return new Promise((resolve, reject) => { + result.value.once("ready", resolve); + result.value.once("error", reject); + }); + }), + ); + // Ready handlers reconcile synchronously before these promises resolve. + // Unchanged empty inventory suppresses public events, but discovery still invalidates. + const sourceVersion = getSkillsSourceVersion(workspaceDir); + await fs.symlink(outside, link, "dir"); + await expect + .poll(() => getSkillsSourceVersion(workspaceDir), { timeout: 3_000 }) + .toBeGreaterThan(sourceVersion); + expect( + watch.mock.calls.some( + ([watched]) => + typeof watched === "string" && (watched === link || watched.startsWith(`${link}/`)), + ), + ).toBe(false); + await fs.unlink(link); + const skillDir = path.join(sourceRoot, "ordinary-proof"); + await fs.mkdir(skillDir, { recursive: true }); + await fs.writeFile( + path.join(skillDir, "SKILL.md"), + "---\nname: ordinary-proof\ndescription: Ordinary replacement\n---\n", + ); + await expect.poll(read, { timeout: 3_000 }).toContain("ordinary-proof"); + }, + ); +}); diff --git a/src/skills/runtime/refresh.subscriptions.test.ts b/src/skills/runtime/refresh.subscriptions.test.ts new file mode 100644 index 000000000000..133341ff9c3b --- /dev/null +++ b/src/skills/runtime/refresh.subscriptions.test.ts @@ -0,0 +1,257 @@ +import fs from "node:fs/promises"; +import path from "node:path"; +import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { withEnvAsync } from "../../test-utils/env.js"; +import { + bumpSkillsSnapshotVersion, + getSkillsSnapshotVersion, + getSkillsSourceVersion, +} from "./refresh-state.js"; +import { + createSkillsWatcherMock, + useSkillsWatcherFixture, +} from "./refresh.watcher.test-support.js"; + +type SkillsChangeEvent = NonNullable[0]>; +const { createdWatchers, watchMock, watchForSkillRoot } = createSkillsWatcherMock(); +let refreshModule: typeof import("./refresh.js"); +let fixtureWorkspaceDir: string; + +vi.mock("chokidar", () => ({ default: { watch: watchMock } })); +vi.mock("../loading/plugin-skills.js", () => ({ + resolvePluginSkillRoots: () => [], + resolvePluginSkillRootsFromMetadata: () => [], +})); + +describe("skills watcher subscription lifecycle", () => { + const fixture = useSkillsWatcherFixture(); + const { createFixtureDirectory } = fixture; + beforeAll(async () => { + refreshModule = await import("./refresh.js"); + }); + beforeEach(() => { + vi.stubEnv("CHOKIDAR_USEPOLLING", "false"); + watchMock.mockClear(); + createdWatchers.length = 0; + fixtureWorkspaceDir = fixture.workspaceDir; + }); + + it("stops fanning a shared-directory change to a workspace after it unsubscribes", async () => { + vi.useFakeTimers(); + const secondWorkspace = await createFixtureDirectory("second-workspace"); + const sharedRoot = await createFixtureDirectory("shared"); + const config = { skills: { load: { extraDirs: [sharedRoot] } } }; + const seen: SkillsChangeEvent[] = []; + refreshModule.registerSkillsChangeListener((change) => { + seen.push(change); + }); + refreshModule.ensureSkillsWatcher({ workspaceDir: fixtureWorkspaceDir, config }); + refreshModule.ensureSkillsWatcher({ workspaceDir: secondWorkspace, config }); + const sharedWatcher = watchForSkillRoot(sharedRoot).watcher; + + refreshModule.ensureSkillsWatcher({ + workspaceDir: fixtureWorkspaceDir, + config: { skills: { load: { extraDirs: [sharedRoot], watch: false } } }, + }); + seen.length = 0; + expect(sharedWatcher.close).not.toHaveBeenCalled(); + const changedPath = path.join(sharedRoot, "demo", "SKILL.md"); + sharedWatcher.emit("all", "change", changedPath); + await vi.advanceTimersByTimeAsync(250); + + expect(seen).toEqual([{ workspaceDir: secondWorkspace, reason: "watch", changedPath }]); + }); + + it("preserves workspace invalidation on watch disable without changing other workspaces", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + const workspaceDir = fixtureWorkspaceDir; + const otherWorkspace = await createFixtureDirectory("other-workspace"); + refreshModule.ensureSkillsWatcher({ workspaceDir: otherWorkspace }); + const otherVersion = getSkillsSnapshotVersion(otherWorkspace); + const globalVersion = getSkillsSnapshotVersion(); + refreshModule.ensureSkillsWatcher({ + workspaceDir, + config: { skills: { load: {} } }, + }); + + const firstVersion = bumpSkillsSnapshotVersion({ + workspaceDir, + reason: "watch", + changedPath: `${workspaceDir}/skills/demo/SKILL.md`, + }); + refreshModule.ensureSkillsWatcher({ + workspaceDir, + config: { skills: { load: { watch: false } } }, + }); + + const nextVersion = getSkillsSnapshotVersion(workspaceDir); + expect(nextVersion).toBe(firstVersion); + expect(getSkillsSnapshotVersion(otherWorkspace)).toBe(otherVersion); + expect(getSkillsSnapshotVersion()).toBe(globalVersion); + vi.setSystemTime(new Date(nextVersion)); + refreshModule.ensureSkillsWatcher({ workspaceDir }); + expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(nextVersion); + }); + + it("evicts idle workspace subscriptions on a later ensure call", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + const idleWorkspaceDir = fixtureWorkspaceDir; + const activeWorkspaceDir = await createFixtureDirectory("workspace-active"); + refreshModule.ensureSkillsWatcher({ + workspaceDir: idleWorkspaceDir, + config: { skills: { load: {} } }, + }); + const idleSkillsWatcher = watchForSkillRoot(path.join(idleWorkspaceDir, "skills")).watcher; + const firstVersion = bumpSkillsSnapshotVersion({ + workspaceDir: idleWorkspaceDir, + reason: "watch", + }); + const globalVersion = getSkillsSnapshotVersion(); + + vi.advanceTimersByTime(60 * 60_000 + 1_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir: activeWorkspaceDir, + config: { skills: { load: {} } }, + }); + + expect(idleSkillsWatcher.close).toHaveBeenCalledTimes(1); + const evictedVersion = getSkillsSnapshotVersion(idleWorkspaceDir); + expect(evictedVersion).toBe(firstVersion); + expect(getSkillsSnapshotVersion()).toBe(globalVersion); + vi.setSystemTime(new Date(evictedVersion)); + refreshModule.ensureSkillsWatcher({ workspaceDir: idleWorkspaceDir }); + expect(getSkillsSnapshotVersion(idleWorkspaceDir)).toBeGreaterThan(evictedVersion); + }); + + it("keeps another execution subscription for the workspace alive after disposal", async () => { + vi.useFakeTimers(); + const workspaceDir = fixtureWorkspaceDir; + const executionWorkspaceDir = await createFixtureDirectory("remaining-worktree"); + refreshModule.ensureSkillsWatcher({ workspaceDir }); + refreshModule.ensureSkillsWatcher({ workspaceDir, executionWorkspaceDir }); + const version = getSkillsSnapshotVersion(workspaceDir); + const globalVersion = getSkillsSnapshotVersion(); + const watcher = watchForSkillRoot(path.join(workspaceDir, "skills")).watcher; + const seen: SkillsChangeEvent[] = []; + refreshModule.registerSkillsChangeListener((change) => seen.push(change)); + + refreshModule.ensureSkillsWatcher({ + workspaceDir, + config: { skills: { load: { watch: false } } }, + }); + + expect(watcher.close).not.toHaveBeenCalled(); + expect(getSkillsSnapshotVersion(workspaceDir)).toBe(version); + expect(getSkillsSnapshotVersion()).toBe(globalVersion); + expect(seen).toEqual([]); + const changedPath = path.join(workspaceDir, "skills", "demo", "SKILL.md"); + watcher.emit("all", "change", changedPath); + await vi.advanceTimersByTimeAsync(250); + expect(seen).toEqual([{ workspaceDir, reason: "watch", changedPath }]); + }); + + it("keeps refreshed workspace subscriptions within the idle TTL", async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + const activeWorkspaceDir = fixtureWorkspaceDir; + const otherWorkspaceDir = await createFixtureDirectory("workspace-other"); + refreshModule.ensureSkillsWatcher({ + workspaceDir: activeWorkspaceDir, + config: { skills: { load: {} } }, + }); + const activeSkillsWatcher = watchForSkillRoot(path.join(activeWorkspaceDir, "skills")).watcher; + + vi.advanceTimersByTime(30 * 60_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir: activeWorkspaceDir, + config: { skills: { load: {} } }, + }); + vi.advanceTimersByTime(31 * 60_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir: otherWorkspaceDir, + config: { skills: { load: {} } }, + }); + + expect(activeSkillsWatcher.close).not.toHaveBeenCalled(); + }); + + it.each(["execution", "base"] as const)( + "keeps an idle %s source active while another consumer remains", + async (scope) => { + vi.useFakeTimers(); + const { loadWorkspaceSkills } = await import("../loading/workspace-skill-loader.js"); + const workspaceDir = fixtureWorkspaceDir; + const executionWorkspaceDir = await createFixtureDirectory("shared-execution"); + const idleScope = { + executionWorkspaceDir: scope === "execution" ? executionWorkspaceDir : undefined, + }; + const skillDir = await createFixtureDirectory( + scope === "execution" ? "shared-execution/skills/demo" : "workspace/skills/demo", + ); + const skillFile = path.join(skillDir, "SKILL.md"); + await fs.writeFile( + skillFile, + "---\nname: demo\ndescription: Demo\n---\nOriginal instructions\n", + ); + await withEnvAsync({ OPENCLAW_STATE_DIR: workspaceDir }, async () => { + const options = { + ...idleScope, + agentId: "agent-b", + bundledSkillsDir: "", + managedSkillsDir: path.join(workspaceDir, "missing-managed"), + }; + const original = loadWorkspaceSkills(workspaceDir, options)[0]!.skill.contentHash; + refreshModule.ensureSkillsWatcher({ + workspaceDir, + ...idleScope, + agentId: "agent-a", + }); + refreshModule.ensureSkillsWatcher({ + workspaceDir, + executionWorkspaceDir, + agentId: "agent-b", + }); + vi.advanceTimersByTime(30 * 60_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir, + executionWorkspaceDir, + agentId: "agent-b", + }); + const sourceVersion = getSkillsSourceVersion(workspaceDir, idleScope); + vi.advanceTimersByTime(31 * 60_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir, + executionWorkspaceDir, + agentId: "agent-b", + }); + expect(getSkillsSourceVersion(workspaceDir, idleScope)).toBe(sourceVersion); + + const version = getSkillsSnapshotVersion(workspaceDir); + await fs.appendFile(skillFile, "\nUpdated instructions\n"); + bumpSkillsSnapshotVersion({ reason: "workshop" }); + expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(version); + expect(loadWorkspaceSkills(workspaceDir, options)[0]!.skill.contentHash).not.toBe(original); + + vi.advanceTimersByTime(60 * 60_000 + 1_000); + refreshModule.ensureSkillsWatcher({ + workspaceDir: await createFixtureDirectory("other-active-workspace"), + }); + const retiredVersion = getSkillsSnapshotVersion(workspaceDir); + await fs.appendFile(skillFile, "\nInstructions changed while retired\n"); + bumpSkillsSnapshotVersion({ reason: "workshop" }); + expect(getSkillsSnapshotVersion(workspaceDir)).toBe(retiredVersion); + + refreshModule.ensureSkillsWatcher({ + workspaceDir, + executionWorkspaceDir: + scope === "base" + ? await createFixtureDirectory("new-execution-workspace") + : executionWorkspaceDir, + }); + expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(retiredVersion); + }); + }, + ); +}); diff --git a/src/skills/runtime/refresh.test.ts b/src/skills/runtime/refresh.test.ts index 865f0713c9fa..41f3bf4250e7 100644 --- a/src/skills/runtime/refresh.test.ts +++ b/src/skills/runtime/refresh.test.ts @@ -68,7 +68,10 @@ describe("ensureSkillsWatcher", () => { expect(options.depth).toBe(7); const projectWatch = watchForSkillRoot(projectSkillsRoot); expect(projectWatch.watchRoot).toBe(workspaceDir.replaceAll("\\", "/")); - expect(projectWatch.options.depth).toBe(9); + expect(projectWatch.options.depth).toBe(0); + await fs.mkdir(projectSkillsRoot, { recursive: true }); + projectWatch.watcher.emit("all", "addDir", path.dirname(projectSkillsRoot)); + expect(watchForSkillRoot(projectSkillsRoot).options.depth).toBe(7); expect( watchForSkillRoot(path.join(os.homedir(), ".agents", "skills")).options.followSymlinks, ).toBe(false); @@ -164,11 +167,16 @@ describe("ensureSkillsWatcher", () => { const watched = watchForSkillRoot(path.join(workspaceDir, "skills")); expect(watched.watchRoot).toBe(workspaceDir.replaceAll("\\", "/")); - expect(watched.options.depth).toBe(8); + expect(watched.options.depth).toBe(0); const changedPath = path.join(workspaceDir, "skills", "group", "demo", "SKILL.md"); + await fs.mkdir(path.dirname(changedPath), { recursive: true }); + watched.watcher.emit("all", "addDir", path.join(workspaceDir, "skills")); + const promoted = watchForSkillRoot(path.join(workspaceDir, "skills")); + expect(promoted.options.depth).toBe(7); + await vi.advanceTimersByTimeAsync(250); seen.length = 0; - watched.watcher.emit("all", "change", changedPath); + promoted.watcher.emit("all", "change", changedPath); await vi.advanceTimersByTimeAsync(250); expect(seen).toEqual([ @@ -371,11 +379,11 @@ describe("ensureSkillsWatcher", () => { expect(rootWatch.watchRoot).toBe(repoDir.replaceAll("\\", "/")); expect(companionWatch.watchRoot).toBe(rootWatch.watchRoot); expect(rootWatch.options.depth).toBe(3); - expect(companionWatch.options.depth).toBe(8); + expect(companionWatch.options.depth).toBe(0); expect(companionWatch.options.ignored(path.join(repoDir, "other"))).toBe(true); }); - it("bumps missing configured root depth for first nested skill creation", async () => { + it("promotes missing configured roots for first nested skill creation", async () => { const parentDir = await createFixtureDirectory("missing-skill-root"); const missingRoot = path.join(parentDir, "repo"); refreshModule.ensureSkillsWatcher({ @@ -387,8 +395,12 @@ describe("ensureSkillsWatcher", () => { const companionWatch = watchForSkillRoot(path.join(missingRoot, "skills")); expect(rootWatch.watchRoot).toBe(parentDir.replaceAll("\\", "/")); expect(companionWatch.watchRoot).toBe(rootWatch.watchRoot); - expect(rootWatch.options.depth).toBe(4); - expect(companionWatch.options.depth).toBe(9); + expect(rootWatch.options.depth).toBe(0); + expect(companionWatch.watcher).toBe(rootWatch.watcher); + await fs.mkdir(path.join(missingRoot, "skills", "group", "demo"), { recursive: true }); + rootWatch.watcher.emit("all", "addDir", missingRoot); + expect(watchForSkillRoot(missingRoot).options.depth).toBe(3); + expect(watchForSkillRoot(path.join(missingRoot, "skills")).options.depth).toBe(7); }); it("watches configured roots named skills at grouped depth", async () => { @@ -443,6 +455,170 @@ describe("ensureSkillsWatcher", () => { expect(createdWatchers[firstIndex]?.close).not.toHaveBeenCalled(); }); + it.each(["error", "error-then-ready", "ready-then-error", "null-error"] as const)( + "retries a failed shared ancestor for a new root after %s", + async (scan) => { + vi.useFakeTimers(); + const ancestor = await createFixtureDirectory("failed-ancestor"); + const firstRoot = path.join(ancestor, "first"); + const secondRoot = path.join(ancestor, "second"); + const secondWorkspace = await createFixtureDirectory("second-workspace"); + refreshModule.ensureSkillsWatcher({ + workspaceDir: fixtureWorkspaceDir, + config: { skills: { load: { extraDirs: [firstRoot] } } }, + }); + const failed = watchForSkillRoot(firstRoot).watcher; + const lateReady = failed.on.mock.calls.find(([event]) => event === "ready")![1]; + const lateError = failed.on.mock.calls.find(([event]) => event === "error")![1]; + if (scan === "ready-then-error") { + for (const watcher of createdWatchers) { + watcher.emit("ready"); + } + } + failed.emit( + "error", + scan === "null-error" + ? null + : Object.assign(new Error("native watch failed"), { code: "EIO" }), + ); + if (scan === "error-then-ready") { + // Chokidar can finish scanning after native watch installation failed. + failed.emit("ready"); + } + refreshModule.ensureSkillsWatcher({ + workspaceDir: secondWorkspace, + config: { skills: { load: { extraDirs: [secondRoot] } } }, + }); + const replacement = watchForSkillRoot(secondRoot).watcher; + expect(replacement).not.toBe(failed); + expect(watchForSkillRoot(firstRoot).watcher).toBe(replacement); + expect(failed.close).toHaveBeenCalledOnce(); + for (const watcher of createdWatchers) { + if (watcher !== replacement) { + watcher.emit("ready"); + } + } + const firstBeforeReady = getSkillsSourceVersion(fixtureWorkspaceDir); + const secondBeforeReady = getSkillsSourceVersion(secondWorkspace); + lateReady(); + lateError(new Error("retired scan failed")); + expect(getSkillsSourceVersion(fixtureWorkspaceDir)).toBe(firstBeforeReady); + expect(getSkillsSourceVersion(secondWorkspace)).toBe(secondBeforeReady); + replacement.emit("ready"); + if (scan === "ready-then-error") { + expect(getSkillsSourceVersion(fixtureWorkspaceDir)).toBe(firstBeforeReady); + } else { + expect(getSkillsSourceVersion(fixtureWorkspaceDir)).toBeGreaterThan(firstBeforeReady); + } + expect(getSkillsSourceVersion(secondWorkspace)).toBeGreaterThan(secondBeforeReady); + + const firstBeforeCreation = getSkillsSourceVersion(fixtureWorkspaceDir); + await fs.mkdir(firstRoot); + replacement.emit("all", "addDir", firstRoot); + await vi.advanceTimersByTimeAsync(250); + expect(getSkillsSourceVersion(fixtureWorkspaceDir)).toBeGreaterThan(firstBeforeCreation); + expect(replacement.close).not.toHaveBeenCalled(); + const secondBeforeCreation = getSkillsSourceVersion(secondWorkspace); + await fs.mkdir(secondRoot); + replacement.emit("all", "addDir", secondRoot); + await vi.advanceTimersByTimeAsync(250); + expect(getSkillsSourceVersion(secondWorkspace)).toBeGreaterThan(secondBeforeCreation); + expect(replacement.close).not.toHaveBeenCalled(); + const first = watchForSkillRoot(firstRoot).watcher; + const second = watchForSkillRoot(secondRoot).watcher; + refreshModule.ensureSkillsWatcher({ + workspaceDir: fixtureWorkspaceDir, + config: { skills: { load: { watch: false } } }, + }); + expect(first.close).toHaveBeenCalledOnce(); + expect(second.close).not.toHaveBeenCalled(); + expect(replacement.close).not.toHaveBeenCalled(); + refreshModule.ensureSkillsWatcher({ + workspaceDir: secondWorkspace, + config: { skills: { load: { watch: false } } }, + }); + expect(second.close).toHaveBeenCalledOnce(); + expect(replacement.close).toHaveBeenCalledOnce(); + }, + ); + + it.each(["before-content", "around-content"] as const)( + "discovers initial content when ancestor scans fail %s", + async (ordering) => { + const { loadWorkspaceSkills } = await import("../loading/workspace-skill-loader.js"); + const ancestor = await createFixtureDirectory("ancestor-error"); + const intermediate = path.join(ancestor, "nested"); + const innerAncestor = path.join(intermediate, "inner"); + const logicalRoot = path.join(innerAncestor, "skills"); + const config = { skills: { load: { extraDirs: [logicalRoot] } } }; + const read = () => + loadWorkspaceSkills(fixtureWorkspaceDir, { + config, + bundledSkillsDir: "", + managedSkillsDir: path.join(ancestor, "unused"), + }).map((entry) => entry.skill.name); + refreshModule.ensureSkillsWatcher({ workspaceDir: fixtureWorkspaceDir, config }); + const initialAncestor = watchForSkillRoot(logicalRoot).watcher; + await fs.mkdir(logicalRoot, { recursive: true }); + // Promotion through initial readiness creates no addDir debounce that could + // later invalidate the empty cache independently of content readiness. + initialAncestor.emit("ready"); + const content = watchForSkillRoot(logicalRoot).watcher; + expect(content).not.toBe(initialAncestor); + const failedIndex = watchMock.mock.calls.findLastIndex( + ([watchRoot, options], index) => + watchRoot === intermediate.replaceAll("\\", "/") && + options.depth === 0 && + !createdWatchers[index]?.closed, + ); + expect(failedIndex).toBeGreaterThanOrEqual(0); + const failedAncestor = createdWatchers[failedIndex]!; + const lastFailedIndex = + ordering === "around-content" + ? watchMock.mock.calls.findLastIndex( + ([watchRoot, options], index) => + watchRoot === innerAncestor.replaceAll("\\", "/") && + options.depth === 0 && + !createdWatchers[index]?.closed, + ) + : -1; + if (ordering === "around-content") { + expect(lastFailedIndex).toBeGreaterThanOrEqual(0); + } + const lastFailedAncestor = createdWatchers[lastFailedIndex]; + for (const watcher of createdWatchers) { + if (watcher !== content && watcher !== failedAncestor && watcher !== lastFailedAncestor) { + watcher.emit("ready"); + } + } + // Existing shared ancestors notify late subscriptions in a microtask. Finish + // those callbacks before error delivery; they must not repair the cache later. + await Promise.resolve(); + failedAncestor.emit( + "error", + Object.assign(new Error("ancestor scan failed"), { code: "EIO" }), + ); + expect(read()).toEqual([]); + const skillDir = path.join(logicalRoot, "ancestor-error-proof"); + await fs.mkdir(skillDir); + await fs.writeFile( + path.join(skillDir, "SKILL.md"), + "---\nname: ancestor-error-proof\ndescription: Discovered by healthy content scan\n---\n", + ); + // ignoreInitial may suppress all/change events for content found by this scan. + expect(read()).toEqual([]); + content.emit("ready"); + if (lastFailedAncestor) { + expect(read()).toEqual([]); + lastFailedAncestor.emit( + "error", + Object.assign(new Error("last ancestor scan failed"), { code: "EIO" }), + ); + } + expect(read()).toEqual(["ancestor-error-proof"]); + }, + ); + it.each(["ready", "error-then-ready"] as const)( "preserves shared coverage when a replaced ancestor scan is %s", async (scan) => { @@ -459,7 +635,7 @@ describe("ensureSkillsWatcher", () => { }); const shallow = watchForSkillRoot(logicalRoot); expect(shallow.watchRoot).toBe(ancestor.replaceAll("\\", "/")); - expect(shallow.options.depth).toBe(8); + expect(shallow.options.depth).toBe(0); for (const watcher of createdWatchers) { watcher.emit("ready"); } @@ -473,15 +649,14 @@ describe("ensureSkillsWatcher", () => { ); refreshModule.ensureSkillsWatcher({ workspaceDir: secondWorkspace }); const deeper = watchForSkillRoot(logicalRoot); - expect(shallow.watcher.close).toHaveBeenCalledOnce(); expect(deeper.watchRoot).toBe(logicalRoot.replaceAll("\\", "/")); - // Eight physical levels cover two skill levels, metadata, and five missing ancestors. expect(deeper.options.depth).toBe(7); for (const watcher of createdWatchers) { if (watcher !== deeper.watcher) { watcher.emit("ready"); } } + expect(shallow.watcher.close).not.toHaveBeenCalled(); if (scan === "error-then-ready") { deeper.watcher.emit("error", new Error("initial scan interrupted")); expect(getSkillsSourceVersion(fixtureWorkspaceDir)).toBeGreaterThan(beforeReplacement); @@ -494,8 +669,15 @@ describe("ensureSkillsWatcher", () => { refreshModule.registerSkillsChangeListener((change) => seen.push(change)); const versionBefore = getSkillsSnapshotVersion(fixtureWorkspaceDir); const removedParent = path.dirname(logicalRoot); + const parentWatchRoot = path.dirname(removedParent).replaceAll("\\", "/"); + const parentWatchIndex = watchMock.mock.calls.findLastIndex( + ([watchRoot, options], index) => + watchRoot === parentWatchRoot && options.depth === 0 && !createdWatchers[index]?.closed, + ); + expect(parentWatchIndex).toBeGreaterThanOrEqual(0); + const parentWatcher = createdWatchers[parentWatchIndex]!; await fs.rm(removedParent, { recursive: true }); - deeper.watcher.emit("all", "unlinkDir", logicalRoot); + parentWatcher.emit("all", "unlinkDir", removedParent); await vi.advanceTimersByTimeAsync(250); expect(getSkillsSnapshotVersion(fixtureWorkspaceDir)).toBeGreaterThan(versionBefore); @@ -509,10 +691,16 @@ describe("ensureSkillsWatcher", () => { const rebuilt = watchForSkillRoot(logicalRoot); expect(deeper.watcher.close).toHaveBeenCalledOnce(); expect(rebuilt.watchRoot).toBe(path.dirname(removedParent).replaceAll("\\", "/")); - expect(rebuilt.options.depth).toBe(9); - seen.length = 0; + expect(rebuilt.options.depth).toBe(0); const changedPath = path.join(logicalRoot, "group", "nested", "demo", "SKILL.md"); - rebuilt.watcher.emit("all", "change", changedPath); + await fs.mkdir(path.dirname(changedPath), { recursive: true }); + rebuilt.watcher.emit("all", "addDir", removedParent); + const promoted = watchForSkillRoot(logicalRoot); + expect(promoted.watchRoot).toBe(logicalRoot.replaceAll("\\", "/")); + expect(promoted.options.depth).toBe(7); + await vi.advanceTimersByTimeAsync(250); + seen.length = 0; + promoted.watcher.emit("all", "change", changedPath); await vi.advanceTimersByTimeAsync(250); expect(seen).toEqual([ { workspaceDir: fixtureWorkspaceDir, reason: "watch", changedPath }, @@ -531,8 +719,13 @@ describe("ensureSkillsWatcher", () => { const nestedRoot = path.join(repoDir, "skills"); const watched = watchForSkillRoot(nestedRoot); expect(watched.watchRoot).toBe(repoDir.replaceAll("\\", "/")); - expect(watched.options.depth).toBe(8); + expect(watched.options.depth).toBe(0); expect(watched.options.ignored(path.join(repoDir, "unrelated"))).toBe(true); + await fs.mkdir(path.join(nestedRoot, "group", "demo"), { recursive: true }); + watched.watcher.emit("all", "addDir", nestedRoot); + const promoted = watchForSkillRoot(nestedRoot); + expect(promoted.watchRoot).toBe(nestedRoot.replaceAll("\\", "/")); + expect(promoted.options.depth).toBe(7); }); it("watches nested skills roots for plugin skill dirs", async () => { @@ -571,8 +764,13 @@ describe("ensureSkillsWatcher", () => { const nestedRoot = path.join(pluginDir, "skills"); const watched = watchForSkillRoot(nestedRoot); expect(watched.watchRoot).toBe(pluginDir.replaceAll("\\", "/")); - expect(watched.options.depth).toBe(8); + expect(watched.options.depth).toBe(0); expect(watched.options.ignored(path.join(pluginDir, "unrelated"))).toBe(true); + await fs.mkdir(path.join(nestedRoot, "group", "demo"), { recursive: true }); + watched.watcher.emit("all", "addDir", nestedRoot); + const promoted = watchForSkillRoot(nestedRoot); + expect(promoted.watchRoot).toBe(nestedRoot.replaceAll("\\", "/")); + expect(promoted.options.depth).toBe(7); }); it.runIf(process.platform !== "win32")( @@ -764,7 +962,7 @@ describe("ensureSkillsWatcher", () => { const second = watchForSkillRoot(secondRoot); expect(first.watchRoot).toBe(ancestor.replaceAll("\\", "/")); expect(second.watchRoot).toBe(first.watchRoot); - expect(first.watcher).not.toBe(second.watcher); + expect(first.watcher).toBe(second.watcher); for (const [watched, root, sibling] of [ [first, firstRoot, secondRoot], [second, secondRoot, firstRoot], @@ -772,13 +970,16 @@ describe("ensureSkillsWatcher", () => { for (const included of [ancestor, path.dirname(root), root, path.join(root, "group")]) { expect(watched.options.ignored(included, { isDirectory: () => true })).toBe(false); } - for (const excluded of [sibling, path.join(path.dirname(root), "unrelated")]) { - expect(watched.options.ignored(excluded, { isDirectory: () => true })).toBe(true); - } + expect(watched.options.ignored(sibling, { isDirectory: () => true })).toBe(false); + expect( + watched.options.ignored(path.join(path.dirname(root), "unrelated"), { + isDirectory: () => true, + }), + ).toBe(true); } const seen: SkillsChangeEvent[] = []; refreshModule.registerSkillsChangeListener((change) => seen.push(change)); - const watchers = [first.watcher, second.watcher]; + const watchers = [...new Set([first.watcher, second.watcher])]; for (const watcher of watchers) { watcher.emit("all", "addDir", path.join(ancestor, "unrelated")); watcher.emit("raw", "change", "SKILL.md", { watchedPath: ancestor }); @@ -787,8 +988,17 @@ describe("ensureSkillsWatcher", () => { await vi.advanceTimersByTimeAsync(500); expect(seen).toEqual([]); + await fs.mkdir(firstRoot, { recursive: true }); + first.watcher.emit("all", "addDir", path.dirname(firstRoot)); + expect(first.watcher.close).not.toHaveBeenCalled(); + await fs.mkdir(secondRoot, { recursive: true }); + second.watcher.emit("all", "addDir", path.dirname(secondRoot)); + expect(first.watcher.close).not.toHaveBeenCalled(); + const promoted = [watchForSkillRoot(firstRoot).watcher, watchForSkillRoot(secondRoot).watcher]; + await vi.advanceTimersByTimeAsync(250); + seen.length = 0; const firstChanged = path.join(firstRoot, "demo", "SKILL.md"); - for (const watcher of watchers) { + for (const watcher of promoted) { watcher.emit("all", "change", firstChanged); } await vi.advanceTimersByTimeAsync(250); @@ -797,7 +1007,7 @@ describe("ensureSkillsWatcher", () => { ]); seen.length = 0; const secondChanged = path.join(secondRoot, "demo", "SKILL.md"); - for (const watcher of watchers) { + for (const watcher of promoted) { watcher.emit("raw", "change", "SKILL.md", { watchedPath: path.dirname(secondChanged) }); } await vi.advanceTimersByTimeAsync(500); @@ -805,7 +1015,7 @@ describe("ensureSkillsWatcher", () => { { workspaceDir: secondWorkspace, reason: "watch", changedPath: secondChanged }, ]); seen.length = 0; - for (const watcher of watchers) { + for (const watcher of promoted) { watcher.emit("raw", "rename", undefined, { watchedPath: firstRoot }); } await vi.advanceTimersByTimeAsync(250); @@ -857,223 +1067,4 @@ describe("ensureSkillsWatcher", () => { ]); }, ); - - it("stops fanning a shared-directory change to a workspace after it unsubscribes", async () => { - vi.useFakeTimers(); - const secondWorkspace = await createFixtureDirectory("second-workspace"); - const sharedRoot = await createFixtureDirectory("shared"); - const config = { skills: { load: { extraDirs: [sharedRoot] } } }; - const seen: SkillsChangeEvent[] = []; - refreshModule.registerSkillsChangeListener((change) => { - seen.push(change); - }); - refreshModule.ensureSkillsWatcher({ workspaceDir: fixtureWorkspaceDir, config }); - refreshModule.ensureSkillsWatcher({ workspaceDir: secondWorkspace, config }); - const sharedWatcher = watchForSkillRoot(sharedRoot).watcher; - - refreshModule.ensureSkillsWatcher({ - workspaceDir: fixtureWorkspaceDir, - config: { skills: { load: { extraDirs: [sharedRoot], watch: false } } }, - }); - seen.length = 0; - expect(sharedWatcher.close).not.toHaveBeenCalled(); - const changedPath = path.join(sharedRoot, "demo", "SKILL.md"); - sharedWatcher.emit("all", "change", changedPath); - await vi.advanceTimersByTimeAsync(250); - - expect(seen).toEqual([{ workspaceDir: secondWorkspace, reason: "watch", changedPath }]); - }); - - it("preserves workspace invalidation on watch disable without changing other workspaces", async () => { - vi.useFakeTimers(); - vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); - const workspaceDir = fixtureWorkspaceDir; - const otherWorkspace = await createFixtureDirectory("other-workspace"); - refreshModule.ensureSkillsWatcher({ workspaceDir: otherWorkspace }); - const otherVersion = getSkillsSnapshotVersion(otherWorkspace); - const globalVersion = getSkillsSnapshotVersion(); - refreshModule.ensureSkillsWatcher({ - workspaceDir, - config: { skills: { load: {} } }, - }); - - const firstVersion = bumpSkillsSnapshotVersion({ - workspaceDir, - reason: "watch", - changedPath: `${workspaceDir}/skills/demo/SKILL.md`, - }); - refreshModule.ensureSkillsWatcher({ - workspaceDir, - config: { skills: { load: { watch: false } } }, - }); - - const nextVersion = getSkillsSnapshotVersion(workspaceDir); - expect(nextVersion).toBe(firstVersion); - expect(getSkillsSnapshotVersion(otherWorkspace)).toBe(otherVersion); - expect(getSkillsSnapshotVersion()).toBe(globalVersion); - vi.setSystemTime(new Date(nextVersion)); - refreshModule.ensureSkillsWatcher({ workspaceDir }); - expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(nextVersion); - }); - - it("evicts idle workspace subscriptions on a later ensure call", async () => { - vi.useFakeTimers(); - vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); - const idleWorkspaceDir = fixtureWorkspaceDir; - const activeWorkspaceDir = await createFixtureDirectory("workspace-active"); - refreshModule.ensureSkillsWatcher({ - workspaceDir: idleWorkspaceDir, - config: { skills: { load: {} } }, - }); - const idleSkillsWatcher = watchForSkillRoot(path.join(idleWorkspaceDir, "skills")).watcher; - const firstVersion = bumpSkillsSnapshotVersion({ - workspaceDir: idleWorkspaceDir, - reason: "watch", - }); - const globalVersion = getSkillsSnapshotVersion(); - - vi.advanceTimersByTime(60 * 60_000 + 1_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir: activeWorkspaceDir, - config: { skills: { load: {} } }, - }); - - expect(idleSkillsWatcher.close).toHaveBeenCalledTimes(1); - const evictedVersion = getSkillsSnapshotVersion(idleWorkspaceDir); - expect(evictedVersion).toBe(firstVersion); - expect(getSkillsSnapshotVersion()).toBe(globalVersion); - vi.setSystemTime(new Date(evictedVersion)); - refreshModule.ensureSkillsWatcher({ workspaceDir: idleWorkspaceDir }); - expect(getSkillsSnapshotVersion(idleWorkspaceDir)).toBeGreaterThan(evictedVersion); - }); - - it("keeps another execution subscription for the workspace alive after disposal", async () => { - vi.useFakeTimers(); - const workspaceDir = fixtureWorkspaceDir; - const executionWorkspaceDir = await createFixtureDirectory("remaining-worktree"); - refreshModule.ensureSkillsWatcher({ workspaceDir }); - refreshModule.ensureSkillsWatcher({ workspaceDir, executionWorkspaceDir }); - const version = getSkillsSnapshotVersion(workspaceDir); - const globalVersion = getSkillsSnapshotVersion(); - const watcher = watchForSkillRoot(path.join(workspaceDir, "skills")).watcher; - const seen: SkillsChangeEvent[] = []; - refreshModule.registerSkillsChangeListener((change) => seen.push(change)); - - refreshModule.ensureSkillsWatcher({ - workspaceDir, - config: { skills: { load: { watch: false } } }, - }); - - expect(watcher.close).not.toHaveBeenCalled(); - expect(getSkillsSnapshotVersion(workspaceDir)).toBe(version); - expect(getSkillsSnapshotVersion()).toBe(globalVersion); - expect(seen).toEqual([]); - const changedPath = path.join(workspaceDir, "skills", "demo", "SKILL.md"); - watcher.emit("all", "change", changedPath); - await vi.advanceTimersByTimeAsync(250); - expect(seen).toEqual([{ workspaceDir, reason: "watch", changedPath }]); - }); - - it("keeps refreshed workspace subscriptions within the idle TTL", async () => { - vi.useFakeTimers(); - vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); - const activeWorkspaceDir = fixtureWorkspaceDir; - const otherWorkspaceDir = await createFixtureDirectory("workspace-other"); - refreshModule.ensureSkillsWatcher({ - workspaceDir: activeWorkspaceDir, - config: { skills: { load: {} } }, - }); - const activeSkillsWatcher = watchForSkillRoot(path.join(activeWorkspaceDir, "skills")).watcher; - - vi.advanceTimersByTime(30 * 60_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir: activeWorkspaceDir, - config: { skills: { load: {} } }, - }); - vi.advanceTimersByTime(31 * 60_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir: otherWorkspaceDir, - config: { skills: { load: {} } }, - }); - - expect(activeSkillsWatcher.close).not.toHaveBeenCalled(); - }); - - it.each(["execution", "base"] as const)( - "keeps an idle %s source active while another consumer remains", - async (scope) => { - vi.useFakeTimers(); - const { loadWorkspaceSkills } = await import("../loading/workspace-skill-loader.js"); - const workspaceDir = fixtureWorkspaceDir; - const executionWorkspaceDir = await createFixtureDirectory("shared-execution"); - const idleScope = { - executionWorkspaceDir: scope === "execution" ? executionWorkspaceDir : undefined, - }; - const skillDir = await createFixtureDirectory( - scope === "execution" ? "shared-execution/skills/demo" : "workspace/skills/demo", - ); - const skillFile = path.join(skillDir, "SKILL.md"); - await fs.writeFile( - skillFile, - "---\nname: demo\ndescription: Demo\n---\nOriginal instructions\n", - ); - await withEnvAsync({ OPENCLAW_STATE_DIR: workspaceDir }, async () => { - const options = { - ...idleScope, - agentId: "agent-b", - bundledSkillsDir: "", - managedSkillsDir: path.join(workspaceDir, "missing-managed"), - }; - const original = loadWorkspaceSkills(workspaceDir, options)[0]!.skill.contentHash; - refreshModule.ensureSkillsWatcher({ - workspaceDir, - ...idleScope, - agentId: "agent-a", - }); - refreshModule.ensureSkillsWatcher({ - workspaceDir, - executionWorkspaceDir, - agentId: "agent-b", - }); - vi.advanceTimersByTime(30 * 60_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir, - executionWorkspaceDir, - agentId: "agent-b", - }); - const sourceVersion = getSkillsSourceVersion(workspaceDir, idleScope); - vi.advanceTimersByTime(31 * 60_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir, - executionWorkspaceDir, - agentId: "agent-b", - }); - expect(getSkillsSourceVersion(workspaceDir, idleScope)).toBe(sourceVersion); - - const version = getSkillsSnapshotVersion(workspaceDir); - await fs.appendFile(skillFile, "\nUpdated instructions\n"); - bumpSkillsSnapshotVersion({ reason: "workshop" }); - expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(version); - expect(loadWorkspaceSkills(workspaceDir, options)[0]!.skill.contentHash).not.toBe(original); - - vi.advanceTimersByTime(60 * 60_000 + 1_000); - refreshModule.ensureSkillsWatcher({ - workspaceDir: await createFixtureDirectory("other-active-workspace"), - }); - const retiredVersion = getSkillsSnapshotVersion(workspaceDir); - await fs.appendFile(skillFile, "\nInstructions changed while retired\n"); - bumpSkillsSnapshotVersion({ reason: "workshop" }); - expect(getSkillsSnapshotVersion(workspaceDir)).toBe(retiredVersion); - - refreshModule.ensureSkillsWatcher({ - workspaceDir, - executionWorkspaceDir: - scope === "base" - ? await createFixtureDirectory("new-execution-workspace") - : executionWorkspaceDir, - }); - expect(getSkillsSnapshotVersion(workspaceDir)).toBeGreaterThan(retiredVersion); - }); - }, - ); }); diff --git a/src/skills/runtime/refresh.ts b/src/skills/runtime/refresh.ts index 50dcaf43edc1..ba1c8888cf4c 100644 --- a/src/skills/runtime/refresh.ts +++ b/src/skills/runtime/refresh.ts @@ -1,12 +1,10 @@ import { AsyncLocalStorage } from "node:async_hooks"; -import fs from "node:fs"; import os from "node:os"; import path from "node:path"; import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; -import chokidar, { type FSWatcher } from "chokidar"; +import chokidar from "chokidar"; import { isDefaultStateDir } from "../../config/paths.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; -import { resolveRealpathOrAbsolute } from "../../infra/boundary-path.js"; import { getFileWatchCapacityCode } from "../../infra/fs-watch-errors.js"; import { isPathInside } from "../../infra/path-guards.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; @@ -16,16 +14,14 @@ import { resolvePluginSkillRoots, resolvePluginSkillRootsFromMetadata, } from "../loading/plugin-skills.js"; -import { - resolveAllowedSkillSymlinkTargetRealPaths, - tryRealpath, -} from "../loading/symlink-targets.js"; +import { resolveAllowedSkillSymlinkTargetRealPaths } from "../loading/symlink-targets.js"; import { normalizeWorkspaceSkillRoots, resolveWorkspaceSkillDirectories, } from "../loading/workspace-skill-roots.js"; import { resolveWorkshopWatchRoots } from "../workshop/skills-root.js"; import { areOrderedArraysEqual } from "./ordered-array-equality.js"; +import { acquireSkillsAncestorWatcher } from "./refresh-ancestor-watch.js"; import { createRawSkillFileScheduler } from "./refresh-file-stability.js"; import { bumpSkillsSnapshotVersion, @@ -38,23 +34,28 @@ import { import { joinSkillsWatcherCloses, teardownSkillsPathWatcher } from "./refresh-watch-close.js"; import { createSkillsWatchPathFilter, - DEFAULT_SKILLS_WATCH_IGNORED, getRawWatchedPath, isSkillDiscoveryFileWatchPath, - isTrustedSymlinkSkillTarget, rawPathToString, - readBudgetedDirEntries, resolveRawSkillsWatchPath, makeSkillsWatchTarget, resolveSkillsWatcherUsePolling, toWatchRoot, } from "./refresh-watch-path.js"; +import { + addSkillSourceWatchTargets, + GROUPED_SKILLS_WATCH_DEPTH, + type WatchTarget, +} from "./refresh-watch-targets.js"; export { registerSkillsChangeListener } from "./refresh-state.js"; type SkillsWatchChange = "skills" | "supporting"; type SkillsPathWatchState = { - watcher: FSWatcher; + closed: boolean; + close: () => void; + schedule: (path?: string) => void; watchRoot: string; + ancestorRoot: string; depth: number; initialScan: "pending" | "ready" | "error"; timer?: ReturnType; @@ -63,13 +64,6 @@ type SkillsPathWatchState = { readonly subscribers: Set; }; -type WatchTarget = { - path: string; - watchRoot: string; - depth: number; - executionOnly?: true; -}; - type WatchTargetCacheEntry = { signature: string; targets: WatchTarget[]; @@ -79,11 +73,6 @@ const log = createSubsystemLogger("gateway/skills"); // Gateway startup imports this owner before serving turns. Shared watcher handles, // including later rebuilds, must inherit that lifetime rather than the triggering turn. const runInSkillsWatcherContext = AsyncLocalStorage.snapshot(); -const GROUPED_SKILLS_WATCH_DEPTH = 6; -const CONFIGURED_ROOT_WATCH_DEPTH = 2; -const MAX_SYMLINK_WATCH_TARGETS_PER_ROOT = 100; -const MAX_SYMLINK_WATCH_DIRECTORY_SCANS_PER_ROOT = 200; -const MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT = 2_000; const SKILLS_WATCH_DEBOUNCE_MS = 250; // One watcher per unique watched directory. Agent workspaces that include the // same shared skill root (the global skills dir, the home skills dir, or a @@ -205,171 +194,97 @@ function resolveWatchTargets( return sortedTargets; } -function addWatchTarget(targets: Map, raw: string, depth: number): void { - const target = makeSkillsWatchTarget(raw, depth); - target.depth = Math.max(target.depth, targets.get(target.path)?.depth ?? 0); - targets.set(target.path, target); -} - -function addSkillRootWatchTargets( - targets: Map, - root: string, - rootDepth: number, -): string { - addWatchTarget(targets, root, rootDepth); - const companionSkillsRoot = path.join(root, "skills"); - addWatchTarget(targets, companionSkillsRoot, GROUPED_SKILLS_WATCH_DEPTH); - return companionSkillsRoot; -} - -function addSkillSourceWatchTargets( - targets: Map, - root: string, - source: string, - allowedSymlinkTargetRealPaths: readonly string[], - rootDepth = path.basename(root) === "skills" - ? GROUPED_SKILLS_WATCH_DEPTH - : CONFIGURED_ROOT_WATCH_DEPTH, -): void { - const companionSkillsRoot = addSkillRootWatchTargets(targets, root, rootDepth); - // Both bounded scans share the source's containment identity for this preparation. - // Trusted symlink leaves below remain registration-only, never recursive scans. - const rootRealPath = resolveRealpathOrAbsolute(root); - addTrustedSymlinkSkillWatchTargets( - targets, - root, - source, - allowedSymlinkTargetRealPaths, - rootDepth, - rootRealPath, - rootRealPath, - ); - addTrustedSymlinkSkillWatchTargets( - targets, - companionSkillsRoot, - source, - allowedSymlinkTargetRealPaths, - GROUPED_SKILLS_WATCH_DEPTH, - rootRealPath, - resolveRealpathOrAbsolute(companionSkillsRoot), - ); -} - -function addTrustedSymlinkSkillWatchTargets( - targets: Map, - root: string, - source: string, - allowedSymlinkTargetRealPaths: readonly string[], - maxDepth: number, - containmentRootRealPath: string, - rootRealPath: string, -): void { - try { - if ( - fs.lstatSync(root).isSymbolicLink() && - isTrustedSymlinkSkillTarget( - source, - containmentRootRealPath, - rootRealPath, - allowedSymlinkTargetRealPaths, - ) - ) { - addSkillRootWatchTargets(targets, rootRealPath, maxDepth); - } - } catch { - return; - } - const queue: Array<{ dir: string; depth: number }> = [{ dir: root, depth: 0 }]; - let watched = 0; - let directoryScans = 0; - let rawEntries = 0; - for (const queued of queue) { - if ( - watched >= MAX_SYMLINK_WATCH_TARGETS_PER_ROOT || - directoryScans >= MAX_SYMLINK_WATCH_DIRECTORY_SCANS_PER_ROOT || - rawEntries >= MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT - ) { - break; - } - const current = queued; - if (!current) { - continue; - } - const scan = readBudgetedDirEntries( - current.dir, - MAX_SYMLINK_WATCH_RAW_ENTRIES_PER_ROOT - rawEntries, - ); - directoryScans += 1; - rawEntries += scan.scannedEntryCount; - if (!scan.ok) { - continue; - } - for (const entry of scan.entries.toSorted((a, b) => a.name.localeCompare(b.name))) { - if (watched >= MAX_SYMLINK_WATCH_TARGETS_PER_ROOT) { - break; - } - if (entry.name.startsWith(".") || entry.name === "node_modules") { - continue; - } - const childPath = path.join(current.dir, entry.name); - if (DEFAULT_SKILLS_WATCH_IGNORED.some((re) => re.test(childPath))) { - continue; - } - if (entry.isSymbolicLink()) { - const targetRealPath = tryRealpath(childPath); - if ( - targetRealPath && - isTrustedSymlinkSkillTarget( - source, - containmentRootRealPath, - targetRealPath, - allowedSymlinkTargetRealPaths, - ) - ) { - addSkillRootWatchTargets(targets, targetRealPath, GROUPED_SKILLS_WATCH_DEPTH); - watched += 1; - } - continue; - } - if (entry.isDirectory() && current.depth < maxDepth) { - queue.push({ dir: childPath, depth: current.depth + 1 }); - } - } - } -} - -function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { +function createSkillsPathWatcher( + target: WatchTarget, + previousAncestorRoot = target.watchRoot, +): SkillsPathWatchState { const usePolling = resolveSkillsWatcherUsePolling(); const pathFilter = createSkillsWatchPathFilter(target.path, usePolling); - // Chokidar's missing-root fallback retains only the final basename, so it - // misses creation through multiple absent parents. Watch the existing prefix - // and restrict traversal to the logical root and its ancestor chain. - const watcher = runInSkillsWatcherContext(() => - chokidar.watch(target.watchRoot, { - ignoreInitial: true, - followSymlinks: false, - usePolling, - // Observe every discovered skill plus its identity metadata directory, - // which sits one level below the deepest admitted skill directory. - depth: - target.depth + - 1 + - path.relative(target.watchRoot, target.path).split(path.sep).filter(Boolean).length, - awaitWriteFinish: { - stabilityThreshold: SKILLS_WATCH_DEBOUNCE_MS, - pollInterval: 100, - }, - ignored: pathFilter.ignored, - }), - ); - + // Descendant native watches do not report ancestor moves. Keep shallow + // observation along the original path even after its content watch promotes. + const ancestorRoot = isPathInside(previousAncestorRoot, target.watchRoot) + ? previousAncestorRoot + : target.watchRoot; + const ancestorRoots: string[] = []; + let currentRoot = target.watchRoot; + while (isPathInside(ancestorRoot, currentRoot)) { + if (currentRoot !== target.path) { + ancestorRoots.push(currentRoot); + } + const parent = toWatchRoot(path.dirname(currentRoot)); + if (parent === currentRoot) { + break; + } + currentRoot = parent; + } + const pendingAncestors = new Set(ancestorRoots); + let contentReady = target.path !== target.watchRoot; + const watcher = + target.path === target.watchRoot + ? runInSkillsWatcherContext(() => + chokidar.watch(target.path, { + ignoreInitial: true, + followSymlinks: false, + usePolling, + // Identity metadata sits one level below the deepest admitted skill. + depth: target.depth + 1, + awaitWriteFinish: { + stabilityThreshold: SKILLS_WATCH_DEBOUNCE_MS, + pollInterval: 100, + }, + ignored: pathFilter.ignored, + }), + ) + : undefined; + const releaseAncestors: (() => void)[] = []; const state: SkillsPathWatchState = { - watcher, + closed: false, + close: () => { + if (state.closed) { + return; + } + state.closed = true; + clearTimeout(state.timer); + if (watcher) { + void teardownSkillsPathWatcher({ watcher }); + } + for (const release of releaseAncestors) { + release(); + } + }, + schedule: (changedPath) => schedule(changedPath), watchRoot: target.watchRoot, + ancestorRoot, depth: target.depth, initialScan: "pending", subscribers: new Set(), }; + const isCurrent = () => !state.closed && pathWatchers.get(target.path) === state; + const reconcileRoot = (changedPath?: string) => { + if (!isCurrent()) { + return true; + } + const nextTarget = makeSkillsWatchTarget(target.path, state.depth, state.ancestorRoot); + if (nextTarget.watchRoot === state.watchRoot) { + return false; + } + for (const subscriber of state.subscribers) { + workspaceWatchTargetCache.delete(subscriber); + for (const entry of workspaceWatchTargets.get(subscriber) ?? []) { + if (entry.path === target.path) { + entry.watchRoot = nextTarget.watchRoot; + } + } + } + const subscriber = state.subscribers.values().next().value; + if (subscriber !== undefined) { + subscribeWorkspaceToPath(subscriber, nextTarget); + if (changedPath) { + pathWatchers.get(target.path)?.schedule(changedPath); + } + } + return true; + }; const publishChanges = ( watcherKeys: Iterable, @@ -415,12 +330,7 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { } }; const settleInitialScan = (result: "ready" | "error") => { - if ( - watcher.closed || - pathWatchers.get(target.path) !== state || - state.initialScan === "ready" || - state.initialScan === result - ) { + if (!isCurrent() || state.initialScan === "ready" || state.initialScan === result) { return; } state.initialScan = result; @@ -430,7 +340,7 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { if ( targets?.every((entry) => { const current = pathWatchers.get(entry.path); - return current && !current.watcher.closed && current.initialScan !== "pending"; + return current && !current.closed && current.initialScan !== "pending"; }) ) { readySubscribers.push(watcherKey); @@ -441,13 +351,16 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { const schedule = (changedPath?: string, change: SkillsWatchChange = "skills") => { // File-stability work may finish after this subscription has been closed. - if (watcher.closed || (change === "supporting" && state.pendingChange === "skills")) { + if (!isCurrent() || (change === "supporting" && state.pendingChange === "skills")) { return; } state.pendingPath = changedPath ?? state.pendingPath; state.pendingChange = change; clearTimeout(state.timer); state.timer = setTimeout(() => { + if (!isCurrent()) { + return; + } const pendingPath = state.pendingPath; const pendingChange = state.pendingChange; state.pendingPath = undefined; @@ -459,7 +372,7 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { }, SKILLS_WATCH_DEBOUNCE_MS); }; const scheduleRawSkillFile = createRawSkillFileScheduler({ - watcher, + watcher: state, stabilityMs: SKILLS_WATCH_DEBOUNCE_MS, schedule, onError: (changedPath, err) => { @@ -470,14 +383,32 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { // ignoreInitial suppresses writes discovered before native watches are ready. // Reconcile the whole workspace once its initial scans finish, rather than // rebuilding metadata for every root that becomes ready. - watcher.on("ready", () => settleInitialScan("ready")); - watcher.on("all", (event, changedPath) => { + const ready = () => { + if (!reconcileRoot() && contentReady && pendingAncestors.size === 0) { + settleInitialScan("ready"); + } + }; + watcher?.on("ready", () => { + contentReady = true; + ready(); + }); + const onChange = (event: string, changedPath: string) => { + if ( + !isCurrent() || + ((!watcher || event === "addDir" || event === "unlinkDir") && reconcileRoot(changedPath)) + ) { + return; + } const skillsRelevant = pathFilter.isRelevant(event, changedPath); if (skillsRelevant || pathFilter.isSupportingPath(changedPath)) { schedule(changedPath, skillsRelevant ? "skills" : "supporting"); } - }); - watcher.on("raw", (_eventName, rawPath, details) => { + }; + watcher?.on("all", onChange); + const onRaw = (_eventName: string, rawPath: unknown, details: unknown) => { + if (!isCurrent()) { + return; + } const rawPathText = rawPathToString(rawPath); if (!rawPathText) { const watchedPath = getRawWatchedPath(details); @@ -500,9 +431,10 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { } else if (changedPath && pathFilter.isSupportingPath(changedPath)) { schedule(changedPath, "supporting"); } - }); - watcher.on("error", (err) => { - if (watcher.closed) { + }; + watcher?.on("raw", onRaw); + const onError = (err: unknown) => { + if (!isCurrent()) { return; } const capacityCode = usePolling ? undefined : getFileWatchCapacityCode(err); @@ -513,7 +445,7 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { `skills native watcher capacity exhausted (${capacityCode}); refreshing skills during agent preparation`, ); for (const active of pathWatchers.values()) { - void teardownSkillsPathWatcher(active); + active.close(); } } return; @@ -522,7 +454,30 @@ function createSkillsPathWatcher(target: WatchTarget): SkillsPathWatchState { // A failed scan may never emit ready. Let healthy roots reconcile; if the // failed scan continues, its eventual ready still closes that read gap. settleInitialScan("error"); - }); + }; + watcher?.on("error", onError); + for (const root of ancestorRoots) { + releaseAncestors.push( + acquireSkillsAncestorWatcher(root, usePolling, { + path: target.path, + ignored: pathFilter.ignored, + ready: () => { + pendingAncestors.delete(root); + ready(); + }, + changed: onChange, + raw: onRaw, + error: (error) => { + pendingAncestors.delete(root); + onError(error); + // Missing roots still need the ancestor's recovery scan. + if (watcher) { + ready(); + } + }, + }).release, + ); + } return state; } @@ -539,10 +494,13 @@ function subscribeWorkspaceToPath(workspaceDir: string, watchTarget: WatchTarget } if (existing) { // A changed ancestor or deeper target needs a rebuilt watcher, preserving subscribers. - const next = createSkillsPathWatcher({ - ...watchTarget, - depth: Math.max(existing.depth, watchTarget.depth), - }); + const next = createSkillsPathWatcher( + { + ...watchTarget, + depth: Math.max(existing.depth, watchTarget.depth), + }, + existing.ancestorRoot, + ); for (const subscriber of existing.subscribers) { next.subscribers.add(subscriber); const owner = workspaceWatchOwners.get(subscriber); @@ -556,7 +514,7 @@ function subscribeWorkspaceToPath(workspaceDir: string, watchTarget: WatchTarget } } next.subscribers.add(workspaceDir); - void teardownSkillsPathWatcher(existing); + existing.close(); pathWatchers.set(watchTarget.path, next); return; } @@ -572,7 +530,7 @@ function unsubscribeWorkspaceFromPath(workspaceDir: string, watchTarget: WatchTa } state.subscribers.delete(workspaceDir); if (state.subscribers.size === 0) { - void teardownSkillsPathWatcher(state); + state.close(); pathWatchers.delete(watchTarget.path); } } @@ -752,7 +710,7 @@ export async function closeSkillsWatchers(resetState = false): Promise { workspaceWatchTargetCache.clear(); workspaceWatchLastEnsuredAt.clear(); for (const state of active) { - void teardownSkillsPathWatcher(state); + state.close(); } await joinSkillsWatcherCloses(); } diff --git a/src/skills/runtime/refresh.watcher.test-support.ts b/src/skills/runtime/refresh.watcher.test-support.ts index 51d56da59a69..aad0dcb0031b 100644 --- a/src/skills/runtime/refresh.watcher.test-support.ts +++ b/src/skills/runtime/refresh.watcher.test-support.ts @@ -80,13 +80,28 @@ export function createSkillsWatcherMock() { return watcher; }); function watchForSkillRoot(root: string) { - // Distinguish logical subscriptions that share one physical ancestor by - // their public traversal filter, rather than depending on watcher order. - const index = watchMock.mock.calls.findLastIndex( - ([, options]) => - !options.ignored(path.join(root, "SKILL.md")) && - options.ignored(path.join(path.dirname(root), "SKILL.md")), + // Existing roots have their own recursive watcher. Missing roots share a + // shallow ancestor whose public traversal filter admits the logical path. + const normalizedRoot = root.replaceAll("\\", "/"); + let index = watchMock.mock.calls.findLastIndex( + ([watchRoot, options], candidate) => + watchRoot === normalizedRoot && options.depth > 0 && !createdWatchers[candidate]?.closed, ); + if (index < 0) { + let closest = -1; + for (const [candidate, [watchRoot, options]] of watchMock.mock.calls.entries()) { + if ( + !createdWatchers[candidate]?.closed && + options.depth === 0 && + normalizedRoot.startsWith(watchRoot.endsWith("/") ? watchRoot : `${watchRoot}/`) && + !options.ignored(path.join(root, "SKILL.md")) && + watchRoot.length >= closest + ) { + index = candidate; + closest = watchRoot.length; + } + } + } expect(index, `watch subscription for ${root}`).toBeGreaterThanOrEqual(0); const [watchRoot, options] = watchMock.mock.calls[index]!; return { watchRoot, options, watcher: createdWatchers[index]! }; diff --git a/src/skills/runtime/refresh.windows.test.ts b/src/skills/runtime/refresh.windows.test.ts index d12d6dbcb3e5..fd2967fafbdc 100644 --- a/src/skills/runtime/refresh.windows.test.ts +++ b/src/skills/runtime/refresh.windows.test.ts @@ -87,7 +87,7 @@ describe("Windows skills watcher paths", () => { normalized(scenario === "existing" ? skillsRoot : path.dirname(skillsRoot)), ); expect(skillsWatch.options).toMatchObject({ - depth: scenario === "existing" ? 7 : 8, + depth: scenario === "existing" ? 7 : 0, followSymlinks: false, }); const repoSkillsRoot = path.join(longRoot, "repo", "skills");