Skip to main content

mssf_core/client/
query_client.rs

1// ------------------------------------------------------------
2// Copyright (c) Microsoft Corporation.  All rights reserved.
3// Licensed under the MIT License (MIT). See License.txt in the repo root for license information.
4// ------------------------------------------------------------
5use std::time::Duration;
6
7use mssf_com::{
8    FabricClient::{
9        IFabricGetApplicationListResult2, IFabricGetDeployedServiceReplicaDetailResult,
10        IFabricGetNodeListResult2, IFabricGetPartitionListResult2,
11        IFabricGetPartitionLoadInformationResult, IFabricGetReplicaListResult2,
12        IFabricQueryClient13,
13    },
14    FabricTypes::{
15        FABRIC_APPLICATION_QUERY_DESCRIPTION,
16        FABRIC_DEPLOYED_SERVICE_REPLICA_DETAIL_QUERY_DESCRIPTION, FABRIC_NODE_QUERY_DESCRIPTION,
17        FABRIC_PARTITION_LOAD_INFORMATION_QUERY_DESCRIPTION,
18        FABRIC_SERVICE_PARTITION_QUERY_DESCRIPTION, FABRIC_SERVICE_QUERY_DESCRIPTION,
19        FABRIC_SERVICE_REPLICA_QUERY_DESCRIPTION,
20    },
21};
22
23use crate::mem::{BoxPool, GetRawWithBoxPool};
24
25use crate::types::{
26    DeployedServiceReplicaDetailQueryDescription, DeployedServiceReplicaDetailQueryResult,
27    GetPartitionLoadInformationResult, NodeListResult, NodeQueryDescription,
28    PartitionLoadInformationQueryDescription, ServicePartitionList,
29    ServicePartitionQueryDescription, ServiceReplicaList, ServiceReplicaQueryDescription,
30};
31use crate::{
32    runtime::executor::BoxedCancelToken,
33    sync::{FabricReceiver, fabric_begin_end_proxy},
34    types::ServiceQueryDescription,
35};
36
37#[derive(Debug, Clone)]
38pub struct QueryClient {
39    com: IFabricQueryClient13,
40}
41
42// Internal implementation block
43// Internal functions focuses on changing SF callback to async future,
44// while the public apis impl focuses on type conversion.
45
46impl QueryClient {
47    pub fn get_node_list_internal(
48        &self,
49        query_description: &FABRIC_NODE_QUERY_DESCRIPTION,
50        timeout_milliseconds: u32,
51        cancellation_token: Option<BoxedCancelToken>,
52    ) -> FabricReceiver<crate::Result<IFabricGetNodeListResult2>> {
53        let com1 = &self.com;
54        let com2 = self.com.clone();
55
56        fabric_begin_end_proxy(
57            move |callback| unsafe {
58                com1.BeginGetNodeList(query_description, timeout_milliseconds, callback)
59            },
60            move |ctx| unsafe { com2.EndGetNodeList2(ctx) },
61            cancellation_token,
62        )
63    }
64
65    pub fn get_application_list_internal(
66        &self,
67        query_description: &FABRIC_APPLICATION_QUERY_DESCRIPTION,
68        timeout_milliseconds: u32,
69        cancellation_token: Option<BoxedCancelToken>,
70    ) -> FabricReceiver<crate::Result<IFabricGetApplicationListResult2>> {
71        let com1 = &self.com;
72        let com2 = self.com.clone();
73        fabric_begin_end_proxy(
74            move |callback| unsafe {
75                com1.BeginGetApplicationList(query_description, timeout_milliseconds, callback)
76            },
77            move |ctx| unsafe { com2.EndGetApplicationList2(ctx) },
78            cancellation_token,
79        )
80    }
81
82    fn get_service_list_internal(
83        &self,
84        desc: &FABRIC_SERVICE_QUERY_DESCRIPTION,
85        timeout_milliseconds: u32,
86        cancellation_token: Option<BoxedCancelToken>,
87    ) -> FabricReceiver<crate::Result<mssf_com::FabricClient::IFabricGetServiceListResult2>> {
88        let com1 = &self.com;
89        let com2 = self.com.clone();
90        fabric_begin_end_proxy(
91            move |callback| unsafe {
92                com1.BeginGetServiceList(desc, timeout_milliseconds, callback)
93            },
94            move |ctx| unsafe { com2.EndGetServiceList2(ctx) },
95            cancellation_token,
96        )
97    }
98
99    fn get_partition_list_internal(
100        &self,
101        desc: &FABRIC_SERVICE_PARTITION_QUERY_DESCRIPTION,
102        timeout_milliseconds: u32,
103        cancellation_token: Option<BoxedCancelToken>,
104    ) -> FabricReceiver<crate::Result<IFabricGetPartitionListResult2>> {
105        let com1 = &self.com;
106        let com2 = self.com.clone();
107        fabric_begin_end_proxy(
108            move |callback| unsafe {
109                com1.BeginGetPartitionList(desc, timeout_milliseconds, callback)
110            },
111            move |ctx| unsafe { com2.EndGetPartitionList2(ctx) },
112            cancellation_token,
113        )
114    }
115
116    fn get_replica_list_internal(
117        &self,
118        desc: &FABRIC_SERVICE_REPLICA_QUERY_DESCRIPTION,
119        timeout_milliseconds: u32,
120        cancellation_token: Option<BoxedCancelToken>,
121    ) -> FabricReceiver<crate::Result<IFabricGetReplicaListResult2>> {
122        let com1 = &self.com;
123        let com2 = self.com.clone();
124        fabric_begin_end_proxy(
125            move |callback| unsafe {
126                com1.BeginGetReplicaList(desc, timeout_milliseconds, callback)
127            },
128            move |ctx| unsafe { com2.EndGetReplicaList2(ctx) },
129            cancellation_token,
130        )
131    }
132
133    fn get_partition_load_information_internal(
134        &self,
135        desc: &FABRIC_PARTITION_LOAD_INFORMATION_QUERY_DESCRIPTION,
136        timeout_milliseconds: u32,
137        cancellation_token: Option<BoxedCancelToken>,
138    ) -> FabricReceiver<crate::Result<IFabricGetPartitionLoadInformationResult>> {
139        let com1 = &self.com;
140        let com2 = self.com.clone();
141        fabric_begin_end_proxy(
142            move |callback| unsafe {
143                com1.BeginGetPartitionLoadInformation(desc, timeout_milliseconds, callback)
144            },
145            move |ctx| unsafe { com2.EndGetPartitionLoadInformation(ctx) },
146            cancellation_token,
147        )
148    }
149
150    fn get_deployed_replica_detail_internal(
151        &self,
152        desc: &FABRIC_DEPLOYED_SERVICE_REPLICA_DETAIL_QUERY_DESCRIPTION,
153        timeout_milliseconds: u32,
154        cancellation_token: Option<BoxedCancelToken>,
155    ) -> FabricReceiver<crate::Result<IFabricGetDeployedServiceReplicaDetailResult>> {
156        let com1 = &self.com;
157        let com2 = self.com.clone();
158        fabric_begin_end_proxy(
159            move |callback| unsafe {
160                com1.BeginGetDeployedReplicaDetail(desc, timeout_milliseconds, callback)
161            },
162            move |ctx| unsafe { com2.EndGetDeployedReplicaDetail(ctx) },
163            cancellation_token,
164        )
165    }
166
167    fn get_deployed_application_list_internal(
168        &self,
169        desc: &mssf_com::FabricTypes::FABRIC_DEPLOYED_APPLICATION_QUERY_DESCRIPTION,
170        timeout_milliseconds: u32,
171        cancellation_token: Option<BoxedCancelToken>,
172    ) -> FabricReceiver<
173        crate::Result<mssf_com::FabricClient::IFabricGetDeployedApplicationListResult>,
174    > {
175        let com1 = &self.com;
176        let com2 = self.com.clone();
177        fabric_begin_end_proxy(
178            move |callback| unsafe {
179                com1.BeginGetDeployedApplicationList(desc, timeout_milliseconds, callback)
180            },
181            move |ctx| unsafe { com2.EndGetDeployedApplicationList(ctx) },
182            cancellation_token,
183        )
184    }
185
186    fn get_deployed_service_package_list_internal(
187        &self,
188        desc: &mssf_com::FabricTypes::FABRIC_DEPLOYED_SERVICE_PACKAGE_QUERY_DESCRIPTION,
189        timeout_milliseconds: u32,
190        cancellation_token: Option<BoxedCancelToken>,
191    ) -> FabricReceiver<
192        crate::Result<mssf_com::FabricClient::IFabricGetDeployedServicePackageListResult>,
193    > {
194        let com1 = &self.com;
195        let com2 = self.com.clone();
196        fabric_begin_end_proxy(
197            move |callback| unsafe {
198                com1.BeginGetDeployedServicePackageList(desc, timeout_milliseconds, callback)
199            },
200            move |ctx| unsafe { com2.EndGetDeployedServicePackageList(ctx) },
201            cancellation_token,
202        )
203    }
204}
205
206impl From<IFabricQueryClient13> for QueryClient {
207    fn from(com: IFabricQueryClient13) -> Self {
208        Self { com }
209    }
210}
211
212impl From<QueryClient> for IFabricQueryClient13 {
213    fn from(value: QueryClient) -> Self {
214        value.com
215    }
216}
217
218impl QueryClient {
219    // List nodes in the cluster
220    pub async fn get_node_list(
221        &self,
222        desc: &NodeQueryDescription,
223        timeout: Duration,
224        cancellation_token: Option<BoxedCancelToken>,
225    ) -> crate::Result<NodeListResult> {
226        let com = {
227            let mut pool = BoxPool::new();
228            let arg = desc.get_raw_with_pool(&mut pool);
229            self.get_node_list_internal(
230                &arg,
231                timeout.as_millis().try_into().unwrap(),
232                cancellation_token,
233            )
234        }
235        .await??;
236        Ok(NodeListResult::from(&com))
237    }
238
239    pub async fn get_application_list(
240        &self,
241        desc: &crate::types::ApplicationQueryDescription,
242        timeout: Duration,
243        cancellation_token: Option<BoxedCancelToken>,
244    ) -> crate::Result<crate::types::ApplicationListResult> {
245        let com = {
246            let mut pool = BoxPool::new();
247            let arg = desc.get_raw_with_pool(&mut pool);
248            self.get_application_list_internal(
249                &arg,
250                timeout.as_millis().try_into().unwrap(),
251                cancellation_token,
252            )
253        }
254        .await??;
255        Ok(crate::types::ApplicationListResult::from(&com))
256    }
257    pub async fn get_service_list(
258        &self,
259        desc: &ServiceQueryDescription,
260        timeout: Duration,
261        cancellation_token: Option<BoxedCancelToken>,
262    ) -> crate::Result<crate::types::ServiceListResult> {
263        let com = {
264            let mut pool = BoxPool::new();
265            let arg = desc.get_raw_with_pool(&mut pool);
266            self.get_service_list_internal(&arg, timeout.as_millis() as u32, cancellation_token)
267        }
268        .await??;
269        Ok(crate::types::ServiceListResult::from(&com))
270    }
271
272    pub async fn get_partition_list(
273        &self,
274        desc: &ServicePartitionQueryDescription,
275        timeout: Duration,
276        cancellation_token: Option<BoxedCancelToken>,
277    ) -> crate::Result<ServicePartitionList> {
278        let com = {
279            let raw: FABRIC_SERVICE_PARTITION_QUERY_DESCRIPTION = desc.into();
280            let mili = timeout.as_millis() as u32;
281            self.get_partition_list_internal(&raw, mili, cancellation_token)
282        }
283        .await??;
284        Ok(ServicePartitionList::from(&com))
285    }
286
287    pub async fn get_replica_list(
288        &self,
289        desc: &ServiceReplicaQueryDescription,
290        timeout: Duration,
291        cancellation_token: Option<BoxedCancelToken>,
292    ) -> crate::Result<ServiceReplicaList> {
293        let com = {
294            let raw: FABRIC_SERVICE_REPLICA_QUERY_DESCRIPTION = desc.into();
295            let mili = timeout.as_millis() as u32;
296            self.get_replica_list_internal(&raw, mili, cancellation_token)
297        }
298        .await??;
299        Ok(ServiceReplicaList::from(&com))
300    }
301
302    pub async fn get_partition_load_information(
303        &self,
304        desc: &PartitionLoadInformationQueryDescription,
305        timeout: Duration,
306        cancellation_token: Option<BoxedCancelToken>,
307    ) -> crate::Result<GetPartitionLoadInformationResult> {
308        let com = {
309            let raw: FABRIC_PARTITION_LOAD_INFORMATION_QUERY_DESCRIPTION = desc.into();
310            let timeout_ms = timeout.as_micros() as u32;
311            self.get_partition_load_information_internal(&raw, timeout_ms, cancellation_token)
312        }
313        .await??;
314        Ok(GetPartitionLoadInformationResult::from(&com))
315    }
316
317    pub async fn get_deployed_replica_detail(
318        &self,
319        desc: &DeployedServiceReplicaDetailQueryDescription,
320        timeout: Duration,
321        cancellation_token: Option<BoxedCancelToken>,
322    ) -> crate::Result<DeployedServiceReplicaDetailQueryResult> {
323        let com = {
324            let raw: FABRIC_DEPLOYED_SERVICE_REPLICA_DETAIL_QUERY_DESCRIPTION = desc.into();
325            let timeout_ms = timeout.as_micros() as u32;
326            self.get_deployed_replica_detail_internal(&raw, timeout_ms, cancellation_token)
327        }
328        .await??;
329        Ok(DeployedServiceReplicaDetailQueryResult::new(com))
330    }
331
332    /// Lists the applications deployed on a node.
333    pub async fn get_deployed_application_list(
334        &self,
335        desc: &crate::types::DeployedApplicationQueryDescription,
336        timeout: Duration,
337        cancellation_token: Option<BoxedCancelToken>,
338    ) -> crate::Result<crate::types::DeployedApplicationList> {
339        let com = {
340            let raw: mssf_com::FabricTypes::FABRIC_DEPLOYED_APPLICATION_QUERY_DESCRIPTION =
341                desc.into();
342            let timeout_ms = timeout.as_millis() as u32;
343            self.get_deployed_application_list_internal(&raw, timeout_ms, cancellation_token)
344        }
345        .await??;
346        Ok(crate::types::DeployedApplicationList::from(&com))
347    }
348
349    /// Lists the service packages deployed for an application on a node.
350    pub async fn get_deployed_service_package_list(
351        &self,
352        desc: &crate::types::DeployedServicePackageQueryDescription,
353        timeout: Duration,
354        cancellation_token: Option<BoxedCancelToken>,
355    ) -> crate::Result<crate::types::DeployedServicePackageList> {
356        let com = {
357            let raw: mssf_com::FabricTypes::FABRIC_DEPLOYED_SERVICE_PACKAGE_QUERY_DESCRIPTION =
358                desc.into();
359            let timeout_ms = timeout.as_millis() as u32;
360            self.get_deployed_service_package_list_internal(&raw, timeout_ms, cancellation_token)
361        }
362        .await??;
363        Ok(crate::types::DeployedServicePackageList::from(&com))
364    }
365}