1use std::io::Write;
4
5use serde::Deserialize;
6
7use crate::batch::ClusterManagementOutcome;
8use crate::cli::{
9 ClusterCommand, ClusterMembersCommand, ClusterMetadataCommand, ClusterMutationRequest,
10 ClusterOperationRequest, ClusterPlacementCommand, ClusterRangesCommand,
11 ClusterReadPolicyCommand, ClusterRecoveryCommand, ClusterSchemaCommand,
12 ClusterSchemaOwnerCommand, ClusterSchemaRolloutCommand, ClusterTargetedReadRequest,
13 ClusterUpgradeCommand, CompactionCommand, OutputFormat, ServerCommand,
14};
15use crate::client::admin_resources::{
16 invoke_cluster_management, ClusterManagementOperation, ClusterManagementRequest,
17 ClusterManagementResponse,
18};
19use crate::client::http::{ClientError, HttpClient};
20use crate::error::{CliError, Result};
21use crate::models::{Column, Row};
22use crate::output::server as server_output;
23use crate::output::table::TableFormatter;
24use crate::output::Formatter;
25use crate::tui::renderer::render_output;
26
27#[derive(Debug, Deserialize)]
28struct ServerStatusResponse {
29 version: Option<String>,
30 uptime_secs: Option<u64>,
31 connections: Option<u64>,
32 queries_per_second: Option<f64>,
33 cluster: Option<ServerClusterStatus>,
34}
35
36#[derive(Debug, Deserialize)]
37struct ServerMetricsResponse {
38 qps: Option<f64>,
39 avg_latency_ms: Option<f64>,
40 p99_latency_ms: Option<f64>,
41 memory_usage_mb: Option<u64>,
42 active_connections: Option<u64>,
43}
44
45#[derive(Debug, Deserialize)]
46struct ServerHealthResponse {
47 status: Option<String>,
48 message: Option<String>,
49 degraded: Option<bool>,
50 cluster: Option<ServerClusterStatus>,
51}
52
53#[derive(Debug, Deserialize)]
54struct ServerCompactionResponse {
55 success: Option<bool>,
56 message: Option<String>,
57}
58
59#[derive(Debug, Deserialize)]
60struct ServerClusterOperationResponse {
61 action: Option<String>,
62 cluster: Option<ServerClusterStatus>,
63}
64
65#[derive(Debug, Deserialize)]
66struct ServerClusterStatus {
67 schema_version: Option<u32>,
68 mode: Option<String>,
69 identity: Option<ServerClusterIdentity>,
70 routing_capabilities: Option<ServerRoutingCapabilities>,
71 degraded: Option<bool>,
72 diagnostics: Option<Vec<ServerClusterDiagnostic>>,
73}
74
75#[derive(Debug, Deserialize)]
76struct ServerClusterIdentity {
77 node_id: Option<String>,
78 lifecycle_state: Option<String>,
79}
80
81#[derive(Debug, Deserialize)]
82struct ServerRoutingCapabilities {
83 local_only: Option<bool>,
84 future_distributed_execution_required: Option<bool>,
85 scatter_gather_simulated: Option<bool>,
86}
87
88#[derive(Debug, Deserialize)]
89struct ServerClusterDiagnostic {
90 code: Option<String>,
91}
92
93#[allow(dead_code)] pub async fn execute_remote<W: Write>(
96 client: &HttpClient,
97 cmd: &ServerCommand,
98 writer: &mut W,
99 quiet: bool,
100) -> Result<()> {
101 execute_remote_with_format(client, cmd, writer, quiet, OutputFormat::Table).await
102}
103
104pub async fn execute_remote_with_format<W: Write>(
107 client: &HttpClient,
108 cmd: &ServerCommand,
109 writer: &mut W,
110 quiet: bool,
111 output_format: OutputFormat,
112) -> Result<()> {
113 match cmd {
114 ServerCommand::Status => {
115 let response: ServerStatusResponse = client
116 .get_json("api/admin/status")
117 .await
118 .map_err(map_client_error)?;
119 if quiet {
120 return Ok(());
121 }
122 render_table(
123 writer,
124 server_output::status_columns(),
125 vec![status_row_from_response(&response)],
126 )
127 }
128 ServerCommand::Metrics => {
129 let response: ServerMetricsResponse = client
130 .get_json("api/admin/metrics")
131 .await
132 .map_err(map_client_error)?;
133 if quiet {
134 return Ok(());
135 }
136 render_table(
137 writer,
138 server_output::metrics_columns(),
139 vec![server_output::metrics_row(
140 response.qps,
141 response.avg_latency_ms,
142 response.p99_latency_ms,
143 response.memory_usage_mb,
144 response.active_connections,
145 )],
146 )
147 }
148 ServerCommand::Health => {
149 let response: ServerHealthResponse = client
150 .get_json("api/admin/health")
151 .await
152 .map_err(map_client_error)?;
153 if quiet {
154 return Ok(());
155 }
156 render_table(
157 writer,
158 server_output::health_columns(),
159 vec![health_row_from_response(&response)],
160 )
161 }
162 ServerCommand::Join => {
163 let response = execute_cluster_operation(client, "join").await?;
164 if quiet {
165 return Ok(());
166 }
167 render_cluster_operation(writer, &response)
168 }
169 ServerCommand::Leave => {
170 let response = execute_cluster_operation(client, "leave").await?;
171 if quiet {
172 return Ok(());
173 }
174 render_cluster_operation(writer, &response)
175 }
176 ServerCommand::Compaction { command } => match command {
177 CompactionCommand::Trigger => {
178 let request = serde_json::json!({});
179 let response: ServerCompactionResponse = client
180 .post_json("api/admin/compaction", &request)
181 .await
182 .map_err(map_client_error)?;
183 if quiet {
184 return Ok(());
185 }
186 render_table(
187 writer,
188 server_output::compaction_columns(),
189 vec![server_output::compaction_row(
190 response.success,
191 response.message.as_deref(),
192 )],
193 )
194 }
195 },
196 ServerCommand::Cluster { command } => {
197 let response = execute_cluster_management_command(client, command).await?;
198 if quiet {
199 return ensure_cluster_management_succeeded(&response);
200 }
201 render_cluster_management(writer, &response, output_format)?;
202 ensure_cluster_management_succeeded(&response)
203 }
204 }
205}
206
207pub async fn execute_remote_tui(
208 client: &HttpClient,
209 cmd: &ServerCommand,
210 quiet: bool,
211 connection_label: impl Into<String>,
212 output_format: OutputFormat,
213 admin_launcher: Option<Box<dyn FnMut() -> Result<()> + '_>>,
214) -> Result<()> {
215 match cmd {
216 ServerCommand::Status => {
217 let response: ServerStatusResponse = client
218 .get_json("api/admin/status")
219 .await
220 .map_err(map_client_error)?;
221 if quiet {
222 return Ok(());
223 }
224 render_output(
225 server_output::status_columns(),
226 vec![status_row_from_response(&response)],
227 connection_label,
228 Some(server_command_context(cmd)),
229 true,
230 None,
231 output_format,
232 admin_launcher,
233 )
234 }
235 ServerCommand::Metrics => {
236 let response: ServerMetricsResponse = client
237 .get_json("api/admin/metrics")
238 .await
239 .map_err(map_client_error)?;
240 if quiet {
241 return Ok(());
242 }
243 render_output(
244 server_output::metrics_columns(),
245 vec![server_output::metrics_row(
246 response.qps,
247 response.avg_latency_ms,
248 response.p99_latency_ms,
249 response.memory_usage_mb,
250 response.active_connections,
251 )],
252 connection_label,
253 Some(server_command_context(cmd)),
254 true,
255 None,
256 output_format,
257 admin_launcher,
258 )
259 }
260 ServerCommand::Health => {
261 let response: ServerHealthResponse = client
262 .get_json("api/admin/health")
263 .await
264 .map_err(map_client_error)?;
265 if quiet {
266 return Ok(());
267 }
268 render_output(
269 server_output::health_columns(),
270 vec![health_row_from_response(&response)],
271 connection_label,
272 Some(server_command_context(cmd)),
273 true,
274 None,
275 output_format,
276 admin_launcher,
277 )
278 }
279 ServerCommand::Join => {
280 let response = execute_cluster_operation(client, "join").await?;
281 if quiet {
282 return Ok(());
283 }
284 render_output(
285 server_output::cluster_operation_columns(),
286 vec![cluster_operation_row(&response)],
287 connection_label,
288 Some(server_command_context(cmd)),
289 true,
290 None,
291 output_format,
292 admin_launcher,
293 )
294 }
295 ServerCommand::Leave => {
296 let response = execute_cluster_operation(client, "leave").await?;
297 if quiet {
298 return Ok(());
299 }
300 render_output(
301 server_output::cluster_operation_columns(),
302 vec![cluster_operation_row(&response)],
303 connection_label,
304 Some(server_command_context(cmd)),
305 true,
306 None,
307 output_format,
308 admin_launcher,
309 )
310 }
311 ServerCommand::Compaction { command } => match command {
312 CompactionCommand::Trigger => {
313 let request = serde_json::json!({});
314 let response: ServerCompactionResponse = client
315 .post_json("api/admin/compaction", &request)
316 .await
317 .map_err(map_client_error)?;
318 if quiet {
319 return Ok(());
320 }
321 render_output(
322 server_output::compaction_columns(),
323 vec![server_output::compaction_row(
324 response.success,
325 response.message.as_deref(),
326 )],
327 connection_label,
328 Some(server_command_context(cmd)),
329 true,
330 None,
331 output_format,
332 admin_launcher,
333 )
334 }
335 },
336 ServerCommand::Cluster { command } => {
337 let response = execute_cluster_management_command(client, command).await?;
338 if quiet {
339 return ensure_cluster_management_succeeded(&response);
340 }
341 render_output(
342 server_output::cluster_management_columns(),
343 vec![cluster_management_row(&response)],
344 connection_label,
345 Some(server_command_context(cmd)),
346 true,
347 None,
348 output_format,
349 admin_launcher,
350 )?;
351 ensure_cluster_management_succeeded(&response)
352 }
353 }
354}
355
356fn server_command_context(cmd: &ServerCommand) -> String {
357 match cmd {
358 ServerCommand::Status => "server status".to_string(),
359 ServerCommand::Metrics => "server metrics".to_string(),
360 ServerCommand::Health => "server health".to_string(),
361 ServerCommand::Join => "server join".to_string(),
362 ServerCommand::Leave => "server leave".to_string(),
363 ServerCommand::Compaction { .. } => "server compaction trigger".to_string(),
364 ServerCommand::Cluster { .. } => "server cluster management".to_string(),
365 }
366}
367
368fn render_table<W: Write>(writer: &mut W, columns: Vec<Column>, rows: Vec<Row>) -> Result<()> {
369 let mut formatter = TableFormatter::new();
370 formatter.write_header(writer, &columns)?;
371 for row in rows {
372 formatter.write_row(writer, &row)?;
373 }
374 formatter.write_footer(writer)
375}
376
377fn status_row_from_response(response: &ServerStatusResponse) -> Row {
378 let cluster = ClusterDisplayFields::from(response.cluster.as_ref());
379 server_output::status_row(server_output::StatusRowFields {
380 version: response.version.as_deref(),
381 uptime_secs: response.uptime_secs,
382 connections: response.connections,
383 qps: response.queries_per_second,
384 cluster_schema_version: cluster.schema_version,
385 cluster_mode: cluster.mode,
386 node_id: cluster.node_id,
387 lifecycle_state: cluster.lifecycle_state,
388 degraded: cluster.degraded,
389 local_only: cluster.local_only,
390 future_distributed: cluster.future_distributed,
391 scatter_gather: cluster.scatter_gather,
392 diagnostics: cluster.diagnostics.as_deref(),
393 })
394}
395
396fn health_row_from_response(response: &ServerHealthResponse) -> Row {
397 let cluster = ClusterDisplayFields::from(response.cluster.as_ref());
398 server_output::health_row(
399 response.status.as_deref(),
400 response.message.as_deref(),
401 response.degraded.or(cluster.degraded),
402 cluster.mode,
403 cluster.node_id,
404 )
405}
406
407async fn execute_cluster_operation(
408 client: &HttpClient,
409 action: &str,
410) -> Result<ServerClusterOperationResponse> {
411 let request = serde_json::json!({});
412 let path = format!("api/admin/cluster/{action}");
413 client
414 .post_json(&path, &request)
415 .await
416 .map_err(map_client_error)
417}
418
419fn render_cluster_operation<W: Write>(
420 writer: &mut W,
421 response: &ServerClusterOperationResponse,
422) -> Result<()> {
423 render_table(
424 writer,
425 server_output::cluster_operation_columns(),
426 vec![cluster_operation_row(response)],
427 )
428}
429
430async fn execute_cluster_management_command(
431 client: &HttpClient,
432 command: &ClusterCommand,
433) -> Result<ClusterManagementResponse> {
434 let request = cluster_management_request(command)?;
435 invoke_cluster_management(client, &request)
436 .await
437 .map_err(map_client_error)
438}
439
440struct ClusterInvocation<'a> {
441 operation: ClusterManagementOperation,
442 request: &'a ClusterOperationRequest,
443 target: Option<&'a str>,
444 confirmed: bool,
445}
446
447impl<'a> ClusterInvocation<'a> {
448 fn read(operation: ClusterManagementOperation, request: &'a ClusterOperationRequest) -> Self {
449 Self {
450 operation,
451 request,
452 target: None,
453 confirmed: false,
454 }
455 }
456
457 fn targeted_read(
458 operation: ClusterManagementOperation,
459 request: &'a ClusterTargetedReadRequest,
460 ) -> Self {
461 Self {
462 operation,
463 request: &request.operation,
464 target: Some(&request.target),
465 confirmed: false,
466 }
467 }
468
469 fn mutation(
470 operation: ClusterManagementOperation,
471 request: &'a ClusterMutationRequest,
472 ) -> Self {
473 Self {
474 operation,
475 request: &request.operation,
476 target: Some(&request.target),
477 confirmed: request.confirm,
478 }
479 }
480}
481
482fn cluster_management_request(command: &ClusterCommand) -> Result<ClusterManagementRequest> {
483 let invocation = match command {
484 ClusterCommand::Metadata {
485 command: ClusterMetadataCommand::Show { request },
486 } => ClusterInvocation::read(ClusterManagementOperation::MetadataShow, request),
487 ClusterCommand::Members {
488 command: ClusterMembersCommand::List { request },
489 } => ClusterInvocation::read(ClusterManagementOperation::MembersList, request),
490 ClusterCommand::Members {
491 command: ClusterMembersCommand::Replace { request },
492 } => ClusterInvocation::mutation(ClusterManagementOperation::MembersReplace, request),
493 ClusterCommand::Ranges {
494 command: ClusterRangesCommand::List { request },
495 } => ClusterInvocation::read(ClusterManagementOperation::RangesList, request),
496 ClusterCommand::Ranges {
497 command: ClusterRangesCommand::Show { request },
498 } => ClusterInvocation::targeted_read(ClusterManagementOperation::RangesList, request),
499 ClusterCommand::Ranges {
500 command: ClusterRangesCommand::Register { request },
501 } => ClusterInvocation::mutation(ClusterManagementOperation::RangesRegister, request),
502 ClusterCommand::Ranges {
503 command: ClusterRangesCommand::Update { request },
504 } => ClusterInvocation::mutation(ClusterManagementOperation::RangesUpdate, request),
505 ClusterCommand::Ranges {
506 command: ClusterRangesCommand::Retire { request },
507 } => ClusterInvocation::mutation(ClusterManagementOperation::RangesRetire, request),
508 ClusterCommand::Placement {
509 command: ClusterPlacementCommand::Get { request },
510 } => ClusterInvocation::targeted_read(ClusterManagementOperation::PlacementGet, request),
511 ClusterCommand::Placement {
512 command: ClusterPlacementCommand::Set { request },
513 } => ClusterInvocation::mutation(ClusterManagementOperation::PlacementSet, request),
514 ClusterCommand::Placement {
515 command: ClusterPlacementCommand::Replace { request },
516 } => ClusterInvocation::mutation(ClusterManagementOperation::PlacementReplace, request),
517 ClusterCommand::ReadPolicy {
518 command: ClusterReadPolicyCommand::Get { request },
519 } => ClusterInvocation::read(ClusterManagementOperation::ReadPolicyGet, request),
520 ClusterCommand::ReadPolicy {
521 command: ClusterReadPolicyCommand::Set { request },
522 } => ClusterInvocation::mutation(ClusterManagementOperation::ReadPolicySet, request),
523 ClusterCommand::Schema {
524 command:
525 ClusterSchemaCommand::Owner {
526 command: ClusterSchemaOwnerCommand::Get { request },
527 },
528 } => ClusterInvocation::read(ClusterManagementOperation::SchemaOwnerGet, request),
529 ClusterCommand::Schema {
530 command:
531 ClusterSchemaCommand::Owner {
532 command: ClusterSchemaOwnerCommand::Set { request },
533 },
534 } => ClusterInvocation::mutation(ClusterManagementOperation::SchemaOwnerSet, request),
535 ClusterCommand::Schema {
536 command:
537 ClusterSchemaCommand::Rollout {
538 command: ClusterSchemaRolloutCommand::Start { request },
539 },
540 } => ClusterInvocation::mutation(ClusterManagementOperation::SchemaRolloutStart, request),
541 ClusterCommand::Schema {
542 command:
543 ClusterSchemaCommand::Rollout {
544 command: ClusterSchemaRolloutCommand::Status { request },
545 },
546 } => ClusterInvocation::read(ClusterManagementOperation::SchemaRolloutStatus, request),
547 ClusterCommand::Recovery {
548 command: ClusterRecoveryCommand::Status { request },
549 } => ClusterInvocation::read(ClusterManagementOperation::RecoveryStatus, request),
550 ClusterCommand::Recovery {
551 command: ClusterRecoveryCommand::Restore { request },
552 } => ClusterInvocation::mutation(ClusterManagementOperation::RecoveryRestore, request),
553 ClusterCommand::Upgrade {
554 command: ClusterUpgradeCommand::Status { request },
555 } => ClusterInvocation::read(ClusterManagementOperation::UpgradeStatus, request),
556 ClusterCommand::Upgrade {
557 command: ClusterUpgradeCommand::Start { request },
558 } => ClusterInvocation::mutation(ClusterManagementOperation::UpgradeStart, request),
559 };
560
561 let target = invocation
562 .target
563 .map(|target| {
564 serde_json::from_str(target).map_err(|err| {
565 CliError::InvalidArgument(format!(
566 "cluster management target must be valid JSON: {err}"
567 ))
568 })
569 })
570 .transpose()?;
571 Ok(ClusterManagementRequest::new(
572 &invocation.request.request_id,
573 invocation.operation,
574 invocation.request.expected_version,
575 target,
576 invocation.confirmed,
577 ))
578}
579
580fn render_cluster_management<W: Write>(
581 writer: &mut W,
582 response: &ClusterManagementResponse,
583 output_format: OutputFormat,
584) -> Result<()> {
585 let mut formatter = crate::output::create_formatter(output_format);
586 let columns = server_output::cluster_management_columns();
587 let row = cluster_management_row(response);
588 formatter.write_header(writer, &columns)?;
589 formatter.write_row(writer, &row)?;
590 formatter.write_footer(writer)
591}
592
593fn cluster_management_row(response: &ClusterManagementResponse) -> Row {
594 let prerequisites = response
595 .control
596 .missing_prerequisites
597 .iter()
598 .map(serde_json::Value::to_string)
599 .collect::<Vec<_>>()
600 .join(",");
601 server_output::cluster_management_row(server_output::ClusterManagementRowFields {
602 operation_id: &response.operation_id,
603 operation: &response.operation,
604 outcome_class: &response.outcome_class,
605 reason: &response.reason,
606 state_version: response.state_version,
607 control_available: response.control.available,
608 control_mode: &response.control.mode,
609 control_reason: &response.control.reason,
610 missing_prerequisites: (!prerequisites.is_empty()).then_some(&prerequisites),
611 actor: response.actor.as_deref(),
612 })
613}
614
615fn ensure_cluster_management_succeeded(response: &ClusterManagementResponse) -> Result<()> {
616 let outcome = ClusterManagementOutcome::from_wire(&response.outcome_class, &response.reason);
617 if outcome.is_success() {
618 return Ok(());
619 }
620 Err(CliError::ClusterManagementOutcome {
621 outcome: response.outcome_class.clone(),
622 reason: response.reason.clone(),
623 exit_code: outcome.exit_code(),
624 })
625}
626
627fn cluster_operation_row(response: &ServerClusterOperationResponse) -> Row {
628 let cluster = response.cluster.as_ref();
629 let identity = cluster.and_then(|cluster| cluster.identity.as_ref());
630 server_output::cluster_operation_row(
631 response.action.as_deref(),
632 cluster.and_then(|cluster| cluster.mode.as_deref()),
633 identity.and_then(|identity| identity.node_id.as_deref()),
634 identity.and_then(|identity| identity.lifecycle_state.as_deref()),
635 cluster.and_then(|cluster| cluster.degraded),
636 )
637}
638
639struct ClusterDisplayFields<'a> {
640 mode: Option<&'a str>,
641 schema_version: Option<u32>,
642 node_id: Option<&'a str>,
643 lifecycle_state: Option<&'a str>,
644 degraded: Option<bool>,
645 local_only: Option<bool>,
646 future_distributed: Option<bool>,
647 scatter_gather: Option<bool>,
648 diagnostics: Option<String>,
649}
650
651impl<'a> From<Option<&'a ServerClusterStatus>> for ClusterDisplayFields<'a> {
652 fn from(cluster: Option<&'a ServerClusterStatus>) -> Self {
653 let identity = cluster.and_then(|cluster| cluster.identity.as_ref());
654 let routing = cluster.and_then(|cluster| cluster.routing_capabilities.as_ref());
655 Self {
656 mode: cluster.and_then(|cluster| cluster.mode.as_deref()),
657 schema_version: cluster.and_then(|cluster| cluster.schema_version),
658 node_id: identity.and_then(|identity| identity.node_id.as_deref()),
659 lifecycle_state: identity.and_then(|identity| identity.lifecycle_state.as_deref()),
660 degraded: cluster.and_then(|cluster| cluster.degraded),
661 local_only: routing.and_then(|routing| routing.local_only),
662 future_distributed: routing
663 .and_then(|routing| routing.future_distributed_execution_required),
664 scatter_gather: routing.and_then(|routing| routing.scatter_gather_simulated),
665 diagnostics: cluster.and_then(diagnostic_codes),
666 }
667 }
668}
669
670fn diagnostic_codes(cluster: &ServerClusterStatus) -> Option<String> {
671 let diagnostics = cluster.diagnostics.as_ref()?;
672 let codes = diagnostics
673 .iter()
674 .filter_map(|diagnostic| diagnostic.code.as_deref())
675 .collect::<Vec<_>>();
676 if codes.is_empty() {
677 None
678 } else {
679 Some(codes.join(","))
680 }
681}
682
683fn map_client_error(err: ClientError) -> CliError {
684 match err {
685 ClientError::Request { source, .. } => {
686 CliError::ServerConnection(format!("request failed: {source}"))
687 }
688 ClientError::InvalidUrl(message) => CliError::InvalidArgument(message),
689 ClientError::Build(message) => CliError::InvalidArgument(message),
690 ClientError::Auth(err) => CliError::InvalidArgument(err.to_string()),
691 ClientError::HttpStatus { status, body } => {
692 if status == reqwest::StatusCode::NOT_IMPLEMENTED {
693 CliError::ServerUnsupported(format!(
694 "server returned HTTP {}: {}",
695 status.as_u16(),
696 body
697 ))
698 } else {
699 CliError::InvalidArgument(format!(
700 "Server error: HTTP {} - {}",
701 status.as_u16(),
702 body
703 ))
704 }
705 }
706 }
707}
708
709#[cfg(test)]
710mod tests {
711 use clap::Parser;
712
713 use super::{
714 cluster_management_request, ensure_cluster_management_succeeded, render_cluster_management,
715 ClusterManagementOperation, ClusterManagementResponse,
716 };
717 use crate::batch::ExitCode;
718 use crate::cli::{Cli, ClusterCommand, Command, OutputFormat, ServerCommand};
719 use crate::client::admin_resources::ClusterControlAvailability;
720 use crate::error::CliError;
721
722 fn cluster_command(args: Vec<&str>) -> ClusterCommand {
723 let cli = Cli::try_parse_from(args).unwrap();
724 match cli.command {
725 Some(Command::Server {
726 command: Some(ServerCommand::Cluster { command }),
727 }) => command,
728 other => panic!("expected server cluster command, got {other:?}"),
729 }
730 }
731
732 #[test]
733 fn every_public_cluster_grammar_maps_to_its_management_operation() {
734 let cases = vec![
735 (
736 vec![
737 "alopex",
738 "server",
739 "cluster",
740 "metadata",
741 "show",
742 "--request-id",
743 "metadata-show",
744 ],
745 ClusterManagementOperation::MetadataShow,
746 ),
747 (
748 vec![
749 "alopex",
750 "server",
751 "cluster",
752 "members",
753 "list",
754 "--request-id",
755 "members-list",
756 ],
757 ClusterManagementOperation::MembersList,
758 ),
759 (
760 vec![
761 "alopex",
762 "server",
763 "cluster",
764 "members",
765 "replace",
766 "--request-id",
767 "members-replace",
768 "--target",
769 "{}",
770 "--confirm",
771 ],
772 ClusterManagementOperation::MembersReplace,
773 ),
774 (
775 vec![
776 "alopex",
777 "server",
778 "cluster",
779 "ranges",
780 "list",
781 "--request-id",
782 "ranges-list",
783 ],
784 ClusterManagementOperation::RangesList,
785 ),
786 (
787 vec![
788 "alopex",
789 "server",
790 "cluster",
791 "ranges",
792 "show",
793 "--request-id",
794 "ranges-show",
795 "--target",
796 "{}",
797 ],
798 ClusterManagementOperation::RangesList,
799 ),
800 (
801 vec![
802 "alopex",
803 "server",
804 "cluster",
805 "ranges",
806 "register",
807 "--request-id",
808 "ranges-register",
809 "--target",
810 "{}",
811 "--confirm",
812 ],
813 ClusterManagementOperation::RangesRegister,
814 ),
815 (
816 vec![
817 "alopex",
818 "server",
819 "cluster",
820 "ranges",
821 "update",
822 "--request-id",
823 "ranges-update",
824 "--target",
825 "{}",
826 "--confirm",
827 ],
828 ClusterManagementOperation::RangesUpdate,
829 ),
830 (
831 vec![
832 "alopex",
833 "server",
834 "cluster",
835 "ranges",
836 "retire",
837 "--request-id",
838 "ranges-retire",
839 "--target",
840 "{}",
841 "--confirm",
842 ],
843 ClusterManagementOperation::RangesRetire,
844 ),
845 (
846 vec![
847 "alopex",
848 "server",
849 "cluster",
850 "placement",
851 "get",
852 "--request-id",
853 "placement-get",
854 "--target",
855 "{}",
856 ],
857 ClusterManagementOperation::PlacementGet,
858 ),
859 (
860 vec![
861 "alopex",
862 "server",
863 "cluster",
864 "placement",
865 "set",
866 "--request-id",
867 "placement-set",
868 "--target",
869 "{}",
870 "--confirm",
871 ],
872 ClusterManagementOperation::PlacementSet,
873 ),
874 (
875 vec![
876 "alopex",
877 "server",
878 "cluster",
879 "placement",
880 "replace",
881 "--request-id",
882 "placement-replace",
883 "--target",
884 "{}",
885 "--confirm",
886 ],
887 ClusterManagementOperation::PlacementReplace,
888 ),
889 (
890 vec![
891 "alopex",
892 "server",
893 "cluster",
894 "read-policy",
895 "get",
896 "--request-id",
897 "read-policy-get",
898 ],
899 ClusterManagementOperation::ReadPolicyGet,
900 ),
901 (
902 vec![
903 "alopex",
904 "server",
905 "cluster",
906 "read-policy",
907 "set",
908 "--request-id",
909 "read-policy-set",
910 "--target",
911 "{}",
912 "--confirm",
913 ],
914 ClusterManagementOperation::ReadPolicySet,
915 ),
916 (
917 vec![
918 "alopex",
919 "server",
920 "cluster",
921 "schema",
922 "owner",
923 "get",
924 "--request-id",
925 "schema-owner-get",
926 ],
927 ClusterManagementOperation::SchemaOwnerGet,
928 ),
929 (
930 vec![
931 "alopex",
932 "server",
933 "cluster",
934 "schema",
935 "owner",
936 "set",
937 "--request-id",
938 "schema-owner-set",
939 "--target",
940 "{}",
941 "--confirm",
942 ],
943 ClusterManagementOperation::SchemaOwnerSet,
944 ),
945 (
946 vec![
947 "alopex",
948 "server",
949 "cluster",
950 "schema",
951 "rollout",
952 "start",
953 "--request-id",
954 "schema-rollout-start",
955 "--target",
956 "{}",
957 "--confirm",
958 ],
959 ClusterManagementOperation::SchemaRolloutStart,
960 ),
961 (
962 vec![
963 "alopex",
964 "server",
965 "cluster",
966 "schema",
967 "rollout",
968 "status",
969 "--request-id",
970 "schema-rollout-status",
971 ],
972 ClusterManagementOperation::SchemaRolloutStatus,
973 ),
974 (
975 vec![
976 "alopex",
977 "server",
978 "cluster",
979 "recovery",
980 "status",
981 "--request-id",
982 "recovery-status",
983 ],
984 ClusterManagementOperation::RecoveryStatus,
985 ),
986 (
987 vec![
988 "alopex",
989 "server",
990 "cluster",
991 "recovery",
992 "restore",
993 "--request-id",
994 "recovery-restore",
995 "--target",
996 "{}",
997 "--confirm",
998 ],
999 ClusterManagementOperation::RecoveryRestore,
1000 ),
1001 (
1002 vec![
1003 "alopex",
1004 "server",
1005 "cluster",
1006 "upgrade",
1007 "status",
1008 "--request-id",
1009 "upgrade-status",
1010 ],
1011 ClusterManagementOperation::UpgradeStatus,
1012 ),
1013 (
1014 vec![
1015 "alopex",
1016 "server",
1017 "cluster",
1018 "upgrade",
1019 "start",
1020 "--request-id",
1021 "upgrade-start",
1022 "--target",
1023 "{}",
1024 "--confirm",
1025 ],
1026 ClusterManagementOperation::UpgradeStart,
1027 ),
1028 ];
1029
1030 for (args, expected_operation) in cases {
1031 let command = cluster_command(args);
1032 let request = cluster_management_request(&command).unwrap();
1033 assert_eq!(request.operation, expected_operation);
1034 }
1035 }
1036
1037 #[test]
1038 fn mutation_target_and_confirmation_are_preserved_for_the_http_contract() {
1039 let command = cluster_command(vec![
1040 "alopex",
1041 "server",
1042 "cluster",
1043 "ranges",
1044 "register",
1045 "--request-id",
1046 "range-register-8",
1047 "--expected-version",
1048 "7",
1049 "--target",
1050 r#"{"range_id":"primary/0"}"#,
1051 "--confirm",
1052 ]);
1053 let request = cluster_management_request(&command).unwrap();
1054
1055 assert_eq!(request.request_id, "range-register-8");
1056 assert_eq!(request.expected_version, Some(7));
1057 assert_eq!(
1058 request.target,
1059 Some(serde_json::json!({"range_id": "primary/0"}))
1060 );
1061 assert!(request.confirmed);
1062 }
1063
1064 #[test]
1065 fn invalid_target_json_is_rejected_before_http_invocation() {
1066 let command = cluster_command(vec![
1067 "alopex",
1068 "server",
1069 "cluster",
1070 "placement",
1071 "set",
1072 "--request-id",
1073 "placement-set",
1074 "--target",
1075 "not-json",
1076 "--confirm",
1077 ]);
1078
1079 assert!(cluster_management_request(&command).is_err());
1080 }
1081
1082 fn cluster_response(outcome_class: &str, reason: &str) -> ClusterManagementResponse {
1083 ClusterManagementResponse {
1084 operation_id: "operation-17".to_string(),
1085 operation: "ranges_register".to_string(),
1086 outcome_class: outcome_class.to_string(),
1087 reason: reason.to_string(),
1088 state_version: Some(12),
1089 control: ClusterControlAvailability {
1090 available: true,
1091 mode: "cluster_aware".to_string(),
1092 reason: "ready".to_string(),
1093 missing_prerequisites: Vec::new(),
1094 },
1095 actor: Some("operator-a".to_string()),
1096 }
1097 }
1098
1099 #[test]
1100 fn cluster_management_json_output_retains_the_response_classification() {
1101 let response = cluster_response("pending", "metadata_consensus_adapter_not_attached");
1102 let mut output = Vec::new();
1103 render_cluster_management(&mut output, &response, OutputFormat::Json).unwrap();
1104
1105 let row = &serde_json::from_slice::<serde_json::Value>(&output).unwrap()[0];
1106 assert_eq!(row["Operation ID"], "operation-17");
1107 assert_eq!(row["Operation"], "ranges_register");
1108 assert_eq!(row["Outcome"], "pending");
1109 assert_eq!(row["Reason"], "metadata_consensus_adapter_not_attached");
1110 assert_eq!(row["Control Available"], true);
1111 }
1112
1113 #[test]
1114 fn cluster_management_non_success_outcomes_have_distinct_exit_codes() {
1115 let cases = [
1116 ("pending", "waiting_for_quorum", ExitCode::Warning),
1117 ("retryable_failure", "not_leader", ExitCode::Retryable),
1118 ("terminal_failure", "stale_version", ExitCode::Error),
1119 (
1120 "terminal_failure",
1121 "authorization_denied",
1122 ExitCode::Authorization,
1123 ),
1124 ];
1125
1126 for (outcome_class, reason, exit_code) in cases {
1127 let response = cluster_response(outcome_class, reason);
1128 let err = ensure_cluster_management_succeeded(&response).unwrap_err();
1129 assert!(matches!(
1130 err,
1131 CliError::ClusterManagementOutcome { exit_code: actual, .. } if actual == exit_code
1132 ));
1133 }
1134 assert!(
1135 ensure_cluster_management_succeeded(&cluster_response("succeeded", "committed"))
1136 .is_ok()
1137 );
1138 }
1139}