ossa_core/protocol/
mod.rs1use serde::{Deserialize, Serialize};
2use tokio::{
3 net::TcpStream,
4 sync::{mpsc::UnboundedReceiver, watch},
5};
6
7use crate::{
8 auth::DeviceId,
9 core::{OssaType, StoreStatuses},
10 protocol::manager::v0::PeerManagerCommand,
11 store::ecg::ECGHeader,
12};
13
14pub mod heartbeat;
15pub mod manager;
16pub mod store_peer;
17pub mod v0;
18
19pub(crate) struct MiniProtocolArgs<StoreId, Hash, HeaderId, Header> {
20 peer_id: DeviceId,
21 active_stores: watch::Receiver<StoreStatuses<StoreId, Hash, HeaderId, Header>>,
22 manager_channel: UnboundedReceiver<PeerManagerCommand<StoreId>>,
23}
24
25impl<StoreId, Hash, HeaderId, Header> MiniProtocolArgs<StoreId, Hash, HeaderId, Header> {
26 pub(crate) fn new(
27 peer_id: DeviceId,
28 active_stores: watch::Receiver<StoreStatuses<StoreId, Hash, HeaderId, Header>>,
29 manager_channel: UnboundedReceiver<PeerManagerCommand<StoreId>>,
30 ) -> Self {
31 Self {
32 peer_id,
33 active_stores,
34 manager_channel,
35 }
36 }
37}
38
39#[derive(Clone, Copy, Debug, Deserialize, Serialize)]
41pub enum Version {
42 V0 = 0,
43}
44
45impl Version {
46 pub fn as_byte(&self) -> u8 {
47 *self as u8
48 }
49
50 pub async fn run_miniprotocols_server<O: OssaType>(
51 &self,
52 stream: TcpStream,
53 args: MiniProtocolArgs<
54 O::StoreId,
55 O::Hash,
56 <O::ECGHeader as ECGHeader>::HeaderId,
57 O::ECGHeader,
58 >,
59 ) {
60 match self {
61 Version::V0 => v0::run_miniprotocols_server::<O>(stream, args).await,
62 }
63 }
64
65 pub async fn run_miniprotocols_client<O: OssaType>(
66 &self,
67 stream: TcpStream,
68 args: MiniProtocolArgs<
69 O::StoreId,
70 O::Hash,
71 <O::ECGHeader as ECGHeader>::HeaderId,
72 O::ECGHeader,
73 >,
74 ) {
75 match self {
76 Version::V0 => v0::run_miniprotocols_client::<O>(stream, args).await,
77 }
78 }
79}
80
81pub(crate) const LATEST_VERSION: Version = Version::V0;