Introduction

WebSocket is a protocol providing full-duplex communication channels over a single TCP connection that makes more interaction between a browser and a web server possible, facilitating the real-time data transfer from and to the server.
In this article, I will show you a sample that uses WebSocket to build a real-time application about distributing tasks.
Suppose we have a task center, there are some tasks will be published to the task center, and when the center receives some tasks, it should distribute the tasks to workers.
Let's take a look at the result at first.
Using WebSocket To Build Real-Time Application Via ASP.NET Core
There are four clients connected to the WebSocket Server, and after publishing messages via RabbitMQ management, the messages will be distributed to different clients in time.
Note
Clients means the workers, WebSocket Server means the task center and RabbitMQ management simulates sending task to the task center, and the messages means the tasks that should be handled.
And here is the architecture diagram that shows how it works.
Using WebSocket To Build Real-Time Application Via ASP.NET Core
Let's take a look on how to accomplish it.

Setup RabbitMQ

Before writing some code, we should run up the RabbitMQ server at first. The fastest way is to use Docker.
  1. docker run -p 5672:5672 -p 15672:15672 rabbitmq:management

WebSocket Server

This is the most important section!
At first, we should configure the WebSocket in Startup class.
  1. public class Startup
  2. {
  3. // other ...
  4. public void ConfigureServices(IServiceCollection services)
  5. {
  6. services.AddSingleton<Handlers.IDisHandler, Handlers.DisHandler>();
  7. services.AddControllers();
  8. }
  9. public void Configure(IApplicationBuilder app, IWebHostEnvironment env)
  10. {
  11. // other ..
  12. var webSocketOptions = new WebSocketOptions()
  13. {
  14. KeepAliveInterval = TimeSpan.FromSeconds(120),
  15. ReceiveBufferSize = 4 * 1024
  16. };
  17. app.UseWebSockets(webSocketOptions);
  18. app.Use(async (context, next) =>
  19. {
  20. // ws://www.yourdomian.com/push
  21. if (context.Request.Path == "/push")
  22. {
  23. if (context.WebSockets.IsWebSocketRequest)
  24. {
  25. WebSocket webSocket = await context.WebSockets.AcceptWebSocketAsync();
  26. try
  27. {
  28. var handler = app.ApplicationServices.GetRequiredService<Handlers.IDisHandler>();
  29. await handler.PushAsync(context, webSocket);
  30. }
  31. catch (Exception ex)
  32. {
  33. Console.WriteLine(ex.Message);
  34. }
  35. }
  36. else
  37. {
  38. context.Response.StatusCode = 400;
  39. }
  40. }
  41. else
  42. {
  43. await next();
  44. }
  45. });
  46. }
  47. }
Next step is to handle the logic about how to push messages to clients. Here create a class named DisHandler to do it.
Creating RabbitMQ connection, channel and consumer on the constructor, and the important part of the Received event of consumer.
Let's take a look at the Received event of consumer.
  1. consumer.Received += async (ch, ea) =>
  2. {
  3. var content = Encoding.UTF8.GetString(ea.Body);
  4. Console.WriteLine($"received content = {content}");
  5. // {"TaskName":"Demo", "TaskType":1}
  6. var msg = Newtonsoft.Json.JsonConvert.DeserializeObject<MqMsg>(content);
  7. var workIds = Worker.GetByTaskType(msg.TaskType);
  8. var onlineWorkerIds = _sockets.Keys.Intersect(workIds).ToList();
  9. if (onlineWorkerIds == null || !onlineWorkerIds.Any())
  10. {
  11. if (!ea.Redelivered)
  12. {
  13. Console.WriteLine("No online worker, reject the message and requeue");
  14. // should requeue here
  15. _channel.BasicReject(ea.DeliveryTag, true);
  16. }
  17. else
  18. {
  19. // should not requeue here, but this message will be discarded
  20. _channel.BasicReject(ea.DeliveryTag, false);
  21. }
  22. }
  23. else
  24. {
  25. // free or busy
  26. var randomNumberBuffer = new byte[10];
  27. new RNGCryptoServiceProvider().GetBytes(randomNumberBuffer);
  28. var rd = new Random(BitConverter.ToInt32(randomNumberBuffer, 0));
  29. var index = rd.Next(0, 9999) % onlineWorkerIds.Count;
  30. var workerId = onlineWorkerIds[index];
  31. if (_sockets.TryGetValue(workerId, out var ws) && ws.State == WebSocketState.Open)
  32. {
  33. // simulating handle the message an get the result.
  34. // put your own logic here
  35. var val = msg.TaskName;
  36. if (msg.TaskType != 1) val = $"Special-{msg.TaskName}";
  37. var task = Encoding.UTF8.GetBytes(val);
  38. Console.WriteLine($"send to {workerId}-{val}");
  39. // should ack here? or when to ack is better?
  40. _channel.BasicAck(ea.DeliveryTag, false);
  41. // sending message to specify client
  42. await ws.SendAsync(
  43. new ArraySegment<byte>(task, 0, task.Length),
  44. WebSocketMessageType.Text,
  45. true,
  46. CancellationToken.None);
  47. }
  48. else
  49. {
  50. Console.WriteLine("Not found a worker");
  51. }
  52. }
  53. };
When we receive a message, we should select a worker that can handle this task.
If there are no online workers, the server should reject the message and make it re-queue one more time.
If there are some online workers, the server will select a worker that can handle this task, here use a random number to simulate this scenario.
We also should ensure the connection of this worker is still open, so that the server can send messages to it.
The entry of handling WebSocket request is PushAsync method, it just maintains connections of clients, when client sends its client id to the server, it will record them, and when the client disconnects from the server, it will remove the client.
  1. public async Task PushAsync(HttpContext context, WebSocket webSocket)
  2. {
  3. var buffer = new byte[1024 * 4];
  4. WebSocketReceiveResult result =
  5. await webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None);
  6. string clientId = Encoding.UTF8.GetString(buffer, 0, result.Count);
  7. // record the client id and it's websocket instance
  8. if (_sockets.TryGetValue(clientId, out var wsi))
  9. {
  10. if (wsi.State == WebSocketState.Open)
  11. {
  12. Console.WriteLine($"abort the before clientId named {clientId}");
  13. await wsi.CloseAsync(WebSocketCloseStatus.InternalServerError,
  14. "A new client with same id was connected!",
  15. CancellationToken.None);
  16. }
  17. _sockets.AddOrUpdate(clientId, webSocket, (x, y) => webSocket);
  18. }
  19. else
  20. {
  21. Console.WriteLine($"add or update {clientId}");
  22. _sockets.AddOrUpdate(clientId, webSocket, (x, y) => webSocket);
  23. }
  24. while (!result.CloseStatus.HasValue)
  25. {
  26. result = await webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None);
  27. }
  28. await webSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None);
  29. Console.WriteLine("close=" + clientId);
  30. _sockets.TryRemove(clientId, out _);
  31. }

Client

The code of client just uses the sample in the document of ASP.NET Core WebSocket.
Here is the source code you can find in my GitHub page.

Summary

This article showed you a sample that uses WebSocket to build a real-time application via ASP.NET Core.
I hope this will help you!