RedisChatConnectionTracker.cs 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103
  1. using Application.Abstractions.Cache;
  2. using Application.Abstractions.Chat;
  3. using StackExchange.Redis;
  4. using System.Text.Json;
  5. namespace Infrastructure.Chat;
  6. public class RedisChatConnectionTracker(IConnectionMultiplexer redis) : IChatConnectionTracker
  7. {
  8. private readonly IDatabase _db = redis.GetDatabase();
  9. private static readonly JsonSerializerOptions _jsonOptions = new()
  10. {
  11. PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
  12. WriteIndented = false
  13. };
  14. public async Task AddAsync(string channelSID, ConnectedUser user)
  15. {
  16. var json = JsonSerializer.Serialize(user, _jsonOptions);
  17. await _db.HashSetAsync(CacheKeys.ChatChannelConnections(channelSID), user.ConnectionId, json);
  18. }
  19. public async Task RemoveAsync(string channelSID, string connectionId)
  20. {
  21. await _db.HashDeleteAsync(CacheKeys.ChatChannelConnections(channelSID), connectionId);
  22. }
  23. public async Task<IReadOnlyList<ConnectedUser>> GetByChannelAsync(string channelSID)
  24. {
  25. var entries = await _db.HashGetAllAsync(CacheKeys.ChatChannelConnections(channelSID));
  26. if (entries.Length <= 0)
  27. {
  28. return [];
  29. }
  30. var list = new List<ConnectedUser>(entries.Length);
  31. foreach (var entry in entries)
  32. {
  33. if (!entry.Value.IsNullOrEmpty)
  34. {
  35. var user = JsonSerializer.Deserialize<ConnectedUser>((string)entry.Value!, _jsonOptions);
  36. if (user is not null)
  37. {
  38. list.Add(user);
  39. }
  40. }
  41. }
  42. return list.AsReadOnly();
  43. }
  44. public async Task<int> GetCountByChannelAsync(string channelSID)
  45. {
  46. return (int)await _db.HashLengthAsync(CacheKeys.ChatChannelConnections(channelSID));
  47. }
  48. public async Task<IReadOnlyList<(
  49. string ChannelSID,
  50. ConnectedUser User
  51. )>> GetAllAsync()
  52. {
  53. // 모든 채널 키를 스캔 후 병합 (관리자 집계용)
  54. var server = _db.Multiplexer.GetServer(_db.Multiplexer.GetEndPoints().First());
  55. var list = new List<(
  56. string ChannelSID,
  57. ConnectedUser User
  58. )>();
  59. await foreach (var key in server.KeysAsync(pattern: "chat:channel:*:connections"))
  60. {
  61. // 키에서 channelSID 추출 — chat:channel:{sid}:connections
  62. var keyStr = (string)key!;
  63. var channelSID = keyStr["chat:channel:".Length..^":connections".Length];
  64. var entries = await _db.HashGetAllAsync(key);
  65. foreach (var entry in entries)
  66. {
  67. if (!entry.Value.IsNullOrEmpty)
  68. {
  69. var user = JsonSerializer.Deserialize<ConnectedUser>((string)entry.Value!, _jsonOptions);
  70. if (user is not null)
  71. {
  72. list.Add((channelSID, user));
  73. }
  74. }
  75. }
  76. }
  77. return list.AsReadOnly();
  78. }
  79. public async Task ClearAllAsync()
  80. {
  81. // 서버 시작 시 모든 채널 연결 초기화 (stale 제거)
  82. var server = _db.Multiplexer.GetServer(_db.Multiplexer.GetEndPoints().First());
  83. await foreach (var key in server.KeysAsync(pattern: "chat:channel:*:connections"))
  84. {
  85. await _db.KeyDeleteAsync(key);
  86. }
  87. }
  88. }