Skip to main content

admin/
service.rs

1use crate::error::Result as AdminResult;
2use crate::metrics::collect_system_metrics;
3use crate::realm::realm_to_proto;
4use actrix_proto::NodeAdminService;
5use actrix_proto::{
6    ConfigOverrideEntry as ProtoConfigOverrideEntry, ConfigType, CreateRealmRequest,
7    CreateRealmResponse, DeleteConfigOverrideRequest, DeleteConfigOverrideResponse,
8    DeleteRealmRequest, DeleteRealmResponse, GetConfigRequest, GetConfigResponse,
9    GetNodeInfoRequest, GetNodeInfoResponse, GetRealmRequest, GetRealmResponse,
10    ListConfigOverridesRequest, ListConfigOverridesResponse, ListRealmsRequest, ListRealmsResponse,
11    RealmInfo, ServiceStatus, SetConfigOverrideRequest, SetConfigOverrideResponse, ShutdownRequest,
12    ShutdownResponse, SystemMetrics, UpdateConfigRequest, UpdateConfigResponse, UpdateRealmRequest,
13    UpdateRealmResponse,
14};
15use chrono::Utc;
16use platform::ServiceCollector;
17use platform::config::ActrixConfig;
18use platform::config::config_store::{ConfigOverride, ConfigOverrideStore};
19use platform::config::registry;
20use platform::config::resolver::{self, ResolvedField};
21use platform::realm::{
22    DEFAULT_REALM_SECRET_PREVIOUS_GRACE_SECS, hash_realm_secret, rotate_realm_secret,
23};
24use platform::realm::{Realm, RealmStatus};
25use serde::Serialize;
26use serde_json::Value as JsonValue;
27use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
28use std::collections::HashMap;
29use std::future::Future;
30use std::path::PathBuf;
31use std::pin::Pin;
32use std::str::FromStr;
33use std::sync::Arc;
34use std::time::Instant;
35use tokio::sync::RwLock;
36use tonic::{Request, Response, Status};
37
38type MetricsFuture = Pin<Box<dyn Future<Output = AdminResult<SystemMetrics>> + Send>>;
39type MetricsProvider = Arc<dyn Fn() -> MetricsFuture + Send + Sync>;
40type ShutdownFuture = Pin<Box<dyn Future<Output = AdminResult<()>> + Send>>;
41type ShutdownHandler =
42    Arc<dyn Fn(bool, Option<i32>, Option<String>) -> ShutdownFuture + Send + Sync>;
43type ReloadFuture = Pin<Box<dyn Future<Output = AdminResult<bool>> + Send>>;
44type ReloadHandler = Arc<dyn Fn() -> ReloadFuture + Send + Sync>;
45type GrpcResult<T> = std::result::Result<T, Status>;
46
47#[derive(Hash, Eq, PartialEq, Clone)]
48struct ConfigKey {
49    config_type: i32,
50    key: String,
51}
52
53/// AdminApiService — canonical business-logic layer for node administration.
54///
55/// Both the REST/BFF (admin_api.rs) and the gRPC NodeAdminService trait call
56/// into the `*_direct()` methods defined here.  REST handlers should be thin
57/// wrappers that parse the request and format the JSON response.
58#[derive(Clone)]
59pub struct AdminApiService {
60    node_id: String,
61    name: String,
62    location_tag: String,
63    version: String,
64    config_store: Arc<RwLock<HashMap<ConfigKey, String>>>,
65    /// SQLite-backed config override store (L2 dynamic overrides).
66    /// Present when running in AdminUi head mode.
67    override_store: Option<Arc<ConfigOverrideStore>>,
68    metrics_provider: MetricsProvider,
69    shutdown_handler: Option<ShutdownHandler>,
70    reload_handler: Option<ReloadHandler>,
71    service_collector: ServiceCollector,
72    started_at: Instant,
73    /// Runtime config snapshot (set in AdminUi head mode).
74    running_config: Option<ActrixConfig>,
75    /// Raw TOML content for L1 detection (resolver).
76    toml_content: Option<String>,
77    /// Path to config.toml on disk.
78    config_path: Option<PathBuf>,
79}
80
81impl AdminApiService {
82    /// Create a new control gRPC service instance.
83    pub fn new(
84        node_id: impl Into<String>,
85        name: impl Into<String>,
86        location_tag: impl Into<String>,
87        version: impl Into<String>,
88        service_collector: ServiceCollector,
89    ) -> AdminResult<Self> {
90        Ok(Self {
91            node_id: node_id.into(),
92            name: name.into(),
93            location_tag: location_tag.into(),
94            version: version.into(),
95            config_store: Arc::new(RwLock::new(HashMap::new())),
96            override_store: None,
97            metrics_provider: Arc::new(|| Box::pin(async { collect_system_metrics().await })),
98            shutdown_handler: None,
99            reload_handler: None,
100            service_collector,
101            started_at: Instant::now(),
102            running_config: None,
103            toml_content: None,
104            config_path: None,
105        })
106    }
107
108    /// Override the metrics provider used by GetNodeInfo.
109    pub fn with_metrics_provider<F, Fut>(mut self, provider: F) -> Self
110    where
111        F: Fn() -> Fut + Send + Sync + 'static,
112        Fut: Future<Output = AdminResult<SystemMetrics>> + Send + 'static,
113    {
114        self.metrics_provider = Arc::new(move || {
115            let fut = provider();
116            Box::pin(fut)
117        });
118        self
119    }
120
121    /// Attach a shutdown handler invoked when Shutdown is accepted.
122    pub fn with_shutdown_handler<F, Fut>(mut self, handler: F) -> Self
123    where
124        F: Fn(bool, Option<i32>, Option<String>) -> Fut + Send + Sync + 'static,
125        Fut: Future<Output = AdminResult<()>> + Send + 'static,
126    {
127        self.shutdown_handler = Some(Arc::new(move |graceful, timeout, reason| {
128            let fut = handler(graceful, timeout, reason);
129            Box::pin(fut)
130        }));
131        self
132    }
133
134    /// Attach a SQLite-backed config override store for L2 dynamic overrides.
135    pub fn with_override_store(mut self, store: Arc<ConfigOverrideStore>) -> Self {
136        self.override_store = Some(store);
137        self
138    }
139
140    /// Attach a reload handler invoked when reload is requested.
141    pub fn with_reload_handler<F, Fut>(mut self, handler: F) -> Self
142    where
143        F: Fn() -> Fut + Send + Sync + 'static,
144        Fut: Future<Output = AdminResult<bool>> + Send + 'static,
145    {
146        self.reload_handler = Some(Arc::new(move || {
147            let fut = handler();
148            Box::pin(fut)
149        }));
150        self
151    }
152
153    /// Attach the runtime config snapshot (AdminUi head mode).
154    pub fn with_running_config(mut self, config: ActrixConfig) -> Self {
155        self.running_config = Some(config);
156        self
157    }
158
159    /// Attach the raw TOML content for L1 detection.
160    pub fn with_toml_content(mut self, content: String) -> Self {
161        self.toml_content = Some(content);
162        self
163    }
164
165    /// Attach the config file path.
166    pub fn with_config_path(mut self, path: PathBuf) -> Self {
167        self.config_path = Some(path);
168        self
169    }
170
171    /// Get the override store reference.
172    pub fn override_store(&self) -> Option<&Arc<ConfigOverrideStore>> {
173        self.override_store.as_ref()
174    }
175
176    fn build_config_key(config_type: ConfigType, key: String) -> ConfigKey {
177        ConfigKey {
178            config_type: config_type as i32,
179            key,
180        }
181    }
182
183    async fn load_realm(&self, realm_id: u32) -> GrpcResult<Realm> {
184        let realm = Realm::get(realm_id)
185            .await
186            .map_err(|e| Status::internal(format!("Failed to load realm: {e}")))?;
187
188        realm.ok_or_else(|| Status::not_found(format!("Realm not found: {realm_id}")))
189    }
190
191    async fn collect_metrics(&self) -> GrpcResult<SystemMetrics> {
192        (self.metrics_provider)()
193            .await
194            .map_err(|e| Status::internal(format!("Failed to collect metrics: {e}")))
195    }
196
197    /// Collect current service statuses for this node.
198    ///
199    /// Returns all service statuses from the service registry.
200    pub async fn service_statuses(&self) -> Vec<ServiceStatus> {
201        self.service_collector.all_statuses().await
202    }
203
204    // ── Helper: get overrides list ──────────────────────────────────
205
206    async fn overrides_list(&self) -> Vec<ConfigOverride> {
207        if let Some(store) = &self.override_store {
208            store.list_all().await.unwrap_or_default()
209        } else {
210            Vec::new()
211        }
212    }
213
214    fn toml_content_ref(&self) -> &str {
215        self.toml_content.as_deref().unwrap_or("")
216    }
217
218    fn require_running_config(&self) -> GrpcResult<&ActrixConfig> {
219        self.running_config
220            .as_ref()
221            .ok_or_else(|| Status::failed_precondition("Running config not available"))
222    }
223
224    // ── Direct methods (transport-agnostic) ─────────────────────────
225
226    /// Resolve the effective node name (L2 override → L1 config → running).
227    pub async fn resolve_node_name(&self) -> String {
228        let overrides = self.overrides_list().await;
229        let fields = resolver::resolve_for_service("platform", self.toml_content_ref(), &overrides);
230
231        fields
232            .iter()
233            .find(|f| f.key == "name")
234            .map(|f| f.effective_value.trim().to_string())
235            .filter(|n| !n.is_empty())
236            .unwrap_or_else(|| {
237                self.running_config
238                    .as_ref()
239                    .map(|c| c.name.clone())
240                    .unwrap_or_else(|| self.name.clone())
241            })
242    }
243
244    /// Health check info.
245    pub async fn health_info(&self) -> HealthInfo {
246        let node = self.resolve_node_name().await;
247        HealthInfo {
248            status: "healthy".to_string(),
249            node,
250            version: self.version.clone(),
251        }
252    }
253
254    /// List L2 dynamic config overrides.
255    pub async fn list_overrides_direct(&self) -> GrpcResult<Vec<ConfigOverride>> {
256        let store = self
257            .override_store
258            .as_ref()
259            .ok_or_else(|| Status::failed_precondition("Override store not available"))?;
260        store
261            .list_all()
262            .await
263            .map_err(|e| Status::internal(format!("Failed to list overrides: {e}")))
264    }
265
266    /// Set a dynamic config override.
267    pub async fn set_override_direct(&self, key: &str, value: &str, by: &str) -> GrpcResult<()> {
268        let store = self
269            .override_store
270            .as_ref()
271            .ok_or_else(|| Status::failed_precondition("Override store not available"))?;
272        store
273            .set(key, value, by)
274            .await
275            .map_err(|e| Status::invalid_argument(e.to_string()))
276    }
277
278    /// Delete a dynamic config override.
279    pub async fn delete_override_direct(&self, key: &str) -> GrpcResult<bool> {
280        let store = self
281            .override_store
282            .as_ref()
283            .ok_or_else(|| Status::failed_precondition("Override store not available"))?;
284        store
285            .delete(key)
286            .await
287            .map_err(|e| Status::internal(format!("Failed to delete override: {e}")))
288    }
289
290    /// Get platform detail with resolved config fields.
291    pub async fn get_platform_detail_direct(&self) -> GrpcResult<PlatformDetail> {
292        let cfg = self.require_running_config()?;
293
294        let http_bind = cfg.bind.http.as_ref().map(|h| {
295            serde_json::json!({
296                "ip": h.ip, "port": h.port,
297                "domain_name": h.domain_name, "advertised_ip": h.advertised_ip,
298                "advertised_port": h.effective_advertised_port(),
299                "cert": h.cert, "key": h.key,
300                "is_tls": h.is_tls(),
301            })
302        });
303
304        let config = serde_json::json!({
305            "enable": cfg.enable,
306            "name": cfg.name,
307            "env": cfg.env,
308            "location_tag": cfg.location_tag,
309            "sqlite_path": cfg.sqlite_path.display().to_string(),
310            "bind": {
311                "http": http_bind,
312                "ice": {
313                    "ip": cfg.bind.ice.ip.to_string(),
314                    "port": cfg.bind.ice.port,
315                    "advertised_ip": cfg.bind.ice.advertised_ip.to_string(),
316                    "advertised_port": cfg.bind.ice.advertised_port,
317                },
318            },
319            "turn": {
320                "relay_port_range": cfg.turn.relay_port_range,
321            },
322            "recording": {
323                "service_name": cfg.recording.service_name,
324            },
325            "control": {
326                "head": cfg.control.head,
327                "admin_ui": {
328                    "session_expiry_secs": cfg.control.admin_ui.session_expiry_secs,
329                },
330            },
331        });
332
333        let overrides = self.overrides_list().await;
334        let config_fields =
335            resolver::resolve_for_service("platform", self.toml_content_ref(), &overrides);
336
337        Ok(PlatformDetail {
338            config,
339            config_fields,
340        })
341    }
342
343    /// Get service detail with resolved config fields.
344    pub async fn get_service_detail_direct(&self, name: &str) -> GrpcResult<ServiceDetail> {
345        let type_id = service_name_to_type_id(name)
346            .ok_or_else(|| Status::not_found(format!("Unknown service: {name}")))?;
347
348        let cfg = self.require_running_config()?;
349
350        let enabled = match name {
351            "stun" => cfg.is_stun_enabled(),
352            "turn" => cfg.is_turn_enabled(),
353            "signaling" => cfg.is_signaling_enabled(),
354            "ais" => cfg.is_ais_enabled(),
355            "signer" => cfg.is_signer_enabled(),
356            _ => false,
357        };
358
359        let statuses = self.service_statuses().await;
360        let status = statuses.iter().find(|s| s.r#type == type_id).map(|s| {
361            serde_json::json!({
362                "name": s.name,
363                "type": s.r#type,
364                "is_healthy": s.is_healthy,
365                "active_connections": s.active_connections,
366                "total_requests": s.total_requests,
367                "failed_requests": s.failed_requests,
368                "average_latency_ms": s.average_latency_ms,
369                "url": s.url,
370                "port": s.port,
371                "domain": s.domain,
372            })
373        });
374
375        let config: Option<JsonValue> = match name {
376            "stun" => Some(serde_json::json!({
377                "bind_ip": cfg.bind.ice.ip,
378                "bind_port": cfg.bind.ice.port,
379                "advertised_ip": cfg.bind.ice.advertised_ip,
380                "advertised_port": cfg.bind.ice.advertised_port,
381            })),
382            "turn" => Some(serde_json::json!({
383                "bind_ip": cfg.bind.ice.ip,
384                "bind_port": cfg.bind.ice.port,
385                "advertised_ip": cfg.bind.ice.advertised_ip,
386                "advertised_port": cfg.bind.ice.advertised_port,
387                "relay_port_range": cfg.turn.relay_port_range,
388                "realm": cfg.turn.realm,
389            })),
390            "signaling" => cfg.services.signaling.as_ref().map(|s| {
391                serde_json::json!({
392                    "ws_path": s.server.ws_path,
393                    "rate_limit": {
394                        "connection": {
395                            "enabled": s.server.rate_limit.connection.enabled,
396                            "per_minute": s.server.rate_limit.connection.per_minute,
397                            "burst_size": s.server.rate_limit.connection.burst_size,
398                            "max_concurrent_per_ip": s.server.rate_limit.connection.max_concurrent_per_ip,
399                        },
400                        "message": {
401                            "enabled": s.server.rate_limit.message.enabled,
402                            "per_second": s.server.rate_limit.message.per_second,
403                            "burst_size": s.server.rate_limit.message.burst_size,
404                        }
405                    },
406                    "dependencies": {
407                        "signer": s.dependencies.signer.as_ref().map(|k| serde_json::json!({"endpoint": k.endpoint})),
408                        "ais": s.dependencies.ais.as_ref().map(|a| serde_json::json!({"endpoint": a.endpoint})),
409                    }
410                })
411            }),
412            "ais" => cfg.services.ais.as_ref().map(|a| {
413                serde_json::json!({
414                    "token_ttl_secs": a.server.token_ttl_secs,
415                    "signaling_heartbeat_interval_secs": a.server.signaling_heartbeat_interval_secs,
416                    "dependencies": {
417                        "signer": a.dependencies.signer.as_ref().map(|k| serde_json::json!({"endpoint": k.endpoint})),
418                    }
419                })
420            }),
421            "signer" => cfg.services.signer.as_ref().map(|k| {
422                serde_json::json!({
423                    "storage_backend": format!("{:?}", k.storage.backend),
424                    "key_ttl_seconds": k.storage.key_ttl_seconds,
425                    "tolerance_seconds": k.tolerance_seconds,
426                })
427            }),
428            _ => None,
429        };
430
431        let config_fields = if !name.is_empty() {
432            let overrides = self.overrides_list().await;
433            Some(resolver::resolve_for_service(
434                name,
435                self.toml_content_ref(),
436                &overrides,
437            ))
438        } else {
439            None
440        };
441
442        Ok(ServiceDetail {
443            enabled,
444            status,
445            config,
446            config_fields,
447        })
448    }
449
450    /// Get config registry (all field definitions).
451    pub fn get_registry_direct(&self) -> &'static [registry::ConfigFieldDef] {
452        registry::all_fields()
453    }
454
455    /// Read config file from disk.
456    pub async fn get_config_file_direct(&self) -> GrpcResult<ConfigFileContent> {
457        let path = self
458            .config_path
459            .as_ref()
460            .ok_or_else(|| Status::failed_precondition("Config path not available"))?;
461        let path_display = path.display().to_string();
462        let content = tokio::fs::read_to_string(path)
463            .await
464            .map_err(|e| Status::internal(format!("Failed to read config file: {e}")))?;
465        Ok(ConfigFileContent {
466            content,
467            path: path_display,
468        })
469    }
470
471    /// Write config file to disk (with TOML validation).
472    pub async fn save_config_file_direct(&self, content: &str) -> GrpcResult<()> {
473        let path = self
474            .config_path
475            .as_ref()
476            .ok_or_else(|| Status::failed_precondition("Config path not available"))?;
477
478        // Parse as TOML to validate syntax
479        let new_config = ActrixConfig::from_toml(content)
480            .map_err(|e| Status::invalid_argument(format!("Invalid TOML: {e}")))?;
481
482        // Validate config semantics
483        if let Err(errors) = new_config.validate() {
484            let non_warnings: Vec<_> = errors
485                .iter()
486                .filter(|e| !e.starts_with("Warning:"))
487                .cloned()
488                .collect();
489            if !non_warnings.is_empty() {
490                return Err(Status::invalid_argument(format!(
491                    "Validation failed: {}",
492                    non_warnings.join("; ")
493                )));
494            }
495        }
496
497        tokio::fs::write(path, content.as_bytes())
498            .await
499            .map_err(|e| Status::internal(format!("Failed to write: {e}")))?;
500
501        platform::recording::info!("Config file saved via admin API");
502        Ok(())
503    }
504
505    /// Trigger reload (SIGHUP).
506    pub async fn reload_direct(&self) -> GrpcResult<bool> {
507        platform::recording::info!("Reload requested via admin API, sending SIGHUP to self");
508        #[cfg(unix)]
509        {
510            let pid = std::process::id() as i32;
511            // SAFETY: sending a signal to our own process is safe
512            let ret = unsafe { nix::libc::kill(pid, nix::libc::SIGHUP) };
513            if ret == 0 {
514                Ok(true)
515            } else {
516                Err(Status::internal("Failed to send SIGHUP"))
517            }
518        }
519        #[cfg(not(unix))]
520        {
521            Err(Status::unimplemented(
522                "Reload not supported on this platform",
523            ))
524        }
525    }
526
527    /// Query Signer keys from the SQLite database.
528    pub async fn get_signer_keys_direct(&self) -> GrpcResult<KsKeysResult> {
529        let cfg = self.require_running_config()?;
530        let db_path = cfg.sqlite_path.join("signer_keys.db");
531        if !db_path.exists() {
532            return Ok(KsKeysResult {
533                keys: vec![],
534                total_count: 0,
535            });
536        }
537
538        let url = format!("sqlite:{}?mode=ro", db_path.display());
539        let options = SqliteConnectOptions::from_str(&url)
540            .map_err(|e| Status::internal(format!("DB connect error: {e}")))?
541            .read_only(true);
542        let pool = SqlitePoolOptions::new()
543            .max_connections(1)
544            .connect_with(options)
545            .await
546            .map_err(|e| Status::internal(format!("DB pool error: {e}")))?;
547
548        let now = std::time::SystemTime::now()
549            .duration_since(std::time::UNIX_EPOCH)
550            .unwrap()
551            .as_secs() as i64;
552
553        let rows = sqlx::query_as::<_, (i64, i32, i64, i64)>(
554            "SELECT key_id, length(public_key) as pk_size, created_at, expires_at FROM keys ORDER BY created_at DESC LIMIT 5",
555        )
556        .fetch_all(&pool)
557        .await
558        .map_err(|e| Status::internal(format!("Query error: {e}")))?;
559
560        let total_count = sqlx::query_as::<_, (i64,)>("SELECT COUNT(*) FROM keys")
561            .fetch_one(&pool)
562            .await
563            .map(|r| r.0)
564            .unwrap_or(0);
565
566        let keys = rows
567            .iter()
568            .map(|(key_id, pk_size, created_at, expires_at)| {
569                let is_expired = *expires_at > 0 && *expires_at < now;
570                KeyInfo {
571                    key_id: *key_id,
572                    pk_size: *pk_size,
573                    created_at: Some(*created_at),
574                    fetched_at: None,
575                    expires_at: *expires_at,
576                    tolerance_seconds: None,
577                    is_expired,
578                }
579            })
580            .collect();
581
582        Ok(KsKeysResult { keys, total_count })
583    }
584
585    /// Cleanup expired Signer keys.
586    pub async fn cleanup_signer_keys_direct(&self) -> GrpcResult<KsCleanupResult> {
587        let cfg = self.require_running_config()?;
588        let db_path = cfg.sqlite_path.join("signer_keys.db");
589        if !db_path.exists() {
590            return Ok(KsCleanupResult {
591                deleted: 0,
592                remaining: 0,
593                tolerance_seconds: 0,
594            });
595        }
596
597        let url = format!("sqlite:{}", db_path.display());
598        let options = SqliteConnectOptions::from_str(&url)
599            .map_err(|e| Status::internal(format!("DB connect error: {e}")))?;
600        let pool = SqlitePoolOptions::new()
601            .max_connections(1)
602            .connect_with(options)
603            .await
604            .map_err(|e| Status::internal(format!("DB pool error: {e}")))?;
605
606        let tolerance = cfg
607            .services
608            .signer
609            .as_ref()
610            .map(|k| k.tolerance_seconds)
611            .unwrap_or(3600);
612
613        let now = std::time::SystemTime::now()
614            .duration_since(std::time::UNIX_EPOCH)
615            .unwrap()
616            .as_secs() as i64;
617        let cutoff = now - tolerance as i64;
618
619        let result = sqlx::query("DELETE FROM keys WHERE expires_at > 0 AND expires_at < ?")
620            .bind(cutoff)
621            .execute(&pool)
622            .await
623            .map_err(|e| Status::internal(format!("Cleanup error: {e}")))?;
624
625        let remaining = sqlx::query_as::<_, (i64,)>("SELECT COUNT(*) FROM keys")
626            .fetch_one(&pool)
627            .await
628            .map(|r| r.0)
629            .unwrap_or(0);
630
631        Ok(KsCleanupResult {
632            deleted: result.rows_affected(),
633            remaining,
634            tolerance_seconds: tolerance,
635        })
636    }
637
638    /// Query AIS keys from the SQLite database.
639    pub async fn get_ais_keys_direct(&self) -> GrpcResult<Vec<KeyInfo>> {
640        let cfg = self.require_running_config()?;
641        let db_path = cfg.sqlite_path.join("ais_keys.db");
642        if !db_path.exists() {
643            return Ok(vec![]);
644        }
645
646        let url = format!("sqlite:{}?mode=ro", db_path.display());
647        let options = SqliteConnectOptions::from_str(&url)
648            .map_err(|e| Status::internal(format!("DB connect error: {e}")))?
649            .read_only(true);
650        let pool = SqlitePoolOptions::new()
651            .max_connections(1)
652            .connect_with(options)
653            .await
654            .map_err(|e| Status::internal(format!("DB pool error: {e}")))?;
655
656        let now = std::time::SystemTime::now()
657            .duration_since(std::time::UNIX_EPOCH)
658            .unwrap()
659            .as_secs() as i64;
660
661        let row = sqlx::query_as::<_, (i64, i32, i64, i64, i64)>(
662            "SELECT key_id, length(public_key) as pk_size, fetched_at, expires_at, tolerance_seconds FROM current_key WHERE id = 1",
663        )
664        .fetch_optional(&pool)
665        .await
666        .map_err(|e| Status::internal(format!("Query error: {e}")))?;
667
668        match row {
669            Some((key_id, pk_size, fetched_at, expires_at, tolerance_seconds)) => {
670                let is_expired = expires_at > 0 && expires_at < now;
671                Ok(vec![KeyInfo {
672                    key_id,
673                    pk_size,
674                    created_at: None,
675                    fetched_at: Some(fetched_at),
676                    expires_at,
677                    tolerance_seconds: Some(tolerance_seconds),
678                    is_expired,
679                }])
680            }
681            None => Ok(vec![]),
682        }
683    }
684
685    /// Get node info including metrics and service statuses.
686    pub async fn node_info_direct(&self) -> GrpcResult<GetNodeInfoResponse> {
687        platform::recording::debug!("GetNodeInfo request received");
688
689        let uptime_secs = self.started_at.elapsed().as_secs() as i64;
690        let metrics = self.collect_metrics().await?;
691        let services = self.service_statuses().await;
692
693        Ok(GetNodeInfoResponse {
694            success: true,
695            error_message: None,
696            node_id: self.node_id.clone(),
697            name: self.name.clone(),
698            version: self.version.clone(),
699            location_tag: self.location_tag.clone(),
700            uptime_secs,
701            current_metrics: Some(metrics),
702            services,
703        })
704    }
705
706    /// List all realms.
707    pub async fn list_realms_direct(&self) -> GrpcResult<ListRealmsResponse> {
708        platform::recording::debug!("ListRealms request received");
709
710        let realms = Realm::get_all()
711            .await
712            .map_err(|e| Status::internal(format!("Failed to load realm list: {e}")))?;
713
714        let realms_info: Vec<RealmInfo> = realms.iter().map(realm_to_proto).collect();
715        let total_count = realms_info.len() as u32;
716
717        Ok(ListRealmsResponse {
718            success: true,
719            error_message: None,
720            realms: realms_info,
721            next_page_token: None,
722            total_count,
723        })
724    }
725
726    async fn create_local_realm_internal(
727        &self,
728        req: CreateRealmRequest,
729        expose_plain_secret: bool,
730    ) -> GrpcResult<(CreateRealmResponse, Option<String>)> {
731        platform::recording::info!("Create local realm request received: name={}", req.name);
732
733        // Generate secret at creation time
734        let plain_secret = platform::realm::secret::generate_realm_secret();
735        let secret_hash = hash_realm_secret(&plain_secret);
736
737        let mut realm = match Realm::create(req.name.clone(), secret_hash).await {
738            Ok(r) => r,
739            Err(err) => {
740                return Ok((
741                    CreateRealmResponse {
742                        success: false,
743                        error_message: Some(format!("Failed to create realm: {err}")),
744                        realm: None,
745                    },
746                    None,
747                ));
748            }
749        };
750
751        // Apply optional fields
752        realm.enabled = req.enabled;
753        if req.expires_at > 0 {
754            realm.expires_at = Some(req.expires_at);
755        }
756        if let Some(status) = req.status.as_deref() {
757            match parse_realm_status(status) {
758                Ok(status) => realm.status = status,
759                Err(err) => {
760                    return Ok((
761                        CreateRealmResponse {
762                            success: false,
763                            error_message: Some(err),
764                            realm: None,
765                        },
766                        None,
767                    ));
768                }
769            }
770        }
771        if let Err(err) = realm.save().await {
772            // Rollback
773            let _ = Realm::delete(realm.id).await;
774            return Ok((
775                CreateRealmResponse {
776                    success: false,
777                    error_message: Some(format!("Failed to save realm settings: {err}")),
778                    realm: None,
779                },
780                None,
781            ));
782        }
783
784        platform::recording::info!("Realm created: id={}, name={}", realm.id, realm.name);
785
786        let realm_info = realm_to_proto(&realm);
787
788        Ok((
789            CreateRealmResponse {
790                success: true,
791                error_message: None,
792                realm: Some(realm_info),
793            },
794            if expose_plain_secret {
795                Some(plain_secret)
796            } else {
797                None
798            },
799        ))
800    }
801
802    async fn create_managed_realm_internal(
803        &self,
804        req: CreateRealmRequest,
805    ) -> GrpcResult<CreateRealmResponse> {
806        let Some(realm_id) = req.realm_id else {
807            return Ok(CreateRealmResponse {
808                success: false,
809                error_message: Some("realm_id is required for managed CreateRealm".to_string()),
810                realm: None,
811            });
812        };
813        let Some(secret_current) = req.secret_current_hash else {
814            return Ok(CreateRealmResponse {
815                success: false,
816                error_message: Some(
817                    "secret_current_hash is required for managed CreateRealm".to_string(),
818                ),
819                realm: None,
820            });
821        };
822
823        let status = match req.status.as_deref() {
824            Some(status) => match parse_realm_status(status) {
825                Ok(status) => status,
826                Err(err) => {
827                    return Ok(CreateRealmResponse {
828                        success: false,
829                        error_message: Some(err),
830                        realm: None,
831                    });
832                }
833            },
834            None => RealmStatus::Active,
835        };
836        let secret_previous =
837            build_secret_previous(req.secret_previous_hash, req.secret_previous_valid_until);
838
839        platform::recording::info!(
840            "Create managed realm request received: id={}, name={}",
841            realm_id,
842            req.name
843        );
844
845        match Realm::upsert_managed(
846            realm_id,
847            req.name,
848            status,
849            req.enabled,
850            (req.expires_at > 0).then_some(req.expires_at),
851            secret_current,
852            secret_previous,
853        )
854        .await
855        {
856            Ok(realm) => Ok(CreateRealmResponse {
857                success: true,
858                error_message: None,
859                realm: Some(realm_to_proto(&realm)),
860            }),
861            Err(err) => Ok(CreateRealmResponse {
862                success: false,
863                error_message: Some(format!("Failed to upsert managed realm: {err}")),
864                realm: None,
865            }),
866        }
867    }
868
869    /// Create a realm and return the plaintext realm secret once.
870    pub async fn create_realm_with_secret_direct(
871        &self,
872        req: CreateRealmRequest,
873    ) -> GrpcResult<(CreateRealmResponse, Option<String>)> {
874        self.create_local_realm_internal(req, true).await
875    }
876
877    /// Create a realm (gRPC path, does not expose plaintext secret in response).
878    pub async fn create_realm_direct(
879        &self,
880        req: CreateRealmRequest,
881    ) -> GrpcResult<CreateRealmResponse> {
882        self.create_managed_realm_internal(req).await
883    }
884
885    /// Get a single realm by ID.
886    pub async fn get_realm_direct(&self, realm_id: u32) -> GrpcResult<GetRealmResponse> {
887        platform::recording::debug!("GetRealm request received: realm_id={}", realm_id);
888
889        match self.load_realm(realm_id).await {
890            Ok(realm) => Ok(GetRealmResponse {
891                success: true,
892                error_message: None,
893                realm: Some(realm_to_proto(&realm)),
894            }),
895            Err(status) if status.code() == tonic::Code::NotFound => Ok(GetRealmResponse {
896                success: false,
897                error_message: Some(status.message().to_string()),
898                realm: None,
899            }),
900            Err(e) => Err(e),
901        }
902    }
903
904    /// Update an existing realm.
905    pub async fn update_realm_direct(
906        &self,
907        req: UpdateRealmRequest,
908    ) -> GrpcResult<UpdateRealmResponse> {
909        let mut realm = match self.load_realm(req.realm_id).await {
910            Ok(r) => r,
911            Err(status) if status.code() == tonic::Code::NotFound => {
912                return Ok(UpdateRealmResponse {
913                    success: false,
914                    error_message: Some(status.message().to_string()),
915                    realm: None,
916                });
917            }
918            Err(e) => return Err(e),
919        };
920
921        if let Some(name) = req.name {
922            realm.name = name;
923        }
924        if let Some(enabled) = req.enabled {
925            realm.enabled = enabled;
926        }
927        if let Some(status) = req.status.as_deref() {
928            match parse_realm_status(status) {
929                Ok(status) => realm.status = status,
930                Err(err) => {
931                    return Ok(UpdateRealmResponse {
932                        success: false,
933                        error_message: Some(err),
934                        realm: None,
935                    });
936                }
937            }
938        }
939        if let Some(expires_at) = req.expires_at {
940            realm.expires_at = (expires_at > 0).then_some(expires_at);
941        }
942        if let Some(secret_current_hash) = req.secret_current_hash {
943            if secret_current_hash.trim().is_empty() {
944                return Ok(UpdateRealmResponse {
945                    success: false,
946                    error_message: Some("secret_current_hash must not be empty".to_string()),
947                    realm: None,
948                });
949            }
950            realm.secret_current = secret_current_hash;
951        }
952        if req.secret_previous_hash.is_some() || req.secret_previous_valid_until.is_some() {
953            realm.secret_previous =
954                build_secret_previous(req.secret_previous_hash, req.secret_previous_valid_until);
955        }
956
957        if let Err(err) = realm.save().await {
958            return Ok(UpdateRealmResponse {
959                success: false,
960                error_message: Some(format!("Failed to update realm: {err}")),
961                realm: None,
962            });
963        }
964
965        Ok(UpdateRealmResponse {
966            success: true,
967            error_message: None,
968            realm: Some(realm_to_proto(&realm)),
969        })
970    }
971
972    /// Delete a realm by ID.
973    pub async fn delete_realm_direct(&self, realm_id: u32) -> GrpcResult<DeleteRealmResponse> {
974        match Realm::soft_delete(realm_id).await {
975            Ok(true) => Ok(DeleteRealmResponse {
976                success: true,
977                error_message: None,
978            }),
979            Ok(false) => Ok(DeleteRealmResponse {
980                success: false,
981                error_message: Some("Realm not found".to_string()),
982            }),
983            Err(err) => Ok(DeleteRealmResponse {
984                success: false,
985                error_message: Some(format!("Failed to delete realm: {err}")),
986            }),
987        }
988    }
989
990    /// Hard-delete a realm for local Admin UI standalone mode.
991    pub async fn delete_realm_hard_direct(&self, realm_id: u32) -> GrpcResult<DeleteRealmResponse> {
992        match Realm::delete(realm_id).await {
993            Ok(affected) if affected > 0 => Ok(DeleteRealmResponse {
994                success: true,
995                error_message: None,
996            }),
997            Ok(_) => Ok(DeleteRealmResponse {
998                success: false,
999                error_message: Some("Realm not found".to_string()),
1000            }),
1001            Err(err) => Ok(DeleteRealmResponse {
1002                success: false,
1003                error_message: Some(format!("Failed to delete realm: {err}")),
1004            }),
1005        }
1006    }
1007
1008    /// Rotate realm secret and return plaintext once.
1009    pub async fn rotate_realm_secret_direct(
1010        &self,
1011        realm_id: u32,
1012    ) -> GrpcResult<RealmSecretRotationResult> {
1013        let rotated = rotate_realm_secret(realm_id, Some(DEFAULT_REALM_SECRET_PREVIOUS_GRACE_SECS))
1014            .await
1015            .map_err(|e| Status::internal(format!("Failed to rotate realm secret: {e}")))?;
1016
1017        Ok(RealmSecretRotationResult {
1018            realm_id,
1019            realm_secret: rotated.new_secret,
1020            previous_valid_until: rotated.previous_valid_until,
1021            grace_seconds: DEFAULT_REALM_SECRET_PREVIOUS_GRACE_SECS,
1022        })
1023    }
1024
1025    /// Get a config value.
1026    pub async fn get_config_direct(
1027        &self,
1028        config_type: ConfigType,
1029        config_key: String,
1030    ) -> GrpcResult<GetConfigResponse> {
1031        let key = Self::build_config_key(config_type, config_key);
1032        let store = self.config_store.read().await;
1033
1034        if let Some(value) = store.get(&key) {
1035            Ok(GetConfigResponse {
1036                success: true,
1037                error_message: None,
1038                config_value: Some(value.clone()),
1039            })
1040        } else {
1041            Ok(GetConfigResponse {
1042                success: false,
1043                error_message: Some("Config not found".to_string()),
1044                config_value: None,
1045            })
1046        }
1047    }
1048
1049    /// Update a config value.
1050    pub async fn update_config_direct(
1051        &self,
1052        config_type: ConfigType,
1053        config_key: String,
1054        config_value: String,
1055    ) -> GrpcResult<UpdateConfigResponse> {
1056        let key = Self::build_config_key(config_type, config_key);
1057        let mut store = self.config_store.write().await;
1058        let old_value = store.insert(key, config_value);
1059
1060        Ok(UpdateConfigResponse {
1061            success: true,
1062            error_message: None,
1063            old_value,
1064        })
1065    }
1066
1067    /// Request node shutdown.
1068    pub async fn shutdown_direct(
1069        &self,
1070        graceful: bool,
1071        timeout_secs: Option<i32>,
1072        reason: Option<String>,
1073    ) -> GrpcResult<ShutdownResponse> {
1074        if let Some(handler) = &self.shutdown_handler {
1075            if let Err(e) = handler(graceful, timeout_secs, reason.clone()).await {
1076                return Ok(ShutdownResponse {
1077                    accepted: false,
1078                    error_message: Some(format!("Shutdown handler failed: {e}")),
1079                    estimated_shutdown_time: None,
1080                });
1081            }
1082        } else {
1083            platform::recording::warn!("Shutdown requested but no handler registered");
1084        }
1085
1086        let estimated = if graceful {
1087            timeout_secs.map(|v| Utc::now().timestamp() + v as i64)
1088        } else {
1089            Some(Utc::now().timestamp())
1090        };
1091
1092        Ok(ShutdownResponse {
1093            accepted: true,
1094            error_message: None,
1095            estimated_shutdown_time: estimated,
1096        })
1097    }
1098}
1099
1100// ── Response structs (transport-agnostic, Serialize) ─────────────
1101
1102/// Health check response.
1103#[derive(Debug, Clone, Serialize)]
1104pub struct HealthInfo {
1105    pub status: String,
1106    pub node: String,
1107    pub version: String,
1108}
1109
1110/// Platform detail with resolved config fields.
1111#[derive(Debug, Clone, Serialize)]
1112pub struct PlatformDetail {
1113    pub config: JsonValue,
1114    pub config_fields: Vec<ResolvedField>,
1115}
1116
1117/// Service detail with resolved config fields.
1118#[derive(Debug, Clone, Serialize)]
1119pub struct ServiceDetail {
1120    pub enabled: bool,
1121    pub status: Option<JsonValue>,
1122    pub config: Option<JsonValue>,
1123    pub config_fields: Option<Vec<ResolvedField>>,
1124}
1125
1126/// Config file content with path.
1127#[derive(Debug, Clone, Serialize)]
1128pub struct ConfigFileContent {
1129    pub content: String,
1130    pub path: String,
1131}
1132
1133/// Key info (shared by Signer and AIS).
1134#[derive(Debug, Clone, Serialize)]
1135pub struct KeyInfo {
1136    pub key_id: i64,
1137    pub pk_size: i32,
1138    #[serde(skip_serializing_if = "Option::is_none")]
1139    pub created_at: Option<i64>,
1140    #[serde(skip_serializing_if = "Option::is_none")]
1141    pub fetched_at: Option<i64>,
1142    pub expires_at: i64,
1143    #[serde(skip_serializing_if = "Option::is_none")]
1144    pub tolerance_seconds: Option<i64>,
1145    pub is_expired: bool,
1146}
1147
1148/// Realm secret 轮转结果(明文 secret 仅返回一次)。
1149#[derive(Debug, Clone, Serialize)]
1150pub struct RealmSecretRotationResult {
1151    pub realm_id: u32,
1152    pub realm_secret: String,
1153    #[serde(skip_serializing_if = "Option::is_none")]
1154    pub previous_valid_until: Option<u64>,
1155    pub grace_seconds: u64,
1156}
1157
1158/// Signer keys query result.
1159#[derive(Debug, Clone, Serialize)]
1160pub struct KsKeysResult {
1161    pub keys: Vec<KeyInfo>,
1162    pub total_count: i64,
1163}
1164
1165/// Signer cleanup result.
1166#[derive(Debug, Clone, Serialize)]
1167pub struct KsCleanupResult {
1168    pub deleted: u64,
1169    pub remaining: i64,
1170    pub tolerance_seconds: u64,
1171}
1172
1173/// Map service name string to the ResourceType integer used in ServiceStatus.type
1174fn service_name_to_type_id(name: &str) -> Option<i32> {
1175    match name {
1176        "stun" => Some(1),
1177        "turn" => Some(2),
1178        "signaling" => Some(3),
1179        "ais" => Some(4),
1180        "signer" => Some(5),
1181        _ => None,
1182    }
1183}
1184
1185fn parse_realm_status(value: &str) -> std::result::Result<RealmStatus, String> {
1186    RealmStatus::from_str(value).map_err(|_| {
1187        format!("Invalid realm status '{value}' (expected Active, Inactive, Suspended)")
1188    })
1189}
1190
1191fn build_secret_previous(hash: Option<String>, valid_until: Option<u64>) -> Option<(String, u64)> {
1192    match (hash, valid_until) {
1193        (Some(hash), Some(valid_until)) if !hash.trim().is_empty() && valid_until > 0 => {
1194            Some((hash, valid_until))
1195        }
1196        _ => None,
1197    }
1198}
1199
1200#[tonic::async_trait]
1201impl NodeAdminService for AdminApiService {
1202    async fn update_config(
1203        &self,
1204        request: Request<UpdateConfigRequest>,
1205    ) -> GrpcResult<Response<UpdateConfigResponse>> {
1206        let req = request.into_inner();
1207        self.update_config_direct(req.config_type(), req.config_key, req.config_value)
1208            .await
1209            .map(Response::new)
1210    }
1211
1212    async fn get_config(
1213        &self,
1214        request: Request<GetConfigRequest>,
1215    ) -> GrpcResult<Response<GetConfigResponse>> {
1216        let req = request.into_inner();
1217        self.get_config_direct(req.config_type(), req.config_key)
1218            .await
1219            .map(Response::new)
1220    }
1221
1222    async fn create_realm(
1223        &self,
1224        request: Request<CreateRealmRequest>,
1225    ) -> GrpcResult<Response<CreateRealmResponse>> {
1226        self.create_realm_direct(request.into_inner())
1227            .await
1228            .map(Response::new)
1229    }
1230
1231    async fn get_realm(
1232        &self,
1233        request: Request<GetRealmRequest>,
1234    ) -> GrpcResult<Response<GetRealmResponse>> {
1235        let req = request.into_inner();
1236        self.get_realm_direct(req.realm_id).await.map(Response::new)
1237    }
1238
1239    async fn update_realm(
1240        &self,
1241        request: Request<UpdateRealmRequest>,
1242    ) -> GrpcResult<Response<UpdateRealmResponse>> {
1243        self.update_realm_direct(request.into_inner())
1244            .await
1245            .map(Response::new)
1246    }
1247
1248    async fn delete_realm(
1249        &self,
1250        request: Request<DeleteRealmRequest>,
1251    ) -> GrpcResult<Response<DeleteRealmResponse>> {
1252        let req = request.into_inner();
1253        self.delete_realm_direct(req.realm_id)
1254            .await
1255            .map(Response::new)
1256    }
1257
1258    async fn list_realms(
1259        &self,
1260        request: Request<ListRealmsRequest>,
1261    ) -> GrpcResult<Response<ListRealmsResponse>> {
1262        let _req = request.into_inner();
1263        self.list_realms_direct().await.map(Response::new)
1264    }
1265
1266    async fn get_node_info(
1267        &self,
1268        request: Request<GetNodeInfoRequest>,
1269    ) -> GrpcResult<Response<GetNodeInfoResponse>> {
1270        let _req = request.into_inner();
1271        self.node_info_direct().await.map(Response::new)
1272    }
1273
1274    async fn shutdown(
1275        &self,
1276        request: Request<ShutdownRequest>,
1277    ) -> GrpcResult<Response<ShutdownResponse>> {
1278        let req = request.into_inner();
1279        self.shutdown_direct(req.graceful, req.timeout_secs, req.reason)
1280            .await
1281            .map(Response::new)
1282    }
1283
1284    async fn list_config_overrides(
1285        &self,
1286        _request: Request<ListConfigOverridesRequest>,
1287    ) -> GrpcResult<Response<ListConfigOverridesResponse>> {
1288        match self.list_overrides_direct().await {
1289            Ok(overrides) => {
1290                let entries = overrides
1291                    .into_iter()
1292                    .map(|o| ProtoConfigOverrideEntry {
1293                        key_path: o.key_path,
1294                        value: o.value,
1295                        updated_at: o.updated_at,
1296                        updated_by: o.updated_by,
1297                    })
1298                    .collect();
1299                Ok(Response::new(ListConfigOverridesResponse {
1300                    success: true,
1301                    error_message: None,
1302                    overrides: entries,
1303                }))
1304            }
1305            Err(e) => Ok(Response::new(ListConfigOverridesResponse {
1306                success: false,
1307                error_message: Some(e.message().to_string()),
1308                overrides: vec![],
1309            })),
1310        }
1311    }
1312
1313    async fn set_config_override(
1314        &self,
1315        request: Request<SetConfigOverrideRequest>,
1316    ) -> GrpcResult<Response<SetConfigOverrideResponse>> {
1317        let req = request.into_inner();
1318        let by = req.updated_by.as_deref().unwrap_or("admin");
1319        match self.set_override_direct(&req.key, &req.value, by).await {
1320            Ok(()) => Ok(Response::new(SetConfigOverrideResponse {
1321                success: true,
1322                error_message: None,
1323            })),
1324            Err(e) => Ok(Response::new(SetConfigOverrideResponse {
1325                success: false,
1326                error_message: Some(e.message().to_string()),
1327            })),
1328        }
1329    }
1330
1331    async fn delete_config_override(
1332        &self,
1333        request: Request<DeleteConfigOverrideRequest>,
1334    ) -> GrpcResult<Response<DeleteConfigOverrideResponse>> {
1335        let req = request.into_inner();
1336        match self.delete_override_direct(&req.key).await {
1337            Ok(deleted) => Ok(Response::new(DeleteConfigOverrideResponse {
1338                success: true,
1339                error_message: None,
1340                deleted,
1341            })),
1342            Err(e) => Ok(Response::new(DeleteConfigOverrideResponse {
1343                success: false,
1344                error_message: Some(e.message().to_string()),
1345                deleted: false,
1346            })),
1347        }
1348    }
1349}