Skip to main content

risingwave_meta_model_migration/
m20260805_000000_object_belong_to_oid.rs

1use sea_orm_migration::prelude::*;
2
3use crate::m20230908_072257_init::Object;
4use crate::sea_orm::{ConnectionTrait, DatabaseBackend, Statement};
5
6#[derive(DeriveMigrationName)]
7pub struct Migration;
8
9const FK_NAME: &str = "FK_object_belong_to_oid";
10const INDEX_NAME: &str = "IDX_object_belong_to_oid";
11
12#[derive(DeriveIden)]
13enum ObjectColumn {
14    BelongToOid,
15}
16
17#[async_trait::async_trait]
18impl MigrationTrait for Migration {
19    async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> {
20        let backend = manager.get_database_backend();
21        // The migrator records the version only after `up` returns, while MySQL DDL implicitly
22        // commits. Guard every DDL step so a new Meta leader can retry a partially run migration.
23        if !manager.has_column("object", "belong_to_oid").await? {
24            match backend {
25                DatabaseBackend::MySql | DatabaseBackend::Postgres => {
26                    manager
27                        .alter_table(
28                            Table::alter()
29                                .table(Object::Table)
30                                .add_column(
31                                    ColumnDef::new(ObjectColumn::BelongToOid).integer().null(),
32                                )
33                                .to_owned(),
34                        )
35                        .await?;
36                }
37                DatabaseBackend::Sqlite => {
38                    // SQLite cannot add a foreign key constraint to an existing table separately,
39                    // but it allows REFERENCES on a newly added nullable column.
40                    manager
41                        .get_connection()
42                        .execute(Statement::from_string(
43                            backend,
44                            r#"ALTER TABLE "object" ADD COLUMN "belong_to_oid" INTEGER REFERENCES "object" ("oid") ON DELETE CASCADE"#,
45                        ))
46                        .await?;
47                }
48            }
49        }
50
51        if !manager.has_index("object", INDEX_NAME).await? {
52            manager
53                .create_index(
54                    Index::create()
55                        .name(INDEX_NAME)
56                        .table(Object::Table)
57                        .col(ObjectColumn::BelongToOid)
58                        .to_owned(),
59                )
60                .await?;
61        }
62
63        if matches!(backend, DatabaseBackend::MySql | DatabaseBackend::Postgres)
64            && !has_foreign_key(manager).await?
65        {
66            manager
67                .alter_table(
68                    Table::alter()
69                        .table(Object::Table)
70                        .add_foreign_key(
71                            TableForeignKey::new()
72                                .name(FK_NAME)
73                                .from_tbl(Object::Table)
74                                .from_col(ObjectColumn::BelongToOid)
75                                .to_tbl(Object::Table)
76                                .to_col(Object::Oid)
77                                .on_delete(ForeignKeyAction::Cascade),
78                        )
79                        .to_owned(),
80                )
81                .await?;
82        }
83
84        backfill_belong_to_oid(manager).await?;
85
86        Ok(())
87    }
88
89    async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> {
90        if matches!(
91            manager.get_database_backend(),
92            DatabaseBackend::MySql | DatabaseBackend::Postgres
93        ) {
94            manager
95                .alter_table(
96                    Table::alter()
97                        .table(Object::Table)
98                        .drop_foreign_key(Alias::new(FK_NAME))
99                        .to_owned(),
100                )
101                .await?;
102        }
103
104        manager
105            .drop_index(
106                Index::drop()
107                    .name(INDEX_NAME)
108                    .table(Object::Table)
109                    .to_owned(),
110            )
111            .await?;
112        manager
113            .alter_table(
114                Table::alter()
115                    .table(Object::Table)
116                    .drop_column(ObjectColumn::BelongToOid)
117                    .to_owned(),
118            )
119            .await?;
120
121        Ok(())
122    }
123}
124
125async fn has_foreign_key(manager: &SchemaManager<'_>) -> Result<bool, DbErr> {
126    let backend = manager.get_database_backend();
127    let statement = match backend {
128        DatabaseBackend::MySql => format!(
129            "SELECT 1 FROM information_schema.TABLE_CONSTRAINTS \
130             WHERE CONSTRAINT_SCHEMA = DATABASE() AND TABLE_NAME = 'object' \
131             AND CONSTRAINT_NAME = '{FK_NAME}' AND CONSTRAINT_TYPE = 'FOREIGN KEY' LIMIT 1"
132        ),
133        DatabaseBackend::Postgres => format!(
134            "SELECT 1 FROM information_schema.table_constraints \
135             WHERE constraint_schema = current_schema() AND table_name = 'object' \
136             AND constraint_name = '{FK_NAME}' AND constraint_type = 'FOREIGN KEY' LIMIT 1"
137        ),
138        DatabaseBackend::Sqlite => unreachable!("SQLite defines the foreign key with the column"),
139    };
140    Ok(manager
141        .get_connection()
142        .query_one(Statement::from_string(backend, statement))
143        .await?
144        .is_some())
145}
146
147async fn backfill_belong_to_oid(manager: &SchemaManager<'_>) -> Result<(), DbErr> {
148    let backend = manager.get_database_backend();
149    // Start with namespace parents, then overwrite them with more specific belong-to relations:
150    // internal table -> job and associated source -> table. Indexes and subscriptions are named
151    // objects, so they continue to belong to their schemas; their references to upstream tables
152    // remain represented by the `index`, `subscription`, and `object_dependency` tables.
153    //
154    // For implicit Iceberg objects, reconstruct the associations that catalog operations used
155    // before `belong_to_oid` was persisted. A sink with the literal `__iceberg_sink_` prefix belongs
156    // to the unique Iceberg table in its object dependencies. A generated source belongs to the
157    // same-schema Iceberg table
158    // identified by its `__iceberg_source_<table>` name, because no source-to-table dependency is
159    // persisted.
160    let statements = match backend {
161        DatabaseBackend::MySql => vec![
162            "UPDATE `object` SET `belong_to_oid` = COALESCE(`schema_id`, `database_id`)",
163            "UPDATE `object` AS o JOIN `table` AS t ON t.`table_id` = o.`oid` SET o.`belong_to_oid` = t.`belongs_to_job_id` WHERE t.`belongs_to_job_id` IS NOT NULL",
164            "UPDATE `object` AS o JOIN `table` AS t ON t.`optional_associated_source_id` = o.`oid` SET o.`belong_to_oid` = t.`table_id`",
165            "UPDATE `object` AS o JOIN `sink` AS s ON s.`sink_id` = o.`oid` JOIN (SELECT d.`used_by`, MIN(t.`table_id`) AS `table_id` FROM `object_dependency` AS d JOIN `table` AS t ON t.`table_id` = d.`oid` AND t.`engine` = 'ICEBERG' GROUP BY d.`used_by` HAVING COUNT(*) = 1) AS parent ON parent.`used_by` = o.`oid` SET o.`belong_to_oid` = parent.`table_id` WHERE s.`name` LIKE '!_!_iceberg!_sink!_%' ESCAPE '!'",
166            "UPDATE `object` AS o JOIN `source` AS s ON s.`source_id` = o.`oid` JOIN `object` AS parent ON parent.`database_id` = o.`database_id` AND parent.`schema_id` = o.`schema_id` JOIN `table` AS t ON t.`table_id` = parent.`oid` AND t.`engine` = 'ICEBERG' SET o.`belong_to_oid` = parent.`oid` WHERE s.`name` = CONCAT('__iceberg_source_', t.`name`)",
167        ],
168        DatabaseBackend::Postgres => vec![
169            r#"UPDATE "object" SET "belong_to_oid" = COALESCE("schema_id", "database_id")"#,
170            r#"UPDATE "object" AS o SET "belong_to_oid" = t."belongs_to_job_id" FROM "table" AS t WHERE t."table_id" = o."oid" AND t."belongs_to_job_id" IS NOT NULL"#,
171            r#"UPDATE "object" AS o SET "belong_to_oid" = t."table_id" FROM "table" AS t WHERE t."optional_associated_source_id" = o."oid""#,
172            r#"UPDATE "object" AS o SET "belong_to_oid" = parent."table_id" FROM "sink" AS s, (SELECT d."used_by", MIN(t."table_id") AS "table_id" FROM "object_dependency" AS d JOIN "table" AS t ON t."table_id" = d."oid" AND t."engine" = 'ICEBERG' GROUP BY d."used_by" HAVING COUNT(*) = 1) AS parent WHERE s."sink_id" = o."oid" AND parent."used_by" = o."oid" AND s."name" LIKE '!_!_iceberg!_sink!_%' ESCAPE '!'"#,
173            r#"UPDATE "object" AS o SET "belong_to_oid" = parent."oid" FROM "source" AS s, "object" AS parent, "table" AS t WHERE s."source_id" = o."oid" AND parent."database_id" = o."database_id" AND parent."schema_id" = o."schema_id" AND t."table_id" = parent."oid" AND t."engine" = 'ICEBERG' AND s."name" = '__iceberg_source_' || t."name""#,
174        ],
175        DatabaseBackend::Sqlite => vec![
176            r#"UPDATE "object" SET "belong_to_oid" = COALESCE("schema_id", "database_id")"#,
177            r#"UPDATE "object" SET "belong_to_oid" = (SELECT t."belongs_to_job_id" FROM "table" AS t WHERE t."table_id" = "object"."oid") WHERE EXISTS (SELECT 1 FROM "table" AS t WHERE t."table_id" = "object"."oid" AND t."belongs_to_job_id" IS NOT NULL)"#,
178            r#"UPDATE "object" SET "belong_to_oid" = (SELECT t."table_id" FROM "table" AS t WHERE t."optional_associated_source_id" = "object"."oid") WHERE EXISTS (SELECT 1 FROM "table" AS t WHERE t."optional_associated_source_id" = "object"."oid")"#,
179            r#"UPDATE "object" SET "belong_to_oid" = (SELECT MIN(t."table_id") FROM "object_dependency" AS d JOIN "table" AS t ON t."table_id" = d."oid" AND t."engine" = 'ICEBERG' WHERE d."used_by" = "object"."oid") WHERE EXISTS (SELECT 1 FROM "sink" AS s WHERE s."sink_id" = "object"."oid" AND s."name" LIKE '!_!_iceberg!_sink!_%' ESCAPE '!') AND (SELECT COUNT(*) FROM "object_dependency" AS d JOIN "table" AS t ON t."table_id" = d."oid" AND t."engine" = 'ICEBERG' WHERE d."used_by" = "object"."oid") = 1"#,
180            r#"UPDATE "object" SET "belong_to_oid" = (SELECT parent."oid" FROM "source" AS s JOIN "object" AS parent ON parent."database_id" = "object"."database_id" AND parent."schema_id" = "object"."schema_id" JOIN "table" AS t ON t."table_id" = parent."oid" AND t."engine" = 'ICEBERG' WHERE s."source_id" = "object"."oid" AND s."name" = '__iceberg_source_' || t."name") WHERE EXISTS (SELECT 1 FROM "source" AS s JOIN "object" AS parent ON parent."database_id" = "object"."database_id" AND parent."schema_id" = "object"."schema_id" JOIN "table" AS t ON t."table_id" = parent."oid" AND t."engine" = 'ICEBERG' WHERE s."source_id" = "object"."oid" AND s."name" = '__iceberg_source_' || t."name")"#,
181        ],
182    };
183
184    for statement in statements {
185        manager
186            .get_connection()
187            .execute(Statement::from_string(backend, statement))
188            .await?;
189    }
190
191    Ok(())
192}
193
194#[cfg(test)]
195mod tests {
196    use sea_orm::{Database, TryGetable};
197
198    use super::*;
199
200    #[tokio::test]
201    async fn test_sqlite_backfill_and_cascade() {
202        let db = Database::connect("sqlite::memory:").await.unwrap();
203        for sql in [
204            r#"CREATE TABLE "object" ("oid" INTEGER PRIMARY KEY, "schema_id" INTEGER, "database_id" INTEGER)"#,
205            r#"CREATE TABLE "table" ("table_id" INTEGER PRIMARY KEY, "name" TEXT, "engine" TEXT, "belongs_to_job_id" INTEGER, "optional_associated_source_id" INTEGER)"#,
206            r#"CREATE TABLE "index" ("index_id" INTEGER PRIMARY KEY, "primary_table_id" INTEGER)"#,
207            r#"CREATE TABLE "subscription" ("subscription_id" INTEGER PRIMARY KEY, "dependent_table_id" INTEGER)"#,
208            r#"CREATE TABLE "sink" ("sink_id" INTEGER PRIMARY KEY, "name" TEXT)"#,
209            r#"CREATE TABLE "source" ("source_id" INTEGER PRIMARY KEY, "name" TEXT)"#,
210            r#"CREATE TABLE "object_dependency" ("oid" INTEGER, "used_by" INTEGER)"#,
211            r#"INSERT INTO "object" ("oid", "schema_id", "database_id") VALUES (1, NULL, NULL), (2, NULL, 1), (3, 2, 1), (4, 2, 1), (5, 2, 1), (6, 2, 1), (7, 2, 1), (8, 2, 1), (9, 2, 1), (10, 2, 1), (11, 2, 1), (12, 2, 1)"#,
212            r#"INSERT INTO "table" ("table_id", "name", "engine", "belongs_to_job_id", "optional_associated_source_id") VALUES (3, 'iceberg_table', 'ICEBERG', NULL, 5), (4, 'table_internal', 'HUMMOCK', 3, NULL), (8, 'sink_internal', 'HUMMOCK', 6, NULL), (10, 'index_internal', 'HUMMOCK', 9, NULL)"#,
213            r#"INSERT INTO "index" ("index_id", "primary_table_id") VALUES (9, 3)"#,
214            r#"INSERT INTO "subscription" ("subscription_id", "dependent_table_id") VALUES (11, 3)"#,
215            r#"INSERT INTO "source" ("source_id", "name") VALUES (5, 'table_source'), (7, '__iceberg_source_iceberg_table')"#,
216            r#"INSERT INTO "sink" ("sink_id", "name") VALUES (6, '__iceberg_sink_renamed'), (12, 'myiceberg_sink_output')"#,
217            r#"INSERT INTO "object_dependency" ("oid", "used_by") VALUES (3, 6), (3, 12)"#,
218        ] {
219            db.execute(Statement::from_string(DatabaseBackend::Sqlite, sql))
220                .await
221                .unwrap();
222        }
223
224        Migration.up(&SchemaManager::new(&db)).await.unwrap();
225
226        let rows = db
227            .query_all(Statement::from_string(
228                DatabaseBackend::Sqlite,
229                r#"SELECT "oid", "belong_to_oid" FROM "object" ORDER BY "oid""#,
230            ))
231            .await
232            .unwrap();
233        let belong_to_oids = rows
234            .into_iter()
235            .map(|row| {
236                (
237                    i32::try_get(&row, "", "oid").unwrap(),
238                    Option::<i32>::try_get(&row, "", "belong_to_oid").unwrap(),
239                )
240            })
241            .collect::<Vec<_>>();
242        assert_eq!(
243            belong_to_oids,
244            vec![
245                (1, None),
246                (2, Some(1)),
247                (3, Some(2)),
248                (4, Some(3)),
249                (5, Some(3)),
250                (6, Some(3)),
251                (7, Some(3)),
252                (8, Some(6)),
253                (9, Some(2)),
254                (10, Some(9)),
255                (11, Some(2)),
256                (12, Some(2)),
257            ]
258        );
259
260        db.execute(Statement::from_string(
261            DatabaseBackend::Sqlite,
262            r#"DELETE FROM "object" WHERE "oid" = 3"#,
263        ))
264        .await
265        .unwrap();
266        let remaining = db
267            .query_one(Statement::from_string(
268                DatabaseBackend::Sqlite,
269                r#"SELECT COUNT(*) AS "count" FROM "object" WHERE "oid" IN (3, 4, 5, 6, 7, 8, 9, 10, 11, 12)"#,
270            ))
271            .await
272            .unwrap()
273            .unwrap();
274        assert_eq!(i64::try_get(&remaining, "", "count").unwrap(), 4);
275    }
276
277    #[tokio::test]
278    async fn test_sqlite_partial_run_retry() {
279        let db = Database::connect("sqlite::memory:").await.unwrap();
280        for sql in [
281            r#"CREATE TABLE "object" ("oid" INTEGER PRIMARY KEY, "schema_id" INTEGER, "database_id" INTEGER)"#,
282            r#"CREATE TABLE "table" ("table_id" INTEGER PRIMARY KEY, "name" TEXT, "engine" TEXT, "belongs_to_job_id" INTEGER, "optional_associated_source_id" INTEGER)"#,
283            r#"CREATE TABLE "sink" ("sink_id" INTEGER PRIMARY KEY, "name" TEXT)"#,
284            r#"CREATE TABLE "source" ("source_id" INTEGER PRIMARY KEY, "name" TEXT)"#,
285            r#"CREATE TABLE "object_dependency" ("oid" INTEGER, "used_by" INTEGER)"#,
286            r#"INSERT INTO "object" ("oid", "schema_id", "database_id") VALUES (1, NULL, NULL), (2, NULL, 1), (3, 2, 1), (4, 2, 1)"#,
287            r#"INSERT INTO "table" ("table_id", "name", "engine", "belongs_to_job_id", "optional_associated_source_id") VALUES (3, 'table', 'HUMMOCK', NULL, NULL), (4, 'internal', 'HUMMOCK', 3, NULL)"#,
288            // Simulate Meta exiting after the first DDL statement but before recording the
289            // migration version.
290            r#"ALTER TABLE "object" ADD COLUMN "belong_to_oid" INTEGER REFERENCES "object" ("oid") ON DELETE CASCADE"#,
291        ] {
292            db.execute(Statement::from_string(DatabaseBackend::Sqlite, sql))
293                .await
294                .unwrap();
295        }
296
297        let manager = SchemaManager::new(&db);
298        Migration.up(&manager).await.unwrap();
299        // Simulate another exit after all migration statements but before recording the version.
300        Migration.up(&manager).await.unwrap();
301
302        assert!(manager.has_column("object", "belong_to_oid").await.unwrap());
303        assert!(manager.has_index("object", INDEX_NAME).await.unwrap());
304        let internal = db
305            .query_one(Statement::from_string(
306                DatabaseBackend::Sqlite,
307                r#"SELECT "belong_to_oid" FROM "object" WHERE "oid" = 4"#,
308            ))
309            .await
310            .unwrap()
311            .unwrap();
312        assert_eq!(i32::try_get(&internal, "", "belong_to_oid").unwrap(), 3);
313
314        db.execute(Statement::from_string(
315            DatabaseBackend::Sqlite,
316            r#"DELETE FROM "object" WHERE "oid" = 3"#,
317        ))
318        .await
319        .unwrap();
320        assert!(
321            db.query_one(Statement::from_string(
322                DatabaseBackend::Sqlite,
323                r#"SELECT "oid" FROM "object" WHERE "oid" = 4"#,
324            ))
325            .await
326            .unwrap()
327            .is_none()
328        );
329    }
330}