Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions io/js/src/main/scala/fs2/io/internal/facade/events.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
24 changes: 20 additions & 4 deletions io/js/src/main/scala/fs2/io/ioplatform.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
45 changes: 18 additions & 27 deletions io/js/src/main/scala/fs2/io/net/tls/TLSContextPlatform.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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 =
Expand All @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions io/js/src/main/scala/fs2/io/net/tls/TLSSocketPlatform.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down