how to use actors in akka client side websockets

Viewed 276

i am using akka client websockets https://doc.akka.io/docs/akka-http/current/client-side/websocket-support.html

i have a server which takes requests and respond in json this is the pattern of my request and response

request #1

{
        "janus" : "create",
        "transaction" : "<random alphanumeric string>"
}
response #1 
{
   "janus": "success",
   "session_id": 2630959283560140,
   "transaction": "asqeasd4as3d4asdasddas",
   "data": {
        "id": 4574061985075210
   }
}

then based on response #1 i need to initiate request #2 and upon receiving response #2 i need to initiate request #3 and so on

for example

then based on id 4574061985075210 i will send request #2 and receive it response

request # 2 {
 }

response # 2 {
}
---- 

how can i use the actors with source and sink and re use the flow here is my initial code

import akka.http.scaladsl.model.ws._

import scala.concurrent.Future

object WebSocketClientFlow {
  def main(args: Array[String]) = {
    implicit val system = ActorSystem()
    implicit val materializer = ActorMaterializer()
    import system.dispatcher

    val incoming: Sink[Message, Future[Done]] =
      Sink.foreach[Message] {
        case message: TextMessage.Strict =>
          println(message.text)
//suppose here based on the server response i need to send another message to the server and so on do i need to repeat this same code here again ?????   

      }

    val outgoing = Source.single(TextMessage("hello world!"))

    val webSocketFlow = Http().webSocketClientFlow(WebSocketRequest("ws://echo.websocket.org"))

    val (upgradeResponse, closed) =
      outgoing
        .viaMat(webSocketFlow)(Keep.right) // keep the materialized Future[WebSocketUpgradeResponse]
        .toMat(incoming)(Keep.both) // also keep the Future[Done]
        .run()

    val connected = upgradeResponse.flatMap { upgrade =>
      if (upgrade.response.status == StatusCodes.SwitchingProtocols) {
        Future.successful(Done)
      } else {
        throw new RuntimeException(s"Connection failed: ${upgrade.response.status}")
      }
    }
connected.onComplete(println)
    closed.foreach(_ => println("closed"))
  }
}

and here i used Source.ActorRef

   val url = "ws://0.0.0.0:8188"
    val req = WebSocketRequest(url, Nil, Option("janus-protocol"))

    implicit val system = ActorSystem()
    implicit val materializer = ActorMaterializer()

    import system.dispatcher

    val webSocketFlow = Http().webSocketClientFlow(req)

    val messageSource: Source[Message, ActorRef] =
      Source.actorRef[TextMessage.Strict](bufferSize = 10, OverflowStrategy.fail)

    val messageSink: Sink[Message, NotUsed] =
      Flow[Message]
        .map(message => println(s"Received text message: [$message]"))
        .to(Sink.ignore)

    val ((ws, upgradeResponse), closed) =
      messageSource
        .viaMat(webSocketFlow)(Keep.both)
        .toMat(messageSink)(Keep.both)
        .run()

    val connected = upgradeResponse.flatMap { upgrade =>
      if (upgrade.response.status == StatusCodes.SwitchingProtocols) {
        Future.successful(Done)
      } else {
        throw new RuntimeException(s"Connection failed: ${upgrade.response.status}")
      }
    }

    val source =
      """{ "janus": "create", "transaction":"d1403sa54a5s3d4as3das"}"""
    val jsonAst = source.parseJson


    ws ! TextMessage.Strict(jsonAst.toString())

now i need help in how can i initiate the second request here because i need the "id" returned from the server to initiate request #2

0 Answers
Related