Просмотр исходного кода

fix(core): bound project filesystem watches (#41096)

Kit Langton 1 неделя назад
Родитель
Сommit
e5ef00b8b8

+ 12 - 13
packages/app/src/pages/session.tsx

@@ -688,6 +688,8 @@ export default function Page() {
     return {
       queryKey: [...vcsKey(), mode] as const,
       enabled,
+      refetchOnMount: "always" as const,
+      refetchOnWindowFocus: true,
       queryFn: mode
         ? () =>
             sdk()
@@ -701,6 +703,16 @@ export default function Page() {
     }
   })
   const refreshVcs = debounce(() => void queryClient.invalidateQueries({ queryKey: vcsKey() }), 100)
+  createEffect(
+    on(
+      () => desktopReviewOpen() || mobileChanges(),
+      (open, previous) => {
+        if (!open || previous || !desktopFileTreeOpen() || vcsQuery.isFetching) return
+        refreshVcs()
+      },
+      { defer: true },
+    ),
+  )
   const reviewDiffs = () => {
     if (reviewMode() === "git" || reviewMode() === "branch")
       // avoids suspense
@@ -947,19 +959,6 @@ export default function Page() {
     ),
   )
 
-  const stopVcs = sdk().event.listen((evt) => {
-    const details = evt.details as { type: string; properties?: unknown }
-    if (details.type !== "file.watcher.updated" && details.type !== "filesystem.changed") return
-    const props =
-      typeof details.properties === "object" && details.properties
-        ? (details.properties as Record<string, unknown>)
-        : undefined
-    const file = typeof props?.file === "string" ? props.file : undefined
-    if (!file || file.startsWith(".git/")) return
-    refreshVcs()
-  })
-  onCleanup(stopVcs)
-
   createEffect(
     on(
       () => sdk().directory,

+ 1 - 26
packages/core/src/filesystem/location-watcher.ts

@@ -11,15 +11,6 @@ import { FSUtil } from "@opencode-ai/util/fs-util"
 import { Git } from "../git"
 import { Location } from "../location"
 import { Watcher } from "./watcher"
-import { Ignore } from "./ignore"
-import { Protected } from "./protected"
-
-function protecteds(dir: string) {
-  return Protected.paths().filter((item) => {
-    const relative = path.relative(dir, item)
-    return relative !== "" && !relative.startsWith("..") && !path.isAbsolute(relative)
-  })
-}
 
 export interface Interface {}
 
@@ -44,19 +35,6 @@ const layer = Layer.effect(
       const config = (yield* configService.entries())
         .filter((entry): entry is Document => entry.type === "document")
         .flatMap((item) => item.info.watcher?.ignore ?? [])
-      const home = Protected.isHome(location.directory)
-
-      if (!home && location.vcs) {
-        const updates = yield* watcher.subscribe({
-          path: location.directory,
-          type: "directory",
-          ignore: [...Ignore.PATTERNS, ...config, ...protecteds(location.directory)],
-        })
-        yield* updates.pipe(Stream.runForEach(publish), Effect.forkScoped)
-      }
-      if (home) {
-        yield* Effect.logInfo("location watcher skipped home directory", { directory: location.directory })
-      }
 
       if (location.vcs?.type === "git") {
         const resolved = (yield* git.repo.discover(location.directory))?.gitDirectory
@@ -64,10 +42,7 @@ const layer = Layer.effect(
           ? yield* fs.realPath(resolved).pipe(Effect.catch(() => Effect.succeed(resolved)))
           : undefined
         if (vcs && !config.includes(".git") && !config.includes(vcs) && (!resolved || !config.includes(resolved))) {
-          const ignore = (yield* fs.readDirectoryEntries(vcs).pipe(Effect.catch(() => Effect.succeed([])))).flatMap(
-            (entry) => (entry.name === "HEAD" ? [] : [entry.name]),
-          )
-          const updates = yield* watcher.subscribe({ path: vcs, type: "directory", ignore })
+          const updates = yield* watcher.subscribe({ path: path.join(vcs, "HEAD"), type: "file" })
           yield* updates.pipe(Stream.runForEach(publish), Effect.forkScoped)
         }
       }

+ 71 - 40
packages/core/src/skill.ts

@@ -2,8 +2,7 @@ export * as Skill from "./skill"
 
 import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
 import path from "path"
-import { Context, Effect, Layer, Schema, Scope, Stream, Types } from "effect"
-import { FileSystem } from "@opencode-ai/schema/filesystem"
+import { Context, Effect, FiberMap, Layer, PubSub, Schema, Semaphore, Stream, Types } from "effect"
 import { Skill } from "@opencode-ai/schema/skill"
 import { Agent } from "./agent"
 import { ConfigMarkdown } from "./config/markdown"
@@ -83,47 +82,78 @@ const layer = Layer.effect(
     const fs = yield* FSUtil.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 watches = yield* FiberMap.make<string>()
+    const lock = Semaphore.makeUnsafe(1)
+    const changes = yield* PubSub.unbounded<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)),
+      const changed = yield* lock.withPermit(
+        Effect.gen(function* () {
+          const invalidated = Array.from(cache.entries()).filter(([, loaded]) =>
+            loaded.paths.some((item) => FSUtil.overlaps(item, file)),
+          )
+          if (invalidated.length === 0) return false
+          cache.clear()
+          yield* FiberMap.clear(watches)
+          yield* Effect.logInfo("skill cache invalidated", {
+            file,
+            sources: invalidated.map(([key]) => key),
+            skills: invalidated.flatMap(([, loaded]) => loaded.skills.map((skill) => skill.id)),
+          })
+          return true
+        }),
       )
-      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)),
-      })
+      if (!changed) return
       yield* bus.publish(Skill.Event.Updated, {}).pipe(Effect.asVoid)
     })
 
