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 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 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 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 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 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}