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`