| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171 |
- using Application.Abstractions.Data;
- using Application.Abstractions.Hub;
- using Application.Abstractions.YouTube;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Hosting;
- using Microsoft.Extensions.Logging;
- using StackExchange.Redis;
- namespace Infrastructure.YouTube;
- /// <summary>
- /// 30초 주기로 라이브 중인 채널의 concurrentViewers 를 일괄 조회하여 Redis 에 갱신.
- ///
- /// Quota 예산:
- /// - videos.list (id=v1,v2,...50) = 1 unit / 50 videos
- /// - 30초 주기 + batch 50 = 분당 최대 2 call = 하루 최대 2,880 unit (기본 quota 10,000 대비 29%)
- /// - 라이브 채널 50개 미만이면 비용 동일 (1 batch 내 처리)
- ///
- /// 분산 락:
- /// - Redis SETNX 기반 락으로 멀티 인스턴스 중복 실행 방지 (quota 2배 소모 방지)
- /// - 락 TTL 은 PollInterval 보다 약간 짧게 설정
- ///
- /// SignalR 연동:
- /// - 시청자 수 갱신마다 ReceiveChannelStatus 브로드캐스트 → ChannelSidebar / WatchView 실시간 반영
- /// - concurrentViewers 가 응답에 없으면 라이브 종료로 판단 → ClearLiveAsync + IsLive=false 브로드캐스트
- /// </summary>
- internal sealed class YouTubeLiveViewerPoller(
- IYouTubeLiveStateStore liveStateStore,
- IYouTubeApiService youTubeApi,
- IChannelStatusBroadcaster statusBroadcaster,
- IServiceScopeFactory scopeFactory,
- IConnectionMultiplexer redis,
- ILogger<YouTubeLiveViewerPoller> logger
- ) : BackgroundService
- {
- private const string LockKey = "youtube:viewer-poller:lock";
- private static readonly TimeSpan InitialDelay = TimeSpan.FromSeconds(30);
- private static readonly TimeSpan PollInterval = TimeSpan.FromSeconds(30);
- private static readonly TimeSpan LockDuration = TimeSpan.FromSeconds(25);
- private readonly string _instanceId = Guid.NewGuid().ToString("N");
- protected override async Task ExecuteAsync(CancellationToken stoppingToken)
- {
- await Task.Delay(InitialDelay, stoppingToken);
- logger.LogInformation("[ViewerPoller] 서비스 시작 — 폴링 주기: {Interval}초, 인스턴스: {InstanceId}",
- PollInterval.TotalSeconds, _instanceId);
- while (!stoppingToken.IsCancellationRequested)
- {
- try
- {
- await PollOnceAsync(stoppingToken);
- }
- catch (OperationCanceledException)
- {
- break;
- }
- catch (Exception ex)
- {
- logger.LogError(ex, "[ViewerPoller] Loop 예외");
- }
- try
- {
- await Task.Delay(PollInterval, stoppingToken);
- }
- catch (OperationCanceledException)
- {
- break;
- }
- }
- }
- private async Task PollOnceAsync(CancellationToken ct)
- {
- var db = redis.GetDatabase();
- // 분산 락 획득 (멀티 인스턴스 환경에서 중복 실행 방지)
- var acquired = await db.StringSetAsync(LockKey, _instanceId, LockDuration, When.NotExists);
- if (!acquired)
- {
- logger.LogDebug("[ViewerPoller] 다른 인스턴스가 락 보유 중 — skip");
- return;
- }
- try
- {
- var allLive = await liveStateStore.GetAllLiveAsync();
- var liveList = allLive.Where(l => l.IsLive).ToList();
- if (liveList.Count == 0)
- {
- logger.LogDebug("[ViewerPoller] 라이브 채널 없음 — skip");
- return;
- }
- var videoIds = liveList.Select(l => l.VideoId).ToList();
- var viewers = await youTubeApi.GetConcurrentViewersAsync(videoIds, ct);
- // 🛡️ Quota 초과 방어: viewers 응답이 통째로 비어있으면 API 실패(quotaExceeded 등)로 판단,
- // ended 오인 처리 안 하고 이번 cycle skip. 다음 cycle 또는 다음 PubSub에서 재처리됨.
- // 단일 video만 missing이면 진짜 ended로 처리 (정상 분기).
- if (viewers.Count == 0)
- {
- logger.LogWarning("[ViewerPoller] viewers 응답 empty (live={LiveCount}) — API 실패 추정 (quota?). ended 오인 방지 위해 skip",
- liveList.Count);
- return;
- }
- // YouTube channelId → 내부 Channel.SID 매핑 (SignalR 브로드캐스트용)
- var ytChannelIds = liveList.Select(l => l.ChannelId).Distinct().ToList();
- using var scope = scopeFactory.CreateScope();
- var dbCtx = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- var sidMap = await dbCtx.Channel.AsNoTracking()
- .Where(c => c.YouTubeChannelID != null && ytChannelIds.Contains(c.YouTubeChannelID))
- .Select(c => new { c.YouTubeChannelID, c.SID })
- .ToDictionaryAsync(x => x.YouTubeChannelID!, x => x.SID, ct);
- var now = DateTime.UtcNow;
- var updated = 0;
- var ended = 0;
- foreach (var info in liveList)
- {
- if (ct.IsCancellationRequested)
- {
- break;
- }
- if (!viewers.TryGetValue(info.VideoId, out var cv))
- {
- // concurrentViewers 미제공 = 라이브 종료로 판단 (Redis 정리 + 브로드캐스트)
- // 만약 일시적 0명 상태였다면 다음 PubSub/FeedPoll 사이클에서 다시 마킹됨
- await liveStateStore.ClearLiveAsync(info.ChannelId);
- if (sidMap.TryGetValue(info.ChannelId, out var endedSid))
- {
- await statusBroadcaster.BroadcastAsync(endedSid, isLive: false, viewerCount: 0, videoId: null, ct);
- }
- ended++;
- continue;
- }
- await liveStateStore.UpdateViewerCountAsync(info.ChannelId, cv, now);
- if (sidMap.TryGetValue(info.ChannelId, out var sid))
- {
- await statusBroadcaster.BroadcastAsync(sid, isLive: true, viewerCount: cv, videoId: info.VideoId, ct);
- }
- updated++;
- }
- logger.LogInformation("[ViewerPoller] 완료 — live={Live}, updated={Updated}, ended={Ended}",
- liveList.Count, updated, ended);
- }
- finally
- {
- // 내 인스턴스가 잡은 락만 해제 (만료된 락을 다른 인스턴스가 잡았을 경우 보호)
- var current = await db.StringGetAsync(LockKey);
- if (current == _instanceId)
- {
- await db.KeyDeleteAsync(LockKey);
- }
- }
- }
- }
|