Skip to main content

ossa_core/store/
mod.rs

1use itertools::Itertools;
2use ossa_crdt::CRDT;
3use ossa_typeable::{TypeId, Typeable};
4use rand::{seq::SliceRandom as _, thread_rng};
5use replace_with::replace_with_or_abort;
6use serde::{Deserialize, Serialize};
7use std::fmt::Debug;
8use std::{
9    collections::{BTreeMap, BTreeSet},
10    ops::Range,
11};
12use tokio::{
13    sync::{
14        mpsc::{UnboundedReceiver, UnboundedSender},
15        oneshot::{self, Sender},
16    },
17    task::JoinHandle,
18};
19use tracing::{debug, error, warn};
20
21use crate::store::v0::BLOCK_SIZE;
22use crate::time::ConcretizeTime;
23use crate::util::merkle_tree::{MerkleTree, Potential};
24use crate::{
25    auth::DeviceId,
26    core::{OssaType, SharedState},
27    network::{
28        multiplexer::{run_miniprotocol_async, SpawnMultiplexerTask},
29        protocol::MiniProtocol,
30    },
31    protocol::{
32        manager::v0::PeerManagerCommand,
33        store_peer::v0::{MsgStoreSyncRequest, StoreSync, StoreSyncCommand},
34    },
35    store::{
36        ecg::{ECGBody, ECGHeader, RawECGBody},
37        v0::{BLOCK_REQUEST_LIMIT, MERKLE_REQUEST_LIMIT},
38    },
39    util::{self, compress_consecutive_into_ranges},
40};
41
42pub mod ecg;
43pub mod v0; // TODO: Move this to network::protocol
44
45pub use v0::{MetadataBody, MetadataHeader, Nonce};
46
47pub struct State<StoreId, Header: ecg::ECGHeader, T: CRDT, Hash> {
48    // Peers that also have this store (that we are potentially connected to?).
49    peers: BTreeMap<DeviceId, PeerInfo<Header::HeaderId, Header>>, // BTreeSet<DeviceId>,
50    state_machine: StateMachine<StoreId, Header, T, Hash>,
51    metadata_subscribers: BTreeMap<DeviceId, oneshot::Sender<Option<v0::MetadataHeader<Hash>>>>,
52    merkle_subscribers: BTreeMap<DeviceId, (Vec<Range<u64>>, oneshot::Sender<Option<Vec<Hash>>>)>,
53    block_subscribers: BTreeMap<
54        DeviceId,
55        (
56            Vec<Range<u64>>,
57            oneshot::Sender<Option<Vec<Option<Vec<u8>>>>>,
58        ),
59    >,
60    ecg_subscribers:
61        BTreeMap<DeviceId, oneshot::Sender<ecg::UntypedState<Header::HeaderId, Header>>>,
62    // listeners: Vec<UnboundedSender<StateUpdate<Header, T>>>,
63}
64
65// States are:
66// - Initializing - Setting up the thread that owns the store (not defined here).
67// - DownloadingMetadata - Don't have the header so we're downloading it.
68// - Syncing - Have the header and syncing updates between peers.
69pub enum StateMachine<StoreId, Header: ecg::ECGHeader, T: CRDT, Hash> {
70    DownloadingMetadata {
71        store_id: StoreId,
72    },
73    /// Invariant: At least one partial_merkle_tree node is None.
74    DownloadingMerkle {
75        metadata: MetadataHeader<Hash>,
76        partial_merkle_tree: MerkleTree<Potential<Hash>>,
77    },
78    /// Invariant: At least one initial_state block is None.
79    DownloadingInitialState {
80        metadata: MetadataHeader<Hash>,
81        merkle_tree: MerkleTree<Hash>,
82        initial_state: Vec<Option<Vec<u8>>>,
83    },
84    Syncing {
85        metadata: MetadataHeader<Hash>,
86        merkle_tree: MerkleTree<Hash>,
87        initial_state: Vec<u8>, // Or just T?
88        ecg_state: ecg::State<Header, T>,
89        decrypted_state: DecryptedState<Header, T>, // Temporary
90                                                    // decrypted_state: Option<DecryptedState<Header, T>>, // JP: Is this actually used?
91                                                    // Does it make sense?
92    },
93}
94
95pub struct DecryptedState<Header: ecg::ECGHeader, T: CRDT> {
96    /// Latest ECG application state we've seen.
97    latest_state: T,
98
99    /// Headers corresponding to the latest ECG application state.
100    // TODO: Remove this.
101    latest_headers: BTreeSet<Header::HeaderId>,
102}
103
104/// Information about a peer.
105#[derive(Debug)]
106struct PeerInfo<HeaderId, Header> {
107    /// Status of incoming sync status from peer.
108    incoming_status: PeerStatus<()>,
109    /// Status of outgoing sync status to peer.
110    outgoing_status: PeerStatus<OutgoingPeerStatus<HeaderId, Header>>,
111    // ecg_status: ECGStatus<HeaderId>,
112}
113
114// #[derive(Clone, Debug)]
115// /// Information about a peer's ECG status.
116// pub(crate) struct ECGStatus<HeaderId> {
117//     /// Greatest common ancestor between our ECG graphs.
118//     pub(crate) meet: Vec<HeaderId>,
119//     /// Whether we need to update the meet between us and this peer.
120//     pub(crate) meet_needs_update: bool,
121//     // JP: Track their_tip?
122// }
123
124impl<Hash, Header> PeerInfo<Hash, Header> {
125    /// Checks if the peer is ready for a sync request.
126    /// This means the peer is syncing and does not have an outstanding request.
127    fn is_ready_for_sync(&self) -> bool {
128        if let PeerStatus::Syncing(s) = &self.outgoing_status {
129            !s.is_outstanding
130        } else {
131            false
132        }
133    }
134}
135
136#[derive(Debug)]
137/// Outgoing information about a syncing peer.
138struct OutgoingPeerStatus<HeaderId, Header> {
139    /// Sender channel for requests to the outgoing peer store.
140    sender_peer: UnboundedSender<StoreSyncCommand<HeaderId, Header>>,
141    /// Whether we have an outgoing peer request that is outstanding.
142    is_outstanding: bool,
143}
144
145/// Status of peers who we are potentially syncing this store with.
146#[derive(Debug)]
147pub(crate) enum PeerStatus<T> {
148    /// Peer is known and likely connected to, but is not syncing this store.
149    Known, // JP: Inactive?
150    /// Setting up the thread that syncs the store with the peer. It's possible that the peer will reject the sync request.
151    Initializing,
152    // {
153    //     task: JoinHandle<()>,
154    // },
155    /// The thread that syncs the store with the peer is syncing.
156    Syncing(T), // JP: Running instead?
157
158                //     /// Peers that we are connected to, but are not syncing this store. It's possible that these connections have dropped.
159                //     Known(), // TODO: Last known IP address, port, statistics (latency, bandwidth, ...). JP: Should some of this be stored globally?
160                //     /// Peers that we are connected to, but are not syncing this store. It's possible that these connections have dropped.
161                //     Connected(),
162                //     /// Peers that we are connected to and are syncing this store. It's possible that these connections have dropped.
163                //     Active(),
164}
165
166impl<T> PeerStatus<T> {
167    fn is_known(&self) -> bool {
168        if let PeerStatus::Known = self {
169            true
170        } else {
171            false
172        }
173    }
174
175    fn is_initializing(&self) -> bool {
176        if let PeerStatus::Initializing = self {
177            true
178        } else {
179            false
180        }
181    }
182}
183
184impl<
185        StoreId: Copy + Eq,
186        Header: ecg::ECGHeader + Clone + Debug,
187        T: CRDT + Clone,
188        Hash: util::Hash + Debug + Into<StoreId>,
189    > State<StoreId, Header, T, Hash>
190{
191    /// Initialize a new store with the given state. This initializes the header, including
192    /// generating a random nonce.
193    pub fn new_syncing(initial_state: T) -> State<StoreId, Header, T, Hash>
194    where
195        T: Serialize + Typeable,
196    {
197        let init_body = MetadataBody::new(&initial_state);
198        debug!("Initialized body: {:?}", init_body);
199        let store_header = MetadataHeader::generate::<T>(&init_body);
200        let decrypted_state = DecryptedState {
201            latest_state: initial_state,
202            latest_headers: BTreeSet::new(),
203        };
204
205        let (merkle_tree, initial_state) = init_body.build();
206
207        let state_machine = StateMachine::Syncing {
208            metadata: store_header,
209            merkle_tree,
210            initial_state,
211            ecg_state: ecg::State::new(),
212            decrypted_state, // : Some(decrypted_state),
213        };
214        State {
215            peers: BTreeMap::new(),
216            state_machine,
217            metadata_subscribers: BTreeMap::new(),
218            merkle_subscribers: BTreeMap::new(),
219            block_subscribers: BTreeMap::new(),
220            ecg_subscribers: BTreeMap::new(),
221        }
222    }
223
224    /// Create a new store with the given store id that is downloading the store's header.
225    pub(crate) fn new_downloading(store_id: StoreId) -> Self {
226        let state_machine = StateMachine::DownloadingMetadata { store_id };
227
228        State {
229            peers: BTreeMap::new(),
230            state_machine,
231            metadata_subscribers: BTreeMap::new(),
232            merkle_subscribers: BTreeMap::new(),
233            block_subscribers: BTreeMap::new(),
234            ecg_subscribers: BTreeMap::new(),
235        }
236    }
237
238    pub fn store_id(&self) -> StoreId {
239        match &self.state_machine {
240            StateMachine::DownloadingMetadata { store_id } => *store_id,
241            StateMachine::DownloadingMerkle { metadata, .. } => metadata.store_id(),
242            StateMachine::DownloadingInitialState { metadata, .. } => metadata.store_id(),
243            StateMachine::Syncing { metadata, .. } => metadata.store_id(),
244        }
245    }
246
247    /// Insert a peer as known if its status isn't already tracked by the store.
248    fn insert_known_peer(&mut self, peer: DeviceId) {
249        // let ecg_status = ECGStatus { meet: vec![], meet_needs_update: true };
250        self.peers
251            .entry(peer)
252            // .and_modify(|s| {
253            //     match s {
254            //         PeerStatus::Known => *s = PeerStatus::Known,
255            //         PeerStatus::Initializing => (),
256            //         PeerStatus::Syncing => (),
257            //     }
258            // })
259            // .or_insert(PeerStatus::Known);
260            .or_insert(PeerInfo {
261                incoming_status: PeerStatus::Known,
262                outgoing_status: PeerStatus::Known,
263            }); // , ecg_status});
264    }
265
266    /// Helper to update a known peer to initializing.
267    fn update_peer_to_initializing<A>(
268        &mut self,
269        peer: &DeviceId,
270        direction_lambda: fn(&mut PeerInfo<Header::HeaderId, Header>) -> &mut PeerStatus<A>,
271    ) where
272        A: Debug,
273    {
274        let Some(info) = self.peers.get_mut(peer) else {
275            error!(
276                "Invariant violated. Attempted to initialize an unknown peer: {}",
277                peer
278            );
279            panic!();
280        };
281        let status = direction_lambda(info);
282        if status.is_known() {
283            *status = PeerStatus::Initializing;
284        } else {
285            error!("Invariant violated. Attempted to initialize an already initialized peer: {} - {:?}", peer, status);
286            panic!();
287        }
288    }
289
290    /// Update a known peer's outgoing status to initializing.
291    fn update_peer_to_initializing_outgoing(&mut self, peer: &DeviceId) {
292        self.update_peer_to_initializing(peer, |info| &mut info.outgoing_status);
293    }
294
295    /// Update a known peer's incoming status to initializing.
296    fn update_peer_to_initializing_incoming(&mut self, peer: &DeviceId) {
297        self.update_peer_to_initializing(peer, |info| &mut info.incoming_status);
298    }
299
300    /// Helper to update a known peer to syncing.
301    fn update_peer_to_syncing<A>(
302        &mut self,
303        peer: &DeviceId,
304        direction_lambda: fn(&mut PeerInfo<Header::HeaderId, Header>) -> &mut PeerStatus<A>,
305        sender_m: A,
306    ) where
307        A: Debug,
308    {
309        let Some(info) = self.peers.get_mut(peer) else {
310            error!(
311                "Invariant violated. Attempted to initialize an unknown peer: {}",
312                peer
313            );
314            panic!();
315        };
316        let status = direction_lambda(info);
317        match status {
318            PeerStatus::Initializing => {
319                *status = PeerStatus::Syncing(sender_m);
320            }
321            PeerStatus::Known => {
322                error!(
323                    "Invariant violated. Attempted to initialize an unknown peer: {} - {:?}",
324                    peer, status
325                );
326                panic!();
327            }
328            PeerStatus::Syncing(_) => {
329                error!(
330                    "Invariant violated. Attempted to sync an already syncing peer: {} - {:?}",
331                    peer, status
332                );
333                panic!();
334            }
335        }
336    }
337
338    fn update_peer_to_syncing_incoming(&mut self, peer: &DeviceId) {
339        self.update_peer_to_syncing(peer, |info| &mut info.incoming_status, ());
340    }
341
342    fn update_peer_to_syncing_outgoing(
343        &mut self,
344        peer: &DeviceId,
345        sender: OutgoingPeerStatus<Header::HeaderId, Header>,
346    ) {
347        self.update_peer_to_syncing(peer, |info| &mut info.outgoing_status, sender);
348    }
349
350    fn update_outgoing_peer_to_ready(&mut self, peer: &DeviceId) {
351        let Some(info) = self.peers.get_mut(peer) else {
352            error!(
353                "Invariant violated. Attempted to update an unknown peer: {}",
354                peer
355            );
356            panic!();
357        };
358        match info.outgoing_status {
359            PeerStatus::Initializing => {
360                error!("Invariant violated. Attempted to update a peer that is initializing: {} - {:?}", peer, info.outgoing_status);
361                panic!();
362            }
363            PeerStatus::Known => {
364                error!(
365                    "Invariant violated. Attempted to update a peer that is not syncing: {} - {:?}",
366                    peer, info.outgoing_status
367                );
368                panic!();
369            }
370            PeerStatus::Syncing(ref mut status) => {
371                status.is_outstanding = false;
372            }
373        }
374    }
375
376    /// Send sync requests to peers.
377    fn send_sync_requests(&mut self) {
378        fn send_command<Hash, Header>(
379            i: &mut PeerInfo<Hash, Header>,
380            message: StoreSyncCommand<Hash, Header>,
381        ) {
382            let PeerStatus::Syncing(ref mut s) = i.outgoing_status else {
383                unreachable!("Already checked that the peer is ready.");
384            };
385
386            // Mark as outstanding.
387            s.is_outstanding = true;
388            s.sender_peer.send(message).expect("TODO");
389        }
390
391        // Get peers (of this store) without outstanding requests.
392        let mut peers: Vec<_> = self
393            .peers
394            .iter_mut()
395            .filter(|(_, i)| i.is_ready_for_sync())
396            .collect();
397        // TODO: Rank and (weighted) randomize the peers.
398        let mut rng = thread_rng();
399        peers.shuffle(&mut rng);
400
401        // Send requests to peers based on what we need.
402        match &self.state_machine {
403            StateMachine::DownloadingMetadata { .. } => {
404                // Tell the first peer store task to request the metadata.
405                peers.iter_mut().take(1).for_each(|(_, i)| {
406                    let message = StoreSyncCommand::MetadataHeaderRequest;
407                    send_command(i, message);
408                });
409            }
410            StateMachine::DownloadingMerkle {
411                partial_merkle_tree,
412                ..
413            } => {
414                debug!(
415                    "send_sync_requests: DownloadingMerkle: {:?}",
416                    partial_merkle_tree
417                );
418
419                // TODO: Keep track of (and filter out) which ones are currently requested.
420                let needed_hashes = partial_merkle_tree
421                    .missing_indices()
422                    .chunks(MERKLE_REQUEST_LIMIT as usize);
423                let mut needed_hashes: Vec<_> = needed_hashes.into_iter().collect();
424                if needed_hashes.is_empty() {
425                    panic!("Invariant violated: We have all merkle nodes but are in state DownloadingMerkle.");
426                }
427
428                // Randomize which peer to request the hashes from.
429                needed_hashes.shuffle(&mut rng);
430
431                peers.iter_mut().zip(needed_hashes).for_each(|(p, hashes)| {
432                    let message = StoreSyncCommand::MerkleRequest(
433                        compress_consecutive_into_ranges(hashes).collect(),
434                    );
435                    send_command(p.1, message);
436                });
437            }
438            StateMachine::DownloadingInitialState { initial_state, .. } => {
439                let needed_blocks = initial_state
440                    .iter()
441                    .enumerate()
442                    .filter_map(|h| {
443                        if h.1.is_none() {
444                            Some(h.0 as u64)
445                        } else {
446                            None
447                        }
448                    })
449                    .chunks(BLOCK_REQUEST_LIMIT as usize);
450                let mut needed_blocks: Vec<_> = needed_blocks.into_iter().collect();
451                if needed_blocks.is_empty() {
452                    panic!("Invariant violated: We have all blocks but are in state DownloadingInitialState.");
453                }
454
455                // Randomize which peer to request the blocks from.
456                needed_blocks.shuffle(&mut rng);
457
458                peers.iter_mut().zip(needed_blocks).for_each(|(p, blocks)| {
459                    let message = StoreSyncCommand::InitialStateBlockRequest(
460                        compress_consecutive_into_ranges(blocks).collect(),
461                    );
462                    send_command(p.1, message);
463                });
464            }
465            StateMachine::Syncing { ecg_state, .. } => {
466                debug!("Sending ECG sync requests to peers.");
467                // Request ECG updates from peers
468                peers.iter_mut().for_each(|p| {
469                    // let ecg_status = p.1.ecg_status.clone();
470                    let ecg_state = ecg_state.state().clone();
471                    let message = StoreSyncCommand::ECGSyncRequest { ecg_state };
472                    debug!("Sending ECG sync request to peer ({})", p.0);
473                    send_command(p.1, message)
474                });
475            }
476        }
477    }
478
479    // Handle a sync request for this store from a peer.
480    fn handle_metadata_peer_request(
481        &mut self,
482        peer: DeviceId,
483        response_chan: Sender<HandlePeerResponse<MetadataHeader<Hash>>>,
484    ) {
485        if let Some(metadata) = self.metadata() {
486            // We have the metadata so share it with the peer.
487            response_chan.send(Ok(*metadata)).expect("TODO");
488        } else {
489            // We don't have the metadata so tell them to wait.
490            let (send_chan, recv_chan) = oneshot::channel();
491            response_chan.send(Err(recv_chan)).expect("TODO");
492
493            // Register the wait channel.
494            self.metadata_subscribers.insert(peer, send_chan); // JP: Safe to drop old one?
495        }
496    }
497
498    fn handle_merkle_peer_request(
499        &mut self,
500        peer: DeviceId,
501        node_ids: Vec<Range<u64>>,
502        response_chan: Sender<HandlePeerResponse<Vec<Hash>>>,
503    ) {
504        debug!("Received merkle peer request for node_ids: {node_ids:?}");
505        let hashes = self
506            .merkle_tree()
507            .map(|merkle_tree| handle_merkle_peer_request_helper(merkle_tree, &node_ids));
508
509        if let Some(node_hashes) = hashes {
510            // We have the hashes so share it with the peer.
511            response_chan.send(Ok(node_hashes)).expect("TODO");
512        } else {
513            // We don't have the hashes so tell them to wait.
514            let (send_chan, recv_chan) = oneshot::channel();
515            response_chan.send(Err(recv_chan)).expect("TODO");
516
517            // Register the wait channel.
518            self.merkle_subscribers.insert(peer, (node_ids, send_chan)); // JP: Safe to drop old one?
519        }
520    }
521
522    fn handle_block_peer_request(
523        &mut self,
524        peer: DeviceId,
525        block_ids: Vec<Range<u64>>,
526        response_chan: Sender<HandlePeerResponse<Vec<Option<Vec<u8>>>>>,
527    ) {
528        let blocks = handle_block_peer_request_helper(&self.state_machine, &block_ids);
529
530        if let Some(blocks) = blocks {
531            // We have the blocks so share it with the peer.
532            response_chan.send(Ok(blocks)).expect("TODO");
533        } else {
534            // We don't have the blocks so tell them to wait.
535            let (send_chan, recv_chan) = oneshot::channel();
536            response_chan.send(Err(recv_chan)).expect("TODO");
537
538            // Register the wait channel.
539            self.block_subscribers.insert(peer, (block_ids, send_chan)); // JP: Safe to drop old one?
540        }
541    }
542
543    fn handle_ecg_subscribe(
544        &mut self,
545        peer: DeviceId,
546        tips: Option<BTreeSet<Header::HeaderId>>,
547        response_chan: oneshot::Sender<ecg::UntypedState<Header::HeaderId, Header>>,
548    ) {
549        // Respond immediately if peer thread is stale (or they requested it immediately with None).
550        if let StateMachine::Syncing { ecg_state, .. } = &self.state_machine {
551            let respond_immediately = if let Some(tips) = tips {
552                debug!("our_tips: {:?}", ecg_state.tips());
553                debug!("their_tips: {:?}", tips);
554                !ecg_state.tips().eq(&tips)
555            } else {
556                true
557            };
558
559            if respond_immediately {
560                debug!("Responding immediately with ECG state.");
561                response_chan.send(ecg_state.state.clone()).expect("TODO");
562
563                return;
564            }
565        };
566
567        // Register subscriber.
568        debug!("Registering subscriber for ECG state for peer: {peer}");
569        self.ecg_subscribers.insert(peer, response_chan);
570    }
571
572    // fn handle_ecg_sync_request(&mut self, peer: DeviceId, request: (Vec<Header::HeaderId>, Vec<Header::HeaderId>), response_chan: Sender<HandlePeerResponse<Vec<(Header, RawECGBody)>>>) {
573    //     // JP: Instead of receiving the meet, receive the frontier of what they need? Or have an
574    //     // enum where that's one option (but what if a new branch from the root is added)??
575    //     let (meet, tips) = request;
576    //     // Send everything after the meet (that they don't have). Update meet?
577
578    //     // Traverse backwards from our tip + stop once we get to something they have
579    //     todo!();
580    // }
581
582    fn handle_received_metadata(
583        &mut self,
584        peer: DeviceId,
585        metadata: MetadataHeader<Hash>,
586        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
587    ) where
588        T: for<'d> Deserialize<'d>,
589    {
590        debug!("Recieved metadata from peer ({peer}): {metadata:?}");
591
592        // Mark peer as ready.
593        self.update_outgoing_peer_to_ready(&peer);
594
595        // Validate metadata.
596        let is_valid = metadata.validate_store_id(self.store_id());
597        warn!("TODO: Validate signature.");
598        warn!("TODO: Validate type id.");
599
600        if !is_valid {
601            // TODO: Penalize/blacklist/disconnect peer?
602            warn!("TODO: Peer provided invalid metadata.");
603            return;
604        }
605
606        // Update state.
607        self.state_machine = StateMachine::DownloadingMerkle {
608            partial_merkle_tree: MerkleTree::new_with_capacity(
609                metadata.merkle_root,
610                metadata.block_count(),
611            ),
612            metadata,
613        };
614        debug!("Updated state machine"); // : {:?}", self.state_machine);
615
616        // Send metadata to any peers that are waiting.
617        let subs = std::mem::take(&mut self.metadata_subscribers);
618        for (sub_peer, sub) in subs {
619            let peer_knows = sub_peer == peer;
620            let msg = if peer_knows { None } else { Some(metadata) };
621            sub.send(msg).expect("TODO");
622        }
623
624        // If there's only one block, we already know the entire merkle tree so move onto downloading blocks.
625        if metadata.block_count() <= 1 {
626            self.update_state_to_downloading_initial_state(peer, listeners);
627        }
628    }
629
630    fn metadata(&self) -> Option<&MetadataHeader<Hash>> {
631        match &self.state_machine {
632            StateMachine::DownloadingMetadata { .. } => None,
633            StateMachine::DownloadingMerkle { metadata, .. } => Some(metadata),
634            StateMachine::DownloadingInitialState { metadata, .. } => Some(metadata),
635            StateMachine::Syncing { metadata, .. } => Some(metadata),
636        }
637    }
638
639    fn merkle_tree(&self) -> Option<&MerkleTree<Hash>> {
640        match &self.state_machine {
641            StateMachine::DownloadingMetadata { .. } => None,
642            StateMachine::DownloadingMerkle { .. } => {
643                // We shouldn't share hashes until we've verified them.
644                // TODO: If we've verified some of the hashes, we can send those back.
645                None
646            }
647            StateMachine::DownloadingInitialState { merkle_tree, .. } => Some(merkle_tree),
648            StateMachine::Syncing { merkle_tree, .. } => Some(merkle_tree),
649        }
650    }
651
652    fn handle_received_merkle_hashes(
653        &mut self,
654        peer: DeviceId,
655        node_ids: Vec<Range<u64>>,
656        their_node_hashes: Vec<Hash>,
657        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
658    ) where
659        T: for<'d> Deserialize<'d>,
660    {
661        warn!("TODO: Keep track if you received different hashes from different peers.");
662
663        // Mark peer as ready.
664        self.update_outgoing_peer_to_ready(&peer);
665
666        let node_ids: Vec<_> = node_ids.into_iter().flatten().collect();
667        if node_ids.len() != their_node_hashes.len() {
668            warn!("TODO: Peer provided an invalid response");
669            return;
670        }
671
672        // Update state.
673        let partial_merkle_tree = if let StateMachine::DownloadingMerkle {
674            ref mut partial_merkle_tree,
675            ..
676        } = &mut self.state_machine
677        {
678            partial_merkle_tree
679        } else {
680            // We're not downloading anymore so we're done.
681            return;
682        };
683        node_ids
684            .into_iter()
685            .zip(their_node_hashes)
686            .for_each(|(i, hash)| {
687                let valid = partial_merkle_tree.set(i, hash);
688                if !valid {
689                    warn!("TODO: Peer sent us an invalid merkle hash.");
690                } // TODO: Else credit peer with sending us these merkle nodes.
691            });
692
693        // Update state.
694        self.update_state_to_downloading_initial_state(peer, listeners);
695    }
696
697    fn handle_received_initial_state_blocks(
698        &mut self,
699        peer: DeviceId,
700        block_ids: Vec<Range<u64>>,
701        their_blocks: Vec<Option<Vec<u8>>>,
702        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
703    ) where
704        T: for<'d> Deserialize<'d>,
705    {
706        // Mark peer as ready.
707        self.update_outgoing_peer_to_ready(&peer);
708
709        let block_ids: Vec<_> = block_ids.into_iter().flatten().collect();
710        if block_ids.len() != their_blocks.len() {
711            warn!("TODO: Peer provided an invalid response");
712            return;
713        }
714
715        // Update state.
716        let (initial_state, merkle_tree) = if let StateMachine::DownloadingInitialState {
717            ref mut initial_state,
718            ref merkle_tree,
719            ..
720        } = &mut self.state_machine
721        {
722            (initial_state, merkle_tree)
723        } else {
724            return;
725        };
726        block_ids
727            .into_iter()
728            .zip(their_blocks)
729            .for_each(|(i, their_block)| {
730                if let Some(their_block) = their_block {
731                    let block = &mut initial_state[i as usize];
732                    // Only set block if it's currently None and if it validates.
733                    if block.is_none() {
734                        // TODO: Credit peer with this block. If not valid, penalize peer.
735                        if merkle_tree.validate_chunk(i, &their_block) {
736                            *block = Some(their_block);
737                        } else {
738                            warn!("TODO: Peer sent us an invalid block");
739                        }
740                    }
741                }
742            });
743
744        // Check if we're done. Exit if we're not.
745        if initial_state.iter().any(|o| o.is_none()) {
746            return;
747        }
748
749        // Update state.
750        self.update_state_to_syncing(peer, listeners);
751    }
752
753    fn handle_received_ecg_operations<OT>(
754        &mut self,
755        peer: DeviceId,
756        operations: Vec<(Header, RawECGBody)>,
757        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
758    ) where
759        OT: OssaType<ECGHeader = Header>,
760        T: CRDT<Time = OT::Time> + Debug,
761        OT::ECGBody<T>: for<'d> Deserialize<'d>
762            + Debug
763            + ECGBody<
764                T::Op,
765                <T::Op as ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>>::Serialized,
766                Header = OT::ECGHeader,
767            >, // ECGBody<T, Header = OT::ECGHeader> +
768        // T::Op: ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>,
769        T::Op: ConcretizeTime<<Header as ECGHeader>::HeaderId>,
770    {
771        // Mark peer as ready.
772        self.update_outgoing_peer_to_ready(&peer);
773
774        warn!("TODO: Validate operations from peer");
775
776        let StateMachine::Syncing {
777            ref mut ecg_state,
778            ref mut decrypted_state,
779            ..
780        } = &mut self.state_machine
781        else {
782            unreachable!("We must be syncing");
783        };
784
785        // Parse and apply all operations.
786        operations.into_iter().for_each(|(header, raw_operations)| {
787            let operations = serde_cbor::from_slice(&raw_operations)
788                .expect("TODO: Peer gave us improperly serialized operations");
789            debug!("Applying operations {operations:?}");
790
791            // TODO: Get rid of this clone.
792            let success = ecg_state.insert_header(header.clone(), raw_operations);
793            if !success {
794                debug!("Failed to insert operations from peer.");
795            } else {
796                apply_operations::<OT, _>(decrypted_state, ecg_state, &header, operations);
797            }
798        });
799        debug!("New decrypted state {:?}", decrypted_state.latest_state);
800
801        // Update listeners (except peer).
802        update_listeners(
803            &mut self.ecg_subscribers,
804            listeners,
805            &decrypted_state.latest_state,
806            ecg_state,
807            Some(peer),
808        );
809    }
810
811    // Precondition: State is StateMachine::DownloadingMerkle.
812    fn update_state_to_downloading_initial_state(
813        &mut self,
814        peer: DeviceId,
815        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
816    ) where
817        T: for<'d> Deserialize<'d>,
818    {
819        let StateMachine::DownloadingMerkle {
820            ref partial_merkle_tree,
821            ..
822        } = &self.state_machine
823        else {
824            panic!("Precondition violated. State mush be StateMachine::DownloadingMerkle");
825        };
826        // Check if we're done.
827        let merkle_tree_m = partial_merkle_tree.try_complete();
828        let Some(merkle_tree) = merkle_tree_m else {
829            return;
830        };
831
832        // // Validate piece hashes.
833        // let is_valid = self.metadata().unwrap().merkle_root == util::merkle_root(&piece_hashes);
834        // if !is_valid {
835        //     // TODO: Penalize/blacklist/disconnect peer?
836        //     warn!("TODO: Peer provided invalid piece hashes.");
837        //     return;
838        // }
839
840        self.state_machine = match self.state_machine {
841            StateMachine::DownloadingMerkle { metadata, .. } => {
842                let initial_state = vec![None; metadata.block_count() as usize];
843                StateMachine::DownloadingInitialState {
844                    metadata,
845                    merkle_tree,
846                    initial_state,
847                }
848            }
849            _ => unreachable!("We already checked that we're downloading merkle"),
850        };
851
852        // Send node hashes to any peers that are waiting.
853        let subs = std::mem::take(&mut self.merkle_subscribers);
854        let merkle_tree = self
855            .merkle_tree()
856            .expect("Unreachable: We just set the merkle tree");
857        for (sub_peer, (node_ids, sub)) in subs {
858            let peer_knows = sub_peer == peer;
859            let msg = if peer_knows {
860                None
861            } else {
862                Some(handle_merkle_peer_request_helper(merkle_tree, &node_ids))
863            };
864            sub.send(msg).expect("TODO");
865        }
866
867        // If there's no initial state, we already know all the blocks so move onto syncing.
868        if self.metadata().unwrap().initial_state_size == 0 {
869            self.update_state_to_syncing(peer, listeners);
870        }
871    }
872
873    fn update_state_to_syncing(
874        &mut self,
875        peer: DeviceId,
876        listeners: &[UnboundedSender<StateUpdate<Header, T>>],
877    ) where
878        T: for<'d> Deserialize<'d>,
879    {
880        replace_with_or_abort(&mut self.state_machine, |sm| match sm {
881            StateMachine::DownloadingInitialState {
882                metadata,
883                merkle_tree,
884                initial_state,
885            } => {
886                let ecg_state = ecg::State::new();
887                let initial_state: Vec<u8> = initial_state
888                    .into_iter()
889                    .flatten()
890                    .flatten()
891                    .collect::<Vec<u8>>();
892                let Ok(latest_state) = serde_cbor::de::from_slice::<T>(&initial_state) else {
893                    todo!("TODO: The store is invalid. Initial state does not parse.");
894                };
895                let decrypted_state = DecryptedState {
896                    latest_state,
897                    latest_headers: BTreeSet::new(),
898                };
899                StateMachine::Syncing {
900                    metadata,
901                    merkle_tree,
902                    initial_state,
903                    ecg_state,
904                    decrypted_state,
905                }
906            }
907            _ => unreachable!("We already checked that we're downloading the initial state"),
908        });
909
910        // Update listeners.
911        let StateMachine::Syncing {
912            ecg_state,
913            decrypted_state,
914            ..
915        } = &self.state_machine
916        else {
917            unreachable!("We just set our state to syncing")
918        };
919        update_listeners(
920            &mut self.ecg_subscribers,
921            listeners,
922            &decrypted_state.latest_state,
923            ecg_state,
924            Some(peer),
925        );
926
927        // Send blocks to any peers that are waiting.
928        let subs = std::mem::take(&mut self.block_subscribers);
929        for (_sub_peer, (block_ids, sub)) in subs {
930            // JP: Do we need to keep sub_peer here?
931            let msg = handle_block_peer_request_helper(&self.state_machine, &block_ids)
932                .expect("Unreachable: We just set our state to syncing");
933            sub.send(Some(msg)).expect("TODO");
934        }
935    }
936}
937
938fn update_listeners<Header: ecg::ECGHeader + Clone + Debug, T: CRDT + Clone>(
939    ecg_subscribers: &mut BTreeMap<
940        DeviceId,
941        oneshot::Sender<ecg::UntypedState<Header::HeaderId, Header>>,
942    >,
943    listeners: &[UnboundedSender<StateUpdate<Header, T>>],
944    latest_state: &T,
945    ecg_state: &ecg::State<Header, T>,
946    from_peer: Option<DeviceId>,
947) {
948    for l in listeners {
949        let snapshot: StateUpdate<Header, T> = StateUpdate::Snapshot {
950            snapshot: latest_state.clone(),
951            ecg_state: ecg_state.clone(),
952        };
953        l.send(snapshot).expect("TODO");
954    }
955
956    // Send updated state to one-time subscribers.
957    // warn!("TODO: Do we always want to update ECG subscribers here? Ex: We may not want to when transitioning from downloading to syncing"); JP: Maybe this is ok since our peer_store won't have anything to share and will resubscribe.
958    let subs = std::mem::take(ecg_subscribers);
959    for (sub_peer, sub) in subs {
960        // Skip notifying subscriber if they told us about this update.
961        if Some(sub_peer) != from_peer {
962            sub.send(ecg_state.state.clone()).expect("TODO");
963        } else {
964            warn!("TODO: Add headers that they sent us to their_known.");
965            // Need to add back subscriber.
966            ecg_subscribers.insert(sub_peer, sub);
967        }
968    }
969}
970
971// JP: Or should Ossa own this/peers?
972/// Manage peers by ranking them, randomize, potentially connecting to some of them, etc.
973async fn manage_peers<OT: OssaType, T: CRDT<Time = OT::Time> + Clone + Send + 'static>(
974    store: &mut State<OT::StoreId, OT::ECGHeader, T, OT::Hash>,
975    shared_state: &SharedState<OT::StoreId>,
976    send_commands: &UnboundedSender<
977        UntypedStoreCommand<OT::Hash, <OT::ECGHeader as ECGHeader>::HeaderId, OT::ECGHeader>,
978    >,
979) where
980    // T::Op<CausalTime<T::Time>>: Serialize,
981    OT::ECGHeader: Clone + Serialize + for<'d> Deserialize<'d> + Send + Sync,
982    <OT::ECGHeader as ECGHeader>::HeaderId: Serialize + for<'d> Deserialize<'d> + Send,
983    //OT::ECGHeader<T>::HeaderId : Send,
984    //T: Send,
985{
986    // For now, sync with all (connected?) peers.
987    // Don't connect to peers we're already syncing with.
988    let peers: Vec<_> = store
989        .peers
990        .iter()
991        .filter(|p| p.1.outgoing_status.is_known())
992        .collect();
993    let peers: Vec<_> = {
994        // Acquire lock on shared state.
995        let peer_states = shared_state.peer_state.read().await;
996        peers
997            .into_iter()
998            .filter_map(|(peer_id, _)| {
999                let chan = peer_states.get(peer_id)?;
1000                Some((*peer_id, chan.clone()))
1001            })
1002            .collect()
1003    };
1004    for (peer_id, command_chan) in peers {
1005        // Mark task as initializing.
1006        store.update_peer_to_initializing_outgoing(&peer_id);
1007
1008        // Create closure that spawns task to sync store with peer.
1009        let send_commands = send_commands.clone();
1010        let spawn_task = Box::new(move |_party, stream_id, sender, receiver| {
1011            // Create miniprotocol
1012            // Spawn task that syncs store with peer.
1013            // JP: Should run with initiative?
1014            tokio::spawn(async move {
1015                // Tell store we're running and send it our channel.
1016                let (send_peer, recv_peer) = tokio::sync::mpsc::unbounded_channel::<
1017                    StoreSyncCommand<<OT::ECGHeader as ECGHeader>::HeaderId, OT::ECGHeader>,
1018                >();
1019
1020                let register_cmd = UntypedStoreCommand::RegisterOutgoingPeerSyncing {
1021                    peer: peer_id,
1022                    send_peer,
1023                };
1024                send_commands.send(register_cmd).expect("TODO");
1025
1026                // Start miniprotocol as server.
1027                let mp = StoreSync::<OT::Hash, _, _>::new_server(peer_id, recv_peer, send_commands);
1028                run_miniprotocol_async::<_, OT>(mp, false, stream_id, sender, receiver).await;
1029
1030                debug!("Store sync with peer (with initiative) exited.")
1031
1032                // JP: This requires StorePeer protocols in both directions (if both sides want updates from the other party).
1033                // This has the downside that we may run ECG sync in both directions.. Does it make sense to store StorePeer state in a shared Arc<RWLock>?
1034                // Separate thread for SCSync?
1035                // TODO: Setup in both directions, remove initialized check.
1036                //
1037                // In the store task, sync state:
1038                // Based on status,
1039                //   if we're still downloading the metadata, request the peer to send it
1040                //   if we are downloading the initial state, send a download request for a piece(s) we need. How do we do timeouts here though? Perhaps in the StorePeer task?
1041                //   if we're syncing:
1042                //      run ECG sync to find the meet
1043                //      Request all operations after the meet
1044            })
1045        });
1046
1047        // Send request to peer's manager for stream.
1048        let store_id = store.store_id();
1049        let cmd = PeerManagerCommand::RequestStoreSync {
1050            store_id,
1051            spawn_task,
1052        };
1053        command_chan.send(cmd).expect("TODO");
1054    }
1055}
1056
1057fn apply_operations<OT: OssaType, T>(
1058    decrypted_state: &mut DecryptedState<OT::ECGHeader, T>,
1059    ecg_state: &ecg::State<OT::ECGHeader, T>,
1060    operation_header: &OT::ECGHeader,
1061    operation_body: OT::ECGBody<T>,
1062) where
1063    T: CRDT<Time = OT::Time>,
1064    // T::Op<CausalTime<T::Time>>: Serialize,
1065    T::Op: ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>,
1066    OT::ECGBody<T>: ECGBody<
1067        T::Op,
1068        <T::Op as ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>>::Serialized,
1069        Header = OT::ECGHeader,
1070    >,
1071{
1072    let causal_state = OT::to_causal_state(ecg_state);
1073    for operation in operation_body.operations(operation_header.get_header_id()) {
1074        replace_with_or_abort(&mut decrypted_state.latest_state, |s| {
1075            s.apply(causal_state, operation)
1076        });
1077    }
1078}
1079
1080/// Run the handler that owns this store and manages its state. This handler is typically run in
1081/// its own tokio thread.
1082pub(crate) async fn run_handler<OT: OssaType, T>(
1083    mut store: State<OT::StoreId, OT::ECGHeader, T, OT::Hash>,
1084    mut recv_commands: UnboundedReceiver<StoreCommand<OT::ECGHeader, OT::ECGBody<T>, T>>,
1085    send_commands_untyped: UnboundedSender<
1086        UntypedStoreCommand<OT::Hash, <OT::ECGHeader as ECGHeader>::HeaderId, OT::ECGHeader>,
1087    >,
1088    mut recv_commands_untyped: UnboundedReceiver<
1089        UntypedStoreCommand<OT::Hash, <OT::ECGHeader as ECGHeader>::HeaderId, OT::ECGHeader>,
1090    >,
1091    shared_state: SharedState<OT::StoreId>,
1092) where
1093    <OT as OssaType>::ECGHeader:
1094        Send + Sync + Clone + Serialize + for<'d> Deserialize<'d> + 'static,
1095    // <<OT as OssaType>::ECGHeader as ECGHeader>::Body: ECGBody<T> + Send,
1096    T::Op: ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>,
1097    OT::ECGBody<T>: Serialize
1098        + for<'d> Deserialize<'d>
1099        + Debug
1100        + ECGBody<
1101            T::Op,
1102            <T::Op as ConcretizeTime<<OT::ECGHeader as ECGHeader>::HeaderId>>::Serialized,
1103            Header = OT::ECGHeader,
1104        >,
1105    //     ECGBody<T, Header = OT::ECGHeader> + Send + Serialize + for<'d> Deserialize<'d> + Debug,
1106    <<OT as OssaType>::ECGHeader as ECGHeader>::HeaderId:
1107        Send + Serialize + for<'d> Deserialize<'d>,
1108    // T::Op<CausalTime<T::Time>>: Serialize,
1109    T: CRDT<Time = OT::Time> + Debug + Clone + Send + 'static + for<'d> Deserialize<'d>,
1110{
1111    let mut listeners: Vec<UnboundedSender<StateUpdate<OT::ECGHeader, T>>> = vec![];
1112
1113    // TODO: Check when done
1114    loop {
1115        tokio::select! {
1116            cmd_m = recv_commands.recv() => {
1117                let Some(cmd) = cmd_m else {
1118                    error!("Failed to receive StoreCommand");
1119                    return;
1120                };
1121
1122                // // Rank and connect to a few peers.
1123                // manage_peers::<OT,T>(&mut store, &shared_state).await;
1124
1125                match cmd {
1126                    StoreCommand::Apply {
1127                        operation_header,
1128                        operation_body,
1129                    } => {
1130                        store.state_machine = match store.state_machine {
1131                            // StateMachine::DownloadingMetadata { store_id } => {
1132                            //     // JP: Should we ever apply an operation if we're still downloading the store??
1133                            //     warn!("Is this unreachable?");
1134
1135                            //     // Rank and connect to a few peers.
1136                            //     manage_peers::<OT,T>(&mut store, &shared_state).await;
1137
1138                            //     StateMachine::DownloadingMetadata { store_id }
1139                            // }
1140                            // StateMachine::DownloadingMerkle { metadata, piece_hashes } => {
1141                            //     // JP: Should we ever apply an operation if we're still downloading the store??
1142                            //     warn!("Is this unreachable?");
1143
1144                            //     // Rank and connect to a few peers.
1145                            //     manage_peers::<OT,T>(&mut store, &shared_state).await;
1146
1147                            //     StateMachine::DownloadingMerkle { metadata, piece_hashes }
1148                            // }
1149                            // StateMachine::DownloadingInitialState { metadata, piece_hashes, initial_state } => {
1150                            //     // JP: Should we ever apply an operation if we're still downloading the store??
1151                            //     warn!("Is this unreachable?");
1152
1153                            //     // Rank and connect to a few peers.
1154                            //     manage_peers::<OT,T>(&mut store, &shared_state).await;
1155
1156                            //     StateMachine::DownloadingInitialState { metadata, piece_hashes, initial_state }
1157                            // }
1158                            StateMachine::Syncing { metadata, merkle_tree, initial_state , ecg_state, decrypted_state } => {
1159                                let mut ecg_state = ecg_state;
1160                                let mut decrypted_state = decrypted_state;
1161
1162                                // Update ECG state.
1163                                let serialized_operations = serde_cbor::to_vec(&operation_body).expect("TODO");
1164                                let success = ecg_state.insert_header(operation_header.clone(), serialized_operations);
1165                                if !success {
1166                                    todo!("Invalid header"); // : {:?}", operation_header);
1167                                }
1168
1169                                // Update state.
1170                                // TODO: Get new time?.. Or take it as an argument
1171                                // Operation ID/time is function of tips, current operation, ...? How do we
1172                                // do batching? (HeaderId(h) | Self, Index(u8)) ? This requires having all
1173                                // the batched operations?
1174                                apply_operations::<OT, _>(&mut decrypted_state, &ecg_state, &operation_header, operation_body);
1175
1176                                // Send state to subscribers.
1177                                update_listeners(&mut store.ecg_subscribers, &listeners, &decrypted_state.latest_state, &ecg_state, None);
1178
1179                                StateMachine::Syncing { metadata, merkle_tree, initial_state, ecg_state, decrypted_state }
1180                            }
1181                            _ => {
1182                                warn!("JP: Does this ever happen?");
1183                                store.state_machine
1184                            }
1185                        };
1186                    }
1187                    StoreCommand::SubscribeState { send_state } => {
1188                        // Send current state.
1189                        let snapshot = match &store.state_machine {
1190                            StateMachine::DownloadingMetadata { .. } => {
1191                                StateUpdate::Downloading { percent: 0 }
1192                            }
1193                            StateMachine::DownloadingMerkle { .. } => {
1194                                StateUpdate::Downloading { percent: 0 }
1195                            }
1196                            StateMachine::DownloadingInitialState { metadata, initial_state, .. } => {
1197                                let percent = if metadata.initial_state_size == 0 {
1198                                    0
1199                                } else {
1200                                    let downloaded = initial_state.iter().filter(|p| p.is_some()).count() as u64;
1201                                    100 * downloaded * BLOCK_SIZE / metadata.initial_state_size
1202                                };
1203                                StateUpdate::Downloading { percent }
1204                            }
1205                            StateMachine::Syncing { ref ecg_state, ref decrypted_state, .. } => {
1206                                StateUpdate::Snapshot {
1207                                    snapshot: decrypted_state.latest_state.clone(),
1208                                    ecg_state: ecg_state.clone(),
1209                                }
1210                            }
1211                        };
1212                        send_state.send(snapshot).expect("TODO");
1213
1214                        // Register this subscriber.
1215                        listeners.push(send_state);
1216                    }
1217                }
1218            }
1219            cmd_m = recv_commands_untyped.recv() => {
1220                let Some(cmd) = cmd_m else {
1221                    error!("Failed to receive UntypedStoreCommand");
1222                    return;
1223                };
1224                match cmd {
1225                    // Called when:
1226                    // - Peer manager threads have this thread as a mutual store.
1227                    UntypedStoreCommand::RegisterPeers { peers } => {
1228                        debug!("Received UntypedStoreCommand::RegisterPeers: {:?}", peers);
1229
1230                        // Add peer to known peers.
1231                        for peer in peers {
1232                            store.insert_known_peer(peer);
1233                        }
1234
1235                        debug!("Peer statuses: {:?}", store.peers);
1236
1237                        // Spawn sync threads for each shared store.
1238                        // TODO: Only do this if server?
1239                        // Check if we already are syncing these.
1240                        manage_peers::<OT,T>(&mut store, &shared_state, &send_commands_untyped).await;
1241                    }
1242                    // Sets up task to respond to a request to sync this store from a peer (without initiative).
1243                    // Called when:
1244                    // - The peer requests we sync this store with them
1245                    UntypedStoreCommand::SyncWithPeer { peer, response_chan } => {
1246                        debug!("Received UntypedStoreCommand::SyncWithPeer: {:?}", peer);
1247
1248                        // Insert peer as known if we don't know them (since they're requesting the store).
1249                        store.insert_known_peer(peer);
1250
1251                        let response = {
1252                            // Check if already syncing with this peer. (JP: What if they're both already "Initializing"? Potential race condition where they don't sync)
1253                            if let Some(status) = store.peers.get(&peer) {
1254                                if status.incoming_status.is_known() {
1255                                    // Mark task as initializing.
1256                                    store.update_peer_to_initializing_incoming(&peer);
1257
1258                                    // Create closure that spawns task to sync store with peer.
1259                                    let send_commands_untyped = send_commands_untyped.clone();
1260                                    let spawn_task: Box<SpawnMultiplexerTask> = Box::new(move |party, stream_id, sender, receiver| {
1261                                        // Create miniprotocol
1262                                        // Spawn task that syncs store with peer.
1263                                        // JP: Should run without initiative so that other peer can setup their handler?
1264                                        tokio::spawn(async move {
1265                                            debug!("Sync with peer (without initiative).");
1266
1267                                            // Tell store we're running.
1268                                            let register_cmd = UntypedStoreCommand::RegisterIncomingPeerSyncing {
1269                                                peer,
1270                                            };
1271                                            send_commands_untyped.send(register_cmd).expect("TODO");
1272
1273                                            // Start miniprotocol as client.
1274                                            let mp = StoreSync::<OT::Hash, _, _>::new_client(peer, send_commands_untyped);
1275                                            run_miniprotocol_async::<_, OT>(mp, true, stream_id, sender, receiver).await;
1276                                            debug!("Store sync with peer (without initiative) exited.")
1277                                        })
1278                                    });
1279                                    Some(spawn_task)
1280                                } else {
1281                                    debug!("Store is already running");
1282                                    None
1283                                }
1284                            } else {
1285                                unreachable!("Don't know this peer.");
1286                                // JP: This should be impossible now.
1287                                None
1288                            }
1289                        };
1290
1291                        response_chan.send(response).or(Err(())).expect("TODO");
1292                    }
1293                    UntypedStoreCommand::RegisterOutgoingPeerSyncing{ peer, send_peer } => {
1294                        // JP: Maybe send_peer actually isn't needed??? We could construct oneshots???
1295                        // Update peer's state to syncing and register channel.
1296                        let outgoing_status = OutgoingPeerStatus {
1297                            sender_peer: send_peer,
1298                            is_outstanding: false,
1299                        };
1300                        store.update_peer_to_syncing_outgoing(&peer, outgoing_status);
1301
1302                        // Sync with peer(s). Do this for all commands??
1303                        store.send_sync_requests();
1304                    }
1305                    UntypedStoreCommand::RegisterIncomingPeerSyncing{ peer } => {
1306                        // JP: Maybe this actually isn't needed??? We could construct oneshots for every request..
1307
1308                        // Update peer's state to syncing and register channel.
1309                        store.update_peer_to_syncing_incoming(&peer);
1310                    }
1311                    UntypedStoreCommand::HandleMetadataPeerRequest(HandlePeerRequest { peer, request, response_chan }) => {
1312                        store.handle_metadata_peer_request(peer, response_chan);
1313                    }
1314                    UntypedStoreCommand::HandleMerklePeerRequest(HandlePeerRequest { peer, request, response_chan }) => {
1315                        store.handle_merkle_peer_request(peer, request, response_chan);
1316                    }
1317                    UntypedStoreCommand::HandleBlockPeerRequest(HandlePeerRequest { peer, request, response_chan }) => {
1318                        store.handle_block_peer_request(peer, request, response_chan);
1319                    }
1320                    UntypedStoreCommand::ReceivedMetadata { peer, metadata } => {
1321                        store.handle_received_metadata(peer, metadata, &listeners);
1322                        store.send_sync_requests();
1323                    }
1324                    UntypedStoreCommand::ReceivedMerkleHashes { peer, ranges, nodes } => {
1325                        store.handle_received_merkle_hashes(peer, ranges, nodes, &listeners);
1326                        store.send_sync_requests();
1327                    }
1328                    UntypedStoreCommand::ReceivedInitialStateBlocks { peer, ranges, blocks } => {
1329                        store.handle_received_initial_state_blocks(peer, ranges, blocks, &listeners);
1330                        store.send_sync_requests();
1331                    }
1332                    UntypedStoreCommand::ReceivedECGOperations { peer, operations } => {
1333                        store.handle_received_ecg_operations::<OT>(peer, operations, &listeners);
1334                        store.send_sync_requests();
1335                    }
1336                    UntypedStoreCommand::SubscribeECG { peer, tips, response_chan } => {
1337                        store.handle_ecg_subscribe(peer, tips, response_chan);
1338                    }
1339                }
1340            }
1341        }
1342    }
1343    debug!("Store thread exiting.");
1344}
1345
1346pub(crate) enum StoreCommand<Header: ECGHeader, Body, T> {
1347    Apply {
1348        operation_header: Header, // <Hash, T>,
1349        operation_body: Body,     // <Hash, T>,
1350    },
1351    // TODO: Support unsubscribe.
1352    SubscribeState {
1353        send_state: UnboundedSender<StateUpdate<Header, T>>,
1354    },
1355}
1356
1357pub enum StateUpdate<Header: ECGHeader, T> {
1358    Downloading {
1359        // Percent of the state that we've downloaded (0 - 100).
1360        percent: u64,
1361    },
1362    Snapshot {
1363        snapshot: T,
1364        ecg_state: ecg::State<Header, T>,
1365        // TODO: ECG DAG
1366    },
1367}
1368
1369// trait UntypedCRDT: CRDT<Op = dyn Any, Time = dyn Any> {} // Any + Sized + 'static +
1370// trait UntypedECGHeader: ECGHeader<dyn Any> {} // Any + Sized + 'static +
1371//
1372// struct UntypedStoreCommand(StoreCommand<dyn UntypedECGHeader, dyn UntypedCRDT>);
1373// struct UntypedStoreCommand(StoreCommand<dyn UntypedECGHeader, dyn Any>);
1374
1375type HandlePeerResponse<Response> = Result<Response, oneshot::Receiver<Option<Response>>>;
1376
1377/// Untyped variant of `StoreCommand` since existentials don't work.
1378// #[derive(Debug)]
1379pub(crate) enum UntypedStoreCommand<Hash, HeaderId, Header> {
1380    /// Register the discovered peers.
1381    RegisterPeers {
1382        peers: Vec<DeviceId>,
1383    },
1384    /// Request store to sync with peer. Store can refuse.
1385    SyncWithPeer {
1386        peer: DeviceId,
1387        response_chan: oneshot::Sender<Option<Box<SpawnMultiplexerTask>>>,
1388    },
1389    RegisterOutgoingPeerSyncing {
1390        peer: DeviceId,
1391        send_peer: UnboundedSender<StoreSyncCommand<HeaderId, Header>>,
1392    },
1393    HandleMetadataPeerRequest(HandlePeerRequest<(), v0::MetadataHeader<Hash>>),
1394    HandleMerklePeerRequest(HandlePeerRequest<Vec<Range<u64>>, Vec<Hash>>),
1395    HandleBlockPeerRequest(HandlePeerRequest<Vec<Range<u64>>, Vec<Option<Vec<u8>>>>),
1396    // HandleECGSyncRequest(HandlePeerRequest<(Vec<HeaderId>, Vec<HeaderId>), Vec<(Header, RawECGBody)>>), // (Meet, Tips)
1397    RegisterIncomingPeerSyncing {
1398        peer: DeviceId,
1399    },
1400    ReceivedMetadata {
1401        peer: DeviceId,
1402        metadata: MetadataHeader<Hash>,
1403    },
1404    ReceivedMerkleHashes {
1405        peer: DeviceId,
1406        ranges: Vec<Range<u64>>,
1407        nodes: Vec<Hash>,
1408    },
1409    ReceivedInitialStateBlocks {
1410        peer: DeviceId,
1411        ranges: Vec<Range<u64>>,
1412        blocks: Vec<Option<Vec<u8>>>,
1413    },
1414    ReceivedECGOperations {
1415        peer: DeviceId,
1416        operations: Vec<(Header, RawECGBody)>,
1417    },
1418    SubscribeECG {
1419        peer: DeviceId,
1420        tips: Option<BTreeSet<HeaderId>>,
1421        response_chan: oneshot::Sender<ecg::UntypedState<HeaderId, Header>>,
1422    },
1423}
1424
1425pub(crate) struct HandlePeerRequest<Request, Response> {
1426    pub(crate) peer: DeviceId,
1427    pub(crate) request: Request, // MsgStoreSyncRequest,
1428    /// Return either the result or a channel to wait on for the response.
1429    pub(crate) response_chan: oneshot::Sender<HandlePeerResponse<Response>>,
1430}
1431
1432fn handle_merkle_peer_request_helper<H: Copy>(
1433    merkle_tree: &MerkleTree<H>,
1434    node_ids: &[Range<u64>],
1435) -> Vec<H> {
1436    warn!("TODO: check ranges are in bounds or return error");
1437    let hashes: Vec<_> = node_ids
1438        .iter()
1439        .cloned()
1440        .flatten()
1441        .map(|i| {
1442            merkle_tree
1443                .get(i)
1444                .expect("TODO: Properly handle invalid requests")
1445                .clone()
1446        })
1447        .collect();
1448    hashes
1449}
1450
1451fn handle_block_peer_request_helper<StoreId, Header: ecg::ECGHeader, T: CRDT, Hash>(
1452    state_machine: &StateMachine<StoreId, Header, T, Hash>,
1453    block_ids: &[Range<u64>],
1454) -> Option<Vec<Option<Vec<u8>>>> {
1455    // TODO: Can we avoid these clones?
1456    match state_machine {
1457        StateMachine::DownloadingMetadata { .. } => None,
1458        StateMachine::DownloadingMerkle { .. } => None,
1459        StateMachine::DownloadingInitialState { initial_state, .. } => {
1460            Some(handle_peer_request_range_helper(initial_state, block_ids))
1461        }
1462        StateMachine::Syncing { initial_state, .. } => {
1463            let blocks: Vec<_> = block_ids
1464                .iter()
1465                .cloned()
1466                .flatten()
1467                .map(|i| {
1468                    warn!("TODO: Properly handle invalid requests"); // Return None if i >= metadata.block_count()?
1469                    let start: usize = (i * BLOCK_SIZE) as usize;
1470                    let end = std::cmp::min(((i + 1) * BLOCK_SIZE) as usize, initial_state.len());
1471                    Some(initial_state[start..end].to_vec())
1472                })
1473                .collect();
1474            Some(blocks)
1475        }
1476    }
1477}
1478
1479fn handle_peer_request_range_helper<T: Clone>(slice: &[T], ids: &[Range<u64>]) -> Vec<T> {
1480    warn!("TODO: check ranges are in bounds or return error");
1481    let hashes: Vec<_> = ids
1482        .iter()
1483        .cloned()
1484        .flatten()
1485        .map(|i| {
1486            slice
1487                .get(i as usize)
1488                .expect("TODO: Properly handle invalid requests")
1489                .clone()
1490        })
1491        .collect();
1492    hashes
1493}