From b4955e26197438b1a8aea9ac8d832f597a1c5cbb Mon Sep 17 00:00:00 2001 From: Guntis Smaukstelis Date: Sat, 20 May 2023 21:53:03 +0300 Subject: [PATCH] Implement FileFetchScheduler and refactor Main class --- src/main/scala/fetch/FileFetchScheduler.scala | 29 +++++++++++++++++++ .../scala/{server => fetch}/Scheduler.scala | 3 +- src/main/scala/server/Main.scala | 20 ++++--------- 3 files changed, 36 insertions(+), 16 deletions(-) create mode 100644 src/main/scala/fetch/FileFetchScheduler.scala rename src/main/scala/{server => fetch}/Scheduler.scala (96%) diff --git a/src/main/scala/fetch/FileFetchScheduler.scala b/src/main/scala/fetch/FileFetchScheduler.scala new file mode 100644 index 0000000..1c773c9 --- /dev/null +++ b/src/main/scala/fetch/FileFetchScheduler.scala @@ -0,0 +1,29 @@ +package fetch + +import cats.effect.IO +import db.DBService +import fs2.Stream +import org.typelevel.log4cats.Logger +import org.typelevel.log4cats.slf4j.Slf4jLogger + + +object FileFetchScheduler { + def of(dbService: DBService): IO[FileFetchScheduler] = { + Slf4jLogger.create[IO].map { + new FileFetchScheduler(dbService, _) + } + } +} +class FileFetchScheduler(dbService: DBService, log: Logger[IO]) { + def run: Stream[IO, Unit] = { + val fetchTask = FileNameService.generateCurrentHour.flatMap(FetchService.fetchSingleFile) + Scheduler.scheduleTask(fetchTask) + .evalMap { case (name, content) => + log.info(s"fetched: $name") *> + dbService.save(name, content).attempt.flatMap { + case Right(savedName) => log.info(s"saved: $savedName") + case Left(err) => log.error(s"error: $err") + } + } + } +} \ No newline at end of file diff --git a/src/main/scala/server/Scheduler.scala b/src/main/scala/fetch/Scheduler.scala similarity index 96% rename from src/main/scala/server/Scheduler.scala rename to src/main/scala/fetch/Scheduler.scala index 834b29f..62283a7 100644 --- a/src/main/scala/server/Scheduler.scala +++ b/src/main/scala/fetch/Scheduler.scala @@ -1,9 +1,8 @@ -package server +package fetch import cats.effect._ import cats.effect.unsafe.implicits.global import cats.implicits.catsSyntaxApply -import fetch.{FetchService, FileNameService} import fs2.Stream import java.time.{Duration, LocalTime} diff --git a/src/main/scala/server/Main.scala b/src/main/scala/server/Main.scala index c06e682..ab339a8 100644 --- a/src/main/scala/server/Main.scala +++ b/src/main/scala/server/Main.scala @@ -1,26 +1,18 @@ package server import cats.effect._ import cats.implicits.catsSyntaxTuple2Parallel + import db.DBService -import fetch.{FetchService, FileNameService} +import fetch.FileFetchScheduler object Main extends IOApp { def run(args: List[String]): IO[ExitCode] = { for { dbService <- DBService.of - fetchTask = FileNameService.generateCurrentHour.flatMap(FetchService.fetchSingleFile) - scheduler = Scheduler.scheduleTask(fetchTask) - .evalMap { case (name, content) => - IO(println(s"fetched: $name")) *> - dbService.save(name, content).attempt.flatMap { - case Right(savedName) => IO(println(s"File saved: $savedName")) - case Left(err) => IO(println(s"Error: $err")) - } - }.compile.drain // Convert Stream[IO, Unit] to IO[Unit] - - server = Server.run // This is an IO[Server] - - exitCode <- (server, scheduler).parMapN((_, _) => ExitCode.Success) + fileFetchScheduler <- FileFetchScheduler.of(dbService) + schedulerTask = fileFetchScheduler.run.compile.drain + serverTask = Server.run + exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success) } yield exitCode } } \ No newline at end of file