Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 14 additions & 3 deletions js/module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1302,9 +1302,20 @@ export interface ITransition extends ISource {

export interface IConfigurable {
/**
* Update the settings of the source instance
* correlating to the values held within the
* object passed.
* Merge the supplied values into this instance's settings.
* Supported settings depend on the source or encoder type.
* For video sources, OBS applies changes during video-frame processing,
* so their effects may still be pending when this method returns.
* For media-file sources (`ffmpeg_source`), OSN owns the internal `caching`
* setting. Updates reset it to false while OSN reevaluates cache eligibility
* and memory use; supplying `caching: true` does not force caching.
*
* @param settings JSON-serializable settings to merge; omitted keys are preserved
* except for the OSN-managed media caching flag described above.
* @returns Nothing. Success does not guarantee that the changes have taken effect.
* @throws {TypeError} If settings cannot be converted to an object or JSON serialized.
* @throws {Error} If the source or encoder no longer exists, or the request to OBS fails.
* If communication fails after the update is sent, the settings may already have changed.
*/
update(settings: ISettings): void;

Expand Down
108 changes: 103 additions & 5 deletions obs-studio-server/source/memory-manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,15 @@ struct MediaCacheManager::SourceEntry {
std::atomic<uint64_t> revision{1};
std::atomic<bool> removed{false};

// Serializes guard entry with cache-write validation and the OBS settings
// patch. Acquire only without m_mutex; graphics uses try_lock to avoid waiting.
// Guards release this before returning to the external settings writer.
std::mutex cacheWriteMutex;

// Settings guards and the graphics callback access these under the queue mutex.
unsigned settingsUpdatesInProgress = 0;
uint64_t queryAllowedFromTick = 0;

// Only the worker changes these fields, under the queue mutex.
uint64_t scheduledRevision = 0;
uint64_t reservedBytes = 0; // Includes an enable operation awaiting completion.
Expand All @@ -83,6 +92,24 @@ MediaCacheManager::~MediaCacheManager()
shutdown();
}

MediaCacheManager::SourceSettingsUpdate::SourceSettingsUpdate(MediaCacheManager *manager, std::shared_ptr<SourceEntry> source) noexcept
: m_manager(manager), m_sourceEntry(std::move(source))
{
}

MediaCacheManager::SourceSettingsUpdate::SourceSettingsUpdate(SourceSettingsUpdate &&other) noexcept
: m_manager(std::exchange(other.m_manager, nullptr)), m_sourceEntry(std::move(other.m_sourceEntry))
{
}

MediaCacheManager::SourceSettingsUpdate::~SourceSettingsUpdate()
{
if (m_manager)
m_manager->finishSourceSettingsUpdate(m_sourceEntry);
// Dropping m_sourceEntry may release the last OBS source reference, so it must
// happen after the queue mutex is unlocked.
}

void MediaCacheManager::initialize()
{
// Initialization and shutdown are serialized by the OBS API lifecycle.
Expand Down Expand Up @@ -146,6 +173,59 @@ void MediaCacheManager::requestCacheUpdate(obs_source_t *source)
m_changed.notify_all();
}

MediaCacheManager::SourceSettingsUpdate MediaCacheManager::trackSourceSettingsUpdate(obs_source_t *source, obs_data_t *pendingSettings)
{
std::shared_ptr<SourceEntry> entry;
{
std::lock_guard lock(m_mutex);
auto it = m_sources.find(source);
if (!m_accepting || it == m_sources.end())
return {nullptr, {}};
entry = it->second;
}
// A cache write already past validation must finish before the caller can
// change live settings. Never wait for that write while holding m_mutex.
std::lock_guard cacheWriteLock(entry->cacheWriteMutex);
Comment thread
aleksandr-voitenko marked this conversation as resolved.
{
std::lock_guard lock(m_mutex);
if (!m_accepting || entry->removed)
return {nullptr, {}};
++entry->revision;
++entry->settingsUpdatesInProgress;
m_notified = true;
}
// A partial update would otherwise inherit the old player's cache enable.
// Clear it before the caller can mutate local_file, and also sanitize any
// supplied settings copy. Do not call obs_source_update here: that would
// schedule the plugin before the caller has finished its settings changes.
OBSDataAutoRelease settings = obs_source_get_settings(source);
obs_data_set_bool(settings, "caching", false);
if (pendingSettings)
obs_data_set_bool(pendingSettings, "caching", false);
// Keep the old reservation until a fresh query after the guarded update.
// An older SetCaching completion may still be waiting for the worker, and
// the old player can still be alive until OBS applies the external update.
m_changed.notify_all();
return {this, std::move(entry)};
}

void MediaCacheManager::finishSourceSettingsUpdate(const std::shared_ptr<SourceEntry> &source)
{
{
std::lock_guard lock(m_mutex);
if (!m_accepting || source->removed)
return;
--source->settingsUpdatesInProgress;
++source->revision;
// Tick callbacks run before deferred source updates. The first callback
// after this write must skip querying; the following callback is after
// OBS has had a source-update phase, even if the write finished mid-frame.
source->queryAllowedFromTick = m_graphicsTick + 2;
m_notified = true;
}
m_changed.notify_all();
}

void MediaCacheManager::requestAllCacheUpdates()
{
{
Expand Down Expand Up @@ -223,7 +303,7 @@ void MediaCacheManager::setCaching(obs_source_t *source, bool caching)
// Apply only our setting; never write back an old copy of the source's
// unrelated settings after the user has edited them.
OBSDataAutoRelease patch = obs_data_create();
// "caching" is a custom Streamlabs setting, OBS does not use it
// Streamlabs setting consumed by ffmpeg_source to enable media caching.
obs_data_set_bool(patch, "caching", caching);
obs_source_update(source, patch);
}
Expand All @@ -237,10 +317,20 @@ void MediaCacheManager::graphicsTick(void *param, float)
if (!manager.m_accepting)
return;

// Take one batch per tick.
while (!manager.m_graphicsJobs.empty() && results.size() < JOBS_PER_TICK) {
results.push_back({std::move(manager.m_graphicsJobs.front())});
++manager.m_graphicsTick;
// Scan one bounded batch. Keep delayed queries queued without consuming
// readiness retries or blocking work for another source behind them.
const auto count = std::min(manager.m_graphicsJobs.size(), JOBS_PER_TICK);
for (size_t i = 0; i < count; ++i) {
auto job = std::move(manager.m_graphicsJobs.front());
manager.m_graphicsJobs.pop_front();
const auto &entry = *job.source;
if (job.type == JobType::QuerySource && !entry.removed && entry.revision == job.revision &&
(entry.settingsUpdatesInProgress || manager.m_graphicsTick < entry.queryAllowedFromTick)) {
Comment thread
Copilot marked this conversation as resolved.
manager.m_graphicsJobs.push_back(std::move(job));
continue;
}
results.push_back({std::move(job)});
}
}
if (results.empty())
Expand All @@ -251,6 +341,14 @@ void MediaCacheManager::graphicsTick(void *param, float)
// its deferred settings update first.
for (auto &result : results) {
auto &job = result.job;
std::unique_lock<std::mutex> cacheWriteLock;
if (job.type == JobType::SetCaching) {
// Hold through revision/settings validation and the write. Otherwise a
// guard could change local_file after we validate the old reservation.
cacheWriteLock = std::unique_lock(job.source->cacheWriteMutex, std::try_to_lock);
if (!cacheWriteLock.owns_lock())
continue; // An unapplied completion schedules fresh evaluation.
}
if (job.source->removed || job.source->revision != job.revision)
continue;
result.valid = true;
Expand All @@ -259,7 +357,7 @@ void MediaCacheManager::graphicsTick(void *param, float)
queryMediaOnGraphicsThread(job.source->source, result.snapshot);
} else if (!job.targetCachingEnabled || (result.snapshot.eligible && result.snapshot.file == job.file)) {
if (result.snapshot.caching != job.targetCachingEnabled)
setCaching(job.source->source, job.targetCachingEnabled);
manager.m_setCaching(job.source->source, job.targetCachingEnabled);
result.applied = true;
}
}
Expand Down
66 changes: 62 additions & 4 deletions obs-studio-server/source/memory-manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,16 +35,46 @@
// their estimated total size within a shared memory budget.
//
// One worker owns cache decisions and accounting. OBS callbacks only invalidate
// entries; media queries and settings changes run in the graphics tick callback.
// entries; media queries and cache-enable writes run in the graphics tick callback.
//
// The ffmpeg_source metadata handlers execute synchronously and access its
// current media player. OBS source updates and video ticks can destroy or
// replace that player on the graphics thread. Retaining the OBS source does
// not keep that particular player alive, and our mutex does not protect it.
// Run those queries in the graphics tick to serialize them with replacement,
// then return copied metadata to the worker.
// Settings become visible before the plugin applies them. External writers must
// use trackSourceSettingsUpdate() so queries wait for a source-update phase and
// cannot associate a new filename with the previous player's metadata.
// Guard entry also serializes with cache-setting writes: validation and the
// caching patch must finish before an external writer can change the file.
// The guard then clears the cache flag in live and incoming settings so the
// changed player starts uncached until its metadata has been reevaluated.
class MediaCacheManager {
struct SourceEntry;

public:
// Keeps cache queries pending while the caller changes a source's settings.
class SourceSettingsUpdate {
public:
// Transfers responsibility for finishing the update; the moved-from guard
// is inactive. The manager must outlive the guard.
SourceSettingsUpdate(SourceSettingsUpdate &&other) noexcept;
// Finishes tracking and schedules evaluation after OBS can apply the update.
// Safe after source unregistration or manager shutdown, but must run before
// obs_shutdown(): the guard still retains its source reference.
~SourceSettingsUpdate();
SourceSettingsUpdate(const SourceSettingsUpdate &) = delete;
SourceSettingsUpdate &operator=(const SourceSettingsUpdate &) = delete;
SourceSettingsUpdate &operator=(SourceSettingsUpdate &&) = delete;

private:
friend class MediaCacheManager;
SourceSettingsUpdate(MediaCacheManager *manager, std::shared_ptr<SourceEntry> source) noexcept;
MediaCacheManager *m_manager;
std::shared_ptr<SourceEntry> m_sourceEntry;
};

// Returns a non-owning reference to the process-wide singleton. Access is
// thread-safe, but does not initialize OBS or start the cache worker.
static MediaCacheManager &GetInstance();
Expand Down Expand Up @@ -75,7 +105,29 @@ class MediaCacheManager {
// obs_source_remove() or explicitly disable the source's cache setting.
void unregisterSource(obs_source_t *source);

// Requests reevaluation after a registered source's settings or activity change.
// Invalidates older work before external code changes a source's live settings.
// Keep the returned guard alive through obs_source_update(), including any
// property callbacks that mutate those settings. Retains the original
// registration so finishing this guard cannot affect a later registration
// of the same source.
// Clears the manager-owned "caching" flag in live settings before returning.
// Pass any separate settings object that will be applied as pendingSettings;
// its caching flag is also cleared so a copied enable cannot be written back.
// Both pointers are borrowed. Other settings are preserved. The caller must
// not set caching again; the worker reevaluates it after the guarded update.
// The existing reservation is reconciled by that evaluation, not guard entry.
// Nested updates are supported; queries wait until all guards for this source
// finish and OBS has had a source-update phase. May wait for an executing cache
// setting write before returning; no queue lock is held while waiting. Does not
// wait for queries; their results are invalidated. Does not serialize external
// settings writers or hold a lock for the returned guard's lifetime.
// No queue lock is held across the caller's OBS operations. Null/unregistered
// sources and calls while stopped or stopping return an inactive guard.
// The manager and OBS runtime must outlive the guard.
[[nodiscard]] SourceSettingsUpdate trackSourceSettingsUpdate(obs_source_t *source, obs_data_t *pendingSettings = nullptr);

// Requests reevaluation after a registered source's activity change.
// Use trackSourceSettingsUpdate() around settings writes instead.
// Safe from OBS callbacks; repeated requests are coalesced, and this call does
// not wait for media queries or cache-setting changes to complete.
// Null/unregistered sources and calls while stopped or stopping are ignored.
Expand All @@ -87,7 +139,9 @@ class MediaCacheManager {
void requestAllCacheUpdates();

// Stops accepting requests, removes the tick callback, cancels queued work,
// joins the worker and drops all retained source references.
// joins the worker and drops its queued and tracked source references.
// Outstanding SourceSettingsUpdate guards retain their own references and must
// finish before obs_shutdown(), while this manager is still alive.
// Call before obs_shutdown(), from the serialized OBS lifecycle, on neither
// the graphics thread nor this manager's worker. Do not overlap another
// shutdown() or initialize(). Waits for any executing tick callback and worker,
Expand All @@ -100,7 +154,6 @@ class MediaCacheManager {
private:
friend class MediaCacheManagerTestAccess;
using Clock = std::chrono::steady_clock;
struct SourceEntry;
struct Snapshot {
std::string file;
bool eligible = false;
Expand Down Expand Up @@ -134,6 +187,7 @@ class MediaCacheManager {
void complete(Completion &result);
void queueCacheSettingUpdate(const std::shared_ptr<SourceEntry> &source, uint64_t revision, bool targetCachingEnabled);
void releaseBudget(SourceEntry &source);
void finishSourceSettingsUpdate(const std::shared_ptr<SourceEntry> &source);

// Protects only queues and bookkeeping. Never held during an OBS call,
// source release, wait for graphics, or thread join.
Expand All @@ -148,8 +202,12 @@ class MediaCacheManager {
bool m_stopping = false;
bool m_notified = false;
bool m_rebalance = false;
// Counts every graphics callback, including empty ticks; survives video reset.
uint64_t m_graphicsTick = 0;
uint64_t m_reservedCacheBytes = 0;
uint64_t m_cacheBudgetBytes;
// Injectable clock for deterministic retry tests.
std::function<Clock::time_point()> m_now = Clock::now;
// Injectable OBS settings write for deterministic cache-write race tests.
std::function<void(obs_source_t *, bool)> m_setCaching = setCaching;
};
8 changes: 6 additions & 2 deletions obs-studio-server/source/osn-source.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,8 @@ void osn::Source::GetProperties(void *data, const int64_t id, const std::vector<

rval.push_back(ipc::value((uint64_t)ErrorCode::Ok));

// Property callbacks may change live settings before obs_source_update().
auto settingsUpdate = MediaCacheManager::GetInstance().trackSourceSettingsUpdate(src);
obs_properties_t *prp = obs_source_properties(src);
obs_data *settings = obs_source_get_settings(src);

Expand Down Expand Up @@ -354,8 +356,10 @@ void osn::Source::Update(void *data, const int64_t id, const std::vector<ipc::va
}
}

obs_source_update(src, sets);
MediaCacheManager::GetInstance().requestCacheUpdate(src);
{
auto settingsUpdate = MediaCacheManager::GetInstance().trackSourceSettingsUpdate(src, sets);
obs_source_update(src, sets);
}
obs_data_release(sets);

obs_data_t *updatedSettings = obs_source_get_settings(src);
Expand Down
Loading
Loading