|
| 1 | +use crate::protocol::{ |
| 2 | + Envelope, EnvelopeError, OrderedEnvelopeCollection, ResolveDependencies, Sort, sort, |
| 3 | +}; |
| 4 | +use derive_builder::Builder; |
| 5 | +use xmtp_proto::types::TopicCursor; |
| 6 | + |
| 7 | +#[derive(Debug, Clone, Builder)] |
| 8 | +#[builder(setter(strip_option), build_fn(error = "EnvelopeError"))] |
| 9 | +pub struct Ordered<T, R> { |
| 10 | + envelopes: Vec<T>, |
| 11 | + resolver: R, |
| 12 | +} |
| 13 | + |
| 14 | +impl<T: Clone, R: Clone> Ordered<T, R> { |
| 15 | + pub fn builder() -> OrderedBuilder<T, R> { |
| 16 | + OrderedBuilder::default() |
| 17 | + } |
| 18 | +} |
| 19 | + |
| 20 | +#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)] |
| 21 | +#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))] |
| 22 | +impl<T, R> OrderedEnvelopeCollection for Ordered<T, R> |
| 23 | +where |
| 24 | + T: Envelope<'static>, |
| 25 | + R: ResolveDependencies<ResolvedEnvelope = T>, |
| 26 | +{ |
| 27 | + /// Sort dependencies of `Self` according to XIP |
| 28 | + /// TODO fill in docs |
| 29 | + async fn sort(&mut self) -> Result<(), EnvelopeError> { |
| 30 | + let mut topic_cursor = TopicCursor::default(); |
| 31 | + { |
| 32 | + sort::timestamp(&mut self.envelopes).sort()?; |
| 33 | + } |
| 34 | + while let Some(missing) = sort::causal(&mut self.envelopes, &mut topic_cursor).sort()? { |
| 35 | + let cursors = missing |
| 36 | + .iter() |
| 37 | + .map(|e| e.cursor()) |
| 38 | + .collect::<Result<Vec<_>, _>>()?; |
| 39 | + // try to resolve the missing dependencies |
| 40 | + let resolved = self.resolver.resolve(cursors).await?; |
| 41 | + // re-apply at start of vec |
| 42 | + self.envelopes.splice(0..0, resolved.into_iter()); |
| 43 | + sort::timestamp(&mut self.envelopes).sort()?; |
| 44 | + } |
| 45 | + Ok(()) |
| 46 | + } |
| 47 | +} |
0 commit comments