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