Add support for on-demand involvement

Fixes https://github.com/etkecc/baibot/issues/15
This commit is contained in:
Slavi Pantaleev
2024-10-01 21:06:54 +03:00
parent eae6472c7a
commit 9908512968
23 changed files with 832 additions and 329 deletions

View File

@@ -5,7 +5,9 @@ use std::sync::Arc;
use mxlink::matrix_sdk::ruma::{OwnedEventId, OwnedUserId};
use mxlink::matrix_sdk::{
deserialized_responses::TimelineEvent,
ruma::events::{
relation::Thread,
room::message::{
MessageType, OriginalSyncRoomMessageEvent, Relation, RoomMessageEventContent,
},
@@ -16,7 +18,12 @@ use mxlink::matrix_sdk::{
use mxlink::{MatrixLink, ThreadGetMessagesParams, ThreadInfo};
use super::{MatrixMessage, MatrixMessageProcessingParams, MatrixMessageType, RoomEventFetcher};
use crate::entity::{MessagePayload, ThreadContext, ThreadContextFirstMessage};
use crate::entity::{InteractionContext, InteractionTrigger, MessagePayload};
struct DetailedMessagePayload {
is_mentioning_bot: bool,
message_payload: MessagePayload,
}
pub async fn get_matrix_messages_in_thread(
matrix_link: MatrixLink,
@@ -42,23 +49,123 @@ pub async fn get_matrix_messages_in_thread(
Ok(messages)
}
pub async fn process_matrix_messages_in_thread(
pub async fn get_matrix_messages_in_reply_chain(
event_fetcher: &Arc<RoomEventFetcher>,
room: &Room,
event_id: OwnedEventId,
) -> Result<Vec<MatrixMessage>, mxlink::matrix_sdk::Error> {
let messages_native =
get_matrix_messages_in_reply_chain_native(event_fetcher, room, event_id).await?;
let mut messages: Vec<MatrixMessage> = Vec::new();
for matrix_native_message in messages_native {
let Some(message) = convert_matrix_native_event_to_matrix_message(&matrix_native_message)
else {
continue;
};
messages.push(message);
}
Ok(messages)
}
async fn get_matrix_messages_in_reply_chain_native(
event_fetcher: &Arc<RoomEventFetcher>,
room: &Room,
event_id: OwnedEventId,
) -> Result<Vec<AnyMessageLikeEvent>, mxlink::matrix_sdk::Error> {
let mut next_event_id = Some(event_id.clone());
let mut messages: Vec<AnyMessageLikeEvent> = Vec::new();
let mut handled_event_ids: Vec<OwnedEventId> = Vec::new();
while let Some(next_event_id_in_loop) = next_event_id {
let event = event_fetcher
.fetch_event_in_room(&next_event_id_in_loop, room)
.await
.unwrap();
if handled_event_ids.contains(&next_event_id_in_loop) {
tracing::warn!(
"Not following loop-causing event: {}",
next_event_id_in_loop
);
break;
}
handled_event_ids.push(next_event_id_in_loop.clone());
let event_deserialized = event.event.deserialize()?;
let AnyTimelineEvent::MessageLike(message_like_event) = event_deserialized else {
tracing::warn!(
"Not proceeding past non-MessageLike event: {:?}",
event_deserialized
);
break;
};
next_event_id = match message_like_event.clone() {
AnyMessageLikeEvent::RoomEncrypted(_) => None,
AnyMessageLikeEvent::RoomMessage(room_message) => {
if let MessageLikeEvent::Original(room_message_original) = room_message {
match room_message_original.content.relates_to {
Some(Relation::Reply { in_reply_to }) => Some(in_reply_to.event_id.clone()),
_ => None,
}
} else {
None
}
}
_ => None,
};
messages.push(message_like_event);
}
messages.reverse();
Ok(messages)
}
pub async fn process_matrix_messages(
messages: &[MatrixMessage],
params: &MatrixMessageProcessingParams,
) -> Vec<MatrixMessage> {
let mut messages_filtered: Vec<MatrixMessage> = Vec::new();
for (i, message) in messages.iter().enumerate() {
if !is_message_from_allowed_sender(message, &params.bot_user_id, &params.allowed_users) {
if !is_message_from_allowed_sender(
message,
&params.bot_user_id,
params.allowed_users.as_deref(),
) {
continue;
}
let mut message = message.clone();
if i == 0 && !params.first_message_stripped_prefixes.is_empty() {
if i == 0 && !params.first_message_prefixes_to_strip.is_empty() {
let mut message_text = message.message_text.clone();
for prefix in &params.first_message_stripped_prefixes {
for prefix in &params.first_message_prefixes_to_strip {
if let Some(message_text_stripped) = message_text.strip_prefix(prefix) {
message_text = message_text_stripped.to_owned();
}
}
message.message_text = message_text.trim().to_owned();
}
// We only strip `bot_user_prefixes_to_strip`-defined prefixes from messages that mention the bot user.
if !params.bot_user_prefixes_to_strip.is_empty()
&& message.mentioned_users.contains(&params.bot_user_id)
{
let mut message_text = message.message_text.clone();
for prefix in &params.bot_user_prefixes_to_strip {
if let Some(message_text_stripped) = message_text.strip_prefix(prefix) {
message_text = message_text_stripped.to_owned();
}
@@ -73,16 +180,25 @@ pub async fn process_matrix_messages_in_thread(
messages_filtered
}
/// Tells if the given message is from an allowed sender.
///
/// If allowed_users is None, all messages are allowed.
/// If allowed_users is Some, only messages from the allowed users (and the `bot_user_id`) are allowed.
fn is_message_from_allowed_sender(
matrix_message: &MatrixMessage,
bot_user_id: &str,
allowed_users: &[regex::Regex],
bot_user_id: &OwnedUserId,
allowed_users: Option<&[regex::Regex]>,
) -> bool {
if matrix_message.sender_id == bot_user_id {
if matrix_message.sender_id == *bot_user_id {
return true;
}
if mxidwc::match_user_id(&matrix_message.sender_id, allowed_users) {
if let Some(allowed_users) = allowed_users {
if mxidwc::match_user_id(matrix_message.sender_id.as_str(), allowed_users) {
return true;
}
} else {
// No allowed users configured, so all messages are allowed
return true;
}
@@ -108,30 +224,58 @@ pub fn convert_matrix_native_event_to_matrix_message(
_ => return None,
};
let is_reply = matches!(room_message.relates_to, Some(Relation::Reply { .. }));
let text = if is_reply {
// For regular replies, we need to strip the fallback-for-rich replies part.
// See: https://spec.matrix.org/v1.11/client-server-api/#fallbacks-for-rich-replies
strip_rich_reply_fallback_text(&text)
} else {
text
};
let mentioned_users = room_message
.mentions
.map(|m| m.user_ids.iter().map(|u| u.to_owned()).collect())
.unwrap_or(vec![]);
Some(MatrixMessage {
sender_id: matrix_native_event.sender().to_string(),
sender_id: matrix_native_event.sender().to_owned(),
message_type: if is_notice {
MatrixMessageType::Notice
} else {
MatrixMessageType::Text
},
message_text: text,
mentioned_users,
})
}
/// Determines the thread context (relationship within the thread + first thread message payload) for an incoming (new) room event.
/// This room event is assumed to be the "newest message" in the thread (or a top-level message).
/// If the given event is a regular reply (not a thread reply), this function will return `None`.
/// If the given event is a top-level message, this function will consider this event as the start of the thread.
/// If the given event is a thread reply, this function will inspect the thread root event and will return the thread context.
/// If the thread root event is not found, is redacted, or is of some unsupported MessagePayload type, this function will return `None`.
pub async fn determine_thread_context_for_room_event(
/// Determines the interaction context for an incoming (new) room event.
///
/// This context is created based on the "newest message" (`current_event`), which is:
/// - either a top-level message, which may or may not be mentioning the bot
/// - this function will inspect the event and will likely start a new threaded conversation
///
/// - or a thread reply
/// - this function will inspect the thread root event and will return the interaction context
/// - if the bot only reacts to prefixed messsages (or mentions), this function may ignore the given thread reply, unless it mentions the bot (which causes a synthetic "first message" to be produced)
/// - if the thread root event is not found, is redacted, or is of some unsupported MessagePayload type, this function will return `None`
///
/// - or an in-room (non-threaded) reply to a room message, which may or may not be mentioning the bot
/// - replies that do not mention the bot cause this function to return `None`
/// - other replies create a interaction context which points to a "first message" which is synthetic
#[tracing::instrument(name = "determine_interaction_context_for_room_event", skip_all, fields(room_id = room.room_id().as_str(), event_id = current_event.event_id.as_str()))]
pub async fn determine_interaction_context_for_room_event(
bot_user_id: &OwnedUserId,
room: &Room,
current_event: &OriginalSyncRoomMessageEvent,
current_event_payload: &MessagePayload,
event_fetcher: &Arc<RoomEventFetcher>,
) -> anyhow::Result<Option<ThreadContext>> {
) -> anyhow::Result<Option<InteractionContext>> {
let current_event_is_mentioning_bot =
is_event_mentioning_bot(&current_event.content, bot_user_id);
let Some(relation) = &current_event.content.relates_to else {
// This is a top-level message. We consider it the start of the thread.
let thread_info = ThreadInfo::new(
@@ -139,25 +283,73 @@ pub async fn determine_thread_context_for_room_event(
current_event.event_id.clone(),
);
let is_mentioning_bot = is_event_mentioning_bot(&current_event.content, bot_user_id);
return Ok(Some(ThreadContext {
info: thread_info,
first_message: ThreadContextFirstMessage {
is_mentioning_bot,
return Ok(Some(InteractionContext {
thread_info,
trigger: InteractionTrigger {
is_mentioning_bot: current_event_is_mentioning_bot,
payload: current_event_payload.clone(),
},
}));
};
let Relation::Thread(thread) = relation else {
// This is a reply or a replacement, etc. It's not a thread.
// We don't care about this.
return Ok(None);
};
match relation {
Relation::Thread(thread) => {
determine_interaction_context_for_room_event_related_to_thread(
bot_user_id,
room,
current_event,
event_fetcher,
current_event_is_mentioning_bot,
thread,
)
.await
}
Relation::Reply { in_reply_to } => {
determine_interaction_context_for_room_event_related_to_reply(
current_event,
current_event_is_mentioning_bot,
in_reply_to.event_id.clone(),
)
.await
}
// This is a replacement or something else. It's not something we support.
_ => return Ok(None),
}
}
async fn determine_interaction_context_for_room_event_related_to_thread(
bot_user_id: &OwnedUserId,
room: &Room,
current_event: &OriginalSyncRoomMessageEvent,
event_fetcher: &Arc<RoomEventFetcher>,
current_event_is_mentioning_bot: bool,
thread: &Thread,
) -> anyhow::Result<Option<InteractionContext>> {
let thread_info = ThreadInfo::new(thread.event_id.clone(), current_event.event_id.clone());
tracing::trace!(
?current_event_is_mentioning_bot,
is_thread_root_only = thread_info.is_thread_root_only(),
"Dealing with a thread reply",
);
if current_event_is_mentioning_bot && !thread_info.is_thread_root_only() {
// If the current event is a thread reply and is mentioning the bot,
// it's probably someone trying to involve us in the threaded conversation.
// See: https://github.com/etkecc/baibot/issues/15
//
// In such cases, we don't care what the thread root event is like or what the current event is like,
// we want text-generation to be triggered for this whole thread regardless.
return Ok(Some(InteractionContext {
thread_info,
trigger: InteractionTrigger {
is_mentioning_bot: true,
payload: MessagePayload::SynthethicChatCompletionTriggerInThread,
},
}));
}
let start_time = std::time::Instant::now();
let thread_start_timeline_event = event_fetcher
@@ -183,82 +375,46 @@ pub async fn determine_thread_context_for_room_event(
"Fetched thread start event"
);
let thread_start_timeline_event_deserialized =
match thread_start_timeline_event.event.deserialize() {
Ok(value) => value,
Err(err) => {
return Err(anyhow::format_err!(
"Failed to deserialize thread start event {}: {:?}",
thread.event_id,
err
));
}
};
let thread_start_detailed_message_payload = timeline_event_to_detailed_message_payload(
&thread.event_id,
thread_start_timeline_event,
thread_info.clone(),
bot_user_id,
)?;
let AnyTimelineEvent::MessageLike(thread_start_message_like_event) =
thread_start_timeline_event_deserialized
else {
tracing::trace!(
"Ignoring non-MessageLike thread start event: {:?}",
thread_start_timeline_event_deserialized
);
let Some(detailed_message_payload) = thread_start_detailed_message_payload else {
return Ok(None);
};
let (thread_start_message_is_mentioning_bot, thread_start_message_payload) =
match thread_start_message_like_event {
AnyMessageLikeEvent::RoomEncrypted(room_message) => {
tracing::warn!(
"Could not inspect thread start event {} because it failed to decrypt: {:?}",
thread.event_id.clone(),
room_message
);
Ok(Some(InteractionContext {
thread_info,
trigger: InteractionTrigger {
is_mentioning_bot: detailed_message_payload.is_mentioning_bot,
payload: detailed_message_payload.message_payload,
},
}))
}
// There's no way to know and it doesn't matter anyway.
let is_mentioning_bot = false;
async fn determine_interaction_context_for_room_event_related_to_reply(
current_event: &OriginalSyncRoomMessageEvent,
current_event_is_mentioning_bot: bool,
reply_to_event_id: OwnedEventId,
) -> anyhow::Result<Option<InteractionContext>> {
tracing::trace!(?current_event_is_mentioning_bot, "Dealing with a reply");
(
is_mentioning_bot,
MessagePayload::Encrypted(thread_info.clone()),
)
}
AnyMessageLikeEvent::RoomMessage(room_message) => {
if let MessageLikeEvent::Original(room_message_original) = room_message {
let room_message_payload: Result<MessagePayload, String> =
room_message_original.content.msgtype.clone().try_into();
if !current_event_is_mentioning_bot {
// If the current event is not mentioning the bot, we don't care about it.
tracing::trace!("Ignoring reply event which does not mention the bot");
return Ok(None);
}
let Ok(room_message_payload) = room_message_payload else {
tracing::debug!(
msg_type = room_message_original.content.msgtype(),
"Ignoring thread start message of unknown type",
);
return Ok(None);
};
let thread_info = ThreadInfo::new(reply_to_event_id.clone(), current_event.event_id.clone());
let is_mentioning_bot =
is_event_mentioning_bot(&room_message_original.content, bot_user_id);
(is_mentioning_bot, room_message_payload)
} else {
tracing::error!("Ignoring thread start message which appears to be redacted");
return Ok(None);
}
}
other => {
tracing::trace!(
"Ignoring unknown MessageLike thread start event: {:?}",
other
);
return Ok(None);
}
};
Ok(Some(ThreadContext {
info: thread_info,
first_message: ThreadContextFirstMessage {
is_mentioning_bot: thread_start_message_is_mentioning_bot,
payload: thread_start_message_payload,
Ok(Some(InteractionContext {
thread_info,
trigger: InteractionTrigger {
is_mentioning_bot: true,
payload: MessagePayload::SynthethicChatCompletionTriggerForReply,
},
}))
}
@@ -267,21 +423,157 @@ fn is_event_mentioning_bot(
event_content: &RoomMessageEventContent,
bot_user_id: &OwnedUserId,
) -> bool {
if let Some(mentions) = &event_content.mentions {
mentions
.user_ids
.iter()
.any(|user_id| user_id == bot_user_id)
} else {
// For compatibility with clients that do not support the new Mentions specification
// (see https://spec.matrix.org/latest/client-server-api/#user-and-room-mentions),
// we also do string matching here.
//
// It may be even better to match not only against the MXID, but also against the bot's
// room-specific display name.
//
// We may consider dropping this string-matching behavior altogether in the future,
// so improving this compatibility block is not a high priority.
event_content.body().contains(bot_user_id.as_str())
}
// As a fallback, we used to do string matching (`event_content.body().contains(bot_user_id.as_str())`) here as well.
// However, this is unreliable. In 2024+, clients that do not have proper mentions support should get fixed,
// instead of us having to deal with the possibility of false positives.
//
let Some(mentions) = &event_content.mentions else {
return false;
};
mentions
.user_ids
.iter()
.any(|user_id| user_id == bot_user_id)
}
/// Strips the rich reply fallback text from the given text.
/// See: https://spec.matrix.org/v1.11/client-server-api/#fallbacks-for-rich-replies
///
/// Example:
/// ```rust,ignore
/// let text = "> <@admin:example.com> What's the difference between Matrix and XMPP?\n\nAnswer me";
/// let stripped_text = strip_rich_reply_fallback_text(text);
/// assert_eq!(stripped_text, "Answer me");
/// ```
fn strip_rich_reply_fallback_text(text: &str) -> String {
let lines = text.lines();
let mut stripped_lines = Vec::new();
let mut encountered_non_prefix = false;
for line in lines {
if !encountered_non_prefix && line.starts_with("> ") {
continue;
} else {
encountered_non_prefix = true;
stripped_lines.push(line);
}
}
stripped_lines.join("\n").trim().to_owned()
}
fn timeline_event_to_detailed_message_payload(
timeline_event_id: &OwnedEventId,
timeline_event: TimelineEvent,
thread_info: ThreadInfo,
bot_user_id: &OwnedUserId,
) -> anyhow::Result<Option<DetailedMessagePayload>> {
let timeline_event_deserialized = match timeline_event.event.deserialize() {
Ok(value) => value,
Err(err) => {
return Err(anyhow::format_err!(
"Failed to deserialize timeline event {}: {:?}",
timeline_event_id,
err
));
}
};
let AnyTimelineEvent::MessageLike(thread_start_message_like_event) =
timeline_event_deserialized
else {
tracing::trace!(
"Ignoring non-MessageLike timeline event: {:?}",
timeline_event_deserialized
);
return Ok(None);
};
let (is_mentioning_bot, message_payload) = match thread_start_message_like_event {
AnyMessageLikeEvent::RoomEncrypted(room_message) => {
tracing::warn!(
"Could not inspect event {} because it failed to decrypt: {:?}",
timeline_event_id.clone(),
room_message
);
// There's no way to know and it doesn't matter anyway.
let is_mentioning_bot = false;
(
is_mentioning_bot,
MessagePayload::Encrypted(thread_info.clone()),
)
}
AnyMessageLikeEvent::RoomMessage(room_message) => {
if let MessageLikeEvent::Original(room_message_original) = room_message {
let room_message_payload: Result<MessagePayload, String> =
room_message_original.content.msgtype.clone().try_into();
let Ok(room_message_payload) = room_message_payload else {
tracing::debug!(
msg_type = room_message_original.content.msgtype(),
"Ignoring event message of unknown type",
);
return Ok(None);
};
let is_mentioning_bot =
is_event_mentioning_bot(&room_message_original.content, bot_user_id);
(is_mentioning_bot, room_message_payload)
} else {
tracing::error!("Ignoring event message which appears to be redacted");
return Ok(None);
}
}
other => {
tracing::trace!("Ignoring unknown MessageLike event: {:?}", other);
return Ok(None);
}
};
Ok(Some(DetailedMessagePayload {
is_mentioning_bot,
message_payload,
}))
}
/// Creates a list of prefixes to strip from the beginning of message texts that mention the bot user.
///
/// Different clients do mentions differently.
/// The body text containing the mention usually contains one of:
/// - the full user ID (includes a @ prefix by default)
/// - the localpart (with a @ prefix)
/// - the localpart (without a @ prefix)
/// - the display name (with a @ prefix)
/// - the display name (without a @ prefix)
///
/// Some add a `: ` suffix after the mention.
///
/// There's no guarantee that the mention is at the start even.
/// It being there is most common and we try to strip it from there
/// as best as we can.
pub fn create_list_of_bot_user_prefixes_to_strip(
bot_user_id: &OwnedUserId,
bot_display_name: &Option<String>,
) -> Vec<String> {
let bot_user_id_localpart = bot_user_id.localpart();
let mut prefixes_to_strip = vec![
bot_user_id.as_str().to_owned(),
format!("@{}", bot_user_id_localpart),
bot_user_id_localpart.to_owned(),
];
if let Some(bot_display_name) = bot_display_name {
prefixes_to_strip.push(format!("@{}", bot_display_name));
prefixes_to_strip.push(bot_display_name.to_owned());
}
prefixes_to_strip.push(":".to_owned());
prefixes_to_strip
}