Page MenuHomePhorge

No OneTemporary

Size
142 KB
Referenced Files
None
Subscribers
None
diff --git a/src/base/kazvevents.hpp b/src/base/kazvevents.hpp
index 8f713aa..a4c6742 100644
--- a/src/base/kazvevents.hpp
+++ b/src/base/kazvevents.hpp
@@ -1,393 +1,414 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020 Tusooa Zhu
* 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>;
}
diff --git a/src/base/serialization/immer-set.hpp b/src/base/serialization/immer-set.hpp
new file mode 100644
index 0000000..32c6188
--- /dev/null
+++ b/src/base/serialization/immer-set.hpp
@@ -0,0 +1,58 @@
+/*
+ * This file is part of libkazv.
+ * SPDX-FileCopyrightText: 2021-2025 tusooa <tusooa@kazv.moe>
+ * SPDX-License-Identifier: AGPL-3.0-or-later
+ */
+
+
+#pragma once
+
+#include <libkazv-config.hpp>
+
+#include <boost/serialization/nvp.hpp>
+#include <boost/serialization/split_free.hpp>
+
+#include <immer/set.hpp>
+#include <immer/set_transient.hpp>
+
+namespace boost::serialization
+{
+
+ template <class Archive, class T, class H, class E, class MP>
+ void save(Archive &ar, const immer::set<T, H, E, MP> &s, const unsigned int /* version */)
+ {
+ auto size = s.size();
+ ar << BOOST_SERIALIZATION_NVP(size);
+ for (const auto &k : s) {
+ ar << k;
+ }
+ }
+
+ template <class Archive, class T, class H, class E, class MP>
+ void load(Archive &ar, immer::set<T, H, E, MP> &s, const unsigned int /* version */)
+ {
+ using TransientT = decltype(s.transient());
+ using SizeT = decltype(s.size());
+
+ TransientT t{};
+
+ SizeT size{};
+ ar >> BOOST_SERIALIZATION_NVP(size);
+
+ for (auto i = SizeT{}; i < size; ++i) {
+ T k;
+ ar >> k;
+ t.insert(std::move(k));
+ }
+
+ assert(size == t.size());
+ s = t.persistent();
+ }
+
+ template<class Archive, class T, class H, class E, class MP>
+ inline void serialize(Archive &ar, immer::set<T, H, E, MP> &s, const unsigned int version)
+ {
+ boost::serialization::split_free(ar, s, version);
+ }
+
+}
diff --git a/src/client/client-model.cpp b/src/client/client-model.cpp
index 2d75184..497d18a 100644
--- a/src/client/client-model.cpp
+++ b/src/client/client-model.cpp
@@ -1,426 +1,494 @@
/*
* 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 <immer/algorithm.hpp>
#include <lager/util.hpp>
#include <lager/context.hpp>
#include <functional>
#include <zug/transducer/filter.hpp>
#include <immer/flex_vector_transient.hpp>
#include "debug.hpp"
#include "immer-utils.hpp"
#include "json-utils.hpp"
#include "client-model.hpp"
#include "actions/states.hpp"
#include "actions/auth.hpp"
#include "actions/membership.hpp"
#include "actions/paginate.hpp"
#include "actions/send.hpp"
#include "actions/states.hpp"
#include "actions/account-data.hpp"
#include "actions/sync.hpp"
#include "actions/ephemeral.hpp"
#include "actions/content.hpp"
#include "actions/encryption.hpp"
#include "actions/profile.hpp"
#include "actions/storage.hpp"
namespace Kazv
{
auto ClientModel::update(ClientModel m, Action a) -> Result
{
auto oldClient = m;
auto oldDeviceLists = m.deviceLists;
+ auto actionIsStorage = std::holds_alternative<LoadEventsFromStorageAction>(a) || std::holds_alternative<PurgeRoomTimelineAction>(a);
+
auto [newClient, effect] = lager::match(std::move(a))(
[&](RoomListAction a) -> Result {
m.roomList = RoomListModel::update(std::move(m.roomList), a);
return {std::move(m), lager::noop};
},
[&](ResubmitJobAction a) -> Result {
m.addJob(std::move(a.job));
return { std::move(m), lager::noop };
},
[&](auto a) -> decltype(updateClient(m, a)) {
return updateClient(m, a);
},
#define RESPONSE_FOR(_jobId) \
if (r.jobId() == #_jobId) { \
return processResponse(m, _jobId##Response{std::move(r)}); \
}
[&](ProcessResponseAction a) -> Result {
auto r = std::move(a.response);
// auth
RESPONSE_FOR(Login);
RESPONSE_FOR(GetWellknown);
RESPONSE_FOR(GetVersions);
RESPONSE_FOR(Logout);
// paginate
RESPONSE_FOR(GetRoomEvents);
// sync
RESPONSE_FOR(Sync);
RESPONSE_FOR(DefineFilter);
// membership
RESPONSE_FOR(CreateRoom);
RESPONSE_FOR(InviteUser);
RESPONSE_FOR(JoinRoomById);
RESPONSE_FOR(JoinRoom);
RESPONSE_FOR(LeaveRoom);
RESPONSE_FOR(ForgetRoom);
RESPONSE_FOR(Kick);
RESPONSE_FOR(Ban);
RESPONSE_FOR(Unban);
// send
RESPONSE_FOR(SendMessage);
RESPONSE_FOR(SendToDevice);
RESPONSE_FOR(RedactEvent);
// states
RESPONSE_FOR(GetRoomState);
RESPONSE_FOR(SetRoomStateWithKey);
RESPONSE_FOR(GetRoomStateWithKey);
// account data
RESPONSE_FOR(SetAccountData);
RESPONSE_FOR(SetAccountDataPerRoom);
// ephemeral
RESPONSE_FOR(SetTyping);
RESPONSE_FOR(PostReceipt);
RESPONSE_FOR(SetReadMarker);
// content
RESPONSE_FOR(UploadContent);
RESPONSE_FOR(GetContent);
RESPONSE_FOR(GetContentThumbnail);
// encryption
RESPONSE_FOR(UploadKeys);
RESPONSE_FOR(QueryKeys);
RESPONSE_FOR(ClaimKeys);
// profile
RESPONSE_FOR(GetUserProfile);
RESPONSE_FOR(SetAvatarUrl);
RESPONSE_FOR(SetDisplayName);
m.addTrigger(UnrecognizedResponse{std::move(r)});
return { std::move(m), lager::noop };
}
#undef RESPONSE_FOR
);
newClient.maybeRotateSessions(oldClient);
+ // if it is an storage action, do not add it back because the things
+ // should be the same as the ones in storage
+ if (!actionIsStorage) {
+ newClient.maybeAddSaveEventsTrigger(oldClient);
+ }
+
return { std::move(newClient), std::move(effect) };
}
std::pair<Event, std::optional<std::string>> ClientModel::megOlmEncrypt(
Event e, std::string roomId, Timestamp timeMs, RandomData random)
{
if (!crypto) {
kzo.client.dbg() << "We do not have e2ee, so do not encrypt events" << std::endl;
return { e, std::nullopt };
}
if (e.encrypted()) {
kzo.client.dbg() << "The event is already encrypted. Ignoring it." << std::endl;
return { e, std::nullopt };
}
auto j = e.originalJson().get();
auto r = roomList[roomId];
if (! r.encrypted) {
kzo.client.dbg() << "The room " << roomId
<< " is not encrypted, so do not encrypt events" << std::endl;
return { e, std::nullopt };
}
auto desc = r.sessionRotateDesc();
auto keyOpt = std::optional<std::string>{};
if (r.shouldRotateSessionKey) {
kzo.client.dbg() << "We should rotate this session." << std::endl;
keyOpt = withCrypto([&](auto &c) { return c.rotateMegOlmSessionWithRandom(random, timeMs, roomId); });
} else {
keyOpt = withCrypto([&](auto &c) { return c.rotateMegOlmSessionWithRandomIfNeeded(random, timeMs, roomId, desc); });
}
// we no longer need to rotate session
// until next time a device change happens
roomList.rooms = std::move(roomList.rooms)
.update(roomId, [](auto r) { r.shouldRotateSessionKey = false; return r; });
auto relation = hasAtThat(j["content"], "m.relates_to", &json::is_object) ? j["content"]["m.relates_to"] : json(nullptr);
// so that Crypto::encryptMegOlm() can find room id
j["room_id"] = roomId;
auto content = withCrypto([&](auto &c) { return c.encryptMegOlm(j); });
j["type"] = "m.room.encrypted";
j["content"] = std::move(content);
j["content"]["device_id"] = deviceId;
// add relationship to plaintext
if (relation.is_object()) {
j["content"]["m.relates_to"] = relation;
}
return { Event(JsonWrap(j)), keyOpt };
}
immer::map<std::string, immer::map<std::string, Event>> ClientModel::olmEncryptSplit(
Event e,
immer::map<std::string, immer::flex_vector<std::string>> userIdToDeviceIdMap, RandomData random)
{
using ResT = immer::map<std::string, immer::map<std::string, Event>>;
if (!crypto) {
kzo.client.dbg() << "We do not have e2ee, so do not encrypt events" << std::endl;
return ResT{};
}
if (e.encrypted()) {
kzo.client.dbg() << "The event is already encrypted. Ignoring it." << std::endl;
return ResT{};
}
auto origJson = e.originalJson().get();
auto encJson = json::object();
encJson["content"] = json{
{"algorithm", CryptoConstants::olmAlgo},
{"ciphertext", json::object()},
{"sender_key", constCrypto().curve25519IdentityKey()},
};
encJson["type"] = "m.room.encrypted";
ResT messages;
for (auto [userId, devices] : userIdToDeviceIdMap) {
messages = std::move(messages).set(userId, immer::map<std::string, Event>());
for (auto dev : devices) {
auto devInfoOpt = deviceLists.get(userId, dev);
if (! devInfoOpt) {
continue;
}
auto devInfo = devInfoOpt.value();
auto jsonForThisDevice = origJson;
jsonForThisDevice["sender"] = this->userId;
jsonForThisDevice["recipient"] = userId;
jsonForThisDevice["recipient_keys"] = json{
{CryptoConstants::ed25519, devInfo.ed25519Key}
};
jsonForThisDevice["keys"] = json{
{CryptoConstants::ed25519, constCrypto().ed25519IdentityKey()}
};
auto thisEventJson = encJson;
thisEventJson["content"]["ciphertext"]
.merge_patch(withCrypto([&](auto &c) { return c.encryptOlmWithRandom(random, jsonForThisDevice, devInfo.curve25519Key); }));
random.erase(0, Crypto::encryptOlmMaxRandomSize());
messages = setIn(std::move(messages), Event(thisEventJson), userId, dev);
}
}
return messages;
}
immer::flex_vector<std::string /* deviceId */> ClientModel::devicesToSendKeys(std::string userId) const
{
auto trustLevelNeeded = this->trustLevelNeededToSendKeys;
// XXX: preliminary approach
auto shouldSendP = [=](auto deviceInfo, auto /* deviceMap */) {
return deviceInfo.trustLevel >= trustLevelNeeded;
};
auto devices = deviceLists.devicesFor(userId);
return intoImmer(
immer::flex_vector<std::string>{},
zug::filter([=](auto n) {
auto [id, dev] = n;
return shouldSendP(dev, devices);
})
| zug::map([=](auto n) {
return n.first;
}),
devices);
}
std::size_t ClientModel::numOneTimeKeysNeeded() const
{
const auto &crypto = constCrypto();
// Keep half of max supported number of keys
int numUploadedKeys = crypto.uploadedOneTimeKeysCount(CryptoConstants::signedCurve25519);
int numKeysNeeded = crypto.maxNumberOfOneTimeKeys() / 2
- numUploadedKeys;
// Subtract the number of existing one-time keys, in case
// the previous upload was not successful.
int numKeysToGenerate = numKeysNeeded - crypto.numUnpublishedOneTimeKeys();
if (numKeysToGenerate < 0) {
numKeysToGenerate = 0;
}
return numKeysToGenerate;
}
std::size_t EncryptMegOlmEventAction::maxRandomSize()
{
return Crypto::rotateMegOlmSessionRandomSize();
}
std::size_t EncryptMegOlmEventAction::minRandomSize()
{
return 0;
}
std::size_t PrepareForSharingRoomKeyAction::randomSize(PrepareForSharingRoomKeyAction::UserIdToDeviceIdMap devices)
{
auto singleRandomSize = Crypto::encryptOlmMaxRandomSize();
auto deviceNum = accumulate(devices, std::size_t{},
[](auto counter, auto pair) { return counter + pair.second.size(); });
return deviceNum * singleRandomSize;
}
std::size_t GenerateAndUploadOneTimeKeysAction::randomSize(std::size_t numToGen)
{
return Crypto::genOneTimeKeysRandomSize(numToGen);
}
std::size_t ClaimKeysAction::randomSize(immer::map<std::string, immer::flex_vector<std::string>> devicesToSend)
{
auto singleRandomSize = Crypto::createOutboundSessionRandomSize();
auto deviceNum = accumulate(devicesToSend, std::size_t{},
[](auto counter, auto pair) { return counter + pair.second.size(); });
return deviceNum * singleRandomSize;
}
void ClientModel::maybeRotateSessions(ClientModel oldClient)
{
auto roomIds = intoImmer(
immer::flex_vector<std::string>{},
zug::filter([](const auto &pair) {
return pair.second.encrypted && !pair.second.shouldRotateSessionKey;
})
| zug::map([](const auto &pair) { return pair.first; }),
roomList.rooms
);
auto markRotate = [this](const auto &roomId) {
roomList.rooms =
std::move(roomList.rooms)
.update(roomId, [](auto room) {
room.shouldRotateSessionKey = true;
return room;
});
};
// Rotate megolm keys for rooms whose users' device list has changed
auto changedUsers = deviceLists.diff(oldClient.deviceLists);
if (! changedUsers.empty()) {
for (auto roomId : roomIds) {
auto it = std::find_if(changedUsers.begin(), changedUsers.end(),
[=](auto userId) { return roomList.rooms[roomId].hasUser(userId); });
if (it != changedUsers.end()) {
kzo.client.dbg() << "rotate keys for room " << roomId << std::endl;
markRotate(roomId);
}
}
}
roomIds = intoImmer(
immer::flex_vector<std::string>{},
zug::filter([](const auto &pair) {
return pair.second.encrypted && !pair.second.shouldRotateSessionKey;
})
| zug::map([](const auto &pair) { return pair.first; }),
roomList.rooms
);
for (auto roomId : roomIds) {
auto userIds = roomList.rooms[roomId].joinedMemberIds();
auto devicesNotChanged = [oldClient, this](const auto &userId) {
return oldClient.devicesToSendKeys(userId) == devicesToSendKeys(userId);
};
// if any user has the device changes
if (!immer::all_of(userIds, devicesNotChanged)) {
kzo.client.dbg() << "rotate keys for room " << roomId << std::endl;
markRotate(roomId);
}
}
}
auto ClientModel::directRoomMap() const -> immer::map<std::string, std::string>
{
auto directs = accountData["m.direct"].content().get();
auto directItems = directs.items();
return std::accumulate(directItems.begin(), directItems.end(), immer::map<std::string, std::string>(),
[](auto acc, const auto &cur) {
auto [userId, roomIds] = cur;
if (!roomIds.is_array()) {
return acc;
}
for (auto roomId : roomIds) {
if (roomId.is_string()) {
acc = std::move(acc).set(roomId.template get<std::string>(), userId);
}
}
return acc;
}
);
}
auto ClientModel::roomIdsUnderTag(std::string tagId) const -> immer::map<std::string, double>
{
return std::accumulate(
roomList.rooms.begin(), roomList.rooms.end(),
immer::map<std::string, double>{},
[tagId](auto acc, auto cur) {
auto [roomId, room] = cur;
auto tags = room.tags();
if (tags.count(tagId)) {
acc = std::move(acc).set(roomId, tags[tagId]);
}
return acc;
}
);
}
auto ClientModel::roomIdsByTagId() const -> immer::map<std::string, immer::map<std::string, double>>
{
return std::accumulate(
roomList.rooms.begin(), roomList.rooms.end(),
immer::map<std::string, immer::map<std::string, double>>{},
[](auto acc, auto cur) {
auto [roomId, room] = cur;
auto tags = room.tags();
if (tags.empty()) {
acc = setIn(std::move(acc), ROOM_TAG_DEFAULT_ORDER, "", roomId);
} else {
for (const auto &[tagId, order] : tags) {
acc = setIn(std::move(acc), order, tagId, roomId);
}
}
return acc;
}
);
}
const Crypto &ClientModel::constCrypto() const
{
return crypto.value().get();
}
+
+ void ClientModel::maybeAddSaveEventsTrigger(const ClientModel &old)
+ {
+ SaveEventsRequested trigger;
+ using ELT = immer::flex_vector_transient<Event>;
+ auto addToTrigger = [&trigger](const std::string &roomId, ELT tl, ELT nonTl) {
+ if (!tl.empty()) {
+ trigger.timelineEvents = std::move(trigger.timelineEvents).set(roomId, std::move(tl).persistent());
+ }
+ if (!nonTl.empty()) {
+ trigger.nonTimelineEvents = std::move(trigger.nonTimelineEvents).set(roomId, std::move(nonTl).persistent());
+ }
+ };
+ immer::diff(old.roomList.rooms, roomList.rooms, immer::make_differ(
+ /* addedFn = */ [&addToTrigger](const auto &p) {
+ const auto &[roomId, room] = p;
+
+ ELT tl;
+ ELT nonTl;
+ for (const auto &[id, e] : room.messages) {
+ if (room.isInTimeline(id)) {
+ tl.push_back(e);
+ } else {
+ nonTl.push_back(e);
+ }
+ }
+
+ addToTrigger(roomId, std::move(tl), std::move(nonTl));
+ },
+ /* removedFn = */ [](auto &&) {},
+ /* changedFn = */ [&addToTrigger](const auto &p1, const auto &p2) {
+ const auto &roomId = p2.first;
+ const auto &newRoom = p2.second;
+ const auto &oldRoom = p1.second;
+
+ ELT tl;
+ ELT nonTl;
+ auto addEvent = [&tl, &nonTl, &newRoom](const auto &newPair) {
+ const auto &[eventId, event] = newPair;
+ if (newRoom.isInTimeline(eventId)) {
+ tl.push_back(event);
+ } else {
+ nonTl.push_back(event);
+ }
+ };
+ immer::diff(oldRoom.messages, newRoom.messages, immer::make_differ(
+ /* addedFn = */ addEvent,
+ /* removedFn = */ [](auto &&) {},
+ /* changedFn = */ [&addEvent](auto &&, auto &&newPair) {
+ addEvent(std::forward<decltype(newPair)>(newPair));
+ }
+ ));
+
+ addToTrigger(roomId, std::move(tl), std::move(nonTl));
+ }
+ ));
+ if (!trigger.timelineEvents.empty() || !trigger.nonTimelineEvents.empty()) {
+ nextTriggers = std::move(nextTriggers).push_back(trigger);
+ }
+ }
}
diff --git a/src/client/client-model.hpp b/src/client/client-model.hpp
index c4a9cdc..7e5adec 100644
--- a/src/client/client-model.hpp
+++ b/src/client/client-model.hpp
@@ -1,658 +1,660 @@
/*
* 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;
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) {
addTriggers({t});
}
inline void addTriggers(immer::flex_vector<KazvEvent> 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;
};
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/room/room-model.cpp b/src/client/room/room-model.cpp
index 2756017..4d6ca72 100644
--- a/src/client/room/room-model.cpp
+++ b/src/client/room/room-model.cpp
@@ -1,750 +1,764 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020-2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <lager/util.hpp>
#include <zug/sequence.hpp>
#include <zug/transducer/map.hpp>
#include <zug/transducer/filter.hpp>
#include "debug.hpp"
#include "room-model.hpp"
#include "cursorutil.hpp"
#include "immer-utils.hpp"
inline const auto receiptTypes = immer::flex_vector<std::string>{"m.read", "m.read.private"};
template<class Func>
static std::string getMaxInTimeline(std::string a, std::string b, Func sortKey)
{
if (a.empty() && b.empty()) {
return std::string();
} else {
// for an unexisting event, event id is empty and timestamp is 0
// for an existing event, event id is not empty and timestamp >= 0
// so this handles all cases even when we do not have the corresponding event
return std::max(a, b, [=](const std::string &x, const std::string &y) {
return sortKey(x) < sortKey(y);
});
}
}
namespace Kazv
{
PendingRoomKeyEvent makePendingRoomKeyEventV0(std::string txnId, Event event, immer::map<std::string, immer::flex_vector<std::string>> devices)
{
immer::map<std::string, immer::map<std::string, Event>> messages;
for (auto [userId, deviceIds] : devices) {
messages = setIn(std::move(messages), immer::map<std::string, Event>(), userId);
for (auto deviceId : deviceIds) {
messages = setIn(
std::move(messages),
event,
userId, deviceId
);
}
}
return PendingRoomKeyEvent{txnId, messages};
}
auto sortKeyForTimelineEvent(Event e) -> std::tuple<Timestamp, std::string>
{
return std::make_tuple(e.originServerTs(), e.id());
}
RoomModel RoomModel::update(RoomModel r, Action a)
{
return lager::match(std::move(a))(
[&](AddStateEventsAction a) {
r.stateEvents = merge(std::move(r.stateEvents), a.stateEvents, keyOfState);
// If m.room.encryption state event appears,
// configure the room to use encryption.
if (r.stateEvents.find(KeyOfState{"m.room.encryption", ""})) {
auto newRoom = update(std::move(r), SetRoomEncryptionAction{});
r = std::move(newRoom);
}
return r;
},
[&](MaybeAddStateEventsAction a) {
for (auto it = a.stateEvents.rbegin();
it != a.stateEvents.rend();
++it) {
const auto &e = *it;
auto k = keyOfState(e);
if (!r.stateEvents.count(k)) {
r.stateEvents = std::move(r.stateEvents).set(k, e);
}
}
return r;
},
[&](AddMessagesAction a) {
r.messages = merge(std::move(r.messages), a.events, keyOfTimeline);
+ if (!a.alsoInTimeline) {
+ for (const auto &e : a.events) {
+ r.nonTimelineEvents = std::move(r.nonTimelineEvents).insert(keyOfTimeline(e));
+ }
+ }
auto handleRedaction =
[&r](const auto &event) {
if (event.type() == "m.room.redaction") {
auto origJson = event.originalJson().get();
if (origJson.contains("redacts") && origJson.at("redacts").is_string()) {
auto redactedEventId = origJson.at("redacts").template get<std::string>();
if (r.messages.find(redactedEventId)) {
r.messages = std::move(r.messages).update(redactedEventId, [&origJson](const auto &eventToBeRedacted) {
auto newJson = eventToBeRedacted.originalJson().get();
newJson.merge_patch(json{
{"unsigned", {{"redacted_because", std::move(origJson)}}},
});
newJson["content"] = json::object();
return Event(newJson);
});
}
}
}
return event;
};
immer::for_each(a.events, handleRedaction);
// remove all local echoes that are received
for (const auto &e : a.events) {
auto jw = e.originalJson();
const auto &json = jw.get();
if (json.contains("unsigned")
&& json["unsigned"].contains("transaction_id")
&& json["unsigned"]["transaction_id"].is_string()) {
r = update(std::move(r), RemoveLocalEchoAction{json["unsigned"]["transaction_id"].template get<std::string>()});
}
}
// calculate event relationships
r.generateRelationships(a.events);
r.addToUndecryptedEvents(a.events);
return r;
},
[&](AddToTimelineAction a) {
auto eventIds = intoImmer(immer::flex_vector<std::string>(),
zug::map(keyOfTimeline), a.events);
auto oldMessages = r.messages;
- auto next = RoomModel::update(std::move(r), AddMessagesAction{a.events});
+ auto next = RoomModel::update(std::move(r), AddMessagesAction{a.events, /* alsoInTimeline = */ true});
r = std::move(next);
auto key =
[=](auto eventId) {
// sort first by timestamp, then by id
return sortKeyForTimelineEvent(r.messages[eventId]);
};
// Things in messages do not always appear in the timeline.
// Let `exists` function to always return false, so that it is checked for duplicates automatically.
r.timeline = sortedUniqueMerge(r.timeline, eventIds, [](auto &&) { return false; }, key);
+ for (const auto &e : eventIds) {
+ r.nonTimelineEvents = std::move(r.nonTimelineEvents).erase(e);
+ }
+
// We have 3 possibilities for the source of calling this action:
// pagination, sync, or load from storage.
//
// In pagination: `gapEventId.has_value() && !limited.has_value()`
// If we can paginate back: `prevBatch.has_value()`
//
// In sync: `!gapEventId.has_value()`
// If limited: `limited.has_value() && limited.value()`
// If we can paginate back: `prevBatch.has_value()` (should always be present, but we do not add a Gap if it is not limited)
//
// In load from storage: `!gapEventId.has_value() && !limited.has_value() && !prevBatch.has_value()` (Because of actions/storage.cpp)
// Only pagination and sync can add a Gap
if (((a.limited.has_value() && a.limited.value())
|| a.gapEventId.has_value())
&& a.prevBatch.has_value()) {
// this sync is limited, add a Gap here
if (!eventIds.empty()) {
r.timelineGaps = std::move(r.timelineGaps).set(eventIds[0], a.prevBatch.value());
}
}
// Only pagination can remove Gaps
// remove the original Gap, as it is resolved
if (a.gapEventId.has_value()) {
r.timelineGaps = std::move(r.timelineGaps).erase(a.gapEventId.value());
}
// remove all Gaps between the gapped event and the first event in this batch
if (!eventIds.empty() && a.gapEventId.has_value()) {
auto cmp = [=](auto a, auto b) {
return key(a) < key(b);
};
auto thisBatchStart = std::equal_range(r.timeline.begin(), r.timeline.end(), eventIds[0], cmp).first;
auto origBatchStart = std::equal_range(thisBatchStart, r.timeline.end(), a.gapEventId.value(), cmp).first;
// Safety assert: we do not want to execute the for_each if the range is empty,
// or it will go out of bounds.
if (thisBatchStart.index() < origBatchStart.index()) {
std::for_each(thisBatchStart + 1, origBatchStart,
[&](auto eventId) {
r.timelineGaps = std::move(r.timelineGaps).erase(eventId);
});
}
}
return r;
},
[&](AddAccountDataAction a) {
r.accountData = merge(std::move(r.accountData), a.events, keyOfAccountData);
return r;
},
[&](ChangeMembershipAction a) {
r.membership = a.membership;
return r;
},
[&](ChangeInviteStateAction a) {
r.inviteState = merge(immer::map<KeyOfState, Event>{}, a.events, keyOfState);
return r;
},
[&](AddEphemeralAction a) {
auto processReceipt = [&](Event e) {
const auto content = e.content().get();
for (auto [eventId, receipts] : content.items()) {
if (!receipts.is_object()) {
continue;
}
for (auto receiptType : receiptTypes) {
if (!(receipts.contains(receiptType)
&& receipts[receiptType].is_object())) {
continue;
}
for (auto [user, receipt]: receipts[receiptType].items()) {
ReadReceipt readReceipt{
eventId,
0,
};
if (receipt.is_object() && receipt.contains("ts")
&& receipt["ts"].is_number()) {
readReceipt.timestamp = receipt["ts"].template get<Timestamp>();
}
// Remove old receipts
if (r.readReceipts.count(user)) {
auto oldReceiptEventId = r.readReceipts[user].eventId;
if (r.eventReadUsers.count(oldReceiptEventId)) {
auto remaining =
intoImmer(
immer::flex_vector<std::string>{},
zug::filter([user=user](auto userId) {
return userId != user;
}),
r.eventReadUsers[oldReceiptEventId]
);
if (remaining.empty()) {
r.eventReadUsers = std::move(r.eventReadUsers).erase(oldReceiptEventId);
} else {
r.eventReadUsers = std::move(r.eventReadUsers).set(oldReceiptEventId, remaining);
}
}
}
// Add new receipt
r.readReceipts = std::move(r.readReceipts).set(user, readReceipt);
auto oldReadUsers = r.eventReadUsers[eventId];
r.eventReadUsers = std::move(r.eventReadUsers).set(eventId, oldReadUsers.push_back(user));
}
}
}
};
for (auto e : a.events) {
if (e.type() == "m.receipt") {
processReceipt(e);
}
}
r.ephemeral = merge(std::move(r.ephemeral), a.events, keyOfEphemeral);
return r;
},
[&](SetLocalDraftAction a) {
r.localDraft = a.localDraft;
return r;
},
[&](SetRoomEncryptionAction) {
r.encrypted = true;
return r;
},
[&](MarkMembersFullyLoadedAction) {
r.membersFullyLoaded = true;
return r;
},
[&](SetHeroIdsAction a) {
r.heroIds = a.heroIds;
return r;
},
[&](AddLocalEchoAction a) {
auto it = std::find_if(r.localEchoes.begin(), r.localEchoes.end(), [a](const auto &desc) {
return desc.txnId == a.localEcho.txnId;
});
if (it == r.localEchoes.end()) {
r.localEchoes = std::move(r.localEchoes).push_back(a.localEcho);
} else {
r.localEchoes = std::move(r.localEchoes).set(it.index(), a.localEcho);
}
return r;
},
[&](RemoveLocalEchoAction a) {
auto it = std::find_if(r.localEchoes.begin(), r.localEchoes.end(), [a](const auto &desc) {
return desc.txnId == a.txnId;
});
if (it != r.localEchoes.end()) {
r.localEchoes = std::move(r.localEchoes).erase(it.index());
}
return r;
},
[&](AddPendingRoomKeyAction a) {
auto it = std::find_if(r.pendingRoomKeyEvents.begin(), r.pendingRoomKeyEvents.end(), [a](const auto &p) {
return p.txnId == a.pendingRoomKeyEvent.txnId;
});
if (it == r.pendingRoomKeyEvents.end()) {
r.pendingRoomKeyEvents = std::move(r.pendingRoomKeyEvents).push_back(a.pendingRoomKeyEvent);
} else {
r.pendingRoomKeyEvents = std::move(r.pendingRoomKeyEvents).set(it.index(), a.pendingRoomKeyEvent);
}
return r;
},
[&](RemovePendingRoomKeyAction a) {
auto it = std::find_if(r.pendingRoomKeyEvents.begin(), r.pendingRoomKeyEvents.end(), [a](const auto &desc) {
return desc.txnId == a.txnId;
});
if (it != r.pendingRoomKeyEvents.end()) {
r.pendingRoomKeyEvents = std::move(r.pendingRoomKeyEvents).erase(it.index());
}
return r;
},
[&](UpdateJoinedMemberCountAction a) {
r.joinedMemberCount = a.joinedMemberCount;
return r;
},
[&](UpdateInvitedMemberCountAction a) {
r.invitedMemberCount = a.invitedMemberCount;
return r;
},
[&](AddLocalNotificationsAction a) {
auto k = [r](const auto &id) {
return sortKeyForTimelineEvent(r.messages[id]);
};
auto readReceiptForCurrentUser = getMaxInTimeline(r.readReceipts[a.myUserId].eventId, r.localReadMarker, k);
auto newEventIds = intoImmer(
immer::flex_vector<std::string>{},
zug::map(&Event::id),
a.newEvents
);
auto needToAddPredicate = [a, r, k, readReceiptForCurrentUser](const auto &eid) {
if (!readReceiptForCurrentUser.empty()
&& k(readReceiptForCurrentUser) >= k(eid)) {
// this means this event is already read
return false;
}
auto e = r.messages[eid];
return e.sender() != a.myUserId
&& a.pushRulesDesc.handle(e, r).shouldNotify;
};
r.unreadNotificationEventIds = sortedUniqueMerge(std::move(r.unreadNotificationEventIds), newEventIds, [needToAddPredicate](const auto &e) { return !needToAddPredicate(e); }, k);
return r;
},
[&](RemoveReadLocalNotificationsAction a) {
auto k = [r](const auto &id) {
return sortKeyForTimelineEvent(r.messages[id]);
};
auto cmp = [k](const auto &a, const auto &b) {
return k(a) < k(b);
};
auto rr = getMaxInTimeline(
r.readReceipts[a.myUserId].eventId,
r.localReadMarker,
k
);
if (rr.empty()) {
return r;
}
auto it = std::upper_bound(
r.unreadNotificationEventIds.begin(),
r.unreadNotificationEventIds.end(),
rr,
cmp
);
// *it > rr, *(it - 1) <= rr (if it - 1 is valid)
if (it == r.unreadNotificationEventIds.end()) {
// If it == end(), it means everything is read
r.unreadNotificationEventIds = {};
} else if (it == r.unreadNotificationEventIds.begin()) {
// If it == begin(), it means everything is unread, so nothing to do
} else {
// it is somewhere in the middle, pointing to the first element that is unread
r.unreadNotificationEventIds = std::move(r.unreadNotificationEventIds).erase(0, it.index());
}
return r;
},
[&](UpdateLocalReadMarkerAction a) {
r.localReadMarker = a.localReadMarker;
auto next = RoomModel::update(std::move(r), RemoveReadLocalNotificationsAction{a.myUserId});
return next;
},
[&](PurgeEventsAction a) {
if (r.timeline.size() <= a.maxToKeep) {
return r;
}
auto numToDrop = r.timeline.size() - a.maxToKeep;
auto keepEvents = intoImmer(
EventList{},
zug::map([&r](const auto &eventId) {
return r.messages[eventId];
}),
std::move(r.timeline).drop(numToDrop)
);
auto origMessages = r.messages;
r.reverseEventRelationships = {};
r.timeline = {};
r.messages = {};
r.undecryptedEvents = {};
auto msgs = intoImmer(
r.localReadMarker.empty() ? EventList{} : EventList{origMessages[r.localReadMarker]},
zug::filter([l=r.localReadMarker](const auto &eventId) {
return eventId != l;
})
| zug::map([&origMessages](const auto &eventId) {
return origMessages[eventId];
}),
r.unreadNotificationEventIds
);
auto next = update(std::move(r), AddMessagesAction{msgs});
return update(std::move(next), AddToTimelineAction{
keepEvents,
std::nullopt,
std::nullopt,
std::nullopt,
});
}
);
}
RoomListModel RoomListModel::update(RoomListModel l, Action a)
{
return lager::match(std::move(a))(
[&](UpdateRoomAction a) {
l.rooms = std::move(l.rooms)
.update(a.roomId,
[=](RoomModel oldRoom) {
oldRoom.roomId = a.roomId; // in case it is a new room
return RoomModel::update(std::move(oldRoom), a.roomAction);
});
return l;
}
);
}
static auto membershipTransducer(const std::string &membership)
{
return zug::filter([](auto val) {
auto [k, v] = val;
auto [type, stateKey] = k;
return type == "m.room.member"s;
})
| zug::map([](auto val) {
auto [k, v] = val;
auto [type, stateKey] = k;
return std::pair<std::string, Kazv::Event>{stateKey, v};
})
| zug::filter([&membership](auto val) {
auto [stateKey, ev] = val;
return ev.content().get()
.at("membership"s) == membership;
});
}
static auto memberIdsByMembership(immer::map<KeyOfState, Event> stateEvents, const std::string &membership)
{
return intoImmer(
immer::flex_vector<std::string>{},
membershipTransducer(membership)
| zug::map([](auto val) {
auto [stateKey, ev] = val;
return stateKey;
}),
stateEvents);
}
auto memberEventsByMembership(immer::map<KeyOfState, Event> stateEvents, const std::string &membership)
{
return intoImmer(
EventList{},
membershipTransducer(membership)
| zug::map([](auto val) {
auto [stateKey, ev] = val;
return ev;
}),
stateEvents);
}
immer::flex_vector<std::string> RoomModel::joinedMemberIds() const
{
return memberIdsByMembership(stateEvents, "join"s);
}
immer::flex_vector<std::string> RoomModel::invitedMemberIds() const
{
return memberIdsByMembership(stateEvents, "invite"s);
}
immer::flex_vector<std::string> RoomModel::knockedMemberIds() const
{
return memberIdsByMembership(stateEvents, "knock"s);
}
immer::flex_vector<std::string> RoomModel::leftMemberIds() const
{
return memberIdsByMembership(stateEvents, "leave"s);
}
immer::flex_vector<std::string> RoomModel::bannedMemberIds() const
{
return memberIdsByMembership(stateEvents, "ban"s);
}
EventList RoomModel::joinedMemberEvents() const
{
return memberEventsByMembership(stateEvents, "join"s);
}
EventList RoomModel::invitedMemberEvents() const
{
return memberEventsByMembership(stateEvents, "invite"s);
}
EventList RoomModel::knockedMemberEvents() const
{
return memberEventsByMembership(stateEvents, "knock"s);
}
EventList RoomModel::leftMemberEvents() const
{
return memberEventsByMembership(stateEvents, "leave"s);
}
EventList RoomModel::bannedMemberEvents() const
{
return memberEventsByMembership(stateEvents, "ban"s);
}
EventList RoomModel::heroMemberEvents() const
{
return intoImmer(
EventList{},
zug::filter([heroIds=heroIds](auto val) {
auto [k, ev] = val;
auto [type, stateKey] = k;
return type == "m.room.member"s &&
std::find(heroIds.begin(), heroIds.end(), stateKey) != heroIds.end();
})
| zug::map([](auto val) {
auto [_, ev] = val;
return ev;
}),
stateEvents);
}
static Timestamp defaultRotateMs = 604800000;
static int defaultRotateMsgs = 100;
MegOlmSessionRotateDesc RoomModel::sessionRotateDesc() const
{
auto k = KeyOfState{"m.room.encryption", ""};
auto content = stateEvents[k].content().get();
auto ms = content.contains("rotation_period_ms")
? content["rotation_period_ms"].get<Timestamp>()
: defaultRotateMs;
auto msgs = content.contains("rotation_period_msgs")
? content["rotation_period_msgs"].get<int>()
: defaultRotateMsgs;
return MegOlmSessionRotateDesc{ ms, msgs };
}
bool RoomModel::hasUser(std::string userId) const
{
try {
auto ev = stateEvents.at(KeyOfState{"m.room.member", userId});
if (ev.content().get().at("membership") == "join") {
return true;
}
} catch (const std::exception &) {
return false;
}
return false;
}
std::optional<LocalEchoDesc> RoomModel::getLocalEchoByTxnId(std::string txnId) const
{
auto it = std::find_if(localEchoes.begin(), localEchoes.end(), [txnId](const auto &desc) {
return txnId == desc.txnId;
});
if (it != localEchoes.end()) {
return *it;
} else {
return std::nullopt;
}
}
std::optional<PendingRoomKeyEvent> RoomModel::getPendingRoomKeyEventByTxnId(std::string txnId) const
{
auto it = std::find_if(pendingRoomKeyEvents.begin(), pendingRoomKeyEvents.end(), [txnId](const auto &desc) {
return txnId == desc.txnId;
});
if (it != pendingRoomKeyEvents.end()) {
return *it;
} else {
return std::nullopt;
}
}
static double getTagOrder(const json &tag)
{
// https://spec.matrix.org/v1.7/client-server-api/#events-12
// If a room has a tag without an order key then it should appear after the rooms with that tag that have an order key.
return tag.contains("order") && tag["order"].is_number()
? tag["order"].template get<double>()
: ROOM_TAG_DEFAULT_ORDER;
}
immer::map<std::string, double> RoomModel::tags() const
{
auto content = accountData["m.tag"].content().get();
if (!content.contains("tags") || !content["tags"].is_object()) {
return {};
}
auto tagsObject = content["tags"];
auto tagsItems = tagsObject.items();
return std::accumulate(tagsItems.begin(), tagsItems.end(), immer::map<std::string, double>(),
[=](auto acc, const auto &cur) {
auto [id, tag] = cur;
return std::move(acc).set(id, getTagOrder(tag));
}
);
}
static auto normalizeTagEventJson(Event e)
{
auto content = e.content().get();
if (!content.contains("tags") || !content["tags"].is_object()) {
content["tags"] = json::object();
}
return json{
{"content", content},
{"type", "m.tag"},
};
}
Event RoomModel::makeAddTagEvent(std::string tagId, std::optional<double> order) const
{
auto eventJson = normalizeTagEventJson(accountData["m.tag"]);
auto tag = json::object();
if (order.has_value()) {
tag["order"] = order.value();
}
eventJson["content"]["tags"][tagId] = tag;
return Event(eventJson);
}
Event RoomModel::makeRemoveTagEvent(std::string tagId) const
{
auto eventJson = normalizeTagEventJson(accountData["m.tag"]);
eventJson["content"]["tags"].erase(tagId);
return Event(eventJson);
}
void RoomModel::generateRelationships(EventList newEvents)
{
for (const auto &event: newEvents) {
auto [relType, eventId] = event.relationship();
if (!relType.empty()) {
reverseEventRelationships = updateIn(std::move(reverseEventRelationships), [event](auto &&evs) {
return evs.push_back(event.id());
}, eventId, relType);
}
}
}
void RoomModel::regenerateRelationships()
{
generateRelationships(intoImmer(EventList{}, zug::map([](const auto &kv) {
return kv.second;
}), messages));
}
void RoomModel::addToUndecryptedEvents(EventList newEvents)
{
if (!encrypted) {
return;
}
for (auto event : newEvents) {
if (event.encrypted() && !event.decrypted()) {
auto original = event.originalJson();
const auto &o = original.get();
if (o.contains("content")
&& o["content"].contains("session_id")
&& o["content"]["session_id"].is_string()
) {
auto sessionId = o["content"]["session_id"].template get<std::string>();
undecryptedEvents = std::move(undecryptedEvents)
.update(sessionId, [id=event.id()](const auto &v) {
return v.push_back(id);
});
}
}
}
}
void RoomModel::recalculateUndecryptedEvents()
{
if (!encrypted) {
return;
}
addToUndecryptedEvents(intoImmer(EventList{}, zug::map([](const auto &kv) {
return kv.second;
}), messages));
}
bool RoomModel::checkInvariants() const
{
auto inMessages = [this](const std::string &eventId) {
return !!messages.count(eventId);
};
return immer::all_of(timeline, inMessages)
&& immer::all_of(unreadNotificationEventIds, inMessages)
&& std::all_of(undecryptedEvents.begin(), undecryptedEvents.end(), [&inMessages](const auto &p) {
return immer::all_of(p.second, inMessages);
})
&& (localReadMarker.empty() || messages.count(localReadMarker))
&& std::all_of(reverseEventRelationships.begin(), reverseEventRelationships.end(), [&inMessages](const auto &p) {
const auto &relTypeToEventsMap = p.second;
return std::all_of(relTypeToEventsMap.begin(), relTypeToEventsMap.end(), [&inMessages](const auto &p2) {
const auto &events = p2.second;
return immer::all_of(events, inMessages);
});
});
}
+
+ bool RoomModel::isInTimeline(const std::string &eventId) const
+ {
+ return !nonTimelineEvents.count(eventId) && messages.count(eventId);
+ }
}
diff --git a/src/client/room/room-model.hpp b/src/client/room/room-model.hpp
index 71189e8..b9b0c2f 100644
--- a/src/client/room/room-model.hpp
+++ b/src/client/room/room-model.hpp
@@ -1,485 +1,506 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2021-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <string>
#include <variant>
#include <immer/flex_vector.hpp>
#include <immer/map.hpp>
+#include <immer/set.hpp>
#include <serialization/immer-flex-vector.hpp>
#include <serialization/immer-box.hpp>
#include <serialization/immer-map.hpp>
+#include <serialization/immer-set.hpp>
#include <serialization/immer-array.hpp>
#include <csapi/sync.hpp>
#include <event.hpp>
#include <crypto.hpp>
#include "push-rules-desc.hpp"
#include "local-echo.hpp"
#include "clientutil.hpp"
namespace Kazv
{
struct PendingRoomKeyEvent
{
std::string txnId;
immer::map<std::string, immer::map<std::string, Event>> messages;
friend bool operator==(const PendingRoomKeyEvent &a, const PendingRoomKeyEvent &b) = default;
friend bool operator!=(const PendingRoomKeyEvent &a, const PendingRoomKeyEvent &b) = default;
};
PendingRoomKeyEvent makePendingRoomKeyEventV0(std::string txnId, Event event, immer::map<std::string, immer::flex_vector<std::string>> devices);
struct ReadReceipt
{
std::string eventId;
Timestamp timestamp;
friend bool operator==(const ReadReceipt &a, const ReadReceipt &b) = default;
friend bool operator!=(const ReadReceipt &a, const ReadReceipt &b) = default;
};
template<class Archive>
void serialize(Archive &ar, ReadReceipt &r, std::uint32_t const /* version */)
{
ar & r.eventId & r.timestamp;
}
struct EventReader
{
std::string userId;
Timestamp timestamp;
friend bool operator==(const EventReader &a, const EventReader &b) = default;
friend bool operator!=(const EventReader &a, const EventReader &b) = default;
};
struct AddStateEventsAction
{
immer::flex_vector<Event> stateEvents;
};
/// Go from the back of stateEvents to the beginning,
/// adding the event to room state only if the room
/// has no state event with that state key.
struct MaybeAddStateEventsAction
{
immer::flex_vector<Event> stateEvents;
};
/// Add events to the messages map, but not the timeline.
/// Usually because their position in the timeline is not known.
struct AddMessagesAction
{
EventList events;
+ /// @internal only to be used by AddToTimelineAction
+ bool alsoInTimeline{false};
};
struct AddToTimelineAction
{
/// Events from oldest to latest
immer::flex_vector<Event> events;
std::optional<std::string> prevBatch;
std::optional<bool> limited;
std::optional<std::string> gapEventId;
};
struct AddAccountDataAction
{
immer::flex_vector<Event> events;
};
struct ChangeMembershipAction
{
RoomMembership membership;
};
struct ChangeInviteStateAction
{
immer::flex_vector<Event> events;
};
struct AddEphemeralAction
{
EventList events;
};
struct SetLocalDraftAction
{
std::string localDraft;
};
struct SetRoomEncryptionAction
{
};
struct MarkMembersFullyLoadedAction
{
};
struct SetHeroIdsAction
{
immer::flex_vector<std::string> heroIds;
};
struct AddLocalEchoAction
{
LocalEchoDesc localEcho;
};
struct RemoveLocalEchoAction
{
std::string txnId;
};
struct AddPendingRoomKeyAction
{
PendingRoomKeyEvent pendingRoomKeyEvent;
};
struct RemovePendingRoomKeyAction
{
std::string txnId;
};
struct UpdateJoinedMemberCountAction
{
std::size_t joinedMemberCount;
};
struct UpdateInvitedMemberCountAction
{
std::size_t invitedMemberCount;
};
/// Update local notifications to include the new events
///
/// Precondition: newEvents are already in room.messages
struct AddLocalNotificationsAction
{
EventList newEvents;
PushRulesDesc pushRulesDesc;
std::string myUserId;
};
/// Remove local notifications that are already read
struct RemoveReadLocalNotificationsAction
{
std::string myUserId;
};
/// Update the local read marker, removing any read notifications before it.
struct UpdateLocalReadMarkerAction
{
std::string localReadMarker;
std::string myUserId;
};
/// Remove events from the model, retaining only the latest `maxToKeep` events.
struct PurgeEventsAction
{
std::size_t maxToKeep;
};
inline const double ROOM_TAG_DEFAULT_ORDER = 2;
template<class Archive>
void serialize(Archive &ar, PendingRoomKeyEvent &e, std::uint32_t const version)
{
if (version < 1) {
// loading an older version where there is only one event
std::string txnId;
Event event;
immer::map<std::string, immer::flex_vector<std::string>> devices;
ar & txnId & event & devices;
e = makePendingRoomKeyEventV0(
std::move(txnId), std::move(event), std::move(devices));
} else {
ar & e.txnId & e.messages;
}
}
/**
* Get the sort key for a timeline event.
*
* If the key is larger, the event should be placed
* at the more recent end of the timeline.
*
* @param e The event to get the sort key for.
* @return The sort key. You MUST use `auto` to store the
* result.
*/
auto sortKeyForTimelineEvent(Event e) -> std::tuple<Timestamp, std::string>;
/**
* The model to store information about a room.
*
* Room invariants:
* Any event in timeline is in messages.
* Any event in undecryptedEvents is in messages.
* Any event in unreadNotificationEventIds is in messages.
* Any relater (i.e. child) event in reverseEventRelationships is in messages.
* localReadMarker (if not empty) is in messages.
*/
struct RoomModel
{
using Membership = RoomMembership;
using ReverseEventRelationshipMap = immer::map<
std::string /* related event id */,
immer::map<std::string /* relation type */, immer::flex_vector<std::string /* relater event id */>>>;
std::string roomId;
immer::map<KeyOfState, Event> stateEvents;
immer::map<KeyOfState, Event> inviteState;
// Smaller indices mean earlier events
// (oldest) 0 --------> n (latest)
immer::flex_vector<std::string> timeline;
immer::map<std::string, Event> messages;
immer::map<std::string, Event> accountData;
Membership membership{};
std::string paginateBackToken;
/// whether this room has earlier events to be fetched
bool canPaginateBack{true};
immer::map<std::string /* eventId */, std::string /* prevBatch */> timelineGaps;
immer::map<std::string, Event> ephemeral;
std::string localDraft;
bool encrypted{false};
/// a marker to indicate whether we need to rotate
/// the session key earlier than it expires
/// (e.g. when a user in the room's device list changed
/// or when someone joins or leaves)
bool shouldRotateSessionKey{true};
bool membersFullyLoaded{false};
immer::flex_vector<std::string> heroIds;
immer::flex_vector<LocalEchoDesc> localEchoes;
immer::flex_vector<PendingRoomKeyEvent> pendingRoomKeyEvents;
ReverseEventRelationshipMap reverseEventRelationships;
std::size_t joinedMemberCount{0};
std::size_t invitedMemberCount{0};
/// The local read marker for this room. Indicates that
/// you have read up to this event.
std::string localReadMarker;
/// The local unread count for this room.
std::size_t localUnreadCount{0};
/// The local unread notification count for this room.
/// XXX this is never used.
std::size_t localNotificationCount{0};
/// Read receipts for all users
immer::map<std::string /* userId */, ReadReceipt> readReceipts;
/// A map from event id to a list of users that has read
/// receipt at that point
immer::map<
std::string /* eventId */,
immer::flex_vector<std::string /* userId */>> eventReadUsers;
/// A map from the session id to a list of event ids of events
/// that cannot (yet) be decrypted.
immer::map<
std::string /* sessionId */,
immer::flex_vector<std::string /* eventId */>> undecryptedEvents;
immer::flex_vector<std::string> unreadNotificationEventIds;
+ /// The set of event ids that are in `messages` but not in `timeline`.
+ /// The rationale is that non-timeline events are sparse, so
+ /// instead of recording events that are in the timeline, we
+ /// record those not in the timeline.
+ immer::set<std::string> nonTimelineEvents;
immer::flex_vector<std::string> joinedMemberIds() const;
immer::flex_vector<std::string> invitedMemberIds() const;
immer::flex_vector<std::string> knockedMemberIds() const;
immer::flex_vector<std::string> leftMemberIds() const;
immer::flex_vector<std::string> bannedMemberIds() const;
EventList joinedMemberEvents() const;
EventList invitedMemberEvents() const;
EventList knockedMemberEvents() const;
EventList leftMemberEvents() const;
EventList bannedMemberEvents() const;
EventList heroMemberEvents() const;
MegOlmSessionRotateDesc sessionRotateDesc() const;
bool hasUser(std::string userId) const;
std::optional<LocalEchoDesc> getLocalEchoByTxnId(std::string txnId) const;
std::optional<PendingRoomKeyEvent> getPendingRoomKeyEventByTxnId(std::string txnId) const;
immer::map<std::string, double> tags() const;
Event makeAddTagEvent(std::string tagId, std::optional<double> order) const;
Event makeRemoveTagEvent(std::string tagId) const;
/**
* Fill in reverseEventRelationships by gathering
* the relationships specified in `newEvents`
*
* @param newEvents The events that just came in after last time event relationships
* are gathered.
*/
void generateRelationships(EventList newEvents);
void regenerateRelationships();
/**
* Fill in undecryptedEvents by gathering
* the session ids specified in `newEvents`.
*
* @param newEvents New incoming events.
*/
void addToUndecryptedEvents(EventList newEvents);
void recalculateUndecryptedEvents();
/**
* Check if the invariants in the model are satisfied.
* @return true iff the invariants are satisfied.
*/
bool checkInvariants() const;
+ /**
+ * Check if the event is in the timeline.
+ *
+ * This function takes constant time.
+ * @param eventId The id of the event to check.
+ * @return true iff the event is in the timeline.
+ */
+ bool isInTimeline(const std::string &eventId) const;
+
using Action = std::variant<
AddStateEventsAction,
MaybeAddStateEventsAction,
AddMessagesAction,
AddToTimelineAction,
AddAccountDataAction,
ChangeMembershipAction,
ChangeInviteStateAction,
AddEphemeralAction,
SetLocalDraftAction,
SetRoomEncryptionAction,
MarkMembersFullyLoadedAction,
SetHeroIdsAction,
AddLocalEchoAction,
RemoveLocalEchoAction,
AddPendingRoomKeyAction,
RemovePendingRoomKeyAction,
UpdateJoinedMemberCountAction,
UpdateInvitedMemberCountAction,
AddLocalNotificationsAction,
RemoveReadLocalNotificationsAction,
UpdateLocalReadMarkerAction,
PurgeEventsAction
>;
static RoomModel update(RoomModel r, Action a);
friend bool operator==(const RoomModel &a, const RoomModel &b) = default;
};
using RoomAction = RoomModel::Action;
struct UpdateRoomAction
{
std::string roomId;
RoomAction roomAction;
};
struct RoomListModel
{
immer::map<std::string, RoomModel> rooms;
inline auto at(std::string id) const { return rooms.at(id); }
inline auto operator[](std::string id) const { return rooms[id]; }
inline bool has(std::string id) const { return rooms.find(id); }
using Action = std::variant<
UpdateRoomAction
>;
static RoomListModel update(RoomListModel l, Action a);
friend bool operator==(const RoomListModel &a, const RoomListModel &b) = default;
};
using RoomListAction = RoomListModel::Action;
template<class Archive>
void serialize(Archive &ar, RoomModel &r, std::uint32_t const version)
{
ar
& r.roomId
& r.stateEvents
& r.inviteState
& r.timeline
& r.messages
& r.accountData
& r.membership
& r.paginateBackToken
& r.canPaginateBack
& r.timelineGaps
& r.ephemeral
& r.localDraft
& r.encrypted
& r.shouldRotateSessionKey
& r.membersFullyLoaded
;
if (version >= 1) {
ar
& r.heroIds
;
}
if (version >= 2) {
ar & r.localEchoes;
}
if (version >= 3) {
ar & r.pendingRoomKeyEvents;
}
if (version >= 4) {
ar & r.reverseEventRelationships;
} else { // must be reading from an older version
if constexpr (typename Archive::is_loading()) {
r.regenerateRelationships();
}
}
if (version >= 5) {
ar & r.joinedMemberCount & r.invitedMemberCount;
}
if (version >= 6) {
ar
& r.localReadMarker
& r.localUnreadCount
& r.localNotificationCount
& r.readReceipts
& r.eventReadUsers;
}
if (version >= 7) {
ar & r.undecryptedEvents;
} else {
if constexpr (typename Archive::is_loading()) {
r.recalculateUndecryptedEvents();
}
}
if (version >= 8) {
ar & r.unreadNotificationEventIds;
}
+ if (version >= 9) {
+ ar & r.nonTimelineEvents;
+ }
}
template<class Archive>
void serialize(Archive &ar, RoomListModel &l, std::uint32_t const /*version*/)
{
ar & l.rooms;
}
}
BOOST_CLASS_VERSION(Kazv::PendingRoomKeyEvent, 1)
BOOST_CLASS_VERSION(Kazv::ReadReceipt, 0)
-BOOST_CLASS_VERSION(Kazv::RoomModel, 8)
+BOOST_CLASS_VERSION(Kazv::RoomModel, 9)
BOOST_CLASS_VERSION(Kazv::RoomListModel, 0)
diff --git a/src/tests/CMakeLists.txt b/src/tests/CMakeLists.txt
index 8c4bb61..5222c53 100644
--- a/src/tests/CMakeLists.txt
+++ b/src/tests/CMakeLists.txt
@@ -1,127 +1,128 @@
include(CTest)
set(KAZVTEST_RESPATH ${CMAKE_CURRENT_SOURCE_DIR}/resources)
configure_file(kazvtest-respath.hpp.in kazvtest-respath.hpp)
function(libkazv_add_tests)
set(options "")
set(oneValueArgs "")
set(multiValueArgs EXTRA_LINK_LIBRARIES EXTRA_INCLUDE_DIRECTORIES)
cmake_parse_arguments(PARSE_ARGV 0 libkazv_add_tests "${options}" "${oneValueArgs}" "${multiValueArgs}")
foreach(test_source ${libkazv_add_tests_UNPARSED_ARGUMENTS})
string(REGEX REPLACE "\\.cpp$" "" test_executable "${test_source}")
string(REGEX REPLACE "/|\\\\" "--" test_executable "${test_executable}")
message(STATUS "Test ${test_executable} added")
add_executable("${test_executable}" "${test_source}")
target_link_libraries("${test_executable}"
PRIVATE Catch2::Catch2WithMain
Threads::Threads
${libkazv_add_tests_EXTRA_LINK_LIBRARIES}
)
target_include_directories(
"${test_executable}"
PRIVATE ${CMAKE_CURRENT_BINARY_DIR}
${CMAKE_CURRENT_SOURCE_DIR}
${CMAKE_CURRENT_SOURCE_DIR}/..
${libkazv_add_tests_EXTRA_INCLUDE_DIRECTORIES}
)
target_compile_definitions("${test_executable}" PRIVATE CATCH_CONFIG_ENABLE_ALL_STRINGMAKERS)
add_test(NAME "${test_executable}" COMMAND "${test_executable}" "--allow-running-no-tests" "~[needs-internet]")
endforeach()
endfunction()
libkazv_add_tests(
event-test.cpp
cursorutiltest.cpp
base/serialization-test.cpp
base/types-test.cpp
base/immer-utils-test.cpp
base/json-utils-test.cpp
EXTRA_LINK_LIBRARIES kazvbase
)
add_library(client-test-lib SHARED client/client-test-util.cpp)
target_link_libraries(client-test-lib PUBLIC kazvjob kazvclient)
libkazv_add_tests(
client/discovery-test.cpp
client/sync-test.cpp
client/content-test.cpp
client/paginate-test.cpp
client/storage-actions-test.cpp
client/util-test.cpp
client/serialization-test.cpp
client/encrypted-file-test.cpp
client/sdk-test.cpp
client/thread-safety-test.cpp
client/room-test.cpp
client/random-generator-test.cpp
client/profile-test.cpp
client/kick-test.cpp
client/ban-test.cpp
client/join-test.cpp
client/keys-test.cpp
client/device-ops-test.cpp
client/send-test.cpp
client/encryption-test.cpp
client/redact-test.cpp
client/tagging-test.cpp
client/account-data-test.cpp
client/room/room-actions-test.cpp
client/room/local-echo-test.cpp
client/room/event-relationships-test.cpp
client/room/member-membership-test.cpp
client/room/purge-test.cpp
client/push-rules-desc-test.cpp
client/notification-handler-test.cpp
client/validator-test.cpp
client/power-levels-desc-test.cpp
client/client-test.cpp
client/create-room-test.cpp
client/device-list-tracker-test.cpp
client/device-list-tracker-benchmark-test.cpp
client/room/read-receipt-test.cpp
client/room/undecrypted-events-test.cpp
client/encryption-benchmark-test.cpp
client/logout-test.cpp
client/room/pinned-events-test.cpp
client/get-versions-test.cpp
client/alias-test.cpp
client/encode-test.cpp
+ client/maybe-add-save-events-trigger-benchmark-test.cpp
EXTRA_LINK_LIBRARIES kazvclient kazveventemitter kazvjob client-test-lib kazvtestfixtures
EXTRA_INCLUDE_DIRECTORIES ${CMAKE_CURRENT_SOURCE_DIR}/client
)
libkazv_add_tests(
basejobtest.cpp
kazvjobtest.cpp
file-desc-test.cpp
EXTRA_LINK_LIBRARIES kazvbase kazvjob
)
libkazv_add_tests(
promise-test.cpp
EXTRA_LINK_LIBRARIES kazvbase kazvjob kazvstore
)
libkazv_add_tests(
event-emitter-test.cpp
EXTRA_LINK_LIBRARIES kazvbase kazveventemitter
)
libkazv_add_tests(
crypto-test.cpp
crypto/inbound-group-session-test.cpp
crypto/outbound-group-session-test.cpp
crypto/session-test.cpp
crypto/key-export-test.cpp
EXTRA_LINK_LIBRARIES kazvcrypto
)
libkazv_add_tests(
store-test.cpp
EXTRA_LINK_LIBRARIES kazvstore kazvjob
)
diff --git a/src/tests/base/serialization-test.cpp b/src/tests/base/serialization-test.cpp
index e2fdbf4..6b06d85 100644
--- a/src/tests/base/serialization-test.cpp
+++ b/src/tests/base/serialization-test.cpp
@@ -1,119 +1,128 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2021 Tusooa Zhu <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <catch2/catch_all.hpp>
#include <sstream>
#include <boost/archive/text_iarchive.hpp>
#include <boost/archive/text_oarchive.hpp>
#include <serialization/immer-flex-vector.hpp>
#include <serialization/immer-map.hpp>
+#include <serialization/immer-set.hpp>
#include <serialization/immer-box.hpp>
#include <serialization/immer-array.hpp>
#include <serialization/std-optional.hpp>
#include <event.hpp>
using namespace Kazv;
using IAr = boost::archive::text_iarchive;
using OAr = boost::archive::text_oarchive;
template<class T>
static void serializeTest(const T &in, T &out)
{
std::stringstream stream;
{
auto ar = OAr(stream);
ar << in;
}
{
auto ar = IAr(stream);
ar >> out;
}
REQUIRE(in == out);
}
TEST_CASE("Serialize immer::array", "[base][serialization]")
{
immer::array<int> v{1, 2, 3, 4, 5};
immer::array<int> v2{50};
serializeTest(v, v2);
}
TEST_CASE("Serialize immer::flex_vector", "[base][serialization]")
{
immer::flex_vector<int> v{1, 2, 3, 4, 5};
immer::flex_vector<int> v2{50};
serializeTest(v, v2);
}
TEST_CASE("Serialize immer::map", "[base][serialization]")
{
auto v = immer::map<int, int>{}
.set(1, 6)
.set(2, 7)
.set(3, 8)
.set(4, 9)
.set(5, 10);
auto v2 = immer::map<int, int>{}.set(10, 7);
serializeTest(v, v2);
}
+TEST_CASE("Serialize immer::set", "[base][serialization]")
+{
+ auto v = immer::set<int>{1, 2, 3, 4, 5};
+ auto v2 = immer::set<int>{}.insert(10);
+
+ serializeTest(v, v2);
+}
+
TEST_CASE("Serialize immer::box", "[base][serialization]")
{
auto v = immer::box<int>{42};
immer::box<int> v2;
serializeTest(v, v2);
}
TEST_CASE("Serialize JsonWrap", "[base][serialization]")
{
auto v = JsonWrap{json::object({{"foo", "bar"}})};
auto v2 = JsonWrap{};
serializeTest(v, v2);
v = json::array({"mew"});
serializeTest(v, v2);
}
TEST_CASE("Serialize Event", "[base][serialization]")
{
auto v = Event{R"({
"type": "m.room.encrypted",
"content": {},
"sender": "@example:example.org",
"event_id": "!example:example.org"
})"_json};
auto v2 = Event{};
serializeTest(v, v2);
}
TEST_CASE("Serialize std::optional", "[base][serialization]")
{
std::optional<int> o{20};
std::optional<int> o2{1};
serializeTest(o, o2);
o2.reset();
serializeTest(o, o2);
o.reset();
serializeTest(o, o2);
}
diff --git a/src/tests/client/maybe-add-save-events-trigger-benchmark-test.cpp b/src/tests/client/maybe-add-save-events-trigger-benchmark-test.cpp
new file mode 100644
index 0000000..02167e0
--- /dev/null
+++ b/src/tests/client/maybe-add-save-events-trigger-benchmark-test.cpp
@@ -0,0 +1,121 @@
+/*
+ * 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 "factory.hpp"
+#include <client-model.hpp>
+#include <zug/transducer/repeat.hpp>
+#include <catch2/catch_test_macros.hpp>
+#include <catch2/benchmark/catch_benchmark.hpp>
+#include <iostream>
+
+using namespace Kazv;
+using namespace Kazv::Factory;
+
+static auto generateChanges(std::size_t newCount, std::size_t changedCount, const EventList &events)
+{
+ return intoImmer(
+ EventList{},
+ zug::repeatn(newCount, 0)
+ | zug::map([](auto &&) {
+ return makeEvent();
+ })
+ ) + intoImmer(
+ EventList{},
+ zug::take(changedCount)
+ | zug::map([](auto &&e) {
+ return makeEvent(withEventId(e.id()) | withEventContent(json{{"mew", "mew"}}));
+ }),
+ events
+ );
+}
+
+// Assumption is that new/changed messages take only a small part in the whole room.
+static std::pair<ClientModel, ClientModel> data(
+ std::size_t roomCount,
+ std::size_t timelineMessagesPerRoom,
+ std::size_t nonTimelineMessagesPerRoom,
+ std::size_t newTimelineMessagesPerRoom,
+ std::size_t newNonTimelineMessagesPerRoom,
+ std::size_t changedTimelineMessagesPerRoom,
+ std::size_t changedNonTimelineMessagesPerRoom
+)
+{
+ std::cerr << "Generating data" << std::endl;
+ auto old = makeClient();
+ auto client = makeClient();
+ std::size_t generated = 0;
+
+ for (std::size_t r = 0; r < roomCount; ++r) {
+ auto room = makeRoom();
+ auto events = EventList{};
+ for (std::size_t m = 0; m < timelineMessagesPerRoom; ++m) {
+ auto event = makeEvent();
+ events = std::move(events).push_back(event);
+ ++generated;
+ if (generated % 10000 == 0) {
+ std::cerr << "Generated " << generated << " events" << std::endl;
+ }
+ }
+ withRoomTimeline(events)(room);
+
+ auto origNonTl = EventList{};
+ for (std::size_t m = 0; m < nonTimelineMessagesPerRoom; ++m) {
+ auto event = makeEvent();
+ events = std::move(events).push_back(event);
+ ++generated;
+ if (generated % 10000 == 0) {
+ std::cerr << "Generated " << generated << " events" << std::endl;
+ }
+ }
+
+ room = RoomModel::update(std::move(room), AddMessagesAction{origNonTl});
+ withRoom(room)(old);
+
+ auto tl = generateChanges(newTimelineMessagesPerRoom, changedTimelineMessagesPerRoom, events);
+
+ auto nonTl = generateChanges(newNonTimelineMessagesPerRoom, changedNonTimelineMessagesPerRoom, origNonTl);
+
+ withRoomTimeline(tl)(room);
+ room = RoomModel::update(std::move(room), AddMessagesAction{nonTl});
+ withRoom(room)(client);
+ }
+ return {old, client};
+}
+
+TEST_CASE("maybeAddSaveEventsTrigger() benchmark: small", "[client][storage-actions][!benchmark]")
+{
+ auto [old, client] = data(10, 200, 100, 5, 5, 2, 2);
+ BENCHMARK("maybeAddSaveEventsTrigger()") {
+ return client.maybeAddSaveEventsTrigger(old);
+ };
+}
+
+TEST_CASE("maybeAddSaveEventsTrigger() benchmark: medium", "[client][storage-actions][!benchmark]")
+{
+ auto [old, client] = data(50, 1000, 200, 10, 10, 5, 5);
+ BENCHMARK("maybeAddSaveEventsTrigger()") {
+ return client.maybeAddSaveEventsTrigger(old);
+ };
+}
+
+// unprobable case: all 100 rooms have 80 events added/changed during 1 sync
+TEST_CASE("maybeAddSaveEventsTrigger() benchmark: large", "[client][storage-actions][!benchmark]")
+{
+ auto [old, client] = data(100, 4000, 400, 20, 20, 20, 20);
+ BENCHMARK("maybeAddSaveEventsTrigger()") {
+ return client.maybeAddSaveEventsTrigger(old);
+ };
+}
+
+// unprobable case: all 500 rooms have 20 events added/changed during 1 sync
+TEST_CASE("maybeAddSaveEventsTrigger() benchmark: many rooms with few events", "[client][storage-actions][!benchmark]")
+{
+ auto [old, client] = data(500, 200, 100, 5, 5, 5, 5);
+ BENCHMARK("maybeAddSaveEventsTrigger()") {
+ return client.maybeAddSaveEventsTrigger(old);
+ };
+}
diff --git a/src/tests/client/room/purge-test.cpp b/src/tests/client/room/purge-test.cpp
index f1d5297..b3df291 100644
--- a/src/tests/client/room/purge-test.cpp
+++ b/src/tests/client/room/purge-test.cpp
@@ -1,169 +1,175 @@
/*
* 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 "client-test-util.hpp"
#include "action-mock-utils.hpp"
#include "factory.hpp"
#include <client.hpp>
#include <catch2/catch_test_macros.hpp>
#include <zug/transducer/repeat.hpp>
#include <zug/transducer/enumerate.hpp>
using namespace Kazv;
using namespace Kazv::Factory;
inline static Event makeEncryptedEvent(std::string sessionId = "xxx"s)
{
return makeEvent(
withEventType("m.room.encrypted")
| withEventContent(json{
{"algorithm", "m.megolm.v1.aes-sha2"},
{"sender_key", "yyy"},
{"device_id", "device1"},
{"session_id", sessionId},
{"ciphertext", "AAA"},
})
);
}
TEST_CASE("PurgeEventsAction", "[client][room][purge]")
{
WHEN("simple purging") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
RoomModel r = makeRoom(withRoomTimeline(events));
auto next = RoomModel::update(r, PurgeEventsAction{20});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 20);
}
WHEN("nothing to purge") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
RoomModel r = makeRoom(withRoomTimeline(events));
auto next = RoomModel::update(r, PurgeEventsAction{120});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 100);
}
WHEN("rewrites undecryptedEvents") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::enumerate
| zug::map([](const auto &i, const auto &) {
return i % 2
? makeEvent()
: makeEncryptedEvent("sess" + std::to_string(i));
})
);
RoomModel r = makeRoom(withRoomEncrypted(true) | withRoomTimeline(events));
auto next = RoomModel::update(r, PurgeEventsAction{20});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 20);
REQUIRE(next.undecryptedEvents.size() == 10);
}
WHEN("keep unreadNotificationEventIds") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
RoomModel r = makeRoom(withRoomTimeline(events));
r.unreadNotificationEventIds = {
events[0].id(),
events[25].id(),
events[70].id(),
events[81].id(),
events[95].id(),
};
auto next = RoomModel::update(r, PurgeEventsAction{20});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 20);
REQUIRE(next.messages.size() == 23);
REQUIRE(next.messages.count(events[0].id()));
REQUIRE(next.messages.count(events[25].id()));
REQUIRE(next.messages.count(events[70].id()));
}
WHEN("keep localReadMarker") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
RoomModel r = makeRoom(withRoomTimeline(events));
r.localReadMarker = events[70].id();
auto next = RoomModel::update(r, PurgeEventsAction{20});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 20);
REQUIRE(next.messages.size() == 21);
REQUIRE(next.messages.count(events[70].id()));
}
WHEN("rewrites relationships") {
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
events = events.set(10, makeEvent(withEventId(events[10].id()) | withEventRelationship("moe.kazv.mxc.some-type", events[9].id())));
events = events.set(26, makeEvent(withEventId(events[26].id()) | withEventRelationship("moe.kazv.mxc.some-type", events[8].id())));
events = events.set(85, makeEvent(withEventId(events[85].id()) | withEventRelationship("moe.kazv.mxc.some-type", events[9].id())));
events = events.set(88, makeEvent(withEventId(events[88].id()) | withEventRelationship("moe.kazv.mxc.some-type", events[20].id())));
events = events.set(90, makeEvent(withEventId(events[90].id()) | withEventRelationship("moe.kazv.mxc.some-type", events[9].id())));
RoomModel r = makeRoom(withRoomTimeline(events));
r.localReadMarker = events[70].id();
auto next = RoomModel::update(r, PurgeEventsAction{20});
REQUIRE(next.checkInvariants());
REQUIRE(next.timeline.size() == 20);
REQUIRE(next.messages.size() == 21);
REQUIRE(next.messages.count(events[70].id()));
auto expected = RoomModel::ReverseEventRelationshipMap{
{events[9].id(), {{"moe.kazv.mxc.some-type", {events[85].id(), events[90].id()}}}},
{events[20].id(), {{"moe.kazv.mxc.some-type", {events[88].id()}}}},
};
REQUIRE(next.reverseEventRelationships == expected);
}
}
TEST_CASE("PurgeRoomTimelineAction", "[client][room][purge]")
{
auto events = intoImmer(
EventList{},
zug::repeatn(100, 0)
| zug::map([](const auto &) { return makeEvent(); })
);
auto r1 = makeRoom(withRoomTimeline(events));
auto r2 = makeRoom(withRoomTimeline(events));
auto r3 = makeRoom(withRoomTimeline(events));
auto m = makeClient(withRoom(r1) | withRoom(r2) | withRoom(r3));
auto u = makeMockSdkUtil(m);
auto d = u.getMockDispatcher(passDown<PurgeRoomTimelineAction>());
auto c = u.getClient(d);
+ bool saveEventsRequested = false;
+ auto watchable = u.ee.watchable();
+ watchable.after<SaveEventsRequested>([&saveEventsRequested](auto &&) {
+ saveEventsRequested = true;
+ });
c.purgeRoomEvents({
{r1.roomId, 80},
{r2.roomId, 50},
}).then([&u](const auto &stat) {
REQUIRE(stat);
u.io.stop();
});
u.io.run();
REQUIRE(c.room(r1.roomId).timelineEventIds().get().size() == 80);
REQUIRE(c.room(r2.roomId).timelineEventIds().get().size() == 50);
REQUIRE(c.room(r3.roomId).timelineEventIds().get().size() == 100);
+ REQUIRE(!saveEventsRequested);
}
diff --git a/src/tests/client/storage-actions-test.cpp b/src/tests/client/storage-actions-test.cpp
index 7bb0b21..c3e38d6 100644
--- a/src/tests/client/storage-actions-test.cpp
+++ b/src/tests/client/storage-actions-test.cpp
@@ -1,78 +1,196 @@
/*
* 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 "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) {
+ 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("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 b4740f2..ee344b3 100644
--- a/src/tests/client/sync-test.cpp
+++ b/src/tests/client/sync-test.cpp
@@ -1,745 +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) {
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);
+}

File Metadata

Mime Type
text/x-diff
Expires
Thu, Oct 8, 6:42 PM (2 h, 22 m)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1784199
Default Alt Text
(142 KB)

Event Timeline