From 8c32153355ae0cdd657fa33797add7369496e238 Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 7 Nov 2022 21:02:13 +0000 Subject: [PATCH 1/5] Introduce `raceOutcomeBoth` --- .../scala/cats/effect/kernel/GenSpawn.scala | 49 +++++++++++++------ 1 file changed, 35 insertions(+), 14 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala index b678561ce4..88b42f7d88 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala @@ -300,6 +300,27 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { def racePair[A, B](fa: F[A], fb: F[B]) : F[Either[(Outcome[F, E, A], Fiber[F, E, B]), (Fiber[F, E, A], Outcome[F, E, B])]] + /** + * Races the evaluation of two fibers, cancels the loser, and returns the [[Outcome]] of both. + * If the race is canceled before one or both participants complete, then then whichever ones + * are incomplete are canceled. + * + * @param fa + * the effect for the first racing fiber + * @param fb + * the effect for the second racing fiber + */ + def raceOutcomeBoth[A, B](fa: F[A], fb: F[B]) + : F[Either[(Outcome[F, E, A], Outcome[F, E, B]), (Outcome[F, E, A], Outcome[F, E, B])]] = + uncancelable { poll => + poll(racePair(fa, fb)).flatMap { + case Left((oca, f)) => + f.cancel.whenA(!oca.isCanceled) *> f.join.map(ocb => Left((oca, ocb))) + case Right((f, ocb)) => + f.cancel.whenA(!ocb.isCanceled) *> f.join.map(oca => Right((oca, ocb))) + } + } + /** * Races the evaluation of two fibers that returns the [[Outcome]] of the winner. The winner * of the race is considered to be the first fiber that completes with an outcome. The loser @@ -315,9 +336,9 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { */ def raceOutcome[A, B](fa: F[A], fb: F[B]): F[Either[Outcome[F, E, A], Outcome[F, E, B]]] = uncancelable { poll => - poll(racePair(fa, fb)).flatMap { - case Left((oc, f)) => f.cancel.as(Left(oc)) - case Right((f, oc)) => f.cancel.as(Right(oc)) + poll(raceOutcomeBoth(fa, fb)).map { + case Left((oc, _)) => Left(oc) + case Right((_, oc)) => Right(oc) } } @@ -346,24 +367,24 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { */ def race[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = uncancelable { poll => - poll(racePair(fa, fb)).flatMap { - case Left((oc, f)) => - oc match { - case Outcome.Succeeded(fa) => f.cancel *> fa.map(Left(_)) - case Outcome.Errored(ea) => f.cancel *> raiseError(ea) + poll(raceOutcomeBoth(fa, fb)).flatMap { + case Left((oca, ocb)) => + oca match { + case Outcome.Succeeded(fa) => fa.map(Left(_)) + case Outcome.Errored(ea) => raiseError(ea) case Outcome.Canceled() => - poll(f.join).onCancel(f.cancel).flatMap { + ocb match { case Outcome.Succeeded(fb) => fb.map(Right(_)) case Outcome.Errored(eb) => raiseError(eb) case Outcome.Canceled() => poll(canceled) *> never } } - case Right((f, oc)) => - oc match { - case Outcome.Succeeded(fb) => f.cancel *> fb.map(Right(_)) - case Outcome.Errored(eb) => f.cancel *> raiseError(eb) + case Right((oca, ocb)) => + ocb match { + case Outcome.Succeeded(fb) => fb.map(Right(_)) + case Outcome.Errored(eb) => raiseError(eb) case Outcome.Canceled() => - poll(f.join).onCancel(f.cancel).flatMap { + oca match { case Outcome.Succeeded(fa) => fa.map(Left(_)) case Outcome.Errored(ea) => raiseError(ea) case Outcome.Canceled() => poll(canceled) *> never From 174257033dfd7c4014d2ac4aaa96d1324070633a Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 7 Nov 2022 22:15:06 +0000 Subject: [PATCH 2/5] Introduce `raceBiased` --- .../scala/cats/effect/kernel/GenSpawn.scala | 71 +++++++++++++------ 1 file changed, 48 insertions(+), 23 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala index 88b42f7d88..e8207962a4 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala @@ -366,33 +366,58 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { * [[raceOutcome]] for a variant that returns the outcome of the winner. */ def race[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = + uncancelable(poll => poll(raceOutcomeBoth(fa, fb)).flatMap(embedRaceOutcome(_))) + + /** + * Like [[race]], but in the case that neither fiber completes with [[Outcome.Canceled]], + * priority is given to the outcome of the second fiber. Hence, this method is right-biased. + * + * @param fa + * the effect for the first racing fiber + * @param fb + * the effect for the second racing fiber + * + * @see + * [[race]] for a non-biased variant + */ + def raceBiased[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = uncancelable { poll => - poll(raceOutcomeBoth(fa, fb)).flatMap { - case Left((oca, ocb)) => - oca match { - case Outcome.Succeeded(fa) => fa.map(Left(_)) - case Outcome.Errored(ea) => raiseError(ea) - case Outcome.Canceled() => - ocb match { - case Outcome.Succeeded(fb) => fb.map(Right(_)) - case Outcome.Errored(eb) => raiseError(eb) - case Outcome.Canceled() => poll(canceled) *> never - } - } - case Right((oca, ocb)) => - ocb match { - case Outcome.Succeeded(fb) => fb.map(Right(_)) - case Outcome.Errored(eb) => raiseError(eb) - case Outcome.Canceled() => - oca match { - case Outcome.Succeeded(fa) => fa.map(Left(_)) - case Outcome.Errored(ea) => raiseError(ea) - case Outcome.Canceled() => poll(canceled) *> never - } - } + poll(raceOutcomeBoth(fa, fb)).flatMap { raceOutcome => + raceOutcome.merge._2 match { + case Outcome.Succeeded(fb) => fb.map(Left(_)) + case Outcome.Errored(eb) => raiseError(eb) + case Outcome.Canceled() => embedRaceOutcome(raceOutcome) + } } } + private[this] def embedRaceOutcome[A, B]( + raceOutcome: Either[Outcome[F, E, A], Outcome[F, E, B]]): F[Either[A, B]] = + raceOutcome match { + case Left((oca, ocb)) => + oca match { + case Outcome.Succeeded(fa) => fa.map(Left(_)) + case Outcome.Errored(ea) => raiseError(ea) + case Outcome.Canceled() => + ocb match { + case Outcome.Succeeded(fb) => fb.map(Right(_)) + case Outcome.Errored(eb) => raiseError(eb) + case Outcome.Canceled() => poll(canceled) *> never + } + } + case Right((oca, ocb)) => + ocb match { + case Outcome.Succeeded(fb) => fb.map(Right(_)) + case Outcome.Errored(eb) => raiseError(eb) + case Outcome.Canceled() => + oca match { + case Outcome.Succeeded(fa) => fa.map(Left(_)) + case Outcome.Errored(ea) => raiseError(ea) + case Outcome.Canceled() => poll(canceled) *> never + } + } + } + /** * Races the evaluation of two fibers and returns the [[Outcome]] of both. If the race is * canceled before one or both participants complete, then then whichever ones are incomplete From 9f6bffcbc8070eaabc94a548911739ee4fec89de Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 7 Nov 2022 23:23:43 +0000 Subject: [PATCH 3/5] Fix compile --- .../main/scala/cats/effect/kernel/GenSpawn.scala | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala index e8207962a4..d80c230a30 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala @@ -366,7 +366,9 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { * [[raceOutcome]] for a variant that returns the outcome of the winner. */ def race[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = - uncancelable(poll => poll(raceOutcomeBoth(fa, fb)).flatMap(embedRaceOutcome(_))) + uncancelable { poll => + poll(raceOutcomeBoth(fa, fb)).flatMap(oc => poll(embedRaceOutcome(oc))) + } /** * Like [[race]], but in the case that neither fiber completes with [[Outcome.Canceled]], @@ -384,15 +386,17 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { uncancelable { poll => poll(raceOutcomeBoth(fa, fb)).flatMap { raceOutcome => raceOutcome.merge._2 match { - case Outcome.Succeeded(fb) => fb.map(Left(_)) + case Outcome.Succeeded(fb) => fb.map(Right(_)) case Outcome.Errored(eb) => raiseError(eb) - case Outcome.Canceled() => embedRaceOutcome(raceOutcome) + case Outcome.Canceled() => poll(embedRaceOutcome(raceOutcome)) } } } private[this] def embedRaceOutcome[A, B]( - raceOutcome: Either[Outcome[F, E, A], Outcome[F, E, B]]): F[Either[A, B]] = + raceOutcome: Either[ + (Outcome[F, E, A], Outcome[F, E, B]), + (Outcome[F, E, A], Outcome[F, E, B])]): F[Either[A, B]] = raceOutcome match { case Left((oca, ocb)) => oca match { @@ -402,7 +406,7 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { ocb match { case Outcome.Succeeded(fb) => fb.map(Right(_)) case Outcome.Errored(eb) => raiseError(eb) - case Outcome.Canceled() => poll(canceled) *> never + case Outcome.Canceled() => canceled *> never } } case Right((oca, ocb)) => @@ -413,7 +417,7 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { oca match { case Outcome.Succeeded(fa) => fa.map(Left(_)) case Outcome.Errored(ea) => raiseError(ea) - case Outcome.Canceled() => poll(canceled) *> never + case Outcome.Canceled() => canceled *> never } } } From 12cadeab81cf202696ee1d8aa15c8acf471b47da Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 7 Nov 2022 23:25:36 +0000 Subject: [PATCH 4/5] Remove some `uncancelable` --- .../scala/cats/effect/kernel/GenSpawn.scala | 24 +++++++------------ 1 file changed, 9 insertions(+), 15 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala index d80c230a30..b673ebd448 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala @@ -335,11 +335,9 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { * [[race]] for a simpler variant that returns the successful outcome. */ def raceOutcome[A, B](fa: F[A], fb: F[B]): F[Either[Outcome[F, E, A], Outcome[F, E, B]]] = - uncancelable { poll => - poll(raceOutcomeBoth(fa, fb)).map { - case Left((oc, _)) => Left(oc) - case Right((_, oc)) => Right(oc) - } + raceOutcomeBoth(fa, fb).map { + case Left((oc, _)) => Left(oc) + case Right((_, oc)) => Right(oc) } /** @@ -366,9 +364,7 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { * [[raceOutcome]] for a variant that returns the outcome of the winner. */ def race[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = - uncancelable { poll => - poll(raceOutcomeBoth(fa, fb)).flatMap(oc => poll(embedRaceOutcome(oc))) - } + raceOutcomeBoth(fa, fb).flatMap(embedRaceOutcome(_)) /** * Like [[race]], but in the case that neither fiber completes with [[Outcome.Canceled]], @@ -383,13 +379,11 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { * [[race]] for a non-biased variant */ def raceBiased[A, B](fa: F[A], fb: F[B]): F[Either[A, B]] = - uncancelable { poll => - poll(raceOutcomeBoth(fa, fb)).flatMap { raceOutcome => - raceOutcome.merge._2 match { - case Outcome.Succeeded(fb) => fb.map(Right(_)) - case Outcome.Errored(eb) => raiseError(eb) - case Outcome.Canceled() => poll(embedRaceOutcome(raceOutcome)) - } + raceOutcomeBoth(fa, fb).flatMap { raceOutcome => + raceOutcome.merge._2 match { + case Outcome.Succeeded(fb) => fb.map(Right(_)) + case Outcome.Errored(eb) => raiseError(eb) + case Outcome.Canceled() => embedRaceOutcome(raceOutcome) } } From 499def9cdb15191c1fca932f836614dc4382300a Mon Sep 17 00:00:00 2001 From: Arman Bilge Date: Mon, 7 Nov 2022 23:31:48 +0000 Subject: [PATCH 5/5] Fix unmasking in `raceOutcomeBoth` --- .../main/scala/cats/effect/kernel/GenSpawn.scala | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala index 88b42f7d88..9e67fbbd92 100644 --- a/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala +++ b/kernel/shared/src/main/scala/cats/effect/kernel/GenSpawn.scala @@ -315,9 +315,19 @@ trait GenSpawn[F[_], E] extends MonadCancel[F, E] with Unique[F] { uncancelable { poll => poll(racePair(fa, fb)).flatMap { case Left((oca, f)) => - f.cancel.whenA(!oca.isCanceled) *> f.join.map(ocb => Left((oca, ocb))) + val joined = + if (oca.isCanceled) + poll(f.join) + else + f.cancel *> f.join + joined.map(ocb => Left((oca, ocb))) case Right((f, ocb)) => - f.cancel.whenA(!ocb.isCanceled) *> f.join.map(oca => Right((oca, ocb))) + val joined = + if (ocb.isCanceled) + poll(f.join) + else + f.cancel *> f.join + joined.map(oca => Right((oca, ocb))) } }