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;
}
}
}