Adds initial support for seeking

This commit is contained in:
Dimitris C
2020-11-15 20:14:59 +00:00
parent 101c7ddf34
commit cece65475d
16 changed files with 392 additions and 118 deletions
@@ -85,6 +85,18 @@
isEnabled = "NO">
</EnvironmentVariable>
</EnvironmentVariables>
<AdditionalOptions>
<AdditionalOption
key = "MallocStackLogging"
value = ""
isEnabled = "YES">
</AdditionalOption>
<AdditionalOption
key = "PrefersMallocStackLoggingLite"
value = ""
isEnabled = "YES">
</AdditionalOption>
</AdditionalOptions>
</LaunchAction>
<ProfileAction
buildConfiguration = "Release"
+30 -21
View File
@@ -38,7 +38,7 @@ enum AudioContent: Int, CaseIterable {
var streamUrl: URL {
switch self {
case .enlefko:
return URL(string: "https://ample-02.radiojar.com/srzwv225e3quv?rj-ttl=5&rj-tok=AAABcmnuHngA2PJ4KSBI9k5cCw")!
return URL(string: "https://stream.radiojar.com/srzwv225e3quv")!
case .offradio:
return URL(string: "http://s3.yesstreaming.net:7033/stream")!
case .pepper966:
@@ -46,7 +46,8 @@ enum AudioContent: Int, CaseIterable {
case .radiox:
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://p.scdn.co/mp3-preview/cab4b09c23ffc11774d879977131df9d150fcef4?cid=d8a5ed958d274c2e8ee717e6a4b0971d")!
return URL(string: "https://t4.bcbits.com/stream/8fce516db797c2496c024ef8560f1f29/mp3-128/3790863592?p=0&ts=1604330347&t=6415e93344ffaf498b6b4b24717782a046c214ce&token=1604330347_7161bf285cfd2ad2eed6f000b0b280313229290b")!
case .piano:
return URL(string: "https://www.kozco.com/tech/piano2-CoolEdit.mp3")!
}
@@ -144,6 +145,7 @@ class ViewController: UIViewController {
slider.thumbTintColor = .black
}
slider.isContinuous = true
slider.semanticContentAttribute = .playback
slider.addTarget(self, action: #selector(sliderTouchedDown), for: .touchDown)
slider.addTarget(self, action: #selector(sliderTouchedUp), for: [.touchUpInside, .touchUpOutside])
slider.addTarget(self, action: #selector(sliderValueChanged), for: .valueChanged)
@@ -187,20 +189,24 @@ class ViewController: UIViewController {
return button
}
var seekValue: Float = 0
var isScrubbing: Bool = false
@objc
func sliderTouchedDown() {
stopDisplayLink(resetLabels: false)
isScrubbing = true
}
@objc
func sliderTouchedUp() {
startDisplayLink()
isScrubbing = false
if player.duration() > 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]) {
+4
View File
@@ -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 = "<group>"; };
B5667A912499063D00D93F85 /* AudioPlayerContext.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AudioPlayerContext.swift; sourceTree = "<group>"; };
B5667B3D249BC43000D93F85 /* AudioPlayerRenderProcessor.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AudioPlayerRenderProcessor.swift; sourceTree = "<group>"; };
B573733F254DE43E003DFBEC /* measure.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = measure.swift; sourceTree = "<group>"; };
B57829CE2548B32B00C78D36 /* Lock.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = Lock.swift; sourceTree = "<group>"; };
B58386372544A2C10087A712 /* EntryFrames.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = EntryFrames.swift; sourceTree = "<group>"; };
B583863F254584A50087A712 /* ProcessedPackets.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ProcessedPackets.swift; sourceTree = "<group>"; };
@@ -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 */,
@@ -64,6 +64,23 @@
debugServiceExtension = "internal"
enableGPUValidationMode = "1"
allowLocationSimulation = "YES">
<AdditionalOptions>
<AdditionalOption
key = "MallocStackLogging"
value = ""
isEnabled = "YES">
</AdditionalOption>
<AdditionalOption
key = "PrefersMallocStackLoggingLite"
value = ""
isEnabled = "YES">
</AdditionalOption>
<AdditionalOption
key = "NSZombieEnabled"
value = "YES"
isEnabled = "YES">
</AdditionalOption>
</AdditionalOptions>
</LaunchAction>
<ProfileAction
buildConfiguration = "Release"
+15
View File
@@ -0,0 +1,15 @@
//
// measure.swift
// AudioStreaming
//
// Created by Dimitrios Chatzieleftheriou on 31/10/2020.
// Copyright © 2020 Decimal. All rights reserved.
//
import Foundation
func measure(name: String = "", block: () -> Void) {
let started = ProcessInfo.processInfo.systemUptime
block()
print("diff for \(name): \(String(format: "%.6f", ProcessInfo.processInfo.systemUptime - started))")
}
@@ -6,15 +6,23 @@
import Foundation
internal final class NetworkDataStream {
typealias StreamResult = Result<StreamResponse, Error>
typealias StreamCompletion = (_ event: NetworkDataStream.StreamEvent) -> Void
typealias StreamResult = Result<Response, Error>
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
}
}
}
@@ -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
}
}
@@ -10,6 +10,10 @@ struct BiMap<Left, Right> 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)
}
@@ -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) {
@@ -24,6 +24,9 @@ final class AudioRendererContext {
let framesRequiredToStartPlaying: UInt32
let framesRequiredAfterRebuffering: UInt32
let framesRequiredForDataAfterSeekPlaying: UInt32
var waitingForDataAfterSeekFrameCount = Protected<Int32>(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
}
}
@@ -19,16 +19,24 @@ struct AudioConvertInfo {
let packDescription: UnsafeMutablePointer<AudioStreamPacketDescription>?
}
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()
}
@@ -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 }
}
}
}
@@ -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)
}
@@ -9,6 +9,8 @@
import Foundation
final class SeekRequest {
let lock = UnfairLock()
var requested: Bool = false
var time: Float = 0
var version = Protected<Int>(0)
var time: Double = 0
}
@@ -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 = []
}
}
@@ -11,20 +11,20 @@ enum PlayerQueueType {
/// Handles the buffering and upcoming upcoming `AudioEntries`
/// The underlying objects are defined as `Queue<AudioEntry>`
final class PlayerQueueEntries {
private let syncQueue: DispatchQueue
private let lock = UnfairLock()
private var bufferring: Queue<AudioEntry>
private var upcoming: Queue<AudioEntry>
/// 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<AudioEntry>()
upcoming = Queue<AudioEntry>()
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`