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/ioplatform.scala b/io/js/src/main/scala/fs2/io/ioplatform.scala index 011f770764..ebc7af035a 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.onFinalizeCase[F] { + 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/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..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 @@ -38,13 +38,13 @@ 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 - tlsSockReadable <- suspendReadableAndRead( + tlsSockReadable <- suspendResourceReadableAndRead( destroyIfNotEnded = false, destroyIfCanceled = false )(upgrade(duplex))