using System.Globalization; using System.Security.Cryptography; using System.Text; using System.Xml.Linq; using Application.Abstractions.Cache; using Application.Abstractions.YouTube; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using SharedKernel; using StackExchange.Redis; namespace Infrastructure.YouTube; /// /// YouTube PubSubHubbub(WebSub) 구독 관리 /// - 구독 신청/해지: POST https://pubsubhubbub.appspot.com/subscribe /// - 콜백 검증: hub.challenge GET 요청 응답 /// - Atom Feed 파싱: 새 영상/생방송 알림 수신 /// - Redis Set으로 구독 중인 채널 추적 (자동 재구독용) /// internal sealed class YouTubePubSubService( IHttpClientFactory httpClientFactory, IConnectionMultiplexer redis, IOptions settings, ILogger logger ) : IYouTubePubSubService { private const string TopicBase = "https://www.youtube.com/feeds/videos.xml?channel_id="; private string HubUrl => string.IsNullOrWhiteSpace(settings.Value.YouTube.HubUrl) ? "https://pubsubhubbub.appspot.com/subscribe" : settings.Value.YouTube.HubUrl; public string HmacSecret => settings.Value.YouTube.HmacSecret; public string CallbackUrl => settings.Value.YouTube.CallbackUrl; public async Task SubscribeAsync(string channelId, CancellationToken ct) { var success = await SendHubRequestAsync("subscribe", channelId, ct); if (success) { await redis.GetDatabase().SetAddAsync(CacheKeys.YouTubePubSubChannels, channelId); } return success; } public async Task UnsubscribeAsync(string channelId, CancellationToken ct) { var success = await SendHubRequestAsync("unsubscribe", channelId, ct); // 실패해도 Set에서는 제거 (더 이상 재구독하지 않도록) await redis.GetDatabase().SetRemoveAsync(CacheKeys.YouTubePubSubChannels, channelId); return success; } public async Task> GetSubscribedChannelIdsAsync(CancellationToken ct) { var values = await redis.GetDatabase().SetMembersAsync(CacheKeys.YouTubePubSubChannels); return values.Select(x => x.ToString()).Where(x => !string.IsNullOrEmpty(x)).ToList(); } public async Task TryMarkProcessedAsync(string videoId, CancellationToken ct) { if (string.IsNullOrWhiteSpace(videoId)) { return false; } var key = $"youtube:pubsub:processed:{videoId}"; return await redis.GetDatabase().StringSetAsync(key, "1", TimeSpan.FromHours(24), When.NotExists); } public async Task GetLastProcessedVideoIdAsync(string youtubeChannelID, CancellationToken ct) { if (string.IsNullOrWhiteSpace(youtubeChannelID)) { return null; } var val = await redis.GetDatabase().HashGetAsync(LastProcessedHashKey, youtubeChannelID); return val.IsNullOrEmpty ? null : val.ToString(); } public async Task SetLastProcessedVideoIdAsync(string youtubeChannelID, string videoId, CancellationToken ct) { if (string.IsNullOrWhiteSpace(youtubeChannelID) || string.IsNullOrWhiteSpace(videoId)) { return; } await redis.GetDatabase().HashSetAsync(LastProcessedHashKey, youtubeChannelID, videoId); } public async Task SetLastPolledAtAsync(string youtubeChannelID, DateTime utcNow, CancellationToken ct) { if (string.IsNullOrWhiteSpace(youtubeChannelID)) { return; } await redis.GetDatabase().HashSetAsync( CacheKeys.YouTubeFeedLastPolled, youtubeChannelID, utcNow.ToString("O", CultureInfo.InvariantCulture) ); } public async Task> GetAllLastPolledAtAsync(CancellationToken ct) { var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubeFeedLastPolled); var result = new Dictionary(entries.Length); foreach (var entry in entries) { if (DateTime.TryParse(entry.Value.ToString(), CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dt)) { result[entry.Name.ToString()] = dt; } } return result; } public async Task> GetAllLastProcessedVideoIdsAsync(CancellationToken ct) { var entries = await redis.GetDatabase().HashGetAllAsync(LastProcessedHashKey); var result = new Dictionary(entries.Length); foreach (var entry in entries) { var val = entry.Value.ToString(); if (!string.IsNullOrEmpty(val)) { result[entry.Name.ToString()] = val; } } return result; } public async Task SetLeaseExpiryAsync(string youtubeChannelID, DateTime expiresAtUtc, CancellationToken ct) { if (string.IsNullOrWhiteSpace(youtubeChannelID)) { return; } await redis.GetDatabase().HashSetAsync( CacheKeys.YouTubePubSubLease, youtubeChannelID, expiresAtUtc.ToString("O", CultureInfo.InvariantCulture) ); } public async Task> GetAllLeaseExpiriesAsync(CancellationToken ct) { var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubePubSubLease); var result = new Dictionary(entries.Length); foreach (var entry in entries) { if (DateTime.TryParse(entry.Value.ToString(), CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var dt)) { result[entry.Name.ToString()] = dt; } } return result; } private const string LastProcessedHashKey = "youtube:feed:lastprocessed"; public YouTubePubSubNotification? ParseNotification(string atomXml) { try { var doc = XDocument.Parse(atomXml); XNamespace atom = "http://www.w3.org/2005/Atom"; XNamespace yt = "http://www.youtube.com/xml/schemas/2015"; var entry = doc.Descendants(atom + "entry").FirstOrDefault(); if (entry is null) { return null; } var videoId = entry.Element(yt + "videoId")?.Value; var channelId = entry.Element(yt + "channelId")?.Value; var title = entry.Element(atom + "title")?.Value ?? ""; var published = entry.Element(atom + "published")?.Value; var updated = entry.Element(atom + "updated")?.Value; if (string.IsNullOrEmpty(videoId) || string.IsNullOrEmpty(channelId)) { logger.LogWarning("[PubSub] Atom feed missing videoId or channelId"); return null; } return new YouTubePubSubNotification( videoId, channelId, title, DateTime.TryParse(published, out var pub) ? pub : DateTime.UtcNow, DateTime.TryParse(updated, out var upd) ? upd : DateTime.UtcNow ); } catch (Exception ex) { logger.LogError(ex, "[PubSub] Failed to parse Atom feed"); return null; } } public bool VerifySignature(string payload, string signature) { if (string.IsNullOrEmpty(signature)) { return false; } // X-Hub-Signature 형식: "sha1=hex_digest" 또는 "sha256=hex_digest" var parts = signature.Split('=', 2); if (parts.Length != 2) { return false; } var algorithm = parts[0]; var expectedHex = parts[1]; byte[] hash; var keyBytes = Encoding.UTF8.GetBytes(HmacSecret); var payloadBytes = Encoding.UTF8.GetBytes(payload); if (algorithm is "sha1") { hash = HMACSHA1.HashData(keyBytes, payloadBytes); } else if (algorithm is "sha256") { hash = HMACSHA256.HashData(keyBytes, payloadBytes); } else { logger.LogWarning("[PubSub] Unsupported HMAC algorithm: {Algorithm}", algorithm); return false; } var computedHex = Convert.ToHexString(hash).ToLowerInvariant(); return string.Equals(computedHex, expectedHex, StringComparison.OrdinalIgnoreCase); } // ── Private ────────────────────────────────────────────────────── private async Task SendHubRequestAsync(string mode, string channelId, CancellationToken ct) { try { var client = httpClientFactory.CreateClient("PubSubHub"); var topicUrl = $"{TopicBase}{channelId}"; var content = new FormUrlEncodedContent(new Dictionary { ["hub.mode"] = mode, ["hub.topic"] = topicUrl, ["hub.callback"] = CallbackUrl, ["hub.secret"] = HmacSecret, ["hub.verify"] = "async" }); var response = await client.PostAsync(HubUrl, content, ct); if (response.IsSuccessStatusCode || (int)response.StatusCode == 202) { logger.LogInformation("[PubSub] {Mode} request accepted for channelId={ChannelId}", mode, channelId); return true; } var body = await response.Content.ReadAsStringAsync(ct); logger.LogWarning("[PubSub] {Mode} failed: {StatusCode} {Body}", mode, response.StatusCode, body); return false; } catch (Exception ex) { logger.LogError(ex, "[PubSub] {Mode} error for channelId={ChannelId}", mode, channelId); return false; } } }