1use crate::infrastructure::di::ServiceContainer;
3#[cfg(feature = "multi-tenant")]
4use crate::infrastructure::web::tenant_api::*;
5use crate::{
6 domain::repositories::TenantRepository,
7 infrastructure::{
8 cluster::{
9 ClusterManager, ClusterMember, GeoReplicationManager, GeoSyncRequest, VoteRequest,
10 },
11 replication::{WalReceiver, WalShipper},
12 security::{
13 auth::AuthManager,
14 middleware::{AuthState, RateLimitState, auth_middleware, rate_limit_middleware},
15 rate_limit::RateLimiter,
16 },
17 web::{audit_api::*, auth_api::*, config_api::*},
18 },
19 store::EventStore,
20};
21use axum::{
22 Json, Router,
23 extract::{Path, State},
24 middleware,
25 response::IntoResponse,
26 routing::{delete, get, post, put},
27};
28use std::sync::Arc;
29use tower_http::{
30 cors::{Any, CorsLayer},
31 trace::TraceLayer,
32};
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
36#[serde(rename_all = "lowercase")]
37pub enum NodeRole {
38 Leader,
39 Follower,
40}
41
42impl NodeRole {
43 pub fn from_env() -> Self {
49 if let Ok(role) = std::env::var("ALLSOURCE_ROLE") {
50 match role.to_lowercase().as_str() {
51 "follower" => return NodeRole::Follower,
52 "leader" => return NodeRole::Leader,
53 other => {
54 tracing::warn!(
55 "Unknown ALLSOURCE_ROLE value '{}', defaulting to leader",
56 other
57 );
58 return NodeRole::Leader;
59 }
60 }
61 }
62 if let Ok(read_only) = std::env::var("ALLSOURCE_READ_ONLY")
63 && (read_only == "true" || read_only == "1")
64 {
65 return NodeRole::Follower;
66 }
67 NodeRole::Leader
68 }
69
70 pub fn is_follower(self) -> bool {
71 self == NodeRole::Follower
72 }
73
74 fn to_u8(self) -> u8 {
75 match self {
76 NodeRole::Leader => 0,
77 NodeRole::Follower => 1,
78 }
79 }
80
81 fn from_u8(v: u8) -> Self {
82 match v {
83 1 => NodeRole::Follower,
84 _ => NodeRole::Leader,
85 }
86 }
87}
88
89impl std::fmt::Display for NodeRole {
90 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
91 match self {
92 NodeRole::Leader => write!(f, "leader"),
93 NodeRole::Follower => write!(f, "follower"),
94 }
95 }
96}
97
98#[derive(Clone)]
103pub struct AtomicNodeRole(Arc<std::sync::atomic::AtomicU8>);
104
105impl AtomicNodeRole {
106 pub fn new(role: NodeRole) -> Self {
107 Self(Arc::new(std::sync::atomic::AtomicU8::new(role.to_u8())))
108 }
109
110 pub fn load(&self) -> NodeRole {
111 NodeRole::from_u8(self.0.load(std::sync::atomic::Ordering::Relaxed))
112 }
113
114 pub fn store(&self, role: NodeRole) {
115 self.0
116 .store(role.to_u8(), std::sync::atomic::Ordering::Relaxed);
117 }
118}
119
120#[derive(Clone)]
122pub struct AppState {
123 pub store: Arc<EventStore>,
124 pub auth_manager: Arc<AuthManager>,
125 pub tenant_repo: Arc<dyn TenantRepository>,
126 pub service_container: ServiceContainer,
128 pub role: AtomicNodeRole,
130 pub wal_shipper: Arc<tokio::sync::RwLock<Option<Arc<WalShipper>>>>,
133 pub wal_receiver: Option<Arc<WalReceiver>>,
135 pub replication_port: u16,
137 pub cluster_manager: Option<Arc<ClusterManager>>,
139 pub geo_replication: Option<Arc<GeoReplicationManager>>,
141}
142
143impl axum::extract::FromRef<AppState> for Arc<EventStore> {
146 fn from_ref(state: &AppState) -> Self {
147 state.store.clone()
148 }
149}
150
151pub async fn serve_v1(
152 store: Arc<EventStore>,
153 auth_manager: Arc<AuthManager>,
154 tenant_repo: Arc<dyn TenantRepository>,
155 rate_limiter: Arc<RateLimiter>,
156 service_container: ServiceContainer,
157 addr: &str,
158 role: NodeRole,
159 wal_shipper: Option<Arc<WalShipper>>,
160 wal_receiver: Option<Arc<WalReceiver>>,
161 replication_port: u16,
162 cluster_manager: Option<Arc<ClusterManager>>,
163 geo_replication: Option<Arc<GeoReplicationManager>>,
164) -> anyhow::Result<()> {
165 let app_state = AppState {
166 store,
167 auth_manager: auth_manager.clone(),
168 tenant_repo,
169 service_container,
170 role: AtomicNodeRole::new(role),
171 wal_shipper: Arc::new(tokio::sync::RwLock::new(wal_shipper)),
172 wal_receiver,
173 replication_port,
174 cluster_manager,
175 geo_replication,
176 };
177
178 let auth_state = AuthState {
179 auth_manager: auth_manager.clone(),
180 };
181
182 let rate_limit_state = RateLimitState { rate_limiter };
183
184 let app = Router::new()
185 .route("/health", get(health_v1))
187 .route("/metrics", get(super::api::prometheus_metrics))
188 .route("/api/v1/auth/register", post(register_handler))
190 .route("/api/v1/auth/login", post(login_handler))
191 .route("/api/v1/auth/me", get(me_handler))
192 .route("/api/v1/auth/api-keys", post(create_api_key_handler))
193 .route("/api/v1/auth/api-keys", get(list_api_keys_handler))
194 .route("/api/v1/auth/api-keys/{id}", delete(revoke_api_key_handler))
195 .route("/api/v1/auth/users", get(list_users_handler))
196 .route("/api/v1/auth/users/{id}", delete(delete_user_handler))
197 ;
199 #[cfg(feature = "multi-tenant")]
200 let app = app
201 .route("/api/v1/tenants", post(create_tenant_handler))
202 .route("/api/v1/tenants", get(list_tenants_handler))
203 .route("/api/v1/tenants/{id}", get(get_tenant_handler))
204 .route("/api/v1/tenants/{id}", put(update_tenant_handler))
205 .route(
206 "/api/v1/tenants/{id}/metadata",
207 axum::routing::patch(merge_tenant_metadata_handler),
211 )
212 .route("/api/v1/tenants/{id}/stats", get(get_tenant_stats_handler))
213 .route("/api/v1/tenants/{id}/quotas", put(update_quotas_handler))
214 .route(
215 "/api/v1/tenants/{id}/usage/increment",
216 post(increment_usage_handler),
217 )
218 .route(
219 "/api/v1/tenants/{id}/schema-enforcement",
220 put(update_schema_enforcement_handler),
221 )
222 .route(
223 "/api/v1/tenants/{id}/deactivate",
224 post(deactivate_tenant_handler),
225 )
226 .route(
227 "/api/v1/tenants/{id}/activate",
228 post(activate_tenant_handler),
229 )
230 .route("/api/v1/tenants/{id}", delete(delete_tenant_handler));
231 let app = app
232 .route("/api/v1/audit/events", post(log_audit_event))
234 .route("/api/v1/audit/events", get(query_audit_events))
235 .route("/api/v1/config", get(list_configs))
237 .route("/api/v1/config", post(set_config))
238 .route("/api/v1/config/{key}", get(get_config))
239 .route("/api/v1/config/{key}", put(update_config))
240 .route("/api/v1/config/{key}", delete(delete_config))
241 .route(
243 "/api/v1/demo/seed",
244 post(super::demo_api::demo_seed_handler),
245 )
246 .route("/api/v1/events", post(super::api::ingest_event_v1))
248 .route(
249 "/api/v1/events/batch",
250 post(super::api::ingest_events_batch_v1),
251 )
252 .route("/api/v1/events/query", get(super::api::query_events))
253 .route(
254 "/api/v1/events/{event_id}",
255 get(super::api::get_event_by_id),
256 )
257 .route("/api/v1/events/stream", get(super::api::events_websocket))
258 .route("/api/v1/entities", get(super::api::list_entities))
259 .route(
260 "/api/v1/entities/duplicates",
261 get(super::api::detect_duplicates),
262 )
263 .route(
264 "/api/v1/entities/{entity_id}/state",
265 get(super::api::get_entity_state),
266 )
267 .route(
268 "/api/v1/entities/{entity_id}/snapshot",
269 get(super::api::get_entity_snapshot),
270 )
271 .route("/api/v1/stats", get(super::api::get_stats))
272 .route("/api/v1/streams", get(super::api::list_streams))
274 .route("/api/v1/event-types", get(super::api::list_event_types))
275 .route(
277 "/api/v1/analytics/frequency",
278 get(super::api::analytics_frequency),
279 )
280 .route(
281 "/api/v1/analytics/summary",
282 get(super::api::analytics_summary),
283 )
284 .route(
285 "/api/v1/analytics/correlation",
286 get(super::api::analytics_correlation),
287 )
288 .route("/api/v1/snapshots", post(super::api::create_snapshot))
290 .route("/api/v1/snapshots", get(super::api::list_snapshots))
291 .route(
292 "/api/v1/snapshots/{entity_id}/latest",
293 get(super::api::get_latest_snapshot),
294 )
295 .route(
297 "/api/v1/compaction/trigger",
298 post(super::api::trigger_compaction),
299 )
300 .route(
301 "/api/v1/compaction/stats",
302 get(super::api::compaction_stats),
303 )
304 .route("/api/v1/schemas", post(super::api::register_schema))
306 .route("/api/v1/schemas", get(super::api::list_subjects))
307 .route("/api/v1/schemas/{subject}", get(super::api::get_schema))
308 .route(
309 "/api/v1/schemas/{subject}/versions",
310 get(super::api::list_schema_versions),
311 )
312 .route(
313 "/api/v1/schemas/validate",
314 post(super::api::validate_event_schema),
315 )
316 .route(
317 "/api/v1/schemas/{subject}/compatibility",
318 put(super::api::set_compatibility_mode),
319 )
320 .route("/api/v1/replay", post(super::api::start_replay))
322 .route("/api/v1/replay", get(super::api::list_replays))
323 .route(
324 "/api/v1/replay/{replay_id}",
325 get(super::api::get_replay_progress),
326 )
327 .route(
328 "/api/v1/replay/{replay_id}/cancel",
329 post(super::api::cancel_replay),
330 )
331 .route(
332 "/api/v1/replay/{replay_id}",
333 delete(super::api::delete_replay),
334 )
335 .route("/api/v1/pipelines", post(super::api::register_pipeline))
337 .route("/api/v1/pipelines", get(super::api::list_pipelines))
338 .route(
339 "/api/v1/pipelines/stats",
340 get(super::api::all_pipeline_stats),
341 )
342 .route(
343 "/api/v1/pipelines/{pipeline_id}",
344 get(super::api::get_pipeline),
345 )
346 .route(
347 "/api/v1/pipelines/{pipeline_id}",
348 delete(super::api::remove_pipeline),
349 )
350 .route(
351 "/api/v1/pipelines/{pipeline_id}/stats",
352 get(super::api::get_pipeline_stats),
353 )
354 .route(
355 "/api/v1/pipelines/{pipeline_id}/reset",
356 put(super::api::reset_pipeline),
357 )
358 .route("/api/v1/projections", get(super::api::list_projections))
360 .route(
361 "/api/v1/projections/{name}",
362 get(super::api::get_projection),
363 )
364 .route(
365 "/api/v1/projections/{name}",
366 delete(super::api::delete_projection),
367 )
368 .route(
369 "/api/v1/projections/{name}/state",
370 get(super::api::get_projection_state_summary),
371 )
372 .route(
373 "/api/v1/projections/{name}/reset",
374 post(super::api::reset_projection),
375 )
376 .route(
377 "/api/v1/projections/{name}/pause",
378 post(super::api::pause_projection),
379 )
380 .route(
381 "/api/v1/projections/{name}/start",
382 post(super::api::start_projection),
383 )
384 .route(
385 "/api/v1/projections/{name}/{entity_id}/state",
386 get(super::api::get_projection_state),
387 )
388 .route(
389 "/api/v1/projections/{name}/{entity_id}/state",
390 post(super::api::save_projection_state),
391 )
392 .route(
393 "/api/v1/projections/{name}/{entity_id}/state",
394 put(super::api::save_projection_state),
395 )
396 .route(
397 "/api/v1/projections/{name}/bulk",
398 post(super::api::bulk_get_projection_states),
399 )
400 .route(
401 "/api/v1/projections/{name}/bulk/save",
402 post(super::api::bulk_save_projection_states),
403 )
404 .route("/api/v1/webhooks", post(super::api::register_webhook))
406 .route("/api/v1/webhooks", get(super::api::list_webhooks))
407 .route(
408 "/api/v1/webhooks/{webhook_id}",
409 get(super::api::get_webhook),
410 )
411 .route(
412 "/api/v1/webhooks/{webhook_id}",
413 put(super::api::update_webhook),
414 )
415 .route(
416 "/api/v1/webhooks/{webhook_id}",
417 delete(super::api::delete_webhook),
418 )
419 .route(
420 "/api/v1/webhooks/{webhook_id}/deliveries",
421 get(super::api::list_webhook_deliveries),
422 )
423 .route("/api/v1/consumers", post(super::api::register_consumer))
425 .route(
426 "/api/v1/consumers/{consumer_id}",
427 get(super::api::get_consumer),
428 )
429 .route(
430 "/api/v1/consumers/{consumer_id}/events",
431 get(super::api::poll_consumer_events),
432 )
433 .route(
434 "/api/v1/consumers/{consumer_id}/ack",
435 post(super::api::ack_consumer),
436 )
437 .route("/api/v1/cluster/status", get(cluster_status_handler))
439 .route("/api/v1/cluster/members", get(cluster_list_members_handler))
440 .route("/api/v1/cluster/members", post(cluster_add_member_handler))
441 .route(
442 "/api/v1/cluster/members/{node_id}",
443 delete(cluster_remove_member_handler),
444 )
445 .route(
446 "/api/v1/cluster/members/{node_id}/heartbeat",
447 post(cluster_heartbeat_handler),
448 )
449 .route("/api/v1/cluster/vote", post(cluster_vote_handler))
450 .route("/api/v1/cluster/election", post(cluster_election_handler))
451 .route(
452 "/api/v1/cluster/partitions",
453 get(cluster_partitions_handler),
454 )
455 .route("/api/v1/graphql", post(super::api::graphql_query))
457 .route("/api/v1/geospatial/query", post(super::api::geo_query))
458 .route("/api/v1/geospatial/stats", get(super::api::geo_stats))
459 .route(
460 "/api/v1/exactly-once/stats",
461 get(super::api::exactly_once_stats),
462 )
463 .route(
464 "/api/v1/schema-evolution/history/{event_type}",
465 get(super::api::schema_evolution_history),
466 )
467 .route(
468 "/api/v1/schema-evolution/schema/{event_type}",
469 get(super::api::schema_evolution_schema),
470 )
471 .route(
472 "/api/v1/schema-evolution/stats",
473 get(super::api::schema_evolution_stats),
474 )
475 .route("/api/v1/geo/status", get(geo_status_handler))
477 .route("/api/v1/geo/sync", post(geo_sync_handler))
478 .route("/api/v1/geo/peers", get(geo_peers_handler))
479 .route("/api/v1/geo/failover", post(geo_failover_handler));
480 #[cfg(feature = "replication")]
482 let app = app
483 .route("/internal/promote", post(promote_handler))
484 .route("/internal/repoint", post(repoint_handler));
485 let app = app;
486
487 #[cfg(feature = "embedded-sync")]
489 let app = app
490 .route("/api/v1/sync/pull", post(super::api::sync_pull_handler))
491 .route("/api/v1/sync/push", post(super::api::sync_push_handler));
492
493 let app = app
494 .with_state(app_state.clone())
495 .layer(middleware::from_fn_with_state(
498 app_state,
499 read_only_middleware,
500 ))
501 .layer(middleware::from_fn_with_state(
502 rate_limit_state,
503 rate_limit_middleware,
504 ))
505 .layer(middleware::from_fn_with_state(auth_state, auth_middleware))
506 .layer(
507 CorsLayer::new()
508 .allow_origin(Any)
509 .allow_methods(Any)
510 .allow_headers(Any),
511 )
512 .layer(TraceLayer::new_for_http());
513
514 #[cfg(feature = "prime")]
516 let app = {
517 let data_dir =
518 std::env::var("PRIME_DATA_DIR").unwrap_or_else(|_| "/tmp/prime-data".to_string());
519 match crate::prime::Prime::open(&data_dir).await {
520 Ok(prime) => {
521 let prime_state = Arc::new(super::prime_api::PrimeState { prime });
522 tracing::info!("Prime API enabled at /api/v1/prime/*");
523 app.nest(
524 "/api/v1/prime",
525 super::prime_api::prime_router().with_state(prime_state),
526 )
527 }
528 Err(e) => {
529 tracing::warn!("Prime API disabled: failed to open Prime: {e}");
530 app
531 }
532 }
533 };
534
535 let listener = tokio::net::TcpListener::bind(addr).await?;
536
537 axum::serve(listener, app)
539 .with_graceful_shutdown(shutdown_signal())
540 .await?;
541
542 tracing::info!("🛑 AllSource Core shutdown complete");
543 Ok(())
544}
545
546const WRITE_PATHS: &[&str] = &[
548 "/api/v1/events",
549 "/api/v1/events/batch",
550 "/api/v1/snapshots",
551 "/api/v1/projections/",
552 "/api/v1/schemas",
553 "/api/v1/replay",
554 "/api/v1/pipelines",
555 "/api/v1/compaction/trigger",
556 "/api/v1/audit/events",
557 "/api/v1/config",
558 "/api/v1/webhooks",
559 "/api/v1/demo/seed",
560];
561
562fn is_write_request(method: &axum::http::Method, path: &str) -> bool {
564 use axum::http::Method;
565 if method != Method::POST && method != Method::PUT && method != Method::DELETE {
567 return false;
568 }
569 WRITE_PATHS
570 .iter()
571 .any(|write_path| path.starts_with(write_path))
572}
573
574fn is_internal_request(path: &str) -> bool {
576 path.starts_with("/internal/")
577 || path.starts_with("/api/v1/cluster/")
578 || path.starts_with("/api/v1/geo/")
579}
580
581async fn read_only_middleware(
587 State(state): State<AppState>,
588 request: axum::extract::Request,
589 next: axum::middleware::Next,
590) -> axum::response::Response {
591 let path = request.uri().path();
592 if state.role.load().is_follower()
593 && is_write_request(request.method(), path)
594 && !is_internal_request(path)
595 {
596 return (
597 axum::http::StatusCode::CONFLICT,
598 axum::Json(serde_json::json!({
599 "error": "read_only",
600 "message": "This node is a read-only follower"
601 })),
602 )
603 .into_response();
604 }
605 next.run(request).await
606}
607
608async fn health_v1(State(state): State<AppState>) -> impl IntoResponse {
613 let has_system_repos = state.service_container.has_system_repositories();
614
615 let system_streams = if has_system_repos {
616 let (tenant_count, config_count, total_events) =
617 if let Some(store) = state.service_container.system_store() {
618 use crate::domain::value_objects::system_stream::SystemDomain;
619 (
620 store.count_stream(SystemDomain::Tenant),
621 store.count_stream(SystemDomain::Config),
622 store.total_events(),
623 )
624 } else {
625 (0, 0, 0)
626 };
627
628 serde_json::json!({
629 "status": "healthy",
630 "mode": "event-sourced",
631 "total_events": total_events,
632 "tenant_events": tenant_count,
633 "config_events": config_count,
634 })
635 } else {
636 serde_json::json!({
637 "status": "disabled",
638 "mode": "in-memory",
639 })
640 };
641
642 let replication = {
643 #[cfg(feature = "replication")]
644 {
645 let shipper_guard = state.wal_shipper.read().await;
646 if let Some(ref shipper) = *shipper_guard {
647 serde_json::to_value(shipper.status()).unwrap_or_default()
648 } else if let Some(ref receiver) = state.wal_receiver {
649 serde_json::to_value(receiver.status()).unwrap_or_default()
650 } else {
651 serde_json::json!(null)
652 }
653 }
654 #[cfg(not(feature = "replication"))]
655 serde_json::json!({"edition": "community", "status": "not_available"})
656 };
657
658 let current_role = state.role.load();
659
660 Json(serde_json::json!({
661 "status": "healthy",
662 "service": "allsource-core",
663 "version": env!("CARGO_PKG_VERSION"),
664 "role": current_role,
665 "system_streams": system_streams,
666 "replication": replication,
667 }))
668}
669
670#[cfg(feature = "replication")]
676async fn promote_handler(State(state): State<AppState>) -> impl IntoResponse {
677 let current_role = state.role.load();
678 if current_role == NodeRole::Leader {
679 return (
680 axum::http::StatusCode::OK,
681 Json(serde_json::json!({
682 "status": "already_leader",
683 "message": "This node is already the leader",
684 })),
685 );
686 }
687
688 tracing::info!("PROMOTE: Switching role from follower to leader");
689
690 state.role.store(NodeRole::Leader);
692
693 if let Some(ref receiver) = state.wal_receiver {
695 receiver.shutdown();
696 tracing::info!("PROMOTE: WAL receiver shutdown signalled");
697 }
698
699 let replication_port = state.replication_port;
701 let (mut shipper, tx) = WalShipper::new();
702 state.store.enable_wal_replication(tx);
703 shipper.set_store(Arc::clone(&state.store));
704 shipper.set_metrics(state.store.metrics());
705 let shipper = Arc::new(shipper);
706
707 {
709 let mut shipper_guard = state.wal_shipper.write().await;
710 *shipper_guard = Some(Arc::clone(&shipper));
711 }
712
713 let shipper_clone = Arc::clone(&shipper);
715 tokio::spawn(async move {
716 if let Err(e) = shipper_clone.serve(replication_port).await {
717 tracing::error!("Promoted WAL shipper error: {}", e);
718 }
719 });
720
721 tracing::info!(
722 "PROMOTE: Now accepting writes. WAL shipper listening on port {}",
723 replication_port,
724 );
725
726 (
727 axum::http::StatusCode::OK,
728 Json(serde_json::json!({
729 "status": "promoted",
730 "role": "leader",
731 "replication_port": replication_port,
732 })),
733 )
734}
735
736#[cfg(feature = "replication")]
741async fn repoint_handler(
742 State(state): State<AppState>,
743 axum::extract::Query(params): axum::extract::Query<std::collections::HashMap<String, String>>,
744) -> impl IntoResponse {
745 let current_role = state.role.load();
746 if current_role != NodeRole::Follower {
747 return (
748 axum::http::StatusCode::CONFLICT,
749 Json(serde_json::json!({
750 "error": "not_follower",
751 "message": "Repoint only applies to follower nodes",
752 })),
753 );
754 }
755
756 let new_leader = match params.get("leader") {
757 Some(l) if !l.is_empty() => l.clone(),
758 _ => {
759 return (
760 axum::http::StatusCode::BAD_REQUEST,
761 Json(serde_json::json!({
762 "error": "missing_leader",
763 "message": "Query parameter 'leader' is required (e.g. ?leader=new-leader:3910)",
764 })),
765 );
766 }
767 };
768
769 tracing::info!("REPOINT: Switching replication target to {}", new_leader);
770
771 if let Some(ref receiver) = state.wal_receiver {
772 receiver.repoint(&new_leader);
773 tracing::info!("REPOINT: WAL receiver repointed to {}", new_leader);
774 } else {
775 tracing::warn!("REPOINT: No WAL receiver to repoint");
776 }
777
778 (
779 axum::http::StatusCode::OK,
780 Json(serde_json::json!({
781 "status": "repointed",
782 "new_leader": new_leader,
783 })),
784 )
785}
786
787async fn cluster_status_handler(State(state): State<AppState>) -> impl IntoResponse {
793 let Some(ref cm) = state.cluster_manager else {
794 return (
795 axum::http::StatusCode::SERVICE_UNAVAILABLE,
796 Json(serde_json::json!({
797 "error": "cluster_not_enabled",
798 "message": "Cluster mode is not enabled on this node"
799 })),
800 );
801 };
802
803 let status = cm.status().await;
804 (
805 axum::http::StatusCode::OK,
806 Json(serde_json::to_value(status).unwrap()),
807 )
808}
809
810async fn cluster_list_members_handler(State(state): State<AppState>) -> impl IntoResponse {
812 let Some(ref cm) = state.cluster_manager else {
813 return (
814 axum::http::StatusCode::SERVICE_UNAVAILABLE,
815 Json(serde_json::json!({
816 "error": "cluster_not_enabled",
817 "message": "Cluster mode is not enabled on this node"
818 })),
819 );
820 };
821
822 let members = cm.all_members();
823 (
824 axum::http::StatusCode::OK,
825 Json(serde_json::json!({
826 "members": members,
827 "count": members.len(),
828 })),
829 )
830}
831
832async fn cluster_add_member_handler(
834 State(state): State<AppState>,
835 Json(member): Json<ClusterMember>,
836) -> impl IntoResponse {
837 let Some(ref cm) = state.cluster_manager else {
838 return (
839 axum::http::StatusCode::SERVICE_UNAVAILABLE,
840 Json(serde_json::json!({
841 "error": "cluster_not_enabled",
842 "message": "Cluster mode is not enabled on this node"
843 })),
844 );
845 };
846
847 let node_id = member.node_id;
848 cm.add_member(member).await;
849
850 tracing::info!("Cluster member {} added", node_id);
851 (
852 axum::http::StatusCode::OK,
853 Json(serde_json::json!({
854 "status": "added",
855 "node_id": node_id,
856 })),
857 )
858}
859
860async fn cluster_remove_member_handler(
862 State(state): State<AppState>,
863 Path(node_id): Path<u32>,
864) -> impl IntoResponse {
865 let Some(ref cm) = state.cluster_manager else {
866 return (
867 axum::http::StatusCode::SERVICE_UNAVAILABLE,
868 Json(serde_json::json!({
869 "error": "cluster_not_enabled",
870 "message": "Cluster mode is not enabled on this node"
871 })),
872 );
873 };
874
875 match cm.remove_member(node_id).await {
876 Some(_) => {
877 tracing::info!("Cluster member {} removed", node_id);
878 (
879 axum::http::StatusCode::OK,
880 Json(serde_json::json!({
881 "status": "removed",
882 "node_id": node_id,
883 })),
884 )
885 }
886 None => (
887 axum::http::StatusCode::NOT_FOUND,
888 Json(serde_json::json!({
889 "error": "not_found",
890 "message": format!("Node {} not found in cluster", node_id),
891 })),
892 ),
893 }
894}
895
896#[derive(serde::Deserialize)]
898struct HeartbeatRequest {
899 wal_offset: u64,
900 #[serde(default = "default_true")]
901 healthy: bool,
902}
903
904fn default_true() -> bool {
905 true
906}
907
908async fn cluster_heartbeat_handler(
909 State(state): State<AppState>,
910 Path(node_id): Path<u32>,
911 Json(req): Json<HeartbeatRequest>,
912) -> impl IntoResponse {
913 let Some(ref cm) = state.cluster_manager else {
914 return (
915 axum::http::StatusCode::SERVICE_UNAVAILABLE,
916 Json(serde_json::json!({
917 "error": "cluster_not_enabled",
918 "message": "Cluster mode is not enabled on this node"
919 })),
920 );
921 };
922
923 cm.update_member_heartbeat(node_id, req.wal_offset, req.healthy);
924 (
925 axum::http::StatusCode::OK,
926 Json(serde_json::json!({
927 "status": "updated",
928 "node_id": node_id,
929 })),
930 )
931}
932
933async fn cluster_vote_handler(
935 State(state): State<AppState>,
936 Json(request): Json<VoteRequest>,
937) -> impl IntoResponse {
938 let Some(ref cm) = state.cluster_manager else {
939 return (
940 axum::http::StatusCode::SERVICE_UNAVAILABLE,
941 Json(serde_json::json!({
942 "error": "cluster_not_enabled",
943 "message": "Cluster mode is not enabled on this node"
944 })),
945 );
946 };
947
948 let response = cm.handle_vote_request(&request).await;
949 (
950 axum::http::StatusCode::OK,
951 Json(serde_json::to_value(response).unwrap()),
952 )
953}
954
955async fn cluster_election_handler(State(state): State<AppState>) -> impl IntoResponse {
957 let Some(ref cm) = state.cluster_manager else {
958 return (
959 axum::http::StatusCode::SERVICE_UNAVAILABLE,
960 Json(serde_json::json!({
961 "error": "cluster_not_enabled",
962 "message": "Cluster mode is not enabled on this node"
963 })),
964 );
965 };
966
967 let candidate = cm.select_leader_candidate();
969
970 match candidate {
971 Some(candidate_id) => {
972 let new_term = cm.start_election().await;
973 tracing::info!(
974 "Cluster election started: term={}, candidate={}",
975 new_term,
976 candidate_id,
977 );
978
979 if candidate_id == cm.self_id() {
982 cm.become_leader(new_term).await;
983 tracing::info!("Node {} became leader at term {}", candidate_id, new_term);
984 }
985
986 (
987 axum::http::StatusCode::OK,
988 Json(serde_json::json!({
989 "status": "election_started",
990 "term": new_term,
991 "candidate_id": candidate_id,
992 "self_is_leader": candidate_id == cm.self_id(),
993 })),
994 )
995 }
996 None => (
997 axum::http::StatusCode::CONFLICT,
998 Json(serde_json::json!({
999 "error": "no_candidates",
1000 "message": "No healthy members available for leader election",
1001 })),
1002 ),
1003 }
1004}
1005
1006async fn cluster_partitions_handler(State(state): State<AppState>) -> impl IntoResponse {
1008 let Some(ref cm) = state.cluster_manager else {
1009 return (
1010 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1011 Json(serde_json::json!({
1012 "error": "cluster_not_enabled",
1013 "message": "Cluster mode is not enabled on this node"
1014 })),
1015 );
1016 };
1017
1018 let registry = cm.registry();
1019 let distribution = registry.partition_distribution();
1020 let total_partitions: usize = distribution.values().map(std::vec::Vec::len).sum();
1021
1022 (
1023 axum::http::StatusCode::OK,
1024 Json(serde_json::json!({
1025 "total_partitions": total_partitions,
1026 "node_count": registry.node_count(),
1027 "healthy_node_count": registry.healthy_node_count(),
1028 "distribution": distribution,
1029 })),
1030 )
1031}
1032
1033async fn geo_status_handler(State(state): State<AppState>) -> impl IntoResponse {
1039 let Some(ref geo) = state.geo_replication else {
1040 return (
1041 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1042 Json(serde_json::json!({
1043 "error": "geo_replication_not_enabled",
1044 "message": "Geo-replication is not enabled on this node"
1045 })),
1046 );
1047 };
1048
1049 let status = geo.status();
1050 (
1051 axum::http::StatusCode::OK,
1052 Json(serde_json::to_value(status).unwrap()),
1053 )
1054}
1055
1056async fn geo_sync_handler(
1058 State(state): State<AppState>,
1059 Json(request): Json<GeoSyncRequest>,
1060) -> impl IntoResponse {
1061 let Some(ref geo) = state.geo_replication else {
1062 return (
1063 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1064 Json(serde_json::json!({
1065 "error": "geo_replication_not_enabled",
1066 "message": "Geo-replication is not enabled on this node"
1067 })),
1068 );
1069 };
1070
1071 tracing::info!(
1072 "Geo-sync received from region '{}': {} events",
1073 request.source_region,
1074 request.events.len(),
1075 );
1076
1077 let response = geo.receive_sync(&request);
1078 (
1079 axum::http::StatusCode::OK,
1080 Json(serde_json::to_value(response).unwrap()),
1081 )
1082}
1083
1084async fn geo_peers_handler(State(state): State<AppState>) -> impl IntoResponse {
1086 let Some(ref geo) = state.geo_replication else {
1087 return (
1088 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1089 Json(serde_json::json!({
1090 "error": "geo_replication_not_enabled",
1091 "message": "Geo-replication is not enabled on this node"
1092 })),
1093 );
1094 };
1095
1096 let status = geo.status();
1097 (
1098 axum::http::StatusCode::OK,
1099 Json(serde_json::json!({
1100 "region_id": status.region_id,
1101 "peers": status.peers,
1102 })),
1103 )
1104}
1105
1106async fn geo_failover_handler(State(state): State<AppState>) -> impl IntoResponse {
1108 let Some(ref geo) = state.geo_replication else {
1109 return (
1110 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1111 Json(serde_json::json!({
1112 "error": "geo_replication_not_enabled",
1113 "message": "Geo-replication is not enabled on this node"
1114 })),
1115 );
1116 };
1117
1118 match geo.select_failover_region() {
1119 Some(failover_region) => {
1120 tracing::info!(
1121 "Geo-failover: selected region '{}' as failover target",
1122 failover_region,
1123 );
1124 (
1125 axum::http::StatusCode::OK,
1126 Json(serde_json::json!({
1127 "status": "failover_target_selected",
1128 "failover_region": failover_region,
1129 "message": "Region selected for failover. DNS/routing update required externally.",
1130 })),
1131 )
1132 }
1133 None => (
1134 axum::http::StatusCode::CONFLICT,
1135 Json(serde_json::json!({
1136 "error": "no_healthy_peers",
1137 "message": "No healthy peer regions available for failover",
1138 })),
1139 ),
1140 }
1141}
1142
1143async fn shutdown_signal() {
1145 let ctrl_c = async {
1146 tokio::signal::ctrl_c()
1147 .await
1148 .expect("failed to install Ctrl+C handler");
1149 };
1150
1151 #[cfg(unix)]
1152 let terminate = async {
1153 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
1154 .expect("failed to install SIGTERM handler")
1155 .recv()
1156 .await;
1157 };
1158
1159 #[cfg(not(unix))]
1160 let terminate = std::future::pending::<()>();
1161
1162 tokio::select! {
1163 () = ctrl_c => {
1164 tracing::info!("📤 Received Ctrl+C, initiating graceful shutdown...");
1165 }
1166 () = terminate => {
1167 tracing::info!("📤 Received SIGTERM, initiating graceful shutdown...");
1168 }
1169 }
1170}