new cspot/bell

This commit is contained in:
philippe44
2023-05-06 23:50:26 +02:00
parent e0e7e718ba
commit 8bad480112
163 changed files with 6611 additions and 6739 deletions

View File

@@ -1,18 +1,26 @@
#include "TrackPlayer.h"
#include <mutex> // for mutex, scoped_lock
#include <string> // for string
#include <type_traits> // for remove_extent_t
#include <vector> // for vector, vector<>::value_type
#include <mutex> // for mutex, scoped_lock
#include <string> // for string
#include <type_traits> // for remove_extent_t
#include <vector> // for vector, vector<>::value_type
#include "BellLogger.h" // for AbstractLogger
#include "BellUtils.h" // for BELL_SLEEP_MS
#include "CDNTrackStream.h" // for CDNTrackStream, CDNTrackStream::TrackInfo
#include "Logger.h" // for CSPOT_LOG
#include "Packet.h" // for cspot
#include "TrackProvider.h" // for TrackProvider
#include "TrackQueue.h" // for CDNTrackStream, CDNTrackStream::TrackInfo
#include "WrappedSemaphore.h" // for WrappedSemaphore
#ifdef BELL_VORBIS_FLOAT
#define VORBIS_SEEK(file, position) (ov_time_seek(file, (double)position / 1000))
#define VORBIS_READ(file, buffer, bufferSize, section) (ov_read(file, buffer, bufferSize, 0, 2, 1, section))
#else
#define VORBIS_SEEK(file, position) (ov_time_seek(file, position))
#define VORBIS_READ(file, buffer, bufferSize, section) \
(ov_read(file, buffer, bufferSize, section))
#endif
namespace cspot {
struct Context;
struct TrackReference;
@@ -38,13 +46,14 @@ static long vorbisTellCb(TrackPlayer* self) {
return self->_vorbisTell();
}
TrackPlayer::TrackPlayer(std::shared_ptr<cspot::Context> ctx, isAiringCallback isAiring, EOFCallback eof, TrackLoadedCallback trackLoaded)
: bell::Task("cspot_player", 48 * 1024, 5, 1) {
TrackPlayer::TrackPlayer(std::shared_ptr<cspot::Context> ctx,
std::shared_ptr<cspot::TrackQueue> trackQueue,
EOFCallback eof, TrackLoadedCallback trackLoaded)
: bell::Task("cspot_player", 32 * 1024, 5, 1) {
this->ctx = ctx;
this->isAiring = isAiring;
this->eofCallback = eof;
this->trackLoaded = trackLoaded;
this->trackProvider = std::make_shared<cspot::TrackProvider>(ctx);
this->trackQueue = trackQueue;
this->playbackSemaphore = std::make_unique<bell::WrappedSemaphore>(5);
// Initialize vorbis callbacks
@@ -55,9 +64,6 @@ TrackPlayer::TrackPlayer(std::shared_ptr<cspot::Context> ctx, isAiringCallback i
(decltype(ov_callbacks::close_func))&vorbisCloseCb,
(decltype(ov_callbacks::tell_func))&vorbisTellCb,
};
isRunning = true;
startTask();
}
TrackPlayer::~TrackPlayer() {
@@ -65,123 +71,192 @@ TrackPlayer::~TrackPlayer() {
std::scoped_lock lock(runningMutex);
}
void TrackPlayer::loadTrackFromRef(TrackReference& ref, size_t positionMs,
bool startAutomatically) {
this->playbackPosition = positionMs;
this->autoStart = startAutomatically;
auto nextTrack = trackProvider->loadFromTrackRef(ref);
stopTrack();
this->sequence++;
this->currentTrackStream = nextTrack;
this->playbackSemaphore->give();
void TrackPlayer::start() {
if (!isRunning) {
isRunning = true;
startTask();
}
}
void TrackPlayer::stopTrack() {
void TrackPlayer::stop() {
isRunning = false;
resetState();
std::scoped_lock lock(runningMutex);
}
void TrackPlayer::resetState() {
// Mark for reset
this->pendingReset = true;
this->currentSongPlaying = false;
std::scoped_lock lock(playbackMutex);
std::scoped_lock lock(dataOutMutex);
CSPOT_LOG(info, "Resetting state");
}
void TrackPlayer::seekMs(size_t ms) {
std::scoped_lock lock(seekMutex);
#ifdef BELL_VORBIS_FLOAT
ov_time_seek(&vorbisFile, (double)ms / 1000);
#else
ov_time_seek(&vorbisFile, ms);
#endif
if (inFuture) {
// We're in the middle of the next track, so we need to reset the player in order to seek
resetState();
}
CSPOT_LOG(info, "Seeking...");
this->pendingSeekPositionMs = ms;
}
void TrackPlayer::runTask() {
std::scoped_lock lock(runningMutex);
std::shared_ptr<QueuedTrack> track, newTrack = nullptr;
int trackOffset = 0;
bool eof = false;
bool endOfQueueReached = false;
while (isRunning) {
this->playbackSemaphore->twait(100);
if (this->currentTrackStream == nullptr) {
// Ensure we even have any tracks to play
if (!this->trackQueue->hasTracks() ||
(endOfQueueReached && trackQueue->isFinished())) {
this->trackQueue->playableSemaphore->twait(300);
continue;
}
CSPOT_LOG(info, "Player received a track, waiting for it to be ready...");
// when track changed many times and very quickly, we are stuck on never-given semaphore
while (this->currentTrackStream->trackReady->twait(250));
CSPOT_LOG(info, "Got track");
// Last track was interrupted, reset to default
if (pendingReset) {
track = nullptr;
pendingReset = false;
inFuture = false;
}
if (this->currentTrackStream->status == CDNTrackStream::Status::FAILED) {
CSPOT_LOG(error, "Track failed to load, skipping it");
this->currentTrackStream = nullptr;
this->eofCallback();
endOfQueueReached = false;
// Wait 800ms. If next reset is requested in meantime, restart the queue.
// Gets rid of excess actions during rapid queueing
BELL_SLEEP_MS(50);
if (pendingReset) {
continue;
}
this->currentSongPlaying = true;
newTrack = trackQueue->consumeTrack(track, trackOffset);
this->trackLoaded();
if (newTrack == nullptr) {
if (trackOffset == -1) {
// Reset required
track = nullptr;
}
this->playbackMutex.lock();
int32_t r = ov_open_callbacks(this, &vorbisFile, NULL, 0, vorbisCallbacks);
if (playbackPosition > 0) {
#ifdef BELL_VORBIS_FLOAT
ov_time_seek(&vorbisFile, (double)playbackPosition / 1000);
#else
ov_time_seek(&vorbisFile, playbackPosition);
#endif
BELL_SLEEP_MS(100);
continue;
}
bool eof = false;
track = newTrack;
while (!eof && currentSongPlaying) {
seekMutex.lock();
#ifdef BELL_VORBIS_FLOAT
long ret = ov_read(&vorbisFile, (char*)&pcmBuffer[0], pcmBuffer.size(),
0, 2, 1, &currentSection);
#else
long ret = ov_read(&vorbisFile, (char*)&pcmBuffer[0], pcmBuffer.size(),
&currentSection);
#endif
seekMutex.unlock();
if (ret == 0) {
CSPOT_LOG(info, "EOF");
// and done :)
eof = true;
} else if (ret < 0) {
CSPOT_LOG(error, "An error has occured in the stream %d", ret);
currentSongPlaying = false;
} else {
inFuture = trackOffset > 0;
if (this->dataCallback != nullptr) {
auto toWrite = ret;
if (track->state != QueuedTrack::State::READY) {
track->loadedSemaphore->twait(5000);
while (!eof && currentSongPlaying && toWrite > 0) {
auto written =
dataCallback(pcmBuffer.data() + (ret - toWrite), toWrite,
this->currentTrackStream->trackInfo.trackId, this->sequence);
if (written == 0) {
BELL_SLEEP_MS(50);
if (track->state != QueuedTrack::State::READY) {
CSPOT_LOG(error, "Track failed to load, skipping it");
this->eofCallback();
continue;
}
}
CSPOT_LOG(info, "Got track ID=%s", track->identifier.c_str());
currentSongPlaying = true;
{
std::scoped_lock lock(playbackMutex);
currentTrackStream = track->getAudioFile();
// Open the stream
currentTrackStream->openStream();
if (pendingReset || !currentSongPlaying) {
continue;
}
if (trackOffset == 0 && pendingSeekPositionMs == 0) {
this->trackLoaded(track);
}
int32_t r =
ov_open_callbacks(this, &vorbisFile, NULL, 0, vorbisCallbacks);
if (pendingSeekPositionMs > 0) {
track->requestedPosition = pendingSeekPositionMs;
}
if (track->requestedPosition > 0) {
VORBIS_SEEK(&vorbisFile, track->requestedPosition);
}
eof = false;
CSPOT_LOG(info, "Playing");
while (!eof && currentSongPlaying) {
// Execute seek if needed
if (pendingSeekPositionMs > 0) {
uint32_t seekPosition = pendingSeekPositionMs;
// Reset the pending seek position
pendingSeekPositionMs = 0;
// Seek to the new position
VORBIS_SEEK(&vorbisFile, seekPosition);
}
long ret = VORBIS_READ(&vorbisFile, (char*)&pcmBuffer[0],
pcmBuffer.size(), &currentSection);
if (ret == 0) {
CSPOT_LOG(info, "EOF");
// and done :)
eof = true;
} else if (ret < 0) {
CSPOT_LOG(error, "An error has occured in the stream %d", ret);
currentSongPlaying = false;
} else {
if (this->dataCallback != nullptr) {
auto toWrite = ret;
while (!eof && currentSongPlaying && !pendingReset && toWrite > 0) {
int written = 0;
{
std::scoped_lock dataOutLock(dataOutMutex);
// If reset happened during playback, return
if (!currentSongPlaying || pendingReset)
break;
written =
dataCallback(pcmBuffer.data() + (ret - toWrite), toWrite, track->identifier);
}
if (written == 0) {
BELL_SLEEP_MS(50);
}
toWrite -= written;
}
toWrite -= written;
}
}
}
}
ov_clear(&vorbisFile);
ov_clear(&vorbisFile);
// With very large buffers, track N+1 can be downloaded while N has not aired yet and
// if we continue, the currentTrackStream will be emptied, causing a crash in
// notifyAudioReachedPlayback when it will look for trackInfo. A busy loop is never
// ideal, but this low impact, infrequent and more simple than yet another semaphore
while (currentSongPlaying && !isAiring()) {
BELL_SLEEP_MS(100);
}
CSPOT_LOG(info, "Playing done");
// always move back to LOADING (ensure proper seeking after last track has been loaded)
this->currentTrackStream.reset();
this->playbackMutex.unlock();
// always move back to LOADING (ensure proper seeking after last track has been loaded)
currentTrackStream = nullptr;
}
if (eof) {
if (trackQueue->isFinished()) {
endOfQueueReached = true;
}
this->eofCallback();
}
}
@@ -226,10 +301,6 @@ long TrackPlayer::_vorbisTell() {
return this->currentTrackStream->getPosition();
}
CDNTrackStream::TrackInfo TrackPlayer::getCurrentTrackInfo() {
return this->currentTrackStream->trackInfo;
}
void TrackPlayer::setDataCallback(DataCallback callback) {
this->dataCallback = callback;
}