mirror of
https://github.com/logos-co/logos-protocol.git
synced 2026-08-27 20:11:07 +00:00
258 lines
11 KiB
C++
258 lines
11 KiB
C++
#include "logos_api_consumer.h"
|
||
#include "logos_object.h"
|
||
#include "module_proxy.h"
|
||
#include "token_manager.h"
|
||
#include "logos_mode.h"
|
||
#include "logos_instance.h"
|
||
#include "logos_transport.h"
|
||
#include "logos_transport_factory.h"
|
||
#include <chrono>
|
||
#include <thread>
|
||
#include <QDebug>
|
||
#include <QUrl>
|
||
#include <QMetaObject>
|
||
#include <QTimer>
|
||
#include <QTime>
|
||
#include <QPointer>
|
||
|
||
LogosAPIConsumer::LogosAPIConsumer(const QString& module_to_talk_to,
|
||
const QString& origin_module,
|
||
TokenManager* token_manager,
|
||
const LogosTransportConfig& transport,
|
||
QObject *parent)
|
||
: QObject(parent)
|
||
, m_registryUrl(LogosInstance::id(module_to_talk_to))
|
||
, m_token_manager(token_manager)
|
||
{
|
||
// Single transport-resolution path: the factory combines LogosMode
|
||
// + LogosTransportConfig (mode wins for Mock/Local; transport
|
||
// chooses the wire protocol in Remote mode). The choice scopes to
|
||
// this consumer only — any LogosAPIProvider in the same LogosAPI
|
||
// still constructs its host from the global default.
|
||
m_transport = LogosTransportFactory::createConnection(transport, m_registryUrl);
|
||
|
||
// Initial connect with deadline-driven retry. The target module's
|
||
// listener may not be ready yet — particularly for TCP/TLS, where
|
||
// the child subprocess's QTcpServer::listen() lags the runtime
|
||
// returning from its load callback. QLocalSocket internally
|
||
// tolerates this (it retries connect until a deadline), but
|
||
// boost::asio::connect on TCP fails fast with "connection refused"
|
||
// and we'd surface a warning + return nullptr for any subsequent
|
||
// requestObject before the listener even came up.
|
||
//
|
||
// 50ms × up-to-100 attempts ≈ 5s budget — same shape as
|
||
// logos-liblogos's sendTokenToProcess loop and generous enough to
|
||
// cover cold-start child Qt initialisation under load.
|
||
using clock = std::chrono::steady_clock;
|
||
const auto deadline = clock::now() + std::chrono::milliseconds(5000);
|
||
while (true) {
|
||
if (m_transport->connectToHost()) break;
|
||
if (clock::now() >= deadline) break;
|
||
std::this_thread::sleep_for(std::chrono::milliseconds(50));
|
||
}
|
||
}
|
||
|
||
LogosAPIConsumer::LogosAPIConsumer(const QString& module_to_talk_to,
|
||
const QString& origin_module,
|
||
TokenManager* token_manager,
|
||
QObject *parent)
|
||
: LogosAPIConsumer(module_to_talk_to, origin_module, token_manager,
|
||
LogosTransportConfigGlobal::getDefault(), parent)
|
||
{
|
||
}
|
||
|
||
LogosAPIConsumer::~LogosAPIConsumer()
|
||
{
|
||
}
|
||
|
||
LogosObject* LogosAPIConsumer::requestObject(const QString& objectName, Timeout timeout)
|
||
{
|
||
qDebug() << "LogosAPIConsumer: Requesting object:" << objectName << "at" << QTime::currentTime().toString("hh:mm:ss.zzz");
|
||
|
||
if (objectName.isEmpty()) {
|
||
qWarning() << "LogosAPIConsumer: Object name cannot be empty";
|
||
return nullptr;
|
||
}
|
||
|
||
if (!m_transport->isConnected()) {
|
||
qWarning() << "LogosAPIConsumer: Not connected to registry. Cannot request object:" << objectName;
|
||
return nullptr;
|
||
}
|
||
|
||
LogosObject* object = m_transport->requestObject(objectName, timeout.ms);
|
||
if (object) {
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: acquired LogosObject for:" << objectName << "(id:" << object->id() << ")";
|
||
}
|
||
return object;
|
||
}
|
||
|
||
bool LogosAPIConsumer::isConnected() const
|
||
{
|
||
return m_transport->isConnected();
|
||
}
|
||
|
||
QString LogosAPIConsumer::registryUrl() const
|
||
{
|
||
return m_registryUrl;
|
||
}
|
||
|
||
bool LogosAPIConsumer::reconnect()
|
||
{
|
||
qDebug() << "LogosAPIConsumer: Attempting to reconnect to registry:" << m_registryUrl;
|
||
return m_transport->reconnect();
|
||
}
|
||
|
||
QVariant LogosAPIConsumer::invokeRemoteMethod(const QString& authToken, const QString& objectName, const QString& methodName,
|
||
const QVariantList& args, Timeout timeout)
|
||
{
|
||
return invokeRemoteMethod(authToken, objectName, methodName, args, timeout, nullptr);
|
||
}
|
||
|
||
QVariant LogosAPIConsumer::invokeRemoteMethod(const QString& authToken, const QString& objectName, const QString& methodName,
|
||
const QVariantList& args, Timeout timeout, logos::CallError* err)
|
||
{
|
||
if (err) err->clear();
|
||
qDebug() << "LogosAPIConsumer: Calling invokeRemoteMethod:" << objectName << methodName << "args_count:" << args.size() << "timeout:" << timeout.ms;
|
||
|
||
LogosObject* plugin = m_transport->requestObject(objectName, timeout.ms);
|
||
if (!plugin) {
|
||
qWarning() << "LogosAPIConsumer: Failed to acquire plugin/replica for object:" << objectName;
|
||
if (err) {
|
||
err->code = "object_unavailable";
|
||
err->message = "failed to acquire remote object '"
|
||
+ objectName.toStdString()
|
||
+ "' (module not loaded, not published, or transport failure)";
|
||
err->origin = objectName.toStdString();
|
||
}
|
||
return QVariant();
|
||
}
|
||
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: calling via LogosObject::callMethod" << methodName;
|
||
QVariant result = plugin->callMethod(authToken, methodName, args, timeout.ms);
|
||
plugin->release();
|
||
return result;
|
||
}
|
||
|
||
void LogosAPIConsumer::invokeRemoteMethodAsync(const QString& authToken, const QString& objectName, const QString& methodName,
|
||
const QVariantList& args,
|
||
AsyncResultCallback callback,
|
||
Timeout timeout)
|
||
{
|
||
// Delegate to the CallError-aware overload so there is one acquire/dispatch
|
||
// path — the legacy callback simply drops the error field.
|
||
invokeRemoteMethodAsync(authToken, objectName, methodName, args,
|
||
[cb = std::move(callback)](QVariant r, const logos::CallError&) mutable {
|
||
if (cb) cb(std::move(r));
|
||
},
|
||
timeout);
|
||
}
|
||
|
||
void LogosAPIConsumer::invokeRemoteMethodAsync(const QString& authToken, const QString& objectName, const QString& methodName,
|
||
const QVariantList& args,
|
||
AsyncResultErrorCallback callback,
|
||
Timeout timeout)
|
||
{
|
||
if (!callback) {
|
||
qWarning() << "LogosAPIConsumer: invokeRemoteMethodAsync called with null callback";
|
||
return;
|
||
}
|
||
|
||
LogosObject* plugin = m_transport->requestObject(objectName, timeout.ms);
|
||
if (!plugin) {
|
||
qWarning() << "LogosAPIConsumer: Failed to acquire plugin/replica for object:" << objectName;
|
||
logos::CallError err;
|
||
err.code = "object_unavailable";
|
||
err.message = "failed to acquire remote object '"
|
||
+ objectName.toStdString()
|
||
+ "' (module not loaded, not published, or transport failure)";
|
||
err.origin = objectName.toStdString();
|
||
QTimer::singleShot(0, this, [callback, err]() { callback(QVariant(), err); });
|
||
return;
|
||
}
|
||
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: async calling via LogosObject::callMethodAsync" << methodName;
|
||
// QPointer guards against use-after-free: if the consumer is destroyed
|
||
// before the transport callback fires, the callback is silently dropped.
|
||
// No re-queuing needed -- the transport already delivers on a deferred
|
||
// event-loop iteration (QTimer / QueuedConnection).
|
||
QPointer<LogosAPIConsumer> self(this);
|
||
plugin->callMethodAsync(authToken, methodName, args, timeout.ms,
|
||
[plugin, callback, self](QVariant result) {
|
||
// Deliver the result before release(): destroying the replica can
|
||
// invalidate storage that QVariant still references for some return
|
||
// types (matches sync invokeRemoteMethod: callMethod then release).
|
||
if (!self) {
|
||
plugin->release();
|
||
return;
|
||
}
|
||
callback(result, logos::CallError{});
|
||
plugin->release();
|
||
});
|
||
}
|
||
|
||
void LogosAPIConsumer::onEvent(LogosObject* originObject, const QString& eventName, std::function<void(const QString&, const QVariantList&)> callback)
|
||
{
|
||
qDebug() << "[LogosObject] LogosAPIConsumer::onEvent registering for:" << eventName << "on LogosObject id:" << originObject;
|
||
|
||
if (!originObject) {
|
||
qWarning() << "LogosAPIConsumer: Cannot register event on null object";
|
||
return;
|
||
}
|
||
|
||
originObject->onEvent(eventName, std::move(callback));
|
||
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: event callback registered for:" << eventName;
|
||
}
|
||
|
||
bool LogosAPIConsumer::informModuleToken(const QString& authToken, const QString& moduleName, const QString& token)
|
||
{
|
||
qDebug() << "LogosAPIConsumer: Informing module token for module:" << moduleName << "with token:" << redactToken(token);
|
||
|
||
LogosObject* plugin = m_transport->requestObject("capability_module", 20000);
|
||
if (!plugin) {
|
||
qWarning() << "LogosAPIConsumer: Failed to acquire plugin/replica for object: capability_module";
|
||
return false;
|
||
}
|
||
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: calling LogosObject::informModuleToken for" << moduleName;
|
||
bool result = plugin->informModuleToken(authToken, moduleName, token, 20000);
|
||
qDebug() << "LogosAPIConsumer: informModuleToken completed with result:" << result;
|
||
plugin->release();
|
||
return result;
|
||
}
|
||
|
||
bool LogosAPIConsumer::informModuleToken_module(const QString& authToken, const QString& originModule, const QString& moduleName, const QString& token)
|
||
{
|
||
qDebug() << "LogosAPIConsumer: Informing module token for module:" << moduleName << "with token:" << redactToken(token);
|
||
|
||
LogosObject* plugin = m_transport->requestObject(originModule, 20000);
|
||
if (!plugin) {
|
||
qWarning() << "LogosAPIConsumer: Failed to acquire plugin/replica for object:" << originModule;
|
||
return false;
|
||
}
|
||
|
||
qDebug() << "[LogosObject] LogosAPIConsumer: calling LogosObject::informModuleToken for" << moduleName << "on" << originModule;
|
||
bool result = plugin->informModuleToken(authToken, moduleName, token, 20000);
|
||
qDebug() << "LogosAPIConsumer: informModuleToken completed with result:" << result;
|
||
plugin->release();
|
||
return result;
|
||
}
|
||
|
||
std::string LogosAPIConsumer::requestModule(const std::string& authToken, const std::string& originModule, const std::string& targetModule)
|
||
{
|
||
const QString qOrigin = QString::fromStdString(originModule);
|
||
const QString qTarget = QString::fromStdString(targetModule);
|
||
qDebug() << "LogosAPIConsumer: requestModule for origin:" << qOrigin << "target:" << qTarget;
|
||
|
||
LogosObject* plugin = m_transport->requestObject("capability_module", 20000);
|
||
if (!plugin) {
|
||
qWarning() << "LogosAPIConsumer: Failed to acquire plugin/replica for object: capability_module";
|
||
return {};
|
||
}
|
||
|
||
QVariant result = plugin->callMethod(QString::fromStdString(authToken), QStringLiteral("requestModule"),
|
||
QVariantList() << qOrigin << qTarget, 20000);
|
||
plugin->release();
|
||
return result.toString().toStdString();
|
||
}
|