Browse Source

fix(core): reload changed skill sources (#40954)

Kit Langton 1 week ago
parent
commit
c79ced174e
2 changed files with 224 additions and 30 deletions
  1. 60 20
      packages/core/src/skill.ts
  2. 164 10
      packages/core/test/skill.test.ts

+ 60 - 20
packages/core/src/skill.ts

@@ -2,7 +2,7 @@ export * as Skill from "./skill"
 
 
 import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
 import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
 import path from "path"
 import path from "path"
-import { Context, Effect, Layer, Schema, Stream, Types } from "effect"
+import { Context, Effect, Layer, Schema, Scope, Stream, Types } from "effect"
 import { FileSystem } from "@opencode-ai/schema/filesystem"
 import { FileSystem } from "@opencode-ai/schema/filesystem"
 import { Skill } from "@opencode-ai/schema/skill"
 import { Skill } from "@opencode-ai/schema/skill"
 import { Agent } from "./agent"
 import { Agent } from "./agent"
@@ -13,6 +13,7 @@ import { Permission } from "./permission"
 import { AbsolutePath } from "./schema"
 import { AbsolutePath } from "./schema"
 import { SkillDiscovery } from "./skill/discovery"
 import { SkillDiscovery } from "./skill/discovery"
 import { State } from "./state"
 import { State } from "./state"
+import { Watcher } from "./filesystem/watcher"
 
 
 export const DirectorySource = Skill.DirectorySource
 export const DirectorySource = Skill.DirectorySource
 export type DirectorySource = Skill.DirectorySource
 export type DirectorySource = Skill.DirectorySource
@@ -81,6 +82,51 @@ const layer = Layer.effect(
     const discovery = yield* SkillDiscovery.Service
     const discovery = yield* SkillDiscovery.Service
     const fs = yield* FSUtil.Service
     const fs = yield* FSUtil.Service
     const bus = yield* Bus.Service
     const bus = yield* Bus.Service
+    const watcher = yield* Watcher.Service
+    const scope = yield* Scope.Scope
+    const cache = new Map<string, { skills: Info[]; paths: readonly string[] }>()
+    const watched = new Set<string>()
+
+    const invalidate = Effect.fn("Skill.invalidateFromWatcher")(function* (file: string) {
+      const invalidated = Array.from(cache.entries()).filter(([, loaded]) =>
+        loaded.paths.some((item) => FSUtil.overlaps(item, file)),
+      )
+      if (invalidated.length === 0) return
+      for (const [key] of invalidated) cache.delete(key)
+      yield* Effect.logInfo("skill cache invalidated", {
+        file,
+        sources: invalidated.map(([key]) => key),
+        skills: invalidated.flatMap(([, loaded]) => loaded.skills.map((skill) => skill.id)),
+      })
+      yield* bus.publish(Skill.Event.Updated, {}).pipe(Effect.asVoid)
+    })
+
+    const watch = Effect.fn("Skill.watch")(function* (directory: string) {
+      const target = path.resolve(directory)
+      if (watched.has(target)) return
+      watched.add(target)
+      const updates = yield* watcher.subscribe({ path: target, type: "directory" })
+      yield* updates.pipe(
+        Stream.runForEach((update) => invalidate(update.path)),
+        Effect.forkIn(scope, { startImmediately: true }),
+      )
+    })
+
+    const watchDirectory = Effect.fn("Skill.watchDirectory")(function* (directory: string) {
+      const target = path.resolve(directory)
+      const resolved = yield* fs.realPath(directory).pipe(Effect.catch(() => Effect.succeed(undefined)))
+      if (resolved) {
+        yield* watch(resolved)
+        if (resolved !== target) {
+          yield* watch(path.dirname(target))
+        }
+        return resolved === target ? [target] : [target, resolved]
+      }
+      if (yield* fs.isDir(path.dirname(target))) {
+        yield* watch(path.dirname(target))
+      }
+      return [target]
+    })
 
 
     const state = State.create<Data, Draft>({
     const state = State.create<Data, Draft>({
       name: "skill",
       name: "skill",
@@ -92,7 +138,8 @@ const layer = Layer.effect(
         },
         },
         list: () => draft.sources as Source[],
         list: () => draft.sources as Source[],
       }),
       }),
-      finalize: () => bus.publish(Skill.Event.Updated, {}).pipe(Effect.asVoid),
+      finalize: () =>
+        Effect.sync(() => cache.clear()).pipe(Effect.andThen(bus.publish(Skill.Event.Updated, {})), Effect.asVoid),
     })
     })
 
 
     const load = Effect.fn("Skill.load")(function* (source: Source) {
     const load = Effect.fn("Skill.load")(function* (source: Source) {
@@ -104,14 +151,22 @@ const layer = Layer.effect(
           directories: [],
           directories: [],
           skills: [source.skill.id],
           skills: [source.skill.id],
         })
         })
-        return { skills: [source.skill], directories: [] }
+        return { skills: [source.skill], paths: [] }
       }
       }
       const directories = source.type === "directory" ? [source.path] : yield* discovery.pull(source.url)
       const directories = source.type === "directory" ? [source.path] : yield* discovery.pull(source.url)
