سكالا التعلم: الجزء 4 - مقبس الويب



مرحبا هبر! هذه المرة حاولت إجراء محادثة بسيطة عبر Websockets. لمزيد من التفاصيل ، مرحبا بكم تحت القط.



المحتوى





الروابط



  1. رموز المصدر
  2. صور عامل ميناء
  3. التابير
  4. Http4s
  5. Fs2
  6. دوبي
  7. ScalaTest
  8. ScalaCheck
  9. ScalaTestPlusScalaCheck


في الواقع كل الشفرة موجودة في كائن ChatHub واحد



class ChatHub[F[_]] private(
                             val topic: Topic[F, WebSocketFrame],
                             private val ref: Ref[F, Int]
                           )
                           (
                             implicit concurrent: Concurrent[F],
                             timer: Timer[F]
                           ) extends Http4sDsl[F] {

  val endpointWs: ServerEndpoint[String, Unit, String, Stream[IO, WebSocketFrame], IO] = endpoint
    .get
    .in("chat")
    .tag("WebSockets")
    .summary("   .    : ws://localhost:8080/chat")
    .description("   ")
    .in(
      stringBody
        .description("      ")
        .example("!")
    )
    .out(
      stringBody
        .description("  -   ")
        .example("6 :     Id  f518a53d: !")
    )
    //    . 
    .serverLogic(_ => IO(Left(()): Either[Unit, String]))

  def routeWs: HttpRoutes[F] = {
    HttpRoutes.of[F] {
      case GET -> Root / "chat" => logic()
    }
  }

  private def logic(): F[Response[F]] = {
    val toClient: Stream[F, WebSocketFrame] =
      topic.subscribe(1000)
    val fromClient: Pipe[F, WebSocketFrame, Unit] =
      handle
    WebSocketBuilder[F].build(toClient, fromClient)
  }

  private def handle(s: Stream[F, WebSocketFrame]): Stream[F, Unit] = s
    .collect({
      case WebSocketFrame.Text(text, _) => text
    })
    .evalMap(text => ref.modify(count => (count + 1, WebSocketFrame.Text(s"${count + 1} : $text"))))
    .through(topic.publish)
}

object ChatHub {

  def apply[F[_]]()(implicit concurrent: Concurrent[F], timer: Timer[F]): F[ChatHub[F]] = for {
    ref <- Ref.of[F, Int](0)
    topic <- Topic[F, WebSocketFrame](WebSocketFrame.Text("==="))
  } yield new ChatHub(topic, ref)
}


هنا يجب أن تقول على الفور عن الموضوع - وهو مبدأ مزامنة أولي من Fs2 يسمح لك بإنشاء نموذج ناشر - مشترك ، ويمكن أن يكون لديك العديد من الناشرين والعديد من المشتركين في نفس الوقت. بشكل عام ، من الأفضل إرسال الرسائل إليه من خلال نوع من المخزن المؤقت مثل قائمة الانتظار ، لأنه يحتوي على حد لعدد الرسائل في قائمة الانتظار وينتظر الناشر حتى يتلقى جميع المشتركين الرسائل في قائمة انتظار الرسائل الخاصة بهم وإذا تم تجاوزها فقد يتعطل.



val topic: Topic[F, WebSocketFrame],


هنا أيضًا أحسب عدد الرسائل التي تم إرسالها إلى الدردشة كرقم كل رسالة. نظرًا لأنني بحاجة إلى القيام بذلك من خيوط مختلفة ، فقد استخدمت نظيرًا من Atomic ، والذي يسمى هنا Ref ويضمن ذرية العملية.



  private val ref: Ref[F, Int]


معالجة دفق من الرسائل من المستخدمين.



  private def handle(stream: Stream[F, WebSocketFrame]): Stream[F, Unit] = 
    stream
//       . 
    .collect({
      case WebSocketFrame.Text(text, _) => text
    })
//               .
    .evalMap(text => ref.modify(count => (count + 1, WebSocketFrame.Text(s"${count + 1} : $text"))))
//     
    .through(topic.publish)


في الواقع ، منطق إنشاء المقبس.



private def logic(): F[Response[F]] = {
//    .
    val toClient: Stream[F, WebSocketFrame] =
//        
      topic.subscribe(1000)
//        
    val fromClient: Pipe[F, WebSocketFrame, Unit] =
//      
      handle
//         .
    WebSocketBuilder[F].build(toClient, fromClient)
  }


نربط المقبس الخاص بنا بالمسار الموجود على الخادم (ws: // localhost: 8080 / chat)



def routeWs: HttpRoutes[F] = {
    HttpRoutes.of[F] {
      case GET -> Root / "chat" => logic()
    }
  }


في الواقع ، هذا كل شيء. ثم يمكنك بدء الخادم بهذا المسار. ما زلت أرغب في عمل أي نوع من الوثائق. بشكل عام ، لتوثيق WebSocket والتفاعلات الأخرى القائمة على الأحداث مثل RabbitMQ AMPQ ، هناك AsynAPI ، لكن لا يوجد شيء ضمن Tapir ، لذلك قمت للتو بوضع وصف لنقطة النهاية لـ Swagger كطلب GET. بالطبع ، لن يعمل. بتعبير أدق ، سيتم إرجاع خطأ 501 ، ولكن سيتم عرضه في Swagger



  val endpointWs: Endpoint[String, Unit, String, fs2.Stream[F, Byte]] = endpoint
    .get
    .in("chat")
    .tag("WebSockets")
    .summary("   .    : ws://localhost:8080/chat")
    .description("   ")
    .in(
      stringBody
        .description("      ")
        .example("!")
    )
    .out(
      stringBody
        .description("  -   ")
        .example("6 :     Id  f518a53d: !")
    )


في اختيال نفسه ، يبدو مثل هذا. قم بتوصيل







الدردشة الخاصة بنا بخادم API الخاص بنا



    todosController = new TodosController()
    imagesController = new ImagesController()
//   
    chatHub <- Resource.liftF(ChatHub[IO]())
    endpoints = todosController.endpoints ::: imagesController.endpoints
//     Swagger
    docs = (chatHub.endpointWs :: endpoints).toOpenAPI("The Scala Todo List", "0.0.1")
    yml: String = docs.toYaml
//      
    routes = chatHub.routeWs <+>
      endpoints.toRoutes <+>
      new SwaggerHttp4s(yml, "swagger").routes[IO]
    httpApp = Router(
      "/" -> routes
    ).orNotFound
    blazeServer <- BlazeServerBuilder[IO](serverEc)
      .bindHttp(settings.host.port, settings.host.host)
      .withHttpApp(httpApp)
      .resource


نقوم بالاتصال بالدردشة بنص بسيط للغاية.



    <script>
        const id = `f${(~~(Math.random() * 1e8)).toString(16)}`;
        const webSocket = new WebSocket('ws://localhost:8080/chat');

        webSocket.onopen = event => {
            alert('onopen ');
        };

        webSocket.onmessage = event => {
            console.log(event);
            receive(event.data);
        };

        webSocket.onclose = event => {
            alert('onclose ');
        };

        function send() {
            let text = document.getElementById("message");
            webSocket.send(`    Id  ${id}: ${text.value}`);
            text.value = '';
        }

        function receive(m) {
            let text = document.getElementById("chat");
            text.value = text.value + '\n\r' + m;
        }
    </script>


هذا في الواقع كل شيء. آمل أن يهتم شخص يدرس موسيقى الروك بهذه المقالة وربما يكون مفيدًا أيضًا.



All Articles