-
Notifications
You must be signed in to change notification settings - Fork 0
/
package.scala
35 lines (30 loc) · 1020 Bytes
/
package.scala
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
package forex
package app
package stream
import programs.rates.Algebra
import core.rates.domains.Pair
import core.rates.domains.Rate
import core.rates.domains.Currency
import io.chrisdavenport.log4cats.Logger
import cats.effect.concurrent.Ref
import cats.effect.Timer
import cats.Functor
package object updater {
import scala.concurrent.duration._
def stream[F[_]: Functor: Timer: Logger](
updateEvery: FiniteDuration,
clientRateAlg: Algebra[F],
mapRef: Ref[F, Map[Pair, Rate]]
): fs2.Stream[F, Unit] =
fs2.Stream.eval(Logger[F].info(s"Starting the cache updater. Updating every $updateEvery")) ++
fs2.Stream
.every[F](updateEvery)
.evalMap(_ => Logger[F].info(s"Updating the cache"))
.flatMap(_ => clientRateAlg.allRates)
.chunkN(Currency.allCombinationsLength)
.evalMap { rates =>
val newMap = rates.map(r => r.pair -> r).toList.toMap
mapRef.set(newMap)
}
.flatMap(_ => fs2.Stream.sleep(updateEvery))
}