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 }
937 }
938
939 fn clear_env() {
940 for key in [
941 "BLOCK_URL",
942 "BLOCK_KEY",
943 "XBP_ATHENA_BACKEND",
944 "XBP_ATHENA_URL",
945 "XBP_ATHENA_KEY",
946 "SUPABASE_URL",
947 "SUPABASE_KEY",
948 "XLX_SUPABASE_URL",
949 "XLX_SUPABASE_ANON_KEY",
950 ] {
951 std::env::remove_var(key);
952 }
953 }
954
955 #[test]
956 fn resolve_runtime_config_prefers_database_block_env_binding() {
957 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
958 clear_env();
959 std::env::set_var("BLOCK_URL", "https://block.example");
960 std::env::set_var("BLOCK_KEY", "block-key");
961 std::env::set_var("XBP_ATHENA_URL", "https://xbp.example");
962 std::env::set_var("XBP_ATHENA_KEY", "xbp-key");
963
964 let mut config = base_config();
965 config.database = Some(DatabaseConfig {
966 enabled: Some(true),
967 backend: Some("postgres".to_string()),
968 url_env: Some("BLOCK_URL".to_string()),
969 key_env: Some("BLOCK_KEY".to_string()),
970 schema: Some("custom_schema".to_string()),
971 });
972
973 let runtime = resolve_runtime_config(Some(&config)).expect("runtime");
974 assert_eq!(runtime.backend, "postgres");
975 assert_eq!(runtime.url, "https://block.example");
976 assert_eq!(runtime.key, "block-key");
977 assert_eq!(runtime.schema, "custom_schema");
978 }
979
980 #[test]
981 fn resolve_runtime_config_uses_xbp_env_before_supabase_fallback() {
982 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
983 clear_env();
984 std::env::set_var("XBP_ATHENA_BACKEND", "supabase");
985 std::env::set_var("XBP_ATHENA_URL", "https://xbp.env");
986 std::env::set_var("XBP_ATHENA_KEY", "xbp-env-key");
987 std::env::set_var("SUPABASE_URL", "https://supabase.env");
988 std::env::set_var("SUPABASE_KEY", "supabase-key");
989
990 let runtime = resolve_runtime_config(None).expect("runtime");
991 assert_eq!(runtime.url, "https://xbp.env");
992 assert_eq!(runtime.key, "xbp-env-key");
993 }
994
995 #[test]
996 fn resolve_runtime_config_falls_back_to_supabase_then_xlx_supabase() {
997 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
998 clear_env();
999 std::env::set_var("SUPABASE_URL", "https://supabase.env");
1000 std::env::set_var("SUPABASE_KEY", "supabase-key");
1001 let runtime = resolve_runtime_config(None).expect("runtime");
1002 assert_eq!(runtime.url, "https://supabase.env");
1003 assert_eq!(runtime.key, "supabase-key");
1004
1005 clear_env();
1006 std::env::set_var("XLX_SUPABASE_URL", "https://xlx-supabase.env");
1007 std::env::set_var("XLX_SUPABASE_ANON_KEY", "xlx-key");
1008 let runtime = resolve_runtime_config(None).expect("runtime");
1009 assert_eq!(runtime.url, "https://xlx-supabase.env");
1010 assert_eq!(runtime.key, "xlx-key");
1011 }
1012
1013 #[test]
1014 fn resolve_runtime_config_respects_database_disable_switch() {
1015 let _guard = ENV_TEST_MUTEX.lock().expect("env lock");
1016 clear_env();
1017 std::env::set_var("SUPABASE_URL", "https://supabase.env");
1018 std::env::set_var("SUPABASE_KEY", "supabase-key");
1019
1020 let mut config = base_config();
1021 config.database = Some(DatabaseConfig {
1022 enabled: Some(false),
1023 backend: None,
1024 url_env: None,
1025 key_env: None,
1026 schema: None,
1027 });
1028
1029 assert!(resolve_runtime_config(Some(&config)).is_none());
1030 }
1031
1032 #[test]
1033 fn sql_literal_escapes_quotes_and_json_payloads() {
1034 assert_eq!(sql_literal(&json!("O'Hara")), "'O''Hara'");
1035 assert_eq!(sql_literal(&json!(true)), "true");
1036 assert_eq!(sql_literal(&json!(12.5)), "12.5");
1037 assert_eq!(
1038 sql_literal(&json!({"nested":"quote's"})),
1039 "'{\"nested\":\"quote''s\"}'::jsonb"
1040 );
1041 }
1042
1043 #[test]
1044 fn payload_mappers_produce_expected_shapes() {
1045 let mut config = base_config();
1046 config.monitor_url = Some("https://monitor.example".to_string());
1047 config.kafka_topic = Some("xbp.logs".to_string());
1048 config.systemd_service_name = Some("xbp-api".to_string());
1049 config.services = Some(vec![ServiceConfig {
1050 name: "api".to_string(),
1051 target: "rust".to_string(),
1052 branch: "main".to_string(),
1053 port: 8080,
1054 root_directory: Some("services/api".to_string()),
1055 environment: Some(HashMap::from([(
1056 "RUST_LOG".to_string(),
1057 "info".to_string(),
1058 )])),
1059 url: Some("https://api.example.com".to_string()),
1060 healthcheck_path: Some("/health".to_string()),
1061 restart_policy: Some("always".to_string()),
1062 restart_policy_max_failure_count: Some(5),
1063 start_wrapper: Some("pm2".to_string()),
1064 commands: None,
1065 force_run_from_root: Some(false),
1066 version_targets: None,
1067 watch_paths: None,
1068 systemd_service_name: Some("xbp-api".to_string()),
1069 systemd: None,
1070 openapi: None,
1071 }]);
1072
1073 let project_meta = build_project_metadata(&config);
1074 assert_eq!(project_meta["services_count"], 1);
1075 assert_eq!(project_meta["kafka_topic"], "xbp.logs");
1076
1077 let service_meta =
1078 build_service_metadata(config.services.as_ref().unwrap().first().unwrap());
1079 assert_eq!(service_meta["force_run_from_root"], false);
1080 assert_eq!(service_meta["restart_policy_max_failure_count"], 5);
1081 }
1082
1083 #[tokio::test]
1084 async fn fail_open_wrapper_returns_none_on_error_and_value_on_success() {
1085 let success = with_fail_open("test-success", async { Ok::<_, String>(42) }).await;
1086 assert_eq!(success, Some(42));
1087
1088 let failed = with_fail_open::<i32, _>("test-failure", async {
1089 Err::<i32, _>("forced failure".to_string())
1090 })
1091 .await;
1092 assert_eq!(failed, None);
1093 }
1094
1095 #[test]
1096 fn cron_restart_parser_handles_split_and_equals_forms() {
1097 let args = vec![
1098 "npm".to_string(),
1099 "start".to_string(),
1100 "--cron-restart".to_string(),
1101 "0 */6 * * *".to_string(),
1102 ];
1103 assert_eq!(
1104 extract_cron_restart_expression(&args),
1105 Some("0 */6 * * *".to_string())
1106 );
1107
1108 let args = vec![
1109 "npm".to_string(),
1110 "start".to_string(),
1111 "--cron-restart=*/5 * * * *".to_string(),
1112 ];
1113 assert_eq!(
1114 extract_cron_restart_expression(&args),
1115 Some("*/5 * * * *".to_string())
1116 );
1117 }
1118}