Skip to main content

mail4agent_server/
messaging.rs

1//! Timeline send, redact, and history paging over an open connection.
2//! Entitlements are the builder's. A private room does not consult a paid flag.
3
4use std::collections::HashSet;
5
6use rusqlite::Connection;
7
8use crate::error::MatrixError;
9use crate::events::client_event_json;
10use crate::store::{Membership, PowerAction, Room, TxnDedupEntry};
11
12// ============================================================================
13// Rate limits (plan §3.8) — own copies per this codebase's
14// one-copy-per-route-module convention (mirrors `dm_db::DM_MESSAGE_DAILY_LIMIT`
15// via `routes::dm::create_message`'s own daily check).
16// ============================================================================
17
18/// Per-sender messages-per-day cap in private rooms, mirroring
19/// [`dm_db`](crate::dm_db)'s own `DM_MESSAGE_DAILY_LIMIT` value.
20pub const MATRIX_SEND_DAILY_LIMIT: i64 = 300;
21
22/// Short-window burst cap: at most this many sends/redacts per device.
23pub const BURST_LIMIT_COUNT: i64 = 20;
24
25/// ...within this many seconds.
26pub const BURST_LIMIT_WINDOW_SECS: i64 = 10;
27
28
29pub const ONE_DAY_MS: u64 = 24 * 60 * 60 * 1000;
30
31
32pub fn check_rate_limits(conn: &Connection, user_id: i64, device_id: &str, now: &str) -> Result<(), MatrixError> {
33    let now_dt = chrono::DateTime::parse_from_rfc3339(now).map_err(|_| MatrixError::internal())?;
34
35    let since_ms = (now_dt - chrono::Duration::hours(24)).timestamp_millis();
36    let daily = crate::store::private_room_messages_sent_since(conn, user_id, since_ms)?;
37    if daily >= MATRIX_SEND_DAILY_LIMIT {
38        return Err(MatrixError::limit_exceeded(ONE_DAY_MS));
39    }
40
41    let burst_since = (now_dt - chrono::Duration::seconds(BURST_LIMIT_WINDOW_SECS)).to_rfc3339();
42    let burst = crate::store::txn_dedup_count_since(conn, user_id, device_id, &burst_since)?;
43    if burst >= BURST_LIMIT_COUNT {
44        return Err(MatrixError::limit_exceeded(BURST_LIMIT_WINDOW_SECS as u64 * 1000));
45    }
46    Ok(())
47}
48
49
50// ============================================================================
51// Pagination tokens and limits — shared by `/messages` and `/relations`
52// ============================================================================
53
54pub const MESSAGES_DEFAULT_LIMIT: i64 = 10;
55
56pub const MESSAGES_MAX_LIMIT: i64 = 100;
57
58
59/// A parsed pagination position — see the module doc.
60///
61/// A `t{stream_id}` token (or a bare integer) sits ON event `stream_id`:
62/// a backward page starts strictly before it, a forward page strictly after
63/// it. A sync token `s{stream_id}_{typing_gen}` (a `/sync` `next_batch` or
64/// `prev_batch`) sits AFTER event `stream_id` — that event was already
65/// delivered — so a backward page from it must still include that event. The
66/// typing half of a sync token is irrelevant to history and is ignored.
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68pub struct StreamToken {
69    pub stream_id: i64,
70    pub after_event: bool,
71}
72
73
74impl StreamToken {
75    /// The exclusive upper bound of a page that walks backward from this
76    /// position: events with `stream_id < bound`.
77    pub fn backward_bound(self) -> i64 {
78        if self.after_event {
79            self.stream_id.saturating_add(1)
80        } else {
81            self.stream_id
82        }
83    }
84
85    /// The exclusive lower bound of a page that walks forward from this
86    /// position: events with `stream_id > bound`.
87    pub fn forward_bound(self) -> i64 {
88        self.stream_id
89    }
90
91    /// The `to=` boundary of a page walking `dir`: the far end a page must
92    /// not cross. It is the position's bound as seen from the opposite walk
93    /// direction, since `to` is approached rather than left.
94    pub fn to_bound(self, dir: Direction) -> i64 {
95        match dir {
96            Direction::Backward => self.forward_bound(),
97            Direction::Forward => self.backward_bound(),
98        }
99    }
100}
101
102
103/// `"t{stream_id}"`, a bare `"0"`, or a sync token `"s{stream_id}_{typing_gen}"`
104/// — see the module doc. Anything else, and any negative stream id, is
105/// refused with `M_INVALID_PARAM`.
106pub fn parse_stream_token(raw: &str) -> Result<StreamToken, MatrixError> {
107    let malformed = || MatrixError::invalid_param("malformed pagination token");
108    let token = if raw.starts_with('s') {
109        let sync = crate::sync_token::parse(raw).map_err(|_| malformed())?;
110        StreamToken { stream_id: sync.stream_id, after_event: true }
111    } else {
112        let digits = raw.strip_prefix('t').unwrap_or(raw);
113        StreamToken { stream_id: digits.parse::<i64>().map_err(|_| malformed())?, after_event: false }
114    };
115    if token.stream_id < 0 {
116        return Err(malformed());
117    }
118    Ok(token)
119}
120
121
122/// The page start `GET /messages` walks from: the explicit `from` token when
123/// present, else the far end of the timeline for `dir` — the newest event
124/// going backward, the very beginning going forward (Matrix makes `from`
125/// optional and defines exactly this).
126pub fn resolve_messages_start(conn: &Connection, from: Option<StreamToken>, dir: Direction) -> rusqlite::Result<i64> {
127    Ok(match (from, dir) {
128        (Some(token), Direction::Backward) => token.backward_bound(),
129        (Some(token), Direction::Forward) => token.forward_bound(),
130        (None, Direction::Backward) => crate::store::max_stream_id(conn)?.saturating_add(1),
131        (None, Direction::Forward) => 0,
132    })
133}
134
135
136pub fn format_stream_token(stream_id: i64) -> String {
137    format!("t{stream_id}")
138}
139
140
141pub fn clamp_limit(requested: Option<i64>) -> i64 {
142    requested.unwrap_or(MESSAGES_DEFAULT_LIMIT).clamp(1, MESSAGES_MAX_LIMIT)
143}
144
145
146#[derive(Debug, Clone, Copy, PartialEq, Eq)]
147pub enum Direction {
148    Backward,
149    Forward,
150}
151
152
153impl Direction {
154    pub fn from_query(raw: Option<&str>) -> Result<Self, MatrixError> {
155        match raw.unwrap_or("b") {
156            "b" => Ok(Direction::Backward),
157            "f" => Ok(Direction::Forward),
158            _ => Err(MatrixError::invalid_param("dir must be 'b' or 'f'")),
159        }
160    }
161}
162
163
164#[derive(serde::Deserialize, Default, Clone)]
165pub struct EventFilter {
166    #[serde(default)]
167    pub types: Option<Vec<String>>,
168    #[serde(default)]
169    pub not_types: Option<Vec<String>>,
170    #[serde(default)]
171    pub lazy_load_members: bool,
172}
173
174
175pub fn event_passes_filter(event_type: &str, filter: &EventFilter) -> bool {
176    if let Some(types) = &filter.types {
177        if !types.iter().any(|t| t == event_type) {
178            return false;
179        }
180    }
181    if let Some(not_types) = &filter.not_types {
182        if not_types.iter().any(|t| t == event_type) {
183            return false;
184        }
185    }
186    true
187}
188
189
190pub fn parse_filter(raw: Option<&str>) -> Result<EventFilter, MatrixError> {
191    match raw {
192        Some(raw) => serde_json::from_str(raw).map_err(|_| MatrixError::invalid_param("malformed filter")),
193        None => Ok(EventFilter::default()),
194    }
195}
196
197
198// ============================================================================
199// DB-only gate/read helpers — own copies of `routes::matrix::rooms`'s small
200// checks (per this codebase's one-copy-per-route-module convention); only
201// `member_and_invited_ids` (the wake fan-out set) is reused verbatim — see
202// the module doc.
203// ============================================================================
204
205pub fn power_levels_of(conn: &Connection, room_id: &str) -> Result<serde_json::Value, MatrixError> {
206    match crate::store::current_state_event(conn, room_id, "m.room.power_levels", "")? {
207        Some(event) => Ok(serde_json::from_str(&event.content)?),
208        None => Ok(serde_json::json!({})),
209    }
210}
211
212
213pub fn require_member(membership: Option<Membership>) -> Result<(), MatrixError> {
214    match membership {
215        Some(Membership::Join) => Ok(()),
216        _ => Err(MatrixError::forbidden("not a member of this room")),
217    }
218}
219
220
221pub fn require_power(power_levels: &serde_json::Value, mxid: &str, action: PowerAction) -> Result<(), MatrixError> {
222    if crate::store::can(power_levels, action, mxid) {
223        Ok(())
224    } else {
225        Err(MatrixError::forbidden("insufficient power level for this action"))
226    }
227}
228
229
230// ============================================================================
231// Pure decision functions — no DB, no `MatrixCaller`, unit-tested directly
232// ============================================================================
233
234/// Every state-event type `PUT /state/{eventType}/{stateKey}` owns, refused
235/// outright on `/send` (P6 binding rule) — includes our own
236/// `org.example.legacy_dm_key` alongside the standard `m.room.*` state
237/// types.
238pub const REFUSED_SEND_STATE_TYPES: [&str;
239 10] = [
240    "m.room.create",
241    "m.room.member",
242    "m.room.power_levels",
243    "m.room.join_rules",
244    "m.room.history_visibility",
245    "m.room.name",
246    "m.room.topic",
247    "m.room.avatar",
248    "m.room.encryption",
249    "m.room.pinned_events",
250];
251pub const LEGACY_DM_KEY_STATE_TYPE: &str = "org.example.legacy_dm_key";
252
253
254/// `PUT /send/{eventType}/{txnId}`'s type-refusal rule (P6 binding rule):
255/// state-event types and `org.example.legacy_dm_key` go through
256/// `/state`; `org.example.legacy_dm` is migration-only; `m.room.redaction`
257/// goes through `/redact`; and, in an encrypted room, only `m.room.encrypted`
258/// and `m.reaction` (server-visible metadata, coordinator ruling) are
259/// accepted at all.
260pub fn check_send_event_type_allowed(event_type: &str, room_is_encrypted: bool) -> Result<(), MatrixError> {
261    if REFUSED_SEND_STATE_TYPES.contains(&event_type) || event_type == LEGACY_DM_KEY_STATE_TYPE {
262        return Err(MatrixError::bad_json(format!("{event_type} is a state event; send it via PUT /state/{{eventType}}/{{stateKey}}")));
263    }
264    if event_type == "org.example.legacy_dm" {
265        return Err(MatrixError::forbidden("org.example.legacy_dm is migration-only and cannot be sent by a client"));
266    }
267    if event_type == "m.room.redaction" {
268        return Err(MatrixError::bad_json("use PUT /rooms/{roomId}/redact/{eventId}/{txnId} to redact an event"));
269    }
270    if room_is_encrypted && event_type != "m.room.encrypted" && event_type != "m.reaction" {
271        return Err(MatrixError::forbidden("only m.room.encrypted and m.reaction may be sent in an encrypted room"));
272    }
273    Ok(())
274}
275
276
277/// An `m.replace` must come from the target event's ORIGINAL sender (P6
278/// binding rule) — a no-op for content with no `m.relates_to`, or a
279/// `rel_type` other than `m.replace`.
280pub fn check_replace_target_sender(conn: &Connection, content: &serde_json::Value, sender_user_id: i64) -> Result<(), MatrixError> {
281    let Some(relates_to) = content.get("m.relates_to") else { return Ok(()) };
282    if relates_to.get("rel_type").and_then(|v| v.as_str()) != Some("m.replace") {
283        return Ok(());
284    }
285    let Some(target_id) = relates_to.get("event_id").and_then(|v| v.as_str()) else { return Ok(()) };
286    let target = crate::store::get_event(conn, target_id)?.ok_or_else(|| MatrixError::not_found("relation target not found"))?;
287    if target.sender_user_id != sender_user_id {
288        return Err(MatrixError::invalid_param("m.replace must be sent by the target event's original sender"));
289    }
290    Ok(())
291}
292
293
294// ============================================================================
295// DB-only action cores — `&mut Connection`, no `MatrixCaller`, unit-tested
296// directly against an in-memory `matrix_store` fixture
297// ============================================================================
298
299#[derive(Debug, Clone, PartialEq)]
300pub struct SendOutcome {
301    pub event: crate::store::MatrixEvent,
302    pub is_new: bool,
303    pub wake_ids: HashSet<i64>,
304}
305
306
307/// `PUT /send/{eventType}/{txnId}`'s whole DB-side decision + write (plan §4
308/// `send` row): txn-dedup peek first (a repeat short-circuits everything
309/// else and returns the ORIGINAL event, per the idempotency rule); then
310/// Member, the type-refusal rule, PowerCheck(`events[type]` else
311/// `events_default`), the
312/// send-rate limits, and the `m.replace`-sender rule; then the dedup-checked
313/// insert itself.
314#[allow(clippy::too_many_arguments)]
315pub fn apply_send(
316    conn: &mut Connection,
317    room: &Room,
318    caller_user_id: i64,
319    caller_mxid: &str,
320    device_id: &str,
321    txn_id: &str,
322    event_id: &str,
323    event_type: &str,
324    content_str: &str,
325    now: &str,
326    origin_ts: i64,
327) -> Result<SendOutcome, MatrixError> {
328    if let TxnDedupEntry::Seen(existing_event_id) = crate::store::txn_dedup_lookup(conn, caller_user_id, device_id, txn_id)? {
329        let existing_event_id = existing_event_id.ok_or_else(MatrixError::internal)?;
330        let event = crate::store::get_event(conn, &existing_event_id)?.ok_or_else(MatrixError::internal)?;
331        return Ok(SendOutcome { event, is_new: false, wake_ids: HashSet::new() });
332    }
333
334    if let Some(existing) = crate::public_channels::seen_txn(conn, caller_user_id, device_id, txn_id).map_err(|_| MatrixError::internal())? {
335        return Ok(SendOutcome { event: existing, is_new: false, wake_ids: HashSet::new() });
336    }
337
338    let caller_membership = crate::store::room_member(conn, &room.id, caller_user_id)?.map(|m| m.membership);
339    require_member(caller_membership)?;
340    check_send_event_type_allowed(event_type, room.is_encrypted)?;
341
342    let power_levels = power_levels_of(conn, &room.id)?;
343    if crate::store::user_level(&power_levels, caller_mxid) < crate::store::event_level(&power_levels, event_type, false) {
344        return Err(MatrixError::forbidden("insufficient power level to send this event"));
345    }
346
347    check_rate_limits(conn, caller_user_id, device_id, now)?;
348
349    let content: serde_json::Value = serde_json::from_str(content_str)?;
350    check_replace_target_sender(conn, &content, caller_user_id)?;
351
352    if crate::public_channels::is_public_room(conn, &room.id)? {
353        // Plaintext channel: the post goes to the public store, never to `events`.
354        return match crate::public_channels::insert_event_deduped(conn, device_id, txn_id, event_id, &room.id, caller_user_id, event_type, content_str, origin_ts)
355            .map_err(|_| MatrixError::internal())?
356        {
357            crate::public_channels::PublicWrite::New(event) => {
358                let wake_ids = crate::rooms::member_and_invited_ids(conn, &room.id)?;
359                Ok(SendOutcome { event, is_new: true, wake_ids })
360            }
361            crate::public_channels::PublicWrite::Existing(event) => Ok(SendOutcome { event, is_new: false, wake_ids: HashSet::new() }),
362        };
363    }
364
365    match crate::store::insert_timeline_event_deduped(conn, device_id, txn_id, event_id, &room.id, caller_user_id, event_type, content_str, origin_ts, now)? {
366        crate::store::DedupedWrite::New(event) => {
367            let wake_ids = crate::rooms::member_and_invited_ids(conn, &room.id)?;
368            Ok(SendOutcome { event, is_new: true, wake_ids })
369        }
370        // Cannot happen — the peek above already established `NotSeen`
371        // under the same connection, and this module's single-writer
372        // discipline means nothing else can have raced it — but handled
373        // defensively rather than panicking.
374        crate::store::DedupedWrite::Existing(event) => Ok(SendOutcome { event, is_new: false, wake_ids: HashSet::new() }),
375    }
376}
377
378
379#[derive(Debug, Clone, PartialEq)]
380pub struct RedactOutcome {
381    pub event: crate::store::MatrixEvent,
382    pub is_new: bool,
383    pub wake_ids: HashSet<i64>,
384}
385
386
387/// `PUT /redact/{eventId}/{txnId}`'s whole DB-side decision + write (plan §4
388/// `redact` row): txn-dedup peek first; then Member; then own-event OR
389/// PowerCheck(redact) — redact does NOT require the sender's level to
390/// exceed the target's own level (spec rule, unlike kick/ban/unban); then
391/// the dedup-checked redaction itself (which enforces the v11
392/// unredactable-event-type rule and maps to 403 via
393/// `MatrixStoreError::UnredactableEvent`).
394#[allow(clippy::too_many_arguments)]
395pub fn apply_redact(
396    conn: &mut Connection,
397    room_id: &str,
398    caller_user_id: i64,
399    caller_mxid: &str,
400    device_id: &str,
401    txn_id: &str,
402    target_event_id: &str,
403    redaction_event_id: &str,
404    reason: Option<&str>,
405    now: &str,
406    origin_ts: i64,
407) -> Result<RedactOutcome, MatrixError> {
408    if let TxnDedupEntry::Seen(existing_event_id) = crate::store::txn_dedup_lookup(conn, caller_user_id, device_id, txn_id)? {
409        let existing_event_id = existing_event_id.ok_or_else(MatrixError::internal)?;
410        let event = crate::store::get_event(conn, &existing_event_id)?.ok_or_else(MatrixError::internal)?;
411        return Ok(RedactOutcome { event, is_new: false, wake_ids: HashSet::new() });
412    }
413
414    let caller_membership = crate::store::room_member(conn, room_id, caller_user_id)?.map(|m| m.membership);
415    require_member(caller_membership)?;
416
417    let target = crate::store::get_event(conn, target_event_id)?.ok_or_else(|| MatrixError::not_found("no such event"))?;
418    if target.room_id != room_id {
419        return Err(MatrixError::not_found("no such event"));
420    }
421    if target.sender_user_id != caller_user_id {
422        let power_levels = power_levels_of(conn, room_id)?;
423        require_power(&power_levels, caller_mxid, PowerAction::Redact)?;
424    }
425
426    if crate::public_channels::get_event(conn, target_event_id).map_err(|_| MatrixError::internal())?.is_some() {
427        return match crate::public_channels::redact_deduped(conn, device_id, txn_id, room_id, target_event_id, redaction_event_id, caller_user_id, reason, origin_ts)
428            .map_err(|_| MatrixError::internal())?
429        {
430            crate::public_channels::PublicWrite::New(event) => {
431                let wake_ids = crate::rooms::member_and_invited_ids(conn, room_id)?;
432                Ok(RedactOutcome { event, is_new: true, wake_ids })
433            }
434            crate::public_channels::PublicWrite::Existing(event) => Ok(RedactOutcome { event, is_new: false, wake_ids: HashSet::new() }),
435        };
436    }
437
438    match crate::store::redact_event_deduped(conn, device_id, txn_id, room_id, target_event_id, redaction_event_id, caller_user_id, reason, origin_ts, now)? {
439        crate::store::DedupedWrite::New(event) => {
440            let wake_ids = crate::rooms::member_and_invited_ids(conn, room_id)?;
441            Ok(RedactOutcome { event, is_new: true, wake_ids })
442        }
443        crate::store::DedupedWrite::Existing(event) => Ok(RedactOutcome { event, is_new: false, wake_ids: HashSet::new() }),
444    }
445}
446
447
448/// One page of `/messages`' timeline window — [`get_messages`]'s own DB-only
449/// core, unit-tested directly.
450pub struct MessagesPage {
451    pub chunk: Vec<crate::store::MatrixEvent>,
452    pub start: String,
453    /// `None` when there is nothing further this caller could ever see
454    /// beyond this page (the physical edge of the room, the caller's own
455    /// [`crate::store::HistoryWindow`], or an explicit `to=` boundary) —
456    /// NEVER `None` merely because a `types`/`not_types` filter thinned this
457    /// particular page (more matching events could still exist further on).
458    pub end: Option<String>,
459}
460
461
462#[allow(clippy::too_many_arguments)]
463pub fn paginate_messages(
464    conn: &Connection,
465    room_id: &str,
466    window: crate::store::HistoryWindow,
467    from_stream: i64,
468    to_stream: Option<i64>,
469    dir: Direction,
470    limit: i64,
471    filter: &EventFilter,
472) -> rusqlite::Result<MessagesPage> {
473    let raw_page = match dir {
474        Direction::Backward => crate::store::events_in_room_before(conn, room_id, from_stream, limit)?,
475        Direction::Forward => crate::store::events_in_room_after(conn, room_id, from_stream, limit)?,
476    };
477    let raw_len = raw_page.len();
478
479    let hard_truncated = |e: &crate::store::MatrixEvent| -> bool {
480        if !window.contains(e.stream_id) {
481            return true;
482        }
483        match (to_stream, dir) {
484            (Some(to), Direction::Backward) => e.stream_id <= to,
485            (Some(to), Direction::Forward) => e.stream_id >= to,
486            (None, _) => false,
487        }
488    };
489    let any_hard_truncated = raw_page.iter().any(hard_truncated);
490    let chunk: Vec<crate::store::MatrixEvent> =
491        raw_page.iter().filter(|e| !hard_truncated(e) && event_passes_filter(&e.event_type, filter)).cloned().collect();
492
493    let end = if raw_len < limit as usize || any_hard_truncated {
494        None
495    } else {
496        raw_page.last().map(|e| format_stream_token(e.stream_id))
497    };
498
499    Ok(MessagesPage { chunk, start: format_stream_token(from_stream), end })
500}
501
502
503pub fn lazy_load_member_state(conn: &Connection, room_id: &str, chunk: &[crate::store::MatrixEvent]) -> Result<Vec<serde_json::Value>, MatrixError> {
504    let mut senders: Vec<i64> = chunk.iter().map(|e| e.sender_user_id).collect();
505    senders.sort_unstable();
506    senders.dedup();
507
508    let mut out = Vec::new();
509    for sender in senders {
510        let Some(mxid) = crate::store::mxid_of(conn, sender)? else { continue };
511        if let Some(event) = crate::store::current_state_event(conn, room_id, "m.room.member", &mxid)? {
512            out.push(client_event_json(conn, &event, None)?);
513        }
514    }
515    Ok(out)
516}
517
518
519// ============================================================================
520// PUT /_matrix/client/v3/rooms/{roomId}/redact/{eventId}/{txnId}
521// ============================================================================
522
523#[derive(serde::Deserialize, Default)]
524pub struct RedactBody {
525    #[serde(default)]
526    pub reason: Option<String>,
527}
528
529
530// ============================================================================
531// GET /_matrix/client/v3/rooms/{roomId}/messages
532// ============================================================================
533
534#[derive(serde::Deserialize)]
535pub struct MessagesQuery {
536    /// Optional (Matrix v1.3+): absent means "from the newest event" for
537    /// `dir=b` and "from the beginning" for `dir=f`.
538    #[serde(default)]
539    pub from: Option<String>,
540    #[serde(default)]
541    pub to: Option<String>,
542    #[serde(default)]
543    pub dir: Option<String>,
544    #[serde(default)]
545    pub limit: Option<i64>,
546    #[serde(default)]
547    pub filter: Option<String>,
548}
549
550
551// ============================================================================
552// GET /_matrix/client/v3/rooms/{roomId}/relations/{eventId}[/{relType}[/{eventType}]]
553// ============================================================================
554
555#[derive(serde::Deserialize, Default)]
556pub struct RelationsQuery {
557    #[serde(default)]
558    pub from: Option<String>,
559    #[serde(default)]
560    pub limit: Option<i64>,
561}