Dependency inject logger and replace println
This commit is contained in:
@@ -6,7 +6,8 @@ import org.http4s._
|
||||
import org.http4s.client.Client
|
||||
import org.http4s.headers.Authorization
|
||||
import org.http4s.ember.client.EmberClientBuilder
|
||||
|
||||
import org.typelevel.log4cats.Logger
|
||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||
import pureconfig._
|
||||
import pureconfig.generic.auto._
|
||||
|
||||
@@ -20,6 +21,12 @@ final case class WeatherServerConfig(
|
||||
)
|
||||
|
||||
object FetchService {
|
||||
def of: IO[FetchService] = {
|
||||
Slf4jLogger.create[IO].map(logger => new FetchService(logger))
|
||||
}
|
||||
}
|
||||
|
||||
class FetchService(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")
|
||||
@@ -36,10 +43,10 @@ object FetchService {
|
||||
|
||||
client.expect[String](request).redeemWith(
|
||||
error => IO(Left(error))
|
||||
// .flatTap(_ => IO.println(s"Request failed to url: $url with error: ${error.getMessage}")),
|
||||
// .flatTap(_ => log.error(s"Request failed to url: $url with error: ${error.getMessage}"))
|
||||
,
|
||||
fileContent => IO(Right((fileName, fileContent)))
|
||||
// .flatTap(_ => IO.println(s"Fetched: $fileName"))
|
||||
// .flatTap(_ => log.info(s"Fetched: $fileName"))
|
||||
)
|
||||
}
|
||||
|
||||
@@ -54,11 +61,11 @@ object FetchService {
|
||||
fetchFiles(List(fileName)).map { results =>
|
||||
results.headOption match {
|
||||
case Some(Right(result)) =>
|
||||
IO.println(s"fetched: $fileName").as(Right(result))
|
||||
log.info(s"fetched: $fileName").as(Right(result))
|
||||
case Some(Left(err)) =>
|
||||
IO.println(s"failed fetch: $fileName with error: ${err.getMessage}").as(Left(err))
|
||||
log.error(s"failed fetch: $fileName with error: ${err.getMessage}").as(Left(err))
|
||||
case None =>
|
||||
IO.println(s"failed fetch: $fileName").as(Left(new Exception("No file fetched")))
|
||||
log.error(s"failed fetch: $fileName").as(Left(new Exception("No file fetched")))
|
||||
}
|
||||
}.flatten
|
||||
}
|
||||
|
||||
@@ -8,28 +8,30 @@ import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||
|
||||
|
||||
object FileFetchScheduler {
|
||||
def of(dbService: DBService): IO[FileFetchScheduler] = {
|
||||
def of(dbService: DBService, fetch: FetchService): IO[FileFetchScheduler] = {
|
||||
Slf4jLogger.create[IO].map {
|
||||
new FileFetchScheduler(dbService, _)
|
||||
new FileFetchScheduler(dbService, fetch, _)
|
||||
}
|
||||
}
|
||||
}
|
||||
class FileFetchScheduler(dbService: DBService, log: Logger[IO]) {
|
||||
class FileFetchScheduler(dbService: DBService, fetch: FetchService, log: Logger[IO]) {
|
||||
def run: Stream[IO, Unit] = {
|
||||
val fetchTask = FileNameService.generateCurrentHour.flatMap(FetchService.fetchSingleFile)
|
||||
Scheduler.scheduleTask(fetchTask)
|
||||
.evalMap {
|
||||
case Left(fetchErr) =>
|
||||
log.error(s"Fetch error: $fetchErr")
|
||||
case Right((name, content)) =>
|
||||
val fetchTask = FileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
||||
|
||||
dbService.save(name, content).attempt.flatMap {
|
||||
case Left(err) => log.error(s"error: $err")
|
||||
case Right(saveResult) => saveResult match {
|
||||
Stream.eval(Scheduler.of).flatMap { scheduler =>
|
||||
scheduler.scheduleTask(fetchTask)
|
||||
.evalMap {
|
||||
case Left(fetchErr) =>
|
||||
log.error(s"Fetch error: $fetchErr")
|
||||
case Right((name, content)) =>
|
||||
dbService.save(name, content).attempt.flatMap {
|
||||
case Left(err) => log.error(s"error: $err")
|
||||
case Right(savedName) => log.info(s"saved: $savedName")
|
||||
case Right(saveResult) => saveResult match {
|
||||
case Left(err) => log.error(s"error: $err")
|
||||
case Right(savedName) => log.info(s"saved: $savedName")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import cats.effect.IO
|
||||
import cats.effect.unsafe.implicits.global
|
||||
import cats.implicits.toTraverseOps
|
||||
import db.DBService
|
||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||
|
||||
import java.time.LocalDateTime
|
||||
|
||||
@@ -12,8 +13,10 @@ object Main {
|
||||
val from = LocalDateTime.of(2023, 4, 28, 10, 0)
|
||||
val to = LocalDateTime.of(2023, 4, 28, 13, 30)
|
||||
for {
|
||||
// fetchResultEither <- FetchService.fetchFromDate(LocalDate.of(2023, 4, 28)).attempt
|
||||
fetchResultEither <- FetchService.fetchInRange(from, to).attempt
|
||||
log <- Slf4jLogger.create[IO]
|
||||
fetch <- FetchService.of
|
||||
// fetchResultEither <- fetch.fetchFromDate(LocalDate.of(2023, 4, 28)).attempt
|
||||
fetchResultEither <- fetch.fetchInRange(from, to).attempt
|
||||
fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList
|
||||
fetchResult = fetchResultEither.getOrElse(List.empty)
|
||||
(fetchErrors, successDownloads) = fetchResult.partitionMap(identity)
|
||||
@@ -22,8 +25,8 @@ object Main {
|
||||
(saveErrors, successSaves) = saveResults.partitionMap(identity)
|
||||
successes = successDownloads.map(s => s"fetched: ${s._1}") ++ successSaves.map(s => s"saved: $s")
|
||||
errors = fetchServiceError ++ fetchErrors.map(e => s"FetchError: ${e.getMessage}") ++ saveErrors.map(e => s"SaveError: ${e.getMessage}")
|
||||
_ <- IO.println(s"errors: $errors")
|
||||
_ <- IO.println(s"successes: $successes")
|
||||
_ <- log.info(s"errors: $errors")
|
||||
_ <- log.info(s"successes: $successes")
|
||||
} yield (successes, errors)
|
||||
}
|
||||
|
||||
|
||||
@@ -4,15 +4,27 @@ import cats.effect._
|
||||
import cats.effect.unsafe.implicits.global
|
||||
import cats.implicits.catsSyntaxApply
|
||||
import fs2.Stream
|
||||
import org.typelevel.log4cats.Logger
|
||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||
import server.Server
|
||||
|
||||
import java.time.{Duration, LocalTime}
|
||||
import scala.concurrent.duration._
|
||||
|
||||
|
||||
object Scheduler {
|
||||
def of: IO[Scheduler] = {
|
||||
Slf4jLogger.create[IO].map {
|
||||
new Scheduler(_)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
class Scheduler(log: Logger[IO]) {
|
||||
private val downloadMinute = 31
|
||||
|
||||
// private def testTask: IO[Unit] = {
|
||||
// IO(println("Running task"))
|
||||
// log.info("Running task")
|
||||
// }
|
||||
|
||||
def durationToNextHalfHour(implicit clock: Clock[IO]): IO[FiniteDuration] = {
|
||||
@@ -34,16 +46,18 @@ object Scheduler {
|
||||
|
||||
def scheduleTask(task: IO[Either[Throwable, (String, String)]]): Stream[IO, Either[Throwable, (String, String)]] = {
|
||||
Stream.eval(durationToNextHalfHour).flatMap { delay => {
|
||||
Stream.eval(IO.println(s"Scheduler started with delay: ${delay.toMinutes} min")) *>
|
||||
Stream.eval(log.info(s"Scheduler started with delay: ${delay.toMinutes} min")) *>
|
||||
(Stream.sleep[IO](delay) ++ Stream.awakeEvery[IO](1.hour))
|
||||
.evalMap(_ => task)
|
||||
}}
|
||||
}
|
||||
|
||||
// This is just for testing
|
||||
def main(args: Array[String]): Unit = {
|
||||
// run.compile.drain.unsafeRunSync()
|
||||
|
||||
val fetchTask = FileNameService.generateCurrentHour.flatMap(FetchService.fetchSingleFile)
|
||||
scheduleTask(fetchTask).compile.drain.unsafeRunSync()
|
||||
for {
|
||||
fetch <- FetchService.of
|
||||
fetchTask = FileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
||||
} yield scheduleTask(fetchTask).compile.drain.unsafeRunSync()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user