Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 26 additions & 7 deletions Jobs/TwitchLiveJob.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<TwitchSubscription> subscriptions = await db.TwitchSubscriptions
.Include(s => s.Webhook)
.ToListAsync();
.ToListAsync(cancellationToken);

if (subscriptions.Count == 0)
return;

List<string> 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);
Expand All @@ -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<bool> UpdateSubscriptionAsync(
TwitchSubscription sub,
TwitchService.TwitchStream? stream,
Func<TwitchSubscription, TwitchService.TwitchStream, Task<bool>> announceAsync)
Func<TwitchSubscription, TwitchService.TwitchStream, Task<bool>> announceAsync,
CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
bool changed = false;

if (stream != null)
Expand Down Expand Up @@ -81,15 +91,24 @@ internal static async Task<bool> UpdateSubscriptionAsync(
return changed;
}

private async Task<bool> AnnounceAsync(TwitchSubscription sub, TwitchService.TwitchStream stream)
private async Task<bool> AnnounceAsync(
TwitchSubscription sub,
TwitchService.TwitchStream stream,
CancellationToken cancellationToken)
{
if (sub.Webhook == null)
return false;

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

Expand Down
75 changes: 75 additions & 0 deletions Morpheus.Tests/TwitchLiveJobTests.cs
Original file line number Diff line number Diff line change
@@ -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<DB> options = new DbContextOptionsBuilder<DB>()
.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<OperationCanceledException>(() => job.Execute(context));
}

[Fact]
public async Task UpdateSubscriptionAsync_WhenAnnouncementFails_DoesNotRecordStreamId()
{
Expand Down Expand Up @@ -49,4 +77,51 @@ Task<bool> 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<OperationCanceledException>(() =>
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<IJobExecutionContext, JobExecutionContextProxy>();
}

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