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}