perf(playlists): bound matching and persistence work

This commit is contained in:
joshpatra committed 2026-07-27 08:11:28 -04:00
1 parent e8a2eb8a51
commit a39fff4d78
5 files changed
+517 -175

No files matched your search

@@ -1,10 +1,52 @@
using System.Diagnostics;
using allstarr.Core.Matching;
using allstarr.Core.Playlists;
using Xunit.Abstractions;
namespace allstarr.Tests;
public sealed class PlaylistMaterializationPlannerTests
public sealed class PlaylistMaterializationPlannerTests(ITestOutputHelper output)
{
[Fact]
public void Reconciliation_baseline_is_linear_at_100_1000_and_10000_tracks()
{
var baselines = new List<(int Count, long Allocated, long ElapsedTicks)>();
foreach (var count in new[] { 100, 1_000, 10_000 })
{
var entries = Enumerable.Range(0, count)
.Select(index => Entry(index, $"source-{index}"))
.ToArray();
var source = Source(entries);
var decisions = entries
.Select((entry, index) => Accepted(entry, $"local-{index}"))
.ToArray();
var allocatedBefore = GC.GetAllocatedBytesForCurrentThread();
var timer = Stopwatch.StartNew();
var plan = Planner().Plan(
PlaylistPlanMode.Reconcile, source, decisions, Target(), Rules());
timer.Stop();
var allocated = GC.GetAllocatedBytesForCurrentThread() - allocatedBefore;
Assert.Equal(count, plan.Entries.Count);
Assert.Equal(count, plan.OrderedBackendItemIds.Count);
Assert.Equal(
Enumerable.Range(0, count).Select(index => $"local-{index}"),
plan.OrderedBackendItemIds);
Assert.Equal(
plan.IdempotencyKey,
Planner().Plan(
PlaylistPlanMode.Reconcile, source, decisions, Target(), Rules()).IdempotencyKey);
baselines.Add((count, allocated, timer.ElapsedTicks));
output.WriteLine(
$"reconciliation tracks={count} allocated_bytes={allocated} elapsed_ticks={timer.ElapsedTicks}");
}
Assert.All(baselines.Zip(baselines.Skip(1)), pair =>
Assert.True(pair.Second.Allocated < pair.First.Allocated * 30,
$"Allocation growth from {pair.First.Count} to {pair.Second.Count} tracks was quadratic."));
}
[Fact]
public void Accepted_local_matches_keep_first_source_order_and_deduplicate_backend_items()
{
@@ -1,3 +1,5 @@
using System.Data.Common;
using System.Diagnostics;
using System.Security.Cryptography;
using System.Text;
using allstarr.Core.Identity;
@@ -11,12 +13,14 @@ using allstarr.Core.Playlists.Targets;
using allstarr.Core.Protocols;
using allstarr.Core.Storage;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using Moq;
using Npgsql;
using Xunit.Abstractions;
namespace allstarr.Tests;
public sealed class PlaylistOrchestrationIntegrationTests : IAsyncLifetime
public sealed class PlaylistOrchestrationIntegrationTests(ITestOutputHelper output) : IAsyncLifetime
{
private PostgresTestDatabase _database = null!;
private DbFactory _factory = null!;
@@ -95,6 +99,98 @@ public sealed class PlaylistOrchestrationIntegrationTests : IAsyncLifetime
await db.SaveChangesAsync();
}
[Fact]
public async Task PostgreSql_playlist_baseline_is_chunk_bounded_at_100_1000_and_10000_tracks()
{
await SetLink(mode: PlaylistLinkMode.Virtual);
var commands = new CommandCounter();
var options = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(_database.ConnectionString)
.AddInterceptors(commands)
.Options;
var factory = new DbFactory(options);
var service = new PlaylistOrchestrationService(
factory,
_source,
new FakeTargetResolver(_target),
new PlaylistMaterializationPlanner(),
new TrackMatchDecisionEngine(),
new TrackMatchCommandService(
factory,
new TrackMatchDecisionEngine(),
new ProviderAccountResolver(factory, new ProviderPolicyOptions()),
new Clock(_now)),
new Clock(_now));
var baselines = new List<(int Count, int Commands, long Allocated, long ElapsedTicks)>();
foreach (var count in new[] { 100, 1_000, 10_000 })
{
_source.Snapshot = Snapshot(
$"scale-{count}",
Enumerable.Range(0, count)
.Select(index => Entry(
index,
$"scale-{count}-entry-{index}",
$"scale-{count}-source-{index}",
"One"))
.ToArray());
commands.Reset();
var allocatedBefore = GC.GetTotalAllocatedBytes();
var timer = Stopwatch.StartNew();
var refresh = await service.RefreshAsync(Context(), _link);
timer.Stop();
var elapsedTicks = timer.ElapsedTicks;
var allocated = GC.GetTotalAllocatedBytes() - allocatedBefore;
await using (var db = await factory.CreateDbContextAsync())
{
Assert.Equal(
count,
await db.PlaylistSourceEntries.CountAsync(
item => item.PlaylistSourceSnapshotId == refresh.SnapshotId));
}
var measuredCommands = commands.Count - 1;
Assert.True(
measuredCommands <= 15 +
(int)Math.Ceiling(count / 20d) * 2 +
(int)Math.Ceiling(count / 500d),
$"{measuredCommands} SQL commands exceeded the chunk budget for {count} tracks.");
commands.Reset();
Assert.Equal(
refresh.SnapshotId,
(await service.RefreshAsync(Context(), _link)).SnapshotId);
Assert.InRange(commands.Count, 1, 15);
commands.Reset();
allocatedBefore = GC.GetTotalAllocatedBytes();
timer.Restart();
var run = await service.RunAsync(
Context(), new PlaylistOrchestrationRequest(_link, 1, refresh.SnapshotId));
timer.Stop();
elapsedTicks += timer.ElapsedTicks;
allocated += GC.GetTotalAllocatedBytes() - allocatedBefore;
measuredCommands += commands.Count;
Assert.Equal(count, run.Plan.Entries.Count);
Assert.Equal(PlaylistPreviewEntryStatus.Included, run.Plan.Entries[0].Status);
Assert.All(
run.Plan.Entries.Skip(1),
item => Assert.Equal(PlaylistPreviewEntryStatus.Duplicate, item.Status));
Assert.Equal(["local-1"], run.Plan.OrderedBackendItemIds);
Assert.True(
commands.Count <= 15 + (int)Math.Ceiling(count / 20d),
$"{commands.Count} matching SQL commands exceeded the chunk budget for {count} tracks.");
baselines.Add((count, measuredCommands, allocated, elapsedTicks));
output.WriteLine(
$"postgres-playlist tracks={count} commands={measuredCommands} returned_rows={run.Plan.Entries.Count} accepted_routes={run.Plan.OrderedBackendItemIds.Count} allocated_bytes={allocated} elapsed_ticks={elapsedTicks}");
}
Assert.All(baselines.Zip(baselines.Skip(1)), pair =>
Assert.True(pair.Second.Commands < pair.First.Commands * 30,
$"SQL command growth from {pair.First.Count} to {pair.Second.Count} tracks was quadratic."));
}
[Fact]
public async Task Refresh_persists_duplicate_source_positions_with_one_external_snapshot()
{
@@ -902,6 +998,34 @@ public sealed class PlaylistOrchestrationIntegrationTests : IAsyncLifetime
public AllstarrDbContext CreateDbContext() => new(options);
public Task<AllstarrDbContext> CreateDbContextAsync(CancellationToken cancellationToken = default) => Task.FromResult(CreateDbContext());
}
private sealed class CommandCounter : DbCommandInterceptor
{
private int _count;
public int Count => Volatile.Read(ref _count);
public void Reset() => Interlocked.Exchange(ref _count, 0);
private void Increment() => Interlocked.Increment(ref _count);
public override ValueTask<InterceptionResult<DbDataReader>> ReaderExecutingAsync(
DbCommand command, CommandEventData eventData, InterceptionResult<DbDataReader> result,
CancellationToken cancellationToken = default)
{
Increment();
return base.ReaderExecutingAsync(command, eventData, result, cancellationToken);
}
public override ValueTask<InterceptionResult<int>> NonQueryExecutingAsync(
DbCommand command, CommandEventData eventData, InterceptionResult<int> result,
CancellationToken cancellationToken = default)
{
Increment();
return base.NonQueryExecutingAsync(command, eventData, result, cancellationToken);
}
public override ValueTask<InterceptionResult<object>> ScalarExecutingAsync(
DbCommand command, CommandEventData eventData, InterceptionResult<object> result,
CancellationToken cancellationToken = default)
{
Increment();
return base.ScalarExecutingAsync(command, eventData, result, cancellationToken);
}
}
private sealed class EmptyServices : IServiceProvider
{
public static EmptyServices Instance { get; } = new();
@@ -1,12 +1,59 @@
using System.Text.Json;
using System.Diagnostics;
using allstarr.Core.Matching;
using allstarr.Core.Playlists;
using allstarr.Core.Storage;
using Xunit.Abstractions;
namespace allstarr.Tests;
public sealed class TrackMatchDecisionEngineTests
public sealed class TrackMatchDecisionEngineTests(ITestOutputHelper output)
{
[Fact]
public void Matching_baseline_is_linear_at_100_1000_and_10000_tracks()
{
var baselines = new List<(int Count, long Allocated, long ElapsedTicks)>();
foreach (var count in new[] { 100, 1_000, 10_000 })
{
var scope = Scope();
var candidates = Enumerable.Range(0, count)
.Select(index => Candidate(scope) with
{
LibraryTrackId = Guid.CreateVersion7(),
BackendItemId = $"local-{index}",
Title = $"Track {index}",
Isrc = $"USAAA26{index:D5}"
})
.ToArray();
var sources = candidates.Select((candidate, index) =>
Source() with
{
SnapshotId = index.ToString(),
Title = candidate.Title,
Isrc = candidate.Isrc
}).ToArray();
var engine = new TrackMatchDecisionEngine();
var allocatedBefore = GC.GetAllocatedBytesForCurrentThread();
var timer = Stopwatch.StartNew();
var prepared = engine.PrepareCandidates(candidates);
var decisions = sources.Select(source => engine.Decide(scope, source, prepared)).ToArray();
timer.Stop();
var allocated = GC.GetAllocatedBytesForCurrentThread() - allocatedBefore;
Assert.Equal(count, decisions.Length);
Assert.All(decisions, decision => Assert.Equal(TrackMatchReviewState.Accepted, decision.State));
Assert.All(decisions, decision => Assert.Single(decision.Candidates));
baselines.Add((count, allocated, timer.ElapsedTicks));
output.WriteLine(
$"matching tracks={count} allocated_bytes={allocated} elapsed_ticks={timer.ElapsedTicks}");
}
Assert.All(baselines.Zip(baselines.Skip(1)), pair =>
Assert.True(pair.Second.Allocated < pair.First.Allocated * 30,
$"Allocation growth from {pair.First.Count} to {pair.Second.Count} tracks was quadratic."));
}
[Fact]
public void PreparedCandidatesAndPersistenceInputPreserveOneDecision()
{
@@ -159,6 +159,11 @@ public interface ITrackMatchRepository
MatchDecisionInput input,
CancellationToken cancellationToken = default);
Task<IReadOnlyList<TrackMatchRecord>> RecordDecisionsAsync(
ProtocolExecutionContext context,
IReadOnlyCollection<MatchDecisionInput> inputs,
CancellationToken cancellationToken = default);
Task<ManualTrackOverrideRecord> SetOverrideAsync(
ProtocolExecutionContext context,
ManualOverrideInput input,
@@ -289,88 +294,111 @@ public sealed class TrackMatchCommandService(
public async Task<TrackMatchRecord> RecordDecisionAsync(
ProtocolExecutionContext context,
MatchDecisionInput input,
CancellationToken cancellationToken = default) =>
(await RecordDecisionsAsync(context, [input], cancellationToken)).Single();
public async Task<IReadOnlyList<TrackMatchRecord>> RecordDecisionsAsync(
ProtocolExecutionContext context,
IReadOnlyCollection<MatchDecisionInput> inputs,
CancellationToken cancellationToken = default)
{
var actor = context.RequireActor();
if (input.DecisionVersion <= 0 ||
input.SourceSnapshotVersion <= 0 ||
string.IsNullOrWhiteSpace(input.MatcherVersion) ||
input.Confidence is < 0 or > 1 ||
input.Threshold is < 0 or > 1 ||
string.IsNullOrWhiteSpace(input.PolicyVersion))
throw new ArgumentException("The match decision is incomplete.", nameof(input));
PersistenceGuard.ValidateSafeJson(input.CandidateResultsJson, nameof(input.CandidateResultsJson));
PersistenceGuard.ValidateSafeJson(input.ReasonsJson, nameof(input.ReasonsJson));
PersistenceGuard.ValidateSafeJson(input.WarningsJson, nameof(input.WarningsJson));
var requested = inputs.ToArray();
if (requested.Length == 0)
return [];
foreach (var input in requested)
ValidateDecisionInput(input);
if (requested.Select(item => (item.ExternalSnapshotId, item.DecisionVersion)).Distinct().Count() !=
requested.Length)
throw new ArgumentException("A match decision version may appear only once.", nameof(inputs));
await using var db = await contextFactory.CreateDbContextAsync(cancellationToken);
var snapshot = await OwnedSnapshotAsync(db, actor, input.ExternalSnapshotId, cancellationToken);
PersistenceGuard.RequireLibrary(context, snapshot.LibraryScopeId);
if (input.State is TrackMatchState.Accepted or TrackMatchState.Pinned &&
!input.LibraryTrackId.HasValue ||
input.State is TrackMatchState.Unresolved or TrackMatchState.Suggested or
TrackMatchState.Rejected or TrackMatchState.Ambiguous &&
input.LibraryTrackId.HasValue)
throw new ArgumentException(
"The selected library track does not match the decision state.",
nameof(input));
if (input.State == TrackMatchState.Accepted && input.Confidence < input.Threshold)
throw new ArgumentException(
"A match below its acceptance threshold cannot be accepted for automatic action.",
nameof(input));
if (input.LibraryTrackId.HasValue &&
!await db.LibraryTracks.AnyAsync(item =>
item.Id == input.LibraryTrackId &&
item.TenantId == actor.TenantId &&
item.OwnerUserId == snapshot.OwnerUserId &&
item.LibraryScopeId == snapshot.LibraryScopeId,
cancellationToken))
throw new UnauthorizedAccessException(
"The selected library track is outside the snapshot scope.");
var existing = await db.TrackMatches.AsNoTracking().SingleOrDefaultAsync(item =>
item.TenantId == actor.TenantId &&
item.OwnerUserId == snapshot.OwnerUserId &&
item.LibraryScopeId == snapshot.LibraryScopeId &&
item.ExternalSnapshotId == snapshot.Id &&
item.DecisionVersion == input.DecisionVersion,
cancellationToken);
if (existing != null)
var snapshotIds = requested.Select(item => item.ExternalSnapshotId).Distinct().ToArray();
var snapshots = await db.ExternalMetadataSnapshots
.Where(item => item.TenantId == actor.TenantId && snapshotIds.Contains(item.Id))
.ToDictionaryAsync(item => item.Id, cancellationToken);
if (snapshots.Count != snapshotIds.Length)
throw new UnauthorizedAccessException("A source snapshot is outside the actor scope.");
foreach (var snapshot in snapshots.Values)
{
if (!MatchesImmutableDecision(existing, input))
throw new InvalidOperationException(
"The match decision version already exists with different content.");
return existing;
PersistenceGuard.RequireOwner(actor, snapshot.OwnerUserId);
PersistenceGuard.RequireLibrary(context, snapshot.LibraryScopeId);
}
var libraryTrackIds = requested
.Where(item => item.LibraryTrackId.HasValue)
.Select(item => item.LibraryTrackId!.Value)
.Distinct()
.ToArray();
var libraryTracks = await db.LibraryTracks.AsNoTracking()
.Where(item => item.TenantId == actor.TenantId && libraryTrackIds.Contains(item.Id))
.ToDictionaryAsync(item => item.Id, cancellationToken);
foreach (var input in requested.Where(item => item.LibraryTrackId.HasValue))
{
var snapshot = snapshots[input.ExternalSnapshotId];
if (!libraryTracks.TryGetValue(input.LibraryTrackId!.Value, out var libraryTrack) ||
libraryTrack.OwnerUserId != snapshot.OwnerUserId ||
libraryTrack.LibraryScopeId != snapshot.LibraryScopeId)
throw new UnauthorizedAccessException(
"The selected library track is outside the snapshot scope.");
}
var versions = requested.Select(item => item.DecisionVersion).Distinct().ToArray();
var existing = await db.TrackMatches.AsNoTracking()
.Where(item => item.TenantId == actor.TenantId &&
snapshotIds.Contains(item.ExternalSnapshotId) &&
versions.Contains(item.DecisionVersion))
.ToDictionaryAsync(
item => (item.ExternalSnapshotId, item.DecisionVersion),
cancellationToken);
var now = clock.UtcNow;
var records = new List<TrackMatchRecord>(requested.Length);
foreach (var input in requested)
{
if (existing.TryGetValue((input.ExternalSnapshotId, input.DecisionVersion), out var stored))
{
if (!MatchesImmutableDecision(stored, input))
throw new InvalidOperationException(
"The match decision version already exists with different content.");
records.Add(stored);
continue;
}
var snapshot = snapshots[input.ExternalSnapshotId];
var record = ToRecord(
input, actor.TenantId, snapshot.OwnerUserId, snapshot.LibraryScopeId,
context.CorrelationId, now);
db.TrackMatches.Add(record);
records.Add(record);
}
var record = ToRecord(
input, actor.TenantId, snapshot.OwnerUserId, snapshot.LibraryScopeId,
context.CorrelationId, clock.UtcNow);
db.TrackMatches.Add(record);
try
{
await db.SaveChangesAsync(cancellationToken);
return record;
return records;
}
catch (DbUpdateException)
{
db.ChangeTracker.Clear();
var winner = await db.TrackMatches.AsNoTracking().SingleOrDefaultAsync(item =>
item.TenantId == actor.TenantId &&
item.OwnerUserId == snapshot.OwnerUserId &&
item.LibraryScopeId == snapshot.LibraryScopeId &&
item.ExternalSnapshotId == snapshot.Id &&
item.DecisionVersion == input.DecisionVersion,
cancellationToken);
if (winner is null)
var winners = await db.TrackMatches.AsNoTracking()
.Where(item => item.TenantId == actor.TenantId &&
snapshotIds.Contains(item.ExternalSnapshotId) &&
versions.Contains(item.DecisionVersion))
.ToDictionaryAsync(
item => (item.ExternalSnapshotId, item.DecisionVersion),
cancellationToken);
foreach (var input in requested)
{
throw;
if (!winners.TryGetValue(
(input.ExternalSnapshotId, input.DecisionVersion), out var winner))
throw;
if (!MatchesImmutableDecision(winner, input))
throw new InvalidOperationException(
"A concurrent match decision used the same version with different content.");
}
if (!MatchesImmutableDecision(winner, input))
throw new InvalidOperationException(
"A concurrent match decision used the same version with different content.");
return winner;
return requested
.Select(input => winners[(input.ExternalSnapshotId, input.DecisionVersion)])
.ToArray();
}
}
@@ -816,17 +844,12 @@ public sealed class TrackMatchCommandService(
item.RevokedAt == null &&
ownedSnapshotIds.Contains(item.ExternalSnapshotId))
.ToListAsync(cancellationToken);
var decisions = (await db.TrackMatches.AsNoTracking()
var decisions = await LatestDecisions(db.TrackMatches.AsNoTracking()
.Where(item => item.TenantId == actor.TenantId &&
item.OwnerUserId == ownerUserId &&
(libraryScopeId == null || item.LibraryScopeId == libraryScopeId) &&
ownedSnapshotIds.Contains(item.ExternalSnapshotId))
.OrderByDescending(item => item.DecisionVersion)
.ThenByDescending(item => item.DecidedAt)
.ToListAsync(cancellationToken))
.GroupBy(item => item.ExternalSnapshotId)
.Select(group => group.First())
.ToArray();
ownedSnapshotIds.Contains(item.ExternalSnapshotId)))
.ToArrayAsync(cancellationToken);
return new(snapshots, identities, overrides, decisions);
}
@@ -1076,13 +1099,33 @@ public sealed class TrackMatchCommandService(
ownerIds.Contains(item.OwnerUserId) &&
libraryScopes.Contains(item.LibraryScopeId))
.ToListAsync(cancellationToken);
var latestDecisions = (await db.TrackMatches
.Where(item => snapshotIds.Contains(item.ExternalSnapshotId))
.OrderByDescending(item => item.DecisionVersion)
.ThenByDescending(item => item.DecidedAt)
.ToListAsync(cancellationToken))
.GroupBy(item => item.ExternalSnapshotId)
.ToDictionary(group => group.Key, group => group.First());
var latestDecisions = await LatestDecisions(db.TrackMatches
.Where(item => snapshotIds.Contains(item.ExternalSnapshotId)))
.ToDictionaryAsync(item => item.ExternalSnapshotId, cancellationToken);
var scopedLibraries = libraryTracks
.GroupBy(item => new
{
item.TenantId,
item.OwnerUserId,
item.LibraryScopeId,
item.BackendInstanceId
})
.ToDictionary(
group => group.Key,
group =>
{
var tracks = group.ToArray();
return (
Tracks: tracks,
ById: tracks.ToDictionary(item => item.Id),
PlayableIds: tracks.Select(item => item.Id).ToHashSet(),
Candidates: decisionEngine.PrepareCandidates(tracks.Select(ToLocalCandidate)));
});
var emptyLibrary = (
Tracks: Array.Empty<LibraryTrackRecord>(),
ById: new Dictionary<Guid, LibraryTrackRecord>(),
PlayableIds: new HashSet<Guid>(),
Candidates: decisionEngine.PrepareCandidates([]));
var results = new List<AutomatedSourceMatchResult>(snapshots.Length);
var now = DateTimeOffset.UtcNow;
@@ -1093,22 +1136,23 @@ public sealed class TrackMatchCommandService(
var seed = tracks[SourceKey(identity.ProviderId, identity.ExternalId)];
latestDecisions.TryGetValue(snapshot.Id, out var latest);
activeOverrides.TryGetValue(snapshot.Id, out var manual);
var scopedTracks = libraryTracks
.Where(item =>
item.TenantId == snapshot.TenantId &&
item.OwnerUserId == snapshot.OwnerUserId &&
item.LibraryScopeId == snapshot.LibraryScopeId &&
item.BackendInstanceId == snapshot.BackendInstanceId)
.ToArray();
if (!scopedLibraries.TryGetValue(new
{
snapshot.TenantId,
snapshot.OwnerUserId,
snapshot.LibraryScopeId,
snapshot.BackendInstanceId
}, out var library))
library = emptyLibrary;
if (manual?.Decision == ManualOverrideDecision.Pin ||
manual?.Decision == ManualOverrideDecision.Reject && !manual.LibraryTrackId.HasValue)
{
var classification = TrackClassifier.Classify(
manual,
latest,
playableLibraryTrackIds: scopedTracks.Select(item => item.Id).ToHashSet());
playableLibraryTrackIds: library.PlayableIds);
var protectedLocal = classification.LibraryTrackId is { } protectedId
? scopedTracks.FirstOrDefault(item => item.Id == protectedId)
? library.ById.GetValueOrDefault(protectedId)
: null;
results.Add(ToAutomatedResult(
seed,
@@ -1118,7 +1162,7 @@ public sealed class TrackMatchCommandService(
continue;
}
var candidates = decisionEngine.PrepareCandidates(scopedTracks.Select(ToLocalCandidate));
var candidates = library.Candidates;
var libraryIndexRevision = candidates.Revision;
var scope = new TrackMatchScope(
snapshot.TenantId,
@@ -1156,7 +1200,7 @@ public sealed class TrackMatchCommandService(
var decision = decisionEngine.Decide(
scope, source, candidates, rejectedOverride);
var selected = decision.SelectedLibraryTrackId is { } selectedId
? scopedTracks.Single(item => item.Id == selectedId)
? library.ById[selectedId]
: null;
if (selected != null && !selected.CanonicalRecordingId.HasValue)
{
@@ -1197,6 +1241,23 @@ public sealed class TrackMatchCommandService(
private static string SourceKey(string providerId, string externalId) =>
$"{providerId.Trim().ToLowerInvariant()}:{externalId.Trim()}";
private static IQueryable<TrackMatchRecord> LatestDecisions(
IQueryable<TrackMatchRecord> decisions)
{
var versions = decisions
.GroupBy(item => item.ExternalSnapshotId)
.Select(group => new
{
ExternalSnapshotId = group.Key,
DecisionVersion = group.Max(item => item.DecisionVersion)
});
return from decision in decisions
join version in versions
on new { decision.ExternalSnapshotId, decision.DecisionVersion }
equals new { version.ExternalSnapshotId, version.DecisionVersion }
select decision;
}
private static LocalTrackMatchCandidate ToLocalCandidate(LibraryTrackRecord item) => new(
item.Id,
item.TenantId,
@@ -1616,6 +1677,32 @@ public sealed class TrackMatchCommandService(
private static string Hash(string value) =>
Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(value))).ToLowerInvariant();
private static void ValidateDecisionInput(MatchDecisionInput input)
{
if (input.DecisionVersion <= 0 ||
input.SourceSnapshotVersion <= 0 ||
string.IsNullOrWhiteSpace(input.MatcherVersion) ||
input.Confidence is < 0 or > 1 ||
input.Threshold is < 0 or > 1 ||
string.IsNullOrWhiteSpace(input.PolicyVersion))
throw new ArgumentException("The match decision is incomplete.", nameof(input));
PersistenceGuard.ValidateSafeJson(input.CandidateResultsJson, nameof(input.CandidateResultsJson));
PersistenceGuard.ValidateSafeJson(input.ReasonsJson, nameof(input.ReasonsJson));
PersistenceGuard.ValidateSafeJson(input.WarningsJson, nameof(input.WarningsJson));
if (input.State is TrackMatchState.Accepted or TrackMatchState.Pinned &&
!input.LibraryTrackId.HasValue ||
input.State is TrackMatchState.Unresolved or TrackMatchState.Suggested or
TrackMatchState.Rejected or TrackMatchState.Ambiguous &&
input.LibraryTrackId.HasValue)
throw new ArgumentException(
"The selected library track does not match the decision state.",
nameof(input));
if (input.State == TrackMatchState.Accepted && input.Confidence < input.Threshold)
throw new ArgumentException(
"A match below its acceptance threshold cannot be accepted for automatic action.",
nameof(input));
}
private static bool MatchesImmutableDecision(
TrackMatchRecord record,
MatchDecisionInput input) =>
@@ -564,24 +564,53 @@ public sealed class PlaylistOrchestrationService : IPlaylistOrchestrationService
providerIdentityIds.Contains(item.Id))
.ToDictionaryAsync(item => item.Id, cancellationToken);
var candidates = await db.LibraryTracks.AsNoTracking().Where(item =>
item.TenantId == link.TenantId && item.OwnerUserId == link.OwnerUserId &&
item.LibraryScopeId == link.LibraryScopeId && item.BackendInstanceId == link.TargetBackendInstanceId)
var candidates = await db.LibraryTracks.AsNoTracking()
.Where(item =>
item.TenantId == link.TenantId && item.OwnerUserId == link.OwnerUserId &&
item.LibraryScopeId == link.LibraryScopeId &&
item.BackendInstanceId == link.TargetBackendInstanceId)
.Select(item => new LibraryTrackRecord
{
Id = item.Id,
TenantId = item.TenantId,
OwnerUserId = item.OwnerUserId,
CanonicalRecordingId = item.CanonicalRecordingId,
LibraryScopeId = item.LibraryScopeId,
BackendInstanceId = item.BackendInstanceId,
BackendItemId = item.BackendItemId,
Title = item.Title,
Artist = item.Artist,
Album = item.Album,
AlbumArtist = item.AlbumArtist,
DurationMilliseconds = item.DurationMilliseconds,
Isrc = item.Isrc,
MusicBrainzRecordingId = item.MusicBrainzRecordingId,
ProviderIdsJson = item.ProviderIdsJson
})
.ToListAsync(cancellationToken);
var candidateIds = candidates.Select(item => item.Id).ToHashSet();
var priorAccepted = (await db.TrackMatches.AsNoTracking()
.Where(item =>
item.TenantId == link.TenantId &&
item.OwnerUserId == link.OwnerUserId &&
item.LibraryScopeId == link.LibraryScopeId &&
item.CanonicalRecordingId.HasValue &&
item.LibraryTrackId.HasValue &&
candidateIds.Contains(item.LibraryTrackId.Value))
.OrderByDescending(item => item.DecisionVersion)
.ThenByDescending(item => item.DecidedAt)
.ToListAsync(cancellationToken))
var candidateDecisions = db.TrackMatches.AsNoTracking()
.Where(item =>
item.TenantId == link.TenantId &&
item.OwnerUserId == link.OwnerUserId &&
item.LibraryScopeId == link.LibraryScopeId &&
item.CanonicalRecordingId.HasValue &&
item.LibraryTrackId.HasValue &&
candidateIds.Contains(item.LibraryTrackId.Value));
var latestCandidateVersions = candidateDecisions
.GroupBy(item => item.ExternalSnapshotId)
.Select(group => group.First())
.Select(group => new
{
ExternalSnapshotId = group.Key,
DecisionVersion = group.Max(item => item.DecisionVersion)
});
var latestCandidateDecisions =
from decision in candidateDecisions
join version in latestCandidateVersions
on new { decision.ExternalSnapshotId, decision.DecisionVersion }
equals new { version.ExternalSnapshotId, version.DecisionVersion }
select decision;
var priorAccepted = (await latestCandidateDecisions.ToListAsync(cancellationToken))
.Where(item => item.State is TrackMatchState.Accepted or TrackMatchState.Pinned)
.OrderByDescending(item => item.DecidedAt)
.GroupBy(item => item.LibraryTrackId!.Value)
@@ -635,6 +664,84 @@ public sealed class PlaylistOrchestrationService : IPlaylistOrchestrationService
var allManualOverrides = resolution.ActiveOverrides
.ToDictionary(item => item.ExternalSnapshotId);
var pendingDecisions = new Dictionary<Guid, MatchDecisionInput>();
foreach (var entry in entries)
{
cancellationToken.ThrowIfCancellationRequested();
var external = externals[entry.ExternalMetadataSnapshotId];
if (pendingDecisions.ContainsKey(external.Id))
continue;
allManualOverrides.TryGetValue(external.Id, out var manual);
storedByExternalId.TryGetValue(external.Id, out var stored);
if (stored != null &&
stored.SourceSnapshotVersion == external.SnapshotVersion &&
stored.LibraryIndexRevision == libraryIndexRevision &&
stored.MatcherVersion == TrackMatchDecisionEngine.AlgorithmVersion &&
stored.PolicyVersion == link.PolicyVersion)
continue;
using var payload = JsonDocument.Parse(external.PayloadJson);
var root = payload.RootElement;
var artists = root.GetProperty("Artists").EnumerateArray().Select(item => item.GetString()).Where(item => item != null).ToArray();
var canonicalRecordingId = external.ProviderTrackIdentityId.HasValue &&
providerIdentities.TryGetValue(
external.ProviderTrackIdentityId.Value, out var providerIdentity)
? providerIdentity.CanonicalRecordingId
: root.TryGetProperty("CanonicalRecordingId", out var canonical) &&
canonical.ValueKind == JsonValueKind.String &&
canonical.TryGetGuid(out var parsedCanonical)
? parsedCanonical
: (Guid?)null;
var source = new ExternalTrackMatchSnapshot(external.Id.ToString("N"), link.SourceProviderId,
external.ExternalIdHash, root.TryGetProperty("Title", out var title) ? title.GetString() ?? "Unknown" : "Unknown",
artists.Length > 0 ? string.Join(", ", artists) : "Unknown",
root.TryGetProperty("Album", out var album) ? album.GetString() : null, null,
ReadDurationMilliseconds(root),
root.TryGetProperty("Isrc", out var isrc) ? isrc.GetString() : null, null,
root.TryGetProperty("IsExplicit", out var explicitValue) && explicitValue.ValueKind is JsonValueKind.True or JsonValueKind.False ? explicitValue.GetBoolean() : null,
canonicalRecordingId);
var rejectedOverride =
manual?.Decision == ManualOverrideDecision.Reject &&
manual.LibraryTrackId.HasValue &&
manual.MatcherVersion == TrackMatchDecisionEngine.AlgorithmVersion
? new ScopedTrackMatchOverride(
link.TenantId,
link.OwnerUserId,
link.LibraryScopeId,
source.ProviderId,
source.ExternalId,
null,
new HashSet<Guid> { manual.LibraryTrackId.Value })
: null;
var match = _matcher.Decide(
new TrackMatchScope(link.TenantId, link.OwnerUserId, link.TargetBackendInstanceId, link.LibraryScopeId, link.ProviderAccountId, 1, snapshot.SnapshotVersion),
source,
candidateSet,
rejectedOverride);
var matchedCanonicalRecordingId = match.SelectedLibraryTrackId.HasValue &&
candidateById.TryGetValue(match.SelectedLibraryTrackId.Value, out var matchedCandidate)
? canonicalRecordingId ?? matchedCandidate.CanonicalRecordingId
: null;
pendingDecisions[external.Id] = MatchDecisionInput.FromDecision(
external.Id,
matchedCanonicalRecordingId,
match,
(stored?.DecisionVersion ?? 0) + 1,
external.SnapshotVersion,
libraryIndexRevision,
link.PolicyVersion);
}
if (pendingDecisions.Count > 0)
{
var storedDecisions = await _trackMatches.RecordDecisionsAsync(
execution, pendingDecisions.Values, cancellationToken);
foreach (var stored in storedDecisions)
storedByExternalId[stored.ExternalSnapshotId] = stored;
}
var decisions = new List<PersistedPlaylistMatchDecision>(entries.Count);
var decisionIds = new Dictionary<Guid, Guid?>(entries.Count);
@@ -643,72 +750,7 @@ public sealed class PlaylistOrchestrationService : IPlaylistOrchestrationService
cancellationToken.ThrowIfCancellationRequested();
var external = externals[entry.ExternalMetadataSnapshotId];
allManualOverrides.TryGetValue(external.Id, out var manual);
storedByExternalId.TryGetValue(external.Id, out var stored);
if (stored == null ||
stored.SourceSnapshotVersion != external.SnapshotVersion ||
stored.LibraryIndexRevision != libraryIndexRevision ||
stored.MatcherVersion != TrackMatchDecisionEngine.AlgorithmVersion ||
stored.PolicyVersion != link.PolicyVersion)
{
using var payload = JsonDocument.Parse(external.PayloadJson);
var root = payload.RootElement;
var artists = root.GetProperty("Artists").EnumerateArray().Select(item => item.GetString()).Where(item => item != null).ToArray();
var canonicalRecordingId = external.ProviderTrackIdentityId.HasValue &&
providerIdentities.TryGetValue(
external.ProviderTrackIdentityId.Value, out var providerIdentity)
? providerIdentity.CanonicalRecordingId
: root.TryGetProperty("CanonicalRecordingId", out var canonical) &&
canonical.ValueKind == JsonValueKind.String &&
canonical.TryGetGuid(out var parsedCanonical)
? parsedCanonical
: (Guid?)null;
var source = new ExternalTrackMatchSnapshot(external.Id.ToString("N"), link.SourceProviderId,
external.ExternalIdHash, root.TryGetProperty("Title", out var title) ? title.GetString() ?? "Unknown" : "Unknown",
artists.Length > 0 ? string.Join(", ", artists) : "Unknown",
root.TryGetProperty("Album", out var album) ? album.GetString() : null, null,
ReadDurationMilliseconds(root),
root.TryGetProperty("Isrc", out var isrc) ? isrc.GetString() : null, null,
root.TryGetProperty("IsExplicit", out var explicitValue) && explicitValue.ValueKind is JsonValueKind.True or JsonValueKind.False ? explicitValue.GetBoolean() : null,
canonicalRecordingId);
var rejectedOverride =
manual?.Decision == ManualOverrideDecision.Reject &&
manual.LibraryTrackId.HasValue &&
manual.MatcherVersion == TrackMatchDecisionEngine.AlgorithmVersion
? new ScopedTrackMatchOverride(
link.TenantId,
link.OwnerUserId,
link.LibraryScopeId,
source.ProviderId,
source.ExternalId,
null,
new HashSet<Guid> { manual.LibraryTrackId.Value })
: null;
var match = _matcher.Decide(
new TrackMatchScope(link.TenantId, link.OwnerUserId, link.TargetBackendInstanceId, link.LibraryScopeId, link.ProviderAccountId, 1, snapshot.SnapshotVersion),
source,
candidateSet,
rejectedOverride);
var matchedCanonicalRecordingId = match.SelectedLibraryTrackId.HasValue &&
candidateById.TryGetValue(match.SelectedLibraryTrackId.Value, out var matchedCandidate)
? canonicalRecordingId ?? matchedCandidate.CanonicalRecordingId
: null;
stored = await _trackMatches.RecordDecisionAsync(
execution,
MatchDecisionInput.FromDecision(
external.Id,
matchedCanonicalRecordingId,
match,
(stored?.DecisionVersion ?? 0) + 1,
external.SnapshotVersion,
libraryIndexRevision,
link.PolicyVersion),
cancellationToken);
storedByExternalId[external.Id] = stored;
}
var stored = storedByExternalId[external.Id];
var classification = TrackClassifier.Classify(
manual,