Skip to content

Instantly share code, notes, and snippets.

@djspiewak
Created September 3, 2021 13:49
Show Gist options
  • Select an option

  • Save djspiewak/006c65d9622958a01faf45eb1675f97c to your computer and use it in GitHub Desktop.

Select an option

Save djspiewak/006c65d9622958a01faf45eb1675f97c to your computer and use it in GitHub Desktop.
final override def map2[A, B, Z](fa: ParallelF[F, A], fb: ParallelF[F, B])(
f: (A, B) => Z): ParallelF[F, Z] =
ParallelF(
F.uncancelable { poll =>
for {
fiberA <- F.start(ParallelF.value(fa))
fiberB <- F.start(ParallelF.value(fb))
// start a pair of supervisors to ensure that the opposite is canceled on error
_ <- F start {
fiberB.join flatMap {
case Outcome.Succeeded(_) => F.unit
case _ => fiberA.cancel
}
}
_ <- F start {
fiberA.join flatMap {
case Outcome.Succeeded(_) => F.unit
case _ => fiberB.cancel
}
}
a <- F
.onCancel(poll(fiberA.join), F.both(fiberA.cancel, fiberB.cancel).void)
.flatMap[A] {
case Outcome.Succeeded(fa) =>
fa
case Outcome.Errored(e) =>
fiberB.cancel *> F.raiseError(e)
case Outcome.Canceled() =>
fiberB.cancel *> poll {
fiberB.join flatMap {
case Outcome.Succeeded(_) | Outcome.Canceled() =>
F.canceled *> F.never
case Outcome.Errored(e) =>
F.raiseError(e)
}
}
}
z <- F.onCancel(poll(fiberB.join), fiberB.cancel).flatMap[Z] {
case Outcome.Succeeded(fb) =>
fb.map(b => f(a, b))
case Outcome.Errored(e) =>
F.raiseError(e)
case Outcome.Canceled() =>
poll {
fiberA.join flatMap {
case Outcome.Succeeded(_) | Outcome.Canceled() =>
F.canceled *> F.never
case Outcome.Errored(e) =>
F.raiseError(e)
}
}
}
} yield z
}
)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment