| 250 | } |
| 251 | |
| 252 | function wrapModuleHandle<E extends ExecEncoding>( |
| 253 | runtime: WorkspaceModuleBackendHandle, |
| 254 | backend: string, |
| 255 | id: string, |
| 256 | source: ReadableStream<WorkspaceRuntimeEvent>, |
| 257 | encoding: E | undefined, |
| 258 | resultMayUseSource = true, |
| 259 | ): WorkspaceRuntimeExecHandle<E> { |
| 260 | let claimed: "result" | "stream" | undefined; |
| 261 | let sourceCancelled = false; |
| 262 | let reader: ReadableStreamDefaultReader<WorkspaceRuntimeEvent<E>> | undefined; |
| 263 | let resultReader: ReadableStreamDefaultReader<WorkspaceRuntimeEvent> | undefined; |
| 264 | let resultPromise: Promise<WorkspaceRuntimeResult<E>> | undefined; |
| 265 | const stream = new ReadableStream<WorkspaceRuntimeEvent<E>>( |
| 266 | { |
| 267 | async pull(controller) { |
| 268 | if (claimed === "result") { |
| 269 | controller.error(new Error("runtime handle already consumed by result()")); |
| 270 | return; |
| 271 | } |
| 272 | claimed = "stream"; |
| 273 | reader ??= transformModuleEvents(source, encoding).getReader(); |
| 274 | try { |
| 275 | const next = await reader.read(); |
| 276 | if (next.done) { |
| 277 | reader.releaseLock(); |
| 278 | reader = undefined; |
| 279 | controller.close(); |
| 280 | } else controller.enqueue(next.value); |
| 281 | } catch (error) { |
| 282 | reader?.releaseLock(); |
| 283 | reader = undefined; |
| 284 | controller.error(error); |
| 285 | } |
| 286 | }, |
| 287 | async cancel(reason) { |
| 288 | sourceCancelled = true; |
| 289 | if (reader) { |
| 290 | try { |
| 291 | await reader.cancel(reason); |
| 292 | } finally { |
| 293 | reader.releaseLock(); |
| 294 | reader = undefined; |
| 295 | } |
| 296 | } else await source.cancel(reason); |
| 297 | }, |
| 298 | }, |
| 299 | { highWaterMark: 0 }, |
| 300 | ) as WorkspaceRuntimeExecHandle<E>; |
| 301 | Object.defineProperties(stream, { |
| 302 | id: { value: id, enumerable: false }, |
| 303 | backend: { value: backend, enumerable: false }, |
| 304 | result: { |
| 305 | value: (): Promise<WorkspaceRuntimeResult<E>> => { |
| 306 | if (claimed === "stream") { |
| 307 | throw new Error("runtime handle already streaming: result() and streaming are exclusive"); |
| 308 | } |
| 309 | claimed = "result"; |