diff --git a/packages/reflex-base/news/7053.performance.md b/packages/reflex-base/news/7053.performance.md new file mode 100644 index 00000000000..f3e3ffeb473 --- /dev/null +++ b/packages/reflex-base/news/7053.performance.md @@ -0,0 +1 @@ +Reduce event-queue processing overhead when prepending pending events. diff --git a/packages/reflex-base/src/reflex_base/.templates/web/utils/state.js b/packages/reflex-base/src/reflex_base/.templates/web/utils/state.js index 8ba6d00509c..3cf490e75ca 100644 --- a/packages/reflex-base/src/reflex_base/.templates/web/utils/state.js +++ b/packages/reflex-base/src/reflex_base/.templates/web/utils/state.js @@ -506,15 +506,10 @@ export const queueEvents = async ( params, ) => { if (prepend) { - // Drain the existing queue and place it after the given events. - events = [ - ...events, - ...Array.from({ length: event_queue.length }).map(() => - event_queue.shift(), - ), - ]; + event_queue.unshift(...events.filter((e) => e !== undefined && e !== null)); + } else { + event_queue.push(...events.filter((e) => e !== undefined && e !== null)); } - event_queue.push(...events.filter((e) => e !== undefined && e !== null)); await processEvent(resolveSocket(socket), navigate, params); }; diff --git a/tests/units/reflex_base/client_event_queue.mjs b/tests/units/reflex_base/client_event_queue.mjs new file mode 100644 index 00000000000..6c4a1d24f54 --- /dev/null +++ b/tests/units/reflex_base/client_event_queue.mjs @@ -0,0 +1,416 @@ +import assert from "node:assert/strict"; +import fs from "node:fs"; +import { test } from "node:test"; +import { createQueueRuntime } from "./client_event_queue_runtime.mjs"; + +const source = fs.readFileSync(process.argv[2], "utf8"); +const createRuntime = (options) => createQueueRuntime(source, options); +const params = { current: {} }; +const stateful = (id) => ({ + name: "reflex___state.test.event", + payload: { id }, +}); +const local = (id, output) => ({ + name: "_call_function", + payload: { function: () => output.push(id) }, +}); +const socketFor = (output, connected = true) => ({ + connected, + emit: (_, event) => output.push(event.payload.id), +}); +const deferred = () => { + let resolve; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +}; + +const loaderFixture = ` +const event_queue = []; +let backend_state_mismatch = false; +export const isStateful = () => { + return false; +}; +export const queueEventIfSocketExists = async () => { +}; +export const applyEvent = async () => { +}; +export const applyRestEvent = async () => { +}; +const resolveSocket = (socket) => { + return socket; +}; +export const processEvent = async () => { +}; +function urlFrom(string) { + return new URL(string); +} +`; + +for (const [name, declaration] of [ + [ + "normally formatted declarations", + `export const queueEvents = async (events) => { + event_queue.push(...events); + return event_queue.length; +};`, + ], + [ + "compact declarations", + "export const queueEvents=async(events)=>{event_queue.push(...events);return event_queue.length};", + ], + [ + "function declarations", + `export async function queueEvents(events) { + event_queue.push(...events); + return event_queue.length; +}`, + ], + [ + "nested closures with unindented braces", + `export const queueEvents = async (events) => { +const enqueue = (event) => { +event_queue.push(event); +}; +events.forEach(enqueue); +return event_queue.length; +};`, + ], +]) { + test(`runtime loads ${name} without rewriting functions`, async () => { + const runtime = await createQueueRuntime(loaderFixture + declaration); + const event = { name: "test.event" }; + assert.equal(await runtime.queueEvents([event]), 1); + assert.equal(runtime.event_queue[0], event); + }); +} + +test("FIFO filtering and both raw and reference sockets", async () => { + for (const ref of [false, true]) { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output); + await runtime.queueEvents( + [null, stateful(1), undefined, local(2, output), stateful(3)], + ref ? { current: socket } : socket, + false, + () => {}, + params, + ); + assert.deepEqual(output, [1, 2, 3]); + assert.equal(runtime.event_queue.length, 0); + assert.equal(runtime.isStateful(), false); + } +}); + +test("offline stateful events hold the entire queue until reconnect", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false); + await runtime.queueEvents( + [local(1, output), stateful(2), local(3, output)], + socket, + false, + () => {}, + params, + ); + assert.deepEqual(output, []); + assert.equal(runtime.event_queue.length, 3); + assert.equal(runtime.isStateful(), true); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, [1, 2, 3]); + assert.equal(runtime.event_queue.length, 0); +}); + +test("local events run without a socket, including an empty queue", async () => { + const runtime = await createRuntime(), + output = []; + await runtime.queueEvents([], null, false, () => {}, params); + await runtime.queueEvents( + [local(1, output), local(2, output)], + null, + false, + () => {}, + params, + ); + assert.deepEqual(output, [1, 2]); +}); + +test("prepend preserves new and pending order and does not mutate the input", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false); + await runtime.queueEvents( + [stateful(3), local(4, output)], + socket, + false, + () => {}, + params, + ); + const front = [stateful(1), null, local(2, output), undefined]; + const original = front.slice(); + await runtime.queueEvents(front, socket, true, () => {}, params); + assert.deepEqual(front, original); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, [1, 2, 3, 4]); + assert.equal(runtime.event_queue.length, 0); +}); + +test("disconnect during an event pauses remaining stateful events", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output); + const disconnect = { + name: "_call_function", + payload: { + function: () => { + output.push(1); + socket.connected = false; + }, + }, + }; + await runtime.queueEvents( + [disconnect, local(2, output), stateful(3)], + socket, + false, + () => {}, + params, + ); + assert.deepEqual(output, [1]); + assert.equal(runtime.event_queue.length, 2); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, [1, 2, 3]); +}); + +test("reentrant enqueue and prepend retain FIFO semantics", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output); + let nested; + const enqueue = { + name: "_call_function", + payload: { + function: () => { + output.push(1); + nested = runtime.queueEvents( + [local(2, output)], + socket, + true, + () => {}, + params, + ); + }, + }, + }; + await runtime.queueEvents( + [enqueue, local(3, output)], + socket, + false, + () => {}, + params, + ); + await nested; + assert.deepEqual(output, [1, 2, 3]); + assert.equal(runtime.event_queue.length, 0); +}); + +test("overlapping calls and promise settling preserve async handler behavior", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output), + wait = deferred(); + let firstSettled = false; + const asynchronous = { + name: "_call_function", + payload: { + function: () => { + output.push(1); + return wait.promise; + }, + callback: () => output.push(4), + }, + }; + const first = runtime + .queueEvents( + [asynchronous, local(2, output)], + socket, + false, + () => {}, + params, + ) + .then(() => { + firstSettled = true; + }); + const second = runtime.queueEvents( + [local(3, output)], + socket, + false, + () => {}, + params, + ); + await second; + assert.deepEqual(output, [1, 2, 3]); + assert.equal(firstSettled, false); + wait.resolve(); + await first; + assert.deepEqual(output, [1, 2, 3, 4]); + assert.equal(firstSettled, true); +}); + +test("fatal mismatch clears pending events and lets drain promises settle", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false); + await runtime.queueEvents( + [stateful(1), local(2, output)], + socket, + false, + () => {}, + params, + ); + runtime.setMismatch(true); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, []); + assert.equal(runtime.event_queue.length, 0); +}); + +test("redirect and REST events retain the ordering of pending work", async () => { + const output = []; + const runtime = await createRuntime({ + uploadFiles: () => output.push("upload"), + }), + socket = socketFor(output); + await runtime.queueEvents( + [ + { name: "_redirect", payload: { path: "/next", replace: true } }, + { + name: "reflex___state.test.upload", + handler: "uploadFiles", + payload: { files: [] }, + }, + local("last", output), + ], + socket, + false, + (path, options) => output.push([path, options]), + params, + ); + assert.deepEqual(output, [["/next", { replace: true }], "upload", "last"]); + assert.equal(runtime.event_queue.length, 0); +}); + +test("prepend shifts pending events only when they are dispatched", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false); + await runtime.queueEvents( + [stateful(2), local(3, output)], + socket, + false, + () => {}, + params, + ); + let shifts = 0; + runtime.event_queue.shift = () => { + shifts++; + return Array.prototype.shift.call(runtime.event_queue); + }; + socket.connected = true; + await runtime.queueEvents([local(1, output)], socket, true, () => {}, params); + assert.deepEqual(output, [1, 2, 3]); + assert.equal(shifts, 3); +}); + +test("dispatch rejection leaves pending work available to a later drain", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output); + socket.emit = () => { + throw new Error("socket write failed"); + }; + await assert.rejects( + runtime.queueEvents( + [stateful(1), local(2, output)], + socket, + false, + () => {}, + params, + ), + /socket write failed/, + ); + assert.equal(runtime.event_queue.length, 1); + socket.emit = (_, event) => output.push(event.payload.id); + await runtime.queueEvents([stateful(3)], socket, false, () => {}, params); + assert.deepEqual(output, [2, 3]); + assert.equal(runtime.event_queue.length, 0); +}); + +test("stateful arrival during an offline local await pauses the later drain", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false), + wait = deferred(); + const asynchronous = { + name: "_call_function", + payload: { function: () => wait.promise, callback: () => output.push(1) }, + }; + const first = runtime.queueEvents( + [asynchronous], + socket, + false, + () => {}, + params, + ); + await runtime.queueEvents( + [local(2, output), stateful(3)], + socket, + false, + () => {}, + params, + ); + assert.deepEqual(output, []); + wait.resolve(); + await first; + assert.deepEqual(output, [1]); + assert.equal(runtime.event_queue.length, 2); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, [1, 2, 3]); +}); + +test("mismatch retains the existing offline stateful guard until reconnect", async () => { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output, false); + runtime.setMismatch(true); + await runtime.queueEvents([stateful(1)], socket, false, () => {}, params); + assert.equal(runtime.event_queue.length, 1); + assert.deepEqual(output, []); + socket.connected = true; + await runtime.processEvent(socket, () => {}, params); + assert.equal(runtime.event_queue.length, 0); + assert.deepEqual(output, []); +}); + +test("one-event queues drain exactly once for both local and stateful handlers", async () => { + for (const makeEvent of [local, (id) => stateful(id)]) { + const runtime = await createRuntime(), + output = [], + socket = socketFor(output); + await runtime.queueEvents( + [makeEvent(1, output)], + socket, + false, + () => {}, + params, + ); + await runtime.processEvent(socket, () => {}, params); + assert.deepEqual(output, [1]); + assert.equal(runtime.event_queue.length, 0); + } +}); diff --git a/tests/units/reflex_base/client_event_queue_runtime.mjs b/tests/units/reflex_base/client_event_queue_runtime.mjs new file mode 100644 index 00000000000..6a1258d8664 --- /dev/null +++ b/tests/units/reflex_base/client_event_queue_runtime.mjs @@ -0,0 +1,79 @@ +import { SourceTextModule, SyntheticModule } from "node:vm"; + +/** Evaluate the complete frontend module with isolated dependency stubs. */ +export async function createQueueRuntime(source, options = {}) { + const unused = () => { + throw new Error("Unexpected frontend dependency in queue test"); + }; + const dependencies = { + "test:browser": { + window: options.window ?? { + location: { host: "localhost", pathname: "/", search: "", hash: "" }, + }, + document: options.document ?? {}, + localStorage: options.localStorage ?? { clear() {}, removeItem() {} }, + sessionStorage: options.sessionStorage ?? { clear() {}, removeItem() {} }, + }, + "socket.io-client": { default: unused }, + "$/env.json": { default: {} }, + "$/reflex.json": { default: {} }, + "universal-cookie": { + default: class { + constructor() { + return options.cookies ?? { remove() {} }; + } + }, + }, + react: { + useCallback: unused, + useEffect: unused, + useRef: unused, + useState: unused, + }, + "react-router": { + useLocation: unused, + useNavigate: unused, + useSearchParams: unused, + useParams: unused, + }, + "$/utils/context": { + initialEvents: options.initialEvents ?? (() => []), + initialState: {}, + onLoadInternalEvent: unused, + state_name: "test_state", + exception_state_name: "test_exception_state", + }, + "$/utils/helpers/debounce": { default: unused }, + "$/utils/helpers/throttle": { default: unused }, + "$/utils/helpers/upload": { + uploadFiles: options.uploadFiles ?? unused, + }, + }; + // Let Node parse the unchanged module; only expose private state to tests. + const module = new SourceTextModule( + `import { window, document, localStorage, sessionStorage } from "test:browser"; +${source} +export { event_queue }; +export function setMismatch(value) { backend_state_mismatch = value; } +`, + ); + const linked = new Map(); + await module.link((specifier) => { + if (!linked.has(specifier)) { + const exports = dependencies[specifier]; + if (!exports) + throw new Error(`Unexpected import in queue test: ${specifier}`); + linked.set( + specifier, + new SyntheticModule(Object.keys(exports), function () { + for (const [name, value] of Object.entries(exports)) { + this.setExport(name, value); + } + }), + ); + } + return linked.get(specifier); + }); + await module.evaluate(); + return module.namespace; +} diff --git a/tests/units/reflex_base/test_client_event_queue.py b/tests/units/reflex_base/test_client_event_queue.py new file mode 100644 index 00000000000..ff026557aee --- /dev/null +++ b/tests/units/reflex_base/test_client_event_queue.py @@ -0,0 +1,30 @@ +"""Execute behavioral regressions against the actual frontend event queue.""" + +import shutil +import subprocess +from pathlib import Path + +import pytest + + +@pytest.mark.skipif(shutil.which("node") is None, reason="Node.js is unavailable") +def test_client_event_queue() -> None: + """Verify queue ordering, reconnect handling, and prepend processing cost.""" + tests = Path(__file__).parent + source = ( + tests.parents[2] + / "packages/reflex-base/src/reflex_base/.templates/web/utils/state.js" + ) + result = subprocess.run( + [ + "node", + "--experimental-vm-modules", + str(tests / "client_event_queue.mjs"), + str(source), + ], + capture_output=True, + text=True, + check=False, + timeout=30, + ) + assert result.returncode == 0, result.stdout + result.stderr