refactor(core): manage watcher lifecycle with RcMap (#39203)

This commit is contained in:
Kit Langton 2026-07-27 20:01:17 -04:00 committed by GitHub
commit 4333a44e65
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 251 additions and 107 deletions

View file

@ -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<void>
/** 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<Subscription | undefined>
}
/**
* 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<Native, NativeInterface>()("@opencode/Watcher/Native") {}
export interface Interface {
readonly subscribe: (input: WatchInput) => Stream.Stream<Update>
}
@ -65,113 +87,135 @@ export interface TestInterface extends Interface {
export class Test extends Context.Service<Test, TestInterface>()("@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<Update>()
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<Update>
readonly subscription: { readonly unsubscribe: () => Promise<void> }
refs: number
}
const entries = new Map<string, Entry>()
const locks = KeyedMutex.makeUnsafe<string>()
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<Update>()
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<Update>(), (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<Update>,
) {
ignore: readonly string[],
publish: (update: Update) => void,
): Effect.Effect<Subscription | undefined> {
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)),
),
)
}

View file

@ -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<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(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,