Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85629776
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/client/room/room-model.cpp b/src/client/room/room-model.cpp
index 4d6ca72..fc3b25a 100644
--- a/src/client/room/room-model.cpp
+++ b/src/client/room/room-model.cpp
@@ -1,764 +1,771 @@
/*
* 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, /* 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);
}
+
+ bool RoomModel::isReadBy(const std::string &eventId, const std::string &userId) const
+ {
+ auto readReceiptAt = messages[readReceipts[userId].eventId];
+ auto event = messages[eventId];
+ return sortKeyForTimelineEvent(readReceiptAt) >= sortKeyForTimelineEvent(event);
+ }
}
diff --git a/src/client/room/room-model.hpp b/src/client/room/room-model.hpp
index b9b0c2f..8e14717 100644
--- a/src/client/room/room-model.hpp
+++ b/src/client/room/room-model.hpp
@@ -1,506 +1,516 @@
/*
* 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;
+
+ /**
+ * Check if an event has been read by the given user
+ *
+ * @param eventId eventId of the event to check
+ * @param userId userId of the given user
+ * @return true iff the event is read by the user
+ */
+ bool isReadBy(const std::string &eventId, const std::string &userId) 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, 9)
BOOST_CLASS_VERSION(Kazv::RoomListModel, 0)
diff --git a/src/client/room/room.cpp b/src/client/room/room.cpp
index 2aef687..ef2d02d 100644
--- a/src/client/room/room.cpp
+++ b/src/client/room/room.cpp
@@ -1,947 +1,952 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2021-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <zug/into_vector.hpp>
#include <immer/algorithm.hpp>
#include <debug.hpp>
#include "room.hpp"
namespace Kazv
{
Room::Room(lager::reader<SdkModel> sdk,
lager::reader<std::string> roomId,
Context<ClientAction> ctx,
DepsT deps)
: m_sdk(sdk)
, m_roomId(roomId)
, m_ctx(ctx)
, m_deps(deps)
, m_roomCursor(makeRoomCursor())
{
}
Room::Room(lager::reader<SdkModel> sdk,
lager::reader<std::string> roomId,
Context<ClientAction> ctx)
: m_sdk(sdk)
, m_roomId(roomId)
, m_ctx(ctx)
, m_deps(std::nullopt)
, m_roomCursor(makeRoomCursor())
{
}
Room::Room(InEventLoopTag, std::string roomId, ContextT ctx, DepsT deps)
: m_sdk(std::nullopt)
, m_roomId(roomId)
, m_ctx(ctx)
, m_deps(deps)
, m_roomCursor(makeRoomCursor())
#ifdef KAZV_USE_THREAD_SAFETY_HELPER
, KAZV_ON_EVENT_LOOP_VAR(true)
#endif
{
}
Room Room::toEventLoop() const
{
assert(m_deps.has_value());
return Room(InEventLoopTag{}, currentRoomId(), m_ctx, m_deps.value());
}
const lager::reader<SdkModel> &Room::sdkCursor() const
{
if (m_sdk.has_value()) { return m_sdk.value(); }
assert(m_deps.has_value());
return *lager::get<SdkModelCursorKey>(m_deps.value());
}
const lager::reader<RoomModel> &Room::roomCursor() const
{
return m_roomCursor;
}
lager::reader<RoomModel> Room::makeRoomCursor() const
{
return lager::match(m_roomId)(
[&](const lager::reader<std::string> &roomId) -> lager::reader<RoomModel> {
return lager::with(sdkCursor().map(&SdkModel::c)[&ClientModel::roomList], roomId)
.map([](auto rooms, auto id) {
return rooms[id];
}).make();
},
[&](const std::string &roomId) -> lager::reader<RoomModel> {
return sdkCursor().map(&SdkModel::c)[&ClientModel::roomList].map([roomId](auto rooms) { return rooms[roomId]; }).make();
});
}
std::string Room::currentRoomId() const
{
return lager::match(m_roomId)(
[&](const lager::reader<std::string> &roomId) {
return roomId.get();
},
[&](const std::string &roomId) {
return roomId;
});
}
auto Room::message(lager::reader<std::string> eventId) const -> lager::reader<Event>
{
return lager::with(roomCursor()[&RoomModel::messages], eventId)
.map([](const auto &msgs, const auto &id) {
return msgs[id];
});
}
auto Room::localEcho(lager::reader<std::string> txnId) const -> lager::reader<LocalEchoDesc>
{
return lager::with(roomCursor()[&RoomModel::localEchoes], txnId)
.map([](const auto localEchoes, const auto &id) {
auto localEchoIt = std::find_if(
localEchoes.begin(), localEchoes.end(),
[id](const auto localEcho) {
return localEcho.txnId == id;
});
if (localEchoIt == localEchoes.end()) {
return LocalEchoDesc {};
}
return *localEchoIt;
});
}
auto Room::members() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](auto room) {
return room.joinedMemberIds();
});
}
auto Room::invitedMembers() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](auto room) {
return room.invitedMemberIds();
});
}
auto Room::knockedMembers() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](auto room) {
return room.knockedMemberIds();
});
}
auto Room::leftMembers() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](auto room) {
return room.leftMemberIds();
});
}
auto Room::bannedMembers() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](auto room) {
return room.bannedMemberIds();
});
}
auto Room::joinedMemberEvents() const -> lager::reader<EventList>
{
return roomCursor().map([](auto room) {
return room.joinedMemberEvents();
});
}
auto Room::invitedMemberEvents() const -> lager::reader<EventList>
{
return roomCursor().map([](auto room) {
return room.invitedMemberEvents();
});
}
auto Room::knockedMemberEvents() const -> lager::reader<EventList>
{
return roomCursor().map([](auto room) {
return room.knockedMemberEvents();
});
}
auto Room::leftMemberEvents() const -> lager::reader<EventList>
{
return roomCursor().map([](auto room) {
return room.leftMemberEvents();
});
}
auto Room::bannedMemberEvents() const -> lager::reader<EventList>
{
return roomCursor().map([](auto room) {
return room.bannedMemberEvents();
});
}
auto Room::memberEventByCursor(lager::reader<std::string> userId) const -> lager::reader<Event>
{
return inviteStateOrStateEvent(userId.map([](auto id) {
return KeyOfState{"m.room.member", id};
}));
}
auto Room::memberEventFor(std::string userId) const -> lager::reader<Event>
{
return memberEventByCursor(lager::make_constant(userId));
}
lager::reader<immer::map<KeyOfState, Event>> Room::inviteStateOrState() const
{
return lager::with(stateEvents(), inviteState(), membership())
.map([](const auto &stateEv, const auto &inviteSt, const auto &mem) {
if (mem == RoomMembership::Invite) {
return inviteSt;
} else {
return stateEv;
}
});
}
lager::reader<Event> Room::inviteStateOrStateEvent(lager::reader<KeyOfState> key) const
{
return lager::with(stateEvents(), inviteState(), membership(), key)
.map([](const auto &stateEv, const auto &inviteSt, const auto &mem, const auto &k) {
if (mem == RoomMembership::Invite) {
auto maybePtr = inviteSt.find(k);
return maybePtr ? *maybePtr : stateEv[k];
} else {
return stateEv[k];
}
});
}
auto Room::inviteState() const -> lager::reader<immer::map<KeyOfState, Event>>
{
return roomCursor()[&RoomModel::inviteState];
}
auto Room::heroMemberEvents() const -> lager::reader<immer::flex_vector<Event>>
{
using namespace lager::lenses;
return roomCursor().map(&RoomModel::heroMemberEvents);
}
auto Room::heroDisplayNames() const
-> lager::reader<immer::flex_vector<std::string>>
{
return heroMemberEvents()
.xform(zug::map([](const auto &events) {
return intoImmer(
immer::flex_vector<std::string>{},
zug::map([](const auto &event) {
auto content = event.content();
return content.get().contains("displayname")
? content.get()["displayname"].template get<std::string>()
: std::string();
}),
events);
}));
}
auto Room::nameOpt() const -> lager::reader<std::optional<std::string>>
{
using namespace lager::lenses;
return inviteStateOrState()
[KeyOfState{"m.room.name", ""}]
[or_default]
.xform(eventContent)
.xform(zug::map([](const JsonWrap &content) {
return content.get().contains("name")
? std::optional<std::string>(content.get()["name"])
: std::nullopt;
}));
}
auto Room::name() const -> lager::reader<std::string>
{
using namespace lager::lenses;
return nameOpt()[value_or("<no name>")];
}
lager::reader<bool> Room::encrypted() const
{
return roomCursor()[&RoomModel::encrypted];
}
auto Room::localReadMarker() const
-> lager::reader<std::string>
{
return roomCursor()[&RoomModel::localReadMarker];
}
auto Room::heroIds() const
-> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor().map([](const auto &room) {
return room.heroIds;
});
}
auto Room::joinedMemberCount() const -> lager::reader<std::size_t>
{
return roomCursor()[&RoomModel::joinedMemberCount];
}
auto Room::invitedMemberCount() const -> lager::reader<std::size_t>
{
return roomCursor()[&RoomModel::invitedMemberCount];
}
auto Room::avatarMxcUri() const -> lager::reader<std::string>
{
using namespace lager::lenses;
return inviteStateOrState()
[KeyOfState{"m.room.avatar", ""}]
[or_default]
.map([](Event ev) {
auto content = ev.content().get();
return content.contains("url")
? std::string(content["url"])
: "";
});
}
auto Room::setLocalDraft(std::string localDraft) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(UpdateRoomAction{+roomId(), SetLocalDraftAction{localDraft}});
}
auto Room::sendMessage(Event msg) const
-> PromiseT
{
using namespace CursorOp;
auto hasCrypto = ~sdkCursor().map([](const auto &sdk) -> bool {
return sdk.c().crypto.has_value();
});
auto roomEncrypted = ~roomCursor()[&RoomModel::encrypted];
auto noFullMembers = ~roomCursor()[&RoomModel::membersFullyLoaded]
.map([](auto b) { return !b; });
auto rid = +roomId();
// Don't use m_ctx directly in the callbacks
// as `this` may have been destroyed when
// the callbacks are called.
auto ctx = m_ctx;
auto promise = ctx.createResolvedPromise(true);
// If we do not need encryption just send it as-is
if (! +allCursors(hasCrypto, roomEncrypted)) {
return promise
.then([ctx, rid, msg](auto succ) {
if (! succ) {
return ctx.createResolvedPromise(false);
}
return ctx.dispatch(SendMessageAction{rid, msg});
});
}
if (! m_deps) {
return ctx.createResolvedPromise({false, json{{"error", "missing-deps"}}});
}
auto deps = m_deps.value();
promise = promise.then([ctx, rid, msg](auto) {
return ctx.dispatch(SaveLocalEchoAction{rid, msg});
});
auto saveLocalEchoPromise = promise;
// If the room member list is not complete, load it fully first.
if (+allCursors(hasCrypto, roomEncrypted, noFullMembers)) {
kzo.client.dbg() << "The members of " << rid
<< " are not fully loaded." << std::endl;
promise = promise
.then([ctx, rid](auto) { // SaveLocalEchoAction can't fail
return ctx.dispatch(GetRoomStatesAction{rid});
})
.then([ctx, rid](auto succ) {
if (! succ) {
kzo.client.warn() << "Loading members of " << rid
<< " failed." << std::endl;
return ctx.createResolvedPromise(false);
} else {
// XXX remove the hard-coded initialSync parameter
return ctx.dispatch(QueryKeysAction{true})
.then([](auto succ) {
if (! succ) {
kzo.client.warn() << "Query keys failed" << std::endl;
}
return succ;
});
}
});
}
auto encryptEvent = [ctx, rid, msg, deps](auto &&status) {
if (! status) { return ctx.createResolvedPromise(status); }
kzo.client.dbg() << "encrypting megolm" << std::endl;
auto &rg = lager::get<RandomInterface &>(deps);
return ctx.dispatch(EncryptMegOlmEventAction{
rid,
msg,
currentTimeMs(),
rg.generateRange<RandomData>(EncryptMegOlmEventAction::maxRandomSize())
});
};
auto saveEncryptedLocalEcho = [saveLocalEchoPromise, msg, rid, ctx](auto &&status) mutable {
if (! status) { return ctx.createResolvedPromise(status); }
return saveLocalEchoPromise
.then([status, msg, rid, ctx](auto &&st) {
auto txnId = st.dataStr("txnId");
kzo.client.dbg() << "saving encrypted local echo with txn id " << txnId << std::endl;
auto encrypted = status.dataJson("encrypted");
auto decrypted = msg.originalJson();
auto event = Event(encrypted).setDecryptedJson(decrypted, Event::Decrypted);
return ctx.dispatch(SaveLocalEchoAction{rid, event, txnId});
})
.then([status](auto &&) {
return status;
});
};
auto maybeSendKeys = [ctx, rid, r=toEventLoop(), deps](auto status) {
if (! status) { return ctx.createResolvedPromise(status); }
auto encryptedEvent = status.dataJson("encrypted");
auto content = encryptedEvent.at("content");
auto ret = ctx.createResolvedPromise({});
if (status.data().get().contains("key")) {
kzo.client.dbg() << "megolm session rotated, sending session key" << std::endl;
auto key = status.dataStr("key");
ret = ret
.then([rid, r, key, ctx, deps, sessionId=content.at("session_id")](auto &&) {
auto members = (+r.roomCursor()).joinedMemberIds();
auto client = (+r.sdkCursor()).c();
using DeviceMapT = immer::map<std::string, immer::flex_vector<std::string>>;
auto devicesToSend = accumulate(
members, DeviceMapT{}, [client](auto map, auto uid) {
return std::move(map)
.set(uid, client.devicesToSendKeys(uid));
});
auto &rg = lager::get<RandomInterface &>(deps);
return ctx.dispatch(ClaimKeysAction{
rid, sessionId, key, devicesToSend,
rg.generateRange<RandomData>(ClaimKeysAction::randomSize(devicesToSend))
})
.then([rid, ctx, deps, devicesToSend](auto status) {
if (! status) { return ctx.createResolvedPromise({}); }
kzo.client.dbg() << "olm-encrypting key event" << std::endl;
auto keyEv = status.dataJson("keyEvent");
auto &rg = lager::get<RandomInterface &>(deps);
return ctx.dispatch(PrepareForSharingRoomKeyAction{
rid,
devicesToSend, keyEv,
rg.generateRange<RandomData>(PrepareForSharingRoomKeyAction::randomSize(devicesToSend))
});
})
.then([ctx, r, devicesToSend, rid](auto status) {
if (! status) { return ctx.createResolvedPromise({}); }
auto txnId = status.dataStr("txnId");
return r.sendPendingKeyEvent(txnId);
});
});
}
return ret
.then([ctx, prevStatus=status](auto status) {
if (! status) { return status; }
return prevStatus;
});
};
auto sendEncryptedEvent = [ctx, rid, saveLocalEchoPromise, msg](auto status) mutable {
if (! status) { return ctx.createResolvedPromise(status); }
return saveLocalEchoPromise.then([ctx, rid, status, msg](auto st) {
auto txnId = st.dataStr("txnId");
kzo.client.dbg() << "sending encrypted message with txn id " << txnId << std::endl;
auto encrypted = status.dataJson("encrypted");
auto decrypted = msg.originalJson();
auto event = Event(encrypted).setDecryptedJson(decrypted, Event::Decrypted);
return ctx.dispatch(SendMessageAction{rid, event, txnId});
});
};
auto maybeSetStatusToFailed = [ctx, saveLocalEchoPromise, rid](const auto &st) mutable {
if (st.success()) {
return ctx.createResolvedPromise(st);
}
return saveLocalEchoPromise
.then([rid, ctx](const auto &saveLocalEchoStatus) {
auto txnId = saveLocalEchoStatus.dataStr("txnId");
return ctx.dispatch(UpdateLocalEchoStatusAction{rid, txnId, LocalEchoDesc::Failed});
})
.then([st](const auto &) { // can't fail
return st;
});
};
return promise
.then(encryptEvent)
.then(saveEncryptedLocalEcho)
.then(maybeSendKeys)
.then(sendEncryptedEvent)
.then(maybeSetStatusToFailed);
}
auto Room::sendTextMessage(std::string text) const
-> PromiseT
{
json j{
{"type", "m.room.message"},
{"content", {
{"msgtype", "m.text"},
{"body", text}
}
}
};
Event e{j};
return sendMessage(e);
}
auto Room::resendMessage(std::string txnId) const
-> PromiseT
{
return m_ctx.createResolvedPromise({})
.then([that=toEventLoop()](const auto &) {
kzo.client.dbg() << "resending all pending key events" << std::endl;
return that.sendAllPendingKeyEvents();
})
.then([that=toEventLoop(), txnId](const auto &sendKeyStatus) {
if (!sendKeyStatus.success()) {
kzo.client.warn() << "resending all pending key events failed" << std::endl;
return that.m_ctx.createResolvedPromise(sendKeyStatus);
}
auto maybeLocalEcho = that.roomCursor().map([txnId](const auto &room) {
return room.getLocalEchoByTxnId(txnId);
}).make();
auto roomEncrypted = that.roomCursor()[&RoomModel::encrypted].make();
if (!maybeLocalEcho.get()) {
return that.m_ctx.createResolvedPromise({false, json{
{"errorCode", "MOE_KAZV_MXC_NO_SUCH_TXNID"},
{"error", "No such txn id"}
}});
}
auto localEcho = maybeLocalEcho.get().value();
auto event = localEcho.event;
auto rid = that.roomId().make().get();
if ((roomEncrypted.get() && event.encrypted()) || !roomEncrypted.get()) {
kzo.client.info() << "resending message with txn id " << localEcho.txnId << std::endl;
return that.m_ctx.dispatch(SendMessageAction{rid, event, localEcho.txnId});
} else {
kzo.client.info() << "room encrypted but event isn't, just resend everything" << std::endl;
return that.m_ctx.dispatch(RoomListAction{UpdateRoomAction{rid, RemoveLocalEchoAction{localEcho.txnId}}})
.then([localEcho, that](auto) { // Can't fail
return that.sendMessage(localEcho.event);
});
}
});
}
auto Room::redactEvent(std::string eventId, std::optional<std::string> reason) const -> PromiseT
{
return m_ctx.dispatch(RedactEventAction{
roomId().make().get(),
eventId,
reason
});
}
auto Room::sendPendingKeyEvent(std::string txnId) const -> PromiseT
{
return m_ctx.createResolvedPromise({})
.then([txnId, r=toEventLoop()](auto &&) {
auto ctx = r.m_ctx;
auto rid = r.roomId().make().get();
kzo.client.dbg() << "sending key event as to-device message" << std::endl;
kzo.client.dbg() << "txnId of key event: " << txnId << std::endl;
auto maybePending = r.roomCursor().make().get().getPendingRoomKeyEventByTxnId(txnId);
if (!maybePending.has_value()) {
kzo.client.warn() << "No such pending room key event";
return ctx.createResolvedPromise(EffectStatus(/* succ = */ false, json::object(
{{"errorCode", "MOE_KAZV_MXC_NO_SUCH_PENDING_ROOM_KEY_EVENT"},
{"error", "No such pending room key event"}}
)));
}
auto pending = maybePending.value();
return ctx.dispatch(SendMultipleToDeviceMessagesAction{pending.messages})
.then([ctx, rid, txnId](const auto &sendToDeviceStatus) {
if (!sendToDeviceStatus.success()) {
return ctx.createResolvedPromise(sendToDeviceStatus);
} else {
return ctx.dispatch(UpdateRoomAction{rid, RemovePendingRoomKeyAction{txnId}});
}
});
});
}
auto Room::sendAllPendingKeyEvents() const -> PromiseT
{
return m_ctx.createResolvedPromise({})
.then([r=toEventLoop()](auto &&) {
auto pendingEvents = r.pendingRoomKeyEvents().make().get();
if (pendingEvents.empty()) {
return r.m_ctx.createResolvedPromise(EffectStatus(/* succ = */ true));
} else {
auto txnId = pendingEvents[0].txnId;
return r.sendPendingKeyEvent(txnId)
.then([r, txnId](const auto &stat) {
if (!stat.success()) {
kzo.client.warn() << "Can't send pending key event of txnId " << txnId << std::endl;
return r.m_ctx.createResolvedPromise(stat);
}
return r.sendAllPendingKeyEvents();
});
}
});
}
auto Room::refreshRoomState() const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(GetRoomStatesAction{+roomId()});
}
auto Room::getStateEvent(std::string type, std::string stateKey) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(GetStateEventAction{+roomId(), type, stateKey});
}
auto Room::sendStateEvent(Event state) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(SendStateEventAction{+roomId(), state});
}
auto Room::setName(std::string name) const
-> PromiseT
{
json j{
{"type", "m.room.name"},
{"content", {
{"name", name}
}
}
};
Event e{j};
return sendStateEvent(e);
}
auto Room::setTopic(std::string topic) const
-> PromiseT
{
json j{
{"type", "m.room.topic"},
{"content", {
{"topic", topic}
}
}
};
Event e{j};
return sendStateEvent(e);
}
auto Room::invite(std::string userId) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(InviteToRoomAction{+roomId(), userId});
}
auto Room::typingUsers() const -> lager::reader<immer::flex_vector<std::string>>
{
using namespace lager::lenses;
return ephemeral("m.typing")
.xform(eventContent
| jsonAtOr("user_ids", immer::flex_vector<std::string>{}));
}
auto Room::typingMemberEvents() const -> lager::reader<EventList>
{
return lager::with(typingUsers(), roomCursor()[&RoomModel::stateEvents])
.map([](const auto &userIds, const auto &events) {
return intoImmer(EventList{}, zug::map([events](const auto &id) {
return events[KeyOfState{"m.room.member", id}];
}), userIds);
});
}
auto Room::setTyping(bool typing, std::optional<int> timeoutMs) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(SetTypingAction{+roomId(), typing, timeoutMs});
}
auto Room::leave() const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(LeaveRoomAction{+roomId()});
}
auto Room::forget() const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(ForgetRoomAction{+roomId()});
}
auto Room::kick(std::string userId, std::optional<std::string> reason) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(KickAction{+roomId(), userId, reason});
}
auto Room::ban(std::string userId, std::optional<std::string> reason) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(BanAction{+roomId(), userId, reason});
}
auto Room::unban(std::string userId/*, std::optional<std::string> reason*/) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(UnbanAction{+roomId(), userId});
}
auto Room::setAccountData(Event accountDataEvent) const -> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(SetAccountDataPerRoomAction{+roomId(), accountDataEvent});
}
auto Room::tags() const -> lager::reader<immer::map<std::string, double>>
{
return roomCursor().map(&RoomModel::tags);
}
auto Room::addOrSetTag(std::string tagId, std::optional<double> order) const -> PromiseT
{
using namespace CursorOp;
return m_ctx.createResolvedPromise({})
.then([that=toEventLoop(), tagId, order](auto) {
return that.setAccountData((+that.roomCursor()).makeAddTagEvent(tagId, order));
});
}
auto Room::removeTag(std::string tagId) const -> PromiseT
{
using namespace CursorOp;
return m_ctx.createResolvedPromise({})
.then([that=toEventLoop(), tagId](auto) {
return that.setAccountData((+that.roomCursor()).makeRemoveTagEvent(tagId));
});
}
auto Room::pinnedEvents() const -> lager::reader<immer::flex_vector<std::string>>
{
return state(KeyOfState{"m.room.pinned_events", ""})
.xform(eventContent | jsonAtOr("pinned", immer::flex_vector<std::string>{}));
}
auto Room::setPinnedEvents(immer::flex_vector<std::string> eventIds) const
-> PromiseT
{
json j{
{"type", "m.room.pinned_events"},
{"content", {
{"pinned", eventIds}
}
}
};
Event e{j};
return sendStateEvent(e);
}
auto Room::pinEvents(immer::flex_vector<std::string> eventIds) const -> PromiseT
{
return m_ctx.createResolvedPromise({})
.then([room=toEventLoop(), eventIds](auto) {
auto curPinnedEvents = room.pinnedEvents().get();
auto newPinnedEvents = intoImmer(
curPinnedEvents,
zug::filter([curPinnedEvents](const auto &eventId) {
return std::find(curPinnedEvents.begin(), curPinnedEvents.end(), eventId) == curPinnedEvents.end();
}),
eventIds
);
return room.setPinnedEvents(newPinnedEvents);
});
}
auto Room::unpinEvents(immer::flex_vector<std::string> eventIds) const -> PromiseT
{
return m_ctx.createResolvedPromise({})
.then([room=toEventLoop(), eventIds](auto) {
auto curPinnedEvents = room.pinnedEvents().get();
auto newPinnedEvents = intoImmer(
immer::flex_vector<std::string>(),
zug::filter([eventIds](const auto &eventId) {
return std::find(eventIds.begin(), eventIds.end(), eventId) == eventIds.end();
}),
curPinnedEvents
);
return room.setPinnedEvents(newPinnedEvents);
});
}
auto Room::timelineEventIds() const -> lager::reader<immer::flex_vector<std::string>>
{
return roomCursor()[&RoomModel::timeline];
}
auto Room::messagesMap() const -> lager::reader<immer::map<std::string, Event>>
{
return roomCursor()[&RoomModel::messages];
}
auto Room::timelineEvents() const -> lager::reader<immer::flex_vector<Event>>
{
return roomCursor()
.xform(zug::map([](auto r) {
auto messages = r.messages;
auto timeline = r.timeline;
return intoImmer(
immer::flex_vector<Event>{},
zug::map([=](auto eventId) {
return messages[eventId];
}),
timeline);
}));
}
auto Room::timelineGaps() const
-> lager::reader<immer::map<std::string /* eventId */,
std::string /* prevBatch */>>
{
return roomCursor()[&RoomModel::timelineGaps];
}
auto Room::paginateBackFromEvent(std::string eventId) const
-> PromiseT
{
using namespace CursorOp;
return m_ctx.dispatch(PaginateTimelineAction{
+roomId(), eventId, std::nullopt});
}
auto Room::localEchoes() const -> lager::reader<immer::flex_vector<LocalEchoDesc>>
{
return roomCursor()[&RoomModel::localEchoes];
}
auto Room::removeLocalEcho(std::string txnId) const -> PromiseT
{
return m_ctx.dispatch(RoomListAction{UpdateRoomAction{
currentRoomId(),
RemoveLocalEchoAction{txnId},
}});
}
auto Room::pendingRoomKeyEvents() const -> lager::reader<immer::flex_vector<PendingRoomKeyEvent>>
{
return roomCursor()[&RoomModel::pendingRoomKeyEvents];
}
auto Room::powerLevels() const -> lager::reader<PowerLevelsDesc>
{
return state({"m.room.power_levels", ""})
.map([](const Event &event) {
return PowerLevelsDesc(event);
});
}
auto Room::relatedEvents(lager::reader<std::string> eventId, std::string relType) const -> lager::reader<EventList>
{
return lager::with(
eventId,
roomCursor()[&RoomModel::reverseEventRelationships],
roomCursor()[&RoomModel::messages]
).map([relType](const auto &eid, const auto &rels, const auto &msgs) {
auto relatedEvents = zug::into_vector(
zug::map([msgs](const auto &id) {
return msgs[id];
}),
rels[eid][relType]
);
std::sort(relatedEvents.begin(), relatedEvents.end(), [](const Event &a, const Event &b) {
return std::make_tuple(a.originServerTs(), a.id()) < std::make_tuple(b.originServerTs(), b.id());
});
return EventList(relatedEvents.begin(), relatedEvents.end());
}).make();
}
auto Room::eventReaders(lager::reader<std::string> eventId) const -> lager::reader<immer::flex_vector<EventReader>>
{
return lager::with(
roomCursor()[&RoomModel::readReceipts],
roomCursor()[&RoomModel::eventReadUsers],
eventId
).map([](const auto &receipts, const auto &readUsers, const auto &eid) {
auto userIds = readUsers[eid];
return intoImmer(
immer::flex_vector<EventReader>{},
zug::map([receipts](const auto &userId) {
return EventReader{userId, receipts[userId].timestamp};
}),
userIds
);
});
}
auto Room::postReceipt(std::string eventId) const -> PromiseT
{
return m_ctx.dispatch(PostReceiptAction{roomId().make().get(), eventId});
}
auto Room::unreadNotificationEventIds() const -> lager::reader<immer::flex_vector<std::string>>
{
return m_roomCursor[&RoomModel::unreadNotificationEventIds];
}
+
+ bool Room::isReadBy(const std::string &eventId, const std::string &userId)
+ {
+ return m_roomCursor.get().isReadBy(eventId, userId);
+ }
}
diff --git a/src/client/room/room.hpp b/src/client/room/room.hpp
index b7fd521..d6bd3b7 100644
--- a/src/client/room/room.hpp
+++ b/src/client/room/room.hpp
@@ -1,834 +1,843 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2020-2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <lager/reader.hpp>
#include <lager/context.hpp>
#include <lager/with.hpp>
#include <lager/constant.hpp>
#include <lager/lenses/optional.hpp>
#include <zug/transducer/map.hpp>
#include <zug/transducer/filter.hpp>
#include <zug/sequence.hpp>
#include <immer/flex_vector_transient.hpp>
#include "debug.hpp"
#include "sdk-model.hpp"
#include "client-model.hpp"
#include "room-model.hpp"
#include <cursorutil.hpp>
#include "sdk-model-cursor-tag.hpp"
#include "random-generator.hpp"
#include "power-levels-desc.hpp"
namespace Kazv
{
/**
* Represent a Matrix room.
*
* This class has the same constraints as Client.
*/
class Room
{
public:
using PromiseT = SingleTypePromise<DefaultRetType>;
using DepsT = lager::deps<SdkModelCursorKey, RandomInterface &
#ifdef KAZV_USE_THREAD_SAFETY_HELPER
, EventLoopThreadIdKeeper &
#endif
>;
using ContextT = Context<ClientAction>;
struct InEventLoopTag {};
/**
* Constructor.
*
* Construct the room with @c roomId .
*
* `sdk` and `roomId` must be cursors in the same thread.
*
* The constructed room will be in the same thread as `sdk` and `roomId`.
*
* @warning Do not use this directly. Use `Client::room()` and
* `Client::roomBycursor()` instead.
*/
Room(lager::reader<SdkModel> sdk,
lager::reader<std::string> roomId,
ContextT ctx);
/**
* Constructor.
*
* Construct the room with @c roomId and with Deps support.
*
* `sdk` and `roomId` must be cursors in the same thread.
*
* The constructed room will be in the same thread as `sdk` and `roomId`.
*
* @warning Do not use this directly. Use `Client::room()` and
* `Client::roomBycursor()` instead.
*/
Room(lager::reader<SdkModel> sdk,
lager::reader<std::string> roomId,
ContextT ctx, DepsT deps);
/**
* Construct a Room in the same thread as the event loop.
*
* The constructed Room is not constructed from a cursor,
* and thus copying-constructing from that is thread-safe as long as each thread
* calls with different objects.
*
* this must have Deps support.
*
* @warning Do not use this directly. Use `Client::room()` and
* `Client::roomBycursor()` instead.
*/
Room(InEventLoopTag, std::string roomId, ContextT ctx, DepsT deps);
/**
* Return a Room that represents the room *currently represented* by this,
* but suitable for use in the event loop of the context.
*
* This function can only be called from the thread where this belongs.
*
* Example:
*
* ```
* auto ctx = sdk.context();
* auto client = sdk.clientFromSecondaryRoot(sr);
* auto room = client.room("!room-id:domain.name");
* room.sendTextMessage("test")
* .then([r=room.toEventLoop(), ctx](auto &&st) {
* if (!st) {
* std::cerr << "Cannot send message" << std::endl;
* return ctx.createResolvedPromise(st);
* }
* return r.sendTextMessage("follow-up");
* });
* ```
*
* @sa Sdk::clientFromSecondaryRoot , Client::room
*/
Room toEventLoop() const;
/* lager::reader<MapT<KeyOfState, Event>> */
inline auto stateEvents() const {
return roomCursor()
[&RoomModel::stateEvents];
}
/**
* Get the invite_state of this room.
*
* @return A lager::reader containing a map from KeyOfState to the event.
*/
auto inviteState() const -> lager::reader<immer::map<KeyOfState, Event>>;
/* lager::reader<std::optional<Event>> */
inline auto stateOpt(KeyOfState k) const {
return stateEvents()
[std::move(k)];
}
/* lager::reader<Event> */
inline auto state(KeyOfState k) const {
return stateOpt(k)
[lager::lenses::or_default];
}
/**
* Get the timeline event ids of this room in ascending timestamp order.
* It takes constant time for the cursor to be updated.
*
* @return A lager::reader of a RangeT of event ids.
*/
auto timelineEventIds() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get a map from event ids to events. It takes constant time for the
* cursor to be updated.
*
* @return A lager::reader of a map for all timeline events.
*/
auto messagesMap() const -> lager::reader<immer::map<std::string, Event>>;
/**
* Get a list of timeline events in this room. It takes O(timeline.size())
* time for the cursor to be updated.
*
* @return A lager::reader<RangeT<Event>> of timeline events from
* the oldest to the latest.
*/
auto timelineEvents() const -> lager::reader<immer::flex_vector<Event>>;
/**
* Get one message by the event id.
*
* @param A lager::reader of the event id.
*
* @return A lager::reader of the event.
*/
auto message(lager::reader<std::string> eventId) const -> lager::reader<Event>;
/**
* Get one local echo by the txnId.
*
* @param A lager::reader of the txnId.
*
* @return A lager::reader of the local echo.
*/
auto localEcho(lager::reader<std::string> txnId) const -> lager::reader<LocalEchoDesc>;
/**
* Get the member events of heroes this room.
*
* @return a lager::reader of a RangeT of Event containing the member events.
*/
auto heroMemberEvents() const -> lager::reader<immer::flex_vector<Event>>;
/**
* Get the member events of heroes this room.
*
* @return a lager::reader of a RangeT of std::string containing the member events.
*/
auto heroDisplayNames() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the name of this room.
*
* If there is a m.room.name state event, the name in it is used.
* If there is none, the returned cursor will hold std::nullopt.
*
* @return a lager::reader of an optional std::string.
*/
auto nameOpt() const -> lager::reader<std::optional<std::string>>;
/**
* Get the name of this room.
*
* If there is a m.room.name state event, the name in it is used.
* If there is none, the returned cursor will hold a placeholder string.
*
* @return a lager::reader of an std::string.
*/
auto name() const -> lager::reader<std::string>;
/**
* Get the avatar mxc uri of this room.
*
* @return A lager::reader containing the mxc uri of this room.
*/
auto avatarMxcUri() const -> lager::reader<std::string>;
/**
* Get the list of joined member ids.
*
* @return A lager::reader containing an RangeT of the joined members' ids.
*/
auto members() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the list of invited member ids.
*
* @return A lager::reader containing an RangeT of the invited members' ids.
*/
auto invitedMembers() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the list of knocked member ids.
*
* @return A lager::reader containing an RangeT of the knocked members' ids.
*/
auto knockedMembers() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the list of left member ids.
*
* @return A lager::reader containing an RangeT of the left members' ids.
*/
auto leftMembers() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the list of banned member ids.
*
* @return A lager::reader containing an RangeT of the banned members' id.
*/
auto bannedMembers() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the list of joined member events.
*
* @return A lager::reader containing an EventList of the joined members' events.
*/
auto joinedMemberEvents() const -> lager::reader<EventList>;
/**
* Get the list of invited member events.
*
* @return A lager::reader containing an EventList of the invited members' events.
*/
auto invitedMemberEvents() const -> lager::reader<EventList>;
/**
* Get the list of knocked member events.
*
* @return A lager::reader containing an EventList of the knocked members' events.
*/
auto knockedMemberEvents() const -> lager::reader<EventList>;
/**
* Get the list of left member events.
*
* @return A lager::reader containing an EventList of the left members' events.
*/
auto leftMemberEvents() const -> lager::reader<EventList>;
/**
* Get the list of banned member events.
*
* @return A lager::reader containing an EventList of the banned members' events.
*/
auto bannedMemberEvents() const -> lager::reader<EventList>;
/**
* Get the member event for userId.
*
* If membership of the current user is Invite, it prefers
* the event in inviteState to the one in stateEvents.
*
* @return A lager::reader containing the state event.
*/
auto memberEventByCursor(lager::reader<std::string> userId) const -> lager::reader<Event>;
/**
* Get the member event for userId.
*
* If membership of the current user is Invite, it prefers
* the event in inviteState to the one in stateEvents.
*
* @return A lager::reader containing the state event.
*/
auto memberEventFor(std::string userId) const -> lager::reader<Event>;
/**
* Get whether this room is encrypted.
*
* The encryption status is changed to true if the client
* receives a state event that turns on encryption.
* If that state event is removed later, the status will
* not be changed.
*
* @return A lager::reader<bool> that contains
* whether this room is encrypted.
*/
lager::reader<bool> encrypted() const;
/*lager::reader<std::string>*/
KAZV_WRAP_ATTR(RoomModel, roomCursor(), roomId);
/*lager::reader<RoomMembership>*/
KAZV_WRAP_ATTR(RoomModel, roomCursor(), membership);
/*lager::reader<std::string>*/
KAZV_WRAP_ATTR(RoomModel, roomCursor(), localDraft);
/* lager::reader<bool> */
KAZV_WRAP_ATTR(RoomModel, roomCursor(), membersFullyLoaded);
/**
* Get the local read marker in this room.
*
* @return a lager::reader of a std::string of the local read marker.
*/
auto localReadMarker() const -> lager::reader<std::string>;
/**
* Get the ids of the heroes of the room.
*
* @return a lager::reader of a RangeT<std::string> containing
* the ids of the heroes of the room.
*/
auto heroIds() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the joined member count of this room.
*
* @return A lager::reader of std::size_t containing the joined
* member count for this room.
*/
auto joinedMemberCount() const -> lager::reader<std::size_t>;
/**
* Get the invited member count of this room.
*
* @return A lager::reader of std::size_t containing the invited
* member count for this room.
*/
auto invitedMemberCount() const -> lager::reader<std::size_t>;
/**
* Set local draft for this room.
*
* After the returned Promise is resolved,
* @c localDraft() will contain @c localDraft .
*
* @param localDraft The local draft to send.
* @return A Promise that resolves when the local draft
* has been set, or when there is an error.
*/
PromiseT setLocalDraft(std::string localDraft) const;
/**
* Send an event to this room.
*
* @param msg The message to send
* @return A Promise that resolves when the event has been sent,
* or when there is an error.
*/
PromiseT sendMessage(Event msg) const;
/**
* Send a text message to this room.
*
* @param text The text
* @return A Promise that resolves when the text message has
* been sent, or when there is an error.
*/
PromiseT sendTextMessage(std::string text) const;
/**
* Resend an event to this room.
*
* @param txnId The transaction id of the unsent message
* @return A Promise that resolves when the event has been sent,
* or when there is an error.
*/
PromiseT resendMessage(std::string txnId) const;
/**
* Redact an event.
*
* @param eventId The event id of the event you want to redact
* @param reason The reason of redaction
* @return A Promise that resolves when the event has been sent,
* or when there is an error.
*/
PromiseT redactEvent(std::string eventId, std::optional<std::string> reason) const;
/**
* Send one pending key event in this room.
*
* @param txnId The transaction id of the pending key event
* @return A Promise that resolves when the event has been sent,
* or when there is an error.
*/
PromiseT sendPendingKeyEvent(std::string txnId) const;
/**
* Send all pending key events in this room.
*
* @return A Promise that resolves when all events has been sent,
* or when there is an error.
*/
PromiseT sendAllPendingKeyEvents() const;
/**
* Get the full state of this room.
*
* This method will update the Client as needed.
*
* After the returned Promise resolves successfully,
* @c stateEvents() will contain the fetched state.
*
* @return A Promise that resolves when the room state
* has been fetched, or when there is an error.
*/
PromiseT refreshRoomState() const;
/**
* Get one state event with @c type and @c stateKey .
*
* This method will update the Client as needed.
*
* After the returned Promise resolves successfully,
* @c state({type,stateKey}) will contain the fetched
* state event.
*
* @return A Promise that resolves when the state
* event has been fetched, or when there is an error.
*/
PromiseT getStateEvent(std::string type, std::string stateKey) const;
/**
* Send a state event to this room.
*
* @param state The state event to send.
* @return A Promise that resolves when the state event
* has been sent, or when there is an error.
*/
PromiseT sendStateEvent(Event state) const;
/**
* Set the room name.
*
* @param name The new name for this room.
* @return A Promise that resolves when the state event
* for the name change has been sent, or when there is an error.
*/
PromiseT setName(std::string name) const;
// lager::reader<std::string>
inline auto topic() const {
using namespace lager::lenses;
return stateEvents()
[KeyOfState{"m.room.topic", ""}]
[or_default]
.xform(eventContent
| jsonAtOr("topic"s, ""s));
}
/**
* Set the room topic.
*
* @param topic The new topic for this room.
* @return A Promise that resolves when the state event
* for the topic change has been sent, or when there is an error.
*/
PromiseT setTopic(std::string topic) const;
/**
* Invite a user to this room
*
* @param userId The user id for the user to invite.
* @return A Promise that resolves when the state event
* for the invite has been sent, or when there is an error.
*/
PromiseT invite(std::string userId) const;
/* lager::reader<MapT<std::string, Event>> */
inline auto ephemeralEvents() const {
return roomCursor()
[&RoomModel::ephemeral];
}
/* lager::reader<std::optional<Event>> */
inline auto ephemeralOpt(std::string type) const {
return roomCursor()
[&RoomModel::ephemeral]
[type];
}
/* lager::reader<Event> */
inline auto ephemeral(std::string type) const {
return roomCursor()
[&RoomModel::ephemeral]
[type]
[lager::lenses::or_default];
}
/**
* Get the ids of all typing users in this room.
*
* @return A lager::reader of an RangeT of std::string of ids of all typing users.
*/
auto typingUsers() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Get the member events of all typing users in this room.
*
* @return A lager::reader of an EventList of member events of all typing users.
*/
auto typingMemberEvents() const -> lager::reader<EventList>;
/**
* Set the typing status of the current user in this room.
*
* @param typing Whether the user is now typing.
* @param timeoutMs How long this typing status should last,
* in milliseconds.
* @return A Promise that resolves when the typing status
* has been sent, or when there is an error.
*/
PromiseT setTyping(bool typing, std::optional<int> timeoutMs) const;
/* lager::reader<MapT<std::string, Event>> */
inline auto accountDataEvents() const {
return roomCursor()
[&RoomModel::accountData];
}
/* lager::reader<std::optional<Event>> */
inline auto accountDataOpt(std::string type) const {
return roomCursor()
[&RoomModel::accountData]
[type];
}
/* lager::reader<Event> */
inline auto accountData(std::string type) const {
return roomCursor()
[&RoomModel::accountData]
[type]
[lager::lenses::or_default];
}
/* lager::reader<std::string> */
inline auto readMarker() const {
using namespace lager::lenses;
return accountData("m.fully_read")
.xform(eventContent
| jsonAtOr("event_id", std::string{}));
}
/**
* Set the account data for this room.
*
* @return A Promise that resolves when the account data
* has been set, or when there is an error.
*/
PromiseT setAccountData(Event accountDataEvent) const;
/**
* Get the tags of the current room.
*
* @return A lager::reader of a map from tag id of this room to
* the corresponding order.
*/
auto tags() const -> lager::reader<immer::map<std::string, double>>;
/**
* Add or set a tag to this room.
*
* @param tagId The tag id to add or set.
* @param order The order to specify.
*
* @return A Promise that resolves when the tag
* has been set, or when there is an error.
*/
PromiseT addOrSetTag(std::string tagId, std::optional<double> order = std::nullopt) const;
/**
* Remove a tag from this room.
*
* @param tagId The tag id to remove.
*
* @return A Promise that resolves when the tag
* has been removed, or when there is an error.
*/
PromiseT removeTag(std::string tagId) const;
/**
* Leave this room.
*
* @return A Promise that resolves when the state event
* for the leaving has been sent, or when there is an error.
*/
PromiseT leave() const;
/**
* Forget this room.
*
* One can only forget a room when they have already left it.
*
* @return A Promise that resolves when the room has been
* forgot, or when there is an error.
*/
PromiseT forget() const;
/**
* Kick a user from this room.
*
* You must have enough power levels in this room to do so.
*
* @param userId The id of the user that will be kicked.
* @param reason The reason to explain this kick.
*
* @return A Promise that resolves when the kick is done,
* or when there is an error.
*/
PromiseT kick(std::string userId, std::optional<std::string> reason = std::nullopt) const;
/**
* Ban a user from this room.
*
* You must have enough power levels in this room to do so.
*
* @param userId The id of the user that will be banned.
* @param reason The reason to explain this ban.
*
* @return A Promise that resolves when the ban is done,
* or when there is an error.
*/
PromiseT ban(std::string userId, std::optional<std::string> reason = std::nullopt) const;
// TODO: v1.1 adds reason field
/**
* Unban a user from this room.
*
* You must have enough power levels in this room to do so.
*
* @param userId The id of the user that will be unbanned.
*
* @return A Promise that resolves when the unban is done,
* or when there is an error.
*/
PromiseT unban(std::string userId/*, std::optional<std::string> reason = std::nullopt*/) const;
/* lager::reader<JsonWrap> */
inline auto avatar() const {
return state(KeyOfState{"m.room.avatar", ""})
.xform(eventContent);
}
/**
* Get pinned events of this room
*
* @return A lager::reader of RangeT of the pinned event ids.
*/
auto pinnedEvents() const -> lager::reader<immer::flex_vector<std::string>>;
/**
* Set pinned events of this room
*
* @param eventIds The event ids of the new pinned events
*
* @return A Promise that resolves when the state event
* for the pinned events change has been sent, or when there is an error.
*/
PromiseT setPinnedEvents(immer::flex_vector<std::string> eventIds) const;
/**
* Add eventIds to the pinned events of this room
*
* @param eventIds The ids of events you want to add to the pinned events.
* @return A Promise that resolves when the state event
* for the pinned events change has been sent, or when there is an error.
*/
PromiseT pinEvents(immer::flex_vector<std::string> eventIds) const;
/**
* Remove eventIds from the pinned events of this room
*
* @param eventIds The ids of events you want to remove from the pinned events.
* @return A Promise that resolves when the state event
* for the pinned events change has been sent, or when there is an error.
*/
PromiseT unpinEvents(immer::flex_vector<std::string> eventIds) const;
/**
* Get the Gaps in the timeline for this room.
*
* Any key of the map in the returned reader can be send as
* an argument of paginateBackFromEvent() to try to fill the Gap
* at that event.
*
* @return A lager::reader that contains an evnetId-to-prevBatch map.
*/
lager::reader<immer::map<std::string /* eventId */, std::string /* prevBatch */>> timelineGaps() const;
/**
* Try to paginate back from @c eventId.
*
* @param eventId An event id that is in the key of `+timelineGaps()`.
*
* @return A Promise that resolves when the pagination is
* successful, or when there is an error. If it is successful,
* `+timelineGaps()` will no longer contain eventId as key, and
* `timeline()` will contain the events before eventId in the
* full event chain on the homeserver.
* If `eventId` is not in `+timelineGaps()`, it is considered
* to be failed.
*/
PromiseT paginateBackFromEvent(std::string eventId) const;
/**
* Get the list of local echoes in this room.
*
* @return A lager::reader that contains the list of all local echoes in the room.
* The event type will always be m.room.message and event id will not be meaningful.
*/
auto localEchoes() const -> lager::reader<immer::flex_vector<LocalEchoDesc>>;
/**
* Remove a local echo from this room.
*
* @param txnId The transaction id associated with that local echo.
*
* @return A Promise that resolves when the local echo is removed.
*/
PromiseT removeLocalEcho(std::string txnId) const;
/**
* Get the list of pending room key events in this room.
*
* @return A lager::reader that contains the list of all pending room key events in the room.
* The event type will always be m.room.encrypted and event id will not be meaningful.
*/
auto pendingRoomKeyEvents() const -> lager::reader<immer::flex_vector<PendingRoomKeyEvent>>;
/**
* Get the power levels of this room.
*
* @return A lager::reader that contains a PowerLevelDesc describing the power levels of this room.
*/
auto powerLevels() const -> lager::reader<PowerLevelsDesc>;
/**
* Get a list of child events of a specified event.
*
* @param eventId The id of the event.
* @param relType The type of the relationship.
* @return A lager::reader of a RangeT of events that are related to the event via `relType`. It is sorted in the same order as the timeline.
*/
auto relatedEvents(lager::reader<std::string> eventId, std::string relType) const -> lager::reader<EventList>;
/**
* Get a list of read receipts of some event in this room.
*
* @param eventId The id of the event.
* @return A lager::reader of a RangeT of matrix id and timestamp of receipts of the event.
*/
auto eventReaders(lager::reader<std::string> eventId) const -> lager::reader<immer::flex_vector<EventReader>>;
/**
* Post a read receipt for this room.
*
* @param eventId The event id the user has read up to.
*
* @return A Promise that resolves when the read receipt is sent.
*/
PromiseT postReceipt(std::string eventId) const;
/**
* Get a list of event ids of unread notifications.
*
* @return A lager::reader of a RangeT of event ids that
* has an unread notification.
*/
auto unreadNotificationEventIds() const -> lager::reader<immer::flex_vector<std::string>>;
+ /**
+ * Check if an event has been read by the given user
+ *
+ * @param eventId eventId of the event to check
+ * @param userId userId of the given user
+ * @return true iff the event is read by the user
+ */
+ bool isReadBy(const std::string &eventId, const std::string &userId);
+
private:
const lager::reader<SdkModel> &sdkCursor() const;
const lager::reader<RoomModel> &roomCursor() const;
lager::reader<RoomModel> makeRoomCursor() const;
std::string currentRoomId() const;
/// if membership is invite, return inviteState; otherwise return stateEvents
lager::reader<immer::map<KeyOfState, Event>> inviteStateOrState() const;
/// if membership is invite and inviteState[key] exists, return inviteState[key], otherwise return stateEvents[key]
lager::reader<Event> inviteStateOrStateEvent(lager::reader<KeyOfState> key) const;
std::optional<lager::reader<SdkModel>> m_sdk;
std::variant<lager::reader<std::string>, std::string> m_roomId;
ContextT m_ctx;
std::optional<DepsT> m_deps;
lager::reader<RoomModel> m_roomCursor;
KAZV_DECLARE_THREAD_ID();
KAZV_DECLARE_EVENT_LOOP_THREAD_ID_KEEPER(m_deps.has_value() ? &lager::get<EventLoopThreadIdKeeper &>(m_deps.value()) : 0);
};
}
diff --git a/src/testfixtures/factory.cpp b/src/testfixtures/factory.cpp
index 0b183df..5b6a4b0 100644
--- a/src/testfixtures/factory.cpp
+++ b/src/testfixtures/factory.cpp
@@ -1,321 +1,326 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include "factory.hpp"
namespace Kazv::Factory
{
static std::string generateUserId()
{
static std::size_t next = 0;
++next;
return "@u" + std::to_string(next) + ":example.com";
}
static std::string generateEventId()
{
static std::size_t next = 0;
++next;
std::stringstream s;
s << std::setw(40) << std::setfill('0') << next;
return "$" + s.str();
}
static std::string generateRoomId()
{
static std::size_t next = 0;
++next;
return "!" + std::to_string(next) + ":example.com";
}
static std::string generateDeviceId()
{
static std::size_t next = 0;
++next;
return "dev" + std::to_string(next);
}
ClientModel makeClient(const ComposedModifier<ClientModel> &mod)
{
ClientModel m;
m.serverUrl = "example.com";
m.userId = "@bob:example.com";
m.token = "exampletoken";
m.deviceId = "exampledevice";
m.loggedIn = true;
mod(m);
return m;
}
ComposedModifier<ClientModel> withRoom(RoomModel room)
{
return [room](ClientModel &m) {
m.roomList.rooms = m.roomList.rooms.set(room.roomId, room);
};
}
ComposedModifier<ClientModel> withAccountData(immer::flex_vector<Event> accountDataEvent)
{
return [accountDataEvent](ClientModel &m) {
m.accountData = merge(std::move(m.accountData), accountDataEvent, keyOfAccountData);
};
}
ComposedModifier<ClientModel> withCrypto(const Crypto &crypto)
{
return [crypto](ClientModel &m) {
m.crypto = crypto;
};
}
ComposedModifier<ClientModel> withDevice(std::string userId, DeviceKeyInfo info)
{
return [userId, info](ClientModel &m) {
m.deviceLists.deviceLists = std::move(m.deviceLists.deviceLists)
.update(userId, [=](auto deviceMap) {
return std::move(deviceMap).set(info.deviceId, info);
});
};
}
DeviceKeyInfo makeDeviceKeyInfo(const ComposedModifier<DeviceKeyInfo> &mod)
{
DeviceKeyInfo m;
mod(m);
if (m.deviceId.empty()) {
withDeviceId(generateDeviceId())(m);
}
if (m.ed25519Key.empty()) {
m.ed25519Key = m.deviceId + "ed25519";
}
if (m.curve25519Key.empty()) {
m.curve25519Key = m.deviceId + "curve25519";
}
return m;
}
ComposedModifier<DeviceKeyInfo> withDeviceId(std::string deviceId)
{
return [deviceId](DeviceKeyInfo &m) { m.deviceId = deviceId; };
}
ComposedModifier<DeviceKeyInfo> withDeviceDisplayName(std::string displayName)
{
return [displayName](DeviceKeyInfo &m) { m.displayName = displayName; };
}
ComposedModifier<DeviceKeyInfo> withDeviceTrustLevel(DeviceTrustLevel trustLevel)
{
return [trustLevel](DeviceKeyInfo &m) { m.trustLevel = trustLevel; };
}
Crypto makeCrypto(const ComposedModifier<Crypto> &mod)
{
auto m = Crypto(RandomTag{}, genRandomData(Crypto::constructRandomSize()));
mod(m);
return m;
}
RoomModel makeRoom(const ComposedModifier<RoomModel> &mod)
{
RoomModel room;
room.roomId = generateRoomId();
mod(room);
return room;
}
ComposedModifier<RoomModel> withRoomId(std::string id)
{
return [id](RoomModel &m) { m.roomId = id; };
}
ComposedModifier<RoomModel> withRoomAccountData(immer::flex_vector<Event> accountDataEvent)
{
return [accountDataEvent](RoomModel &m) {
m.accountData = merge(std::move(m.accountData), accountDataEvent, keyOfAccountData);
};
}
ComposedModifier<RoomModel> withRoomState(immer::flex_vector<Event> stateEvent)
{
return [stateEvent](RoomModel &m) {
m.stateEvents = merge(std::move(m.stateEvents), stateEvent, keyOfState);
};
}
ComposedModifier<RoomModel> withRoomTimeline(immer::flex_vector<Event> timelineEvents)
{
return [timelineEvents](RoomModel &m) {
m = RoomModel::update(m, AddToTimelineAction{timelineEvents, std::nullopt, std::nullopt, std::nullopt});
};
}
ComposedModifier<RoomModel> withRoomTimelineGaps(immer::map<std::string, std::string> timelineGaps)
{
return [timelineGaps](RoomModel &m) {
m.timelineGaps = timelineGaps;
};
}
ComposedModifier<RoomModel> withRoomMembership(RoomMembership membership)
{
return [membership](RoomModel &m) {
m.membership = membership;
};
}
ComposedModifier<RoomModel> withRoomEncrypted(bool encrypted)
{
return [encrypted](RoomModel &m) {
m.encrypted = encrypted;
};
}
Event makeEvent(const ComposedModifier<Event> &mod)
{
Event e = json{
{"type", "m.room.message"},
{"sender", "@foo:tusooa.xyz"},
{"content", {
{"msgtype", "m.text"},
{"body", "test"},
}},
{"origin_server_ts", 1000},
{"event_id", generateEventId()},
};
mod(e);
return e;
}
Event makeMemberEvent(const ComposedModifier<Event> &mod)
{
auto userId = generateUserId();
auto event = makeEvent(
withEventType("m.room.member")
| withStateKey(userId)
| withEventContent(json::object())
| withMembership("join")
| mod
);
if (event.sender().empty()) {
withEventSenderId(event.stateKey());
}
return event;
}
ComposedModifier<Event> withEventJson(const json &j)
{
return [j](Event &e) {
e = Event(j);
};
}
ComposedModifier<Event> withEventKV(const json::json_pointer &k, const json &v)
{
return [k, v](Event &e) {
auto j = e.originalJson().get();
j[k] = v;
withEventJson(j)(e);
};
}
+ ComposedModifier<Event> withEventTs(Timestamp ts)
+ {
+ return withEventKV("/origin_server_ts"_json_pointer, ts);
+ }
+
ComposedModifier<Event> withEventId(std::string id)
{
return withEventKV("/event_id"_json_pointer, id);
}
ComposedModifier<Event> withEventType(std::string type)
{
return withEventKV("/type"_json_pointer, type);
}
ComposedModifier<Event> withEventContent(const json &content)
{
return withEventKV("/content"_json_pointer, content);
}
ComposedModifier<Event> withStateKey(std::string id)
{
return withEventKV("/state_key"_json_pointer, id);
}
ComposedModifier<Event> withMembership(std::string membership)
{
return withEventKV("/content/membership"_json_pointer, membership);
}
ComposedModifier<Event> withMemberDisplayName(std::string displayName)
{
return withEventKV("/content/displayname"_json_pointer, displayName);
}
ComposedModifier<Event> withMemberAvatarUrl(std::string avatarUrl)
{
return withEventKV("/content/avatar_url"_json_pointer, avatarUrl);
}
ComposedModifier<Event> withEventSenderId(std::string sender)
{
return withEventKV("/sender"_json_pointer, sender);
}
ComposedModifier<Event> withEventRelationship(std::string relType, std::string eventId)
{
return withEventKV("/content/m.relates_to"_json_pointer, {
{"rel_type", relType},
{"event_id", eventId},
});
}
ComposedModifier<Event> withEventReplyTo(std::string eventId)
{
return withEventKV("/content/m.relates_to"_json_pointer, {
{"m.in_reply_to", {{"event_id", eventId}}},
});
}
Response makeResponse(std::string jobId, const ComposedModifier<Response> &mod)
{
Response r;
(withResponseStatusCode(200)
| withResponseJsonBody(json::object())
| withResponseDataKV("-job-id", jobId)
| mod)(r);
return r;
}
ComposedModifier<Response> withResponseStatusCode(int code)
{
return withAttr(&Response::statusCode, code);
}
ComposedModifier<Response> withResponseJsonBody(const json &body)
{
return withAttr(&Response::body, JsonWrap(body));
}
ComposedModifier<Response> withResponseBytesBody(const Bytes &body)
{
return withAttr(&Response::body, body);
}
ComposedModifier<Response> withResponseFileBody(const FileDesc &body)
{
return withAttr(&Response::body, body);
}
ComposedModifier<Response> withResponseDataKV(std::string k, const json &v)
{
return [k, v](Response &m) {
auto data = m.extraData.get();
if (!data.is_object()) {
data = json::object();
}
data[k] = v;
m.extraData = data;
};
}
}
diff --git a/src/testfixtures/factory.hpp b/src/testfixtures/factory.hpp
index 9997404..e5990cc 100644
--- a/src/testfixtures/factory.hpp
+++ b/src/testfixtures/factory.hpp
@@ -1,105 +1,106 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2023 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#pragma once
#include <libkazv-config.hpp>
#include <client-model.hpp>
namespace Kazv
{
namespace Factory
{
template<class ModelType>
struct ComposedModifier : std::function<void(ModelType &)>
{
using BaseT = std::function<void(ModelType &)>;
using BaseT::BaseT;
using BaseT::operator();
ComposedModifier() : BaseT([](ModelType &) {}) {}
ComposedModifier operator|(const ComposedModifier &that) const
{
return ComposedModifier([that, *this](ModelType &model) {
(*this)(model);
that(model);
});
}
};
template<class PointerToMember>
struct AttrModifier
{};
template<class ModelType, class DataType>
struct AttrModifier<DataType ModelType::*> : public ComposedModifier<ModelType>
{
using BaseT = ComposedModifier<ModelType>;
using PtrT = DataType ModelType::*;
using ModelT = ModelType;
using DataT = DataType;
AttrModifier(PtrT ptr, const DataType &data)
: BaseT([ptr, data](ModelType &m) {
m.*ptr = data;
})
{}
};
template<class PointerToMember>
AttrModifier<PointerToMember> withAttr(PointerToMember p, const typename AttrModifier<PointerToMember>::DataT &data)
{
return AttrModifier<PointerToMember>(p, data);
}
ClientModel makeClient(const ComposedModifier<ClientModel> &mod = {});
ComposedModifier<ClientModel> withRoom(RoomModel room);
ComposedModifier<ClientModel> withAccountData(immer::flex_vector<Event> accountDataEvent);
ComposedModifier<ClientModel> withCrypto(const Crypto &crypto);
ComposedModifier<ClientModel> withDevice(std::string userId, DeviceKeyInfo info);
DeviceKeyInfo makeDeviceKeyInfo(const ComposedModifier<DeviceKeyInfo> &mod = {});
ComposedModifier<DeviceKeyInfo> withDeviceId(std::string deviceId);
ComposedModifier<DeviceKeyInfo> withDeviceDisplayName(std::string displayName);
ComposedModifier<DeviceKeyInfo> withDeviceTrustLevel(DeviceTrustLevel trustLevel);
Crypto makeCrypto(const ComposedModifier<Crypto> &mod = {});
RoomModel makeRoom(const ComposedModifier<RoomModel> &mod = {});
ComposedModifier<RoomModel> withRoomId(std::string id);
ComposedModifier<RoomModel> withRoomAccountData(immer::flex_vector<Event> accountDataEvent);
ComposedModifier<RoomModel> withRoomState(immer::flex_vector<Event> stateEvent);
ComposedModifier<RoomModel> withRoomInviteState(immer::flex_vector<Event> stateEvent);
ComposedModifier<RoomModel> withRoomTimeline(immer::flex_vector<Event> timelineEvents);
ComposedModifier<RoomModel> withRoomTimelineGaps(immer::map<std::string, std::string> timelineGaps);
ComposedModifier<RoomModel> withRoomMembership(RoomMembership membership);
ComposedModifier<RoomModel> withRoomEncrypted(bool encrypted);
Event makeEvent(const ComposedModifier<Event> &mod = {});
Event makeMemberEvent(const ComposedModifier<Event> &mod = {});
ComposedModifier<Event> withEventJson(const json &j);
ComposedModifier<Event> withEventKV(const json::json_pointer &k, const json &v);
+ ComposedModifier<Event> withEventTs(Timestamp ts);
ComposedModifier<Event> withEventId(std::string id);
ComposedModifier<Event> withEventType(std::string type);
ComposedModifier<Event> withEventContent(const json &content);
ComposedModifier<Event> withStateKey(std::string id);
ComposedModifier<Event> withMembership(std::string membership);
ComposedModifier<Event> withMemberDisplayName(std::string displayName);
ComposedModifier<Event> withMemberAvatarUrl(std::string avatarUrl);
ComposedModifier<Event> withEventSenderId(std::string sender);
ComposedModifier<Event> withEventRelationship(std::string relType, std::string eventId);
ComposedModifier<Event> withEventReplyTo(std::string eventId);
Response makeResponse(std::string jobId, const ComposedModifier<Response> &mod = {});
ComposedModifier<Response> withResponseStatusCode(int code);
ComposedModifier<Response> withResponseJsonBody(const json &body);
ComposedModifier<Response> withResponseBytesBody(const Bytes &body);
ComposedModifier<Response> withResponseFileBody(const FileDesc &body);
ComposedModifier<Response> withResponseDataKV(std::string k, const json &v);
}
}
diff --git a/src/tests/client/room/read-receipt-test.cpp b/src/tests/client/room/read-receipt-test.cpp
index da09f83..162d443 100644
--- a/src/tests/client/room/read-receipt-test.cpp
+++ b/src/tests/client/room/read-receipt-test.cpp
@@ -1,410 +1,438 @@
/*
* This file is part of libkazv.
* SPDX-FileCopyrightText: 2024 tusooa <tusooa@kazv.moe>
* SPDX-License-Identifier: AGPL-3.0-or-later
*/
#include <libkazv-config.hpp>
#include <boost/asio.hpp>
#include <catch2/catch_test_macros.hpp>
#include <catch2/matchers/catch_matchers_range_equals.hpp>
#include <lager/event_loop/boost_asio.hpp>
#include <asio-promise-handler.hpp>
#include <room/room-model.hpp>
#include <sdk-model.hpp>
#include <client/client.hpp>
#include <client/actions/ephemeral.hpp>
#include <cprjobhandler.hpp>
#include <lagerstoreeventemitter.hpp>
#include <testfixtures/factory.hpp>
#include "client-test-util.hpp"
#include "action-mock-utils.hpp"
#include "push-rules-test-util.hpp"
using namespace Kazv;
using namespace Kazv::Factory;
TEST_CASE("Adding an m.receipt ephemeral event", "[client][room][receipt]")
{
auto room = makeRoom();
auto receiptEvent = makeEvent(
withEventType("m.receipt")
// https://spec.matrix.org/v1.8/client-server-api/#events-5
| withEventContent(R"({
"$1435641916114394fHBLK:matrix.org": {
"m.read": {
"@rikj:jki.re": {
"ts": 1436451550453
}
},
"m.read.private": {
"@self:example.org": {
"ts": 1661384801651
}
}
}
})"_json));
room = RoomModel::update(room, AddEphemeralAction{{receiptEvent}});
auto receipt = room.readReceipts.at("@rikj:jki.re");
REQUIRE(
room.readReceipts.at("@rikj:jki.re")
==
ReadReceipt{"$1435641916114394fHBLK:matrix.org", 1436451550453}
);
REQUIRE(
room.eventReadUsers.at("$1435641916114394fHBLK:matrix.org")
==
immer::flex_vector<std::string>{"@rikj:jki.re", "@self:example.org"}
);
auto anotherReceiptEvent = makeEvent(
withEventType("m.receipt")
| withEventContent(json{
{"$123", {
{"m.read", {
{"@foo:example.com", {{"ts", 1796451550450}}}
}},
}},
}));
room = RoomModel::update(room, AddEphemeralAction{{anotherReceiptEvent}});
// the receipt that was there should still be there.
REQUIRE(
room.readReceipts.at("@rikj:jki.re")
==
ReadReceipt{"$1435641916114394fHBLK:matrix.org", 1436451550453}
);
REQUIRE(
room.readReceipts.at("@foo:example.com")
==
ReadReceipt{"$123", 1796451550450}
);
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto model = makeClient(withRoom(room));
auto store = createTestClientStoreFrom(model, ph);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }),
store,
std::nullopt);
auto r = client.room(room.roomId);
auto readers1 = r.eventReaders(lager::make_constant<std::string>("$1435641916114394fHBLK:matrix.org")).make().get();
auto expected1 = immer::flex_vector<EventReader>{{"@rikj:jki.re", 1436451550453}, {"@self:example.org", 1661384801651}};
REQUIRE_THAT(readers1, Catch::Matchers::UnorderedRangeEquals(expected1));
auto readers2 = r.eventReaders(lager::make_constant<std::string>("$123")).make().get();
auto expected2 = immer::flex_vector<EventReader>{{"@foo:example.com", 1796451550450}};
REQUIRE(readers2 == expected2);
}
TEST_CASE("Update a receipt for some user", "[client][room][receipt]")
{
auto room = makeRoom();
auto receiptEvent = makeEvent(
withEventType("m.receipt")
// https://spec.matrix.org/v1.8/client-server-api/#events-5
| withEventContent(R"({
"$1435641916114394fHBLK:matrix.org": {
"m.read": {
"@rikj:jki.re": {
"ts": 1436451550453
}
}
}
})"_json));
room = RoomModel::update(room, AddEphemeralAction{{receiptEvent}});
auto anotherReceiptEvent = makeEvent(
withEventType("m.receipt")
| withEventContent(json{
{"$123", {
{"m.read", {
{"@rikj:jki.re", {{"ts", 1796451550450}}}
}},
}},
}));
room = RoomModel::update(room, AddEphemeralAction{{anotherReceiptEvent}});
REQUIRE(
room.readReceipts.at("@rikj:jki.re")
==
ReadReceipt{"$123", 1796451550450}
);
REQUIRE(room.eventReadUsers.count("$1435641916114394fHBLK:matrix.org") == 0);
REQUIRE(room.eventReadUsers.count("$123") == 1);
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto model = makeClient(withRoom(room));
auto store = createTestClientStoreFrom(model, ph);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }),
store,
std::nullopt);
auto r = client.room(room.roomId);
auto readers1 = r.eventReaders(lager::make_constant<std::string>("$1435641916114394fHBLK:matrix.org")).make().get();
auto expected1 = immer::flex_vector<EventReader>{};
REQUIRE(readers1 == expected1);
auto readers2 = r.eventReaders(lager::make_constant<std::string>("$123")).make().get();
auto expected2 = immer::flex_vector<EventReader>{{"@rikj:jki.re", 1796451550450}};
REQUIRE(readers2 == expected2);
}
TEST_CASE("localReadMarker", "[client][room][receipt]")
{
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto tl = EventList{
makeEvent(),
};
auto room = makeRoom(withRoomTimeline(tl));
room.localReadMarker = tl[0].id();
auto model = makeClient(withRoom(room));
auto store = createTestClientStoreFrom(model, ph);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }),
store,
std::nullopt);
auto r = client.room(room.roomId);
REQUIRE(r.localReadMarker().get() == tl[0].id());
}
TEST_CASE("Posting receipts", "[client][room][receipt]")
{
auto r = makeRoom();
auto m = makeClient(withRoom(r));
boost::asio::io_context io;
SingleTypePromiseInterface<EffectStatus> sgph{AsioPromiseHandler{io.get_executor()}};
auto jh = Kazv::CprJobHandler{io.get_executor()};
auto ee = Kazv::LagerStoreEventEmitter(lager::with_boost_asio_event_loop{io.get_executor()});
auto sdk = Kazv::makeSdk(
SdkModel{m},
jh,
ee,
Kazv::AsioPromiseHandler{io.get_executor()},
zug::identity
);
auto ctx = sdk.context();
auto dispatcher = getMockDispatcher(
sgph,
ctx,
returnEmpty<PostReceiptAction>()
);
auto mockContext = getMockContext(sgph, dispatcher);
auto client = Client(Client::InEventLoopTag{}, mockContext, sdk.context());
auto room = client.room(r.roomId);
room.postReceipt("$1")
.then([&io](auto) {
io.stop();
});
io.run();
REQUIRE(dispatcher.template calledTimes<PostReceiptAction>() == 1);
auto action = dispatcher.template of<PostReceiptAction>()[0];
REQUIRE(action.roomId == r.roomId);
REQUIRE(action.eventId == "$1");
}
TEST_CASE("PostReceiptAction", "[client][room][receipt]")
{
auto m = makeClient();
auto [next, _ignore] = updateClient(m, PostReceiptAction{"!someroom:example.com", "$someevent"});
assert1Job(next);
for1stJob(next, [](const BaseJob &job) {
REQUIRE(job.jobId() == "PostReceipt");
REQUIRE(job.url().find("rooms/!someroom:example.com/receipt/m.read/$someevent") != std::string::npos);
REQUIRE(json::parse(std::get<Bytes>(job.requestBody())) == json::object());
});
}
static auto pushRules = PushRulesDesc(Event(R"({
"type": "m.push_rules",
"content": {
"global": {
"override": [{
"rule_id": "moe.kazv.mxc.some-rule",
"default": true,
"enabled": true,
"conditions": [{
"kind": "event_match",
"key": "type",
"pattern": "moe.kazv.mxc.test-event"
}],
"actions": ["notify"]
}]
}
}
})"_json));
TEST_CASE("AddLocalNotificationsAction", "[client][room][receipt]")
{
auto myUserId = "@mew:example.com"s;
auto newEvents = EventList{
makeEvent(withEventType("moe.kazv.mxc.test-event")),
makeEvent(),
makeEvent(withEventType("moe.kazv.mxc.test-event")),
makeEvent(withEventType("moe.kazv.mxc.test-event") | withEventSenderId(myUserId)),
makeEvent(withEventType("moe.kazv.mxc.test-event")),
};
auto r = makeRoom(withRoomTimeline(newEvents));
r.readReceipts = {{myUserId, ReadReceipt{newEvents[0].id(), 0}}};
WHEN("push rules does not notify") {
auto next = RoomModel::update(r, AddLocalNotificationsAction{
newEvents,
PushRulesDesc(Event()),
myUserId,
});
REQUIRE(next.unreadNotificationEventIds == immer::flex_vector<std::string>{});
}
WHEN("push rules notifies some") {
auto next = RoomModel::update(r, AddLocalNotificationsAction{
newEvents,
pushRules,
myUserId,
});
REQUIRE(
next.unreadNotificationEventIds
== immer::flex_vector<std::string>{
newEvents[2].id(),
newEvents[4].id(),
}
);
}
WHEN("push rules notifies some but is before local read marker, only things after local read marker is added") {
r.localReadMarker = newEvents[3].id();
auto next = RoomModel::update(r, AddLocalNotificationsAction{
newEvents,
pushRules,
myUserId,
});
REQUIRE(
next.unreadNotificationEventIds
== immer::flex_vector<std::string>{
newEvents[4].id(),
}
);
}
}
TEST_CASE("PostReceiptAction should update local read marker and remove read notifications")
{
auto myUserId = "@mew:example.com"s;
auto newEvents = EventList{
makeEvent(withEventType("moe.kazv.mxc.test-event")),
makeEvent(withEventType("moe.kazv.mxc.test-event")),
makeEvent(withEventType("moe.kazv.mxc.test-event")),
};
auto r = makeRoom(withRoomTimeline(newEvents));
r.readReceipts = {{myUserId, ReadReceipt{newEvents[0].id(), 0}}};
r = RoomModel::update(r, AddLocalNotificationsAction{
newEvents,
pushRules,
myUserId,
});
auto client = makeClient(withRoom(r));
client.userId = myUserId;
auto receiptPos = newEvents[1].id();
auto [next, _] = ClientModel::update(client, PostReceiptAction{r.roomId, receiptPos});
auto nextRoom = next.roomList.rooms.at(r.roomId);
REQUIRE(nextRoom.localReadMarker == receiptPos);
REQUIRE(nextRoom.unreadNotificationEventIds == immer::flex_vector<std::string>{newEvents[2].id()});
}
TEST_CASE("RemoveReadLocalNotificationsAction", "[client][room][receipt]")
{
auto myUserId = "@mew:example.com"s;
auto oldEvents = EventList{
makeEvent(),
makeEvent(),
};
auto newEvents = EventList{
makeEvent(),
makeEvent(),
makeEvent(),
makeEvent(),
makeEvent(),
makeEvent(),
};
auto timeline = oldEvents + newEvents;
auto r = makeRoom(withRoomTimeline(timeline));
r.unreadNotificationEventIds = intoImmer(
immer::flex_vector<std::string>(),
zug::map(&Event::id),
newEvents.erase(5)
);
WHEN("a read receipt is in the middle") {
r.readReceipts = {{myUserId, ReadReceipt{newEvents[2].id(), 0}}};
auto next = RoomModel::update(r, RemoveReadLocalNotificationsAction{myUserId});
REQUIRE(next.unreadNotificationEventIds == immer::flex_vector<std::string>{newEvents[3].id(), newEvents[4].id()});
}
WHEN("a read receipt is before the beginning") {
r.readReceipts = {{myUserId, ReadReceipt{oldEvents[0].id(), 0}}};
auto next = RoomModel::update(r, RemoveReadLocalNotificationsAction{myUserId});
REQUIRE(next.unreadNotificationEventIds == r.unreadNotificationEventIds);
}
WHEN("a read receipt does not exist") {
r.readReceipts = {};
auto next = RoomModel::update(r, RemoveReadLocalNotificationsAction{myUserId});
REQUIRE(next.unreadNotificationEventIds == r.unreadNotificationEventIds);
}
WHEN("a read receipt is at the end") {
r.readReceipts = {{myUserId, ReadReceipt{newEvents[5].id(), 0}}};
auto next = RoomModel::update(r, RemoveReadLocalNotificationsAction{myUserId});
REQUIRE(next.unreadNotificationEventIds == immer::flex_vector<std::string>{});
}
WHEN("a read receipt is after the end") {
r.readReceipts = {{myUserId, ReadReceipt{newEvents[5].id(), 0}}};
auto next = RoomModel::update(r, RemoveReadLocalNotificationsAction{myUserId});
REQUIRE(next.unreadNotificationEventIds == immer::flex_vector<std::string>{});
}
}
TEST_CASE("Room::unreadNotificationEventIds()", "[client][room][receipt]")
{
boost::asio::io_context io;
AsioPromiseHandler ph{io.get_executor()};
auto room = makeRoom();
room.unreadNotificationEventIds = {"$foo", "$bar"};
auto model = makeClient(withRoom(room));
auto store = createTestClientStoreFrom(model, ph);
auto client = Client(store.reader().map([](auto c) { return SdkModel{c}; }),
store,
std::nullopt);
auto r = client.room(room.roomId);
REQUIRE(r.unreadNotificationEventIds().get()
== room.unreadNotificationEventIds);
}
+
+TEST_CASE("RoomModel::isReadBy", "[client][room][receipt]")
+{
+ auto room = makeRoom();
+ auto receiptEvent = makeEvent(
+ withEventType("m.receipt")
+ // https://spec.matrix.org/v1.8/client-server-api/#events-5
+ | withEventContent(R"({
+ "$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z2": {
+ "m.read": {
+ "@rikj:jki.re": {
+ "ts": 1145141919810
+ }
+ }
+ }
+ })"_json));
+ auto timeline = EventList{makeEvent(withEventTs(1436451550452) | withEventId("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z0")),
+ makeEvent(withEventTs(1436451550453) | withEventId("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z7")),
+ makeEvent(withEventTs(1436451550453) | withEventId("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z2")),
+ makeEvent(withEventTs(1436451550454) | withEventId("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z1"))
+ };
+ room = RoomModel::update(std::move(room), AddMessagesAction{timeline});
+ room = RoomModel::update(std::move(room), AddEphemeralAction({receiptEvent}));
+ REQUIRE(room.isReadBy("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z0", "@rikj:jki.re") == true);
+ REQUIRE(room.isReadBy("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z7", "@rikj:jki.re") == false);
+ REQUIRE(room.isReadBy("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z2", "@rikj:jki.re") == true);
+ REQUIRE(room.isReadBy("$Rqnc-F-dvnEYJTyHq_iKxU2bZ1CI92-kuZq3a5lr5Z1", "@rikj:jki.re") == false);
+}
File Metadata
Details
Attached
Mime Type
text/x-diff
Expires
Sun, Aug 9, 5:07 AM (1 d, 14 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1724234
Default Alt Text
(142 KB)
Attached To
Mode
rL libkazv
Attached
Detach File
Event Timeline
Log In to Comment