From 35acc7909577a36df50af8a4072b17af2a138e99 Mon Sep 17 00:00:00 2001 From: Guntis Smaukstelis Date: Wed, 12 Feb 2025 23:18:40 +0200 Subject: [PATCH] Universal scheduler, cleanup task for scheduler --- src/main/scala/app/Main.scala | 20 ++++++++++-- src/main/scala/fetch/Scheduler.scala | 2 +- src/main/scala/fetchDMI/Scheduler.scala | 42 +++++++++++++++++++++++++ 3 files changed, 61 insertions(+), 3 deletions(-) create mode 100644 src/main/scala/fetchDMI/Scheduler.scala diff --git a/src/main/scala/app/Main.scala b/src/main/scala/app/Main.scala index a5ae410..da8c021 100644 --- a/src/main/scala/app/Main.scala +++ b/src/main/scala/app/Main.scala @@ -1,7 +1,8 @@ package app import cats.effect._ -import cats.implicits.catsSyntaxTuple2Parallel +import cats.implicits.catsSyntaxTuple3Parallel +import data.DataService import db.{DBConnection, PostgresService} import fetch.{FetchService, FileFetchScheduler} import server.Server @@ -17,10 +18,25 @@ object Main extends IOApp { fileFetchScheduler <- FileFetchScheduler.of(postgresService, fetch) schedulerTask = fileFetchScheduler.run.compile.drain + +// scheduler <- fetchDMI.Scheduler.of("Grib", List(3, 27, 39, 51)) +// simpleTask = IO.delay { +// val timeNow = LocalTime.now().format(DateTimeFormatter.ofPattern("HH:mm")) +// List(s"Task executed at $timeNow") +// } +// simpleScheduler = scheduler.scheduleTask(simpleTask).compile.drain + + + scheduler <- fetchDMI.Scheduler.of("Cleanup", List(1)) + cleanupTask = DataService.deleteOldForecasts() + cleanupScheduler = scheduler.scheduleTask(cleanupTask).compile.drain + + + server <- Server.of(postgresService, fetch) serverTask = server.run - exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success) + exitCode <- (serverTask, schedulerTask, cleanupScheduler).parMapN((_, _, _) => ExitCode.Success) } yield exitCode } } \ No newline at end of file diff --git a/src/main/scala/fetch/Scheduler.scala b/src/main/scala/fetch/Scheduler.scala index 02b3af3..1e4f5c8 100644 --- a/src/main/scala/fetch/Scheduler.scala +++ b/src/main/scala/fetch/Scheduler.scala @@ -34,7 +34,7 @@ class Scheduler(log: Logger[IO]) { def scheduleTask(task: IO[Either[Throwable, (String, String)]]): Stream[IO, Either[Throwable, (String, String)]] = { Stream.eval(durationToNextHalfHour).flatMap { delay => { - Stream.eval(log.info(s"Scheduler started with delay: ${delay.toMinutes} min")) *> + Stream.eval(log.info(s"CSV scheduler started with delay: ${delay.toMinutes} min")) *> (Stream.sleep[IO](delay) ++ Stream.awakeEvery[IO](1.hour)) .evalMap(_ => task) }} diff --git a/src/main/scala/fetchDMI/Scheduler.scala b/src/main/scala/fetchDMI/Scheduler.scala new file mode 100644 index 0000000..e7b9822 --- /dev/null +++ b/src/main/scala/fetchDMI/Scheduler.scala @@ -0,0 +1,42 @@ +package fetchDMI + +import cats.effect._ +import cats.implicits.catsSyntaxApply +import fs2.Stream +import org.typelevel.log4cats.Logger +import org.typelevel.log4cats.slf4j.Slf4jLogger + +import java.time.{Duration, LocalTime} +import scala.concurrent.duration._ + + +object Scheduler { + def of(name: String, minutes: List[Int]): IO[Scheduler] = { + Slf4jLogger.create[IO].map(logger => new Scheduler(name, minutes, logger)) + } +} + +class Scheduler(name: String, minutes: List[Int], log: Logger[IO]) { + private def durationToNext(targetMinute: Int)(implicit clock: Clock[IO]): IO[FiniteDuration] = { + clock.realTime.map { duration => + val now = LocalTime.ofSecondOfDay((duration.toMillis / 1000) % (24 * 60 * 60)) + val nextTime = + if (now.getMinute < targetMinute) now.withMinute(targetMinute) + else now.plusHours(1).withMinute(targetMinute) + val durationToNext = Duration.between(now, nextTime) + FiniteDuration(durationToNext.toMillis, MILLISECONDS) + } + } + + def scheduleTask(task: IO[List[String]]): Stream[IO, List[String]] = { + def streamForMinute(minute: Int): Stream[IO, List[String]] = { + Stream.eval(durationToNext(minute)).flatMap { delay => + Stream.eval(log.info(s"$name scheduler in: ${delay.toMinutes} min")) *> + (Stream.sleep[IO](delay) ++ Stream.awakeEvery[IO](1.hour)) + .evalMap(_ => task) + } + } + + minutes.map(streamForMinute).reduce(_ merge _) + } +} \ No newline at end of file