From b5bb3c0e488c1c06b0098814f3573e4c1935a7c8 Mon Sep 17 00:00:00 2001 From: Mike Shepherd Date: Wed, 15 Nov 2023 19:04:09 +0000 Subject: [PATCH 1/3] Remove JS TLS socket event listeners on finalise --- .../scala/fs2/io/internal/facade/events.scala | 20 +++++++++ .../fs2/io/net/tls/TLSContextPlatform.scala | 45 ++++++++----------- .../fs2/io/net/tls/TLSSocketPlatform.scala | 5 ++- 3 files changed, 41 insertions(+), 29 deletions(-) diff --git a/io/js/src/main/scala/fs2/io/internal/facade/events.scala b/io/js/src/main/scala/fs2/io/internal/facade/events.scala index 8cdb1041db..d967671f93 100644 --- a/io/js/src/main/scala/fs2/io/internal/facade/events.scala +++ b/io/js/src/main/scala/fs2/io/internal/facade/events.scala @@ -80,6 +80,26 @@ private[io] object EventEmitter { })(fn => F.delay(eventTarget.removeListener(eventName, fn))) .void + def registerOneTimeListener0[F[_]](eventName: String, dispatcher: Dispatcher[F])( + listener: => F[Unit] + )(implicit F: Sync[F]): Resource[F, Unit] = Resource + .make(F.delay { + val fn: js.Function0[Unit] = () => dispatcher.unsafeRunAndForget(listener) + eventTarget.once(eventName, fn) + fn + })(fn => F.delay(eventTarget.removeListener(eventName, fn))) + .void + + def registerOneTimeListener[F[_], E](eventName: String, dispatcher: Dispatcher[F])( + listener: E => F[Unit] + )(implicit F: Sync[F]): Resource[F, Unit] = Resource + .make(F.delay { + val fn: js.Function1[E, Unit] = e => dispatcher.unsafeRunAndForget(listener(e)) + eventTarget.once(eventName, fn) + fn + })(fn => F.delay(eventTarget.removeListener(eventName, fn))) + .void + def registerOneTimeListener[F[_], E](eventName: String)( listener: E => Unit )(implicit F: Sync[F]): F[Option[F[Unit]]] = F.delay { diff --git a/io/js/src/main/scala/fs2/io/net/tls/TLSContextPlatform.scala b/io/js/src/main/scala/fs2/io/net/tls/TLSContextPlatform.scala index a506e02f58..4ac4395303 100644 --- a/io/js/src/main/scala/fs2/io/net/tls/TLSContextPlatform.scala +++ b/io/js/src/main/scala/fs2/io/net/tls/TLSContextPlatform.scala @@ -76,19 +76,16 @@ private[tls] trait TLSContextCompanionPlatform { self: TLSContext.type => options.rejectUnauthorized = false options.enableTrace = logger != TLSLogger.Disabled options.socket = sock - val tlsSock = facade.tls.connect(options) - tlsSock.once( - "secureConnect", - () => seqDispatcher.unsafeRunAndForget(handshake.complete(Either.unit)) - ) - tlsSock.once[js.Error]( - "error", - e => - seqDispatcher.unsafeRunAndForget( - handshake.complete(Left(new js.JavaScriptException(e))) - ) - ) - tlsSock + for { + tlsSock <- Resource.eval(F.delay(facade.tls.connect(options))) + _ <- tlsSock.registerOneTimeListener0[F]("secureConnect", seqDispatcher)( + handshake.complete(Either.unit).void + ) + _ <- tlsSock.registerOneTimeListener[F, js.Error]("error", seqDispatcher)( + e => handshake.complete(Left(new js.JavaScriptException(e))).void + ) + } yield tlsSock + } ) .evalTap(_ => handshake.get.rethrow) @@ -105,10 +102,9 @@ private[tls] trait TLSContextCompanionPlatform { self: TLSContext.type => options.rejectUnauthorized = false options.enableTrace = logger != TLSLogger.Disabled options.isServer = true - val tlsSock = new facade.tls.TLSSocket(sock, options) - tlsSock.once( - "secure", - { () => + for { + tlsSock <- Resource.eval(F.delay(new facade.tls.TLSSocket(sock, options))) + _ <- tlsSock.registerOneTimeListener0[F]("secure", seqDispatcher) { val requestCert = options.requestCert.getOrElse(false) val rejectUnauthorized = options.rejectUnauthorized.getOrElse(true) val result = @@ -117,17 +113,12 @@ private[tls] trait TLSContextCompanionPlatform { self: TLSContext.type => .map(e => new JavaScriptSSLException(js.JavaScriptException(e))) .toLeft(()) else Either.unit - seqDispatcher.unsafeRunAndForget(verifyError.complete(result)) + verifyError.complete(result).void } - ) - tlsSock.once[js.Error]( - "error", - e => - seqDispatcher.unsafeRunAndForget( - verifyError.complete(Left(new js.JavaScriptException(e))) - ) - ) - tlsSock + _ <- tlsSock.registerOneTimeListener[F, js.Error]("error", seqDispatcher)( + e => verifyError.complete(Left(new js.JavaScriptException(e))).void + ) + } yield tlsSock } ) .evalTap(_ => verifyError.get.rethrow) diff --git a/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala b/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala index 050f9d8353..5abc2c94f9 100644 --- a/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala +++ b/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala @@ -38,16 +38,17 @@ private[tls] trait TLSSocketCompanionPlatform { self: TLSSocket.type => private[tls] def forAsync[F[_]]( socket: Socket[F], - upgrade: fs2.io.Duplex => facade.tls.TLSSocket + upgrade: fs2.io.Duplex => Resource[F, facade.tls.TLSSocket] )(implicit F: Async[F]): Resource[F, TLSSocket[F]] = for { duplexOut <- mkDuplex(socket.reads) (duplex, out) = duplexOut _ <- out.through(socket.writes).compile.drain.background + upgraded <- upgrade(duplex) tlsSockReadable <- suspendReadableAndRead( destroyIfNotEnded = false, destroyIfCanceled = false - )(upgrade(duplex)) + )(upgraded) (tlsSock, readable) = tlsSockReadable readStream <- SuspendedStream(readable) } yield new AsyncTLSSocket( From acb97d3d0e97e697adbc3293526588d3f4533f3c Mon Sep 17 00:00:00 2001 From: Mike Shepherd Date: Wed, 15 Nov 2023 21:37:39 +0000 Subject: [PATCH 2/3] Make TLS socket creation still happen on MicrotaskExecutor --- io/js/src/main/scala/fs2/io/ioplatform.scala | 24 +++++++++++++++---- .../fs2/io/net/tls/TLSSocketPlatform.scala | 5 ++-- 2 files changed, 22 insertions(+), 7 deletions(-) diff --git a/io/js/src/main/scala/fs2/io/ioplatform.scala b/io/js/src/main/scala/fs2/io/ioplatform.scala index 011f770764..4faff25019 100644 --- a/io/js/src/main/scala/fs2/io/ioplatform.scala +++ b/io/js/src/main/scala/fs2/io/ioplatform.scala @@ -59,22 +59,38 @@ private[fs2] trait ioplatform { destroyIfNotEnded: Boolean = true, destroyIfCanceled: Boolean = true )(thunk: => R)(implicit F: Async[F]): Resource[F, (R, Stream[F, Byte])] = + suspendResourceReadableAndRead(destroyIfNotEnded, destroyIfCanceled)( + Resource.eval(F.delay(thunk)) + ) + + /** Suspends the creation of a `Readable` and a `Stream` that reads all bytes from that `Readable`. + * + * Accepts a Resource to allow finalizers to be run for the acquired Readable. + * Be aware that the readable may have been destroyed before the finalizers for the resource are run. + * N.B. This has different semantics to running the resource prior to passing it here + * as it will be run on a different executor (see implementation note) + */ + def suspendResourceReadableAndRead[F[_], R <: Readable]( + destroyIfNotEnded: Boolean = true, + destroyIfCanceled: Boolean = true + )(thunk: Resource[F, R])(implicit F: Async[F]): Resource[F, (R, Stream[F, Byte])] = (for { dispatcher <- Dispatcher.sequential[F] channel <- Channel.unbounded[F, Unit].toResource error <- F.deferred[Throwable].toResource readableResource = for { - readable <- Resource.makeCase(F.delay(thunk)) { - case (readable, Resource.ExitCase.Succeeded) => + readable <- thunk + _ <- Resource.makeCase(F.unit) { + case (_, Resource.ExitCase.Succeeded) => F.delay { if (!readable.readableEnded & destroyIfNotEnded) readable.destroy() } - case (readable, Resource.ExitCase.Errored(_)) => + case (_, Resource.ExitCase.Errored(_)) => // tempting, but don't propagate the error! // that would trigger a unhandled Node.js error that circumvents FS2/CE error channels F.delay(readable.destroy()) - case (readable, Resource.ExitCase.Canceled) => + case (_, Resource.ExitCase.Canceled) => if (destroyIfCanceled) F.delay(readable.destroy()) else diff --git a/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala b/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala index 5abc2c94f9..feb4f10f1d 100644 --- a/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala +++ b/io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala @@ -44,11 +44,10 @@ private[tls] trait TLSSocketCompanionPlatform { self: TLSSocket.type => duplexOut <- mkDuplex(socket.reads) (duplex, out) = duplexOut _ <- out.through(socket.writes).compile.drain.background - upgraded <- upgrade(duplex) - tlsSockReadable <- suspendReadableAndRead( + tlsSockReadable <- suspendResourceReadableAndRead( destroyIfNotEnded = false, destroyIfCanceled = false - )(upgraded) + )(upgrade(duplex)) (tlsSock, readable) = tlsSockReadable readStream <- SuspendedStream(readable) } yield new AsyncTLSSocket( From fa6820384b6cbde4a979f1872d8a7dc0d5b294b0 Mon Sep 17 00:00:00 2001 From: Mike Shepherd Date: Thu, 16 Nov 2023 09:46:06 +0000 Subject: [PATCH 3/3] Use Resource.onFinalizeCase --- io/js/src/main/scala/fs2/io/ioplatform.scala | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/io/js/src/main/scala/fs2/io/ioplatform.scala b/io/js/src/main/scala/fs2/io/ioplatform.scala index 4faff25019..ebc7af035a 100644 --- a/io/js/src/main/scala/fs2/io/ioplatform.scala +++ b/io/js/src/main/scala/fs2/io/ioplatform.scala @@ -80,17 +80,17 @@ private[fs2] trait ioplatform { error <- F.deferred[Throwable].toResource readableResource = for { readable <- thunk - _ <- Resource.makeCase(F.unit) { - case (_, Resource.ExitCase.Succeeded) => + _ <- Resource.onFinalizeCase[F] { + case Resource.ExitCase.Succeeded => F.delay { if (!readable.readableEnded & destroyIfNotEnded) readable.destroy() } - case (_, Resource.ExitCase.Errored(_)) => + case Resource.ExitCase.Errored(_) => // tempting, but don't propagate the error! // that would trigger a unhandled Node.js error that circumvents FS2/CE error channels F.delay(readable.destroy()) - case (_, Resource.ExitCase.Canceled) => + case Resource.ExitCase.Canceled => if (destroyIfCanceled) F.delay(readable.destroy()) else