diff --git a/packages/react-server/src/ReactFlightServer.js b/packages/react-server/src/ReactFlightServer.js index ec9d5aa487..f2df73b336 100644 --- a/packages/react-server/src/ReactFlightServer.js +++ b/packages/react-server/src/ReactFlightServer.js @@ -91,6 +91,7 @@ import { getChildFormatContext, initAsyncDebugInfo, markAsyncSequenceRootTask, + markAsyncSequenceRequest, getCurrentAsyncSequence, getAsyncSequenceFromPromise, parseStackTrace, @@ -801,6 +802,11 @@ function RequestInstance( performance.timeOrigin, ); this.abortTime = -0.0; + if (__DEV__ && enableAsyncDebugInfo) { + // Let the async sequence tracking know this request may consume async + // debug info from here on. + markAsyncSequenceRequest(this as any); + } } else { timeOrigin = 0; } @@ -898,6 +904,15 @@ export function resolveRequest(): null | Request { return null; } +// A closed main stream doesn't mean the request is done: a separate debug +// channel can stay open and keep receiving debug info. +export function canRequestStillEmitDebugInfo(request: Request): boolean { + return ( + request.debugDestination !== null || + (request.status !== CLOSING && request.status !== CLOSED) + ); +} + function isTypedArray(value: any): boolean { if (value instanceof ArrayBuffer) { return true; diff --git a/packages/react-server/src/ReactFlightServerConfigDebugNode.js b/packages/react-server/src/ReactFlightServerConfigDebugNode.js index f4030eda5f..e5cabe963e 100644 --- a/packages/react-server/src/ReactFlightServerConfigDebugNode.js +++ b/packages/react-server/src/ReactFlightServerConfigDebugNode.js @@ -26,7 +26,13 @@ import { UNRESOLVED_AWAIT_NODE, } from './ReactFlightAsyncSequence'; import {resolveOwner} from './flight/ReactFlightCurrentOwner'; -import {resolveRequest, isAwaitInUserspace} from './ReactFlightServer'; +import type {Request} from './ReactFlightServer'; + +import { + resolveRequest, + isAwaitInUserspace, + canRequestStillEmitDebugInfo, +} from './ReactFlightServer'; import {createHook, executionAsyncId} from 'async_hooks'; import {promiseHooks} from 'v8'; import {enableAsyncDebugInfo} from 'shared/ReactFeatureFlags'; @@ -72,6 +78,85 @@ const executingPromises: Array> = // Keep the last resolved await as a workaround for async functions missing data. let lastRanAwait: null | AwaitNode = null; +// The requests that might still consume async debug info, held weakly so one +// that never formally closes still expires with GC. +const liveRequests: Set> = + __DEV__ && enableAsyncDebugInfo ? new Set() : (null as any); + +export function markAsyncSequenceRequest(request: Request): void { + if (__DEV__ && enableAsyncDebugInfo) { + liveRequests.add(new WeakRef(request)); + } +} + +function getOldestLiveRequestTimeOrigin(): number { + let oldest = Infinity; + liveRequests.forEach(requestRef => { + const request = requestRef.deref(); + if (request === undefined || !canRequestStillEmitDebugInfo(request)) { + liveRequests.delete(requestRef); + } else if (request.timeOrigin < oldest) { + oldest = request.timeOrigin; + } + }); + return oldest; +} + +// visitAsyncNodeImpl stops at any node that resolved at or before the +// consuming request's time origin, so a node that resolved before every +// possible consumer's origin contributes nothing wherever it appears and can +// drop its own links, releasing the history behind it. Two consumers keep links +// alive: a live request, which pins its own history from its time origin (those +// nodes are re-enqueued, so they don't count against this budget), and a +// request that doesn't exist yet, which may backdate itself with +// options.startTime. Only the latter is unbounded, so it's bounded by count +// rather than by time, which would be multiplied by the tracking rate. Measured +// gaps between capturing a start time and creating the request stay near a +// hundred operations even under load. +const MAX_RETAINED_OPERATIONS = 10000; + +const agingOperations: Array = + __DEV__ && enableAsyncDebugInfo ? [] : (null as any); +let agingHead = 0; +let operationsSinceSweep = 0; + +function trackAgingOperation(node: AsyncSequence): void { + agingOperations.push(node); + if (++operationsSinceSweep >= 128) { + sweepAgingOperations(); + } +} + +function sweepAgingOperations(): void { + operationsSinceSweep = 0; + const cutoff = getOldestLiveRequestTimeOrigin(); + // A node that can't be cleared goes back on the tail, so the queue doesn't + // shorten and the range has to be taken up front. Capped so a backlog drains + // over later sweeps instead of pausing this one. + let remaining = agingOperations.length - agingHead - MAX_RETAINED_OPERATIONS; + if (remaining > 1000) { + remaining = 1000; + } + while (remaining-- > 0) { + const node = agingOperations[agingHead]; + agingOperations[agingHead] = null as any; + agingHead++; + // Unsettled nodes (end < 0) are never cleared. The walk stops at `<=` the + // origin, so requiring strictly before the cutoff is conservative. + if (node.end >= 0 && node.end < cutoff) { + (node as any).awaited = null; + (node as any).previous = null; + } else if (node.awaited !== null || node.previous !== null) { + // Pinned by a live request or not settled yet; reconsider later. + agingOperations.push(node); + } + } + if (agingHead > 1024 && agingHead * 2 >= agingOperations.length) { + agingOperations.splice(0, agingHead); + agingHead = 0; + } +} + function resolvePromiseOrAwaitNode( unresolvedNode: UnresolvedAwaitNode | UnresolvedPromiseNode, endTime: number, @@ -245,6 +330,7 @@ function promiseSettled(node: AsyncSequence, selfResolved: boolean): void { awaited: resolvedNode.awaited, previous: resolvedNode.previous, }; + trackAgingOperation(clonedNode); // We started awaiting on the callback when the original .then() resolved. resolvedNode.start = resolvedNode.end; // It resolved now. We could use the end time of "awaited" maybe. @@ -301,6 +387,7 @@ export function initAsyncDebugInfo(): void { node = createPromiseNode(promise, getCurrentOperation()); } pendingPromises.set(promise, node); + trackAgingOperation(node); }, before(promise: Promise): void { executingPromises.push(promise); @@ -394,6 +481,7 @@ export function initAsyncDebugInfo(): void { awaited: null, previous: null, } as IONode; + trackAgingOperation(node); } else if ( trigger.tag === AWAIT_NODE || trigger.tag === UNRESOLVED_AWAIT_NODE @@ -411,6 +499,7 @@ export function initAsyncDebugInfo(): void { awaited: null, previous: trigger, } as IONode; + trackAgingOperation(node); } else { // Otherwise, this is just a continuation of the same I/O sequence. node = trigger; @@ -452,6 +541,7 @@ export function initAsyncDebugInfo(): void { awaited: ioNode.awaited, previous: ioNode.previous, }; + trackAgingOperation(clonedNode); pendingOperations.set(asyncId, clonedNode); } } else { @@ -473,6 +563,8 @@ export function markAsyncSequenceRootTask(): void { } else { pendingOperations.delete(executionAsyncId()); } + // Renders are also a good time to release expired history. + sweepAgingOperations(); } } diff --git a/packages/react-server/src/ReactFlightServerConfigDebugNoop.js b/packages/react-server/src/ReactFlightServerConfigDebugNoop.js index e435929114..6eb40e585e 100644 --- a/packages/react-server/src/ReactFlightServerConfigDebugNoop.js +++ b/packages/react-server/src/ReactFlightServerConfigDebugNoop.js @@ -12,6 +12,7 @@ import type {AsyncSequence} from './ReactFlightAsyncSequence'; // Exported for runtimes that don't support Promise instrumentation for async debugging. export function initAsyncDebugInfo(): void {} export function markAsyncSequenceRootTask(): void {} +export function markAsyncSequenceRequest(request: any): void {} export function getCurrentAsyncSequence(): null | AsyncSequence { return null; } diff --git a/packages/react-server/src/__tests__/ReactFlightAsyncDebugInfoExpiration-test.js b/packages/react-server/src/__tests__/ReactFlightAsyncDebugInfoExpiration-test.js index 661b2c9989..f1b97ba0ec 100644 --- a/packages/react-server/src/__tests__/ReactFlightAsyncDebugInfoExpiration-test.js +++ b/packages/react-server/src/__tests__/ReactFlightAsyncDebugInfoExpiration-test.js @@ -11,6 +11,8 @@ import {patchSetImmediate} from '../../../../scripts/jest/patchSetImmediate'; let React; +let ReactServer; +let cache; let ReactServerDOMServer; let ReactServerDOMClient; let Stream; @@ -42,8 +44,9 @@ describe('ReactFlightAsyncDebugInfoExpiration', () => { jest.mock('react-server-dom-webpack/server', () => jest.requireActual('react-server-dom-webpack/server.node'), ); - require('react'); + ReactServer = require('react'); ReactServerDOMServer = require('react-server-dom-webpack/server'); + cache = ReactServer.cache; jest.resetModules(); jest.useRealTimers(); @@ -160,4 +163,132 @@ describe('ReactFlightAsyncDebugInfoExpiration', () => { expect(getDebugInfo(result).filter(entry => entry.awaited)).toEqual([]); } }); + + it('keeps history alive while a debug channel is still open', async () => { + // A bidirectional debug channel lets the client ask for more debug info + // after the response is done. The main stream closes but this request can + // still walk the graph, so its history has to stay put. + const debugChannel = new Stream.Duplex({ + ...streamOptions, + read() {}, + write(chunk, encoding, callback) { + callback(); + }, + }); + async function Init() { + await delay(1); + return 'ok'; + } + const initStream = ReactServerDOMServer.renderToPipeableStream( + , + {}, + {filterStackFrame, debugChannel}, + ); + const initReadable = new Stream.PassThrough(streamOptions); + const initResult = ReactServerDOMClient.createFromNodeStream(initReadable, { + moduleMap: {}, + moduleLoading: {}, + }); + initStream.pipe(initReadable); + expect(await initResult).toBe('ok'); + await finishLoadingStream(initReadable); + + async function fetchCachedData() { + await delay(5); + return 'hello'; + } + const cachedStartTime = + // $FlowFixMe[prop-missing] + performance.timeOrigin + performance.now(); + const cachedData = fetchCachedData(); + await cachedData; + + churnPastRetention(); + + async function Component() { + return 'data:' + (await cachedData); + } + const stream = ReactServerDOMServer.renderToPipeableStream( + , + {}, + { + filterStackFrame, + startTime: cachedStartTime, + }, + ); + const readable = new Stream.PassThrough(streamOptions); + const result = ReactServerDOMClient.createFromNodeStream(readable, { + moduleMap: {}, + moduleLoading: {}, + }); + stream.pipe(readable); + expect(await result).toBe('data:hello'); + await finishLoadingStream(readable); + + if ( + __DEV__ && + gate( + flags => + flags.enableComponentPerformanceTrack && flags.enableAsyncDebugInfo, + ) + ) { + // The fetch resolved after the first request started, so that request + // pinned it through the sweeps and it's still here. + const awaitedNames = getDebugInfo(result) + .filter(entry => entry.awaited) + .map(entry => entry.awaited.name); + expect(awaitedNames).toContain('setTimeout'); + } + }); + + it('keeps debug info intact when history expires during the render', async () => { + const getData = cache(async function getData(text) { + await delay(1); + return text.toUpperCase(); + }); + + // History expires while this render is still running. Anything the + // render itself can still emit must survive the sweeps. + async function Child() { + const greeting = await getData('hi'); + return greeting + ', Seb'; + } + + async function Component() { + await getData('hi'); + churnPastRetention(); + return ; + } + + const stream = ReactServerDOMServer.renderToPipeableStream( + , + {}, + { + filterStackFrame, + }, + ); + const readable = new Stream.PassThrough(streamOptions); + const result = ReactServerDOMClient.createFromNodeStream(readable, { + moduleMap: {}, + moduleLoading: {}, + }); + stream.pipe(readable); + expect(await result).toBe('HI, Seb'); + await finishLoadingStream(readable); + + if ( + __DEV__ && + gate( + flags => + flags.enableComponentPerformanceTrack && flags.enableAsyncDebugInfo, + ) + ) { + // The cached entry resolved after this request started, so its awaited + // I/O info must survive the sweeps. + const awaitedNames = getDebugInfo(result) + .filter(entry => entry.awaited) + .map(entry => entry.awaited.name); + expect(awaitedNames).toContain('delay'); + } + }); });