From cece65475dca5001dcfd15c5a4b7faab3e45fe33 Mon Sep 17 00:00:00 2001 From: Dimitris C Date: Sun, 1 Nov 2020 17:19:42 +0000 Subject: [PATCH] Adds initial support for seeking --- .../xcschemes/AudioExample.xcscheme | 12 ++ .../AudioExample/ViewController.swift | 51 +++++---- AudioStreaming.xcodeproj/project.pbxproj | 4 + .../xcschemes/AudioStreaming.xcscheme | 17 +++ AudioStreaming/Core/Helpers/measure.swift | 15 +++ .../Core/Network/NetworkDataStream.swift | 59 ++++++++-- .../Core/Network/NetworkingClient.swift | 29 +++-- AudioStreaming/Core/Structures/BiMap.swift | 4 + .../Streaming/AudioPlayer/AudioPlayer.swift | 105 ++++++++++++++++-- .../AudioPlayer/AudioRendererContext.swift | 7 +- .../Processors/AudioFileStreamProcessor.swift | 88 +++++++++++++-- .../AudioPlayerRenderProcessor.swift | 15 ++- .../Streaming/AudioSource/AudioEntry.swift | 11 +- .../AudioSource/Models/SeekRequest.swift | 4 +- .../AudioSource/RemoteAudioSource.swift | 22 ++-- .../Helpers/PlayerQueueEntries.swift | 67 ++++++----- 16 files changed, 392 insertions(+), 118 deletions(-) create mode 100644 AudioStreaming/Core/Helpers/measure.swift diff --git a/AudioExample/AudioExample.xcodeproj/xcshareddata/xcschemes/AudioExample.xcscheme b/AudioExample/AudioExample.xcodeproj/xcshareddata/xcschemes/AudioExample.xcscheme index 52549b0..9ebe11e 100644 --- a/AudioExample/AudioExample.xcodeproj/xcshareddata/xcschemes/AudioExample.xcscheme +++ b/AudioExample/AudioExample.xcodeproj/xcshareddata/xcschemes/AudioExample.xcscheme @@ -85,6 +85,18 @@ isEnabled = "NO"> + + + + + + 0 { + player.seek(to: Double(slider.value)) + } } @objc func sliderValueChanged() { - - print(slider.value) + seekValue = slider.value } @objc @@ -242,7 +248,8 @@ class ViewController: UIViewController { private func startDisplayLink() { displayLink?.invalidate() displayLink = UIScreen.main.displayLink(withTarget: self, selector: #selector(tick)) - displayLink?.add(to: .current, forMode: .default) + displayLink?.preferredFramesPerSecond = 6 + displayLink?.add(to: .current, forMode: .common) } private func stopDisplayLink(resetLabels: Bool) { @@ -270,7 +277,9 @@ class ViewController: UIViewController { slider.minimumValue = 0 slider.maximumValue = Float(duration) - slider.value = Float(progress) + if !isScrubbing { + slider.value = Float(progress) + } elapsedPlayTimeLabel.text = timeFrom(seconds: elapsed) remainingPlayTimeLabel.text = timeFrom(seconds: remaining) @@ -298,28 +307,28 @@ extension ViewController: AudioPlayerDelegate { print("did cancel items") } - func audioPlayerDidStartPlaying(player _: AudioPlayer, with entryId: AudioEntryId) { - print("did start playing: \(entryId)") + func audioPlayerDidStartPlaying(player _: AudioPlayer, with _: AudioEntryId) { +// print("did start playing: \(entryId)") metadataLabel.text = "" } - func audioPlayerDidFinishBuffering(player _: AudioPlayer, with entryId: AudioEntryId) { - print("did finish buffering: \(entryId)") + func audioPlayerDidFinishBuffering(player _: AudioPlayer, with _: AudioEntryId) { +// print("did finish buffering: \(entryId)") } - func audioPlayerStateChanged(player _: AudioPlayer, with newState: AudioPlayerState, previous: AudioPlayerState) { - print("player state changed from: \(previous) to: \(newState)") + func audioPlayerStateChanged(player _: AudioPlayer, with _: AudioPlayerState, previous _: AudioPlayerState) { +// print("player state changed from: \(previous) to: \(newState)") } - func audioPlayerDidFinishPlaying(player _: AudioPlayer, entryId: AudioEntryId, stopReason: AudioPlayerStopReason, progress: Double, duration: Double) { - print("player finished playing for: \(entryId)") - print("===> stop reason: \(stopReason)") - print("===> progress: \(progress)") - print("===> duration: \(duration)") + func audioPlayerDidFinishPlaying(player _: AudioPlayer, entryId _: AudioEntryId, stopReason _: AudioPlayerStopReason, progress _: Double, duration _: Double) { +// print("player finished playing for: \(entryId)") +// print("===> stop reason: \(stopReason)") +// print("===> progress: \(progress)") +// print("===> duration: \(duration)") } - func audioPlayerUnexpectedError(player _: AudioPlayer, error: AudioPlayerError) { - print("player error'd unexpectedly: \(error)") + func audioPlayerUnexpectedError(player _: AudioPlayer, error _: AudioPlayerError) { +// print("player error'd unexpectedly: \(error)") } func audioPlayerDidReadMetadata(player _: AudioPlayer, metadata: [String: String]) { diff --git a/AudioStreaming.xcodeproj/project.pbxproj b/AudioStreaming.xcodeproj/project.pbxproj index dc9ef26..8a6f3eb 100644 --- a/AudioStreaming.xcodeproj/project.pbxproj +++ b/AudioStreaming.xcodeproj/project.pbxproj @@ -32,6 +32,7 @@ B5667A902499018D00D93F85 /* AudioFileStreamProcessor.swift in Sources */ = {isa = PBXBuildFile; fileRef = B5667A8F2499018D00D93F85 /* AudioFileStreamProcessor.swift */; }; B5667A922499063D00D93F85 /* AudioPlayerContext.swift in Sources */ = {isa = PBXBuildFile; fileRef = B5667A912499063D00D93F85 /* AudioPlayerContext.swift */; }; B5667B3E249BC43100D93F85 /* AudioPlayerRenderProcessor.swift in Sources */ = {isa = PBXBuildFile; fileRef = B5667B3D249BC43000D93F85 /* AudioPlayerRenderProcessor.swift */; }; + B5737340254DE43E003DFBEC /* measure.swift in Sources */ = {isa = PBXBuildFile; fileRef = B573733F254DE43E003DFBEC /* measure.swift */; }; B57829CF2548B32B00C78D36 /* Lock.swift in Sources */ = {isa = PBXBuildFile; fileRef = B57829CE2548B32B00C78D36 /* Lock.swift */; }; B58386382544A2C10087A712 /* EntryFrames.swift in Sources */ = {isa = PBXBuildFile; fileRef = B58386372544A2C10087A712 /* EntryFrames.swift */; }; B5838640254584A50087A712 /* ProcessedPackets.swift in Sources */ = {isa = PBXBuildFile; fileRef = B583863F254584A50087A712 /* ProcessedPackets.swift */; }; @@ -111,6 +112,7 @@ B5667A8F2499018D00D93F85 /* AudioFileStreamProcessor.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AudioFileStreamProcessor.swift; sourceTree = ""; }; B5667A912499063D00D93F85 /* AudioPlayerContext.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AudioPlayerContext.swift; sourceTree = ""; }; B5667B3D249BC43000D93F85 /* AudioPlayerRenderProcessor.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AudioPlayerRenderProcessor.swift; sourceTree = ""; }; + B573733F254DE43E003DFBEC /* measure.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = measure.swift; sourceTree = ""; }; B57829CE2548B32B00C78D36 /* Lock.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Lock.swift; sourceTree = ""; }; B58386372544A2C10087A712 /* EntryFrames.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = EntryFrames.swift; sourceTree = ""; }; B583863F254584A50087A712 /* ProcessedPackets.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ProcessedPackets.swift; sourceTree = ""; }; @@ -274,6 +276,7 @@ B592E13025460883008866FB /* Helpers */ = { isa = PBXGroup; children = ( + B573733F254DE43E003DFBEC /* measure.swift */, B514657E248E3884005C03F7 /* DispatchTimerSource.swift */, B57829CE2548B32B00C78D36 /* Lock.swift */, B500731F24D00BAC00BB4475 /* Logger.swift */, @@ -554,6 +557,7 @@ B59DF1A32493E90C0043C498 /* AudioFileStream+Helpers.swift in Sources */, B54D876D2490E4A000C361A0 /* UnitDescriptions.swift in Sources */, B514657F248E3884005C03F7 /* DispatchTimerSource.swift in Sources */, + B5737340254DE43E003DFBEC /* measure.swift in Sources */, B55CEABC24853CD20001C498 /* AudioPlayer.swift in Sources */, B5667B3E249BC43100D93F85 /* AudioPlayerRenderProcessor.swift in Sources */, B5276B6F247D21A000D2F56A /* NetworkingClient.swift in Sources */, diff --git a/AudioStreaming.xcodeproj/xcshareddata/xcschemes/AudioStreaming.xcscheme b/AudioStreaming.xcodeproj/xcshareddata/xcschemes/AudioStreaming.xcscheme index 2b03310..98e86cc 100644 --- a/AudioStreaming.xcodeproj/xcshareddata/xcschemes/AudioStreaming.xcscheme +++ b/AudioStreaming.xcodeproj/xcshareddata/xcschemes/AudioStreaming.xcscheme @@ -64,6 +64,23 @@ debugServiceExtension = "internal" enableGPUValidationMode = "1" allowLocationSimulation = "YES"> + + + + + + + + Void) { + let started = ProcessInfo.processInfo.systemUptime + block() + print("diff for \(name): \(String(format: "%.6f", ProcessInfo.processInfo.systemUptime - started))") +} diff --git a/AudioStreaming/Core/Network/NetworkDataStream.swift b/AudioStreaming/Core/Network/NetworkDataStream.swift index 4a5e1d1..1f1dcd3 100644 --- a/AudioStreaming/Core/Network/NetworkDataStream.swift +++ b/AudioStreaming/Core/Network/NetworkDataStream.swift @@ -6,15 +6,23 @@ import Foundation internal final class NetworkDataStream { - typealias StreamResult = Result - typealias StreamCompletion = (_ event: NetworkDataStream.StreamEvent) -> Void + typealias StreamResult = Result + typealias StreamCompletion = (_ event: NetworkDataStream.ResponseEvent) -> Void - struct StreamResponse { + enum State { + case initialised + case resumed + case suspended + case cancelled + case finished + } + + struct Response { let response: HTTPURLResponse? let data: Data? } - enum StreamEvent { + enum ResponseEvent { case stream(StreamResult) case complete(Completion) case response(HTTPURLResponse?) @@ -31,8 +39,14 @@ internal final class NetworkDataStream { private let underlyingQueue: DispatchQueue private let id: UUID + private var state: State + + var isCancelled: Bool { + state == .cancelled + } + /// the underlying task of the network request - private var task: URLSessionTask? + var task: URLSessionTask? var urlResponse: HTTPURLResponse? { task?.response as? HTTPURLResponse @@ -41,6 +55,7 @@ internal final class NetworkDataStream { internal init(id: UUID, underlyingQueue: DispatchQueue) { self.id = id self.underlyingQueue = underlyingQueue + state = .initialised } func task(for request: URLRequest, using session: URLSession) -> URLSessionTask { @@ -57,13 +72,16 @@ internal final class NetworkDataStream { @discardableResult func resume() -> Self { - underlyingQueue.async { [weak self] in - self?.task?.resume() - } + guard state.canBecome(.resumed) else { return self } + state = .resumed + task?.resume() return self } func cancel() { + guard state.canBecome(.cancelled) else { return } + state = .cancelled + streamCallback = nil task?.cancel() task = nil } @@ -82,7 +100,7 @@ internal final class NetworkDataStream { underlyingQueue.async { [weak self] in guard let self = self else { return } guard let streamCallback = self.streamCallback else { return } - let streamResponse = StreamResponse(response: response, data: data) + let streamResponse = Response(response: response, data: data) streamCallback(.stream(.success(streamResponse))) } } @@ -114,3 +132,26 @@ extension NetworkDataStream: Hashable { hasher.combine(id) } } + +extension NetworkDataStream.State { + func canBecome(_ state: NetworkDataStream.State) -> Bool { + switch (self, state) { + case (.initialised, _): + return true + case (_, .initialised), + (.cancelled, _), + (.finished, _): + return false + case (.resumed, .cancelled), + (.resumed, .suspended), + (.suspended, .resumed), + (.suspended, .cancelled): + return true + case (.suspended, .suspended), + (.resumed, .resumed): + return false + case (_, .finished): + return true + } + } +} diff --git a/AudioStreaming/Core/Network/NetworkingClient.swift b/AudioStreaming/Core/Network/NetworkingClient.swift index 4038894..11599ef 100644 --- a/AudioStreaming/Core/Network/NetworkingClient.swift +++ b/AudioStreaming/Core/Network/NetworkingClient.swift @@ -15,12 +15,12 @@ public enum NetworkError: Error, Equatable { 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 + case (.failure, failure): + return true + case (.serverError, .serverError): + return true + default: + return false } } } @@ -74,9 +74,9 @@ internal final class NetworkingClient { } internal func remove(task: NetworkDataStream) { - networkQueue.async { [weak self] in - self?.activeTasks.remove(task) - self?.tasks[task] = nil + activeTasks.remove(task) + if tasks.isEmpty { + tasks[task] = nil } } @@ -86,13 +86,10 @@ internal final class NetworkingClient { /// - parameter stream: The `NetworkDataStream` object to be performed /// - parameter request: The `URLRequest` for the `stream` private func setupRequest(_ stream: NetworkDataStream, request: URLRequest) { - networkQueue.async { [weak self] in - guard let self = self else { return } - - self.activeTasks.insert(stream) - let task = stream.task(for: request, using: self.session) - self.tasks[stream] = task - } + guard !stream.isCancelled else { return } + activeTasks.insert(stream) + let task = stream.task(for: request, using: session) + tasks[stream] = task } } diff --git a/AudioStreaming/Core/Structures/BiMap.swift b/AudioStreaming/Core/Structures/BiMap.swift index a6762de..9aafb07 100644 --- a/AudioStreaming/Core/Structures/BiMap.swift +++ b/AudioStreaming/Core/Structures/BiMap.swift @@ -10,6 +10,10 @@ struct BiMap where Left: Hashable, Right: Hashable { private var leftToRight: [Left: Right] = [:] private var rightToLeft: [Right: Left] = [:] + var isEmpty: Bool { + leftValues.isEmpty && rightValues.isEmpty + } + var leftValues: [Left] { leftToRight.lazy.map(\.key) } diff --git a/AudioStreaming/Streaming/AudioPlayer/AudioPlayer.swift b/AudioStreaming/Streaming/AudioPlayer/AudioPlayer.swift index 4e92b11..4cd7b1f 100644 --- a/AudioStreaming/Streaming/AudioPlayer/AudioPlayer.swift +++ b/AudioStreaming/Streaming/AudioPlayer/AudioPlayer.swift @@ -77,7 +77,7 @@ public final class AudioPlayer { let playerRenderProcessor: AudioPlayerRenderProcessor private let audioReadSource: DispatchTimerSource - private let underlyingQueue = DispatchQueue(label: "streaming.core.queue", qos: .userInitiated) + private let underlyingQueue = DispatchQueue(label: "streaming.core.queue", qos: .userInitiated, attributes: .concurrent) private let sourceQueue: DispatchQueue private(set) lazy var networking = NetworkingClient() @@ -93,7 +93,7 @@ public final class AudioPlayer { entriesQueue = PlayerQueueEntries() sourceQueue = DispatchQueue(label: "source.queue", qos: .userInitiated, target: underlyingQueue) - audioReadSource = DispatchTimerSource(interval: .milliseconds(500), queue: underlyingQueue) + audioReadSource = DispatchTimerSource(interval: .milliseconds(200), queue: underlyingQueue) fileStreamProcessor = AudioFileStreamProcessor(playerContext: playerContext, rendererContext: rendererContext, @@ -202,11 +202,38 @@ public final class AudioPlayer { Logger.debug("resuming audio engine failed: %@", category: .generic, args: error.localizedDescription) } - playerContext.audioPlayingEntry?.resume() + if let playingEntry = playerContext.audioReadingEntry { + if playingEntry.seekRequest.requested { + rendererContext.resetBuffers() + } + playingEntry.resume() + } startPlayer(resetBuffers: false) startReadProcessFromSourceIfNeeded() } + public func seek(to time: Double) { + guard let playingEntry = playerContext.audioPlayingEntry else { + return + } + playingEntry.seekRequest.lock.lock() + let alreadyRequestedToSeek = playingEntry.seekRequest.requested + playingEntry.seekRequest.requested = true + playingEntry.seekRequest.time = time + playingEntry.seekRequest.lock.unlock() + + if !alreadyRequestedToSeek { + playingEntry.seekRequest.version.write { version in + version += 1 + } + playingEntry.suspend() + checkRenderWaitingAndNotifyIfNeeded() + sourceQueue.async { [weak self] in + self?.processSource() + } + } + } + /// The duration of the audio, in seconds. /// /// **NOTE** In live audio playback this will be `0.0` @@ -229,12 +256,11 @@ public final class AudioPlayer { /// The progress of the audio playback, in seconds. public func progress() -> Double { - // TODO: account for seek request guard playerContext.internalState != .pendingNext else { return 0 } playerContext.entriesLock.lock() let playingEntry = playerContext.audioPlayingEntry playerContext.entriesLock.unlock() - guard let entry = playingEntry else { return 0 } + guard let entry = playingEntry, !entry.seekRequest.requested else { return 0 } return entry.lock.around { return Double(entry.seekTime) + (Double(entry.framesState.played) / outputAudioFormat.sampleRate) @@ -306,6 +332,18 @@ public final class AudioPlayer { self.processSource() } } + + fileStreamProcessor.effectCallback = { [weak self] effect in + guard let self = self else { return } + switch effect { + case .proccessSource: + self.sourceQueue.async { + self.processSource() + } + case let .raiseError(error): + self.raiseUnxpected(error: error) + } + } } /// Attaches and connect nodes to the `AudioEngine`. @@ -347,7 +385,7 @@ public final class AudioPlayer { /// /// - parameter reason: A value of `AudioPlayerStopReason` indicating the reason the engine stopped. private func stopEngine(reason: AudioPlayerStopReason) { - guard isEngineRunning else { + guard isEngineRunning && player.auAudioUnit.isRunning else { Logger.debug("already already stopped 🛑", category: .generic) return } @@ -406,6 +444,27 @@ public final class AudioPlayer { playerContext.internalState = .waitingForData setCurrentReading(entry: entry, startPlaying: true, shouldClearQueue: true) rendererContext.resetBuffers() + } else if let playingEntry = playerContext.audioPlayingEntry, + playingEntry.seekRequest.requested, + playingEntry != playerContext.audioReadingEntry + { + playingEntry.audioStreamState.processedDataFormat = false + playingEntry.reset() + if let readingEntry = playerContext.audioReadingEntry { + readingEntry.delegate = nil + readingEntry.close() + } + if configuration.flushQueueOnSeek { + playerContext.setInternalState(to: .waitingForDataAfterSeek) + setCurrentReading(entry: playingEntry, startPlaying: true, shouldClearQueue: true) + } else { + entriesQueue.requeueBufferingEntries { audioEntry in + audioEntry.reset() + } + playerContext.setInternalState(to: .waitingForDataAfterSeek) + setCurrentReading(entry: playingEntry, startPlaying: true, shouldClearQueue: true) + } + } else if playerContext.audioReadingEntry == nil { if entriesQueue.count(for: .upcoming) > 0 { let entry = entriesQueue.dequeue(type: .upcoming) @@ -419,6 +478,29 @@ public final class AudioPlayer { } } } + + if let playingEntry = playerContext.audioPlayingEntry, + playingEntry.audioStreamState.processedDataFormat, + playingEntry.calculatedBitrate() > 0.0 + { + let currSeekVersion = playingEntry.seekRequest.version.value + let originalSeekToTimeRequested = playingEntry.seekRequest.requested + + if originalSeekToTimeRequested, playerContext.audioReadingEntry === playingEntry { + proccessSeekTime() + + playingEntry.seekRequest.version.read { version in + if currSeekVersion == version { + playingEntry.seekRequest.requested = false + } + } + } + } + } + + private func proccessSeekTime() { + assert(playerContext.audioReadingEntry === playerContext.audioPlayingEntry, "reading and playing entry must be the same") + fileStreamProcessor.processSeek() } private func setCurrentReading(entry: AudioEntry?, startPlaying: Bool, shouldClearQueue: Bool) { @@ -479,10 +561,12 @@ public final class AudioPlayer { if let nextEntry = nextEntry { if !isPlayingSameItemProbablySeek { - nextEntry.lock.lock() - nextEntry.seekTime = 0 - nextEntry.lock.unlock() - // seek requested no. + nextEntry.lock.around { + nextEntry.seekTime = 0 + } + nextEntry.seekRequest.lock.around { + nextEntry.seekRequest.requested = false + } } playerContext.entriesLock.lock() playerContext.audioPlayingEntry = nextEntry @@ -551,7 +635,6 @@ extension AudioPlayer: AudioStreamSourceDelegate { } } - // TODO: check for discontinuous stream and add flag if fileStreamProcessor.isFileStreamOpen { guard fileStreamProcessor.parseFileStreamBytes(data: data) == noErr else { if let playingEntry = playerContext.audioPlayingEntry, playingEntry.has(same: source) { diff --git a/AudioStreaming/Streaming/AudioPlayer/AudioRendererContext.swift b/AudioStreaming/Streaming/AudioPlayer/AudioRendererContext.swift index a994b61..a04770b 100644 --- a/AudioStreaming/Streaming/AudioPlayer/AudioRendererContext.swift +++ b/AudioStreaming/Streaming/AudioPlayer/AudioRendererContext.swift @@ -24,6 +24,9 @@ final class AudioRendererContext { let framesRequiredToStartPlaying: UInt32 let framesRequiredAfterRebuffering: UInt32 + let framesRequiredForDataAfterSeekPlaying: UInt32 + + var waitingForDataAfterSeekFrameCount = Protected(0) private let configuration: AudioPlayerConfiguration @@ -34,6 +37,7 @@ final class AudioRendererContext { framesRequiredToStartPlaying = UInt32(canonicalStream.mSampleRate) * UInt32(configuration.secondsRequiredToStartPlaying) framesRequiredAfterRebuffering = UInt32(canonicalStream.mSampleRate) * UInt32(configuration.secondsRequiredToStartPlayingAfterBufferUnderun) + framesRequiredForDataAfterSeekPlaying = UInt32(canonicalStream.mSampleRate) * UInt32(configuration.gracePeriodAfterSeekInSeconds) let dataByteSize = Int(canonicalStream.mSampleRate * configuration.bufferSizeInSeconds) * Int(canonicalStream.mBytesPerFrame) inOutAudioBufferList = allocateBufferList(dataByteSize: dataByteSize) @@ -55,7 +59,8 @@ final class AudioRendererContext { /// Resets the `BufferContext` public func resetBuffers() { lock.lock(); defer { lock.unlock() } - bufferContext.reset() + bufferContext.frameStartIndex = 0 + bufferContext.frameUsedCount = 0 } } diff --git a/AudioStreaming/Streaming/AudioPlayer/Processors/AudioFileStreamProcessor.swift b/AudioStreaming/Streaming/AudioPlayer/Processors/AudioFileStreamProcessor.swift index e9d69ff..596bcd3 100644 --- a/AudioStreaming/Streaming/AudioPlayer/Processors/AudioFileStreamProcessor.swift +++ b/AudioStreaming/Streaming/AudioPlayer/Processors/AudioFileStreamProcessor.swift @@ -19,16 +19,24 @@ struct AudioConvertInfo { let packDescription: UnsafeMutablePointer? } +enum FileStreamProccessorEffect { + case proccessSource + case raiseError(AudioPlayerError) +} + /// An object that handles the proccessing of AudioFileStream, its packets etc. final class AudioFileStreamProcessor { private let maxCompressedPacketForBitrate = 4096 + var effectCallback: ((FileStreamProccessorEffect) -> Void)? + private let playerContext: AudioPlayerContext private let rendererContext: AudioRendererContext private let outputAudioFormat: AudioStreamBasicDescription internal var audioFileStream: AudioFileStreamID? internal var audioConverter: AudioConverterRef? + internal var discontinuous: Bool = false internal var inputFormat = AudioStreamBasicDescription() internal var fileFormat: String = "" @@ -74,11 +82,68 @@ final class AudioFileStreamProcessor { func parseFileStreamBytes(data: Data) -> OSStatus { guard let stream = audioFileStream else { return 0 } guard !data.isEmpty else { return 0 } + let flags: AudioFileStreamParseFlags = discontinuous ? .discontinuity : .init() return data.withUnsafeBytes { buffer -> OSStatus in - AudioFileStreamParseBytes(stream, UInt32(buffer.count), buffer.baseAddress, .init()) + AudioFileStreamParseBytes(stream, UInt32(buffer.count), buffer.baseAddress, flags) } } + func processSeek() { + guard let stream = audioFileStream else { return } + guard let readingEntry = playerContext.audioReadingEntry else { + return + } + + guard readingEntry.calculatedBitrate() > 0.0 || (playerContext.audioPlayingEntry?.length ?? 0) > 0 else { + return + } + + let dataOffset = Double(readingEntry.audioStreamState.dataOffset) + let seekTimeToProgress = readingEntry.seekRequest.time / readingEntry.duration() + let dataLengthInBytes = Double(readingEntry.audioDataLengthBytes()) + var seekByteOffset = Int64((dataOffset + seekTimeToProgress) * dataLengthInBytes) + + if seekByteOffset > readingEntry.length - (2 * Int(readingEntry.processedPacketsState.bufferSize)) { + seekByteOffset = Int64(readingEntry.length - (2 * Int(readingEntry.processedPacketsState.bufferSize))) + } + + readingEntry.lock.lock() + readingEntry.seekTime = readingEntry.seekRequest.time + readingEntry.lock.unlock() + + let bitrate = readingEntry.calculatedBitrate() + if readingEntry.processedPacketsState.count > 0, bitrate > 0 { + var ioFlags = AudioFileStreamSeekFlags(rawValue: 0) + var packetsAlignedByteOffset: Int64 = 0 + let seekPacket = Int64(floor(readingEntry.seekRequest.time / readingEntry.packetDuration)) + + guard AudioFileStreamSeek(stream, seekPacket, &packetsAlignedByteOffset, &ioFlags) == noErr else { + Logger.error("seek failed", category: .generic) + return + } + + let dataOffset = Int64(readingEntry.audioStreamState.dataOffset) + seekByteOffset = packetsAlignedByteOffset + dataOffset + if !ioFlags.contains(.offsetIsEstimated) { + let delta = Double((seekByteOffset - dataOffset) - packetsAlignedByteOffset) / bitrate * 8 + + readingEntry.lock.lock() + readingEntry.seekTime -= delta + readingEntry.lock.unlock() + } + } + + if let converted = audioConverter { + AudioConverterReset(converted) + } + + readingEntry.reset() + readingEntry.seek(at: Int(seekByteOffset)) + rendererContext.waitingForDataAfterSeekFrameCount.write { $0 = 0 } + playerContext.internalState = .waitingForDataAfterSeek + rendererContext.resetBuffers() + } + /// Creates an `AudioConverter` instance to be used for converting the remote audio data to the canonical audio format /// /// - parameter fromFormat: An `AudioStreamBasicDescription` indicating the format of the remote audio @@ -100,7 +165,7 @@ final class AudioFileStreamProcessor { if audioConverter == nil { guard AudioConverterNew(&inputFormat, &outputFormat, &audioConverter) == noErr else { - // raise error... + effectCallback?(.raiseError(.audioSystemError(.fileStreamError))) return } } @@ -118,7 +183,7 @@ final class AudioFileStreamProcessor { return } guard AudioFileStreamSetProperty(fileStream, kAudioConverterDecompressionMagicCookie, cookieSize, cookie) == noErr else { - // todo raise error + effectCallback?(.raiseError(.audioSystemError(.fileStreamError))) return } } @@ -175,6 +240,9 @@ final class AudioFileStreamProcessor { var packetCountSize = UInt32(MemoryLayout.size(ofValue: packetCount)) AudioFileStreamGetProperty(fileStream, kAudioFileStreamProperty_AudioDataPacketCount, &packetCountSize, &packetCount) playerContext.audioPlayingEntry?.audioStreamState.dataPacketCount = Double(packetCount) + if playerContext.audioPlayingEntry?.audioStreamFormat.mFormatID != kAudioFormatLinearPCM { + discontinuous = true + } } private func processFileFormat(fileStream: AudioFileStreamID) { @@ -209,9 +277,8 @@ final class AudioFileStreamProcessor { playerContext.audioReadingEntry?.lock.around { playerContext.audioReadingEntry?.processedPacketsState.bufferSize = packetBufferSize } - } - if let readingEntry = playerContext.audioReadingEntry { - createAudioConverter(from: readingEntry.audioStreamFormat, to: outputAudioFormat) + + createAudioConverter(from: audioStreamFormat, to: outputAudioFormat) } } @@ -261,7 +328,7 @@ final class AudioFileStreamProcessor { if let playingEntry = playerContext.audioPlayingEntry, playingEntry.seekRequest.requested, playingEntry.calculatedBitrate() > 0 { - // TODO: call proccess source on player... + effectCallback?(.proccessSource) if rendererContext.waiting.value { rendererContext.packetsSemaphore.signal() } @@ -273,6 +340,9 @@ final class AudioFileStreamProcessor { return } + // reset discontinuity + discontinuous = false + var convertInfo = AudioConvertInfo(done: false, numberOfPackets: inNumberPackets, packDescription: inPacketDescriptions) @@ -314,11 +384,11 @@ final class AudioFileStreamProcessor { { return } - // TODO: check for seek time and proccess + if let playingEntry = playerContext.audioPlayingEntry, playingEntry.seekRequest.requested, playingEntry.calculatedBitrate() > 0 { - // TODO: call proccess source on player... + effectCallback?(.proccessSource) if rendererContext.waiting.value { rendererContext.packetsSemaphore.signal() } diff --git a/AudioStreaming/Streaming/AudioPlayer/Processors/AudioPlayerRenderProcessor.swift b/AudioStreaming/Streaming/AudioPlayer/Processors/AudioPlayerRenderProcessor.swift index a086f2d..f880740 100644 --- a/AudioStreaming/Streaming/AudioPlayer/Processors/AudioPlayerRenderProcessor.swift +++ b/AudioStreaming/Streaming/AudioPlayer/Processors/AudioPlayerRenderProcessor.swift @@ -187,14 +187,25 @@ final class AudioPlayerRenderProcessor: NSObject { offset: Int(totalFramesCopied * frameSizeInBytes)) if playingEntry != nil || AudioPlayer.InternalState.waiting.contains(state) { - // buffering if playerContext.internalState != .rebuffering { playerContext.setInternalState(to: .rebuffering, when: { state -> Bool in state.contains(.running) && state != .paused }) } } else if state == .waitingForDataAfterSeek { - // TODO: implement this + if totalFramesCopied == 0 { + rendererContext.waitingForDataAfterSeekFrameCount.write { $0 += Int32(inNumberFrames - totalFramesCopied) } + if rendererContext.waitingForDataAfterSeekFrameCount.value > rendererContext.framesRequiredForDataAfterSeekPlaying { + if playerContext.internalState != .playing { + playerContext.setInternalState(to: .playing) { state -> Bool in + state.contains(.running) && state != .playing + } + } + rendererContext.waitingForDataAfterSeekFrameCount.write { $0 = 0 } + } + } else { + rendererContext.waitingForDataAfterSeekFrameCount.write { $0 = 0 } + } } } diff --git a/AudioStreaming/Streaming/AudioSource/AudioEntry.swift b/AudioStreaming/Streaming/AudioSource/AudioEntry.swift index 776c6d9..029aaa4 100644 --- a/AudioStreaming/Streaming/AudioSource/AudioEntry.swift +++ b/AudioStreaming/Streaming/AudioSource/AudioEntry.swift @@ -29,16 +29,21 @@ internal class AudioEntry { source.audioFileHint } + var length: Int { + source.length + } + var audioStreamFormat = AudioStreamBasicDescription() - var seekTime: Float + /// Hold the seek time, if a seek was requested + var seekTime: Double private(set) var seekRequest: SeekRequest private(set) var audioStreamState: AudioStreamState private(set) var framesState: EntryFramesState private(set) var processedPacketsState: ProcessedPacketsState - private var packetDuration: Double { + var packetDuration: Double { return Double(audioStreamFormat.mFramesPerPacket) / Double(sampleRate) } @@ -124,7 +129,7 @@ internal class AudioEntry { return Double(audioDataLengthBytes()) / (calculatedBitrate / 8) } - private func audioDataLengthBytes() -> UInt { + func audioDataLengthBytes() -> UInt { if let byteCount = audioStreamState.dataByteCount { return UInt(byteCount) } diff --git a/AudioStreaming/Streaming/AudioSource/Models/SeekRequest.swift b/AudioStreaming/Streaming/AudioSource/Models/SeekRequest.swift index 18eacdd..b265164 100644 --- a/AudioStreaming/Streaming/AudioSource/Models/SeekRequest.swift +++ b/AudioStreaming/Streaming/AudioSource/Models/SeekRequest.swift @@ -9,6 +9,8 @@ import Foundation final class SeekRequest { + let lock = UnfairLock() var requested: Bool = false - var time: Float = 0 + var version = Protected(0) + var time: Double = 0 } diff --git a/AudioStreaming/Streaming/AudioSource/RemoteAudioSource.swift b/AudioStreaming/Streaming/AudioSource/RemoteAudioSource.swift index 73be749..a67ae5e 100644 --- a/AudioStreaming/Streaming/AudioSource/RemoteAudioSource.swift +++ b/AudioStreaming/Streaming/AudioSource/RemoteAudioSource.swift @@ -39,6 +39,7 @@ public class RemoteAudioSource: AudioStreamSource { internal let underlyingQueue: DispatchQueue internal let streamOperationQueue: OperationQueue + private var streamOperations: [BlockOperation] = [] init(networking: NetworkingClient, metadataStreamSource: MetadataStreamSource, @@ -85,12 +86,13 @@ public class RemoteAudioSource: AudioStreamSource { } func close() { + streamOperationQueue.isSuspended = true + streamOperationQueue.cancelAllOperations() streamRequest?.cancel() if let streamTask = streamRequest { networkingClient.remove(task: streamTask) } streamRequest = nil - streamOperationQueue.cancelAllOperations() } func seek(at offset: Int) { @@ -105,7 +107,6 @@ public class RemoteAudioSource: AudioStreamSource { return } - resume() performOpen(seek: offset) } @@ -135,12 +136,11 @@ public class RemoteAudioSource: AudioStreamSource { // MARK: - Network Handle Methods - private func handleResponse(event: NetworkDataStream.StreamEvent) { + private func handleResponse(event: NetworkDataStream.ResponseEvent) { switch event { case let .response(urlResponse): - addStreamOperation { [weak self] in - self?.parseResponseHeader(response: urlResponse) - } + parseResponseHeader(response: urlResponse) + resume() case let .stream(event): addStreamOperation { [weak self] in self?.handleStreamEvent(event: event) @@ -205,9 +205,8 @@ public class RemoteAudioSource: AudioStreamSource { urlRequest.addValue("1", forHTTPHeaderField: "Icy-MetaData") if let supportsSeek = parsedHeaderOutput?.supportsSeek, supportsSeek, seekOffset > 0 { - urlRequest.addValue("bytes=\(seekOffset)", forHTTPHeaderField: "Range") + urlRequest.addValue("bytes=\(seekOffset)-", forHTTPHeaderField: "Range") } - return urlRequest } @@ -218,7 +217,11 @@ public class RemoteAudioSource: AudioStreamSource { /// - Parameter block: A closure to be executed private func addStreamOperation(_ block: @escaping () -> Void) { let operation = BlockOperation(block: block) + if let lastOp = streamOperations.last { + operation.addDependency(lastOp) + } streamOperationQueue.addOperation(operation) + streamOperations.append(operation) } /// Schedules the given block on the stream operation queue as a completion @@ -226,10 +229,11 @@ public class RemoteAudioSource: AudioStreamSource { /// - Parameter block: A closure to be executed private func addCompletionOperation(_ block: @escaping () -> Void) { let operation = BlockOperation(block: block) - if let lastOperation = streamOperationQueue.operations.last { + if let lastOperation = streamOperations.last { operation.addDependency(lastOperation) } streamOperationQueue.addOperation(operation) + streamOperations = [] } } diff --git a/AudioStreaming/Streaming/Helpers/PlayerQueueEntries.swift b/AudioStreaming/Streaming/Helpers/PlayerQueueEntries.swift index 69bbb0a..3e482d4 100644 --- a/AudioStreaming/Streaming/Helpers/PlayerQueueEntries.swift +++ b/AudioStreaming/Streaming/Helpers/PlayerQueueEntries.swift @@ -11,20 +11,20 @@ enum PlayerQueueType { /// Handles the buffering and upcoming upcoming `AudioEntries` /// The underlying objects are defined as `Queue` final class PlayerQueueEntries { - private let syncQueue: DispatchQueue + private let lock = UnfairLock() private var bufferring: Queue private var upcoming: Queue /// Returns `true` when both underlying entries are empty var isEmpty: Bool { - syncQueue.sync { + lock.around { bufferring.isEmpty && upcoming.isEmpty } } /// Returns the count of both underlying entries var count: Int { - syncQueue.sync { + lock.around { bufferring.count + upcoming.count } } @@ -32,17 +32,14 @@ final class PlayerQueueEntries { init() { bufferring = Queue() upcoming = Queue() - syncQueue = DispatchQueue(label: "sync.qe", qos: .default, attributes: .concurrent) } /// Adds the `item` to the underlying queue for the specified `type` /// - parameter item: An `AudioEntry` object to be added /// - parameter type: The type fo the underlying queue as expressed by `PlayerQueueType` func enqueue(item: AudioEntry, type: PlayerQueueType) { - syncQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.queue(for: type).enqueue(item: item) - } + lock.lock(); defer { lock.unlock() } + queue(for: type).enqueue(item: item) } /// Returns and removes the `item` to the underlying queue for the specified `type` @@ -50,62 +47,60 @@ final class PlayerQueueEntries { /// - parameter type: The type fo the underlying queue as expressed by `PlayerQueueType` /// - returns: An `AudioEntry` if found func dequeue(type: PlayerQueueType) -> AudioEntry? { - syncQueue.sync(flags: .barrier) { - queue(for: type).dequeue() - } + lock.lock(); defer { lock.unlock() } + return queue(for: type).dequeue() } /// Appends (skips) the `items` to the underlying queue for the specified `type` /// - parameter item: An `AudioEntry` object to be added /// - parameter type: The type fo the underlying queue as expressed by `PlayerQueueType` func skip(items: [AudioEntry], type: PlayerQueueType) { - syncQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.queue(for: type).skip(items: items) - } + lock.lock(); defer { lock.unlock() } + queue(for: type).skip(items: items) } /// Append (skip) the `item` to the underlying queue for the specified `type` /// - parameter item: An `AudioEntry` object to be added /// - parameter type: The type fo the underlying queue as expressed by `PlayerQueueType` func skip(item: AudioEntry, type: PlayerQueueType) { - syncQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.queue(for: type).skip(item: item) - } + lock.lock(); defer { lock.unlock() } + queue(for: type).skip(item: item) } func count(for type: PlayerQueueType) -> Int { - syncQueue.sync { - queue(for: type).count - } + lock.lock(); defer { lock.unlock() } + return queue(for: type).count } /// Removes all elements from the specified queue type func removeAll(for type: PlayerQueueType) { - syncQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.queue(for: type).removeAll() - } + lock.lock(); defer { lock.unlock() } + queue(for: type).removeAll() } /// Removes all elements from all queue type func removeAll() { - syncQueue.async(flags: .barrier) { [weak self] in - guard let self = self else { return } - self.queue(for: .buffering).removeAll() - self.queue(for: .upcoming).removeAll() - } + lock.lock(); defer { lock.unlock() } + queue(for: .buffering).removeAll() + queue(for: .upcoming).removeAll() } /// Returns an array of `AudioEntryId` of both underlying queues /// - returns: The newly constructed array of `AudioEntryId` objects func pendingEntriesId() -> [AudioEntryId] { - syncQueue.sync { - let upcomingIds = upcoming.map { $0.id } - let bufferingIds = bufferring.map { $0.id } - return upcomingIds + bufferingIds - } + lock.lock(); defer { lock.unlock() } + let upcomingIds = upcoming.map { $0.id } + let bufferingIds = bufferring.map { $0.id } + return upcomingIds + bufferingIds + } + + func requeueBufferingEntries(block: (AudioEntry) -> Void) { + lock.lock(); defer { lock.unlock() } + let bufferring = queue(for: .buffering) + bufferring.forEach(block) + let bufferringItems = bufferring.map { $0 } + queue(for: .upcoming).skip(items: bufferringItems) + queue(for: .buffering).removeAll() } /// - parameter type: A `PlayerQueueType`