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