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 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 &None::<()>,
104 )
105 .await?;
106 Ok(req_no)
107 }
108
109 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}