Skip to main content

mail4agent_server/
sync.rs

1//! `/sync` snapshot. Returns JSON. Does not wait, poll, or wake.
2
3use std::time::Instant;
4
5use rusqlite::Connection;
6
7use crate::error::MatrixError;
8use crate::events::client_event_json;
9use crate::store::{Membership, ReceiptType, Room};
10use crate::sync_token;
11
12/// `timeout` (ms) is clamped to this ceiling regardless of what the client
13/// asks for (plan §4 `/sync` row).
14pub const SYNC_MAX_TIMEOUT_MS: u64 = 30_000;
15
16/// `filter.room.timeline.limit` default and ceiling (P10 brief — overrides
17/// plan §3.2's own "default 20" wherever they differ).
18pub const SYNC_TIMELINE_DEFAULT_LIMIT: i64 = 10;
19
20pub const SYNC_TIMELINE_MAX_LIMIT: i64 = 50;
21
22/// `to_device.events` cap per response (P10 brief).
23pub const SYNC_TO_DEVICE_CAP: i64 = 100;
24
25/// `m.heroes` cap (plan §3.4).
26pub const SYNC_HERO_LIMIT: i64 = 5;
27
28
29// ============================================================================
30// Filter — `room.timeline.limit/types/not_types`, `room.state.
31// lazy_load_members`, `room.include_leave`, `account_data.types/not_types`.
32// `presence`/`event_fields` are accepted (serde ignores unrecognized JSON
33// keys by default — no `deny_unknown_fields`) and never interpreted.
34// ============================================================================
35
36#[derive(serde::Deserialize, Default, Clone)]
37pub struct RoomTimelineFilter {
38    #[serde(default)]
39    pub limit: Option<i64>,
40    #[serde(default)]
41    pub types: Option<Vec<String>>,
42    #[serde(default)]
43    pub not_types: Option<Vec<String>>,
44}
45
46
47#[derive(serde::Deserialize, Default, Clone)]
48pub struct RoomStateFilter {
49    #[serde(default)]
50    pub lazy_load_members: Option<bool>,
51}
52
53
54#[derive(serde::Deserialize, Default, Clone)]
55pub struct RoomFilter {
56    #[serde(default)]
57    pub timeline: RoomTimelineFilter,
58    #[serde(default)]
59    pub state: RoomStateFilter,
60    #[serde(default)]
61    pub include_leave: bool,
62}
63
64
65#[derive(serde::Deserialize, Default, Clone)]
66pub struct AccountDataFilter {
67    #[serde(default)]
68    pub types: Option<Vec<String>>,
69    #[serde(default)]
70    pub not_types: Option<Vec<String>>,
71}
72
73
74#[derive(serde::Deserialize, Default, Clone)]
75pub struct SyncFilter {
76    #[serde(default)]
77    pub room: RoomFilter,
78    #[serde(default)]
79    pub account_data: AccountDataFilter,
80    /// Build room blocks only for these rooms (sliding sync: the rooms inside the asked ranges).
81    #[serde(skip)]
82    pub only_rooms: Option<std::collections::HashSet<String>>,
83}
84
85
86/// `types`/`not_types` gate shared by the timeline and account-data filter
87/// sections — own copy per this codebase's one-copy-per-route-module
88/// convention (mirrors `routes::matrix::messaging::event_passes_filter`).
89pub fn passes_type_filter(value: &str, types: &Option<Vec<String>>, not_types: &Option<Vec<String>>) -> bool {
90    if let Some(types) = types {
91        if !types.iter().any(|t| t == value) {
92            return false;
93        }
94    }
95    if let Some(not_types) = not_types {
96        if not_types.iter().any(|t| t == value) {
97            return false;
98        }
99    }
100    true
101}
102
103
104impl SyncFilter {
105    fn timeline_limit(&self) -> i64 {
106        self.room.timeline.limit.unwrap_or(SYNC_TIMELINE_DEFAULT_LIMIT).clamp(1, SYNC_TIMELINE_MAX_LIMIT)
107    }
108
109    /// Default `true` — every client this server ships sets it, and it is
110    /// the single biggest `/sync` payload-size win for anything but a tiny
111    /// room (plan §3.3).
112    pub fn lazy_load_members(&self) -> bool {
113        self.room.state.lazy_load_members.unwrap_or(true)
114    }
115
116    fn include_leave(&self) -> bool {
117        self.room.include_leave
118    }
119
120    fn timeline_event_passes(&self, event_type: &str) -> bool {
121        passes_type_filter(event_type, &self.room.timeline.types, &self.room.timeline.not_types)
122    }
123
124    fn account_data_passes(&self, data_type: &str) -> bool {
125        passes_type_filter(data_type, &self.account_data.types, &self.account_data.not_types)
126    }
127}
128
129
130/// Resolve the `filter` query param: a stored `filter_id` (looked up via
131/// [`crate::store::get_filter`], scoped to `user_id` — the same isolation
132/// `routes::matrix::account::get_filter` enforces) or inline JSON (a leading
133/// `{`). Absent/empty is [`SyncFilter::default`].
134pub fn parse_filter_param(conn: &Connection, user_id: i64, raw: Option<&str>) -> Result<SyncFilter, MatrixError> {
135    let Some(raw) = raw.map(str::trim).filter(|s| !s.is_empty()) else {
136        return Ok(SyncFilter::default());
137    };
138    let definition = if raw.starts_with('{') {
139        raw.to_string()
140    } else {
141        let filter_id: i64 = raw.parse().map_err(|_| MatrixError::invalid_param("filter must be a filter id or inline JSON"))?;
142        crate::store::get_filter(conn, user_id, filter_id)?.ok_or_else(|| MatrixError::not_found("no such filter"))?
143    };
144    serde_json::from_str(&definition).map_err(|_| MatrixError::invalid_param("malformed filter"))
145}
146
147
148// ============================================================================
149// Response emptiness — what makes a long-poll keep waiting
150// ============================================================================
151
152fn events_array_is_empty(value: &serde_json::Value, key: &str) -> bool {
153    value.get(key).and_then(|v| v.get("events")).and_then(|v| v.as_array()).map(|a| a.is_empty()).unwrap_or(true)
154}
155
156
157/// "No rooms/account-data/to-device/device-lists changes, no typing/receipt
158/// change" (P10 brief) — `device_one_time_keys_count`/
159/// `device_unused_fallback_key_types` are deliberately NOT part of this
160/// check: they are static per-device facts always present, never a "new
161/// since last time" signal.
162pub fn sync_response_is_empty(value: &serde_json::Value) -> bool {
163    let rooms_empty = value
164        .get("rooms")
165        .map(|rooms| {
166            ["join", "invite", "leave"]
167                .iter()
168                .all(|kind| rooms.get(kind).and_then(|v| v.as_object()).map(|o| o.is_empty()).unwrap_or(true))
169        })
170        .unwrap_or(true);
171    let device_lists_empty = value
172        .get("device_lists")
173        .map(|dl| {
174            ["changed", "left"]
175                .iter()
176                .all(|k| dl.get(k).and_then(|v| v.as_array()).map(|a| a.is_empty()).unwrap_or(true))
177        })
178        .unwrap_or(true);
179    rooms_empty && events_array_is_empty(value, "account_data") && events_array_is_empty(value, "to_device") && device_lists_empty
180}
181
182
183// ============================================================================
184// Event formatting helpers
185// ============================================================================
186
187pub fn format_event(conn: &Connection, caller_user_id: i64, caller_device_id: &str, event: &crate::store::MatrixEvent) -> Result<serde_json::Value, MatrixError> {
188    let own_txn_id = crate::store::txn_id_for_event(conn, caller_user_id, caller_device_id, &event.event_id)?;
189    client_event_json(conn, event, own_txn_id.as_deref())
190}
191
192
193pub fn format_pagination_token(stream_id: i64) -> String {
194    format!("t{stream_id}")
195}
196
197
198/// One coalesced `m.receipt` ephemeral event (plan §3.6): every `m.read`
199/// receipt (visible to the whole room by definition) plus ONLY the caller's
200/// own `m.read.private` receipt — Matrix's own privacy rule that a private
201/// receipt must never reach another user's `/sync`. `None` when `receipts`
202/// contained only OTHER users' private receipts — after the privacy filter
203/// there is nothing left worth an (otherwise empty-content) event at all,
204/// and an empty-but-present `m.receipt` event would wrongly count as "this
205/// room changed" for the emptiness check that gates a room's inclusion.
206pub fn build_receipt_event(conn: &Connection, caller_user_id: i64, receipts: &[crate::store::ReceiptRow]) -> Result<Option<serde_json::Value>, MatrixError> {
207    let mut by_event: serde_json::Map<String, serde_json::Value> = serde_json::Map::new();
208    for receipt in receipts {
209        if receipt.receipt_type == ReceiptType::ReadPrivate && receipt.user_id != caller_user_id {
210            continue;
211        }
212        let Some(user_mxid) = crate::store::mxid_of(conn, receipt.user_id)? else { continue };
213
214        let event_entry = by_event.entry(receipt.event_id.clone()).or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
215        let event_map = event_entry.as_object_mut().ok_or_else(MatrixError::internal)?;
216        let type_entry = event_map
217            .entry(receipt.receipt_type.as_str().to_string())
218            .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
219        let type_map = type_entry.as_object_mut().ok_or_else(MatrixError::internal)?;
220        type_map.insert(user_mxid, serde_json::json!({ "ts": receipt.ts }));
221    }
222    if by_event.is_empty() {
223        return Ok(None);
224    }
225    Ok(Some(serde_json::json!({ "type": "m.receipt", "content": serde_json::Value::Object(by_event) })))
226}
227
228
229// ============================================================================
230// Per-room block builders
231// ============================================================================
232
233/// One `rooms.join.{roomId}` block, or `None` when this is an incremental
234/// sync AND nothing changed for this room at all (a room is only listed
235/// when something happened — see the module doc's emptiness rule).
236#[allow(clippy::too_many_arguments)]
237pub fn build_join_room_block(
238    conn: &Connection,
239    typing: &crate::typing::TypingRegistry,
240    room: &Room,
241    caller_user_id: i64,
242    caller_mxid: &str,
243    caller_device_id: &str,
244    since_stream: i64,
245    upto: i64,
246    is_initial_or_full: bool,
247    typing_gen_since: u64,
248    filter: &SyncFilter,
249    now: Instant,
250) -> Result<Option<serde_json::Value>, MatrixError> {
251    let member_event = crate::store::current_state_event(conn, &room.id, "m.room.member", caller_mxid)?;
252    let joined_stream = member_event.as_ref().map(|e| e.stream_id).unwrap_or(0);
253    let newly_joined = !is_initial_or_full && joined_stream > since_stream;
254    let full_view = is_initial_or_full || newly_joined;
255    let effective_since = if full_view { 0 } else { since_stream };
256
257    let limit = filter.timeline_limit();
258    let (mut timeline_events, limited) = if full_view {
259        let mut page = crate::store::events_in_room_before(conn, &room.id, upto.saturating_add(1), limit + 1)?;
260        let limited = page.len() as i64 > limit;
261        if limited {
262            page.pop(); // drop the oldest of the N+1 (page is newest-first)
263        }
264        page.reverse(); // ascending
265        (page, limited)
266    } else {
267        let mut page = crate::store::events_in_room_after(conn, &room.id, since_stream, limit + 1)?;
268        let limited = page.len() as i64 > limit;
269        if limited {
270            let drop_count = page.len() - limit as usize;
271            page.drain(0..drop_count); // keep only the newest `limit`
272        }
273        (page, limited)
274    };
275    // Computed from the RAW (pre-type-filter) page — `limited` signals a gap
276    // in raw stream continuity, independent of what a content filter thins
277    // out (matches `routes::matrix::messaging::paginate_messages`'s own
278    // precedent).
279    let prev_batch = timeline_events.first().map(|e| format_pagination_token(e.stream_id - 1));
280    timeline_events.retain(|e| filter.timeline_event_passes(&e.event_type));
281
282    // STATE
283    let mut state_events: Vec<crate::store::MatrixEvent> = Vec::new();
284    if filter.lazy_load_members() {
285        // Clause (a): member events of every DISTINCT sender in the
286        // returned (capped) timeline window.
287        let mut senders: Vec<i64> = timeline_events.iter().map(|e| e.sender_user_id).collect();
288        senders.sort_unstable();
289        senders.dedup();
290        for sender in senders {
291            if let Some(mxid) = crate::store::mxid_of(conn, sender)? {
292                if let Some(ev) = crate::store::current_state_event(conn, &room.id, "m.room.member", &mxid)? {
293                    state_events.push(ev);
294                }
295            }
296        }
297        // Clause (b): the matrix-spec#942 gap rule — only when there is
298        // actually a gap the timeline did not cover (incremental + limited;
299        // a full/initial view has no prior client state to keep consistent
300        // with, so nothing to backfill).
301        if !full_view && limited {
302            for ev in crate::store::member_state_changed_in_window(conn, &room.id, since_stream, upto)? {
303                if !state_events.iter().any(|existing| existing.event_id == ev.event_id) {
304                    state_events.push(ev);
305                }
306            }
307        }
308    } else {
309        state_events.extend(crate::store::state_events_of_type_at(conn, &room.id, "m.room.member", upto)?);
310    }
311    // Every OTHER state type: the delta since `effective_since` (never
312    // gated on `limited` — see the module doc), MINUS anything already
313    // present in the returned `timeline` (manager review, 2026-09-24):
314    // `state` is the state at the START of the timeline window, not a
315    // second copy of state changes the client is ALREADY receiving as
316    // timeline entries. On a non-limited incremental sync this makes
317    // non-member `state` collapse to empty (every state change in the
318    // window is, by definition, already IN that unlimited timeline); on a
319    // limited one, only the gap's changes — the ones the truncated
320    // timeline did NOT cover — survive the exclusion below.
321    state_events.extend(crate::store::non_member_state_changed_in_window(conn, &room.id, effective_since, upto)?);
322    let timeline_event_ids: std::collections::HashSet<&str> = timeline_events.iter().map(|e| e.event_id.as_str()).collect();
323    state_events.retain(|e| !timeline_event_ids.contains(e.event_id.as_str()));
324
325    // EPHEMERAL: typing + receipts
326    let mut ephemeral_events = Vec::new();
327    let room_typing_serial = typing.typing_serial(&room.id, now);
328    if full_view || room_typing_serial > typing_gen_since {
329        let user_ids: Vec<String> =
330            typing.typing_users(&room.id, now).into_iter().filter_map(|uid| crate::store::mxid_of(conn, uid).ok().flatten()).collect();
331        ephemeral_events.push(serde_json::json!({ "type": "m.typing", "content": { "user_ids": user_ids } }));
332    }
333    let receipts = crate::store::receipts_changed_in_room(conn, &room.id, effective_since, upto)?;
334    if !receipts.is_empty() {
335        if let Some(receipt_event) = build_receipt_event(conn, caller_user_id, &receipts)? {
336            ephemeral_events.push(receipt_event);
337        }
338    }
339
340    // ACCOUNT DATA
341    let mut account_data_events = Vec::new();
342    for row in crate::store::account_data_since(conn, caller_user_id, &room.id, effective_since)? {
343        if !filter.account_data_passes(&row.data_type) {
344            continue;
345        }
346        let content: serde_json::Value = serde_json::from_str(&row.content).unwrap_or_else(|_| serde_json::json!({}));
347        account_data_events.push(serde_json::json!({ "type": row.data_type, "content": content }));
348    }
349
350    let is_empty_incremental =
351        !full_view && timeline_events.is_empty() && state_events.is_empty() && ephemeral_events.is_empty() && account_data_events.is_empty();
352    if is_empty_incremental {
353        return Ok(None);
354    }
355
356    let mut timeline_json = Vec::with_capacity(timeline_events.len());
357    for event in &timeline_events {
358        timeline_json.push(format_event(conn, caller_user_id, caller_device_id, event)?);
359    }
360    let mut state_json = Vec::with_capacity(state_events.len());
361    for event in &state_events {
362        state_json.push(format_event(conn, caller_user_id, caller_device_id, event)?);
363    }
364
365    // `summary` is always attached whenever the room block itself is sent
366    // (never gated on detecting a change against a prior sync — this
367    // server does not persist a "last reported summary" to compare
368    // against; re-sending unchanged, cheap, idempotent metadata is spec-
369    // permitted and harmless, unlike re-sending the whole timeline would
370    // be).
371    let heroes: Vec<String> = crate::store::room_heroes(conn, &room.id, caller_user_id, SYNC_HERO_LIMIT)?
372        .into_iter()
373        .filter_map(|uid| crate::store::mxid_of(conn, uid).ok().flatten())
374        .collect();
375    let joined_count = crate::store::room_members(conn, &room.id, Some(Membership::Join))?.len();
376    let invited_count = crate::store::room_members(conn, &room.id, Some(Membership::Invite))?.len();
377    let notification_count = crate::store::notification_count(conn, room, caller_user_id)?;
378
379    let mut timeline = serde_json::json!({ "events": timeline_json, "limited": limited });
380    if let Some(prev_batch) = prev_batch {
381        timeline["prev_batch"] = serde_json::Value::String(prev_batch);
382    }
383
384    Ok(Some(serde_json::json!({
385        "summary": {
386            "m.heroes": heroes,
387            "m.joined_member_count": joined_count,
388            "m.invited_member_count": invited_count,
389        },
390        "state": { "events": state_json },
391        "timeline": timeline,
392        "ephemeral": { "events": ephemeral_events },
393        "account_data": { "events": account_data_events },
394        // `highlight_count` is always 0 — no push-rule/keyword-highlight
395        // engine in this server (plan §3.5).
396        "unread_notifications": { "notification_count": notification_count, "highlight_count": 0 },
397    })))
398}
399
400
401/// One `rooms.invite.{roomId}` block, or `None` when the invite is not new
402/// (an incremental sync whose invite predates `since_stream`).
403pub fn build_invite_room_block(conn: &Connection, room_id: &str, caller_mxid: &str, is_initial: bool, since_stream: i64) -> Result<Option<serde_json::Value>, MatrixError> {
404    let Some(member_event) = crate::store::current_state_event(conn, room_id, "m.room.member", caller_mxid)? else {
405        return Ok(None);
406    };
407    if !(is_initial || member_event.stream_id > since_stream) {
408        return Ok(None);
409    }
410    let mut events = crate::store::stripped_invite_state(conn, room_id, member_event.sender_user_id)?;
411    events.push(crate::store::stripped_state_json(conn, &member_event)?);
412    Ok(Some(serde_json::json!({ "invite_state": { "events": events } })))
413}
414
415
416/// One `rooms.leave.{roomId}` block: timeline up to and including the
417/// leave/kick/ban event itself (naturally bounded — nothing after that
418/// point is ever queried), `state` minimal (v1 simplification: the client
419/// was previously joined, so it already has the room's bootstrap state from
420/// before it left — nothing new is reconstructed here).
421pub fn build_leave_room_block(
422    conn: &Connection,
423    room: &Room,
424    caller_user_id: i64,
425    caller_device_id: &str,
426    member_event: &crate::store::MatrixEvent,
427    limit: i64,
428) -> Result<serde_json::Value, MatrixError> {
429    let mut page = crate::store::events_in_room_before(conn, &room.id, member_event.stream_id.saturating_add(1), limit + 1)?;
430    let limited = page.len() as i64 > limit;
431    if limited {
432        page.pop();
433    }
434    page.reverse();
435
436    let mut events_json = Vec::with_capacity(page.len());
437    for event in &page {
438        events_json.push(format_event(conn, caller_user_id, caller_device_id, event)?);
439    }
440
441    Ok(serde_json::json!({
442        "timeline": { "events": events_json, "limited": limited },
443        "state": { "events": [] },
444    }))
445}
446
447
448// ============================================================================
449// The whole-response builder — one consistent cut (see the module doc)
450// ============================================================================
451
452#[allow(clippy::too_many_arguments)]
453pub fn build_sync_response(
454    conn: &Connection,
455    typing: &crate::typing::TypingRegistry,
456    caller_user_id: i64,
457    caller_mxid: &str,
458    caller_device_id: &str,
459    since: Option<sync_token::SyncToken>,
460    filter: &SyncFilter,
461    full_state: bool,
462    now: Instant,
463) -> Result<serde_json::Value, MatrixError> {
464    let upto = crate::store::max_stream_id(conn)?;
465    let since_stream = since.map(|t| t.stream_id).unwrap_or(0);
466    let typing_gen_since = since.map(|t| t.typing_gen).unwrap_or(0);
467    // Snapshotted BEFORE any per-room typing read (mirrors `live.rs`'s own
468    // "register before read" discipline): any typing change racing in
469    // AFTER this point still stamps a room serial strictly greater than
470    // what we are about to hand back as `next_batch`, so the client's NEXT
471    // poll (using that value as its own `typing_gen_since`) is guaranteed
472    // to observe it. Capturing this AFTER building rooms would risk the
473    // reverse: a change landing after a room's own (already-negative) check
474    // but before this snapshot would be stamped INTO next_batch without
475    // ever having been reported — silently lost until a later change
476    // happens to exceed it.
477    let typing_gen_now = typing.current_typing_gen();
478    let is_initial = since.is_none() || full_state;
479
480    // Delete-after-ack: the client asking for `since_stream` is itself the
481    // proof it already durably received everything up to that point
482    // (plan §3.7). A no-op for an initial sync (`since_stream == 0`).
483    crate::keys::delete_to_device_up_to(conn, caller_user_id, caller_device_id, since_stream)?;
484
485    let timeline_limit = filter.timeline_limit();
486
487    let mut joined_room_ids = crate::store::rooms_for_user(conn, caller_user_id, Some(Membership::Join))?;
488    if let Some(only) = &filter.only_rooms {
489        joined_room_ids.retain(|r| only.contains(r));
490    }
491    // Changed-room prefilter (manager review, 2026-09-24): on an initial/
492    // `full_state` sync every joined room is built regardless (there is no
493    // `since` boundary to prefilter against). On an incremental sync,
494    // building EVERY joined room's block just to discover most did not
495    // change costs a room block builder's own dozen-odd queries per room —
496    // with a few hundred rooms that is thousands of queries per wake, all
497    // held under the single `messenger.db` mutex, blocking every other
498    // user. `rooms_changed_in_window` answers "did anything happen here at
499    // all" in three cheap, indexed queries total (not per room), and the
500    // in-memory typing check costs nothing; only rooms in the resulting set
501    // ever reach [`build_join_room_block`].
502    let rooms_to_build: Vec<String> = if is_initial {
503        joined_room_ids.clone()
504    } else {
505        let mut changed = crate::store::rooms_changed_in_window(conn, &joined_room_ids, caller_user_id, since_stream, upto)?;
506        for room_id in &joined_room_ids {
507            if typing.typing_serial(room_id, now) > typing_gen_since {
508                changed.insert(room_id.clone());
509            }
510        }
511        joined_room_ids.iter().filter(|room_id| changed.contains(*room_id)).cloned().collect()
512    };
513
514    let mut rooms_join = serde_json::Map::new();
515    for room_id in rooms_to_build {
516        let Some(room) = crate::store::get_room(conn, &room_id)? else { continue };
517        if let Some(block) = build_join_room_block(
518            conn,
519            typing,
520            &room,
521            caller_user_id,
522            caller_mxid,
523            caller_device_id,
524            since_stream,
525            upto,
526            is_initial,
527            typing_gen_since,
528            filter,
529            now,
530        )? {
531            rooms_join.insert(room_id, block);
532        }
533    }
534
535    let mut rooms_invite = serde_json::Map::new();
536    for room_id in crate::store::rooms_for_user(conn, caller_user_id, Some(Membership::Invite))? {
537        if filter.only_rooms.as_ref().is_some_and(|o| !o.contains(&room_id)) {
538            continue;
539        }
540        if let Some(block) = build_invite_room_block(conn, &room_id, caller_mxid, is_initial, since_stream)? {
541            rooms_invite.insert(room_id, block);
542        }
543    }
544
545    // Left rooms: an initial sync omits them entirely (there is no boundary
546    // to test "left inside the window" against), and even on an incremental
547    // sync only when the filter opts in (`room.include_leave`, default
548    // false — P10 brief).
549    let mut rooms_leave = serde_json::Map::new();
550    if since.is_some() && filter.include_leave() {
551        for membership in [Membership::Leave, Membership::Ban] {
552            for room_id in crate::store::rooms_for_user(conn, caller_user_id, Some(membership))? {
553                let Some(room) = crate::store::get_room(conn, &room_id)? else { continue };
554                let Some(member_event) = crate::store::current_state_event(conn, &room_id, "m.room.member", caller_mxid)? else { continue };
555                if member_event.stream_id > since_stream {
556                    let block = build_leave_room_block(conn, &room, caller_user_id, caller_device_id, &member_event, timeline_limit)?;
557                    rooms_leave.insert(room_id, block);
558                }
559            }
560        }
561    }
562
563    let account_data_since_global = if is_initial { 0 } else { since_stream };
564    let mut global_account_data = Vec::new();
565    for row in crate::store::account_data_since(conn, caller_user_id, crate::store::GLOBAL_ACCOUNT_DATA_ROOM, account_data_since_global)? {
566        if !filter.account_data_passes(&row.data_type) {
567            continue;
568        }
569        let content: serde_json::Value = serde_json::from_str(&row.content).unwrap_or_else(|_| serde_json::json!({}));
570        global_account_data.push(serde_json::json!({ "type": row.data_type, "content": content }));
571    }
572
573    // Cap to-device delivery at `SYNC_TO_DEVICE_CAP`. If capped, lower ONLY
574    // `next_batch`'s own stream position to the last DELIVERED message's
575    // `stream_id` (P10 brief's "deliver in order and let the next sync
576    // continue" option) — every other field in THIS response still reflects
577    // the full `upto` snapshot; anything past the lowered watermark that
578    // would otherwise have been reported now is simply reported again on
579    // the next call, which is safe (event ids are stable, receipts/account
580    // data are last-write-wins).
581    let mut to_device_raw = crate::keys::to_device_for(conn, caller_user_id, caller_device_id, since_stream, SYNC_TO_DEVICE_CAP + 1)?;
582    let to_device_capped = to_device_raw.len() as i64 > SYNC_TO_DEVICE_CAP;
583    if to_device_capped {
584        to_device_raw.truncate(SYNC_TO_DEVICE_CAP as usize);
585    }
586    let effective_upto = if to_device_capped { to_device_raw.last().map(|m| m.stream_id).unwrap_or(upto) } else { upto };
587    let mut to_device_events = Vec::with_capacity(to_device_raw.len());
588    for message in &to_device_raw {
589        let Some(sender_mxid) = crate::store::mxid_of(conn, message.sender_user_id)? else { continue };
590        let content: serde_json::Value = serde_json::from_str(&message.content).unwrap_or_else(|_| serde_json::json!({}));
591        to_device_events.push(serde_json::json!({ "sender": sender_mxid, "type": message.event_type, "content": content }));
592    }
593
594    // Device lists: omitted entirely on an initial/`full_state` sync (a
595    // fresh client queries `/keys/query` for every room member outright —
596    // an incremental DELTA has no meaning against no prior state).
597    //
598    // Incrementally, `changed` also carries users who newly share an
599    // encrypted room with the caller (a join/invite since `since`), and `left`
600    // users who no longer share any room with it — see
601    // [`super::keys::device_list_delta`], which `/keys/changes` shares.
602    let (device_lists_changed, device_lists_left) = if is_initial {
603        (Vec::<String>::new(), Vec::<String>::new())
604    } else {
605        let delta = crate::key_ops::device_list_delta(conn, caller_user_id, since_stream, upto)?;
606        (delta.changed, delta.left)
607    };
608
609    let otk_counts = crate::keys::count_one_time_keys(conn, caller_user_id, caller_device_id)?;
610    let unused_fallback = crate::keys::unused_fallback_key_types(conn, caller_user_id, caller_device_id)?;
611
612    let next_batch = sync_token::format(effective_upto, typing_gen_now);
613
614    Ok(serde_json::json!({
615        "next_batch": next_batch,
616        "account_data": { "events": global_account_data },
617        "to_device": { "events": to_device_events },
618        "device_lists": { "changed": device_lists_changed, "left": device_lists_left },
619        "device_one_time_keys_count": otk_counts,
620        "device_unused_fallback_key_types": unused_fallback,
621        "presence": { "events": crate::http::presence::sync_events(conn, caller_user_id, since_stream, upto)? },
622        "rooms": {
623            "join": serde_json::Value::Object(rooms_join),
624            "invite": serde_json::Value::Object(rooms_invite),
625            "leave": serde_json::Value::Object(rooms_leave),
626        },
627    }))
628}
629#[derive(serde::Deserialize, Default)]
630pub struct SyncQuery {
631    #[serde(default)]
632    pub since: Option<String>,
633    #[serde(default)]
634    pub timeout: Option<u64>,
635    #[serde(default)]
636    pub filter: Option<String>,
637    #[serde(default)]
638    pub full_state: Option<bool>,
639}