+      const roots = (yield* Effect.forEach(directories, watchDirectory)).flat()
+      const paths = [...roots]
       for (const directory of directories) {
       for (const directory of directories) {
         const files = yield* fs
         const files = yield* fs
           .scan("{*.md,**/SKILL.md}", { cwd: directory, absolute: true, include: "file", symlink: true, dot: true })
           .scan("{*.md,**/SKILL.md}", { cwd: directory, absolute: true, include: "file", symlink: true, dot: true })
           .pipe(Effect.catch(() => Effect.succeed([] as string[])))
           .pipe(Effect.catch(() => Effect.succeed([] as string[])))
         for (const filepath of files.toSorted()) {
         for (const filepath of files.toSorted()) {
+          const resolved = yield* fs.realPath(filepath).pipe(Effect.catch(() => Effect.succeed(filepath)))
+          if (!roots.some((root) => FSUtil.contains(root, resolved))) {
+            const external = path.dirname(resolved)
+            paths.push(external)
+            yield* watch(external)
+          }
           const content = yield* fs.readFileStringSafe(filepath).pipe(Effect.catch(() => Effect.succeed(undefined)))
           const content = yield* fs.readFileStringSafe(filepath).pipe(Effect.catch(() => Effect.succeed(undefined)))
           if (!content) continue
           if (!content) continue
           const markdown = ConfigMarkdown.parseOption(content)
           const markdown = ConfigMarkdown.parseOption(content)
@@ -139,22 +194,7 @@ const layer = Layer.effect(
         directories,
         directories,
         skills: skills.map((skill) => skill.id),
         skills: skills.map((skill) => skill.id),
       })
       })
-      return { skills, directories }
-    })
-
-    const cache = new Map<string, { skills: Info[]; directories: readonly string[] }>()
-    const invalidate = Effect.fn("Skill.invalidateFromWatcher")(function* (file: string) {
-      const invalidated = Array.from(cache.entries()).filter(([, loaded]) =>
-        loaded.directories.some((directory) => FSUtil.contains(directory, file)),
-      )
-      if (invalidated.length === 0) return
-      for (const [key] of invalidated) cache.delete(key)
-      yield* Effect.logInfo("skill cache invalidated", {
-        file,
-        sources: invalidated.map(([key]) => key),
-        skills: invalidated.flatMap(([, loaded]) => loaded.skills.map((skill) => skill.id)),
-      })
-      yield* bus.publish(Skill.Event.Updated, {}).pipe(Effect.asVoid)
+      return { skills, paths }
     })
     })
 
 
     yield* bus.subscribe(FileSystem.Event.Changed).pipe(
     yield* bus.subscribe(FileSystem.Event.Changed).pipe(
@@ -187,5 +227,5 @@ const layer = Layer.effect(
 export const node = makeLocationNode({
 export const node = makeLocationNode({
   service: Service,
   service: Service,
   layer,
   layer,
-  deps: [SkillDiscovery.node, FSUtil.node, Bus.node],
+  deps: [SkillDiscovery.node, FSUtil.node, Bus.node, Watcher.node],
 })
 })

+ 164 - 10
packages/core/test/skill.test.ts

@@ -6,11 +6,11 @@ import { Agent } from "@opencode-ai/core/agent"
 import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
 import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
 import { LayerNode } from "@opencode-ai/util/effect/layer-node"
 import { LayerNode } from "@opencode-ai/util/effect/layer-node"
 import { Bus } from "@opencode-ai/core/bus"
 import { Bus } from "@opencode-ai/core/bus"
-import { FSUtil } from "@opencode-ai/util/fs-util"
 import { AbsolutePath } from "@opencode-ai/core/schema"
 import { AbsolutePath } from "@opencode-ai/core/schema"
 import { Skill } from "@opencode-ai/core/skill"
 import { Skill } from "@opencode-ai/core/skill"
 import { SkillDiscovery } from "@opencode-ai/core/skill/discovery"
 import { SkillDiscovery } from "@opencode-ai/core/skill/discovery"
 import { FileSystem } from "@opencode-ai/schema/filesystem"
 import { FileSystem } from "@opencode-ai/schema/filesystem"
+import { Watcher } from "@opencode-ai/core/filesystem/watcher"
 import { tmpdir } from "./fixture/tmpdir"
 import { tmpdir } from "./fixture/tmpdir"
 import { testEffect } from "./lib/effect"
 import { testEffect } from "./lib/effect"
 
 
@@ -25,8 +25,15 @@ const discovery = Layer.succeed(
     },
     },
   }),
   }),
 )
 )
