diff --git a/forex-mtl/build.sbt b/forex-mtl/build.sbt index 8994026f..12c057fc 100644 --- a/forex-mtl/build.sbt +++ b/forex-mtl/build.sbt @@ -67,3 +67,10 @@ libraryDependencies ++= Seq( Libraries.scalaCheck % Test, Libraries.catsScalaCheck % Test ) + +libraryDependencies ++= Seq( + "net.debasishg" %% "redisclient" % "3.41" +) + +libraryDependencies += "com.lihaoyi" %% "requests" % "0.8.0" + diff --git a/forex-mtl/src/main/resources/application.conf b/forex-mtl/src/main/resources/application.conf index b2af6efd..5cb5e54f 100644 --- a/forex-mtl/src/main/resources/application.conf +++ b/forex-mtl/src/main/resources/application.conf @@ -1,8 +1,11 @@ app { http { host = "0.0.0.0" - port = 8080 + port = 8081 timeout = 40 seconds } + oneframe { + token = "10dc303535874aeccc86a8251e6992f5" + } } diff --git a/forex-mtl/src/main/scala/forex/Main.scala b/forex-mtl/src/main/scala/forex/Main.scala index 6dda10a7..950d096f 100644 --- a/forex-mtl/src/main/scala/forex/Main.scala +++ b/forex-mtl/src/main/scala/forex/Main.scala @@ -8,9 +8,14 @@ import org.http4s.blaze.server.BlazeServerBuilder object Main extends IOApp { - override def run(args: List[String]): IO[ExitCode] = - new Application[IO].stream(executionContext).compile.drain.as(ExitCode.Success) - + override def run(args: List[String]): IO[ExitCode] = { + val startupTask = IO {StartupTasks.process()} + val app: IO[ExitCode] = new Application[IO].stream(executionContext).compile.drain.as(ExitCode.Success) + for { + _ <- startupTask.start //This will run async and starts another task after 4minutes to refresh cache... + exitCode <- app //This will run the app in current thread + } yield { exitCode} + } } class Application[F[_]: ConcurrentEffect: Timer] { @@ -24,5 +29,4 @@ class Application[F[_]: ConcurrentEffect: Timer] { .withHttpApp(module.httpApp) .serve } yield () - } diff --git a/forex-mtl/src/main/scala/forex/StartupTasks.scala b/forex-mtl/src/main/scala/forex/StartupTasks.scala new file mode 100644 index 00000000..6eb34b95 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/StartupTasks.scala @@ -0,0 +1,12 @@ +package forex + +import cats.effect.{ExitCode, IO} +import forex.components.scheduler.tasks.TaskScheduler +import forex.thirdparties.oneframae.OneFrameForexRatesHandler + +object StartupTasks { + def process(): IO[ExitCode] = { + TaskScheduler.scheduleForexCacheRefreshJob + IO(OneFrameForexRatesHandler.handleCache()).as(ExitCode.Success) + } +} diff --git a/forex-mtl/src/main/scala/forex/components/cache/Algebra.scala b/forex-mtl/src/main/scala/forex/components/cache/Algebra.scala new file mode 100644 index 00000000..ea8e8f16 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/cache/Algebra.scala @@ -0,0 +1,11 @@ +package forex.components.cache + +import forex.services.rates.errors._ + + +trait Algebra { + + def put[A](key:String, value: A) : Boolean + + def get(key:String) : Either[Error, String] +} diff --git a/forex-mtl/src/main/scala/forex/components/cache/Protocol.scala b/forex-mtl/src/main/scala/forex/components/cache/Protocol.scala new file mode 100644 index 00000000..f4b69c45 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/cache/Protocol.scala @@ -0,0 +1,7 @@ +package forex.components.cache + +import forex.components.cache.redis.Redis + +object Protocol { + def redis(segment: String) : Algebra = new Redis(segment) +} diff --git a/forex-mtl/src/main/scala/forex/components/cache/redis/Redis.scala b/forex-mtl/src/main/scala/forex/components/cache/redis/Redis.scala new file mode 100644 index 00000000..5907df2f --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/cache/redis/Redis.scala @@ -0,0 +1,24 @@ +package forex.components.cache.redis + +import com.redis.RedisClient +import forex.components.cache.Algebra +import forex.programs.rates.ErrorCodes +import forex.services.rates.errors._ + +class Redis(segment: String) extends Algebra { + + private[redis] lazy val redisClient: RedisClient = new RedisClient("localhost", 6379) + + override def put[A](key: String, value: A): Boolean = { + redisClient.set(key = ForexRedisHelper.getSegmentPrefixedKey(segment, key), value = value) + } + + override def get(key:String): Either[Error, String] = { + redisClient.get(ForexRedisHelper.getSegmentPrefixedKey(segment, key)) match { + case Some(value) => Right(value).withLeft[Error] + case None => Left(Error.RateLookupFailed(ErrorCodes.cacheFetchFailed, "Value Not Found In Cache")) + } + + } + +} diff --git a/forex-mtl/src/main/scala/forex/components/cache/redis/RedisHelper.scala b/forex-mtl/src/main/scala/forex/components/cache/redis/RedisHelper.scala new file mode 100644 index 00000000..ea1e2471 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/cache/redis/RedisHelper.scala @@ -0,0 +1,17 @@ +package forex.components.cache.redis + +abstract class RedisHelper { + def getSegmentPrefixedKey(segement: String, key:String): String +} + +object Segments { + final val forexRates = "FOREX_RATES" +} + +object ForexRedisHelper extends RedisHelper { + override def getSegmentPrefixedKey(segment: String, key: String): String = segment ++ "::" ++ key + + def getFormattedKey(from: String, to: String): String = { + from + "&" + to + } +} diff --git a/forex-mtl/src/main/scala/forex/components/package.scala b/forex-mtl/src/main/scala/forex/components/package.scala new file mode 100644 index 00000000..f3323028 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/package.scala @@ -0,0 +1,8 @@ +package forex + +import forex.components.cache.Algebra + +package object components { + type Cache = Algebra + final val CacheAPI = cache.Protocol +} diff --git a/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TaskScheduler.scala b/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TaskScheduler.scala new file mode 100644 index 00000000..b99ad6a1 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TaskScheduler.scala @@ -0,0 +1,45 @@ +package forex.components.scheduler.tasks + +import cats.effect.{Concurrent, IO, IOApp, Timer} +import cats.implicits.toFlatMapOps +import forex.domain.Tasks +import forex.thirdparties.oneframae.OneFrameForexRatesHandler +import org.log4s.{Logger, getLogger} + +import scala.concurrent.duration._ + +class TaskScheduler(task: Tasks, f: IO[Unit], frequency: FiniteDuration) extends IOApp .Simple { + + val logger: Logger = getLogger(getClass) + + private def schedule[F[_] : Concurrent : Timer](task: F[Unit], frequency: FiniteDuration) : F[Unit] = { + Timer[F].sleep(frequency).flatMap(_ => task) + } + + private[tasks] def scheduleJob(): Unit = { + if (TasksLibrary.isTaskSpawned(this.task)) { + logger.info("Job Already Spawned...") + } else { + schedule[IO](this.f, this.frequency) + TasksLibrary.addSpawnedTask(task) match { + case Some(_) => logger.info("Task Successfully Added To Task Library.") + case None => logger.warn("Task Addition to Library Failed. Audit and fix the issue.") + } + } + } + + override def run: IO[Unit] = { + IO(scheduleJob()) + } +} + +object TaskScheduler { + def apply(task: Tasks, f: IO[Unit], frequency: FiniteDuration): TaskScheduler = { + TaskScheduler(task, f, frequency) + } + + def scheduleForexCacheRefreshJob = { + TaskScheduler(Tasks.FOREX_JOB, OneFrameForexRatesHandler.handleCache(), 4.minutes) + } +} + diff --git a/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TasksLibrary.scala b/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TasksLibrary.scala new file mode 100644 index 00000000..1e15abee --- /dev/null +++ b/forex-mtl/src/main/scala/forex/components/scheduler/tasks/TasksLibrary.scala @@ -0,0 +1,13 @@ +package forex.components.scheduler.tasks + +import forex.domain.Tasks + +import scala.collection.mutable + +object TasksLibrary { + private val tasksSpawned: mutable.Map[String, Boolean] = mutable.Map[String, Boolean]() + + def isTaskSpawned(task: Tasks) : Boolean = tasksSpawned(task.name()) + + private[tasks] def addSpawnedTask(task: Tasks) : Option[Boolean] = tasksSpawned.put(task.name(), true) +} diff --git a/forex-mtl/src/main/scala/forex/domain/Currency.scala b/forex-mtl/src/main/scala/forex/domain/Currency.scala index a6f2857d..1f726138 100644 --- a/forex-mtl/src/main/scala/forex/domain/Currency.scala +++ b/forex-mtl/src/main/scala/forex/domain/Currency.scala @@ -15,6 +15,8 @@ object Currency { case object SGD extends Currency case object USD extends Currency + val supportedCurrencies : Set[Currency] = Set(AUD, CAD, CHF, EUR, GBP, NZD, JPY, SGD, USD) + implicit val show: Show[Currency] = Show.show { case AUD => "AUD" case CAD => "CAD" diff --git a/forex-mtl/src/main/scala/forex/domain/Tasks.scala b/forex-mtl/src/main/scala/forex/domain/Tasks.scala new file mode 100644 index 00000000..955a6d0b --- /dev/null +++ b/forex-mtl/src/main/scala/forex/domain/Tasks.scala @@ -0,0 +1,14 @@ +package forex.domain + +sealed trait Tasks { + protected var taskName: String = null + protected def setName() : Unit + def name() : String = this.taskName +} + +object Tasks { + case object FOREX_JOB extends Tasks { + override def setName(): Unit = this.taskName = "FOREX_JOB" + } +} + diff --git a/forex-mtl/src/main/scala/forex/http/rates/Protocol.scala b/forex-mtl/src/main/scala/forex/http/rates/Protocol.scala index 75391f9d..2eb9b295 100644 --- a/forex-mtl/src/main/scala/forex/http/rates/Protocol.scala +++ b/forex-mtl/src/main/scala/forex/http/rates/Protocol.scala @@ -4,6 +4,7 @@ package rates import forex.domain.Currency.show import forex.domain.Rate.Pair import forex.domain._ +import forex.programs.rates.errors import io.circe._ import io.circe.generic.extras.Configuration import io.circe.generic.extras.semiauto.deriveConfiguredEncoder @@ -36,4 +37,10 @@ object Protocol { implicit val responseEncoder: Encoder[GetApiResponse] = deriveConfiguredEncoder[GetApiResponse] + implicit val errorEncoder: Encoder[errors.Error] = + Encoder.instance { + case error: errors.Error.RateLookupFailed => rateLookupFailedEncoder(error) + } + implicit val rateLookupFailedEncoder: Encoder[errors.Error.RateLookupFailed] = + deriveConfiguredEncoder[errors.Error.RateLookupFailed] _ } diff --git a/forex-mtl/src/main/scala/forex/http/rates/RatesHttpRoutes.scala b/forex-mtl/src/main/scala/forex/http/rates/RatesHttpRoutes.scala index d91dcffb..b3ea8623 100644 --- a/forex-mtl/src/main/scala/forex/http/rates/RatesHttpRoutes.scala +++ b/forex-mtl/src/main/scala/forex/http/rates/RatesHttpRoutes.scala @@ -3,11 +3,15 @@ package rates import cats.effect.Sync import cats.syntax.flatMap._ +import forex.domain.Rate import forex.programs.RatesProgram -import forex.programs.rates.{ Protocol => RatesProgramProtocol } +import forex.programs.rates.{ErrorCodes, errors, Protocol => RatesProgramProtocol} +import io.circe.syntax.EncoderOps import org.http4s.HttpRoutes import org.http4s.dsl.Http4sDsl import org.http4s.server.Router +import org.log4s.{Logger, getLogger} + class RatesHttpRoutes[F[_]: Sync](rates: RatesProgram[F]) extends Http4sDsl[F] { @@ -15,10 +19,21 @@ class RatesHttpRoutes[F[_]: Sync](rates: RatesProgram[F]) extends Http4sDsl[F] { private[http] val prefixPath = "/rates" + val logger: Logger = getLogger(getClass) + private val httpRoutes: HttpRoutes[F] = HttpRoutes.of[F] { case GET -> Root :? FromQueryParam(from) +& ToQueryParam(to) => - rates.get(RatesProgramProtocol.GetRatesRequest(from, to)).flatMap(Sync[F].fromEither).flatMap { rate => - Ok(rate.asGetApiResponse) + try { + val response: F[Either[errors.Error, Rate]] = rates.get(RatesProgramProtocol.GetRatesRequest(from, to)) + response.flatMap { + case Right(_) => response.flatMap(Sync[F].fromEither).flatMap { rate => Ok(rate.asGetApiResponse) } + case Left(value) => BadRequest(value.asJson) + } + } catch { + case exception: Exception => + logger.error(exception)("Exception Occured in Forex Service.") + val error: errors.Error = errors.Error.RateLookupFailed(ErrorCodes.internalError, "Rate Lookup Failed due to " + exception.getMessage()) + BadRequest(error.asJson) } } diff --git a/forex-mtl/src/main/scala/forex/programs/package.scala b/forex-mtl/src/main/scala/forex/programs/package.scala index 9d381587..f68ef47d 100644 --- a/forex-mtl/src/main/scala/forex/programs/package.scala +++ b/forex-mtl/src/main/scala/forex/programs/package.scala @@ -3,4 +3,4 @@ package forex package object programs { type RatesProgram[F[_]] = rates.Algebra[F] final val RatesProgram = rates.Program -} +} \ No newline at end of file diff --git a/forex-mtl/src/main/scala/forex/programs/rates/ErrorCodes.scala b/forex-mtl/src/main/scala/forex/programs/rates/ErrorCodes.scala new file mode 100644 index 00000000..f02b14ff --- /dev/null +++ b/forex-mtl/src/main/scala/forex/programs/rates/ErrorCodes.scala @@ -0,0 +1,10 @@ +package forex.programs.rates + +object ErrorCodes { + final val cachePopulationFailed: String = "FX_CacheNotPopulated" + final val cacheFetchFailed: String = "FX_CacheFetchFailed" + final val cacheInitFailed: String = "FX_CachecInitializationFailed" + final val cacheRefreshFailed: String = "FX_CacheRefreshFailed" + final val fxRateLookUpFailed: String = "FX_RatesFetchFailed" + final val internalError: String = "FX_InternalError" +} diff --git a/forex-mtl/src/main/scala/forex/programs/rates/Program.scala b/forex-mtl/src/main/scala/forex/programs/rates/Program.scala index 528ee1f9..26642e19 100644 --- a/forex-mtl/src/main/scala/forex/programs/rates/Program.scala +++ b/forex-mtl/src/main/scala/forex/programs/rates/Program.scala @@ -11,7 +11,7 @@ class Program[F[_]: Functor]( ) extends Algebra[F] { override def get(request: Protocol.GetRatesRequest): F[Error Either Rate] = - EitherT(ratesService.get(Rate.Pair(request.from, request.to))).leftMap(toProgramError(_)).value + EitherT(ratesService.get(Rate.Pair(request.from, request.to))).leftMap(toProgramError).value } diff --git a/forex-mtl/src/main/scala/forex/programs/rates/errors.scala b/forex-mtl/src/main/scala/forex/programs/rates/errors.scala index 39496b13..4e4569c6 100644 --- a/forex-mtl/src/main/scala/forex/programs/rates/errors.scala +++ b/forex-mtl/src/main/scala/forex/programs/rates/errors.scala @@ -4,12 +4,15 @@ import forex.services.rates.errors.{ Error => RatesServiceError } object errors { - sealed trait Error extends Exception + class Error(code: String, message: String) extends Exception { + override def toString(): String = s"Code: $code, Message: $message" + } object Error { - final case class RateLookupFailed(msg: String) extends Error + final case class RateLookupFailed(code:String, msg: String) extends Error(code, msg) } def toProgramError(error: RatesServiceError): Error = error match { - case RatesServiceError.OneFrameLookupFailed(msg) => Error.RateLookupFailed(msg) + case RatesServiceError.OneFrameLookupFailed(code, msg) => Error.RateLookupFailed(code, msg) + case RatesServiceError.RateLookupFailed(code, msg) => Error.RateLookupFailed(code, msg) } } diff --git a/forex-mtl/src/main/scala/forex/services/rates/errors.scala b/forex-mtl/src/main/scala/forex/services/rates/errors.scala index 0584dcf4..2aab35bb 100644 --- a/forex-mtl/src/main/scala/forex/services/rates/errors.scala +++ b/forex-mtl/src/main/scala/forex/services/rates/errors.scala @@ -4,7 +4,8 @@ object errors { sealed trait Error object Error { - final case class OneFrameLookupFailed(msg: String) extends Error + final case class OneFrameLookupFailed(code:String, msg: String) extends Error + final case class RateLookupFailed(code:String, msg: String) extends Error } } diff --git a/forex-mtl/src/main/scala/forex/services/rates/interpreters/OneFrameDummy.scala b/forex-mtl/src/main/scala/forex/services/rates/interpreters/OneFrameDummy.scala index 37a3f50c..29548845 100644 --- a/forex-mtl/src/main/scala/forex/services/rates/interpreters/OneFrameDummy.scala +++ b/forex-mtl/src/main/scala/forex/services/rates/interpreters/OneFrameDummy.scala @@ -4,12 +4,29 @@ import forex.services.rates.Algebra import cats.Applicative import cats.syntax.applicative._ import cats.syntax.either._ -import forex.domain.{ Price, Rate, Timestamp } +import forex.domain.{Price, Rate, Timestamp} import forex.services.rates.errors._ +import forex.components._ +import forex.components.cache.redis.Segments + class OneFrameDummy[F[_]: Applicative] extends Algebra[F] { - override def get(pair: Rate.Pair): F[Error Either Rate] = - Rate(pair, Price(BigDecimal(100)), Timestamp.now).asRight[Error].pure[F] + private final val redisAPI: Cache = CacheAPI.redis(Segments.forexRates) + + override def get(pair: Rate.Pair): F[Error Either Rate] = { + val formattedCacheKey: String = getFormattedCacheKey(pair) + redisAPI.get(formattedCacheKey) match { + case Right(value) => Rate(pair, Price(BigDecimal(value.toDouble)), Timestamp.now) + .asRight[Error] + .pure[F] + case Left(value) => value + .asLeft[Rate] + .pure[F] + } + } + private def getFormattedCacheKey(pair: Rate.Pair) : String = { + pair.from.toString + "&" + pair.to.toString + } } diff --git a/forex-mtl/src/main/scala/forex/thirdparties/oneframae/OneFrameForexRatesHandler.scala b/forex-mtl/src/main/scala/forex/thirdparties/oneframae/OneFrameForexRatesHandler.scala new file mode 100644 index 00000000..fb97e506 --- /dev/null +++ b/forex-mtl/src/main/scala/forex/thirdparties/oneframae/OneFrameForexRatesHandler.scala @@ -0,0 +1,83 @@ +package forex.thirdparties.oneframae + +import cats.effect.IO +import com.typesafe.config.{ConfigFactory, ConfigRenderOptions} +import forex.components.CacheAPI +import forex.components.cache.redis.{ForexRedisHelper, Segments} +import forex.domain.Currency +import forex.programs.rates.ErrorCodes +import io.circe.{Json, ParsingFailure} +import io.circe.parser._ +import requests.Response +import forex.programs.rates.errors._ +import org.log4s.{Logger, getLogger} + +object OneFrameForexRatesHandler { + + val logger: Logger = getLogger(getClass) + + def handleCache(): IO[Unit] = { + populateCacheFromOneFrameService() match { + case Right(_) => IO(logger.info("Successfully Refreshed Cache.")) + case Left(value) => IO(logger.warn(s"Cache Refresh Failed due to ${value.getMessage}. Have a quick audit and fix it.")) + } + } + + private def populateCacheFromOneFrameService(): Either[Error, Boolean] = { + val redis = CacheAPI.redis(Segments.forexRates) + + fetchAllExchangeRates match { + case Right(value) => + value.foreach(json => { + + redis.put(ForexRedisHelper.getFormattedKey(json.hcursor.get[String]("from").toOption.get + , json.hcursor.get[String]("to").toOption.get) + , json.hcursor.get[BigDecimal]("price").toOption.get) + }) + Option(true).toRight(Error.RateLookupFailed(ErrorCodes.cacheInitFailed, "Cache Init Failed. Aborting Service Start!")) + case Left(value) => Option(value).toLeft(false) + } + } + + private def fetchAllExchangeRates: Either[Error, Vector[Json]] = { + val queryParam = constructCurrencyCombinations + .map(x => s"pair=$x") + .reduce((x,y) => x + "&" + y) + + parse(callService(queryParam)) match { + case Right(json) => + json.asArray match { + case Some(array) => Option(array).toRight(Error.RateLookupFailed(ErrorCodes.cacheInitFailed, "Cache Init Failed. Aborting Service Start!")) + case None => Option(Error.RateLookupFailed(ErrorCodes.fxRateLookUpFailed, "Cache Init Failed. Aborting Service Start!")).toLeft(Vector[Json]()) + } + case Left(_) => Option(Error.RateLookupFailed(ErrorCodes.fxRateLookUpFailed, "Cache Init Failed. Aborting Service Start!")).toLeft(Vector[Json]()) + } + } + + + private def constructCurrencyCombinations : Set[String] = { + for { + x <- Currency.supportedCurrencies + y <- Currency.supportedCurrencies + } yield { + val pair = Currency.show.show(x) + Currency.show.show(y) + pair + } + } + + private def callService(queryParam: String) : String = { + val config = ConfigFactory.load() + val appConfig = config.getConfig("app") + val appConfigJson: Either[ParsingFailure, Json] = parse(appConfig.root().render(ConfigRenderOptions.concise())) + appConfigJson match { + case Right(value) => { + val response: Response = requests.get("http://localhost:8080/rates?" + queryParam, headers = Map("token" -> value.hcursor.downField("oneframe").downField("token").as[String].toString)) + response.text() + } + case Left(_) => "" + } + + + } + +}