| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400 |
- import { $ } from "bun"
- import { describe, expect } from "bun:test"
- import fs from "fs/promises"
- import path from "path"
- import { Deferred, Duration, Effect, Fiber, Layer, Option, Schedule, Stream } from "effect"
- import { Config } from "@opencode-ai/core/config"
- import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
- import { LayerNode } from "@opencode-ai/util/effect/layer-node"
- import { Bus } from "@opencode-ai/core/bus"
- import { FSUtil } from "@opencode-ai/util/fs-util"
- import { LocationWatcher } from "@opencode-ai/core/filesystem/location-watcher"
- import { Watcher } from "@opencode-ai/core/filesystem/watcher"
- import { FileSystem } from "@opencode-ai/schema/filesystem"
- import { Location } from "@opencode-ai/core/location"
- import { AbsolutePath } from "@opencode-ai/core/schema"
- import { location } from "../fixture/location"
- import { tmpdir } from "../fixture/tmpdir"
- import { testEffect } from "../lib/effect"
- 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])))
- const configLayer = Config.testLayer()
- describe("Watcher.testLayer", () => {
- it.effect("records subscriptions and broadcasts emitted updates through the service", () =>
- Effect.gen(function* () {
- const watcher = yield* Watcher.Service
- const test = yield* Watcher.Test
- const updates = yield* watcher.subscribe({ path: "/root", type: "directory" })
- const received = yield* updates.pipe(
- Stream.take(1),
- Stream.runCollect,
- Effect.forkScoped({ startImmediately: true }),
- )
- yield* Effect.yieldNow
- yield* test.emit({ type: "update", path: "/root/file.md" })
- expect(Array.from(yield* Fiber.join(received))).toEqual([{ type: "update", path: "/root/file.md" }])
- // subscriptions() reports acquired watches, so paths come back resolved.
- expect(yield* test.subscriptions()).toEqual([{ path: path.resolve("/root"), type: "directory" }])
- }).pipe(Effect.provide(Watcher.testLayer)),
- )
- })
- function withNative(native: Watcher.NativeInterface) {
- return Effect.provide(Watcher.layer().pipe(Layer.provide(Layer.succeed(Watcher.Native, native))))
- }
- function countingNative() {
- const counts = { subscribes: 0, unsubscribes: 0 }
- const native: Watcher.NativeInterface = {
- subscribe: () =>
- Effect.sync(() => {
- counts.subscribes++
- return {
- unsubscribe: () => {
- counts.unsubscribes++
- return Promise.resolve()
- },
- }
- }),
- }
- return { native, counts }
- }
- describe("Watcher lifecycle", () => {
- it.effect("interrupting a consumer interrupts a pending acquisition", () =>
- Effect.gen(function* () {
- const started = yield* Deferred.make<void>()
- 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 }))
- yield* Deferred.await(started)
- yield* Fiber.interrupt(consumer)
- expect(yield* Deferred.isDone(interrupted)).toBe(true)
- }).pipe(
- withNative({
- subscribe: () =>
- Deferred.succeed(started, undefined).pipe(
- Effect.andThen(Effect.never),
- Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)),
- ),
- }),
- )
- }),
- )
- it.effect("shares one subscription and releases exactly once after the final consumer", () => {
- const { native, counts } = countingNative()
- 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 }))
- const first = yield* consume()
- const second = yield* consume()
- yield* Effect.yieldNow
- expect(counts.subscribes).toBe(1)
- yield* Fiber.interrupt(first)
- expect(counts.unsubscribes).toBe(0)
- yield* Fiber.interrupt(second)
- expect(counts.subscribes).toBe(1)
- expect(counts.unsubscribes).toBe(1)
- }).pipe(withNative(native))
- })
- it.effect("scope shutdown releases an active subscription exactly once", () => {
- const { native, counts } = countingNative()
- return Effect.gen(function* () {
- const consumer = yield* Effect.gen(function* () {
- const watcher = yield* Watcher.Service
- const updates = yield* watcher.subscribe({ path: "/active", type: "directory" })
- const consumer = yield* updates.pipe(Stream.runDrain, Effect.forkScoped({ startImmediately: true }))
- yield* Effect.yieldNow
- expect(counts.subscribes).toBe(1)
- expect(counts.unsubscribes).toBe(0)
- return consumer
- }).pipe(withNative(native))
- // Closing the layer scope tears the native subscription down while the
- // consumer still holds a reference; the consumer's own release as its
- // stream ends must not tear it down a second time.
- yield* Fiber.join(consumer)
- expect(counts.unsubscribes).toBe(1)
- })
- })
- })
- 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 })),
- )
- 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>
- watcher?: Layer.Layer<Watcher.Service>
- },
- ) {
- return Effect.acquireRelease(
- Effect.promise(async () => {
- const tmp = await tmpdir()
- if (options?.vcs === "hg") {
- await fs.mkdir(path.join(tmp.path, ".hg"))
- return { tmp, vcs: { type: "hg" as const, store: AbsolutePath.make(path.join(tmp.path, ".hg")) } }
- }
- if (options?.vcs !== "git") return { tmp, vcs: undefined }
- await $`git init`.cwd(tmp.path).quiet()
- await $`git config core.fsmonitor false`.cwd(tmp.path).quiet()
- await $`git config commit.gpgsign false`.cwd(tmp.path).quiet()
- await $`git config user.email test@opencode.test`.cwd(tmp.path).quiet()
- await $`git config user.name Test`.cwd(tmp.path).quiet()
- await $`git commit --allow-empty -m root`.cwd(tmp.path).quiet()
- await options.init?.(tmp.path)
- 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, 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
- const deferred = yield* Deferred.make<WatcherEvent>()
- const fiber = yield* bus.subscribe(FileSystem.Event.Changed).pipe(
- Stream.runForEach((event) => {
- if (!check(event.data)) return Effect.void
- return Deferred.succeed(deferred, event.data).pipe(Effect.asVoid)
- }),
- Effect.forkScoped,
- )
- yield* Effect.yieldNow
- return { deferred, fiber }
- })
- }
- function maybeNextUpdate<E>(
- check: (event: WatcherEvent) => boolean,
- trigger: Effect.Effect<void, E>,
- timeout: Duration.Input = "5 seconds",
- ) {
- return Effect.acquireUseRelease(
- wait(check),
- ({ deferred }) => trigger.pipe(Effect.andThen(Deferred.await(deferred)), Effect.timeoutOption(timeout)),
- ({ fiber }) => Fiber.interrupt(fiber),
- )
- }
- function nextUpdate<E>(check: (event: WatcherEvent) => boolean, trigger: Effect.Effect<void, E>) {
- return Effect.gen(function* () {
- const result = yield* maybeNextUpdate(check, trigger)
- if (Option.isSome(result)) return result.value
- return yield* Effect.fail(new Error("timed out waiting for file watcher update"))
- })
- }
- function eventuallyUpdate<E>(check: (event: WatcherEvent) => boolean, trigger: () => Effect.Effect<void, E>) {
- return Effect.gen(function* () {
- while (true) {
- const result = yield* maybeNextUpdate(check, trigger(), "250 millis")
- if (Option.isSome(result)) return result.value
- }
- }).pipe(
- Effect.timeoutOrElse({
- duration: "5 seconds",
- orElse: () => Effect.fail(new Error("timed out waiting for file watcher readiness")),
- }),
- )
- }
- 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 === eventFile,
- () => fs.writeFileString(file, content),
- ).pipe(Effect.asVoid)
- })
- }
- describeNative("LocationWatcher", () => {
- it.live("limits file watches to the exact target", () =>
- withTmp((directory) =>
- Effect.gen(function* () {
- const fs = yield* FSUtil.Service
- const watcher = yield* Watcher.Service
- const target = path.join(directory, "opencode.json")
- const sibling = path.join(directory, "other.json")
- const updates = yield* watcher.subscribe({ path: target, type: "file" })
- const update = yield* updates.pipe(
- Stream.take(1),
- Stream.runHead,
- Effect.forkScoped({ startImmediately: true }),
- )
- yield* fs.writeFileString(sibling, "sibling")
- const writes = yield* Effect.suspend(() => fs.writeFileString(target, `target-${Math.random()}`)).pipe(
- Effect.repeat(Schedule.spaced("10 millis")),
- Effect.forkScoped,
- )
- const event = yield* Fiber.join(update).pipe(Effect.ensuring(Fiber.interrupt(writes)))
- expect(event.valueOrUndefined?.path).toBe(target)
- }).pipe(Effect.provide(AppNodeBuilder.build(Watcher.node))),
- ),
- )
- it.live("detects creation of a missing directory target", () =>
- withTmp((directory) =>
- Effect.gen(function* () {
- const fs = yield* FSUtil.Service
- 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)))
- expect(event.valueOrUndefined?.path).toBe(target)
- }).pipe(Effect.provide(AppNodeBuilder.build(Watcher.node))),
- ),
- )
- it.live("publishes .git/HEAD events", () =>
- withTmp(
- (directory) =>
- Effect.gen(function* () {
- const fs = yield* FSUtil.Service
- const head = path.join(directory, ".git", "HEAD")
- const branch = `watch-${Math.random().toString(36).slice(2)}`
- 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`)),
- ).toEqual({ file: head, event: "change" })
- }),
- { vcs: "git" },
- ),
- )
- const describeSymlink = process.platform !== "win32" ? describe : describe.skip
- describeSymlink("symlinked .git", () => {
- it.live("publishes .git/HEAD events through a symlinked .git directory", () =>
- withTmp(
- (directory) =>
- Effect.gen(function* () {
- 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 })))
- 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(
- yield* nextUpdate(
- (event) => event.file === path.join(actual, "HEAD"),
- afs.writeFileString(head, `ref: refs/heads/${branch}\n`),
- ),
- ).toEqual({ file: path.join(actual, "HEAD"), event: "change" })
- }),
- {
- vcs: "git",
- init: async (directory) => {
- const actual = path.join(directory, "..", `actual_${path.basename(directory)}`)
- await fs.rename(path.join(directory, ".git"), actual)
- await fs.symlink(actual, path.join(directory, ".git"))
- },
- },
- ),
- )
- })
- it.live("publishes .hg/branch events", () =>
- withTmp(
- (directory) =>
- Effect.gen(function* () {
- const fs = yield* FSUtil.Service
- const branch = path.join(directory, ".hg", "branch")
- yield* ready(branch)
- expect(
- yield* nextUpdate((event) => event.file === branch, fs.writeFileString(branch, "feature\n")),
- ).toMatchObject({ file: branch })
- }),
- { vcs: "hg" },
- ),
- )
- })
|