Refactor FileFetchScheduler
This commit is contained in:
@@ -9,16 +9,18 @@ import org.typelevel.log4cats.slf4j.Slf4jLogger
|
|||||||
|
|
||||||
object FileFetchScheduler {
|
object FileFetchScheduler {
|
||||||
def of(dbService: DBService, fetch: FetchService): IO[FileFetchScheduler] = {
|
def of(dbService: DBService, fetch: FetchService): IO[FileFetchScheduler] = {
|
||||||
|
Scheduler.of.flatMap { scheduler =>
|
||||||
Slf4jLogger.create[IO].map {
|
Slf4jLogger.create[IO].map {
|
||||||
new FileFetchScheduler(dbService, fetch, _)
|
new FileFetchScheduler(dbService, fetch, scheduler, _)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
class FileFetchScheduler(dbService: DBService, fetch: FetchService, log: Logger[IO]) {
|
|
||||||
|
class FileFetchScheduler(dbService: DBService, fetch: FetchService, scheduler: Scheduler, log: Logger[IO]) {
|
||||||
def run: Stream[IO, Unit] = {
|
def run: Stream[IO, Unit] = {
|
||||||
val fetchTask = FileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
val fetchTask = FileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
||||||
|
|
||||||
Stream.eval(Scheduler.of).flatMap { scheduler =>
|
|
||||||
scheduler.scheduleTask(fetchTask)
|
scheduler.scheduleTask(fetchTask)
|
||||||
.evalMap {
|
.evalMap {
|
||||||
case Left(fetchErr) =>
|
case Left(fetchErr) =>
|
||||||
@@ -33,5 +35,4 @@ class FileFetchScheduler(dbService: DBService, fetch: FetchService, log: Logger[
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user