diff --git a/Sources/AudioTeeCLI/AudioTee.swift b/Sources/AudioTeeCLI/AudioTee.swift index cb9d1a9..e7f8973 100644 --- a/Sources/AudioTeeCLI/AudioTee.swift +++ b/Sources/AudioTeeCLI/AudioTee.swift @@ -9,7 +9,6 @@ struct AudioTee { var stereo: Bool = false var sampleRate: Double? var chunkDuration: Double = 0.2 - var flush: Bool = false init() {} @@ -33,7 +32,6 @@ struct AudioTee { audiotee --include-processes 1234 5678 9012 # Tap only these processes audiotee --exclude-processes 1234 5678 # Tap everything except these audiotee --mute # Mute processes being tapped - audiotee --flush # Flush stdout after each chunk """ ) @@ -45,8 +43,6 @@ struct AudioTee { name: "exclude-processes", help: "Process IDs to exclude (space-separated)") parser.addFlag(name: "mute", help: "Mute processes being tapped") parser.addFlag(name: "stereo", help: "Records in stereo") - parser.addFlag( - name: "flush", help: "Flush stdout after each audio chunk (reduces latency when piping)") parser.addOption( name: "sample-rate", help: "Target sample rate (8000, 16000, 22050, 24000, 32000, 44100, 48000)") @@ -64,7 +60,6 @@ struct AudioTee { audioTee.excludeProcesses = try parser.getArrayValue("exclude-processes", as: Int32.self) audioTee.mute = parser.getFlag("mute") audioTee.stereo = parser.getFlag("stereo") - audioTee.flush = parser.getFlag("flush") audioTee.sampleRate = try parser.getOptionalValue("sample-rate", as: Double.self) audioTee.chunkDuration = try parser.getValue("chunk-duration", as: Double.self) @@ -142,7 +137,7 @@ struct AudioTee { throw ExitCode.failure } - let outputHandler = BinaryAudioOutputHandler(flushAfterWrite: flush) + let outputHandler = BinaryAudioOutputHandler() let recorder = try AudioRecorder( deviceID: deviceID, outputHandler: outputHandler, convertToSampleRate: sampleRate, chunkDuration: chunkDuration) diff --git a/Sources/AudioTeeCLI/BinaryOutputHandler.swift b/Sources/AudioTeeCLI/BinaryOutputHandler.swift index 9f8a5b6..732b837 100644 --- a/Sources/AudioTeeCLI/BinaryOutputHandler.swift +++ b/Sources/AudioTeeCLI/BinaryOutputHandler.swift @@ -4,17 +4,19 @@ import Foundation /// CLI-specific output handler that writes raw PCM audio to stdout /// and lifecycle messages to stderr via the logger. class BinaryAudioOutputHandler: AudioOutputHandler { - private let flushAfterWrite: Bool + private let fd = STDOUT_FILENO - init(flushAfterWrite: Bool = false) { - self.flushAfterWrite = flushAfterWrite - } - - func handleAudioPacket(_ packet: AudioPacket) { - // Write raw binary audio data directly to stdout - FileHandle.standardOutput.write(packet.data) - if flushAfterWrite { - fflush(stdout) + func handleAudioData(_ pointer: UnsafeRawPointer, count: Int) { + var written = 0 + while written < count { + let result = write(fd, pointer.advanced(by: written), count - written) + if result >= 0 { + written += result + } else if errno == EINTR { + continue + } else { + break // EPIPE, EIO, etc — consumer gone or real error + } } } diff --git a/Sources/AudioTeeCore/Core/AudioBuffer.swift b/Sources/AudioTeeCore/Core/AudioBuffer.swift index 636cd9c..7b5c781 100644 --- a/Sources/AudioTeeCore/Core/AudioBuffer.swift +++ b/Sources/AudioTeeCore/Core/AudioBuffer.swift @@ -10,20 +10,21 @@ import Foundation public class AudioBuffer { /// Raw heap-allocated ring buffer backing store. private let buffer: UnsafeMutableRawPointer + /// Pre-allocated buffer for linearizing chunks that straddle the ring + /// buffer boundary. Avoids a heap allocation on the wrap-around path. + private let linearizationBuffer: UnsafeMutableRawPointer private var writeIndex: Int = 0 private var readIndex: Int = 0 private var availableBytes: Int = 0 private let maxBufferSize: Int - private let bytesPerChunk: Int - private let chunkDuration: Double + public let bytesPerChunk: Int public init(format: AudioStreamBasicDescription, chunkDuration: Double = 0.2) { // Pre-calculate chunk parameters let bytesPerFrame = Int(format.mBytesPerFrame) let samplesPerChunk = Int(format.mSampleRate * chunkDuration) self.bytesPerChunk = samplesPerChunk * bytesPerFrame - self.chunkDuration = Double(samplesPerChunk) / format.mSampleRate // Calculate max buffer size to hold ~10 seconds of audio (safety limit) let bytesPerSecond = Int(format.mSampleRate) * bytesPerFrame @@ -36,10 +37,16 @@ public class AudioBuffer { alignment: MemoryLayout.alignment ) buffer.initializeMemory(as: UInt8.self, repeating: 0, count: maxBufferSize) + + self.linearizationBuffer = UnsafeMutableRawPointer.allocate( + byteCount: bytesPerChunk, + alignment: MemoryLayout.alignment + ) } deinit { buffer.deallocate() + linearizationBuffer.deallocate() } /// Appends audio data directly from a raw pointer into the ring buffer. @@ -82,51 +89,33 @@ public class AudioBuffer { availableBytes += count } - /// Extracts all complete chunks currently available in the buffer. - public func processChunks() -> [AudioPacket] { - var packets: [AudioPacket] = [] + /// Calls `handler` once for each complete chunk available in the buffer. + /// The pointer passed to the handler is valid only for the duration of + /// that call. In the common (contiguous) case this points directly into + /// the ring buffer — zero copies. In the wrap-around case the chunk is + /// linearized into a pre-allocated scratch buffer — one memcpy, zero + /// heap allocations. + public func processChunks(_ handler: (UnsafeRawPointer, Int) -> Void) { + while availableBytes >= bytesPerChunk { + if readIndex + bytesPerChunk <= maxBufferSize { + // Contiguous: point directly into the ring buffer + handler(buffer.advanced(by: readIndex), bytesPerChunk) + readIndex = (readIndex + bytesPerChunk) % maxBufferSize + } else { + // Wrap-around: linearize into the pre-allocated scratch buffer + let firstChunkSize = maxBufferSize - readIndex + let secondChunkSize = bytesPerChunk - firstChunkSize - while let packet = nextChunk() { - packets.append(packet) + linearizationBuffer.copyMemory( + from: buffer.advanced(by: readIndex), byteCount: firstChunkSize) + linearizationBuffer.advanced(by: firstChunkSize).copyMemory( + from: buffer, byteCount: secondChunkSize) + + handler(linearizationBuffer, bytesPerChunk) + readIndex = secondChunkSize + } + + availableBytes -= bytesPerChunk } - - return packets - } - - private func nextChunk() -> AudioPacket? { - // Check if we have enough data for a complete chunk - guard availableBytes >= bytesPerChunk else { return nil } - - let chunkData: Data - - // Check if we can copy in one block (no wrap-around) - if readIndex + bytesPerChunk <= maxBufferSize { - // one copy needed - chunkData = Data(bytes: buffer.advanced(by: readIndex), count: bytesPerChunk) - readIndex = (readIndex + bytesPerChunk) % maxBufferSize - } else { - // two copies needed due to wrap-around - let firstChunkSize = maxBufferSize - readIndex - let secondChunkSize = bytesPerChunk - firstChunkSize - - var assembled = Data(capacity: bytesPerChunk) - assembled.append( - buffer.advanced(by: readIndex).assumingMemoryBound(to: UInt8.self), - count: firstChunkSize) - assembled.append( - buffer.assumingMemoryBound(to: UInt8.self), - count: secondChunkSize) - chunkData = assembled - - readIndex = secondChunkSize - } - - availableBytes -= bytesPerChunk - - return AudioPacket( - timestamp: Date(), - duration: chunkDuration, - data: chunkData - ) } } diff --git a/Sources/AudioTeeCore/Core/AudioFormatConverter.swift b/Sources/AudioTeeCore/Core/AudioFormatConverter.swift index 7bcc95f..5a52284 100644 --- a/Sources/AudioTeeCore/Core/AudioFormatConverter.swift +++ b/Sources/AudioTeeCore/Core/AudioFormatConverter.swift @@ -126,27 +126,28 @@ public class AudioFormatConverter { return (inputBuf, outputBuf) } - public func transform(_ packet: AudioPacket) -> AudioPacket { - let inputData = packet.data - - // Calculate frame count from the input data size + /// Converts audio data in-place through the pre-allocated converter buffers. + /// Calls `handler` with a pointer to the converted output, valid only for + /// the duration of that call. Returns false on failure (caller should + /// pass through the original data or drop it). + @discardableResult + public func transform( + from source: UnsafeRawPointer, count: Int, + handler: (UnsafeRawPointer, Int) -> Void + ) -> Bool { let bytesPerFrame = Int(sourceFormat.streamDescription.pointee.mBytesPerFrame) - let inputFrameCount = AVAudioFrameCount(inputData.count / bytesPerFrame) + let inputFrameCount = AVAudioFrameCount(count / bytesPerFrame) - // Get or create pre-allocated buffers guard let (inputBuffer, outputBuffer) = getBuffers(inputFrameCount: inputFrameCount) else { - return packet + return false } - // Copy input data into the reusable input buffer - inputData.withUnsafeBytes { bytes in - let dest = inputBuffer.audioBufferList.pointee.mBuffers.mData! - dest.copyMemory(from: bytes.baseAddress!, byteCount: inputData.count) - } + // Copy source data into the reusable input buffer + let dest = inputBuffer.audioBufferList.pointee.mBuffers.mData! + dest.copyMemory(from: source, byteCount: count) inputBuffer.frameLength = inputFrameCount - // Perform conversion — the block-based API lets AVAudioConverter pull - // input data as needed. We do NOT call avConverter.reset() between + // Perform conversion — we do NOT call avConverter.reset() between // calls because the resampler maintains internal state for continuity // across chunks (avoiding discontinuity artifacts). var error: NSError? @@ -157,7 +158,6 @@ public class AudioFormatConverter { return inputBuffer } - // Check if conversion produced output (regardless of status code) guard outputBuffer.frameLength > 0 else { AudioTeeLogging.logger.error( "Audio conversion produced no output", @@ -167,19 +167,13 @@ public class AudioFormatConverter { "input_frames": String(inputBuffer.frameLength), "output_capacity": String(outputBuffer.frameCapacity), ]) - return packet + return false } - // Extract converted data from the reusable output buffer - let outputData = Data( - bytes: outputBuffer.audioBufferList.pointee.mBuffers.mData!, - count: Int(outputBuffer.frameLength * targetFormat.streamDescription.pointee.mBytesPerFrame)) - - return AudioPacket( - timestamp: packet.timestamp, - duration: packet.duration, - data: outputData - ) + let outputCount = Int( + outputBuffer.frameLength * targetFormat.streamDescription.pointee.mBytesPerFrame) + handler(outputBuffer.audioBufferList.pointee.mBuffers.mData!, outputCount) + return true } public static func toSampleRate( diff --git a/Sources/AudioTeeCore/Core/AudioPacket.swift b/Sources/AudioTeeCore/Core/AudioPacket.swift deleted file mode 100644 index 8db7d13..0000000 --- a/Sources/AudioTeeCore/Core/AudioPacket.swift +++ /dev/null @@ -1,17 +0,0 @@ -import Foundation - -public struct AudioPacket { - public let timestamp: Date - public let duration: Double - public let data: Data - - public init( - timestamp: Date, - duration: Double, - data: Data - ) { - self.timestamp = timestamp - self.duration = duration - self.data = data - } -} diff --git a/Sources/AudioTeeCore/Core/AudioRecorder.swift b/Sources/AudioTeeCore/Core/AudioRecorder.swift index d0564c6..6d8dce2 100644 --- a/Sources/AudioTeeCore/Core/AudioRecorder.swift +++ b/Sources/AudioTeeCore/Core/AudioRecorder.swift @@ -132,10 +132,17 @@ public class AudioRecorder { } private func processAudioBuffer() { - // Process and send complete chunks, applying conversion if needed - audioBuffer?.processChunks().forEach { packet in - let processedPacket = converter?.transform(packet) ?? packet - outputHandler.handleAudioPacket(processedPacket) + audioBuffer?.processChunks { pointer, count in + if let converter = self.converter { + if !converter.transform(from: pointer, count: count, handler: { outPtr, outCount in + self.outputHandler.handleAudioData(outPtr, count: outCount) + }) { + // Conversion failed — pass through unconverted audio + self.outputHandler.handleAudioData(pointer, count: count) + } + } else { + self.outputHandler.handleAudioData(pointer, count: count) + } } } diff --git a/Sources/AudioTeeCore/Output/AudioOutputProtocol.swift b/Sources/AudioTeeCore/Output/AudioOutputProtocol.swift index 6175a8b..5a100f5 100644 --- a/Sources/AudioTeeCore/Output/AudioOutputProtocol.swift +++ b/Sources/AudioTeeCore/Output/AudioOutputProtocol.swift @@ -2,7 +2,9 @@ import Foundation /// Protocol for handling audio output in different formats public protocol AudioOutputHandler { - func handleAudioPacket(_ packet: AudioPacket) + /// Called with a pointer to raw PCM audio data. The pointer is only + /// valid for the duration of this call. + func handleAudioData(_ pointer: UnsafeRawPointer, count: Int) func handleMetadata(_ metadata: AudioStreamMetadata) func handleStreamStart() func handleStreamStop() diff --git a/Tests/AudioTeeCoreTests/AudioBufferTests.swift b/Tests/AudioTeeCoreTests/AudioBufferTests.swift index d2153d1..2268e12 100644 --- a/Tests/AudioTeeCoreTests/AudioBufferTests.swift +++ b/Tests/AudioTeeCoreTests/AudioBufferTests.swift @@ -44,6 +44,15 @@ final class AudioBufferTests: XCTestCase { } } + /// Collects chunks from the buffer as Data objects for test verification. + private func collectChunks(from buffer: AudioBuffer) -> [Data] { + var chunks: [Data] = [] + buffer.processChunks { pointer, count in + chunks.append(Data(bytes: pointer, count: count)) + } + return chunks + } + // MARK: - Basic append + processChunks func testSingleChunkExtraction() { @@ -55,10 +64,10 @@ final class AudioBufferTests: XCTestCase { let data = makeData(byte: 0xAB, count: chunkSize) appendData(data, to: buffer) - let packets = buffer.processChunks() - XCTAssertEqual(packets.count, 1) - XCTAssertEqual(packets[0].data.count, chunkSize) - XCTAssertEqual(packets[0].data, data) + let chunks = collectChunks(from: buffer) + XCTAssertEqual(chunks.count, 1) + XCTAssertEqual(chunks[0].count, chunkSize) + XCTAssertEqual(chunks[0], data) } func testMultipleChunksExtracted() { @@ -69,11 +78,11 @@ final class AudioBufferTests: XCTestCase { // Append 2.5 chunks worth appendData(makeData(byte: 0x01, count: chunkSize * 2 + chunkSize / 2), to: buffer) - let packets = buffer.processChunks() + let chunks = collectChunks(from: buffer) // Should get 2 complete chunks, remainder stays in buffer - XCTAssertEqual(packets.count, 2) - XCTAssertEqual(packets[0].data.count, chunkSize) - XCTAssertEqual(packets[1].data.count, chunkSize) + XCTAssertEqual(chunks.count, 2) + XCTAssertEqual(chunks[0].count, chunkSize) + XCTAssertEqual(chunks[1].count, chunkSize) } func testInsufficientDataReturnsNoChunks() { @@ -84,8 +93,8 @@ final class AudioBufferTests: XCTestCase { // Append less than one chunk appendData(makeData(byte: 0xFF, count: chunkSize - 1), to: buffer) - let packets = buffer.processChunks() - XCTAssertEqual(packets.count, 0) + let chunks = collectChunks(from: buffer) + XCTAssertEqual(chunks.count, 0) } // MARK: - Wrap-around @@ -103,7 +112,7 @@ final class AudioBufferTests: XCTestCase { for _ in 0..<33 { appendData(makeData(byte: 0x00, count: chunkSize), to: buffer) } - let drained = buffer.processChunks() + let drained = collectChunks(from: buffer) XCTAssertEqual(drained.count, 33) // Next write of 4800 bytes starts at 158400. 158400 + 4800 = 163200 > 160000. @@ -117,9 +126,9 @@ final class AudioBufferTests: XCTestCase { XCTAssertEqual(wrappingData.count, chunkSize) appendData(wrappingData, to: buffer) - let packets = buffer.processChunks() - XCTAssertEqual(packets.count, 1) - XCTAssertEqual(packets[0].data, wrappingData) + let chunks = collectChunks(from: buffer) + XCTAssertEqual(chunks.count, 1) + XCTAssertEqual(chunks[0], wrappingData) } func testWrapAroundRead() { @@ -133,7 +142,7 @@ final class AudioBufferTests: XCTestCase { for _ in 0..<33 { appendData(makeData(byte: 0x00, count: chunkSize), to: buffer) } - _ = buffer.processChunks() + _ = collectChunks(from: buffer) // Write one chunk starting at 158400. The write itself wraps (tested above), // but crucially the READ will also wrap: readIndex = 158400, @@ -145,9 +154,9 @@ final class AudioBufferTests: XCTestCase { crossBoundaryData.append(makeData(byte: 0xDD, count: 3200)) appendData(crossBoundaryData, to: buffer) - let packets = buffer.processChunks() - XCTAssertEqual(packets.count, 1) - XCTAssertEqual(packets[0].data, crossBoundaryData) + let chunks = collectChunks(from: buffer) + XCTAssertEqual(chunks.count, 1) + XCTAssertEqual(chunks[0], crossBoundaryData) } // MARK: - Overflow guard @@ -164,13 +173,13 @@ final class AudioBufferTests: XCTestCase { appendData(makeData(byte: 0x02, count: 100), to: buffer) // Drain and verify we only got the original data - let packets = buffer.processChunks() - let totalBytes = packets.reduce(0) { $0 + $1.data.count } + let chunks = collectChunks(from: buffer) + let totalBytes = chunks.reduce(0) { $0 + $1.count } XCTAssertEqual(totalBytes, maxBuffer) // Every byte should be 0x01, not 0x02 - for packet in packets { - XCTAssertTrue(packet.data.allSatisfy { $0 == 0x01 }) + for chunk in chunks { + XCTAssertTrue(chunk.allSatisfy { $0 == 0x01 }) } } @@ -187,26 +196,24 @@ final class AudioBufferTests: XCTestCase { appendData(makeData(byte: UInt8(i), count: callbackSize), to: buffer) } - let packets = buffer.processChunks() - XCTAssertEqual(packets.count, 1) - XCTAssertEqual(packets[0].data.count, chunkSize) + let chunks = collectChunks(from: buffer) + XCTAssertEqual(chunks.count, 1) + XCTAssertEqual(chunks[0].count, chunkSize) // Verify the data is in the correct order for i in 0..<10 { - let slice = packets[0].data.subdata(in: (i * callbackSize)..<((i + 1) * callbackSize)) + let slice = chunks[0].subdata(in: (i * callbackSize)..<((i + 1) * callbackSize)) XCTAssertTrue(slice.allSatisfy { $0 == UInt8(i) }) } } - // MARK: - Packet metadata + // MARK: - Chunk size - func testChunkDurationIsCorrect() { + func testBytesPerChunkIsCorrect() { let format = makeFormat() let buffer = AudioBuffer(format: format, chunkDuration: 0.1) - appendData(makeData(byte: 0x00, count: 3200), to: buffer) - let packets = buffer.processChunks() - - XCTAssertEqual(packets[0].duration, 0.1, accuracy: 0.001) + // 16kHz * 0.1s * 2 bytes/frame = 3200 + XCTAssertEqual(buffer.bytesPerChunk, 3200) } } diff --git a/Tests/AudioTeeCoreTests/AudioPacketTests.swift b/Tests/AudioTeeCoreTests/AudioPacketTests.swift deleted file mode 100644 index e93b388..0000000 --- a/Tests/AudioTeeCoreTests/AudioPacketTests.swift +++ /dev/null @@ -1,31 +0,0 @@ -import XCTest - -@testable import AudioTeeCore - -final class AudioPacketTests: XCTestCase { - func testPacketCreation() { - let timestamp = Date() - let duration = 1.0 - let data = Data([0x01, 0x02, 0x03, 0x04]) - - let packet = AudioPacket( - timestamp: timestamp, - duration: duration, - data: data - ) - - XCTAssertEqual(packet.timestamp, timestamp) - XCTAssertEqual(packet.duration, duration) - XCTAssertEqual(packet.data, data) - } - - func testPacketDataSize() { - let packet = AudioPacket( - timestamp: Date(), - duration: 0.5, - data: Data(repeating: 0xFF, count: 1024) - ) - - XCTAssertEqual(packet.data.count, 1024) - } -}