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}/usage/queries/admit",
220 post(super::tenant_query_usage_api::admit_query_usage_handler)
221 .layer(axum::extract::DefaultBodyLimit::max(2_048)),
222 )
223 .route(
224 "/api/v1/tenants/{id}/usage/queries",
225 get(super::tenant_query_usage_api::query_usage_handler),
226 )
227 .route(
228 "/api/v1/tenants/{id}/usage/queries/reset",
229 post(super::tenant_query_usage_api::reset_query_usage_handler)
230 .layer(axum::extract::DefaultBodyLimit::max(2_048)),
231 )
232 .route(
233 "/api/v1/tenants/{id}/schema-enforcement",
234 put(update_schema_enforcement_handler),
235 )
236 .route(
237 "/api/v1/tenants/{id}/deactivate",
238 post(deactivate_tenant_handler),
239 )
240 .route(
241 "/api/v1/tenants/{id}/activate",
242 post(activate_tenant_handler),
243 )
244 .route("/api/v1/tenants/{id}", delete(delete_tenant_handler));
245 let app = app
246 .route("/api/v1/audit/events", post(log_audit_event))
248 .route("/api/v1/audit/events", get(query_audit_events))
249 .route("/api/v1/config", get(list_configs))
251 .route("/api/v1/config", post(set_config))
252 .route(
253 "/api/v1/config/conditional/set",
254 post(set_config_conditionally),
255 )
256 .route("/api/v1/config/{key}", get(get_config))
257 .route("/api/v1/config/{key}", put(update_config))
258 .route("/api/v1/config/{key}", delete(delete_config))
259 .route(
261 "/api/v1/demo/seed",
262 post(super::demo_api::demo_seed_handler),
263 )
264 .route("/api/v1/events", post(super::api::ingest_event_v1))
266 .route(
267 "/api/v1/events/batch",
268 post(super::api::ingest_events_batch_v1),
269 )
270 .route("/api/v1/events/query", get(super::api::query_events))
271 .route(
272 "/api/v1/events/{event_id}",
273 get(super::api::get_event_by_id),
274 )
275 .route("/api/v1/events/stream", get(super::api::events_websocket))
276 .route("/api/v1/entities", get(super::api::list_entities))
277 .route(
278 "/api/v1/entities/duplicates",
279 get(super::api::detect_duplicates),
280 )
281 .route(
282 "/api/v1/entities/{entity_id}/state",
283 get(super::api::get_entity_state),
284 )
285 .route(
286 "/api/v1/entities/{entity_id}/snapshot",
287 get(super::api::get_entity_snapshot),
288 )
289 .route("/api/v1/stats", get(super::api::get_stats))
290 .route("/api/v1/streams", get(super::api::list_streams))
292 .route("/api/v1/event-types", get(super::api::list_event_types))
293 .route(
295 "/api/v1/analytics/frequency",
296 get(super::api::analytics_frequency),
297 )
298 .route(
299 "/api/v1/analytics/summary",
300 get(super::api::analytics_summary),
301 )
302 .route(
303 "/api/v1/analytics/correlation",
304 get(super::api::analytics_correlation),
305 )
306 .route("/api/v1/snapshots", post(super::api::create_snapshot))
308 .route("/api/v1/snapshots", get(super::api::list_snapshots))
309 .route(
310 "/api/v1/snapshots/{entity_id}/latest",
311 get(super::api::get_latest_snapshot),
312 )
313 .route(
315 "/api/v1/compaction/trigger",
316 post(super::api::trigger_compaction),
317 )
318 .route(
319 "/api/v1/compaction/stats",
320 get(super::api::compaction_stats),
321 )
322 .route("/api/v1/schemas", post(super::api::register_schema))
324 .route("/api/v1/schemas", get(super::api::list_subjects))
325 .route("/api/v1/schemas/{subject}", get(super::api::get_schema))
326 .route(
327 "/api/v1/schemas/{subject}/versions",
328 get(super::api::list_schema_versions),
329 )
330 .route(
331 "/api/v1/schemas/validate",
332 post(super::api::validate_event_schema),
333 )
334 .route(
335 "/api/v1/schemas/{subject}/compatibility",
336 put(super::api::set_compatibility_mode),
337 )
338 .route("/api/v1/replay", post(super::api::start_replay))
340 .route("/api/v1/replay", get(super::api::list_replays))
341 .route(
342 "/api/v1/replay/{replay_id}",
343 get(super::api::get_replay_progress),
344 )
345 .route(
346 "/api/v1/replay/{replay_id}/cancel",
347 post(super::api::cancel_replay),
348 )
349 .route(
350 "/api/v1/replay/{replay_id}",
351 delete(super::api::delete_replay),
352 )
353 .route("/api/v1/pipelines", post(super::api::register_pipeline))
355 .route("/api/v1/pipelines", get(super::api::list_pipelines))
356 .route(
357 "/api/v1/pipelines/stats",
358 get(super::api::all_pipeline_stats),
359 )
360 .route(
361 "/api/v1/pipelines/{pipeline_id}",
362 get(super::api::get_pipeline),
363 )
364 .route(
365 "/api/v1/pipelines/{pipeline_id}",
366 delete(super::api::remove_pipeline),
367 )
368 .route(
369 "/api/v1/pipelines/{pipeline_id}/stats",
370 get(super::api::get_pipeline_stats),
371 )
372 .route(
373 "/api/v1/pipelines/{pipeline_id}/reset",
374 put(super::api::reset_pipeline),
375 )
376 .route("/api/v1/projections", get(super::api::list_projections))
378 .route(
379 "/api/v1/projections/{name}",
380 get(super::api::get_projection),
381 )
382 .route(
383 "/api/v1/projections/{name}",
384 delete(super::api::delete_projection),
385 )
386 .route(
387 "/api/v1/projections/{name}/state",
388 get(super::api::get_projection_state_summary),
389 )
390 .route(
391 "/api/v1/projections/{name}/reset",
392 post(super::api::reset_projection),
393 )
394 .route(
395 "/api/v1/projections/{name}/pause",
396 post(super::api::pause_projection),
397 )
398 .route(
399 "/api/v1/projections/{name}/start",
400 post(super::api::start_projection),
401 )
402 .route(
403 "/api/v1/projections/{name}/{entity_id}/state",
404 get(super::api::get_projection_state),
405 )
406 .route(
407 "/api/v1/projections/{name}/{entity_id}/state",
408 post(super::api::save_projection_state),
409 )
410 .route(
411 "/api/v1/projections/{name}/{entity_id}/state",
412 put(super::api::save_projection_state),
413 )
414 .route(
415 "/api/v1/projections/{name}/bulk",
416 post(super::api::bulk_get_projection_states),
417 )
418 .route(
419 "/api/v1/projections/{name}/bulk/save",
420 post(super::api::bulk_save_projection_states),
421 )
422 .route("/api/v1/webhooks", post(super::api::register_webhook))
424 .route("/api/v1/webhooks", get(super::api::list_webhooks))
425 .route(
426 "/api/v1/webhooks/{webhook_id}",
427 get(super::api::get_webhook),
428 )
429 .route(
430 "/api/v1/webhooks/{webhook_id}",
431 put(super::api::update_webhook),
432 )
433 .route(
434 "/api/v1/webhooks/{webhook_id}",
435 delete(super::api::delete_webhook),
436 )
437 .route(
438 "/api/v1/webhooks/{webhook_id}/deliveries",
439 get(super::api::list_webhook_deliveries),
440 )
441 .route("/api/v1/consumers", post(super::api::register_consumer))
443 .route(
444 "/api/v1/consumers/{consumer_id}",
445 get(super::api::get_consumer),
446 )
447 .route(
448 "/api/v1/consumers/{consumer_id}/events",
449 get(super::api::poll_consumer_events),
450 )
451 .route(
452 "/api/v1/consumers/{consumer_id}/ack",
453 post(super::api::ack_consumer),
454 )
455 .route("/api/v1/cluster/status", get(cluster_status_handler))
457 .route("/api/v1/cluster/members", get(cluster_list_members_handler))
458 .route("/api/v1/cluster/members", post(cluster_add_member_handler))
459 .route(
460 "/api/v1/cluster/members/{node_id}",
461 delete(cluster_remove_member_handler),
462 )
463 .route(
464 "/api/v1/cluster/members/{node_id}/heartbeat",
465 post(cluster_heartbeat_handler),
466 )
467 .route("/api/v1/cluster/vote", post(cluster_vote_handler))
468 .route("/api/v1/cluster/election", post(cluster_election_handler))
469 .route(
470 "/api/v1/cluster/partitions",
471 get(cluster_partitions_handler),
472 )
473 .route("/api/v1/graphql", post(super::api::graphql_query))
475 .route("/api/v1/geospatial/query", post(super::api::geo_query))
476 .route("/api/v1/geospatial/stats", get(super::api::geo_stats))
477 .route(
478 "/api/v1/exactly-once/stats",
479 get(super::api::exactly_once_stats),
480 )
481 .route(
482 "/api/v1/schema-evolution/history/{event_type}",
483 get(super::api::schema_evolution_history),
484 )
485 .route(
486 "/api/v1/schema-evolution/schema/{event_type}",
487 get(super::api::schema_evolution_schema),
488 )
489 .route(
490 "/api/v1/schema-evolution/stats",
491 get(super::api::schema_evolution_stats),
492 )
493 .route("/api/v1/geo/status", get(geo_status_handler))
495 .route("/api/v1/geo/sync", post(geo_sync_handler))
496 .route("/api/v1/geo/peers", get(geo_peers_handler))
497 .route("/api/v1/geo/failover", post(geo_failover_handler));
498 #[cfg(feature = "replication")]
500 let app = app
501 .route("/internal/promote", post(promote_handler))
502 .route("/internal/repoint", post(repoint_handler));
503 let app = app;
504
505 #[cfg(feature = "embedded-sync")]
507 let app = app
508 .route("/api/v1/sync/pull", post(super::api::sync_pull_handler))
509 .route("/api/v1/sync/push", post(super::api::sync_push_handler));
510
511 let app = app
512 .with_state(app_state.clone())
513 .layer(middleware::from_fn_with_state(
516 app_state,
517 read_only_middleware,
518 ))
519 .layer(middleware::from_fn_with_state(
520 rate_limit_state,
521 rate_limit_middleware,
522 ))
523 .layer(middleware::from_fn_with_state(auth_state, auth_middleware))
524 .layer(
525 CorsLayer::new()
526 .allow_origin(Any)
527 .allow_methods(Any)
528 .allow_headers(Any),
529 )
530 .layer(TraceLayer::new_for_http());
531
532 #[cfg(feature = "prime")]
534 let app = {
535 let data_dir =
536 std::env::var("PRIME_DATA_DIR").unwrap_or_else(|_| "/tmp/prime-data".to_string());
537 match crate::prime::Prime::open(&data_dir).await {
538 Ok(prime) => {
539 let prime_state = Arc::new(super::prime_api::PrimeState { prime });
540 tracing::info!("Prime API enabled at /api/v1/prime/*");
541 app.nest(
542 "/api/v1/prime",
543 super::prime_api::prime_router().with_state(prime_state),
544 )
545 }
546 Err(e) => {
547 tracing::warn!("Prime API disabled: failed to open Prime: {e}");
548 app
549 }
550 }
551 };
552
553 let listener = tokio::net::TcpListener::bind(addr).await?;
554
555 axum::serve(listener, app)
557 .with_graceful_shutdown(shutdown_signal())
558 .await?;
559
560 tracing::info!("🛑 AllSource Core shutdown complete");
561 Ok(())
562}
563
564const WRITE_PATHS: &[&str] = &[
566 "/api/v1/events",
567 "/api/v1/events/batch",
568 "/api/v1/snapshots",
569 "/api/v1/projections/",
570 "/api/v1/schemas",
571 "/api/v1/replay",
572 "/api/v1/pipelines",
573 "/api/v1/compaction/trigger",
574 "/api/v1/audit/events",
575 "/api/v1/config",
576 "/api/v1/tenants",
577 "/api/v1/webhooks",
578 "/api/v1/demo/seed",
579];
580
581fn is_write_request(method: &axum::http::Method, path: &str) -> bool {
583 use axum::http::Method;
584 if method != Method::POST && method != Method::PUT && method != Method::DELETE {
586 return false;
587 }
588 WRITE_PATHS
589 .iter()
590 .any(|write_path| path.starts_with(write_path))
591}
592
593fn is_internal_request(path: &str) -> bool {
595 path.starts_with("/internal/")
596 || path.starts_with("/api/v1/cluster/")
597 || path.starts_with("/api/v1/geo/")
598}
599
600async fn read_only_middleware(
606 State(state): State<AppState>,
607 request: axum::extract::Request,
608 next: axum::middleware::Next,
609) -> axum::response::Response {
610 let path = request.uri().path();
611 if state.role.load().is_follower()
612 && is_write_request(request.method(), path)
613 && !is_internal_request(path)
614 {
615 return (
616 axum::http::StatusCode::CONFLICT,
617 axum::Json(serde_json::json!({
618 "error": "read_only",
619 "message": "This node is a read-only follower"
620 })),
621 )
622 .into_response();
623 }
624 next.run(request).await
625}
626
627async fn health_v1(State(state): State<AppState>) -> impl IntoResponse {
632 let has_system_repos = state.service_container.has_system_repositories();
633
634 let system_streams = if has_system_repos {
635 let (tenant_count, config_count, total_events) =
636 if let Some(store) = state.service_container.system_store() {
637 use crate::domain::value_objects::system_stream::SystemDomain;
638 (
639 store.count_stream(SystemDomain::Tenant),
640 store.count_stream(SystemDomain::Config),
641 store.total_events(),
642 )
643 } else {
644 (0, 0, 0)
645 };
646
647 serde_json::json!({
648 "status": "healthy",
649 "mode": "event-sourced",
650 "total_events": total_events,
651 "tenant_events": tenant_count,
652 "config_events": config_count,
653 })
654 } else {
655 serde_json::json!({
656 "status": "disabled",
657 "mode": "in-memory",
658 })
659 };
660
661 let replication = {
662 #[cfg(feature = "replication")]
663 {
664 let shipper_guard = state.wal_shipper.read().await;
665 if let Some(ref shipper) = *shipper_guard {
666 serde_json::to_value(shipper.status()).unwrap_or_default()
667 } else if let Some(ref receiver) = state.wal_receiver {
668 serde_json::to_value(receiver.status()).unwrap_or_default()
669 } else {
670 serde_json::json!(null)
671 }
672 }
673 #[cfg(not(feature = "replication"))]
674 serde_json::json!({"edition": "community", "status": "not_available"})
675 };
676
677 let current_role = state.role.load();
678
679 Json(serde_json::json!({
680 "status": "healthy",
681 "service": "allsource-core",
682 "version": env!("CARGO_PKG_VERSION"),
683 "role": current_role,
684 "system_streams": system_streams,
685 "replication": replication,
686 }))
687}
688
689#[cfg(feature = "replication")]
695async fn promote_handler(State(state): State<AppState>) -> impl IntoResponse {
696 let current_role = state.role.load();
697 if current_role == NodeRole::Leader {
698 return (
699 axum::http::StatusCode::OK,
700 Json(serde_json::json!({
701 "status": "already_leader",
702 "message": "This node is already the leader",
703 })),
704 );
705 }
706
707 tracing::info!("PROMOTE: Switching role from follower to leader");
708
709 state.role.store(NodeRole::Leader);
711
712 if let Some(ref receiver) = state.wal_receiver {
714 receiver.shutdown();
715 tracing::info!("PROMOTE: WAL receiver shutdown signalled");
716 }
717
718 let replication_port = state.replication_port;
720 let (mut shipper, tx) = WalShipper::new();
721 state.store.enable_wal_replication(tx);
722 shipper.set_store(Arc::clone(&state.store));
723 shipper.set_metrics(state.store.metrics());
724 let shipper = Arc::new(shipper);
725
726 {
728 let mut shipper_guard = state.wal_shipper.write().await;
729 *shipper_guard = Some(Arc::clone(&shipper));
730 }
731
732 let shipper_clone = Arc::clone(&shipper);
734 tokio::spawn(async move {
735 if let Err(e) = shipper_clone.serve(replication_port).await {
736 tracing::error!("Promoted WAL shipper error: {}", e);
737 }
738 });
739
740 tracing::info!(
741 "PROMOTE: Now accepting writes. WAL shipper listening on port {}",
742 replication_port,
743 );
744
745 (
746 axum::http::StatusCode::OK,
747 Json(serde_json::json!({
748 "status": "promoted",
749 "role": "leader",
750 "replication_port": replication_port,
751 })),
752 )
753}
754
755#[cfg(feature = "replication")]
760async fn repoint_handler(
761 State(state): State<AppState>,
762 axum::extract::Query(params): axum::extract::Query<std::collections::HashMap<String, String>>,
763) -> impl IntoResponse {
764 let current_role = state.role.load();
765 if current_role != NodeRole::Follower {
766 return (
767 axum::http::StatusCode::CONFLICT,
768 Json(serde_json::json!({
769 "error": "not_follower",
770 "message": "Repoint only applies to follower nodes",
771 })),
772 );
773 }
774
775 let new_leader = match params.get("leader") {
776 Some(l) if !l.is_empty() => l.clone(),
777 _ => {
778 return (
779 axum::http::StatusCode::BAD_REQUEST,
780 Json(serde_json::json!({
781 "error": "missing_leader",
782 "message": "Query parameter 'leader' is required (e.g. ?leader=new-leader:3910)",
783 })),
784 );
785 }
786 };
787
788 tracing::info!("REPOINT: Switching replication target to {}", new_leader);
789
790 if let Some(ref receiver) = state.wal_receiver {
791 receiver.repoint(&new_leader);
792 tracing::info!("REPOINT: WAL receiver repointed to {}", new_leader);
793 } else {
794 tracing::warn!("REPOINT: No WAL receiver to repoint");
795 }
796
797 (
798 axum::http::StatusCode::OK,
799 Json(serde_json::json!({
800 "status": "repointed",
801 "new_leader": new_leader,
802 })),
803 )
804}
805
806async fn cluster_status_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 status = cm.status().await;
823 (
824 axum::http::StatusCode::OK,
825 Json(serde_json::to_value(status).unwrap()),
826 )
827}
828
829async fn cluster_list_members_handler(State(state): State<AppState>) -> impl IntoResponse {
831 let Some(ref cm) = state.cluster_manager else {
832 return (
833 axum::http::StatusCode::SERVICE_UNAVAILABLE,
834 Json(serde_json::json!({
835 "error": "cluster_not_enabled",
836 "message": "Cluster mode is not enabled on this node"
837 })),
838 );
839 };
840
841 let members = cm.all_members();
842 (
843 axum::http::StatusCode::OK,
844 Json(serde_json::json!({
845 "members": members,
846 "count": members.len(),
847 })),
848 )
849}
850
851async fn cluster_add_member_handler(
853 State(state): State<AppState>,
854 Json(member): Json<ClusterMember>,
855) -> impl IntoResponse {
856 let Some(ref cm) = state.cluster_manager else {
857 return (
858 axum::http::StatusCode::SERVICE_UNAVAILABLE,
859 Json(serde_json::json!({
860 "error": "cluster_not_enabled",
861 "message": "Cluster mode is not enabled on this node"
862 })),
863 );
864 };
865
866 let node_id = member.node_id;
867 cm.add_member(member).await;
868
869 tracing::info!("Cluster member {} added", node_id);
870 (
871 axum::http::StatusCode::OK,
872 Json(serde_json::json!({
873 "status": "added",
874 "node_id": node_id,
875 })),
876 )
877}
878
879async fn cluster_remove_member_handler(
881 State(state): State<AppState>,
882 Path(node_id): Path<u32>,
883) -> impl IntoResponse {
884 let Some(ref cm) = state.cluster_manager else {
885 return (
886 axum::http::StatusCode::SERVICE_UNAVAILABLE,
887 Json(serde_json::json!({
888 "error": "cluster_not_enabled",
889 "message": "Cluster mode is not enabled on this node"
890 })),
891 );
892 };
893
894 match cm.remove_member(node_id).await {
895 Some(_) => {
896 tracing::info!("Cluster member {} removed", node_id);
897 (
898 axum::http::StatusCode::OK,
899 Json(serde_json::json!({
900 "status": "removed",
901 "node_id": node_id,
902 })),
903 )
904 }
905 None => (
906 axum::http::StatusCode::NOT_FOUND,
907 Json(serde_json::json!({
908 "error": "not_found",
909 "message": format!("Node {} not found in cluster", node_id),
910 })),
911 ),
912 }
913}
914
915#[derive(serde::Deserialize)]
917struct HeartbeatRequest {
918 wal_offset: u64,
919 #[serde(default = "default_true")]
920 healthy: bool,
921}
922
923fn default_true() -> bool {
924 true
925}
926
927async fn cluster_heartbeat_handler(
928 State(state): State<AppState>,
929 Path(node_id): Path<u32>,
930 Json(req): Json<HeartbeatRequest>,
931) -> impl IntoResponse {
932 let Some(ref cm) = state.cluster_manager else {
933 return (
934 axum::http::StatusCode::SERVICE_UNAVAILABLE,
935 Json(serde_json::json!({
936 "error": "cluster_not_enabled",
937 "message": "Cluster mode is not enabled on this node"
938 })),
939 );
940 };
941
942 cm.update_member_heartbeat(node_id, req.wal_offset, req.healthy);
943 (
944 axum::http::StatusCode::OK,
945 Json(serde_json::json!({
946 "status": "updated",
947 "node_id": node_id,
948 })),
949 )
950}
951
952async fn cluster_vote_handler(
954 State(state): State<AppState>,
955 Json(request): Json<VoteRequest>,
956) -> 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 response = cm.handle_vote_request(&request).await;
968 (
969 axum::http::StatusCode::OK,
970 Json(serde_json::to_value(response).unwrap()),
971 )
972}
973
974async fn cluster_election_handler(State(state): State<AppState>) -> impl IntoResponse {
976 let Some(ref cm) = state.cluster_manager else {
977 return (
978 axum::http::StatusCode::SERVICE_UNAVAILABLE,
979 Json(serde_json::json!({
980 "error": "cluster_not_enabled",
981 "message": "Cluster mode is not enabled on this node"
982 })),
983 );
984 };
985
986 let candidate = cm.select_leader_candidate();
988
989 match candidate {
990 Some(candidate_id) => {
991 let new_term = cm.start_election().await;
992 tracing::info!(
993 "Cluster election started: term={}, candidate={}",
994 new_term,
995 candidate_id,
996 );
997
998 if candidate_id == cm.self_id() {
1001 cm.become_leader(new_term).await;
1002 tracing::info!("Node {} became leader at term {}", candidate_id, new_term);
1003 }
1004
1005 (
1006 axum::http::StatusCode::OK,
1007 Json(serde_json::json!({
1008 "status": "election_started",
1009 "term": new_term,
1010 "candidate_id": candidate_id,
1011 "self_is_leader": candidate_id == cm.self_id(),
1012 })),
1013 )
1014 }
1015 None => (
1016 axum::http::StatusCode::CONFLICT,
1017 Json(serde_json::json!({
1018 "error": "no_candidates",
1019 "message": "No healthy members available for leader election",
1020 })),
1021 ),
1022 }
1023}
1024
1025async fn cluster_partitions_handler(State(state): State<AppState>) -> impl IntoResponse {
1027 let Some(ref cm) = state.cluster_manager else {
1028 return (
1029 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1030 Json(serde_json::json!({
1031 "error": "cluster_not_enabled",
1032 "message": "Cluster mode is not enabled on this node"
1033 })),
1034 );
1035 };
1036
1037 let registry = cm.registry();
1038 let distribution = registry.partition_distribution();
1039 let total_partitions: usize = distribution.values().map(std::vec::Vec::len).sum();
1040
1041 (
1042 axum::http::StatusCode::OK,
1043 Json(serde_json::json!({
1044 "total_partitions": total_partitions,
1045 "node_count": registry.node_count(),
1046 "healthy_node_count": registry.healthy_node_count(),
1047 "distribution": distribution,
1048 })),
1049 )
1050}
1051
1052async fn geo_status_handler(State(state): State<AppState>) -> impl IntoResponse {
1058 let Some(ref geo) = state.geo_replication else {
1059 return (
1060 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1061 Json(serde_json::json!({
1062 "error": "geo_replication_not_enabled",
1063 "message": "Geo-replication is not enabled on this node"
1064 })),
1065 );
1066 };
1067
1068 let status = geo.status();
1069 (
1070 axum::http::StatusCode::OK,
1071 Json(serde_json::to_value(status).unwrap()),
1072 )
1073}
1074
1075async fn geo_sync_handler(
1077 State(state): State<AppState>,
1078 Json(request): Json<GeoSyncRequest>,
1079) -> impl IntoResponse {
1080 let Some(ref geo) = state.geo_replication else {
1081 return (
1082 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1083 Json(serde_json::json!({
1084 "error": "geo_replication_not_enabled",
1085 "message": "Geo-replication is not enabled on this node"
1086 })),
1087 );
1088 };
1089
1090 tracing::info!(
1091 "Geo-sync received from region '{}': {} events",
1092 request.source_region,
1093 request.events.len(),
1094 );
1095
1096 let response = geo.receive_sync(&request);
1097 (
1098 axum::http::StatusCode::OK,
1099 Json(serde_json::to_value(response).unwrap()),
1100 )
1101}
1102
1103async fn geo_peers_handler(State(state): State<AppState>) -> impl IntoResponse {
1105 let Some(ref geo) = state.geo_replication else {
1106 return (
1107 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1108 Json(serde_json::json!({
1109 "error": "geo_replication_not_enabled",
1110 "message": "Geo-replication is not enabled on this node"
1111 })),
1112 );
1113 };
1114
1115 let status = geo.status();
1116 (
1117 axum::http::StatusCode::OK,
1118 Json(serde_json::json!({
1119 "region_id": status.region_id,
1120 "peers": status.peers,
1121 })),
1122 )
1123}
1124
1125async fn geo_failover_handler(State(state): State<AppState>) -> impl IntoResponse {
1127 let Some(ref geo) = state.geo_replication else {
1128 return (
1129 axum::http::StatusCode::SERVICE_UNAVAILABLE,
1130 Json(serde_json::json!({
1131 "error": "geo_replication_not_enabled",
1132 "message": "Geo-replication is not enabled on this node"
1133 })),
1134 );
1135 };
1136
1137 match geo.select_failover_region() {
1138 Some(failover_region) => {
1139 tracing::info!(
1140 "Geo-failover: selected region '{}' as failover target",
1141 failover_region,
1142 );
1143 (
1144 axum::http::StatusCode::OK,
1145 Json(serde_json::json!({
1146 "status": "failover_target_selected",
1147 "failover_region": failover_region,
1148 "message": "Region selected for failover. DNS/routing update required externally.",
1149 })),
1150 )
1151 }
1152 None => (
1153 axum::http::StatusCode::CONFLICT,
1154 Json(serde_json::json!({
1155 "error": "no_healthy_peers",
1156 "message": "No healthy peer regions available for failover",
1157 })),
1158 ),
1159 }
1160}
1161
1162async fn shutdown_signal() {
1164 let ctrl_c = async {
1165 tokio::signal::ctrl_c()
1166 .await
1167 .expect("failed to install Ctrl+C handler");
1168 };
1169
1170 #[cfg(unix)]
1171 let terminate = async {
1172 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
1173 .expect("failed to install SIGTERM handler")
1174 .recv()
1175 .await;
1176 };
1177
1178 #[cfg(not(unix))]
1179 let terminate = std::future::pending::<()>();
1180
1181 tokio::select! {
1182 () = ctrl_c => {
1183 tracing::info!("📤 Received Ctrl+C, initiating graceful shutdown...");
1184 }
1185 () = terminate => {
1186 tracing::info!("📤 Received SIGTERM, initiating graceful shutdown...");
1187 }
1188 }
1189}