| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160 |
- using Application.Abstractions.Data;
- using Application.Abstractions.YouTube;
- using Application.Abstractions.Messaging;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Hosting;
- using Microsoft.Extensions.Logging;
- using Microsoft.Extensions.Options;
- using SharedKernel;
- namespace Infrastructure.YouTube;
- /// <summary>
- /// YouTube 공개 Atom Feed를 주기 폴링하여 신규 비디오를 감지.
- /// - YouTube PubSub이 라이브 이벤트를 누락하는 문제 보완
- /// - Feed 다운로드 자체는 0 unit (공개 XML)
- /// - 신규 videoId 감지 시 기존 PubSubNotify.Command 재사용 (downstream 파이프라인 통일)
- /// - 첫 실행 시 기존 영상들은 bootstrap으로 "마지막 처리"에만 기록 (무더기 처리 방지)
- /// </summary>
- internal sealed class YouTubeFeedPollingService(
- IServiceScopeFactory scopeFactory,
- IYouTubeFeedPoller feedPoller,
- IYouTubePubSubService pubSubService,
- IOptions<AppSettings> settings,
- ILogger<YouTubeFeedPollingService> logger
- ) : BackgroundService
- {
- private static readonly TimeSpan InitialDelay = TimeSpan.FromSeconds(60);
- private static readonly TimeSpan InterChannelDelay = TimeSpan.FromMilliseconds(500);
- protected override async Task ExecuteAsync(CancellationToken stoppingToken)
- {
- await Task.Delay(InitialDelay, stoppingToken);
- var intervalMinutes = settings.Value.YouTube.FeedPollingIntervalMinutes;
- if (intervalMinutes <= 0)
- {
- intervalMinutes = 3;
- }
- var interval = TimeSpan.FromMinutes(intervalMinutes);
- logger.LogInformation("[FeedPoller] 서비스 시작 — 폴링 주기: {Interval}분", intervalMinutes);
- var bootstrapped = false;
- while (!stoppingToken.IsCancellationRequested)
- {
- try
- {
- await PollOnceAsync(bootstrapped, stoppingToken);
- bootstrapped = true;
- }
- catch (OperationCanceledException)
- {
- break;
- }
- catch (Exception ex)
- {
- logger.LogError(ex, "[FeedPoller] Loop 예외");
- }
- try
- {
- await Task.Delay(interval, stoppingToken);
- }
- catch (OperationCanceledException)
- {
- break;
- }
- }
- }
- private async Task PollOnceAsync(bool bootstrapped, CancellationToken ct)
- {
- using var scope = scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- var mediator = scope.ServiceProvider.GetRequiredService<IMediator>();
- var channels = await db.Channel.AsNoTracking()
- .Where(c => c.IsActive && c.YouTubeChannelID != null && c.YouTubeChannelID != "")
- .Select(c => c.YouTubeChannelID!)
- .ToListAsync(ct);
- if (channels.Count == 0)
- {
- logger.LogDebug("[FeedPoller] 활성 채널 없음");
- return;
- }
- logger.LogDebug("[FeedPoller] Loop start — {Count} channels, bootstrapped={Bootstrapped}",
- channels.Count, bootstrapped);
- var detected = 0;
- var skipped = 0;
- foreach (var ytChannelId in channels)
- {
- if (ct.IsCancellationRequested)
- {
- break;
- }
- try
- {
- // 신규 감지 여부와 무관하게 "이 채널을 지금 폴링했음"을 기록
- await pubSubService.SetLastPolledAtAsync(ytChannelId, DateTime.UtcNow, ct);
- var latest = await feedPoller.GetLatestVideoIdAsync(ytChannelId, ct);
- if (string.IsNullOrEmpty(latest))
- {
- continue;
- }
- var lastProcessed = await pubSubService.GetLastProcessedVideoIdAsync(ytChannelId, ct);
- if (lastProcessed == latest)
- {
- skipped++;
- continue;
- }
- if (!bootstrapped && lastProcessed is null)
- {
- // 첫 루프 + Redis에 기록 없음 = 배포 직후 상태
- // 기존 영상을 "신규"로 오인하지 않도록 기록만 하고 trigger skip
- await pubSubService.SetLastProcessedVideoIdAsync(ytChannelId, latest, ct);
- logger.LogInformation("[FeedPoller] Bootstrap — channelId={ChannelId}, videoId={VideoId} recorded without trigger",
- ytChannelId, latest);
- continue;
- }
- // 신규 videoId 감지 → PubSubNotify 파이프라인 재사용
- logger.LogInformation("[FeedPoller] 신규 비디오 감지 — channelId={ChannelId}, videoId={VideoId} (prev={Prev})",
- ytChannelId, latest, lastProcessed ?? "(none)");
- await mediator.Send(new Application.Features.Api.YouTube.PubSubNotify.Command(
- latest, ytChannelId, string.Empty, DateTime.UtcNow, DateTime.UtcNow
- ), ct);
- // Handler 내부에서도 SetLastProcessed 호출하지만 여기서도 명시적으로 업데이트
- // (Handler가 라이브 아니라고 판단해 일찍 return해도 이 videoId는 더 안 봄)
- await pubSubService.SetLastProcessedVideoIdAsync(ytChannelId, latest, ct);
- detected++;
- }
- catch (OperationCanceledException)
- {
- throw;
- }
- catch (Exception ex)
- {
- logger.LogError(ex, "[FeedPoller] 채널 처리 예외 — channelId={ChannelId}", ytChannelId);
- }
- await Task.Delay(InterChannelDelay, ct);
- }
- logger.LogInformation("[FeedPoller] Loop 완료 — total={Total}, detected={Detected}, skipped={Skipped}",
- channels.Count, detected, skipped);
- }
- }
|