From 5d13bf65b98e9418106daed678c7b63f29c6666d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Luis=20Miguel=20Mej=C3=ADa=20Su=C3=A1rez?= Date: Mon, 3 Mar 2025 20:04:23 -0500 Subject: [PATCH 1/4] Add JobManager API --- .../scala/cats/effect/std/JobManager.scala | 63 +++++++++++++++++++ 1 file changed, 63 insertions(+) create mode 100644 std/shared/src/main/scala/cats/effect/std/JobManager.scala diff --git a/std/shared/src/main/scala/cats/effect/std/JobManager.scala b/std/shared/src/main/scala/cats/effect/std/JobManager.scala new file mode 100644 index 0000000000..77d44e977f --- /dev/null +++ b/std/shared/src/main/scala/cats/effect/std/JobManager.scala @@ -0,0 +1,63 @@ +/* + * Copyright 2020-2025 Typelevel + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package cats.effect.std + +import cats.effect.kernel._ + +/** + * A `JobManager` allows you to launch `Jobs` in the background using a unique identifier. Then + * you can use the identifier to query for the status of the job or cancel it. + */ +trait JobManager[F[_], Id, S] { + + /** + * Creates and launches the given `Job` in the background. If another Job with the same id was + * already running, it will be cancelled before starting this one. + */ + def startJob(id: Id, job: Resource[F, JobManager.Job[F, S]]): F[Unit] + + /** + * Gets the status of the `Job` associated with the given `id`. If `id` doesn't exists or the + * `Job` already finished then the returned value will be a `None`. + */ + def getJobStatus(id: Id): F[Option[S]] + + /** + * Signals cancellation of the `Job` associated with the given `id`, and waits for its + * completion. + */ + def cancelJob(id: Id): F[Unit] +} + +object JobManager { + + /** + * Represents a job managed by a `JobManager`. + */ + trait Job[F[_], S] { + + /** + * Starts the logic of this `Job`. + */ + def run: F[Unit] + + /** + * Gets the status of this `Job`. + */ + def getStatus: F[S] + } +} From 5381f44f5e99093e5095bebd34465a943a94fd13 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Luis=20Miguel=20Mej=C3=ADa=20Su=C3=A1rez?= Date: Mon, 3 Mar 2025 20:40:51 -0500 Subject: [PATCH 2/4] Implement JobManager --- .../scala/cats/effect/std/JobManager.scala | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/std/shared/src/main/scala/cats/effect/std/JobManager.scala b/std/shared/src/main/scala/cats/effect/std/JobManager.scala index 77d44e977f..8949469960 100644 --- a/std/shared/src/main/scala/cats/effect/std/JobManager.scala +++ b/std/shared/src/main/scala/cats/effect/std/JobManager.scala @@ -17,6 +17,7 @@ package cats.effect.std import cats.effect.kernel._ +import cats.syntax.all._ /** * A `JobManager` allows you to launch `Jobs` in the background using a unique identifier. Then @@ -60,4 +61,39 @@ object JobManager { */ def getStatus: F[S] } + + private final case class RunningJob[F[_], S]( + status: F[S], + cancel: F[Unit] + ) + + def apply[F[_], Id, S](implicit F: Concurrent[F]): Resource[F, JobManager[F, Id, S]] = + for { + supervisor <- Supervisor[F](await = true) + jobsMap <- Resource.eval(MapRef[F, Id, RunningJob[F, S]]) + } yield new JobManager[F, Id, S] { + override def startJob(id: Id, jobR: Resource[F, Job[F, S]]): F[Unit] = { + val runJob = jobR.use { job => + supervisor.supervise(job.run).flatMap { fiber => + jobsMap(id) + .getAndSet( + RunningJob( + status = job.getStatus, + cancel = fiber.cancel + ).some + ) + .flatMap(_.traverse_(_.cancel)) >> + fiber.join + } + } + + supervisor.supervise(runJob).void + } + + override def getJobStatus(id: Id): F[Option[S]] = + jobsMap(id).get.flatMap(_.traverse(_.status)) + + override def cancelJob(id: Id): F[Unit] = + jobsMap(id).getAndSet(None).flatMap(_.traverse_(_.cancel)) + } } From 3ebd753e086236a3588a3bab951fa720e846e03b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Luis=20Miguel=20Mej=C3=ADa=20Su=C3=A1rez?= Date: Sun, 23 Mar 2025 15:58:15 -0500 Subject: [PATCH 3/4] startJob waits until the job has been successfully registered --- .../scala/cats/effect/std/JobManager.scala | 43 +++++++++++++------ 1 file changed, 30 insertions(+), 13 deletions(-) diff --git a/std/shared/src/main/scala/cats/effect/std/JobManager.scala b/std/shared/src/main/scala/cats/effect/std/JobManager.scala index 8949469960..9135fe59cf 100644 --- a/std/shared/src/main/scala/cats/effect/std/JobManager.scala +++ b/std/shared/src/main/scala/cats/effect/std/JobManager.scala @@ -28,6 +28,9 @@ trait JobManager[F[_], Id, S] { /** * Creates and launches the given `Job` in the background. If another Job with the same id was * already running, it will be cancelled before starting this one. + * + * Note: This waits until the actual job's `run` process has started. In other words, you can + * query the `status` of the associated id immediately after. */ def startJob(id: Id, job: Resource[F, JobManager.Job[F, S]]): F[Unit] @@ -73,21 +76,35 @@ object JobManager { jobsMap <- Resource.eval(MapRef[F, Id, RunningJob[F, S]]) } yield new JobManager[F, Id, S] { override def startJob(id: Id, jobR: Resource[F, Job[F, S]]): F[Unit] = { - val runJob = jobR.use { job => - supervisor.supervise(job.run).flatMap { fiber => - jobsMap(id) - .getAndSet( - RunningJob( - status = job.getStatus, - cancel = fiber.cancel - ).some - ) - .flatMap(_.traverse_(_.cancel)) >> - fiber.join + Deferred[F, Unit].flatMap { waitForJobRegistration => + // Running a job has the following steps: + // 1. Run the job setup (Resource acquire). + // 2. Launch the job in the background. + // 3. Register the job in the jobsMap. + // 4. In case there was another job already register with the same id, cancel it. + // 5. Wait for the job to finish. + // 6. Run the job cleanup (Resource release). + val runJob = jobR.use { job => + supervisor.supervise(job.run).flatMap { jobFiber => + val registerJob = + F.guarantee( + jobsMap(id).getAndSet( + RunningJob( + status = job.getStatus, + cancel = jobFiber.cancel + ).some + ), + fin = waitForJobRegistration.complete(()).void + ) + + registerJob.flatMap(_.traverse_(_.cancel)) >> jobFiber.join + } } - } - supervisor.supervise(runJob).void + // Starts the run job process in the background, + // but wait until it has been registered in the jobsMap. + supervisor.supervise(runJob) >> waitForJobRegistration.get + } } override def getJobStatus(id: Id): F[Option[S]] = From 4fb97c7f44871702fc0ad2053f56572d8a83e8a0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Luis=20Miguel=20Mej=C3=ADa=20Su=C3=A1rez?= Date: Tue, 25 Aug 2026 10:59:54 -0500 Subject: [PATCH 4/4] Do not swallow `Job` results and ensure correct cleanup --- .../scala/cats/effect/std/JobManager.scala | 118 ++++++++++++------ 1 file changed, 79 insertions(+), 39 deletions(-) diff --git a/std/shared/src/main/scala/cats/effect/std/JobManager.scala b/std/shared/src/main/scala/cats/effect/std/JobManager.scala index 9135fe59cf..8c156c1946 100644 --- a/std/shared/src/main/scala/cats/effect/std/JobManager.scala +++ b/std/shared/src/main/scala/cats/effect/std/JobManager.scala @@ -21,9 +21,10 @@ import cats.syntax.all._ /** * A `JobManager` allows you to launch `Jobs` in the background using a unique identifier. Then - * you can use the identifier to query for the status of the job or cancel it. + * you can use the identifier to query for the status of the job, wait for its completion, or + * cancel it. */ -trait JobManager[F[_], Id, S] { +trait JobManager[F[_], Id, S, R] { /** * Creates and launches the given `Job` in the background. If another Job with the same id was @@ -31,12 +32,25 @@ trait JobManager[F[_], Id, S] { * * Note: This waits until the actual job's `run` process has started. In other words, you can * query the `status` of the associated id immediately after. + * + * @return + * the [[cats.effect.kernel.Fiber]] of the started `Job`. + */ + def startJob(id: Id, job: Resource[F, JobManager.Job[F, S, R]]): F[Fiber[F, Throwable, R]] + + /** + * Waits until the given `Job` finishes and returns its result. + * + * Note: If `id` does not exists or the `Job` already finished then the returned value will be + * a `None`. */ - def startJob(id: Id, job: Resource[F, JobManager.Job[F, S]]): F[Unit] + def join(id: Id): F[Option[Outcome[F, Throwable, R]]] /** - * Gets the status of the `Job` associated with the given `id`. If `id` doesn't exists or the - * `Job` already finished then the returned value will be a `None`. + * Gets the status of the `Job` associated with the given `id`. + * + * Note: If `id` does not exists or the `Job` already finished then the returned value will be + * a `None`. */ def getJobStatus(id: Id): F[Option[S]] @@ -52,12 +66,12 @@ object JobManager { /** * Represents a job managed by a `JobManager`. */ - trait Job[F[_], S] { + trait Job[F[_], S, R] { /** * Starts the logic of this `Job`. */ - def run: F[Unit] + def run: F[R] /** * Gets the status of this `Job`. @@ -65,52 +79,78 @@ object JobManager { def getStatus: F[S] } - private final case class RunningJob[F[_], S]( + private final case class RunningJob[F[_], S, R]( + token: Unique.Token, status: F[S], - cancel: F[Unit] + fiber: Fiber[F, Throwable, R] ) - def apply[F[_], Id, S](implicit F: Concurrent[F]): Resource[F, JobManager[F, Id, S]] = + def apply[F[_], Id, S, R]( + implicit F: Concurrent[F] + ): Resource[F, JobManager[F, Id, S, R]] = for { supervisor <- Supervisor[F](await = true) - jobsMap <- Resource.eval(MapRef[F, Id, RunningJob[F, S]]) - } yield new JobManager[F, Id, S] { - override def startJob(id: Id, jobR: Resource[F, Job[F, S]]): F[Unit] = { - Deferred[F, Unit].flatMap { waitForJobRegistration => - // Running a job has the following steps: - // 1. Run the job setup (Resource acquire). - // 2. Launch the job in the background. - // 3. Register the job in the jobsMap. - // 4. In case there was another job already register with the same id, cancel it. - // 5. Wait for the job to finish. - // 6. Run the job cleanup (Resource release). - val runJob = jobR.use { job => - supervisor.supervise(job.run).flatMap { jobFiber => - val registerJob = - F.guarantee( - jobsMap(id).getAndSet( - RunningJob( - status = job.getStatus, - cancel = jobFiber.cancel - ).some - ), - fin = waitForJobRegistration.complete(()).void + jobsMap <- Resource.eval(MapRef[F, Id, RunningJob[F, S, R]]) + } yield new JobManager[F, Id, S, R] { + override def startJob( + id: Id, + jobR: Resource[F, Job[F, S, R]] + ): F[Fiber[F, Throwable, R]] = { + ( + F.unique, + Deferred[F, Fiber[F, Throwable, R]] + ).flatMapN { + case (jobToken, waitForJobRegistration) => + // Running a job has the following steps: + // 1. Run the job setup (Resource acquire). + // 2. Launch the job in the background. + // 3. Register the job in the jobsMap. + // 4. In case there was another job already register with the same id, cancel it. + // 5. Wait for the job to finish. + // 6. Run the job cleanup (Resource release). + val runJob = jobR.use { job => + supervisor + .supervise( + // Ensure that no matter what happens with the job, + // we will clean up it from the jobsMap. + F.guarantee( + job.run, + fin = jobsMap(id).update { runningJob => + // Guard to avoid removing a different job from the jobsMap. + runningJob.filterNot(_.token == jobToken) + } + ) ) - - registerJob.flatMap(_.traverse_(_.cancel)) >> jobFiber.join + .flatMap { jobFiber => + val registerJob = + F.guarantee( + jobsMap(id).getAndSet( + RunningJob( + token = jobToken, + status = job.getStatus, + fiber = jobFiber + ).some + ), + fin = waitForJobRegistration.complete(jobFiber).void + ) + + registerJob.flatMap(_.traverse_(_.fiber.cancel)) >> jobFiber.join + } } - } - // Starts the run job process in the background, - // but wait until it has been registered in the jobsMap. - supervisor.supervise(runJob) >> waitForJobRegistration.get + // Starts the run job process in the background, + // but wait until it has been registered in the jobsMap. + supervisor.supervise(runJob) >> waitForJobRegistration.get } } + override def join(id: Id): F[Option[Outcome[F, Throwable, R]]] = + jobsMap(id).get.flatMap(_.traverse(_.fiber.join)) + override def getJobStatus(id: Id): F[Option[S]] = jobsMap(id).get.flatMap(_.traverse(_.status)) override def cancelJob(id: Id): F[Unit] = - jobsMap(id).getAndSet(None).flatMap(_.traverse_(_.cancel)) + F.uncancelable(_ => jobsMap(id).get.flatMap(_.traverse_(_.fiber.cancel))) } }