+const watcherLayer = Watcher.testLayer
 const it = testEffect(
 const it = testEffect(
-  AppNodeBuilder.build(LayerNode.group([Skill.node, Agent.node, Bus.node]), [[SkillDiscovery.node, discovery]]),
+  Layer.mergeAll(
+    AppNodeBuilder.build(LayerNode.group([Skill.node, Agent.node, Bus.node]), [
+      [SkillDiscovery.node, discovery],
+      [Watcher.node, watcherLayer],
+    ]),
+    watcherLayer,
+  ),
 )
 )
 
 
 function write(directory: string, name: string, description: string) {
 function write(directory: string, name: string, description: string) {
@@ -53,6 +60,24 @@ function waitForSkillUpdate() {
   })
   })
 }
 }
 
 
+function expectSubscription(check: (input: Watcher.WatchInput) => boolean) {
+  return Effect.gen(function* () {
+    const watcher = yield* Watcher.Test
+    expect((yield* watcher.subscriptions()).some(check)).toBe(true)
+  })
+}
+
+function emitAndWait(update: Watcher.Update) {
+  return Effect.gen(function* () {
+    const watcher = yield* Watcher.Test
+    yield* Effect.acquireUseRelease(
+      waitForSkillUpdate(),
+      ({ deferred }) => watcher.emit(update).pipe(Effect.andThen(Deferred.await(deferred)), Effect.timeout("1 second")),
+      ({ fiber }) => Fiber.interrupt(fiber),
+    )
+  })
+}
+
 describe("Skill", () => {
 describe("Skill", () => {
   it.live("publishes updates when skill sources change", () =>
   it.live("publishes updates when skill sources change", () =>
     Effect.gen(function* () {
     Effect.gen(function* () {
@@ -198,7 +223,7 @@ metadata:
     ),
     ),
   )
   )
 
 
-  it.live("invalidates cached skills and publishes updates for watcher changes", () =>
+  it.live("clears cached skills when sources reload", () =>
     Effect.acquireRelease(
     Effect.acquireRelease(
       Effect.promise(() => tmpdir()),
       Effect.promise(() => tmpdir()),
       (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
       (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
@@ -210,26 +235,155 @@ metadata:
             await write(tmp.path, "deploy", "Initial deploy")
             await write(tmp.path, "deploy", "Initial deploy")
           })
           })
 
 
-          const bus = yield* Bus.Service
           const skill = yield* Skill.Service
           const skill = yield* Skill.Service
           yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(tmp.path) }))
           yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(tmp.path) }))
+          expect((yield* skill.list()).find((item) => item.id === "deploy")?.description).toBe("Initial deploy")
 
 
-          expect((yield* skill.list()).find((item) => item.name === "deploy")?.description).toBe("Initial deploy")
-
-          const file = path.join(tmp.path, "deploy", "SKILL.md")
           yield* Effect.promise(() => write(tmp.path, "deploy", "Updated deploy"))
           yield* Effect.promise(() => write(tmp.path, "deploy", "Updated deploy"))
-          expect((yield* skill.list()).find((item) => item.name === "deploy")?.description).toBe("Initial deploy")
+          yield* skill.reload()
 
 
+          expect((yield* skill.list()).find((item) => item.id === "deploy")?.description).toBe("Updated deploy")
+        }),
+      ),
+    ),
+  )
+
+  it.live("reloads project sources created after their missing parent", () =>
+    Effect.acquireRelease(
+      Effect.promise(() => tmpdir()),
+      (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
+    ).pipe(
+      Effect.flatMap((tmp) =>
+        Effect.gen(function* () {
+          const source = path.join(tmp.path, "generated", "skills")
+          const file = path.join(source, "deploy", "SKILL.md")
+          const skill = yield* Skill.Service
+          const bus = yield* Bus.Service
+          yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(source) }))
+          expect(yield* skill.list()).toEqual([])
+
+          yield* Effect.promise(async () => {
+            await fs.mkdir(path.dirname(file), { recursive: true })
+            await write(source, "deploy", "Deploy production")
+          })
           yield* Effect.acquireUseRelease(
           yield* Effect.acquireUseRelease(
             waitForSkillUpdate(),
             waitForSkillUpdate(),
             ({ deferred }) =>
             ({ deferred }) =>
               bus
               bus
-                .publish(FileSystem.Event.Changed, { file, event: "change" })
+                .publish(FileSystem.Event.Changed, { file, event: "add" })
                 .pipe(Effect.andThen(Deferred.await(deferred)), Effect.timeout("1 second")),
                 .pipe(Effect.andThen(Deferred.await(deferred)), Effect.timeout("1 second")),
             ({ fiber }) => Fiber.interrupt(fiber),
             ({ fiber }) => Fiber.interrupt(fiber),
           )
           )
 
 
-          expect((yield* skill.list()).find((item) => item.name === "deploy")?.description).toBe("Updated deploy")
+          expect((yield* skill.list()).map((item) => item.id)).toEqual([Skill.ID.make("deploy")])
+        }),
+      ),
+    ),
+  )
+
+  it.live("watches directory sources for added and changed skills", () =>
+    Effect.acquireRelease(
+      Effect.promise(() => tmpdir()),
+      (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
+    ).pipe(
+      Effect.flatMap((tmp) =>
+        Effect.gen(function* () {
+          yield* Effect.promise(async () => {
+            await fs.mkdir(path.join(tmp.path, "deploy"), { recursive: true })
+            await write(tmp.path, "deploy", "Initial deploy")
+          })
+
+          const skill = yield* Skill.Service
+          yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(tmp.path) }))
+          expect((yield* skill.list()).map((item) => item.id)).toEqual([Skill.ID.make("deploy")])
+          yield* expectSubscription((input) => input.type === "directory" && input.path === tmp.path)
+
+          const deploy = path.join(tmp.path, "deploy", "SKILL.md")
+          yield* Effect.promise(() => write(tmp.path, "deploy", "Updated deploy"))
+          yield* emitAndWait({ type: "update", path: deploy })
+          expect((yield* skill.list()).find((item) => item.id === "deploy")?.description).toBe("Updated deploy")
+
+          yield* Effect.promise(async () => {
+            await fs.mkdir(path.join(tmp.path, "review"), { recursive: true })
+            await write(tmp.path, "review", "Review changes")
+          })
+          const review = path.join(tmp.path, "review", "SKILL.md")
+          yield* emitAndWait({ type: "create", path: review })
+          expect((yield* skill.list()).map((item) => item.id)).toEqual([
+            Skill.ID.make("deploy"),
+            Skill.ID.make("review"),
+          ])
+
+          yield* Effect.promise(() => fs.rm(path.join(tmp.path, "review"), { recursive: true }))
+          yield* emitAndWait({ type: "delete", path: review })
+          expect((yield* skill.list()).map((item) => item.id)).toEqual([Skill.ID.make("deploy")])
+        }),
+      ),
+    ),
+  )
+
+  it.live("watches canonical directories behind symlinked skills", () =>
+    Effect.acquireRelease(
+      Effect.promise(() => tmpdir()),
+      (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
+    ).pipe(
+      Effect.flatMap((tmp) =>
+        Effect.gen(function* () {
+          const source = path.join(tmp.path, "source")
+          const target = path.join(tmp.path, "target", "bro")
+          const file = path.join(target, "SKILL.md")
+          yield* Effect.promise(async () => {
+            await fs.mkdir(source, { recursive: true })
+            await fs.mkdir(target, { recursive: true })
+            await fs.writeFile(file, "---\nname: bro\ndescription: Initial\n---\n# bro")
+            await fs.symlink(target, path.join(source, "bro"))
+          })
+
+          const skill = yield* Skill.Service
+          yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(source) }))
+          expect((yield* skill.list()).find((item) => item.id === "bro")?.description).toBe("Initial")
+          yield* expectSubscription((input) => input.type === "directory" && input.path === target)
+
+          yield* Effect.promise(() => fs.writeFile(file, "---\nname: bro\ndescription: Updated\n---\n# bro"))
+          yield* emitAndWait({ type: "update", path: file })
+          expect((yield* skill.list()).find((item) => item.id === "bro")?.description).toBe("Updated")
+        }),
+      ),
+    ),
+  )
+
+  it.live("invalidates symlinked sources when their target changes", () =>
+    Effect.acquireRelease(
+      Effect.promise(() => tmpdir()),
+      (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
+    ).pipe(
+      Effect.flatMap((tmp) =>
+        Effect.gen(function* () {
+          const source = path.join(tmp.path, "source")
+          const first = path.join(tmp.path, "first")
+          const second = path.join(tmp.path, "second")
+          yield* Effect.promise(async () => {
+            await fs.mkdir(path.join(first, "bro"), { recursive: true })
+            await fs.mkdir(path.join(second, "bro"), { recursive: true })
+            await write(first, "bro", "First")
+            await write(second, "bro", "Second")
+            await fs.symlink(first, source)
+          })
+
+          const skill = yield* Skill.Service
+          yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(source) }))
+          expect((yield* skill.list()).find((item) => item.id === "bro")?.description).toBe("First")
+          yield* expectSubscription((input) => input.type === "directory" && input.path === first)
+          yield* expectSubscription((input) => input.type === "directory" && input.path === tmp.path)
+
+          yield* Effect.promise(async () => {
+            await fs.unlink(source)
+            await fs.symlink(second, source)
+          })
+          yield* emitAndWait({ type: "update", path: source })
+
+          expect((yield* skill.list()).find((item) => item.id === "bro")?.description).toBe("Second")
+          yield* expectSubscription((input) => input.type === "directory" && input.path === second)
         }),
         }),
       ),
       ),
     ),
     ),