| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285 |
- 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;
- /// <summary>
- /// YouTube PubSubHubbub(WebSub) 구독 관리
- /// - 구독 신청/해지: POST https://pubsubhubbub.appspot.com/subscribe
- /// - 콜백 검증: hub.challenge GET 요청 응답
- /// - Atom Feed 파싱: 새 영상/생방송 알림 수신
- /// - Redis Set으로 구독 중인 채널 추적 (자동 재구독용)
- /// </summary>
- internal sealed class YouTubePubSubService(
- IHttpClientFactory httpClientFactory,
- IConnectionMultiplexer redis,
- IOptions<AppSettings> settings,
- ILogger<YouTubePubSubService> 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<bool> 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<bool> 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<IReadOnlyList<string>> 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<bool> 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<string?> 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<IReadOnlyDictionary<string, DateTime>> GetAllLastPolledAtAsync(CancellationToken ct)
- {
- var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubeFeedLastPolled);
- var result = new Dictionary<string, DateTime>(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<IReadOnlyDictionary<string, string>> GetAllLastProcessedVideoIdsAsync(CancellationToken ct)
- {
- var entries = await redis.GetDatabase().HashGetAllAsync(LastProcessedHashKey);
- var result = new Dictionary<string, string>(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<IReadOnlyDictionary<string, DateTime>> GetAllLeaseExpiriesAsync(CancellationToken ct)
- {
- var entries = await redis.GetDatabase().HashGetAllAsync(CacheKeys.YouTubePubSubLease);
- var result = new Dictionary<string, DateTime>(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<bool> SendHubRequestAsync(string mode, string channelId, CancellationToken ct)
- {
- try
- {
- var client = httpClientFactory.CreateClient("PubSubHub");
- var topicUrl = $"{TopicBase}{channelId}";
- var content = new FormUrlEncodedContent(new Dictionary<string, string>
- {
- ["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;
- }
- }
- }
|