مرحبا هبر! هذه المرة حاولت إجراء محادثة بسيطة عبر Websockets. لمزيد من التفاصيل ، مرحبا بكم تحت القط.
المحتوى
- سكالا التعليمية: الجزء 1 - لعبة الأفعى
- سكالا التعلم: الجزء 2 - صحيفة تودو مع إمكانية تحميل الصور
- مقياس التعلم: الجزء 3 - اختبارات الوحدة
- سكالا التعلم: الجزء 4 - مقبس الويب
الروابط
في الواقع كل الشفرة موجودة في كائن 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>
هذا في الواقع كل شيء. آمل أن يهتم شخص يدرس موسيقى الروك بهذه المقالة وربما يكون مفيدًا أيضًا.