mirror of
https://github.com/logos-co/logos-protocol.git
synced 2026-08-27 12:01:15 +00:00
lp_client_create() made the CALLING thread the client's owner thread. Callers reach it through a lazily-created wrapper (the generated bind_<iface>() -> LpClient::ensure()), so the first thread to make an outbound call captured the whole transport for the life of the process. For the qt_remote transport that thread also ends up owning the QRemoteObjectNode and its QLocalSocket, which are only serviced by a thread running a Qt event loop. A module whose first call came from a worker — an HTTP handler, a timer thread — bound its transport to a thread that only pumps events while it is already blocked inside a call. Replica acquisition then never completed: every requestObject() burned its full 20s timeout and returned nullptr, and since a failed acquire yields an empty result the data loss was silent. openmetrics-module hit exactly this: one GET /metrics took 40s (2 x 20s) and came back missing a module, /health went unanswered behind the wedged libmicrohttpd thread, and the follow-up stop RPC failed. Construct the client on the Qt main thread when the transport needs a Qt event loop, so the per-call marshal that already exists (logos::runOnOwnerThread) lands somewhere that can actually service it. This is the anchor the Qt path always had — LogosAPI::getClient marshals construction to the LogosAPI's thread — given to the lp path. Plain (tcp/tcp_ssl) and mock transports are Qt-free and thread-agnostic, so they keep the calling thread: a worker-thread consumer stays off the main thread's back. LogosTransportFactory::needsQtEventLoop() carries that rule next to the createConnection resolution it mirrors. When there is nothing to anchor to (a Qt-affine transport with no QCoreApplication) we now warn instead of letting it surface as a mute timeout. Tests: a worker thread creates an lp client over qt_remote and calls a provider published on the main thread; passes in ~0.15s, and with the construction hop reverted fails after 24.8s / 49.9s — the acquire timeouts themselves. Plus a truth table for needsQtEventLoop. 183/183 protocol tests pass. Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
79 lines
3.3 KiB
C++
79 lines
3.3 KiB
C++
#ifndef LOGOS_THREAD_MARSHAL_H
|
|
#define LOGOS_THREAD_MARSHAL_H
|
|
|
|
#include <type_traits>
|
|
#include <utility>
|
|
|
|
#include <QCoreApplication>
|
|
#include <QMetaObject>
|
|
#include <QObject>
|
|
#include <QThread>
|
|
|
|
namespace logos {
|
|
|
|
// Run `fn` on `obj`'s (owner) thread, blocking the caller until it completes,
|
|
// and forward the return value. If already on that thread, runs directly with
|
|
// no marshaling and no overhead (the common case).
|
|
//
|
|
// Why: Logos inter-module calls go over Qt Remote Objects, whose replicas only
|
|
// work on the thread that owns them (the module's main/event-loop thread). This
|
|
// lets a module call other modules from a worker thread (e.g. an HTTP server
|
|
// thread) without the module touching Qt — the SDK transparently marshals the
|
|
// call onto the owner thread.
|
|
//
|
|
// Requirements:
|
|
// - `obj`'s thread must be running an event loop (it is — the module's main
|
|
// thread runs QCoreApplication::exec()). The same-thread guard avoids the
|
|
// BlockingQueuedConnection self-deadlock.
|
|
// - The return type must be void or default-constructible (the marshaled
|
|
// branch holds the result in a local before assigning it), and must not be
|
|
// a reference (there'd be nothing to bind the local to). Both are satisfied
|
|
// by the SDK's uses here (void, QVariant, LogosObject*, LogosAPIClient*).
|
|
template <typename Fn>
|
|
auto runOnOwnerThread(QObject* obj, Fn&& fn) -> decltype(fn())
|
|
{
|
|
using Ret = decltype(fn());
|
|
static_assert(!std::is_reference_v<Ret>,
|
|
"runOnOwnerThread does not support reference return types");
|
|
if (QThread::currentThread() == obj->thread()) {
|
|
return fn();
|
|
}
|
|
if constexpr (std::is_void_v<Ret>) {
|
|
QMetaObject::invokeMethod(obj, [&]() { fn(); }, Qt::BlockingQueuedConnection);
|
|
return;
|
|
} else {
|
|
Ret ret{};
|
|
QMetaObject::invokeMethod(obj, [&]() { ret = fn(); }, Qt::BlockingQueuedConnection);
|
|
return ret;
|
|
}
|
|
}
|
|
|
|
// Run `fn` on the process's Qt main thread — the QCoreApplication's thread,
|
|
// the one thread guaranteed to be running an event loop for the life of a
|
|
// module — and forward the return value.
|
|
//
|
|
// Why this exists separately from runOnOwnerThread: that one marshals to an
|
|
// object's *existing* owner thread, which presupposes the object was created
|
|
// somewhere sane. This one is for deciding where to create it in the first
|
|
// place. A Qt-affine transport (see LogosTransportFactory::needsQtEventLoop)
|
|
// binds its node/socket to whichever thread constructs it, so construction on
|
|
// a worker thread — an HTTP handler making the module's first outbound call,
|
|
// say — permanently binds the client to a thread that only pumps events while
|
|
// it happens to be blocked inside a call. Replica acquisition then never
|
|
// completes and every call burns its full timeout.
|
|
//
|
|
// Falls back to running inline when there is no QCoreApplication (a pure-lp
|
|
// host with no Qt loop — there is no better thread to pick) or when already on
|
|
// the main thread (the common case: module init, and any call made from it).
|
|
template <typename Fn>
|
|
auto runOnQtMainThread(Fn&& fn) -> decltype(fn())
|
|
{
|
|
QCoreApplication* app = QCoreApplication::instance();
|
|
if (!app) return fn();
|
|
return runOnOwnerThread(app, std::forward<Fn>(fn));
|
|
}
|
|
|
|
} // namespace logos
|
|
|
|
#endif // LOGOS_THREAD_MARSHAL_H
|