matrix_sdk_common/linked_chunk/updates.rs
1// Copyright 2024 The Matrix.org Foundation C.I.C.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::{
16 collections::HashMap,
17 pin::Pin,
18 sync::{Arc, RwLock, Weak},
19 task::{Context, Poll, Waker},
20};
21
22use futures_core::Stream;
23
24use super::{ChunkIdentifier, Position};
25
26/// Represent the updates that have happened inside a [`LinkedChunk`].
27///
28/// To retrieve the updates, use [`LinkedChunk::updates`].
29///
30/// These updates are useful to store a `LinkedChunk` in another form of
31/// storage, like a database or something similar.
32///
33/// [`LinkedChunk`]: super::LinkedChunk
34/// [`LinkedChunk::updates`]: super::LinkedChunk::updates
35#[derive(Debug, Clone, PartialEq)]
36pub enum Update<Item, Gap> {
37 /// A new chunk of kind Items has been created.
38 NewItemsChunk {
39 /// The identifier of the previous chunk of this new chunk.
40 previous: Option<ChunkIdentifier>,
41
42 /// The identifier of the new chunk.
43 new: ChunkIdentifier,
44
45 /// The identifier of the next chunk of this new chunk.
46 next: Option<ChunkIdentifier>,
47 },
48
49 /// A new chunk of kind Gap has been created.
50 NewGapChunk {
51 /// The identifier of the previous chunk of this new chunk.
52 previous: Option<ChunkIdentifier>,
53
54 /// The identifier of the new chunk.
55 new: ChunkIdentifier,
56
57 /// The identifier of the next chunk of this new chunk.
58 next: Option<ChunkIdentifier>,
59
60 /// The content of the chunk.
61 gap: Gap,
62 },
63
64 /// A chunk has been removed.
65 RemoveChunk(ChunkIdentifier),
66
67 /// Items are pushed inside a chunk of kind Items.
68 PushItems {
69 /// The [`Position`] of the items.
70 ///
71 /// This value is given to prevent the need for position computations by
72 /// the update readers. Items are pushed, so the positions should be
73 /// incrementally computed from the previous items, which requires the
74 /// reading of the last previous item. With `at`, the update readers no
75 /// longer need to do so.
76 at: Position,
77
78 /// The items.
79 items: Vec<Item>,
80 },
81
82 /// An item has been replaced in the linked chunk.
83 ///
84 /// The `at` position MUST resolve to the actual position an existing _item_
85 /// (not a gap).
86 ReplaceItem {
87 /// The position of the item that's being replaced.
88 at: Position,
89
90 /// The new value for the item.
91 item: Item,
92 },
93
94 /// An item has been removed inside a chunk of kind Items.
95 RemoveItem {
96 /// The [`Position`] of the item.
97 at: Position,
98 },
99
100 /// The last items of a chunk have been detached, i.e. the chunk has been
101 /// truncated.
102 DetachLastItems {
103 /// The split position. Before this position (`..position`), items are
104 /// kept, from this position (`position..`), items are detached.
105 at: Position,
106 },
107
108 /// Detached items (see [`Self::DetachLastItems`]) starts being reattached.
109 StartReattachItems,
110
111 /// Reattaching items (see [`Self::StartReattachItems`]) is finished.
112 EndReattachItems,
113
114 /// All chunks have been cleared, i.e. all items and all gaps have been
115 /// dropped.
116 Clear,
117}
118
119impl<Item, Gap> Update<Item, Gap> {
120 /// Get the items from the [`Update`] if any.
121 ///
122 /// This function is useful if you only care about the items from the
123 /// [`Update`] and not what kind of update it was and where the items should
124 /// be placed.
125 ///
126 /// [`Update`] variants which don't contain any items will return an empty
127 /// [`Vec`].
128 pub fn into_items(self) -> Vec<Item> {
129 match self {
130 Update::NewItemsChunk { .. }
131 | Update::NewGapChunk { .. }
132 | Update::RemoveChunk(_)
133 | Update::RemoveItem { .. }
134 | Update::DetachLastItems { .. }
135 | Update::StartReattachItems
136 | Update::EndReattachItems
137 | Update::Clear => vec![],
138 Update::PushItems { items, .. } => items,
139 Update::ReplaceItem { item, .. } => vec![item],
140 }
141 }
142}
143
144/// A collection of [`Update`]s that can be observed.
145///
146/// Get a value for this type with [`LinkedChunk::updates`].
147///
148/// All clones of this type share the same data.
149///
150/// [`LinkedChunk::updates`]: super::LinkedChunk::updates
151#[derive(Debug)]
152pub struct ObservableUpdates<Item, Gap> {
153 pub(super) inner: Arc<RwLock<UpdatesInner<Item, Gap>>>,
154}
155
156impl<Item, Gap> ObservableUpdates<Item, Gap> {
157 /// Create a new [`ObservableUpdates`].
158 pub(super) fn new() -> Self {
159 Self { inner: Arc::new(RwLock::new(UpdatesInner::new())) }
160 }
161
162 /// Push a new update.
163 pub(super) fn push(&mut self, update: Update<Item, Gap>) {
164 self.inner.write().unwrap().push(update);
165 }
166
167 /// Clear all pending updates.
168 pub(super) fn clear_pending(&mut self) {
169 self.inner.write().unwrap().clear_pending();
170 }
171
172 /// Take new updates.
173 ///
174 /// Updates that have been taken will not be read again.
175 pub fn take(&mut self) -> Vec<Update<Item, Gap>>
176 where
177 Item: Clone,
178 Gap: Clone,
179 {
180 self.inner.write().unwrap().take().to_owned()
181 }
182
183 /// Subscribe to updates by using a [`Stream`].
184 pub fn subscribe(&mut self) -> UpdatesSubscriber<Item, Gap> {
185 // A subscriber is a new update reader, it needs its own token.
186 let token = self.new_reader_token();
187
188 UpdatesSubscriber::new(Arc::downgrade(&self.inner), token)
189 }
190
191 /// Generate a new [`ReaderToken`].
192 pub(super) fn new_reader_token(&mut self) -> ReaderToken {
193 let mut inner = self.inner.write().unwrap();
194
195 // Add 1 before reading the `last_token`, in this particular order,
196 // because the 0 token is reserved by `MAIN_READER_TOKEN`.
197 inner.last_token += 1;
198 let last_token = inner.last_token;
199
200 inner.last_index_per_reader.insert(last_token, 0);
201
202 last_token
203 }
204
205 /// Create a new [`ObservableUpdatesPusher`], privately.
206 pub(super) fn new_pusher(&self) -> ObservableUpdatesPusher<Item, Gap> {
207 ObservableUpdatesPusher { inner: self.inner.clone() }
208 }
209}
210
211/// This type is similar to [`ObservableUpdates`] except it has a single `push`
212/// method which takes a `&self` instead of a `&mut self` to accommodate a
213/// particular need in `Ends` for lazily get the first chunk.
214pub(super) struct ObservableUpdatesPusher<Item, Gap> {
215 inner: Arc<RwLock<UpdatesInner<Item, Gap>>>,
216}
217
218impl<Item, Gap> ObservableUpdatesPusher<Item, Gap> {
219 /// Push a new update, even if `&self` while we could expect a `&mut self`.
220 pub fn push(&self, update: Update<Item, Gap>) {
221 self.inner.write().unwrap().push(update);
222 }
223}
224
225/// A token used to represent readers that read the updates in [`UpdatesInner`].
226pub(super) type ReaderToken = usize;
227
228/// Inner type for [`ObservableUpdates`].
229///
230/// The particularity of this type is that multiple readers can read the
231/// updates. A reader has a [`ReaderToken`]. The public API (i.e.
232/// [`ObservableUpdates`]) is considered to be the _main reader_ (it has the
233/// token [`Self::MAIN_READER_TOKEN`]).
234///
235/// An update that have been read by all readers are garbage collected to be
236/// removed from the memory. An update will never be read twice by the same
237/// reader.
238///
239/// Why do we need multiple readers? The public API reads the updates with
240/// [`ObservableUpdates::take`], but the private API must also read the updates
241/// for example with [`UpdatesSubscriber`]. Of course, they can be multiple
242/// `UpdatesSubscriber`s at the same time. Hence the need of supporting multiple
243/// readers.
244#[derive(Debug)]
245pub(super) struct UpdatesInner<Item, Gap> {
246 /// All the updates that have not been read by all readers.
247 updates: Vec<Update<Item, Gap>>,
248
249 /// Updates are stored in [`Self::updates`]. Multiple readers can read them.
250 /// A reader is identified by a [`ReaderToken`].
251 ///
252 /// To each reader token is associated an index that represents the index of
253 /// the last reading. It is used to never return the same update twice.
254 last_index_per_reader: HashMap<ReaderToken, usize>,
255
256 /// The last generated token. This is useful to generate new token.
257 last_token: ReaderToken,
258
259 /// Pending wakers for [`UpdateSubscriber`]s. A waker is removed every time
260 /// it is called.
261 wakers: Vec<Waker>,
262}
263
264impl<Item, Gap> UpdatesInner<Item, Gap> {
265 /// The token used by the main reader. See [`Self::take`] to learn more.
266 const MAIN_READER_TOKEN: ReaderToken = 0;
267
268 /// Create a new [`Self`].
269 fn new() -> Self {
270 Self {
271 updates: Vec::with_capacity(8),
272 last_index_per_reader: {
273 let mut map = HashMap::with_capacity(2);
274 map.insert(Self::MAIN_READER_TOKEN, 0);
275
276 map
277 },
278 last_token: Self::MAIN_READER_TOKEN,
279 wakers: Vec::with_capacity(2),
280 }
281 }
282
283 /// Push a new update.
284 fn push(&mut self, update: Update<Item, Gap>) {
285 self.updates.push(update);
286
287 // Wake them up \o/.
288 for waker in self.wakers.drain(..) {
289 waker.wake();
290 }
291 }
292
293 /// Clear all pending updates.
294 fn clear_pending(&mut self) {
295 self.updates.clear();
296
297 // Reset all the per-reader indices.
298 for idx in self.last_index_per_reader.values_mut() {
299 *idx = 0;
300 }
301
302 // No need to wake the wakers; they're waiting for a new update, and we
303 // just made them all disappear.
304 }
305
306 /// Take new updates; it considers the caller is the main reader, i.e. it
307 /// will use the [`Self::MAIN_READER_TOKEN`].
308 ///
309 /// Updates that have been read will never be read again by the current
310 /// reader.
311 ///
312 /// Learn more by reading [`Self::take_with_token`].
313 fn take(&mut self) -> &[Update<Item, Gap>] {
314 self.take_with_token(Self::MAIN_READER_TOKEN)
315 }
316
317 /// Take new updates with a particular reader token.
318 ///
319 /// Updates are stored in [`Self::updates`]. Multiple readers can read them.
320 /// A reader is identified by a [`ReaderToken`]. Every reader can
321 /// take/read/consume each update only once. An internal index is stored per
322 /// reader token to know where to start reading updates next time this
323 /// method is called.
324 pub(super) fn take_with_token(&mut self, token: ReaderToken) -> &[Update<Item, Gap>] {
325 // Let's garbage collect unused updates.
326 self.garbage_collect();
327
328 let index = self
329 .last_index_per_reader
330 .get_mut(&token)
331 .expect("Given `UpdatesToken` does not map to any index");
332
333 // Read new updates, and update the index.
334 let slice = &self.updates[*index..];
335 *index = self.updates.len();
336
337 slice
338 }
339
340 /// Has the given reader, identified by its [`ReaderToken`], some pending
341 /// updates, or has it consumed all the pending updates?
342 pub(super) fn is_reader_up_to_date(&self, token: ReaderToken) -> bool {
343 *self.last_index_per_reader.get(&token).expect("unknown reader token") == self.updates.len()
344 }
345
346 /// Return the number of updates in the buffer.
347 #[cfg(test)]
348 fn len(&self) -> usize {
349 self.updates.len()
350 }
351
352 /// Garbage collect unused updates. An update is considered unused when it's
353 /// been read by all readers.
354 ///
355 /// Basically, it reduces to finding the smallest last index for all
356 /// readers, and clear from 0 to that index.
357 fn garbage_collect(&mut self) {
358 let min_index = self.last_index_per_reader.values().min().copied().unwrap_or(0);
359
360 if min_index > 0 {
361 let _ = self.updates.drain(0..min_index);
362
363 // Let's shift the indices to the left by `min_index` to preserve
364 // them.
365 for index in self.last_index_per_reader.values_mut() {
366 *index -= min_index;
367 }
368 }
369 }
370}
371
372/// A subscriber to [`ObservableUpdates`]. It is helpful to receive updates via
373/// a [`Stream`].
374#[derive(Debug)]
375pub struct UpdatesSubscriber<Item, Gap> {
376 /// Weak reference to [`UpdatesInner`].
377 ///
378 /// Using a weak reference allows [`ObservableUpdates`] to be dropped freely
379 /// even if a subscriber exists.
380 updates: Weak<RwLock<UpdatesInner<Item, Gap>>>,
381
382 /// The token to read the updates.
383 token: ReaderToken,
384}
385
386impl<Item, Gap> UpdatesSubscriber<Item, Gap> {
387 /// Create a new [`Self`].
388 fn new(updates: Weak<RwLock<UpdatesInner<Item, Gap>>>, token: ReaderToken) -> Self {
389 Self { updates, token }
390 }
391}
392
393impl<Item, Gap> Stream for UpdatesSubscriber<Item, Gap>
394where
395 Item: Clone,
396 Gap: Clone,
397{
398 type Item = Vec<Update<Item, Gap>>;
399
400 fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
401 let Some(updates) = self.updates.upgrade() else {
402 // The `ObservableUpdates` has been dropped. It's time to close this
403 // stream.
404 return Poll::Ready(None);
405 };
406
407 let mut updates = updates.write().unwrap();
408 let the_updates = updates.take_with_token(self.token);
409
410 // No updates.
411 if the_updates.is_empty() {
412 // Let's register the waker.
413 updates.wakers.push(context.waker().clone());
414
415 // The stream is pending.
416 return Poll::Pending;
417 }
418
419 // There is updates! Let's forward them in this stream.
420 Poll::Ready(Some(the_updates.to_owned()))
421 }
422}
423
424impl<Item, Gap> Drop for UpdatesSubscriber<Item, Gap> {
425 fn drop(&mut self) {
426 // Remove `Self::token` from `UpdatesInner::last_index_per_reader`. This
427 // is important so that the garbage collector can do its jobs correctly
428 // without a dead dangling reader token.
429 if let Some(updates) = self.updates.upgrade() {
430 let mut updates = updates.write().unwrap();
431
432 // Remove the reader token from `UpdatesInner`. It's safe to ignore
433 // the result of `remove` here: `None` means the token was already
434 // removed (note: it should be unreachable).
435 let _ = updates.last_index_per_reader.remove(&self.token);
436 }
437 }
438}
439
440#[cfg(test)]
441mod tests {
442 use std::{
443 sync::{Arc, Mutex},
444 task::{Context, Poll, Wake},
445 };
446
447 use assert_matches::assert_matches;
448 use futures_core::Stream;
449 use futures_util::pin_mut;
450
451 use super::{super::LinkedChunk, ChunkIdentifier, Position, UpdatesInner};
452 use crate::linked_chunk::Update;
453
454 #[test]
455 fn test_updates_take_and_garbage_collector() {
456 use super::Update::*;
457
458 let mut linked_chunk = LinkedChunk::<10, char, ()>::new_with_update_history();
459
460 // Simulate another updates “reader”, it can a subscriber.
461 let main_token = UpdatesInner::<char, ()>::MAIN_READER_TOKEN;
462 let other_token = {
463 let updates = linked_chunk.updates().unwrap();
464 let mut inner = updates.inner.write().unwrap();
465 inner.last_token += 1;
466
467 let other_token = inner.last_token;
468 inner.last_index_per_reader.insert(other_token, 0);
469
470 other_token
471 };
472
473 // Let's trigger the chunk creation to simplify the test.
474 let _ = linked_chunk.first_chunk();
475
476 // There is an update.
477 {
478 let updates = linked_chunk.updates().unwrap();
479
480 assert_eq!(
481 updates.take(),
482 &[NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None }],
483 );
484 assert_eq!(
485 updates.inner.write().unwrap().take_with_token(other_token),
486 &[NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None }],
487 );
488 }
489
490 // No new update.
491 {
492 let updates = linked_chunk.updates().unwrap();
493
494 assert!(updates.take().is_empty());
495 assert!(updates.inner.write().unwrap().take_with_token(other_token).is_empty());
496 }
497
498 linked_chunk.push_items_back(['a']);
499 linked_chunk.push_items_back(['b']);
500 linked_chunk.push_items_back(['c']);
501
502 // Scenario 1: “main” takes the new updates, “other” doesn't take the
503 // new updates.
504 //
505 // 0 1 2 3
506 // +---+---+---+
507 // | a | b | c |
508 // +---+---+---+
509 //
510 // “main” will move its index from 0 to 3. “other” won't move its index.
511 {
512 let updates = linked_chunk.updates().unwrap();
513
514 {
515 // Inspect number of updates in memory.
516 assert_eq!(updates.inner.read().unwrap().len(), 3);
517 }
518
519 assert_eq!(
520 updates.take(),
521 &[
522 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] },
523 PushItems { at: Position(ChunkIdentifier(0), 1), items: vec!['b'] },
524 PushItems { at: Position(ChunkIdentifier(0), 2), items: vec!['c'] },
525 ]
526 );
527
528 {
529 let inner = updates.inner.read().unwrap();
530
531 // Inspect number of updates in memory. It must be the same
532 // number as before as the garbage collector weren't not able to
533 // remove any unused updates.
534 assert_eq!(inner.len(), 3);
535
536 // Inspect the indices.
537 let indices = &inner.last_index_per_reader;
538
539 assert_eq!(indices.get(&main_token), Some(&3));
540 assert_eq!(indices.get(&other_token), Some(&0));
541 }
542 }
543
544 linked_chunk.push_items_back(['d']);
545 linked_chunk.push_items_back(['e']);
546 linked_chunk.push_items_back(['f']);
547
548 // Scenario 2: “other“ takes the new updates, “main” doesn't take the
549 // new updates.
550 //
551 // 0 1 2 3 4 5 6
552 // +---+---+---+---+---+---+
553 // | a | b | c | d | e | f |
554 // +---+---+---+---+---+---+
555 //
556 // “main” won't move its index. “other” will move its index from 0 to 6.
557 {
558 let updates = linked_chunk.updates().unwrap();
559
560 assert_eq!(
561 updates.inner.write().unwrap().take_with_token(other_token),
562 &[
563 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] },
564 PushItems { at: Position(ChunkIdentifier(0), 1), items: vec!['b'] },
565 PushItems { at: Position(ChunkIdentifier(0), 2), items: vec!['c'] },
566 PushItems { at: Position(ChunkIdentifier(0), 3), items: vec!['d'] },
567 PushItems { at: Position(ChunkIdentifier(0), 4), items: vec!['e'] },
568 PushItems { at: Position(ChunkIdentifier(0), 5), items: vec!['f'] },
569 ]
570 );
571
572 {
573 let inner = updates.inner.read().unwrap();
574
575 // Inspect number of updates in memory. It must be the same
576 // number as before as the garbage collector will be able to
577 // remove unused updates but at the next call…
578 assert_eq!(inner.len(), 6);
579
580 // Inspect the indices.
581 let indices = &inner.last_index_per_reader;
582
583 assert_eq!(indices.get(&main_token), Some(&3));
584 assert_eq!(indices.get(&other_token), Some(&6));
585 }
586 }
587
588 // Scenario 3: “other” take new updates, but there is none, “main”
589 // doesn't take new updates. The garbage collector will run and collect
590 // unused updates.
591 //
592 // 0 1 2 3
593 // +---+---+---+
594 // | d | e | f |
595 // +---+---+---+
596 //
597 // “main” will have its index updated from 3 to 0. “other” will have its
598 // index updated from 6 to 3.
599 {
600 let updates = linked_chunk.updates().unwrap();
601
602 assert!(updates.inner.write().unwrap().take_with_token(other_token).is_empty());
603
604 {
605 let inner = updates.inner.read().unwrap();
606
607 // Inspect number of updates in memory. The garbage collector
608 // has removed unused updates.
609 assert_eq!(inner.len(), 3);
610
611 // Inspect the indices. They must have been adjusted.
612 let indices = &inner.last_index_per_reader;
613
614 assert_eq!(indices.get(&main_token), Some(&0));
615 assert_eq!(indices.get(&other_token), Some(&3));
616 }
617 }
618
619 linked_chunk.push_items_back(['g']);
620 linked_chunk.push_items_back(['h']);
621 linked_chunk.push_items_back(['i']);
622
623 // Scenario 4: both “main” and “other” take the new updates.
624 //
625 // 0 1 2 3 4 5 6
626 // +---+---+---+---+---+---+
627 // | d | e | f | g | h | i |
628 // +---+---+---+---+---+---+
629 //
630 // “main” will have its index updated from 0 to 3. “other” will have its
631 // index updated from 6 to 3.
632 {
633 let updates = linked_chunk.updates().unwrap();
634
635 assert_eq!(
636 updates.take(),
637 &[
638 PushItems { at: Position(ChunkIdentifier(0), 3), items: vec!['d'] },
639 PushItems { at: Position(ChunkIdentifier(0), 4), items: vec!['e'] },
640 PushItems { at: Position(ChunkIdentifier(0), 5), items: vec!['f'] },
641 PushItems { at: Position(ChunkIdentifier(0), 6), items: vec!['g'] },
642 PushItems { at: Position(ChunkIdentifier(0), 7), items: vec!['h'] },
643 PushItems { at: Position(ChunkIdentifier(0), 8), items: vec!['i'] },
644 ]
645 );
646 assert_eq!(
647 updates.inner.write().unwrap().take_with_token(other_token),
648 &[
649 PushItems { at: Position(ChunkIdentifier(0), 6), items: vec!['g'] },
650 PushItems { at: Position(ChunkIdentifier(0), 7), items: vec!['h'] },
651 PushItems { at: Position(ChunkIdentifier(0), 8), items: vec!['i'] },
652 ]
653 );
654
655 {
656 let inner = updates.inner.read().unwrap();
657
658 // Inspect number of updates in memory. The garbage collector
659 // had a chance to collect the first 3 updates.
660 assert_eq!(inner.len(), 3);
661
662 // Inspect the indices.
663 let indices = &inner.last_index_per_reader;
664
665 assert_eq!(indices.get(&main_token), Some(&3));
666 assert_eq!(indices.get(&other_token), Some(&3));
667 }
668 }
669
670 // Scenario 5: no more updates but they both try to take new updates.
671 // The garbage collector will collect all updates as all of them as been
672 // read already.
673 //
674 // “main” will have its index updated from 0 to 0. “other” will have its
675 // index updated from 3 to 0.
676 {
677 let updates = linked_chunk.updates().unwrap();
678
679 assert!(updates.take().is_empty());
680 assert!(updates.inner.write().unwrap().take_with_token(other_token).is_empty());
681
682 {
683 let inner = updates.inner.read().unwrap();
684
685 // Inspect number of updates in memory. The garbage collector
686 // had a chance to collect all updates.
687 assert_eq!(inner.len(), 0);
688
689 // Inspect the indices.
690 let indices = &inner.last_index_per_reader;
691
692 assert_eq!(indices.get(&main_token), Some(&0));
693 assert_eq!(indices.get(&other_token), Some(&0));
694 }
695 }
696 }
697
698 struct CounterWaker {
699 number_of_wakeup: Mutex<usize>,
700 }
701
702 impl Wake for CounterWaker {
703 fn wake(self: Arc<Self>) {
704 *self.number_of_wakeup.lock().unwrap() += 1;
705 }
706 }
707
708 #[test]
709 fn test_updates_stream() {
710 use super::Update::*;
711
712 let counter_waker = Arc::new(CounterWaker { number_of_wakeup: Mutex::new(0) });
713 let waker = counter_waker.clone().into();
714 let mut context = Context::from_waker(&waker);
715
716 let mut linked_chunk = LinkedChunk::<3, char, ()>::new_with_update_history();
717
718 let updates_subscriber = linked_chunk.updates().unwrap().subscribe();
719 pin_mut!(updates_subscriber);
720
721 // No initial update, stream is pending.
722 assert_matches!(updates_subscriber.as_mut().poll_next(&mut context), Poll::Pending);
723 assert_eq!(*counter_waker.number_of_wakeup.lock().unwrap(), 0);
724
725 // Let's generate an update.
726 linked_chunk.push_items_back(['a']);
727
728 // The waker must have been called.
729 assert_eq!(*counter_waker.number_of_wakeup.lock().unwrap(), 1);
730
731 // There is an update! Right after that, the stream is pending again.
732 assert_matches!(
733 updates_subscriber.as_mut().poll_next(&mut context),
734 Poll::Ready(Some(items)) => {
735 assert_eq!(
736 items,
737 &[
738 NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None },
739 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] }
740 ]
741 );
742 }
743 );
744 assert_matches!(updates_subscriber.as_mut().poll_next(&mut context), Poll::Pending);
745
746 // Let's generate two other updates.
747 linked_chunk.push_items_back(['b']);
748 linked_chunk.push_items_back(['c']);
749
750 // The waker must have been called only once for the two updates.
751 assert_eq!(*counter_waker.number_of_wakeup.lock().unwrap(), 2);
752
753 // We can consume the updates without the stream, but the stream
754 // continues to know it has updates.
755 assert_eq!(
756 linked_chunk.updates().unwrap().take(),
757 &[
758 NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None },
759 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] },
760 PushItems { at: Position(ChunkIdentifier(0), 1), items: vec!['b'] },
761 PushItems { at: Position(ChunkIdentifier(0), 2), items: vec!['c'] },
762 ]
763 );
764 assert_matches!(
765 updates_subscriber.as_mut().poll_next(&mut context),
766 Poll::Ready(Some(items)) => {
767 assert_eq!(
768 items,
769 &[
770 PushItems { at: Position(ChunkIdentifier(0), 1), items: vec!['b'] },
771 PushItems { at: Position(ChunkIdentifier(0), 2), items: vec!['c'] },
772 ]
773 );
774 }
775 );
776 assert_matches!(updates_subscriber.as_mut().poll_next(&mut context), Poll::Pending);
777
778 // When dropping the `LinkedChunk`, it closes the stream.
779 drop(linked_chunk);
780 assert_matches!(updates_subscriber.as_mut().poll_next(&mut context), Poll::Ready(None));
781
782 // Wakers calls have not changed.
783 assert_eq!(*counter_waker.number_of_wakeup.lock().unwrap(), 2);
784 }
785
786 #[test]
787 fn test_updates_multiple_streams() {
788 use super::Update::*;
789
790 let counter_waker1 = Arc::new(CounterWaker { number_of_wakeup: Mutex::new(0) });
791 let counter_waker2 = Arc::new(CounterWaker { number_of_wakeup: Mutex::new(0) });
792
793 let waker1 = counter_waker1.clone().into();
794 let waker2 = counter_waker2.clone().into();
795
796 let mut context1 = Context::from_waker(&waker1);
797 let mut context2 = Context::from_waker(&waker2);
798
799 let mut linked_chunk = LinkedChunk::<3, char, ()>::new_with_update_history();
800
801 let updates_subscriber1 = linked_chunk.updates().unwrap().subscribe();
802 pin_mut!(updates_subscriber1);
803
804 // Scope for `updates_subscriber2`.
805 let updates_subscriber2_token = {
806 let updates_subscriber2 = linked_chunk.updates().unwrap().subscribe();
807 pin_mut!(updates_subscriber2);
808
809 // No initial updates, streams are pending.
810 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Pending);
811 assert_eq!(*counter_waker1.number_of_wakeup.lock().unwrap(), 0);
812
813 assert_matches!(updates_subscriber2.as_mut().poll_next(&mut context2), Poll::Pending);
814 assert_eq!(*counter_waker2.number_of_wakeup.lock().unwrap(), 0);
815
816 // Let's generate an update.
817 linked_chunk.push_items_back(['a']);
818
819 // The wakers must have been called.
820 assert_eq!(*counter_waker1.number_of_wakeup.lock().unwrap(), 1);
821 assert_eq!(*counter_waker2.number_of_wakeup.lock().unwrap(), 1);
822
823 // There is an update! Right after that, the streams are pending
824 // again.
825 assert_matches!(
826 updates_subscriber1.as_mut().poll_next(&mut context1),
827 Poll::Ready(Some(items)) => {
828 assert_eq!(
829 items,
830 &[
831 NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None },
832 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] }
833 ]
834 );
835 }
836 );
837 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Pending);
838 assert_matches!(
839 updates_subscriber2.as_mut().poll_next(&mut context2),
840 Poll::Ready(Some(items)) => {
841 assert_eq!(
842 items,
843 &[
844 NewItemsChunk { previous: None, new: ChunkIdentifier(0), next: None },
845 PushItems { at: Position(ChunkIdentifier(0), 0), items: vec!['a'] }
846 ]
847 );
848 }
849 );
850 assert_matches!(updates_subscriber2.as_mut().poll_next(&mut context2), Poll::Pending);
851
852 // Let's generate two other updates.
853 linked_chunk.push_items_back(['b']);
854 linked_chunk.push_items_back(['c']);
855
856 // A waker is consumed when called. The first call to
857 // `push_items_back` will call and consume the wakers. The second
858 // call to `push_items_back` will do nothing as the wakers have been
859 // consumed. New wakers will be registered on polling.
860 //
861 // So, the waker must have been called only once for the two
862 // updates.
863 assert_eq!(*counter_waker1.number_of_wakeup.lock().unwrap(), 2);
864 assert_eq!(*counter_waker2.number_of_wakeup.lock().unwrap(), 2);
865
866 // Let's poll `updates_subscriber1` only.
867 assert_matches!(
868 updates_subscriber1.as_mut().poll_next(&mut context1),
869 Poll::Ready(Some(items)) => {
870 assert_eq!(
871 items,
872 &[
873 PushItems { at: Position(ChunkIdentifier(0), 1), items: vec!['b'] },
874 PushItems { at: Position(ChunkIdentifier(0), 2), items: vec!['c'] },
875 ]
876 );
877 }
878 );
879 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Pending);
880
881 // For the sake of this test, we also need to advance the main
882 // reader token.
883 let _ = linked_chunk.updates().unwrap().take();
884 let _ = linked_chunk.updates().unwrap().take();
885
886 // If we inspect the garbage collector state, `a`, `b` and `c`
887 // should still be present because not all of them have been
888 // consumed by `updates_subscriber2` yet.
889 {
890 let updates = linked_chunk.updates().unwrap();
891
892 let inner = updates.inner.read().unwrap();
893
894 // Inspect number of updates in memory. We get 2 because the
895 // garbage collector runs before data are taken, not after:
896 // `updates_subscriber2` has read `a` only, so `b` and `c`
897 // remain.
898 assert_eq!(inner.len(), 2);
899
900 // Inspect the indices.
901 let indices = &inner.last_index_per_reader;
902
903 assert_eq!(indices.get(&updates_subscriber1.token), Some(&2));
904 assert_eq!(indices.get(&updates_subscriber2.token), Some(&0));
905 }
906
907 // Poll `updates_subscriber1` again: there is no new update so it
908 // must be pending.
909 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Pending);
910
911 // The state of the garbage collector is unchanged: `a`, `b` and `c`
912 // are still in memory.
913 {
914 let updates = linked_chunk.updates().unwrap();
915
916 let inner = updates.inner.read().unwrap();
917
918 // Inspect number of updates in memory. Value is unchanged.
919 assert_eq!(inner.len(), 2);
920
921 // Inspect the indices. They are unchanged.
922 let indices = &inner.last_index_per_reader;
923
924 assert_eq!(indices.get(&updates_subscriber1.token), Some(&2));
925 assert_eq!(indices.get(&updates_subscriber2.token), Some(&0));
926 }
927
928 updates_subscriber2.token
929 // Drop `updates_subscriber2`!
930 };
931
932 // `updates_subscriber2` has been dropped. Poll `updates_subscriber1`
933 // again: still no new update, but it will run the garbage collector
934 // again, and this time `updates_subscriber2` is not “retaining” `b` and
935 // `c`. The garbage collector must be empty.
936 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Pending);
937
938 // Inspect the garbage collector.
939 {
940 let updates = linked_chunk.updates().unwrap();
941
942 let inner = updates.inner.read().unwrap();
943
944 // Inspect number of updates in memory.
945 assert_eq!(inner.len(), 0);
946
947 // Inspect the indices.
948 let indices = &inner.last_index_per_reader;
949
950 assert_eq!(indices.get(&updates_subscriber1.token), Some(&0));
951 assert_eq!(indices.get(&updates_subscriber2_token), None); // token is unknown!
952 }
953
954 // When dropping the `LinkedChunk`, it closes the stream.
955 drop(linked_chunk);
956 assert_matches!(updates_subscriber1.as_mut().poll_next(&mut context1), Poll::Ready(None));
957 }
958
959 #[test]
960 fn test_update_into_items() {
961 let updates: Update<_, u32> =
962 Update::PushItems { at: Position::new(ChunkIdentifier(0), 0), items: vec![1, 2, 3] };
963
964 assert_eq!(updates.into_items(), vec![1, 2, 3]);
965
966 let updates: Update<u32, u32> = Update::Clear;
967 assert!(updates.into_items().is_empty());
968
969 let updates: Update<u32, u32> =
970 Update::RemoveItem { at: Position::new(ChunkIdentifier(0), 0) };
971 assert!(updates.into_items().is_empty());
972
973 let updates: Update<u32, u32> =
974 Update::ReplaceItem { at: Position::new(ChunkIdentifier(0), 0), item: 42 };
975 assert_eq!(updates.into_items(), vec![42]);
976 }
977}