mirror of
https://github.com/FluidInference/FluidAudio.git
synced 2026-06-11 20:24:36 +00:00
This PR addresses three high-priority consistency improvements in the Parakeet ASR folder from issue #457. ## Summary - ✅ **Task 1:** Standardized lifecycle method names across all managers (13 files) - ✅ **Task 2:** Consolidated ~230 lines of duplicate token deduplication logic - ✅ **Task 3:** Extracted shared streaming code into reusable utilities ## Changes ### 1. Lifecycle Method Standardization Unified naming conventions to eliminate confusion: | Manager | Old Method | New Method | |---------|-----------|------------| | `AsrManager` | `loadModels(_:)` | `configure(models:)` | | `SlidingWindowAsrSession` | `initialize()` | `loadModels()` | | `SlidingWindowAsrManager` | `start()` | `startStreaming()` | | `StreamingEouAsrManager` | `loadModelsFromHuggingFace()` | `loadModels()` | **Files updated:** 5 managers + 8 CLI commands ### 2. Token Deduplication Consolidation Extracted duplicate matching algorithms into generic, type-safe utilities: **New Files:** - `SequenceMatch.swift` - Data structure for sequence matches - `SequenceMatcher.swift` - 5 reusable matching algorithms: - `findSuffixPrefixMatch()` - O(n) greedy boundary detection - `findBoundedSubstringMatch()` - Windowed search - `findLongestCommonSubsequence()` - O(n²) LCS via DP - `findContiguousMatches()` - Longest consecutive run - `consolidateMatches()` - Merge adjacent matches - `TokenDeduplicationRegressionTests.swift` - 12 comprehensive tests **Refactored:** - `AsrManager+TokenProcessing.swift` - Reduced from ~65 to ~40 lines (-38%) - `ChunkProcessor.swift` - Removed ~77 lines of duplicate code ### 3. Streaming Code Extraction Created utilities for common patterns in both `StreamingEouAsrManager` and `StreamingNemotronAsrManager`: **New Utilities:** - `EncoderCacheManager` - Cache initialization and extraction - `StreamingAsrUtils` - Audio buffering, state reset, token decoding ## Impact | Metric | Result | |--------|--------| | **Duplicate code eliminated** | ~230 lines | | **New reusable utilities** | 430 lines | | **Test coverage** | +12 regression tests | | **API consistency** | Unified lifecycle naming | | **Performance** | No regression ✅ | | **WER** | 0.4% (verified) ✅ | | **RTFx** | 43.3x (verified) ✅ | | **Tests** | 25/25 passing ✅ | ## Testing ```bash # Token deduplication regression tests swift test --filter TokenDeduplicationRegressionTests # ✅ 12/12 tests passing # Nemotron streaming tests swift test --filter StreamingNemotronAsrManagerTests # ✅ 16/16 tests passing # ASR benchmark (no WER regression) swift run -c release fluidaudiocli asr-benchmark --max-files 10 # ✅ WER: 0.4%, RTFx: 43.3x ``` ## Breaking Changes ⚠️ This PR contains breaking API changes: - Renamed lifecycle methods (no deprecation wrappers) - All call sites updated in this PR Closes #457 <!-- devin-review-badge-begin --> --- <a href="https://app.devin.ai/review/fluidinference/fluidaudio/pull/494" target="_blank"> <picture> <source media="(prefers-color-scheme: dark)" srcset="https://static.devin.ai/assets/gh-open-in-devin-review-dark.svg?v=1"> <img src="https://static.devin.ai/assets/gh-open-in-devin-review-light.svg?v=1" alt="Open with Devin"> </picture> </a> <!-- devin-review-badge-end --> ---------
261 lines
9.1 KiB
Swift
261 lines
9.1 KiB
Swift
#if os(macOS)
|
|
@preconcurrency import AVFoundation
|
|
import FluidAudio
|
|
import Foundation
|
|
|
|
/// Command to demonstrate multi-stream ASR with shared model loading
|
|
enum MultiStreamCommand {
|
|
private static let logger = AppLogger(category: "MultiStream")
|
|
|
|
static func run(arguments: [String]) async {
|
|
// Parse arguments
|
|
guard !arguments.isEmpty else {
|
|
logger.error("No audio files specified")
|
|
printUsage()
|
|
exit(1)
|
|
}
|
|
|
|
let audioFile1 = arguments[0]
|
|
var audioFile2: String? = nil
|
|
|
|
// Parse options
|
|
var i = 1
|
|
while i < arguments.count {
|
|
switch arguments[i] {
|
|
case "--help", "-h":
|
|
printUsage()
|
|
exit(0)
|
|
default:
|
|
// Check if it's a second audio file
|
|
if audioFile2 == nil && arguments[i].hasSuffix(".wav") {
|
|
audioFile2 = arguments[i]
|
|
} else {
|
|
logger.warning("Unknown option: \(arguments[i])")
|
|
}
|
|
}
|
|
i += 1
|
|
}
|
|
|
|
// Use same file for both streams if only one provided
|
|
let micAudioFile = audioFile1
|
|
let systemAudioFile = audioFile2 ?? audioFile1
|
|
|
|
logger.info("🎤 Multi-Stream ASR Test\n========================\n")
|
|
|
|
if audioFile2 != nil {
|
|
logger.info(
|
|
"📁 Processing two different files:\n Microphone: \(micAudioFile)\n System: \(systemAudioFile)\n")
|
|
} else {
|
|
logger.info("📁 Processing single file on both streams: \(audioFile1)\n")
|
|
}
|
|
|
|
do {
|
|
// Load first audio file (microphone)
|
|
let micFileURL = URL(fileURLWithPath: micAudioFile)
|
|
let micFileHandle = try AVAudioFile(forReading: micFileURL)
|
|
let micFormat = micFileHandle.processingFormat
|
|
let micFrameCount = AVAudioFrameCount(micFileHandle.length)
|
|
|
|
guard
|
|
let micBuffer = AVAudioPCMBuffer(
|
|
pcmFormat: micFormat, frameCapacity: micFrameCount)
|
|
else {
|
|
logger.error("Failed to create microphone audio buffer")
|
|
return
|
|
}
|
|
try micFileHandle.read(into: micBuffer)
|
|
|
|
// Load second audio file (system)
|
|
let systemFileURL = URL(fileURLWithPath: systemAudioFile)
|
|
let systemFileHandle = try AVAudioFile(forReading: systemFileURL)
|
|
let systemFormat = systemFileHandle.processingFormat
|
|
let systemFrameCount = AVAudioFrameCount(systemFileHandle.length)
|
|
|
|
guard
|
|
let systemBuffer = AVAudioPCMBuffer(
|
|
pcmFormat: systemFormat, frameCapacity: systemFrameCount)
|
|
else {
|
|
logger.error("Failed to create system audio buffer")
|
|
return
|
|
}
|
|
try systemFileHandle.read(into: systemBuffer)
|
|
|
|
logger.info(
|
|
"""
|
|
📊 Audio file info:
|
|
🎙️ Microphone file:
|
|
Sample rate: \(micFormat.sampleRate) Hz
|
|
Channels: \(micFormat.channelCount)
|
|
Duration: \(String(format: "%.2f", Double(micFileHandle.length) / micFormat.sampleRate)) seconds
|
|
|
|
System audio file:
|
|
Sample rate: \(systemFormat.sampleRate) Hz
|
|
Channels: \(systemFormat.channelCount)
|
|
Duration: \(String(format: "%.2f", Double(systemFileHandle.length) / systemFormat.sampleRate)) seconds
|
|
"""
|
|
)
|
|
|
|
// Create a streaming session
|
|
logger.info("Creating streaming session...")
|
|
let session = SlidingWindowAsrSession()
|
|
|
|
// Initialize models once
|
|
logger.info("Loading ASR models (shared across streams)...")
|
|
let startTime = Date()
|
|
try await session.loadModels()
|
|
let loadTime = Date().timeIntervalSince(startTime)
|
|
logger.info("Models loaded in \(String(format: "%.2f", loadTime))s\n")
|
|
|
|
// Create streams for different sources
|
|
logger.info("Creating streams for different audio sources...")
|
|
let micStream = try await session.createStream(
|
|
source: .microphone,
|
|
config: .default
|
|
)
|
|
logger.info("Created microphone stream")
|
|
|
|
let systemStream = try await session.createStream(
|
|
source: .system,
|
|
config: .default
|
|
)
|
|
logger.info("Created system audio stream\n")
|
|
|
|
// Listen for updates from both streams (only if debug enabled)
|
|
let micTask = Task {
|
|
for await update in await micStream.transcriptionUpdates {
|
|
logger.info("[MIC] \(update.isConfirmed ? "✓" : "~") \(update.text)")
|
|
}
|
|
}
|
|
|
|
let systemTask = Task {
|
|
for await update in await systemStream.transcriptionUpdates {
|
|
logger.info("[SYS] \(update.isConfirmed ? "✓" : "~") \(update.text)")
|
|
}
|
|
}
|
|
|
|
logger.info("Streaming audio files in parallel...\n Both streams using default config (10.0s chunks)\n")
|
|
|
|
// Process both files in parallel
|
|
let micProcessingTask = Task {
|
|
await streamAudioFile(
|
|
buffer: micBuffer,
|
|
format: micFormat,
|
|
to: micStream,
|
|
label: "MIC"
|
|
)
|
|
}
|
|
|
|
let systemProcessingTask = Task {
|
|
await streamAudioFile(
|
|
buffer: systemBuffer,
|
|
format: systemFormat,
|
|
to: systemStream,
|
|
label: "SYS"
|
|
)
|
|
}
|
|
|
|
// Wait for both to complete
|
|
await micProcessingTask.value
|
|
await systemProcessingTask.value
|
|
|
|
logger.info("Finalizing transcriptions...")
|
|
|
|
// Get final results
|
|
let micFinal = try await micStream.finish()
|
|
let systemFinal = try await systemStream.finish()
|
|
|
|
// Cancel update tasks
|
|
micTask.cancel()
|
|
systemTask.cancel()
|
|
|
|
// Print results
|
|
logger.info(
|
|
"" + String(repeating: "=", count: 60) + "TRANSCRIPTION RESULTS\n"
|
|
+ String(repeating: "=", count: 60) + "")
|
|
logger.info("MICROPHONE STREAM:\n\(micFinal)")
|
|
logger.info("SYSTEM AUDIO STREAM:\n\(systemFinal)")
|
|
logger.info("Session info:")
|
|
let activeStreams = await session.activeStreams
|
|
logger.info("Active streams: \(activeStreams.count)")
|
|
for (source, stream) in activeStreams {
|
|
logger.info(" - \(source): \(await stream.source)")
|
|
}
|
|
|
|
await session.cleanup()
|
|
} catch {
|
|
logger.error("Error: \(error)")
|
|
}
|
|
}
|
|
|
|
/// Helper function to stream an audio file to a stream
|
|
private static func streamAudioFile(
|
|
buffer: AVAudioPCMBuffer,
|
|
format: AVAudioFormat,
|
|
to stream: SlidingWindowAsrManager,
|
|
label: String
|
|
) async {
|
|
let chunkDuration = 0.5 // 500ms chunks
|
|
let samplesPerChunk = Int(chunkDuration * format.sampleRate)
|
|
var position = 0
|
|
|
|
while position < Int(buffer.frameLength) {
|
|
let remainingSamples = Int(buffer.frameLength) - position
|
|
let chunkSize = min(samplesPerChunk, remainingSamples)
|
|
|
|
// Create chunk buffer
|
|
guard
|
|
let chunkBuffer = AVAudioPCMBuffer(
|
|
pcmFormat: format,
|
|
frameCapacity: AVAudioFrameCount(chunkSize)
|
|
)
|
|
else {
|
|
break
|
|
}
|
|
|
|
for channel in 0..<Int(format.channelCount) {
|
|
if let sourceData = buffer.floatChannelData?[channel],
|
|
let destData = chunkBuffer.floatChannelData?[channel]
|
|
{
|
|
for i in 0..<chunkSize {
|
|
destData[i] = sourceData[position + i]
|
|
}
|
|
}
|
|
}
|
|
chunkBuffer.frameLength = AVAudioFrameCount(chunkSize)
|
|
|
|
await stream.streamAudio(chunkBuffer)
|
|
|
|
position += chunkSize
|
|
}
|
|
|
|
logger.info("[\(label)] Streaming complete")
|
|
}
|
|
|
|
private static func printUsage() {
|
|
logger.info(
|
|
"""
|
|
|
|
Multi-Stream Command Usage:
|
|
fluidaudio multi-stream <audio_file1> [audio_file2] [options]
|
|
|
|
Options:
|
|
--help, -h Show this help message
|
|
|
|
Examples:
|
|
# Process same file on both streams
|
|
fluidaudio multi-stream audio.wav
|
|
|
|
# Process two different files in parallel
|
|
fluidaudio multi-stream mic_audio.wav system_audio.wav
|
|
|
|
|
|
This command demonstrates:
|
|
- Loading ASR models once and sharing across streams
|
|
- Creating separate streams for microphone and system audio
|
|
- Parallel transcription with shared resources
|
|
"""
|
|
)
|
|
}
|
|
}
|
|
#endif
|