Remove StatefulFetchService as we have already DataService with lst 24h state
This commit is contained in:
@@ -20,19 +20,13 @@ final case class WeatherServerConfig(
|
|||||||
url: String,
|
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 {
|
object FetchService {
|
||||||
def of: IO[FetchService] = {
|
def of: IO[FetchService] = {
|
||||||
Slf4jLogger.create[IO].map(logger => new FetchService(new FileNameService, logger))
|
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 {
|
private val weatherServerConfig: WeatherServerConfig = ConfigSource.default.load[WeatherServerConfig] match {
|
||||||
case Right(config) => config
|
case Right(config) => config
|
||||||
case Left(errors) => throw new RuntimeException(s"Unable to load config: $errors")
|
case Left(errors) => throw new RuntimeException(s"Unable to load config: $errors")
|
||||||
|
|||||||
@@ -1,14 +1,14 @@
|
|||||||
package fetch
|
package fetch
|
||||||
|
|
||||||
import cats.effect.IO
|
import cats.effect.IO
|
||||||
import db.{DBService, DataServiceTrait}
|
import db.DataServiceTrait
|
||||||
import fs2.Stream
|
import fs2.Stream
|
||||||
import org.typelevel.log4cats.Logger
|
import org.typelevel.log4cats.Logger
|
||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
|
|
||||||
|
|
||||||
object FileFetchScheduler {
|
object FileFetchScheduler {
|
||||||
def of(dataService: DataServiceTrait, fetch: FetchServiceTrait): IO[FileFetchScheduler] = {
|
def of(dataService: DataServiceTrait, fetch: FetchService): IO[FileFetchScheduler] = {
|
||||||
Scheduler.of.flatMap { scheduler =>
|
Scheduler.of.flatMap { scheduler =>
|
||||||
Slf4jLogger.create[IO].map {
|
Slf4jLogger.create[IO].map {
|
||||||
new FileFetchScheduler(dataService, fetch, new FileNameService(), scheduler, _)
|
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] = {
|
def run: Stream[IO, Unit] = {
|
||||||
val fetchTask = fileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
val fetchTask = fileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
||||||
scheduler.scheduleTask(fetchTask)
|
scheduler.scheduleTask(fetchTask)
|
||||||
|
|||||||
@@ -33,9 +33,8 @@ object Main {
|
|||||||
def run: IO[Unit] = {
|
def run: IO[Unit] = {
|
||||||
for {
|
for {
|
||||||
fetch <- FetchService.of
|
fetch <- FetchService.of
|
||||||
statefulFetch <- StatefulFetchService.of(fetch)
|
fetchResultEither <- fetch.fetchSingleFile("20230524_0030.csv").attempt
|
||||||
fetchResultEither <- statefulFetch.fetchSingleFile("20230524_0030.csv").attempt
|
fetchResultEither <- fetch.fetchSingleFile("20230522_0130.csv").attempt
|
||||||
fetchResultEither <- statefulFetch.fetchSingleFile("20230522_0130.csv").attempt
|
|
||||||
fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList
|
fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList
|
||||||
fetchResult = fetchResultEither.flatMap(res => res.flatMap(aaa => {
|
fetchResult = fetchResultEither.flatMap(res => res.flatMap(aaa => {
|
||||||
println(s"fffffff: ${aaa._1}")
|
println(s"fffffff: ${aaa._1}")
|
||||||
|
|||||||
@@ -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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@@ -2,7 +2,7 @@ package server
|
|||||||
import cats.effect._
|
import cats.effect._
|
||||||
import cats.implicits.catsSyntaxTuple2Parallel
|
import cats.implicits.catsSyntaxTuple2Parallel
|
||||||
import db.{DBService, DataService}
|
import db.{DBService, DataService}
|
||||||
import fetch.{FetchService, FileFetchScheduler, StatefulFetchService}
|
import fetch.{FetchService, FileFetchScheduler}
|
||||||
|
|
||||||
object Main extends IOApp {
|
object Main extends IOApp {
|
||||||
def run(args: List[String]): IO[ExitCode] = {
|
def run(args: List[String]): IO[ExitCode] = {
|
||||||
@@ -10,10 +10,9 @@ object Main extends IOApp {
|
|||||||
dbService <- DBService.of
|
dbService <- DBService.of
|
||||||
dataService <- DataService.of(dbService)
|
dataService <- DataService.of(dbService)
|
||||||
fetch <- FetchService.of
|
fetch <- FetchService.of
|
||||||
statefulFetch <- StatefulFetchService.of(fetch)
|
fileFetchScheduler <- FileFetchScheduler.of(dataService, fetch)
|
||||||
fileFetchScheduler <- FileFetchScheduler.of(dataService, statefulFetch)
|
|
||||||
schedulerTask = fileFetchScheduler.run.compile.drain
|
schedulerTask = fileFetchScheduler.run.compile.drain
|
||||||
server <- Server.of(dataService, statefulFetch)
|
server <- Server.of(dataService, fetch)
|
||||||
serverTask = server.run
|
serverTask = server.run
|
||||||
exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success)
|
exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success)
|
||||||
} yield exitCode
|
} yield exitCode
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import cats.effect._
|
|||||||
import cats.implicits.toTraverseOps
|
import cats.implicits.toTraverseOps
|
||||||
import com.comcast.ip4s.IpLiteralSyntax
|
import com.comcast.ip4s.IpLiteralSyntax
|
||||||
import db.DataService
|
import db.DataService
|
||||||
import fetch.FetchServiceTrait
|
import fetch.FetchService
|
||||||
import parse.{Parser, WeatherData}
|
import parse.{Parser, WeatherData}
|
||||||
import server.ValidateRoutes.{AggKey, CityList, DateTimeRange, ValidDate}
|
import server.ValidateRoutes.{AggKey, CityList, DateTimeRange, ValidDate}
|
||||||
import io.circe.{Json, Printer}
|
import io.circe.{Json, Printer}
|
||||||
@@ -27,14 +27,14 @@ import scala.concurrent.duration.DurationInt
|
|||||||
|
|
||||||
|
|
||||||
object Server {
|
object Server {
|
||||||
def of(dataService: DataService, fetch: FetchServiceTrait): IO[Server] = {
|
def of(dataService: DataService, fetch: FetchService): IO[Server] = {
|
||||||
Slf4jLogger.create[IO].map {
|
Slf4jLogger.create[IO].map {
|
||||||
new Server(dataService, fetch, _)
|
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
|
// Define the extension method `pretty` for Json
|
||||||
implicit class JsonPrettyPrinter(json: Json) {
|
implicit class JsonPrettyPrinter(json: Json) {
|
||||||
|
|||||||
Reference in New Issue
Block a user