Skip to main content

xbp_cli/data/
athena.rs

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}
235/// Cap log details persisted to Athena. Full `cargo publish` transcripts can be
236/// multi-megabyte and explode SQL payload size when escaped into INSERT strings.
237const 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(&current_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}