diff --git a/src/main/scala/Main.scala b/src/main/scala/Main.scala index 711b537..d1167eb 100644 --- a/src/main/scala/Main.scala +++ b/src/main/scala/Main.scala @@ -5,6 +5,7 @@ import data.DataService import db.{DBConnection, PostgresService} import fetch.csv.{FetchService, FileNameService} import fetch.dmi +import fetch.lvgmc import scheduler.Scheduler import server.Server @@ -22,25 +23,28 @@ object Main extends IOApp { scheduler <- Scheduler.of dataService <- DataService.of - fetchService <- FetchService.of + fetchLegacyService <- FetchService.of + fetchLvgmcService <- lvgmc.FetchService.of - fetchCsvList = new FileNameService().generateCurrentHour - .flatMap(fetchService.fetchSingleFile) - .flatMap { - case Right((name, content)) => postgresService.save(name, content) - case Left(_) => IO.unit // ignore error - } - fetchCsvTask = scheduler.scheduleTask("Fetch CSV", List(31), fetchCsvList).compile.drain - cleanupTask = scheduler.scheduleTask("Cleanup", List(1), dataService.deleteOldForecasts()).compile.drain + fetchWeatherStations = for { + fileName <- new FileNameService().generateCurrentHour + stationDataStr <- fetchLvgmcService.fetchWeatherStations() + _ <- IO.println(fileName) + // TODO refactor after legacy deleted. No need for fileName, instead pass only year + _ <- postgresService.save(fileName, stationDataStr) + } yield () + fetchStationsTask = scheduler.scheduleTask("Fetch Weather Stations", List(11,13,23,30), fetchWeatherStations).compile.drain + + cleanupTask = scheduler.scheduleTask("Cleanup old Grib", List(41), dataService.deleteOldForecasts()).compile.drain fetchGrib <- dmi.FetchService.of(dataService) - fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(3), fetchGrib.fetchRecentForecasts()).compile.drain + fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(43), fetchGrib.fetchRecentForecasts()).compile.drain - server <- Server.of(postgresService, dataService, fetchService) + server <- Server.of(postgresService, dataService, fetchLegacyService) serverTask = server.run - exitCode <- (serverTask, fetchCsvTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ExitCode.Success) + exitCode <- (serverTask, fetchStationsTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ExitCode.Success) } yield exitCode program.handleErrorWith { error => diff --git a/src/main/scala/fetch/lvgmc/FetchService.scala b/src/main/scala/fetch/lvgmc/FetchService.scala index a94ca3d..56e22a8 100644 --- a/src/main/scala/fetch/lvgmc/FetchService.scala +++ b/src/main/scala/fetch/lvgmc/FetchService.scala @@ -6,6 +6,7 @@ import org.typelevel.log4cats.Logger import org.typelevel.log4cats.slf4j.Slf4jLogger import java.io.ByteArrayOutputStream +import java.nio.charset.StandardCharsets final case class LVGMCServerConfig( username: String, @@ -85,4 +86,8 @@ class FetchService(log: Logger[IO]) { } yield content } } + + def fetchWeatherStations(): IO[String] = { + fetchFile("Latvija_faktiskais_laiks.csv").map(result => new String(result, StandardCharsets.UTF_8)) + } } \ No newline at end of file diff --git a/src/main/scala/fetch/lvgmc/FetchServiceTest.scala b/src/main/scala/fetch/lvgmc/FetchServiceTest.scala index ffa4f9e..09b8330 100644 --- a/src/main/scala/fetch/lvgmc/FetchServiceTest.scala +++ b/src/main/scala/fetch/lvgmc/FetchServiceTest.scala @@ -2,12 +2,18 @@ package fetch.lvgmc import cats.effect.IO import cats.effect.unsafe.implicits.global +import db.{DBConnection, PostgresService} +import fetch.csv.FileNameService + +import java.nio.charset.StandardCharsets object FetchServiceTest { def main(args: Array[String]): Unit = { - fetchFile().unsafeRunSync() +// fetchFile().unsafeRunSync() + fetchAndSave().unsafeRunSync() } + // Eiropa_LTV_pilsetas_nakama_dn.csv // Eiropa_LTV_pilsetas_tekosa_dn.csv // Latvija_LTV_pilsetas_nakama_dnn.csv @@ -22,4 +28,25 @@ object FetchServiceTest { } yield () program } + + private def fetchAndSave (): IO[Unit] = { + var program = for { + fileName <- new FileNameService().generateCurrentHour + _ <- IO.println(fileName) + + fetch <- FetchService.of + result <- fetch.fetchFile("Latvija_faktiskais_laiks.csv") + resultStr = new String(result, StandardCharsets.UTF_8) + _ <- IO.println(s"fetched file, size: ${result.length}") +// _ <- IO.println(resultStr) + + transactor <- DBConnection.transactor[IO] + postgresService <- PostgresService.of(transactor) + + _ <- postgresService.save(fileName, resultStr) + } yield () + program + } + + }