diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 549f84d8d..2349dfd8a 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -7068,7 +7068,6 @@ it.layer(NodeServices.layer)("server router seam", (it) => { projectionSnapshotQuery: { getThreadDetailSnapshot: () => Effect.gen(function* () { - yield* Effect.sleep("25 millis"); yield* PubSub.publish(liveEvents, messageEvent); return Option.some({ snapshotSequence: 1, thread }); }), @@ -7081,14 +7080,19 @@ it.layer(NodeServices.layer)("server router seam", (it) => { withWsRpcClient(wsUrl, (client) => client[ORCHESTRATION_WS_METHODS.subscribeThread]({ threadId: defaultThreadId, - }).pipe(Stream.take(2), Stream.runCollect), + requestCompletionMarker: true, + }).pipe( + Stream.takeUntil((item) => item.kind === "synchronized"), + Stream.runCollect, + ), ), - ).pipe(Effect.timeout("2 seconds")); + ); assert.equal(items[0]?.kind, "snapshot"); assert.equal(items[1]?.kind, "event"); assert.equal(items[1]?.kind === "event" ? items[1].event.sequence : null, 2); - }).pipe(Effect.provide(NodeHttpServer.layerTest), TestClock.withLive), + assert.equal(items[2]?.kind, "synchronized"); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); it.effect("coalesces buffered live tool updates to the latest state", () => @@ -7503,12 +7507,15 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }, } satisfies Extract; + const readyForLiveEvents = yield* Deferred.make(); yield* buildAppUnderTest({ layers: { orchestrationEngine: { latestSequence: Effect.succeed(3), readEvents: () => Stream.empty, - streamDomainEvents: Stream.make(rollbackEvent, laterEvent), + streamDomainEvents: Stream.fromEffect(Deferred.await(readyForLiveEvents)).pipe( + Stream.flatMap(() => Stream.make(rollbackEvent, laterEvent)), + ), }, }, }); @@ -7519,7 +7526,15 @@ it.layer(NodeServices.layer)("server router seam", (it) => { threadId: defaultThreadId, afterSequence: 3, requestCompletionMarker: true, - }).pipe(Stream.take(2), Stream.runCollect), + }).pipe( + Stream.tap((item) => + item.kind === "synchronized" + ? Deferred.succeed(readyForLiveEvents, undefined) + : Effect.void, + ), + Stream.take(2), + Stream.runCollect, + ), ), ); assert.deepEqual( diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 8ed4900fa..df07314ed 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1571,7 +1571,9 @@ const makeWsRpcLayer = ( // Attach live delivery before reading either replay or snapshot state. // Otherwise an event published while the snapshot is loading is lost. const liveBuffer = yield* makeThreadLiveEventCoalescer(); - yield* Effect.forkScoped(liveStream.pipe(Stream.runForEach(liveBuffer.offer))); + yield* Effect.forkScoped(liveStream.pipe(Stream.runForEach(liveBuffer.offer)), { + startImmediately: true, + }); const bufferedLiveStream = liveBuffer.stream; // When the client already loaded the snapshot over HTTP it passes