From e9408391ac874f633aa88c6692da18163d6981fd Mon Sep 17 00:00:00 2001 From: Guntis Smaukstelis Date: Wed, 5 Mar 2025 23:25:42 +0200 Subject: [PATCH] Rest api for deleting old forecasts --- src/main/scala/data/DataService.scala | 21 +++++++++++++++++---- src/main/scala/data/DataServiceTest.scala | 3 ++- src/main/scala/server/Server.scala | 4 +++- 3 files changed, 22 insertions(+), 6 deletions(-) diff --git a/src/main/scala/data/DataService.scala b/src/main/scala/data/DataService.scala index c28a60b..d15b9ba 100644 --- a/src/main/scala/data/DataService.scala +++ b/src/main/scala/data/DataService.scala @@ -6,7 +6,10 @@ import cats.effect.kernel.Resource import cats.effect.std.Semaphore import cats.effect.unsafe.implicits.global import cats.implicits._ +import data.DataService.DeletionResult import fs2.io.file.{Files, Path} +import io.circe.Encoder +import io.circe.generic.semiauto.deriveEncoder import org.typelevel.log4cats.Logger import org.typelevel.log4cats.slf4j.Slf4jLogger import parse.grib.{Grib, GribParser} @@ -19,6 +22,10 @@ import scala.concurrent.ExecutionContext import scala.util.Try object DataService { + case class DeletionResult(fileName: String, success: Boolean, errorMessage: Option[String]) + + implicit val deletionResultEncoder: Encoder[DeletionResult] = deriveEncoder[DeletionResult] + def of: IO[DataService] = { for { logger <- Slf4jLogger.create[IO] @@ -107,7 +114,7 @@ class DataService(log: Logger[IO]) { } } - def deleteOldForecasts(maxHours: Int = 9): IO[List[String]] = { + def deleteOldForecasts(maxHours: Int = 9): IO[List[DeletionResult]] = { val nowUTC = ZonedDateTime.now(ZoneOffset.UTC) val ageThreshold = nowUTC.minusHours(maxHours) @@ -119,9 +126,15 @@ class DataService(log: Logger[IO]) { ) keepList = fileDateList.filter(_._2.isAfter(ageThreshold)).map(_._1) deleteList = fileList.filter(!keepList.contains(_)) - _ <- deleteList.traverse(name => Files[IO].delete(Path(s"$GRIB_FOLDER/${name}"))) - _ <- deleteList.traverse(name => log.info(s"delete: $name")) - } yield deleteList + results <- deleteList.traverse { name => + val path = Path(s"$GRIB_FOLDER/${name}") + Files[IO].delete(path).attempt.flatMap { + case Right(_) => log.info(s"delete: $name").as(DeletionResult(name, true, None)) + case Left(error) => log.error(s"Failed to delete $name: ${error.getMessage}") + .as(DeletionResult(name, false, Some(error.getMessage))) + } + } + } yield results } private def getTimeFromName(filename: String): Option[(ZonedDateTime, ZonedDateTime)] = { diff --git a/src/main/scala/data/DataServiceTest.scala b/src/main/scala/data/DataServiceTest.scala index b5c0128..eb87e9f 100644 --- a/src/main/scala/data/DataServiceTest.scala +++ b/src/main/scala/data/DataServiceTest.scala @@ -2,7 +2,8 @@ package data import cats.effect.IO import cats.effect.unsafe.implicits.global -import io.circe.syntax.EncoderOps +import data.DataService.deletionResultEncoder +import io.circe.syntax._ object DataServiceTest { def main(args: Array[String]): Unit = { diff --git a/src/main/scala/server/Server.scala b/src/main/scala/server/Server.scala index 071f765..f5f10dd 100644 --- a/src/main/scala/server/Server.scala +++ b/src/main/scala/server/Server.scala @@ -24,7 +24,6 @@ import parse.csv.Aggregate.{AggregateKey, UserQuery} import org.http4s.circe.jsonEncoder import org.typelevel.log4cats.Logger import org.typelevel.log4cats.slf4j.Slf4jLogger -import fs2.io.file.Path import parse.csv.Aggregate import scala.concurrent.duration.DurationInt @@ -68,6 +67,9 @@ class Server(postgresService: PostgresService, dataService: DataService, fetch: case GET -> Root / "grib" / "binary-chunk" / ValidateInt(binaryOffset) / ValidateInt(binaryLength) / fileName => dataService.getBinaryChunk(binaryOffset, binaryLength, fileName).flatMap(buffer => Ok(buffer)) + // http://0.0.0.0:8080/api/grib/delete-old-forecasts + case GET -> Root / "grib" / "delete-old-forecasts" => + dataService.deleteOldForecasts().flatMap(result => Ok(result.asJson))