refactor(skill): scope provider invalidation
This commit is contained in:
@@ -26,6 +26,7 @@ import {
|
||||
type SkillDefinition,
|
||||
type SkillLookupOptions,
|
||||
type SkillProvider,
|
||||
type SkillProviderControl,
|
||||
type SkillSource,
|
||||
} from '@deepseek-ai/dsh-skill'
|
||||
|
||||
@@ -119,8 +120,11 @@ interface ResolvedWatchConfig {
|
||||
|
||||
/** Register the local filesystem skill provider on `ctx.skills`. */
|
||||
export function apply(ctx: Context, config: Config = {}): void {
|
||||
const provider = new LocalSkillProvider(ctx, config)
|
||||
ctx.skills.registerProvider(provider)
|
||||
let provider!: LocalSkillProvider
|
||||
ctx.skills.registerProvider((control) => {
|
||||
provider = new LocalSkillProvider(ctx, control, config)
|
||||
return provider
|
||||
})
|
||||
ctx.effect(function* () {
|
||||
yield async () => { await provider.dispose() }
|
||||
}, 'skill-local watcher')
|
||||
@@ -138,12 +142,18 @@ export class LocalSkillProvider implements SkillProvider {
|
||||
private readonly customSkillDirs: string[]
|
||||
private readonly watchManager: SkillWatchManager
|
||||
private readonly bundledSkillDir: string | undefined
|
||||
private disposal: Promise<void> | undefined
|
||||
|
||||
constructor(private readonly ctx: Context, config: Config = {}) {
|
||||
constructor(
|
||||
private readonly ctx: Context,
|
||||
control: SkillProviderControl,
|
||||
config: Config = {},
|
||||
) {
|
||||
this.dshHome = resolveDshHome(config.dshHome)
|
||||
this.agentsHome = resolve(config.agentsHome ?? process.env.DSH_AGENTS_HOME ?? join(homedir(), '.agents'))
|
||||
this.customSkillDirs = (config.customSkillDirs ?? []).map(root => resolve(root))
|
||||
this.watchManager = new SkillWatchManager(ctx, this, resolveWatchConfig(config))
|
||||
this.watchManager = new SkillWatchManager(ctx, control.invalidate, resolveWatchConfig(config))
|
||||
control.signal.addEventListener('abort', () => { void this.dispose() }, { once: true })
|
||||
const bundledSkillDir = config.bundledSkillDir ?? process.env.DSH_BUNDLED_SKILL_DIR
|
||||
this.bundledSkillDir = bundledSkillDir === undefined ? undefined : resolve(bundledSkillDir)
|
||||
}
|
||||
@@ -197,9 +207,13 @@ export class LocalSkillProvider implements SkillProvider {
|
||||
this.watchManager.observeHostMutation(path)
|
||||
}
|
||||
|
||||
/** Close every host watcher and contain late filesystem callbacks. */
|
||||
async dispose(): Promise<void> {
|
||||
await this.watchManager.dispose()
|
||||
/**
|
||||
* Close every host watcher and contain late filesystem callbacks.
|
||||
* @returns a shared promise that settles when every watcher reaches quiescence.
|
||||
*/
|
||||
dispose(): Promise<void> {
|
||||
this.disposal ??= this.watchManager.dispose()
|
||||
return this.disposal
|
||||
}
|
||||
|
||||
private async roots(cwd: string | undefined): Promise<SkillRoot[]> {
|
||||
@@ -250,7 +264,7 @@ class SkillWatchManager {
|
||||
|
||||
constructor(
|
||||
private readonly ctx: Context,
|
||||
private readonly provider: SkillProvider,
|
||||
private readonly invalidate: () => void,
|
||||
private readonly config: ResolvedWatchConfig,
|
||||
) {}
|
||||
|
||||
@@ -286,18 +300,17 @@ class SkillWatchManager {
|
||||
evictedProject = true
|
||||
}
|
||||
await Promise.all(pending)
|
||||
if (evictedProject) this.ctx.skills.invalidateProvider(this.provider)
|
||||
if (evictedProject) this.invalidate()
|
||||
}
|
||||
|
||||
observeHostMutation(path: string): void {
|
||||
if (this.closing) return
|
||||
const normalized = resolve(path)
|
||||
if (![...this.roots.values()].some(state => isPotentialSkillPath(state.root, normalized))) return
|
||||
this.ctx.skills.invalidateProvider(this.provider)
|
||||
this.invalidate()
|
||||
}
|
||||
|
||||
async dispose(): Promise<void> {
|
||||
if (this.closing) return
|
||||
this.closing = true
|
||||
const states = [...this.roots.values()]
|
||||
this.roots.clear()
|
||||
@@ -376,6 +389,8 @@ class SkillWatchManager {
|
||||
}
|
||||
}
|
||||
|
||||
// FIXME(file-watch-service): Extract Chokidar and missing-root observation below into a Cordis
|
||||
// service; keep skill filtering and invalidation here.
|
||||
private async openStableWatcher(state: RootWatchState): Promise<WatchHandle | undefined> {
|
||||
while (!this.closing && state.owners.size > 0) {
|
||||
const mode = await resolveRootWatchMode(state.root.path)
|
||||
@@ -489,7 +504,7 @@ class SkillWatchManager {
|
||||
queueMicrotask(() => {
|
||||
this.invalidationQueued = false
|
||||
if (this.closing) return
|
||||
this.ctx.skills.invalidateProvider(this.provider)
|
||||
this.invalidate()
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -123,12 +123,8 @@ describe('skill-local watcher failures', () => {
|
||||
watchStabilityThresholdMs: 20,
|
||||
})
|
||||
expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['watched-skill'])
|
||||
const invalidateProvider = ctx.skills.invalidateProvider.bind(ctx.skills)
|
||||
let invalidations = 0
|
||||
ctx.skills.invalidateProvider = (provider) => {
|
||||
invalidations += 1
|
||||
invalidateProvider(provider)
|
||||
}
|
||||
ctx.on('skills/change', () => { invalidations += 1 })
|
||||
const first = watcherHarness.watchers[0]
|
||||
if (first === undefined) throw new Error('expected a root watcher')
|
||||
|
||||
@@ -169,14 +165,17 @@ describe('skill-local watcher failures', () => {
|
||||
watcherHarness.deferredReady = 1
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SkillService)
|
||||
const provider = new SkillLocal.LocalSkillProvider(ctx, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
watch: true,
|
||||
watchPollIntervalMs: 10,
|
||||
watchStabilityThresholdMs: 20,
|
||||
let provider!: InstanceType<typeof SkillLocal.LocalSkillProvider>
|
||||
const disposeProvider = ctx.skills.registerProvider((control) => {
|
||||
provider = new SkillLocal.LocalSkillProvider(ctx, control, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
watch: true,
|
||||
watchPollIntervalMs: 10,
|
||||
watchStabilityThresholdMs: 20,
|
||||
})
|
||||
return provider
|
||||
})
|
||||
ctx.skills.registerProvider(provider)
|
||||
|
||||
const discovery = provider.list({})
|
||||
await settle()
|
||||
@@ -187,6 +186,7 @@ describe('skill-local watcher failures', () => {
|
||||
first.emitter.emit('ready')
|
||||
|
||||
await Promise.all([discovery, disposal])
|
||||
disposeProvider()
|
||||
await settle()
|
||||
expect(first.closeCalls).toBeGreaterThan(0)
|
||||
})
|
||||
@@ -198,14 +198,17 @@ describe('skill-local watcher failures', () => {
|
||||
watcherHarness.deferredReady = 1
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SkillService)
|
||||
const provider = new SkillLocal.LocalSkillProvider(ctx, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
watch: true,
|
||||
watchPollIntervalMs: 10,
|
||||
watchStabilityThresholdMs: 20,
|
||||
let provider!: InstanceType<typeof SkillLocal.LocalSkillProvider>
|
||||
const disposeProvider = ctx.skills.registerProvider((control) => {
|
||||
provider = new SkillLocal.LocalSkillProvider(ctx, control, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
watch: true,
|
||||
watchPollIntervalMs: 10,
|
||||
watchStabilityThresholdMs: 20,
|
||||
})
|
||||
return provider
|
||||
})
|
||||
ctx.skills.registerProvider(provider)
|
||||
|
||||
const discovery = provider.list({})
|
||||
await settle()
|
||||
@@ -216,5 +219,6 @@ describe('skill-local watcher failures', () => {
|
||||
|
||||
await expect(discovery).rejects.toThrow('opening failed during disposal')
|
||||
await disposal
|
||||
disposeProvider()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -572,12 +572,8 @@ describe('LocalSkillProvider', () => {
|
||||
const root = join(home, '.agents/skills')
|
||||
const ctx = await setupLocal(home)
|
||||
expect(await ctx.skills.list()).toEqual([])
|
||||
const invalidateProvider = ctx.skills.invalidateProvider.bind(ctx.skills)
|
||||
let invalidations = 0
|
||||
ctx.skills.invalidateProvider = (provider) => {
|
||||
invalidations += 1
|
||||
invalidateProvider(provider)
|
||||
}
|
||||
ctx.on('skills/change', () => { invalidations += 1 })
|
||||
|
||||
await writeSkill(root, 'observed-skill', 'Observed skill')
|
||||
const path = join(root, 'observed-skill/SKILL.md')
|
||||
@@ -657,15 +653,18 @@ describe('LocalSkillProvider', () => {
|
||||
await writeSkill(join(home, '.agents/skills'), 'disposed-skill', 'Disposed skill')
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SkillService)
|
||||
const provider = new SkillLocal.LocalSkillProvider(ctx, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
customSkillDirs: [nonDirectoryRoot],
|
||||
watch: true,
|
||||
watchStabilityThresholdMs: 20,
|
||||
watchPollIntervalMs: 10,
|
||||
let provider!: SkillLocal.LocalSkillProvider
|
||||
const disposeProvider = ctx.skills.registerProvider((control) => {
|
||||
provider = new SkillLocal.LocalSkillProvider(ctx, control, {
|
||||
dshHome: join(home, '.dsh'),
|
||||
agentsHome: join(home, '.agents'),
|
||||
customSkillDirs: [nonDirectoryRoot],
|
||||
watch: true,
|
||||
watchStabilityThresholdMs: 20,
|
||||
watchPollIntervalMs: 10,
|
||||
})
|
||||
return provider
|
||||
})
|
||||
ctx.skills.registerProvider(provider)
|
||||
expect((await provider.list({})).map(skill => skill.name)).toEqual(['disposed-skill'])
|
||||
|
||||
await provider.dispose()
|
||||
@@ -673,6 +672,7 @@ describe('LocalSkillProvider', () => {
|
||||
provider.observeHostMutation(join(home, '.agents/skills/disposed-skill/SKILL.md'))
|
||||
|
||||
expect((await provider.list({})).map(skill => skill.name)).toEqual(['disposed-skill'])
|
||||
disposeProvider()
|
||||
})
|
||||
|
||||
it('refreshes frontmatter through a followed skill symlink', { timeout: 10000 }, async () => {
|
||||
@@ -740,7 +740,10 @@ describe('LocalSkillProvider', () => {
|
||||
expect(await empty.skills.list()).toEqual([])
|
||||
|
||||
delete process.env.DSH_AGENTS_HOME
|
||||
expect(new SkillLocal.LocalSkillProvider(empty, { dshHome: join(envHome, 'empty-dsh') }).name).toBe('local')
|
||||
expect(new SkillLocal.LocalSkillProvider(empty, {
|
||||
signal: new AbortController().signal,
|
||||
invalidate() {},
|
||||
}, { dshHome: join(envHome, 'empty-dsh') }).name).toBe('local')
|
||||
} finally {
|
||||
if (previousDshHome === undefined) {
|
||||
delete process.env.DSH_HOME
|
||||
|
||||
Reference in New Issue
Block a user