diff --git a/.env-sample b/.env-sample index 8de0d55..5b6bf86 100644 --- a/.env-sample +++ b/.env-sample @@ -4,4 +4,9 @@ POSTGRES_PASSWORD=aaa METEO_USER=aaa METEO_PASSWORD=aaa -METEO_URL=aaa \ No newline at end of file +METEO_URL=aaa + +HARMONIE_EDR_API_KEY=aaa +HARMONIE_EDR_URL=aaa +HARMONIE_STAC_API_KEY=aaa +HARMONIE_STAC_URL=aaa \ No newline at end of file diff --git a/Dockerfile b/Dockerfile index 3ac6458..7011088 100644 --- a/Dockerfile +++ b/Dockerfile @@ -3,6 +3,6 @@ FROM amazoncorretto:17-alpine COPY ./web/dist/ /web/dist/ COPY ./target/scala-2.13/WeatherTool-assembly-0.1.1-SNAPSHOT.jar /app.jar -VOLUME /data +VOLUME /app/data CMD ["java", "-jar", "/app.jar"] \ No newline at end of file diff --git a/Dockerfile.local b/Dockerfile.local index 950515c..bd1c4e5 100644 --- a/Dockerfile.local +++ b/Dockerfile.local @@ -10,4 +10,6 @@ FROM amazoncorretto:17-alpine COPY --from=build /app/web/dist/ /web/dist/ COPY --from=build /app/target/scala-2.13/WeatherTool-assembly-0.1.1-SNAPSHOT.jar /app.jar +VOLUME /app/data + CMD ["java", "-jar", "/app.jar"] diff --git a/fly.toml b/fly.toml index b2438a7..3beb42e 100644 --- a/fly.toml +++ b/fly.toml @@ -13,8 +13,9 @@ primary_region = "waw" auto_start_machines = true [[mounts]] - source = "data_volume" - destination = "/data" + source = "file_volume" + destination = "/app/data" + size_gb = 10 [[vm]] memory = 512 @@ -22,5 +23,5 @@ primary_region = "waw" cpus = 1 [scale] - count = 0 + count = 1 idle_timeout = 0 diff --git a/src/main/scala/app/Main.scala b/src/main/scala/app/Main.scala index 3ef82a6..203b895 100644 --- a/src/main/scala/app/Main.scala +++ b/src/main/scala/app/Main.scala @@ -19,12 +19,14 @@ object Main extends IOApp { fetchCsvTask = fileFetchScheduler.run.compile.drain scheduler <- fetchDMI.Scheduler.of - cleanupTask = scheduler.scheduleTask("Cleanup", List(1), DataService.deleteOldForecasts()).compile.drain + dataService <- DataService.of - fetchGrib <- fetchDMI.FetchService.of + cleanupTask = scheduler.scheduleTask("Cleanup", List(1), dataService.deleteOldForecasts()).compile.drain + + fetchGrib <- fetchDMI.FetchService.of(dataService) fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(2), fetchGrib.fetchRecentForecasts()).compile.drain - server <- Server.of(postgresService, fetch) + server <- Server.of(postgresService, dataService, fetch) serverTask = server.run exitCode <- (serverTask, fetchCsvTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ExitCode.Success) diff --git a/src/main/scala/data/DataService.scala b/src/main/scala/data/DataService.scala index 6ef12af..547b555 100644 --- a/src/main/scala/data/DataService.scala +++ b/src/main/scala/data/DataService.scala @@ -4,26 +4,49 @@ import cats.effect.IO import cats.implicits.toTraverseOps import fs2.io.file.{Files, Path} import grib.{Grib, GribParser} +import org.typelevel.log4cats.Logger +import org.typelevel.log4cats.slf4j.Slf4jLogger +import java.nio.file.Paths import java.time.{ZoneId, ZoneOffset, ZonedDateTime} import java.time.format.DateTimeFormatter import scala.io.Source import scala.util.Try object DataService { - val FOLDER = "data" + def of: IO[DataService] = { + for { + logger <- Slf4jLogger.create[IO] + service = new DataService(logger) + _ <- service.init + } yield service + } +} + + +class DataService(log: Logger[IO]) { + val BASE_FOLDER = "data" + val GRIB_FOLDER = s"$BASE_FOLDER/grib" + + def init: IO[Unit] = { + for { + _ <- fs2.io.file.Files[IO].createDirectories(fs2.io.file.Path(BASE_FOLDER)) + _ <- fs2.io.file.Files[IO].createDirectories(fs2.io.file.Path(GRIB_FOLDER)) + _ <- log.info(s"Created directories: $BASE_FOLDER and $GRIB_FOLDER") + } yield () + } def getFileList(): IO[List[String]] = Files[IO] - .list(Path(FOLDER)) + .list(Path(GRIB_FOLDER)) .map(_.toString) .filter(_.endsWith(".grib")) - .map(_.replace(s"$FOLDER/", "")) + .map(_.replace(s"$GRIB_FOLDER/", "")) .compile .toList def getGribStucture(fileName: String): IO[List[Grib]] = { - val filePath = Path(s"$FOLDER/$fileName") + val filePath = Path(s"$GRIB_FOLDER/$fileName") Files[IO].exists(filePath).flatMap { case true => GribParser.parseFile(filePath) @@ -33,7 +56,7 @@ object DataService { def getBinaryChunk(offset: Int, length: Int, fileName: String): IO[Array[Byte]] = { IO { - val source = Source.fromFile(s"$FOLDER/$fileName", "ISO-8859-1") + val source = Source.fromFile(s"$GRIB_FOLDER/$fileName", "ISO-8859-1") try { source.slice(offset, offset + length).map(_.toByte).toArray } finally { @@ -54,14 +77,14 @@ object DataService { val oldThreshold = nowUTC.minusHours(maxHours) for { - _ <- IO.println("start cleanup") + _ <- log.info("start cleanup") fileList <- getFileList() fileDateList = fileList.flatMap(fileName => getTimeFromName(fileName).map(extracted => (fileName, extracted._1)) ) deleteList = fileDateList.filter(_._2.isBefore(oldThreshold)).map(_._1) - _ <- deleteList.traverse(name => Files[IO].delete(Path(s"$FOLDER/${name}"))) - _ <- deleteList.traverse(name => IO.println(s"delete: $name")) + _ <- deleteList.traverse(name => Files[IO].delete(Path(s"$GRIB_FOLDER/${name}"))) + _ <- deleteList.traverse(name => log.info(s"delete: $name")) } yield deleteList } diff --git a/src/main/scala/data/DataServiceTest.scala b/src/main/scala/data/DataServiceTest.scala index aa56bff..b5c0128 100644 --- a/src/main/scala/data/DataServiceTest.scala +++ b/src/main/scala/data/DataServiceTest.scala @@ -12,8 +12,9 @@ object DataServiceTest { private def deleteOldForecasts(): IO[Unit] = { val program = for { - deleteList <- DataService.deleteOldForecasts() -// _ <- IO.println(deleteList.asJson) + dataService <- DataService.of + deleteList <- dataService.deleteOldForecasts() + _ <- IO.println(deleteList.asJson) } yield () program @@ -21,7 +22,8 @@ object DataServiceTest { private def getForecasts(): IO[Unit] = { val program = for { - forecasts <- DataService.getForecasts() + dataService <- DataService.of + forecasts <- dataService.getForecasts() _ <- IO.println(forecasts.asJson) } yield () diff --git a/src/main/scala/fetchDMI/FetchService.scala b/src/main/scala/fetchDMI/FetchService.scala index d8a0609..766317c 100644 --- a/src/main/scala/fetchDMI/FetchService.scala +++ b/src/main/scala/fetchDMI/FetchService.scala @@ -22,12 +22,12 @@ final case class HarmonieServerConfig( ) object FetchService { - def of: IO[FetchService] = { - Slf4jLogger.create[IO].map(logger => new FetchService(logger)) + def of(dataService: DataService): IO[FetchService] = { + Slf4jLogger.create[IO].map(logger => new FetchService(dataService, logger)) } } -class FetchService(log: Logger[IO]) { +class FetchService(dataService: DataService, log: Logger[IO]) { private val edrConfig: HarmonieServerConfig = ( sys.env.get("HARMONIE_EDR_API_KEY"), sys.env.get("HARMONIE_EDR_URL"), @@ -35,7 +35,7 @@ class FetchService(log: Logger[IO]) { case (Some(api_key), Some(url)) => HarmonieServerConfig(api_key, url) case _ => - throw new RuntimeException("Unable to load harmonie config: Missing required environment variables") + throw new RuntimeException("Unable to load harmonie edr config: Missing required environment variables") } private val stacConfig: HarmonieServerConfig = ( @@ -45,7 +45,7 @@ class FetchService(log: Logger[IO]) { case (Some(api_key), Some(url)) => HarmonieServerConfig(api_key, url) case _ => - throw new RuntimeException("Unable to load harmonie config: Missing required environment variables") + throw new RuntimeException("Unable to load harmonie stac config: Missing required environment variables") } /* @@ -77,7 +77,8 @@ class FetchService(log: Logger[IO]) { urlWithParams = Uri.unsafeFromString(s"${base.toString}?${queryParams.toString}") request = Request[IO](Method.GET, urlWithParams) - tmpPath = Path(s"${DataService.FOLDER}/tmp.grib") + // TODO proly better to call dataService method than property + tmpPath = Path(s"${dataService.GRIB_FOLDER}/tmp.grib") _ <- client.stream(request) .flatMap(_.body) .through(Files[IO].writeAll(tmpPath)) @@ -85,7 +86,7 @@ class FetchService(log: Logger[IO]) { .drain gribList <- GribParser.parseFile(tmpPath) gribTime = gribList.head.time - fileName = Path(s"${DataService.FOLDER}/harmonie_${gribTime.referenceTime}_${gribTime.forecastTime}.grib".replace(":", "")) + fileName = Path(s"${dataService.GRIB_FOLDER}/harmonie_${gribTime.referenceTime}_${gribTime.forecastTime}.grib".replace(":", "")) _ <- Files[IO].move(tmpPath, fileName, CopyFlags.apply(CopyFlag.ReplaceExisting)) fileSizeBytes <- Files[IO].size(fileName) fileSizeMB = fileSizeBytes.toDouble / (1024 * 1024) @@ -111,7 +112,7 @@ class FetchService(log: Logger[IO]) { for { availableResult <- fetchAvailableForecasts() (modelRun, forecastDateList) = availableResult - localForecasts <- DataService.getForecasts() + localForecasts <- dataService.getForecasts() toFetchList = forecastDateList.filter(dateTime => !localForecasts.contains((modelRun, dateTime))) } yield toFetchList } diff --git a/src/main/scala/fetchDMI/FetchServiceTest.scala b/src/main/scala/fetchDMI/FetchServiceTest.scala index b3dde79..9182d79 100644 --- a/src/main/scala/fetchDMI/FetchServiceTest.scala +++ b/src/main/scala/fetchDMI/FetchServiceTest.scala @@ -2,6 +2,7 @@ package fetchDMI; import cats.effect.IO import cats.effect.unsafe.implicits.global +import data.DataService import java.time.{ZoneOffset, ZonedDateTime} @@ -15,7 +16,8 @@ object FetchServiceTest { private def fetchRecentForecasts(): IO[Unit] = { val program = for { - fetch <- FetchService.of + dataService <- DataService.of + fetch <- FetchService.of(dataService) result <- fetch.fetchRecentForecasts() _ <- IO.println("-=fetch finished=-") } yield () @@ -24,7 +26,8 @@ object FetchServiceTest { private def generateFetchList(): IO[Unit] = { val program = for { - fetch <- FetchService.of + dataService <- DataService.of + fetch <- FetchService.of(dataService) list <- fetch.generateFetchList() _ <- IO.println(list) } yield () @@ -33,7 +36,8 @@ object FetchServiceTest { private def fetchAvailableForecasts(): IO[Unit] = { val program = for { - fetch <- FetchService.of + dataService <- DataService.of + fetch <- FetchService.of(dataService) result <- fetch.fetchAvailableForecasts() (modelRun, forecastTimes) = result _ <- IO.println(modelRun) @@ -48,7 +52,8 @@ object FetchServiceTest { nowUTC <- IO(ZonedDateTime.now(ZoneOffset.UTC)) referenceTime = FileName.getClosestReferenceTime(nowUTC) timeList = FileName.generateTimeList(referenceTime) - fetch <- FetchService.of + dataService <- DataService.of + fetch <- FetchService.of(dataService) nameList <- fetch.fetchFromList(timeList) _ <- IO.println(nameList) } yield () diff --git a/src/main/scala/grib/GribParserTest.scala b/src/main/scala/grib/GribParserTest.scala index 11f9fe6..f0f9e93 100644 --- a/src/main/scala/grib/GribParserTest.scala +++ b/src/main/scala/grib/GribParserTest.scala @@ -12,12 +12,10 @@ object GribParserTest { val gribTitle = Codes.codesToString(0, 0, 2) println(gribTitle) - -// val fileName = s"${DataService.FOLDER}/HARMONIE_DINI_SF_2025-01-24T030000Z_2025-01-26T010000Z.grib" - val fileName = s"${DataService.FOLDER}/harmonie_2025-02-01T1500Z_2025-02-01T180000Z.grib" - val path = Path(fileName) - val program = for { + dataService <- DataService.of + fileName = s"${dataService.GRIB_FOLDER}/harmonie_2025-02-01T1500Z_2025-02-01T180000Z.grib" + path = Path(fileName) gribList <- GribParser.parseFile(path) json = gribList.asJson _ <- IO.println(json.spaces2) diff --git a/src/main/scala/server/Server.scala b/src/main/scala/server/Server.scala index 1dbef52..680748a 100644 --- a/src/main/scala/server/Server.scala +++ b/src/main/scala/server/Server.scala @@ -31,14 +31,14 @@ import scala.concurrent.duration.DurationInt object Server { - def of(postgresService: PostgresService, fetch: FetchService): IO[Server] = { + def of(postgresService: PostgresService, dataService: DataService, fetch: FetchService): IO[Server] = { Slf4jLogger.create[IO].map { - new Server(postgresService, fetch, _) + new Server(postgresService, dataService, fetch, _) } } } -class Server(postgresService: PostgresService, fetch: FetchService, log: Logger[IO]) { +class Server(postgresService: PostgresService, dataService: DataService, fetch: FetchService, log: Logger[IO]) { // Define the extension method `pretty` for Json implicit class JsonPrettyPrinter(json: Json) { @@ -54,17 +54,17 @@ class Server(postgresService: PostgresService, fetch: FetchService, log: Logger[ // http://0.0.0.0:8080/api/show/grib-name/harmonie_2025-02-01T1500Z_2025-02-01T180000Z.grib case GET -> Root / "show" / "grib" / fileName => // val fileName = "data/HARMONIE_DINI_SF_2025-01-24T030000Z_2025-01-26T010000Z.grib" - DataService.getGribStucture(fileName).flatMap(response => Ok(response.asJson.pretty)) + dataService.getGribStucture(fileName).flatMap(response => Ok(response.asJson.pretty)) case GET -> Root / "show" / "gribName" => Ok("{\"fileName\":\"TODO replace this fake name\"}") // http://0.0.0.0:8080/api/show/grib-list case GET -> Root / "show" / "grib-list" => - DataService.getFileList().flatMap(fileList => Ok(fileList.asJson)) + dataService.getFileList().flatMap(fileList => Ok(fileList.asJson)) // TODO implement binary-chunk/ get request case GET -> Root / "grib" / "binary-chunk" / ValidateInt(binaryOffset) / ValidateInt(binaryLength) / fileName => - DataService.getBinaryChunk(binaryOffset, binaryLength, fileName).flatMap(buffer => Ok(buffer)) + dataService.getBinaryChunk(binaryOffset, binaryLength, fileName).flatMap(buffer => Ok(buffer))