edc_connector_client/api/
transfer_process.rs1use crate::{
2 client::EdcConnectorClientInternal,
3 types::{
4 context::WithContext,
5 query::Query,
6 response::IdResponse,
7 transfer_process::{
8 SuspendTransfer, TerminateTransfer, TransferProcess, TransferProcessState,
9 TransferRequest, TransferState,
10 },
11 },
12 EdcConnectorApiVersion, EdcResult,
13};
14
15const TRANSFER_PROCESSES_PATH: &str = "transferprocesses";
16
17pub struct TransferProcessApi<'a> {
18 client: &'a EdcConnectorClientInternal,
19 version: EdcConnectorApiVersion,
20}
21
22impl<'a> TransferProcessApi<'a> {
23 pub(crate) fn new(
24 client: &'a EdcConnectorClientInternal,
25 version: EdcConnectorApiVersion,
26 ) -> TransferProcessApi<'a> {
27 TransferProcessApi { client, version }
28 }
29
30 pub async fn initiate(
31 &self,
32 transfer_request: &TransferRequest,
33 ) -> EdcResult<IdResponse<String>> {
34 let url = self
35 .client
36 .path_for(self.version, &[TRANSFER_PROCESSES_PATH]);
37 self.client
38 .post::<_, WithContext<IdResponse<String>>>(
39 url,
40 &self.client.context_for(self.version, &transfer_request),
41 )
42 .await
43 .map(|ctx| ctx.inner)
44 }
45
46 pub async fn get(&self, id: &str) -> EdcResult<TransferProcess> {
47 let url = self
48 .client
49 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, id]);
50 self.client
51 .get::<WithContext<TransferProcess>>(url)
52 .await
53 .map(|ctx| ctx.inner)
54 }
55
56 pub async fn get_state(&self, id: &str) -> EdcResult<TransferProcessState> {
57 let url = self
58 .client
59 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, id]);
60 self.client
61 .get::<WithContext<TransferState>>(url)
62 .await
63 .map(|ctx| ctx.inner.state().clone())
64 }
65
66 pub async fn query(&self, query: Query) -> EdcResult<Vec<TransferProcess>> {
67 let url = self
68 .client
69 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, "request"]);
70
71 self.client
72 .post::<_, Vec<WithContext<TransferProcess>>>(
73 url,
74 &self.client.context_for(self.version, &query),
75 )
76 .await
77 .map(|results| results.into_iter().map(|ctx| ctx.inner).collect())
78 }
79
80 pub async fn terminate(&self, id: &str, reason: &str) -> EdcResult<()> {
81 let url = self
82 .client
83 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, id, "terminate"]);
84
85 let request = TerminateTransfer::builder()
86 .id(id.to_string())
87 .reason(reason.to_string())
88 .build();
89
90 self.client
91 .post_no_response(url, &self.client.context_for(self.version, &request))
92 .await
93 .map(|_| ())
94 }
95
96 pub async fn suspend(&self, id: &str, reason: &str) -> EdcResult<()> {
97 let url = self
98 .client
99 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, id, "suspend"]);
100
101 let request = SuspendTransfer::builder()
102 .id(id.to_string())
103 .reason(reason.to_string())
104 .build();
105
106 self.client
107 .post_no_response(url, &self.client.context_for(self.version, &request))
108 .await
109 .map(|_| ())
110 }
111
112 pub async fn resume(&self, id: &str) -> EdcResult<()> {
113 let url = self
114 .client
115 .path_for(self.version, &[TRANSFER_PROCESSES_PATH, id, "resume"]);
116 self.client
117 .post_no_response(url, &Option::<()>::None)
118 .await
119 .map(|_| ())
120 }
121}