1use std::path::Path;
2use std::sync::Arc;
3
4use alopex_cluster::ClusterStatusSnapshot;
5use alopex_core::kv::any::AnyKV;
6use axum::extract::{Extension, Path as AxumPath};
7use axum::http::StatusCode;
8use axum::response::{IntoResponse, Response};
9use axum::Json;
10use serde::{Deserialize, Serialize};
11use uuid::Uuid;
12
13use crate::auth::AuthMode;
14use crate::http::{error_response, RequestContext};
15use crate::metrics::ClusterMetricsSurface;
16use crate::ops::backup::{export_snapshot, BackupHandle};
17use crate::ops::restore::{RestoreHandle, RestoreSource};
18use crate::ops::state::{OperationState, RestoreMetadata};
19use crate::ops::status::StatusReporter;
20use crate::ops::status::StatusView;
21use crate::server::ServerState;
22
23#[derive(Serialize)]
24struct AdminCapabilitiesResponse {
25 scope: &'static str,
26 allowed_actions: Vec<&'static str>,
27 unsupported_actions: Vec<&'static str>,
28}
29
30#[derive(Serialize)]
31struct AdminStatusResponse {
32 version: Option<String>,
33 uptime_secs: Option<u64>,
34 connections: Option<u64>,
35 queries_per_second: Option<f64>,
36 cluster: ClusterStatusSnapshot,
37 #[serde(flatten)]
38 status: StatusView,
39}
40
41#[derive(Serialize)]
42struct AdminMetricsResponse {
43 qps: Option<f64>,
44 avg_latency_ms: Option<f64>,
45 p99_latency_ms: Option<f64>,
46 memory_usage_mb: Option<u64>,
47 active_connections: Option<u64>,
48 cluster: ClusterStatusSnapshot,
49 cluster_metrics: ClusterMetricsSurface,
50}
51
52#[derive(Serialize)]
53struct AdminHealthResponse {
54 status: &'static str,
55 message: &'static str,
56 degraded: bool,
57 cluster: ClusterStatusSnapshot,
58}
59
60#[derive(Serialize)]
61struct AdminClusterOperationResponse {
62 action: &'static str,
63 cluster: ClusterStatusSnapshot,
64}
65
66#[derive(Deserialize)]
67pub struct AdminLifecycleRequest {
68 action: String,
69}
70
71#[derive(Deserialize)]
72pub struct AdminRestoreRequest {
73 #[serde(default)]
74 source: Option<String>,
75}
76
77#[derive(Serialize)]
78struct AdminLifecycleResponse {
79 status: &'static str,
80 message: String,
81}
82
83#[derive(Serialize)]
84struct AdminExportResponse {
85 status: &'static str,
86 location: String,
87}
88
89#[derive(Serialize)]
90struct AdminBackupResponse {
91 handle: String,
92 location: String,
93 state: OperationState,
94}
95
96#[derive(Serialize)]
97struct AdminRestoreResponse {
98 handle: String,
99 state: OperationState,
100 metadata: Option<RestoreMetadata>,
101}
102
103pub async fn capabilities(Extension(state): Extension<Arc<ServerState>>) -> impl IntoResponse {
104 let (scope, allowed_actions, unsupported_actions) = capabilities_for_auth(&state.auth);
105 Json(AdminCapabilitiesResponse {
106 scope,
107 allowed_actions,
108 unsupported_actions,
109 })
110}
111
112pub async fn status(
113 Extension(state): Extension<Arc<ServerState>>,
114 Extension(ctx): Extension<RequestContext>,
115) -> Response {
116 let uptime = state.start_time.elapsed().as_secs();
117 let reporter = StatusReporter::new(state.lifecycle_state.clone(), state.recovery_info.clone());
118 let status = reporter.status_view();
119 let cluster = match state.cluster_status_snapshot() {
120 Ok(snapshot) => snapshot,
121 Err(err) => return error_response(err, &ctx),
122 };
123 state.metrics.record_cluster_status(&cluster);
124 Json(AdminStatusResponse {
125 version: Some(env!("CARGO_PKG_VERSION").to_string()),
126 uptime_secs: Some(uptime),
127 connections: None,
128 queries_per_second: None,
129 cluster,
130 status,
131 })
132 .into_response()
133}
134
135pub async fn metrics(
136 Extension(state): Extension<Arc<ServerState>>,
137 Extension(ctx): Extension<RequestContext>,
138) -> Response {
139 let cluster = match state.cluster_status_snapshot() {
140 Ok(snapshot) => snapshot,
141 Err(err) => return error_response(err, &ctx),
142 };
143 state.metrics.record_cluster_status(&cluster);
144 Json(AdminMetricsResponse {
145 qps: None,
146 avg_latency_ms: None,
147 p99_latency_ms: None,
148 memory_usage_mb: None,
149 active_connections: None,
150 cluster_metrics: ClusterMetricsSurface::from(&cluster),
151 cluster,
152 })
153 .into_response()
154}
155
156pub async fn health(
157 Extension(state): Extension<Arc<ServerState>>,
158 Extension(ctx): Extension<RequestContext>,
159) -> Response {
160 let cluster = match state.cluster_status_snapshot() {
161 Ok(snapshot) => snapshot,
162 Err(err) => return error_response(err, &ctx),
163 };
164 state.metrics.record_cluster_status(&cluster);
165 let (status, message) = if cluster.degraded {
166 ("degraded", "cluster status degraded")
167 } else {
168 ("ok", "ready")
169 };
170 Json(AdminHealthResponse {
171 status,
172 message,
173 degraded: cluster.degraded,
174 cluster,
175 })
176 .into_response()
177}
178
179pub async fn cluster_join(
180 Extension(state): Extension<Arc<ServerState>>,
181 Extension(ctx): Extension<RequestContext>,
182) -> Response {
183 cluster_operation_response(&state, &ctx, "join")
184}
185
186pub async fn cluster_leave(
187 Extension(state): Extension<Arc<ServerState>>,
188 Extension(ctx): Extension<RequestContext>,
189) -> Response {
190 cluster_operation_response(&state, &ctx, "leave")
191}
192
193pub async fn compaction(
194 Extension(_state): Extension<Arc<ServerState>>,
195 Extension(ctx): Extension<RequestContext>,
196) -> Response {
197 error_response(
198 crate::error::ServerError::NotImplemented(
199 "manual compaction is not available for the server's LSM storage engine".into(),
200 ),
201 &ctx,
202 )
203}
204
205pub async fn start_backup(
206 Extension(state): Extension<Arc<ServerState>>,
207 Extension(ctx): Extension<RequestContext>,
208) -> Response {
209 match state.backup_coordinator.start_backup().await {
210 Ok(handle) => match backup_response(&state, &handle) {
211 Ok(response) => Json(response).into_response(),
212 Err(err) => error_response(err, &ctx),
213 },
214 Err(err) => error_response(err, &ctx),
215 }
216}
217
218pub async fn export(
219 Extension(state): Extension<Arc<ServerState>>,
220 Extension(ctx): Extension<RequestContext>,
221) -> Response {
222 let export_state = state.clone();
223 let result = tokio::task::spawn_blocking(move || perform_export(export_state.as_ref()))
224 .await
225 .map_err(|err| crate::error::ServerError::Internal(err.to_string()))
226 .and_then(|res| res);
227
228 match result {
229 Ok(location) => Json(AdminExportResponse {
230 status: "OK",
231 location,
232 })
233 .into_response(),
234 Err(err) => error_response(err, &ctx),
235 }
236}
237
238pub async fn backup_status(
239 AxumPath(id): AxumPath<String>,
240 Extension(state): Extension<Arc<ServerState>>,
241 Extension(ctx): Extension<RequestContext>,
242) -> Response {
243 let handle = match parse_backup_handle(&id) {
244 Ok(handle) => handle,
245 Err(err) => return error_response(err, &ctx),
246 };
247 match backup_response(&state, &handle) {
248 Ok(response) => Json(response).into_response(),
249 Err(err) => error_response(err, &ctx),
250 }
251}
252
253pub async fn start_restore(
254 Extension(state): Extension<Arc<ServerState>>,
255 Extension(ctx): Extension<RequestContext>,
256 Json(request): Json<AdminRestoreRequest>,
257) -> Response {
258 let source_path = match request.source {
259 Some(source) => source.into(),
260 None => match crate::ops::restore::resolve_default_source(&state.config.data_dir) {
261 Ok(path) => path,
262 Err(crate::error::ServerError::NotFound(_)) => {
263 match state.backup_coordinator.latest_location() {
264 Some(path) => path,
265 None => {
266 let data_dir = state.config.data_dir.clone();
267 let archive_result = tokio::task::spawn_blocking(move || {
268 perform_lifecycle_action("archive", Path::new(&data_dir))
269 })
270 .await
271 .map_err(|err| crate::error::ServerError::Internal(err.to_string()))
272 .and_then(|res| res.map_err(crate::error::ServerError::BadRequest));
273 if let Err(err) = archive_result {
274 return error_response(err, &ctx);
275 }
276 match crate::ops::restore::resolve_default_source(&state.config.data_dir) {
277 Ok(path) => path,
278 Err(err) => return error_response(err, &ctx),
279 }
280 }
281 }
282 }
283 Err(err) => return error_response(err, &ctx),
284 },
285 };
286 let source = RestoreSource { path: source_path };
287 match state.restore_coordinator.start_restore(source).await {
288 Ok(handle) => match restore_response(&state, &handle) {
289 Ok(response) => Json(response).into_response(),
290 Err(err) => error_response(err, &ctx),
291 },
292 Err(err) => error_response(err, &ctx),
293 }
294}
295
296pub async fn restore_status(
297 AxumPath(id): AxumPath<String>,
298 Extension(state): Extension<Arc<ServerState>>,
299 Extension(ctx): Extension<RequestContext>,
300) -> Response {
301 let handle = match parse_restore_handle(&id) {
302 Ok(handle) => handle,
303 Err(err) => return error_response(err, &ctx),
304 };
305 match restore_response(&state, &handle) {
306 Ok(response) => Json(response).into_response(),
307 Err(err) => error_response(err, &ctx),
308 }
309}
310
311pub async fn lifecycle(
312 Extension(state): Extension<Arc<ServerState>>,
313 Json(request): Json<AdminLifecycleRequest>,
314) -> impl IntoResponse {
315 let data_dir = state.config.data_dir.clone();
316 let action = request.action;
317 let result = tokio::task::spawn_blocking(move || {
318 perform_lifecycle_action(action.as_str(), Path::new(&data_dir))
319 })
320 .await
321 .map_err(|err| err.to_string())
322 .and_then(|res| res.map_err(|err| err.to_string()));
323
324 match result {
325 Ok(message) => (
326 StatusCode::OK,
327 Json(AdminLifecycleResponse {
328 status: "OK",
329 message,
330 }),
331 )
332 .into_response(),
333 Err(err) => (
334 StatusCode::BAD_REQUEST,
335 Json(AdminLifecycleResponse {
336 status: "Error",
337 message: err,
338 }),
339 )
340 .into_response(),
341 }
342}
343
344fn parse_backup_handle(id: &str) -> crate::error::Result<BackupHandle> {
345 let id = Uuid::parse_str(id)
346 .map_err(|_| crate::error::ServerError::BadRequest("invalid backup handle".into()))?;
347 Ok(BackupHandle { id })
348}
349
350fn parse_restore_handle(id: &str) -> crate::error::Result<RestoreHandle> {
351 let id = Uuid::parse_str(id)
352 .map_err(|_| crate::error::ServerError::BadRequest("invalid restore handle".into()))?;
353 Ok(RestoreHandle { id })
354}
355
356fn backup_response(
357 state: &ServerState,
358 handle: &BackupHandle,
359) -> crate::error::Result<AdminBackupResponse> {
360 let location = state.backup_coordinator.location(handle)?;
361 let status = state.backup_coordinator.status(handle)?;
362 Ok(AdminBackupResponse {
363 handle: handle.id.to_string(),
364 location: location.display().to_string(),
365 state: status,
366 })
367}
368
369fn restore_response(
370 state: &ServerState,
371 handle: &RestoreHandle,
372) -> crate::error::Result<AdminRestoreResponse> {
373 let status = state.restore_coordinator.status(handle)?;
374 let metadata = state.restore_coordinator.metadata(handle)?;
375 Ok(AdminRestoreResponse {
376 handle: handle.id.to_string(),
377 state: status,
378 metadata,
379 })
380}
381
382fn capabilities_for_auth(
383 auth: &crate::auth::AuthMiddleware,
384) -> (&'static str, Vec<&'static str>, Vec<&'static str>) {
385 match auth.mode() {
386 AuthMode::None => ("full", Vec::new(), unsupported_actions()),
387 AuthMode::Dev { .. } => ("restricted", all_actions(), unsupported_actions()),
388 }
389}
390
391fn unsupported_actions() -> Vec<&'static str> {
392 vec!["compaction"]
393}
394
395fn all_actions() -> Vec<&'static str> {
396 vec![
397 "read", "create", "update", "delete", "archive", "restore", "backup", "export", "join",
398 "leave",
399 ]
400}
401
402fn cluster_operation_response(
403 state: &Arc<ServerState>,
404 ctx: &RequestContext,
405 action: &'static str,
406) -> Response {
407 let cluster = match action {
408 "join" => state.cluster_join(),
409 "leave" => state.cluster_leave(),
410 _ => unreachable!("cluster membership action is fixed by route"),
411 };
412 let cluster = match cluster {
413 Ok(snapshot) => snapshot,
414 Err(err) => return error_response(err, ctx),
415 };
416 state.metrics.record_cluster_status(&cluster);
417 Json(AdminClusterOperationResponse { action, cluster }).into_response()
418}
419
420fn perform_lifecycle_action(action: &str, data_dir: &Path) -> Result<String, String> {
421 if !data_dir.exists() {
422 return Err(format!(
423 "Data directory does not exist: {}",
424 data_dir.display()
425 ));
426 }
427 if !data_dir.is_dir() {
428 return Err(format!(
429 "Data directory is not a directory: {}",
430 data_dir.display()
431 ));
432 }
433
434 let lifecycle_root = data_dir.join(".lifecycle");
435 std::fs::create_dir_all(&lifecycle_root).map_err(|err| err.to_string())?;
436
437 match action {
438 "archive" => {
439 let dest = lifecycle_root.join("archive").join(timestamp_dir());
440 copy_data_dir(data_dir, &dest)?;
441 write_latest_marker(&lifecycle_root.join("archive"), &dest)?;
442 Ok(format!("Archived data to {}", dest.display()))
443 }
444 "export" => {
445 let dest = lifecycle_root.join("export").join(timestamp_dir());
446 copy_data_dir(data_dir, &dest)?;
447 write_latest_marker(&lifecycle_root.join("export"), &dest)?;
448 Ok(format!("Exported data to {}", dest.display()))
449 }
450 _ => Err("Unknown lifecycle action.".to_string()),
451 }
452}
453
454fn perform_export(state: &ServerState) -> crate::error::Result<String> {
455 match state.store.as_ref() {
456 AnyKV::Lsm(kv) => {
457 let _ = kv.checkpoint()?;
458 }
459 _ => {
460 return Err(crate::error::ServerError::BadRequest(
461 "checkpoint unsupported for current storage engine".to_string(),
462 ));
463 }
464 }
465 let data_dir = state.config.data_dir.as_path();
466 let lifecycle_root = data_dir.join(".lifecycle");
467 std::fs::create_dir_all(&lifecycle_root)?;
468 let dest = lifecycle_root.join("export").join(timestamp_dir());
469 std::fs::create_dir_all(&dest)?;
470 export_snapshot(data_dir, &dest)?;
471 write_latest_marker(&lifecycle_root.join("export"), &dest)
472 .map_err(crate::error::ServerError::Internal)?;
473 Ok(dest.display().to_string())
474}
475
476fn timestamp_dir() -> String {
477 let seconds = std::time::SystemTime::now()
478 .duration_since(std::time::UNIX_EPOCH)
479 .unwrap_or_default()
480 .as_secs();
481 format!("ts-{seconds}")
482}
483
484fn copy_data_dir(src: &Path, dest: &Path) -> Result<(), String> {
485 std::fs::create_dir_all(dest).map_err(|err| err.to_string())?;
486 copy_dir_filtered(src, dest)
487}
488
489fn copy_dir_filtered(src: &Path, dest: &Path) -> Result<(), String> {
490 for entry in std::fs::read_dir(src).map_err(|err| err.to_string())? {
491 let entry = entry.map_err(|err| err.to_string())?;
492 let file_type = entry.file_type().map_err(|err| err.to_string())?;
493 let name = entry.file_name();
494 if name == ".lifecycle" {
495 continue;
496 }
497 let dest_path = dest.join(name);
498 if file_type.is_dir() {
499 copy_data_dir(&entry.path(), &dest_path)?;
500 } else {
501 std::fs::copy(entry.path(), &dest_path).map_err(|err| err.to_string())?;
502 }
503 }
504 Ok(())
505}
506
507fn write_latest_marker(root: &Path, dest: &Path) -> Result<(), String> {
508 let marker = root.join("latest");
509 std::fs::create_dir_all(root).map_err(|err| err.to_string())?;
510 std::fs::write(&marker, dest.to_string_lossy().as_bytes()).map_err(|err| err.to_string())?;
511 Ok(())
512}