Skip to content
File

Blob: src/node/internal/internal_tls_jsstream.ts

typescript220 lines
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 
28import { notStrictEqual } from 'node-internal:internal_assert';
29import { Socket } from 'node-internal:internal_net';
30import { ERR_STREAM_WRAP } from 'node-internal:internal_errors';
31import { Duplex, toBYOBWeb } from 'node-internal:streams_duplex';
32import type {
33 SocketInfo,
34 Writer,
35 Socket as CloudflareSocket,
36} from 'node-internal:sockets';
37 
38const kCurrentWriteRequest = Symbol('kCurrentWriteRequest');
39const kCurrentShutdownRequest = Symbol('kCurrentShutdownRequest');
40const kPendingShutdownRequest = Symbol('kPendingShutdownRequest');
41const 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 */
55export 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}