fix(playlists): publish rematched decisions

This commit is contained in:
joshpatra committed 2026-08-05 22:06:45 -04:00
1 parent cb81b0fd3f
commit ff4875fecb
4 files changed
+189 -18

No files matched your search

@@ -581,6 +581,9 @@ public sealed class PlaylistOrchestrationIntegrationTests(ITestOutputHelper outp
Entry(1, "entry-protected", "source-alias", "One"));
var refresh = await _service.RefreshAsync(Context(), _link);
await _service.RunAsync(Context(), new(_link, 1, refresh.SnapshotId));
Guid rematchedExternalId;
Guid originalPublishedMatchId;
Guid rematchedDecisionId;
await using (var setup = await _factory.CreateDbContextAsync())
{
var snapshots = await setup.PlaylistSourceEntries
@@ -588,6 +591,10 @@ public sealed class PlaylistOrchestrationIntegrationTests(ITestOutputHelper outp
.OrderBy(item => item.SourcePosition)
.Select(item => item.ExternalMetadataSnapshotId)
.ToArrayAsync();
rematchedExternalId = snapshots[0];
originalPublishedMatchId = (await setup.PlaylistSourceEntries.SingleAsync(item =>
item.PlaylistSourceSnapshotId == refresh.SnapshotId && item.SourcePosition == 0))
.PublishedTrackMatchId!.Value;
setup.ManualTrackOverrides.Add(Override(snapshots[1], ManualOverrideDecision.Pin, _trackOne));
var decisions = await setup.TrackMatches.Where(item => snapshots.Contains(item.ExternalSnapshotId)).ToListAsync();
Assert.Equal(2, decisions.Count);
@@ -611,7 +618,9 @@ public sealed class PlaylistOrchestrationIntegrationTests(ITestOutputHelper outp
1,
PlaylistRematchJobHandler.Type,
JsonSerializer.SerializeToElement(new PlaylistRematchJobPayload(
preview.ConfirmationId, preview.ScopeFingerprint)),
preview.ConfirmationId,
preview.ScopeFingerprint,
preview.Targets)),
_tenant,
_user,
null,
@@ -621,24 +630,57 @@ public sealed class PlaylistOrchestrationIntegrationTests(ITestOutputHelper outp
"controlled-rematch",
"worker",
_now.AddMinutes(1));
var handler = new PlaylistRematchJobHandler(_factory, rematches, _trackMatches, new Clock(_now));
var handler = new PlaylistRematchJobHandler(
_factory, rematches, _trackMatches, _service, new Clock(_now));
var completion = await handler.ExecuteAsync(
new DurableJobExecutionContext(claim, EmptyServices.Instance), default);
Assert.Equal(DurableJobCompletionKind.Succeeded, completion.Kind);
var after = await rematches.PreviewAsync(_tenant, _user);
Assert.False(after.CanApply);
await using (var verify = await _factory.CreateDbContextAsync())
{
var latest = await verify.TrackMatches
.Where(item => item.ExternalSnapshotId == rematchedExternalId)
.OrderByDescending(item => item.DecisionVersion)
.FirstAsync();
rematchedDecisionId = latest.Id;
var entry = await verify.PlaylistSourceEntries.SingleAsync(item =>
item.PlaylistSourceSnapshotId == refresh.SnapshotId && item.SourcePosition == 0);
Assert.Equal(latest.Id, entry.PublishedTrackMatchId);
Assert.Single(await verify.ManualTrackOverrides.Where(item => item.RevokedAt == null).ToListAsync());
Assert.Single(await verify.AuditEvents.Where(item => item.Category == "playlist-rematch").ToListAsync());
foreach (var identity in await verify.ProviderTrackIdentities.ToListAsync())
identity.ProviderId = "detached";
await verify.SaveChangesAsync();
}
var after = await rematches.PreviewAsync(_tenant, _user);
Assert.False(after.CanApply);
await using (var interrupted = await _factory.CreateDbContextAsync())
{
var entry = await interrupted.PlaylistSourceEntries.SingleAsync(item =>
item.PlaylistSourceSnapshotId == refresh.SnapshotId && item.SourcePosition == 0);
entry.PublishedTrackMatchId = originalPublishedMatchId;
await interrupted.SaveChangesAsync();
}
var resumed = await handler.ExecuteAsync(
new DurableJobExecutionContext(claim with { AttemptNumber = 2 }, EmptyServices.Instance), default);
Assert.Equal(DurableJobCompletionKind.Succeeded, resumed.Kind);
await using var final = await _factory.CreateDbContextAsync();
Assert.Single(await final.AuditEvents.Where(item => item.Category == "playlist-rematch").ToListAsync());
Assert.Equal(
["rematched", "republished"],
await final.AuditEvents.Where(item => item.Category == "playlist-rematch")
.OrderBy(item => item.CreatedAt)
.Select(item => item.Outcome)
.ToArrayAsync());
var finalLatest = await final.TrackMatches
.Where(item => item.ExternalSnapshotId == rematchedExternalId)
.OrderByDescending(item => item.DecisionVersion)
.FirstAsync();
Assert.Equal(rematchedDecisionId, finalLatest.Id);
Assert.Equal(finalLatest.Id, (await final.PlaylistSourceEntries.SingleAsync(item =>
item.PlaylistSourceSnapshotId == refresh.SnapshotId && item.SourcePosition == 0))
.PublishedTrackMatchId);
}
[Fact]
@@ -741,7 +741,7 @@ public sealed class PlaylistLinksController(
var queued = await jobs.EnqueueAsync(new DurableJobEnqueueRequest<PlaylistRematchJobPayload>(
PlaylistRematchJobHandler.Type,
$"playlist-rematch:{session.AllstarrUserId:N}:{preview.ConfirmationId}",
new(preview.ConfirmationId, preview.ScopeFingerprint),
new(preview.ConfirmationId, preview.ScopeFingerprint, preview.Targets),
session.TenantId,
session.AllstarrUserId,
CorrelationId: HttpContext.TraceIdentifier), cancellationToken);
@@ -16,8 +16,10 @@ public sealed record PlaylistRematchTarget(
Guid ExternalSnapshotId,
int SnapshotVersion,
int DecisionVersion,
bool RequiresDecision,
Guid PlaylistLinkId,
long PlaylistLinkRevision,
Guid PlaylistSourceSnapshotId,
string PolicyVersion,
string LibraryScopeId,
string BackendInstanceId,
@@ -81,6 +83,22 @@ public sealed class PlaylistRematchService(
externalIds.Contains(item.ExternalSnapshotId) &&
item.RevokedAt == null)
.ToDictionaryAsync(item => item.ExternalSnapshotId, cancellationToken);
var snapshotLinks = byLink.ToDictionary(item => item.Value.SnapshotId, item => item.Key);
var publishedRows = await db.PlaylistSourceEntries.AsNoTracking()
.Where(item => snapshotLinks.Keys.Contains(item.PlaylistSourceSnapshotId))
.Select(item => new
{
item.PlaylistSourceSnapshotId,
item.ExternalMetadataSnapshotId,
item.PublishedTrackMatchId
})
.ToListAsync(cancellationToken);
var publishedByLink = publishedRows.GroupBy(item => (
LinkId: snapshotLinks[item.PlaylistSourceSnapshotId],
item.ExternalMetadataSnapshotId))
.ToDictionary(group => group.Key, group => group
.Select(item => item.PublishedTrackMatchId)
.ToHashSet());
var libraryTracks = await db.LibraryTracks.AsNoTracking()
.Where(item => item.TenantId == tenantId && item.OwnerUserId == ownerUserId)
.ToListAsync(cancellationToken);
@@ -119,17 +137,23 @@ public sealed class PlaylistRematchService(
decision.MatcherVersion != TrackMatchDecisionEngine.AlgorithmVersion ||
decision.PolicyVersion != row.Link.PolicyVersion);
if (stale) staleIds.Add(group.Key);
var missingProviderIdentity = group.All(item => item.Entry.ProviderRoutes.Count == 0);
var genericProviderIdentity = group.Any(item =>
item.Entry.RouteKind == "external" && !IsExactProvider(item.Entry.RouteProviderId));
var publicationOutdated = decision != null && group.Any(item =>
!publishedByLink.TryGetValue((item.Link.Id, group.Key), out var published) ||
published.Any(id => id != decision.Id));
if (overrides.ContainsKey(group.Key) || contextConflicts.Contains(group.Key) ||
decision != null && !stale && !missingProviderIdentity)
decision != null && !stale && !genericProviderIdentity && !publicationOutdated)
continue;
targetIds.Add(group.Key);
targets.Add(new(
group.Key,
snapshot.SnapshotVersion,
decision?.DecisionVersion ?? 0,
decision == null || stale || genericProviderIdentity,
row.Link.Id,
row.Link.Revision,
byLink[row.Link.Id].SnapshotId,
row.Link.PolicyVersion,
row.Link.LibraryScopeId,
row.Link.TargetBackendInstanceId,
@@ -143,7 +167,11 @@ public sealed class PlaylistRematchService(
.Concat(externalIds.Order().Select(id =>
$"row:{id:N}:{latest.GetValueOrDefault(id)?.DecisionVersion ?? 0}:" +
$"{latest.GetValueOrDefault(id)?.Revision ?? 0}:" +
$"{overrides.GetValueOrDefault(id)?.Revision ?? 0}:{targetIds.Contains(id)}"))));
$"{overrides.GetValueOrDefault(id)?.Revision ?? 0}:{targetIds.Contains(id)}:" +
string.Join(',', publishedByLink
.Where(item => item.Key.ExternalMetadataSnapshotId == id)
.OrderBy(item => item.Key.LinkId)
.SelectMany(item => item.Value.OrderBy(value => value))))) ));
var conflictingIds = contextConflicts
.Concat(latest.Where(item => item.Value.State == TrackMatchState.Ambiguous).Select(item => item.Key))
.ToHashSet();
@@ -210,12 +238,16 @@ public sealed class PlaylistRematchService(
DurablePlaylistEntryProjection Entry);
}
public sealed record PlaylistRematchJobPayload(string ConfirmationId, string ScopeFingerprint);
public sealed record PlaylistRematchJobPayload(
string ConfirmationId,
string ScopeFingerprint,
IReadOnlyList<PlaylistRematchTarget>? ApprovedTargets = null);
public sealed class PlaylistRematchJobHandler(
IDbContextFactory<AllstarrDbContext> contextFactory,
PlaylistRematchService rematches,
ITrackMatchRepository trackMatches,
PlaylistOrchestrationService orchestration,
IPlatformClock clock) : IDurableJobHandler
{
public const string Type = "playlist.rematch";
@@ -243,7 +275,7 @@ public sealed class PlaylistRematchJobHandler(
!preview.ConfirmationId.Equals(payload.ConfirmationId, StringComparison.Ordinal))
return DurableJobCompletion.Failure(
"playlist_rematch_preview_changed", "The match state changed. Review the rematch preview again.");
if (!preview.CanApply) return DurableJobCompletion.Success();
var approvedTargets = payload.ApprovedTargets ?? preview.Targets;
var runtime = await LoadRuntimeAsync(context, cancellationToken);
var started = Stopwatch.GetTimestamp();
@@ -264,7 +296,15 @@ public sealed class PlaylistRematchJobHandler(
continue;
}
var execution = CreateExecution(context, runtime, target, cancellationToken);
if (!target.RequiresDecision)
{
audits.Add(Audit(context, target.ExternalSnapshotId, "republished", target.DecisionVersion));
completed++;
continue;
}
var execution = CreateExecution(
context, runtime, target.TargetProtocol, target.BackendInstanceId,
target.LibraryScopeId, cancellationToken);
var result = await trackMatches.RematchSnapshotAsync(
execution,
target.ExternalSnapshotId,
@@ -288,6 +328,32 @@ public sealed class PlaylistRematchJobHandler(
ThroughputPerSecond: completed / Math.Max(Stopwatch.GetElapsedTime(started).TotalSeconds, .001)),
cancellationToken);
}
var current = await rematches.PreviewAsync(
context.Claim.TenantId.Value, context.Claim.OwnerUserId.Value, cancellationToken);
if (!current.ScopeFingerprint.Equals(payload.ScopeFingerprint, StringComparison.Ordinal))
return DurableJobCompletion.Failure(
"playlist_rematch_scope_changed", "A playlist or library changed. Review the rematch preview again.");
foreach (var publication in approvedTargets.GroupBy(item => new
{
item.PlaylistLinkId,
item.PlaylistLinkRevision,
item.PlaylistSourceSnapshotId,
item.LibraryScopeId,
item.BackendInstanceId,
item.TargetProtocol
}))
{
var execution = CreateExecution(
context, runtime, publication.Key.TargetProtocol, publication.Key.BackendInstanceId,
publication.Key.LibraryScopeId, cancellationToken);
await orchestration.PublishExistingDecisionsAsync(
execution,
publication.Key.PlaylistLinkId,
publication.Key.PlaylistLinkRevision,
publication.Key.PlaylistSourceSnapshotId,
publication.Select(item => item.ExternalSnapshotId).Distinct().ToArray(),
cancellationToken);
}
return failures == 0
? DurableJobCompletion.Success()
: DurableJobCompletion.Failure(
@@ -349,25 +415,27 @@ public sealed class PlaylistRematchJobHandler(
private ProtocolExecutionContext CreateExecution(
DurableJobExecutionContext context,
RuntimeIdentity runtime,
PlaylistRematchTarget target,
string targetProtocol,
string backendInstanceId,
string libraryScopeId,
CancellationToken cancellationToken)
{
var tenantId = context.Claim.TenantId!.Value;
var ownerUserId = context.Claim.OwnerUserId!.Value;
var protocol = target.TargetProtocol == "jellyfin" ? ProtocolKind.Jellyfin : ProtocolKind.Subsonic;
var protocol = targetProtocol == "jellyfin" ? ProtocolKind.Jellyfin : ProtocolKind.Subsonic;
var backendType = protocol.ToString().ToLowerInvariant();
if (!runtime.PrincipalIds.TryGetValue($"{backendType}\n{target.BackendInstanceId}", out var principalId))
if (!runtime.PrincipalIds.TryGetValue($"{backendType}\n{backendInstanceId}", out var principalId))
throw new UnauthorizedAccessException("The target backend identity is unavailable.");
return new(
protocol,
target.BackendInstanceId,
backendInstanceId,
principalId,
new AllstarrPrincipal(tenantId, ownerUserId, backendType, target.BackendInstanceId,
new AllstarrPrincipal(tenantId, ownerUserId, backendType, backendInstanceId,
principalId, runtime.DisplayName, false),
context.Claim.CorrelationId,
clock.UtcNow.AddMinutes(2),
cancellationToken,
libraryScopeId: target.LibraryScopeId);
libraryScopeId: libraryScopeId);
}
private AuditEventRecord Audit(
@@ -363,6 +363,67 @@ public sealed class PlaylistOrchestrationService : IPlaylistOrchestrationService
return new PlaylistRefreshResult(snapshot.Id, snapshot.SnapshotVersion, snapshot.ProviderRevision);
}
public async Task PublishExistingDecisionsAsync(
ProtocolExecutionContext execution,
Guid playlistLinkId,
long expectedLinkRevision,
Guid sourceSnapshotId,
IReadOnlyCollection<Guid> externalSnapshotIds,
CancellationToken cancellationToken = default)
{
var actor = execution.RequireActor();
await using var db = await _factory.CreateDbContextAsync(cancellationToken);
await using var transaction = await db.Database.BeginTransactionAsync(cancellationToken);
var link = await db.PlaylistLinks.SingleOrDefaultAsync(item =>
item.Id == playlistLinkId && item.TenantId == actor.TenantId,
cancellationToken) ?? throw new KeyNotFoundException("Playlist link not found.");
PersistenceGuard.RequireOwner(actor, link.OwnerUserId);
PersistenceGuard.RequireLibrary(execution, link.LibraryScopeId);
if (link.Revision != expectedLinkRevision)
throw new DbUpdateConcurrencyException("The playlist changed before rematch publication.");
var snapshot = await LoadSnapshotAsync(db, link, sourceSnapshotId, cancellationToken);
var currentId = await db.PlaylistSourceSnapshots.AsNoTracking()
.Where(item => item.TenantId == link.TenantId &&
item.PlaylistLinkId == link.Id &&
item.PublishedAt.HasValue)
.OrderByDescending(item => item.SnapshotVersion)
.ThenByDescending(item => item.RetrievedAt)
.Select(item => (Guid?)item.Id)
.FirstOrDefaultAsync(cancellationToken);
if (currentId != snapshot.Id)
throw new DbUpdateConcurrencyException("The published playlist changed before rematch publication.");
var externalIds = externalSnapshotIds.Distinct().ToArray();
var entries = await db.PlaylistSourceEntries
.Where(item => item.TenantId == link.TenantId &&
item.PlaylistSourceSnapshotId == snapshot.Id &&
externalIds.Contains(item.ExternalMetadataSnapshotId))
.ToListAsync(cancellationToken);
if (entries.Select(item => item.ExternalMetadataSnapshotId).Distinct().Count() != externalIds.Length)
throw new DbUpdateConcurrencyException("The rematched playlist rows changed before publication.");
var decisions = db.TrackMatches.AsNoTracking()
.Where(item => item.TenantId == link.TenantId &&
item.OwnerUserId == link.OwnerUserId &&
item.LibraryScopeId == link.LibraryScopeId &&
externalIds.Contains(item.ExternalSnapshotId));
var versions = decisions.GroupBy(item => item.ExternalSnapshotId).Select(group => new
{
ExternalSnapshotId = group.Key,
DecisionVersion = group.Max(item => item.DecisionVersion)
});
var latest = await (from decision in decisions
join version in versions
on new { decision.ExternalSnapshotId, decision.DecisionVersion }
equals new { version.ExternalSnapshotId, version.DecisionVersion }
select decision)
.ToDictionaryAsync(item => item.ExternalSnapshotId, cancellationToken);
if (latest.Count != externalIds.Length)
throw new InvalidOperationException("A rematched playlist cannot publish without one decision per source entry.");
foreach (var entry in entries)
entry.PublishedTrackMatchId = latest[entry.ExternalMetadataSnapshotId].Id;
await db.SaveChangesAsync(cancellationToken);
await transaction.CommitAsync(cancellationToken);
}
private async Task LogReconciliationAsync(
PlaylistLinkRecord link,
CancellationToken cancellationToken)