diff --git a/data_connector/README.md b/data_connector/README.md index aca18367c5..c9594d2159 100644 --- a/data_connector/README.md +++ b/data_connector/README.md @@ -153,6 +153,8 @@ postgres: owner: "myschema" # Oracle: schema prefix (MYSCHEMA."TABLE") # Redis: key prefix ("myschema:conversation:{id}") # Postgres: ignored (use search_path for schema control) + version: 2 # Skip migrations 1–2 (database already at v2) + auto_migrate: true # Opt in to automatic migration (default: false) conversations: table: "my_conversations" # Overrides default "conversations" @@ -176,7 +178,7 @@ postgres: | Type | Fields | Purpose | |------|--------|---------| -| `SchemaConfig` | `owner`, `conversations`, `responses`, `conversation_items`, `conversation_item_links` | Top-level config with an optional owner/prefix and per-table settings | +| `SchemaConfig` | `owner`, `version`, `auto_migrate`, `conversations`, `responses`, `conversation_items`, `conversation_item_links` | Top-level config with an optional owner/prefix, migration control, and per-table settings | | `TableConfig` | `table`, `columns` | Per-table config: physical table name and a map of logical-to-physical column name overrides | Key behaviors: @@ -390,11 +392,68 @@ PostgreSQL additionally creates an index on `conversation_item_links(conversation_id, added_at)` for efficient cursor-based listing. -Column names within each table can also be overridden via `SchemaConfig`. The -config describes the existing database schema — it does not perform migrations. +Column names within each table can also be overridden via `SchemaConfig`. If you rename a column in config, the corresponding database column must already exist with that name. +### Schema Versioning + +The data connector includes a built-in migration system that tracks applied +schema changes in a `_schema_versions` table. Each backend (Oracle, Postgres) +defines its own migration list with backend-specific DDL. Redis has no +structural schema and does not use migrations. + +**Safe by default**: `auto_migrate` defaults to `false`. When pending +migrations are detected, startup **fails with the exact SQL statements** +needed so you can review and apply them manually. Set `auto_migrate: true` +to opt in to automatic migration. + +On startup: +1. Tables are created if they don't exist (`CREATE TABLE` / `CREATE TABLE IF NOT EXISTS`) +2. The `_schema_versions` tracking table is created +3. Pending migrations are checked: + - If `auto_migrate: true` → migrations are applied automatically + - If `auto_migrate: false` (default) → startup fails with actionable SQL if migrations are pending + +Current migrations: + +| Version | Description | +|---------|-------------| +| 1 | Add `safety_identifier` column to responses | +| 2 | Remove legacy `user_id` column from responses | + +#### Controlling migrations + +| Config field | Type | Default | Description | +|---|---|---|---| +| `auto_migrate` | `bool` | `false` | Set to `true` to apply migrations automatically on startup. When `false`, startup fails with the exact SQL if migrations are pending. | +| `version` | `u32` (optional) | `None` | Starting version — migrations up to this number are skipped. Use when your database is already at a known version. | + +Example: opt in to automatic migration: + +```yaml +oracle: + schema: + auto_migrate: true +``` + +Example: database already at version 2, skip those migrations: + +```yaml +oracle: + schema: + version: 2 + auto_migrate: true +``` + +#### Concurrency safety + +- **Postgres**: Uses `pg_advisory_lock` to serialize migrations across + concurrent application instances. +- **Oracle**: DDL statements use PL/SQL exception handling for idempotency. + Duplicate version records from concurrent instances are detected and skipped + (ORA-00001). + ## Testing Run the unit tests (Memory and NoOp backends, config validation, ID generation): diff --git a/data_connector/src/factory.rs b/data_connector/src/factory.rs index 08622356aa..1c24dc7773 100644 --- a/data_connector/src/factory.rs +++ b/data_connector/src/factory.rs @@ -178,10 +178,19 @@ async fn create_postgres_storage(postgres_cfg: &PostgresConfig) -> Result Result<(), String> { - let s = &schema.responses; - let col_safety = s.col("safety_identifier"); - // Table and column names are already uppercased by OracleStore::new(). - let col_upper = col_safety.to_uppercase(); - let table = s.qualified_table(schema.owner.as_deref()); - - let present: i64 = conn - .query_row_as( - &format!( - "SELECT COUNT(*) FROM user_tab_columns \ - WHERE table_name = '{}' AND column_name = '{col_upper}'", - s.table - ), - &[], - ) - .map_err(map_oracle_error)?; - - if present == 0 { - if let Err(err) = conn.execute( - &format!("ALTER TABLE {table} ADD ({col_safety} VARCHAR2(128))"), - &[], - ) { - let present_after: i64 = conn - .query_row_as( - &format!( - "SELECT COUNT(*) FROM user_tab_columns \ - WHERE table_name = '{}' AND column_name = '{col_upper}'", - s.table - ), - &[], - ) - .map_err(map_oracle_error)?; - if present_after == 0 { - return Err(map_oracle_error(err)); - } - } - } - - Ok(()) - } - - fn remove_user_id_column_if_exists( - conn: &Connection, - schema: &SchemaConfig, - ) -> Result<(), String> { - // Table and column names are already uppercased by OracleStore::new(). - let s = &schema.responses; - let table = s.qualified_table(schema.owner.as_deref()); - - let present: i64 = conn - .query_row_as( - &format!( - "SELECT COUNT(*) FROM user_tab_columns \ - WHERE table_name = '{}' AND column_name = 'USER_ID'", - s.table - ), - &[], - ) - .map_err(map_oracle_error)?; - - if present > 0 { - if let Err(err) = conn.execute(&format!("ALTER TABLE {table} DROP COLUMN USER_ID"), &[]) - { - let present_after: i64 = conn - .query_row_as( - &format!( - "SELECT COUNT(*) FROM user_tab_columns \ - WHERE table_name = '{}' AND column_name = 'USER_ID'", - s.table - ), - &[], - ) - .map_err(map_oracle_error)?; - if present_after > 0 { - return Err(map_oracle_error(err)); - } - } - } - - Ok(()) - } - fn build_response_from_row(row: &Row, schema: &SchemaConfig) -> Result { let s = &schema.responses; let col_id = s.col("id"); @@ -1625,7 +1555,8 @@ fn create_index_if_missing( if let Some(db_err) = err.db_error() { // ORA-00955: name is already used by an existing object // ORA-01408: such column list already indexed - if db_err.code() != 955 && db_err.code() != 1408 { + // ORA-00904: invalid identifier (column not yet added by migration) + if db_err.code() != 955 && db_err.code() != 1408 && db_err.code() != 904 { return Err(map_oracle_error(err)); } } else { diff --git a/data_connector/src/oracle_migrations.rs b/data_connector/src/oracle_migrations.rs new file mode 100644 index 0000000000..13cfd7e242 --- /dev/null +++ b/data_connector/src/oracle_migrations.rs @@ -0,0 +1,130 @@ +//! Oracle-specific schema migrations. +//! +//! Each migration is a function that generates Oracle DDL from [`SchemaConfig`], +//! so it respects custom table/column names. PL/SQL exception handling ensures +//! idempotency (safe to re-run if a previous attempt partially completed). + +use crate::{schema::SchemaConfig, versioning::Migration}; + +/// Oracle migration list. Append new migrations here. +pub(crate) static ORACLE_MIGRATIONS: [Migration; 2] = [ + Migration { + version: 1, + description: "Add safety_identifier column to responses", + up: oracle_v1_up, + }, + Migration { + version: 2, + description: "Remove legacy user_id column from responses", + up: oracle_v2_up, + }, +]; + +fn oracle_v1_up(schema: &SchemaConfig) -> Vec { + let s = &schema.responses; + if s.is_skipped("safety_identifier") { + return vec![]; + } + let table = s.qualified_table(schema.owner.as_deref()); + let col = s.col("safety_identifier"); + // PL/SQL block: ORA-01430 = "column already exists" (idempotent) + vec![format!( + "BEGIN EXECUTE IMMEDIATE 'ALTER TABLE {table} ADD ({col} VARCHAR2(128))'; \ + EXCEPTION WHEN OTHERS THEN IF SQLCODE != -1430 THEN RAISE; END IF; END;" + )] +} + +fn oracle_v2_up(schema: &SchemaConfig) -> Vec { + let s = &schema.responses; + // Don't drop USER_ID if a configured column maps to that name + // or if it's defined as an extra column. + if s.columns + .values() + .any(|v| v.eq_ignore_ascii_case("USER_ID")) + || s.extra_columns + .keys() + .any(|k| k.eq_ignore_ascii_case("USER_ID")) + { + return vec![]; + } + let table = s.qualified_table(schema.owner.as_deref()); + // PL/SQL block: ORA-00904 = "invalid identifier" (column doesn't exist) + vec![format!( + "BEGIN EXECUTE IMMEDIATE 'ALTER TABLE {table} DROP (USER_ID)'; \ + EXCEPTION WHEN OTHERS THEN IF SQLCODE != -904 THEN RAISE; END IF; END;" + )] +} + +// ── Tests ────────────────────────────────────────────────────────────────── + +#[cfg(test)] +mod tests { + use super::*; + use crate::schema::TableConfig; + + #[test] + fn oracle_migrations_are_sequential() { + for (i, m) in ORACLE_MIGRATIONS.iter().enumerate() { + assert_eq!(m.version, (i + 1) as u32, "migration {i} has wrong version"); + } + } + + #[test] + fn oracle_v1_up_generates_plsql_add_column() { + let schema = SchemaConfig::default(); + let stmts = oracle_v1_up(&schema); + assert_eq!(stmts.len(), 1); + assert!(stmts[0].contains("ADD"), "got: {}", stmts[0]); + assert!(stmts[0].contains("SQLCODE"), "got: {}", stmts[0]); + } + + #[test] + fn oracle_v1_up_skipped_returns_empty() { + let schema = SchemaConfig { + responses: TableConfig { + skip_columns: ["safety_identifier".to_string()].into_iter().collect(), + ..TableConfig::with_table("responses") + }, + ..Default::default() + }; + let stmts = oracle_v1_up(&schema); + assert!(stmts.is_empty()); + } + + #[test] + fn oracle_v2_up_generates_plsql_drop_column() { + let schema = SchemaConfig::default(); + let stmts = oracle_v2_up(&schema); + assert_eq!(stmts.len(), 1); + assert!(stmts[0].contains("DROP"), "got: {}", stmts[0]); + assert!(stmts[0].contains("USER_ID"), "got: {}", stmts[0]); + } + + #[test] + fn oracle_v2_up_skipped_when_column_maps_to_user_id() { + let mut schema = SchemaConfig::default(); + schema + .responses + .columns + .insert("safety_identifier".to_string(), "USER_ID".to_string()); + let stmts = oracle_v2_up(&schema); + assert!(stmts.is_empty(), "should skip drop when USER_ID is mapped"); + } + + #[test] + fn oracle_v2_up_skipped_when_extra_column_is_user_id() { + let mut schema = SchemaConfig::default(); + schema.responses.extra_columns.insert( + "USER_ID".to_string(), + crate::schema::ColumnDef { + sql_type: "VARCHAR2(128)".to_string(), + default_value: None, + }, + ); + let stmts = oracle_v2_up(&schema); + assert!( + stmts.is_empty(), + "should skip drop when USER_ID is an extra column" + ); + } +} diff --git a/data_connector/src/postgres.rs b/data_connector/src/postgres.rs index bdebe2b4b9..c4f4a0d058 100644 --- a/data_connector/src/postgres.rs +++ b/data_connector/src/postgres.rs @@ -28,9 +28,12 @@ use crate::{ ListParams, NewConversation, NewConversationItem, ResponseId, ResponseResult, ResponseStorage, ResponseStorageError, SortOrder, StoredResponse, }, + postgres_migrations::POSTGRES_MIGRATIONS, schema::SchemaConfig, }; +// ── Store ──────────────────────────────────────────────────────────────── + pub(crate) struct PostgresStore { pool: Pool, pub(crate) schema: Arc, @@ -55,6 +58,46 @@ impl PostgresStore { Ok(Self { pool, schema }) } + + /// Run versioned schema migrations after tables have been created. + pub(crate) async fn run_migrations(&self) -> Result, String> { + let mut client = self + .pool + .get() + .await + .map_err(|e| format!("failed to get connection for migrations: {e}"))?; + crate::versioning::run_postgres_migrations( + &mut client, + &self.schema, + &POSTGRES_MIGRATIONS, + self.schema.version, + self.schema.auto_migrate, + ) + .await + } + + /// Create indexes that may have been deferred during init because + /// migration-added columns did not yet exist. + pub(crate) async fn ensure_response_indexes(&self) -> Result<(), String> { + let s = &self.schema.responses; + if s.is_skipped("safety_identifier") { + return Ok(()); + } + let table = s.qualified_table(self.schema.owner.as_deref()); + let col = s.col("safety_identifier"); + let idx_ddl = + format!("CREATE INDEX IF NOT EXISTS responses_safety_idx ON {table} ({col});"); + let client = self + .pool + .get() + .await + .map_err(|e| format!("failed to get connection for index creation: {e}"))?; + client + .batch_execute(&idx_ddl) + .await + .map_err(|e| format!("failed to create response index: {e}"))?; + Ok(()) + } } impl Clone for PostgresStore { @@ -848,16 +891,10 @@ impl PostgresResponseStorage { } col_defs.extend(extra_column_defs(s)); - let mut ddl = format!( + let table_ddl = format!( "CREATE TABLE IF NOT EXISTS {table} ({});", col_defs.join(", "), ); - if !s.is_skipped("safety_identifier") { - ddl.push_str(&format!( - "\nCREATE INDEX IF NOT EXISTS responses_safety_idx ON {table} ({});", - s.col("safety_identifier") - )); - } let client = store .pool @@ -865,10 +902,23 @@ impl PostgresResponseStorage { .await .map_err(|e| ResponseStorageError::StorageError(e.to_string()))?; client - .batch_execute(&ddl) + .batch_execute(&table_ddl) .await .map_err(|e| ResponseStorageError::StorageError(e.to_string()))?; + // Index creation is separate so legacy tables missing migrated + // columns (e.g. safety_identifier) don't block startup. Indexes + // are retried after migrations via ensure_response_indexes(). + if !s.is_skipped("safety_identifier") { + let idx_ddl = format!( + "CREATE INDEX IF NOT EXISTS responses_safety_idx ON {table} ({});", + s.col("safety_identifier") + ); + if let Err(e) = client.batch_execute(&idx_ddl).await { + tracing::debug!("deferred response index creation (column may not exist yet): {e}"); + } + } + let select_base = build_response_select_base(&store.schema); Ok(Self { store, select_base }) } diff --git a/data_connector/src/postgres_migrations.rs b/data_connector/src/postgres_migrations.rs new file mode 100644 index 0000000000..e7bca4f62e --- /dev/null +++ b/data_connector/src/postgres_migrations.rs @@ -0,0 +1,122 @@ +//! Postgres-specific schema migrations. +//! +//! Each migration is a function that generates Postgres DDL from [`SchemaConfig`], +//! so it respects custom table/column names. `IF NOT EXISTS` / `IF EXISTS` +//! clauses ensure idempotency. + +use crate::{schema::SchemaConfig, versioning::Migration}; + +/// Postgres migration list. Append new migrations here. +pub(crate) static POSTGRES_MIGRATIONS: [Migration; 2] = [ + Migration { + version: 1, + description: "Add safety_identifier column to responses", + up: pg_v1_up, + }, + Migration { + version: 2, + description: "Remove legacy user_id column from responses", + up: pg_v2_up, + }, +]; + +fn pg_v1_up(schema: &SchemaConfig) -> Vec { + let s = &schema.responses; + if s.is_skipped("safety_identifier") { + return vec![]; + } + let table = s.qualified_table(schema.owner.as_deref()); + let col = s.col("safety_identifier"); + vec![format!( + "ALTER TABLE {table} ADD COLUMN IF NOT EXISTS {col} VARCHAR(128)" + )] +} + +fn pg_v2_up(schema: &SchemaConfig) -> Vec { + let s = &schema.responses; + // Don't drop user_id if a configured column maps to that name + // or if it's defined as an extra column. + if s.columns + .values() + .any(|v| v.eq_ignore_ascii_case("user_id")) + || s.extra_columns + .keys() + .any(|k| k.eq_ignore_ascii_case("user_id")) + { + return vec![]; + } + let table = s.qualified_table(schema.owner.as_deref()); + vec![format!("ALTER TABLE {table} DROP COLUMN IF EXISTS user_id")] +} + +// ── Tests ────────────────────────────────────────────────────────────────── + +#[cfg(test)] +mod tests { + use super::*; + use crate::schema::TableConfig; + + #[test] + fn postgres_migrations_are_sequential() { + for (i, m) in POSTGRES_MIGRATIONS.iter().enumerate() { + assert_eq!(m.version, (i + 1) as u32, "migration {i} has wrong version"); + } + } + + #[test] + fn pg_v1_up_generates_add_column_if_not_exists() { + let schema = SchemaConfig::default(); + let stmts = pg_v1_up(&schema); + assert_eq!(stmts.len(), 1); + assert!(stmts[0].contains("IF NOT EXISTS"), "got: {}", stmts[0]); + } + + #[test] + fn pg_v1_up_skipped_returns_empty() { + let schema = SchemaConfig { + responses: TableConfig { + skip_columns: ["safety_identifier".to_string()].into_iter().collect(), + ..TableConfig::with_table("responses") + }, + ..Default::default() + }; + let stmts = pg_v1_up(&schema); + assert!(stmts.is_empty()); + } + + #[test] + fn pg_v2_up_generates_drop_column_if_exists() { + let schema = SchemaConfig::default(); + let stmts = pg_v2_up(&schema); + assert_eq!(stmts.len(), 1); + assert!(stmts[0].contains("IF EXISTS"), "got: {}", stmts[0]); + } + + #[test] + fn pg_v2_up_skipped_when_column_maps_to_user_id() { + let mut schema = SchemaConfig::default(); + schema + .responses + .columns + .insert("safety_identifier".to_string(), "user_id".to_string()); + let stmts = pg_v2_up(&schema); + assert!(stmts.is_empty(), "should skip drop when user_id is mapped"); + } + + #[test] + fn pg_v2_up_skipped_when_extra_column_is_user_id() { + let mut schema = SchemaConfig::default(); + schema.responses.extra_columns.insert( + "user_id".to_string(), + crate::schema::ColumnDef { + sql_type: "VARCHAR(128)".to_string(), + default_value: None, + }, + ); + let stmts = pg_v2_up(&schema); + assert!( + stmts.is_empty(), + "should skip drop when user_id is an extra column" + ); + } +} diff --git a/data_connector/src/schema.rs b/data_connector/src/schema.rs index 77210a19fe..2e9df0e915 100644 --- a/data_connector/src/schema.rs +++ b/data_connector/src/schema.rs @@ -25,6 +25,20 @@ pub struct SchemaConfig { #[serde(skip_serializing_if = "Option::is_none")] pub owner: Option, + /// Starting schema version. Set this when your database already has + /// migrations applied (e.g. `version: 3` skips migrations 1–3). + /// `None` means start from 0 (apply all migrations). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub version: Option, + + /// Whether to run schema migrations automatically on startup. + /// Defaults to `false` (safe by default). When `false` and pending + /// migrations are detected, startup fails with the exact SQL statements + /// needed so you can review and apply them manually. + /// Set to `true` to opt in to automatic migration. + #[serde(default = "default_auto_migrate")] + pub auto_migrate: bool, + pub conversations: TableConfig, pub responses: TableConfig, pub conversation_items: TableConfig, @@ -72,10 +86,19 @@ pub struct ColumnDef { // Defaults // ──────────────────────────────────────────────────────────────────────────── +fn default_auto_migrate() -> bool { + std::env::var("DB_AUTO_MIGRATE") + .ok() + .map(|v| v.eq_ignore_ascii_case("true") || v == "1") + .unwrap_or(false) +} + impl Default for SchemaConfig { fn default() -> Self { Self { owner: None, + version: None, + auto_migrate: default_auto_migrate(), conversations: TableConfig::with_table("conversations"), responses: TableConfig::with_table("responses"), conversation_items: TableConfig::with_table("conversation_items"), @@ -537,6 +560,38 @@ mod tests { assert_eq!(cfg, SchemaConfig::default()); } + // ── version / auto_migrate serde ──────────────────────────────────── + + #[test] + fn serde_roundtrip_with_version_and_auto_migrate() { + let cfg = SchemaConfig { + version: Some(3), + auto_migrate: false, + ..Default::default() + }; + let json = serde_json::to_string(&cfg).expect("serialize"); + let restored: SchemaConfig = serde_json::from_str(&json).expect("deserialize"); + assert_eq!(restored.version, Some(3)); + assert!(!restored.auto_migrate); + } + + #[test] + fn serde_defaults_version_none_and_auto_migrate_false() { + let cfg: SchemaConfig = serde_json::from_str("{}").expect("deserialize empty"); + assert_eq!(cfg.version, None); + assert!(!cfg.auto_migrate); + } + + #[test] + fn serde_version_none_is_omitted_from_json() { + let cfg = SchemaConfig::default(); + let json = serde_json::to_string(&cfg).expect("serialize"); + assert!( + !json.contains("version"), + "version:None should be skipped: {json}" + ); + } + // ── is_skipped() ───────────────────────────────────────────────────── #[test] diff --git a/data_connector/src/versioning.rs b/data_connector/src/versioning.rs new file mode 100644 index 0000000000..f5bdb32df3 --- /dev/null +++ b/data_connector/src/versioning.rs @@ -0,0 +1,596 @@ +//! Schema versioning and migration infrastructure. +//! +//! Replaces ad-hoc `ALTER TABLE` / column-existence checks with a tracked, +//! per-backend migration system. Each backend defines its own migration list +//! (DDL syntax differs across databases) and migrations reference +//! [`SchemaConfig`] so they work correctly even with custom table/column names. +//! +//! # Version tracking +//! +//! A `_schema_versions` table records which migrations have been applied. +//! On startup the runner reads the current version and applies any pending +//! migrations in order. +//! +//! # Safety +//! +//! `auto_migrate` defaults to `false`. When pending migrations are detected +//! and auto-migration is off, startup **fails** with the exact SQL statements +//! needed so the operator can review and apply them manually. +//! +//! # Configuration +//! +//! ```yaml +//! oracle: +//! schema: +//! version: 2 # "my schema is already at v2, skip 1-2" +//! auto_migrate: true # opt in to automatic migration +//! ``` + +use crate::schema::SchemaConfig; + +// ── Types ────────────────────────────────────────────────────────────────── + +/// A single schema migration. +/// +/// Migrations are functions (not static SQL strings) so they can reference +/// [`SchemaConfig`] for table/column names. `up` returns a `Vec` +/// because some migrations require multiple DDL statements (e.g. ALTER TABLE +/// followed by CREATE INDEX). +pub struct Migration { + /// Monotonically increasing version number (1, 2, 3, …). + pub version: u32, + /// Human-readable description for the version log. + pub description: &'static str, + /// Generate the forward-migration DDL statements. + pub up: fn(&SchemaConfig) -> Vec, +} + +// ── Versions table DDL ───────────────────────────────────────────────────── + +/// Name of the schema-versions tracking table. +pub const VERSIONS_TABLE: &str = "_schema_versions"; + +/// Oracle-qualified name for the versions table (always quoted since `_` prefix +/// is invalid for unquoted Oracle identifiers). +fn oracle_versions_table(schema: &SchemaConfig) -> String { + match &schema.owner { + Some(owner) => format!("{owner}.\"{VERSIONS_TABLE}\""), + None => format!("\"{VERSIONS_TABLE}\""), + } +} + +/// Oracle DDL for creating the versions tracking table. +pub fn oracle_create_versions_table(schema: &SchemaConfig) -> String { + let table = oracle_versions_table(schema); + format!( + "CREATE TABLE {table} (\ + version NUMBER(10) NOT NULL PRIMARY KEY, \ + description VARCHAR2(512), \ + applied_at TIMESTAMP DEFAULT SYSTIMESTAMP NOT NULL)" + ) +} + +/// Postgres DDL for creating the versions tracking table. +pub fn postgres_create_versions_table() -> String { + format!( + "CREATE TABLE IF NOT EXISTS {VERSIONS_TABLE} (\ + version INTEGER NOT NULL PRIMARY KEY, \ + description VARCHAR(512), \ + applied_at TIMESTAMPTZ NOT NULL DEFAULT NOW())" + ) +} + +// ── Pending-migration error ─────────────────────────────────────────────── + +/// Build an actionable error message listing pending migrations and their SQL. +fn pending_migrations_error( + backend: &str, + current: u32, + pending: &[&Migration], + schema: &SchemaConfig, +) -> String { + let mut msg = format!( + "Schema migration required (current version: {current}, \ + latest version: {}).\n\n\ + The following migrations need to be applied:\n", + pending.last().map_or(current, |m| m.version) + ); + + let versions_insert = if backend == "oracle" { + format!( + "INSERT INTO {} (version, description) VALUES", + oracle_versions_table(schema), + ) + } else { + format!("INSERT INTO {VERSIONS_TABLE} (version, description) VALUES") + }; + + for m in pending { + msg.push_str(&format!("\n v{}: {}\n", m.version, m.description)); + let stmts = (m.up)(schema); + for stmt in &stmts { + if !stmt.is_empty() { + msg.push_str(&format!(" {stmt}\n")); + } + } + // Include the version-tracking INSERT so operators record the + // migration after applying the DDL manually. + msg.push_str(&format!( + " {versions_insert} ({}, '{}');\n", + m.version, m.description, + )); + } + + msg.push_str(&format!( + "\nTo apply automatically, set `auto_migrate: true` in your {backend} schema config.\n\ + To skip, set `version: {}` to mark your schema as already up to date.", + pending.last().map_or(current, |m| m.version) + )); + + msg +} + +// ── Oracle helpers ───────────────────────────────────────────────────────── + +/// Run pending Oracle migrations on a synchronous `oracle::Connection`. +/// +/// Returns the list of newly applied version numbers. +pub fn run_oracle_migrations( + conn: &oracle::Connection, + schema: &SchemaConfig, + migrations: &[Migration], + config_version: Option, + auto_migrate: bool, +) -> Result, String> { + // Ensure the versions table exists (needed to check current version) + ensure_oracle_versions_table(conn, schema)?; + + let current = oracle_current_version(conn, schema)?; + let skip_up_to = config_version.unwrap_or(0); + let effective_start = current.max(skip_up_to); + + let pending: Vec<&Migration> = migrations + .iter() + .filter(|m| m.version > effective_start) + .collect(); + + if pending.is_empty() { + tracing::debug!(current_version = effective_start, "schema is up to date"); + return Ok(Vec::new()); + } + + // When auto_migrate is off, fail with actionable info + if !auto_migrate { + return Err(pending_migrations_error( + "oracle", + effective_start, + &pending, + schema, + )); + } + + tracing::info!( + current_version = effective_start, + pending = pending.len(), + "applying schema migrations" + ); + + let versions_table = oracle_versions_table(schema); + + let mut applied = Vec::new(); + for migration in pending { + // NOTE: Oracle DDL implicitly commits. If the DDL below succeeds but + // the subsequent INSERT into _schema_versions fails (for a reason + // other than ORA-00001), the schema change is committed without a + // version record. Next startup will re-attempt the migration, so + // all Oracle migrations MUST be idempotent (e.g. use PL/SQL + // EXCEPTION handlers to tolerate "column already exists"). + let stmts = (migration.up)(schema); + for stmt in &stmts { + if stmt.is_empty() { + continue; + } + conn.execute(stmt, &[]).map_err(|e| { + format!( + "migration v{} ({}) failed: {}", + migration.version, migration.description, e + ) + })?; + } + // Record the applied migration. + // ORA-00001 (unique constraint) means another instance already + // applied this migration concurrently — safe to skip since the + // DDL statements above are idempotent. + match conn.execute( + &format!("INSERT INTO {versions_table} (version, description) VALUES (:1, :2)"), + &[&migration.version, &migration.description], + ) { + Ok(_) => { + conn.commit().map_err(|e| format!("commit failed: {e}"))?; + } + Err(e) if e.db_error().is_some_and(|de| de.code() == 1) => { + tracing::info!( + version = migration.version, + "migration already applied by another instance, skipping" + ); + continue; + } + Err(e) => { + return Err(format!( + "failed to record migration v{}: {}", + migration.version, e + )); + } + } + + tracing::info!( + version = migration.version, + description = migration.description, + "applied migration" + ); + applied.push(migration.version); + } + + let final_version = oracle_current_version(conn, schema)?; + tracing::info!( + schema_version = final_version, + "schema version after migrations" + ); + + Ok(applied) +} + +/// Ensure the `_schema_versions` table exists in Oracle. +/// +/// Uses `all_tables` with an owner filter when `schema.owner` is set, +/// falling back to `user_tables` for the current user's schema. +/// +/// The table is always created with a quoted identifier since `_` is not +/// valid as the first character of an unquoted Oracle identifier. This +/// preserves the lowercase name in the catalog. +fn ensure_oracle_versions_table( + conn: &oracle::Connection, + schema: &SchemaConfig, +) -> Result<(), String> { + // The table is always created with a quoted identifier (preserves + // lowercase in Oracle's catalog) since `_` is invalid as the first + // character of an unquoted Oracle identifier. + let check_sql = match &schema.owner { + Some(owner) => format!( + "SELECT COUNT(*) FROM all_tables WHERE owner = '{}' AND table_name = '{VERSIONS_TABLE}'", + owner.to_ascii_uppercase() + ), + None => { + format!("SELECT COUNT(*) FROM user_tables WHERE table_name = '{VERSIONS_TABLE}'") + } + }; + let present: i64 = conn + .query_row_as(&check_sql, &[]) + .map_err(|e| format!("failed to check for {VERSIONS_TABLE} table: {e}"))?; + + if present == 0 { + let ddl = oracle_create_versions_table(schema); + if let Err(err) = conn.execute(&ddl, &[]) { + // ORA-00955: name is already used — another instance created + // the table between our check and this CREATE. Safe to ignore. + if err.db_error().is_some_and(|de| de.code() == 955) { + tracing::debug!("versions table created by concurrent instance, proceeding"); + } else { + return Err(format!("failed to create {VERSIONS_TABLE} table: {err}")); + } + } + conn.commit().map_err(|e| format!("commit failed: {e}"))?; + } + Ok(()) +} + +/// Read the current (highest) schema version from the Oracle versions table. +fn oracle_current_version(conn: &oracle::Connection, schema: &SchemaConfig) -> Result { + let versions_table = oracle_versions_table(schema); + let row: Option = conn + .query_row_as_named(&format!("SELECT MAX(version) FROM {versions_table}"), &[]) + .map_err(|e| format!("failed to read current schema version: {e}"))?; + + Ok(row.unwrap_or(0) as u32) +} + +// ── Postgres helpers ─────────────────────────────────────────────────────── + +/// Run pending Postgres migrations on a `tokio_postgres::Client`. +/// +/// Uses a transaction per migration for atomicity. Returns the list of +/// newly applied version numbers. +pub async fn run_postgres_migrations( + client: &mut tokio_postgres::Client, + schema: &SchemaConfig, + migrations: &[Migration], + config_version: Option, + auto_migrate: bool, +) -> Result, String> { + // Ensure the versions table exists (needed to check current version) + client + .batch_execute(&postgres_create_versions_table()) + .await + .map_err(|e| format!("failed to create {VERSIONS_TABLE} table: {e}"))?; + + let current = postgres_current_version(client).await?; + let skip_up_to = config_version.unwrap_or(0); + let effective_start = current.max(skip_up_to); + + let pending: Vec<&Migration> = migrations + .iter() + .filter(|m| m.version > effective_start) + .collect(); + + if pending.is_empty() { + tracing::debug!(current_version = effective_start, "schema is up to date"); + return Ok(Vec::new()); + } + + // When auto_migrate is off, fail with actionable info + if !auto_migrate { + return Err(pending_migrations_error( + "postgres", + effective_start, + &pending, + schema, + )); + } + + // Acquire session-level advisory lock to serialize migrations across + // concurrent application instances. Released explicitly below (or + // automatically when the connection is closed). + const MIGRATION_LOCK_ID: i64 = 0x736D675F6D696772; // "smg_migr" + client + .execute("SELECT pg_advisory_lock($1)", &[&MIGRATION_LOCK_ID]) + .await + .map_err(|e| format!("failed to acquire migration lock: {e}"))?; + + let result = + apply_postgres_migrations(client, schema, migrations, config_version, effective_start) + .await; + + // Always release the lock, even on error + let _ = client + .execute("SELECT pg_advisory_unlock($1)", &[&MIGRATION_LOCK_ID]) + .await; + + result +} + +/// Inner migration logic, separated so the caller can manage the advisory lock. +/// +/// Re-reads the current version under the advisory lock to handle the case +/// where another instance applied migrations between our initial check and +/// acquiring the lock. +async fn apply_postgres_migrations( + client: &mut tokio_postgres::Client, + schema: &SchemaConfig, + migrations: &[Migration], + config_version: Option, + pre_lock_start: u32, +) -> Result, String> { + // Re-read under lock — another instance may have migrated since our + // initial check (before lock acquisition). + let current = postgres_current_version(client).await?; + let skip_up_to = config_version.unwrap_or(0); + let effective_start = current.max(skip_up_to).max(pre_lock_start); + + let pending: Vec<&Migration> = migrations + .iter() + .filter(|m| m.version > effective_start) + .collect(); + + if pending.is_empty() { + tracing::debug!(current_version = effective_start, "schema is up to date"); + return Ok(Vec::new()); + } + + tracing::info!( + current_version = effective_start, + pending = pending.len(), + "applying schema migrations" + ); + + let mut applied = Vec::new(); + for migration in pending { + // Run each migration in a transaction for atomicity. + // Note: DDL in Postgres is transactional (unlike Oracle). + let tx = client + .transaction() + .await + .map_err(|e| format!("failed to begin transaction: {e}"))?; + + let stmts = (migration.up)(schema); + for stmt in &stmts { + if stmt.is_empty() { + continue; + } + tx.batch_execute(stmt).await.map_err(|e| { + format!( + "migration v{} ({}) failed: {}", + migration.version, migration.description, e + ) + })?; + } + + // Record the applied migration + tx.execute( + &format!("INSERT INTO {VERSIONS_TABLE} (version, description) VALUES ($1, $2)"), + &[&(migration.version as i32), &migration.description], + ) + .await + .map_err(|e| format!("failed to record migration v{}: {}", migration.version, e))?; + + tx.commit() + .await + .map_err(|e| format!("commit failed: {e}"))?; + + tracing::info!( + version = migration.version, + description = migration.description, + "applied migration" + ); + applied.push(migration.version); + } + + let final_version = postgres_current_version(client).await?; + tracing::info!( + schema_version = final_version, + "schema version after migrations" + ); + + Ok(applied) +} + +/// Read the current (highest) schema version from the Postgres versions table. +async fn postgres_current_version(client: &tokio_postgres::Client) -> Result { + let row = client + .query_one( + &format!("SELECT COALESCE(MAX(version), 0) FROM {VERSIONS_TABLE}"), + &[], + ) + .await + .map_err(|e| format!("failed to read current schema version: {e}"))?; + + let version: i32 = row.get(0); + Ok(version as u32) +} + +// ── Tests ────────────────────────────────────────────────────────────────── + +#[cfg(test)] +mod tests { + use super::*; + use crate::schema::TableConfig; + + #[test] + fn oracle_versions_table_name_is_always_quoted() { + let no_owner = SchemaConfig::default(); + assert_eq!(oracle_versions_table(&no_owner), "\"_schema_versions\""); + + let with_owner = SchemaConfig { + owner: Some("ADMIN".to_string()), + ..Default::default() + }; + assert_eq!( + oracle_versions_table(&with_owner), + "ADMIN.\"_schema_versions\"" + ); + } + + #[test] + fn oracle_versions_table_ddl_without_owner() { + let schema = SchemaConfig::default(); + let ddl = oracle_create_versions_table(&schema); + assert!( + ddl.contains("\"_schema_versions\""), + "must be quoted for Oracle: {ddl}" + ); + assert!(ddl.contains("PRIMARY KEY")); + } + + #[test] + fn oracle_versions_table_ddl_with_owner() { + let schema = SchemaConfig { + owner: Some("ADMIN".to_string()), + ..Default::default() + }; + let ddl = oracle_create_versions_table(&schema); + assert!(ddl.contains("ADMIN.\"_schema_versions\""), "got: {ddl}"); + } + + #[test] + fn postgres_versions_table_ddl() { + let ddl = postgres_create_versions_table(); + assert!(ddl.contains("IF NOT EXISTS")); + assert!(ddl.contains("_schema_versions")); + assert!(ddl.contains("PRIMARY KEY")); + } + + #[test] + fn migration_up_respects_schema_config() { + let schema = SchemaConfig { + owner: Some("ADMIN".to_string()), + responses: TableConfig { + table: "MY_RESPONSES".to_string(), + ..Default::default() + }, + ..Default::default() + }; + + let m = Migration { + version: 1, + description: "test", + up: |s| { + let t = s.responses.qualified_table(s.owner.as_deref()); + vec![format!("ALTER TABLE {t} ADD COLUMN x INT")] + }, + }; + let stmts = (m.up)(&schema); + assert!( + stmts[0].contains("ADMIN.\"MY_RESPONSES\""), + "got: {}", + stmts[0] + ); + } + + #[test] + fn pending_migrations_error_includes_sql_and_hints() { + let schema = SchemaConfig::default(); + let migrations = [ + Migration { + version: 1, + description: "add col_x", + up: |_| vec!["ALTER TABLE t ADD COLUMN x INT".to_string()], + }, + Migration { + version: 2, + description: "drop col_y", + up: |_| vec!["ALTER TABLE t DROP COLUMN y".to_string()], + }, + ]; + let pending: Vec<&Migration> = migrations.iter().collect(); + let err = pending_migrations_error("postgres", 0, &pending, &schema); + + assert!(err.contains("v1: add col_x"), "should list v1: {err}"); + assert!(err.contains("v2: drop col_y"), "should list v2: {err}"); + assert!( + err.contains("ALTER TABLE t ADD COLUMN x INT"), + "should include SQL: {err}" + ); + assert!( + err.contains("INSERT INTO _schema_versions"), + "should include version INSERT: {err}" + ); + assert!( + err.contains("auto_migrate: true"), + "should hint auto_migrate: {err}" + ); + assert!( + err.contains("version: 2"), + "should hint version skip: {err}" + ); + } + + #[test] + fn pending_migrations_error_shows_current_version() { + let schema = SchemaConfig::default(); + let migrations = [Migration { + version: 3, + description: "test", + up: |_| vec!["SELECT 1".to_string()], + }]; + let pending: Vec<&Migration> = migrations.iter().collect(); + let err = pending_migrations_error("oracle", 2, &pending, &schema); + + assert!( + err.contains("current version: 2"), + "should show current version: {err}" + ); + assert!( + err.contains("latest version: 3"), + "should show target version: {err}" + ); + } +} diff --git a/scripts/ci_agentic_svc_deps.sh b/scripts/ci_agentic_svc_deps.sh index 296b9174d3..15ae9d4618 100755 --- a/scripts/ci_agentic_svc_deps.sh +++ b/scripts/ci_agentic_svc_deps.sh @@ -109,6 +109,7 @@ PYEOF echo "ATP_USER=$TEST_USER" >> "$GITHUB_ENV" echo "ATP_PASSWORD=$TEST_PASS" >> "$GITHUB_ENV" echo "ATP_DSN=$oracle_dsn" >> "$GITHUB_ENV" + echo "DB_AUTO_MIGRATE=true" >> "$GITHUB_ENV" } cmd_cleanup_oracle_user() {