FetchSingleFile with Either
This commit is contained in:
@@ -50,14 +50,17 @@ object FetchService {
|
||||
}
|
||||
}
|
||||
|
||||
def fetchSingleFile(fileName: String): IO[(String, String)] = {
|
||||
fetchFiles(List(fileName)).flatMap { results =>
|
||||
def fetchSingleFile(fileName: String): IO[Either[Throwable, (String, String)]] = {
|
||||
fetchFiles(List(fileName)).map { results =>
|
||||
results.headOption match {
|
||||
case Some(Right(result)) => IO.pure(result)
|
||||
case Some(Left(err)) => IO.raiseError(err)
|
||||
case None => IO.raiseError(new Exception("No file fetched"))
|
||||
case Some(Right(result)) =>
|
||||
IO.println(s"fetched: $fileName").as(Right(result))
|
||||
case Some(Left(err)) =>
|
||||
IO.println(s"failed fetch: $fileName with error: ${err.getMessage}").as(Left(err))
|
||||
case None =>
|
||||
IO.println(s"failed fetch: $fileName").as(Left(new Exception("No file fetched")))
|
||||
}
|
||||
}
|
||||
}.flatten
|
||||
}
|
||||
|
||||
def fetchInRange(from: LocalDateTime, to: LocalDateTime): IO[List[Either[Throwable, (String, String)]]] = {
|
||||
|
||||
@@ -18,11 +18,17 @@ class FileFetchScheduler(dbService: DBService, log: Logger[IO]) {
|
||||
def run: Stream[IO, Unit] = {
|
||||
val fetchTask = FileNameService.generateCurrentHour.flatMap(FetchService.fetchSingleFile)
|
||||
Scheduler.scheduleTask(fetchTask)
|
||||
.evalMap { case (name, content) =>
|
||||
log.info(s"fetched: $name") *>
|
||||
.evalMap {
|
||||
case Left(fetchErr) =>
|
||||
log.error(s"Fetch error: $fetchErr")
|
||||
case Right((name, content)) =>
|
||||
|
||||
dbService.save(name, content).attempt.flatMap {
|
||||
case Right(savedName) => log.info(s"saved: $savedName")
|
||||
case Left(err) => log.error(s"error: $err")
|
||||
case Right(saveResult) => saveResult match {
|
||||
case Left(err) => log.error(s"error: $err")
|
||||
case Right(savedName) => log.info(s"saved: $savedName")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,7 +32,7 @@ object Scheduler {
|
||||
// }
|
||||
// }
|
||||
|
||||
def scheduleTask(task: IO[(String, String)]): Stream[IO, (String, String)] = {
|
||||
def scheduleTask(task: IO[Either[Throwable, (String, String)]]): Stream[IO, Either[Throwable, (String, String)]] = {
|
||||
Stream.eval(durationToNextHalfHour).flatMap { delay => {
|
||||
Stream.eval(IO.println(s"Scheduler started with delay: ${delay.toMinutes} min")) *>
|
||||
(Stream.sleep[IO](delay) ++ Stream.awakeEvery[IO](1.hour))
|
||||
|
||||
Reference in New Issue
Block a user