diff --git a/apps/web/src/components/preview/PreviewAutomationHosts.test.tsx b/apps/web/src/components/preview/PreviewAutomationHosts.test.tsx index 542b84907b6e..d78dac503107 100644 --- a/apps/web/src/components/preview/PreviewAutomationHosts.test.tsx +++ b/apps/web/src/components/preview/PreviewAutomationHosts.test.tsx @@ -90,9 +90,9 @@ const snapshot: PreviewSessionSnapshot = { }; const emptyList = { sessions: [], serverEpoch: "test-server", revision: 0 }; const listAtom = Atom.make(AsyncResult.success(emptyList)); -const requestsAtom = Atom.make>( - AsyncResult.initial(false), -); +const requestsAtom = Atom.make< + AsyncResult.AsyncResult, Error> +>(AsyncResult.initial(false)); const requestEvent: PreviewAutomationStreamEvent = { type: "request", connectionId: "automation-connection", @@ -164,7 +164,7 @@ describe("PreviewAutomationHosts open", () => { mocks.respond.mockImplementationOnce(async ({ input }) => response.resolve(input)); await act(async () => { - appAtomRegistry.set(requestsAtom, AsyncResult.success(requestEvent)); + appAtomRegistry.set(requestsAtom, AsyncResult.success([requestEvent])); await readStarted.promise; }); expect(mocks.open).not.toHaveBeenCalled(); @@ -195,7 +195,7 @@ describe("PreviewAutomationHosts open", () => { mocks.respond.mockImplementationOnce(async ({ input }) => response.resolve(input)); await act(async () => { - appAtomRegistry.set(requestsAtom, AsyncResult.success(requestEvent)); + appAtomRegistry.set(requestsAtom, AsyncResult.success([requestEvent])); await response.promise; }); @@ -216,7 +216,7 @@ describe("PreviewAutomationHosts ownership", () => { await act(() => { appAtomRegistry.set( requestsAtom, - AsyncResult.success({ type: "connected", connectionId: "automation-connection" }), + AsyncResult.success([{ type: "connected", connectionId: "automation-connection" }]), ); applyPreviewServerSnapshot(threadRef, snapshot); }); @@ -281,7 +281,7 @@ describe("PreviewAutomationHosts ownership", () => { await act(() => { appAtomRegistry.set( requestsAtom, - AsyncResult.success({ type: "connected", connectionId: "reconnected" }), + AsyncResult.success([{ type: "connected", connectionId: "reconnected" }]), ); }); expect(mocks.focus).toHaveBeenLastCalledWith( @@ -308,14 +308,14 @@ describe("PreviewAutomationHosts ownership", () => { await act(() => { appAtomRegistry.set( requestsAtom, - AsyncResult.success({ type: "connected", connectionId: "first" }), + AsyncResult.success([{ type: "connected", connectionId: "first" }]), ); }); if (reconnect) { await act(() => { appAtomRegistry.set( requestsAtom, - AsyncResult.success({ type: "connected", connectionId: "second" }), + AsyncResult.success([{ type: "connected", connectionId: "second" }]), ); }); } @@ -343,15 +343,17 @@ describe("PreviewAutomationHosts ownership", () => { applyPreviewServerSnapshot(threadRef, snapshot); appAtomRegistry.set( requestsAtom, - AsyncResult.success({ - ...requestEvent, - request: { - ...requestEvent.request, - operation: "status", - tabId: snapshot.tabId, - input: {}, + AsyncResult.success([ + { + ...requestEvent, + request: { + ...requestEvent.request, + operation: "status", + tabId: snapshot.tabId, + input: {}, + }, }, - }), + ]), ); await response.promise; }); diff --git a/apps/web/src/components/preview/previewAutomationRequestConsumer.test.ts b/apps/web/src/components/preview/previewAutomationRequestConsumer.test.ts index af3a95c32c78..57123e53bc95 100644 --- a/apps/web/src/components/preview/previewAutomationRequestConsumer.test.ts +++ b/apps/web/src/components/preview/previewAutomationRequestConsumer.test.ts @@ -6,6 +6,7 @@ import { PreviewTabId, ThreadId, } from "@t3tools/contracts"; +import * as Stream from "effect/Stream"; import { AsyncResult, Atom, AtomRegistry } from "effect/unstable/reactivity"; import { describe, expect, it, vi } from "vite-plus/test"; @@ -55,10 +56,9 @@ const consumerState = (handleRequest: (request: PreviewAutomationRequest) => Pro describe("previewAutomationRequestConsumer", () => { it("acknowledges a replacement stream before consuming requests from it", async () => { const requestsAtom = Atom.make( - AsyncResult.success({ - type: "connected", - connectionId, - }), + AsyncResult.success, Error>([ + { type: "connected", connectionId }, + ]), ); const handleRequest = vi.fn(async () => undefined); const respond = vi.fn(async () => undefined); @@ -75,7 +75,7 @@ describe("previewAutomationRequestConsumer", () => { const registry = AtomRegistry.make(); registry.mount(consumerAtom); - registry.set(requestsAtom, AsyncResult.success(requestEvent("request-after-connect"))); + registry.set(requestsAtom, AsyncResult.success([requestEvent("request-after-connect")])); await vi.waitFor(() => expect(registry.get(state.connectionAtom)).toBe(connectionId)); await vi.waitFor(() => expect(respond).toHaveBeenCalledTimes(1)); @@ -85,10 +85,9 @@ describe("previewAutomationRequestConsumer", () => { it("drops late requests from an older stream generation", async () => { const requestsAtom = Atom.make( - AsyncResult.success({ - type: "connected", - connectionId: "connection-2", - }), + AsyncResult.success, Error>([ + { type: "connected", connectionId: "connection-2" }, + ]), ); const handleRequest = vi.fn(async () => undefined); const respond = vi.fn(async () => undefined); @@ -107,7 +106,7 @@ describe("previewAutomationRequestConsumer", () => { registry.mount(consumerAtom); registry.set( requestsAtom, - AsyncResult.success(requestEvent("request-stale", {}, "connection-1")), + AsyncResult.success([requestEvent("request-stale", {}, "connection-1")]), ); await vi.waitFor(() => expect(registry.get(state.connectionAtom)).toBe("connection-2")); @@ -117,9 +116,9 @@ describe("previewAutomationRequestConsumer", () => { }); it("consumes every request emitted before React can render", async () => { - const requestsAtom = Atom.make>( - AsyncResult.initial(false), - ); + const requestsAtom = Atom.make< + AsyncResult.AsyncResult, Error> + >(AsyncResult.initial, Error>(false)); const handleRequest = vi.fn(async (value: PreviewAutomationRequest) => ({ requestId: value.requestId, })); @@ -140,8 +139,8 @@ describe("previewAutomationRequestConsumer", () => { const registry = AtomRegistry.make(); registry.mount(consumerAtom); - registry.set(requestsAtom, AsyncResult.success(requestEvent("request-1"))); - registry.set(requestsAtom, AsyncResult.success(requestEvent("request-2"))); + registry.set(requestsAtom, AsyncResult.success([requestEvent("request-1")])); + registry.set(requestsAtom, AsyncResult.success([requestEvent("request-2")])); await vi.waitFor(() => expect(respond).toHaveBeenCalledTimes(2)); expect(handleRequest.mock.calls.map(([value]) => value.requestId)).toEqual([ @@ -152,10 +151,45 @@ describe("previewAutomationRequestConsumer", () => { registry.dispose(); }); - it("uses the latest request handler without rebuilding the stream consumer", async () => { - const requestsAtom = Atom.make>( - AsyncResult.initial(false), + it("handles every request that arrives in one batch of the request stream", async () => { + // Requests from several threads reach the host together. The production atom delivers each + // batch whole, because a stream atom keeps only the last element of a batch. + const requestsAtom = Atom.make( + Stream.fromIterable([ + { type: "connected", connectionId }, + requestEvent("request-1"), + requestEvent("request-2"), + requestEvent("request-3"), + ]).pipe(Stream.chunks), ); + const handleRequest = vi.fn(async () => undefined); + const respond = vi.fn(async (_response: PreviewAutomationResponse) => undefined); + const state = consumerState(handleRequest); + const consumerAtom = createPreviewAutomationRequestConsumerAtom({ + requestsAtom, + clientId, + connectionAtom: state.connectionAtom, + environmentId, + requestHandlerAtom: state.requestHandlerAtom, + respond, + label: "test:preview-automation-batch", + }); + const registry = AtomRegistry.make(); + registry.mount(consumerAtom); + + await vi.waitFor(() => expect(respond).toHaveBeenCalledTimes(3)); + expect(respond.mock.calls.map(([response]) => response.requestId)).toEqual([ + "request-1", + "request-2", + "request-3", + ]); + registry.dispose(); + }); + + it("uses the latest request handler without rebuilding the stream consumer", async () => { + const requestsAtom = Atom.make< + AsyncResult.AsyncResult, Error> + >(AsyncResult.initial, Error>(false)); const firstHandler = vi.fn(async () => "first"); const secondHandler = vi.fn(async () => "second"); const respond = vi.fn(async (_response: PreviewAutomationResponse) => undefined); @@ -172,10 +206,10 @@ describe("previewAutomationRequestConsumer", () => { const registry = AtomRegistry.make(); registry.mount(consumerAtom); - registry.set(requestsAtom, AsyncResult.success(requestEvent("request-first"))); + registry.set(requestsAtom, AsyncResult.success([requestEvent("request-first")])); await vi.waitFor(() => expect(respond).toHaveBeenCalledTimes(1)); registry.set(state.requestHandlerAtom, { handle: secondHandler }); - registry.set(requestsAtom, AsyncResult.success(requestEvent("request-second"))); + registry.set(requestsAtom, AsyncResult.success([requestEvent("request-second")])); await vi.waitFor(() => expect(respond).toHaveBeenCalledTimes(2)); expect(firstHandler).toHaveBeenCalledTimes(1); @@ -186,7 +220,9 @@ describe("previewAutomationRequestConsumer", () => { it("consumes a request that arrived immediately before the consumer mounted", async () => { const requestsAtom = Atom.make( - AsyncResult.success(requestEvent("request-ready")), + AsyncResult.success, Error>([ + requestEvent("request-ready"), + ]), ); const respond = vi.fn(async (_response: PreviewAutomationResponse) => undefined); const state = consumerState(async () => undefined); @@ -351,9 +387,9 @@ describe("previewAutomationRequestConsumer", () => { }); it("sanitizes unexpected handler failures at the response boundary", async () => { - const requestsAtom = Atom.make>( - AsyncResult.initial(false), - ); + const requestsAtom = Atom.make< + AsyncResult.AsyncResult, Error> + >(AsyncResult.initial, Error>(false)); const responses: PreviewAutomationResponse[] = []; const state = consumerState(async () => { throw new Error("desktop IPC secret: do-not-return"); @@ -374,12 +410,12 @@ describe("previewAutomationRequestConsumer", () => { registry.set( requestsAtom, - AsyncResult.success( + AsyncResult.success([ requestEvent("request-failed", { operation: "click", tabId, }), - ), + ]), ); await vi.waitFor(() => expect(responses).toHaveLength(1)); diff --git a/apps/web/src/components/preview/previewAutomationRequestConsumer.ts b/apps/web/src/components/preview/previewAutomationRequestConsumer.ts index 89a9387e4af6..dbe0d7056ed6 100644 --- a/apps/web/src/components/preview/previewAutomationRequestConsumer.ts +++ b/apps/web/src/components/preview/previewAutomationRequestConsumer.ts @@ -12,7 +12,11 @@ import { serializePreviewAutomationHostError, } from "./previewAutomationErrors"; -type AutomationStreamResult = AsyncResult.AsyncResult; +/** The events the host's request stream delivered together, in order. */ +type AutomationStreamResult = AsyncResult.AsyncResult< + ReadonlyArray, + E +>; export function serializePreviewAutomationError( error: unknown, @@ -45,7 +49,10 @@ export function createPreviewAutomationRequestConsumerAtom(options: { const consume = (result: AutomationStreamResult) => { if (!AsyncResult.isSuccess(result)) return; - const event = result.value; + for (const event of result.value) consumeEvent(event); + }; + + const consumeEvent = (event: PreviewAutomationStreamEvent) => { if (event.type === "connected") { activeConnectionId = event.connectionId; connectionExplicitlyAnnounced = true; @@ -96,12 +103,15 @@ export function createPreviewAutomationRequestConsumerAtom(options: { disposed = true; }); const initialRequest = get.once(options.requestsAtom); - if (AsyncResult.isSuccess(initialRequest)) { - activeConnectionId = initialRequest.value.connectionId; - connectionExplicitlyAnnounced = initialRequest.value.type === "connected"; - if (initialRequest.value.type === "connected") { - reportedConnectionId = initialRequest.value.connectionId; - get.set(options.connectionAtom, initialRequest.value.connectionId); + const initialEvent = AsyncResult.isSuccess(initialRequest) + ? initialRequest.value.at(-1) + : undefined; + if (initialEvent !== undefined) { + activeConnectionId = initialEvent.connectionId; + connectionExplicitlyAnnounced = initialEvent.type === "connected"; + if (initialEvent.type === "connected") { + reportedConnectionId = initialEvent.connectionId; + get.set(options.connectionAtom, initialEvent.connectionId); } } get.subscribe(options.requestsAtom, (result) => { @@ -110,9 +120,9 @@ export function createPreviewAutomationRequestConsumerAtom(options: { }); queueMicrotask(() => { const initialConnectionWasSkipped = - AsyncResult.isSuccess(initialRequest) && - initialRequest.value.connectionId === activeConnectionId && - initialRequest.value.connectionId !== reportedConnectionId; + initialEvent !== undefined && + initialEvent.connectionId === activeConnectionId && + initialEvent.connectionId !== reportedConnectionId; if (!disposed && (requestsVersion === 0 || initialConnectionWasSkipped)) { consume(initialRequest); } diff --git a/packages/client-runtime/src/state/preview.ts b/packages/client-runtime/src/state/preview.ts index 86ca157047ba..03a117f7b157 100644 --- a/packages/client-runtime/src/state/preview.ts +++ b/packages/client-runtime/src/state/preview.ts @@ -1,4 +1,5 @@ import { WS_METHODS } from "@t3tools/contracts"; +import * as Stream from "effect/Stream"; import { Atom } from "effect/unstable/reactivity"; import type { EnvironmentRegistry } from "../connection/registry.ts"; @@ -52,6 +53,9 @@ export function createPreviewEnvironmentAtoms( // stream immediately with its owner so stale requests cannot replay when // a thread remounts and the server can clear disconnected hosts promptly. idleTtlMs: 0, + // A stream atom keeps only the last event of each batch it receives, and requests from + // several threads arrive together. Each value is the whole batch, so none is dropped. + transform: Stream.chunks, }), open: createEnvironmentRpcCommand(runtime, { label: "environment-data:preview:open",