using Application.Abstractions.Data;
using Application.Abstractions.YouTube;
using Domain.Entities.Members;
using Domain.Entities.Members.ValueObject;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
namespace Infrastructure.YouTube;
///
/// 매일 자정 00:00(KST) 1회 실행 — YouTube 방송 지표 집계.
/// - 각 활성 채널의 uploads 플레이리스트에서 최근 50개 videoId 조회
/// - IsFinalized=false 세션 또는 최근 7일 내 업로드 videoId만 대상
/// - videos.list 로 liveStreamingDetails + statistics 배치 조회
/// - BroadcastSession 생성/갱신, End+24h 경과 시 IsFinalized=true
/// - ChannelBroadcastStats 누적 재계산
///
/// API quota: 채널당 1(playlistItems) + 3(videos.list) = 4 units/일
/// 1,000채널까지 안전 (기본 10,000 quota).
///
internal sealed class YouTubeDailyAggregatorService(
IServiceScopeFactory scopeFactory,
IYouTubeApiService youTubeApi,
ILogger logger
) : BackgroundService
{
private const int BatchSize = 50; // YouTube API 최대 50개/요청
private const int UploadsLookback = 50; // 채널당 조회할 최근 업로드 수
private static readonly TimeSpan InitialDelay = TimeSpan.FromMinutes(2);
// 한국 표준시(KST) 자정 기준으로 실행
private static readonly TimeZoneInfo Kst = TimeZoneInfo.FindSystemTimeZoneById("Korea Standard Time");
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
await Task.Delay(InitialDelay, stoppingToken);
logger.LogInformation("[YouTubeDailyAggregator] 서비스 시작 — 매일 KST 00:00 실행");
while (!stoppingToken.IsCancellationRequested)
{
var delay = CalculateDelayUntilNextMidnight();
logger.LogInformation("[YouTubeDailyAggregator] 다음 실행까지 대기: {Hours}시간 {Minutes}분", delay.Hours, delay.Minutes);
try
{
await Task.Delay(delay, stoppingToken);
}
catch (TaskCanceledException)
{
break;
}
try
{
await AggregateAllAsync(stoppingToken);
}
catch (Exception ex)
{
logger.LogError(ex, "[YouTubeDailyAggregator] 집계 실행 중 오류");
}
}
}
private static TimeSpan CalculateDelayUntilNextMidnight()
{
var nowKst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, Kst);
var nextMidnightKst = nowKst.Date.AddDays(1);
var delay = nextMidnightKst - nowKst;
if (delay <= TimeSpan.Zero)
{
delay = TimeSpan.FromMinutes(1);
}
return delay;
}
private async Task AggregateAllAsync(CancellationToken ct)
{
logger.LogInformation("[YouTubeDailyAggregator] 집계 시작");
using var scope = scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
// YouTubeChannelID 가 설정된 활성 채널만 대상
var channels = await db.Channel.AsNoTracking()
.Where(c => c.IsActive && c.YouTubeChannelID != null && c.YouTubeChannelID != "")
.Select(c => new { c.ID, c.YouTubeChannelID })
.ToListAsync(ct);
if (channels.Count == 0)
{
logger.LogInformation("[YouTubeDailyAggregator] 대상 채널 없음");
return;
}
var processed = 0;
var failed = 0;
foreach (var channel in channels)
{
if (ct.IsCancellationRequested)
{
break;
}
try
{
await AggregateChannelAsync(channel.ID, channel.YouTubeChannelID!, ct);
processed++;
}
catch (Exception ex)
{
failed++;
logger.LogError(ex, "[YouTubeDailyAggregator] 채널 {ChannelID} 집계 실패", channel.ID);
}
// API rate limit 방어
await Task.Delay(TimeSpan.FromMilliseconds(500), ct);
}
logger.LogInformation("[YouTubeDailyAggregator] 집계 완료 — 성공: {Processed}, 실패: {Failed}", processed, failed);
}
private async Task AggregateChannelAsync(int channelID, string youtubeChannelID, CancellationToken ct)
{
// 1. 최근 업로드 videoId 조회 (1 unit)
var recentVideoIds = await youTubeApi.GetRecentUploadsAsync(youtubeChannelID, UploadsLookback, ct);
if (recentVideoIds.Count == 0)
{
return;
}
using var scope = scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
// 2. 기존 세션 중 IsFinalized=false 또는 최근 7일 내 세션 videoId 추가
var now = DateTime.UtcNow;
var lookbackStart = now.AddDays(-7);
var pendingSessions = await db.BroadcastSession
.Where(s => s.ChannelID == channelID && (!s.IsFinalized || (s.EndAt != null && s.EndAt > lookbackStart)))
.ToListAsync(ct);
var targetVideoIds = new HashSet(recentVideoIds);
foreach (var session in pendingSessions)
{
targetVideoIds.Add(session.VideoID);
}
// 3. videos.list 배치 조회 (50개씩, 3 units/배치)
var targetList = targetVideoIds.ToList();
var allStats = new List();
for (var i = 0; i < targetList.Count; i += BatchSize)
{
var batch = targetList.Skip(i).Take(BatchSize).ToList();
var stats = await youTubeApi.GetVideoStatsBatchAsync(batch, ct);
allStats.AddRange(stats);
if (i + BatchSize < targetList.Count)
{
await Task.Delay(TimeSpan.FromMilliseconds(250), ct);
}
}
// 4. 세션 upsert
var existingByVideoId = pendingSessions.ToDictionary(s => s.VideoID, s => s);
var sessionsToAdd = new List();
foreach (var stat in allStats)
{
// 라이브 방송이 아닌 일반 업로드는 제외
if (!stat.IsLiveStream || stat.ActualStartTime == null)
{
continue;
}
if (existingByVideoId.TryGetValue(stat.VideoID, out var existing))
{
existing.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
// End + 24h 경과 시 finalize
if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
{
existing.MarkFinalized();
}
}
else
{
var session = BroadcastSession.Create(channelID, BroadcastPlatform.YouTube, stat.VideoID);
session.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
{
session.MarkFinalized();
}
sessionsToAdd.Add(session);
}
}
if (sessionsToAdd.Count > 0)
{
await db.BroadcastSession.AddRangeAsync(sessionsToAdd, ct);
}
await db.SaveChangesAsync(ct);
// 5. ChannelBroadcastStats 누적 재계산 (모든 세션 스캔)
await RecalculateChannelStatsAsync(db, channelID, ct);
}
private static async Task RecalculateChannelStatsAsync(IAppDbContext db, int channelID, CancellationToken ct)
{
var aggregate = await db.BroadcastSession
.Where(s => s.ChannelID == channelID && s.StartAt != null && s.EndAt != null)
.GroupBy(s => 1)
.Select(g => new
{
TotalDurationSec = g.Sum(s => (long)s.DurationSec),
CumulativeViews = g.Sum(s => s.FinalViews),
TotalLikes = g.Sum(s => s.FinalLikes),
SessionCount = g.Count(),
LastBroadcastAt = g.Max(s => s.EndAt)
})
.FirstOrDefaultAsync(ct);
var stats = await db.ChannelBroadcastStats.FirstOrDefaultAsync(s => s.ChannelID == channelID, ct);
if (aggregate == null)
{
return;
}
if (stats == null)
{
stats = ChannelBroadcastStats.Create(channelID);
stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
await db.ChannelBroadcastStats.AddAsync(stats, ct);
}
else
{
stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
}
await db.SaveChangesAsync(ct);
}
}