Skip to content
Merged
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
101 changes: 79 additions & 22 deletions packages/runtime-playground/src/preview-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,18 @@ function previewProxyServer(target: URL, routes: InternalPreviewRouteRegistry):
return
}

upstreamQueue(() => proxyPreviewRequest(target, incoming, outgoing)).catch((error: Error) => writeProxyError(outgoing, error))
upstreamQueue(
() => proxyPreviewRequest(target, incoming, outgoing),
() => incoming.aborted || incoming.destroyed || outgoing.destroyed,
(cancel) => {
incoming.once("aborted", cancel)
outgoing.once("close", cancel)
return () => {
incoming.off("aborted", cancel)
outgoing.off("close", cancel)
}
},
).catch((error: Error) => writeProxyError(outgoing, error))
})
}

Expand Down Expand Up @@ -203,6 +214,8 @@ function proxyPreviewRequest(target: URL, incoming: IncomingMessage, outgoing: S
outgoing.destroy(error)
settle()
})
response.on("end", settle)
response.on("close", settle)
outgoing.on("finish", settle)
outgoing.on("close", settle)
if (bodyTransform) {
Expand All @@ -218,6 +231,7 @@ function proxyPreviewRequest(target: URL, incoming: IncomingMessage, outgoing: S
writeProxyError(outgoing, error)
settle()
})
incoming.on("aborted", abortUpstream)
incoming.on("error", () => {
abortUpstream()
})
Expand Down Expand Up @@ -258,37 +272,80 @@ function previewProxyRequestTarget(incoming: IncomingMessage, target: URL): Prev
return { upstreamHost: target.host, visibleHost: authority.host, path: rawUrl, port: authority.port || "80", protocol: "http:", rewriteTargetOrigin: true }
}

function createPreviewProxyQueue(): (task: () => Promise<void>) => Promise<void> {
function createPreviewProxyQueue(): (
task: () => Promise<void>,
isCanceled: () => boolean,
observeCancellation: (cancel: () => void) => () => void,
) => Promise<void> {
let active = false
const pending: Array<() => void> = []

const acquire = async () => {
if (!active) {
active = true
return
}

await new Promise<void>((resolve) => pending.push(resolve))
interface PendingRequest {
task: () => Promise<void>
isCanceled: () => boolean
stopObservingCancellation: () => void
resolve: () => void
reject: (error: unknown) => void
}

const release = () => {
const next = pending.shift()
if (next) {
next()
return
const pending: PendingRequest[] = []

function release(): void {
let next = pending.shift()
while (next) {
next.stopObservingCancellation()
if (!next.isCanceled()) {
run(next)
return
}
next.resolve()
next = pending.shift()
}

active = false
}

return async (task) => {
await acquire()
function run(request: PendingRequest): void {
request.stopObservingCancellation()
if (request.isCanceled()) {
request.resolve()
release()
return
}
let task: Promise<void>
try {
await task()
} finally {
task = request.task()
} catch (error) {
request.reject(error)
release()
return
}
task.then(request.resolve, request.reject).finally(release)
}

return (task, isCanceled, observeCancellation) => new Promise<void>((resolve, reject) => {
const request: PendingRequest = {
task,
isCanceled,
stopObservingCancellation: () => {},
resolve,
reject,
}
const cancel = () => {
const index = pending.indexOf(request)
if (index === -1) {
return
}
pending.splice(index, 1)
request.stopObservingCancellation()
resolve()
}
request.stopObservingCancellation = observeCancellation(cancel)

if (active) {
pending.push(request)
return
}

active = true
run(request)
})
}

async function listenPreviewProxy(proxy: PreviewProxyServer, port: number, bind: string): Promise<void> {
Expand Down
59 changes: 59 additions & 0 deletions tests/browser-callback-materialization-contracts.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -259,6 +259,65 @@ try {
await closeHttpServer(hangingTargetServer)
}

let releaseFirstQueuedRequest: (() => void) | undefined
let firstQueuedRequestReached: (() => void) | undefined
const firstQueuedRequest = new Promise<void>((resolve) => {
firstQueuedRequestReached = resolve
})
const queuedTargetRequests: string[] = []
const queuedTargetServer = createServer((request, response) => {
queuedTargetRequests.push(request.url ?? "")
if (request.url === "/first") {
firstQueuedRequestReached?.()
new Promise<void>((resolve) => {
releaseFirstQueuedRequest = resolve
}).then(() => response.end("first"))
return
}
response.end("third")
})
const queuedTargetServerUrl = await listenLocalHttpServer(queuedTargetServer)
const queuedAbortProxy = await withPreviewProxy({
playground: { async run() { return { text: "" } } },
serverUrl: queuedTargetServerUrl,
async [Symbol.asyncDispose]() {},
} satisfies PlaygroundCliServer, 0)
let secondQueuedRequestReached: (() => void) | undefined
const secondQueuedRequest = new Promise<void>((resolve) => {
secondQueuedRequestReached = resolve
})
let secondQueuedRequestCanceled: (() => void) | undefined
const secondQueuedCancellation = new Promise<void>((resolve) => {
secondQueuedRequestCanceled = resolve
})
const stopObservingQueuedRequest = queuedAbortProxy.previewRoutes?.add((request, response) => {
if (request.url === "/second") {
secondQueuedRequestReached?.()
response.once("close", () => secondQueuedRequestCanceled?.())
}
return false
})
try {
const firstRequest = fetch(`${queuedAbortProxy.serverUrl}/first`)
await firstQueuedRequest

const controller = new AbortController()
const canceledRequest = fetch(`${queuedAbortProxy.serverUrl}/second`, { signal: controller.signal }).catch((error) => error)
await secondQueuedRequest
controller.abort()
await canceledRequest
await secondQueuedCancellation

releaseFirstQueuedRequest?.()
assert.equal(await (await firstRequest).text(), "first")
assert.equal(await (await fetch(`${queuedAbortProxy.serverUrl}/third`)).text(), "third")
assert.deepEqual(queuedTargetRequests, ["/first", "/third"], "a canceled queued request must not open an upstream request")
} finally {
stopObservingQueuedRequest?.()
await queuedAbortProxy[Symbol.asyncDispose]()
await closeHttpServer(queuedTargetServer)
}

const phase = materializationPhaseResult({
phase: "persist-browser-artifacts",
status: "completed",
Expand Down
Loading