Skip to main content

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}