Fixes and cleanup

This commit is contained in:
Dimitris C
2020-10-20 23:41:25 +01:00
parent 48007f7a07
commit 32ea41c4d7
20 changed files with 233 additions and 245 deletions
@@ -50,7 +50,7 @@ enum AudioContent: Int, CaseIterable {
return URL(string: "https://media-ssl.musicradio.com/RadioXLondon")!
case .khruangbin:
// return URL(string: "https://p.scdn.co/mp3-preview/cab4b09c23ffc11774d879977131df9d150fcef4?cid=d8a5ed958d274c2e8ee717e6a4b0971d")!
return URL(string: "https://t4.bcbits.com/stream/8fce516db797c2496c024ef8560f1f29/mp3-128/3790863592?p=0&ts=1603183493&t=facdb0dc4754388596665b27aabf33da88e23882&token=1603183493_06268f83a569966a14237624a3dc1b1ca9dfb5e0")!
return URL(string: "https://t4.bcbits.com/stream/8fce516db797c2496c024ef8560f1f29/mp3-128/3790863592?p=0&ts=1603280742&t=9437510cf44c5621623d38959c409d059e2b99cc&token=1603280742_5e747a280086bc47077abeba8f2eea17694882e6")!
case .flac:
return URL(string: "http://www.lindberg.no/hires/test/2L-145_01_stereo_01.cd.flac")!
case .piano:
+4 -4
View File
@@ -40,25 +40,25 @@ final class DispatchTimerSource {
/// - parameter handler: A closure for the event handler
func add(handler: @escaping () -> Void) {
let handler = handler
self.timer.setEventHandler(handler: handler)
timer.setEventHandler(handler: handler)
}
/// Removes the added event handler from the timer.
func removeHandler() {
self.timer.setEventHandler(handler: nil)
timer.setEventHandler(handler: nil)
}
/// Activates the timer, if needed
func activate() {
if state == .activated { return }
state = .activated
self.timer.activate()
timer.activate()
}
/// Suspends the timer, if needed.
func suspend() {
if state == .suspended { return }
state = .suspended
self.timer.suspend()
timer.suspend()
}
}
@@ -11,7 +11,7 @@ extension AVAudioFormat {
///
/// This exposes the `pointee` value of the `UsafePointer<AudioStreamBasicDescription>`
var basicStreamDescription: AudioStreamBasicDescription {
return self.streamDescription.pointee
return streamDescription.pointee
}
}
@@ -5,16 +5,6 @@
import Foundation
extension UnsafeMutablePointer where Pointee == UInt8 {
/// Allocates and performs binding to memory of an `UnsafeMutableRawPointer` to `UnsafeMutablePointer<UInt8>`
static func uint8pointer(of size: Int) -> UnsafeMutablePointer<UInt8> {
let alignment = MemoryLayout<UInt8>.alignment
return UnsafeMutableRawPointer
.allocate(byteCount: size, alignment: alignment)
.bindMemory(to: UInt8.self, capacity: size)
}
}
extension UnsafeMutableRawPointer {
/// Converts an UnsafeMutableRawPointer to the given Object type
func to<Object: AnyObject>(type: Object.Type) -> Object {
@@ -7,7 +7,7 @@ import Foundation
internal final class NetworkDataStream {
typealias StreamResult = Result<StreamResponse, Error>
typealias StreamCompletion = (queue: OperationQueue, event: (_ event: NetworkDataStream.StreamEvent) -> Void)
typealias StreamCompletion = (_ event: NetworkDataStream.StreamEvent) -> Void
struct StreamResponse {
let response: HTTPURLResponse?
@@ -50,9 +50,8 @@ internal final class NetworkDataStream {
}
@discardableResult
func responseStream(on queue: OperationQueue,
completion: @escaping (_ event: NetworkDataStream.StreamEvent) -> Void) -> Self {
self.streamCallback = (queue, completion)
func responseStream(completion: @escaping (_ event: NetworkDataStream.StreamEvent) -> Void) -> Self {
self.streamCallback = completion
return self
}
@@ -74,42 +73,30 @@ internal final class NetworkDataStream {
internal func didReceive(response: HTTPURLResponse?) {
underlyingQueue.async { [weak self] in
guard let self = self else { return }
guard let stream = self.streamCallback else { return }
stream.queue.addOperation {
stream.event(.response(response))
}
guard let streamCallback = self.streamCallback else { return }
streamCallback(.response(response))
}
}
internal func didReceive(data: Data, response: HTTPURLResponse?) {
underlyingQueue.async { [weak self] in
guard let self = self else { return }
guard let stream = self.streamCallback else { return }
let operation = BlockOperation {
let streamResponse = StreamResponse(response: response, data: data)
stream.event(.stream(.success(streamResponse)))
}
stream.queue.addOperation(operation)
guard let streamCallback = self.streamCallback else { return }
let streamResponse = StreamResponse(response: response, data: data)
streamCallback(.stream(.success(streamResponse)))
}
}
internal func didComplete(with error: Error?) {
internal func didComplete(with error: Error?, response: HTTPURLResponse?) {
underlyingQueue.async { [weak self] in
guard let self = self else { return }
guard let stream = self.streamCallback else { return }
let operation = BlockOperation {
if let error = error {
stream.event(.stream(.failure(error)))
} else {
let completion = Completion(response: self.task?.response as? HTTPURLResponse,
error: error)
stream.event(.complete(completion))
}
if let error = error {
stream(.stream(.failure(error)))
} else {
let completion = Completion(response: response, error: error)
stream(.complete(completion))
}
if let lastOp = stream.queue.operations.last {
operation.addDependency(lastOp)
}
stream.queue.addOperation(operation)
}
}
@@ -26,7 +26,7 @@ internal final class NetworkSessionDelegate: NSObject, URLSessionDataDelegate {
internal func urlSession(_ session: URLSession, task: URLSessionTask, didCompleteWithError error: Error?) {
if let stream = self.stream(for: task) {
stream.didComplete(with: error)
stream.didComplete(with: error, response: task.response as? HTTPURLResponse)
}
}
@@ -10,6 +10,21 @@ enum DataStreamError: Error {
case sessionDeinit
}
public enum NetworkError: Error, Equatable {
case failure(Error)
case serverError
public static func == (lhs: NetworkError, rhs: NetworkError) -> Bool {
switch (lhs, rhs) {
case (.failure, failure):
return true
case (.serverError, .serverError):
return true
default:
return false
}
}
}
protocol StreamTaskProvider: class {
func dataStream(for request: URLSessionTask) -> NetworkDataStream?
}
+1 -11
View File
@@ -23,7 +23,7 @@ extension Lock {
}
/// A wrapper for `os_unfair_lock`
final public class UnfairLock {
final public class UnfairLock: Lock {
private let unfairLock: os_unfair_lock_t
public init() {
@@ -45,16 +45,6 @@ final public class UnfairLock {
}
}
extension UnfairLock: Lock { }
func setLock(_ lock: os_unfair_lock_t) {
os_unfair_lock_lock(lock)
}
func setUnlock(_ lock: os_unfair_lock_t) {
os_unfair_lock_unlock(lock)
}
@propertyWrapper
final class Protected<Value> {
private let lock = UnfairLock()
@@ -20,8 +20,8 @@ public final class AudioPlayer {
/// Defaults to 1.0. Valid ranges are 0.0 to 1.0
/// The value is restricted from 0.0 to 1.0
public var volume: Float32 {
get { self.audioEngine.mainMixerNode.outputVolume }
set { self.audioEngine.mainMixerNode.outputVolume = min(1.0, max(0.0, newValue)) }
get { audioEngine.mainMixerNode.outputVolume }
set { audioEngine.mainMixerNode.outputVolume = min(1.0, max(0.0, newValue)) }
}
/// The playback rate of the player.
///
@@ -30,8 +30,8 @@ public final class AudioPlayer {
/// **NOTE:** Setting this to a value of more than `1.0` while playing a live broadcast stream would
/// result in the audio being exhausted before it could fetch new data.
public var rate: Float {
get { self.rateNode.rate }
set { self.rateNode.rate = newValue }
get { rateNode.rate }
set { rateNode.rate = newValue }
}
/// The player's current state.
@@ -53,7 +53,7 @@ public final class AudioPlayer {
}()
/// Keeps track of the player's state before being paused.
private var stateBeforePaused: PlayerInternalState = .initial
private var stateBeforePaused: InternalState = .initial
/// The underlying `AVAudioEngine` object
let audioEngine = AVAudioEngine()
@@ -128,13 +128,13 @@ public final class AudioPlayer {
/// - parameter url: A `URL` specifying the audio context to be played.
/// - parameter headers: A `Dictionary` specifying any additional headers to be pass to the network request.
public func play(url: URL, headers: [String: String]) {
let audioSource = RemoteAudioSource(networking: self.networking,
let audioSource = RemoteAudioSource(networking: networking,
url: url,
underlyingQueue: sourceQueue,
httpHeaders: headers)
let entry = AudioEntry(source: audioSource,
entryId: AudioEntryId(id: url.absoluteString))
audioSource.delegate = self
entry.delegate = self
clearQueue()
entriesQueue.enqueue(item: entry, type: .upcoming)
playerContext.internalState = .pendingNext
@@ -160,8 +160,8 @@ public final class AudioPlayer {
stopReadProccessFromSource()
sourceQueue.async { [weak self] in
guard let self = self else { return }
self.playerContext.audioReadingEntry?.source.delegate = nil
self.playerContext.audioReadingEntry?.source.close()
self.playerContext.audioReadingEntry?.delegate = nil
self.playerContext.audioReadingEntry?.close()
if let playingEntry = self.playerContext.audioPlayingEntry {
self.processFinishPlaying(entry: playingEntry, with: nil)
}
@@ -183,7 +183,7 @@ public final class AudioPlayer {
pauseEngine()
stopReadProccessFromSource()
playerContext.audioPlayingEntry?.source.suspend()
playerContext.audioPlayingEntry?.suspend()
sourceQueue.async { [weak self] in
self?.processSource()
}
@@ -200,7 +200,7 @@ public final class AudioPlayer {
Logger.debug("resuming audio engine failed: %@", category: .generic, args: error.localizedDescription)
}
playerContext.audioPlayingEntry?.source.resume()
playerContext.audioPlayingEntry?.resume()
startPlayer(resetBuffers: false)
startReadProcessFromSourceIfNeeded()
}
@@ -284,12 +284,12 @@ public final class AudioPlayer {
/// Attaches callbacks to the `playerContext` and `renderProcessor`.
private func configPlayerContext() {
self.playerContext.stateChanged = { [weak self] oldValue, newValue in
playerContext.stateChanged = { [weak self] oldValue, newValue in
guard let self = self else { return }
self.delegate?.audioPlayerStateChanged(player: self, with: newValue, previous: oldValue)
}
self.playerRenderProcessor.audioFinished = { [weak self] entry in
playerRenderProcessor.audioFinished = { [weak self] entry in
guard let self = self else { return }
self.sourceQueue.async {
let nextEntry = self.entriesQueue.dequeue(type: .buffering)
@@ -426,15 +426,15 @@ public final class AudioPlayer {
fileStreamProcessor.closeFileStreamIfNeeded()
if let readingEntry = playerContext.audioReadingEntry {
readingEntry.source.delegate = nil
readingEntry.source.close()
readingEntry.delegate = nil
readingEntry.close()
}
playerContext.entriesLock.around {
playerContext.audioReadingEntry = entry
}
playerContext.audioReadingEntry?.source.delegate = self
playerContext.audioReadingEntry?.source.seek(at: 0)
playerContext.audioReadingEntry?.delegate = self
playerContext.audioReadingEntry?.seek(at: 0)
if startPlaying {
if shouldClearQueue {
@@ -531,7 +531,8 @@ public final class AudioPlayer {
extension AudioPlayer: AudioStreamSourceDelegate {
func dataAvailable(source: AudioStreamSource, data: Data) {
guard playerContext.audioReadingEntry?.source === source else { return }
guard let readingEntry = playerContext.audioReadingEntry,
readingEntry.has(same: source) else { return }
if !fileStreamProcessor.isFileStreamOpen {
guard fileStreamProcessor.openFileStream(with: source.audioFileHint) == noErr else {
@@ -543,7 +544,8 @@ extension AudioPlayer: AudioStreamSourceDelegate {
// TODO: check for discontinuous stream and add flag
if fileStreamProcessor.isFileStreamOpen {
guard fileStreamProcessor.parseFileStreamBytes(data: data) == noErr else {
if source === playerContext.audioPlayingEntry?.source {
if let playingEntry = playerContext.audioPlayingEntry,
playingEntry.has(same: source) {
raiseUnxpected(error: .streamParseBytesFailure)
}
return
@@ -557,13 +559,14 @@ extension AudioPlayer: AudioStreamSourceDelegate {
}
}
func errorOccured(source: AudioStreamSource) {
guard let entry = playerContext.audioReadingEntry, entry.source === source else { return }
raiseUnxpected(error: .dataNotFound)
func errorOccured(source: AudioStreamSource, error: Error) {
guard let entry = playerContext.audioReadingEntry, entry.has(same: source) else { return }
raiseUnxpected(error: .networkError(.failure(error)))
}
func endOfFileOccured(source: AudioStreamSource) {
if playerContext.audioReadingEntry == nil || playerContext.audioReadingEntry?.source !== source {
let hasSameSource = playerContext.audioReadingEntry?.has(same: source) ?? false
guard playerContext.audioReadingEntry == nil || hasSameSource else {
source.delegate = nil
source.close()
return
@@ -583,8 +586,8 @@ extension AudioPlayer: AudioStreamSourceDelegate {
readingEntry.framesState.lastFrameQueued = readingEntry.framesState.queued
readingEntry.source.delegate = nil
readingEntry.source.close()
readingEntry.delegate = nil
readingEntry.close()
playerContext.entriesLock.lock()
playerContext.audioReadingEntry = nil
@@ -25,9 +25,9 @@ internal final class AudioPlayerContext {
/// - NOTE: Do not use directly instead use the `internalState` to set and get the property
/// or the `setInternalState(to:when:)`method
@Protected
private var __playerInternalState: PlayerInternalState = .initial
private var __playerInternalState: AudioPlayer.InternalState = .initial
var internalState: PlayerInternalState {
var internalState: AudioPlayer.InternalState {
get { __playerInternalState }
set { setInternalState(to: newValue) }
}
@@ -49,8 +49,8 @@ internal final class AudioPlayerContext {
/// - parameter state: The new `PlayerInternalState`
/// - parameter inState: If the `inState` expression is not nil, the internalState will be set if the evaluated expression is `true`
/// - NOTE: This sets the underlying `__playerInternalState` variable
internal func setInternalState(to state: PlayerInternalState,
when inState: ((PlayerInternalState) -> Bool)? = nil) {
internal func setInternalState(to state: AudioPlayer.InternalState,
when inState: ((AudioPlayer.InternalState) -> Bool)? = nil) {
let newValues = playerStateAndStopReason(for: state)
$stopReason.write { reason in
reason = newValues.stopReason
@@ -6,30 +6,32 @@
import Foundation
// MARK: Internal State
internal struct PlayerInternalState: OptionSet {
var rawValue: Int
static let initial = PlayerInternalState([])
static let running = PlayerInternalState(rawValue: 1)
static let playing = PlayerInternalState(rawValue: 1 << 1 | PlayerInternalState.running.rawValue)
static let rebuffering = PlayerInternalState(rawValue: 1 << 2 | PlayerInternalState.running.rawValue)
static let startingThread = PlayerInternalState(rawValue: 1 << 3 | PlayerInternalState.running.rawValue)
static let waitingForData = PlayerInternalState(rawValue: 1 << 4 | PlayerInternalState.running.rawValue)
static let waitingForDataAfterSeek = PlayerInternalState(rawValue: 1 << 5 | PlayerInternalState.running.rawValue)
static let paused = PlayerInternalState(rawValue: 1 << 6 | PlayerInternalState.running.rawValue)
static let stopped = PlayerInternalState(rawValue: 1 << 9)
static let pendingNext = PlayerInternalState(rawValue: 1 << 10)
static let disposed = PlayerInternalState(rawValue: 1 << 30)
static let error = PlayerInternalState(rawValue: 1 << 31)
static let isPlaying: PlayerInternalState =
[.running, .startingThread, .playing, .waitingForDataAfterSeek]
static let isBuffering: PlayerInternalState =
[.pendingNext, .rebuffering, .waitingForData]
extension AudioPlayer {
internal struct InternalState: OptionSet {
var rawValue: Int
static let initial = InternalState([])
static let running = InternalState(rawValue: 1)
static let playing = InternalState(rawValue: 1 << 1 | InternalState.running.rawValue)
static let rebuffering = InternalState(rawValue: 1 << 2 | InternalState.running.rawValue)
static let startingThread = InternalState(rawValue: 1 << 3 | InternalState.running.rawValue)
static let waitingForData = InternalState(rawValue: 1 << 4 | InternalState.running.rawValue)
static let waitingForDataAfterSeek = InternalState(rawValue: 1 << 5 | InternalState.running.rawValue)
static let paused = InternalState(rawValue: 1 << 6 | InternalState.running.rawValue)
static let stopped = InternalState(rawValue: 1 << 9)
static let pendingNext = InternalState(rawValue: 1 << 10)
static let disposed = InternalState(rawValue: 1 << 30)
static let error = InternalState(rawValue: 1 << 31)
static let isPlaying: InternalState =
[.running, .startingThread, .playing, .waitingForDataAfterSeek]
static let isBuffering: InternalState =
[.pendingNext, .rebuffering, .waitingForData]
}
}
func playerStateAndStopReason(for internalState: PlayerInternalState) -> (state: AudioPlayerState, stopReason: AudioPlayerStopReason) {
func playerStateAndStopReason(for internalState: AudioPlayer.InternalState) -> (state: AudioPlayerState, stopReason: AudioPlayerStopReason) {
var playerNewState: AudioPlayerState
var stopReason: AudioPlayerStopReason = .none
@@ -89,6 +91,7 @@ public enum AudioPlayerError: LocalizedError, Equatable {
case audioSystemError(AudioSystemError)
case codecError
case dataNotFound
case networkError(NetworkError)
case other
public var errorDescription: String? {
@@ -101,6 +104,8 @@ public enum AudioPlayerError: LocalizedError, Equatable {
return "Codec error while parsing data packets"
case .dataNotFound:
return "No data supplied from network stream"
case .networkError(let error):
return error.localizedDescription
case .other:
return "Audio Player error"
}
@@ -104,7 +104,7 @@ final class AudioFileStreamProcessor {
self.inputFormat = inputFormat
// magic cookie info
let fileHint = playerContext.audioReadingEntry?.source.audioFileHint
let fileHint = playerContext.audioReadingEntry?.audioFileHint
if let fileStream = audioFileStream, fileHint != kAudioFileAAC_ADTSType {
var cookieSize: UInt32 = 0
guard AudioFileStreamGetPropertyInfo(fileStream, kAudioFileStreamProperty_MagicCookieData, &cookieSize, nil) == noErr else {
@@ -163,15 +163,15 @@ final class AudioFileStreamProcessor {
private func processDataOffset(fileStream: AudioFileStreamID) {
var offset: UInt64 = 0
fileStreamGetProperty(value: &offset, fileStream: fileStream, propertyId: kAudioFileStreamProperty_DataOffset)
playerContext.audioReadingEntry?.parsedHeader = true
playerContext.audioReadingEntry?.audioDataOffset = offset
playerContext.audioReadingEntry?.audioStreamState.processedDataFormat = true
playerContext.audioReadingEntry?.audioStreamState.dataOffset = offset
}
private func processReadyToProducePackets(fileStream: AudioFileStreamID) {
var packetCount: UInt64 = 0
var packetCountSize = UInt32(MemoryLayout.size(ofValue: packetCount))
AudioFileStreamGetProperty(fileStream, kAudioFileStreamProperty_AudioDataPacketCount, &packetCountSize, &packetCount)
playerContext.audioPlayingEntry?.packetCount = Double(packetCount)
playerContext.audioPlayingEntry?.audioStreamState.dataPacketCount = Double(packetCount)
}
private func processFileFormat(fileStream: AudioFileStreamID) {
@@ -186,7 +186,7 @@ final class AudioFileStreamProcessor {
private func processDataFormat(fileStream: AudioFileStreamID) {
var audioStreamFormat = AudioStreamBasicDescription()
guard let entry = playerContext.audioReadingEntry else { return }
if !entry.parsedHeader {
if !entry.audioStreamState.processedDataFormat {
fileStreamGetProperty(value: &audioStreamFormat, fileStream: fileStream, propertyId: kAudioFileStreamProperty_DataFormat)
var packetBufferSize: UInt32 = 0
var status = fileStreamGetProperty(value: &packetBufferSize, fileStream: fileStream, propertyId: kAudioFileStreamProperty_PacketSizeUpperBound)
@@ -202,7 +202,7 @@ final class AudioFileStreamProcessor {
}
playerContext.entriesLock.unlock()
playerContext.audioReadingEntry?.lock.around {
playerContext.audioPlayingEntry?.processedPacketsState.buferSize = packetBufferSize
playerContext.audioReadingEntry?.processedPacketsState.bufferSize = packetBufferSize
}
}
if let readingEntry = playerContext.audioReadingEntry {
@@ -214,7 +214,7 @@ final class AudioFileStreamProcessor {
guard let entry = playerContext.audioReadingEntry else { return }
var audioDataByteCount: UInt64 = 0
fileStreamGetProperty(value: &audioDataByteCount, fileStream: fileStream, propertyId: kAudioFileStreamProperty_AudioDataByteCount)
entry.audioDataByteCount = audioDataByteCount
entry.audioStreamState.dataByteCount = audioDataByteCount
}
@@ -222,7 +222,7 @@ final class AudioFileStreamProcessor {
guard let entry = playerContext.audioReadingEntry else { return }
var audioDataPacketCount: UInt64 = 0
fileStreamGetProperty(value: &audioDataPacketCount, fileStream: fileStream, propertyId: kAudioFileStreamProperty_AudioDataPacketCount)
entry.audioDataPacketOffset = audioDataPacketCount
entry.audioStreamState.dataPacketOffset = audioDataPacketCount
}
private func processFormatList(fileStream: AudioFileStreamID) {
@@ -251,7 +251,7 @@ final class AudioFileStreamProcessor {
inInputData: UnsafeRawPointer,
inPacketDescriptions: UnsafeMutablePointer<AudioStreamPacketDescription>?) {
guard let entry = playerContext.audioReadingEntry,
entry.parsedHeader && !playerContext.disposedRequested else { return }
entry.audioStreamState.processedDataFormat && !playerContext.disposedRequested else { return }
if let playingEntry = playerContext.audioPlayingEntry,
rendererContext.seekRequest.requested && playingEntry.calculatedBitrate() > 0 {
@@ -473,6 +473,7 @@ private func _converterCallback(inAudioConverter: AudioConverterRef,
inUserData: UnsafeMutableRawPointer?) -> OSStatus {
guard let convertInfo = inUserData?.assumingMemoryBound(to: AudioConvertInfo.self) else { return 0 }
// we need to tell the converter to stop converting after it should stop converting
if convertInfo.pointee.done {
ioNumberDataPackets.pointee = 0
return AudioConvertStatus.done.rawValue
@@ -73,10 +73,9 @@ final class MetadataStreamProcessor: MetadataStreamSource {
@inline(__always)
func proccessMetadata(data: Data) -> Data {
var audioData = Data()
data.withUnsafeBytes { buffer in
guard buffer.count > 0 else { return }
data.withUnsafeBytes { buffer -> Data in
guard buffer.count > 0 else { return data }
var audioData = Data()
var bytesRead = 0
let bytes = buffer.baseAddress!.assumingMemoryBound(to: UInt8.self)
while bytesRead < buffer.count {
@@ -113,9 +112,8 @@ final class MetadataStreamProcessor: MetadataStreamSource {
bytesRead += audioBytesToRead
}
}
return audioData
}
return audioData
}
}
@@ -5,7 +5,6 @@
import AVFoundation
private let outputChannels: UInt32 = 2
struct UnitDescriptions {
@@ -21,48 +21,59 @@ final class EntryFramesState {
}
final class ProcessedPacketsState {
var buferSize: UInt32 = 0
var bufferSize: UInt32 = 0
var count: UInt32 = 0
var sizeTotal: UInt32 = 0
}
final class AudioStreamState {
var processedDataFormat: Bool = false
var dataOffset: UInt64 = 0
var dataByteCount: UInt64? = nil
var dataPacketOffset: UInt64? = nil
var dataPacketCount: Double = 0
var streamFormat = AudioStreamBasicDescription()
}
public class AudioEntry {
private let estimationMinPackets = 2
private let estimationMinPacketsPreferred = 64
let lock = UnfairLock()
let source: AudioStreamSource
weak var delegate: AudioStreamSourceDelegate?
let id: AudioEntryId
var seekTime: Float
var parsedHeader: Bool = false
var packetCount: Double = 0
var packetDuration: Double {
return Double(audioStreamFormat.mFramesPerPacket) / Double(sampleRate)
}
/// The sample rate from the `audioStreamFormat`
var sampleRate: Float {
Float(audioStreamFormat.mSampleRate)
}
var framesState: EntryFramesState
var processedPacketsState: ProcessedPacketsState
var audioDataOffset: UInt64 = 0
var audioDataByteCount: UInt64?
var audioDataPacketOffset: UInt64?
var audioFileHint: AudioFileTypeID {
source.audioFileHint
}
var audioStreamFormat = AudioStreamBasicDescription()
private(set) var audioStreamState: AudioStreamState
private(set) var framesState: EntryFramesState
private(set) var processedPacketsState: ProcessedPacketsState
private var packetDuration: Double {
return Double(audioStreamFormat.mFramesPerPacket) / Double(sampleRate)
}
private var avaragePacketByteSize: Double {
let packets = processedPacketsState
guard packets.count > 0 else { return 0 }
return Double(packets.sizeTotal / packets.count)
}
private let source: AudioStreamSource
init(source: AudioStreamSource, entryId: AudioEntryId) {
self.source = source
self.id = entryId
@@ -71,11 +82,34 @@ public class AudioEntry {
self.processedPacketsState = ProcessedPacketsState()
self.framesState = EntryFramesState()
self.audioStreamState = AudioStreamState()
}
func close() {
source.delegate = nil
source.close()
}
func suspend() {
source.suspend()
}
func resume() {
source.resume()
}
func seek(at offset: Int) {
source.delegate = self
source.seek(at: offset)
}
func reset() {
lock.lock(); defer { lock.unlock() }
self.framesState = EntryFramesState()
framesState = EntryFramesState()
}
func has(same source: AudioStreamSource) -> Bool {
source === self.source
}
func calculatedBitrate() -> Double {
@@ -98,7 +132,7 @@ public class AudioEntry {
func duration() -> Double {
guard sampleRate > 0 else { return 0 }
if let audioDataPacketOffset = audioDataPacketOffset {
if let audioDataPacketOffset = audioStreamState.dataPacketOffset {
let franesPerPacket = UInt64(audioStreamFormat.mFramesPerPacket)
if audioDataPacketOffset > 0 && franesPerPacket > 0 {
return Double(audioDataPacketOffset * franesPerPacket) / audioStreamFormat.mSampleRate
@@ -113,15 +147,33 @@ public class AudioEntry {
}
private func audioDataLengthBytes() -> UInt {
if let byteCount = audioDataByteCount {
if let byteCount = audioStreamState.dataByteCount {
return UInt(byteCount)
}
guard source.length > 0 else { return 0 }
return UInt(source.length) - UInt(audioDataOffset)
return UInt(source.length) - UInt(audioStreamState.dataOffset)
}
}
extension AudioEntry: AudioStreamSourceDelegate {
func dataAvailable(source: AudioStreamSource, data: Data) {
delegate?.dataAvailable(source: source, data: data)
}
func errorOccured(source: AudioStreamSource, error: Error) {
delegate?.errorOccured(source: source, error: error)
}
func endOfFileOccured(source: AudioStreamSource) {
delegate?.endOfFileOccured(source: source)
}
func metadataReceived(data: [String : String]) {
delegate?.metadataReceived(data: data)
}
}
extension AudioEntry: Equatable {
public static func == (lhs: AudioEntry, rhs: AudioEntry) -> Bool {
lhs.id == rhs.id
@@ -6,18 +6,18 @@
import Foundation
import AudioToolbox
protocol AudioStreamSourceDelegate: class {
protocol AudioStreamSourceDelegate: AnyObject {
/// Indicates that there's data available
func dataAvailable(source: AudioStreamSource, data: Data)
/// Indicates an error occurred
func errorOccured(source: AudioStreamSource)
func errorOccured(source: AudioStreamSource, error: Error)
/// Indicates end of file has occurred
func endOfFileOccured(source: AudioStreamSource)
/// Indicates metadata read from stream
func metadataReceived(data: [String: String])
}
protocol CoreAudioStreamSource: class {
protocol CoreAudioStreamSource: AnyObject {
/// An `Int` that represents the position of the audio
var position: Int { get }
/// The length of the audio in bytes
@@ -21,19 +21,16 @@ public class RemoteAudioSource: AudioStreamSource {
private let url: URL
private let networking: NetworkingClient
internal var metadataStreamProccessor: MetadataStreamSource
private var streamRequest: NetworkDataStream?
private var additionalRequestHeaders: [String: String]
private var httpHeaderParsed: Bool
private var httpResponse: HTTPURLResponse? {
streamRequest?.urlResponse
}
private var parsedHeaderOutput: HTTPHeaderParserOutput?
private var relativePosition: Int
private var seekOffset: Int
internal var metadataStreamProccessor: MetadataStreamSource
internal var audioFileHint: AudioFileTypeID {
if let output = parsedHeaderOutput {
return output.typeId
@@ -53,7 +50,6 @@ public class RemoteAudioSource: AudioStreamSource {
self.metadataStreamProccessor = metadataStreamSource
self.url = url
self.additionalRequestHeaders = httpHeaders
self.httpHeaderParsed = false
self.relativePosition = 0
self.seekOffset = 0
self.underlyingQueue = underlyingQueue
@@ -100,7 +96,7 @@ public class RemoteAudioSource: AudioStreamSource {
relativePosition = 0
seekOffset = offset
if let supportsSeek = self.parsedHeaderOutput?.supportsSeek,
if let supportsSeek = parsedHeaderOutput?.supportsSeek,
!supportsSeek && offset != relativePosition {
return
}
@@ -123,7 +119,7 @@ public class RemoteAudioSource: AudioStreamSource {
let urlRequest = buildUrlRequest(with: url, seekIfNeeded: seekOffset)
streamRequest = networking.stream(request: urlRequest)
.responseStream(on: networkStreamQueue) { [weak self] event in
.responseStream { [weak self] event in
guard let self = self else { return }
self.handleResponse(event: event)
}
@@ -131,32 +127,39 @@ public class RemoteAudioSource: AudioStreamSource {
metadataStreamProccessor.delegate = self
}
// MARK: - Network Handle Methods
private func handleResponse(event: NetworkDataStream.StreamEvent) {
switch event {
case .response(let urlResponse):
self.parseResponseHeader(response: urlResponse)
parseResponseHeader(response: urlResponse)
case .stream(let event):
self.handleStreamEvent(event: event)
addStreamOperation { [weak self] in
self?.handleStreamEvent(event: event)
}
case .complete:
self.delegate?.endOfFileOccured(source: self)
addCompletionOperation { [weak self] in
guard let self = self else { return }
self.delegate?.endOfFileOccured(source: self)
}
}
}
private func handleStreamEvent(event: NetworkDataStream.StreamResult) {
switch event {
case .success(let responseValue):
if let data = responseValue.data {
case .success(let value):
if let data = value.data {
if metadataStreamProccessor.canProccessMetadata {
let extractedAudioData = metadataStreamProccessor.proccessMetadata(data: data)
self.delegate?.dataAvailable(source: self, data: extractedAudioData)
delegate?.dataAvailable(source: self, data: extractedAudioData)
} else {
self.delegate?.dataAvailable(source: self, data: data)
delegate?.dataAvailable(source: self, data: data)
}
relativePosition += data.count
}
case .failure(let error):
print(error)
self.delegate?.errorOccured(source: self)
delegate?.errorOccured(source: self, error: error)
break
}
}
@@ -178,13 +181,13 @@ public class RemoteAudioSource: AudioStreamSource {
delegate?.endOfFileOccured(source: self)
}
else if httpStatusCode >= 300 {
delegate?.errorOccured(source: self)
delegate?.errorOccured(source: self, error: NetworkError.serverError)
}
}
private func buildUrlRequest(with url: URL, seekIfNeeded seekOffset: Int) -> URLRequest {
var urlRequest = URLRequest(url: self.url)
var urlRequest = URLRequest(url: url)
urlRequest.networkServiceType = .avStreaming
urlRequest.cachePolicy = .reloadIgnoringLocalCacheData
@@ -201,11 +204,28 @@ public class RemoteAudioSource: AudioStreamSource {
return urlRequest
}
// MARK: - Network Stream Operation Queue
private func addStreamOperation(_ block: @escaping () -> Void) {
let operation = BlockOperation(block: block)
networkStreamQueue.addOperation(operation)
}
private func addCompletionOperation(_ block: @escaping () -> Void) {
let operation = BlockOperation(block: block)
operation.qualityOfService = .background
operation.queuePriority = .veryLow
if let lastOperation = networkStreamQueue.operations.last {
operation.addDependency(lastOperation)
}
networkStreamQueue.addOperation(operation)
}
}
extension RemoteAudioSource: MetadataStreamSourceDelegate {
func didReceiveMetadata(metadata: Result<[String : String], MetadataParsingError>) {
guard case let .success(data) = metadata else { return }
self.delegate?.metadataReceived(data: data)
delegate?.metadataReceived(data: data)
}
}
@@ -72,52 +72,4 @@ struct HTTPHeaderParser: Parser {
metadataStep: metadataStep)
}
}
struct CFHTTPResponseParser: Parser {
typealias Input = CFHTTPMessage
typealias Output = HTTPHeaderParserOutput
func parse(input: CFHTTPMessage) -> HTTPHeaderParserOutput {
let headers = CFHTTPMessageCopyAllHeaderFields(input)?.takeRetainedValue() as? [String: Any]
let statusCode = CFHTTPMessageGetResponseStatusCode(input)
var supportsSeek = false
if let acceptRanges = headers?[HeaderField.acceptRanges] as? String, acceptRanges != "none" {
supportsSeek = true
}
var typeId: UInt32 = 0
if let contentType = headers?["Content-Type"] as? String {
typeId = audioFileType(mimeType: contentType)
}
var fileLength: Int = 0
if statusCode == 200 {
if let contentLength = headers?[HeaderField.contentLength] as? String,
let length = Int(contentLength) {
fileLength = length
}
} else if statusCode == 206 {
if let contentLength = headers?[HeaderField.contentRange] as? String {
let components = contentLength.components(separatedBy: "/")
if components.count == 2 {
if let last = components.last, let length = Int(last) {
fileLength = length
}
}
}
}
var metadataStep = 0
if let icyMetaint = headers?[IcyHeaderField.icyMentaint] as? String,
let intValue = Int(icyMetaint) {
metadataStep = intValue
}
return HTTPHeaderParserOutput(supportsSeek: supportsSeek,
fileLength: fileLength,
typeId: typeId,
metadataStep: metadataStep)
}
}
@@ -53,6 +53,8 @@ class NetworkingClientTests: XCTestCase {
case .complete(let completion):
responseCompletion = completion
expectation.fulfill()
case .response:
break
}
}
.resume()
@@ -64,35 +66,6 @@ class NetworkingClientTests: XCTestCase {
XCTAssertNotNil(receivedData)
}
func testThatStreamCanProduceAnInputStream() {
let expect = expectation(description: "stream complete")
let networking = NetworkingClient()
let url = URL(string: "https://httpbin.org/xml")!
var request = URLRequest(url: url)
request.addValue("application/xml", forHTTPHeaderField: "Content-Type")
let inputStream = networking
.stream(request: request)
.responseStream { event in
switch event {
case .complete:
expect.fulfill()
default: break
}
}
.asInputStream()
wait(for: [expect], timeout: 10)
let xmlParser = XMLParser(stream: inputStream!)
let xmlParsed = xmlParser.parse()
XCTAssertTrue(xmlParsed)
XCTAssertNil(xmlParser.parserError)
}
func testThatStreamCanBeCalledAndCompleteAtAGivenThread() {
let networking = NetworkingClient()
@@ -103,11 +76,9 @@ class NetworkingClientTests: XCTestCase {
var responseCompletion: NetworkDataStream.Completion?
var receivedData: Data?
let receivedDataQueue = DispatchQueue(label: "received.data.queue")
networking.stream(request: request)
.responseStream(on: receivedDataQueue) { event in
.responseStream { event in
switch event {
case .stream(let result):
XCTAssertFalse(Thread.current.isMainThread)
@@ -120,6 +91,8 @@ class NetworkingClientTests: XCTestCase {
XCTAssertFalse(Thread.current.isMainThread)
responseCompletion = completion
expectation.fulfill()
case .response:
XCTAssertFalse(Thread.current.isMainThread)
}
}
.resume()
@@ -140,7 +140,10 @@ class PlayerQueueEntriesTest: XCTestCase {
private let networkingClient = NetworkingClient(configuration: .ephemeral)
private func audioEntry(id: String) -> AudioEntry {
let source =
RemoteAudioSource(networking: networkingClient, url: URL(string: "www.a-url.com")!, sourceQueue: .main, readBufferSize: 1024)
RemoteAudioSource(networking: networkingClient,
url: URL(string: "www.a-url.com")!,
underlyingQueue: DispatchQueue(label: "some-queue"),
httpHeaders: [:])
return AudioEntry(source: source, entryId: AudioEntryId(id: id))
}