Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 19 additions & 17 deletions apps/web/src/components/preview/PreviewAutomationHosts.test.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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.AsyncResult<PreviewAutomationStreamEvent, Error>>(
AsyncResult.initial(false),
);
const requestsAtom = Atom.make<
AsyncResult.AsyncResult<ReadonlyArray<PreviewAutomationStreamEvent>, Error>
>(AsyncResult.initial(false));
const requestEvent: PreviewAutomationStreamEvent = {
type: "request",
connectionId: "automation-connection",
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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;
});

Expand All @@ -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);
});
Expand Down Expand Up @@ -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(
Expand All @@ -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" }]),
);
});
}
Expand Down Expand Up @@ -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;
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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<PreviewAutomationStreamEvent, Error>({
type: "connected",
connectionId,
}),
AsyncResult.success<ReadonlyArray<PreviewAutomationStreamEvent>, Error>([
{ type: "connected", connectionId },
]),
);
const handleRequest = vi.fn(async () => undefined);
const respond = vi.fn(async () => undefined);
Expand All @@ -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));
Expand All @@ -85,10 +85,9 @@ describe("previewAutomationRequestConsumer", () => {

it("drops late requests from an older stream generation", async () => {
const requestsAtom = Atom.make(
AsyncResult.success<PreviewAutomationStreamEvent, Error>({
type: "connected",
connectionId: "connection-2",
}),
AsyncResult.success<ReadonlyArray<PreviewAutomationStreamEvent>, Error>([
{ type: "connected", connectionId: "connection-2" },
]),
);
const handleRequest = vi.fn(async () => undefined);
const respond = vi.fn(async () => undefined);
Expand All @@ -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"));
Expand All @@ -117,9 +116,9 @@ describe("previewAutomationRequestConsumer", () => {
});

it("consumes every request emitted before React can render", async () => {
const requestsAtom = Atom.make<AsyncResult.AsyncResult<PreviewAutomationStreamEvent, Error>>(
AsyncResult.initial<PreviewAutomationStreamEvent, Error>(false),
);
const requestsAtom = Atom.make<
AsyncResult.AsyncResult<ReadonlyArray<PreviewAutomationStreamEvent>, Error>
>(AsyncResult.initial<ReadonlyArray<PreviewAutomationStreamEvent>, Error>(false));
const handleRequest = vi.fn(async (value: PreviewAutomationRequest) => ({
requestId: value.requestId,
}));
Expand All @@ -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([
Expand All @@ -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.AsyncResult<PreviewAutomationStreamEvent, Error>>(
AsyncResult.initial<PreviewAutomationStreamEvent, Error>(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<PreviewAutomationStreamEvent>([
{ 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<ReadonlyArray<PreviewAutomationStreamEvent>, Error>
>(AsyncResult.initial<ReadonlyArray<PreviewAutomationStreamEvent>, Error>(false));
const firstHandler = vi.fn(async () => "first");
const secondHandler = vi.fn(async () => "second");
const respond = vi.fn(async (_response: PreviewAutomationResponse) => undefined);
Expand All @@ -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);
Expand All @@ -186,7 +220,9 @@ describe("previewAutomationRequestConsumer", () => {

it("consumes a request that arrived immediately before the consumer mounted", async () => {
const requestsAtom = Atom.make(
AsyncResult.success<PreviewAutomationStreamEvent, Error>(requestEvent("request-ready")),
AsyncResult.success<ReadonlyArray<PreviewAutomationStreamEvent>, Error>([
requestEvent("request-ready"),
]),
);
const respond = vi.fn(async (_response: PreviewAutomationResponse) => undefined);
const state = consumerState(async () => undefined);
Expand Down Expand Up @@ -351,9 +387,9 @@ describe("previewAutomationRequestConsumer", () => {
});

it("sanitizes unexpected handler failures at the response boundary", async () => {
const requestsAtom = Atom.make<AsyncResult.AsyncResult<PreviewAutomationStreamEvent, Error>>(
AsyncResult.initial<PreviewAutomationStreamEvent, Error>(false),
);
const requestsAtom = Atom.make<
AsyncResult.AsyncResult<ReadonlyArray<PreviewAutomationStreamEvent>, Error>
>(AsyncResult.initial<ReadonlyArray<PreviewAutomationStreamEvent>, Error>(false));
const responses: PreviewAutomationResponse[] = [];
const state = consumerState(async () => {
throw new Error("desktop IPC secret: do-not-return");
Expand All @@ -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));
Expand Down
32 changes: 21 additions & 11 deletions apps/web/src/components/preview/previewAutomationRequestConsumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,11 @@ import {
serializePreviewAutomationHostError,
} from "./previewAutomationErrors";

type AutomationStreamResult<E> = AsyncResult.AsyncResult<PreviewAutomationStreamEvent, E>;
/** The events the host's request stream delivered together, in order. */
type AutomationStreamResult<E> = AsyncResult.AsyncResult<
ReadonlyArray<PreviewAutomationStreamEvent>,
E
>;

export function serializePreviewAutomationError(
error: unknown,
Expand Down Expand Up @@ -45,7 +49,10 @@ export function createPreviewAutomationRequestConsumerAtom<E>(options: {

const consume = (result: AutomationStreamResult<E>) => {
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;
Expand Down Expand Up @@ -96,12 +103,15 @@ export function createPreviewAutomationRequestConsumerAtom<E>(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) => {
Expand All @@ -110,9 +120,9 @@ export function createPreviewAutomationRequestConsumerAtom<E>(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);
}
Expand Down
4 changes: 4 additions & 0 deletions packages/client-runtime/src/state/preview.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -52,6 +53,9 @@ export function createPreviewEnvironmentAtoms<R, E>(
// 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",
Expand Down
Loading