diff --git a/Jobs/TwitchLiveJob.cs b/Jobs/TwitchLiveJob.cs index c010e1b..2c0e02a 100644 --- a/Jobs/TwitchLiveJob.cs +++ b/Jobs/TwitchLiveJob.cs @@ -17,18 +17,21 @@ public class TwitchLiveJob(DB db, TwitchService twitch, DiscordWebhookService di { public async Task Execute(IJobExecutionContext context) { + CancellationToken cancellationToken = context.CancellationToken; + cancellationToken.ThrowIfCancellationRequested(); + if (!twitch.IsConfigured) return; List subscriptions = await db.TwitchSubscriptions .Include(s => s.Webhook) - .ToListAsync(); + .ToListAsync(cancellationToken); if (subscriptions.Count == 0) return; List userIds = subscriptions.Select(s => s.TwitchUserId).Distinct().ToList(); - TwitchService.LiveStreamsResult result = await twitch.GetLiveStreamsResultAsync(userIds, context.CancellationToken); + TwitchService.LiveStreamsResult result = await twitch.GetLiveStreamsResultAsync(userIds, cancellationToken); if (!result.Succeeded) { logsService.Log("TwitchLiveJob: live-status request failed; preserving existing subscription state.", LogSeverity.Warning); @@ -41,19 +44,26 @@ public async Task Execute(IJobExecutionContext context) foreach (TwitchSubscription sub in subscriptions) { + cancellationToken.ThrowIfCancellationRequested(); live.TryGetValue(sub.TwitchUserId, out TwitchService.TwitchStream? stream); - changed |= await UpdateSubscriptionAsync(sub, stream, AnnounceAsync); + changed |= await UpdateSubscriptionAsync( + sub, + stream, + (subscription, liveStream) => AnnounceAsync(subscription, liveStream, cancellationToken), + cancellationToken); } if (changed) - await db.SaveChangesAsync(); + await db.SaveChangesAsync(cancellationToken); } internal static async Task UpdateSubscriptionAsync( TwitchSubscription sub, TwitchService.TwitchStream? stream, - Func> announceAsync) + Func> announceAsync, + CancellationToken cancellationToken = default) { + cancellationToken.ThrowIfCancellationRequested(); bool changed = false; if (stream != null) @@ -81,7 +91,10 @@ internal static async Task UpdateSubscriptionAsync( return changed; } - private async Task AnnounceAsync(TwitchSubscription sub, TwitchService.TwitchStream stream) + private async Task AnnounceAsync( + TwitchSubscription sub, + TwitchService.TwitchStream stream, + CancellationToken cancellationToken) { if (sub.Webhook == null) return false; @@ -89,7 +102,13 @@ private async Task AnnounceAsync(TwitchSubscription sub, TwitchService.Twi string title = string.IsNullOrWhiteSpace(stream.Title) ? string.Empty : $"\n{stream.Title}"; string content = $"🔴 **{sub.TwitchDisplayName}** is now live!{title}\nhttps://www.twitch.tv/{sub.TwitchLogin}"; - bool ok = await discordWebhook.SendAsync(sub.Webhook.WebhookId, sub.Webhook.Token, content, sub.TwitchDisplayName, sub.AvatarUrl); + bool ok = await discordWebhook.SendAsync( + sub.Webhook.WebhookId, + sub.Webhook.Token, + content, + sub.TwitchDisplayName, + sub.AvatarUrl, + cancellationToken); if (!ok) logsService.Log($"TwitchLiveJob: failed to post go-live for {sub.TwitchLogin} to channel {sub.ChannelDiscordId}", LogSeverity.Warning); diff --git a/Morpheus.Tests/TwitchLiveJobTests.cs b/Morpheus.Tests/TwitchLiveJobTests.cs index c4d3f75..2c745d5 100644 --- a/Morpheus.Tests/TwitchLiveJobTests.cs +++ b/Morpheus.Tests/TwitchLiveJobTests.cs @@ -1,11 +1,39 @@ +using Microsoft.Data.Sqlite; +using Microsoft.EntityFrameworkCore; +using Morpheus.Database; using Morpheus.Database.Models; using Morpheus.Jobs; using Morpheus.Services; +using Quartz; +using System.Reflection; namespace Morpheus.Tests; public class TwitchLiveJobTests { + [Fact] + public async Task ExecuteAsync_WhenCanceledBeforeLoadingSubscriptions_PropagatesCancellation() + { + await using SqliteConnection connection = new("Data Source=:memory:"); + await connection.OpenAsync(); + DbContextOptions options = new DbContextOptionsBuilder() + .UseSqlite(connection) + .Options; + await using DB db = new(options); + await db.Database.EnsureCreatedAsync(); + + LogsService logsService = new(new LogQueue()); + using HttpClient httpClient = new(); + TwitchService twitch = new(logsService, httpClient, "test-client", "test-secret"); + TwitchLiveJob job = new(db, twitch, new DiscordWebhookService(logsService), logsService); + using CancellationTokenSource cancellation = new(); + await cancellation.CancelAsync(); + + IJobExecutionContext context = CreateContext(cancellation.Token); + + await Assert.ThrowsAnyAsync(() => job.Execute(context)); + } + [Fact] public async Task UpdateSubscriptionAsync_WhenAnnouncementFails_DoesNotRecordStreamId() { @@ -49,4 +77,51 @@ Task AnnounceAsync(TwitchSubscription _, TwitchService.TwitchStream __) Assert.Equal(2, attempts); Assert.Equal("current-stream", subscription.LastAnnouncedStreamId); } + + [Fact] + public async Task UpdateSubscriptionAsync_WhenCanceled_DoesNotAnnounceOrMutate() + { + TwitchSubscription subscription = new(); + TwitchService.TwitchStream stream = new("current-stream", "Test stream"); + using CancellationTokenSource cancellation = new(); + await cancellation.CancelAsync(); + bool announced = false; + + await Assert.ThrowsAnyAsync(() => + TwitchLiveJob.UpdateSubscriptionAsync( + subscription, + stream, + (_, _) => + { + announced = true; + return Task.FromResult(true); + }, + cancellation.Token)); + + Assert.False(announced); + Assert.False(subscription.IsLive); + Assert.Null(subscription.LastAnnouncedStreamId); + } + + private static IJobExecutionContext CreateContext(CancellationToken cancellationToken) + { + JobExecutionContextProxy.CurrentCancellationToken = cancellationToken; + return DispatchProxy.Create(); + } + + private class JobExecutionContextProxy : DispatchProxy + { + public static CancellationToken CurrentCancellationToken { get; set; } + + protected override object? Invoke(MethodInfo? targetMethod, object?[]? args) + { + if (targetMethod?.ReturnType == typeof(CancellationToken)) + return CurrentCancellationToken; + + Type returnType = targetMethod?.ReturnType ?? typeof(void); + return returnType == typeof(void) || !returnType.IsValueType + ? null + : Activator.CreateInstance(returnType); + } + } }