diff --git a/packages/core/src/plugin/supervisor.ts b/packages/core/src/plugin/supervisor.ts index 59c9ebd798d..9343424ac25 100644 --- a/packages/core/src/plugin/supervisor.ts +++ b/packages/core/src/plugin/supervisor.ts @@ -233,7 +233,7 @@ const layer = Layer.effect( const bus = yield* Bus.Service const watcher = yield* Watcher.Service const fs = yield* FSUtil.Service - const ready = yield* Deferred.make() + const ready = { current: yield* Deferred.make() } let observed = 0 // Configured local plugin files can live outside config roots, where the @@ -291,7 +291,13 @@ const layer = Layer.effect( bus.subscribe([Event.Updated, SdkPlugins.Updated]), ).pipe( // Make accepted work visible to flush before coalescing the burst. - Stream.mapEffect(() => Effect.sync(() => ++observed)), + Stream.mapEffect(() => + Effect.gen(function* () { + observed++ + if (yield* Deferred.isDone(ready.current)) ready.current = yield* Deferred.make() + return observed + }), + ), ) yield* Stream.concat(Stream.succeed(0), updates).pipe( // Keep observing updates while activation runs, retaining only the latest generation request. @@ -300,12 +306,12 @@ const layer = Layer.effect( Stream.runForEach((target) => Effect.gen(function* () { yield* activate() - if (observed === target) yield* Deferred.succeed(ready, undefined) + if (observed === target) yield* Deferred.succeed(ready.current, undefined) }).pipe(Effect.catchCause((cause) => Effect.logError("failed to reload plugins", { cause }))), ), Effect.forkScoped({ startImmediately: true }), ) - return Service.of({ flush: Deferred.await(ready) }) + return Service.of({ flush: Effect.suspend(() => Deferred.await(ready.current)) }) }), ) diff --git a/packages/core/test/location-layer.test.ts b/packages/core/test/location-layer.test.ts index a70810d0622..da9e73c3c76 100644 --- a/packages/core/test/location-layer.test.ts +++ b/packages/core/test/location-layer.test.ts @@ -305,11 +305,13 @@ describe("LocationServiceMap", () => { ) yield* Deferred.await(started) - yield* PluginSupervisor.Service.use((supervisor) => supervisor.flush).pipe( + const flushFiber = yield* PluginSupervisor.Service.use((supervisor) => supervisor.flush).pipe( Effect.provide(context), - Effect.timeout("1 second"), + Effect.forkChild({ startImmediately: true }), ) + expect(flushFiber.pollUnsafe()).toBeUndefined() yield* Deferred.succeed(release, undefined) + yield* Fiber.join(flushFiber) yield* Deferred.await(completed) }), ),