Page MenuHomePhorge

No OneTemporary

Size
63 KB
Referenced Files
None
Subscribers
None
diff --git a/src/db-store.cpp b/src/db-store.cpp
index c8ca77d..4e98d58 100644
--- a/src/db-store.cpp
+++ b/src/db-store.cpp
@@ -1,211 +1,299 @@
/*
* This file is part of kazv.
* SPDX-FileCopyrightText: 2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include "db-store.hpp"
#include "matrix-sdk.hpp"
#include "event.hpp"
#include "kazv-log.hpp"
#include <QPromise>
#include <QSqlError>
#include <QSqlQuery>
+#include <QSqlRecord>
#include <QCoroFuture>
using namespace Qt::Literals::StringLiterals;
using namespace Kazv;
QString DbStore::getHandleFor(std::string userId, std::string deviceId)
{
return QString::fromStdString(userId + "/" + deviceId);
}
+std::optional<Kazv::Event> DbStore::resultToEvent(bool encrypted, bool decrypted, QString originalJson, QString decryptedJson)
+{
+ Event e;
+ try {
+ auto oj = json::parse(std::move(originalJson).toStdString());
+ e = Event(JsonWrap(oj));
+ } catch (const json::parse_error &) {
+ return std::nullopt;
+ }
+ try {
+ if (encrypted) {
+ auto dj = json::parse(std::move(decryptedJson).toStdString());
+ e = std::move(e).setDecryptedJson(
+ JsonWrap(dj),
+ decrypted ? Event::Decrypted : Event::NotDecrypted
+ );
+ }
+ return e;
+ } catch (const json::parse_error &) {
+ return e;
+ }
+}
+
DbStore::DbStore()
: m_thread(new QThread)
, m_obj(new QObject)
, m_d(std::nullopt)
, m_handle()
{
m_obj->moveToThread(m_thread);
}
DbStore::~DbStore()
{
cleanup();
}
QCoro::Task<std::pair<bool, QString>> DbStore::setup(std::string userDataDir, std::string userId, std::string deviceId)
{
m_handle = getHandleFor(userId, deviceId);
m_d = QSqlDatabase::addDatabase(u"QSQLITE"_s, m_handle);
auto dbFullDir = sessionDirForUserAndDeviceId(userDataDir, userId, deviceId) / dbDirName;
std::error_code err;
if ((! std::filesystem::create_directories(dbFullDir, err))
&& err) {
qCWarning(kazvLog) << "DbStore: setup: cannot create db dir:" << dbFullDir.string() << ", error:" << err.message();
co_return {false, u"Cannot create db dir"_s};
}
auto dbName = QString::fromStdString((dbFullDir / dbBaseName).string());
qCDebug(kazvLog) << "DbStore: setup: data file name" << dbName;
m_d->setDatabaseName(dbName);
bool status = m_d->open();
if (!status) {
qCWarning(kazvLog) << "DbStore: Cannot open database";
std::pair<bool, QString> res = {status, m_d->lastError().text()};
m_d.reset();
QSqlDatabase::removeDatabase(m_handle);
co_return res;
}
qCDebug(kazvLog) << "DbStore: Opened database";
m_d->moveToThread(m_thread);
m_thread->start();
qCDebug(kazvLog) << "DbStore: Running database migrations";
auto res = co_await migrate();
co_return res;
}
QCoro::Task<std::pair<bool, QString>> DbStore::importAllFrom(const Kazv::SdkModel &model)
{
std::size_t count = 0;
for (const auto &[roomId, room] : model.c().roomList.rooms) {
auto roomIdQs = QString::fromStdString(roomId);
for (const auto &[eventId, event] : room.messages) {
auto res = co_await importOne(roomIdQs, event, /* isInTimeline = */ true);
if (!res.first) {
qCWarning(kazvLog) << "Error inserting value:" << res.second;
co_return res;
}
++count;
if (count % 100 == 0) {
qCInfo(kazvLog) << "Imported" << count << "events";
}
}
}
co_return {true, QString()};
}
+QCoro::Task<std::pair<bool, QString>> DbStore::saveEvents(Kazv::SaveEventsRequested s)
+{
+ for (const auto &[roomId, el] : s.timelineEvents) {
+ auto roomIdQs = QString::fromStdString(roomId);
+ for (const auto &event : el) {
+ auto res = co_await importOne(roomIdQs, event, /* isInTimeline = */ true);
+ if (!res.first) {
+ qCWarning(kazvLog) << "Error inserting value:" << res.second;
+ co_return res;
+ }
+ }
+ }
+
+ for (const auto &[roomId, el] : s.nonTimelineEvents) {
+ auto roomIdQs = QString::fromStdString(roomId);
+ for (const auto &event : el) {
+ auto res = co_await importOne(roomIdQs, event, /* isInTimeline = */ false);
+ if (!res.first) {
+ qCWarning(kazvLog) << "Error inserting value:" << res.second;
+ co_return res;
+ }
+ }
+ }
+ co_return {true, QString()};
+}
+
+inline const auto getEventByIdQuery = uR"xxx(SELECT
+ encrypted, decrypted,
+ original_json, decrypted_json, is_in_timeline
+FROM events WHERE
+ room_id = :room_id AND event_id = :event_id
+;)xxx"_s;
+QCoro::Task<std::optional<std::pair<Kazv::Event, bool>>> DbStore::getEventById(const QString &roomId, const QString &eventId)
+{
+ QPromise<std::optional<std::pair<Kazv::Event, bool>>> p;
+ auto fut = p.future();
+ QMetaObject::invokeMethod(m_obj, [this, p=std::move(p), roomId, eventId]() mutable {
+ QSqlQuery q(m_d.value());
+ auto res = q.prepare(getEventByIdQuery);
+ if (!res) {
+ p.addResult(std::nullopt);
+ p.finish();
+ return;
+ }
+ q.bindValue(u":room_id"_s, roomId);
+ q.bindValue(u":event_id"_s, eventId);
+ res = q.exec();
+ if (!res || !q.next()) {
+ p.addResult(std::nullopt);
+ p.finish();
+ return;
+ }
+ auto encrypted = q.value(0).toBool();
+ auto decrypted = q.value(1).toBool();
+ auto originalJson = q.value(2).toString();
+ auto decryptedJson = q.value(3).toString();
+ auto isInTimeline = q.value(4).toBool();
+ auto eventOpt = resultToEvent(encrypted, decrypted, std::move(originalJson), std::move(decryptedJson));
+ auto ret = eventOpt
+ ? std::optional<std::pair<Event, bool>>(
+ {eventOpt.value(), isInTimeline}
+ ) : std::nullopt;
+ p.addResult(ret);
+ p.finish();
+ });
+ co_return co_await fut;
+}
+
bool DbStore::valid() const
{
return m_d.has_value();
}
void DbStore::cleanup()
{
if (m_d.has_value()) {
m_d->close();
m_d.reset();
QSqlDatabase::removeDatabase(m_handle);
m_handle = QString();
}
QMetaObject::invokeMethod(m_obj, [thread=m_thread]() {
thread->quit();
});
m_thread->wait();
m_thread->deleteLater();
m_obj->deleteLater();
}
QCoro::Task<std::pair<bool, QString>> DbStore::migrate()
{
auto res = co_await asyncQuery(uR"xxx(CREATE TABLE IF NOT EXISTS events (
id INTEGER PRIMARY KEY NOT NULL,
room_id TEXT NOT NULL,
event_id TEXT NOT NULL,
encrypted INTEGER NOT NULL,
decrypted INTEGER NOT NULL,
original_json TEXT NOT NULL,
decrypted_json TEXT,
is_in_timeline INTEGER NOT NULL,
origin_server_ts INTEGER NOT NULL GENERATED ALWAYS AS (coalesce(CAST(original_json->>'origin_server_ts' AS INTEGER), 0)) STORED
);)xxx"_s);
if (!res.first) { co_return res; }
res = co_await asyncQuery(uR"xxx(CREATE UNIQUE INDEX IF NOT EXISTS events_room_id_event_id_index
ON events (
room_id,
event_id
);)xxx"_s);
if (!res.first) { co_return res; }
res = co_await asyncQuery(uR"xxx(CREATE INDEX IF NOT EXISTS events_room_id_index
ON events (
room_id
);)xxx"_s);
if (!res.first) { co_return res; }
res = co_await asyncQuery(uR"xxx(CREATE INDEX IF NOT EXISTS events_order_index
ON events (origin_server_ts, event_id)
)xxx"_s);
co_return res;
}
static const auto importEventQuery = uR"xxx(INSERT INTO events (
room_id, event_id, encrypted, decrypted,
original_json, decrypted_json, is_in_timeline
) VALUES (
:room_id, :event_id, :encrypted, :decrypted,
:original_json, :decrypted_json, :is_in_timeline
) ON CONFLICT (room_id, event_id) DO UPDATE SET
encrypted = :encrypted, decrypted = :decrypted,
original_json = :original_json, decrypted_json = :decrypted_json,
is_in_timeline = is_in_timeline OR :is_in_timeline
;)xxx"_s;
QCoro::Task<std::pair<bool, QString>> DbStore::importOne(const QString &roomId, const Kazv::Event &event, bool isInTimeline)
{
auto eventIdQs = QString::fromStdString(event.id());
auto encrypted = event.encrypted();
auto decrypted = event.decrypted();
auto originalJson = QString::fromStdString(event.originalJson().get().dump());
QVariant decryptedJson = event.encrypted() ? QString::fromStdString(event.decryptedJson().get().dump()) : QVariant(QMetaType::fromType<QString>());
auto tf = [=](QSqlQuery &q) {
q.bindValue(u":room_id"_s, roomId);
q.bindValue(u":event_id"_s, eventIdQs);
q.bindValue(u":encrypted"_s, encrypted);
q.bindValue(u":decrypted"_s, decrypted);
q.bindValue(u":original_json"_s, originalJson);
q.bindValue(u":decrypted_json"_s, decryptedJson);
q.bindValue(u":is_in_timeline"_s, isInTimeline);
};
co_return co_await asyncQuery(importEventQuery, tf);
}
QCoro::Task<std::pair<bool, QString>> DbStore::asyncQuery(const QString &qs, std::function<void(QSqlQuery &)> tf)
{
// This function exists because QSqlQuery must be used in the thread where
// the QSqlDatabase is in.
// We invoke the query inside the thread of the database, then use the
// promise-future protocol to pass the result back as a coroutine.
QPromise<std::pair<bool, QString>> p;
- qCDebug(kazvLog) << "DbStore: Async Query:" << qs;
auto fut = p.future();
QMetaObject::invokeMethod(m_obj, [this, p=std::move(p), s=std::move(qs), tf]() mutable {
p.start();
QSqlQuery q(m_d.value());
- qCDebug(kazvLog) << "DbStore: Invoking Query:" << s;
auto res = q.prepare(std::move(s));
if (!res) {
- qCDebug(kazvLog) << "DbStore: Async Query Prepare Failed: Driver:" << q.lastError().driverText();
- qCDebug(kazvLog) << "DbStore: Async Query Prepare Failed: Db:" << q.lastError().databaseText();
+ qCWarning(kazvLog) << "DbStore: Async Query Prepare Failed: Driver:" << q.lastError().driverText();
+ qCWarning(kazvLog) << "DbStore: Async Query Prepare Failed: Db:" << q.lastError().databaseText();
p.addResult({false, q.lastError().text()});
p.finish();
return;
}
tf(q);
- qCDebug(kazvLog) << "DbStore: where:" << q.boundValueNames() << q.boundValues();
res = q.exec();
- qCDebug(kazvLog) << "DbStore: Async Query Done:" << res;
if (!res) {
- qCDebug(kazvLog) << "DbStore: Async Query Failed: Driver:" << q.lastError().driverText();
- qCDebug(kazvLog) << "DbStore: Async Query Failed: Db:" << q.lastError().databaseText();
+ qCWarning(kazvLog) << "DbStore: Async Query Failed: Driver:" << q.lastError().driverText();
+ qCWarning(kazvLog) << "DbStore: Async Query Failed: Db:" << q.lastError().databaseText();
p.addResult({false, q.lastError().text()});
} else {
p.addResult({true, QString()});
}
p.finish();
});
co_return co_await fut;
}
diff --git a/src/db-store.hpp b/src/db-store.hpp
index 587308c..009014c 100644
--- a/src/db-store.hpp
+++ b/src/db-store.hpp
@@ -1,65 +1,83 @@
/*
* This file is part of kazv.
* SPDX-FileCopyrightText: 2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <kazv-defs.hpp>
#include <tuple>
#include <functional>
#include <QString>
#include <QThread>
#include <QFuture>
#include <QSqlDatabase>
#include <QCoroTask>
#include <sdk-model.hpp>
+#include <kazvevents.hpp>
/**
* This implements the event storage system using a relational database.
*
* The functions in this class are implemented as async because we do not
* want to block the calling thread. Just like in MatrixSdk, we create a
* dedicated thread for the database, and perform all work there.
* Results are sent back using coroutines as they are easier to use than QFuture.
*/
class DbStore
{
public:
static inline constexpr std::string dbDirName{"database"};
static inline constexpr std::string dbBaseName{"db.sqlite3"};
static QString getHandleFor(std::string userId, std::string deviceId);
+ static std::optional<Kazv::Event> resultToEvent(bool encrypted, bool decrypted, QString originalJson, QString decryptedJson);
DbStore();
~DbStore();
/**
* Open the database for a specific session.
* @param userDataDir The user data dir from MatrixSdk.
* @param userId The user id.
* @param deviceId The device id.
* @return A pair of (status, error text). status is true iff the setup is successful. error text contains the error text reported by the database.
*/
QCoro::Task<std::pair<bool, QString>> setup(std::string userDataDir, std::string userId, std::string deviceId);
/// @return Whether the DbStore is valid
bool valid() const;
public Q_SLOTS:
/**
* Import all events from a model.
* @param model The model.
*/
QCoro::Task<std::pair<bool, QString>> importAllFrom(const Kazv::SdkModel &model);
+ /**
+ * Save events from a SaveEventsRequested trigger.
+ *
+ * @param s The SaveEventsRequested trigger.
+ */
+ QCoro::Task<std::pair<bool, QString>> saveEvents(Kazv::SaveEventsRequested s);
+
+ /**
+ * Get one event from the database.
+ *
+ * @param roomId The room id of the event.
+ * @param eventId The id of the event.
+ * @return A pair of (Event, isInTimeline) associated with the ids or std::nullopt if it is not found.
+ */
+ QCoro::Task<std::optional<std::pair<Kazv::Event, bool>>> getEventById(const QString &roomId, const QString &eventId);
+
private:
void cleanup();
QCoro::Task<std::pair<bool, QString>> migrate();
QCoro::Task<std::pair<bool, QString>> importOne(const QString &roomId, const Kazv::Event &event, bool isInTimeline);
QCoro::Task<std::pair<bool, QString>> asyncQuery(const QString &qs, std::function<void(QSqlQuery &)> tf = [](auto &) {});
QThread *m_thread;
QObject *m_obj;
std::optional<QSqlDatabase> m_d;
QString m_handle;
};
diff --git a/src/matrix-sdk.cpp b/src/matrix-sdk.cpp
index cbdf35e..3a4e3fa 100644
--- a/src/matrix-sdk.cpp
+++ b/src/matrix-sdk.cpp
@@ -1,972 +1,1059 @@
/*
* This file is part of kazv.
* SPDX-FileCopyrightText: 2020-2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <kazv-defs.hpp>
#include <boost/archive/text_oarchive.hpp>
#include <boost/archive/text_iarchive.hpp>
#include <fstream>
#include <filesystem>
#include <chrono>
#include <QMutex>
#include <QMutexLocker>
#include <QtConcurrent>
#include <QThreadPool>
#include <KConfig>
#include <KConfigGroup>
#include <QCoroTask>
#include <eventemitter/lagerstoreeventemitter.hpp>
#include <client/sdk.hpp>
#include <client/notification-handler.hpp>
#include <crypto/base64.hpp>
#include <csapi/directory.hpp>
#include <client/alias.hpp>
#include <zug/util.hpp>
#include <lager/event_loop/qt.hpp>
#include "matrix-sdk.hpp"
#include "matrix-room-list.hpp"
#include "matrix-promise.hpp"
#include "matrix-event.hpp"
#include "helper.hpp"
#include "kazv-path-config.hpp"
#include "kazv-version.hpp"
#include "qt-json.hpp"
#include "qt-rand-adapter.hpp"
#include "qt-promise-handler.hpp"
#include "qt-job-handler.hpp"
#include "device-mgmt/matrix-device-list.hpp"
#include "matrix-sticker-pack-list.hpp"
#include "matrix-sticker-pack-list-p.hpp"
#include "matrix-user-given-attrs-map.hpp"
#include "kazv-log.hpp"
#include "matrix-utils.hpp"
#include "kazv-session-lock-guard.hpp"
#include "db-store.hpp"
using namespace Qt::Literals::StringLiterals;
using namespace Kazv;
static const std::string clientName = "kazv";
// Sdk with qt event loop, identity transform and no enhancers
using SdkT =
decltype(makeSdk(
SdkModel{},
detail::declref<JobInterface>(),
detail::declref<EventInterface>(),
QtPromiseHandler(detail::declref<QObject>()),
zug::identity,
withRandomGenerator(detail::declref<RandomInterface>())));
struct QtEventLoop
{
QObject *m_obj;
template<class Fn>
void async(Fn &&) { throw std::runtime_error{"not implemented!"}; }
template<class Fn>
void post(Fn &&fn)
{
QMetaObject::invokeMethod(
m_obj, std::forward<Fn>(fn), Qt::QueuedConnection
);
}
void finish() {}
void pause() { throw std::runtime_error{"not implemented!"}; }
void resume() { throw std::runtime_error{"not implemented!"}; }
};
+QCoro::Task<std::pair<bool, QString>> setupAndImport(SdkModel model, DbStore *dbStore, std::string userDataDir)
+{
+ auto res = co_await dbStore->setup(
+ userDataDir,
+ model.c().userId,
+ model.c().deviceId
+ );
+ if (!res.first) {
+ co_return res;
+ }
+ co_return co_await dbStore->importAllFrom(model);
+}
+
std::filesystem::path sessionDirForUserAndDeviceId(std::filesystem::path userDataDir, std::string userId, std::string deviceId)
{
auto encodedUserId = encodeBase64(userId, Base64Opts::urlSafe);
auto sessionDir = userDataDir / "sessions"
/ encodedUserId / deviceId;
return sessionDir;
}
+struct PendingSaveEvents : public SaveEventsRequested
+{
+ using DataT = immer::map<std::string, EventList>;
+ void add(SaveEventsRequested s)
+ {
+ timelineEvents = addOneType(std::move(timelineEvents), s.timelineEvents);
+ nonTimelineEvents = addOneType(std::move(nonTimelineEvents), s.nonTimelineEvents);
+ }
+
+ static DataT addOneType(DataT orig, DataT addon)
+ {
+ for (const auto &[roomId, el]: addon) {
+ orig = std::move(orig).update(roomId, [&el](EventList origEl) {
+ // This works without explicit dedupe because we do not do parallel
+ // inserts and the insert order is the same as list order.
+ return origEl + el;
+ });
+ }
+ return orig;
+ }
+};
+
+enum DbStatus
+{
+ Pending,
+ Ready,
+ Failed,
+};
+
struct MatrixSdkPrivate
{
MatrixSdkPrivate(MatrixSdk *q, bool testing, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore);
MatrixSdkPrivate(MatrixSdk *q, bool testing, SdkModel model, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore);
bool testing;
std::string userDataDir;
std::unique_ptr<KazvSessionLockGuard> lockGuard;
RandomInterface randomGenerator;
std::unique_ptr<DbStore> dbStore;
QThread *thread;
QObject *obj;
QtJobHandler *jobHandler;
LagerStoreEventEmitter ee;
LagerStoreEventEmitter::Watchable watchable;
SdkT sdk;
QTimer saveTimer;
using SecondaryRootT = decltype(sdk.createSecondaryRoot(std::declval<lager::with_qt_event_loop>()));
SecondaryRootT secondaryRoot;
Client clientOnSecondaryRoot;
NotificationHandler notificationHandler;
+ DbStatus dbStatus;
+ PendingSaveEvents pendingSaveEvents{};
+
void runIoContext() {
thread->start();
}
void stopIoContext() {
thread->quit();
}
void maybeSerialize()
{
if (!testing) {
serializeClientToFile(clientOnSecondaryRoot);
}
}
void serializeClientToFile(Client c);
+
+ void saveOrQueueEvents(SaveEventsRequested s)
+ {
+ if (dbStatus == Ready) {
+ saveEvents(std::move(s));
+ } else if (dbStatus == Pending) {
+ pendingSaveEvents.add(std::move(s));
+ }
+ // Database open failed, do not do anything
+ }
+
+ void saveEvents(SaveEventsRequested s)
+ {
+ dbStore->saveEvents(std::move(s))
+ .then([](const auto &stat) {
+ if (!stat.first) {
+ qCWarning(kazvLog) << "Cannot save events:" << stat.second;
+ }
+ });
+ }
};
// Cleaning up notes:
// 0. Callback functions may store the context for an indefinite time
// 1. The QThread event loop can be stopped
// 2. QtJobHandler::submit() should only happen in the primary event loop thread
// 3. QtJobHandler lives in the primary event loop thread
// 4. QtJobs live in the primary event loop thread
// 5. Job callbacks are called in the primary event loop thread
// 6. QtPromise::then() callbacks are called in the primary event loop thread
// 7. When the QThread event loop stops, no more callbacks will be executed (there is nothing to post to)
// 8. The QThread should stop before obj is deleted
class CleanupHelper : public QObject
{
Q_OBJECT
public:
explicit CleanupHelper(std::unique_ptr<MatrixSdkPrivate> d)
: oldD(std::move(d))
{
}
std::unique_ptr<MatrixSdkPrivate> oldD;
void cleanup()
{
qCInfo(kazvLog) << "start to clean up everything";
oldD->clientOnSecondaryRoot.stopSyncing()
.then([obj=oldD->obj, thread=oldD->thread](auto &&) {
qCDebug(kazvLog) << "stopped syncing";
QMetaObject::invokeMethod(obj, [thread]() {
thread->quit();
});
});
oldD->thread->wait();
oldD->thread->deleteLater();
// After the thread's event loop is finished, we can delete the root object
oldD->obj->deleteLater();
this->deleteLater();
qCInfo(kazvLog) << "thread is done";
}
};
void MatrixSdkPrivate::serializeClientToFile(Client c)
{
using namespace Kazv::CursorOp;
auto userId = +c.userId();
auto deviceId = +c.deviceId();
if (userId.empty() || deviceId.empty()) {
qDebug() << "Not logged in, nothing to serialize";
return;
}
using StdPath = std::filesystem::path;
auto userDataDir = StdPath(this->userDataDir);
auto sessionDir = sessionDirForUserAndDeviceId(userDataDir, userId, deviceId);
auto storeFile = sessionDir / "store";
auto storeFileNew = sessionDir / "store.new";
auto metadataFile = sessionDir / "metadata";
qDebug() << "storeFile=" << QString::fromStdString(storeFile.string());
std::error_code err;
if ((! std::filesystem::create_directories(sessionDir, err))
&& err) {
qDebug() << "Unable to create sessionDir";
return;
}
if (!lockGuard) {
try {
lockGuard = std::make_unique<KazvSessionLockGuard>(sessionDir);
} catch (const std::runtime_error &e) {
qCWarning(kazvLog) << "Error locking session: " << e.what();
return;
}
}
try {
auto storeStream = std::ofstream(storeFileNew);
if (! storeStream) {
qCWarning(kazvLog) << "Unable to open storeFile";
return;
}
using OAr = boost::archive::text_oarchive;
auto archive = OAr{storeStream};
c.serializeTo(archive);
} catch (const std::exception &e) {
qCWarning(kazvLog) << "Cannot write to store file: " << e.what();
return;
}
err.clear();
std::filesystem::rename(storeFileNew, storeFile, err);
if (err) {
qCWarning(kazvLog) << "Cannot move storeFile into place: " << QString::fromStdString(err.message());
return;
}
qDebug() << "Serialization done";
// store metadata
{
KConfig metadata(QString::fromStdString(metadataFile.string()));
KConfigGroup mdGroup(&metadata, u"Metadata"_s);
mdGroup.writeEntry("kazvVersion", QString::fromStdString(kazvVersionString()));
mdGroup.writeEntry("archiveFormat", "text-with-sqlite");
}
}
MatrixSdkPrivate::MatrixSdkPrivate(MatrixSdk *q, bool testing, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore)
: testing(testing)
, userDataDir{kazvUserDataDir().toStdString()}
, lockGuard(std::move(lockGuard))
, randomGenerator(QtRandAdapter{})
, dbStore(std::move(dbStore))
, thread(new QThread())
, obj(new QObject())
, jobHandler(new QtJobHandler(obj))
, ee{QtEventLoop{obj}}
, watchable(ee.watchable())
, sdk(makeDefaultSdkWithCryptoRandom(
randomGenerator.generateRange<std::string>(makeDefaultSdkWithCryptoRandomSize()),
static_cast<JobInterface &>(*jobHandler),
static_cast<EventInterface &>(ee),
QtPromiseHandler(*obj),
zug::identity,
withRandomGenerator(randomGenerator)))
, secondaryRoot(sdk.createSecondaryRoot(QtEventLoop{q}))
, clientOnSecondaryRoot(sdk.clientFromSecondaryRoot(secondaryRoot))
, notificationHandler(clientOnSecondaryRoot.notificationHandler())
+ , dbStatus(this->dbStore ? Ready : Pending)
{
obj->moveToThread(thread);
}
MatrixSdkPrivate::MatrixSdkPrivate(MatrixSdk *q, bool testing, SdkModel model, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore)
: testing(testing)
, userDataDir{kazvUserDataDir().toStdString()}
, lockGuard(std::move(lockGuard))
, randomGenerator(QtRandAdapter{})
, dbStore(std::move(dbStore))
, thread(new QThread())
, obj(new QObject())
, jobHandler(new QtJobHandler(obj))
, ee{QtEventLoop{obj}}
, watchable(ee.watchable())
, sdk(makeSdk(
model,
static_cast<JobInterface &>(*jobHandler),
static_cast<EventInterface &>(ee),
QtPromiseHandler(*obj),
zug::identity,
withRandomGenerator(randomGenerator)))
, secondaryRoot(sdk.createSecondaryRoot(QtEventLoop{q}, std::move(model)))
, clientOnSecondaryRoot(sdk.clientFromSecondaryRoot(secondaryRoot))
, notificationHandler(clientOnSecondaryRoot.notificationHandler())
+ , dbStatus(this->dbStore ? Ready : Pending)
{
obj->moveToThread(thread);
}
MatrixSdk::MatrixSdk(std::unique_ptr<MatrixSdkPrivate> d, QObject *parent)
: QObject(parent)
, m_d(std::move(d))
{
init();
connect(this, &MatrixSdk::trigger,
- this, [](KazvEvent e) {
- qDebug() << "receiving trigger:";
- if (std::holds_alternative<LoginSuccessful>(e)) {
- qDebug() << "Login successful";
- }
- });
+ this, [this](KazvEvent e) {
+ if (std::holds_alternative<SaveEventsRequested>(e)) {
+ auto s = std::get<SaveEventsRequested>(e);
+ m_d->saveOrQueueEvents(s);
+ }
+ });
+
+ // loginSuccessful will be emitted only for new sessions (not for loaded
+ // sessions),
+ connect(this, &MatrixSdk::loginSuccessful, this, [this]() {
+ if (m_d->dbStore) {
+ qCWarning(kazvLog) << "Trying to replace current db store, ignoring";
+ return;
+ }
+ m_d->dbStore = std::make_unique<DbStore>();
+ // This must be called from the GUI thread
+ auto model = m_d->secondaryRoot.get();
+ setupAndImport(std::move(model), m_d->dbStore.get(), m_d->userDataDir)
+ .then([this](const auto &stat) {
+ if (!stat.first) {
+ qCWarning(kazvLog) << "Cannot set up database:" << stat.second;
+ }
+ Q_EMIT dbResolved(stat.first);
+ });
+ });
+
+ connect(this, &MatrixSdk::dbResolved, this, [this](bool success) {
+ if (!success) {
+ m_d->dbStatus = Failed;
+ } else {
+ qCDebug(kazvLog) << "db loading succeeded";
+ auto pendingSaveEvents = std::move(m_d->pendingSaveEvents);
+ m_d->pendingSaveEvents = {};
+ m_d->saveEvents(std::move(pendingSaveEvents));
+ m_d->dbStatus = Ready;
+ }
+ });
}
void MatrixSdk::init()
{
LAGER_QT(serverUrl) = m_d->clientOnSecondaryRoot.serverUrl().xform(strToQt); Q_EMIT serverUrlChanged(serverUrl());
LAGER_QT(userId) = m_d->clientOnSecondaryRoot.userId().xform(strToQt); Q_EMIT userIdChanged(userId());
LAGER_QT(token) = m_d->clientOnSecondaryRoot.token().xform(strToQt); Q_EMIT tokenChanged(token());
LAGER_QT(deviceId) = m_d->clientOnSecondaryRoot.deviceId().xform(strToQt); Q_EMIT deviceIdChanged(deviceId());
LAGER_QT(specVersions) = m_d->clientOnSecondaryRoot.supportVersions(); Q_EMIT specVersionsChanged(specVersions());
m_d->watchable.afterAll(
[this](KazvEvent e) {
Q_EMIT this->trigger(e);
});
m_d->watchable.after<LoginSuccessful>(
[this](LoginSuccessful e) {
Q_EMIT this->loginSuccessful(e);
});
m_d->watchable.after<LoginFailed>(
[this](LoginFailed e) {
Q_EMIT this->loginFailed(
QString::fromStdString(e.errorCode),
QString::fromStdString(e.error)
);
});
m_d->watchable.after<ReceivingRoomTimelineEvent>(
[this](ReceivingRoomTimelineEvent e) {
Q_EMIT this->receivedMessage(
QString::fromStdString(e.roomId),
QString::fromStdString(e.event.id())
);
});
connect(&m_d->saveTimer, &QTimer::timeout, &m_d->saveTimer, [m_d=m_d.get()]() {
m_d->maybeSerialize();
});
const int saveIntervalMs = 1000 * 60 * 5;
m_d->saveTimer.start(std::chrono::milliseconds{saveIntervalMs});
}
MatrixSdk::MatrixSdk(QObject *parent)
: MatrixSdk(std::make_unique<MatrixSdkPrivate>(
this,
/* testing = */ false,
std::unique_ptr<KazvSessionLockGuard>(),
std::unique_ptr<DbStore>()
), parent)
{
}
MatrixSdk::MatrixSdk(SdkModel model, bool testing, QObject *parent)
: MatrixSdk(std::make_unique<MatrixSdkPrivate>(
this,
testing,
std::move(model),
std::unique_ptr<KazvSessionLockGuard>(),
std::unique_ptr<DbStore>()
), parent)
{
}
static void cleanupDPointer(std::unique_ptr<MatrixSdkPrivate> oldD)
{
oldD->saveTimer.disconnect();
oldD->saveTimer.stop();
auto helper = new CleanupHelper(std::move(oldD));
helper->cleanup();
}
MatrixSdk::~MatrixSdk()
{
if (m_d) {
serializeToFile();
cleanupDPointer(std::move(m_d));
}
}
QString MatrixSdk::mxcUriToHttp(QString mxcUri) const
{
return QString::fromStdString(m_d->clientOnSecondaryRoot.mxcUriToHttp(mxcUri.toStdString()));
}
QString MatrixSdk::mxcUriToHttpAuthenticatedV1(QString mxcUri) const
{
return QString::fromStdString(m_d->clientOnSecondaryRoot.mxcUriToHttpV1(mxcUri.toStdString()));
}
MatrixDeviceList *MatrixSdk::devicesOfUser(QString userId) const
{
return new MatrixDeviceList(m_d->clientOnSecondaryRoot.devicesOfUser(userId.toStdString()));
}
bool isIllFormatSpecVersion(const QString &version)
{
const bool isLegacy = version.startsWith(u"r"_s);
if (isLegacy && version.split(u'.').size() == 3) {
return false;
}
if (!isLegacy && version.split(u'.').size() == 2) {
return false;
}
return true;
}
// Return true if v1 is at least as new as v2, false if v1 is older than v2
// Return false if v1 or v2 is ill-format
bool compareSpecVersion(QString v1, QString v2)
{
// Check parameters format
if (isIllFormatSpecVersion(v1) || isIllFormatSpecVersion(v2)) {
return false;
}
const bool v1IsLegacy = v1.startsWith(u"r"_s);
const bool v2IsLegacy = v2.startsWith(u"r"_s);
if (v1IsLegacy != v2IsLegacy) {
return v2IsLegacy;
}
v1.remove(0, 1);
v2.remove(0, 1);
auto v1VersionNumbers = v1.split(u'.');
auto v2VersionNumbers = v2.split(u'.');
for (int i = 0; i < v1VersionNumbers.size(); i++) {
auto v1VerNum = v1VersionNumbers[i].toInt();
auto v2VerNum = v2VersionNumbers[i].toInt();
if (v1VerNum != v2VerNum) {
return v1VerNum > v2VerNum;
}
}
// v1 is equal to v2
return true;
}
// Return true if version in the range [minVer, maxVer]
// Return false if any parameter is ill-format
bool compareSpecVersionRange(const QString &version,
const QString &minVer, const QString &maxVer)
{
if (isIllFormatSpecVersion(version)
|| isIllFormatSpecVersion(minVer)
|| isIllFormatSpecVersion(maxVer)) {
return false;
}
if (compareSpecVersion(version, minVer) && compareSpecVersion(maxVer, version)) {
return true;
}
return false;
}
bool MatrixSdk::checkSpecVersion(QString version) const
{
return std::find_if(specVersions().begin(), specVersions().end(),
[&version](auto v) {
return compareSpecVersion(QString::fromStdString(v), version);
}) != specVersions().end();
}
bool MatrixSdk::checkSpecVersionRange(QString minVer, QString maxVer) const
{
return std::find_if(specVersions().begin(), specVersions().end(),
[&minVer, maxVer](auto v) {
return compareSpecVersionRange(QString::fromStdString(v), minVer, maxVer);
}) != specVersions().end();
}
QStringList MatrixSdk::directRoomIds(QString userId) const
{
auto content = m_d->clientOnSecondaryRoot.accountData().get()["m.direct"].content().get();
auto roomIds = QStringList{};
for (auto i : content[userId.toStdString()]) {
roomIds.push_back(QString::fromStdString(i.get<std::string>()));
}
return roomIds;
}
std::string MatrixSdk::validateHomeserverUrl(const QString &url)
{
if (url.isEmpty()) {
return std::string();
}
auto u = QUrl::fromUserInput(url);
if (!u.isValid()) {
return std::string();
}
if (u.scheme() == u"http"_s) {
qCInfo(kazvLog) << "url" << u << "is http. Force switching to https.";
u.setScheme(u"https"_s);
} else if (u.scheme() != u"https"_s) {
qCWarning(kazvLog) << "url" << u << "is not http/https.";
return std::string();
}
return u.toString().toStdString();
}
void MatrixSdk::login(const QString &userId, const QString &password, const QString &homeserverUrl)
{
auto loginFunc = [userId, password](const Client &client, const std::string &serverUrl) {
client.passwordLogin(
serverUrl,
userId.toStdString(),
password.toStdString(),
clientName
);
};
auto validated = validateHomeserverUrl(homeserverUrl);
if (!validated.empty()) {
loginFunc(m_d->clientOnSecondaryRoot, validated);
} else {
m_d->clientOnSecondaryRoot
.autoDiscover(userId.toStdString()) // autoDiscover() will dispatch GetVersionAction to get supported versions of the server
.then([
this,
client=m_d->clientOnSecondaryRoot.toEventLoop(),
userId,
password,
loginFunc
](auto res) {
if (!res.success()) {
// FIXME use real error codes and msgs when available in libkazv
Q_EMIT this->discoverFailed(u""_s, u""_s);
return res;
}
auto serverUrl = res.dataStr("homeserverUrl");
loginFunc(client, serverUrl);
return res;
});
}
m_d->runIoContext();
}
void MatrixSdk::logout()
{
m_d->clientOnSecondaryRoot.logout()
.then([&] (EffectStatus stat) {
if (stat.success()) {
m_d->stopIoContext();
Q_EMIT this->logoutSuccessful();
} else {
Q_EMIT this->logoutFailed(QString::fromStdString(stat.dataStr("errorCode")), QString::fromStdString(stat.dataStr("error")));
}
});
}
MatrixRoomList *MatrixSdk::roomList() const
{
return new MatrixRoomList(m_d->clientOnSecondaryRoot);
}
void MatrixSdk::emplace(std::optional<SdkModel> model, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore)
{
auto testing = m_d->testing;
auto userDataDir = m_d->userDataDir;
if (m_d) {
cleanupDPointer(std::move(m_d));
}
m_d = model.has_value()
? std::make_unique<MatrixSdkPrivate>(
this, testing, std::move(model.value()),
std::move(lockGuard), std::move(dbStore))
: std::make_unique<MatrixSdkPrivate>(
this, testing, std::move(lockGuard), std::move(dbStore)
);
m_d->userDataDir = userDataDir;
// Re-initialize lager-qt cursors and watchable connections
init();
m_d->runIoContext();
m_d->clientOnSecondaryRoot.startSyncing();
Q_EMIT sessionChanged();
}
QStringList MatrixSdk::allSessions() const
{
using StdPath = std::filesystem::path;
auto userDataDir = StdPath(m_d->userDataDir);
auto allSessionsDir = userDataDir / "sessions";
QStringList sessionNames;
try {
for (const auto &p : std::filesystem::directory_iterator(allSessionsDir)) {
if (p.is_directory()) {
auto maybeEncodedUserId = p.path().filename().string();
auto userId = decodeBase64(maybeEncodedUserId, Base64Opts::urlSafe);
if (userId.empty() || userId[0] != '@') {
continue;
}
for (const auto &q : std::filesystem::directory_iterator(p.path())) {
auto path = q.path();
auto deviceId = path.filename().string();
std::error_code err;
if (std::filesystem::exists(path / "store", err)) {
sessionNames.append(QString::fromStdString(userId + "/" + deviceId));
}
}
}
}
} catch (const std::filesystem::filesystem_error &) {
qDebug() << "sessionDir not available, ignoring";
}
return sessionNames;
}
void MatrixSdk::serializeToFile() const
{
m_d->maybeSerialize();
}
void MatrixSdk::loadSession(QString sessionName)
{
using StdPath = std::filesystem::path;
auto loadFromSession = [this, sessionName](StdPath sessionDir, std::unique_ptr<KazvSessionLockGuard> lockGuard) {
auto storeFile = sessionDir / "store";
auto metadataFile = sessionDir / "metadata";
if (! std::filesystem::exists(storeFile)) {
qDebug() << "storeFile does not exist, skip loading session " << sessionName;
Q_EMIT loadSessionFinished(sessionName, SessionNotFound);
return;
}
if (std::filesystem::exists(metadataFile)) {
KConfig metadata(QString::fromStdString(metadataFile.string()));
KConfigGroup mdGroup(&metadata, u"Metadata"_s);
auto format = mdGroup.readEntry(u"archiveFormat"_s);
if (format != u"text"_s && format != u"text-with-sqlite"_s) {
qDebug() << "Unknown archive format:" << format;
Q_EMIT loadSessionFinished(sessionName, SessionFormatUnknown);
return;
}
auto version = mdGroup.readEntry(u"kazvVersion"_s);
auto curVersion = kazvVersionString();
if (version != QString::fromStdString(curVersion)) {
qDebug() << "A different version from the current one, making a backup";
std::error_code err;
auto now = std::chrono::system_clock::now();
auto backupName =
std::to_string(std::chrono::duration_cast<std::chrono::seconds>(now.time_since_epoch()).count());
auto backupDir = sessionDir / "backup" / backupName;
if (! std::filesystem::create_directories(backupDir, err)
&& err) {
qDebug() << "Cannot create backup directory";
Q_EMIT loadSessionFinished(sessionName, SessionCannotBackup);
return;
}
std::filesystem::copy_file(storeFile, backupDir / "store");
std::filesystem::copy_file(metadataFile, backupDir / "metadata");
}
}
SdkModel model;
try {
auto storeStream = std::ifstream(storeFile);
if (! storeStream) {
qDebug() << "Unable to open storeFile";
Q_EMIT loadSessionFinished(sessionName, SessionCannotOpenFile);
return;
}
using IAr = boost::archive::text_iarchive;
auto archive = IAr{storeStream};
archive >> model;
qDebug() << "Finished loading session store file";
} catch (const std::exception &e) {
qDebug() << "Error when loading session:" << QString::fromStdString(e.what());
Q_EMIT loadSessionFinished(sessionName, SessionDeserializeFailed);
return;
}
auto dbStore = std::make_unique<DbStore>();
- auto stat = QCoro::waitFor([](SdkModel model, DbStore *dbStore, std::string userDataDir) -> QCoro::Task<std::pair<bool, QString>> {
- qCDebug(kazvLog) << "open db:" << userDataDir << model.c().userId << model.c().deviceId;
- auto res = co_await dbStore->setup(
- userDataDir,
- model.c().userId,
- model.c().deviceId
- );
- if (!res.first) {
- co_return res;
- }
- co_return co_await dbStore->importAllFrom(model);
- }(model, dbStore.get(), m_d->userDataDir));
+ auto stat = QCoro::waitFor(setupAndImport(model, dbStore.get(), m_d->userDataDir));
if (!stat.first) {
qCWarning(kazvLog) << "Unable to open db store:" << stat.second;
Q_EMIT loadSessionFinished(sessionName, SessionDeserializeFailed);
return;
}
QMetaObject::invokeMethod(this, [this, model=std::move(model), lockGuard=std::move(lockGuard), dbStore=std::move(dbStore)]() mutable {
emplace(std::move(model), std::move(lockGuard), std::move(dbStore));
});
Q_EMIT loadSessionFinished(sessionName, SessionLoadSuccess);
};
qDebug() << "in loadSession(), sessionName=" << sessionName;
auto userDataDir = StdPath(m_d->userDataDir);
auto parts = sessionName.split(u'/');
if (parts.size() == 2) {
auto userId = parts[0].toStdString();
auto deviceId = parts[1].toStdString();
auto sessionDir = sessionDirForUserAndDeviceId(userDataDir, userId, deviceId);
std::unique_ptr<KazvSessionLockGuard> lockGuard;
try {
lockGuard = std::make_unique<KazvSessionLockGuard>(sessionDir);
} catch (const std::runtime_error &e) {
qCWarning(kazvLog) << "Error locking session: " << e.what();
Q_EMIT loadSessionFinished(sessionName, SessionLockFailed);
return;
}
QThreadPool::globalInstance()->start([loadFromSession, sessionDir, lockGuard=std::move(lockGuard)]() mutable {
loadFromSession(sessionDir, std::move(lockGuard));
});
return;
}
qDebug(kazvLog) << "no session found for" << sessionName;
Q_EMIT loadSessionFinished(sessionName, SessionNotFound);
}
bool MatrixSdk::deleteSession(QString sessionName) {
using StdPath = std::filesystem::path;
qDebug() << "in deleteSession(), sessionName=" << sessionName;
auto userDataDir = StdPath(m_d->userDataDir);
auto parts = sessionName.split(u'/');
if (parts.size() == 2) {
auto userId = parts[0].toStdString();
auto deviceId = parts[1].toStdString();
auto sessionDir = sessionDirForUserAndDeviceId(userDataDir, userId, deviceId);
if (std::filesystem::exists(sessionDir)) {
qCDebug(kazvLog) << "new path works";
return std::filesystem::remove_all(sessionDir);
}
qCDebug(kazvLog) << "trying legacy path";
auto legacySessionDir = userDataDir / "sessions" / userId / deviceId;
if (std::filesystem::exists(legacySessionDir)) {
qCDebug(kazvLog) << "legacy path works";
return std::filesystem::remove_all(legacySessionDir);
}
}
qDebug(kazvLog) << "no session found for" << sessionName;
return false;
}
bool MatrixSdk::startNewSession()
{
emplace(std::nullopt, std::unique_ptr<KazvSessionLockGuard>(), std::unique_ptr<DbStore>());
return true;
}
static std::optional<std::string> optMaybe(QString s)
{
if (s.isEmpty()) {
return std::nullopt;
} else {
return s.toStdString();
}
}
MatrixPromise *MatrixSdk::createRoom(
bool isPrivate,
const QString &name,
const QString &alias,
const QStringList &invite,
bool isDirect,
bool allowFederate,
const QString &topic,
const QJsonValue &powerLevelContentOverride,
CreateRoomPreset preset,
bool encrypted
)
{
immer::array<Event> initialState;
if (encrypted) {
initialState = {Event{json{
{"type", "m.room.encryption"},
{"state_key", ""},
{"content", {
{"algorithm", "m.megolm.v1.aes-sha2"},
}},
}}};
}
return new MatrixPromise(m_d->clientOnSecondaryRoot.createRoom(
isPrivate ? Kazv::RoomVisibility::Private : Kazv::RoomVisibility::Public,
optMaybe(name),
optMaybe(alias),
qStringListToStdF(invite),
isDirect,
allowFederate,
optMaybe(topic),
nlohmann::json(powerLevelContentOverride),
static_cast<Kazv::CreateRoomPreset>(preset),
initialState
));
}
MatrixPromise *MatrixSdk::joinRoom(const QString &idOrAlias, const QStringList &servers)
{
return new MatrixPromise(m_d->clientOnSecondaryRoot.joinRoom(
idOrAlias.toStdString(),
qStringListToStdF(servers)
));
}
MatrixPromise *MatrixSdk::setDeviceTrustLevel(QString userId, QString deviceId, QString trustLevel)
{
return new MatrixPromise(
m_d->clientOnSecondaryRoot.setDeviceTrustLevel(
userId.toStdString(),
deviceId.toStdString(),
qStringToTrustLevelFunc(trustLevel)
)
);
}
MatrixPromise *MatrixSdk::getSelfProfile()
{
return new MatrixPromise(
m_d->clientOnSecondaryRoot.getProfile(userId().toStdString())
);
}
MatrixPromise *MatrixSdk::setDisplayName(QString displayName)
{
return new MatrixPromise(
m_d->clientOnSecondaryRoot.setDisplayName(
displayName.isEmpty() ? std::nullopt : std::optional<std::string>(displayName.toStdString())
)
);
}
MatrixPromise *MatrixSdk::setAvatarUrl(QString avatarUrl)
{
return new MatrixPromise(
m_d->clientOnSecondaryRoot.setAvatarUrl(
avatarUrl.isEmpty() ? std::nullopt : std::optional<std::string>(avatarUrl.toStdString())
)
);
}
void MatrixSdk::startThread()
{
m_d->runIoContext();
}
RandomInterface &MatrixSdk::randomGenerator() const
{
return m_d->randomGenerator;
}
bool MatrixSdk::shouldNotify(MatrixEvent *event) const
{
// Do not notify own event
if (event->sender() == userId()) {
return false;
}
return m_d->notificationHandler.handleNotification(event->underlyingEvent()).shouldNotify;
}
bool MatrixSdk::shouldPlaySound(MatrixEvent *event) const
{
return m_d->notificationHandler.handleNotification(event->underlyingEvent()).sound.has_value();
}
MatrixStickerPackList *MatrixSdk::stickerPackList() const
{
return new MatrixStickerPackList(m_d->clientOnSecondaryRoot);
}
MatrixEvent *MatrixSdk::stickerRoomsEvent() const
{
return new MatrixEvent(m_d->clientOnSecondaryRoot.accountData()[imagePackRoomsEventType][lager::lenses::or_default]);
}
MatrixPromise *MatrixSdk::updateStickerPack(MatrixStickerPackSource source)
{
if (source.source == MatrixStickerPackSource::AccountData) {
auto eventJson = std::move(source.event).raw().get();
eventJson["type"] = source.eventType;
return sendAccountDataImpl(Event(std::move(eventJson)));
} else if (source.source == MatrixStickerPackSource::RoomState) {
auto eventJson = std::move(source.event).raw().get();
eventJson["type"] = source.eventType;
eventJson["state_key"] = source.stateKey;
return new MatrixPromise(
m_d->clientOnSecondaryRoot
.room(source.roomId)
.sendStateEvent(Event(std::move(eventJson))));
} else {
return 0;
}
}
MatrixUserGivenAttrsMap *MatrixSdk::userGivenNicknameMap() const
{
return new MatrixUserGivenAttrsMap(
userGivenNicknameMapFor(m_d->clientOnSecondaryRoot),
[client=m_d->clientOnSecondaryRoot](json content) {
return client.setAccountData(json{
{"type", USER_GIVEN_NICKNAME_EVENT_TYPES[0]},
{"content", std::move(content)},
});
}
);
}
MatrixPromise *MatrixSdk::sendAccountData(const QString &type, const QJsonObject &content)
{
Event e = json{
{"type", type.toStdString()},
{"content", content},
};
return sendAccountDataImpl(std::move(e));
}
MatrixPromise *MatrixSdk::sendAccountDataImpl(Event event)
{
return new MatrixPromise(m_d->clientOnSecondaryRoot.setAccountData(event));
}
MatrixPromise *MatrixSdk::getSpecVersions()
{
return new MatrixPromise(m_d->clientOnSecondaryRoot.getVersions(LAGER_QT(serverUrl).get().toStdString()));
}
MatrixPromise *MatrixSdk::addDirectRoom(const QString &userId, const QString &roomId)
{
return new MatrixPromise(m_d->clientOnSecondaryRoot.addDirectRoom(userId.toStdString(), roomId.toStdString()));
}
MatrixPromise *MatrixSdk::getRoomIdByAlias(const QString &roomAlias)
{
auto job = m_d->clientOnSecondaryRoot
.getRoomIdByAliasJob(roomAlias.toStdString());
auto jh = m_d->jobHandler;
return new MatrixPromise(
m_d->sdk.context().createPromise([jh, job](auto resolve) {
jh->submit(job, [resolve](GetRoomIdByAliasResponse r) {
resolve(parseGetRoomIdByAliasResponse(r));
});
}));
}
void MatrixSdk::setUserDataDir(const std::string &userDataDir)
{
m_d->userDataDir = userDataDir;
}
#include "matrix-sdk.moc"
diff --git a/src/matrix-sdk.hpp b/src/matrix-sdk.hpp
index 5698e3f..2f5219a 100644
--- a/src/matrix-sdk.hpp
+++ b/src/matrix-sdk.hpp
@@ -1,325 +1,326 @@
/*
* This file is part of kazv.
* SPDX-FileCopyrightText: 2020-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <kazv-defs.hpp>
#include <QObject>
#include <QQmlEngine>
#include <QString>
#include <string>
#include <memory>
#include <filesystem>
#include <lager/extra/qt.hpp>
#include <immer/map.hpp>
#include <sdk-model.hpp>
#include <random-generator.hpp>
#include "meta-types.hpp"
Q_MOC_INCLUDE("matrix-room-list.hpp")
Q_MOC_INCLUDE("matrix-device-list.hpp")
Q_MOC_INCLUDE("matrix-promise.hpp")
Q_MOC_INCLUDE("matrix-event.hpp")
Q_MOC_INCLUDE("matrix-sticker-pack-list.hpp")
Q_MOC_INCLUDE("matrix-user-given-attrs-map.hpp")
class MatrixRoomList;
class MatrixDeviceList;
class MatrixPromise;
class MatrixSdkTest;
class MatrixSdkSessionsTest;
class MatrixEvent;
class MatrixStickerPackList;
class MatrixUserGivenAttrsMap;
class KazvSessionLockGuard;
class DbStore;
struct MatrixSdkPrivate;
std::filesystem::path sessionDirForUserAndDeviceId(std::filesystem::path userDataDir, std::string userId, std::string deviceId);
class MatrixSdk : public QObject
{
Q_OBJECT
QML_ELEMENT
std::unique_ptr<MatrixSdkPrivate> m_d;
/// @param d A dynamically allocated d-pointer, whose ownership
/// will be transferred to this MatrixSdk.
explicit MatrixSdk(std::unique_ptr<MatrixSdkPrivate> d, QObject *parent);
void init();
public:
enum CreateRoomPreset {
PrivateChat = Kazv::CreateRoomPreset::PrivateChat,
PublicChat = Kazv::CreateRoomPreset::PublicChat,
TrustedPrivateChat = Kazv::CreateRoomPreset::TrustedPrivateChat,
};
Q_ENUM(CreateRoomPreset);
enum LoadSessionResult {
/// Successfully loaded the session
SessionLoadSuccess,
/// There is no store file
SessionNotFound,
/// The format of the store file is not supported
SessionFormatUnknown,
/// The store file cannot be backed up
SessionCannotBackup,
/// Cannot grab the lock on the session file
SessionLockFailed,
/// Cannot open store file
SessionCannotOpenFile,
/// Cannot deserialize the store file
SessionDeserializeFailed,
};
Q_ENUM(LoadSessionResult);
enum MediaDownloadEndpointVersion {
UnauthenticatedMediaV3, // Classical endpoint series
AuthenticatedMediaV1, // Endpoints that needs authorization
};
Q_ENUM(MediaDownloadEndpointVersion);
explicit MatrixSdk(QObject *parent = 0);
~MatrixSdk() override;
LAGER_QT_READER(QString, serverUrl);
LAGER_QT_READER(QString, userId);
LAGER_QT_READER(QString, token);
LAGER_QT_READER(QString, deviceId);
LAGER_QT_READER(immer::array<std::string>, specVersions); // The versions of the Matrix Spec supported by the server.
Q_INVOKABLE MatrixRoomList *roomList() const;
Q_INVOKABLE QString mxcUriToHttp(QString mxcUri) const;
Q_INVOKABLE QString mxcUriToHttpAuthenticatedV1(QString mxcUri) const;
Q_INVOKABLE MatrixDeviceList *devicesOfUser(QString userId) const;
// Return true if version is at least as new as the spec version of server
Q_INVOKABLE bool checkSpecVersion(QString version) const;
// Return true if the spec version of server in the range [minVer, maxVer]
Q_INVOKABLE bool checkSpecVersionRange(QString minVer, QString maxVer) const;
Q_INVOKABLE QStringList directRoomIds(QString userId) const;
Kazv::RandomInterface &randomGenerator() const;
private:
// Replaces the store with another one
void emplace(std::optional<Kazv::SdkModel> model, std::unique_ptr<KazvSessionLockGuard> lockGuard, std::unique_ptr<DbStore> dbStore);
static std::string validateHomeserverUrl(const QString &url);
Q_SIGNALS:
void trigger(Kazv::KazvEvent e);
void loginSuccessful(Kazv::KazvEvent e);
void loginFailed(QString errorCode, QString errorMsg);
void discoverFailed(QString errorCode, QString errorMsg);
void logoutSuccessful();
void logoutFailed(QString errorCode, QString errorMsg);
void receivedMessage(QString roomId, QString eventId);
void sessionChanged();
void loadSessionFinished(QString sessionName, MatrixSdk::LoadSessionResult result);
+ void dbResolved(bool success);
public Q_SLOTS:
void login(const QString &userId, const QString &password, const QString &homeserverUrl);
void logout();
/**
* Serialize data to <AppDataDir>/sessions/<userid>/<deviceid>/
*
* If not logged in, do nothing.
*/
void serializeToFile() const;
/**
* Load session at <AppDataDir>/sessions/<sessionName> asynchronously.
*
* When the session is loaded or when there is an error, loadSessionFinished
* will be called with the result.
*
* @param sessionName A string in the form of <userid>/<deviceid> .
*/
void loadSession(QString sessionName);
/**
* Delete session at <AppDataDir>/sessions/<sessionName> .
*
* @param sessionName A string in the form of <userid>/<deviceid> .
*
* @return true if successful, false otherwise.
*/
bool deleteSession(QString sessionName);
/**
* Start an empty session.
*
* The new session is not logged in, and need to call login().
*
* @return true if successful, false otherwise.
*/
bool startNewSession();
/**
* Get all saved sessions.
*
* @return A list of session names in the form of <userid>/<deviceid> .
*/
QStringList allSessions() const;
/**
* Create a new room.
*
* @param isPrivate Whether the room is private.
* @param name The room's name.
* @param alias The alias of the room.
* @param invite List of matrix ids of users to invite.
* @param isDirect Whether it is a direct message room.
* @param allowFederate Whether to allow users on other servers to join.
* @param topic The topic of the room.
* @param powerLevelContentOverride The content to override m.room.power_levels event.
* @param preset The preset to create the room with.
* @param encrypted Whether to enable encryption for this room.
*/
MatrixPromise *createRoom(
bool isPrivate,
const QString &name,
const QString &alias,
const QStringList &invite,
bool isDirect,
bool allowFederate,
const QString &topic,
const QJsonValue &powerLevelContentOverride,
CreateRoomPreset preset,
bool encrypted
);
/**
* Join a room.
* @param idOrAlias The id or alias of the room to join.
* @param servers The servers to use when joining the room.
*/
MatrixPromise *joinRoom(
const QString &idOrAlias,
const QStringList &servers
);
/**
* Change the trust level of a device.
*
* @param userId The user id that owns the device.
* @param deviceId The device id to set the trust level.
* @param trustLevel The trust level.
*
* @return A MatrixPromise representing the progress.
*/
MatrixPromise *setDeviceTrustLevel(QString userId, QString deviceId, QString trustLevel);
/**
* Get the profile of the current user.
*
* @return A MatrixPromise representing the progress.
*/
MatrixPromise *getSelfProfile();
/**
* Set the display name of the current user.
*
* @return A MatrixPromise representing the progress.
*/
MatrixPromise *setDisplayName(QString displayName);
/**
* Set the avatar url of the current user.
*
* @return A MatrixPromise representing the progress.
*/
MatrixPromise *setAvatarUrl(QString avatarUrl);
/**
* Check if an event should be notified.
*
* @param event The event to check.
* @return Whether `event` should be notified.
*/
bool shouldNotify(MatrixEvent *event) const;
/**
* Check if an event should be notified with sound.
*
* You should only call this method when `shouldNotify(event)`
* returns true.
*
* @param event The event to check.
* @return Whether `event` should be notified with sound.
*/
bool shouldPlaySound(MatrixEvent *event) const;
/**
* Get the sticker pack list for the current account.
*
* @return A list of sticker packs associated with the current account.
*/
MatrixStickerPackList *stickerPackList() const;
/**
* Get the sticker rooms account data event for the current account.
*
* @return A MatrixEvent representing the sticker rooms account data event.
*/
MatrixEvent *stickerRoomsEvent() const;
/**
* Update the sticker pack from source.
*
* @param source The source of the sticker pack to update.
* @return A promise that resolves when the sticker pack is updated,
* or when there is an error.
*/
MatrixPromise *updateStickerPack(MatrixStickerPackSource source);
MatrixUserGivenAttrsMap *userGivenNicknameMap() const;
MatrixPromise *sendAccountData(const QString &type, const QJsonObject &content);
/**
* Get all Matrix Spec versions supported by the server.
* Use MatrixSdk::supportSpecVersion() to check if a version is supported.
*/
MatrixPromise *getSpecVersions();
MatrixPromise *addDirectRoom(const QString &userId, const QString &roomId);
MatrixPromise *getRoomIdByAlias(const QString &roomAlias);
private:
MatrixPromise *sendAccountDataImpl(Kazv::Event event);
private: // Testing
friend MatrixSdkTest;
friend MatrixSdkSessionsTest;
friend MatrixSdk *makeTestSdk(Kazv::SdkModel model);
void setUserDataDir(const std::string &userDataDir);
explicit MatrixSdk(Kazv::SdkModel model, bool testing = false, QObject *parent = 0);
void startThread();
};
diff --git a/src/tests/db-store-test.cpp b/src/tests/db-store-test.cpp
index 1002678..bb1a4d2 100644
--- a/src/tests/db-store-test.cpp
+++ b/src/tests/db-store-test.cpp
@@ -1,67 +1,118 @@
/*
* This file is part of kazv.
* SPDX-FileCopyrightText: 2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <kazv-defs.hpp>
#include "db-store.hpp"
#include <QtTest>
#include <QTemporaryDir>
#include <QCoroCore>
#include <factory.hpp>
using namespace Kazv::Factory;
using namespace Kazv;
+using namespace Qt::Literals::StringLiterals;
class DbStoreTest : public QObject
{
Q_OBJECT
private Q_SLOTS:
void init();
void testSetup();
void testImportAllFrom();
+ void testResultToEvent();
+ void testSaveEvents();
void cleanup();
private:
QTemporaryDir m_tempDir;
};
void DbStoreTest::init()
{
m_tempDir = QTemporaryDir();
}
void DbStoreTest::cleanup()
{
}
void DbStoreTest::testSetup()
{
DbStore dbStore;
auto stat = QCoro::waitFor(dbStore.setup(m_tempDir.path().toStdString(), "@userid:example.com", "device1"));
QVERIFY(stat.first);
}
void DbStoreTest::testImportAllFrom()
{
DbStore dbStore;
auto r1 = makeRoom(withRoomTimeline({makeEvent(), makeEvent()}));
auto r2 = makeRoom(withRoomTimeline({makeEvent(), makeEvent(), makeEvent()}));
auto m = SdkModel(makeClient(
withRoom(r1) | withRoom(r2)
));
auto stat = QCoro::waitFor([](
std::string tempDir, DbStore &dbStore, SdkModel m
) -> QCoro::Task<std::pair<bool, QString>> {
auto res = co_await dbStore.setup(tempDir, "@userid:example.com", "device1");
if (!res.first) { co_return res; }
co_return co_await dbStore.importAllFrom(m);
}(m_tempDir.path().toStdString(), dbStore, m));
QVERIFY(stat.first);
}
+void DbStoreTest::testResultToEvent()
+{
+ {
+ Event e = makeEvent(withEventType("m.room.encrypted"))
+ .setDecryptedJson(json{{"content", {{"text", "mew"}}}}, Event::Decrypted);
+ QVERIFY(e == DbStore::resultToEvent(
+ e.encrypted(),
+ e.decrypted(),
+ QString::fromStdString(e.originalJson().get().dump()),
+ QString::fromStdString(e.decryptedJson().get().dump())
+ ));
+ }
+ {
+ Event e = makeEvent();
+ QVERIFY(e == DbStore::resultToEvent(
+ e.encrypted(),
+ e.decrypted(),
+ QString::fromStdString(e.originalJson().get().dump()),
+ QVariant().toString()
+ ));
+ }
+}
+
+void DbStoreTest::testSaveEvents()
+{
+ DbStore dbStore;
+ {
+ auto stat = QCoro::waitFor(dbStore.setup(m_tempDir.path().toStdString(), "@userid:example.com", "device1"));
+ QVERIFY(stat.first);
+ }
+ auto e1 = makeEvent();
+ auto e2 = makeEvent();
+ auto e3 = makeEvent();
+ SaveEventsRequested s{
+ {{"!room1:example.com", {e1, e2}}},
+ {{"!room2:example.com", {e3}}},
+ };
+ {
+ auto stat = QCoro::waitFor(dbStore.saveEvents(s));
+ QVERIFY(stat.first);
+ }
+ auto res = QCoro::waitFor(dbStore.getEventById(u"!room1:example.com"_s, QString::fromStdString(e1.id())));
+ QVERIFY(res.has_value() && res.value().first == e1 && res.value().second);
+
+ res = QCoro::waitFor(dbStore.getEventById(u"!room2:example.com"_s, QString::fromStdString(e3.id())));
+ QVERIFY(res.has_value() && res.value().first == e3 && !res.value().second);
+}
+
QTEST_MAIN(DbStoreTest)
#include "db-store-test.moc"

File Metadata

Mime Type
text/x-diff
Expires
Thu, Oct 8, 6:39 PM (6 h, 26 m)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1784198
Default Alt Text
(63 KB)

Event Timeline