using MeetingAssistant.MeetingNotes; using MeetingAssistant.Transcription; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; namespace MeetingAssistant.Speakers; public sealed class ResemblyzerSpeakerIdentificationService : ISpeakerIdentificationService { private readonly IDbContextFactory dbContextFactory; private readonly ISpeakerSnippetExtractor snippetExtractor; private readonly IResemblyzerVoiceEncoder encoder; private readonly ResemblyzerVoiceClusterMatcher clusterMatcher; private readonly ResemblyzerVoiceVectorOutlierPruner outlierPruner; private readonly SpeakerIdentificationOptions options; private readonly ResemblyzerSpeakerRecognitionOptions resemblyzerOptions; private readonly ILogger logger; public ResemblyzerSpeakerIdentificationService( IDbContextFactory dbContextFactory, ISpeakerSnippetExtractor snippetExtractor, IResemblyzerVoiceEncoder encoder, ResemblyzerVoiceClusterMatcher clusterMatcher, ResemblyzerVoiceVectorOutlierPruner outlierPruner, IOptions options, ILogger logger) { this.dbContextFactory = dbContextFactory; this.snippetExtractor = snippetExtractor; this.encoder = encoder; this.clusterMatcher = clusterMatcher; this.outlierPruner = outlierPruner; this.options = options.Value.SpeakerIdentification; resemblyzerOptions = this.options.Resemblyzer; this.logger = logger; } public Task IdentifyKnownSpeakersAsync( SpeakerIdentificationRequest request, CancellationToken cancellationToken) { return ProcessTranscriptAsync(request, final: false, allowAudioFallback: false, cancellationToken); } public Task IdentifyFinishedSpeakersAsync( SpeakerIdentificationRequest request, CancellationToken cancellationToken) { return ProcessTranscriptAsync(request, final: false, allowAudioFallback: true, cancellationToken); } public Task ProcessFinishedTranscriptAsync( SpeakerIdentificationRequest request, CancellationToken cancellationToken) { return ProcessTranscriptAsync(request, final: true, allowAudioFallback: true, cancellationToken); } public async Task ApplySpeakerOverrideAsync( SpeakerIdentificationRequest request, string sourceSpeaker, string targetSpeaker, CancellationToken cancellationToken) { if (!options.Enabled || string.IsNullOrWhiteSpace(sourceSpeaker) || string.IsNullOrWhiteSpace(targetSpeaker) || string.Equals(sourceSpeaker, targetSpeaker, StringComparison.OrdinalIgnoreCase)) { return; } var sourceLabel = sourceSpeaker.Trim(); var targetName = targetSpeaker.Trim(); await using var context = await dbContextFactory.CreateDbContextAsync(cancellationToken); await SpeakerIdentitySchema.EnsureCreatedOrUpdatedAsync(context, cancellationToken); var identities = await LoadIdentities(context).ToListAsync(cancellationToken); var target = identities .Where(identity => SpeakerIdentityNaming.GetAcceptedNames(identity).Contains(targetName)) .OrderBy(identity => string.Equals(identity.CanonicalName, targetName, StringComparison.OrdinalIgnoreCase) ? 0 : 1) .ThenBy(identity => identity.Id) .FirstOrDefault(); var reference = CreateReference(request.MeetingNote, DateTimeOffset.UtcNow); var sourceCandidate = identities .Where(identity => string.IsNullOrWhiteSpace(identity.CanonicalName)) .Where(identity => identity.References.Any(existing => SpeakerIdentityReferences.IsSame(existing, reference))) .Where(identity => SpeakerIdentityNaming.GetAcceptedNames(identity).Contains(targetName)) .OrderBy(identity => identity.Id) .FirstOrDefault(); var vectors = await ResolveAvailableVectorsAsync(request, sourceLabel, cancellationToken); if (target is null && sourceCandidate is null && vectors.Count == 0) { logger.LogWarning( "Skipping Resemblyzer speaker override from {SourceSpeaker} to {TargetSpeaker} because no source evidence was available", sourceLabel, targetName); return; } var now = DateTimeOffset.UtcNow; if (target is null) { target = sourceCandidate ?? new SpeakerIdentity { CreatedAt = now }; if (target.Id == 0) { context.SpeakerIdentities.Add(target); } } else if (sourceCandidate is not null && sourceCandidate.Id != target.Id) { SpeakerIdentityMerger.MergeIntoAndPrune( target, sourceCandidate, options.MaxSnippetsPerSpeaker, resemblyzerOptions.MaxVectorsPerIdentity, outlierPruner); context.SpeakerIdentities.Remove(sourceCandidate); } target.CanonicalName = targetName; SpeakerIdentityNaming.SetCandidates(target, [targetName]); SpeakerIdentityReferences.AddIfMissing(target, reference, now); outlierPruner.AddAndPrune( target, vectors, now); target.UpdatedAt = now; await context.SaveChangesAsync(cancellationToken); await SpeakerIdentityTranscriptAudit.AppendIdentifiedAsync( target.References, sourceLabel, targetName, cancellationToken); } public async Task DeleteSpeakerIdentityAsync( string identity, CancellationToken cancellationToken) { if (!options.Enabled || string.IsNullOrWhiteSpace(identity)) { return; } await using var context = await dbContextFactory.CreateDbContextAsync(cancellationToken); await SpeakerIdentitySchema.EnsureCreatedOrUpdatedAsync(context, cancellationToken); var identities = await LoadIdentities(context).ToListAsync(cancellationToken); var target = identities.FirstOrDefault(candidate => SpeakerIdentityNaming.GetAcceptedNames(candidate).Contains(identity.Trim())); if (target is null) { return; } context.SpeakerIdentities.Remove(target); await context.SaveChangesAsync(cancellationToken); } private async Task ProcessTranscriptAsync( SpeakerIdentificationRequest request, bool final, bool allowAudioFallback, CancellationToken cancellationToken) { if (!options.Enabled || request.Segments.Count == 0) { return EmptyResult(request.Segments); } await using var context = await dbContextFactory.CreateDbContextAsync(cancellationToken); await SpeakerIdentitySchema.EnsureCreatedOrUpdatedAsync(context, cancellationToken); var attendees = SpeakerIdentityNaming.NormalizeAttendees(request.MeetingNote.Frontmatter.Attendees); var knownMappings = request.KnownSpeakerMappings ?? new Dictionary(StringComparer.OrdinalIgnoreCase); var knownLabels = knownMappings.Keys.ToHashSet(StringComparer.OrdinalIgnoreCase); if (final && knownMappings.Count > 0) { await PersistMappedSpeakerEvidenceAsync( context, request, knownMappings, cancellationToken); } var identifiedNames = knownMappings.Values .Where(name => !string.IsNullOrWhiteSpace(name)) .Select(name => name.Trim()) .ToHashSet(StringComparer.OrdinalIgnoreCase); foreach (var segmentSpeaker in request.Segments .Select(segment => segment.Speaker) .Where(speaker => !string.IsNullOrWhiteSpace(speaker) && !SpeakerIdentityNaming.IsDiarizedSpeakerLabel(speaker))) { identifiedNames.Add(segmentSpeaker.Trim()); } var mappings = new Dictionary(StringComparer.OrdinalIgnoreCase); var attendeeMatches = new List(); var pendingAudits = new List(); var matchedAcceptedNames = identifiedNames.ToHashSet(StringComparer.OrdinalIgnoreCase); var unmatchedSpeakers = new List<(string Speaker, IReadOnlyList Vectors)>(); foreach (var speaker in request.Segments .Select(segment => segment.Speaker) .Where(speaker => !string.IsNullOrWhiteSpace(speaker)) .Distinct(StringComparer.OrdinalIgnoreCase)) { if (knownLabels.Contains(speaker) || identifiedNames.Contains(speaker) || !SpeakerIdentityNaming.IsDiarizedSpeakerLabel(speaker)) { continue; } var vectors = await ResolveAutomaticVectorsAsync( request, speaker, allowAudioFallback, cancellationToken); if (vectors.Count < resemblyzerOptions.RequiredVectorsPerSpeaker) { logger.LogInformation( "Resemblyzer matching waits for more vectors for {Speaker}: {VectorCount}/{RequiredVectorCount}", speaker, vectors.Count, resemblyzerOptions.RequiredVectorsPerSpeaker); continue; } var (identity, decision) = await FindMatchAsync( context, attendees, identifiedNames, vectors, cancellationToken); if (identity is null) { if (final && decision.Cohesion >= resemblyzerOptions.MinimumClusterCohesion) { unmatchedSpeakers.Add((speaker, vectors)); } continue; } var now = DateTimeOffset.UtcNow; var previousCanonicalName = identity.CanonicalName; var previousReferenceCount = identity.References.Count; outlierPruner.AddAndPrune( identity, vectors, now); SpeakerIdentityReferences.AddIfMissing( identity, CreateReference(request.MeetingNote, now), now); if (identity.References.Count != previousReferenceCount) { identity.UpdatedAt = now; } if (final) { UpdateMatchedIdentity(identity, attendees); if (string.IsNullOrWhiteSpace(previousCanonicalName) && !string.IsNullOrWhiteSpace(identity.CanonicalName)) { pendingAudits.Add(new PendingIdentificationAudit( identity, speaker, identity.CanonicalName)); } foreach (var acceptedName in SpeakerIdentityNaming.GetAcceptedNames(identity)) { matchedAcceptedNames.Add(acceptedName); } } var displayName = identity.GetDisplayName(); if (!string.IsNullOrWhiteSpace(displayName)) { mappings[speaker] = displayName; identifiedNames.Add(displayName); foreach (var acceptedName in SpeakerIdentityNaming.GetAcceptedNames(identity)) { identifiedNames.Add(acceptedName); } attendeeMatches.Add(new SpeakerIdentityAttendeeMatch( displayName, SpeakerIdentityNaming.GetAcceptedNames(identity).ToList())); } } if (final) { LearnUnmatchedSpeakers( context, request.MeetingNote, attendees, matchedAcceptedNames, unmatchedSpeakers, pendingAudits); } await context.SaveChangesAsync(cancellationToken); await AppendAuditsAsync(pendingAudits, cancellationToken); var relabeled = request.Segments .Select(segment => mappings.TryGetValue(segment.Speaker, out var name) ? segment with { Speaker = name } : segment) .ToList(); return new SpeakerIdentificationResult(relabeled, mappings, attendeeMatches); } private async Task<(SpeakerIdentity? Identity, ResemblyzerVoiceClusterMatchResult Decision)> FindMatchAsync( SpeakerIdentityDbContext context, IReadOnlyList attendees, IReadOnlySet identifiedNames, IReadOnlyList queryVectors, CancellationToken cancellationToken) { var activeCutoff = DateTimeOffset.UtcNow - options.MatchIdentityActiveAge; var identities = await LoadIdentities(context) .OrderByDescending(identity => identity.References.Count) .ThenBy(identity => identity.Id) .ToListAsync(cancellationToken); var candidates = identities .Select(identity => new { Identity = identity, IsAttendee = SpeakerIdentityNaming.MatchesAttendees(identity, attendees), IsActive = identity.UpdatedAt >= activeCutoff }) .Where(candidate => candidate.IsAttendee || candidate.IsActive) .Where(candidate => !SpeakerIdentityNaming.MatchesAnyAcceptedName(candidate.Identity, identifiedNames)) .Where(candidate => candidate.Identity.VoiceVectors.Any(vector => string.Equals(vector.ModelId, resemblyzerOptions.ModelId, StringComparison.Ordinal))) .OrderByDescending(candidate => candidate.IsAttendee) .ThenByDescending(candidate => candidate.Identity.ReferenceCount) .ThenBy(candidate => candidate.Identity.Id) .Take(Math.Max(1, options.MaxMatchCandidates)) .Select(candidate => new ResemblyzerVoiceVectorCandidate( candidate.Identity.Id, SpeakerVoiceVectors.DecodeCompatible( candidate.Identity, resemblyzerOptions.ModelId, logger))) .Where(candidate => candidate.Vectors.Count > 0) .ToList(); var match = clusterMatcher.Match(queryVectors, candidates); return ( match.IdentityId is { } identityId ? identities.Single(identity => identity.Id == identityId) : null, match); } private void LearnUnmatchedSpeakers( SpeakerIdentityDbContext context, MeetingNote meetingNote, IReadOnlyList attendees, IReadOnlySet matchedAcceptedNames, IReadOnlyList<(string Speaker, IReadOnlyList Vectors)> unmatchedSpeakers, ICollection pendingAudits) { var candidates = attendees .Except(matchedAcceptedNames, StringComparer.OrdinalIgnoreCase) .Order(StringComparer.OrdinalIgnoreCase) .ToList(); if (candidates.Count == 0) { return; } foreach (var (speaker, vectors) in unmatchedSpeakers) { var now = DateTimeOffset.UtcNow; var identity = new SpeakerIdentity { CanonicalName = candidates.Count == 1 ? candidates[0] : null, CreatedAt = now, UpdatedAt = now, CandidateNames = candidates .Select(name => new SpeakerCandidateName { Name = name }) .ToList(), References = [CreateReference(meetingNote, now)] }; outlierPruner.AddAndPrune( identity, vectors, now); context.SpeakerIdentities.Add(identity); logger.LogInformation( "Created Resemblyzer identity candidate for {Speaker} with {VectorCount} vector(s) and candidates {Candidates}", speaker, identity.VoiceVectors.Count, string.Join(", ", candidates)); if (!string.IsNullOrWhiteSpace(identity.CanonicalName)) { pendingAudits.Add(new PendingIdentificationAudit( identity, speaker, identity.CanonicalName)); } } } private async Task AppendAuditsAsync( IEnumerable pendingAudits, CancellationToken cancellationToken) { foreach (var audit in pendingAudits) { try { await SpeakerIdentityTranscriptAudit.AppendIdentifiedAsync( audit.Identity.References, audit.Speaker, audit.Name, cancellationToken); } catch (Exception exception) when (exception is not OperationCanceledException) { logger.LogError( exception, "Resemblyzer identity {IdentityId} was saved, but its transcript identification audit could not be written", audit.Identity.Id); } } } private async Task> ResolveAutomaticVectorsAsync( SpeakerIdentificationRequest request, string speaker, bool allowAudioFallback, CancellationToken cancellationToken) { var requiredCount = resemblyzerOptions.RequiredVectorsPerSpeaker; var samples = await ResolveWavSamplesAsync( request, speaker, resemblyzerOptions.MaxVectorsPerIdentity, allowAudioFallback, cancellationToken); if (samples.Count < requiredCount) { return []; } return await encoder.EncodeAsync(samples, cancellationToken); } private async Task PersistMappedSpeakerEvidenceAsync( SpeakerIdentityDbContext context, SpeakerIdentificationRequest request, IReadOnlyDictionary knownMappings, CancellationToken cancellationToken) { var identities = await LoadIdentities(context).ToListAsync(cancellationToken); foreach (var (speaker, mappedName) in knownMappings) { var identity = identities .Where(candidate => SpeakerIdentityNaming.GetAcceptedNames(candidate).Contains(mappedName)) .OrderBy(candidate => string.Equals(candidate.CanonicalName, mappedName, StringComparison.OrdinalIgnoreCase) ? 0 : 1) .ThenBy(candidate => candidate.Id) .FirstOrDefault(); if (identity is null) { logger.LogWarning( "Could not retain final Resemblyzer evidence for mapped speaker {Speaker}: identity {MappedName} was not found", speaker, mappedName); continue; } var vectors = await ResolveAvailableVectorsAsync(request, speaker, cancellationToken); if (vectors.Count == 0) { continue; } var now = DateTimeOffset.UtcNow; var previousReferenceCount = identity.References.Count; var vectorUpdate = outlierPruner.AddAndPrune( identity, vectors, now); SpeakerIdentityReferences.AddIfMissing( identity, CreateReference(request.MeetingNote, now), now); if (identity.References.Count != previousReferenceCount) { identity.UpdatedAt = now; } logger.LogInformation( "Retained {AddedVectorCount} new final Resemblyzer vector(s) for mapped speaker {Speaker} as identity {IdentityId}", vectorUpdate.AddedCount, speaker, identity.Id); } } private async Task> ResolveAvailableVectorsAsync( SpeakerIdentificationRequest request, string speaker, CancellationToken cancellationToken) { var samples = await ResolveWavSamplesAsync( request, speaker, resemblyzerOptions.MaxVectorsPerIdentity, allowAudioFallback: true, cancellationToken); if (samples.Count == 0) { return []; } return await encoder.EncodeAsync(samples, cancellationToken); } private async Task> ResolveWavSamplesAsync( SpeakerIdentificationRequest request, string speaker, int maxSamples, bool allowAudioFallback, CancellationToken cancellationToken) { var suppliedSamples = request.Samples? .Where(sample => string.Equals(sample.Speaker, speaker, StringComparison.OrdinalIgnoreCase)) .Where(sample => sample.WavBytes.Length > 0) .OrderByDescending(sample => sample.Score) .Take(maxSamples) .ToList() ?? []; var wavSamples = suppliedSamples .Select(sample => sample.WavBytes) .ToList(); if (!allowAudioFallback || wavSamples.Count >= maxSamples) { return wavSamples; } var spans = SpeakerSampleSpanSelector.SelectBestSameSpeakerSpans( request.Segments, speaker, options.MinimumSampleSpeechDuration, options.MaximumSampleDuration, request.Segments.Count); foreach (var span in spans.Where(span => !OverlapsSuppliedSample(span, suppliedSamples))) { var wavBytes = await snippetExtractor.ExtractSnippetAsync( request.AudioPath, span, cancellationToken); if (wavBytes.Length > 0) { wavSamples.Add(wavBytes); } if (wavSamples.Count >= maxSamples) { break; } } logger.LogInformation( "Resolved {SampleCount}/{RequestedSampleCount} Resemblyzer WAV samples for {Speaker}: {SuppliedSampleCount} supplied, {ExtractedSampleCount} extracted from completed audio", wavSamples.Count, maxSamples, speaker, suppliedSamples.Count, wavSamples.Count - suppliedSamples.Count); return wavSamples; } private static bool OverlapsSuppliedSample( IReadOnlyList span, IReadOnlyList suppliedSamples) { if (span.Count == 0) { return false; } var start = span[0].Start; var end = span[^1].End; return suppliedSamples.Any(sample => sample.Segment.Start < end && sample.Segment.End > start); } private static IQueryable LoadIdentities(SpeakerIdentityDbContext context) { return context.SpeakerIdentities .AsSplitQuery() .Include(identity => identity.Aliases) .Include(identity => identity.CandidateNames) .Include(identity => identity.Snippets) .Include(identity => identity.VoiceVectors) .Include(identity => identity.References); } private void UpdateMatchedIdentity(SpeakerIdentity identity, IReadOnlyList attendees) { identity.UpdatedAt = DateTimeOffset.UtcNow; if (!string.IsNullOrWhiteSpace(identity.CanonicalName) || attendees.Count == 0) { return; } SpeakerIdentityNaming.UpdateCandidateNames(identity, attendees); } private static SpeakerIdentityReference CreateReference(MeetingNote meetingNote, DateTimeOffset timestamp) { return SpeakerIdentityReferences.Create(meetingNote.Path, meetingNote.Frontmatter.Transcript, timestamp); } private static SpeakerIdentificationResult EmptyResult(IReadOnlyList segments) { return new SpeakerIdentificationResult(segments, new Dictionary()); } private sealed record PendingIdentificationAudit( SpeakerIdentity Identity, string Speaker, string Name); }