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#[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 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 running_config: Option<ActrixConfig>,
75 toml_content: Option<String>,
77 config_path: Option<PathBuf>,
79}
80
81impl AdminApiService {
82 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 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 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 pub fn with_override_store(mut self, store: Arc<ConfigOverrideStore>) -> Self {
136 self.override_store = Some(store);
137 self
138 }
139
140 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 pub fn with_running_config(mut self, config: ActrixConfig) -> Self {
155 self.running_config = Some(config);
156 self
157 }
158
159 pub fn with_toml_content(mut self, content: String) -> Self {
161 self.toml_content = Some(content);
162 self
163 }
164
165 pub fn with_config_path(mut self, path: PathBuf) -> Self {
167 self.config_path = Some(path);
168 self
169 }
170
171 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 pub async fn service_statuses(&self) -> Vec<ServiceStatus> {
201 self.service_collector.all_statuses().await
202 }
203
204 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 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 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 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 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 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 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 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 pub fn get_registry_direct(&self) -> &'static [registry::ConfigFieldDef] {
452 registry::all_fields()
453 }
454
455 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 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 let new_config = ActrixConfig::from_toml(content)
480 .map_err(|e| Status::invalid_argument(format!("Invalid TOML: {e}")))?;
481
482 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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#[derive(Debug, Clone, Serialize)]
1104pub struct HealthInfo {
1105 pub status: String,
1106 pub node: String,
1107 pub version: String,
1108}
1109
1110#[derive(Debug, Clone, Serialize)]
1112pub struct PlatformDetail {
1113 pub config: JsonValue,
1114 pub config_fields: Vec<ResolvedField>,
1115}
1116
1117#[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#[derive(Debug, Clone, Serialize)]
1128pub struct ConfigFileContent {
1129 pub content: String,
1130 pub path: String,
1131}
1132
1133#[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#[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#[derive(Debug, Clone, Serialize)]
1160pub struct KsKeysResult {
1161 pub keys: Vec<KeyInfo>,
1162 pub total_count: i64,
1163}
1164
1165#[derive(Debug, Clone, Serialize)]
1167pub struct KsCleanupResult {
1168 pub deleted: u64,
1169 pub remaining: i64,
1170 pub tolerance_seconds: u64,
1171}
1172
1173fn 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}