1use serde::{Deserialize, Serialize};
2use serde_json::Value;
3
4use crate::client::http::{ClientResult, HttpClient};
5
6#[derive(Debug, Default)]
7pub struct AdminResourcesRequest {
8 pub limit: Option<usize>,
9 pub include_columnar_columns: Option<bool>,
10 pub columnar_column_limit: Option<usize>,
11 pub kv_prefix: Option<String>,
12}
13
14#[derive(Debug, Deserialize)]
15pub struct AdminResourcesResponse {
16 pub sql_tables: Vec<SqlTableResource>,
17 pub columnar_segments: Vec<ColumnarSegmentResource>,
18 pub kv_keys: Vec<String>,
19 pub truncated: TruncatedSections,
20}
21
22#[derive(Debug, Deserialize)]
23pub struct TruncatedSections {
24 pub sql_tables: bool,
25 pub columnar_segments: bool,
26 pub kv_keys: bool,
27}
28
29#[derive(Debug, Deserialize)]
30pub struct SqlTableResource {
31 pub name: String,
32 pub columns: Vec<SqlColumnResource>,
33}
34
35#[derive(Debug, Deserialize)]
36pub struct SqlColumnResource {
37 pub name: String,
38 #[allow(dead_code)]
39 pub data_type: String,
40}
41
42#[derive(Debug, Deserialize)]
43pub struct ColumnarSegmentResource {
44 pub id: String,
45 pub columns: Option<Vec<String>>,
46}
47
48#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
52#[serde(rename_all = "snake_case")]
53pub enum ClusterManagementOperation {
54 MetadataShow,
55 MembersList,
56 MembersReplace,
57 RangesList,
58 RangesRegister,
59 RangesUpdate,
60 RangesRetire,
61 PlacementGet,
62 PlacementSet,
63 PlacementReplace,
64 ReadPolicyGet,
65 ReadPolicySet,
66 SchemaOwnerGet,
67 SchemaOwnerSet,
68 SchemaRolloutStart,
69 SchemaRolloutStatus,
70 RecoveryStatus,
71 RecoveryRestore,
72 UpgradeStatus,
73 UpgradeStart,
74}
75
76#[derive(Debug, Serialize)]
80pub struct ClusterManagementRequest {
81 pub request_id: String,
82 pub operation: ClusterManagementOperation,
83 #[serde(skip_serializing_if = "Option::is_none")]
84 pub expected_version: Option<u64>,
85 #[serde(skip_serializing_if = "Option::is_none")]
86 pub target: Option<Value>,
87 pub confirmed: bool,
88}
89
90impl ClusterManagementRequest {
91 pub fn new(
92 request_id: impl Into<String>,
93 operation: ClusterManagementOperation,
94 expected_version: Option<u64>,
95 target: Option<Value>,
96 confirmed: bool,
97 ) -> Self {
98 Self {
99 request_id: request_id.into(),
100 operation,
101 expected_version,
102 target,
103 confirmed,
104 }
105 }
106}
107
108#[derive(Debug, Deserialize)]
110pub struct ClusterManagementResponse {
111 pub operation_id: String,
112 #[allow(dead_code)]
113 pub operation: String,
114 pub outcome_class: String,
115 pub reason: String,
116 pub state_version: Option<u64>,
117 pub control: ClusterControlAvailability,
118 pub actor: Option<String>,
119}
120
121#[derive(Debug, Deserialize)]
123pub struct ClusterControlAvailability {
124 pub available: bool,
125 pub mode: String,
126 pub reason: String,
127 pub missing_prerequisites: Vec<Value>,
128}
129
130pub async fn fetch_admin_resources(
131 client: &HttpClient,
132 request: &AdminResourcesRequest,
133) -> ClientResult<AdminResourcesResponse> {
134 let path = build_query_path(request);
135 client.get_json(&path).await
136}
137
138pub async fn invoke_cluster_management(
142 client: &HttpClient,
143 request: &ClusterManagementRequest,
144) -> ClientResult<ClusterManagementResponse> {
145 client
146 .post_json("api/admin/cluster/operations", request)
147 .await
148}
149
150fn build_query_path(request: &AdminResourcesRequest) -> String {
151 let mut params = Vec::new();
152 if let Some(limit) = request.limit {
153 params.push(format!("limit={limit}"));
154 }
155 if let Some(include) = request.include_columnar_columns {
156 params.push(format!("include_columnar_columns={include}"));
157 }
158 if let Some(columnar_column_limit) = request.columnar_column_limit {
159 params.push(format!("columnar_column_limit={columnar_column_limit}"));
160 }
161 if let Some(prefix) = request.kv_prefix.as_deref() {
162 params.push(format!("kv_prefix={}", encode_query_component(prefix)));
163 }
164
165 if params.is_empty() {
166 "api/admin/resources".to_string()
167 } else {
168 format!("api/admin/resources?{}", params.join("&"))
169 }
170}
171
172fn encode_query_component(value: &str) -> String {
173 let mut out = String::with_capacity(value.len());
174 for b in value.bytes() {
175 match b {
176 b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => {
177 out.push(b as char)
178 }
179 _ => out.push_str(&format!("%{b:02X}")),
180 }
181 }
182 out
183}
184
185#[cfg(test)]
186mod tests {
187 use super::{
188 build_query_path, AdminResourcesRequest, ClusterManagementOperation,
189 ClusterManagementRequest,
190 };
191 use serde_json::json;
192
193 #[test]
194 fn build_query_path_includes_params() {
195 let request = AdminResourcesRequest {
196 limit: Some(10),
197 include_columnar_columns: Some(true),
198 columnar_column_limit: Some(5),
199 kv_prefix: Some("app/".to_string()),
200 };
201 let path = build_query_path(&request);
202 assert!(path.starts_with("api/admin/resources?"));
203 assert!(path.contains("limit=10"));
204 assert!(path.contains("include_columnar_columns=true"));
205 assert!(path.contains("columnar_column_limit=5"));
206 assert!(path.contains("kv_prefix=app%2F"));
207 }
208
209 #[test]
210 fn cluster_management_request_preserves_idempotency_and_confirmation() {
211 let request = ClusterManagementRequest::new(
212 "operation-42",
213 ClusterManagementOperation::RangesRegister,
214 Some(9),
215 Some(json!({"range_id": "primary/0"})),
216 true,
217 );
218
219 assert_eq!(
220 serde_json::to_value(request).unwrap(),
221 json!({
222 "request_id": "operation-42",
223 "operation": "ranges_register",
224 "expected_version": 9,
225 "target": {"range_id": "primary/0"},
226 "confirmed": true,
227 })
228 );
229 }
230}