mirror of
https://github.com/logos-co/logos-protocol.git
synced 2026-08-27 20:11:07 +00:00
* fix: make isConnected() mean connected, and stop the log claiming it QRemoteObjectNode::connectToNode() returns false only when the URL SCHEME is unregistered -- it never contacts the peer. Our registry URLs are COMPUTED rather than discovered (logos_instance.h: local:logos_<module>_<instanceId>), so they are identical whether or not the module exists. Latching m_connected from that return therefore made isConnected() answer "yes" for modules that were never loaded, which made every `if (!client->isConnected()) return;` guard in the codebase DEAD CODE. Callers then paid a 20 s waitForSource per call, twice over, because the token handshake tries capability_module first. Measured in Basecamp with package_manager absent: ~417 s of blocked GUI thread on macOS and 361 s on Linux before the window appeared, and over 900 s under load. Not a Windows bug -- the Windows port merely exposed it. isConnected() now also requires a listener at the endpoint. For `local:` that is a direct socket / named-pipe probe, which costs microseconds precisely in the case that used to cost 20 seconds; any other scheme keeps its previous behaviour. Two logging changes, because the diagnostics cost more than the defect: "Successfully connected to registry" asserted a connection that often did not exist and sent three separate investigations to the wrong place -- it now says a connect attempt started and makes no claim about the peer. And requestObject warns BEFORE a doomed wait instead of going silent for 20 s and then reporting failure. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: let event subscriptions survive a module that is not reachable yet requestObject() answers "is the module there RIGHT NOW", and every subscriber in this codebase asks at the one moment the answer is no: a module's init(), a UI backend's onContextReady(), a QML view's Component.onCompleted. All of those run while the dependency's host process has been spawned but has not called listen() yet. The subscriber then gave up permanently -- lp_subscribe returned nullptr with no log at all, and callers turned that into a `false` the documented example discards. Method calls kept working through the same window because acquireCachedObject() reaches the replica by a path that never asks, so the symptom was "events are broken", not "the subscription never happened".1238316(isConnected() means connected) is what made this deterministic rather than lucky, and it must not be reverted -- it removed ~417 s (macOS) / 361 s (Linux) of blocked GUI thread at Basecamp startup. So the subscription becomes deferrable instead. - LogosTransportAsyncAcquire: a sibling interface (dynamic_cast, like LogosObjectErrorChannel) so LogosTransportConnection's installed vtable is unchanged. requestObjectWhenAvailable() registers interest and returns; it never blocks and never spins a nested event loop. - qt_remote implements it by acquiring a dynamic replica before the peer exists -- legal, free, and armed by the node's existing 250 ms reconnect loop, so it adds no polling. Delivery is deferred one event-loop turn because stateChanged fires from inside onClientRead (the refresh_balances re-entrancy SIGSEGV). - LogosAPIConsumer::onEventWhenAvailable() holds the pending subscriptions, arms them when the object appears, shares ONE handle per object (separate from the call cache, so a call re-acquiring a stale handle cannot silently kill a live subscription), and re-arms them after reconnect(). Unbounded in time on purpose -- a module can be installed mid-session -- but bounded in noise: one warning at 3 s, one at 60 s, a log line when it arms, and a loud abandon when the transport proves it impossible. - lp_subscribe routes through it, which fixes the same defect for every C++/Nim/Rust module and UI backend without touching qt-sdk or any generated code. tests/protocol/test_deferred_subscription.cpp pins all three layers, each with a published-first control so a red case cannot be a mis-wired fixture. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: close the remaining silent-failure holes in deferred event subscriptions The deferred-subscription registry from the previous commit fixed the reported defect, but review found six ways it could still lose a subscription without saying so — five in the registry itself, one in the plain transport's host — and every one of them lived in a cell with no test. All of its tests ran in Remote mode; three of the four transports had none at all. Registry (cpp/logos_api_consumer.cpp): * An already-present module was deferred to the first 250 ms tick on every transport without a deferred acquire, and every event emitted in that window was dropped. lp_subscribe used to attach synchronously and deliver them, so this relocated the silent event loss rather than removing it. startAcquire() now reports which of three answers the transport gave, and only an Unsupported answer takes the one synchronous requestObject() — which is also what keeps that call structurally away from qt_remote, whose requestObject() enters waitForSource()'s nested event loop even at timeout 0. Previously that invariant lived in a comment, and tick() could reach it whenever acquireDynamic() returned null. * reconnected() put every armed subscription back in the pending set but never restarted the timer, which takeMatching() had stopped when they armed. Since tick() is the sole driver of both the retry and the watchdog, a reconnect left the subscription dead AND silent — quieter than the "not connected" warning it replaced. * armAgainst() released a stale handle while entries were still attached to its event helper. Those entries stayed in m_armed, never fired again, and reported as healthy. They are now revived and re-armed against the new handle. * The retry timer ran forever at the 5 s cap with nothing to do. It now stops once every pending entry has an acquire in flight and has said everything it will say, and restarts when that changes. * A cancelled subscription had no way to leave the registry, so lp_unsubscribe left it holding the timer up and warning about a subscription nobody wanted. onEventWhenAvailable() now returns an id; cancelEventSubscription() and eventSubscriptionState() are its counterparts, and lp_unsubscribe uses them. Plain transport (cpp/implementations/plain/plain_transport_host.cpp): * onSubscribe() dropped a Subscribe for an object that was not published YET — which is exactly when consumers subscribe — and the consumer could not know, because requestObject() had already succeeded. Publishing also overwrote the sink table wholesale, so a republish took every subscriber down with it. The sinks now live in a table keyed independently of publication. Also adds lp_pending_subscriptions() to the C ABI. The Qt consumer has had this visibility all along and the C ABI had none, which is why a subscription that silently never armed was undetectable from Rust, Nim or a universal C++ module. tests/protocol/test_event_delivery_matrix.cpp pins the product rather than a sample of it: 3 transports x 2 provider kinds (Qt-native and universal/std, which reach the wire by different conversions) x 2 consumer paths (onEventWhenAvailable and lp_subscribe) x 6 timings, plus mock and the non-blocking guard. Every delivery case has a control that is green independently of these fixes. One thing that is NOT fixed and is now stated in the contract: arming is not retroactive and no transport buffers, so a module that emits a one-shot "ready" event synchronously inside its own init() can still be missed. That window is inherent to the transport — the blocking requestObject() this replaced had it too — but "subscriptions survive a late module" is not "no event can be missed". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * docs: name the QtRO invariant the stale-handle revive rests on * test(events): state what the non-blocking guard can and cannot catch The acquireCount assertion catches a retry that polls qt_remote's blocking requestObject() in the ordinary case. It cannot reach the narrow one -- the poll is only reachable when the transport declines a deferred acquire while still reporting connected, which needs acquireDynamic() to return null and is not forcible from outside. That case is held shut by control flow instead, and saying so is better than leaving a reader to assume the test covers it. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: make the async-acquire contract and lp_subscribe's return honest Both from review on #47, both real. The LogosTransportAsyncAcquire contract promised that a true return means onReady "WILL be invoked exactly once". It will not: RemoteTransportConnection parents every in-flight PendingAcquire to m_pendingAcquires, which is reset at the top of the destructor and rebuilt on reconnect, so an accepted request is cancelled silently with no callback whenever the connection it belongs to goes away. The contract now says AT MOST once, names both cancellation triggers, and states what a caller has to do about them — re-issue after a reconnect, or carry its own deadline. It also records that the layer above already does the first, which is why a subscription made through onEventWhenAvailable() survives something the raw transport call does not. That asymmetry is the reason to prefer the consumer API, and it was previously implicit. lp_subscribe returned a non-null lp_subscription even when onEventWhenAvailable refused and returned 0, leaving the caller with a handle that can never fire while the ABI documents NULL as the one signal that the arguments were refused. It now checks sub->id and returns nullptr. That second one is defensive rather than a live bug, and the code says so: the guard at the top of lp_subscribe already rejects an empty event name and a null callback, and lp_client_create rejects an empty target, so the three inputs that make onEventWhenAvailable() return 0 cannot all arrive there today. No test drives it. The two contracts simply have to agree, and one of them changing is how they would stop agreeing. 374/374 green. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: stop lp_unsubscribe deadlocking, without dereferencing a freed client lp_unsubscribe took ownerGuard->mutex and, while holding it, called cancelEventSubscription(), which marshals to the owner thread with a BLOCKING queued connection. The delivery callback lp_subscribe installs runs ON that thread and takes subGuard->mutex then clientGuard->mutex — and clientGuard IS ownerGuard, both assigned from client->guard. Lock-order inversion. It also hung outright once the owner's event loop had stopped, which is exactly when a language binding drops its subscription handle. The first attempt at this dropped the guard entirely and checked `alive` inside the posted lambda. That was a use-after-free: QMetaObject::invokeMethod dereferences the target (it reads object->thread()) before the lambda can run, and lp_client_destroy sets alive=false and deletes the client synchronously — so the check was unreachable on the exact ordering lp_subscription's own comment documents as supported. Proven rather than argued: with MallocScribble=1, a test that destroys the client before unsubscribing segfaulted 6/6 with the guard removed and passed 6/6 with it restored. So the guard is held across the POST and not across the cancel. Both halves are load-bearing, and the distinction is the whole fix: posting never waits on the owner thread, so holding the mutex across it cannot invert; only the blocking marshal ever had to move. Consequence, now stated in the ABI header: un-registration is EVENTUAL. The callback-will-not-fire guarantee stays synchronous and unconditional, but lp_pending_subscriptions() may still list a just-cancelled subscription until the owner thread runs, and if the client is destroyed first the cancellation never runs at all — correct, since the registry died with it. The matrix test now pumps for the drain instead of asserting it happened synchronously. 374/374 green. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: arm a subscription immediately when the module is already reachable Deferral introduced a narrower version of the loss it removed. The common consumer shape is a call followed by a subscription in the same function -- wallet-ui's backend calls get_chains() and subscribes on the next line, the tutorial's C++ UI backend does the same. Before deferral the generated Qt wrapper acquired synchronously, so the subscription was live before on() returned and an event emitted straight after was delivered. Holding it until the next event-loop turn silently drops that event. Measured on the generated-wrapper harness: 1/1 delivered pre-migration, 0/1 after, over 3 runs. LogosTransportAsyncAcquire gains tryAcquireNow(): hand back a handle ONLY if that costs nothing -- for qt_remote, a replica that is already Valid, which is exactly the state a prior call leaves behind since QtRO shares one replica implementation per object name on a node. It must never block, never spin a nested event loop and never wait on a peer; "not immediately available" is an answer and the caller falls back to the deferred path. Default returns nullptr, so a transport that cannot answer cheaply simply does not. Delivering inline here is safe for the reason the never-synchronous rule exists: that rule protects against re-entering the transport's READ stack from a stateChanged callback. tryAcquireNow runs on the subscriber's own stack. The new matrix case fires ONCE, synchronously, with no pumping in between -- re-firing would hide the exact gap under test -- and states the transport difference rather than papering over it. Subscription registration is local on qt_remote (attach to a held replica) and qt_local (connect an in-process signal), so delivery there must be instant. On plain it is a wire frame to the host, so instant delivery was never on offer and never was before this change either; that leg asserts it still arms and delivers. Also de-flaked EventDeliveryNonBlocking: its heartbeat COUNT over a fixed wall-clock window measures the machine, not the code. The gap assertion is the one that means something; the count is now only a floor. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: stop tryAcquireNow leaving a dangling facade in QtRO's connect liste9f82acintroduced a use-after-free. tryAcquireNow() acquired a dynamic replica and, when it was not already Valid, deleted it. That is not safe: QtRO shares one replica IMPLEMENTATION per object name per node, and while that implementation is still waiting for the source's metaobject it records every facade built on it as a RAW pointer in QConnectedReplicaImplementation::m_parentsNeedingConnect. ~QRemoteObjectReplica is an empty body, so destroying a facade never deregisters it, and the implementation dereferences the whole list when the class definition arrives. So each probe of an unreachable module left one dangling pointer behind. WHY IT HID. The first probe owns the only implementation and takes it down with itself, so a single subscription is harmless. It needs a second subscription whose implementation is pinned by an in-flight PendingAcquire before a freed facade can outlive its implementation. A consumer subscribing once sees nothing; the QML plugin shape -- a view registering every event it cares about up front -- dies. REPRODUCED, 4 runs of 4, serially as well as in parallel, in logos-view-module-runtime's existing suite (unchanged from master, and green there against this same protocol checkout): LogosQmlBridge: subscription accepted for "echo_module" :: "ev13" Received signal 10 (SIGBUS), code 1, for address 0x5a SIGBUS code 1 is BUS_ADRALN -- a misaligned atomic access on a garbage base read out of a recycled heap block, in the event loop rather than at the call site, which is why it reads as a mystery crash rather than as a subscription bug. PROVEN, before writing this fix, by commenting out that single `delete replica`: the same suite went 4 failures -> 6/6 with no other change. With this fix: 6/6. THE FIX IS TO PARK, NOT TO FREE. One probe per object name, parented to m_pendingAcquires -- which both the destructor and reconnect() already destroy BEFORE the node, so the implementations die in the same breath and freeing them there is safe. Ownership transfers out only when the replica reaches Valid, by which point the implementation is configured and is no longer holding the facade. It costs one idle replica per name until it goes Valid or the connection dies. AND REMOVE THE MULTIPLIER: beginAcquire() probed on EVERY add(), ahead of startAcquire() and therefore ahead of the m_acquiring one-acquire-per-object guard. tick() already applies that filter; beginAcquire() was the one caller that did not, which is what turned one probe per module into one per subscription. While an acquire is in flight its PendingAcquire already holds a replica and will arm every waiting entry at once, so the probe buys nothing there. Not QML-specific: lp_subscribe reaches the same entry point. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
894 lines
39 KiB
C++
894 lines
39 KiB
C++
#include "remote_transport.h"
|
|
#include "../../logos_async_dispatch.h"
|
|
#include "../../logos_socket_paths.h"
|
|
#include "qt_socket_path.h"
|
|
#include <QRemoteObjectRegistryHost>
|
|
#include <QRemoteObjectNode>
|
|
#include <QRemoteObjectReplica>
|
|
#include <QRemoteObjectPendingCall>
|
|
#include <QRemoteObjectPendingCallWatcher>
|
|
#include <QTimer>
|
|
#include <QEventLoop>
|
|
#include <QDebug>
|
|
#include <QUrl>
|
|
#include <QLocalSocket>
|
|
#include <QMetaObject>
|
|
#include <QTime>
|
|
#include <QJsonArray>
|
|
#include <QVariantMap>
|
|
#include <atomic>
|
|
|
|
// Process-wide count of replicas acquired by requestObject() — a test hook to
|
|
// prove the consumer reuses one cached handle instead of re-acquiring per call.
|
|
static std::atomic<long> g_acquireCount{0};
|
|
|
|
using logos::qtremote::localSocketFilePath;
|
|
|
|
// ── RemoteLogosObject ────────────────────────────────────────────────────────
|
|
|
|
namespace {
|
|
|
|
class RemoteEventHelper : public QObject {
|
|
Q_OBJECT
|
|
public:
|
|
explicit RemoteEventHelper(QObject* parent = nullptr) : QObject(parent) {}
|
|
|
|
void addCallback(const QString& eventName, LogosObject::EventCallback cb) {
|
|
m_callbacks[eventName].append(std::move(cb));
|
|
}
|
|
|
|
public slots:
|
|
void onEventResponse(const QString& eventName, const QVariantList& data) {
|
|
// Dispatch to callbacks registered for this specific event name,
|
|
// plus any wildcard subscribers (callbacks registered with an
|
|
// empty event name, meaning "receive every event").
|
|
auto cbs = m_callbacks.value(eventName);
|
|
cbs.append(m_callbacks.value(QString()));
|
|
if (!cbs.isEmpty()) {
|
|
qDebug() << "[LogosObject] Remote EventHelper: dispatching event" << eventName << "to" << cbs.size() << "callback(s) (via IPC)";
|
|
}
|
|
for (const auto& cb : cbs) {
|
|
try { cb(eventName, data); } catch (...) {}
|
|
}
|
|
}
|
|
|
|
private:
|
|
QHash<QString, QList<LogosObject::EventCallback>> m_callbacks;
|
|
};
|
|
|
|
} // anonymous namespace
|
|
|
|
class RemoteLogosObject : public LogosObject, public LogosObjectErrorChannel {
|
|
public:
|
|
// objectName is carried purely so a failure can name the module it belongs
|
|
// to — logos::CallError::origin, the same field the acquire-time error and
|
|
// lp_invoke's out_error_json already fill in.
|
|
RemoteLogosObject(QObject* replica, QString objectName)
|
|
: m_replica(replica), m_helper(nullptr), m_objectName(std::move(objectName))
|
|
{
|
|
qDebug() << "[LogosObject] Created RemoteLogosObject wrapping QRemoteObjectReplica" << reinterpret_cast<quintptr>(replica);
|
|
if (m_replica) {
|
|
// Eager event wiring — a deferred ("multi") call's result arrives as a
|
|
// completion event, so the channel must be live even when the caller
|
|
// never subscribes to a user event. onEvent() reuses this same helper.
|
|
m_helper = new RemoteEventHelper();
|
|
QObject::connect(m_replica, SIGNAL(eventResponse(QString,QVariantList)),
|
|
m_helper, SLOT(onEventResponse(QString,QVariantList)));
|
|
m_helper->addCallback(logos::callCompleteEvent(),
|
|
[this](const QString&, const QVariantList& data) {
|
|
if (data.size() != 2) return;
|
|
const QString id = data.at(0).toString();
|
|
const QVariant result = data.at(1);
|
|
m_completions.insert(id, result);
|
|
if (QEventLoop* loop = m_completionWaiters.value(id, nullptr))
|
|
loop->quit(); // wake a sync waiter
|
|
if (m_asyncCompletionCbs.contains(id)) { // fire an async waiter
|
|
auto cb = m_asyncCompletionCbs.take(id);
|
|
m_completions.remove(id);
|
|
// Deliver on the NEXT event-loop turn, never inline. This
|
|
// runs from RemoteEventHelper::onEventResponse, which —
|
|
// cross-process — is on QtRO's onClientRead read stack. The
|
|
// user callback (LogosAPIConsumer's async lambda → the
|
|
// module's completion handler) routinely emits a module
|
|
// event (host ModuleProxy → QtRO source serialization) and
|
|
// then release()s this object. Doing either while
|
|
// onClientRead is still unwinding re-enters QtRO and
|
|
// corrupts the node — the refresh_balances SIGSEGV
|
|
// (KERN_INVALID_ADDRESS in onClientRead). Deferring runs the
|
|
// user code after the read fully unwinds. m_helper is the
|
|
// context so the callback is dropped if we're torn down
|
|
// first (mirrors the QPointer guard in invokeRemoteMethodAsync).
|
|
if (cb)
|
|
QTimer::singleShot(0, m_helper,
|
|
[cb = std::move(cb), result]() {
|
|
cb(result, logos::CallError{});
|
|
});
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
~RemoteLogosObject() override {
|
|
qDebug() << "[LogosObject] Destroying RemoteLogosObject" << reinterpret_cast<quintptr>(m_replica);
|
|
// release() normally clears m_helper first (deferred). If we get here on a
|
|
// direct delete, defer the helper too: a direct delete can still be reached
|
|
// from within the helper's own slot dispatch. See disconnectEvents().
|
|
if (m_helper) {
|
|
if (m_replica)
|
|
QObject::disconnect(m_replica, nullptr, m_helper, nullptr);
|
|
m_helper->deleteLater();
|
|
m_helper = nullptr;
|
|
}
|
|
}
|
|
|
|
// Adapter over the error-carrying implementation: discards the diagnosis,
|
|
// which is exactly what this entry point has always done.
|
|
QVariant callMethod(const QString& authToken,
|
|
const QString& methodName,
|
|
const QVariantList& args,
|
|
int timeoutMs) override
|
|
{
|
|
return callMethodWithError(authToken, methodName, args, timeoutMs, nullptr);
|
|
}
|
|
|
|
QVariant callMethodWithError(const QString& authToken,
|
|
const QString& methodName,
|
|
const QVariantList& args,
|
|
int timeoutMs,
|
|
logos::CallError* err) override
|
|
{
|
|
if (err) err->clear();
|
|
const std::string origin = m_objectName.toStdString();
|
|
const std::string method = methodName.toStdString();
|
|
|
|
if (!m_replica) {
|
|
qWarning() << "RemoteLogosObject: Cannot call method on null replica";
|
|
if (err)
|
|
*err = logos::callErrorObjectUnavailable(
|
|
origin, "replica for '" + origin + "' is gone (module unloaded"
|
|
" or transport dropped)");
|
|
return QVariant();
|
|
}
|
|
qDebug() << "[LogosObject] RemoteLogosObject::callMethod" << methodName << "args:" << args.size();
|
|
|
|
QRemoteObjectPendingCall pendingCall;
|
|
bool success = QMetaObject::invokeMethod(
|
|
m_replica,
|
|
"callRemoteMethod",
|
|
Qt::DirectConnection,
|
|
Q_RETURN_ARG(QRemoteObjectPendingCall, pendingCall),
|
|
Q_ARG(QString, authToken),
|
|
Q_ARG(QString, methodName),
|
|
Q_ARG(QVariantList, args)
|
|
);
|
|
|
|
if (!success) {
|
|
qWarning() << "RemoteLogosObject: Failed to invoke callRemoteMethod on replica";
|
|
if (err)
|
|
*err = logos::callErrorCallFailed(
|
|
origin, "replica did not accept callRemoteMethod for '"
|
|
+ method + "'");
|
|
return QVariant();
|
|
}
|
|
|
|
pendingCall.waitForFinished(timeoutMs);
|
|
|
|
// Two distinct outcomes that used to collapse into one empty QVariant:
|
|
// the deadline elapsed with the call still in flight (timeout), and QtRO
|
|
// itself failed the call (transport).
|
|
if (!pendingCall.isFinished()) {
|
|
qWarning() << "RemoteLogosObject: callRemoteMethod timed out";
|
|
if (err) *err = logos::callErrorTimeout(origin, method, timeoutMs);
|
|
return QVariant();
|
|
}
|
|
if (pendingCall.error() != QRemoteObjectPendingCall::NoError) {
|
|
qWarning() << "RemoteLogosObject: callRemoteMethod failed:" << pendingCall.error();
|
|
if (err)
|
|
*err = logos::callErrorTransport(
|
|
origin, "QtRO call to '" + origin + "." + method
|
|
+ "' failed with error "
|
|
+ std::to_string(static_cast<int>(pendingCall.error())));
|
|
return QVariant();
|
|
}
|
|
|
|
// A "multi" provider may have deferred the result (returned a pending
|
|
// sentinel); resolveDeferred waits for the completion event, or returns
|
|
// the value unchanged for an ordinary (synchronous) result.
|
|
return resolveDeferred(pendingCall.returnValue(), timeoutMs, methodName, err);
|
|
}
|
|
|
|
// Adapter over the error-carrying implementation: discards the diagnosis,
|
|
// which is exactly what this entry point has always done.
|
|
void callMethodAsync(const QString& authToken,
|
|
const QString& methodName,
|
|
const QVariantList& args,
|
|
int timeoutMs,
|
|
AsyncResultCallback callback) override
|
|
{
|
|
if (!callback) return;
|
|
callMethodAsyncWithError(authToken, methodName, args, timeoutMs,
|
|
[cb = std::move(callback)](QVariant v, const logos::CallError&) mutable {
|
|
cb(std::move(v));
|
|
});
|
|
}
|
|
|
|
void callMethodAsyncWithError(const QString& authToken,
|
|
const QString& methodName,
|
|
const QVariantList& args,
|
|
int timeoutMs,
|
|
AsyncResultErrorCallback callback) override
|
|
{
|
|
if (!callback) return;
|
|
const std::string origin = m_objectName.toStdString();
|
|
const std::string method = methodName.toStdString();
|
|
|
|
if (!m_replica) {
|
|
const logos::CallError e = logos::callErrorObjectUnavailable(
|
|
origin, "replica for '" + origin + "' is gone (module unloaded"
|
|
" or transport dropped)");
|
|
QTimer::singleShot(0, [callback, e]() { callback(QVariant(), e); });
|
|
return;
|
|
}
|
|
|
|
qDebug() << "[LogosObject] RemoteLogosObject::callMethodAsync" << methodName << "args:" << args.size();
|
|
|
|
QRemoteObjectPendingCall pendingCall;
|
|
bool success = QMetaObject::invokeMethod(
|
|
m_replica,
|
|
"callRemoteMethod",
|
|
Qt::DirectConnection,
|
|
Q_RETURN_ARG(QRemoteObjectPendingCall, pendingCall),
|
|
Q_ARG(QString, authToken),
|
|
Q_ARG(QString, methodName),
|
|
Q_ARG(QVariantList, args)
|
|
);
|
|
|
|
if (!success) {
|
|
qWarning() << "RemoteLogosObject: Failed to invoke callRemoteMethod on replica (async)";
|
|
const logos::CallError e = logos::callErrorCallFailed(
|
|
origin, "replica did not accept callRemoteMethod for '" + method + "'");
|
|
QTimer::singleShot(0, [callback, e]() { callback(QVariant(), e); });
|
|
return;
|
|
}
|
|
|
|
auto* watcher = new QRemoteObjectPendingCallWatcher(pendingCall);
|
|
|
|
// Timeout timer -- parented to the watcher so it is auto-deleted
|
|
// when the watcher is destroyed, preventing late timeout callbacks.
|
|
auto* timer = new QTimer(watcher);
|
|
timer->setSingleShot(true);
|
|
|
|
// Exactly-once gate. The finished handler and the timeout timer can
|
|
// both be queued around the same moment; without a guard that races
|
|
// into a double callback, which violates callMethodAsyncWithError's
|
|
// contract and is a latent double-free for every consumer. The same
|
|
// gate wraps the deferred-completion path so a late completion cannot
|
|
// deliver after the initial timeout already has (or vice versa).
|
|
auto delivered = std::make_shared<std::atomic_bool>(false);
|
|
AsyncResultErrorCallback deliverOnce =
|
|
[delivered, callback = std::move(callback)](QVariant result,
|
|
const logos::CallError& err) mutable {
|
|
if (delivered->exchange(true))
|
|
return;
|
|
if (callback)
|
|
callback(std::move(result), err);
|
|
};
|
|
|
|
// Success handler -- delivers result on the consumer's thread
|
|
QObject::connect(watcher, &QRemoteObjectPendingCallWatcher::finished,
|
|
watcher, [this, deliverOnce, timer, timeoutMs, origin, method, delivered](QRemoteObjectPendingCallWatcher* w) {
|
|
timer->stop(); // cancel timeout
|
|
// Timeout may already have won the race and deleteLater'd us; if
|
|
// the slot still runs, do not enter the deferred-completion path
|
|
// or we would arm a second delivery after the caller already saw
|
|
// a timeout.
|
|
if (delivered->load()) {
|
|
w->deleteLater();
|
|
return;
|
|
}
|
|
QVariant result;
|
|
logos::CallError err;
|
|
if (w->error() == QRemoteObjectPendingCall::NoError) {
|
|
result = w->returnValue();
|
|
} else {
|
|
qWarning() << "RemoteLogosObject: async callMethod error:" << w->error();
|
|
err = logos::callErrorTransport(
|
|
origin, "QtRO call to '" + origin + "." + method
|
|
+ "' failed with error "
|
|
+ std::to_string(static_cast<int>(w->error())));
|
|
}
|
|
w->deleteLater();
|
|
// A "multi" provider may have deferred the result: wait for the
|
|
// completion event instead of delivering the pending sentinel.
|
|
if (err.ok()) {
|
|
QString callId;
|
|
if (logos::isPendingCallSentinel(result, &callId)) {
|
|
if (m_completions.contains(callId)) {
|
|
deliverOnce(m_completions.take(callId), logos::CallError{});
|
|
return;
|
|
}
|
|
m_asyncCompletionCbs.insert(callId, deliverOnce);
|
|
// Bound the wait: a completion that never lands is a timeout,
|
|
// reported as one instead of as an empty result.
|
|
QTimer::singleShot(timeoutMs, m_helper, [this, callId, origin, method, timeoutMs]() {
|
|
if (m_asyncCompletionCbs.contains(callId)) {
|
|
auto cb = m_asyncCompletionCbs.take(callId);
|
|
if (cb)
|
|
cb(QVariant(),
|
|
logos::callErrorTimeout(origin, method, timeoutMs));
|
|
}
|
|
});
|
|
return;
|
|
}
|
|
}
|
|
deliverOnce(result, err);
|
|
}, Qt::QueuedConnection);
|
|
|
|
// Timeout handler -- stops the watcher and reports the elapsed deadline
|
|
QObject::connect(timer, &QTimer::timeout, watcher, [watcher, deliverOnce, origin, method, timeoutMs]() {
|
|
qWarning() << "RemoteLogosObject: async callMethod timed out";
|
|
deliverOnce(QVariant(), logos::callErrorTimeout(origin, method, timeoutMs));
|
|
watcher->deleteLater(); // also destroys the timer (child)
|
|
});
|
|
|
|
timer->start(timeoutMs);
|
|
}
|
|
|
|
bool informModuleToken(const QString& authToken,
|
|
const QString& moduleName,
|
|
const QString& token,
|
|
int timeoutMs) override
|
|
{
|
|
if (!m_replica) {
|
|
qWarning() << "RemoteLogosObject: Cannot call informModuleToken on null replica";
|
|
return false;
|
|
}
|
|
|
|
QRemoteObjectPendingCall pendingCall;
|
|
bool success = QMetaObject::invokeMethod(
|
|
m_replica,
|
|
"informModuleToken",
|
|
Qt::DirectConnection,
|
|
Q_RETURN_ARG(QRemoteObjectPendingCall, pendingCall),
|
|
Q_ARG(QString, authToken),
|
|
Q_ARG(QString, moduleName),
|
|
Q_ARG(QString, token)
|
|
);
|
|
|
|
if (!success) {
|
|
qWarning() << "RemoteLogosObject: Failed to invoke informModuleToken on replica";
|
|
return false;
|
|
}
|
|
|
|
pendingCall.waitForFinished(timeoutMs);
|
|
|
|
if (!pendingCall.isFinished() || pendingCall.error() != QRemoteObjectPendingCall::NoError) {
|
|
qWarning() << "RemoteLogosObject: informModuleToken failed or timed out:" << pendingCall.error();
|
|
return false;
|
|
}
|
|
|
|
return pendingCall.returnValue().toBool();
|
|
}
|
|
|
|
void onEvent(const QString& eventName, EventCallback callback) override
|
|
{
|
|
if (!m_replica) return;
|
|
|
|
qDebug() << "[LogosObject] RemoteLogosObject::onEvent subscribing to event:" << eventName;
|
|
if (!m_helper) {
|
|
m_helper = new RemoteEventHelper();
|
|
QObject::connect(m_replica, SIGNAL(eventResponse(QString,QVariantList)),
|
|
m_helper, SLOT(onEventResponse(QString,QVariantList)));
|
|
qDebug() << "[LogosObject] RemoteLogosObject: connected EventHelper to QRemoteObjectReplica signals (IPC)";
|
|
}
|
|
m_helper->addCallback(eventName, std::move(callback));
|
|
}
|
|
|
|
void disconnectEvents() override
|
|
{
|
|
if (!m_helper) return;
|
|
// Defer the helper's destruction instead of deleting it inline.
|
|
//
|
|
// release()/disconnectEvents() can be reached *synchronously from inside
|
|
// the helper's own onEventResponse() slot*: a deferred ("multi") call's
|
|
// result is delivered as a completion event that fires onEventResponse,
|
|
// whose callback (LogosAPIConsumer::invokeRemoteMethodAsync) runs the
|
|
// user callback and then calls plugin->release(). With the qt_remote
|
|
// transport the event arrives cross-process, so onEventResponse runs on
|
|
// the QtRO read stack (QRemoteObjectNodePrivate::onClientRead) while the
|
|
// replica is still emitting. Deleting the helper (the signal receiver)
|
|
// here — and the replica (the sender) in release() — corrupts the
|
|
// connection list Qt is still iterating, a use-after-free that crashes in
|
|
// QMetaObjectPrivate::signal / onClientRead (the refresh_balances SIGSEGV).
|
|
//
|
|
// Stop further dispatch now by disconnecting, then let the event loop
|
|
// delete the helper once the current emission has fully unwound.
|
|
if (m_replica)
|
|
QObject::disconnect(m_replica, nullptr, m_helper, nullptr);
|
|
m_helper->deleteLater();
|
|
m_helper = nullptr;
|
|
}
|
|
|
|
void emitEvent(const QString& eventName, const QVariantList& data) override
|
|
{
|
|
if (!m_replica) return;
|
|
qDebug() << "[LogosObject] RemoteLogosObject::emitEvent" << eventName << "data:" << data.size() << "items (via IPC)";
|
|
QMetaObject::invokeMethod(m_replica, "eventResponse",
|
|
Qt::QueuedConnection,
|
|
Q_ARG(QString, eventName),
|
|
Q_ARG(QVariantList, data));
|
|
}
|
|
|
|
QJsonArray getMethods() override
|
|
{
|
|
// Remote introspection not implemented — callers should use
|
|
// the local module inspection tools (lm) instead.
|
|
return QJsonArray();
|
|
}
|
|
|
|
void release() override
|
|
{
|
|
disconnectEvents(); // defers the helper (signal receiver)
|
|
// The replica is the QtRO signal *sender* whose eventResponse() may be the
|
|
// very emission that re-entered release() (a deferred completion event).
|
|
// Deleting the sender mid-emit corrupts QtRO's read path
|
|
// (QRemoteObjectNodePrivate::onClientRead) — defer it too. deleteLater also
|
|
// keeps any QVariant return storage backed by the replica alive until the
|
|
// consumer callback (which runs before this release) has consumed it.
|
|
if (m_replica) {
|
|
m_replica->deleteLater();
|
|
m_replica = nullptr;
|
|
}
|
|
delete this;
|
|
}
|
|
|
|
quintptr id() const override { return reinterpret_cast<quintptr>(m_replica); }
|
|
|
|
// Valid only while the underlying replica is synced to its source. When the
|
|
// target module unloads, the source drops and state() leaves Valid, so a
|
|
// cached handle knows to re-acquire instead of calling into a dead replica.
|
|
bool isValid() const override
|
|
{
|
|
auto* r = qobject_cast<QRemoteObjectReplica*>(m_replica);
|
|
return r && r->state() == QRemoteObjectReplica::Valid;
|
|
}
|
|
|
|
private:
|
|
// Resolve a possibly-deferred result. If `rv` is a pending sentinel from a
|
|
// "multi" provider, wait (up to timeoutMs) for the completion event keyed by
|
|
// callId, pumping the consumer event loop; otherwise return `rv` unchanged.
|
|
QVariant resolveDeferred(const QVariant& rv, int timeoutMs,
|
|
const QString& methodName = QString(),
|
|
logos::CallError* err = nullptr)
|
|
{
|
|
QString callId;
|
|
if (!logos::isPendingCallSentinel(rv, &callId)) return rv;
|
|
// The completion can arrive during the call's own waitForFinished above.
|
|
if (m_completions.contains(callId)) return m_completions.take(callId);
|
|
QEventLoop loop;
|
|
m_completionWaiters.insert(callId, &loop);
|
|
QTimer timer;
|
|
timer.setSingleShot(true);
|
|
QObject::connect(&timer, &QTimer::timeout, &loop, &QEventLoop::quit);
|
|
timer.start(timeoutMs > 0 ? timeoutMs : 30000);
|
|
loop.exec();
|
|
m_completionWaiters.remove(callId);
|
|
if (m_completions.contains(callId)) return m_completions.take(callId);
|
|
qWarning() << "RemoteLogosObject: deferred call" << callId << "timed out";
|
|
if (err)
|
|
*err = logos::callErrorTimeout(m_objectName.toStdString(),
|
|
methodName.toStdString(),
|
|
timeoutMs > 0 ? timeoutMs : 30000);
|
|
return QVariant();
|
|
}
|
|
|
|
QObject* m_replica;
|
|
RemoteEventHelper* m_helper;
|
|
// Deferred ("multi") completions delivered over the event channel, keyed by
|
|
// callId: buffered results + sync (QEventLoop) and async (callback) waiters.
|
|
// Touched only on the consumer event-loop thread.
|
|
QHash<QString, QVariant> m_completions;
|
|
QHash<QString, QEventLoop*> m_completionWaiters;
|
|
QHash<QString, AsyncResultErrorCallback> m_asyncCompletionCbs;
|
|
QString m_objectName;
|
|
};
|
|
|
|
// ── PendingAcquire ───────────────────────────────────────────────────────────
|
|
|
|
namespace {
|
|
|
|
// One outstanding non-blocking acquire.
|
|
//
|
|
// Acquiring a dynamic replica before the peer exists is legal and free:
|
|
// QRemoteObjectNodePrivate::handleNewAcquire registers the replica in the
|
|
// node's `replicas` map with no connection, and when the peer finally shows up
|
|
// the ObjectList packet handler calls handleReplicaConnection for every waiting
|
|
// replica, which drives the replica to Valid. There is no polling here — the
|
|
// node's own 250 ms reconnect timer (onShouldReconnect) is running for this
|
|
// endpoint regardless of whether anyone is waiting on it.
|
|
//
|
|
// Deletes itself after delivering exactly once. If it is destroyed before
|
|
// delivering (transport torn down), it owns and deletes the replica.
|
|
class PendingAcquire : public QObject {
|
|
Q_OBJECT
|
|
public:
|
|
PendingAcquire(QRemoteObjectReplica* replica,
|
|
const QString& objectName,
|
|
LogosTransportAsyncAcquire::AcquireCallback cb,
|
|
QObject* parent)
|
|
: QObject(parent)
|
|
, m_replica(replica)
|
|
, m_objectName(objectName)
|
|
, m_cb(std::move(cb))
|
|
{
|
|
connect(m_replica, &QRemoteObjectReplica::stateChanged,
|
|
this, &PendingAcquire::onStateChanged);
|
|
// Already live (module was up before we asked) — still deliver via the
|
|
// event loop, so the callback contract is "never synchronously".
|
|
if (m_replica->state() == QRemoteObjectReplica::Valid)
|
|
schedule(true);
|
|
}
|
|
|
|
~PendingAcquire() override
|
|
{
|
|
// Only reached when we never handed the replica over.
|
|
delete m_replica;
|
|
}
|
|
|
|
private slots:
|
|
void onStateChanged(QRemoteObjectReplica::State state,
|
|
QRemoteObjectReplica::State /*old*/)
|
|
{
|
|
if (m_done) return;
|
|
if (state == QRemoteObjectReplica::Valid) {
|
|
schedule(true);
|
|
} else if (state == QRemoteObjectReplica::SignatureMismatch) {
|
|
// Terminal: handleReplicaConnection warns and never connects, and
|
|
// no further state change will arrive. Report it rather than
|
|
// sitting Pending forever in silence — a subscription that can
|
|
// never arm has to SAY so.
|
|
qWarning() << "RemoteTransportConnection: source signature mismatch for"
|
|
<< m_objectName << "-- this acquire can never succeed";
|
|
schedule(false);
|
|
}
|
|
// Suspect / Default are NOT terminal: the source can come back and the
|
|
// node will re-drive the replica to Valid.
|
|
}
|
|
|
|
private:
|
|
void schedule(bool ok)
|
|
{
|
|
if (m_done) return;
|
|
m_done = true;
|
|
// NEVER deliver inline. stateChanged is emitted from inside
|
|
// QRemoteObjectNodePrivate::onClientRead; constructing a
|
|
// RemoteLogosObject (which connects to the same replica) and then
|
|
// running caller code there re-enters QtRO while the replica is still
|
|
// emitting. That is precisely the refresh_balances SIGSEGV documented
|
|
// on RemoteLogosObject::disconnectEvents(). Hop to the next event-loop
|
|
// turn; `this` as context cancels the call if we are destroyed first.
|
|
QTimer::singleShot(0, this, [this, ok]() { deliver(ok); });
|
|
}
|
|
|
|
void deliver(bool ok)
|
|
{
|
|
QRemoteObjectReplica* replica = m_replica;
|
|
m_replica = nullptr; // ownership leaves us either way
|
|
disconnect(replica, nullptr, this, nullptr);
|
|
|
|
LogosObject* obj = nullptr;
|
|
if (ok) {
|
|
g_acquireCount.fetch_add(1, std::memory_order_relaxed);
|
|
obj = new RemoteLogosObject(replica, m_objectName);
|
|
} else {
|
|
delete replica;
|
|
}
|
|
auto cb = std::move(m_cb);
|
|
deleteLater();
|
|
cb(obj); // may re-enter the transport safely
|
|
}
|
|
|
|
QRemoteObjectReplica* m_replica;
|
|
QString m_objectName;
|
|
LogosTransportAsyncAcquire::AcquireCallback m_cb;
|
|
bool m_done = false;
|
|
};
|
|
|
|
} // anonymous namespace
|
|
|
|
// ── RemoteTransportHost ──────────────────────────────────────────────────────
|
|
|
|
RemoteTransportHost::RemoteTransportHost(const QString& registryUrl)
|
|
: m_registryHost(nullptr)
|
|
, m_registryUrl(registryUrl)
|
|
{
|
|
}
|
|
|
|
RemoteTransportHost::~RemoteTransportHost()
|
|
{
|
|
delete m_registryHost;
|
|
}
|
|
|
|
bool RemoteTransportHost::publishObject(const QString& name, QObject* object)
|
|
{
|
|
if (!m_registryHost) {
|
|
// Construct WITHOUT a URL and listen via setRegistryUrl() so we can
|
|
// observe the result. The QUrl-taking ctor swallows the listen bool: on
|
|
// a failed bind (stale socket, permission denied, sun_path overflow) it
|
|
// leaves a half-built host that silently rejects every enableRemoting()
|
|
// with OperationNotValidOnClientNode — the failure only surfaced as
|
|
// clients hanging on a dead endpoint.
|
|
auto* host = new QRemoteObjectRegistryHost();
|
|
if (!host->setRegistryUrl(QUrl(m_registryUrl))) {
|
|
qCritical() << "RemoteTransportHost: failed to listen on" << m_registryUrl
|
|
<< "- error" << host->lastError()
|
|
<< "- socket path" << localSocketFilePath(m_registryUrl);
|
|
delete host;
|
|
return false;
|
|
}
|
|
m_registryHost = host;
|
|
|
|
// Make the freshly-bound socket reachable by a co-resident client per
|
|
// the LOGOS_SOCKET_GROUP / LOGOS_SOCKET_MODE policy. No-op when those
|
|
// env vars are unset (socket keeps its owner-only default). Only `local:`
|
|
// registries have a socket file — a tcp:// one does not.
|
|
if (QUrl(m_registryUrl).scheme() == QLatin1String("local")) {
|
|
std::string permErr;
|
|
if (!logos::applySocketPerms(
|
|
localSocketFilePath(m_registryUrl).toStdString(), &permErr)) {
|
|
qWarning() << "RemoteTransportHost: could not apply socket perms:"
|
|
<< QString::fromStdString(permErr);
|
|
}
|
|
}
|
|
qDebug() << "RemoteTransportHost: Created registry host with URL:" << m_registryUrl;
|
|
}
|
|
|
|
bool success = m_registryHost->enableRemoting(object, name);
|
|
if (success) {
|
|
qDebug() << "RemoteTransportHost: Published object:" << name;
|
|
} else {
|
|
qCritical() << "RemoteTransportHost: Failed to publish object:" << name;
|
|
}
|
|
return success;
|
|
}
|
|
|
|
void RemoteTransportHost::unpublishObject(const QString& /*name*/)
|
|
{
|
|
}
|
|
|
|
// ── RemoteTransportConnection ────────────────────────────────────────────────
|
|
|
|
RemoteTransportConnection::RemoteTransportConnection(const QString& registryUrl)
|
|
: m_node(new QRemoteObjectNode())
|
|
, m_pendingAcquires(new QObject())
|
|
, m_registryUrl(registryUrl)
|
|
, m_connected(false)
|
|
{
|
|
}
|
|
|
|
RemoteTransportConnection::~RemoteTransportConnection()
|
|
{
|
|
// Pending replicas and parked probes belong to the node; kill them FIRST so
|
|
// none outlives it.
|
|
m_probes.clear();
|
|
delete m_pendingAcquires;
|
|
m_pendingAcquires = nullptr;
|
|
delete m_node;
|
|
}
|
|
|
|
bool RemoteTransportConnection::connectToHost()
|
|
{
|
|
return connectToRegistry();
|
|
}
|
|
|
|
bool RemoteTransportConnection::endpointHasListener() const
|
|
{
|
|
const QUrl url(m_registryUrl);
|
|
// Only `local:` can be probed cheaply and without side effects worth
|
|
// worrying about. Anything else keeps the previous semantics.
|
|
if (url.scheme() != QLatin1String("local"))
|
|
return true;
|
|
|
|
const QString serverName = url.path().isEmpty() ? url.host() : url.path();
|
|
if (serverName.isEmpty())
|
|
return true;
|
|
|
|
// A client connect is exactly what QtRO itself does, so a transient one is
|
|
// benign against a live registry, and against a dead endpoint it fails
|
|
// immediately -- ERROR_FILE_NOT_FOUND on Windows, ENOENT/ECONNREFUSED on
|
|
// Unix. That is the whole point: cost is microseconds when the answer is
|
|
// "no", which is the case that currently costs 20 seconds.
|
|
QLocalSocket probe;
|
|
probe.connectToServer(serverName, QIODevice::ReadOnly);
|
|
const bool alive = probe.waitForConnected(250);
|
|
probe.abort();
|
|
return alive;
|
|
}
|
|
|
|
bool RemoteTransportConnection::isConnected() const
|
|
{
|
|
// m_connected alone is NOT an answer. QRemoteObjectNode::connectToNode()
|
|
// returns false only when the URL scheme is unregistered -- it never
|
|
// contacts the peer -- and our registry URLs are COMPUTED rather than
|
|
// discovered (logos_instance.h: local:logos_<module>_<instanceId>), so they
|
|
// are identical whether or not the module exists. Returning the raw latch
|
|
// made every `if (!client->isConnected()) return;` guard in the codebase
|
|
// dead code, and callers then paid a 20 s waitForSource per call against
|
|
// modules that were never loaded -- measured at ~417 s of blocked GUI
|
|
// thread in Basecamp on macOS, 361 s on Linux, before its window appeared.
|
|
if (!m_connected)
|
|
return false;
|
|
return endpointHasListener();
|
|
}
|
|
|
|
bool RemoteTransportConnection::reconnect()
|
|
{
|
|
qDebug() << "RemoteTransportConnection: Attempting to reconnect to registry:" << m_registryUrl;
|
|
|
|
if (m_connected) {
|
|
// In-flight acquires hold replicas belonging to the node about to die.
|
|
// Drop them first; their callbacks are cancelled with them, so the
|
|
// consumer's subscription registry re-arms against the new node
|
|
// (LogosAPIConsumer::reconnect -> reconnected()).
|
|
m_probes.clear();
|
|
delete m_pendingAcquires;
|
|
m_pendingAcquires = new QObject();
|
|
delete m_node;
|
|
m_node = new QRemoteObjectNode();
|
|
m_connected = false;
|
|
}
|
|
|
|
return connectToRegistry();
|
|
}
|
|
|
|
bool RemoteTransportConnection::connectToRegistry()
|
|
{
|
|
if (!m_node) {
|
|
qWarning() << "RemoteTransportConnection: Remote object node is null";
|
|
return false;
|
|
}
|
|
|
|
if (m_registryUrl.isEmpty()) {
|
|
qWarning() << "RemoteTransportConnection: Registry URL is empty";
|
|
return false;
|
|
}
|
|
|
|
qDebug() << "RemoteTransportConnection: Connecting to registry:" << m_registryUrl
|
|
<< "at" << QTime::currentTime().toString("hh:mm:ss.zzz");
|
|
|
|
QUrl url(m_registryUrl);
|
|
bool success = m_node->connectToNode(url);
|
|
|
|
if (success) {
|
|
m_connected = true;
|
|
// Deliberately NOT "Successfully connected". connectToNode() only
|
|
// accepted the URL scheme; it never contacted a peer, and this line
|
|
// previously asserted a connection that frequently did not exist. That
|
|
// false claim sent three separate investigations to the wrong place --
|
|
// it cost considerably more than the bug it hid.
|
|
qDebug() << "RemoteTransportConnection: Registry connect attempt started (no peer contact yet):"
|
|
<< m_registryUrl;
|
|
} else {
|
|
m_connected = false;
|
|
qWarning() << "RemoteTransportConnection: Registry URL scheme rejected:" << m_registryUrl;
|
|
}
|
|
|
|
return m_connected;
|
|
}
|
|
|
|
LogosObject* RemoteTransportConnection::requestObject(const QString& objectName, int timeoutMs)
|
|
{
|
|
if (!m_connected) {
|
|
qWarning() << "RemoteTransportConnection: Not connected. Cannot request object:" << objectName;
|
|
return nullptr;
|
|
}
|
|
|
|
// Warn BEFORE a doomed wait rather than after it. Without this the log went
|
|
// silent for the full timeout and then reported failure, which reads as a
|
|
// hang with no cause attached.
|
|
if (!endpointHasListener()) {
|
|
qWarning() << "RemoteTransportConnection: no listener at" << m_registryUrl
|
|
<< "-- request for" << objectName << "will block up to" << timeoutMs
|
|
<< "ms and then fail. Is the module loaded?";
|
|
}
|
|
|
|
qDebug() << "RemoteTransportConnection: Requesting object:" << objectName
|
|
<< "at" << QTime::currentTime().toString("hh:mm:ss.zzz");
|
|
|
|
QRemoteObjectReplica* replica = m_node->acquireDynamic(objectName);
|
|
if (!replica) {
|
|
qWarning() << "RemoteTransportConnection: Failed to acquire replica for:" << objectName;
|
|
return nullptr;
|
|
}
|
|
|
|
if (!replica->waitForSource(timeoutMs)) {
|
|
qWarning() << "RemoteTransportConnection: Timeout waiting for replica:" << objectName;
|
|
delete replica;
|
|
return nullptr;
|
|
}
|
|
|
|
qDebug() << "[LogosObject] RemoteTransportConnection: returning RemoteLogosObject for:" << objectName;
|
|
g_acquireCount.fetch_add(1, std::memory_order_relaxed);
|
|
return new RemoteLogosObject(replica, objectName);
|
|
}
|
|
|
|
bool RemoteTransportConnection::requestObjectWhenAvailable(const QString& objectName,
|
|
AcquireCallback onReady)
|
|
{
|
|
if (objectName.isEmpty() || !onReady) return false;
|
|
if (!m_node || !m_pendingAcquires) return false;
|
|
// m_connected is the raw "connectToNode() accepted the URL" latch, and here
|
|
// that is exactly the right question: we are NOT asking whether the peer is
|
|
// up (it deliberately need not be), only whether this node is wired to the
|
|
// endpoint so QtRO's reconnect loop is running for it. Calling the truthful
|
|
// isConnected() here would reintroduce the very bug this exists to fix.
|
|
if (!m_connected) return false;
|
|
|
|
QRemoteObjectReplica* replica = m_node->acquireDynamic(objectName);
|
|
if (!replica) {
|
|
qWarning() << "RemoteTransportConnection: Failed to acquire replica for:" << objectName;
|
|
return false;
|
|
}
|
|
|
|
new PendingAcquire(replica, objectName, std::move(onReady), m_pendingAcquires);
|
|
return true;
|
|
}
|
|
|
|
LogosObject* RemoteTransportConnection::tryAcquireNow(const QString& objectName)
|
|
{
|
|
if (objectName.isEmpty() || !m_node) return nullptr;
|
|
// Raw latch, same reasoning as requestObjectWhenAvailable: the question is
|
|
// whether this node is wired to the endpoint, not whether the peer is up.
|
|
if (!m_connected) return nullptr;
|
|
|
|
// PARK the probe; never free it while it is unconfigured. QtRO shares one
|
|
// replica implementation per object name on a node, and while that
|
|
// implementation is still waiting for the source's metaobject it records
|
|
// every facade built on it as a RAW pointer in m_parentsNeedingConnect.
|
|
// ~QRemoteObjectReplica is an empty body, so destroying a facade does not
|
|
// deregister it — and the implementation dereferences the whole list when
|
|
// the class definition finally arrives.
|
|
//
|
|
// The earlier acquire-then-delete form was therefore a use-after-free that
|
|
// left one dangling pointer per probe. It survived a single subscription
|
|
// (the probe owned the only implementation and took it down with itself)
|
|
// and crashed once a second subscription shared an implementation pinned by
|
|
// an in-flight PendingAcquire: SIGBUS with BUS_ADRALN, in the consumer's
|
|
// event loop rather than at the call site.
|
|
//
|
|
// Parking costs one idle replica per name until it goes Valid or the
|
|
// connection dies. Parented to m_pendingAcquires, which both ~RemoteTransportConnection
|
|
// and reconnect() destroy BEFORE the node — the ordering is what makes
|
|
// freeing them safe, since the implementations die in the same breath.
|
|
QPointer<QRemoteObjectReplica>& probe = m_probes[objectName];
|
|
if (!probe) {
|
|
QRemoteObjectReplica* fresh = m_node->acquireDynamic(objectName);
|
|
if (!fresh) {
|
|
m_probes.remove(objectName);
|
|
return nullptr;
|
|
}
|
|
if (m_pendingAcquires) fresh->setParent(m_pendingAcquires);
|
|
probe = fresh;
|
|
}
|
|
|
|
// Not Valid means "would have to be waited on", and waiting is exactly what
|
|
// this must not do — requestObjectWhenAvailable owns that. Leave the probe
|
|
// parked and answer no.
|
|
if (probe->state() != QRemoteObjectReplica::Valid)
|
|
return nullptr;
|
|
|
|
// Valid: the implementation is configured, so it is no longer holding this
|
|
// facade in m_parentsNeedingConnect and handing ownership over is safe.
|
|
QRemoteObjectReplica* replica = probe;
|
|
m_probes.remove(objectName);
|
|
replica->setParent(nullptr);
|
|
|
|
g_acquireCount.fetch_add(1, std::memory_order_relaxed);
|
|
return new RemoteLogosObject(replica, objectName);
|
|
}
|
|
|
|
long RemoteTransportConnection::acquireCount() { return g_acquireCount.load(std::memory_order_relaxed); }
|
|
void RemoteTransportConnection::resetAcquireCount() { g_acquireCount.store(0, std::memory_order_relaxed); }
|
|
|
|
#include "remote_transport.moc"
|