-    const watch = Effect.fn("Skill.watch")(function* (directory: string) {
+    yield* Stream.fromPubSub(changes).pipe(Stream.runForEach(invalidate), Effect.forkScoped({ startImmediately: true }))
+
+    const watch = Effect.fn("Skill.watch")(function* (directory: string, type: Watcher.WatchInput["type"]) {
       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 updates = yield* watcher.subscribe(
+        type === "file" ? { path: target, type: "file" } : { path: target, type: "directory" },
+      )
+      yield* FiberMap.run(
+        watches,
+        `${type}:${target}`,
+        updates.pipe(Stream.runForEach((update) => PubSub.publish(changes, update.path).pipe(Effect.asVoid))),
+        {
+          onlyIfMissing: true,
+          startImmediately: true,
+        },
       )
     })
 
-    const watchDirectory = Effect.fn("Skill.watchDirectory")(function* (directory: string) {
+    function firstMissing(target: string): Effect.Effect<string | undefined> {
+      const parent = path.dirname(target)
+      if (parent === target) return Effect.succeed(undefined)
+      return fs.isDir(parent).pipe(Effect.flatMap((exists) => (exists ? Effect.succeed(target) : firstMissing(parent))))
+    }
+
+    const watchDirectory: (directory: string) => Effect.Effect<string[]> = 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)
+        yield* watch(resolved, "directory")
         if (resolved !== target) {
-          yield* watch(path.dirname(target))
+          yield* watch(target, "file")
         }
         return resolved === target ? [target] : [target, resolved]
       }
-      if (yield* fs.isDir(path.dirname(target))) {
-        yield* watch(path.dirname(target))
+      const missing = yield* firstMissing(target)
+      if (missing) yield* watch(missing, "file")
+      if (
+        yield* fs.realPath(directory).pipe(
+          Effect.as(true),
+          Effect.catch(() => Effect.succeed(false)),
+        )
+      ) {
+        if (missing) yield* FiberMap.remove(watches, `file:${path.resolve(missing)}`)
+        return yield* watchDirectory(directory)
       }
       return [target]
     })
@@ -139,7 +169,9 @@ const layer = Layer.effect(
         list: () => draft.sources as Source[],
       }),
       finalize: () =>
