Skip to main content

edc_connector_client/api/
transfer_process.rs

1use 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}