Rest api for deleting old forecasts
This commit is contained in:
@@ -6,7 +6,10 @@ import cats.effect.kernel.Resource
|
|||||||
import cats.effect.std.Semaphore
|
import cats.effect.std.Semaphore
|
||||||
import cats.effect.unsafe.implicits.global
|
import cats.effect.unsafe.implicits.global
|
||||||
import cats.implicits._
|
import cats.implicits._
|
||||||
|
import data.DataService.DeletionResult
|
||||||
import fs2.io.file.{Files, Path}
|
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.Logger
|
||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
import parse.grib.{Grib, GribParser}
|
import parse.grib.{Grib, GribParser}
|
||||||
@@ -19,6 +22,10 @@ import scala.concurrent.ExecutionContext
|
|||||||
import scala.util.Try
|
import scala.util.Try
|
||||||
|
|
||||||
object DataService {
|
object DataService {
|
||||||
|
case class DeletionResult(fileName: String, success: Boolean, errorMessage: Option[String])
|
||||||
|
|
||||||
|
implicit val deletionResultEncoder: Encoder[DeletionResult] = deriveEncoder[DeletionResult]
|
||||||
|
|
||||||
def of: IO[DataService] = {
|
def of: IO[DataService] = {
|
||||||
for {
|
for {
|
||||||
logger <- Slf4jLogger.create[IO]
|
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 nowUTC = ZonedDateTime.now(ZoneOffset.UTC)
|
||||||
val ageThreshold = nowUTC.minusHours(maxHours)
|
val ageThreshold = nowUTC.minusHours(maxHours)
|
||||||
|
|
||||||
@@ -119,9 +126,15 @@ class DataService(log: Logger[IO]) {
|
|||||||
)
|
)
|
||||||
keepList = fileDateList.filter(_._2.isAfter(ageThreshold)).map(_._1)
|
keepList = fileDateList.filter(_._2.isAfter(ageThreshold)).map(_._1)
|
||||||
deleteList = fileList.filter(!keepList.contains(_))
|
deleteList = fileList.filter(!keepList.contains(_))
|
||||||
_ <- deleteList.traverse(name => Files[IO].delete(Path(s"$GRIB_FOLDER/${name}")))
|
results <- deleteList.traverse { name =>
|
||||||
_ <- deleteList.traverse(name => log.info(s"delete: $name"))
|
val path = Path(s"$GRIB_FOLDER/${name}")
|
||||||
} yield deleteList
|
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)] = {
|
private def getTimeFromName(filename: String): Option[(ZonedDateTime, ZonedDateTime)] = {
|
||||||
|
|||||||
@@ -2,7 +2,8 @@ package data
|
|||||||
|
|
||||||
import cats.effect.IO
|
import cats.effect.IO
|
||||||
import cats.effect.unsafe.implicits.global
|
import cats.effect.unsafe.implicits.global
|
||||||
import io.circe.syntax.EncoderOps
|
import data.DataService.deletionResultEncoder
|
||||||
|
import io.circe.syntax._
|
||||||
|
|
||||||
object DataServiceTest {
|
object DataServiceTest {
|
||||||
def main(args: Array[String]): Unit = {
|
def main(args: Array[String]): Unit = {
|
||||||
|
|||||||
@@ -24,7 +24,6 @@ import parse.csv.Aggregate.{AggregateKey, UserQuery}
|
|||||||
import org.http4s.circe.jsonEncoder
|
import org.http4s.circe.jsonEncoder
|
||||||
import org.typelevel.log4cats.Logger
|
import org.typelevel.log4cats.Logger
|
||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
import fs2.io.file.Path
|
|
||||||
import parse.csv.Aggregate
|
import parse.csv.Aggregate
|
||||||
|
|
||||||
import scala.concurrent.duration.DurationInt
|
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 =>
|
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))
|
||||||
|
|
||||||
|
// 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))
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user