Page MenuHomePhorge

No OneTemporary

Size
155 KB
Referenced Files
None
Subscribers
None
diff --git a/src/base/eventinterface.hpp b/src/base/eventinterface.hpp
index 054f166..2f82689 100644
--- a/src/base/eventinterface.hpp
+++ b/src/base/eventinterface.hpp
@@ -1,20 +1,20 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include "libkazv-config.hpp"
-#include "kazvevents.hpp"
+#include "kazv-triggers.hpp"
namespace Kazv
{
class EventInterface
{
public:
virtual ~EventInterface() = default;
- virtual void emit(KazvEvent e) = 0;
+ virtual void emit(KazvTrigger e) = 0;
};
}
diff --git a/src/base/kazv-triggers.hpp b/src/base/kazv-triggers.hpp
new file mode 100644
index 0000000..56a146e
--- /dev/null
+++ b/src/base/kazv-triggers.hpp
@@ -0,0 +1,104 @@
+/*
+ * This file is part of libkazv.
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
+ * SPDX-License-Identifier: AGPL-3.0-or-later
+ */
+
+
+#pragma once
+#include "libkazv-config.hpp"
+#include <variant>
+#include "types.hpp"
+#include "event.hpp"
+#include "basejob.hpp"
+
+namespace Kazv
+{
+ struct LoginSuccessful {};
+
+ struct LoginFailed
+ {
+ std::string errorCode;
+ std::string error;
+ };
+
+ struct ReceivingPresenceEvent { Event event; };
+
+ struct ReceivingAccountDataEvent { Event event; };
+
+ struct ReceivingRoomStateEvent {
+ Event event;
+ std::string roomId;
+ };
+
+ struct ReceivingRoomTimelineEvent {
+ Event event;
+ std::string roomId;
+ };
+
+ struct ReceivingRoomAccountDataEvent {
+ Event event;
+ std::string roomId;
+ };
+
+ struct ReceivingToDeviceMessage
+ {
+ Event event;
+ };
+
+ struct RoomMembershipChanged {
+ RoomMembership membership;
+ std::string roomId;
+ };
+
+ /**
+ * Indicate that there are events to be saved.
+ */
+ struct SaveEventsRequested
+ {
+ /**
+ * The events to be saved in the timeline.
+ * The events saved in the timeline part of the storage can be considered continuous, i.e. they can be
+ * loaded to fully re-construct the timeline. There can be gaps,
+ * but all gaps are properly recorded in the client state so that
+ * it can be easily paginated.
+ */
+ immer::map<std::string /* roomId */, EventList> timelineEvents;
+ /**
+ * The events that should be saved, but not in the timeline.
+ */
+ immer::map<std::string /* roomId */, EventList> nonTimelineEvents;
+ };
+
+ struct UnrecognizedResponse
+ {
+ Response response;
+ };
+
+ using KazvTrigger = std::variant<
+ // use this for placeholder of "no events yet"
+ // otherwise the first LoginSuccessful event cannot be detected
+ std::monostate,
+
+ // matrix events
+ ReceivingPresenceEvent,
+ ReceivingAccountDataEvent,
+ ReceivingRoomTimelineEvent,
+ ReceivingRoomStateEvent,
+ RoomMembershipChanged,
+ ReceivingRoomAccountDataEvent,
+ ReceivingToDeviceMessage,
+
+ // auth
+ LoginSuccessful, LoginFailed,
+ // storage
+ SaveEventsRequested,
+
+ // general
+ UnrecognizedResponse
+ >;
+ using KazvEvent [[deprecated("renamed to KazvTrigger")]] = KazvTrigger;
+
+ using KazvTriggerList = immer::flex_vector<KazvTrigger>;
+ using KazvEventList [[deprecated("renamed to KazvTriggerList")]] = KazvTriggerList;
+}
diff --git a/src/base/kazvevents.hpp b/src/base/kazvevents.hpp
index a4c6742..3a8b8e1 100644
--- a/src/base/kazvevents.hpp
+++ b/src/base/kazvevents.hpp
@@ -1,414 +1,7 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu
+ * SPDX-FileCopyrightText: 2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
-
-#pragma once
-#include "libkazv-config.hpp"
-#include <variant>
-#include "types.hpp"
-#include "event.hpp"
-#include "basejob.hpp"
-
-namespace Kazv
-{
- struct LoginSuccessful {};
-
- struct LoginFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- struct SyncSuccessful
- {
- std::string nextToken;
- };
-
- struct SyncFailed
- {
- };
-
- struct PostInitialFiltersSuccessful
- {
- };
-
- struct PostInitialFiltersFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- struct ReceivingPresenceEvent { Event event; };
-
- struct ReceivingAccountDataEvent { Event event; };
-
- struct ReceivingRoomStateEvent {
- Event event;
- std::string roomId;
- };
-
- struct ReceivingRoomTimelineEvent {
- Event event;
- std::string roomId;
- };
-
- struct ReceivingRoomAccountDataEvent {
- Event event;
- std::string roomId;
- };
-
- struct ReceivingToDeviceMessage
- {
- Event event;
- };
-
- struct RoomMembershipChanged {
- RoomMembership membership;
- std::string roomId;
- };
-
- struct PaginateSuccessful
- {
- std::string roomId;
- };
-
- struct PaginateFailed
- {
- std::string roomId;
- };
-
- struct CreateRoomSuccessful
- {
- std::string roomId;
- };
-
- struct CreateRoomFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- struct InviteUserSuccessful
- {
- std::string roomId;
- std::string userId;
- };
-
- struct InviteUserFailed
- {
- std::string roomId;
- std::string userId;
- std::string errorCode;
- std::string error;
- };
-
- struct JoinRoomSuccessful
- {
- std::string roomIdOrAlias;
- };
-
- struct JoinRoomFailed
- {
- std::string roomIdOrAlias;
- std::string errorCode;
- std::string error;
- };
-
- struct LeaveRoomSuccessful
- {
- std::string roomId;
- };
-
- struct LeaveRoomFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct ForgetRoomSuccessful
- {
- std::string roomId;
- };
-
- struct ForgetRoomFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct SendMessageSuccessful
- {
- std::string roomId;
- std::string eventId;
- };
-
- struct SendMessageFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct SendToDeviceMessageSuccessful
- {
- immer::map<std::string, immer::flex_vector<std::string>> devicesToSend;
- std::string txnId;
- };
-
- struct SendToDeviceMessageFailed
- {
- immer::map<std::string, immer::flex_vector<std::string>> devicesToSend;
- std::string txnId;
- std::string errorCode;
- std::string error;
- };
-
- struct InvalidMessageFormat
- {
- };
-
- struct GetRoomStatesSuccessful
- {
- std::string roomId;
- };
-
- struct GetRoomStatesFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct GetStateEventSuccessful
- {
- std::string roomId;
- JsonWrap content;
- };
-
- struct GetStateEventFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct SendStateEventSuccessful
- {
- std::string roomId;
- std::string eventId;
- std::string eventType;
- std::string stateKey;
- };
-
- struct SendStateEventFailed
- {
- std::string roomId;
- std::string eventType;
- std::string stateKey;
- std::string errorCode;
- std::string error;
- };
-
- struct SetTypingSuccessful
- {
- std::string roomId;
- };
-
- struct SetTypingFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct PostReceiptSuccessful
- {
- std::string roomId;
- };
-
- struct PostReceiptFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct SetReadMarkerSuccessful
- {
- std::string roomId;
- };
-
- struct SetReadMarkerFailed
- {
- std::string roomId;
- std::string errorCode;
- std::string error;
- };
-
- struct UploadContentSuccessful
- {
- std::string mxcUri;
- std::string uploadId;
- };
-
- struct UploadContentFailed
- {
- std::string uploadId;
- std::string errorCode;
- std::string error;
- };
-
- struct DownloadContentSuccessful
- {
- std::string mxcUri;
- immer::box<Bytes> content;
- std::optional<std::string> filename;
- std::optional<std::string> contentType;
- };
-
- struct DownloadContentFailed
- {
- std::string mxcUri;
- std::string errorCode;
- std::string error;
- };
-
- struct DownloadThumbnailSuccessful
- {
- std::string mxcUri;
- immer::box<Bytes> content;
- std::optional<std::string> contentType;
- };
-
- struct DownloadThumbnailFailed
- {
- std::string mxcUri;
- std::string errorCode;
- std::string error;
- };
-
- struct UploadIdentityKeysSuccessful
- {
- };
-
- struct UploadIdentityKeysFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- struct UploadOneTimeKeysSuccessful
- {
- };
-
- struct UploadOneTimeKeysFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- struct ClaimKeysSuccessful
- {
- Event eventToSend;
- immer::map<std::string, immer::flex_vector<std::string>> devicesToSend;
- };
-
- struct ClaimKeysFailed
- {
- std::string errorCode;
- std::string error;
- };
-
- /**
- * Indicate that there are events to be saved.
- */
- struct SaveEventsRequested
- {
- /**
- * The events to be saved in the timeline.
- * The events saved in the timeline part of the storage can be considered continuous, i.e. they can be
- * loaded to fully re-construct the timeline. There can be gaps,
- * but all gaps are properly recorded in the client state so that
- * it can be easily paginated.
- */
- immer::map<std::string /* roomId */, EventList> timelineEvents;
- /**
- * The events that should be saved, but not in the timeline.
- */
- immer::map<std::string /* roomId */, EventList> nonTimelineEvents;
- };
-
- struct UnrecognizedResponse
- {
- Response response;
- };
-
-
- struct ShouldQueryKeys
- {
- bool isInitialSync;
- };
-
- using KazvEvent = std::variant<
- // use this for placeholder of "no events yet"
- // otherwise the first LoginSuccessful event cannot be detected
- std::monostate,
-
- // matrix events
- ReceivingPresenceEvent,
- ReceivingAccountDataEvent,
- ReceivingRoomTimelineEvent,
- ReceivingRoomStateEvent,
- RoomMembershipChanged,
- ReceivingRoomAccountDataEvent,
- ReceivingToDeviceMessage,
-
- // auth
- LoginSuccessful, LoginFailed,
- // sync
- SyncSuccessful, SyncFailed,
- PostInitialFiltersSuccessful, PostInitialFiltersFailed,
- // paginate
- PaginateSuccessful, PaginateFailed,
- // membership
- CreateRoomSuccessful, CreateRoomFailed,
- InviteUserSuccessful, InviteUserFailed,
- JoinRoomSuccessful, JoinRoomFailed,
- LeaveRoomSuccessful, LeaveRoomFailed,
- ForgetRoomSuccessful, ForgetRoomFailed,
- // send
- SendMessageSuccessful, SendMessageFailed,
- SendToDeviceMessageSuccessful, SendToDeviceMessageFailed,
- InvalidMessageFormat,
- // states
- GetRoomStatesSuccessful, GetRoomStatesFailed,
- GetStateEventSuccessful, GetStateEventFailed,
- SendStateEventSuccessful, SendStateEventFailed,
- // ephemeral
- SetTypingSuccessful, SetTypingFailed,
- PostReceiptSuccessful, PostReceiptFailed,
- SetReadMarkerSuccessful, SetReadMarkerFailed,
- // content
- UploadContentSuccessful, UploadContentFailed,
- DownloadContentSuccessful, DownloadContentFailed,
- DownloadThumbnailSuccessful, DownloadThumbnailFailed,
- // encryption
- UploadIdentityKeysSuccessful, UploadIdentityKeysFailed,
- UploadOneTimeKeysSuccessful, UploadOneTimeKeysFailed,
- ClaimKeysSuccessful, ClaimKeysFailed,
- // storage
- SaveEventsRequested,
-
- // general
- UnrecognizedResponse,
- ShouldQueryKeys
-
- >;
-
- using KazvEventList = immer::flex_vector<KazvEvent>;
-}
+#include "kazv-triggers.hpp"
diff --git a/src/client/actions/content.cpp b/src/client/actions/content.cpp
index 7ae7bb3..c969d67 100644
--- a/src/client/actions/content.cpp
+++ b/src/client/actions/content.cpp
@@ -1,136 +1,132 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020 Tusooa Zhu <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <boost/algorithm/string.hpp>
#include "status-utils.hpp"
#include "content.hpp"
namespace Kazv
{
std::pair<std::string, std::string> mxcUriToMediaDesc(std::string mxcUri)
{
using namespace boost::algorithm;
if (mxcUri.find("mxc://") != 0) {
return {"", ""};
}
mxcUri.erase(0, 6); // remove "mxc://"
std::vector<std::string> a;
split(a, mxcUri, is_any_of("/"));
if (a.size() != 2) {
return {"", ""};
}
return { a[0], a[1] };
}
ClientResult updateClient(ClientModel m, UploadContentAction a)
{
auto job = m.job<UploadContentJob>()
- .make(a.content, a.filename, a.contentType)
- .withData(json{ {"uploadId", a.uploadId} });
+ .make(a.content, a.filename, a.contentType);
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, UploadContentResponse r)
{
- auto uploadId = r.dataStr("uploadId");
if (! r.success()) {
- m.addTrigger(UploadContentFailed{uploadId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
auto mxcUri = r.contentUri();
- m.addTrigger(UploadContentSuccessful{mxcUri, uploadId});
return {
std::move(m),
[=](auto &&) {
return EffectStatus{/* succ = */ true, json{{"mxcUri", mxcUri}}};
}
};
}
ClientResult updateClient(ClientModel m, DownloadContentAction a)
{
auto [serverName, mediaId] = mxcUriToMediaDesc(a.mxcUri);
auto data = json{
{"mxcUri", a.mxcUri},
{"streaming", false},
};
if (a.downloadTo.has_value() && a.downloadTo.value().name().has_value()) {
data["streaming"] = true;
}
auto job = m.job<GetContentJob>()
.make(serverName, mediaId, /* allowRemote = */ true, a.downloadTo)
.withData(std::move(data));
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, GetContentResponse r)
{
auto mxcUri = r.dataStr("mxcUri");
if (! r.success()) {
return { std::move(m), failWithResponse(r) };
}
auto isStreaming = r.dataJson("streaming").template get<bool>();
return {
std::move(m),
[=](auto &&) {
auto data = json::object();
if (! isStreaming) {
data["content"] = std::get<Bytes>(std::move(r).data());
}
return EffectStatus{/* succ = */ true, data};
}
};
}
ClientResult updateClient(ClientModel m, DownloadThumbnailAction a)
{
auto [serverName, mediaId] = mxcUriToMediaDesc(a.mxcUri);
std::optional<std::string> method;
if (a.method) {
method = a.method.value() == Crop ? "crop" : "scale";
}
auto job = m.job<GetContentThumbnailJob>()
.make(serverName, mediaId, a.width, a.height, method, /* allowRemote = */ true, a.downloadTo)
.withData(json{ {"mxcUri", a.mxcUri},
{"streaming", a.downloadTo.has_value() && a.downloadTo.value().name().has_value()} });
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, GetContentThumbnailResponse r)
{
auto mxcUri = r.dataStr("mxcUri");
if (! r.success()) {
return { std::move(m), failWithResponse(r) };
}
auto isStreaming = r.dataJson("streaming").template get<bool>();
return {
std::move(m),
[=](auto &&) {
auto data = json::object();
if (! isStreaming) {
data["content"] = std::get<Bytes>(std::move(r).data());
}
return EffectStatus{/* succ = */ true, data};
}
};
}
}
diff --git a/src/client/actions/encryption.cpp b/src/client/actions/encryption.cpp
index dd5d3ac..3545596 100644
--- a/src/client/actions/encryption.cpp
+++ b/src/client/actions/encryption.cpp
@@ -1,718 +1,712 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2021-2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <zug/transducer/filter.hpp>
#include <zug/transducer/cat.hpp>
#include "encryption.hpp"
#include <immer-utils.hpp>
#include <debug.hpp>
#include "cursorutil.hpp"
#include "status-utils.hpp"
#include "key-export.hpp"
namespace Kazv
{
using namespace CryptoConstants;
static json convertSignature(const ClientModel &m, std::string signature)
{
auto j = json::object();
j[m.userId] = json::object();
j[m.userId][ed25519 + ":" + m.deviceId] = signature;
return j;
}
ClientResult updateClient(ClientModel m, UploadIdentityKeysAction)
{
if (! m.crypto) {
kzo.client.warn() << "Client::crypto is invalid, ignoring it." << std::endl;
return { std::move(m), lager::noop };
}
auto keys =
immer::map<std::string, std::string>{}
.set(ed25519 + ":" + m.deviceId, m.constCrypto().ed25519IdentityKey())
.set(curve25519 + ":" + m.deviceId, m.constCrypto().curve25519IdentityKey());
DeviceKeys k {
m.userId,
m.deviceId,
{olmAlgo, megOlmAlgo},
keys,
{} // signatures to be added soon
};
auto j = json(k);
auto sig = m.withCrypto([&](auto &crypto) { return crypto.sign(j); });
k.signatures = convertSignature(m, sig);
auto job = m.job<UploadKeysJob>()
.make(k)
.withData(json{{"is", "identityKeys"}});
kzo.client.dbg() << "Uploading identity keys" << std::endl;
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, GenerateAndUploadOneTimeKeysAction a)
{
if (! m.crypto) {
kzo.client.warn() << "Client::crypto is invalid, ignoring it." << std::endl;
return { std::move(m), simpleFail };
}
kzo.client.dbg() << "Generating " << a.numToGen << " one-time keys..." << std::endl;
auto maxNumKeys = m.constCrypto().maxNumberOfOneTimeKeys();
auto numLocalKeys = m.constCrypto().numUnpublishedOneTimeKeys();
auto numStoredKeys = m.constCrypto().uploadedOneTimeKeysCount(signedCurve25519) + numLocalKeys;
auto numKeysToGenerate = a.numToGen;
auto genKeysLimit = maxNumKeys - numStoredKeys;
if (numKeysToGenerate > genKeysLimit) {
numKeysToGenerate = genKeysLimit;
}
if (numLocalKeys <= 0 && numKeysToGenerate <= 0) { // we have enough already
kzo.client.dbg() << "We have enough one-time keys. Ignoring this." << std::endl;
return { std::move(m), lager::noop };
}
if (numKeysToGenerate > 0) {
m.withCrypto([&](auto &c) { c.genOneTimeKeysWithRandom(a.random, numKeysToGenerate); });
}
kzo.client.dbg() << "Generating done." << std::endl;
auto keys = m.constCrypto().unpublishedOneTimeKeys();
auto cv25519Keys = keys.at(curve25519);
json oneTimeKeys = json::object();
for (auto [id, keyStr] : cv25519Keys.items()) {
json keyObject = json::object();
keyObject["key"] = keyStr;
keyObject["signatures"] = convertSignature(m, m.withCrypto([&](auto &c) { return c.sign(keyObject); }));
oneTimeKeys[signedCurve25519 + ":" + id] = keyObject;
}
auto job = m.job<UploadKeysJob>()
.make(
std::nullopt, // deviceKeys
oneTimeKeys)
.withData(json{{"is", "oneTimeKeys"}});
kzo.client.dbg() << "Uploading one time keys" << std::endl;
m.addJob(std::move(job));
return { std::move(m), lager::noop };
};
ClientResult processResponse(ClientModel m, UploadKeysResponse r)
{
if (! m.crypto) {
kzo.client.warn() << "Client::crypto is invalid, ignoring it." << std::endl;
return { std::move(m), lager::noop };
}
auto is = r.dataStr("is");
if (is == "identityKeys") {
if (! r.success()) {
kzo.client.dbg() << "Uploading identity keys failed" << std::endl;
- m.addTrigger(UploadIdentityKeysFailed{r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
kzo.client.dbg() << "Uploading identity keys successful" << std::endl;
- m.addTrigger(UploadIdentityKeysSuccessful{});
m.identityKeysUploaded = true;
} else {
if (! r.success()) {
kzo.client.dbg() << "Uploading one-time keys failed" << std::endl;
- m.addTrigger(UploadOneTimeKeysFailed{r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
kzo.client.dbg() << "Uploading one-time keys successful" << std::endl;
- m.addTrigger(UploadOneTimeKeysSuccessful{});
-
m.withCrypto([&](auto &c) { c.markOneTimeKeysAsPublished(); });
}
m.withCrypto([&](auto &c) { c.setUploadedOneTimeKeysCount(r.oneTimeKeyCounts()); });
return { std::move(m), lager::noop };
}
static JsonWrap cannotDecryptEvent(
const std::string &reason,
const std::string &errcode,
const json &raw)
{
return json{
{"type", "m.room.message"},
{"content", {
{"msgtype","moe.kazv.mxc.cannot.decrypt"},
{"body", "**This message cannot be decrypted due to " + reason + ".**"},
{"moe.kazv.mxc.error", reason},
{"moe.kazv.mxc.errcode", errcode},
{"moe.kazv.mxc.raw", raw},
}},
};
}
// returns std::nullopt on success, and an error event on failure
static std::optional<JsonWrap> verifyEvent(ClientModel &m, Event e, const json &plainJson)
{
try {
std::string algo = e.originalJson().get().at("content").at("algorithm");
if (algo == olmAlgo) {
std::string senderCurve25519Key = e.originalJson().get()
.at("content").at("sender_key");
auto deviceInfoOpt = m.deviceLists.findByCurve25519Key(e.sender(), senderCurve25519Key);
if (! deviceInfoOpt) {
kzo.client.dbg() << "Device key " << senderCurve25519Key
<< " unknown, thus invalid" << std::endl;
return cannotDecryptEvent(
"device key unknown",
"MOE.KAZV.MXC_DEVICE_KEY_UNKNOWN",
plainJson
);
}
auto deviceInfo = deviceInfoOpt.value();
if (! (plainJson.at("sender") == e.sender())) {
kzo.client.dbg() << "Sender does not match, thus invalid" << std::endl;
return cannotDecryptEvent(
"sender does not match",
"MOE.KAZV.MXC_BAD_SENDER",
plainJson
);
}
if (! (plainJson.at("recipient") == m.userId)) {
kzo.client.dbg() << "Recipient does not match, thus invalid" << std::endl;
return cannotDecryptEvent(
"recipient does not match",
"MOE.KAZV.MXC_BAD_RECIPIENT",
plainJson
);
}
if (! (plainJson.at("recipient_keys").at(ed25519) == m.constCrypto().ed25519IdentityKey())) {
kzo.client.dbg() << "Recipient key does not match, thus invalid" << std::endl;
return cannotDecryptEvent(
"recipient keys do not match",
"MOE.KAZV.MXC_BAD_RECIPIENT_KEYS",
plainJson
);
}
auto thisEd25519Key = plainJson.at("keys").at(ed25519).get<std::string>();
if (thisEd25519Key != deviceInfo.ed25519Key) {
kzo.client.dbg() << "Sender ed25519 key does not match, thus invalid" << std::endl;
return cannotDecryptEvent(
"sender keys do not match",
"MOE.KAZV.MXC_BAD_SENDER_KEYS",
plainJson
);
}
} else if (algo == megOlmAlgo) {
auto roomId = plainJson.at("room_id").get<std::string>();
if (roomId.empty() ||
roomId != e.originalJson().get().at("room_id").template get<std::string>()) {
kzo.client.dbg() << "Room id does not match, thus invalid" << std::endl;
return cannotDecryptEvent(
"room id does not match",
"MOE.KAZV.MXC_BAD_ROOM_ID",
plainJson
);
}
} else {
kzo.client.dbg() << "Unknown algorithm, thus invalid" << std::endl;
return cannotDecryptEvent(
"unknown algorithm",
"MOE.KAZV.MXC_UNKNOWN_ALGORITHM",
plainJson
);
}
} catch (const std::exception &exception) {
kzo.client.dbg() << "json format is not correct, thus invalid" << std::endl;
return cannotDecryptEvent(
exception.what(),
"M_BAD_JSON",
plainJson
);
}
return std::nullopt;
}
static Event decryptEvent(ClientModel &m, Event e)
{
// no need for decryption
if (e.decrypted() || (! e.encrypted())) {
return e;
}
kzo.client.dbg() << "About to decrypt event: "
<< e.id() << std::endl;
auto maybePlainText = m.withCrypto([&](Crypto &c) {
return c.decrypt(e.originalJson().get());
});
if (! maybePlainText) {
kzo.client.dbg() << "Cannot decrypt: " << maybePlainText.reason() << std::endl;
return e.setDecryptedJson(
cannotDecryptEvent(
maybePlainText.reason(),
"MOE.KAZV.MXC_DECRYPT_ERROR",
json(nullptr)
),
Event::NotDecrypted);
} else {
try {
auto plainJson = json::parse(maybePlainText.value());
auto error = verifyEvent(m, e, plainJson);
auto valid = !error.has_value();
if (valid) {
kzo.client.dbg() << "The decrypted event is valid." << std::endl;
}
return valid
? e.setDecryptedJson(plainJson, Event::Decrypted)
: e.setDecryptedJson(
error.value(),
Event::NotDecrypted);
} catch (const std::exception &exception) {
return e.setDecryptedJson(
cannotDecryptEvent(
exception.what(),
"M_NOT_JSON",
maybePlainText.value()
),
Event::NotDecrypted
);
}
}
}
ClientModel tryDecryptEvents(ClientModel m)
{
if (! m.crypto) {
kzo.client.dbg() << "We have no encryption enabled--ignoring decryption request" << std::endl;
return m;
}
kzo.client.dbg() << "Trying to decrypt events..." << std::endl;
auto decryptFunc = [&](auto e) { return decryptEvent(m, e); };
auto takeOutRoomKeyEvents =
[&](auto e) {
if (e.type() != "m.room_key") {
// Leave it as it is
return true;
}
// It is a room key event, but unencrypted.
// Per matrix spec, we should not trust it as a E2EE key.
// matrix spec also says we should make sure it's Olm-encrypted.
// This is realized by verifying all MegOlm-encrypted events have
// a room_id.
if (!e.encrypted()) {
kzo.client.warn() << "Received an unencrypted room key event. Ignoring." << std::endl;
return false;
}
try {
auto content = e.content();
std::string roomId = content.get().at("room_id");
std::string sessionId = content.get().at("session_id");
kzo.client.dbg() << "Got a room key for room " << roomId
<< ", session id: " << sessionId << std::endl;
std::string sessionKey = content.get().at("session_key");
auto k = KeyOfGroupSession{roomId, sessionId};
std::string ed25519Key = e.decryptedJson().get().at("keys").at(ed25519);
if (m.withCrypto([&](auto &c) { return c.createInboundGroupSession(k, sessionKey, ed25519Key); })) {
return false; // such that this event is removed
} else {
kzo.client.warn() << "The session exists and cannot be merged. Someone is trying to do a session-replace attack." << std::endl;
kzo.client.dbg() << "sender key is " << ed25519Key << std::endl;
return true;
}
} catch (...) {
kzo.client.dbg() << "cannot create group session";
return false;
}
return true;
};
m.toDevice = intoImmer(
EventList{},
zug::map(decryptFunc)
| zug::filter(takeOutRoomKeyEvents),
std::move(m.toDevice));
auto decryptEventInRoom =
[&](auto id, auto room) {
if (! room.encrypted) {
return;
} else {
auto messages = room.messages;
auto undecryptedEvents = room.undecryptedEvents;
for (auto [sessionId, eventIds] : undecryptedEvents) {
if (m.constCrypto().hasInboundGroupSession({
room.roomId,
sessionId,
})) {
auto nextEventIds = intoImmer(
immer::flex_vector<std::string>{},
zug::filter([&](auto eventId) {
auto event = room.messages[eventId];
auto decrypted = decryptFunc(event);
room.messages = std::move(room.messages)
.set(eventId, decrypted);
return !decrypted.decrypted();
}),
eventIds
);
if (nextEventIds.empty()) {
room.undecryptedEvents = std::move(room.undecryptedEvents).erase(sessionId);
} else {
room.undecryptedEvents = std::move(room.undecryptedEvents).set(sessionId, nextEventIds);
}
}
}
m.roomList.rooms = std::move(m.roomList.rooms).set(id, room);
}
};
auto rooms = m.roomList.rooms;
for (auto [id, room]: rooms) {
decryptEventInRoom(id, room);
}
return m;
}
std::optional<BaseJob> clientPerform(ClientModel m, QueryKeysAction a)
{
if (! m.crypto) {
kzo.client.dbg() << "We have no encryption enabled--ignoring this" << std::endl;
return std::nullopt;
}
immer::map<std::string, immer::array<std::string>> deviceKeys;
auto encryptedUsers = m.deviceLists.outdatedUsers();
if (encryptedUsers.empty()) {
kzo.client.dbg() << "Keys are up-to-date." << std::endl;
return std::nullopt;
}
kzo.client.dbg() << "We need to query keys for: " << std::endl;
for (auto userId: encryptedUsers) {
kzo.client.dbg() << userId << std::endl;
deviceKeys = std::move(deviceKeys).set(userId, {});
}
kzo.client.dbg() << "^" << std::endl;
auto job = m.job<QueryKeysJob>()
.make(std::move(deviceKeys),
std::nullopt, // timeout
a.isInitialSync ? std::nullopt : m.syncToken
);
return job;
}
ClientResult updateClient(ClientModel m, QueryKeysAction a)
{
auto jobOpt = clientPerform(m, a);
if (jobOpt) {
m.addJob(jobOpt.value());
}
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, QueryKeysResponse r)
{
if (! m.crypto) {
kzo.client.dbg() << "We have no encryption enabled--ignoring this" << std::endl;
return { std::move(m), simpleFail };
}
if (! r.success()) {
kzo.client.dbg() << "query keys failed: " << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
kzo.client.dbg() << "Received a query key response" << std::endl;
auto usersMap = r.deviceKeys();
for (auto [userId, deviceMap] : usersMap) {
for (auto [deviceId, deviceInfo] : deviceMap) {
kzo.client.dbg() << "Key for " << userId
<< "/" << deviceId
<< ": " << json(deviceInfo).dump()
<< std::endl;
m.withCrypto([&](Crypto &c) {
m.deviceLists.addDevice(userId, deviceId, deviceInfo, c);
});
}
m.deviceLists.markUpToDate(userId);
}
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, ClaimKeysAction a)
{
if (! m.crypto) {
kzo.client.dbg() << "We have no encryption enabled--ignoring this" << std::endl;
return { std::move(m), lager::noop };
}
kzo.client.dbg() << "claim keys for: " << json(a.devicesToSend).dump() << std::endl;
auto keyMap = immer::map<std::string, immer::map<std::string /* deviceId */,
std::string /* curve25519IdentityKey */>>{};
for (auto [userId, devices] : a.devicesToSend) {
kzo.client.dbg() << "Iterating through user " << userId << std::endl;
auto deviceToKey = immer::map<std::string, std::string>{};
for (auto deviceId : devices) {
kzo.client.dbg() << "Device: " << deviceId << std::endl;
auto infoOpt = m.deviceLists.get(userId, deviceId);
if (infoOpt) {
kzo.client.dbg() << "Got device info, curve25519 key is: " << infoOpt.value().curve25519Key << std::endl;
deviceToKey = std::move(deviceToKey)
.set(deviceId, infoOpt.value().curve25519Key);
} else {
kzo.client.dbg() << "Did not get device info" << std::endl;
}
}
keyMap = std::move(keyMap).set(userId, deviceToKey);
}
auto devicesToClaimKeys = m.withCrypto([&](auto &c) { return c.devicesMissingOutboundSessionKey(keyMap); });
kzo.client.dbg() << "Really claim keys for: " << json(devicesToClaimKeys).dump() << std::endl;
auto oneTimeKeys = immer::map<std::string, immer::map<std::string, std::string>>{};
for (auto [userId, devices] : devicesToClaimKeys) {
auto devKeys = immer::map<std::string, std::string>{};
for (auto deviceId: devices) {
devKeys = std::move(devKeys).set(deviceId, signedCurve25519);
}
oneTimeKeys = std::move(oneTimeKeys).set(userId, devKeys);
}
auto job = m.job<ClaimKeysJob>()
.make(std::move(oneTimeKeys))
.withData(json{
{"roomId", a.roomId},
{"sessionId", a.sessionId},
{"sessionKey", a.sessionKey},
{"devicesToSend", a.devicesToSend},
{"random", a.random}
});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, ClaimKeysResponse r)
{
if (! m.crypto) {
kzo.client.dbg() << "We have no encryption enabled--ignoring this" << std::endl;
return { std::move(m), simpleFail };
}
if (! r.success()) {
kzo.client.dbg() << "claim keys failed" << std::endl;
- m.addTrigger(ClaimKeysFailed{r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
kzo.client.dbg() << "claim keys successful" << std::endl;
kzo.client.dbg() << "Json body: " << r.jsonBody().get().dump() << std::endl;
auto roomId = r.dataStr("roomId");
auto sessionKey = r.dataStr("sessionKey");
auto sessionId = r.dataStr("sessionId");
auto devicesToSend =
immer::map<std::string, immer::flex_vector<std::string>>(r.dataJson("devicesToSend"));
auto random = r.dataJson("random").template get<RandomData>();
// create outbound sessions for those devices
auto oneTimeKeys = r.oneTimeKeys();
for (auto [userId, deviceMap] : oneTimeKeys) {
for (auto [deviceId, keyVar] : deviceMap) {
auto keys = keyVar.get();
for (auto [keyId, key] : keys.items()) {
auto deviceInfoOpt = m.deviceLists.get(userId, deviceId);
if (deviceInfoOpt) {
auto deviceInfo = deviceInfoOpt.value();
kzo.client.dbg() << "Verifying key for " << userId
<< "/" << deviceId
<< key.dump()
<< " with ed25519 key "
<< deviceInfo.ed25519Key << std::endl;
auto verified = m.withCrypto([&](auto &c) { return c.verify(key, userId, deviceId, deviceInfo.ed25519Key); });
kzo.client.dbg() << (verified ? "passed" : "did not pass") << std::endl;
if (verified && key.contains("key")) {
auto theirOneTimeKey = key.at("key");
kzo.client.dbg() << "creating outbound session for it" << std::endl;
m.withCrypto([&](auto &c) { c.createOutboundSessionWithRandom(random, deviceInfo.curve25519Key, theirOneTimeKey); });
random.erase(0, Crypto::createOutboundSessionRandomSize());
kzo.client.dbg() << "done" << std::endl;
}
}
}
}
}
auto eventJson = json{
{"content", {{"algorithm", megOlmAlgo},
{"room_id", roomId},
{"session_id", sessionId},
{"session_key", sessionKey}}},
{"type", "m.room_key"}
};
auto event = Event(JsonWrap(eventJson));
return {
std::move(m),
[event](auto &&) { return EffectStatus{ /* success = */ true, json{{ "keyEvent", event.originalJson() }} }; }
};
}
ClientResult updateClient(ClientModel m, EncryptMegOlmEventAction a)
{
auto [encryptedEvent, maybeKey] = m.megOlmEncrypt(a.e, a.roomId, a.timeMs, a.random);
return {
std::move(m),
[=](auto && /* ctx */) {
auto retJson = json::object({
{"encrypted", encryptedEvent.originalJson()},
});
if (maybeKey.has_value()) {
retJson["key"] = maybeKey.value();
}
return EffectStatus(/* succ = */ true, retJson);
}
};
}
ClientResult updateClient(ClientModel m, SetDeviceTrustLevelAction a)
{
auto maybeOldInfo = m.deviceLists.get(a.userId, a.deviceId);
if (!maybeOldInfo) {
return {
std::move(m),
[=](auto && /* ctx */) {
auto retJson = json::object({
{"error", "No such device"},
{"errorCode", "MOE_KAZV_MXC_KAZV_NO_SUCH_DEVICE"},
});
return EffectStatus(/* succ = */ false, retJson);
}
};
}
m.deviceLists.deviceLists = updateIn(
std::move(m.deviceLists.deviceLists),
[a](auto device) {
device.trustLevel = a.trustLevel;
return device;
},
a.userId,
a.deviceId
);
return { m, lager::noop };
}
ClientResult updateClient(ClientModel m, SetTrustLevelNeededToSendKeysAction a)
{
m.trustLevelNeededToSendKeys = a.trustLevel;
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, PrepareForSharingRoomKeyAction a)
{
auto messages = m.olmEncryptSplit(a.e, a.devices, a.random);
auto txnId = getTxnId(Event(), m);
m.roomList = RoomListModel::update(
std::move(m.roomList),
UpdateRoomAction{
a.roomId,
AddPendingRoomKeyAction{
PendingRoomKeyEvent{txnId, messages}
}
}
);
return { std::move(m), [txnId](auto &&) {
return EffectStatus(/* succ = */ true, json::object({{"txnId", txnId}}));
} };
}
ClientResult updateClient(ClientModel m, ImportFromKeyBackupFileAction a)
{
auto maybeExportFile = decryptKeyExport(std::move(a.fileContent), std::move(a.password));
if (!maybeExportFile) {
return {std::move(m), Kazv::detail::ReturnEffectStatusT{{
/* succ = */ false,
json{
{"errorCode", maybeExportFile.reason()},
{"error", maybeExportFile.reason()},
},
}}};
}
std::size_t imported = 0;
m.withCrypto([&maybeExportFile, &imported](Crypto &c) {
imported = c.importInboundGroupSessions(std::move(maybeExportFile).value());
});
return {std::move(m), Kazv::detail::ReturnEffectStatusT{{
/* succ = */ true,
json{
{"imported", imported},
},
}}};
};
}
diff --git a/src/client/actions/ephemeral.cpp b/src/client/actions/ephemeral.cpp
index 2b38820..021960a 100644
--- a/src/client/actions/ephemeral.cpp
+++ b/src/client/actions/ephemeral.cpp
@@ -1,101 +1,87 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu <tusooa@kazv.moe>
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <status-utils.hpp>
#include "ephemeral.hpp"
namespace Kazv
{
ClientResult updateClient(ClientModel m, SetTypingAction a)
{
auto job = m.job<SetTypingJob>().make(
m.userId,
a.roomId,
a.typing,
a.timeoutMs)
.withData(json{{"roomId", a.roomId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, SetTypingResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(SetTypingFailed{roomId, r.errorCode(), r.errorMessage()});
-
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(SetTypingSuccessful{roomId});
-
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, PostReceiptAction a)
{
auto job = m.job<PostReceiptJob>().make(
a.roomId,
/* receiptType = */ "m.read"s,
a.eventId,
/* receipt = */ json::object())
.withData(json{{"roomId", a.roomId}});
m.addJob(std::move(job));
auto userId = m.userId;
auto [next, _] = ClientModel::update(std::move(m), UpdateRoomAction{
a.roomId,
UpdateLocalReadMarkerAction{a.eventId, userId},
});
return { std::move(next), lager::noop };
}
ClientResult processResponse(ClientModel m, PostReceiptResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(PostReceiptFailed{roomId, r.errorCode(), r.errorMessage()});
-
return { std::move(m), lager::noop };
}
- m.addTrigger(PostReceiptSuccessful{roomId});
-
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SetReadMarkerAction a)
{
auto job = m.job<SetReadMarkerJob>().make(
a.roomId,
/* mFullyRead = */ a.eventId)
.withData(json{{"roomId", a.roomId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, SetReadMarkerResponse r)
{
- auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(SetReadMarkerFailed{roomId, r.errorCode(), r.errorMessage()});
-
return { std::move(m), lager::noop };
}
-
- m.addTrigger(SetReadMarkerSuccessful{roomId});
-
return { std::move(m), lager::noop };
}
}
diff --git a/src/client/actions/membership.cpp b/src/client/actions/membership.cpp
index 46971d7..aa3fe77 100644
--- a/src/client/actions/membership.cpp
+++ b/src/client/actions/membership.cpp
@@ -1,297 +1,268 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu <tusooa@kazv.moe>
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <csapi/create_room.hpp>
#include <csapi/inviting.hpp>
#include <csapi/joining.hpp>
#include <debug.hpp>
#include "client-model.hpp"
#include "clientutil.hpp"
#include "cursorutil.hpp"
#include "status-utils.hpp"
#include "membership.hpp"
namespace Kazv
{
static std::string visibilityToStr(CreateRoomAction::Visibility v)
{
using V = CreateRoomAction::Visibility;
switch (v) {
case V::Private:
return "private";
case V::Public:
return "public";
default:
// should not happen
return "";
}
}
static std::string presetToStr(CreateRoomAction::Preset p)
{
using P = CreateRoomAction::Preset;
switch (p) {
case P::PrivateChat:
return "private_chat";
case P::PublicChat:
return "public_chat";
case P::TrustedPrivateChat:
return "trusted_private_chat";
default:
// should not happen
return "";
}
}
ClientResult updateClient(ClientModel m, CreateRoomAction a)
{
auto visibility = visibilityToStr(a.visibility);
std::optional<std::string> preset = a.preset
? std::optional<std::string>(presetToStr(a.preset.value()))
: std::nullopt;
using StateEvT = Kazv::CreateRoomJob::StateEvent;
auto initialState = intoImmer(
immer::array<StateEvT>{},
zug::map(
[](Event e) {
return StateEvT{e.type(), e.content(), e.stateKey()};
}),
a.initialState);
auto job = m.job<CreateRoomJob>().make(
visibility,
a.roomAliasName,
a.name,
a.topic,
a.invite,
DEFVAL, // invite3pid, not supported yet
a.roomVersion,
a.creationContent,
initialState,
preset,
a.isDirect,
a.powerLevelContentOverride);
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, CreateRoomResponse r)
{
if (! r.success()) {
kzo.client.dbg() << "Create room failed" << std::endl;
- m.addTrigger(CreateRoomFailed{r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(CreateRoomSuccessful{r.roomId()});
return { std::move(m), [=](auto &&) {
return EffectStatus(/* succ = */ true, json{
{"roomId", r.roomId()},
});
}};
}
ClientResult updateClient(ClientModel m, InviteToRoomAction a)
{
auto job = m.job<InviteUserJob>()
.make(a.roomId, a.userId)
.withData(json{
{"roomId", a.roomId},
{"userId", a.userId},
});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, InviteUserResponse r)
{
auto roomId = r.dataStr("roomId");
auto userId = r.dataStr("userId");
if (! r.success()) {
// Error
kzo.client.warn() << "Error inviting user " << r.errorCode() << r.errorMessage() << std::endl;
- m.addTrigger(InviteUserFailed{roomId, userId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
kzo.client.info() << "Inviting user successful" << std::endl;
- m.addTrigger(InviteUserSuccessful{roomId, userId});
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, JoinRoomAction a)
{
/// HACK: wait for real boost::url
auto encodedId = a.roomIdOrAlias;
if (!encodedId.empty() && encodedId[0] == '#') {
encodedId.replace(0, 1, "%23");
}
auto job = m.job<JoinRoomJob>()
.make(encodedId, a.serverName)
.withData(json{{"roomIdOrAlias", a.roomIdOrAlias}});
m.addJob(std::move(job));
return { m, lager::noop };
}
ClientResult processResponse(ClientModel m, JoinRoomResponse r)
{
auto roomIdOrAlias = r.dataStr("roomIdOrAlias");
if (! r.success()) {
- m.addTrigger(JoinRoomFailed{
- roomIdOrAlias,
- r.errorCode(),
- r.errorMessage()
- });
kzo.client.warn() << "Error joining room" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
kzo.client.info() << "Successfully joined room" << std::endl;
- m.addTrigger(JoinRoomSuccessful{roomIdOrAlias});
return { std::move(m), lager::noop};
}
ClientResult updateClient(ClientModel m, JoinRoomByIdAction a)
{
auto job = m.job<JoinRoomByIdJob>()
.make(a.roomId)
.withData(json{{"roomIdOrAlias", a.roomId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, JoinRoomByIdResponse r)
{
auto roomIdOrAlias = r.dataStr("roomIdOrAlias");
if (! r.success()) {
- m.addTrigger(JoinRoomFailed{
- roomIdOrAlias,
- r.errorCode(),
- r.errorMessage()
- });
-
kzo.client.warn() << "Error joining room" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
kzo.client.info() << "Successfully joined room" << std::endl;
- m.addTrigger(JoinRoomSuccessful{roomIdOrAlias});
return { std::move(m), lager::noop};
}
ClientResult updateClient(ClientModel m, LeaveRoomAction a)
{
auto job = m.job<LeaveRoomJob>()
.make(a.roomId)
.withData(json{{"roomId", a.roomId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, LeaveRoomResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(LeaveRoomFailed{
- roomId,
- r.errorCode(),
- r.errorMessage()
- });
kzo.client.warn() << "Error leaving room" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
kzo.client.info() << "Successfully left room" << std::endl;
- m.addTrigger(LeaveRoomSuccessful{roomId});
return { std::move(m), lager::noop};
}
ClientResult updateClient(ClientModel m, ForgetRoomAction a)
{
auto job = m.job<ForgetRoomJob>()
.make(a.roomId)
.withData(json{{"roomId", a.roomId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, ForgetRoomResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(ForgetRoomFailed{
- roomId,
- r.errorCode(),
- r.errorMessage()
- });
kzo.client.warn() << "Error forgetting room" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
kzo.client.info() << "Successfully forgot room" << std::endl;
- m.addTrigger(ForgetRoomSuccessful{roomId});
return { std::move(m), lager::noop};
}
ClientResult updateClient(ClientModel m, KickAction a)
{
m.addJob(m.job<KickJob>().make(a.roomId, a.userId, a.reason));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, KickResponse r)
{
if (!r.success()) {
kzo.client.warn() << "Error kicking" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, BanAction a)
{
m.addJob(m.job<BanJob>().make(a.roomId, a.userId, a.reason));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, BanResponse r)
{
if (!r.success()) {
kzo.client.warn() << "Error banning" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, UnbanAction a)
{
m.addJob(m.job<UnbanJob>().make(a.roomId, a.userId));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, UnbanResponse r)
{
if (!r.success()) {
kzo.client.warn() << "Error unbanning" << r.errorCode() << r.errorMessage() << std::endl;
return { std::move(m), failWithResponse(r) };
}
return { std::move(m), lager::noop };
}
};
diff --git a/src/client/actions/paginate.cpp b/src/client/actions/paginate.cpp
index 54591ff..df598a5 100644
--- a/src/client/actions/paginate.cpp
+++ b/src/client/actions/paginate.cpp
@@ -1,107 +1,104 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
// eager requires tuplify but did not include it.
// This is a bug.
#include <zug/tuplify.hpp>
#include <zug/transducer/eager.hpp>
#include <zug/into.hpp>
#include <lager/util.hpp>
#include <types.hpp>
#include <debug.hpp>
#include <cursorutil.hpp>
#include <status-utils.hpp>
#include "paginate.hpp"
namespace Kazv
{
ClientResult updateClient(ClientModel m, PaginateTimelineAction a)
{
auto roomId = a.roomId;
auto room = m.roomList.at(roomId);
if (! room.timelineGaps.find(a.fromEventId)) {
kzo.client.dbg() << "There is no such gap at " << a.fromEventId << std::endl;
return { m, simpleFail };
}
auto paginateBackToken = room.timelineGaps[a.fromEventId];
auto job = m.job<GetRoomEventsJob>().make(
roomId,
"b"s, // dir
paginateBackToken, // from
std::nullopt, // to
a.limit).withData(json{{"roomId", roomId},
{"gapEventId", a.fromEventId}});
m.addJob(std::move(job));
return { m, lager::noop };
}
ClientResult processResponse(ClientModel m, GetRoomEventsResponse r)
{
auto roomId = r.dataStr("roomId");
auto gapEventId = r.dataStr("gapEventId");
if (! r.success()) {
kzo.client.dbg() << "Get room events failed" << std::endl;
- m.addTrigger(PaginateFailed{roomId});
return { m, failWithResponse(r) };
}
auto paginateBackToken = r.end();
auto chunk = r.chunk();
try {
// kzo.client.dbg() << "We got " << chunk.size() << " events here" << std::endl;
auto room = m.roomList.at(roomId);
// The timeline from paginate backwards is returned
// in reversed order, so restore the order.
auto events = intoImmer(EventList{}, zug::reversed, chunk);
AddToTimelineAction action
{events, paginateBackToken, std::nullopt, gapEventId};
m.roomList = RoomListModel::update(
std::move(m.roomList),
UpdateRoomAction{roomId, action});
m.roomList = RoomListModel::update(
std::move(m.roomList),
UpdateRoomAction{roomId, MaybeAddStateEventsAction{
intoImmer(EventList{},
zug::filter(&Event::isState),
events)
}});
// r.state() contains state events prior to the timeline
// we add them later than the timeline, so the states in the
// timeline takes priority
m.roomList = RoomListModel::update(
std::move(m.roomList),
UpdateRoomAction{roomId, MaybeAddStateEventsAction{
r.state()
}});
- m.addTrigger(PaginateSuccessful{roomId});
-
return { std::move(m), lager::noop };
} catch (const std::out_of_range &e) {
// Ignore it, the client model is modified
// such that it knows nothing about this room.
// May happen in debugger.
return { std::move(m), simpleFail };
}
}
}
diff --git a/src/client/actions/send.cpp b/src/client/actions/send.cpp
index 9f3a692..2162101 100644
--- a/src/client/actions/send.cpp
+++ b/src/client/actions/send.cpp
@@ -1,283 +1,281 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu <tusooa@kazv.moe>
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <debug.hpp>
#include <types.hpp>
#include <immer-utils.hpp>
#include "send.hpp"
#include "status-utils.hpp"
namespace Kazv
{
ClientResult updateClient(ClientModel m, SendMessageAction a)
{
auto event = std::move(a.event);
auto roomId = a.roomId;
auto origJson = event.originalJson().get();
if (!origJson.contains("type") || !origJson.contains("content")) {
- m.addTrigger(InvalidMessageFormat{});
- return { std::move(m), lager::noop };
+ return { std::move(m), failEffect(
+ "MOE.KAZV.MXC.INVALID_MESSAGE_FORMAT",
+ "Invalid message format"
+ ) };
}
if (m.roomList.rooms[a.roomId].encrypted && !event.encrypted()) {
if (!event.isState()) {
return { std::move(m), [](auto &&) {
return EffectStatus{
/* succ = */ false,
json{
{"errorCode", "MOE_KAZV_MXC_SENDING_UNENCRYPTED_EVENT_TO_ENCRYPTED_ROOM"},
{"error", "Cannot send unencrypted event to encrypted room"},
}
};
}};
}
}
// We do not use event.type() etc. because we want
// encrypted events stay encrypted.
auto type = origJson["type"];
auto content = origJson["content"];
kzo.client.dbg() << "Sending message of type " << type
<< " with content " << content.dump()
<< " to " << a.roomId
<< " as #" << m.nextTxnId << std::endl;
// We combine the hash of json, the timestamp,
// and a numeric count in the client to avoid collision.
auto txnId = a.txnId.has_value() ? a.txnId.value() : getTxnId(event, m);
m.roomList = RoomListModel::update(std::move(m.roomList),
UpdateRoomAction{
roomId,
AddLocalEchoAction{{txnId, event}},
}
);
auto job = m.job<SendMessageJob>()
.make(a.roomId, type, txnId, content)
.withData(json{
{"roomId", a.roomId},
{"txnId", txnId},
});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, SendMessageResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
auto txnId = r.dataStr("txnId");
kzo.client.dbg() << "Send message failed" << std::endl;
m.roomList.rooms = std::move(m.roomList.rooms).update(roomId, [txnId](auto room) {
auto maybeLocalEcho = room.getLocalEchoByTxnId(txnId);
if (!maybeLocalEcho.has_value()) {
kzo.client.warn() << "We do not have local echo with txnId " << txnId << " . Have we time-travelled?" << std::endl;
return room;
}
auto localEcho = std::move(maybeLocalEcho).value();
localEcho.status = LocalEchoDesc::Failed;
return RoomModel::update(std::move(room), AddLocalEchoAction{localEcho});
});
- m.addTrigger(SendMessageFailed{roomId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(SendMessageSuccessful{roomId, r.eventId()});
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SendToDeviceMessageAction a)
{
auto origJson = a.event.originalJson().get();
if (!origJson.contains("type") || !origJson.contains("content")) {
return { std::move(m), simpleFail };
}
// We do not use event.type() etc. because we want
// encrypted events stay encrypted.
auto type = origJson["type"];
auto content = origJson["content"];
if (type == "m.room_key" && !a.event.encrypted()) {
kzo.client.err() << "Trying to send room key event unencrypted! Rejecting." << std::endl;
return { std::move(m), failEffect(
"MOE_KAZV_MXC_SENDING_ROOM_KEY_EVENT_UNENCRYPTED",
"Cannot send room key event unencrypted"
)};
}
auto txnId = a.txnId.has_value() ? a.txnId.value() : getTxnId(a.event, m);
kzo.client.info() << "sending to-device message with txnId " << txnId;
auto messages =
immer::map<std::string, immer::map<std::string, JsonWrap>>{};
for (auto [userId, devices] : a.devicesToSend) {
auto deviceIdToContentMap = immer::map<std::string, JsonWrap>{};
for (auto deviceId : devices) {
deviceIdToContentMap = std::move(deviceIdToContentMap).set(deviceId, content);
}
messages = std::move(messages).set(userId, deviceIdToContentMap);
}
auto job = m.job<SendToDeviceJob>()
.make(type, txnId, messages)
.withData(json{{"devicesToSend", a.devicesToSend},
{"txnId", txnId}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SendMultipleToDeviceMessagesAction a)
{
std::optional<std::string> type;
std::optional<std::string> txnId;
immer::map<std::string, immer::map<std::string, JsonWrap>> userToDeviceToContentMap;
for (auto [userId, deviceToEventMap] : a.userToDeviceToEventMap) {
userToDeviceToContentMap = std::move(userToDeviceToContentMap).set(userId, immer::map<std::string, JsonWrap>());
for (auto [deviceId, event] : deviceToEventMap) {
auto origJson = event.originalJson().get();
if (!origJson.contains("type") || !origJson.contains("content")) {
return { std::move(m), failEffect(
"MOE_KAZV_MXC_INVALID_EVENT_FORMAT",
"Invalid event format"
)};
}
if (!type.has_value()) {
type = origJson["type"];
} else {
if (origJson["type"] != type.value()) {
return { std::move(m), failEffect(
"MOE_KAZV_MXC_TYPE_NOT_SAME",
"The to-device messages' type must be the same"
)};
}
}
if (origJson["type"] == "m.room_key" && !event.encrypted()) {
kzo.client.err() << "Trying to send room key event unencrypted! Rejecting." << std::endl;
return { std::move(m), failEffect(
"MOE_KAZV_MXC_SENDING_ROOM_KEY_EVENT_UNENCRYPTED",
"Cannot send room key event unencrypted"
)};
}
if (!txnId.has_value()) {
txnId = getTxnId(event, m);
}
auto content = origJson["content"];
userToDeviceToContentMap = setIn(
std::move(userToDeviceToContentMap),
std::move(content),
userId, deviceId
);
}
}
if (!type) {
// nothing to send
return { std::move(m), lager::noop };
}
auto job = m.job<SendToDeviceJob>()
.make(type.value(), txnId.value(), userToDeviceToContentMap)
// XXX not extracting devicesToSend from the map
// those triggers should be deprecated in favour of
// then()-continuation anyway
.withData(json{{"devicesToSend", json::object()},
{"txnId", txnId.value()}});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, SendToDeviceResponse r)
{
auto devicesToSend = r.dataJson("devicesToSend");
auto txnId = r.dataStr("txnId");
if (! r.success()) {
- m.addTrigger(SendToDeviceMessageFailed{devicesToSend, txnId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(SendToDeviceMessageSuccessful{devicesToSend, txnId});
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SaveLocalEchoAction a)
{
auto txnId = a.txnId.has_value() ? a.txnId.value() : getTxnId(a.event, m);
m.roomList = RoomListModel::update(m.roomList, UpdateRoomAction{
a.roomId,
AddLocalEchoAction{{txnId, a.event}},
});
return { std::move(m), [txnId](auto) {
return EffectStatus(true, json::object({{"txnId", txnId}}));
}};
}
ClientResult updateClient(ClientModel m, UpdateLocalEchoStatusAction a)
{
auto txnId = a.txnId;
auto maybeLocalEcho = m.roomList.rooms[a.roomId].getLocalEchoByTxnId(txnId);
if (!maybeLocalEcho.has_value()) {
return { std::move(m), [](const auto &) {
return EffectStatus(/* succ = */ false, json{
{"errorCode", "MOE_KAZV_MXC_LOCAL_ECHO_NOT_FOUND"},
{"error", "Local echo not found"},
});
} };
}
auto localEcho = maybeLocalEcho.value();
m.roomList = RoomListModel::update(m.roomList, UpdateRoomAction{
a.roomId,
AddLocalEchoAction{{txnId, localEcho.event, a.status}},
});
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, RedactEventAction a)
{
auto txnId = getTxnId(Event{json{
{"room_id", a.roomId},
{"event_id", a.eventId},
}}, m);
auto job = m.job<RedactEventJob>()
.make(a.roomId, a.eventId, txnId, a.reason);
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, RedactEventResponse r)
{
if (!r.success()) {
return { std::move(m), failWithResponse(r) };
}
return { std::move(m), lager::noop };
}
}
diff --git a/src/client/actions/states.cpp b/src/client/actions/states.cpp
index 03fd014..8590b11 100644
--- a/src/client/actions/states.cpp
+++ b/src/client/actions/states.cpp
@@ -1,165 +1,156 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <debug.hpp>
#include "clientutil.hpp"
#include "cursorutil.hpp"
#include <status-utils.hpp>
#include "states.hpp"
namespace Kazv
{
BaseJob clientPerform(ClientModel m, GetRoomStatesAction a)
{
auto roomId = a.roomId;
return m.job<GetRoomStateJob>()
.make(roomId)
.withData(json{{"roomId", roomId}});
}
ClientResult updateClient(ClientModel m, GetRoomStatesAction a)
{
m.addJob(clientPerform(m, a));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, GetRoomStateResponse r)
{
auto roomId = r.dataStr("roomId");
if (! r.success()) {
- m.addTrigger(GetRoomStatesFailed{roomId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(GetRoomStatesSuccessful{roomId});
-
auto events = r.data();
auto action = UpdateRoomAction{
roomId,
AddStateEventsAction{std::move(events)}
};
m.roomList = RoomListModel::update(std::move(m.roomList), action);
if (m.crypto && m.roomList[roomId].encrypted) {
m.deviceLists.track(m.roomList[roomId].joinedMemberIds());
}
m.roomList = RoomListModel::update(std::move(m.roomList),
UpdateRoomAction{roomId, MarkMembersFullyLoadedAction{}});
- //m.addTrigger(ShouldQueryKeys{false});
-
return {std::move(m), lager::noop};
}
ClientResult updateClient(ClientModel m, GetStateEventAction a)
{
auto job = m.job<GetRoomStateWithKeyJob>()
.make(a.roomId, a.type, a.stateKey)
.withData(json{
{"roomId", a.roomId},
{"type", a.type},
{"stateKey", a.stateKey}
});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, GetRoomStateWithKeyResponse r)
{
auto roomId = r.dataStr("roomId");
auto type = r.dataStr("type");
auto stateKey = r.dataStr("stateKey");
if (! r.success()) {
- m.addTrigger(GetStateEventFailed{roomId, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
auto content = r.jsonBody();
- m.addTrigger(GetStateEventSuccessful{roomId, content});
-
auto k = KeyOfState{type, stateKey};
auto eventJson = m.roomList[roomId].stateEvents[k]
.originalJson().get();
eventJson["content"] = content;
eventJson["type"] = type;
eventJson["state_key"] = stateKey;
auto a = AddStateEventsAction{EventList{Event{eventJson}}};
auto l = RoomListModel::update(
std::move(m.roomList),
UpdateRoomAction{roomId, a});
m.roomList = std::move(l);
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SendStateEventAction a)
{
auto event = a.event;
if (event.type() == ""s) {
- m.addTrigger(InvalidMessageFormat{});
- return { std::move(m), lager::noop };
+ return { std::move(m), failEffect(
+ "MOE.KAZV.MXC.INVALID_MESSAGE_FORMAT",
+ "Invalid message format"
+ ) };
}
auto type = event.type();
auto content = event.content();
auto stateKey = event.stateKey();
kzo.client.dbg() << "Sending state event of type " << type
<< " with content " << content.get().dump()
<< " to " << a.roomId
<< " with state key #" << stateKey << std::endl;
auto job = m.job<SetRoomStateWithKeyJob>().make(
a.roomId,
type,
stateKey,
content)
.withData(json{
{"roomId", a.roomId},
{"eventType", type},
{"stateKey", stateKey},
});
m.addJob(std::move(job));
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, SetRoomStateWithKeyResponse r)
{
auto roomId = r.dataStr("roomId");
auto eventType = r.dataStr("eventType");
auto stateKey = r.dataStr("stateKey");
if (! r.success()) {
kzo.client.dbg() << "Send state event failed" << std::endl;
- m.addTrigger(SendStateEventFailed{roomId, eventType, stateKey, r.errorCode(), r.errorMessage()});
return { std::move(m), failWithResponse(r) };
}
- m.addTrigger(SendStateEventSuccessful{roomId, r.eventId(), eventType, stateKey});
-
return { std::move(m), lager::noop };
}
}
diff --git a/src/client/actions/sync.cpp b/src/client/actions/sync.cpp
index e49e3dd..404c1a7 100644
--- a/src/client/actions/sync.cpp
+++ b/src/client/actions/sync.cpp
@@ -1,384 +1,379 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2021-2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <lager/util.hpp>
#include <zug/transducer/map.hpp>
#include <zug/transducer/cat.hpp>
#include <zug/transducer/filter.hpp>
#include <zug/sequence.hpp>
#include <jobinterface.hpp>
#include <debug.hpp>
#include "cursorutil.hpp"
#include "sync.hpp"
#include "encryption.hpp"
#include "status-utils.hpp"
namespace Kazv
{
// Atomicity guaranteed: if the sync action is created
// before an action that reasonably changes Client
// (e.g. roll back to an earlier state, obtain other
// events), but executed
// after that action, the sync will still give continuous
// data about the events. (Sync will not "skip" events)
// This is because this function takes the sync token
// from the ClientModel model it is passed.
ClientResult updateClient(ClientModel m, SyncAction)
{
kzo.client.dbg() << "Start syncing with token " <<
(m.syncToken ? m.syncToken.value() : "<null>") << std::endl;
bool isInitialSync = ! m.syncToken;
m.syncing = true;
std::string filter = m.syncToken ? m.incrementalSyncFilterId : m.initialSyncFilterId;
m.addJob(m.job<SyncJob>()
.make(filter,
m.syncToken,
std::nullopt, // fullState
std::nullopt, // setPresence
// Let initial sync return immediately
isInitialSync ? 0 : m.syncTimeoutMs
)
.withData(json{{"is", isInitialSync ? "initial" : "incremental"}}));
return { m, lager::noop };
}
- static KazvEventList loadRoomsFromSyncInPlace(ClientModel &m, SyncJob::Rooms rooms)
+ static KazvTriggerList loadRoomsFromSyncInPlace(ClientModel &m, SyncJob::Rooms rooms)
{
auto l = std::move(m.roomList);
- auto eventsToEmit = KazvEventList{}.transient();
+ auto eventsToEmit = KazvTriggerList{}.transient();
auto pushRules = PushRulesDesc(m.accountData["m.push_rules"]);
auto updateRoomImpl =
[&l](auto id, auto a) {
l = RoomListModel::update(
std::move(l),
UpdateRoomAction{std::move(id), std::move(a)});
};
auto updateSingleRoom =
[&, updateRoomImpl](const auto &id, const auto &room, auto membership) {
if (!l.has(id) || l[id].membership != membership) {
eventsToEmit.push_back(RoomMembershipChanged{membership, id});
}
updateRoomImpl(id, ChangeMembershipAction{membership});
auto timelineEvents =
intoImmer(
EventList{},
zug::map([=](Event e) {
return Event::fromSync(e, id);
}),
room.timeline.events);
eventsToEmit.append(
intoImmer(
- KazvEventList{},
- zug::map([=](Event e) -> KazvEvent {
+ KazvTriggerList{},
+ zug::map([=](Event e) -> KazvTrigger {
return ReceivingRoomTimelineEvent{std::move(e), id};
}),
timelineEvents).transient());
updateRoomImpl(id, AddToTimelineAction{timelineEvents,
room.timeline.prevBatch,
room.timeline.limited,
std::nullopt // we do not have a gapEventId
});
if (room.state) {
eventsToEmit.append(
intoImmer(
- KazvEventList{},
- zug::map([=](Event e) -> KazvEvent {
+ KazvTriggerList{},
+ zug::map([=](Event e) -> KazvTrigger {
return ReceivingRoomStateEvent{std::move(e), id};
}),
room.state.value().events).transient());
updateRoomImpl(id, AddStateEventsAction{room.state.value().events});
}
// Process state events in timeline, which should have arrived later
// than those in room.state .
updateRoomImpl(id, AddStateEventsAction{
intoImmer(EventList{},
zug::filter([=](Event e) {
return e.isState();
}),
timelineEvents)});
if (room.accountData) {
eventsToEmit.append(
intoImmer(
- KazvEventList{},
- zug::map([=](Event e) -> KazvEvent {
+ KazvTriggerList{},
+ zug::map([=](Event e) -> KazvTrigger {
return ReceivingRoomAccountDataEvent{std::move(e), id};
}),
room.accountData.value().events).transient());
updateRoomImpl(id, AddAccountDataAction{room.accountData.value().events});
}
};
auto updateRoomSummary =
[=](const auto &id, const auto &room) {
if (!room.summary.has_value()) {
return;
}
if (!room.summary->mHeroes.empty()) {
auto newHeroes = room.summary->mHeroes;
updateRoomImpl(id, SetHeroIdsAction{immer::flex_vector<std::string>(newHeroes.begin(), newHeroes.end())});
}
if (room.summary->mJoinedMemberCount.has_value()) {
updateRoomImpl(id, UpdateJoinedMemberCountAction{static_cast<std::size_t>(room.summary->mJoinedMemberCount.value())});
}
if (room.summary->mInvitedMemberCount.has_value()) {
updateRoomImpl(id, UpdateInvitedMemberCountAction{static_cast<std::size_t>(room.summary->mInvitedMemberCount.value())});
}
};
auto updateRoomNotifications = [=, &l](const auto &id, const auto &room) {
auto oldRoom = m.roomList.rooms[id];
const auto &newRoom = l.rooms[id];
updateRoomImpl(id, AddLocalNotificationsAction{
room.timeline.events,
pushRules,
m.userId,
});
if (oldRoom.readReceipts[m.userId].eventId
!= newRoom.readReceipts[m.userId].eventId) {
updateRoomImpl(id, RemoveReadLocalNotificationsAction{m.userId});
}
};
auto updateJoinedRoom =
[=](const auto &id, const auto &room) {
updateSingleRoom(id, room, RoomMembership::Join);
if (room.ephemeral) {
updateRoomImpl(id, AddEphemeralAction{room.ephemeral.value().events});
}
updateRoomNotifications(id, room);
updateRoomSummary(id, room);
};
auto updateInvitedRoom =
[=](const auto &id, const auto &room) {
updateRoomImpl(id, ChangeMembershipAction{RoomMembership::Invite});
if (room.inviteState) {
updateRoomImpl(id, ChangeInviteStateAction{room.inviteState.value().events});
}
};
auto updateLeftRoom =
[=](const auto &id, const auto &room) {
updateSingleRoom(id, room, RoomMembership::Leave);
};
for (const auto &[id, room]: rooms.join) {
updateJoinedRoom(id, room);
}
// TODO update info for invited rooms
for (const auto &[id, room]: rooms.invite) {
updateInvitedRoom(id, room);
}
for (const auto &[id, room]: rooms.leave) {
updateLeftRoom(id, room);
}
m.roomList = std::move(l);
return eventsToEmit.persistent();
}
- static KazvEventList loadPresenceFromSyncInPlace(ClientModel &m, EventList presence)
+ static KazvTriggerList loadPresenceFromSyncInPlace(ClientModel &m, EventList presence)
{
auto eventsToEmit = intoImmer(
- KazvEventList{},
+ KazvTriggerList{},
zug::map([](Event e) { return ReceivingPresenceEvent{e}; }),
presence);
m.presence = merge(std::move(m.presence), presence, keyOfPresence);
return eventsToEmit;
}
- static KazvEventList loadAccountDataFromSyncInPlace(ClientModel &m, EventList accountData)
+ static KazvTriggerList loadAccountDataFromSyncInPlace(ClientModel &m, EventList accountData)
{
auto eventsToEmit = intoImmer(
- KazvEventList{},
+ KazvTriggerList{},
zug::map([](Event e) { return ReceivingPresenceEvent{e}; }),
accountData);
m.accountData = merge(std::move(m.accountData), accountData, keyOfAccountData);
return eventsToEmit;
}
- static KazvEventList loadToDeviceFromSyncInPlace(ClientModel &m, JsonWrap toDevice)
+ static KazvTriggerList loadToDeviceFromSyncInPlace(ClientModel &m, JsonWrap toDevice)
{
if (toDevice.get().contains("events")) {
auto events = toDevice.get()["events"];
auto msgs = intoImmer(
EventList{},
zug::map([](json j) {
// Prevent malicious server from injecting a room id
j.erase("room_id");
return Event(j);
}),
events);
m.toDevice = std::move(m.toDevice) + msgs;
return intoImmer(
- KazvEventList{},
+ KazvTriggerList{},
zug::map([](Event e) { return ReceivingToDeviceMessage{e}; }),
msgs);
}
return {};
}
ClientResult processResponse(ClientModel m, SyncResponse r)
{
if (! r.success()) {
- m.addTrigger(SyncFailed{});
kzo.client.dbg() << "Sync failed" << std::endl;
kzo.client.dbg() << r.statusCode << std::endl;
if (isBodyJson(r.body)) {
auto j = r.jsonBody();
kzo.client.dbg() << "Json says: " << j.get().dump() << std::endl;
} else {
kzo.client.dbg() << "Response body: "
<< std::get<BaseJob::BytesBody>(r.body) << std::endl;
}
return { std::move(m), failWithResponse(r) };
}
kzo.client.dbg() << "Sync successful" << std::endl;
auto rooms = r.rooms();
auto accountData = r.accountData();
auto presence = r.presence();
// load the info that has been sync'd
m.syncToken = r.nextBatch();
// Load account data first because it contains push rules
// which can affect the processing of rooms
if (accountData) {
m.addTriggers(loadAccountDataFromSyncInPlace(m, std::move(accountData.value().events)));
}
if (rooms) {
m.addTriggers(loadRoomsFromSyncInPlace(m, std::move(rooms.value())));
}
if (presence) {
m.addTriggers(loadPresenceFromSyncInPlace(m, std::move(presence.value().events)));
}
m.addTriggers(loadToDeviceFromSyncInPlace(m, r.toDevice()));
auto is = r.dataStr("is");
auto isInitialSync = is == "initial";
if (m.crypto) {
kzo.client.dbg() << "E2EE is on. Processing device lists and one-time key counts." << std::endl;
// process deviceLists
if (isInitialSync) {
auto encryptedUsers =
zug::sequence(
zug::map([](auto n) { return n.second; })
| zug::filter([](auto room) { return room.encrypted; })
| zug::map([](auto room) { return room.joinedMemberIds(); })
| zug::cat,
// no need to use distinct here as the map will overwrite
m.roomList.rooms);
m.deviceLists.track(std::move(encryptedUsers));
} else {
auto l = r.deviceLists().get();
if (l.contains("changed")) {
const auto &changed = l.at("changed");
m.deviceLists.track(changed);
}
if (l.contains("left")) {
const auto &left = l.at("left");
m.deviceLists.untrack(left);
}
}
// deviceOneTimeKeysCount
m.withCrypto([&](auto &c) { c.setUploadedOneTimeKeysCount(r.deviceOneTimeKeysCount()); });
auto model = tryDecryptEvents(std::move(m));
m = std::move(model);
}
- m.addTrigger(SyncSuccessful{r.nextBatch()});
-
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, SetShouldSyncAction a)
{
m.shouldSync = a.shouldSync;
return { std::move(m), lager::noop };
}
ClientResult updateClient(ClientModel m, PostInitialFiltersAction)
{
if (m.syncing) {
return { std::move(m), lager::noop };
}
Filter initialSyncFilter;
initialSyncFilter.room.timeline.limit = 1;
initialSyncFilter.room.state.lazyLoadMembers = true;
auto firstJob = m.job<DefineFilterJob>()
.make(m.userId, initialSyncFilter)
.withData(json{{"is", "initialSyncFilter"}})
.withQueue("post-filter", CancelFutureIfFailed);
m.addJob(firstJob);
Filter incrementalSyncFilter;
incrementalSyncFilter.room.timeline.limit = 20;
incrementalSyncFilter.room.state.lazyLoadMembers = true;
m.addJob(m.job<DefineFilterJob>()
.make(m.userId, incrementalSyncFilter)
.withData(json{{"is", "incrementalSyncFilter"}})
.withQueue("post-filter", CancelFutureIfFailed));
m.syncing = true;
return { std::move(m), lager::noop };
}
ClientResult processResponse(ClientModel m, DefineFilterResponse r)
{
auto is = r.dataStr("is");
if (! r.success()) {
m.syncing = false;
kzo.client.dbg() << "posting filter failed: " << r.errorCode() << r.errorMessage() << std::endl;
- m.addTrigger(PostInitialFiltersFailed{r.errorCode(), r.errorMessage()});
return { std::move(m), lager::noop };
}
kzo.client.dbg() << "filter " << is << " is posted" << std::endl;
if (is == "incrementalSyncFilter") {
m.incrementalSyncFilterId = r.filterId();
- m.addTrigger(PostInitialFiltersSuccessful{});
} else {
m.initialSyncFilterId = r.filterId();
}
return { std::move(m), lager::noop };
}
}
diff --git a/src/client/client-model.hpp b/src/client/client-model.hpp
index 097c58e..d0d49a2 100644
--- a/src/client/client-model.hpp
+++ b/src/client/client-model.hpp
@@ -1,675 +1,675 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020-2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <tuple>
#include <variant>
#include <string>
#include <optional>
#include <lager/context.hpp>
#include <boost/hana.hpp>
#include <serialization/std-optional.hpp>
#include <csapi/sync.hpp>
#include <file-desc.hpp>
#include <crypto.hpp>
#include <serialization/immer-flex-vector.hpp>
#include <serialization/immer-box.hpp>
#include <serialization/immer-map.hpp>
#include <serialization/immer-array.hpp>
#include "clientfwd.hpp"
#include "device-list-tracker.hpp"
#include "room/room-model.hpp"
namespace Kazv
{
inline const std::string DEFTXNID{"0"};
enum RoomVisibility
{
Private,
Public,
};
enum CreateRoomPreset
{
PrivateChat,
PublicChat,
TrustedPrivateChat,
};
enum ThumbnailResizingMethod
{
Crop,
Scale,
};
struct ClientModel
{
std::string serverUrl;
std::string userId;
std::string token;
std::string deviceId;
bool loggedIn{false};
bool syncing{false};
bool shouldSync{true};
int firstRetryMs{1000};
int retryTimeFactor{2};
int maxRetryMs{30 * 1000};
int syncTimeoutMs{20000};
std::string initialSyncFilterId;
std::string incrementalSyncFilterId;
std::optional<std::string> syncToken;
RoomListModel roomList;
immer::map<std::string /* sender */, Event> presence;
immer::map<std::string /* type */, Event> accountData;
std::string nextTxnId{DEFTXNID};
immer::flex_vector<BaseJob> nextJobs;
- immer::flex_vector<KazvEvent> nextTriggers;
+ immer::flex_vector<KazvTrigger> nextTriggers;
EventList toDevice;
std::optional<immer::box<Crypto>> crypto;
bool identityKeysUploaded{false};
DeviceListTracker deviceLists;
DeviceTrustLevel trustLevelNeededToSendKeys{DeviceTrustLevel::Unseen};
immer::array<std::string /* version */> versions;
immer::flex_vector<std::string /* deviceId */> devicesToSendKeys(std::string userId) const;
/// rotate sessions for a room if there is a user in the room with
/// devicesToSendKeys changes
void maybeRotateSessions(ClientModel oldClient);
std::pair<Event, std::optional<std::string> /* sessionKey */>
megOlmEncrypt(Event e, std::string roomId, Timestamp timeMs, RandomData random);
/// precondition: the one-time keys for those devices must already be claimed
/// @return A map from user id to device id to encrypted event for that device
immer::map<std::string, immer::map<std::string, Event>> olmEncryptSplit(Event e, immer::map<std::string, immer::flex_vector<std::string>> userIdToDeviceIdMap, RandomData random);
/// @return number of one-time keys we need to generate
std::size_t numOneTimeKeysNeeded() const;
/// @return the mapping from room id to user id of direct rooms
auto directRoomMap() const -> immer::map<std::string, std::string>;
auto roomIdsUnderTag(std::string tagId) const -> immer::map<std::string, double>;
auto roomIdsByTagId() const -> immer::map<std::string, immer::map<std::string, double>>;
/// Get the const reference of crypto of this client.
///
/// `crypto.has_value()` must be true.
const Crypto &constCrypto() const;
/// Do func with crypto, returning its return value.
///
/// `crypto.has_value()` must be true.
template<class Func>
auto withCrypto(Func &&func) -> std::decay_t<std::invoke_result_t<Func &&, Crypto &>>
{
using ResT = std::decay_t<std::invoke_result_t<Func &&, Crypto &>>;
if constexpr (std::is_same_v<ResT, void>) {
crypto = std::move(crypto).value()
.update([f=std::forward<Func>(func)](Crypto c) mutable {
std::forward<Func>(f)(c);
return c;
});
} else {
std::optional<ResT> res;
crypto = std::move(crypto).value()
.update([f=std::forward<Func>(func), &res](Crypto c) mutable {
res = std::forward<Func>(f)(c);
return c;
});
return std::move(res).value();
}
}
// helpers
template<class Job>
struct MakeJobT
{
template<class ...Args>
constexpr auto make(Args &&...args) const {
if constexpr (Job::needsAuth()) {
return Job(
serverUrl,
token,
std::forward<Args>(args)...);
} else {
return Job(
serverUrl,
std::forward<Args>(args)...);
}
}
std::string serverUrl;
std::string token;
};
template<class Job>
constexpr auto job() const {
return MakeJobT<Job>{serverUrl, token};
}
inline void addJob(BaseJob j) {
nextJobs = std::move(nextJobs).push_back(std::move(j));
}
inline auto popAllJobs() {
auto jobs = std::move(nextJobs);
nextJobs = DEFVAL;
return jobs;
};
- inline void addTrigger(KazvEvent t) {
+ inline void addTrigger(KazvTrigger t) {
addTriggers({t});
}
- inline void addTriggers(immer::flex_vector<KazvEvent> c) {
+ inline void addTriggers(immer::flex_vector<KazvTrigger> c) {
nextTriggers = std::move(nextTriggers) + c;
}
inline auto popAllTriggers() {
auto triggers = std::move(nextTriggers);
nextTriggers = DEFVAL;
return triggers;
}
void maybeAddSaveEventsTrigger(const ClientModel &old);
using Action = ClientAction;
using Effect = ClientEffect;
using Result = ClientResult;
static Result update(ClientModel m, Action a);
};
// actions:
struct LoginAction {
std::string serverUrl;
std::string username;
std::string password;
std::optional<std::string> deviceName;
};
struct TokenLoginAction
{
std::string serverUrl;
std::string username;
std::string token;
std::string deviceId;
};
struct LogoutAction {};
struct HardLogoutAction {};
struct GetWellknownAction
{
std::string userId;
};
struct GetVersionsAction
{
std::string serverUrl;
};
struct SyncAction {};
struct SetShouldSyncAction
{
bool shouldSync;
};
struct PaginateTimelineAction
{
std::string roomId;
/// Must be where the Gap is
std::string fromEventId;
std::optional<int> limit;
};
struct SendMessageAction
{
std::string roomId;
Event event;
std::optional<std::string> txnId{std::nullopt};
};
struct SendStateEventAction
{
std::string roomId;
Event event;
};
/**
* Saves an local echo.
*
* After dispatching this action, the result should be such that
* `result.dataStr("txnId")` contains the transaction id to be used
* in SendMessageAction.
*/
struct SaveLocalEchoAction
{
/// The room id
std::string roomId;
/// The event to send
Event event;
/// The chosen txnId for this event. If not specified, generate from the current ClientModel.
std::optional<std::string> txnId{std::nullopt};
};
/**
* Updates the status of an local echo.
*
* After dispatching this action, the local echo's status will be
* set to the one described in the action.
*/
struct UpdateLocalEchoStatusAction
{
/// The room id.
std::string roomId;
/// The chosen txnId for this event.
std::string txnId;
/// The updated status of this local echo.
LocalEchoDesc::Status status;
};
struct RedactEventAction
{
std::string roomId;
std::string eventId;
std::optional<std::string> reason;
};
struct CreateRoomAction
{
using Visibility = RoomVisibility;
using Preset = CreateRoomPreset;
Visibility visibility;
std::optional<std::string> roomAliasName;
std::optional<std::string> name;
std::optional<std::string> topic;
immer::array<std::string> invite;
//immer::array<Invite3pid> invite3pid;
std::optional<std::string> roomVersion;
JsonWrap creationContent;
immer::array<Event> initialState;
std::optional<Preset> preset;
std::optional<bool> isDirect;
JsonWrap powerLevelContentOverride;
};
struct GetRoomStatesAction
{
std::string roomId;
};
struct GetStateEventAction
{
std::string roomId;
std::string type;
std::string stateKey;
};
struct InviteToRoomAction
{
std::string roomId;
std::string userId;
};
struct JoinRoomByIdAction
{
std::string roomId;
};
struct JoinRoomAction
{
std::string roomIdOrAlias;
immer::array<std::string> serverName;
};
struct LeaveRoomAction
{
std::string roomId;
};
struct ForgetRoomAction
{
std::string roomId;
};
struct KickAction
{
std::string roomId;
std::string userId;
std::optional<std::string> reason;
};
struct BanAction
{
std::string roomId;
std::string userId;
std::optional<std::string> reason;
};
struct UnbanAction
{
std::string roomId;
std::string userId;
};
struct SetAccountDataPerRoomAction
{
std::string roomId;
Event accountDataEvent;
};
struct SetTypingAction
{
std::string roomId;
bool typing;
std::optional<int> timeoutMs;
};
struct PostReceiptAction
{
std::string roomId;
std::string eventId;
};
struct SetReadMarkerAction
{
std::string roomId;
std::string eventId;
};
struct UploadContentAction
{
FileDesc content;
std::optional<std::string> filename;
std::optional<std::string> contentType;
std::string uploadId; // to be used by library users
};
struct DownloadContentAction
{
std::string mxcUri;
std::optional<FileDesc> downloadTo;
};
struct DownloadThumbnailAction
{
std::string mxcUri;
int width;
int height;
std::optional<ThumbnailResizingMethod> method;
std::optional<bool> allowRemote;
std::optional<FileDesc> downloadTo;
};
struct ResubmitJobAction
{
BaseJob job;
};
struct ProcessResponseAction
{
Response response;
};
struct PostInitialFiltersAction
{
};
struct SetAccountDataAction
{
Event accountDataEvent;
};
struct SendToDeviceMessageAction
{
Event event;
immer::map<std::string, immer::flex_vector<std::string>> devicesToSend;
std::optional<std::string> txnId{std::nullopt};
};
/**
* Send multiple to device messages.
*
* Due to limitations of the spec, the type of the Events must be the same.
*/
struct SendMultipleToDeviceMessagesAction
{
/// A map from user id to device id to the event.
immer::map<std::string, immer::map<std::string, Event>> userToDeviceToEventMap;
/// An optional transaction id. Will be generated if not provided.
std::optional<std::string> txnId{std::nullopt};
};
struct UploadIdentityKeysAction
{
};
/**
* The action to generate one-time keys.
*
* `random.size()` must be at least `randomSize(numToGen)`.
*
* This action will not generate keys exceeding the local limit of olm.
*/
struct GenerateAndUploadOneTimeKeysAction
{
/// @return The size of random needed to generate
/// `numToGen` one-time keys
static std::size_t randomSize(std::size_t numToGen);
/// The number of keys to generate
std::size_t numToGen;
/// The random data used to generate keys
RandomData random;
};
struct QueryKeysAction
{
bool isInitialSync;
};
struct ClaimKeysAction
{
static std::size_t randomSize(immer::map<std::string, immer::flex_vector<std::string>> devicesToSend);
std::string roomId;
std::string sessionId;
std::string sessionKey;
immer::map<std::string, immer::flex_vector<std::string>> devicesToSend;
RandomData random;
};
/**
* The action to encrypt an megolm event for a room.
*
* If the action is successful, the result `r` will
* be such that `r.dataJson("encrypted")` contains the encrypted event *json*.
*
* If the megolm session is rotated, `r.dataStr("key")` will contain the key
* of the megolm session. Otherwise, `r.data().contains("key")` will be false.
*
* The Action may fail due to insufficient random data,
* when the megolm session needs to be rotated.
* In this case, the reducer for the Action will fail,
* and its result `r` will be such that
* `r.dataStr("reason") == "NotEnoughRandom"`.
* The user needs to provide random data of
* at least size `maxRandomSize()`.
*
*/
struct EncryptMegOlmEventAction
{
static std::size_t maxRandomSize();
static std::size_t minRandomSize();
/// The id of the room to encrypt for.
std::string roomId;
/// The event to encrypt.
Event e;
/// The timestamp, to determine whether the session should expire.
Timestamp timeMs;
/// Random data for the operation. Must be of at least size
/// `minRandomSize()`. If this is a retry of the previous operation
/// due to NotEnoughRandom, it must be of at least size `maxRandomSize()`.
RandomData random;
};
struct SetDeviceTrustLevelAction
{
std::string userId;
std::string deviceId;
DeviceTrustLevel trustLevel;
};
struct SetTrustLevelNeededToSendKeysAction
{
DeviceTrustLevel trustLevel;
};
/// Encrypt room key as olm and add it to the room's
/// pending keyshare slots.
/// This is to ensure atomicity and that we do not lose an olm-encrypted event.
struct PrepareForSharingRoomKeyAction
{
using UserIdToDeviceIdMap = immer::map<std::string, immer::flex_vector<std::string>>;
static std::size_t randomSize(UserIdToDeviceIdMap devices);
/// The room to share the key event in.
std::string roomId;
/// Devices to encrypt for.
UserIdToDeviceIdMap devices;
/// The key event to encrypt.
Event e;
/// The random data for the encryption. Must be of at least
/// size `randomSize(devices)`.
RandomData random;
};
/**
* Import keys from key backup file.
*
* On success, the reducer returns data with `imported` property
* being the number of keys imported. On failure, it returns data with
* `errorCode` and `error` properties set to the error in the process.
*/
struct ImportFromKeyBackupFileAction
{
/// The content of the key backup file.
std::string fileContent;
/// The password.
std::string password;
};
struct GetUserProfileAction
{
std::string userId;
};
struct SetAvatarUrlAction
{
std::optional<std::string> avatarUrl;
};
struct SetDisplayNameAction
{
std::optional<std::string> displayName;
};
/// Load events from the storage into the model
struct LoadEventsFromStorageAction
{
/// Map from room id to a list of
/// loaded events that should be put into the timeline. From oldest to latest.
immer::map<std::string, EventList> timelineEvents;
/// Map from room id to a list of
/// related events that should not be put into the timeline. From oldest to latest.
/// There might be events in the storage that is needed to display
/// existing events or room state (e.g. pinned events), but
/// the storage may not know its place in the timeline.
immer::map<std::string, EventList> relatedEvents;
};
/// Remove events from the model, keeping only the latest `maxToKeep` events.
/// For each room, this takes O(maxToKeep * log(maxToKeep)) time.
struct PurgeRoomTimelineAction
{
/// A map from roomId to maxToKeep
immer::map<std::string, std::size_t> roomIdToMaxToKeepMap;
};
template<class Archive>
void serialize(Archive &ar, ClientModel &m, std::uint32_t const version)
{
bool dummySyncing{false};
ar
& m.serverUrl
& m.userId
& m.token
& m.deviceId
& m.loggedIn
& dummySyncing
& m.firstRetryMs
& m.retryTimeFactor
& m.maxRetryMs
& m.syncTimeoutMs
& m.initialSyncFilterId
& m.incrementalSyncFilterId
& m.syncToken
& m.roomList
& m.presence
& m.accountData
& m.nextTxnId
& m.toDevice;
// version <= 1 uses std::optional<Crypto>
// while version >= 2 uses std::optional<immer::box<Crypto>>
if (version >= 2) {
ar & m.crypto;
} else {
if constexpr (typename Archive::is_loading()) {
std::optional<Crypto> crypto;
ar >> crypto;
if (crypto.has_value()) {
m.crypto = immer::box<Crypto>(std::move(crypto).value());
}
}
// otherwise is_saving, which will always use the latest version
// this is unreachable
}
ar
& m.identityKeysUploaded
& m.deviceLists
;
if (version >= 1) { ar & m.trustLevelNeededToSendKeys; }
}
}
BOOST_CLASS_VERSION(Kazv::ClientModel, 2)
diff --git a/src/client/clientfwd.hpp b/src/client/clientfwd.hpp
index 4f9a5d2..064c16a 100644
--- a/src/client/clientfwd.hpp
+++ b/src/client/clientfwd.hpp
@@ -1,158 +1,157 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <tuple>
#include <variant>
#include <lager/context.hpp>
#include <context.hpp>
#include "room/room-model.hpp"
namespace Kazv
{
using namespace Api;
class JobInterface;
class EventInterface;
struct LoginAction;
struct TokenLoginAction;
struct LogoutAction;
struct HardLogoutAction;
struct GetWellknownAction;
struct GetVersionsAction;
struct SyncAction;
struct SetShouldSyncAction;
struct PostInitialFiltersAction;
struct SetAccountDataAction;
struct PaginateTimelineAction;
struct SendMessageAction;
struct SendStateEventAction;
struct SaveLocalEchoAction;
struct UpdateLocalEchoStatusAction;
struct RedactEventAction;
struct CreateRoomAction;
struct GetRoomStatesAction;
struct GetStateEventAction;
struct InviteToRoomAction;
struct JoinRoomByIdAction;
- struct EmitKazvEventsAction;
struct JoinRoomAction;
struct LeaveRoomAction;
struct ForgetRoomAction;
struct KickAction;
struct BanAction;
struct UnbanAction;
struct SetAccountDataPerRoomAction;
struct ProcessResponseAction;
struct SetTypingAction;
struct PostReceiptAction;
struct SetReadMarkerAction;
struct UploadContentAction;
struct DownloadContentAction;
struct DownloadThumbnailAction;
struct SendToDeviceMessageAction;
struct SendMultipleToDeviceMessagesAction;
struct UploadIdentityKeysAction;
struct GenerateAndUploadOneTimeKeysAction;
struct QueryKeysAction;
struct ClaimKeysAction;
struct EncryptMegOlmEventAction;
struct SetDeviceTrustLevelAction;
struct SetTrustLevelNeededToSendKeysAction;
struct PrepareForSharingRoomKeyAction;
struct ImportFromKeyBackupFileAction;
struct GetUserProfileAction;
struct SetAvatarUrlAction;
struct SetDisplayNameAction;
struct ResubmitJobAction;
struct LoadEventsFromStorageAction;
struct PurgeRoomTimelineAction;
struct ClientModel;
using ClientAction = std::variant<
RoomListAction,
LoginAction,
TokenLoginAction,
LogoutAction,
HardLogoutAction,
GetWellknownAction,
GetVersionsAction,
SyncAction,
SetShouldSyncAction,
PostInitialFiltersAction,
SetAccountDataAction,
PaginateTimelineAction,
SendMessageAction,
SendStateEventAction,
SaveLocalEchoAction,
UpdateLocalEchoStatusAction,
RedactEventAction,
CreateRoomAction,
GetRoomStatesAction,
GetStateEventAction,
InviteToRoomAction,
JoinRoomByIdAction,
JoinRoomAction,
LeaveRoomAction,
ForgetRoomAction,
KickAction,
BanAction,
UnbanAction,
SetAccountDataPerRoomAction,
ProcessResponseAction,
SetTypingAction,
PostReceiptAction,
SetReadMarkerAction,
UploadContentAction,
DownloadContentAction,
DownloadThumbnailAction,
SendToDeviceMessageAction,
SendMultipleToDeviceMessagesAction,
UploadIdentityKeysAction,
GenerateAndUploadOneTimeKeysAction,
QueryKeysAction,
ClaimKeysAction,
EncryptMegOlmEventAction,
SetDeviceTrustLevelAction,
SetTrustLevelNeededToSendKeysAction,
PrepareForSharingRoomKeyAction,
ImportFromKeyBackupFileAction,
GetUserProfileAction,
SetAvatarUrlAction,
SetDisplayNameAction,
ResubmitJobAction,
LoadEventsFromStorageAction,
PurgeRoomTimelineAction
>;
using ClientEffect = Effect<ClientAction, lager::deps<>>;
using ClientResult = std::pair<ClientModel, ClientEffect>;
}
diff --git a/src/eventemitter/lagerstoreeventemitter.hpp b/src/eventemitter/lagerstoreeventemitter.hpp
index 978a6f7..d392add 100644
--- a/src/eventemitter/lagerstoreeventemitter.hpp
+++ b/src/eventemitter/lagerstoreeventemitter.hpp
@@ -1,174 +1,174 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <vector>
#include <memory>
#include <algorithm>
#include <lager/store.hpp>
#include <lager/reader.hpp>
#include <lager/event_loop/manual.hpp>
#include "types.hpp"
#include "eventinterface.hpp"
namespace Kazv
{
class LagerStoreEventEmitter : public EventInterface
{
- struct Model { KazvEvent curEvent; };
- struct Action { KazvEvent nextEvent; };
+ struct Model { KazvTrigger curEvent; };
+ struct Action { KazvTrigger nextEvent; };
struct ListenerHolder;
using Result = std::pair<Model,
lager::effect<Action, lager::deps<ListenerHolder &>>>;
- using SlotT = std::function<void(KazvEvent)>;
+ using SlotT = std::function<void(KazvTrigger)>;
struct Listener
{
- void emit(KazvEvent e) {
+ void emit(KazvTrigger e) {
for (const auto &slot: m_slots) {
slot(e);
}
}
void connect(SlotT slot) {
m_slots.push_back(std::move(slot));
}
std::vector<SlotT> m_slots;
};
using ListenerSP = std::shared_ptr<Listener>;
using ListenerWSP = std::weak_ptr<Listener>;
struct ListenerHolder
{
- void sendToListeners(KazvEvent e) {
+ void sendToListeners(KazvTrigger e) {
bool needsCleanup = false;
for (auto listener : m_listeners) {
auto strongListener = listener.lock();
if (strongListener) {
strongListener->emit(e);
} else {
needsCleanup = true;
}
}
if (needsCleanup) {
std::remove_if(m_listeners.begin(),
m_listeners.end(),
[](auto ptr) {
return ptr.expired();
});
}
}
std::vector<ListenerWSP> m_listeners;
};
inline static Result update(Model, Action a) {
return {
Model{a.nextEvent},
[=](auto &&ctx) {
auto &holder = lager::get<ListenerHolder &>(ctx);
holder.sendToListeners(a.nextEvent);
}
};
}
public:
template<class EventLoop>
LagerStoreEventEmitter(EventLoop loop)
: m_holder{}
, m_store(
lager::make_store<Action>(
Model{},
loop,
lager::with_reducer(&update),
lager::with_deps(std::ref(m_holder))))
, m_postingFunc(
[loop=loop](auto &&func) mutable {
loop.post(std::forward<decltype(func)>(func));
}) {}
~LagerStoreEventEmitter() override = default;
- void emit(KazvEvent e) override {
+ void emit(KazvTrigger e) override {
m_store.dispatch(Action{e});
}
class Watchable
{
public:
Watchable(LagerStoreEventEmitter &ee)
: m_listener(std::make_shared<Listener>()) {
ee.addListener(m_listener);
}
template<class EventType, class Func>
void after(Func &&func) {
m_listener->connect(
- [f=std::forward<Func>(func)](KazvEvent e) {
+ [f=std::forward<Func>(func)](KazvTrigger e) {
if (std::holds_alternative<EventType>(e)) {
f(std::get<EventType>(e));
}
});
}
template<class Func>
void afterAll(Func &&func) {
m_listener->connect(
- [f=std::forward<Func>(func)](KazvEvent e) {
+ [f=std::forward<Func>(func)](KazvTrigger e) {
f(e);
});
}
private:
ListenerSP m_listener;
};
/**
* An object you can watch for events.
* ```
* auto watchable = eventEmitter.watchable();
* watchable.after<ReceivingRoomTimelineEvent>(
* [](auto &&e) { ... });
* ```
*/
Watchable watchable() {
return Watchable(*this);
}
private:
void addListener(ListenerSP listener) {
m_postingFunc([=]() {
m_holder.m_listeners.push_back(listener);
});
}
using StoreT =
decltype(lager::make_store<Action>(
Model{},
lager::with_manual_event_loop{},
lager::with_reducer(&update),
lager::with_deps(std::ref(detail::declref<ListenerHolder>()))));
using PostingFunc = std::function<void(std::function<void()>)>;
ListenerHolder m_holder;
StoreT m_store;
PostingFunc m_postingFunc;
};
}
diff --git a/src/tests/client/storage-actions-test.cpp b/src/tests/client/storage-actions-test.cpp
index 794d2b0..baef7f5 100644
--- a/src/tests/client/storage-actions-test.cpp
+++ b/src/tests/client/storage-actions-test.cpp
@@ -1,215 +1,215 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2025 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include "action-mock-utils.hpp"
#include "actions/storage.hpp"
#include <catch2/catch_test_macros.hpp>
#include <catch2/matchers/catch_matchers_quantifiers.hpp>
#include <catch2/matchers/catch_matchers_range_equals.hpp>
#include <catch2/matchers/catch_matchers_predicate.hpp>
#include <factory.hpp>
using namespace Kazv;
using namespace Kazv::Factory;
using Catch::Matchers::AllMatch;
using Catch::Matchers::NoneMatch;
using Catch::Matchers::Predicate;
using Catch::Matchers::UnorderedRangeEquals;
static const std::string roomId = "!test:example.com";
template<class T>
static auto makeEventInRoom(T modifiers)
{
return makeEvent(withEventKV("/room_id"_json_pointer, roomId) | std::move(modifiers));
};
static auto makeEventInRoom()
{
return makeEvent(withEventKV("/room_id"_json_pointer, roomId));
};
TEST_CASE("LoadEventsFromStorageAction", "[client][storage-actions]")
{
auto existingEvents = EventList{
makeEventInRoom(),
makeEventInRoom(),
};
auto c = makeClient(withRoom(makeRoom(withRoomId(roomId) | withRoomTimeline(existingEvents))));
auto timelineEvents = EventList{
makeEventInRoom(),
makeEventInRoom(),
};
auto relatedEvents = EventList{
makeEventInRoom(withEventRelationship("moe.kazv.mxc.some-rel", timelineEvents[0].id())),
makeEventInRoom(withEventRelationship("moe.kazv.mxc.other-rel", existingEvents[0].id())),
};
auto [next, _] = updateClient(c, LoadEventsFromStorageAction{
{{roomId, timelineEvents}},
{{roomId, relatedEvents}},
});
{
auto room = next.roomList.rooms[roomId];
REQUIRE(room.timeline == intoImmer(immer::flex_vector<std::string>{}, zug::map(&Event::id), existingEvents + timelineEvents));
REQUIRE_THAT(existingEvents + timelineEvents + relatedEvents,
AllMatch(Predicate<Event>([&room](const auto &e) { return !!room.messages.count(e.id()); }, "should be in messages")));
REQUIRE(room.reverseEventRelationships.count(timelineEvents[0].id()));
REQUIRE(room.reverseEventRelationships.count(existingEvents[0].id()));
// should not give out any SaveEventsRequested triggers
- REQUIRE_THAT(next.nextTriggers, NoneMatch(Predicate<KazvEvent>([](const auto &t) {
+ REQUIRE_THAT(next.nextTriggers, NoneMatch(Predicate<KazvTrigger>([](const auto &t) {
return std::holds_alternative<SaveEventsRequested>(t);
})));
}
// Verify the load works when loading into timeline an event in
// messages but not in timeline
auto nextTimelineEvents = relatedEvents + EventList{makeEventInRoom()};
std::tie(next, std::ignore) = updateClient(next, LoadEventsFromStorageAction{
{{roomId, nextTimelineEvents}},
{},
});
{
auto room = next.roomList.rooms[roomId];
REQUIRE(room.timeline == intoImmer(immer::flex_vector<std::string>{}, zug::map(&Event::id), existingEvents + timelineEvents + nextTimelineEvents));
}
}
TEST_CASE("Client::loadEventsFromStorage()", "[client][storage-actions]")
{
auto c = makeClient(withRoom(makeRoom(withRoomId(roomId))));
auto u = makeMockSdkUtil(c);
auto d = u.getMockDispatcher(passDown<LoadEventsFromStorageAction>());
auto client = u.getClient(d);
auto events = EventList{makeEvent(), makeEvent()};
client.loadEventsFromStorage({{roomId, events}}, {}).then([&u](const auto &stat) {
REQUIRE(stat);
u.io.stop();
});
u.io.run();
auto room = client.room(roomId);
REQUIRE(room.timelineEventIds().get().size() == events.size());
REQUIRE(room.messagesMap().get().size() == events.size());
REQUIRE(d.calledTimes<LoadEventsFromStorageAction>() == 1);
}
TEST_CASE("maybeAddSaveEventsTrigger", "[client][storage-actions]")
{
WHEN("compare with self") {
auto c = makeClient(withRoom(makeRoom(withRoomTimeline({
makeEvent(),
makeEvent(),
makeEvent(),
}))));
c.maybeAddSaveEventsTrigger(c);
REQUIRE(c.nextTriggers.empty());
}
WHEN("new room, no events") {
auto old = makeClient(withRoom(makeRoom(withRoomTimeline({
makeEvent(),
makeEvent(),
makeEvent(),
}))));
auto c = old;
withRoom(makeRoom())(c);
c.maybeAddSaveEventsTrigger(old);
REQUIRE(c.nextTriggers.empty());
}
WHEN("new room, with events") {
auto old = makeClient(withRoom(makeRoom(withRoomTimeline({
makeEvent(),
makeEvent(),
makeEvent(),
}))));
auto c = old;
auto tl = EventList{
makeEvent(),
makeEvent(),
};
auto room = makeRoom(withRoomTimeline(tl));
withRoom(room)(c);
c.maybeAddSaveEventsTrigger(old);
REQUIRE(c.nextTriggers.size() == 1);
REQUIRE(std::holds_alternative<SaveEventsRequested>(c.nextTriggers[0]));
auto t = std::get<SaveEventsRequested>(c.nextTriggers[0]);
REQUIRE(t.nonTimelineEvents.empty());
REQUIRE(t.timelineEvents.size() == 1);
REQUIRE_THAT(t.timelineEvents[room.roomId], UnorderedRangeEquals(tl));
}
WHEN("new room, with tl and non-tl events") {
auto old = makeClient(withRoom(makeRoom(withRoomTimeline({
makeEvent(),
makeEvent(),
makeEvent(),
}))));
auto c = old;
auto tl = EventList{
makeEvent(),
makeEvent(),
};
auto nonTl = EventList{makeEvent(), makeEvent()};
auto room = makeRoom(withRoomTimeline(tl));
room = RoomModel::update(std::move(room), AddMessagesAction{nonTl});
withRoom(room)(c);
c.maybeAddSaveEventsTrigger(old);
REQUIRE(c.nextTriggers.size() == 1);
REQUIRE(std::holds_alternative<SaveEventsRequested>(c.nextTriggers[0]));
auto t = std::get<SaveEventsRequested>(c.nextTriggers[0]);
REQUIRE(t.nonTimelineEvents.size() == 1);
REQUIRE_THAT(t.nonTimelineEvents[room.roomId], UnorderedRangeEquals(nonTl));
REQUIRE(t.timelineEvents.size() == 1);
REQUIRE_THAT(t.timelineEvents[room.roomId], UnorderedRangeEquals(tl));
}
WHEN("existing room with new and changed events") {
auto origTl = EventList{
makeEvent(),
makeEvent(),
makeEvent(),
};
auto origNonTl = EventList{makeEvent(), makeEvent()};
auto room = makeRoom(withRoomTimeline(origTl));
room = RoomModel::update(std::move(room), AddMessagesAction{origNonTl});
auto old = makeClient(withRoom(room));
auto c = old;
auto tl = EventList{
makeEvent(withEventId(origTl[0].id()) | withEventContent(json{{"mew", "mew"}})),
makeEvent(),
makeEvent(),
};
withRoomTimeline(tl)(room);
auto nonTl = EventList{
makeEvent(withEventId(origNonTl[0].id()) | withEventContent(json{{"mew", "mew"}})),
makeEvent(),
};
room = RoomModel::update(std::move(room), AddMessagesAction{nonTl});
withRoom(room)(c);
c.maybeAddSaveEventsTrigger(old);
REQUIRE(c.nextTriggers.size() == 1);
REQUIRE(std::holds_alternative<SaveEventsRequested>(c.nextTriggers[0]));
auto t = std::get<SaveEventsRequested>(c.nextTriggers[0]);
REQUIRE(t.timelineEvents.size() == 1);
REQUIRE_THAT(t.timelineEvents[room.roomId], UnorderedRangeEquals(tl));
REQUIRE(t.nonTimelineEvents.size() == 1);
REQUIRE_THAT(t.nonTimelineEvents[room.roomId], UnorderedRangeEquals(nonTl));
}
}
diff --git a/src/tests/client/sync-test.cpp b/src/tests/client/sync-test.cpp
index ee344b3..4d802ce 100644
--- a/src/tests/client/sync-test.cpp
+++ b/src/tests/client/sync-test.cpp
@@ -1,766 +1,766 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <catch2/catch_all.hpp>
#include <boost/asio.hpp>
#include <zug/into_vector.hpp>
#include <asio-promise-handler.hpp>
#include <cursorutil.hpp>
#include <sdk-model.hpp>
#include <client/client.hpp>
#include <client/actions/sync.hpp>
#include <outbound-group-session.hpp>
#include "client-test-util.hpp"
#include "factory.hpp"
using namespace Kazv::Factory;
// The example response is adapted from https://matrix.org/docs/spec/client_server/latest
static json syncResponseJson = R"({
"next_batch": "s72595_4483_1934",
"presence": {
"events": [
{
"content": {
"avatar_url": "mxc://localhost:wefuiwegh8742w",
"last_active_ago": 2478593,
"presence": "online",
"currently_active": false,
"status_msg": "Making cupcakes"
},
"type": "m.presence",
"sender": "@example:localhost"
}
]
},
"account_data": {
"events": [
{
"type": "org.example.custom.config",
"content": {
"custom_config_key": "custom_config_value"
}
}
]
},
"rooms": {
"join": {
"!726s6s6q:example.com": {
"summary": {
"m.heroes": [
"@alice:example.com",
"@bob:example.com"
],
"m.joined_member_count": 2,
"m.invited_member_count": 1
},
"state": {
"events": [
{
"content": {
"membership": "join",
"avatar_url": "mxc://example.org/SEsfnsuifSDFSSEF",
"displayname": "Alice Margatroid"
},
"type": "m.room.member",
"event_id": "$143273582443PhrSn:example.org",
"room_id": "!726s6s6q:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"unsigned": {
"age": 1234
},
"state_key": "@alice:example.org"
}
]
},
"timeline": {
"events": [
{
"content": {
"membership": "join",
"avatar_url": "mxc://example.org/SEsfnsuifSDFSSEF",
"displayname": "Alice Margatroid"
},
"type": "m.room.member",
"event_id": "$143273582443PhrSn:example.org",
"room_id": "!726s6s6q:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"unsigned": {
"age": 1234
},
"state_key": "@alice:example.org"
},
{
"content": {
"body": "This is an example text message",
"msgtype": "m.text",
"format": "org.matrix.custom.html",
"formatted_body": "<b>This is an example text message</b>"
},
"type": "m.room.message",
"event_id": "$anothermessageevent:example.org",
"room_id": "!726s6s6q:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"unsigned": {
"age": 1234
}
}
],
"limited": true,
"prev_batch": "t34-23535_0_0"
},
"ephemeral": {
"events": [
{
"content": {
"user_ids": [
"@alice:matrix.org",
"@bob:example.com"
]
},
"type": "m.typing",
"room_id": "!jEsUZKDJdhlrceRyVU:example.org"
}
]
},
"account_data": {
"events": [
{
"content": {
"tags": {
"u.work": {
"order": 0.9
}
}
},
"type": "m.tag"
},
{
"type": "org.example.custom.room.config",
"content": {
"custom_config_key": "custom_config_value"
}
},
{
"type": "m.fully_read",
"content": {
"event_id": "$anothermessageevent:example.org"
}
}
]
}
}
},
"invite": {
"!696r7674:example.com": {
"invite_state": {
"events": [
{
"sender": "@alice:example.com",
"type": "m.room.name",
"state_key": "",
"content": {
"name": "My Room Name"
}
},
{
"sender": "@alice:example.com",
"type": "m.room.member",
"state_key": "@bob:example.com",
"content": {
"membership": "invite"
}
}
]
}
}
},
"leave": {}
},
"to_device": {
"events": [
{
"sender": "@alice:example.com",
"type": "m.new_device",
"content": {
"device_id": "XYZABCDE",
"rooms": ["!726s6s6q:example.com"]
}
}
]
}
})"_json;
static json stateInTimelineResponseJson = R"({
"next_batch": "some-example-value",
"rooms": {
"join": {
"!exampleroomid:example.com": {
"timeline": {
"events": [
{
"content": { "example": "foo" },
"state_key": "",
"event_id": "$example:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"unsigned": { "age": 1234 },
"type": "moe.kazv.mxc.custom.state.type"
}
],
"limited": false
}
}
}
}
})"_json;
static json txnIdResponseJson = R"({
"next_batch": "some-example-value",
"rooms": {
"join": {
"!exampleroomid:example.com": {
"timeline": {
"events": [
{
"content": { "example": "foo" },
"event_id": "$example:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"unsigned": { "age": 1234, "transaction_id": "some-example-txnid" },
"type": "m.room.message"
}
],
"limited": false
}
}
}
}
})"_json;
static auto addNotificationsJson = R"({
"next_batch": "some-example-value",
"account_data": {
"events": [
{
"type": "m.push_rules",
"content": {
"global": {
"override": [{
"rule_id": "moe.kazv.mxc.catch_all",
"default": true,
"enabled": true,
"conditions": [],
"actions": ["notify"]
}]
}
}
}
]
},
"rooms": {
"join": {
"!exampleroomid:example.com": {
"timeline": {
"events": [
{
"content": { "example": "foo" },
"event_id": "$example:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824653,
"type": "m.room.message"
},
{
"content": { "example": "foo2" },
"event_id": "$example2:example.com",
"sender": "@example:example.org",
"origin_server_ts": 1432735824953,
"type": "m.room.message"
}
],
"limited": false
}
}
}
}
})"_json;
static auto receiptJson = R"({
"type": "m.receipt",
"content": {
"$example:example.com": {
"m.read": {
"@bob:example.com": {
"ts": 1432735824653
}
}
}
}
})"_json;
static auto addAndRemoveNotificationsJson = [](auto a, auto r) {
a["rooms"]["join"]["!exampleroomid:example.com"]["ephemeral"] = {
{"events", {
r,
}}
};
return a;
}(addNotificationsJson, receiptJson);
static auto removeNotificationsJson = [](auto r) {
auto j = R"({
"next_batch": "some-example-value",
"rooms": {
"join": {
"!exampleroomid:example.com": {
"timeline": {
"events": [],
"limited": false
}
}
}
}
})"_json;
j["rooms"]["join"]["!exampleroomid:example.com"]["ephemeral"] = {
{"events", {
r,
}}
};
return j;
}(receiptJson);
TEST_CASE("use sync response to update client model", "[client][sync]")
{
using namespace Kazv::CursorOp;
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto store = createTestClientStore(ph);
auto resp = makeResponse(
"Sync",
withResponseJsonBody(syncResponseJson)
| withResponseDataKV("is", "initial")
);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }), store,
std::nullopt);
store.dispatch(ProcessResponseAction{resp});
io.run();
auto rooms = +client.rooms();
std::string roomId = "!726s6s6q:example.com";
SECTION("rooms should be added") {
REQUIRE(rooms.find(roomId));
}
auto r = client.room(roomId);
SECTION("room members should be updated") {
auto members = +r.members();
auto hasAlice = zug::into_vector(
zug::filter([](auto id) { return id == "@alice:example.org"; }),
members)
.size() > 0;
REQUIRE(hasAlice);
}
SECTION("heroes should be updated") {
auto heroIds = +r.heroIds();
REQUIRE(heroIds == immer::flex_vector<std::string>{"@alice:example.com", "@bob:example.com"});
}
SECTION("joined and invited member counts should be updated") {
REQUIRE(r.joinedMemberCount().get() == 2);
REQUIRE(r.invitedMemberCount().get() == 1);
}
SECTION("ephemeral events should be updated") {
auto users = +r.typingUsers();
REQUIRE((users == immer::flex_vector<std::string>{
"@alice:matrix.org",
"@bob:example.com"
}));
}
auto eventId = "$anothermessageevent:example.org"s;
SECTION("timeline should be updated") {
auto timeline = +r.timelineEvents();
auto filtered = zug::into_vector(
zug::filter([=](auto event) { return event.id() == eventId; }),
timeline);
auto hasEvent = filtered.size() > 0;
REQUIRE(hasEvent);
auto onlyOneEvent = filtered.size() == 1;
REQUIRE(onlyOneEvent);
auto ev = filtered[0];
auto eventHasRoomId = ev.originalJson().get().contains("room_id"s);
REQUIRE(eventHasRoomId);
auto gaps = +r.timelineGaps();
// first event in the batch, correspond to its prevBatch
REQUIRE(gaps.at("$143273582443PhrSn:example.org") == "t34-23535_0_0");
}
SECTION("fully read marker should be updated") {
auto readMarker = +r.readMarker();
REQUIRE(readMarker == eventId);
}
SECTION("toDevice should be updated") {
auto toDevice = +client.toDevice();
REQUIRE(toDevice.size() == 1);
REQUIRE(toDevice[0].sender() == "@alice:example.com");
}
SECTION("emits account data changes") {
auto nextTriggers = store.reader().get().nextTriggers;
auto triggerContains = [=](auto p) {
return std::any_of(
nextTriggers.begin(),
nextTriggers.end(),
- [=](const KazvEvent &t) {
+ [=](const KazvTrigger &t) {
if (!std::holds_alternative<ReceivingRoomAccountDataEvent>(t)) {
return false;
}
auto e = std::get<ReceivingRoomAccountDataEvent>(t);
return p(e);
});
};
REQUIRE(triggerContains([](const auto &e) {
return e.event.type() == "m.tag" && e.roomId == "!726s6s6q:example.com";
}));
REQUIRE(triggerContains([](const auto &e) {
return e.event.type() == "m.fully_read" && e.roomId == "!726s6s6q:example.com";
}));
REQUIRE(triggerContains([](const auto &e) {
return e.event.type() == "org.example.custom.room.config" && e.roomId == "!726s6s6q:example.com";
}));
}
}
TEST_CASE("Sync should record state events in timeline", "[client][sync]")
{
using namespace Kazv::CursorOp;
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto store = createTestClientStore(ph);
auto resp = makeResponse(
"Sync",
withResponseJsonBody(stateInTimelineResponseJson)
| withResponseDataKV("is", "initial")
);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }), store,
std::nullopt);
store.dispatch(ProcessResponseAction{resp});
io.run();
auto r = client.room("!exampleroomid:example.com");
auto stateOpt = +r.stateOpt(KeyOfState{"moe.kazv.mxc.custom.state.type", ""});
REQUIRE(stateOpt.has_value());
REQUIRE(stateOpt.value().content().get().at("example") == "foo");
}
TEST_CASE("Sync should remove already sent local echo", "[client][sync]")
{
using namespace Kazv::CursorOp;
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
ClientModel m = makeClient(
withRoom(makeRoom(
withRoomId("!exampleroomid:example.com")
| withAttr(&RoomModel::localEchoes, {
{"some-example-txnid", json{
{"type", "m.room.message"},
{"content", {{"example", "foo"}}}
}},
{"some-other-txnid", json{
{"type", "m.room.message"},
{"content", {{"example", "foo"}}}
}},
})
))
);
auto store = createTestClientStoreFrom(m, ph);
auto resp = makeResponse(
"Sync",
withResponseJsonBody(txnIdResponseJson)
| withResponseDataKV("is", "initial")
);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }), store,
std::nullopt);
store.dispatch(ProcessResponseAction{resp});
io.run();
auto r = client.room("!exampleroomid:example.com");
auto localEchoes = +r.localEchoes();
REQUIRE(localEchoes.size() == 1);
REQUIRE(localEchoes[0].txnId == "some-other-txnid");
}
TEST_CASE("updating local notifications", "[client][sync]")
{
ClientModel m = makeClient(
withRoom(makeRoom(
withRoomId("!exampleroomid:example.com"))));
WHEN("the receipt for the current user did not change") {
auto resp = makeResponse(
"Sync",
withResponseJsonBody(addNotificationsJson)
| withResponseDataKV("is", "incremental")
);
auto [next, _] = processResponse(m, SyncResponse{resp});
auto room = next.roomList.rooms.at("!exampleroomid:example.com");
REQUIRE(room.unreadNotificationEventIds
== immer::flex_vector<std::string>{
"$example:example.com",
"$example2:example.com"
});
THEN("it changed later") {
auto resp = makeResponse(
"Sync",
withResponseJsonBody(removeNotificationsJson)
| withResponseDataKV("is", "incremental")
);
auto [nextNext, _] = processResponse(next, SyncResponse{resp});
auto room = nextNext.roomList.rooms.at("!exampleroomid:example.com");
REQUIRE(room.unreadNotificationEventIds
== immer::flex_vector<std::string>{
"$example2:example.com"
});
}
}
WHEN("the receipt for the current user changed") {
auto resp = makeResponse(
"Sync",
withResponseJsonBody(addAndRemoveNotificationsJson)
| withResponseDataKV("is", "incremental")
);
auto [next, _] = processResponse(m, SyncResponse{resp});
auto room = next.roomList.rooms.at("!exampleroomid:example.com");
REQUIRE(room.unreadNotificationEventIds
== immer::flex_vector<std::string>{
"$example2:example.com"
});
}
}
TEST_CASE("it does not add a gap when the limited field in the timeline is not present (conduwuit)", "[client][sync]")
{
auto body = R"({
"device_one_time_keys_count": {
"signed_curve25519": 721
},
"device_unused_fallback_key_types": null,
"next_batch": "some",
"rooms": {
"join": {
"!foo:example.com": {
"ephemeral": {
"events": []
},
"timeline": {
"events": [
{
"content": {
},
"event_id": "$1",
"origin_server_ts": 1723379000000,
"sender": "@foo:example.com",
"type": "m.room.message",
"unsigned": {
"age": 1,
"transaction_id": "xxxxxx"
}
}
],
"prev_batch": "prev-batch"
},
"unread_notifications": {
"highlight_count": 0,
"notification_count": 0
}
}
}
}
})"_json;
auto resp = makeResponse(
"Sync",
withResponseJsonBody(body)
| withResponseDataKV("is", "incremental")
);
ClientModel m = makeClient(
withRoom(makeRoom(
withRoomId("!foo:example.com"))));
auto [next, _] = processResponse(m, SyncResponse{resp});
auto room = next.roomList.rooms.at("!foo:example.com");
REQUIRE(room.timelineGaps.size() == 0);
REQUIRE(room.messages.count("$1") != 0);
}
TEST_CASE("it does not crash when receiving a malformed encrypted event")
{
auto crypto = makeCrypto();
auto ogs = OutboundGroupSession(
RandomTag{},
genRandomData(crypto.rotateMegOlmSessionRandomSize()),
0);
REQUIRE(ogs.valid());
crypto.createInboundGroupSession(
KeyOfGroupSession{"!foo:example.com", ogs.sessionId()},
ogs.sessionKey(),
crypto.ed25519IdentityKey()
);
ClientModel m = makeClient(
withCrypto(std::move(crypto))
| withRoom(makeRoom(
withRoomId("!foo:example.com")
| withRoomEncrypted(true))));
auto [validEvent, ignore] = m.megOlmEncrypt(Event(json{
{"content", {{"body", "foo"}}},
{"type", "m.room.message"},
{"sender", "@foo:example.com"},
{"origin_server_ts", 0},
}), "!foo:example.com", 0, genRandomData(EncryptMegOlmEventAction::maxRandomSize()));
auto validEventJson = validEvent.originalJson().get();
validEventJson["event_id"] = "$0";
auto [plainText, exceptedErrorCode] = GENERATE(
table<std::string, std::string>({
{"not-json", "M_NOT_JSON"},
{"null", "M_BAD_JSON"},
{"{}", "M_BAD_JSON"},
{R"_({"room_id": "!other:example.com"})_", "MOE.KAZV.MXC_BAD_ROOM_ID"},
}));
auto encrypted = ogs.encrypt(plainText);
auto encryptedEventJson = json{
{"content", {
{"algorithm", CryptoConstants::megOlmAlgo},
{"ciphertext", encrypted},
{"session_id", ogs.sessionId()},
}},
{"event_id", "$1"},
{"sender", "@foo:example.com"},
{"type", "m.room.encrypted"},
{"origin_server_ts", 0},
};
auto body = json{
{"device_one_time_keys_count", {
{"signed_curve25519", 721},
}},
{"device_unused_fallback_key_types", nullptr},
{"next_batch", "some"},
{"rooms", {
{"join", {
{"!foo:example.com", {
{"timeline", {
{"events", {
validEventJson,
encryptedEventJson,
}},
{"prev_batch", "prev-batch"},
}},
}},
}},
}},
};
auto resp = makeResponse(
"Sync",
withResponseJsonBody(body)
| withResponseDataKV("is", "incremental")
);
auto [next, _] = processResponse(m, SyncResponse{resp});
auto messagesMap = next.roomList.rooms.at("!foo:example.com").messages;
REQUIRE(messagesMap.at("$0").decrypted());
REQUIRE(messagesMap.count("$1"));
auto event = messagesMap.at("$1");
REQUIRE(!event.decrypted());
REQUIRE(event.content().get().at("moe.kazv.mxc.errcode") == exceptedErrorCode);
}
TEST_CASE("emit SaveEventsRequested", "[client][sync]")
{
ClientModel m = makeClient(
withRoom(makeRoom(
withRoomId("!exampleroomid:example.com"))));
auto resp = makeResponse(
"Sync",
withResponseJsonBody(syncResponseJson)
| withResponseDataKV("is", "incremental")
);
auto [next, _] = ClientModel::update(m, ProcessResponseAction{SyncResponse{resp}});
auto it = std::find_if(next.nextTriggers.begin(), next.nextTriggers.end(), [](const auto &trigger) {
return std::holds_alternative<SaveEventsRequested>(trigger);
});
REQUIRE(it != next.nextTriggers.end());
auto t = std::get<SaveEventsRequested>(*it);
REQUIRE(t.timelineEvents["!726s6s6q:example.com"].size() == 2);
}
diff --git a/src/tests/event-emitter-test.cpp b/src/tests/event-emitter-test.cpp
index c1050bd..ccf91b5 100644
--- a/src/tests/event-emitter-test.cpp
+++ b/src/tests/event-emitter-test.cpp
@@ -1,83 +1,83 @@
/*
* This file is part of libkazv.
- * SPDX-FileCopyrightText: 2020 Tusooa Zhu <tusooa@kazv.moe>
+ * SPDX-FileCopyrightText: 2020-2026 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <chrono>
#include <catch2/catch_all.hpp>
#include <eventemitter/lagerstoreeventemitter.hpp>
#include <lager/event_loop/boost_asio.hpp>
using namespace Kazv;
TEST_CASE("Event emitter should work normally", "[eventemitter]")
{
boost::asio::io_context ioContext;
LagerStoreEventEmitter ee{lager::with_boost_asio_event_loop{ioContext.get_executor()}};
auto watchable = ee.watchable();
int counter = 0;
- watchable.after<SyncSuccessful>(
+ watchable.after<ReceivingPresenceEvent>(
[&](auto) {
++counter;
});
SECTION("Event handlers should be run properly") {
- ee.emit(SyncSuccessful{});
- ee.emit(SyncFailed{});
+ ee.emit(ReceivingPresenceEvent{Event()});
+ ee.emit(ReceivingRoomTimelineEvent{Event(), ""});
ioContext.run();
REQUIRE(counter == 1);
}
SECTION("Event handlers should be called when every event fires off") {
for (int i = 0; i < 100; ++i) {
- ee.emit(SyncSuccessful{});
+ ee.emit(ReceivingPresenceEvent{Event()});
}
ioContext.run();
REQUIRE(counter == 100);
}
SECTION("Handler should be disconnected once watchable is destroyed") {
int counter2 = 0;
auto guard = boost::asio::make_work_guard(ioContext.get_executor());
auto thread = std::thread([&] { ioContext.run(); });
- ee.emit(SyncFailed{});
+ ee.emit(ReceivingRoomTimelineEvent{});
std::this_thread::sleep_for(std::chrono::milliseconds{100});
{
auto watchable2 = ee.watchable();
- watchable2.after<SyncFailed>(
+ watchable2.after<ReceivingRoomTimelineEvent>(
[&](auto) {
++counter2;
});
- ee.emit(SyncSuccessful{});
- ee.emit(SyncFailed{});
+ ee.emit(ReceivingRoomTimelineEvent{});
+ ee.emit(ReceivingPresenceEvent{});
std::this_thread::sleep_for(std::chrono::milliseconds{100});
}
- ee.emit(SyncFailed{});
+ ee.emit(ReceivingRoomTimelineEvent{});
guard.reset();
thread.join();
REQUIRE(counter == 1);
REQUIRE(counter2 == 1);
}
}

File Metadata

Mime Type
text/x-diff
Expires
Sun, Aug 9, 5:06 AM (1 d, 15 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1724232
Default Alt Text
(155 KB)

Event Timeline