From be037a311d970cb352a2e48a61a58fc47904f573 Mon Sep 17 00:00:00 2001 From: Guntis Smaukstelis Date: Thu, 25 May 2023 11:53:32 +0300 Subject: [PATCH] Remove StatefulFetchService as we have already DataService with lst 24h state --- src/main/scala/fetch/FetchService.scala | 8 +-- src/main/scala/fetch/FileFetchScheduler.scala | 6 +- src/main/scala/fetch/Main.scala | 5 +- .../scala/fetch/StatefulFetchService.scala | 59 ------------------- src/main/scala/server/Main.scala | 7 +-- src/main/scala/server/Server.scala | 6 +- 6 files changed, 12 insertions(+), 79 deletions(-) delete mode 100644 src/main/scala/fetch/StatefulFetchService.scala diff --git a/src/main/scala/fetch/FetchService.scala b/src/main/scala/fetch/FetchService.scala index 5615250..2ce3818 100644 --- a/src/main/scala/fetch/FetchService.scala +++ b/src/main/scala/fetch/FetchService.scala @@ -20,19 +20,13 @@ final case class WeatherServerConfig( url: String, ) -trait FetchServiceTrait { - def fetchSingleFile(fileName: String): IO[Either[Throwable, (String, String)]] - def fetchInRange(from: LocalDateTime, to: LocalDateTime): IO[List[Either[Throwable, (String, String)]]] - def fetchFromDate(date: LocalDate): IO[List[Either[Throwable, (String, String)]]] -} - object FetchService { def of: IO[FetchService] = { Slf4jLogger.create[IO].map(logger => new FetchService(new FileNameService, logger)) } } -class FetchService(fileNameService: FileNameService, log: Logger[IO]) extends FetchServiceTrait { +class FetchService(fileNameService: FileNameService, log: Logger[IO]) { private val weatherServerConfig: WeatherServerConfig = ConfigSource.default.load[WeatherServerConfig] match { case Right(config) => config case Left(errors) => throw new RuntimeException(s"Unable to load config: $errors") diff --git a/src/main/scala/fetch/FileFetchScheduler.scala b/src/main/scala/fetch/FileFetchScheduler.scala index 1b69f79..e4af8f8 100644 --- a/src/main/scala/fetch/FileFetchScheduler.scala +++ b/src/main/scala/fetch/FileFetchScheduler.scala @@ -1,14 +1,14 @@ package fetch import cats.effect.IO -import db.{DBService, DataServiceTrait} +import db.DataServiceTrait import fs2.Stream import org.typelevel.log4cats.Logger import org.typelevel.log4cats.slf4j.Slf4jLogger object FileFetchScheduler { - def of(dataService: DataServiceTrait, fetch: FetchServiceTrait): IO[FileFetchScheduler] = { + def of(dataService: DataServiceTrait, fetch: FetchService): IO[FileFetchScheduler] = { Scheduler.of.flatMap { scheduler => Slf4jLogger.create[IO].map { new FileFetchScheduler(dataService, fetch, new FileNameService(), scheduler, _) @@ -17,7 +17,7 @@ object FileFetchScheduler { } } -class FileFetchScheduler(dataService: DataServiceTrait, fetch: FetchServiceTrait, fileNameService: FileNameService, scheduler: Scheduler, log: Logger[IO]) { +class FileFetchScheduler(dataService: DataServiceTrait, fetch: FetchService, fileNameService: FileNameService, scheduler: Scheduler, log: Logger[IO]) { def run: Stream[IO, Unit] = { val fetchTask = fileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile) scheduler.scheduleTask(fetchTask) diff --git a/src/main/scala/fetch/Main.scala b/src/main/scala/fetch/Main.scala index d48c7ce..7b07683 100644 --- a/src/main/scala/fetch/Main.scala +++ b/src/main/scala/fetch/Main.scala @@ -33,9 +33,8 @@ object Main { def run: IO[Unit] = { for { fetch <- FetchService.of - statefulFetch <- StatefulFetchService.of(fetch) - fetchResultEither <- statefulFetch.fetchSingleFile("20230524_0030.csv").attempt - fetchResultEither <- statefulFetch.fetchSingleFile("20230522_0130.csv").attempt + fetchResultEither <- fetch.fetchSingleFile("20230524_0030.csv").attempt + fetchResultEither <- fetch.fetchSingleFile("20230522_0130.csv").attempt fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList fetchResult = fetchResultEither.flatMap(res => res.flatMap(aaa => { println(s"fffffff: ${aaa._1}") diff --git a/src/main/scala/fetch/StatefulFetchService.scala b/src/main/scala/fetch/StatefulFetchService.scala deleted file mode 100644 index 6217b68..0000000 --- a/src/main/scala/fetch/StatefulFetchService.scala +++ /dev/null @@ -1,59 +0,0 @@ -package fetch - -import cats.effect.{IO, Ref} -import cats.implicits.toTraverseOps -import org.typelevel.log4cats.Logger -import org.typelevel.log4cats.slf4j.Slf4jLogger - -import java.time.{LocalDate, LocalDateTime} - -object StatefulFetchService { - def of(fetchService: FetchService): IO[StatefulFetchService] = { - Slf4jLogger.create[IO].map(logger => new StatefulFetchService(fetchService, new FileNameService(), logger)) - } -} - -class StatefulFetchService(fetchService: FetchService, fileNameService: FileNameService, log: Logger[IO]) extends FetchServiceTrait { - private val state: Ref[IO, Map[String, String]] = Ref.unsafe(Map.empty) - - private def logState: IO[Unit] = { - state.get.flatMap(currentState => log.info(s"State keys: ${currentState.keys}")) - } - - private def updateState(fileName: String, content: String): IO[Unit] = { - for { - last24Hours <- fileNameService.generateLast24Hours - _ <- state.update(st => (st + (fileName -> content)).filterKeys(last24Hours.contains).toMap) - _ <- logState - } yield () - } - - def fetchSingleFile(fileName: String): IO[Either[Throwable, (String, String)]] = { - fetchService.fetchSingleFile(fileName).flatMap { - case Right((fileName, content)) => - updateState(fileName, content).as(Right((fileName, content))) - case e@Left(_) => IO(e) - } - } - - def fetchInRange(from: LocalDateTime, to: LocalDateTime): IO[List[Either[Throwable, (String, String)]]] = { - fetchService.fetchInRange(from, to).flatMap { results => - val successfulResults = results.collect { case Right(data) => data } - successfulResults.traverse { case (name, content) => updateState(name, content) }.as(results) - } - } - - def fetchFromDate(date: LocalDate): IO[List[Either[Throwable, (String, String)]]] = { - fetchService.fetchFromDate(date).flatMap { results => - val successfulResults = results.collect { case Right(data) => data } - successfulResults.traverse { case (name, content) => updateState(name, content) }.as(results) - } - } -} - - - - - - - diff --git a/src/main/scala/server/Main.scala b/src/main/scala/server/Main.scala index 5fe3911..f0f28bb 100644 --- a/src/main/scala/server/Main.scala +++ b/src/main/scala/server/Main.scala @@ -2,7 +2,7 @@ package server import cats.effect._ import cats.implicits.catsSyntaxTuple2Parallel import db.{DBService, DataService} -import fetch.{FetchService, FileFetchScheduler, StatefulFetchService} +import fetch.{FetchService, FileFetchScheduler} object Main extends IOApp { def run(args: List[String]): IO[ExitCode] = { @@ -10,10 +10,9 @@ object Main extends IOApp { dbService <- DBService.of dataService <- DataService.of(dbService) fetch <- FetchService.of - statefulFetch <- StatefulFetchService.of(fetch) - fileFetchScheduler <- FileFetchScheduler.of(dataService, statefulFetch) + fileFetchScheduler <- FileFetchScheduler.of(dataService, fetch) schedulerTask = fileFetchScheduler.run.compile.drain - server <- Server.of(dataService, statefulFetch) + server <- Server.of(dataService, fetch) serverTask = server.run exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success) } yield exitCode diff --git a/src/main/scala/server/Server.scala b/src/main/scala/server/Server.scala index cda6d4c..ab6b9db 100644 --- a/src/main/scala/server/Server.scala +++ b/src/main/scala/server/Server.scala @@ -4,7 +4,7 @@ import cats.effect._ import cats.implicits.toTraverseOps import com.comcast.ip4s.IpLiteralSyntax import db.DataService -import fetch.FetchServiceTrait +import fetch.FetchService import parse.{Parser, WeatherData} import server.ValidateRoutes.{AggKey, CityList, DateTimeRange, ValidDate} import io.circe.{Json, Printer} @@ -27,14 +27,14 @@ import scala.concurrent.duration.DurationInt object Server { - def of(dataService: DataService, fetch: FetchServiceTrait): IO[Server] = { + def of(dataService: DataService, fetch: FetchService): IO[Server] = { Slf4jLogger.create[IO].map { new Server(dataService, fetch, _) } } } -class Server(dataService: DataService, fetch: FetchServiceTrait, log: Logger[IO]) { +class Server(dataService: DataService, fetch: FetchService, log: Logger[IO]) { // Define the extension method `pretty` for Json implicit class JsonPrettyPrinter(json: Json) {