1use athena_rs::client::backend::QueryResult;
2use athena_rs::AthenaClient;
3use chrono::{DateTime, Local, Utc};
4use once_cell::sync::Lazy;
5use serde_json::{json, Value};
6use std::collections::HashSet;
7use std::env;
8use std::future::Future;
9use std::path::{Path, PathBuf};
10use tokio::sync::Mutex;
11use tracing::warn;
12
13use crate::logging::{LogEntry, LogLevel};
14use crate::strategies::{DatabaseConfig, ServiceConfig, XbpConfig};
15use crate::utils::{find_xbp_config_upwards, parse_config_with_auto_heal};
16
17const DB_NOT_CONFIGURED: &str = "athena database is not configured";
18const DEFAULT_BACKEND: &str = "supabase";
19const DEFAULT_SCHEMA: &str = "public";
20const BOOTSTRAP_SQL: &str = include_str!("../../sql/schema.sql");
21
22static BOOTSTRAPPED_BACKENDS: Lazy<Mutex<HashSet<String>>> =
23 Lazy::new(|| Mutex::new(HashSet::new()));
24
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub struct AthenaRuntimeConfig {
27 pub backend: String,
28 pub url: String,
29 pub key: String,
30 pub schema: String,
31}
32
33#[derive(Debug, Clone)]
34struct ProjectContext {
35 project_root: PathBuf,
36 config: XbpConfig,
37 config_kind: String,
38}
39
40pub async fn persist_project_snapshot(
41 project_root: &Path,
42 config: &XbpConfig,
43 config_kind: Option<&str>,
44) {
45 let _ = with_fail_open(
46 "persist project snapshot",
47 persist_project_snapshot_inner(project_root, config, config_kind),
48 )
49 .await;
50}
51
52pub async fn persist_log_entry(entry: &LogEntry) {
53 let _ = with_fail_open("persist xbp log entry", persist_log_entry_inner(entry)).await;
54}
55
56pub async fn persist_nginx_config_snapshot(
57 domain: &str,
58 config_path: &Path,
59 content: &str,
60 upstream_ports: &[u16],
61 listen_ports: &[u16],
62 source: &str,
63) {
64 let _ = with_fail_open(
65 "persist nginx config snapshot",
66 persist_nginx_config_snapshot_inner(
67 domain,
68 config_path,
69 content,
70 upstream_ports,
71 listen_ports,
72 source,
73 ),
74 )
75 .await;
76}
77
78pub async fn persist_nginx_log(
79 domain: Option<&str>,
80 action: &str,
81 success: bool,
82 message: &str,
83 details: Option<&str>,
84 metadata: Value,
85) {
86 let _ = with_fail_open(
87 "persist nginx log",
88 persist_nginx_log_inner(domain, action, success, message, details, metadata),
89 )
90 .await;
91}
92
93pub async fn persist_nginx_edit_audit_log(
94 domain: Option<&str>,
95 config_path: Option<&Path>,
96 actor: Option<&str>,
97 action: &str,
98 old_content: Option<&str>,
99 new_content: Option<&str>,
100 metadata: Value,
101) {
102 let _ = with_fail_open(
103 "persist nginx edit audit log",
104 persist_nginx_edit_audit_log_inner(
105 domain,
106 config_path,
107 actor,
108 action,
109 old_content,
110 new_content,
111 metadata,
112 ),
113 )
114 .await;
115}
116
117pub async fn persist_docker_container_snapshot(
118 container_id: &str,
119 container_name: &str,
120 status: Option<&str>,
121 ports: Option<&str>,
122 metadata: Value,
123) {
124 let _ = with_fail_open(
125 "persist docker container snapshot",
126 persist_docker_container_snapshot_inner(
127 container_id,
128 container_name,
129 status,
130 ports,
131 metadata,
132 ),
133 )
134 .await;
135}
136
137pub async fn persist_docker_log(
138 container_id: Option<&str>,
139 command: Option<&str>,
140 stream: &str,
141 message: &str,
142 metadata: Value,
143) {
144 let _ = with_fail_open(
145 "persist docker log",
146 persist_docker_log_inner(container_id, command, stream, message, metadata),
147 )
148 .await;
149}
150
151pub async fn persist_schedule(
152 schedule_type: &str,
153 target_kind: &str,
154 target_ref: Option<&str>,
155 expression: &str,
156 enabled: bool,
157 metadata: Value,
158) {
159 let _ = with_fail_open(
160 "persist schedule",
161 persist_schedule_inner(
162 schedule_type,
163 target_kind,
164 target_ref,
165 expression,
166 enabled,
167 metadata,
168 ),
169 )
170 .await;
171}
172
173pub fn extract_cron_restart_expression(args: &[String]) -> Option<String> {
174 for (index, arg) in args.iter().enumerate() {
175 if arg == "--cron-restart" {
176 if let Some(value) = args.get(index + 1) {
177 let value = value.trim();
178 if !value.is_empty() {
179 return Some(value.to_string());
180 }
181 }
182 } else if let Some(value) = arg.strip_prefix("--cron-restart=") {
183 let value = value.trim();
184 if !value.is_empty() {
185 return Some(value.to_string());
186 }
187 }
188 }
189 None
190}
191
192pub fn resolve_runtime_config(config: Option<&XbpConfig>) -> Option<AthenaRuntimeConfig> {
193 let database = config.and_then(|cfg| cfg.database.as_ref());
194 let enabled = database.and_then(|db| db.enabled).unwrap_or(true);
195 if !enabled {
196 return None;
197 }
198
199 let backend = value_or_default(
200 database.and_then(|db| db.backend.as_deref()),
201 env::var("XBP_ATHENA_BACKEND").ok().as_deref(),
202 DEFAULT_BACKEND,
203 );
204
205 let (url, key) = resolve_connection_pair(database)?;
206 let schema = sanitize_identifier(
207 value_or_default(
208 database.and_then(|db| db.schema.as_deref()),
209 None,
210 DEFAULT_SCHEMA,
211 )
212 .as_str(),
213 )
214 .unwrap_or_else(|| DEFAULT_SCHEMA.to_string());
215
216 Some(AthenaRuntimeConfig {
217 backend,
218 url,
219 key,
220 schema,
221 })
222}
223
224async fn persist_project_snapshot_inner(
225 project_root: &Path,
226 config: &XbpConfig,
227 config_kind: Option<&str>,
228) -> Result<(), String> {
229 let (client, runtime) = initialize_client(Some(config)).await?;
230 let _ =
231 upsert_project_snapshot_with_client(&client, &runtime, project_root, config, config_kind)
232 .await?;
233 Ok(())
234}
235const ATHENA_LOG_DETAILS_MAX_CHARS: usize = 32 * 1024;
238
239fn truncate_log_details(details: Option<&str>) -> Option<String> {
240 let details = details?;
241 let char_count = details.chars().count();
242 if char_count <= ATHENA_LOG_DETAILS_MAX_CHARS {
243 return Some(details.to_string());
244 }
245 let omitted = char_count.saturating_sub(ATHENA_LOG_DETAILS_MAX_CHARS);
246 let head: String = details.chars().take(ATHENA_LOG_DETAILS_MAX_CHARS / 2).collect();
247 let tail: String = details
248 .chars()
249 .rev()
250 .take(ATHENA_LOG_DETAILS_MAX_CHARS / 2)
251 .collect::<String>()
252 .chars()
253 .rev()
254 .collect();
255 Some(format!(
256 "{head}\n\n… truncated {omitted} characters for Athena log persistence …\n\n{tail}"
257 ))
258}
259
260async fn persist_log_entry_inner(entry: &LogEntry) -> Result<(), String> {
261 let context = load_current_project_context();
262 let config_ref = context.as_ref().map(|ctx| &ctx.config);
263 let (client, runtime) = initialize_client(config_ref).await?;
264
265 let project_id = if let Some(ctx) = context.as_ref() {
266 Some(
267 upsert_project_snapshot_with_client(
268 &client,
269 &runtime,
270 &ctx.project_root,
271 &ctx.config,
272 Some(ctx.config_kind.as_str()),
273 )
274 .await?,
275 )
276 } else {
277 None
278 };
279
280 let log_table = qualified_table(&runtime.schema, "xbp_logs");
281 let timestamp = to_rfc3339(entry.timestamp);
282 let metadata = json!({
283 "project_name": context.as_ref().map(|ctx| ctx.config.project_name.clone()),
284 "project_path": context.as_ref().map(|ctx| ctx.project_root.display().to_string()),
285 });
286 let details = truncate_log_details(entry.details.as_deref());
287
288 let global_sql = format!(
289 "INSERT INTO {table} (log_level, command, message, details, duration_ms, occurred_at, metadata, updated_at) \
290 VALUES ({log_level}, {command}, {message}, {details}, {duration}, {occurred_at}::timestamptz, {metadata}, now()) \
291 RETURNING log_id",
292 table = log_table,
293 log_level = sql_literal(&Value::String(log_level_label(&entry.level).to_string())),
294 command = sql_literal(&Value::String(entry.command.clone())),
295 message = sql_literal(&Value::String(entry.message.clone())),
296 details = optional_text_literal(details.as_deref()),
297 duration = optional_u64_literal(entry.duration_ms),
298 occurred_at = sql_literal(&Value::String(timestamp)),
299 metadata = sql_literal(&metadata),
300 );
301
302 let global_result = execute_sql(&client, &global_sql).await?;
303 let global_log_id = query_first_column_as_string(&global_result, "log_id");
304
305 if let (Some(project_id), Some(global_log_id)) =
306 (project_id.as_deref(), global_log_id.as_deref())
307 {
308 let project_table = qualified_table(&runtime.schema, "xbp_project_logs");
309 let project_sql = format!(
310 "INSERT INTO {table} (project_id, global_log_id, log_level, command, message, details, duration_ms, occurred_at, metadata, updated_at) \
311 VALUES ({project_id}::uuid, {global_log_id}::uuid, {log_level}, {command}, {message}, {details}, {duration}, {occurred_at}::timestamptz, {metadata}, now())",
312 table = project_table,
313 project_id = sql_literal(&Value::String(project_id.to_string())),
314 global_log_id = sql_literal(&Value::String(global_log_id.to_string())),
315 log_level = sql_literal(&Value::String(log_level_label(&entry.level).to_string())),
316 command = sql_literal(&Value::String(entry.command.clone())),
317 message = sql_literal(&Value::String(entry.message.clone())),
318 details = optional_text_literal(details.as_deref()),
319 duration = optional_u64_literal(entry.duration_ms),
320 occurred_at = sql_literal(&Value::String(to_rfc3339(entry.timestamp))),
321 metadata = sql_literal(&metadata),
322 );
323 let _ = execute_sql(&client, &project_sql).await?;
324 }
325
326 Ok(())
327}
328
329async fn persist_nginx_config_snapshot_inner(
330 domain: &str,
331 config_path: &Path,
332 content: &str,
333 upstream_ports: &[u16],
334 listen_ports: &[u16],
335 source: &str,
336) -> Result<(), String> {
337 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
338 let table = qualified_table(&runtime.schema, "xbp_nginx_configs");
339 let metadata = json!({ "source": source });
340 let upstream = Value::Array(
341 upstream_ports
342 .iter()
343 .map(|port| Value::from(*port as u64))
344 .collect(),
345 );
346 let listen = Value::Array(
347 listen_ports
348 .iter()
349 .map(|port| Value::from(*port as u64))
350 .collect(),
351 );
352
353 let sql = format!(
354 "INSERT INTO {table} (project_id, domain, config_path, content, upstream_ports, listen_ports, metadata, updated_at) \
355 VALUES ({project_id}, {domain}, {config_path}, {content}, {upstream_ports}, {listen_ports}, {metadata}, now()) \
356 ON CONFLICT (domain, config_path) DO UPDATE SET \
357 project_id = EXCLUDED.project_id, \
358 content = EXCLUDED.content, \
359 upstream_ports = EXCLUDED.upstream_ports, \
360 listen_ports = EXCLUDED.listen_ports, \
361 metadata = EXCLUDED.metadata, \
362 updated_at = now()",
363 table = table,
364 project_id = optional_uuid_literal(project_id.as_deref()),
365 domain = sql_literal(&Value::String(domain.to_string())),
366 config_path = sql_literal(&Value::String(config_path.display().to_string())),
367 content = sql_literal(&Value::String(content.to_string())),
368 upstream_ports = sql_literal(&upstream),
369 listen_ports = sql_literal(&listen),
370 metadata = sql_literal(&metadata),
371 );
372
373 let _ = execute_sql(&client, &sql).await?;
374 Ok(())
375}
376
377async fn persist_nginx_log_inner(
378 domain: Option<&str>,
379 action: &str,
380 success: bool,
381 message: &str,
382 details: Option<&str>,
383 metadata: Value,
384) -> Result<(), String> {
385 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
386 let table = qualified_table(&runtime.schema, "xbp_nginx_logs");
387
388 let sql = format!(
389 "INSERT INTO {table} (project_id, domain, action, success, message, details, occurred_at, metadata, updated_at) \
390 VALUES ({project_id}, {domain}, {action}, {success}, {message}, {details}, now(), {metadata}, now())",
391 table = table,
392 project_id = optional_uuid_literal(project_id.as_deref()),
393 domain = optional_text_literal(domain),
394 action = sql_literal(&Value::String(action.to_string())),
395 success = if success { "true" } else { "false" },
396 message = sql_literal(&Value::String(message.to_string())),
397 details = optional_text_literal(details),
398 metadata = sql_literal(&metadata),
399 );
400
401 let _ = execute_sql(&client, &sql).await?;
402 Ok(())
403}
404
405async fn persist_nginx_edit_audit_log_inner(
406 domain: Option<&str>,
407 config_path: Option<&Path>,
408 actor: Option<&str>,
409 action: &str,
410 old_content: Option<&str>,
411 new_content: Option<&str>,
412 metadata: Value,
413) -> Result<(), String> {
414 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
415 let table = qualified_table(&runtime.schema, "xbp_nginx_edit_audit_logs");
416
417 let sql = format!(
418 "INSERT INTO {table} (project_id, domain, config_path, actor, action, old_content, new_content, occurred_at, metadata, updated_at) \
419 VALUES ({project_id}, {domain}, {config_path}, {actor}, {action}, {old_content}, {new_content}, now(), {metadata}, now())",
420 table = table,
421 project_id = optional_uuid_literal(project_id.as_deref()),
422 domain = optional_text_literal(domain),
423 config_path = config_path
424 .map(|path| sql_literal(&Value::String(path.display().to_string())))
425 .unwrap_or_else(|| "NULL".to_string()),
426 actor = optional_text_literal(actor),
427 action = sql_literal(&Value::String(action.to_string())),
428 old_content = optional_text_literal(old_content),
429 new_content = optional_text_literal(new_content),
430 metadata = sql_literal(&metadata),
431 );
432
433 let _ = execute_sql(&client, &sql).await?;
434 Ok(())
435}
436
437async fn persist_docker_container_snapshot_inner(
438 container_id: &str,
439 container_name: &str,
440 status: Option<&str>,
441 ports: Option<&str>,
442 metadata: Value,
443) -> Result<(), String> {
444 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
445 let table = qualified_table(&runtime.schema, "xbp_docker_containers");
446 let sql = format!(
447 "INSERT INTO {table} (project_id, container_id, container_name, status, ports, inspected_at, metadata, updated_at) \
448 VALUES ({project_id}, {container_id}, {container_name}, {status}, {ports}, now(), {metadata}, now()) \
449 ON CONFLICT (container_id) DO UPDATE SET \
450 project_id = EXCLUDED.project_id, \
451 container_name = EXCLUDED.container_name, \
452 status = EXCLUDED.status, \
453 ports = EXCLUDED.ports, \
454 inspected_at = EXCLUDED.inspected_at, \
455 metadata = EXCLUDED.metadata, \
456 updated_at = now()",
457 table = table,
458 project_id = optional_uuid_literal(project_id.as_deref()),
459 container_id = sql_literal(&Value::String(container_id.to_string())),
460 container_name = sql_literal(&Value::String(container_name.to_string())),
461 status = optional_text_literal(status),
462 ports = optional_text_literal(ports),
463 metadata = sql_literal(&metadata),
464 );
465
466 let _ = execute_sql(&client, &sql).await?;
467 Ok(())
468}
469async fn persist_docker_log_inner(
470 container_id: Option<&str>,
471 command: Option<&str>,
472 stream: &str,
473 message: &str,
474 metadata: Value,
475) -> Result<(), String> {
476 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
477 let table = qualified_table(&runtime.schema, "xbp_docker_logs");
478 let sql = format!(
479 "INSERT INTO {table} (project_id, container_id, command, stream, message, occurred_at, metadata, updated_at) \
480 VALUES ({project_id}, {container_id}, {command}, {stream}, {message}, now(), {metadata}, now())",
481 table = table,
482 project_id = optional_uuid_literal(project_id.as_deref()),
483 container_id = optional_text_literal(container_id),
484 command = optional_text_literal(command),
485 stream = sql_literal(&Value::String(stream.to_string())),
486 message = sql_literal(&Value::String(message.to_string())),
487 metadata = sql_literal(&metadata),
488 );
489
490 let _ = execute_sql(&client, &sql).await?;
491 Ok(())
492}
493
494async fn persist_schedule_inner(
495 schedule_type: &str,
496 target_kind: &str,
497 target_ref: Option<&str>,
498 expression: &str,
499 enabled: bool,
500 metadata: Value,
501) -> Result<(), String> {
502 let (client, runtime, project_id) = initialize_client_with_project_context().await?;
503 let table = qualified_table(&runtime.schema, "xbp_schedules");
504 let normalized_target_ref = target_ref.unwrap_or("");
505 let sql = format!(
506 "INSERT INTO {table} (project_id, schedule_type, target_kind, target_ref, expression, timezone, enabled, metadata, occurred_at, updated_at) \
507 VALUES ({project_id}, {schedule_type}, {target_kind}, {target_ref}, {expression}, {timezone}, {enabled}, {metadata}, now(), now()) \
508 ON CONFLICT (schedule_type, target_kind, target_ref, expression) DO UPDATE SET \
509 project_id = EXCLUDED.project_id, \
510 timezone = EXCLUDED.timezone, \
511 enabled = EXCLUDED.enabled, \
512 metadata = EXCLUDED.metadata, \
513 occurred_at = EXCLUDED.occurred_at, \
514 updated_at = now()",
515 table = table,
516 project_id = optional_uuid_literal(project_id.as_deref()),
517 schedule_type = sql_literal(&Value::String(schedule_type.to_string())),
518 target_kind = sql_literal(&Value::String(target_kind.to_string())),
519 target_ref = sql_literal(&Value::String(normalized_target_ref.to_string())),
520 expression = sql_literal(&Value::String(expression.to_string())),
521 timezone = sql_literal(&Value::String("UTC".to_string())),
522 enabled = if enabled { "true" } else { "false" },
523 metadata = sql_literal(&metadata),
524 );
525
526 let _ = execute_sql(&client, &sql).await?;
527 Ok(())
528}
529
530async fn initialize_client_with_project_context(
531) -> Result<(AthenaClient, AthenaRuntimeConfig, Option<String>), String> {
532 let context = load_current_project_context();
533 let config_ref = context.as_ref().map(|ctx| &ctx.config);
534 let (client, runtime) = initialize_client(config_ref).await?;
535 let project_id = if let Some(ctx) = context {
536 Some(
537 upsert_project_snapshot_with_client(
538 &client,
539 &runtime,
540 &ctx.project_root,
541 &ctx.config,
542 Some(ctx.config_kind.as_str()),
543 )
544 .await?,
545 )
546 } else {
547 None
548 };
549 Ok((client, runtime, project_id))
550}
551
552async fn initialize_client(
553 config: Option<&XbpConfig>,
554) -> Result<(AthenaClient, AthenaRuntimeConfig), String> {
555 let runtime = resolve_runtime_config(config).ok_or_else(|| DB_NOT_CONFIGURED.to_string())?;
556
557 let client = AthenaClient::new_with_backend_name(
558 runtime.url.clone(),
559 runtime.key.clone(),
560 "xbp-cli",
561 &runtime.backend,
562 )
563 .await
564 .map_err(|err| format!("failed to initialize Athena client: {err}"))?;
565
566 ensure_schema_bootstrapped(&client, &runtime).await?;
567 Ok((client, runtime))
568}
569
570async fn ensure_schema_bootstrapped(
571 client: &AthenaClient,
572 runtime: &AthenaRuntimeConfig,
573) -> Result<(), String> {
574 let key = format!("{}|{}|{}", runtime.backend, runtime.url, runtime.schema);
575 let mut guard = BOOTSTRAPPED_BACKENDS.lock().await;
576 if guard.contains(&key) {
577 return Ok(());
578 }
579
580 let schema = sanitize_identifier(&runtime.schema).unwrap_or_else(|| DEFAULT_SCHEMA.to_string());
581 let create_schema_sql = format!("CREATE SCHEMA IF NOT EXISTS {}", schema);
582 let _ = execute_sql(client, &create_schema_sql).await?;
583
584 for statement in BOOTSTRAP_SQL.split(';') {
585 let trimmed = statement.trim();
586 if trimmed.is_empty() {
587 continue;
588 }
589 let sql = if schema == DEFAULT_SCHEMA {
590 format!("{};", trimmed)
591 } else {
592 format!("SET search_path TO {}; {};", schema, trimmed)
593 };
594 let _ = execute_sql(client, &sql).await?;
595 }
596
597 guard.insert(key);
598 Ok(())
599}
600
601async fn upsert_project_snapshot_with_client(
602 client: &AthenaClient,
603 runtime: &AthenaRuntimeConfig,
604 project_root: &Path,
605 config: &XbpConfig,
606 config_kind: Option<&str>,
607) -> Result<String, String> {
608 let project_table = qualified_table(&runtime.schema, "xbp_projects");
609 let metadata = build_project_metadata(config);
610 let project_sql = format!(
611 "INSERT INTO {table} (project_name, project_path, version, build_dir, port, app_type, branch, target, config_kind, metadata, updated_at) \
612 VALUES ({project_name}, {project_path}, {version}, {build_dir}, {port}, {app_type}, {branch}, {target}, {config_kind}, {metadata}, now()) \
613 ON CONFLICT (project_path) DO UPDATE SET \
614 project_name = EXCLUDED.project_name, \
615 version = EXCLUDED.version, \
616 build_dir = EXCLUDED.build_dir, \
617 port = EXCLUDED.port, \
618 app_type = EXCLUDED.app_type, \
619 branch = EXCLUDED.branch, \
620 target = EXCLUDED.target, \
621 config_kind = EXCLUDED.config_kind, \
622 metadata = EXCLUDED.metadata, \
623 updated_at = now() \
624 RETURNING project_id",
625 table = project_table,
626 project_name = sql_literal(&Value::String(config.project_name.clone())),
627 project_path = sql_literal(&Value::String(project_root.display().to_string())),
628 version = sql_literal(&Value::String(config.version.clone())),
629 build_dir = sql_literal(&Value::String(config.build_dir.clone())),
630 port = config.port,
631 app_type = optional_text_literal(config.app_type.as_deref()),
632 branch = optional_text_literal(config.branch.as_deref()),
633 target = optional_text_literal(config.target.as_deref()),
634 config_kind = optional_text_literal(config_kind),
635 metadata = sql_literal(&metadata),
636 );
637
638 let project_result = execute_sql(client, &project_sql).await?;
639 let project_id = query_first_column_as_string(&project_result, "project_id")
640 .ok_or_else(|| "failed to resolve project_id from upsert".to_string())?;
641
642 if let Some(services) = config.services.as_ref() {
643 let services_table = qualified_table(&runtime.schema, "xbp_project_services");
644 for service in services {
645 upsert_project_service(client, &services_table, &project_id, service).await?;
646 }
647 }
648
649 Ok(project_id)
650}
651async fn upsert_project_service(
652 client: &AthenaClient,
653 table: &str,
654 project_id: &str,
655 service: &ServiceConfig,
656) -> Result<(), String> {
657 let commands = serde_json::to_value(&service.commands).unwrap_or_else(|_| json!({}));
658 let environment = serde_json::to_value(&service.environment).unwrap_or_else(|_| json!({}));
659 let metadata = build_service_metadata(service);
660 let sql = format!(
661 "INSERT INTO {table} (project_id, service_name, target, branch, port, root_directory, url, healthcheck_path, restart_policy, start_wrapper, systemd_service_name, commands, environment, metadata, updated_at) \
662 VALUES ({project_id}::uuid, {service_name}, {target}, {branch}, {port}, {root_directory}, {url}, {healthcheck_path}, {restart_policy}, {start_wrapper}, {systemd_service_name}, {commands}, {environment}, {metadata}, now()) \
663 ON CONFLICT (project_id, service_name) DO UPDATE SET \
664 target = EXCLUDED.target, \
665 branch = EXCLUDED.branch, \
666 port = EXCLUDED.port, \
667 root_directory = EXCLUDED.root_directory, \
668 url = EXCLUDED.url, \
669 healthcheck_path = EXCLUDED.healthcheck_path, \
670 restart_policy = EXCLUDED.restart_policy, \
671 start_wrapper = EXCLUDED.start_wrapper, \
672 systemd_service_name = EXCLUDED.systemd_service_name, \
673 commands = EXCLUDED.commands, \
674 environment = EXCLUDED.environment, \
675 metadata = EXCLUDED.metadata, \
676 updated_at = now()",
677 table = table,
678 project_id = sql_literal(&Value::String(project_id.to_string())),
679 service_name = sql_literal(&Value::String(service.name.clone())),
680 target = sql_literal(&Value::String(service.target.clone())),
681 branch = sql_literal(&Value::String(service.branch.clone())),
682 port = service.port,
683 root_directory = optional_text_literal(service.root_directory.as_deref()),
684 url = optional_text_literal(service.url.as_deref()),
685 healthcheck_path = optional_text_literal(service.healthcheck_path.as_deref()),
686 restart_policy = optional_text_literal(service.restart_policy.as_deref()),
687 start_wrapper = optional_text_literal(service.start_wrapper.as_deref()),
688 systemd_service_name = optional_text_literal(service.systemd_service_name.as_deref()),
689 commands = sql_literal(&commands),
690 environment = sql_literal(&environment),
691 metadata = sql_literal(&metadata),
692 );
693
694 let _ = execute_sql(client, &sql).await?;
695 Ok(())
696}
697
698fn build_project_metadata(config: &XbpConfig) -> Value {
699 json!({
700 "services_count": config.services.as_ref().map(|services| services.len()).unwrap_or(0),
701 "monitor_url": config.monitor_url,
702 "kafka_topic": config.kafka_topic,
703 "systemd_service_name": config.systemd_service_name,
704 })
705}
706
707fn build_service_metadata(service: &ServiceConfig) -> Value {
708 json!({
709 "force_run_from_root": service.force_run_from_root,
710 "restart_policy_max_failure_count": service.restart_policy_max_failure_count,
711 })
712}
713
714fn resolve_connection_pair(database: Option<&DatabaseConfig>) -> Option<(String, String)> {
715 let block_url = database
716 .and_then(|db| db.url_env.as_deref())
717 .and_then(read_env_nonempty);
718 let block_key = database
719 .and_then(|db| db.key_env.as_deref())
720 .and_then(read_env_nonempty);
721
722 if let (Some(url), Some(key)) = (block_url, block_key) {
723 return Some((url, key));
724 }
725
726 let xbp_url = read_env_nonempty("XBP_ATHENA_URL");
727 let xbp_key = read_env_nonempty("XBP_ATHENA_KEY");
728 if let (Some(url), Some(key)) = (xbp_url, xbp_key) {
729 return Some((url, key));
730 }
731
732 let supabase_url = read_env_nonempty("SUPABASE_URL");
733 let supabase_key = read_env_nonempty("SUPABASE_KEY");
734 if let (Some(url), Some(key)) = (supabase_url, supabase_key) {
735 return Some((url, key));
736 }
737
738 let xlx_supabase_url = read_env_nonempty("XLX_SUPABASE_URL");
739 let xlx_supabase_key = read_env_nonempty("XLX_SUPABASE_ANON_KEY");
740 if let (Some(url), Some(key)) = (xlx_supabase_url, xlx_supabase_key) {
741 return Some((url, key));
742 }
743
744 None
745}
746
747fn read_env_nonempty(name: &str) -> Option<String> {
748 match env::var(name) {
749 Ok(value) if !value.trim().is_empty() => Some(value),
750 _ => None,
751 }
752}
753
754fn value_or_default(primary: Option<&str>, secondary: Option<&str>, default: &str) -> String {
755 if let Some(value) = primary {
756 let trimmed = value.trim();
757 if !trimmed.is_empty() {
758 return trimmed.to_string();
759 }
760 }
761 if let Some(value) = secondary {
762 let trimmed = value.trim();
763 if !trimmed.is_empty() {
764 return trimmed.to_string();
765 }
766 }
767 default.to_string()
768}
769
770async fn execute_sql(client: &AthenaClient, sql: &str) -> Result<QueryResult, String> {
771 client
772 .execute_sql(sql)
773 .await
774 .map_err(|err| format!("athena sql execution failed: {err}"))
775}
776
777fn query_first_column_as_string(result: &QueryResult, key: &str) -> Option<String> {
778 result.rows.first().and_then(|row| {
779 row.get(key).and_then(|value| {
780 value
781 .as_str()
782 .map(|value| value.to_string())
783 .or_else(|| value.as_i64().map(|value| value.to_string()))
784 })
785 })
786}
787
788fn load_current_project_context() -> Option<ProjectContext> {
789 let current_dir = env::current_dir().ok()?;
790 let found = find_xbp_config_upwards(¤t_dir)?;
791 let content = std::fs::read_to_string(&found.config_path).ok()?;
792 let (config, _) = parse_config_with_auto_heal::<XbpConfig>(&content, found.kind).ok()?;
793
794 Some(ProjectContext {
795 project_root: found.project_root,
796 config,
797 config_kind: found.kind.to_string(),
798 })
799}
800
801fn log_level_label(level: &LogLevel) -> &'static str {
802 match level {
803 LogLevel::Info => "INFO",
804 LogLevel::Warning => "WARN",
805 LogLevel::Error => "ERROR",
806 LogLevel::Debug => "DEBUG",
807 LogLevel::Success => "SUCCESS",
808 }
809}
810
811fn to_rfc3339(ts: DateTime<Local>) -> String {
812 ts.with_timezone(&Utc).to_rfc3339()
813}
814
815fn sanitize_identifier(value: &str) -> Option<String> {
816 let trimmed = value.trim();
817 if trimmed.is_empty() {
818 return None;
819 }
820 let mut chars = trimmed.chars();
821 let first = chars.next()?;
822 if !(first.is_ascii_alphabetic() || first == '_') {
823 return None;
824 }
825 if chars.all(|ch| ch.is_ascii_alphanumeric() || ch == '_') {
826 Some(trimmed.to_string())
827 } else {
828 None
829 }
830}
831
832fn qualified_table(schema: &str, table: &str) -> String {
833 let schema = sanitize_identifier(schema).unwrap_or_else(|| DEFAULT_SCHEMA.to_string());
834 format!("{}.{}", schema, table)
835}
836
837pub fn sql_literal(value: &Value) -> String {
838 match value {
839 Value::Null => "NULL".to_string(),
840 Value::Bool(boolean) => boolean.to_string(),
841 Value::Number(number) => number.to_string(),
842 Value::String(text) => format!("'{}'", text.replace('\'', "''")),
843 Value::Array(_) | Value::Object(_) => {
844 format!("'{}'::jsonb", value.to_string().replace('\'', "''"))
845 }
846 }
847}
848
849fn optional_text_literal(value: Option<&str>) -> String {
850 match value {
851 Some(value) if !value.trim().is_empty() => sql_literal(&Value::String(value.to_string())),
852 _ => "NULL".to_string(),
853 }
854}
855
856fn optional_u64_literal(value: Option<u64>) -> String {
857 match value {
858 Some(value) => value.to_string(),
859 None => "NULL".to_string(),
860 }
861}
862
863fn optional_uuid_literal(value: Option<&str>) -> String {
864 match value {
865 Some(value) if !value.trim().is_empty() => {
866 format!("{}::uuid", sql_literal(&Value::String(value.to_string())))
867 }
868 _ => "NULL".to_string(),
869 }
870}
871
872pub async fn with_fail_open<T, F>(operation: &str, fut: F) -> Option<T>
873where
874 F: Future<Output = Result<T, String>>,
875{
876 match fut.await {
877 Ok(result) => Some(result),
878 Err(error) if error == DB_NOT_CONFIGURED => None,
879 Err(error) => {
880 warn!("{} (fail-open): {}", operation, error);
881 None
882 }
883 }
884}
885#[cfg(test)]
886mod tests {
887 use super::{
888 build_project_metadata, build_service_metadata, extract_cron_restart_expression,
889 resolve_runtime_config, sql_literal, with_fail_open,
890 };
891 use crate::strategies::{DatabaseConfig, ServiceConfig, XbpConfig};
892 use serde_json::json;
893 use std::collections::HashMap;
894 use std::sync::Mutex;
895
896 static ENV_TEST_MUTEX: Mutex<()> = Mutex::new(());
897
898 fn base_config() -> XbpConfig {
899 XbpConfig {
900 project_name: "demo".to_string(),
901 version: "0.1.0".to_string(),
902 port: 3000,
903 build_dir: "/tmp/demo".to_string(),
904 app_type: None,
905 build_command: None,
906 start_command: None,
907 install_command: None,
908 environment: None,
909 services: None,
910 openapi: None,
911 workers: None,
912 systemd_service_name: None,
913 systemd: None,
914 kafka_brokers: None,
915 kafka_topic: None,
916 kafka_public_url: None,
917 log_files: None,
918 monitor_url: None,
919 monitor_method: None,
920 monitor_expected_code: None,
921 monitor_interval: None,
922 target: None,
923 branch: None,
924 crate_name: None,
925 npm_script: None,
926 port_storybook: None,
927 url: None,
928 url_storybook: None,
929 database: None,
930 linear: None,
931 github: None,
932 release: None,
933 publish: None,
934 version_targets: Vec::new(),
935 version_domains: Vec::new(),
936 ignore_paths: Vec::new(),
937 }
938 }
939
940 fn clear_env() {
941 for key in [
942 "BLOCK_URL",
943 "BLOCK_KEY",
944 "XBP_ATHENA_BACKEND",
945 "XBP_ATHENA_URL",
946 "XBP_ATHENA_KEY",
947 "SUPABASE_URL",
948 "SUPABASE_KEY",
949 "XLX_SUPABASE_URL",
950 "XLX_SUPABASE_ANON_KEY",
951 ] {
952 std::env::remove_var(key);
953 }
954 }
955
956 #[test]
957 fn resolve_runtime_config_prefers_database_block_env_binding() {
958 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
959 clear_env();
960 std::env::set_var("BLOCK_URL", "https://block.example");
961 std::env::set_var("BLOCK_KEY", "block-key");
962 std::env::set_var("XBP_ATHENA_URL", "https://xbp.example");
963 std::env::set_var("XBP_ATHENA_KEY", "xbp-key");
964
965 let mut config = base_config();
966 config.database = Some(DatabaseConfig {
967 enabled: Some(true),
968 backend: Some("postgres".to_string()),
969 url_env: Some("BLOCK_URL".to_string()),
970 key_env: Some("BLOCK_KEY".to_string()),
971 schema: Some("custom_schema".to_string()),
972 });
973
974 let runtime = resolve_runtime_config(Some(&config)).expect("runtime");
975 assert_eq!(runtime.backend, "postgres");
976 assert_eq!(runtime.url, "https://block.example");
977 assert_eq!(runtime.key, "block-key");
978 assert_eq!(runtime.schema, "custom_schema");
979 }
980
981 #[test]
982 fn resolve_runtime_config_uses_xbp_env_before_supabase_fallback() {
983 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
984 clear_env();
985 std::env::set_var("XBP_ATHENA_BACKEND", "supabase");
986 std::env::set_var("XBP_ATHENA_URL", "https://xbp.env");
987 std::env::set_var("XBP_ATHENA_KEY", "xbp-env-key");
988 std::env::set_var("SUPABASE_URL", "https://supabase.env");
989 std::env::set_var("SUPABASE_KEY", "supabase-key");
990
991 let runtime = resolve_runtime_config(None).expect("runtime");
992 assert_eq!(runtime.url, "https://xbp.env");
993 assert_eq!(runtime.key, "xbp-env-key");
994 }
995
996 #[test]
997 fn resolve_runtime_config_falls_back_to_supabase_then_xlx_supabase() {
998 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
999 clear_env();
1000 std::env::set_var("SUPABASE_URL", "https://supabase.env");
1001 std::env::set_var("SUPABASE_KEY", "supabase-key");
1002 let runtime = resolve_runtime_config(None).expect("runtime");
1003 assert_eq!(runtime.url, "https://supabase.env");
1004 assert_eq!(runtime.key, "supabase-key");
1005
1006 clear_env();
1007 std::env::set_var("XLX_SUPABASE_URL", "https://xlx-supabase.env");
1008 std::env::set_var("XLX_SUPABASE_ANON_KEY", "xlx-key");
1009 let runtime = resolve_runtime_config(None).expect("runtime");
1010 assert_eq!(runtime.url, "https://xlx-supabase.env");
1011 assert_eq!(runtime.key, "xlx-key");
1012 }
1013
1014 #[test]
1015 fn resolve_runtime_config_respects_database_disable_switch() {
1016 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
1017 clear_env();
1018 std::env::set_var("SUPABASE_URL", "https://supabase.env");
1019 std::env::set_var("SUPABASE_KEY", "supabase-key");
1020
1021 let mut config = base_config();
1022 config.database = Some(DatabaseConfig {
1023 enabled: Some(false),
1024 backend: None,
1025 url_env: None,
1026 key_env: None,
1027 schema: None,
1028 });
1029
1030 assert!(resolve_runtime_config(Some(&config)).is_none());
1031 }
1032
1033 #[test]
1034 fn sql_literal_escapes_quotes_and_json_payloads() {
1035 assert_eq!(sql_literal(&json!("O'Hara")), "'O''Hara'");
1036 assert_eq!(sql_literal(&json!(true)), "true");
1037 assert_eq!(sql_literal(&json!(12.5)), "12.5");
1038 assert_eq!(
1039 sql_literal(&json!({"nested":"quote's"})),
1040 "'{\"nested\":\"quote''s\"}'::jsonb"
1041 );
1042 }
1043
1044 #[test]
1045 fn payload_mappers_produce_expected_shapes() {
1046 let mut config = base_config();
1047 config.monitor_url = Some("https://monitor.example".to_string());
1048 config.kafka_topic = Some("xbp.logs".to_string());
1049 config.systemd_service_name = Some("xbp-api".to_string());
1050 config.services = Some(vec![ServiceConfig {
1051 name: "api".to_string(),
1052 target: "rust".to_string(),
1053 branch: "main".to_string(),
1054 port: 8080,
1055 root_directory: Some("services/api".to_string()),
1056 environment: Some(HashMap::from([(
1057 "RUST_LOG".to_string(),
1058 "info".to_string(),
1059 )])),
1060 url: Some("https://api.example.com".to_string()),
1061 healthcheck_path: Some("/health".to_string()),
1062 restart_policy: Some("always".to_string()),
1063 restart_policy_max_failure_count: Some(5),
1064 start_wrapper: Some("pm2".to_string()),
1065 commands: None,
1066 force_run_from_root: Some(false),
1067 version_targets: None,
1068 watch_paths: None,
1069 systemd_service_name: Some("xbp-api".to_string()),
1070 systemd: None,
1071 openapi: None,
1072 }]);
1073
1074 let project_meta = build_project_metadata(&config);
1075 assert_eq!(project_meta["services_count"], 1);
1076 assert_eq!(project_meta["kafka_topic"], "xbp.logs");
1077
1078 let service_meta =
1079 build_service_metadata(config.services.as_ref().unwrap().first().unwrap());
1080 assert_eq!(service_meta["force_run_from_root"], false);
1081 assert_eq!(service_meta["restart_policy_max_failure_count"], 5);
1082 }
1083
1084 #[tokio::test]
1085 async fn fail_open_wrapper_returns_none_on_error_and_value_on_success() {
1086 let success = with_fail_open("test-success", async { Ok::<_, String>(42) }).await;
1087 assert_eq!(success, Some(42));
1088
1089 let failed = with_fail_open::<i32, _>("test-failure", async {
1090 Err::<i32, _>("forced failure".to_string())
1091 })
1092 .await;
1093 assert_eq!(failed, None);
1094 }
1095
1096 #[test]
1097 fn cron_restart_parser_handles_split_and_equals_forms() {
1098 let args = vec![
1099 "npm".to_string(),
1100 "start".to_string(),
1101 "--cron-restart".to_string(),
1102 "0 */6 * * *".to_string(),
1103 ];
1104 assert_eq!(
1105 extract_cron_restart_expression(&args),
1106 Some("0 */6 * * *".to_string())
1107 );
1108
1109 let args = vec![
1110 "npm".to_string(),
1111 "start".to_string(),
1112 "--cron-restart=*/5 * * * *".to_string(),
1113 ];
1114 assert_eq!(
1115 extract_cron_restart_expression(&args),
1116 Some("*/5 * * * *".to_string())
1117 );
1118 }
1119}