Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85802723
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
142 KB
Referenced Files
None
Subscribers
None
View Options
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
Details
Attached
Mime Type
text/x-diff
Expires
Thu, Oct 8, 6:42 PM (3 h, 8 m)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1784199
Default Alt Text
(142 KB)
Attached To
Mode
rL libkazv
Attached
Detach File
Event Timeline
Log In to Comment