Skip to main content

alopex_server/http/
admin_api.rs

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}