diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index a9a2c3fa10d6..5e536b4aa69e 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -6398,7 +6398,10 @@ it.layer(NodeServices.layer)("server router seam", (it) => { projectionSnapshotQuery: { getThreadDetailSnapshot: () => Effect.gen(function* () { - yield* Effect.sleep("25 millis"); + // Publish during the snapshot load, with no sleep. Effect 4 + // schedules forkScoped work on a later tick unless + // startImmediately is set, so a delay here would hide a lost + // event. yield* PubSub.publish(liveEvents, messageEvent); return Option.some({ snapshotSequence: 1, thread }); }), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 226c82cdb1ac..b8c088767b6c 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1465,6 +1465,7 @@ const makeWsRpcLayer = ( const liveBuffer = yield* Queue.unbounded(); yield* Effect.forkScoped( liveStream.pipe(Stream.runForEach((item) => Queue.offer(liveBuffer, item))), + { startImmediately: true }, ); const bufferedLiveStream = Stream.fromQueue(liveBuffer);