Compare commits

...

1 Commits

Author SHA1 Message Date
Kit Langton f0d1c75e38 fix(core): flush plugin reload generations 2026-08-08 13:55:20 -04:00
2 changed files with 14 additions and 6 deletions
+10 -4
View File
@@ -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<void>()
const ready = { current: yield* Deferred.make<void>() }
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<void>()
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)) })
}),
)
+4 -2
View File
@@ -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)
}),
),