From 6333c3088ad849a932894c4d1d47ed128cff5b42 Mon Sep 17 00:00:00 2001 From: mengw15 <125719918+mengw15@users.noreply.github.com> Date: Tue, 18 Aug 2026 00:26:07 -0700 Subject: [PATCH 1/7] feat(amber): decouple the Lakekeeper catalog name from the user-facing name The catalog name was derived from the display name: user--. That one string is also the REST catalog prefix, the S3 key prefix and a component of every result URI an execution wrote into the warehouse, so the display name was frozen at creation and a warehouse could never be renamed -- unlike a computing unit, whose name is pure display metadata because cuid is the identity everywhere else. Derive the catalog name from the row id instead: user--. The id is drawn from the table's sequence before the Lakekeeper call, so the creation order is unchanged (Lakekeeper first, row after, with the compensating delete) and no schema or nullability change is needed. The sequence is resolved through pg_get_serial_sequence rather than named literally, because the generated name is not a stable contract -- the jOOQ output already carries both user_warehouse_whid_seq and ..._seq1. A sequence-derived name also cannot collide, so it needs no retry path. A random suffix would have needed one, and that retry would have to recognise Lakekeeper's name-conflict error -- the same brittle response parsing #7742 just had to harden. Rename user_warehouse.warehouse_name to lakekeeper_warehouse_name to sit beside lakekeeper_warehouse_id; the table already has its own `name` column, and the value is no longer a name in any user-facing sense. The wire DTO keeps warehouseName. The table is empty in every deployment while the flag is off, so the rename carries no data -- but it still needs a schema migration (sql/updates/38.sql): texera_ddl.sql is CREATE TABLE IF NOT EXISTS, and jOOQ generates from the live database. Closes #7753. --- .../user/warehouse/WarehouseResource.scala | 23 ++++++++-- .../texera/web/service/WorkflowService.scala | 2 +- .../warehouse/WarehouseResourceSpec.scala | 46 +++++++++++++++++-- .../WorkflowExecutionsResourceSpec.scala | 2 +- ...ExecutionsMetadataPersistServiceSpec.scala | 4 +- .../WorkflowServiceWarehouseSpec.scala | 2 +- .../apache/texera/dao/UserWarehouseSpec.scala | 4 +- sql/changelog.xml | 5 ++ sql/texera_ddl.sql | 20 ++++---- sql/updates/38.sql | 35 ++++++++++++++ 10 files changed, 117 insertions(+), 26 deletions(-) create mode 100644 sql/updates/38.sql diff --git a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResource.scala b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResource.scala index d4aed3ccdcf..b8dccff8894 100644 --- a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResource.scala +++ b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResource.scala @@ -30,6 +30,7 @@ import org.apache.texera.dao.jooq.generated.enums.UserWarehouseFlavorEnum import org.apache.texera.dao.jooq.generated.tables.records.UserWarehouseRecord import org.apache.texera.web.resource.dashboard.user.warehouse.WarehouseResource._ import org.apache.texera.web.service.LakekeeperClient +import org.jooq.impl.DSL import javax.annotation.security.RolesAllowed import javax.ws.rs._ @@ -60,7 +61,7 @@ object WarehouseResource { DashboardWarehouse( row.getWhid, row.getName, - row.getWarehouseName, + row.getLakekeeperWarehouseName, row.getFlavor.getLiteral, row.getCreatedAt.toInstant.toEpochMilli ) @@ -131,20 +132,34 @@ class WarehouseResource(client: LakekeeperClient, enabled: Boolean) extends Lazy throw new WebApplicationException(s"a warehouse named '$name' already exists", 409) } - val warehouseName = s"user-$uid-$name" + // The catalog name is derived from the row's own id, never from `name`: that one + // string is also the REST catalog prefix, the S3 key prefix and a component of every + // result URI an execution wrote, so deriving it from a user-facing name would freeze + // that name forever (#7753). Take the id from the sequence up front so the creation + // order below is unchanged -- Lakekeeper first, row after, with the compensating + // delete. The sequence is resolved from the catalog rather than named literally, + // because its generated name is not a stable contract. + val whid: Integer = context.fetchValue( + DSL.field( + "nextval(pg_get_serial_sequence('texera_db.user_warehouse','whid'))", + classOf[Integer] + ) + ) + val lakekeeperWarehouseName = s"user-$uid-$whid" // Create in Lakekeeper first, record after: a failed creation leaves no orphaned row. val warehouseId = try { - client.createWarehouse(warehouseName) + client.createWarehouse(lakekeeperWarehouseName) } catch { case e: Exception => throw new WebApplicationException(e.getMessage, 502) } val row = context.newRecord(USER_WAREHOUSE) + row.setWhid(whid) row.setUid(uid) row.setName(name) - row.setWarehouseName(warehouseName) + row.setLakekeeperWarehouseName(lakekeeperWarehouseName) row.setLakekeeperWarehouseId(warehouseId) row.setFlavor(UserWarehouseFlavorEnum.local) row.setS3Bucket(StorageConfig.icebergRESTCatalogS3Bucket) diff --git a/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala b/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala index a1cc08727b4..3693e43ad36 100644 --- a/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala +++ b/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala @@ -99,7 +99,7 @@ object WorkflowService { if (row == null) { throw new IllegalArgumentException(s"no warehouse with id $whid owned by this user") } - row.getWarehouseName + row.getLakekeeperWarehouseName }) } val cleanUpDeadlineInSeconds: Int = ApplicationConfig.executionStateCleanUpInSecs diff --git a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResourceSpec.scala b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResourceSpec.scala index 1d3aeb9eed2..9d5c55857b7 100644 --- a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResourceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/warehouse/WarehouseResourceSpec.scala @@ -22,6 +22,7 @@ package org.apache.texera.web.resource.dashboard.user.warehouse import org.apache.texera.auth.SessionUser import org.apache.texera.common.config.StorageConfig import org.apache.texera.dao.MockTexeraDB +import org.jooq.impl.DSL import org.apache.texera.dao.jooq.generated.Tables.USER_WAREHOUSE import org.apache.texera.dao.jooq.generated.tables.daos.UserDao import org.apache.texera.dao.jooq.generated.tables.pojos.User @@ -125,19 +126,34 @@ class WarehouseResourceSpec // Create / list / delete // --------------------------------------------------------------------------- - "create" should "create in Lakekeeper, record the row, and mint user--" in { + "create" should "create in Lakekeeper, record the row, and mint user--" in { val created = resource.create(CreateWarehouseRequest("mybucket"), sessionUser) created.name shouldBe "mybucket" - created.warehouseName shouldBe s"user-${sessionUser.getUid}-mybucket" + // The catalog name is derived from the row id, never from the display name, so the + // display name stays free to change later (#7753). + created.warehouseName shouldBe s"user-${sessionUser.getUid}-${created.whid}" + created.warehouseName should not include "mybucket" created.flavor shouldBe "local" - createdNames.toList shouldBe List(s"user-${sessionUser.getUid}-mybucket") + createdNames.toList shouldBe List(s"user-${sessionUser.getUid}-${created.whid}") val status = resource.status(sessionUser) status.enabled shouldBe true status.warehouses.map(_.whid) shouldBe List(created.whid) } + it should "mint a fresh catalog name when the same display name is reused" in { + // Deleting a warehouse and creating another with the same display name must not + // reuse the catalog name: stored result URIs embed it, so a reused name would let + // a new warehouse inherit an old one's storage path (#7753). + val first = resource.create(CreateWarehouseRequest("recycled"), sessionUser) + resource.delete(first.whid, sessionUser) + val second = resource.create(CreateWarehouseRequest("recycled"), sessionUser) + + second.name shouldBe first.name + second.warehouseName should not be first.warehouseName + } + it should "reject an unsafe or duplicate name" in { a[BadRequestException] should be thrownBy resource.create(CreateWarehouseRequest("a/b"), sessionUser) @@ -164,7 +180,17 @@ class WarehouseResourceSpec val squatter = getDSLContext.newRecord(USER_WAREHOUSE) squatter.setUid(otherUser.getUid) squatter.setName("unrelated") - squatter.setWarehouseName(s"user-${sessionUser.getUid}-boom") + // The catalog name now comes from the sequence, so claim the id the next create + // will draw: take one number for the squatter itself (set explicitly, so storing it + // consumes nothing further) and squat on the one after it. + val takenWhid = getDSLContext.fetchValue( + DSL.field( + "nextval(pg_get_serial_sequence('texera_db.user_warehouse','whid'))", + classOf[Integer] + ) + ) + squatter.setWhid(takenWhid) + squatter.setLakekeeperWarehouseName(s"user-${sessionUser.getUid}-${takenWhid + 1}") squatter.setLakekeeperWarehouseId(UUID.randomUUID()) squatter.setFlavor( org.apache.texera.dao.jooq.generated.enums.UserWarehouseFlavorEnum.local @@ -203,7 +229,17 @@ class WarehouseResourceSpec val squatter = getDSLContext.newRecord(USER_WAREHOUSE) squatter.setUid(otherUser.getUid) squatter.setName("unrelated-2") - squatter.setWarehouseName(s"user-${sessionUser.getUid}-doublefault") + // The catalog name now comes from the sequence, so claim the id the next create + // will draw: take one number for the squatter itself (set explicitly, so storing it + // consumes nothing further) and squat on the one after it. + val takenWhid = getDSLContext.fetchValue( + DSL.field( + "nextval(pg_get_serial_sequence('texera_db.user_warehouse','whid'))", + classOf[Integer] + ) + ) + squatter.setWhid(takenWhid) + squatter.setLakekeeperWarehouseName(s"user-${sessionUser.getUid}-${takenWhid + 1}") squatter.setLakekeeperWarehouseId(UUID.randomUUID()) squatter.setFlavor( org.apache.texera.dao.jooq.generated.enums.UserWarehouseFlavorEnum.local diff --git a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala index a0e14bcc8dd..7dc5742a184 100644 --- a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala @@ -1136,7 +1136,7 @@ class WorkflowExecutionsResourceSpec val warehouse = getDSLContext.newRecord(USER_WAREHOUSE) warehouse.setUid(testUser.getUid) warehouse.setName("latest-entry-warehouse") - warehouse.setWarehouseName(s"user-${testUser.getUid}-latest-entry-warehouse") + warehouse.setLakekeeperWarehouseName(s"user-${testUser.getUid}-latest-entry-warehouse") warehouse.setLakekeeperWarehouseId(UUID.randomUUID()) warehouse.setFlavor(UserWarehouseFlavorEnum.local) warehouse.store() diff --git a/amber/src/test/scala/org/apache/texera/web/service/ExecutionsMetadataPersistServiceSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/ExecutionsMetadataPersistServiceSpec.scala index bc49d67e56e..620d8312b3c 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/ExecutionsMetadataPersistServiceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/ExecutionsMetadataPersistServiceSpec.scala @@ -227,7 +227,7 @@ class ExecutionsMetadataPersistServiceSpec val row = getDSLContext.newRecord(USER_WAREHOUSE) row.setUid(testUid) row.setName("exec-spec-warehouse") - row.setWarehouseName(s"user-$testUid-exec-spec-warehouse") + row.setLakekeeperWarehouseName(s"user-$testUid-exec-spec-warehouse") row.setLakekeeperWarehouseId(UUID.randomUUID()) row.setFlavor(UserWarehouseFlavorEnum.local) row.store() @@ -256,7 +256,7 @@ class ExecutionsMetadataPersistServiceSpec val row = getDSLContext.newRecord(USER_WAREHOUSE) row.setUid(testUid) row.setName("doomed-warehouse") - row.setWarehouseName(s"user-$testUid-doomed-warehouse") + row.setLakekeeperWarehouseName(s"user-$testUid-doomed-warehouse") row.setLakekeeperWarehouseId(UUID.randomUUID()) row.setFlavor(UserWarehouseFlavorEnum.local) row.store() diff --git a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala index ef7d8cfb7cc..3c9d76167a8 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala @@ -63,7 +63,7 @@ class WorkflowServiceWarehouseSpec val row = getDSLContext.newRecord(USER_WAREHOUSE) row.setUid(ownerUid) row.setName("mybucket") - row.setWarehouseName(s"user-$ownerUid-mybucket") + row.setLakekeeperWarehouseName(s"user-$ownerUid-mybucket") row.setLakekeeperWarehouseId(UUID.randomUUID()) row.setFlavor(UserWarehouseFlavorEnum.local) row.store() diff --git a/common/dao/src/test/scala/org/apache/texera/dao/UserWarehouseSpec.scala b/common/dao/src/test/scala/org/apache/texera/dao/UserWarehouseSpec.scala index 30c09b152ff..6179e9efd6c 100644 --- a/common/dao/src/test/scala/org/apache/texera/dao/UserWarehouseSpec.scala +++ b/common/dao/src/test/scala/org/apache/texera/dao/UserWarehouseSpec.scala @@ -60,7 +60,7 @@ class UserWarehouseSpec extends AnyFlatSpec with Matchers with BeforeAndAfterAll USER_WAREHOUSE, USER_WAREHOUSE.UID, USER_WAREHOUSE.NAME, - USER_WAREHOUSE.WAREHOUSE_NAME, + USER_WAREHOUSE.LAKEKEEPER_WAREHOUSE_NAME, USER_WAREHOUSE.LAKEKEEPER_WAREHOUSE_ID, USER_WAREHOUSE.FLAVOR ) @@ -78,7 +78,7 @@ class UserWarehouseSpec extends AnyFlatSpec with Matchers with BeforeAndAfterAll .where(USER_WAREHOUSE.UID.eq(uid)) .fetchOne() row.getName shouldBe "mybucket" - row.getWarehouseName shouldBe s"user-$uid-mybucket" + row.getLakekeeperWarehouseName shouldBe s"user-$uid-mybucket" row.getFlavor shouldBe UserWarehouseFlavorEnum.local row.getLakekeeperWarehouseId should not be null row.getCreatedAt should not be null diff --git a/sql/changelog.xml b/sql/changelog.xml index d3a7d55228c..c0f735a09c0 100644 --- a/sql/changelog.xml +++ b/sql/changelog.xml @@ -99,6 +99,11 @@ + + + + +