using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; namespace MeetingAssistant.Speakers; public interface ISpeakerIdentityMergeService { Task MergeRecentIdentitiesAsync( TimeSpan? recentIdentityAge, CancellationToken cancellationToken); } public sealed record SpeakerIdentityMergeResult( int RecentIdentityCount, int CandidateIdentityCount, int MatchAttempts, int MergedPairs); public sealed class SpeakerIdentityMergeService : ISpeakerIdentityMergeService { private readonly IDbContextFactory dbContextFactory; private readonly ISpeakerIdentityMatcher matcher; private readonly SpeakerIdentificationOptions options; private readonly ILogger logger; public SpeakerIdentityMergeService( IDbContextFactory dbContextFactory, ISpeakerIdentityMatcher matcher, IOptions options, ILogger logger) { this.dbContextFactory = dbContextFactory; this.matcher = matcher; this.options = options.Value.SpeakerIdentification; this.logger = logger; } public async Task MergeRecentIdentitiesAsync( TimeSpan? recentIdentityAge, CancellationToken cancellationToken) { await using var context = await dbContextFactory.CreateDbContextAsync(cancellationToken); await SpeakerIdentitySchema.EnsureCreatedOrUpdatedAsync(context, cancellationToken); var cutoff = DateTimeOffset.UtcNow - (recentIdentityAge ?? options.MergeRecentIdentityAge); var identities = await context.SpeakerIdentities .Include(identity => identity.Aliases) .Include(identity => identity.CandidateNames) .Include(identity => identity.Snippets) .Include(identity => identity.References) .OrderByDescending(identity => identity.References.Count) .ThenBy(identity => identity.Id) .ToListAsync(cancellationToken); var recentIds = identities .Where(identity => identity.CreatedAt >= cutoff) .Select(identity => identity.Id) .ToHashSet(); var mergedPairs = 0; var attempts = 0; logger.LogInformation( "Speaker identity merge diagnostics started: cutoff {Cutoff}, recent identity ids {RecentIdentityIds}, candidate identity count {CandidateIdentityCount}", cutoff, FormatIds(recentIds), identities.Count); foreach (var sourceId in recentIds.ToList()) { var source = identities.SingleOrDefault(identity => identity.Id == sourceId); if (source is null || source.Snippets.Count == 0) { logger.LogInformation( "Speaker identity merge diagnostics skipped source identity {SourceIdentityId} because it was missing or had no snippets", sourceId); continue; } var targets = identities .Where(identity => identity.Id != source.Id && identity.Snippets.Count > 0) .ToList(); logger.LogInformation( "Speaker identity merge diagnostics evaluating source identity {SourceIdentityId} ({SourceName}) with {SnippetCount} snippet(s) against target ids {TargetIdentityIds}", source.Id, source.GetDisplayName(), source.Snippets.Count, FormatIds(targets.Select(target => target.Id))); foreach (var batch in targets.Chunk(Math.Max(1, options.MatchBatchSize))) { var firstSnippet = SelectSnippet(source, excludedSnippet: null); if (firstSnippet is null) { logger.LogInformation( "Speaker identity merge diagnostics stopped evaluating source identity {SourceIdentityId} because no first snippet was available", source.Id); break; } attempts++; logger.LogInformation( "Speaker identity merge diagnostics round 1 attempt {Attempt} for source identity {SourceIdentityId}: target ids {TargetIdentityIds}, source snippet {SnippetId}", attempts, source.Id, FormatIds(batch.Select(target => target.Id)), firstSnippet.Id); var firstMatch = await matcher.MatchAsync( CreateRequest(source, firstSnippet, batch), cancellationToken); if (firstMatch is null || firstMatch.IdentityId == source.Id) { logger.LogInformation( "Speaker identity merge diagnostics round 1 found no usable target for source identity {SourceIdentityId}", source.Id); continue; } var target = batch.SingleOrDefault(identity => identity.Id == firstMatch.IdentityId); if (target is null) { logger.LogInformation( "Speaker identity merge diagnostics round 1 matched identity {IdentityId}, but it was not in the current target batch", firstMatch.IdentityId); continue; } var secondSnippet = SelectSnippet(source, excludedSnippet: firstSnippet); if (secondSnippet is null) { logger.LogInformation( "Speaker identity merge diagnostics could not validate source identity {SourceIdentityId} against target identity {TargetIdentityId} because no second source snippet was available", source.Id, target.Id); continue; } attempts++; logger.LogInformation( "Speaker identity merge diagnostics round 2 attempt {Attempt} for source identity {SourceIdentityId}: target identity {TargetIdentityId}, source snippet {SnippetId}", attempts, source.Id, target.Id, secondSnippet.Id); var secondMatch = await matcher.MatchAsync( CreateRequest(source, secondSnippet, batch), cancellationToken); if (secondMatch?.IdentityId != target.Id) { logger.LogInformation( "Speaker identity merge diagnostics rejected merge source identity {SourceIdentityId} into target identity {TargetIdentityId}: second match was {SecondMatchIdentityId}", source.Id, target.Id, secondMatch?.IdentityId); continue; } var targetName = target.GetDisplayName() ?? $"identity-{target.Id}"; var sourceName = source.GetDisplayName() ?? $"identity-{source.Id}"; SpeakerIdentityMerger.MergeInto( target, source, options.MaxSnippetsPerSpeaker); logger.LogInformation( "Speaker identity merge diagnostics merging source identity {SourceIdentityId} ({SourceName}) into target identity {TargetIdentityId} ({TargetName})", source.Id, sourceName, target.Id, targetName); await SpeakerIdentityTranscriptAudit.AppendMergedAsync( target.References, targetName, sourceName, cancellationToken); context.SpeakerIdentities.Remove(source); identities.Remove(source); mergedPairs++; break; } } await context.SaveChangesAsync(cancellationToken); logger.LogInformation( "Speaker identity merge diagnostics completed: {MergedPairs} merge(s), {Attempts} match attempt(s)", mergedPairs, attempts); return new SpeakerIdentityMergeResult( recentIds.Count, identities.Count, attempts, mergedPairs); } private static SpeakerIdentityMatchRequest CreateRequest( SpeakerIdentity source, SpeakerSnippet unknownSnippet, IReadOnlyList candidates) { return new SpeakerIdentityMatchRequest( source.GetDisplayName() ?? $"identity-{source.Id}", unknownSnippet.WavBytes, candidates.Select(identity => new SpeakerIdentityMatchCandidate( identity.Id, identity.GetDisplayName(), identity.ReferenceCount, identity.Snippets .OrderBy(snippet => snippet.CreatedAt) .Select(snippet => snippet.WavBytes) .ToList())) .ToList()); } private static SpeakerSnippet? SelectSnippet(SpeakerIdentity identity, SpeakerSnippet? excludedSnippet) { return identity.Snippets .OrderBy(snippet => snippet.CreatedAt) .FirstOrDefault(snippet => excludedSnippet is null || snippet.Id != excludedSnippet.Id); } private static string FormatIds(IEnumerable ids) { var formatted = ids .Distinct() .Order() .Select(id => id.ToString(System.Globalization.CultureInfo.InvariantCulture)) .ToList(); return formatted.Count == 0 ? "" : string.Join(", ", formatted); } }