YouTubePubSubService.cs 9.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285
  1. using System.Globalization;
  2. using System.Security.Cryptography;
  3. using System.Text;
  4. using System.Xml.Linq;
  5. using Application.Abstractions.Cache;
  6. using Application.Abstractions.YouTube;
  7. using Microsoft.Extensions.Logging;
  8. using Microsoft.Extensions.Options;
  9. using SharedKernel;
  10. using StackExchange.Redis;
  11. namespace Infrastructure.YouTube;
  12. /// <summary>
  13. /// YouTube PubSubHubbub(WebSub) 구독 관리
  14. /// - 구독 신청/해지: POST https://pubsubhubbub.appspot.com/subscribe
  15. /// - 콜백 검증: hub.challenge GET 요청 응답
  16. /// - Atom Feed 파싱: 새 영상/생방송 알림 수신
  17. /// - Redis Set으로 구독 중인 채널 추적 (자동 재구독용)
  18. /// </summary>
  19. internal sealed class YouTubePubSubService(
  20. IHttpClientFactory httpClientFactory,
  21. IConnectionMultiplexer redis,
  22. IOptions<AppSettings> settings,
  23. ILogger<YouTubePubSubService> logger
  24. ) : IYouTubePubSubService
  25. {
  26. private const string TopicBase = "https://www.youtube.com/feeds/videos.xml?channel_id=";
  27. private string HubUrl => string.IsNullOrWhiteSpace(settings.Value.YouTube.HubUrl)
  28. ? "https://pubsubhubbub.appspot.com/subscribe"
  29. : settings.Value.YouTube.HubUrl;
  30. public string HmacSecret => settings.Value.YouTube.HmacSecret;
  31. public string CallbackUrl => settings.Value.YouTube.CallbackUrl;
  32. public async Task<bool> SubscribeAsync(string channelId, CancellationToken ct)
  33. {
  34. var success = await SendHubRequestAsync("subscribe", channelId, ct);
  35. if (success)
  36. {
  37. await redis.GetDatabase().SetAddAsync(CacheKeys.YouTubePubSubChannels, channelId);
  38. }
  39. return success;
  40. }
  41. public async Task<bool> UnsubscribeAsync(string channelId, CancellationToken ct)
  42. {
  43. var success = await SendHubRequestAsync("unsubscribe", channelId, ct);
  44. // 실패해도 Set에서는 제거 (더 이상 재구독하지 않도록)
  45. await redis.GetDatabase().SetRemoveAsync(CacheKeys.YouTubePubSubChannels, channelId);
  46. return success;
  47. }
  48. public async Task<IReadOnlyList<string>> GetSubscribedChannelIdsAsync(CancellationToken ct)
  49. {
  50. var values = await redis.GetDatabase().SetMembersAsync(CacheKeys.YouTubePubSubChannels);
  51. return values.Select(x => x.ToString()).Where(x => !string.IsNullOrEmpty(x)).ToList();
  52. }
  53. public async Task<bool> TryMarkProcessedAsync(string videoId, CancellationToken ct)
  54. {
  55. if (string.IsNullOrWhiteSpace(videoId))
  56. {
  57. return false;
  58. }
  59. var key = $"youtube:pubsub:processed:{videoId}";
  60. return await redis.GetDatabase().StringSetAsync(key, "1", TimeSpan.FromHours(24), When.NotExists);
  61. }
  62. public async Task<string?> GetLastProcessedVideoIdAsync(string youtubeChannelID, CancellationToken ct)
  63. {
  64. if (string.IsNullOrWhiteSpace(youtubeChannelID))
  65. {
  66. return null;
  67. }
  68. var val = await redis.GetDatabase().HashGetAsync(LastProcessedHashKey, youtubeChannelID);
  69. return val.IsNullOrEmpty ? null : val.ToString();
  70. }
  71. public async Task SetLastProcessedVideoIdAsync(string youtubeChannelID, string videoId, CancellationToken ct)
  72. {
  73. if (string.IsNullOrWhiteSpace(youtubeChannelID) || string.IsNullOrWhiteSpace(videoId))
  74. {
  75. return;
  76. }
  77. await redis.GetDatabase().HashSetAsync(LastProcessedHashKey, youtubeChannelID, videoId);
  78. }
  79. public async Task SetLastPolledAtAsync(string youtubeChannelID, DateTime utcNow, CancellationToken ct)
  80. {
  81. if (string.IsNullOrWhiteSpace(youtubeChannelID))
  82. {
  83. return;
  84. }
  85. await redis.GetDatabase().HashSetAsync(
  86. CacheKeys.YouTubeFeedLastPolled,
  87. youtubeChannelID,
  88. utcNow.ToString("O", CultureInfo.InvariantCulture)
  89. );
  90. }
  91. public async Task<IReadOnlyDictionary<string, DateTime>> GetAllLastPolledAtAsync(CancellationToken ct)
  92. {
  93. var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubeFeedLastPolled);
  94. var result = new Dictionary<string, DateTime>(entries.Length);
  95. foreach (var entry in entries)
  96. {
  97. if (DateTime.TryParse(entry.Value.ToString(), CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dt))
  98. {
  99. result[entry.Name.ToString()] = dt;
  100. }
  101. }
  102. return result;
  103. }
  104. public async Task<IReadOnlyDictionary<string, string>> GetAllLastProcessedVideoIdsAsync(CancellationToken ct)
  105. {
  106. var entries = await redis.GetDatabase().HashGetAllAsync(LastProcessedHashKey);
  107. var result = new Dictionary<string, string>(entries.Length);
  108. foreach (var entry in entries)
  109. {
  110. var val = entry.Value.ToString();
  111. if (!string.IsNullOrEmpty(val))
  112. {
  113. result[entry.Name.ToString()] = val;
  114. }
  115. }
  116. return result;
  117. }
  118. public async Task SetLeaseExpiryAsync(string youtubeChannelID, DateTime expiresAtUtc, CancellationToken ct)
  119. {
  120. if (string.IsNullOrWhiteSpace(youtubeChannelID))
  121. {
  122. return;
  123. }
  124. await redis.GetDatabase().HashSetAsync(
  125. CacheKeys.YouTubePubSubLease,
  126. youtubeChannelID,
  127. expiresAtUtc.ToString("O", CultureInfo.InvariantCulture)
  128. );
  129. }
  130. public async Task<IReadOnlyDictionary<string, DateTime>> GetAllLeaseExpiriesAsync(CancellationToken ct)
  131. {
  132. var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubePubSubLease);
  133. var result = new Dictionary<string, DateTime>(entries.Length);
  134. foreach (var entry in entries)
  135. {
  136. if (DateTime.TryParse(entry.Value.ToString(), CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dt))
  137. {
  138. result[entry.Name.ToString()] = dt;
  139. }
  140. }
  141. return result;
  142. }
  143. private const string LastProcessedHashKey = "youtube:feed:lastprocessed";
  144. public YouTubePubSubNotification? ParseNotification(string atomXml)
  145. {
  146. try
  147. {
  148. var doc = XDocument.Parse(atomXml);
  149. XNamespace atom = "http://www.w3.org/2005/Atom";
  150. XNamespace yt = "http://www.youtube.com/xml/schemas/2015";
  151. var entry = doc.Descendants(atom + "entry").FirstOrDefault();
  152. if (entry is null)
  153. {
  154. return null;
  155. }
  156. var videoId = entry.Element(yt + "videoId")?.Value;
  157. var channelId = entry.Element(yt + "channelId")?.Value;
  158. var title = entry.Element(atom + "title")?.Value ?? "";
  159. var published = entry.Element(atom + "published")?.Value;
  160. var updated = entry.Element(atom + "updated")?.Value;
  161. if (string.IsNullOrEmpty(videoId) || string.IsNullOrEmpty(channelId))
  162. {
  163. logger.LogWarning("[PubSub] Atom feed missing videoId or channelId");
  164. return null;
  165. }
  166. return new YouTubePubSubNotification(
  167. videoId,
  168. channelId,
  169. title,
  170. DateTime.TryParse(published, out var pub) ? pub : DateTime.UtcNow,
  171. DateTime.TryParse(updated, out var upd) ? upd : DateTime.UtcNow
  172. );
  173. }
  174. catch (Exception ex)
  175. {
  176. logger.LogError(ex, "[PubSub] Failed to parse Atom feed");
  177. return null;
  178. }
  179. }
  180. public bool VerifySignature(string payload, string signature)
  181. {
  182. if (string.IsNullOrEmpty(signature))
  183. {
  184. return false;
  185. }
  186. // X-Hub-Signature 형식: "sha1=hex_digest" 또는 "sha256=hex_digest"
  187. var parts = signature.Split('=', 2);
  188. if (parts.Length != 2)
  189. {
  190. return false;
  191. }
  192. var algorithm = parts[0];
  193. var expectedHex = parts[1];
  194. byte[] hash;
  195. var keyBytes = Encoding.UTF8.GetBytes(HmacSecret);
  196. var payloadBytes = Encoding.UTF8.GetBytes(payload);
  197. if (algorithm is "sha1")
  198. {
  199. hash = HMACSHA1.HashData(keyBytes, payloadBytes);
  200. }
  201. else if (algorithm is "sha256")
  202. {
  203. hash = HMACSHA256.HashData(keyBytes, payloadBytes);
  204. }
  205. else
  206. {
  207. logger.LogWarning("[PubSub] Unsupported HMAC algorithm: {Algorithm}", algorithm);
  208. return false;
  209. }
  210. var computedHex = Convert.ToHexString(hash).ToLowerInvariant();
  211. return string.Equals(computedHex, expectedHex, StringComparison.OrdinalIgnoreCase);
  212. }
  213. // ── Private ──────────────────────────────────────────────────────
  214. private async Task<bool> SendHubRequestAsync(string mode, string channelId, CancellationToken ct)
  215. {
  216. try
  217. {
  218. var client = httpClientFactory.CreateClient("PubSubHub");
  219. var topicUrl = $"{TopicBase}{channelId}";
  220. var content = new FormUrlEncodedContent(new Dictionary<string, string>
  221. {
  222. ["hub.mode"] = mode,
  223. ["hub.topic"] = topicUrl,
  224. ["hub.callback"] = CallbackUrl,
  225. ["hub.secret"] = HmacSecret,
  226. ["hub.verify"] = "async"
  227. });
  228. var response = await client.PostAsync(HubUrl, content, ct);
  229. if (response.IsSuccessStatusCode || (int)response.StatusCode == 202)
  230. {
  231. logger.LogInformation("[PubSub] {Mode} request accepted for channelId={ChannelId}", mode, channelId);
  232. return true;
  233. }
  234. var body = await response.Content.ReadAsStringAsync(ct);
  235. logger.LogWarning("[PubSub] {Mode} failed: {StatusCode} {Body}", mode, response.StatusCode, body);
  236. return false;
  237. }
  238. catch (Exception ex)
  239. {
  240. logger.LogError(ex, "[PubSub] {Mode} error for channelId={ChannelId}", mode, channelId);
  241. return false;
  242. }
  243. }
  244. }