2026-02-14 01:29:34 +01:00
|
|
|
use crate::common::config::AppConfig;
|
2025-03-20 09:22:31 +01:00
|
|
|
use anyhow::Result;
|
2026-02-18 17:56:06 -08:00
|
|
|
use sqlx::{PgPool, postgres::PgPoolOptions};
|
2025-03-20 09:22:31 +01:00
|
|
|
use std::time::Duration;
|
|
|
|
|
|
|
|
|
|
pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
|
2026-02-14 01:29:34 +01:00
|
|
|
tracing::info!(
|
|
|
|
|
"Initializing PostgreSQL connection with URL: {}",
|
|
|
|
|
config
|
|
|
|
|
.database
|
|
|
|
|
.connection_string
|
|
|
|
|
.replace("postgres://", "postgres://[user]:[pass]@")
|
|
|
|
|
);
|
|
|
|
|
|
2025-03-23 22:44:18 +01:00
|
|
|
let mut attempt = 0;
|
2026-02-12 23:20:46 +01:00
|
|
|
const MAX_ATTEMPTS: usize = 5;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-03-23 22:44:18 +01:00
|
|
|
while attempt < MAX_ATTEMPTS {
|
|
|
|
|
attempt += 1;
|
2026-02-14 01:29:34 +01:00
|
|
|
tracing::info!(
|
|
|
|
|
"PostgreSQL connection attempt #{}/{}",
|
|
|
|
|
attempt,
|
|
|
|
|
MAX_ATTEMPTS
|
|
|
|
|
);
|
|
|
|
|
|
2025-03-23 22:44:18 +01:00
|
|
|
match PgPoolOptions::new()
|
|
|
|
|
.max_connections(config.database.max_connections)
|
|
|
|
|
.min_connections(config.database.min_connections)
|
|
|
|
|
.acquire_timeout(Duration::from_secs(config.database.connect_timeout_secs))
|
|
|
|
|
.idle_timeout(Duration::from_secs(config.database.idle_timeout_secs))
|
|
|
|
|
.max_lifetime(Duration::from_secs(config.database.max_lifetime_secs))
|
|
|
|
|
.connect(&config.database.connection_string)
|
2026-02-14 01:29:34 +01:00
|
|
|
.await
|
|
|
|
|
{
|
|
|
|
|
Ok(pool) => {
|
|
|
|
|
match sqlx::query("SELECT 1").execute(&pool).await {
|
|
|
|
|
Ok(_) => {
|
|
|
|
|
tracing::info!("PostgreSQL connection established successfully");
|
|
|
|
|
|
2026-02-18 17:56:06 -08:00
|
|
|
// Always apply schema - it's idempotent (uses IF NOT EXISTS and CREATE OR REPLACE)
|
|
|
|
|
tracing::info!("Applying database schema...");
|
|
|
|
|
if let Err(e) = apply_schema(&pool).await {
|
|
|
|
|
return Err(anyhow::anyhow!(
|
|
|
|
|
"Database schema could not be applied: {}. \
|
|
|
|
|
Run manually: psql -f db/schema.sql",
|
|
|
|
|
e
|
|
|
|
|
));
|
2025-03-23 22:44:18 +01:00
|
|
|
}
|
2026-02-18 17:56:06 -08:00
|
|
|
tracing::info!("Database schema applied successfully");
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
return Ok(pool);
|
2025-03-23 22:44:18 +01:00
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
Err(e) => {
|
|
|
|
|
tracing::error!("Error verifying connection: {}", e);
|
|
|
|
|
if attempt >= MAX_ATTEMPTS {
|
|
|
|
|
return Err(anyhow::anyhow!(
|
|
|
|
|
"Error verifying PostgreSQL connection: {}",
|
|
|
|
|
e
|
|
|
|
|
));
|
|
|
|
|
}
|
2025-03-23 22:44:18 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
Err(e) => {
|
|
|
|
|
tracing::error!(
|
|
|
|
|
"Error connecting to PostgreSQL (attempt {}/{}): {}",
|
|
|
|
|
attempt,
|
|
|
|
|
MAX_ATTEMPTS,
|
|
|
|
|
e
|
|
|
|
|
);
|
|
|
|
|
if attempt >= MAX_ATTEMPTS {
|
|
|
|
|
return Err(anyhow::anyhow!("Error in PostgreSQL connection: {}", e));
|
|
|
|
|
}
|
|
|
|
|
tokio::time::sleep(Duration::from_secs(2)).await;
|
|
|
|
|
}
|
|
|
|
|
}
|
2025-03-23 22:44:18 +01:00
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
Err(anyhow::anyhow!(
|
|
|
|
|
"Could not establish PostgreSQL connection after {} attempts",
|
|
|
|
|
MAX_ATTEMPTS
|
|
|
|
|
))
|
2026-02-12 23:20:46 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Apply the embedded schema.sql to the database.
|
|
|
|
|
/// First tries `raw_sql` (simple query protocol). If that fails, falls back
|
|
|
|
|
/// to splitting the SQL into individual statements and executing them one by one.
|
|
|
|
|
async fn apply_schema(pool: &PgPool) -> Result<()> {
|
|
|
|
|
let schema_sql = include_str!("../../db/schema.sql");
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Attempt 1: raw_sql sends the entire script via the simple query protocol
|
|
|
|
|
match sqlx::raw_sql(schema_sql).execute(pool).await {
|
|
|
|
|
Ok(_) => return Ok(()),
|
|
|
|
|
Err(e) => {
|
2026-02-14 01:29:34 +01:00
|
|
|
tracing::warn!(
|
|
|
|
|
"raw_sql failed ({}), falling back to statement-by-statement execution",
|
|
|
|
|
e
|
|
|
|
|
);
|
2026-02-12 23:20:46 +01:00
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Attempt 2: split into individual statements respecting dollar-quoting
|
|
|
|
|
let statements = split_sql_statements(schema_sql);
|
|
|
|
|
for (i, stmt) in statements.iter().enumerate() {
|
|
|
|
|
let trimmed = stmt.trim();
|
|
|
|
|
if trimmed.is_empty() || trimmed == ";" {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
if let Err(e) = sqlx::raw_sql(trimmed).execute(pool).await {
|
2026-02-14 01:29:34 +01:00
|
|
|
let preview = if trimmed.len() > 200 {
|
|
|
|
|
&trimmed[..200]
|
|
|
|
|
} else {
|
|
|
|
|
trimmed
|
|
|
|
|
};
|
|
|
|
|
tracing::error!(
|
|
|
|
|
"Schema statement {} failed: {}\n--- SQL ---\n{}\n-----------",
|
|
|
|
|
i + 1,
|
|
|
|
|
e,
|
|
|
|
|
preview
|
|
|
|
|
);
|
2026-02-12 23:20:46 +01:00
|
|
|
return Err(anyhow::anyhow!("Schema statement {} failed: {}", i + 1, e));
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
Ok(())
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Split a SQL script into individual statements, correctly handling:
|
|
|
|
|
/// - Dollar-quoted blocks (`$BODY$...$BODY$`, `$$...$$`)
|
|
|
|
|
/// - Single-quoted strings (`'...'`)
|
|
|
|
|
/// - Line comments (`-- ...`)
|
|
|
|
|
/// - Block comments (`/* ... */`)
|
|
|
|
|
fn split_sql_statements(sql: &str) -> Vec<String> {
|
|
|
|
|
let mut statements = Vec::new();
|
|
|
|
|
let mut current = String::new();
|
|
|
|
|
let chars: Vec<char> = sql.chars().collect();
|
|
|
|
|
let len = chars.len();
|
|
|
|
|
let mut i = 0;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
while i < len {
|
|
|
|
|
// Line comment
|
|
|
|
|
if i + 1 < len && chars[i] == '-' && chars[i + 1] == '-' {
|
|
|
|
|
while i < len && chars[i] != '\n' {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Block comment
|
|
|
|
|
if i + 1 < len && chars[i] == '/' && chars[i + 1] == '*' {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
current.push(chars[i + 1]);
|
|
|
|
|
i += 2;
|
|
|
|
|
while i + 1 < len && !(chars[i] == '*' && chars[i + 1] == '/') {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
|
|
|
|
if i + 1 < len {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
current.push(chars[i + 1]);
|
|
|
|
|
i += 2;
|
|
|
|
|
}
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Single-quoted string
|
|
|
|
|
if chars[i] == '\'' {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
while i < len {
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
if chars[i] == '\'' {
|
|
|
|
|
if i + 1 < len && chars[i + 1] == '\'' {
|
|
|
|
|
current.push(chars[i + 1]);
|
|
|
|
|
i += 2;
|
|
|
|
|
} else {
|
|
|
|
|
i += 1;
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Dollar-quoted string ($tag$...$tag$ or $$...$$)
|
|
|
|
|
if chars[i] == '$' {
|
|
|
|
|
let _start = i;
|
|
|
|
|
i += 1;
|
|
|
|
|
let mut tag = String::from("$");
|
|
|
|
|
while i < len && (chars[i].is_alphanumeric() || chars[i] == '_') {
|
|
|
|
|
tag.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
|
|
|
|
if i < len && chars[i] == '$' {
|
|
|
|
|
tag.push('$');
|
|
|
|
|
i += 1;
|
|
|
|
|
// We have a dollar-quote tag, find the closing tag
|
|
|
|
|
current.push_str(&tag);
|
|
|
|
|
loop {
|
|
|
|
|
if i >= len {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
if chars[i] == '$' {
|
|
|
|
|
let remaining = &sql[i..];
|
|
|
|
|
if remaining.starts_with(&tag) {
|
|
|
|
|
current.push_str(&tag);
|
|
|
|
|
i += tag.len();
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
current.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
// Not a valid dollar-quote, push what we consumed
|
|
|
|
|
current.push_str(&tag);
|
|
|
|
|
}
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Statement separator
|
|
|
|
|
if chars[i] == ';' {
|
|
|
|
|
current.push(';');
|
|
|
|
|
let trimmed = current.trim().to_string();
|
|
|
|
|
if !trimmed.is_empty() && trimmed != ";" {
|
|
|
|
|
statements.push(trimmed);
|
|
|
|
|
}
|
|
|
|
|
current.clear();
|
|
|
|
|
i += 1;
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
current.push(chars[i]);
|
|
|
|
|
i += 1;
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
// Trailing statement without semicolon
|
|
|
|
|
let trimmed = current.trim().to_string();
|
|
|
|
|
if !trimmed.is_empty() && trimmed != ";" {
|
|
|
|
|
statements.push(trimmed);
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 23:20:46 +01:00
|
|
|
statements
|
2026-02-14 01:29:34 +01:00
|
|
|
}
|