From 4333a44e652828ed5ad1e3ab3ea4dd4a9ce027f1 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 27 Jul 2026 20:01:17 -0400 Subject: [PATCH] refactor(core): manage watcher lifecycle with RcMap (#39203) --- packages/core/src/filesystem/watcher.ts | 271 +++++++++++------- packages/core/test/filesystem/watcher.test.ts | 93 +++++- 2 files changed, 254 insertions(+), 110 deletions(-) diff --git a/packages/core/src/filesystem/watcher.ts b/packages/core/src/filesystem/watcher.ts index 5da2ef8678..11f8ac7f69 100644 --- a/packages/core/src/filesystem/watcher.ts +++ b/packages/core/src/filesystem/watcher.ts @@ -5,8 +5,7 @@ import { createWrapper } from "@parcel/watcher/wrapper" import type ParcelWatcher from "@parcel/watcher" import { FileSystem } from "@opencode-ai/schema/filesystem" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" -import { Cause, Context, Effect, Layer, PubSub, Schema, Scope, Stream } from "effect" -import { KeyedMutex } from "../effect/keyed-mutex" +import { Cause, Context, Effect, Layer, PubSub, RcMap, Schema, Stream } from "effect" import { lazy } from "../util/lazy" import { watch as watchFileSystem } from "node:fs" import path from "path" @@ -45,6 +44,29 @@ export type WatchInput = | { readonly path: string; readonly type: "file" } | { readonly path: string; readonly type: "directory"; readonly ignore?: readonly string[] } +export type Subscription = { + readonly unsubscribe: () => Promise + /** Backend name for logging, e.g. "node" or "fs-events". */ + readonly backend?: string +} + +export interface NativeInterface { + /** Starts one OS-level watch, reporting events through `publish` until unsubscribed. */ + readonly subscribe: (input: { + readonly type: WatchInput["type"] + readonly target: string + readonly ignore: readonly string[] + readonly publish: (update: Update) => void + }) => Effect.Effect +} + +/** + * The OS-level watch implementation behind the Watcher service. The default + * layer uses `node:fs.watch` for files and `@parcel/watcher` for directories; + * tests provide implementations they can control. + */ +export class Native extends Context.Service()("@opencode/Watcher/Native") {} + export interface Interface { readonly subscribe: (input: WatchInput) => Stream.Stream } @@ -65,113 +87,135 @@ export interface TestInterface extends Interface { export class Test extends Context.Service()("@opencode/Watcher/Test") {} -/** In-memory watcher for tests: records subscribe calls and broadcasts emitted updates. */ -export const testLayer = Layer.effectContext( - Effect.gen(function* () { - const updates = yield* PubSub.unbounded() - const subscriptions: WatchInput[] = [] - const service = Test.of({ - subscribe: (input) => { - subscriptions.push(input) - return Stream.fromPubSub(updates) - }, - emit: (update) => PubSub.publish(updates, update).pipe(Effect.asVoid), - subscriptions: () => Effect.sync(() => [...subscriptions]), - }) - return Context.empty().pipe(Context.add(Service, service), Context.add(Test, service)) - }), -) +export const layer = (options?: Options) => + Layer.effect( + Service, + Effect.gen(function* () { + if (options?.enabled === false) { + return Service.of({ subscribe: () => Stream.empty }) + } + const native = yield* Native -export const layer = (options?: Options) => Layer.effect( - Service, - Effect.gen(function* () { - const backend = getBackend() - const native = watcher() - if (options?.enabled === false) { - return Service.of({ subscribe: () => Stream.empty }) - } - - type Entry = { - readonly pubsub: PubSub.PubSub - readonly subscription: { readonly unsubscribe: () => Promise } - refs: number - } - const entries = new Map() - const locks = KeyedMutex.makeUnsafe() - - const acquire = Effect.fn("Watcher.acquire")(function* (input: WatchInput) { - const scope = yield* Scope.Scope - const target = path.resolve(input.path) - const directory = input.type === "file" ? path.dirname(target) : target - const ignore = [...new Set(input.type === "directory" ? (input.ignore ?? []) : [])].toSorted() - const id = JSON.stringify([input.type, target, ignore]) - const pubsub = yield* locks.withLock(id)( - Effect.gen(function* () { - const existing = entries.get(id) - if (existing) { - existing.refs++ - return existing.pubsub - } - const pubsub = yield* PubSub.unbounded() - const subscription = yield* input.type === "file" - ? Effect.sync(() => { - const subscription = watchFileSystem(directory, { recursive: false }, (_event, file) => { - if (file && path.resolve(directory, file.toString()) !== target) return - PubSub.publishUnsafe(pubsub, { - path: target, - type: "update", - } satisfies Update) - }) - if ("on" in subscription && typeof subscription.on === "function") { - subscription.on("error", (error: unknown) => - Effect.runFork(Effect.logError("watcher callback failed", { path: target, error })), - ) - } - return { unsubscribe: () => Promise.resolve(subscription.close()) } - }) - : subscribeDirectory(native, backend, directory, ignore, pubsub) - if (subscription) { - entries.set(id, { pubsub, subscription, refs: 1 }) + // Keys compare structurally (effect Equal), so equivalent watches share one entry. + type Key = { readonly type: WatchInput["type"]; readonly target: string; readonly ignore: readonly string[] } + const watchers = yield* RcMap.make({ + lookup: (key: Key) => + Effect.gen(function* () { + const pubsub = yield* Effect.acquireRelease(PubSub.unbounded(), (pubsub) => + PubSub.shutdown(pubsub), + ) + const subscription = yield* Effect.acquireRelease( + native.subscribe({ + type: key.type, + target: key.target, + ignore: key.ignore, + publish: (update) => PubSub.publishUnsafe(pubsub, update), + }), + (subscription) => + subscription + ? Effect.promise(() => subscription.unsubscribe()).pipe( + Effect.ignoreCause, + Effect.andThen(Effect.logInfo("watcher stopped", { path: key.target, type: key.type })), + ) + : Effect.void, + // Native subscription may stay pending up to SUBSCRIBE_TIMEOUT_MS; + // scope shutdown must not wait behind an uninterruptible acquisition. + { interruptible: true }, + ) + if (!subscription) { + // Unsupported backend: end subscriber streams instead of hanging them. + yield* PubSub.shutdown(pubsub) + return pubsub + } yield* Effect.logInfo("watcher started", { - path: target, - type: input.type, - backend: input.type === "file" ? "node" : backend, - ignores: ignore.length, + path: key.target, + type: key.type, + backend: subscription.backend, + ignores: key.ignore.length, }) return pubsub - } - yield* PubSub.shutdown(pubsub) - return pubsub - }), - ) - - yield* Scope.addFinalizer( - scope, - locks.withLock(id)( - Effect.gen(function* () { - const entry = entries.get(id) - if (!entry) return - entry.refs-- - if (entry.refs > 0) return - entries.delete(id) - yield* Effect.promise(() => entry.subscription.unsubscribe()).pipe(Effect.ignore) - yield* PubSub.shutdown(entry.pubsub) - yield* Effect.logInfo("watcher stopped", { path: target, type: input.type }) }), - ), - ) - return pubsub + }) + + const subscribe = (input: WatchInput) => { + const target = path.resolve(input.path) + const ignore = [...new Set(input.type === "directory" ? (input.ignore ?? []) : [])].toSorted() + return Stream.unwrap( + RcMap.get(watchers, { type: input.type, target, ignore }).pipe( + Effect.map((pubsub) => Stream.fromPubSub(pubsub)), + ), + ) + } + + return Service.of({ subscribe }) + }), + ) + +/** + * Watcher for tests: the real lifecycle over an in-memory Native that records + * acquired watches and broadcasts emitted updates to every active watch. + */ +export const testLayer = Layer.effectContext( + Effect.gen(function* () { + const subscriptions: WatchInput[] = [] + const active = new Set<(update: Update) => void>() + const native = Native.of({ + subscribe: (input) => + Effect.sync(() => { + subscriptions.push( + input.type === "file" + ? { path: input.target, type: "file" } + : input.ignore.length > 0 + ? { path: input.target, type: "directory", ignore: input.ignore } + : { path: input.target, type: "directory" }, + ) + active.add(input.publish) + return { + unsubscribe: () => { + active.delete(input.publish) + return Promise.resolve() + }, + } + }), }) - - const subscribe = (input: WatchInput) => - Stream.unwrap(acquire(input).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))) - - return Service.of({ subscribe }) + const context = yield* Layer.build(layer().pipe(Layer.provide(Layer.succeed(Native, native)))) + const test = Test.of({ + subscribe: Context.get(context, Service).subscribe, + emit: (update) => Effect.sync(() => active.forEach((publish) => publish(update))), + subscriptions: () => Effect.sync(() => [...subscriptions]), + }) + return Context.empty().pipe(Context.add(Service, test), Context.add(Test, test)) }), ) +export const nativeLayer = Layer.succeed( + Native, + Native.of({ + subscribe: (input) => { + if (input.type === "file") { + return Effect.sync(() => { + const directory = path.dirname(input.target) + const subscription = watchFileSystem(directory, { recursive: false }, (_event, file) => { + if (file && path.resolve(directory, file.toString()) !== input.target) return + input.publish({ path: input.target, type: "update" } satisfies Update) + }) + if ("on" in subscription && typeof subscription.on === "function") { + subscription.on("error", (error: unknown) => + Effect.runFork(Effect.logError("watcher callback failed", { path: input.target, error })), + ) + } + return { unsubscribe: () => Promise.resolve(subscription.close()), backend: "node" } + }) + } + return subscribeDirectory(watcher(), getBackend(), input.target, input.ignore, input.publish) + }, + }), +) + +export const nativeNode = makeGlobalNode({ service: Native, layer: nativeLayer, deps: [] }) + export function configured(options?: Options) { - return makeGlobalNode({ service: Service, layer: layer(options), deps: [] }) + return makeGlobalNode({ service: Service, layer: layer(options), deps: [nativeNode] }) } export const node = configured() @@ -180,9 +224,9 @@ function subscribeDirectory( native: typeof import("@parcel/watcher") | undefined, backend: ParcelWatcher.BackendType | undefined, directory: string, - ignore: string[], - pubsub: PubSub.PubSub, -) { + ignore: readonly string[], + publish: (update: Update) => void, +): Effect.Effect { if (!native || !backend) { return Effect.logError("watcher backend not supported", { directory, platform: process.platform }).pipe( Effect.as(undefined), @@ -190,17 +234,26 @@ function subscribeDirectory( } const callback: ParcelWatcher.SubscribeCallback = (error, updates) => { if (error) Effect.runFork(Effect.logError("watcher callback failed", { error })) - for (const update of updates) PubSub.publishUnsafe(pubsub, update) + for (const update of updates) publish(update) } - const pending = native.subscribe(directory, callback, { ignore, backend }) + // Copy `ignore`: it aliases the RcMap key, whose structural hash is cached, + // so the array handed to native code must never be the mutable original. + const pending = native.subscribe(directory, callback, { ignore: [...ignore], backend }) return Effect.promise(() => pending).pipe( + Effect.map((subscription) => ({ unsubscribe: () => subscription.unsubscribe(), backend })), + // Interruption (including the timeout below) abandons the pending native + // subscription, so close it once it eventually resolves. + Effect.onInterrupt(() => + Effect.sync(() => { + pending.then((subscription) => subscription.unsubscribe()).catch(() => {}) + }), + ), Effect.timeout(SUBSCRIBE_TIMEOUT_MS), - Effect.catchCause((cause) => { - pending.then((subscription) => subscription.unsubscribe()).catch(() => {}) - return Effect.logError("failed to subscribe", { + Effect.catchCause((cause) => + Effect.logError("failed to subscribe", { directory, cause: Cause.pretty(cause), - }).pipe(Effect.as(undefined)) - }), + }).pipe(Effect.as(undefined)), + ), ) } diff --git a/packages/core/test/filesystem/watcher.test.ts b/packages/core/test/filesystem/watcher.test.ts index c2a97740c8..4ff39443e1 100644 --- a/packages/core/test/filesystem/watcher.test.ts +++ b/packages/core/test/filesystem/watcher.test.ts @@ -38,11 +38,102 @@ describe("Watcher.testLayer", () => { yield* test.emit({ type: "update", path: "/root/file.md" }) expect(Array.from(yield* Fiber.join(received))).toEqual([{ type: "update", path: "/root/file.md" }]) - expect(yield* test.subscriptions()).toEqual([{ path: "/root", type: "directory" }]) + // 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() + const interrupted = yield* Deferred.make() + yield* Effect.gen(function* () { + const watcher = yield* Watcher.Service + const consumer = yield* watcher + .subscribe({ path: "/pending", type: "directory" }) + .pipe(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(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 consumer = yield* watcher + .subscribe({ path: "/active", type: "directory" }) + .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"]) { const locationLayer = Layer.succeed( Location.Service,