Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85629774
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
155 KB
Referenced Files
None
Subscribers
None
View Options
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
Details
Attached
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)
Attached To
Mode
rL libkazv
Attached
Detach File
Event Timeline
Log In to Comment