Skip to main content

allsource_core/infrastructure/web/
api_v1.rs

1/// v1.0 API router with authentication and multi-tenancy
2use 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/// Node role for leader-follower replication
35#[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    /// Detect role from environment variables.
44    ///
45    /// Checks `ALLSOURCE_ROLE` ("leader" or "follower") first,
46    /// then falls back to `ALLSOURCE_READ_ONLY` ("true" → follower).
47    /// Defaults to `Leader` if neither is set.
48    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/// Thread-safe, runtime-mutable node role for failover support.
99///
100/// Wraps an `AtomicU8` so the read-only middleware and health endpoint
101/// always see the current role, even after a sentinel-triggered promotion.
102#[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/// Unified application state for all handlers
121#[derive(Clone)]
122pub struct AppState {
123    pub store: Arc<EventStore>,
124    pub auth_manager: Arc<AuthManager>,
125    pub tenant_repo: Arc<dyn TenantRepository>,
126    /// Service container for paywall domain use cases (Creator, Article, Payment, etc.)
127    pub service_container: ServiceContainer,
128    /// Node role for leader-follower replication (runtime-mutable for failover)
129    pub role: AtomicNodeRole,
130    /// WAL shipper for replication status reporting (leader only).
131    /// Wrapped in RwLock so a promoted follower can install a shipper at runtime.
132    pub wal_shipper: Arc<tokio::sync::RwLock<Option<Arc<WalShipper>>>>,
133    /// WAL receiver for replication status reporting (follower only)
134    pub wal_receiver: Option<Arc<WalReceiver>>,
135    /// Replication port used by the WAL shipper (needed for runtime promotion)
136    pub replication_port: u16,
137    /// Cluster manager for multi-node consensus and membership (optional)
138    pub cluster_manager: Option<Arc<ClusterManager>>,
139    /// Geo-replication manager for cross-region replication (optional)
140    pub geo_replication: Option<Arc<GeoReplicationManager>>,
141}
142
143// Enable extracting Arc<EventStore> from AppState
144// This allows handlers that expect State<Arc<EventStore>> to work with AppState
145impl 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        // Public routes (no auth)
186        .route("/health", get(health_v1))
187        .route("/metrics", get(super::api::prometheus_metrics))
188        // Auth routes
189        .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        // Tenant routes (protected — enterprise only)
198        ;
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            // Qualified path: `patch` is only used here, inside the
208            // `multi-tenant`-gated block, so importing it unconditionally would
209            // warn unused under the community edition.
210            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        // Audit endpoints (admin only)
247        .route("/api/v1/audit/events", post(log_audit_event))
248        .route("/api/v1/audit/events", get(query_audit_events))
249        // Config endpoints (admin only)
250        .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        // Demo seeding
260        .route(
261            "/api/v1/demo/seed",
262            post(super::demo_api::demo_seed_handler),
263        )
264        // Event and data routes (protected by auth)
265        .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        // v0.10: Stream and event type discovery endpoints
291        .route("/api/v1/streams", get(super::api::list_streams))
292        .route("/api/v1/event-types", get(super::api::list_event_types))
293        // Analytics
294        .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        // Snapshots
307        .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        // Compaction
314        .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        // Schemas
323        .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        // Replay
339        .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        // Pipelines
354        .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        // v0.7: Projection State API for Query Service integration
377        .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        // v0.11: Webhook management
423        .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        // v0.14: Durable consumer subscriptions
442        .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        // v1.8: Cluster membership management API
456        .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        // v2.0: Advanced query features
474        .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        // v1.9: Geo-replication API
494        .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    // Internal endpoints for sentinel-driven failover (enterprise only)
499    #[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    // v0.11: Bidirectional sync protocol (embedded↔server)
506    #[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        // IMPORTANT: Middleware layers execute from bottom to top in Tower/Axum
514        // Read-only middleware runs after auth (applied before rate limit layer)
515        .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    // Prime API — nested router with its own state (feature-gated)
533    #[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    // Graceful shutdown on SIGTERM (required for serverless platforms)
556    axum::serve(listener, app)
557        .with_graceful_shutdown(shutdown_signal())
558        .await?;
559
560    tracing::info!("🛑 AllSource Core shutdown complete");
561    Ok(())
562}
563
564/// Write paths that should be rejected when running as a follower.
565const 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
581/// Returns true if this request is a write operation that should be blocked on followers.
582fn is_write_request(method: &axum::http::Method, path: &str) -> bool {
583    use axum::http::Method;
584    // Only POST/PUT/DELETE are writes
585    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
593/// Returns true if the request targets an internal or cluster endpoint (not subject to read-only checks).
594fn 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
600/// Middleware that rejects write requests when the node is a follower.
601///
602/// Returns HTTP 409 Conflict with `{"error": "read_only", "message": "..."}`.
603/// Internal endpoints (`/internal/*`) are exempt — they are used by the sentinel
604/// to trigger promotion and repointing during failover.
605async 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
627/// Health endpoint with system stream health reporting.
628///
629/// Reports overall health plus detailed system metadata health when
630/// event-sourced system repositories are configured.
631async 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/// POST /internal/promote — Promote this follower to leader.
690///
691/// Called by the sentinel process during automated failover.
692/// Switches the node role to leader, stops WAL receiving, and starts
693/// a WAL shipper on the replication port so other followers can connect.
694#[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    // 1. Switch role — the read-only middleware will immediately start accepting writes
710    state.role.store(NodeRole::Leader);
711
712    // 2. Signal the WAL receiver to stop (it will stop reconnecting)
713    if let Some(ref receiver) = state.wal_receiver {
714        receiver.shutdown();
715        tracing::info!("PROMOTE: WAL receiver shutdown signalled");
716    }
717
718    // 3. Start a new WAL shipper so remaining followers can connect
719    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    // Install into AppState so health endpoint reports shipper status
727    {
728        let mut shipper_guard = state.wal_shipper.write().await;
729        *shipper_guard = Some(Arc::clone(&shipper));
730    }
731
732    // Spawn the shipper server
733    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/// POST /internal/repoint?leader=host:port — Switch replication target.
756///
757/// Called by the sentinel process to tell a follower to disconnect from
758/// the old leader and reconnect to a newly promoted leader.
759#[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
806// =============================================================================
807// Cluster Management Handlers (v1.8)
808// =============================================================================
809
810/// GET /api/v1/cluster/status — Get cluster status including term, leader, and members
811async 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
829/// GET /api/v1/cluster/members — List all cluster members
830async 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
851/// POST /api/v1/cluster/members — Add a member to the cluster
852async 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
879/// DELETE /api/v1/cluster/members/{node_id} — Remove a member from the cluster
880async 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/// POST /api/v1/cluster/members/{node_id}/heartbeat — Update member heartbeat
916#[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
952/// POST /api/v1/cluster/vote — Handle a vote request (simplified Raft RequestVote)
953async 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
974/// POST /api/v1/cluster/election — Trigger a leader election (manual failover)
975async 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    // Select best candidate deterministically
987    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 this node is the candidate, become leader immediately
999            // (In a full Raft, we'd collect votes from a majority first)
1000            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
1025/// GET /api/v1/cluster/partitions — Get partition distribution across nodes
1026async 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
1052// =============================================================================
1053// Geo-Replication Handlers (v1.9)
1054// =============================================================================
1055
1056/// GET /api/v1/geo/status — Get geo-replication status
1057async 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
1075/// POST /api/v1/geo/sync — Receive replicated events from a peer region
1076async 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
1103/// GET /api/v1/geo/peers — List peer regions and their health
1104async 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
1125/// POST /api/v1/geo/failover — Trigger regional failover
1126async 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
1162/// Listen for shutdown signals (SIGTERM for serverless, SIGINT for local dev)
1163async 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}