Files
libspinscale/probe/probeComponentThread.cpp
T
latentprion 88d6367491 probe: barrier TLS init until make_shared arms shared_from_this.
PuppeteerThread starts its OS thread in the constructor, so initializeTls()
could race and throw bad_weak_ptr; wait for the harness barrier first.
2026-07-12 06:10:41 -04:00

159 lines
4.3 KiB
C++

#include <probe/probeComponentThread.h>
#include <condition_variable>
#include <iostream>
#include <mutex>
#include <spinscale/component.h>
namespace sscl::probe {
namespace {
constexpr sscl::ThreadId PROBE_PUPPETEER_THREAD_ID = 2;
class ProbeDummyPuppeteerComponent
: public sscl::pptr::PuppeteerComponent
{
public:
explicit ProbeDummyPuppeteerComponent(
const std::shared_ptr<sscl::PuppeteerThread>& componentThreadIn)
: sscl::pptr::PuppeteerComponent(componentThreadIn)
{}
void handleLoopExceptionHook() override
{
std::cerr << "ProbeComponentThreadHarness: puppeteer loop exception\n";
}
};
/** EXPLANATION:
* PuppeteerThread starts its std::thread inside the constructor, but
* enable_shared_from_this::weak_this is only armed after make_shared returns.
* Without a barrier, initializeTls()'s shared_from_this() races and throws
* std::bad_weak_ptr. DedicatedIoThread uses the same handshake.
*/
struct ProbeThreadStartupState
{
std::mutex mutex;
std::condition_variable condition;
bool allowInitialization = false;
};
void waitForProbeThreadStartupPermission(
const std::shared_ptr<ProbeThreadStartupState>& startupState)
{
std::unique_lock<std::mutex> lock(startupState->mutex);
startupState->condition.wait(
lock,
[&startupState]() { return startupState->allowInitialization; });
}
void releaseProbeThreadStartupBarrier(
const std::shared_ptr<ProbeThreadStartupState>& startupState)
{
{
std::lock_guard<std::mutex> guard(startupState->mutex);
startupState->allowInitialization = true;
}
startupState->condition.notify_all();
}
void probePuppeteerMain(
const sscl::PuppeteerThread::EntryFnArguments& args,
const std::function<void(
const std::shared_ptr<sscl::ComponentThread>&)>& work,
std::promise<std::exception_ptr>& donePromise,
const std::shared_ptr<ProbeThreadStartupState>& startupState)
{
waitForProbeThreadStartupPermission(startupState);
sscl::PuppeteerThread& thr = args.usableBeforeJolt;
thr.initializeTls();
sscl::ComponentThread::setPuppeteerThreadId(PROBE_PUPPETEER_THREAD_ID);
std::shared_ptr<sscl::PuppeteerThread> thrPtr =
std::static_pointer_cast<sscl::PuppeteerThread>(thr.shared_from_this());
sscl::ComponentThread::setPuppeteerThread(thrPtr);
try {
work(thrPtr);
donePromise.set_value(nullptr);
}
catch (...) {
donePromise.set_value(std::current_exception());
}
thr.getIoContext().stop();
}
} // namespace
ProbeComponentThreadHarness::ProbeComponentThreadHarness(
const char *threadName)
: threadName(threadName),
dummyComponent(std::make_shared<ProbeDummyPuppeteerComponent>(
std::shared_ptr<sscl::PuppeteerThread>()))
{}
ProbeComponentThreadHarness::~ProbeComponentThreadHarness() = default;
std::shared_ptr<sscl::ComponentThread>
ProbeComponentThreadHarness::componentThread() const
{
return lastComponentThread;
}
void ProbeComponentThreadHarness::runSync(
const std::function<void(
const std::shared_ptr<sscl::ComponentThread>&)>& work)
{
std::promise<std::exception_ptr> donePromise;
std::future<std::exception_ptr> doneFuture = donePromise.get_future();
auto startupState = std::make_shared<ProbeThreadStartupState>();
std::shared_ptr<sscl::PuppeteerThread> runThread =
std::make_shared<sscl::PuppeteerThread>(
PROBE_PUPPETEER_THREAD_ID,
threadName,
[&work, &donePromise, startupState](
const sscl::PuppeteerThread::EntryFnArguments& args)
{
probePuppeteerMain(args, work, donePromise, startupState);
},
*dummyComponent,
nullptr);
dummyComponent->thread = runThread;
lastComponentThread = runThread;
releaseProbeThreadStartupBarrier(startupState);
runThread->thread.join();
std::exception_ptr probeException = doneFuture.get();
if (probeException) {
std::rethrow_exception(probeException);
}
}
void runNonViralNurseryOnComponentThread(
const std::shared_ptr<sscl::ComponentThread>& componentThread,
std::function<sscl::co::NonViralNonPostingInvoker(
sscl::co::NonViralTaskNursery::Slot::Lease&)> invokerFactory,
std::chrono::milliseconds timeout)
{
(void)timeout;
sscl::co::NonViralTaskNursery nursery;
nursery.openAdmission();
nursery.launch(
[&invokerFactory](sscl::co::NonViralTaskNursery::Slot::Lease& lease)
{
return invokerFactory(lease);
});
nursery.closeAdmission();
nursery.syncAwaitAllSettlements(componentThread->getIoContext());
}
} // namespace sscl::probe