Skip to content

bug: using structuredClone with ReadableStream prevents process from exiting #44985

Description

@KhafraDev

Version

v18.10.0

Platform

Microsoft Windows NT 10.0.19043.0 x64

Subsystem

No response

What steps will reproduce the bug?

const rs = new ReadableStream({
  start (controller) {
    controller.enqueue(new Uint8Array([65]))
    controller.close()
  }
})

const cloned = structuredClone(rs, { transfer: [rs] })

After this the process will indefinitely stay open.

How often does it reproduce? Is there a required condition?

always

What is the expected behavior?

the process should exit

What do you see instead?

the process stays open indefinitely

Additional information

Exporting readableStreamTee would also solve my use case.

function readableStreamTee(stream, cloneForBranch2) {

Activity

  1. KhafraDev commented on Oct 13, 2022

    @KhafraDev
    MemberAuthor

    A workaround is to tee the cloned stream.

  2. ronag commented on Oct 14, 2022

    @ronag
    Member
  3. ronag commented on Oct 14, 2022

    @ronag
    Member
  4. KhafraDev commented on Oct 18, 2022

    @KhafraDev
    MemberAuthor

    It looks like 2 ports aren't unref'd

    [
      MessagePort [EventTarget] {
        active: true,
        refed: true,
        [Symbol(kEvents)]: SafeMap(4) [Map] {        
          'newListener' => [Object],
          'removeListener' => [Object],
          'message' => [Object],
          'messageerror' => [Object]
        },
        [Symbol(events.maxEventTargetListeners)]: 10,
        [Symbol(events.maxEventTargetListenersWarned)]: false,
        [Symbol(kNewListener)]: [Function (anonymous)],
        [Symbol(kRemoveListener)]: [Function (anonymous)],
        [Symbol(nodejs.internal.kCurrentlyReceivingPorts)]: undefined,
        [Symbol(khandlers)]: SafeMap(2) [Map] {
          'message' => [Function],
          'messageerror' => [Function]
        }
      },
      MessagePort [EventTarget] {
        active: true,
        refed: true,
        [Symbol(kEvents)]: SafeMap(4) [Map] {
          'newListener' => [Object],
          'removeListener' => [Object],
          'message' => [Object],
          'messageerror' => [Object]
        },
        [Symbol(events.maxEventTargetListeners)]: 10,
        [Symbol(events.maxEventTargetListenersWarned)]: false,
        [Symbol(kNewListener)]: [Function (anonymous)],
        [Symbol(kRemoveListener)]: [Function (anonymous)],
        [Symbol(nodejs.internal.kCurrentlyReceivingPorts)]: undefined,
        [Symbol(khandlers)]: SafeMap(2) [Map] {
          'message' => [Function],
          'messageerror' => [Function]
        }
      }
    ]

    which from my best guess are these:

    this[kState].transfer.port1 = port1;
    this[kState].transfer.port2 = port2;

  5. ronag commented on Dec 10, 2023

    @ronag
    Member
  6. tsctx commented on Dec 10, 2023

    @tsctx
    Member

    @ronag
    I have been able to solve this problem.

    const { kTransfer } = require("internal/worker/js_transferable");
    
    const { readableStreamPipeTo } = require("internal/webstreams/readablestream");
    
    const { setPromiseHandled, kState } = require("internal/webstreams/util");
    
    const {
      CrossRealmTransformWritableSink,
    } = require("internal/webstreams/transfer");
    
    function newCrossRealmWritableSink(readable, port) {
      const source = new CrossRealmTransformWritableSink(port);
      const writable = new WritableStream(source);
      const promise = readableStreamPipeTo(readable, writable, false, false, false);
      setPromiseHandled(promise);
      return {
        writable,
        source,
        promise,
      };
    }
    
    const registry = new FinalizationRegistry(({ source }) => {
      source.close();
    });
    
    ReadableStream.prototype[kTransfer] = function () {
      if (this.locked) {
        this[kState].transfer.port1?.close();
        this[kState].transfer.port1 = undefined;
        this[kState].transfer.port2 = undefined;
        throw new DOMException(
          "Cannot transfer a locked ReadableStream",
          "DataCloneError"
        );
      }
    
      const { port1, port2 } = this[kState].transfer;
      this[kState].transfer.port2 = undefined;
    
      const { writable, promise, source } = newCrossRealmWritableSink(this, port1);
    
      this[kState].transfer.writable = writable;
      this[kState].transfer.promise = promise;
    
      registry.register(port2, { source });
    
      return {
        data: { port: port2 },
        deserializeInfo:
          "internal/webstreams/readablestream:TransferredReadableStream",
      };
    };
    
    const stack = [
      new Uint8Array([65]),
      new Uint8Array([65]),
      new Uint8Array([65]),
    ];
    
    const rs = new ReadableStream({
      pull(controller) {
        const data = stack.shift();
        if (data) {
          controller.enqueue(data);
        } else {
          controller.close();
        }
      },
    });
    
    const cloned = structuredClone(rs, { transfer: [rs] });
  7. tsctx commented on Dec 10, 2023

    @tsctx
    Member

    @ronag
    This problem occurs not only with ReadableStream but also with WritableStream.
    I am stuck on this point.
    Could you please work on this instead of me?

  8. ronag commented on Dec 10, 2023

    @ronag
    Member

    Sorry. Webstreams are beyond my expertise and frankly my interest...

  9. tsctx commented on Dec 10, 2023

    @tsctx
    Member

    I apologize for the inconvenience that I have caused you.

  10. ronag commented on Dec 10, 2023

    @ronag
    Member

    No worries. No inconvenience at all. 🤗

  11. added a commit that references this issue on Dec 24, 2023
  12. mcollina commented on Jan 19, 2024

    @mcollina
    SponsorMember

    Reopening after #51491

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions