// UtteranceStreamChunker.swift // OSGKeyboard · Shared // // Splits a Flow utterance PCM stream into ASR-sized chunks. When possible, // extends slightly past the max window to the next pause instead of cutting // mid-word. import Foundation public enum UtteranceStreamChunker { /// Yields chunks as audio arrives; the final chunk is marked `isLast`. public static func chunks( from stream: AsyncStream, config: FlowUtteranceChunkConfig = .flowDefault ) -> AsyncStream { AsyncStream { continuation in let task = Task { var buffer: [Float] = [] let initialCapacity = config.maxChunkSamples(forChunkIndex: 0) + config.pauseExtensionSamples buffer.reserveCapacity(initialCapacity) var chunkIndex = 0 func emit( upTo splitEnd: Int, isLast: Bool, trailingPauseSeconds: Double = 0 ) { guard splitEnd > 0, splitEnd <= buffer.count else { FlowTrace.warn( "pipeline.chunk.emitSkipped", "chunk=\(chunkIndex) splitEnd=\(splitEnd) buffered=\(buffer.count)" ) return } let chunkSamples = Array(buffer[..= buffer.count { buffer.removeAll(keepingCapacity: true) } else { let overlapStart = max(0, splitEnd - config.overlapSamples) buffer = Array(buffer[overlapStart...]) } } var receivedSnapshots = 0 var receivedSamples = 0 for await snap in stream { if Task.isCancelled { break } guard !snap.samples.isEmpty else { continue } receivedSnapshots += 1 receivedSamples += snap.samples.count buffer.append(contentsOf: snap.samples) while buffer.count >= config.maxChunkSamples(forChunkIndex: chunkIndex) { let split = pauseAwareSplit( in: buffer, config: config, chunkIndex: chunkIndex ) emit( upTo: split.index, isLast: false, trailingPauseSeconds: Double(split.pauseSamples) / Double(config.sampleRate) ) } } FlowTrace.pipeline( "chunk.streamEnded", "snapshots=\(receivedSnapshots) samples=\(receivedSamples) " + "seconds=\(FlowTrace.seconds(samples: receivedSamples, sampleRate: config.sampleRate)) " + "chunksEmitted=\(chunkIndex) buffered=\(buffer.count) " + "cancelled=\(Task.isCancelled ? 1 : 0)" ) if !buffer.isEmpty { emit(upTo: buffer.count, isLast: true) } else if chunkIndex == 0 { // Empty utterance — no chunks. The recogniser is never // invoked, so an empty transcript here means the mic stream // itself was empty, not that recognition failed. FlowTrace.warn( "pipeline.chunk.emptyUtterance", "snapshots=\(receivedSnapshots) samples=0 chunksEmitted=0" ) } else { // Stream ended exactly on a chunk boundary; prior emit holds // all tail audio. Marker so FinalChunkRecovery paths run. continuation.yield( UtteranceAudioChunk(index: chunkIndex, samples: [], isLast: true) ) } continuation.finish() } continuation.onTermination = { _ in task.cancel() } } } /// Pick a split index at or after `maxChunkSamples`, preferring a pause. static func pauseAwareSplitIndex( in buffer: [Float], config: FlowUtteranceChunkConfig, chunkIndex: Int = 1 ) -> Int { pauseAwareSplit(in: buffer, config: config, chunkIndex: chunkIndex).index } static func pauseAwareSplit( in buffer: [Float], config: FlowUtteranceChunkConfig, chunkIndex: Int = 1 ) -> (index: Int, pauseSamples: Int) { let minSplit = config.maxChunkSamples(forChunkIndex: chunkIndex) guard buffer.count >= minSplit else { return (buffer.count, 0) } let searchEnd = min(buffer.count, minSplit + config.pauseExtensionSamples) if searchEnd <= minSplit { return (minSplit, 0) } let windowSize = max(config.sampleRate / 50, 160) // ~20 ms let step = max(windowSize / 2, 1) var bestPauseEnd: Int? var bestPauseSamples = 0 var currentPauseStart: Int? var idx = minSplit while idx + windowSize <= searchEnd { if rms(of: buffer, start: idx, count: windowSize) < config.pauseRMSThreshold { if currentPauseStart == nil { currentPauseStart = idx } let pauseSamples = idx + windowSize - (currentPauseStart ?? idx) if pauseSamples > bestPauseSamples { bestPauseSamples = pauseSamples bestPauseEnd = idx + windowSize } } else { currentPauseStart = nil } idx += step } return (bestPauseEnd ?? minSplit, bestPauseSamples) } static func rms(of samples: [Float], start: Int, count: Int) -> Float { guard start >= 0, count > 0, start + count <= samples.count else { return 1 } var sum: Float = 0 for i in start..<(start + count) { let v = samples[i] sum += v * v } return sqrtf(sum / Float(count)) } }