MCPcopy Create free account
hub / github.com/cloudflare/computer / wrapModuleHandle

Function wrapModuleHandle

packages/computer/src/runtime/runtime.ts:252–336  ·  view source on GitHub ↗
(
  runtime: WorkspaceModuleBackendHandle,
  backend: string,
  id: string,
  source: ReadableStream<WorkspaceRuntimeEvent>,
  encoding: E | undefined,
  resultMayUseSource = true,
)

Source from the content-addressed store, hash-verified

250}
251
252function 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";

Callers 2

execMethod · 0.85
getExecMethod · 0.85

Calls 4

drainModuleResultFunction · 0.85
cancelMethod · 0.65
getExecMethod · 0.65
killExecMethod · 0.65

Tested by

no test coverage detected