use anyhow::Result;
use sqlx::{Acquire, Row};
use super::{DatabaseKind, Store};
#[cfg_attr(not(test), allow(dead_code))]
fn is_safe_identifier(s: &str) -> bool {
!s.is_empty() && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
}
impl Store {
pub async fn get_current_version(&self) -> Result<i32> {
let row = sqlx::query("SELECT MAX(version) FROM schema_version WHERE version < $1")
.bind(i32::MAX)
.fetch_optional(&self.pool)
.await?;
match row {
Some(r) => {
let val: Option<i32> = r.try_get(0)?;
Ok(val.unwrap_or(0))
}
None => Ok(0),
}
}
pub async fn is_migration_run(&self, version: i32) -> Result<bool> {
let row = sqlx::query("SELECT 1 FROM schema_version WHERE version = $1")
.bind(version)
.fetch_optional(&self.pool)
.await?;
Ok(row.is_some())
}
pub async fn try_claim_migration(&self, version: i32) -> Result<bool> {
match self.kind {
DatabaseKind::Sqlite | DatabaseKind::Postgres => {
let row = sqlx::query(
"INSERT INTO schema_version (version) VALUES ($1) \
ON CONFLICT (version) DO NOTHING RETURNING version",
)
.bind(version)
.fetch_optional(&self.pool)
.await?;
Ok(row.is_some())
}
DatabaseKind::MySql => {
let result = sqlx::query("INSERT IGNORE INTO schema_version (version) VALUES (?)")
.bind(version)
.execute(&self.pool)
.await?;
if result.rows_affected() == 1 {
return Ok(true);
}
let row = sqlx::query("SELECT 1 FROM schema_version WHERE version = ?")
.bind(version)
.fetch_optional(&self.pool)
.await?;
if row.is_none() {
return Err(anyhow::anyhow!(
"INSERT IGNORE for schema_version v{version}: row absent after insert — \
a non-duplicate error may have been silently suppressed"
));
}
Ok(false)
}
}
}
pub async fn record_migration_run(&self, version: i32) -> Result<()> {
match self.kind {
DatabaseKind::Sqlite | DatabaseKind::Postgres => {
sqlx::query(
"INSERT INTO schema_version (version) VALUES ($1) \
ON CONFLICT (version) DO NOTHING",
)
.bind(version)
.execute(&self.pool)
.await?;
}
DatabaseKind::MySql => {
sqlx::query("INSERT IGNORE INTO schema_version (version) VALUES (?)")
.bind(version)
.execute(&self.pool)
.await?;
}
}
Ok(())
}
pub(super) async fn try_claim_migration_in_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
version: i32,
) -> Result<bool> {
match self.kind {
DatabaseKind::Sqlite | DatabaseKind::Postgres => {
let row = sqlx::query(
"INSERT INTO schema_version (version) VALUES ($1) \
ON CONFLICT (version) DO NOTHING RETURNING version",
)
.bind(version)
.fetch_optional(&mut **tx)
.await?;
Ok(row.is_some())
}
DatabaseKind::MySql => {
let result = sqlx::query("INSERT IGNORE INTO schema_version (version) VALUES (?)")
.bind(version)
.execute(&mut **tx)
.await?;
if result.rows_affected() == 1 {
return Ok(true);
}
let row = sqlx::query("SELECT 1 FROM schema_version WHERE version = ?")
.bind(version)
.fetch_optional(&mut **tx)
.await?;
if row.is_none() {
return Err(anyhow::anyhow!(
"INSERT IGNORE for schema_version v{version}: row absent after insert — \
a non-duplicate error may have been silently suppressed"
));
}
Ok(false)
}
}
}
pub(super) async fn run_migrations(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS schema_version (
version INTEGER PRIMARY KEY,
migrated_at TEXT NOT NULL DEFAULT (CURRENT_TIMESTAMP)
)",
)
.execute(&mut *tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: CREATE TABLE IF NOT EXISTS schema_version\n Error: {e}")
})?;
self.migrate_to_v1_tx(&mut tx).await?;
self.migrate_to_v2_tx(&mut tx).await?;
self.migrate_to_v3_tx(&mut tx).await?;
self.migrate_to_v4_tx(&mut tx).await?;
self.migrate_to_v5_tx(&mut tx).await?;
self.migrate_to_v6_tx(&mut tx).await?;
self.migrate_to_v7_tx(&mut tx).await?;
self.migrate_to_v8_tx(&mut tx).await?;
self.migrate_to_v9_tx(&mut tx).await?;
self.migrate_to_v10_tx(&mut tx).await?;
self.migrate_to_v11_tx(&mut tx).await?;
self.migrate_to_v12_tx(&mut tx).await?;
self.migrate_to_v13_tx(&mut tx).await?;
self.migrate_to_v14_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v1_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 1).await? {
return Ok(());
}
for sql in V1_DDLS {
sqlx::query(sql)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {sql}\n Error: {e}"))?;
}
tracing::info!(version = 1, "migration applied: initial schema");
Ok(())
}
pub(super) async fn migrate_to_v2_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 2).await? {
return Ok(());
}
for sql in V2_DDLS {
sqlx::query(sql)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {sql}\n Error: {e}"))?;
}
tracing::info!(version = 2, "migration applied: backend session tables");
Ok(())
}
pub(super) async fn migrate_to_v3_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 3).await? {
return Ok(());
}
sqlx::query(V3_CREATE_IDX)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {V3_CREATE_IDX}\n Error: {e}"))?;
tracing::info!(version = 3, "migration applied: real_ctx unique index");
Ok(())
}
pub(super) async fn migrate_to_v4_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 4).await? {
return Ok(());
}
let col_exists_in_tx = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM pragma_table_info('context_token_map') \
WHERE name = 'created_at'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.columns \
WHERE table_name = 'context_token_map' AND column_name = 'created_at' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if !col_exists_in_tx {
sqlx::query(V4_ALTER_ADD_CREATED_AT)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V4_ALTER_ADD_CREATED_AT}\n Error: {e}")
})?;
} else {
tracing::debug!(
"v4 migration: created_at column already present (pre-check), skipping ALTER"
);
}
sqlx::query(V4_CREATE_IDX)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {V4_CREATE_IDX}\n Error: {e}"))?;
tracing::info!(
version = 4,
"migration applied: context_token_map created_at column + index"
);
Ok(())
}
pub(super) async fn migrate_to_v5_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 5).await? {
return Ok(());
}
let create_messages = Self::v5_create_messages_sql(self.kind);
sqlx::query(&create_messages)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed (CREATE TABLE messages): {e}"))?;
for sql in V5_INDEXES {
sqlx::query(sql)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {sql}\n Error: {e}"))?;
}
tracing::info!(version = 5, "migration applied: messages table + indexes");
Ok(())
}
pub(super) async fn migrate_to_v6_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 6).await? {
return Ok(());
}
sqlx::query(V6_NORMALIZE_PEER_USER_ID)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V6_NORMALIZE_PEER_USER_ID}\n Error: {e}")
})?;
tracing::info!(
version = 6,
"migration applied: normalize context_token_map.peer_user_id to peer:/group: form"
);
Ok(())
}
pub(super) async fn migrate_to_v7_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 7).await? {
return Ok(());
}
match self.kind {
DatabaseKind::Sqlite | DatabaseKind::Postgres => {
let dedup_sql = match self.kind {
DatabaseKind::Sqlite => V7_DEDUP_PEER_USER_ID_SQLITE,
_ => V7_DEDUP_PEER_USER_ID_POSTGRES,
};
sqlx::query(dedup_sql)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: de-duplicate peer_user_id\n Error: {e}")
})?;
sqlx::query(V7_UNIQUE_IDX_PEER_USER_ID)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V7_UNIQUE_IDX_PEER_USER_ID}\n Error: {e}")
})?;
tracing::info!(
version = 7,
"migration applied: de-dup + unique index on non-empty peer_user_id"
);
}
DatabaseKind::MySql => {
tracing::info!(
version = 7,
"migration skipped on MySQL: partial unique indexes are not supported; \
find_or_create_vctx uses serialised pool instead"
);
}
}
Ok(())
}
pub(super) async fn migrate_to_v8_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 8).await? {
return Ok(());
}
let first_vtoken_row = sqlx::query("SELECT vtoken FROM clients LIMIT 1")
.fetch_optional(&mut **tx)
.await
.or_else(|e| {
let msg = e.to_string();
if msg.contains("no such table") || msg.contains("doesn't exist") {
Ok(None)
} else {
Err(e)
}
})?;
if let Some(row) = first_vtoken_row {
let vtoken: String = row.try_get(0)?;
if crate::hub::is_vtoken_hash(&vtoken) {
tracing::debug!(
"v8 migration: first client vtoken is already hashed, skipping data conversion"
);
return Ok(());
}
} else {
sqlx::query(V8_DDL)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {V8_DDL}\n Error: {e}"))?;
return Ok(());
}
let local_key: Option<crate::runtime::crypto::Key>;
let master_key: &ring::aead::LessSafeKey = if let Some(k) = self.master_key() {
k.as_ref()
} else {
match crate::runtime::crypto::load_or_derive_master_key() {
Ok(k) => {
local_key = Some(k);
local_key.as_ref().expect("just assigned above")
}
Err(e) => {
return Err(anyhow::anyhow!(
"Migration v8 failed: ILINK_HUB_MASTER_KEY is required to encrypt bot credentials, but could not be loaded: {e}"
));
}
}
};
sqlx::query(V8_DDL)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {V8_DDL}\n Error: {e}"))?;
let mut vtokens_to_hash = std::collections::HashSet::new();
let clients_exist = match sqlx::query("SELECT vtoken FROM clients")
.fetch_all(&mut **tx)
.await
{
Ok(rows) => {
for r in rows {
let vt: String = r.try_get(0)?;
if !crate::hub::is_vtoken_hash(&vt) {
vtokens_to_hash.insert(vt);
}
}
true
}
Err(e) => {
let msg = e.to_string();
if msg.contains("no such table") || msg.contains("doesn't exist") {
false
} else {
return Err(e.into());
}
}
};
if clients_exist {
match sqlx::query("SELECT active_vtoken FROM routing_state")
.fetch_all(&mut **tx)
.await
{
Ok(rows) => {
for r in rows {
let vt: String = r.try_get(0)?;
if !crate::hub::is_vtoken_hash(&vt) {
vtokens_to_hash.insert(vt);
}
}
}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") && !msg.contains("doesn't exist") {
return Err(e.into());
}
}
}
match sqlx::query("SELECT DISTINCT vtoken FROM messages WHERE vtoken IS NOT NULL")
.fetch_all(&mut **tx)
.await
{
Ok(rows) => {
for r in rows {
let vt: String = r.try_get(0)?;
if !crate::hub::is_vtoken_hash(&vt) {
vtokens_to_hash.insert(vt);
}
}
}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") && !msg.contains("doesn't exist") {
return Err(e.into());
}
}
}
}
for old_vtoken in vtokens_to_hash {
let new_vtoken = crate::hub::hash_vtoken(&old_vtoken);
if clients_exist {
sqlx::query("UPDATE clients SET vtoken = $1 WHERE vtoken = $2")
.bind(&new_vtoken)
.bind(&old_vtoken)
.execute(&mut **tx)
.await?;
match sqlx::query(
"UPDATE routing_state SET active_vtoken = $1 WHERE active_vtoken = $2",
)
.bind(&new_vtoken)
.bind(&old_vtoken)
.execute(&mut **tx)
.await
{
Ok(_) => {}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") && !msg.contains("doesn't exist") {
return Err(e.into());
}
}
}
match sqlx::query("UPDATE messages SET vtoken = $1 WHERE vtoken = $2")
.bind(&new_vtoken)
.bind(&old_vtoken)
.execute(&mut **tx)
.await
{
Ok(_) => {}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") && !msg.contains("doesn't exist") {
return Err(e.into());
}
}
}
}
}
let cred_rows = match sqlx::query("SELECT id, token FROM bot_credentials")
.fetch_all(&mut **tx)
.await
{
Ok(rows) => Some(rows),
Err(e) => {
let msg = e.to_string();
if msg.contains("no such table") || msg.contains("doesn't exist") {
None
} else {
return Err(e.into());
}
}
};
if let Some(rows) = cred_rows {
for row in rows {
let id: i64 = row.try_get(0)?;
let token: String = row.try_get(1)?;
let is_encrypted =
crate::runtime::crypto::decrypt_token(&token, master_key).is_ok();
if !is_encrypted {
let encrypted = crate::runtime::crypto::encrypt_token(&token, master_key)?;
sqlx::query("UPDATE bot_credentials SET token = $1 WHERE id = $2")
.bind(encrypted)
.bind(id)
.execute(&mut **tx)
.await?;
}
}
}
tracing::info!(
version = 8,
"migration applied: vtoken hashed and bot_credentials encrypted"
);
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v1(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v1_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v2(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v2_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v3(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v3_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v4(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v4_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v5(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v5_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v6(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v6_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v7(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v7_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v8(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v8_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v9_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 9).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => {
sqlx::query("SELECT 1 FROM sqlite_master WHERE type='table' AND name='messages'")
.fetch_optional(&mut **tx)
.await?
.is_some()
}
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables \
WHERE table_name = 'messages' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
sqlx::query(V9_MESSAGES_LOOKUP_IDX)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V9_MESSAGES_LOOKUP_IDX}\n Error: {e}")
})?;
tracing::info!(version = 9, "migration applied: messages lookup index");
} else {
tracing::debug!(
"v9 migration: messages table absent (partial schema), skipping index creation"
);
}
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v9(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v9_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v10_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 10).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => {
sqlx::query("SELECT 1 FROM sqlite_master WHERE type='table' AND name='clients'")
.fetch_optional(&mut **tx)
.await?
.is_some()
}
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables \
WHERE table_name = 'clients' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
for ddl in V10_DDLS {
sqlx::query(ddl)
.execute(&mut **tx)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {ddl}\n Error: {e}"))?;
}
tracing::info!(version = 10, "migration applied: client persona columns");
} else {
tracing::debug!(
"v10 migration: clients table absent (partial schema), skipping persona columns"
);
}
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v10(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v10_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v11_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 11).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name='active_sessions'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables \
WHERE table_name = 'active_sessions' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
let col_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM pragma_table_info('active_sessions') WHERE name = 'a2a_depth'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.columns \
WHERE table_name = 'active_sessions' AND column_name = 'a2a_depth' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if !col_exists {
sqlx::query(V11_ADD_A2A_DEPTH)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V11_ADD_A2A_DEPTH}\n Error: {e}")
})?;
}
tracing::info!(version = 11, "migration applied: active_sessions.a2a_depth");
} else {
tracing::debug!(
"v11 migration: active_sessions table absent (partial schema), skipping a2a_depth column"
);
}
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v11(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v11_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v12_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 12).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => {
sqlx::query("SELECT 1 FROM sqlite_master WHERE type='table' AND name='clients'")
.fetch_optional(&mut **tx)
.await?
.is_some()
}
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables WHERE table_name = 'clients' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
let col_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM pragma_table_info('clients') WHERE name = 'description'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.columns \
WHERE table_name = 'clients' AND column_name = 'description' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if !col_exists {
sqlx::query(V12_ADD_DESCRIPTION)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V12_ADD_DESCRIPTION}\n Error: {e}")
})?;
}
tracing::info!(version = 12, "migration applied: clients.description");
} else {
tracing::debug!(
"v12 migration: clients table absent (partial schema), skipping description column"
);
}
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v12(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v12_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) async fn migrate_to_v13_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 13).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => {
sqlx::query("SELECT 1 FROM sqlite_master WHERE type='table' AND name='messages'")
.fetch_optional(&mut **tx)
.await?
.is_some()
}
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables WHERE table_name = 'messages' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
let col_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM pragma_table_info('messages') WHERE name = 'ilink_msg_id'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.columns \
WHERE table_name = 'messages' AND column_name = 'ilink_msg_id' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if !col_exists {
sqlx::query(V13_ADD_ILINK_MSG_ID)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V13_ADD_ILINK_MSG_ID}\n Error: {e}")
})?;
}
tracing::info!(version = 13, "migration applied: messages.ilink_msg_id");
} else {
tracing::debug!(
"v13 migration: messages table absent (partial schema), skipping ilink_msg_id column"
);
}
Ok(())
}
pub(super) async fn migrate_to_v14_tx(
&self,
tx: &mut sqlx::Transaction<'_, sqlx::Any>,
) -> Result<()> {
if !self.try_claim_migration_in_tx(tx, 14).await? {
return Ok(());
}
let table_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name='backend_sessions_v2'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.tables WHERE table_name = 'backend_sessions_v2' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if table_exists {
let col_exists = match self.kind {
DatabaseKind::Sqlite => sqlx::query(
"SELECT 1 FROM pragma_table_info('backend_sessions_v2') WHERE name = 'last_usage_json'",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
DatabaseKind::Postgres | DatabaseKind::MySql => sqlx::query(
"SELECT 1 FROM information_schema.columns WHERE table_name = 'backend_sessions_v2' AND column_name = 'last_usage_json' LIMIT 1",
)
.fetch_optional(&mut **tx)
.await?
.is_some(),
};
if !col_exists {
sqlx::query(V14_ADD_SESSION_USAGE)
.execute(&mut **tx)
.await
.map_err(|e| {
anyhow::anyhow!("DDL failed: {V14_ADD_SESSION_USAGE}\n Error: {e}")
})?;
}
tracing::info!(
version = 14,
"migration applied: backend_sessions_v2.last_usage_json"
);
} else {
tracing::debug!(
"v14 migration: backend_sessions_v2 absent (partial schema), skipping last_usage_json"
);
}
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v13(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v13_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
#[allow(dead_code)]
pub(super) async fn migrate_to_v14(&self) -> Result<()> {
let mut conn = self.pool.acquire().await?;
let mut tx = conn.begin().await?;
self.migrate_to_v14_tx(&mut tx).await?;
tx.commit().await?;
Ok(())
}
pub(super) fn v5_create_messages_sql(kind: DatabaseKind) -> String {
let id_clause = match kind {
DatabaseKind::Sqlite => "id INTEGER PRIMARY KEY AUTOINCREMENT",
DatabaseKind::Postgres => {
"id INTEGER GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY"
}
DatabaseKind::MySql => "id BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY",
};
format!(
"CREATE TABLE IF NOT EXISTS messages (
{id_clause},
vctx TEXT NOT NULL,
vtoken TEXT,
session_name TEXT NOT NULL DEFAULT 'default',
peer_user_id TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (CURRENT_TIMESTAMP)
)"
)
}
#[cfg_attr(not(test), allow(dead_code))]
pub(super) async fn column_exists(&self, table: &str, column: &str) -> Result<bool> {
if !is_safe_identifier(table) || !is_safe_identifier(column) {
return Ok(false);
}
match self.kind {
DatabaseKind::Sqlite => {
let pragma_sql =
format!("SELECT 1 FROM pragma_table_info('{table}') WHERE name = '{column}'");
let row = sqlx::query(&pragma_sql)
.fetch_optional(&self.pool)
.await
.unwrap_or(None);
Ok(row.is_some())
}
DatabaseKind::Postgres | DatabaseKind::MySql => {
let row = sqlx::query(
"SELECT 1 FROM information_schema.columns \
WHERE table_name = $1 AND column_name = $2 LIMIT 1",
)
.bind(table)
.bind(column)
.fetch_optional(&self.pool)
.await?;
Ok(row.is_some())
}
}
}
#[allow(dead_code)]
pub(super) async fn ddl(&self, sql: &str) -> Result<()> {
let mut conn = self.pool.acquire().await?;
sqlx::query(sql)
.execute(&mut *conn)
.await
.map_err(|e| anyhow::anyhow!("DDL failed: {sql}\n Error: {e}"))?;
Ok(())
}
}
const V1_DDLS: &[&str] = &[include_str!("../../migrations/0001_initial_schema.sql")];
const V2_DDLS: &[&str] = &[include_str!("../../migrations/0002_backend_sessions.sql")];
const V3_CREATE_IDX: &str = include_str!("../../migrations/0003_context_token_real_ctx_index.sql");
const V4_ALTER_ADD_CREATED_AT: &str = "ALTER TABLE context_token_map ADD COLUMN created_at TEXT";
const V4_CREATE_IDX: &str = "CREATE INDEX IF NOT EXISTS idx_context_token_map_created_at \
ON context_token_map (created_at DESC)";
const V5_INDEXES: &[&str] = &[
"CREATE INDEX IF NOT EXISTS idx_messages_vctx_created \
ON messages (vctx, created_at DESC)",
"CREATE INDEX IF NOT EXISTS idx_messages_peer_role_created \
ON messages (peer_user_id, role, created_at DESC)",
];
const V6_NORMALIZE_PEER_USER_ID: &str =
include_str!("../../migrations/0006_normalize_peer_user_id.sql");
const V7_DEDUP_PEER_USER_ID_SQLITE: &str = "DELETE FROM context_token_map \
WHERE peer_user_id != '' \
AND rowid NOT IN ( \
SELECT MIN(rowid) FROM context_token_map \
WHERE peer_user_id != '' \
GROUP BY peer_user_id \
)";
const V7_DEDUP_PEER_USER_ID_POSTGRES: &str = "DELETE FROM context_token_map \
WHERE peer_user_id != '' \
AND ctid NOT IN ( \
SELECT MIN(ctid) FROM context_token_map \
WHERE peer_user_id != '' \
GROUP BY peer_user_id \
)";
const V7_UNIQUE_IDX_PEER_USER_ID: &str =
include_str!("../../migrations/0007_peer_user_id_unique_index.sql");
const V8_DDL: &str = include_str!("../../migrations/0008_vtoken_and_bot_token_hash.sql");
const V9_MESSAGES_LOOKUP_IDX: &str =
include_str!("../../migrations/0009_messages_lookup_index.sql");
const V10_DDLS: &[&str] = &[include_str!("../../migrations/0010_client_persona.sql")];
const V11_ADD_A2A_DEPTH: &str = include_str!("../../migrations/0011_a2a_depth.sql");
const V12_ADD_DESCRIPTION: &str = include_str!("../../migrations/0012_client_description.sql");
const V13_ADD_ILINK_MSG_ID: &str = include_str!("../../migrations/0013_messages_ilink_msg_id.sql");
const V14_ADD_SESSION_USAGE: &str = include_str!("../../migrations/0014_backend_session_usage.sql");