6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
64 lines
2.5 KiB
TypeScript
64 lines
2.5 KiB
TypeScript
import type { Readable } from 'node:stream'
|
|
|
|
/**
|
|
* Bridges a Node `Readable` into a WHATWG `ReadableStream` suitable for a `Response`
|
|
* body. Node's built-in `Readable.toWeb` is NOT used: its adapter throws an unhandled
|
|
* `ERR_INVALID_STATE` ("Controller is already closed") when the web stream is cancelled
|
|
* while the Node stream is still flowing — which happens whenever a consumer aborts a
|
|
* live body (a redirect hop cancelling the previous response, a browser cancelling a
|
|
* download mid-transfer). This bridge instead swallows a late enqueue after close and
|
|
* destroys the source on cancel, so cancelling a live body frees its socket cleanly.
|
|
* Source errors — a size-limit overrun, a storage read failing mid-archive — surface as
|
|
* the source's `error` event and reject the read.
|
|
*/
|
|
export function nodeReadableToWebStream(nodeStream: Readable): ReadableStream<Uint8Array> {
|
|
let settled = false
|
|
return new ReadableStream<Uint8Array>({
|
|
start(controller) {
|
|
nodeStream.on('data', (chunk: Buffer) => {
|
|
try {
|
|
// Copy, not a view: the producer may recycle the pooled buffer backing `chunk`
|
|
// after this handler returns, which would corrupt a chunk still queued for a
|
|
// slow consumer. `new Uint8Array(chunk)` allocates a fresh backing buffer.
|
|
controller.enqueue(new Uint8Array(chunk))
|
|
} catch {
|
|
// Controller already closed (consumer cancelled) — stop the source, drop the chunk.
|
|
nodeStream.destroy()
|
|
return
|
|
}
|
|
if ((controller.desiredSize ?? 1) <= 0) nodeStream.pause()
|
|
})
|
|
nodeStream.once('end', () => {
|
|
settled = true
|
|
try {
|
|
controller.close()
|
|
} catch {}
|
|
})
|
|
nodeStream.once('error', (err) => {
|
|
settled = true
|
|
try {
|
|
controller.error(err)
|
|
} catch {}
|
|
})
|
|
// An abort or upstream reset can `destroy()` the source with no `error` event;
|
|
// without this the reader would hang forever. `close` fires after every terminal
|
|
// path, so only act when `end`/`error` didn't already settle the stream.
|
|
nodeStream.once('close', () => {
|
|
if (settled) return
|
|
settled = true
|
|
try {
|
|
controller.error(new Error('Stream closed before completing'))
|
|
} catch {}
|
|
})
|
|
// Start paused so nothing buffers before the consumer pulls (backpressure).
|
|
nodeStream.pause()
|
|
},
|
|
pull() {
|
|
nodeStream.resume()
|
|
},
|
|
cancel(reason) {
|
|
nodeStream.destroy(reason instanceof Error ? reason : undefined)
|
|
},
|
|
})
|
|
}
|