Skip to main content

kuska_ssb/api/
helper.rs

1use crate::{
2    api::dto::content::{
3        FriendsHops, InviteCreateOptions, RelationshipQuery, SubsetQuery, SubsetQueryOptions,
4        TypedMessage,
5    },
6    feed::Message,
7    rpc::{ArgType, Body, BodyType, RequestNo, RpcType, RpcWriter},
8};
9use async_std::io::Write;
10
11use super::{dto, error::Result};
12
13const MAX_RPC_BODY_LEN: usize = 65536;
14
15#[derive(Debug)]
16pub enum ApiMethod {
17    InviteCreate,
18    InviteUse,
19    FriendsIsFollowing,
20    FriendsIsBlocking,
21    FriendsHops,
22    GetSubset,
23    Publish,
24    WhoAmI,
25    Get,
26    CreateHistoryStream,
27    CreateFeedStream,
28    Latest,
29    BlobsGet,
30    BlobsCreateWants,
31}
32
33impl ApiMethod {
34    pub fn selector(&self) -> &'static [&'static str] {
35        use ApiMethod::*;
36        match self {
37            InviteCreate => &["invite", "create"],
38            InviteUse => &["invite", "use"],
39            FriendsIsFollowing => &["friends", "isFollowing"],
40            FriendsIsBlocking => &["friends", "isBlocking"],
41            FriendsHops => &["friends", "hops"],
42            GetSubset => &["partialReplication", "getSubset"],
43            Publish => &["publish"],
44            WhoAmI => &["whoami"],
45            Get => &["get"],
46            CreateHistoryStream => &["createHistoryStream"],
47            CreateFeedStream => &["createFeedStream"],
48            Latest => &["latest"],
49            BlobsGet => &["blobs", "get"],
50            BlobsCreateWants => &["blobs", "createWants"],
51        }
52    }
53    pub fn from_selector(s: &[&str]) -> Option<Self> {
54        use ApiMethod::*;
55        match s {
56            ["invite", "create"] => Some(InviteCreate),
57            ["invite", "use"] => Some(InviteUse),
58            ["friends", "isFollowing"] => Some(FriendsIsFollowing),
59            ["friends", "isBlocking"] => Some(FriendsIsBlocking),
60            ["friends", "hops"] => Some(FriendsHops),
61            ["partialReplication", "getSubset"] => Some(GetSubset),
62            ["publish"] => Some(Publish),
63            ["whoami"] => Some(WhoAmI),
64            ["get"] => Some(Get),
65            ["createHistoryStream"] => Some(CreateHistoryStream),
66            ["createFeedStream"] => Some(CreateFeedStream),
67            ["latest"] => Some(Latest),
68            ["blobs", "get"] => Some(BlobsGet),
69            ["blobs", "createWants"] => Some(BlobsCreateWants),
70            _ => None,
71        }
72    }
73    pub fn from_rpc_body(body: &Body) -> Option<Self> {
74        let selector = body.name.iter().map(|v| v.as_str()).collect::<Vec<_>>();
75        Self::from_selector(&selector)
76    }
77}
78
79pub struct ApiCaller<W: Write + Unpin> {
80    rpc: RpcWriter<W>,
81}
82
83impl<W: Write + Unpin> ApiCaller<W> {
84    pub fn new(rpc: RpcWriter<W>) -> Self {
85        Self { rpc }
86    }
87
88    pub fn rpc(&mut self) -> &mut RpcWriter<W> {
89        &mut self.rpc
90    }
91
92    /// Send ["invite", "create"] request.
93    pub async fn invite_create_req_send(&mut self, uses: u16) -> Result<RequestNo> {
94        let args = InviteCreateOptions { uses };
95        let req_no = self
96            .rpc
97            .send_request(
98                ApiMethod::InviteCreate.selector(),
99                RpcType::Async,
100                ArgType::Object,
101                &args,
102                // specify None value for `opts`
103                &None::<()>,
104            )
105            .await?;
106        Ok(req_no)
107    }
108
109    /// Send ["invite", "use"] request.
110    pub async fn invite_use_req_send(&mut self, invite_code: &str) -> Result<RequestNo> {
111        let req_no = self
112            .rpc
113            .send_request(
114                ApiMethod::InviteUse.selector(),
115                RpcType::Async,
116                ArgType::Array,
117                &invite_code,
118                &None::<()>,
119            )
120            .await?;
121        Ok(req_no)
122    }
123
124    /// Send ["friends", "isFollowing"] request.
125    pub async fn friends_is_following_req_send(
126        &mut self,
127        args: RelationshipQuery,
128    ) -> Result<RequestNo> {
129        let req_no = self
130            .rpc
131            .send_request(
132                ApiMethod::FriendsIsFollowing.selector(),
133                RpcType::Async,
134                ArgType::Array,
135                &args,
136                &None::<()>,
137            )
138            .await?;
139        Ok(req_no)
140    }
141
142    /// Send ["friends", "isBlocking"] request.
143    pub async fn friends_is_blocking_req_send(
144        &mut self,
145        args: RelationshipQuery,
146    ) -> Result<RequestNo> {
147        let req_no = self
148            .rpc
149            .send_request(
150                ApiMethod::FriendsIsBlocking.selector(),
151                RpcType::Async,
152                ArgType::Array,
153                &args,
154                &None::<()>,
155            )
156            .await?;
157        Ok(req_no)
158    }
159
160    /// Send ["friends", "hops"] request
161    pub async fn friends_hops_req_send(&mut self, args: FriendsHops) -> Result<RequestNo> {
162        let req_no = self
163            .rpc
164            .send_request(
165                ApiMethod::FriendsHops.selector(),
166                RpcType::Source,
167                ArgType::Array,
168                &args,
169                &None::<()>,
170            )
171            .await?;
172        Ok(req_no)
173    }
174
175    /// Send ["partialReplication", "getSubset"] request.
176    pub async fn getsubset_req_send(
177        &mut self,
178        query: SubsetQuery,
179        opts: Option<SubsetQueryOptions>,
180    ) -> Result<RequestNo> {
181        let req_no = self
182            .rpc
183            .send_request(
184                ApiMethod::GetSubset.selector(),
185                RpcType::Source,
186                ArgType::Tuple,
187                &query,
188                &opts,
189            )
190            .await?;
191        Ok(req_no)
192    }
193
194    /// Send ["publish"] request.
195    pub async fn publish_req_send(&mut self, msg: TypedMessage) -> Result<RequestNo> {
196        let req_no = self
197            .rpc
198            .send_request(
199                ApiMethod::Publish.selector(),
200                RpcType::Async,
201                ArgType::Array,
202                &msg,
203                &None::<()>,
204            )
205            .await?;
206        Ok(req_no)
207    }
208
209    /// Send ["publish"] response.
210    pub async fn publish_res_send(&mut self, req_no: RequestNo, msg_ref: String) -> Result<()> {
211        Ok(self
212            .rpc
213            .send_response(req_no, RpcType::Async, BodyType::JSON, msg_ref.as_bytes())
214            .await?)
215    }
216
217    /// Send ["whoami"] request.
218    pub async fn whoami_req_send(&mut self) -> Result<RequestNo> {
219        let args: [&str; 0] = [];
220        let req_no = self
221            .rpc
222            .send_request(
223                ApiMethod::WhoAmI.selector(),
224                RpcType::Async,
225                ArgType::Array,
226                &args,
227                &None::<()>,
228            )
229            .await?;
230        Ok(req_no)
231    }
232
233    /// Send ["whoami"] response.
234    pub async fn whoami_res_send(&mut self, req_no: RequestNo, id: String) -> Result<()> {
235        let body = serde_json::to_string(&dto::WhoAmIOut { id })?;
236        Ok(self
237            .rpc
238            .send_response(req_no, RpcType::Async, BodyType::JSON, body.as_bytes())
239            .await?)
240    }
241
242    /// Send ["get"] request.
243    pub async fn get_req_send(&mut self, msg_id: &str) -> Result<RequestNo> {
244        let req_no = self
245            .rpc
246            .send_request(
247                ApiMethod::Get.selector(),
248                RpcType::Async,
249                ArgType::Array,
250                &msg_id,
251                &None::<()>,
252            )
253            .await?;
254        Ok(req_no)
255    }
256
257    /// Send ["get"] response.
258    pub async fn get_res_send(&mut self, req_no: RequestNo, msg: &Message) -> Result<()> {
259        self.rpc
260            .send_response(
261                req_no,
262                RpcType::Async,
263                BodyType::JSON,
264                msg.to_string().as_bytes(),
265            )
266            .await?;
267        Ok(())
268    }
269
270    /// Send ["createHistoryStream"] request.
271    pub async fn create_history_stream_req_send(
272        &mut self,
273        args: &dto::CreateHistoryStreamIn,
274    ) -> Result<RequestNo> {
275        let req_no = self
276            .rpc
277            .send_request(
278                ApiMethod::CreateHistoryStream.selector(),
279                RpcType::Source,
280                ArgType::Array,
281                &args,
282                &None::<()>,
283            )
284            .await?;
285        Ok(req_no)
286    }
287
288    /// Send ["createFeedStream"] request.
289    pub async fn create_feed_stream_req_send<'a>(
290        &mut self,
291        args: &dto::CreateStreamIn<u64>,
292    ) -> Result<RequestNo> {
293        let req_no = self
294            .rpc
295            .send_request(
296                ApiMethod::CreateFeedStream.selector(),
297                RpcType::Source,
298                ArgType::Array,
299                &args,
300                &None::<()>,
301            )
302            .await?;
303        Ok(req_no)
304    }
305
306    /// Send ["latest"] request.
307    pub async fn latest_req_send(&mut self) -> Result<RequestNo> {
308        let args: [&str; 0] = [];
309        let req_no = self
310            .rpc
311            .send_request(
312                ApiMethod::Latest.selector(),
313                RpcType::Async,
314                ArgType::Array,
315                &args,
316                &None::<()>,
317            )
318            .await?;
319        Ok(req_no)
320    }
321
322    /// Send ["blobs","get"] request.
323    pub async fn blobs_get_req_send(&mut self, args: &dto::BlobsGetIn) -> Result<RequestNo> {
324        let req_no = self
325            .rpc
326            .send_request(
327                ApiMethod::BlobsGet.selector(),
328                RpcType::Source,
329                ArgType::Array,
330                &args,
331                &None::<()>,
332            )
333            .await?;
334        Ok(req_no)
335    }
336
337    /// Send feed response
338    pub async fn feed_res_send(&mut self, req_no: RequestNo, feed: &str) -> Result<()> {
339        self.rpc
340            .send_response(req_no, RpcType::Source, BodyType::JSON, feed.as_bytes())
341            .await?;
342        Ok(())
343    }
344
345    /// Send blob create wants
346    pub async fn blob_create_wants_req_send(&mut self) -> Result<RequestNo> {
347        let args: [&str; 0] = [];
348        let req_no = self
349            .rpc
350            .send_request(
351                ApiMethod::BlobsCreateWants.selector(),
352                RpcType::Source,
353                ArgType::Array,
354                &args,
355                &None::<()>,
356            )
357            .await?;
358        Ok(req_no)
359    }
360
361    /// Send blob response
362    pub async fn blobs_get_res_send<D: AsRef<[u8]>>(
363        &mut self,
364        req_no: RequestNo,
365        data: D,
366    ) -> Result<()> {
367        let mut offset = 0;
368        let data = data.as_ref();
369        while offset < data.len() {
370            let limit = std::cmp::min(data.len(), offset + MAX_RPC_BODY_LEN);
371            self.rpc
372                .send_response(
373                    req_no,
374                    RpcType::Source,
375                    BodyType::Binary,
376                    &data[offset..limit],
377                )
378                .await?;
379            offset += MAX_RPC_BODY_LEN;
380        }
381        self.rpc.send_stream_eof(req_no).await?;
382        Ok(())
383    }
384}