Skip to main content

rhiza_node/
admin.rs

1use std::{
2    collections::HashMap,
3    fmt, fs,
4    io::Write,
5    path::{Path, PathBuf},
6    sync::atomic::{AtomicBool, AtomicUsize, Ordering},
7    sync::Arc,
8};
9
10use axum::{
11    extract::{rejection::JsonRejection, Extension, Request, State},
12    http::StatusCode,
13    middleware::{self, Next},
14    response::{IntoResponse, Response},
15    routing::{get, post},
16    Json, Router,
17};
18use rhiza_core::{ConfigChange, ExecutionProfile, LogAnchor, LogEntry, LogHash, StoredCommand};
19use rhiza_log::{IndexRange, LogStore};
20use rhiza_quepaxa::{Membership, RecorderFileStore};
21use serde::{Deserialize, Serialize};
22use serde_json::Value;
23
24use crate::{
25    client_authenticated, install_successor_recorder, valid_auth_token, ConfigError, NodeError,
26    NodeRuntime, NodeStatus, StopInformation,
27};
28use crate::{CheckpointCoordinator, DurabilityError};
29
30pub const ADMIN_STATUS_PATH: &str = "/v1/admin/membership/status";
31pub const ADMIN_STOP_PATH: &str = "/v1/admin/membership/stop";
32pub const ADMIN_INSTALL_SUCCESSOR_PATH: &str = "/v1/admin/membership/install-successor";
33pub const ADMIN_ACTIVATE_PATH: &str = "/v1/admin/membership/activate";
34pub const ADMIN_COMPACT_PATH: &str = "/v1/admin/checkpoint/compact";
35
36#[derive(Clone, Eq, PartialEq)]
37pub struct AdminConfig {
38    token: String,
39}
40
41impl AdminConfig {
42    pub fn new(token: impl Into<String>) -> Result<Self, ConfigError> {
43        let token = token.into();
44        if !valid_auth_token(&token) {
45            return Err(ConfigError::EmptyAdminToken);
46        }
47        Ok(Self { token })
48    }
49}
50
51impl fmt::Debug for AdminConfig {
52    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
53        formatter
54            .debug_struct("AdminConfig")
55            .field("token", &"[redacted]")
56            .finish()
57    }
58}
59
60#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
61pub struct AdminStatusResponse {
62    pub cluster_id: String,
63    pub execution_profile: ExecutionProfile,
64    pub epoch: u64,
65    pub node: NodeStatus,
66    pub members: Vec<String>,
67    pub recovery_generation: u64,
68    pub qlog_root: LogAnchor,
69    pub checkpoint_root: Option<LogAnchor>,
70    pub stopped_transition: Option<AdminStoppedTransition>,
71}
72
73#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
74#[serde(deny_unknown_fields)]
75pub struct AdminStopRequest {
76    pub operation_id: String,
77    pub expected_config_id: u64,
78    pub successor: AdminSuccessorBundle,
79}
80
81#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
82pub struct AdminStopResponse {
83    pub operation_id: String,
84    pub stop: StopInformation,
85    pub successor: AdminSuccessorBundle,
86}
87
88#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
89pub struct AdminStoppedTransition {
90    pub stop: StopInformation,
91    pub successor: AdminSuccessorBundle,
92}
93
94#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
95#[serde(deny_unknown_fields)]
96pub struct AdminSuccessorBundle {
97    pub config_id: u64,
98    pub members: Vec<String>,
99    pub digest: LogHash,
100}
101
102#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
103#[serde(deny_unknown_fields)]
104pub struct AdminInstallSuccessorRequest {
105    pub operation_id: String,
106    pub expected_config_id: u64,
107    pub expected_stopped_anchor: LogAnchor,
108    pub old_members: Vec<String>,
109    pub stop: StopInformation,
110    pub successor: AdminSuccessorBundle,
111}
112
113#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
114pub struct AdminInstallSuccessorResponse {
115    pub operation_id: String,
116    pub config_id: u64,
117    pub digest: LogHash,
118    pub activated: bool,
119}
120
121#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
122#[serde(deny_unknown_fields)]
123pub struct AdminActivateRequest {
124    pub operation_id: String,
125    pub expected_config_id: u64,
126}
127
128#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
129pub struct AdminActivateResponse {
130    pub operation_id: String,
131    pub entry: LogEntry,
132}
133
134#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
135#[serde(deny_unknown_fields)]
136pub struct AdminCompactRequest {
137    pub operation_id: String,
138    pub expected_config_id: u64,
139    pub expected_recovery_generation: u64,
140    pub expected_root: LogAnchor,
141}
142
143#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
144pub struct AdminCompactResponse {
145    pub operation_id: String,
146    pub anchor: rhiza_core::RecoveryAnchor,
147}
148
149#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
150#[serde(rename_all = "snake_case")]
151pub enum AdminErrorCode {
152    Unauthorized,
153    InvalidRequest,
154    OperationConflict,
155    PreconditionFailed,
156    Unavailable,
157    Internal,
158}
159
160#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
161pub struct AdminErrorResponse {
162    pub code: AdminErrorCode,
163}
164
165#[derive(Clone)]
166struct AdminGateState {
167    token: String,
168    admission: Arc<tokio::sync::Semaphore>,
169    tasks: AdminTaskTracker,
170}
171
172#[derive(Clone)]
173pub struct AdminTaskTracker {
174    state: Arc<AdminTaskState>,
175}
176
177struct AdminTaskState {
178    accepting: AtomicBool,
179    active: AtomicUsize,
180    changed: tokio::sync::watch::Sender<()>,
181}
182
183impl AdminTaskTracker {
184    fn new() -> Self {
185        let (changed, _) = tokio::sync::watch::channel(());
186        Self {
187            state: Arc::new(AdminTaskState {
188                accepting: AtomicBool::new(true),
189                active: AtomicUsize::new(0),
190                changed,
191            }),
192        }
193    }
194
195    fn try_start(&self) -> Option<AdminTaskGuard> {
196        if !self.state.accepting.load(Ordering::Acquire) {
197            return None;
198        }
199        self.state.active.fetch_add(1, Ordering::AcqRel);
200        if self.state.accepting.load(Ordering::Acquire) {
201            Some(AdminTaskGuard {
202                state: Arc::clone(&self.state),
203            })
204        } else {
205            finish_admin_task(&self.state);
206            None
207        }
208    }
209
210    pub fn stop_admission(&self) {
211        self.state.accepting.store(false, Ordering::Release);
212    }
213
214    pub async fn wait_for_idle(&self) {
215        let mut changed = self.state.changed.subscribe();
216        loop {
217            if self.state.active.load(Ordering::Acquire) == 0 {
218                return;
219            }
220            if changed.changed().await.is_err() {
221                return;
222            }
223        }
224    }
225}
226
227struct AdminTaskGuard {
228    state: Arc<AdminTaskState>,
229}
230
231impl Drop for AdminTaskGuard {
232    fn drop(&mut self) {
233        finish_admin_task(&self.state);
234    }
235}
236
237fn finish_admin_task(state: &AdminTaskState) {
238    if state.active.fetch_sub(1, Ordering::AcqRel) == 1 {
239        state.changed.send_replace(());
240    }
241}
242
243struct AdminPermit {
244    _admission: tokio::sync::OwnedSemaphorePermit,
245    _task: AdminTaskGuard,
246}
247
248#[derive(Clone)]
249struct AdminRouteState {
250    runtime: Arc<NodeRuntime>,
251    recorder: RecorderFileStore,
252    coordinator: Option<Arc<CheckpointCoordinator>>,
253    operations: Arc<tokio::sync::Mutex<Option<HashMap<String, OperationRecord>>>>,
254    ledger_path: PathBuf,
255}
256
257#[derive(Clone, Deserialize, Serialize)]
258struct OperationRecord {
259    fingerprint: Vec<u8>,
260    status: u16,
261    body: Value,
262}
263
264#[derive(Default, Deserialize, Serialize)]
265#[serde(deny_unknown_fields)]
266struct OperationLedger {
267    operations: HashMap<String, OperationRecord>,
268}
269
270pub fn node_router_with_admin(
271    runtime: Arc<NodeRuntime>,
272    recorder: RecorderFileStore,
273    admin: AdminConfig,
274) -> Result<Router, ConfigError> {
275    node_router_with_admin_and_tasks(runtime, recorder, admin).map(|(router, _)| router)
276}
277
278pub fn node_router_with_admin_and_tasks(
279    runtime: Arc<NodeRuntime>,
280    recorder: RecorderFileStore,
281    admin: AdminConfig,
282) -> Result<(Router, AdminTaskTracker), ConfigError> {
283    validate_admin_token(&runtime, &admin)?;
284    let (admin_router, tasks) = admin_router(runtime.clone(), recorder.clone(), None, admin);
285    Ok((
286        crate::node_router(runtime, recorder).merge(admin_router),
287        tasks,
288    ))
289}
290
291pub fn node_router_with_checkpoint_and_admin(
292    runtime: Arc<NodeRuntime>,
293    recorder: RecorderFileStore,
294    coordinator: Arc<CheckpointCoordinator>,
295    admin: AdminConfig,
296) -> Result<Router, ConfigError> {
297    node_router_with_checkpoint_and_admin_tasks(runtime, recorder, coordinator, admin)
298        .map(|(router, _)| router)
299}
300
301pub fn node_router_with_checkpoint_and_admin_tasks(
302    runtime: Arc<NodeRuntime>,
303    recorder: RecorderFileStore,
304    coordinator: Arc<CheckpointCoordinator>,
305    admin: AdminConfig,
306) -> Result<(Router, AdminTaskTracker), ConfigError> {
307    validate_admin_token(&runtime, &admin)?;
308    let (admin_router, tasks) = admin_router(
309        runtime.clone(),
310        recorder.clone(),
311        Some(coordinator.clone()),
312        admin,
313    );
314    Ok((
315        crate::node_router_with_checkpoint(runtime, recorder, coordinator).merge(admin_router),
316        tasks,
317    ))
318}
319
320fn validate_admin_token(runtime: &NodeRuntime, admin: &AdminConfig) -> Result<(), ConfigError> {
321    if runtime.config().client_token() == admin.token
322        || runtime
323            .config()
324            .peers()
325            .iter()
326            .any(|peer| peer.token() == admin.token)
327    {
328        return Err(ConfigError::AdminTokenConflictsWithRuntime);
329    }
330    Ok(())
331}
332
333fn admin_router(
334    runtime: Arc<NodeRuntime>,
335    recorder: RecorderFileStore,
336    coordinator: Option<Arc<CheckpointCoordinator>>,
337    admin: AdminConfig,
338) -> (Router, AdminTaskTracker) {
339    let ledger_path = runtime.config().data_dir().join("admin-operations-v1.json");
340    let operations = load_operations(&ledger_path)
341        .map(Some)
342        .unwrap_or_else(|error| {
343            eprintln!("admin operation ledger is unavailable: {error}");
344            None
345        });
346    let state = AdminRouteState {
347        runtime,
348        recorder,
349        coordinator,
350        operations: Arc::new(tokio::sync::Mutex::new(operations)),
351        ledger_path,
352    };
353    let tasks = AdminTaskTracker::new();
354    let router = Router::new()
355        .route(ADMIN_STATUS_PATH, get(handle_status))
356        .route(ADMIN_STOP_PATH, post(handle_stop))
357        .route(ADMIN_INSTALL_SUCCESSOR_PATH, post(handle_install_successor))
358        .route(ADMIN_ACTIVATE_PATH, post(handle_activate))
359        .route(ADMIN_COMPACT_PATH, post(handle_compact))
360        .route_layer(middleware::from_fn_with_state(
361            AdminGateState {
362                token: admin.token,
363                admission: Arc::new(tokio::sync::Semaphore::new(1)),
364                tasks: tasks.clone(),
365            },
366            admin_gate,
367        ))
368        .with_state(state);
369    (router, tasks)
370}
371
372async fn admin_gate(
373    State(state): State<AdminGateState>,
374    mut request: Request,
375    next: Next,
376) -> Response {
377    if !client_authenticated(request.headers(), &state.token) {
378        return admin_error(StatusCode::UNAUTHORIZED, AdminErrorCode::Unauthorized);
379    }
380    let task = match state.tasks.try_start() {
381        Some(task) => task,
382        None => return admin_error(StatusCode::SERVICE_UNAVAILABLE, AdminErrorCode::Unavailable),
383    };
384    let permit = match state.admission.try_acquire_owned() {
385        Ok(admission) => Arc::new(AdminPermit {
386            _admission: admission,
387            _task: task,
388        }),
389        Err(_) => return admin_error(StatusCode::TOO_MANY_REQUESTS, AdminErrorCode::Unavailable),
390    };
391    request.extensions_mut().insert(permit);
392    next.run(request).await
393}
394
395async fn handle_status(
396    State(state): State<AdminRouteState>,
397    Extension(_permit): Extension<Arc<AdminPermit>>,
398) -> Response {
399    let _operations = state.operations.lock().await;
400    match status_response(&state).await {
401        Ok(response) => Json(response).into_response(),
402        Err(error) => node_admin_error(error),
403    }
404}
405
406async fn handle_stop(
407    State(state): State<AdminRouteState>,
408    Extension(permit): Extension<Arc<AdminPermit>>,
409    payload: Result<Json<AdminStopRequest>, JsonRejection>,
410) -> Response {
411    let request = match payload {
412        Ok(Json(request)) => request,
413        Err(_) => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
414    };
415    let runtime = state.runtime.clone();
416    let owned = request.clone();
417    run_async_operation(&state, permit, "stop", &request, async move {
418        tokio::task::spawn_blocking(move || {
419            let successor = validate_successor(
420                &owned.successor,
421                owned.expected_config_id,
422                runtime.config().cluster_id(),
423            )?;
424            runtime
425                .stop_current_configuration_for_successor(&successor)
426                .map(|stop| AdminStopResponse {
427                    operation_id: owned.operation_id,
428                    stop,
429                    successor: owned.successor,
430                })
431                .map_err(OperationError::Node)
432        })
433        .await
434        .unwrap_or(Err(OperationError::Unavailable))
435    })
436    .await
437}
438
439async fn handle_install_successor(
440    State(state): State<AdminRouteState>,
441    Extension(permit): Extension<Arc<AdminPermit>>,
442    payload: Result<Json<AdminInstallSuccessorRequest>, JsonRejection>,
443) -> Response {
444    let request = match payload {
445        Ok(Json(request)) => request,
446        Err(_) => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
447    };
448    let runtime = state.runtime.clone();
449    let recorder = state.recorder.clone();
450    let owned = request.clone();
451    run_async_operation(&state, permit, "install_successor", &request, async move {
452        tokio::task::spawn_blocking(move || {
453            install_successor(&runtime, &recorder, &owned).map_err(OperationError::Node)
454        })
455        .await
456        .unwrap_or(Err(OperationError::Unavailable))
457    })
458    .await
459}
460
461async fn handle_activate(
462    State(state): State<AdminRouteState>,
463    Extension(permit): Extension<Arc<AdminPermit>>,
464    payload: Result<Json<AdminActivateRequest>, JsonRejection>,
465) -> Response {
466    let request = match payload {
467        Ok(Json(request)) => request,
468        Err(_) => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
469    };
470    let runtime = state.runtime.clone();
471    let owned = request.clone();
472    run_async_operation(&state, permit, "activate", &request, async move {
473        tokio::task::spawn_blocking(move || {
474            runtime
475                .activate_successor_if(owned.expected_config_id)
476                .map(|entry| AdminActivateResponse {
477                    operation_id: owned.operation_id,
478                    entry,
479                })
480                .map_err(OperationError::Node)
481        })
482        .await
483        .unwrap_or(Err(OperationError::Unavailable))
484    })
485    .await
486}
487
488async fn handle_compact(
489    State(state): State<AdminRouteState>,
490    Extension(permit): Extension<Arc<AdminPermit>>,
491    payload: Result<Json<AdminCompactRequest>, JsonRejection>,
492) -> Response {
493    let request = match payload {
494        Ok(Json(request)) => request,
495        Err(_) => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
496    };
497    let runtime = state.runtime.clone();
498    let coordinator = state.coordinator.clone();
499    let owned = request.clone();
500    run_async_operation(&state, permit, "compact", &request, async move {
501        match coordinator {
502            Some(coordinator) => coordinator
503                .checkpoint_compact_fenced(
504                    &runtime,
505                    owned.expected_config_id,
506                    owned.expected_recovery_generation,
507                    owned.expected_root,
508                )
509                .await
510                .map(|anchor| AdminCompactResponse {
511                    operation_id: owned.operation_id,
512                    anchor,
513                })
514                .map_err(OperationError::Durability),
515            None => Err(OperationError::Unavailable),
516        }
517    })
518    .await
519}
520
521async fn run_async_operation<T, R, F>(
522    state: &AdminRouteState,
523    permit: Arc<AdminPermit>,
524    kind: &str,
525    request: &T,
526    operation: F,
527) -> Response
528where
529    T: Serialize,
530    R: Serialize + Send + 'static,
531    F: std::future::Future<Output = Result<R, OperationError>> + Send + 'static,
532{
533    let operation_id = match serde_json::to_value(request)
534        .ok()
535        .and_then(|value| value.get("operation_id")?.as_str().map(str::to_owned))
536    {
537        Some(operation_id) => operation_id,
538        None => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
539    };
540    let fingerprint = match operation_fingerprint(kind, request) {
541        Ok(fingerprint) => fingerprint,
542        Err(()) => return admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
543    };
544    {
545        let operations = state.operations.lock().await;
546        let Some(operations) = operations.as_ref() else {
547            return admin_error(StatusCode::SERVICE_UNAVAILABLE, AdminErrorCode::Unavailable);
548        };
549        if let Some(response) = replay(operations, &operation_id, &fingerprint) {
550            return response;
551        }
552    }
553    if let Some(response) = validate_operation_id(&operation_id) {
554        return response;
555    }
556    let detached_state = state.clone();
557    let detached_operation_id = operation_id.clone();
558    let detached_fingerprint = fingerprint.clone();
559    let (completed, mut completion) = tokio::sync::watch::channel(false);
560    tokio::spawn(async move {
561        let result = operation.await;
562        let mut operations = detached_state.operations.lock().await;
563        if let Some(records) = operations.as_mut() {
564            let _ = store_result(records, detached_operation_id, detached_fingerprint, result);
565            if let Err(error) = persist_operations(&detached_state.ledger_path, records) {
566                eprintln!("admin operation ledger persistence failed: {error}");
567                *operations = None;
568            }
569        }
570        drop(operations);
571        drop(detached_state);
572        drop(permit);
573        completed.send_replace(true);
574    });
575    let waited = tokio::time::timeout(std::time::Duration::from_secs(10), async {
576        while !*completion.borrow() {
577            if completion.changed().await.is_err() {
578                break;
579            }
580        }
581    })
582    .await;
583    if waited.is_err() {
584        return admin_error(StatusCode::SERVICE_UNAVAILABLE, AdminErrorCode::Unavailable);
585    }
586    let operations = state.operations.lock().await;
587    let Some(operations) = operations.as_ref() else {
588        return admin_error(StatusCode::SERVICE_UNAVAILABLE, AdminErrorCode::Unavailable);
589    };
590    replay(operations, &operation_id, &fingerprint)
591        .unwrap_or_else(|| admin_error(StatusCode::INTERNAL_SERVER_ERROR, AdminErrorCode::Internal))
592}
593
594fn operation_fingerprint(kind: &str, request: &impl Serialize) -> Result<Vec<u8>, ()> {
595    serde_json::to_vec(&(kind, request)).map_err(|_| ())
596}
597
598fn replay(
599    operations: &HashMap<String, OperationRecord>,
600    operation_id: &str,
601    fingerprint: &[u8],
602) -> Option<Response> {
603    operations.get(operation_id).map(|record| {
604        if record.fingerprint == fingerprint {
605            (
606                StatusCode::from_u16(record.status).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR),
607                Json(record.body.clone()),
608            )
609                .into_response()
610        } else {
611            admin_error(StatusCode::CONFLICT, AdminErrorCode::OperationConflict)
612        }
613    })
614}
615
616fn validate_operation_id(operation_id: &str) -> Option<Response> {
617    (operation_id.trim().is_empty() || operation_id.len() > 256)
618        .then(|| admin_error(StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest))
619}
620
621fn store_result<R: Serialize>(
622    operations: &mut HashMap<String, OperationRecord>,
623    operation_id: String,
624    fingerprint: Vec<u8>,
625    result: Result<R, OperationError>,
626) -> Response {
627    let (status, body) = match result {
628        Ok(response) => match serde_json::to_value(response) {
629            Ok(body) => (StatusCode::OK, body),
630            Err(_) => (
631                StatusCode::INTERNAL_SERVER_ERROR,
632                error_value(AdminErrorCode::Internal),
633            ),
634        },
635        Err(error) => operation_error_value(error),
636    };
637    operations.insert(
638        operation_id,
639        OperationRecord {
640            fingerprint,
641            status: status.as_u16(),
642            body: body.clone(),
643        },
644    );
645    (status, Json(body)).into_response()
646}
647
648async fn status_response(state: &AdminRouteState) -> Result<AdminStatusResponse, NodeError> {
649    let checkpoint_root = match state.coordinator.as_ref() {
650        Some(coordinator) => {
651            let tip = coordinator
652                .refresh_durable_tip()
653                .await
654                .map_err(|error| NodeError::Unavailable(error.to_string()))?;
655            Some(LogAnchor::new(tip.index(), tip.hash()))
656        }
657        None => None,
658    };
659    let _commit = state.runtime.lock_commit()?;
660    let node = state.runtime.status()?;
661    let qlog_root = state.runtime.log_root_unlocked()?;
662    let stopped_transition = stopped_transition(&state.runtime)?;
663    Ok(AdminStatusResponse {
664        cluster_id: state.runtime.config.cluster_id().to_owned(),
665        execution_profile: state.runtime.config.execution_profile(),
666        epoch: state.runtime.config.epoch(),
667        node,
668        members: state.runtime.config.membership().members().to_vec(),
669        recovery_generation: state.runtime.config.recovery_generation(),
670        qlog_root,
671        checkpoint_root,
672        stopped_transition,
673    })
674}
675
676fn stopped_transition(runtime: &NodeRuntime) -> Result<Option<AdminStoppedTransition>, NodeError> {
677    let configuration = runtime.configuration_state()?;
678    let Some(anchor) = configuration.stop().copied() else {
679        return Ok(None);
680    };
681    if configuration.config_id() != runtime.consensus.config_id() {
682        return Ok(None);
683    }
684    let entry = runtime.recover_stop_entry(anchor)?;
685    let successor = successor_from_entry(&entry)?;
686    let proof = runtime
687        .consensus
688        .inspect_decision_proof_at(entry.index)
689        .map_err(|error| NodeError::Unavailable(error.to_string()))?
690        .ok_or_else(|| NodeError::Unavailable("durable Stop proof is unavailable".into()))?;
691    Ok(Some(AdminStoppedTransition {
692        stop: StopInformation { entry, proof },
693        successor,
694    }))
695}
696
697fn successor_from_entry(entry: &LogEntry) -> Result<AdminSuccessorBundle, NodeError> {
698    let command = StoredCommand::new(entry.entry_type, entry.payload.clone());
699    let change = ConfigChange::recognize(&command)
700        .map_err(|_| NodeError::PreconditionFailed("Stop command is not successor-bound".into()))?;
701    let ConfigChange::BoundStop { successor } = change else {
702        return Err(NodeError::PreconditionFailed(
703            "Stop command is not successor-bound".into(),
704        ));
705    };
706    Ok(AdminSuccessorBundle {
707        config_id: successor.config_id(),
708        members: successor.members().to_vec(),
709        digest: successor.digest(),
710    })
711}
712
713fn validate_successor(
714    bundle: &AdminSuccessorBundle,
715    predecessor_config_id: u64,
716    cluster_id: &str,
717) -> Result<Membership, OperationError> {
718    let membership = Membership::from_voters(bundle.members.clone()).map_err(|_| {
719        OperationError::Node(NodeError::InvalidRequest(
720            "successor membership is invalid".into(),
721        ))
722    })?;
723    if predecessor_config_id.checked_add(1) != Some(bundle.config_id)
724        || membership.digest() != bundle.digest
725        || cluster_id.is_empty()
726    {
727        return Err(OperationError::Node(NodeError::PreconditionFailed(
728            "successor descriptor does not match the active configuration".into(),
729        )));
730    }
731    Ok(membership)
732}
733
734fn load_operations(path: &Path) -> Result<HashMap<String, OperationRecord>, std::io::Error> {
735    let bytes = match fs::read(path) {
736        Ok(bytes) => bytes,
737        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
738        Err(error) => return Err(error),
739    };
740    let ledger: OperationLedger = serde_json::from_slice(&bytes)
741        .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
742    Ok(ledger.operations)
743}
744
745fn persist_operations(
746    path: &Path,
747    operations: &HashMap<String, OperationRecord>,
748) -> Result<(), std::io::Error> {
749    let parent = path.parent().ok_or_else(|| {
750        std::io::Error::new(std::io::ErrorKind::InvalidInput, "ledger has no parent")
751    })?;
752    fs::create_dir_all(parent)?;
753    let temporary = path.with_extension(format!("tmp-{}", std::process::id()));
754    let bytes = serde_json::to_vec(&OperationLedger {
755        operations: operations.clone(),
756    })
757    .map_err(std::io::Error::other)?;
758    let mut file = fs::File::create(&temporary)?;
759    file.write_all(&bytes)?;
760    file.sync_all()?;
761    fs::rename(&temporary, path)?;
762    fs::File::open(parent)?.sync_all()
763}
764
765fn install_successor(
766    runtime: &NodeRuntime,
767    recorder: &RecorderFileStore,
768    request: &AdminInstallSuccessorRequest,
769) -> Result<AdminInstallSuccessorResponse, NodeError> {
770    let _commit = runtime.lock_commit()?;
771    runtime.ensure_ready()?;
772    let state = runtime.configuration_state()?;
773    if state.is_active()
774        || state.config_id() != request.expected_config_id
775        || state.stop().copied() != Some(request.expected_stopped_anchor)
776    {
777        return Err(NodeError::PreconditionFailed(
778            "stopped configuration anchor does not match".into(),
779        ));
780    }
781    let old_membership = Membership::from_voters(request.old_members.clone())
782        .map_err(|_| NodeError::InvalidRequest("old membership is invalid".into()))?;
783    if old_membership.digest() != state.digest()
784        || request.stop.entry.cluster_id != runtime.config.cluster_id()
785        || request.stop.entry.epoch != runtime.config.epoch()
786        || request.stop.entry.config_id != request.expected_config_id
787        || request.stop.entry.index != request.expected_stopped_anchor.index()
788        || request.stop.entry.hash != request.expected_stopped_anchor.hash()
789    {
790        return Err(NodeError::PreconditionFailed(
791            "old decision material does not match the stopped runtime".into(),
792        ));
793    }
794    let entries = runtime
795        .log_store
796        .read_range(
797            IndexRange::new(request.stop.entry.index, request.stop.entry.index)
798                .map_err(|error| NodeError::Storage(error.to_string()))?,
799        )
800        .map_err(|error| NodeError::Storage(error.to_string()))?;
801    if entries.as_slice() != [request.stop.entry.clone()] {
802        return Err(NodeError::PreconditionFailed(
803            "old stop entry is not the exact local qlog entry".into(),
804        ));
805    }
806    let successor = Membership::from_voters(request.successor.members.clone())
807        .map_err(|_| NodeError::InvalidRequest("successor membership is invalid".into()))?;
808    if successor.digest() != request.successor.digest
809        || request.expected_config_id.checked_add(1) != Some(request.successor.config_id)
810        || successor_from_entry(&request.stop.entry)? != request.successor
811    {
812        return Err(NodeError::PreconditionFailed(
813            "successor bundle or digest does not match".into(),
814        ));
815    }
816    let installed = install_successor_recorder(
817        recorder,
818        request.successor.config_id,
819        successor,
820        &request.stop,
821    )?;
822    Ok(AdminInstallSuccessorResponse {
823        operation_id: request.operation_id.clone(),
824        config_id: installed.config_id(),
825        digest: installed.config_digest(),
826        activated: installed.is_activated(),
827    })
828}
829
830enum OperationError {
831    Node(NodeError),
832    Durability(DurabilityError),
833    Unavailable,
834}
835
836fn operation_error_value(error: OperationError) -> (StatusCode, Value) {
837    match error {
838        OperationError::Node(error) => {
839            eprintln!("admin operation failed: {error}");
840            let (status, code) = node_admin_status(&error);
841            (status, error_value(code))
842        }
843        OperationError::Durability(DurabilityError::PreconditionFailed) => (
844            StatusCode::CONFLICT,
845            error_value(AdminErrorCode::PreconditionFailed),
846        ),
847        OperationError::Durability(error) => {
848            eprintln!("admin durability operation failed: {error}");
849            (
850                StatusCode::SERVICE_UNAVAILABLE,
851                error_value(AdminErrorCode::Unavailable),
852            )
853        }
854        OperationError::Unavailable => (
855            StatusCode::SERVICE_UNAVAILABLE,
856            error_value(AdminErrorCode::Unavailable),
857        ),
858    }
859}
860
861fn node_admin_error(error: NodeError) -> Response {
862    let (status, code) = node_admin_status(&error);
863    admin_error(status, code)
864}
865
866fn node_admin_status(error: &NodeError) -> (StatusCode, AdminErrorCode) {
867    match error {
868        NodeError::InvalidRequest(_) => (StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest),
869        #[cfg(feature = "sql")]
870        NodeError::InvalidSqlStatement { .. } => {
871            (StatusCode::BAD_REQUEST, AdminErrorCode::InvalidRequest)
872        }
873        NodeError::PreconditionFailed(_) | NodeError::ConfigurationTransition { .. } => {
874            (StatusCode::CONFLICT, AdminErrorCode::PreconditionFailed)
875        }
876        #[cfg(feature = "sql")]
877        NodeError::RequestConflict(_) => (StatusCode::CONFLICT, AdminErrorCode::PreconditionFailed),
878        NodeError::Unavailable(_)
879        | NodeError::ResourceExhausted(_)
880        | NodeError::Contention(_)
881        | NodeError::WinnerLimitExceeded => {
882            (StatusCode::SERVICE_UNAVAILABLE, AdminErrorCode::Unavailable)
883        }
884        NodeError::UnsupportedAckMode(_)
885        | NodeError::ExecutionProfileMismatch { .. }
886        | NodeError::DataRootLocked(_)
887        | NodeError::SnapshotRequired(_)
888        | NodeError::Storage(_)
889        | NodeError::Reconciliation(_)
890        | NodeError::Invariant(_)
891        | NodeError::Fatal(_) => (StatusCode::INTERNAL_SERVER_ERROR, AdminErrorCode::Internal),
892    }
893}
894
895fn admin_error(status: StatusCode, code: AdminErrorCode) -> Response {
896    (status, Json(AdminErrorResponse { code })).into_response()
897}
898
899fn error_value(code: AdminErrorCode) -> Value {
900    serde_json::to_value(AdminErrorResponse { code })
901        .unwrap_or_else(|_| serde_json::json!({"code": "internal"}))
902}
903
904#[cfg(test)]
905mod tests {
906    use std::sync::{
907        atomic::{AtomicUsize, Ordering},
908        Arc,
909    };
910
911    use super::{AdminTaskTracker, OperationLedger};
912
913    #[test]
914    fn operation_ledger_has_one_strict_canonical_shape() {
915        let canonical = serde_json::to_value(OperationLedger::default()).unwrap();
916        assert_eq!(canonical, serde_json::json!({"operations": {}}));
917        assert!(serde_json::from_value::<OperationLedger>(
918            serde_json::json!({"version": 1, "operations": {}})
919        )
920        .is_err());
921    }
922
923    #[tokio::test]
924    async fn shutdown_observes_a_late_admin_mutation_before_sampling_state() {
925        let tasks = AdminTaskTracker::new();
926        let guard = tasks.try_start().unwrap();
927        let committed_tip = Arc::new(AtomicUsize::new(1));
928        let late_tip = Arc::clone(&committed_tip);
929        let (release, wait) = tokio::sync::oneshot::channel();
930        let operation = tokio::spawn(async move {
931            let _ = wait.await;
932            late_tip.store(2, Ordering::Release);
933            drop(guard);
934        });
935
936        tasks.stop_admission();
937        assert!(tasks.try_start().is_none());
938        let mut idle = Box::pin(tasks.wait_for_idle());
939        assert!(
940            tokio::time::timeout(std::time::Duration::from_millis(10), &mut idle)
941                .await
942                .is_err()
943        );
944
945        let _ = release.send(());
946        idle.await;
947        operation.await.unwrap();
948        assert_eq!(committed_tip.load(Ordering::Acquire), 2);
949    }
950
951    #[tokio::test]
952    async fn shutdown_waits_until_every_admitted_admin_task_finishes() {
953        let tasks = AdminTaskTracker::new();
954        let first = tasks.try_start().unwrap();
955        let second = tasks.try_start().unwrap();
956        tasks.stop_admission();
957
958        let mut idle = Box::pin(tasks.wait_for_idle());
959        drop(first);
960        assert!(
961            tokio::time::timeout(std::time::Duration::from_millis(10), &mut idle)
962                .await
963                .is_err()
964        );
965
966        drop(second);
967        tokio::time::timeout(std::time::Duration::from_secs(1), idle)
968            .await
969            .expect("last task completion must wake the shutdown waiter");
970    }
971}