mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
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
This commit is contained in:
parent
b9e18567a6
commit
ca2e2cc268
11 changed files with 1266 additions and 458 deletions
163
src/skills/runtime/refresh-ancestor-watch.ts
Normal file
163
src/skills/runtime/refresh-ancestor-watch.ts
Normal file
|
|
@ -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<typeof createSkillsWatchPathFilter>["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<AncestorSubscription>;
|
||||
};
|
||||
|
||||
// Imported with the refresh owner at Gateway startup, outside turn contexts.
|
||||
const runInWatcherContext = AsyncLocalStorage.snapshot();
|
||||
const ancestorWatchers = new Map<string, AncestorWatcher>();
|
||||
|
||||
function createAncestorWatcher(
|
||||
watchRoot: string,
|
||||
usePolling: boolean,
|
||||
subscriptions: Set<AncestorSubscription>,
|
||||
): 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<AncestorSubscription>([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);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
|
@ -20,7 +20,7 @@ function readFileStabilitySnapshot(filePath: string): FileStabilitySnapshot | un
|
|||
async function waitForStableSkillFile(
|
||||
filePath: string,
|
||||
stabilityMs: number,
|
||||
watcher: FSWatcher,
|
||||
watcher: Pick<FSWatcher, "closed">,
|
||||
readRevision: () => number,
|
||||
): Promise<void> {
|
||||
if (watcher.closed || stabilityMs <= 0) {
|
||||
|
|
@ -66,7 +66,7 @@ export function createRawSkillFileScheduler({
|
|||
schedule,
|
||||
onError,
|
||||
}: {
|
||||
watcher: FSWatcher;
|
||||
watcher: Pick<FSWatcher, "closed">;
|
||||
stabilityMs: number;
|
||||
schedule: (filePath: string) => void;
|
||||
onError: (filePath: string, error: unknown) => void;
|
||||
|
|
|
|||
|
|
@ -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 };
|
||||
}
|
||||
|
||||
|
|
|
|||
156
src/skills/runtime/refresh-watch-targets.ts
Normal file
156
src/skills/runtime/refresh-watch-targets.ts
Normal file
|
|
@ -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<string, WatchTarget>, 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<string, WatchTarget>,
|
||||
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<string, WatchTarget>,
|
||||
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<string, WatchTarget>,
|
||||
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 });
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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[] = [];
|
||||
|
|
|
|||
|
|
@ -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<typeof chokidar.watch>; 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<typeof clearTimeout>[0],
|
||||
{ settled: Promise<void>; 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<void>((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<void>((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");
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
|
|||
257
src/skills/runtime/refresh.subscriptions.test.ts
Normal file
257
src/skills/runtime/refresh.subscriptions.test.ts
Normal file
|
|
@ -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<Parameters<typeof bumpSkillsSnapshotVersion>[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);
|
||||
});
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
@ -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);
|
||||
});
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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<typeof setTimeout>;
|
||||
|
|
@ -63,13 +64,6 @@ type SkillsPathWatchState = {
|
|||
readonly subscribers: Set<string>;
|
||||
};
|
||||
|
||||
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<string, WatchTarget>, 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<string, WatchTarget>,
|
||||
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<string, WatchTarget>,
|
||||
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<string, WatchTarget>,
|
||||
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<string>(),
|
||||
};
|
||||
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<string>,
|
||||
|
|
@ -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<void> {
|
|||
workspaceWatchTargetCache.clear();
|
||||
workspaceWatchLastEnsuredAt.clear();
|
||||
for (const state of active) {
|
||||
void teardownSkillsPathWatcher(state);
|
||||
state.close();
|
||||
}
|
||||
await joinSkillsWatcherCloses();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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]! };
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue