File
Blob: src/node/internal/internal_tls_jsstream.ts
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | // |
| 5 | // Copyright Joyent, Inc. and other Node contributors. |
| 6 | // |
| 7 | // Permission is hereby granted, free of charge, to any person obtaining a |
| 8 | // copy of this software and associated documentation files (the |
| 9 | // "Software"), to deal in the Software without restriction, including |
| 10 | // without limitation the rights to use, copy, modify, merge, publish, |
| 11 | // distribute, sublicense, and/or sell copies of the Software, and to permit |
| 12 | // persons to whom the Software is furnished to do so, subject to the |
| 13 | // following conditions: |
| 14 | // |
| 15 | // The above copyright notice and this permission notice shall be included |
| 16 | // in all copies or substantial portions of the Software. |
| 17 | // |
| 18 | // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS |
| 19 | // OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF |
| 20 | // MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN |
| 21 | // NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, |
| 22 | // DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR |
| 23 | // OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE |
| 24 | // USE OR OTHER DEALINGS IN THE SOFTWARE. |
| 25 | |
| 26 | /* eslint-disable @typescript-eslint/no-redundant-type-constituents */ |
| 27 | |
| 28 | import { notStrictEqual } from 'node-internal:internal_assert'; |
| 29 | import { Socket } from 'node-internal:internal_net'; |
| 30 | import { ERR_STREAM_WRAP } from 'node-internal:internal_errors'; |
| 31 | import { Duplex, toBYOBWeb } from 'node-internal:streams_duplex'; |
| 32 | import type { |
| 33 | SocketInfo, |
| 34 | Writer, |
| 35 | Socket as CloudflareSocket, |
| 36 | } from 'node-internal:sockets'; |
| 37 | |
| 38 | const kCurrentWriteRequest = Symbol('kCurrentWriteRequest'); |
| 39 | const kCurrentShutdownRequest = Symbol('kCurrentShutdownRequest'); |
| 40 | const kPendingShutdownRequest = Symbol('kPendingShutdownRequest'); |
| 41 | const kPendingClose = Symbol('kPendingClose'); |
| 42 | |
| 43 | /* This class serves as a wrapper for when the C++ side of Node wants access |
| 44 | * to a standard JS stream. For example, TLS or HTTP do not operate on network |
| 45 | * resources conceptually, although that is the common case and what we are |
| 46 | * optimizing for; in theory, they are completely composable and can work with |
| 47 | * any stream resource they see. |
| 48 | * |
| 49 | * For the common case, i.e. a TLS socket wrapping around a net.Socket, we |
| 50 | * can skip going through the JS layer and let TLS access the raw C++ handle |
| 51 | * of a net.Socket. The flipside of this is that, to maintain composability, |
| 52 | * we need a way to create "fake" net.Socket instances that call back into a |
| 53 | * "real" JavaScript stream. JSStreamSocket is exactly this. |
| 54 | */ |
| 55 | export class JSStreamSocket extends Socket { |
| 56 | stream: Duplex; |
| 57 | [kCurrentWriteRequest]: null | unknown; |
| 58 | [kCurrentShutdownRequest]: null | unknown; |
| 59 | [kPendingShutdownRequest]: null | unknown; |
| 60 | [kPendingClose]: boolean; |
| 61 | |
| 62 | constructor(stream: Duplex) { |
| 63 | // eslint-disable-next-line @typescript-eslint/no-invalid-void-type |
| 64 | const closePromise = Promise.withResolvers<void>(); |
| 65 | const openPromise = Promise.withResolvers<SocketInfo>(); |
| 66 | |
| 67 | const webStream = toBYOBWeb(stream); |
| 68 | Object.assign(webStream.writable, { |
| 69 | // eslint-disable-next-line @typescript-eslint/require-await |
| 70 | write: async (data: string | ArrayBufferView): Promise<void> => { |
| 71 | stream.write(data); |
| 72 | }, |
| 73 | closed: closePromise.promise, |
| 74 | releaseLock: async (): Promise<void> => {}, |
| 75 | }); |
| 76 | const handle: Socket['_handle'] = { |
| 77 | reading: true, |
| 78 | bytesRead: 0, |
| 79 | bytesWritten: 0, |
| 80 | socket: { |
| 81 | startTls(): CloudflareSocket { |
| 82 | throw new Error( |
| 83 | 'startTls() should not be called for a duplex stream' |
| 84 | ); |
| 85 | }, |
| 86 | upgraded: false, |
| 87 | secureTransport: 'off', |
| 88 | closed: closePromise.promise, |
| 89 | close: async (): Promise<void> => { |
| 90 | queueMicrotask(() => { |
| 91 | closePromise.resolve(); |
| 92 | }); |
| 93 | return closePromise.promise; |
| 94 | }, |
| 95 | opened: openPromise.promise, |
| 96 | readable: webStream.readable, |
| 97 | writable: webStream.writable as unknown as Writer, |
| 98 | }, |
| 99 | // eslint-disable-next-line @typescript-eslint/no-unsafe-argument |
| 100 | reader: new ReadableStreamBYOBReader(webStream.readable), |
| 101 | writer: new WritableStreamDefaultWriter<unknown>(webStream.writable), |
| 102 | options: { |
| 103 | host: '0.0.0.0', |
| 104 | port: 0, |
| 105 | addressType: 4, |
| 106 | }, |
| 107 | }; |
| 108 | |
| 109 | stream.pause(); |
| 110 | stream.on('error', (err) => this.emit('error', err)); |
| 111 | const ondata = (chunk: string | Buffer): void => { |
| 112 | // eslint-disable-next-line @typescript-eslint/no-unnecessary-boolean-literal-compare |
| 113 | if (typeof chunk === 'string' || stream.readableObjectMode === true) { |
| 114 | // Make sure that no further `data` events will happen. |
| 115 | stream.pause(); |
| 116 | stream.removeListener('data', ondata); |
| 117 | |
| 118 | this.emit('error', new ERR_STREAM_WRAP()); |
| 119 | return; |
| 120 | } |
| 121 | |
| 122 | // TODO(soon): We need to trigger read() result for _handle.reader.read(buf) call |
| 123 | // in node:net. |
| 124 | }; |
| 125 | stream.on('data', ondata); |
| 126 | stream.once('end', () => { |
| 127 | closePromise.resolve(); |
| 128 | }); |
| 129 | // Some `Stream` don't pass `hasError` parameters when closed. |
| 130 | stream.once('close', () => { |
| 131 | // Errors emitted from `stream` have also been emitted to this instance |
| 132 | // so that we don't pass errors to `destroy()` again. |
| 133 | this.destroy(); |
| 134 | }); |
| 135 | |
| 136 | super({ handle }); |
| 137 | this.stream = stream; |
| 138 | this[kCurrentWriteRequest] = null; |
| 139 | this[kCurrentShutdownRequest] = null; |
| 140 | this[kPendingShutdownRequest] = null; |
| 141 | this[kPendingClose] = false; |
| 142 | this.readable = stream.readable; |
| 143 | this.writable = stream.writable; |
| 144 | |
| 145 | // eslint-disable-next-line @typescript-eslint/no-floating-promises |
| 146 | handle.socket.closed.then(this.doClose.bind(this)); |
| 147 | |
| 148 | openPromise.resolve({}); |
| 149 | |
| 150 | // Start reading. |
| 151 | this.read(0); |
| 152 | } |
| 153 | |
| 154 | isClosing(): boolean { |
| 155 | return !this.readable || !this.writable; |
| 156 | } |
| 157 | |
| 158 | readStart(): number { |
| 159 | this.stream.resume(); |
| 160 | return 0; |
| 161 | } |
| 162 | |
| 163 | readStop(): number { |
| 164 | this.stream.pause(); |
| 165 | return 0; |
| 166 | } |
| 167 | |
| 168 | doShutdown(req: unknown): number { |
| 169 | // TODO(addaleax): It might be nice if we could get into a state where |
| 170 | // DoShutdown() is not called on streams while a write is still pending. |
| 171 | // |
| 172 | // Currently, the only part of the code base where that happens is the |
| 173 | // TLS implementation, which calls both DoWrite() and DoShutdown() on the |
| 174 | // underlying network stream inside of its own DoShutdown() method. |
| 175 | // Working around that on the native side is not quite trivial (yet?), |
| 176 | // so for now that is supported here. |
| 177 | |
| 178 | if (this[kCurrentWriteRequest] !== null) { |
| 179 | this[kPendingShutdownRequest] = req; |
| 180 | return 0; |
| 181 | } |
| 182 | |
| 183 | this[kCurrentShutdownRequest] = req; |
| 184 | |
| 185 | if (this[kPendingClose]) { |
| 186 | // If doClose is pending, the stream & this._handle are gone. We can't do |
| 187 | // anything. doClose will call finishShutdown with ECANCELED for us shortly. |
| 188 | return 0; |
| 189 | } |
| 190 | |
| 191 | const handle = this._handle; |
| 192 | notStrictEqual(handle, null); |
| 193 | |
| 194 | queueMicrotask(() => { |
| 195 | // Ensure that write is dispatched asynchronously. |
| 196 | this.stream.end(); |
| 197 | }); |
| 198 | return 0; |
| 199 | } |
| 200 | |
| 201 | doClose(): void { |
| 202 | this[kPendingClose] = true; |
| 203 | |
| 204 | const handle = this._handle; |
| 205 | |
| 206 | // When sockets of the "net" module destroyed, they will call |
| 207 | // `this._handle.close()` which will also emit EOF if not emitted before. |
| 208 | // This feature makes sockets on the other side emit "end" and "close" |
| 209 | // even though we haven't called `end()`. As `stream` are likely to be |
| 210 | // instances of `net.Socket`, calling `stream.destroy()` manually will |
| 211 | // avoid issues that don't properly close wrapped connections. |
| 212 | this.stream.destroy(); |
| 213 | |
| 214 | queueMicrotask(() => { |
| 215 | notStrictEqual(handle, null); |
| 216 | this[kPendingClose] = false; |
| 217 | }); |
| 218 | } |
| 219 | } |