1#![cfg_attr(docsrs, feature(doc_cfg))]
16#![doc(html_root_url = "https://docs.rs/ferogram-msgbox/0.6.5")]
17#![deny(unsafe_code)]
91
92mod adaptor;
110pub mod defs;
111
112#[cfg(test)]
113mod tests;
114
115use std::cmp::Ordering;
116use std::time::Duration;
117
118use defs::Instant; use defs::Key;
121pub use defs::{ChannelState, Gap, MessageBoxes, UpdatesLike, UpdatesStateSnap};
122use defs::{
123 LiveEntry, NO_DATE, NO_PTS, NO_SEQ, POSSIBLE_GAP_TIMEOUT, PossibleGap, PtsInfo, UpdateAndPeers,
124};
125use ferogram_tl_types as tl;
126
127pub fn classify_own_response(body: &[u8]) -> Option<UpdatesLike> {
144 use tl::{Deserializable, Identifiable};
145
146 if body.len() < 4 {
147 return None;
148 }
149
150 let mut cur = tl::Cursor::from_slice(body);
152 if let Ok(updates) = tl::enums::Updates::deserialize(&mut cur) {
153 return Some(UpdatesLike::Updates(Box::new(updates)));
154 }
155
156 let mut cur = tl::Cursor::from_slice(body);
161 if let Ok(tl::enums::messages::ChatInviteJoinResult::Ok(ok)) =
162 tl::enums::messages::ChatInviteJoinResult::deserialize(&mut cur)
163 {
164 return Some(UpdatesLike::Updates(Box::new(ok.updates)));
165 }
166
167 let id = u32::from_le_bytes([body[0], body[1], body[2], body[3]]);
170 if id == <tl::types::messages::AffectedMessages as Identifiable>::CONSTRUCTOR_ID {
171 let mut cur = tl::Cursor::from_slice(&body[4..]);
172 if let Ok(affected) = tl::types::messages::AffectedMessages::deserialize(&mut cur) {
173 return Some(UpdatesLike::AffectedMessages(affected));
174 }
175 }
176
177 None
178}
179
180fn next_updates_deadline() -> Instant {
183 Instant::now() + defs::NO_UPDATES_TIMEOUT
184}
185
186fn update_sort_key(update: &tl::enums::Update) -> i32 {
187 match PtsInfo::from_update(update) {
188 Some(info) => info.pts - info.count,
189 None => NO_PTS,
190 }
191}
192
193impl Default for MessageBoxes {
196 fn default() -> Self {
197 Self::new()
198 }
199}
200
201impl MessageBoxes {
202 pub fn new() -> Self {
204 tracing::trace!("[ferogram::msgbox] created new (no prior state)");
205 Self {
206 entries: Vec::new(),
207 date: NO_DATE,
208 seq: NO_SEQ,
209 getting_diff_for: Vec::new(),
210 channel_diff_in_flight: None,
211 next_deadline: next_updates_deadline(),
212 }
213 }
214
215 pub fn load(state: UpdatesStateSnap) -> Self {
217 tracing::trace!("[ferogram::msgbox] loaded from state: {:?}", state);
218 let mut entries = Vec::with_capacity(2 + state.channels.len());
219 let mut getting_diff_for = Vec::with_capacity(2 + state.channels.len());
220 let deadline = next_updates_deadline();
221
222 if state.pts != NO_PTS {
223 entries.push(LiveEntry {
224 key: Key::Common,
225 pts: state.pts,
226 deadline,
227 possible_gap: None,
228 });
229 }
230 if state.qts != NO_PTS {
231 entries.push(LiveEntry {
232 key: Key::Secondary,
233 pts: state.qts,
234 deadline,
235 possible_gap: None,
236 });
237 }
238 entries.extend(state.channels.iter().map(|c| LiveEntry {
239 key: Key::Channel(c.id),
240 pts: c.pts,
241 deadline,
242 possible_gap: None,
243 }));
244 entries.sort_by_key(|e| e.key);
245
246 getting_diff_for.extend(entries.iter().map(|e| e.key));
248
249 Self {
250 entries,
251 date: state.date,
252 seq: state.seq,
253 getting_diff_for,
254 channel_diff_in_flight: None,
255 next_deadline: deadline,
256 }
257 }
258
259 fn entry(&self, key: Key) -> Option<&LiveEntry> {
260 self.entries
261 .binary_search_by_key(&key, |e| e.key)
262 .map(|i| &self.entries[i])
263 .ok()
264 }
265
266 fn update_entry(&mut self, key: Key, f: impl FnOnce(&mut LiveEntry)) -> bool {
267 match self.entries.binary_search_by_key(&key, |e| e.key) {
268 Ok(i) => {
269 f(&mut self.entries[i]);
270 true
271 }
272 Err(_) => false,
273 }
274 }
275
276 fn force_update_entry(&mut self, mut entry: LiveEntry, f: impl FnOnce(&mut LiveEntry)) {
277 match self.entries.binary_search_by_key(&entry.key, |e| e.key) {
278 Ok(i) => f(&mut self.entries[i]),
279 Err(i) => {
280 f(&mut entry);
281 self.entries.insert(i, entry);
282 }
283 }
284 }
285
286 fn set_entry(&mut self, entry: LiveEntry) {
287 match self.entries.binary_search_by_key(&entry.key, |e| e.key) {
288 Ok(i) => self.entries[i] = entry,
289 Err(i) => self.entries.insert(i, entry),
290 }
291 }
292
293 fn set_pts(&mut self, key: Key, pts: i32) {
294 if !self.update_entry(key, |e| e.pts = pts) {
295 self.set_entry(LiveEntry {
296 key,
297 pts,
298 deadline: next_updates_deadline(),
299 possible_gap: None,
300 });
301 }
302 }
303
304 fn pop_entry(&mut self, key: Key) -> Option<LiveEntry> {
305 match self.entries.binary_search_by_key(&key, |e| e.key) {
306 Ok(i) => Some(self.entries.remove(i)),
307 Err(_) => None,
308 }
309 }
310
311 fn reset_deadline(&mut self, key: Key, deadline: Instant) {
312 let mut old_deadline = self.next_deadline;
313 self.update_entry(key, |e| {
314 old_deadline = e.deadline;
315 e.deadline = deadline;
316 });
317 if self.next_deadline == old_deadline {
318 self.next_deadline = self
319 .entries
320 .iter()
321 .fold(deadline, |d, e| d.min(e.effective_deadline()));
322 }
323 }
324
325 fn reset_timeout(&mut self, key: Key, timeout: Option<i32>) {
326 self.reset_deadline(
327 key,
328 timeout
329 .map(|t| Instant::now() + Duration::from_secs(t as _))
330 .unwrap_or_else(next_updates_deadline),
331 );
332 }
333
334 fn try_begin_get_diff(&mut self, key: Key) {
335 if !self.getting_diff_for.contains(&key) {
336 if self.update_entry(key, |e| e.possible_gap = None) {
337 self.getting_diff_for.push(key);
338 } else {
339 tracing::debug!(
340 "[ferogram::msgbox] begin_get_diff skipped for {:?}: no entry exists",
341 key
342 );
343 }
344 }
345 }
346
347 fn try_end_get_diff(&mut self, key: Key) {
348 let i = match self.getting_diff_for.iter().position(|&k| k == key) {
349 Some(i) => i,
350 None => return,
351 };
352 self.getting_diff_for.remove(i);
353 self.reset_deadline(key, next_updates_deadline());
354 debug_assert!(
355 self.entry(key).is_none_or(|e| e.possible_gap.is_none()),
356 "gaps shouldn't be created while getting difference"
357 );
358 }
359
360 pub fn is_empty(&self) -> bool {
362 self.entries.is_empty()
363 }
364
365 pub fn set_state(&mut self, state: tl::types::updates::State) {
367 debug_assert!(self.is_empty());
368 let deadline = next_updates_deadline();
369 self.set_entry(LiveEntry {
370 key: Key::Common,
371 pts: state.pts,
372 deadline,
373 possible_gap: None,
374 });
375 self.set_entry(LiveEntry {
376 key: Key::Secondary,
377 pts: state.qts,
378 deadline,
379 possible_gap: None,
380 });
381 self.date = state.date;
382 self.seq = state.seq;
383 self.next_deadline = deadline;
384 }
385
386 pub fn try_set_channel_state(&mut self, id: i64, pts: i32) {
388 if self.entry(Key::Channel(id)).is_none() {
389 self.set_entry(LiveEntry {
390 key: Key::Channel(id),
391 pts,
392 deadline: next_updates_deadline(),
393 possible_gap: None,
394 });
395 }
396 }
397
398 pub fn session_state(&self) -> UpdatesStateSnap {
400 UpdatesStateSnap {
401 pts: self.entry(Key::Common).map(|e| e.pts).unwrap_or(NO_PTS),
402 qts: self.entry(Key::Secondary).map(|e| e.pts).unwrap_or(NO_PTS),
403 date: self.date,
404 seq: self.seq,
405 channels: self
406 .entries
407 .iter()
408 .filter_map(|e| match e.key {
409 Key::Channel(id) => Some(ChannelState { id, pts: e.pts }),
410 _ => None,
411 })
412 .collect(),
413 }
414 }
415
416 pub fn check_deadlines(&mut self) -> Instant {
422 let now = Instant::now();
423
424 if !self.getting_diff_for.is_empty() {
425 return now; }
427
428 if now >= self.next_deadline {
429 self.getting_diff_for
431 .extend(self.entries.iter().filter_map(|e| {
432 if now >= e.effective_deadline() {
433 tracing::debug!(
434 "[ferogram::msgbox] deadline met for {:?}; forcing diff",
435 e.key
436 );
437 Some(e.key)
438 } else {
439 None
440 }
441 }));
442
443 for i in 0..self.getting_diff_for.len() {
445 self.update_entry(self.getting_diff_for[i], |e| e.possible_gap = None);
446 }
447
448 if self.getting_diff_for.is_empty() {
449 self.next_deadline = next_updates_deadline();
450 }
451 }
452
453 self.next_deadline
454 }
455
456 pub fn get_difference(&self) -> Option<tl::functions::updates::GetDifference> {
461 for key in [Key::Common, Key::Secondary] {
462 if self.getting_diff_for.contains(&key) {
463 let pts = self
464 .entry(Key::Common)
465 .map(|e| e.pts)
466 .expect("Common entry must exist when diffing it");
467
468 return Some(tl::functions::updates::GetDifference {
469 pts,
470 pts_limit: None,
471 pts_total_limit: None,
472 date: self.date.max(1),
473 qts: self.entry(Key::Secondary).map(|e| e.pts).unwrap_or(NO_PTS),
474 qts_limit: None,
475 });
476 }
477 }
478 None
479 }
480
481 pub fn get_channel_difference(
486 &mut self,
487 ) -> Option<(i64, tl::functions::updates::GetChannelDifference)> {
488 let (key, channel_id) = self.getting_diff_for.iter().find_map(|&k| match k {
493 Key::Channel(id) if self.channel_diff_in_flight != Some(id) => Some((k, id)),
494 _ => None,
495 })?;
496
497 let pts = self
498 .entry(key)
499 .map(|e| e.pts)
500 .expect("Channel entry must exist when diffing it");
501
502 self.channel_diff_in_flight = Some(channel_id);
503
504 Some((
505 channel_id,
506 tl::functions::updates::GetChannelDifference {
507 force: false,
508 channel: tl::enums::InputChannel::InputChannel(tl::types::InputChannel {
509 channel_id,
510 access_hash: 0, }),
512 filter: tl::enums::ChannelMessagesFilter::Empty,
513 pts,
514 limit: 0, },
516 ))
517 }
518}
519
520impl MessageBoxes {
523 pub fn process_updates(&mut self, updates: UpdatesLike) -> Result<UpdateAndPeers, Gap> {
528 let deadline = next_updates_deadline();
529
530 let tl::types::UpdatesCombined {
531 date,
532 seq_start,
533 seq,
534 mut updates,
535 users,
536 chats,
537 } = match adaptor::adapt(updates) {
538 Ok(combined) => combined,
539 Err(Gap) => {
540 self.try_begin_get_diff(Key::Common);
541 return Err(Gap);
542 }
543 };
544
545 let new_date = if date == NO_DATE { self.date } else { date };
546 let new_seq = if seq == NO_SEQ { self.seq } else { seq };
547
548 if seq_start != NO_SEQ {
550 match (self.seq + 1).cmp(&seq_start) {
551 Ordering::Equal => {} Ordering::Greater => {
553 tracing::debug!(
554 "[ferogram::msgbox] duplicate seq (local={}, remote={}), skipping",
555 self.seq,
556 seq_start
557 );
558 return Ok((Vec::new(), users, chats));
559 }
560 Ordering::Less => {
561 tracing::debug!(
562 "[ferogram::msgbox] seq gap (local={}, remote={})",
563 self.seq,
564 seq_start
565 );
566 self.try_begin_get_diff(Key::Common);
567 return Err(Gap);
568 }
569 }
570 }
571
572 updates.sort_by_key(update_sort_key);
574
575 let mut result: Vec<tl::enums::Update> = Vec::with_capacity(updates.len());
576 let mut have_unresolved_gaps = false;
577
578 for update in updates {
579 if let tl::enums::Update::ChannelTooLong(ref u) = update {
581 let key = Key::Channel(u.channel_id);
582 if let Some(pts) = u.pts {
583 self.set_entry(LiveEntry {
584 key,
585 pts,
586 deadline,
587 possible_gap: None,
588 });
589 }
590 self.try_begin_get_diff(key);
591 continue;
592 }
593
594 let info = match PtsInfo::from_update(&update) {
595 Some(info) => info,
596 None => {
597 result.push(update);
599 continue;
600 }
601 };
602
603 if self.getting_diff_for.contains(&info.key) {
606 tracing::debug!(
607 "[ferogram::msgbox] update for {:?} suppressed while getDifference is in flight",
608 info.key
609 );
610 self.reset_deadline(info.key, next_updates_deadline());
611 result.push(update);
612 continue;
613 }
614
615 let mut gap_deadline = None;
616
617 self.force_update_entry(
618 LiveEntry {
619 key: info.key,
620 pts: info.pts - info.count,
621 deadline,
622 possible_gap: None,
623 },
624 |entry| {
625 match (entry.pts + info.count).cmp(&info.pts) {
626 Ordering::Equal => {
627 entry.pts = info.pts;
629 entry.deadline = deadline;
630 result.push(update);
631 }
632 Ordering::Greater => {
633 tracing::debug!(
635 "[ferogram::msgbox] duplicate update for {:?} \
636 (local={}, count={}, remote={})",
637 info.key,
638 entry.pts,
639 info.count,
640 info.pts
641 );
642 }
643 Ordering::Less => {
644 tracing::debug!(
646 "[ferogram::msgbox] gap for {:?} \
647 (local={}, count={}, remote={})",
648 info.key,
649 entry.pts,
650 info.count,
651 info.pts
652 );
653 entry
654 .possible_gap
655 .get_or_insert_with(|| PossibleGap {
656 deadline: Instant::now() + POSSIBLE_GAP_TIMEOUT,
657 updates: Vec::new(),
658 })
659 .updates
660 .push(update.clone());
661 }
662 }
663
664 if let Some(mut gap) = entry.possible_gap.take() {
666 gap.updates.sort_by_key(|u| -update_sort_key(u));
667 while let Some(gap_update) = gap.updates.pop() {
668 let gap_info = PtsInfo::from_update(&gap_update)
669 .expect("only updates with pts may be buffered as gaps");
670 match (entry.pts + gap_info.count).cmp(&gap_info.pts) {
671 Ordering::Equal => {
672 entry.pts = gap_info.pts;
673 result.push(gap_update);
674 }
675 Ordering::Greater => {}
676 Ordering::Less => {
677 gap.updates.push(gap_update);
678 break;
679 }
680 }
681 }
682 if !gap.updates.is_empty() {
683 gap_deadline = Some(gap.deadline);
684 entry.possible_gap = Some(gap);
685 have_unresolved_gaps = true;
686 }
687 }
688 },
689 );
690
691 self.next_deadline = self.next_deadline.min(gap_deadline.unwrap_or(deadline));
692 }
693
694 if !result.is_empty() && !have_unresolved_gaps {
695 self.date = new_date;
696 self.seq = new_seq;
697 }
698
699 Ok((result, users, chats))
700 }
701}
702
703impl MessageBoxes {
706 pub fn apply_difference(
708 &mut self,
709 difference: tl::enums::updates::Difference,
710 ) -> UpdateAndPeers {
711 tracing::trace!("[ferogram::msgbox] applying account difference");
712 if !self.getting_diff_for.contains(&Key::Common)
713 && !self.getting_diff_for.contains(&Key::Secondary)
714 {
715 tracing::warn!(
716 "[ferogram::msgbox] apply_difference called but no diff was pending \
717 (concurrent call already completed?); ignoring"
718 );
719 return (Vec::new(), Vec::new(), Vec::new());
720 }
721
722 let finish: bool;
723 let result = match difference {
724 tl::enums::updates::Difference::Empty(e) => {
725 tracing::debug!(
726 "[ferogram::msgbox] difference empty (date={}, seq={})",
727 e.date,
728 e.seq
729 );
730 finish = true;
731 self.date = e.date;
732 self.seq = e.seq;
733 (Vec::new(), Vec::new(), Vec::new())
734 }
735 tl::enums::updates::Difference::Difference(d) => {
736 tracing::debug!(
737 "[ferogram::msgbox] getDifference: received full difference, applying"
738 );
739 finish = true;
740 self.apply_difference_type(d)
741 }
742 tl::enums::updates::Difference::Slice(tl::types::updates::DifferenceSlice {
743 new_messages,
744 new_encrypted_messages,
745 other_updates,
746 chats,
747 users,
748 intermediate_state: state,
749 }) => {
750 tracing::debug!(
751 "[ferogram::msgbox] getDifference: received slice, will request another round"
752 );
753 finish = false;
754 self.apply_difference_type(tl::types::updates::Difference {
755 new_messages,
756 new_encrypted_messages,
757 other_updates,
758 chats,
759 users,
760 state,
761 })
762 }
763 tl::enums::updates::Difference::TooLong(d) => {
764 tracing::warn!(
765 "[ferogram::msgbox] getDifference returned TooLong (pts={}); resetting to server pts",
766 d.pts
767 );
768 finish = true;
769 self.set_pts(Key::Common, d.pts);
770 (Vec::new(), Vec::new(), Vec::new())
771 }
772 };
773
774 if finish {
775 self.try_end_get_diff(Key::Common);
776 self.try_end_get_diff(Key::Secondary);
777 }
778
779 result
780 }
781
782 fn apply_difference_type(
783 &mut self,
784 tl::types::updates::Difference {
785 new_messages,
786 new_encrypted_messages,
787 other_updates: updates,
788 chats,
789 users,
790 state: tl::enums::updates::State::State(state),
791 }: tl::types::updates::Difference,
792 ) -> UpdateAndPeers {
793 self.date = state.date;
794 self.seq = state.seq;
795 self.set_pts(Key::Common, state.pts);
796 self.set_pts(Key::Secondary, state.qts);
797
798 let us = UpdatesLike::Updates(Box::new(tl::enums::Updates::Updates(tl::types::Updates {
800 updates,
801 users,
802 chats,
803 date: NO_DATE,
804 seq: NO_SEQ,
805 })));
806 let (mut result_updates, users, chats) = self
807 .process_updates(us)
808 .expect("gap detected while applying difference - should not happen");
809
810 let msgs: Vec<tl::enums::Update> = new_messages
812 .into_iter()
813 .map(|msg| {
814 tl::enums::Update::NewMessage(tl::types::UpdateNewMessage {
815 message: msg,
816 pts: NO_PTS,
817 pts_count: 0,
818 })
819 })
820 .chain(new_encrypted_messages.into_iter().map(|msg| {
821 tl::enums::Update::NewEncryptedMessage(tl::types::UpdateNewEncryptedMessage {
822 message: msg,
823 qts: NO_PTS,
824 })
825 }))
826 .collect();
827
828 result_updates.splice(0..0, msgs);
829 (result_updates, users, chats)
830 }
831}
832
833impl MessageBoxes {
836 pub fn apply_channel_difference(
838 &mut self,
839 difference: tl::enums::updates::ChannelDifference,
840 ) -> UpdateAndPeers {
841 let Some(channel_id) = self.channel_diff_in_flight.take() else {
842 tracing::warn!(
843 "[ferogram::msgbox] apply_channel_difference called but no channel diff \
844 was in flight (stale/duplicate response?); ignoring"
845 );
846 return (Vec::new(), Vec::new(), Vec::new());
847 };
848
849 let key = Key::Channel(channel_id);
850 if !self.getting_diff_for.contains(&key) {
851 tracing::debug!(
854 "[ferogram::msgbox] apply_channel_difference: channel {} no longer in \
855 getting_diff_for (stale duplicate response); ignoring",
856 channel_id
857 );
858 return (Vec::new(), Vec::new(), Vec::new());
859 }
860
861 tracing::trace!(
862 "[ferogram::msgbox] applying channel {} difference",
863 channel_id
864 );
865 self.update_entry(key, |e| e.possible_gap = None);
866
867 let Ok(tl::types::updates::ChannelDifference {
868 r#final,
869 pts,
870 timeout,
871 new_messages,
872 other_updates: updates,
873 chats,
874 users,
875 }) = adaptor::adapt_channel_difference(difference)
876 else {
877 self.channel_diff_in_flight = Some(channel_id);
881 self.end_channel_difference(PrematureEndReason::TemporaryServerIssues);
882 return (Vec::new(), Vec::new(), Vec::new());
883 };
884
885 if r#final {
886 tracing::debug!(
887 "[ferogram::msgbox] channel {} diff complete (final=true)",
888 channel_id
889 );
890 self.try_end_get_diff(key);
891 } else {
892 tracing::debug!(
893 "[ferogram::msgbox] channel {} diff slice received; requesting next batch",
894 channel_id
895 );
896 }
897
898 self.set_pts(key, pts);
899
900 let us = UpdatesLike::Updates(Box::new(tl::enums::Updates::Updates(tl::types::Updates {
901 updates,
902 users,
903 chats,
904 date: NO_DATE,
905 seq: NO_SEQ,
906 })));
907 let (mut result_updates, users, chats) = self
908 .process_updates(us)
909 .expect("gap detected while applying channel difference");
910
911 let msgs: Vec<tl::enums::Update> = new_messages
913 .into_iter()
914 .map(|msg| {
915 tl::enums::Update::NewChannelMessage(tl::types::UpdateNewChannelMessage {
916 message: msg,
917 pts: NO_PTS,
918 pts_count: 0,
919 })
920 })
921 .collect();
922 result_updates.splice(0..0, msgs);
923
924 self.reset_timeout(key, timeout);
925
926 (result_updates, users, chats)
927 }
928
929 pub fn abort_difference(&mut self) {
940 for key in [Key::Common, Key::Secondary] {
941 self.update_entry(key, |e| e.possible_gap = None);
942 self.try_end_get_diff(key);
943 }
944 tracing::debug!(
945 "[ferogram::msgbox] getDifference aborted; cleared pending state for Common and Secondary"
946 );
947 }
948
949 pub fn force_reset_common_pts(&mut self, pts: i32, qts: i32, date: i32, seq: i32) {
955 self.set_pts(Key::Common, pts);
956 self.set_pts(Key::Secondary, qts);
957 self.date = date;
958 self.seq = seq;
959 tracing::debug!(
960 "[ferogram::msgbox] force_reset_common_pts: pts={pts}, qts={qts}, seq={seq}"
961 );
962 }
963
964 pub fn end_channel_difference(&mut self, reason: PrematureEndReason) {
965 let Some(channel_id) = self.channel_diff_in_flight.take() else {
966 tracing::warn!(
967 "[ferogram::msgbox] end_channel_difference called but no channel diff \
968 was in flight (already ended? duplicate error path)"
969 );
970 return;
971 };
972 let key = Key::Channel(channel_id);
973 if !self.getting_diff_for.contains(&key) {
974 tracing::debug!(
975 "[ferogram::msgbox] end_channel_difference: channel {} no longer in \
976 getting_diff_for (stale duplicate response); ignoring",
977 channel_id
978 );
979 return;
980 }
981
982 tracing::trace!(
983 "[ferogram::msgbox] ending channel {} diff: {:?}",
984 channel_id,
985 reason
986 );
987
988 match reason {
989 PrematureEndReason::TemporaryServerIssues => {
990 self.update_entry(key, |e| e.possible_gap = None);
991 self.try_end_get_diff(key);
992 }
993 PrematureEndReason::Banned => {
994 self.update_entry(key, |e| e.possible_gap = None);
995 self.try_end_get_diff(key);
996 self.pop_entry(key);
997 }
998 }
999 }
1000}
1001
1002#[derive(Debug)]
1004pub enum PrematureEndReason {
1005 TemporaryServerIssues,
1007 Banned,
1009}