#[cfg(feature = "postgres")]
use tokio_postgres::Client;
use crate::config::WaypointConfig;
#[cfg(feature = "postgres")]
use crate::db;
use crate::db::DbClient;
#[cfg(feature = "postgres")]
use crate::db::quote_ident;
#[cfg(feature = "mysql")]
use crate::db::quote_ident_mysql as qi;
use crate::dialect::DialectKind;
use crate::error::{Result, WaypointError};
#[cfg(feature = "postgres")]
#[deprecated(
since = "0.6.0",
note = "Unused PostgreSQL-only entry point superseded by `execute_db`, which handles both engines. Will be removed in 1.0."
)]
pub async fn execute(
client: &Client,
config: &WaypointConfig,
allow_clean: bool,
) -> Result<Vec<String>> {
if !config.migrations.clean_enabled && !allow_clean {
return Err(WaypointError::CleanDisabled);
}
let schema = &config.migrations.schema;
let table = &config.migrations.table;
db::acquire_advisory_lock(client, schema, table).await?;
let result = execute_inner_pg(client, config).await;
if let Err(e) = db::release_advisory_lock(client, schema, table).await {
log::error!("Failed to release advisory lock: {}", e);
}
result
}
pub async fn execute_db(
client: &DbClient,
config: &WaypointConfig,
allow_clean: bool,
) -> Result<Vec<String>> {
if !config.migrations.clean_enabled && !allow_clean {
return Err(WaypointError::CleanDisabled);
}
let schema = client.resolve_schema(&config.migrations.schema).await?;
let table = &config.migrations.table;
client.acquire_lock(&schema, table).await?;
let result = match client.dialect_kind() {
#[cfg(feature = "postgres")]
DialectKind::Postgres => execute_inner_pg(client.as_postgres()?, config).await,
#[cfg(not(feature = "postgres"))]
DialectKind::Postgres => Err(WaypointError::ConfigError(
"PostgreSQL support is not compiled in (enable the `postgres` feature)".into(),
)),
#[cfg(feature = "mysql")]
DialectKind::Mysql => execute_inner_mysql(client, config).await,
#[cfg(not(feature = "mysql"))]
DialectKind::Mysql => Err(WaypointError::ConfigError(
"MySQL support is not compiled in (enable the `mysql` feature)".into(),
)),
};
if let Err(e) = client.release_lock(&schema, table).await {
log::error!("Failed to release advisory lock: {}", e);
}
result
}
#[cfg(feature = "postgres")]
async fn execute_inner_pg(client: &Client, config: &WaypointConfig) -> Result<Vec<String>> {
let schema = &config.migrations.schema;
let schema_q = quote_ident(schema);
let mut dropped = Vec::new();
log::warn!(
"Starting clean — this will DROP all objects in the schema; schema={}",
schema
);
let rows = client
.query(
"SELECT matviewname FROM pg_matviews WHERE schemaname = $1",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let sql = format!(
"DROP MATERIALIZED VIEW IF EXISTS {}.{} CASCADE",
schema_q,
quote_ident(&name)
);
client.batch_execute(&sql).await?;
dropped.push(format!("Materialized view: {}.{}", schema, name));
}
let rows = client
.query(
"SELECT table_name FROM information_schema.views WHERE table_schema = $1",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let sql = format!(
"DROP VIEW IF EXISTS {}.{} CASCADE",
schema_q,
quote_ident(&name)
);
client.batch_execute(&sql).await?;
dropped.push(format!("View: {}.{}", schema, name));
}
let rows = client
.query(
"SELECT tablename FROM pg_tables WHERE schemaname = $1",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let sql = format!(
"DROP TABLE IF EXISTS {}.{} CASCADE",
schema_q,
quote_ident(&name)
);
client.batch_execute(&sql).await?;
dropped.push(format!("Table: {}.{}", schema, name));
}
let rows = client
.query(
"SELECT sequence_name FROM information_schema.sequences WHERE sequence_schema = $1",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let sql = format!(
"DROP SEQUENCE IF EXISTS {}.{} CASCADE",
schema_q,
quote_ident(&name)
);
client.batch_execute(&sql).await?;
dropped.push(format!("Sequence: {}.{}", schema, name));
}
let rows = client
.query(
"SELECT p.proname, pg_get_function_identity_arguments(p.oid) AS args, p.prokind \
FROM pg_proc p \
JOIN pg_namespace n ON p.pronamespace = n.oid \
WHERE n.nspname = $1 \
AND NOT EXISTS ( \
SELECT 1 FROM pg_depend d \
WHERE d.objid = p.oid AND d.deptype = 'e' \
)",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let args: String = row.get(1);
let prokind: i8 = row.get(2);
let (keyword, label) = match prokind as u8 as char {
'p' => ("PROCEDURE", "Procedure"),
'a' => ("AGGREGATE", "Aggregate"),
_ => ("FUNCTION", "Function"),
};
let sql = format!(
"DROP {} IF EXISTS {}.{}({}) CASCADE",
keyword,
schema_q,
quote_ident(&name),
args
);
client.batch_execute(&sql).await?;
dropped.push(format!("{}: {}.{}", label, schema, name));
}
let rows = client
.query(
"SELECT t.typname \
FROM pg_type t \
JOIN pg_namespace n ON t.typnamespace = n.oid \
WHERE n.nspname = $1 \
AND t.typtype IN ('e', 'c') \
AND t.typname NOT LIKE '\\_%'",
&[&schema],
)
.await?;
for row in rows {
let name: String = row.get(0);
let sql = format!(
"DROP TYPE IF EXISTS {}.{} CASCADE",
schema_q,
quote_ident(&name)
);
client.batch_execute(&sql).await?;
dropped.push(format!("Type: {}.{}", schema, name));
}
log::warn!(
"Clean completed; schema={}, objects_dropped={}",
schema,
dropped.len()
);
Ok(dropped)
}
#[cfg(feature = "mysql")]
async fn execute_inner_mysql(client: &DbClient, config: &WaypointConfig) -> Result<Vec<String>> {
use mysql_async::prelude::*;
let pool = client.as_mysql()?;
let schema = client.resolve_schema(&config.migrations.schema).await?;
let mut dropped = Vec::new();
log::warn!(
"Starting clean — this will DROP all objects in the database; database={}",
schema
);
let mut conn = pool.get_conn().await?;
conn.query_drop("SET FOREIGN_KEY_CHECKS = 0").await?;
let views: Vec<String> = conn
.exec(
"SELECT TABLE_NAME FROM information_schema.VIEWS WHERE TABLE_SCHEMA = ?",
(schema.as_str(),),
)
.await?;
for name in views {
let sql = format!("DROP VIEW IF EXISTS {}.{}", qi(&schema), qi(&name));
conn.query_drop(&sql).await?;
dropped.push(format!("View: {}.{}", schema, name));
}
let tables: Vec<String> = conn
.exec(
"SELECT TABLE_NAME FROM information_schema.TABLES \
WHERE TABLE_SCHEMA = ? AND TABLE_TYPE = 'BASE TABLE'",
(schema.as_str(),),
)
.await?;
for name in tables {
let sql = format!("DROP TABLE IF EXISTS {}.{}", qi(&schema), qi(&name));
conn.query_drop(&sql).await?;
dropped.push(format!("Table: {}.{}", schema, name));
}
let routines: Vec<(String, String)> = conn
.exec(
"SELECT ROUTINE_NAME, ROUTINE_TYPE FROM information_schema.ROUTINES \
WHERE ROUTINE_SCHEMA = ?",
(schema.as_str(),),
)
.await?;
for (name, kind) in routines {
let kw = if kind.eq_ignore_ascii_case("PROCEDURE") {
"PROCEDURE"
} else {
"FUNCTION"
};
let sql = format!("DROP {} IF EXISTS {}.{}", kw, qi(&schema), qi(&name));
conn.query_drop(&sql).await?;
dropped.push(format!("{}: {}.{}", kw.to_ascii_lowercase(), schema, name));
}
let events: Vec<String> = conn
.exec(
"SELECT EVENT_NAME FROM information_schema.EVENTS WHERE EVENT_SCHEMA = ?",
(schema.as_str(),),
)
.await?;
for name in events {
let sql = format!("DROP EVENT IF EXISTS {}.{}", qi(&schema), qi(&name));
conn.query_drop(&sql).await?;
dropped.push(format!("Event: {}.{}", schema, name));
}
if let Err(e) = conn.query_drop("SET FOREIGN_KEY_CHECKS = 1").await {
log::warn!(
"Failed to restore FOREIGN_KEY_CHECKS=1 on clean conn: {}",
e
);
}
log::warn!(
"Clean completed; database={}, objects_dropped={}",
schema,
dropped.len()
);
Ok(dropped)
}