Skip to content
Merged
65 changes: 62 additions & 3 deletions data_connector/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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:
Expand Down Expand Up @@ -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):
Expand Down
11 changes: 10 additions & 1 deletion data_connector/src/factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -178,10 +178,19 @@ async fn create_postgres_storage(postgres_cfg: &PostgresConfig) -> Result<Storag
let postgres_conv = PostgresConversationStorage::new(store.clone())
.await
.map_err(|err| format!("failed to initialize Postgres conversation storage: {err}"))?;
let postgres_item = PostgresConversationItemStorage::new(store)
let postgres_item = PostgresConversationItemStorage::new(store.clone())
.await
.map_err(|err| format!("failed to initialize Postgres conversation item storage: {err}"))?;

// Run versioned migrations after all tables are created
let applied = store.run_migrations().await?;

// Re-create indexes that were deferred during init because
// migration-added columns did not yet exist.
if !applied.is_empty() {
store.ensure_response_indexes().await?;
}
Comment on lines +190 to +192

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Retry deferred Postgres index creation unconditionally

PostgresResponseStorage::new now downgrades responses_safety_idx creation failures to a debug log when the migrated column is not present yet, and this function retries index creation only when run_migrations() reports locally applied versions. In concurrent startup, another instance can apply the migrations first (so applied is empty here), which skips this retry path even though this process already deferred index creation earlier; that leaves safety_identifier queries running without the intended index until a later restart.

Useful? React with 👍 / 👎.


Ok((
Arc::new(postgres_resp),
Arc::new(postgres_conv),
Expand Down
3 changes: 3 additions & 0 deletions data_connector/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,12 @@ pub mod hooks;
mod memory;
mod noop;
mod oracle;
mod oracle_migrations;
mod postgres;
mod postgres_migrations;
mod redis;
pub mod schema;
pub(crate) mod versioning;

// Re-export config types
// Re-export core types and traits
Expand Down
115 changes: 23 additions & 92 deletions data_connector/src/oracle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,10 @@ use crate::{
},
config::OracleConfig,
context::current_extra_columns,
oracle_migrations::ORACLE_MIGRATIONS,
schema::SchemaConfig,
};

// ============================================================================
// PART 1: OracleStore Helper + Common Utilities
// ============================================================================
Expand Down Expand Up @@ -79,6 +81,25 @@ impl OracleStore {
for init_schema in init_schemas {
init_schema(&conn, &schema)?;
}

// Run versioned migrations after table creation
let applied = crate::versioning::run_oracle_migrations(
&conn,
&schema,
&ORACLE_MIGRATIONS,
schema.version,
schema.auto_migrate,
)?;
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// Re-run init_schemas when migrations were applied so that indexes
// on newly-added columns (e.g. safety_identifier added by v1) are
// created in this startup cycle rather than requiring a restart.
if !applied.is_empty() {
for init_schema in init_schemas {
init_schema(&conn, &schema)?;
}
}

drop(conn);

// Create connection pool
Expand Down Expand Up @@ -1168,11 +1189,6 @@ impl OracleResponseStorage {
&[],
)
.map_err(map_oracle_error)?;
} else {
if !s.is_skipped("safety_identifier") {
Self::alter_safety_identifier_column(conn, schema)?;
}
Self::remove_user_id_column_if_exists(conn, schema)?;
}

if !s.is_skipped("previous_response_id") {
Expand Down Expand Up @@ -1200,92 +1216,6 @@ impl OracleResponseStorage {
Ok(())
}

fn alter_safety_identifier_column(
conn: &Connection,
schema: &SchemaConfig,
) -> 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<StoredResponse, String> {
let s = &schema.responses;
let col_id = s.col("id");
Expand Down Expand Up @@ -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 {
Comment thread
slin1237 marked this conversation as resolved.
Comment thread
slin1237 marked this conversation as resolved.
return Err(map_oracle_error(err));
}
} else {
Expand Down
Loading