-        Effect.sync(() => cache.clear()).pipe(Effect.andThen(bus.publish(Skill.Event.Updated, {})), Effect.asVoid),
+        lock
+          .withPermit(FiberMap.clear(watches).pipe(Effect.andThen(Effect.sync(() => cache.clear())), Effect.asVoid))
+          .pipe(Effect.andThen(bus.publish(Skill.Event.Updated, {})), Effect.asVoid),
     })
 
     const load = Effect.fn("Skill.load")(function* (source: Source) {
@@ -165,7 +197,7 @@ const layer = Layer.effect(
           if (!roots.some((root) => FSUtil.contains(root, resolved))) {
             const external = path.dirname(resolved)
             paths.push(external)
-            yield* watch(external)
+            yield* watch(external, "directory")
           }
           const content = yield* fs.readFileStringSafe(filepath).pipe(Effect.catch(() => Effect.succeed(undefined)))
           if (!content) continue
@@ -197,20 +229,19 @@ const layer = Layer.effect(
       return { skills, paths }
     })
 
-    yield* bus.subscribe(FileSystem.Event.Changed).pipe(
-      Stream.runForEach((event) => invalidate(event.data.file)),
-      Effect.forkScoped({ startImmediately: true }),
-    )
-
     const list = Effect.fn("Skill.list")(function* () {
-      const skills = new Map<ID, Info>()
-      for (const source of state.get().sources) {
-        const key = Source.key(source)
-        const loaded = cache.get(key) ?? (yield* load(source))
-        cache.set(key, loaded)
-        for (const skill of loaded.skills) skills.set(skill.id, skill)
-      }
-      return Array.from(skills.values())
+      return yield* lock.withPermit(
+        Effect.gen(function* () {
+          const skills = new Map<ID, Info>()
+          for (const source of state.get().sources) {
+            const key = Source.key(source)
+            const loaded = cache.get(key) ?? (yield* load(source))
+            cache.set(key, loaded)
+            for (const skill of loaded.skills) skills.set(skill.id, skill)
+          }
+          return Array.from(skills.values())
+        }),
+      )
     })
 
     return Service.of({

+ 93 - 126
packages/core/test/filesystem/watcher.test.ts

@@ -17,9 +17,8 @@ import { location } from "../fixture/location"
 import { tmpdir } from "../fixture/tmpdir"
 import { testEffect } from "../lib/effect"
 
-const describeWatcher = Watcher.hasNativeBinding() && !process.env.CI ? describe : describe.skip
-
 type WatcherEvent = { file: string; event: "add" | "change" | "unlink" }
+const describeNative = process.env.CI ? describe.skip : describe
 
 const it = testEffect(AppNodeBuilder.build(LayerNode.group([FSUtil.node, Bus.node])))
 
@@ -75,10 +74,9 @@ describe("Watcher lifecycle", () => {
       const interrupted = yield* Deferred.make<void>()
       yield* Effect.gen(function* () {
         const watcher = yield* Watcher.Service
-        const consumer = yield* watcher.subscribe({ path: "/pending", type: "directory" }).pipe(
-          Effect.flatMap(Stream.runDrain),
-          Effect.forkScoped({ startImmediately: true }),
-        )
+        const consumer = yield* watcher
+          .subscribe({ path: "/pending", type: "directory" })
+          .pipe(Effect.flatMap(Stream.runDrain), Effect.forkScoped({ startImmediately: true }))
         yield* Deferred.await(started)
         yield* Fiber.interrupt(consumer)
         expect(yield* Deferred.isDone(interrupted)).toBe(true)
@@ -99,10 +97,9 @@ describe("Watcher lifecycle", () => {
     return Effect.gen(function* () {
       const watcher = yield* Watcher.Service
       const consume = () =>
-        watcher.subscribe({ path: "/shared", type: "directory" }).pipe(
-          Effect.flatMap(Stream.runDrain),
-          Effect.forkScoped({ startImmediately: true }),
-        )
+        watcher
+          .subscribe({ path: "/shared", type: "directory" })
+          .pipe(Effect.flatMap(Stream.runDrain), Effect.forkScoped({ startImmediately: true }))
       const first = yield* consume()
       const second = yield* consume()
       yield* Effect.yieldNow
@@ -138,22 +135,26 @@ describe("Watcher lifecycle", () => {
   })
 })
 
-function provide(directory: string, vcs?: Location.Interface["vcs"]) {
+function provide(directory: string, vcs?: Location.Interface["vcs"], watcher?: Layer.Layer<Watcher.Service>) {
   const locationLayer = Layer.succeed(
     Location.Service,
     Location.Service.of(location({ directory: AbsolutePath.make(directory) }, { vcs })),
   )
-  return Effect.provide(
-    AppNodeBuilder.build(LocationWatcher.node, [
-      [Config.node, configLayer],
-      [Location.node, locationLayer],
-    ]),
-  )
+  const built = AppNodeBuilder.build(LocationWatcher.node, [
+    [Config.node, configLayer],
+    [Location.node, locationLayer],
+    ...(watcher ? ([[Watcher.node, watcher]] as const) : []),
+  ])
+  return Effect.provide(built)
 }
 
 function withTmp<A, E, R>(
   f: (directory: string, vcs?: Location.Interface["vcs"]) => Effect.Effect<A, E, R>,
-  options?: { vcs?: "git" | "hg"; init?: (directory: string) => Promise<void> },
+  options?: {
+    vcs?: "git" | "hg"
+    init?: (directory: string) => Promise<void>
+    watcher?: Layer.Layer<Watcher.Service>
+  },
 ) {
   return Effect.acquireRelease(
     Effect.promise(async () => {
@@ -173,9 +174,57 @@ function withTmp<A, E, R>(
       return { tmp, vcs: { type: "git" as const, store: AbsolutePath.make(path.join(tmp.path, ".git")) } }
     }),
     ({ tmp }) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
-  ).pipe(Effect.flatMap(({ tmp, vcs }) => f(tmp.path, vcs).pipe(provide(tmp.path, vcs))))
+  ).pipe(Effect.flatMap(({ tmp, vcs }) => f(tmp.path, vcs).pipe(provide(tmp.path, vcs, options?.watcher))))
 }
 
+describe("LocationWatcher subscriptions", () => {
+  it.live("watches only exact Git branch metadata", () => {
+    const subscriptions: Watcher.WatchInput[] = []
+    const watcher = Layer.succeed(
+      Watcher.Service,
+      Watcher.Service.of({
+        subscribe: (input) => Effect.sync(() => subscriptions.push(input)).pipe(Effect.as(Stream.empty)),
+      }),
+    )
+    return withTmp(
+      (directory) =>
+        Effect.gen(function* () {
+          yield* LocationWatcher.Service
+          yield* Effect.sync(() => subscriptions.length).pipe(
+            Effect.filterOrFail((count) => count > 0),
+            Effect.retry(Schedule.spaced("10 millis")),
+          )
+          yield* Effect.sleep("10 millis")
+          expect(subscriptions).toEqual([{ path: path.join(directory, ".git", "HEAD"), type: "file" }])
+        }),
+      { vcs: "git", watcher },
+    )
+  })
+
+  it.live("watches only exact Hg branch metadata", () => {
+    const subscriptions: Watcher.WatchInput[] = []
+    const watcher = Layer.succeed(
+      Watcher.Service,
+      Watcher.Service.of({
+        subscribe: (input) => Effect.sync(() => subscriptions.push(input)).pipe(Effect.as(Stream.empty)),
+      }),
+    )
+    return withTmp(
+      (directory) =>
+        Effect.gen(function* () {
+          yield* LocationWatcher.Service
+          yield* Effect.sync(() => subscriptions.length).pipe(
+            Effect.filterOrFail((count) => count > 0),
+            Effect.retry(Schedule.spaced("10 millis")),
+          )
+          yield* Effect.sleep("10 millis")
+          expect(subscriptions).toEqual([{ path: path.join(directory, ".hg", "branch"), type: "file" }])
+        }),
+      { vcs: "hg", watcher },
+    )
+  })
+})
+
 function wait(check: (event: WatcherEvent) => boolean) {
   return Effect.gen(function* () {
     const bus = yield* Bus.Service
@@ -226,31 +275,18 @@ function eventuallyUpdate<E>(check: (event: WatcherEvent) => boolean, trigger: (
   )
 }
 
-function noUpdate<E>(check: (event: WatcherEvent) => boolean, trigger: Effect.Effect<void, E>, timeout = 500) {
-  return Effect.acquireUseRelease(
-    wait(check),
-    ({ deferred }) =>
-      trigger.pipe(
-        Effect.andThen(Deferred.await(deferred)),
-        Effect.timeoutOption(`${timeout} millis`),
-        Effect.tap((result) => Effect.sync(() => expect(result).toEqual(Option.none()))),
-      ),
-    ({ fiber }) => Fiber.interrupt(fiber),
-  )
-}
-
-function ready(directory: string) {
-  const file = path.join(directory, `.watcher-${Math.random().toString(36).slice(2)}`)
+function ready(file: string, eventFile = file) {
   return Effect.gen(function* () {
     const fs = yield* FSUtil.Service
+    const content = (yield* fs.readFileStringSafe(file)) ?? `ready-${Math.random()}`
     yield* eventuallyUpdate(
-      (event) => event.file === file,
-      () => fs.writeFileString(file, `ready-${Math.random()}`),
-    ).pipe(Effect.ensuring(fs.remove(file, { force: true }).pipe(Effect.ignore)), Effect.asVoid)
+      (event) => event.file === eventFile,
+      () => fs.writeFileString(file, content),
+    ).pipe(Effect.asVoid)
   })
 }
 
-describeWatcher("LocationWatcher", () => {
+describeNative("LocationWatcher", () => {
   it.live("limits file watches to the exact target", () =>
     withTmp((directory) =>
       Effect.gen(function* () {
@@ -276,94 +312,25 @@ describeWatcher("LocationWatcher", () => {
     ),
   )
 
-  it.live("publishes root create, update, and delete events", () =>
-    withTmp(
-      (directory) =>
-        Effect.gen(function* () {
-          const fs = yield* FSUtil.Service
-          const file = path.join(directory, "watch.txt")
-          yield* ready(directory)
-          for (const item of [
-            { event: "add" as const, trigger: fs.writeFileString(file, "a") },
-            { event: "change" as const, trigger: fs.writeFileString(file, "b") },
-            { event: "unlink" as const, trigger: fs.remove(file) },
-          ]) {
-            expect(
-              yield* nextUpdate((event) => event.file === file && event.event === item.event, item.trigger),
-            ).toEqual({
-              file,
-              event: item.event,
-            })
-          }
-        }),
-      { vcs: "git" },
-    ),
-  )
-
-  it.live("skips non-git roots", () =>
+  it.live("detects creation of a missing directory target", () =>
     withTmp((directory) =>
       Effect.gen(function* () {
         const fs = yield* FSUtil.Service
-        const file = path.join(directory, "plain.txt")
-        yield* noUpdate((event) => event.file === file, fs.writeFileString(file, "plain"))
-      }),
-    ),
-  )
-
-  it.live("ignores dependency, VCS, and build directories at any depth", () =>
-    withTmp(
-      (directory) =>
-        Effect.gen(function* () {
-          const afs = yield* FSUtil.Service
-          yield* ready(directory)
-          const roots = ["node_modules", ".git", "dist"].map((name) => path.join(directory, "nested", name))
-          const files = roots.map((root) => path.join(root, "package", "index.js"))
-          yield* noUpdate(
-            (event) => roots.some((root) => event.file === root || event.file.startsWith(`${root}${path.sep}`)),
-            Effect.forEach(files, (file) => afs.writeWithDirs(file, "ignored"), {
-              concurrency: "unbounded",
-              discard: true,
-            }),
-          )
-        }),
-      { vcs: "git" },
-    ),
-  )
-
-  it.live("cleanup stops publishing events", () =>
-    Effect.gen(function* () {
-      const bus = yield* Bus.Service
-      const fs = yield* FSUtil.Service
-      const tmp = yield* Effect.acquireRelease(
-        Effect.promise(() => tmpdir()),
-        (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
-      )
-      yield* ready(tmp.path).pipe(
-        provide(tmp.path, { type: "git", store: AbsolutePath.make(path.join(tmp.path, ".git")) }),
-        Effect.scoped,
-      )
-      const file = path.join(tmp.path, "after-dispose.txt")
-      yield* noUpdate((event) => event.file === file, fs.writeFileString(file, "gone")).pipe(
-        Effect.provideService(Bus.Service, bus),
-      )
-    }).pipe(Effect.provide(AppNodeBuilder.build(LayerNode.group([FSUtil.node, Bus.node])))),
-  )
+        const watcher = yield* Watcher.Service
+        const target = path.join(directory, "generated")
+        const updates = yield* watcher.subscribe({ path: target, type: "file" })
+        const update = yield* updates.pipe(
+          Stream.take(1),
+          Stream.runHead,
+          Effect.forkScoped({ startImmediately: true }),
+        )
+        const creates = yield* Effect.suspend(() =>
+          fs.remove(target, { recursive: true, force: true }).pipe(Effect.andThen(fs.ensureDir(target))),
+        ).pipe(Effect.repeat(Schedule.spaced("10 millis")), Effect.forkScoped)
+        const event = yield* Fiber.join(update).pipe(Effect.ensuring(Fiber.interrupt(creates)))
 
-  it.live("ignores .git/index changes", () =>
-    withTmp(
-      (directory) =>
-        Effect.gen(function* () {
-          const fs = yield* FSUtil.Service
-          const index = path.join(directory, ".git", "index")
-          yield* ready(directory)
-          yield* noUpdate(
-            (event) => event.file === index,
-            fs
-              .writeFileString(path.join(directory, "tracked.txt"), "a")
-              .pipe(Effect.andThen(Effect.promise(() => $`git add .`.cwd(directory).quiet())), Effect.asVoid),
-          )
-        }),
-      { vcs: "git" },
+        expect(event.valueOrUndefined?.path).toBe(target)
+      }).pipe(Effect.provide(AppNodeBuilder.build(Watcher.node))),
     ),
   )
 
@@ -374,11 +341,11 @@ describeWatcher("LocationWatcher", () => {
           const fs = yield* FSUtil.Service
           const head = path.join(directory, ".git", "HEAD")
           const branch = `watch-${Math.random().toString(36).slice(2)}`
-          yield* ready(directory)
+          yield* ready(head)
           yield* Effect.promise(() => $`git branch ${branch}`.cwd(directory).quiet())
           expect(
             yield* nextUpdate((event) => event.file === head, fs.writeFileString(head, `ref: refs/heads/${branch}\n`)),
-          ).toMatchObject({ file: head })
+          ).toEqual({ file: head, event: "change" })
         }),
       { vcs: "git" },
     ),
@@ -393,8 +360,8 @@ describeWatcher("LocationWatcher", () => {
             const afs = yield* FSUtil.Service
             const actual = path.join(directory, "..", `actual_${path.basename(directory)}`)
             yield* Effect.addFinalizer(() => Effect.promise(() => fs.rm(actual, { recursive: true, force: true })))
-            yield* ready(directory)
             const head = path.join(directory, ".git", "HEAD")
+            yield* ready(head, path.join(actual, "HEAD"))
             const branch = `watch-${Math.random().toString(36).slice(2)}`
             yield* Effect.promise(() => $`git branch ${branch}`.cwd(directory).quiet())
             expect(
@@ -422,7 +389,7 @@ describeWatcher("LocationWatcher", () => {
         Effect.gen(function* () {
           const fs = yield* FSUtil.Service
           const branch = path.join(directory, ".hg", "branch")
-          yield* ready(directory)
+          yield* ready(branch)
           expect(
             yield* nextUpdate((event) => event.file === branch, fs.writeFileString(branch, "feature\n")),
           ).toMatchObject({ file: branch })

+ 62 - 15
packages/core/test/skill.test.ts

@@ -9,7 +9,6 @@ import { Bus } from "@opencode-ai/core/bus"
 import { AbsolutePath } from "@opencode-ai/core/schema"
 import { Skill } from "@opencode-ai/core/skill"
 import { SkillDiscovery } from "@opencode-ai/core/skill/discovery"
-import { FileSystem } from "@opencode-ai/schema/filesystem"
 import { Watcher } from "@opencode-ai/core/filesystem/watcher"
 import { tmpdir } from "./fixture/tmpdir"
 import { testEffect } from "./lib/effect"
@@ -114,6 +113,7 @@ describe("Skill", () => {
           })
 
           const skill = yield* Skill.Service
+          const watcher = yield* Watcher.Test
           yield* skill.transform((editor) => {
             editor.source({ type: "directory", path: AbsolutePath.make(first) })
             editor.source({ type: "directory", path: AbsolutePath.make(first) })
@@ -144,6 +144,21 @@ describe("Skill", () => {
               content: "# review",
             },
           ])
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: first, type: "directory" },
+            { path: second, type: "directory" },
+          ])
+
+          yield* Effect.promise(() => write(second, "review", "Updated Second"))
+          yield* emitAndWait({ type: "update", path: path.join(second, "review", "SKILL.md") })
+
+          expect((yield* skill.list()).find((item) => item.id === "review")?.description).toBe("Updated Second")
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: first, type: "directory" },
+            { path: second, type: "directory" },
+            { path: first, type: "directory" },
+            { path: second, type: "directory" },
+          ])
         }),
       ),
     ),
@@ -236,13 +251,30 @@ metadata:
           })
 
           const skill = yield* Skill.Service
+          const watcher = yield* Watcher.Test
+          const bus = yield* Bus.Service
           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* watcher.subscriptions()).toEqual([{ path: tmp.path, type: "directory" }])
+
+          let refreshed: Skill.Info[] = []
+          const unsubscribe = yield* bus.listen((event) => {
+            if (event.type !== Skill.Event.Updated.type) return Effect.void
+            return skill.list().pipe(
+              Effect.tap((items) => Effect.sync(() => (refreshed = items))),
+              Effect.asVoid,
+            )
+          })
 
           yield* Effect.promise(() => write(tmp.path, "deploy", "Updated deploy"))
-          yield* skill.reload()
+          yield* skill.reload().pipe(Effect.timeout("1 second"))
+          yield* unsubscribe
 
-          expect((yield* skill.list()).find((item) => item.id === "deploy")?.description).toBe("Updated deploy")
+          expect(refreshed.find((item) => item.id === "deploy")?.description).toBe("Updated deploy")
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: tmp.path, type: "directory" },
+            { path: tmp.path, type: "directory" },
+          ])
         }),
       ),
     ),
@@ -258,24 +290,31 @@ metadata:
           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
+          const watcher = yield* Watcher.Test
           yield* skill.transform((editor) => editor.source({ type: "directory", path: AbsolutePath.make(source) }))
           expect(yield* skill.list()).toEqual([])
+          expect(yield* watcher.subscriptions()).toEqual([{ path: path.join(tmp.path, "generated"), type: "file" }])
+
+          yield* Effect.promise(() => fs.mkdir(path.join(tmp.path, "generated")))
+          yield* emitAndWait({ type: "create", path: path.join(tmp.path, "generated") })
+          expect(yield* skill.list()).toEqual([])
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: path.join(tmp.path, "generated"), type: "file" },
+            { path: source, type: "file" },
+          ])
 
           yield* Effect.promise(async () => {
             await fs.mkdir(path.dirname(file), { recursive: true })
             await write(source, "deploy", "Deploy production")
           })
-          yield* Effect.acquireUseRelease(
-            waitForSkillUpdate(),
-            ({ deferred }) =>
-              bus
-                .publish(FileSystem.Event.Changed, { file, event: "add" })
-                .pipe(Effect.andThen(Deferred.await(deferred)), Effect.timeout("1 second")),
-            ({ fiber }) => Fiber.interrupt(fiber),
-          )
+          yield* emitAndWait({ type: "create", path: source })
 
           expect((yield* skill.list()).map((item) => item.id)).toEqual([Skill.ID.make("deploy")])
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: path.join(tmp.path, "generated"), type: "file" },
+            { path: source, type: "file" },
+            { path: source, type: "directory" },
+          ])
         }),
       ),
     ),
@@ -371,10 +410,13 @@ metadata:
           })
 
           const skill = yield* Skill.Service
+          const watcher = yield* Watcher.Test
           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)
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: first, type: "directory" },
+            { path: source, type: "file" },
+          ])
 
           yield* Effect.promise(async () => {
             await fs.unlink(source)
@@ -383,7 +425,12 @@ metadata:
           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)
+          expect(yield* watcher.subscriptions()).toEqual([
+            { path: first, type: "directory" },
+            { path: source, type: "file" },
+            { path: second, type: "directory" },
+            { path: source, type: "file" },
+          ])
         }),
       ),
     ),