From d31323ae783f474f80b1f87d90cf9270708938a5 Mon Sep 17 00:00:00 2001 From: Andrew Zolotukhin Date: Fri, 2 Oct 2026 22:40:09 +0000 Subject: [PATCH] fix(orm): honor polymorphic write lifecycle (F07) --- .changeset/polymorphic-write-lifecycle.md | 19 + docs/framework-feature-candidates.md | 13 +- libs/knex-schema/MIGRATION-v5.md | 20 + .../integration/polymorphic-writes.test.ts | 618 +++++++++++++++++ .../src/PolymorphicQueryBuilder.ts | 111 ++- libs/knex-schema/src/SchemaQueryBuilder.ts | 13 +- libs/orm/README.md | 56 +- libs/orm/src/change-tracker.ts | 151 +++-- libs/orm/src/dbset.ts | 117 ++-- libs/orm/src/orm.test.ts | 66 +- libs/orm/src/result-types.ts | 19 +- libs/orm/src/variant-write.test-d.ts | 62 ++ libs/orm/src/variant-write.test.ts | 86 +++ libs/orm/src/variant-write.ts | 640 +++++++++--------- 14 files changed, 1461 insertions(+), 530 deletions(-) create mode 100644 .changeset/polymorphic-write-lifecycle.md create mode 100644 libs/knex-schema/integration/polymorphic-writes.test.ts create mode 100644 libs/orm/src/variant-write.test-d.ts create mode 100644 libs/orm/src/variant-write.test.ts diff --git a/.changeset/polymorphic-write-lifecycle.md b/.changeset/polymorphic-write-lifecycle.md new file mode 100644 index 00000000..f9765a37 --- /dev/null +++ b/.changeset/polymorphic-write-lifecycle.md @@ -0,0 +1,19 @@ +--- +"@cleverbrush/orm": major +"@cleverbrush/knex-schema": patch +--- + +Honor hooks, timestamps and base-schema soft deletion in explicit variant writes +and tracked polymorphic saves, for single-table and class-table inheritance. +Add `ofVariant(key).restore()` and `hardDelete()`; CTI soft deletion retains the +child row, while permanent deletion removes both rows atomically. + +**Migration:** Variant `delete()` now respects the base schema's `.softDelete()`. +Use `hardDelete()` for physical removal, and `withDeleted()` to include hidden +rows. Review lifecycle hooks that now run, and remove identity/discriminator/join +keys from update patches. This change is part of the coordinated v5 major release. + +Capture mutation targets inside the write transaction, preserve query restrictions +and transaction bindings, and use savepoints for caller-owned transactions. Fix +base/variant column mapping and visibility of extension-managed deletion columns; +retain exact numeric keys and tracked optimistic-concurrency/rollback semantics. diff --git a/docs/framework-feature-candidates.md b/docs/framework-feature-candidates.md index 5e1282c7..961c5f2d 100644 --- a/docs/framework-feature-candidates.md +++ b/docs/framework-feature-candidates.md @@ -1,6 +1,6 @@ # Framework feature candidates -Status: F01–F05 merged; F06 implemented on the CORS feature branch for PR review; F07 remains proposed. +Status: F01–F06 merged; F07 implemented on the polymorphic-write lifecycle branch for PR review. Assessment date: 2026-10-02. This is an unprioritized list of reusable Framework capabilities and correctness @@ -242,7 +242,7 @@ documentation. See the [CORS guide](../libs/server/README.md#cors). **Priority:** -**Existing support and gap.** Ordinary entities support soft deletion and lifecycle +**Original assessment.** Ordinary entities support soft deletion and lifecycle hooks. Source review of [variant writes](../libs/orm/src/variant-write.ts) found direct Knex updates/deletes: variant deletion physically removes rows and these paths bypass the ordinary write pipeline. This finding needs PostgreSQL @@ -268,6 +268,15 @@ this candidate is pending. - Failures roll back base/variant writes together and respect query predicates. - Type tests describe the supported mutation surface accurately. +**Implementation:** Variant and tracked polymorphic writes share an atomic +lifecycle pipeline for STI and CTI, including hooks, timestamps, soft deletion, +restoration and explicit permanent deletion. Base soft-delete metadata controls +visibility; CTI child rows remain intact until permanent deletion. PostgreSQL +regressions cover scopes, transaction/savepoint rollback, concurrency and exact +keys. The deletion change ships in the coordinated v5 major release; see the +[ORM lifecycle contract](../libs/orm/README.md#variant-deletion-and-lifecycle) and +[v5 migration guide](../libs/knex-schema/MIGRATION-v5.md#7-review-polymorphic-mutation-lifecycle). + **Review notes:** ## Dependencies and review boundaries diff --git a/libs/knex-schema/MIGRATION-v5.md b/libs/knex-schema/MIGRATION-v5.md index 1f188b49..ada6521e 100644 --- a/libs/knex-schema/MIGRATION-v5.md +++ b/libs/knex-schema/MIGRATION-v5.md @@ -198,6 +198,25 @@ with typed child customizers. Use application mappers after decoding for DTO cha Child default scopes and soft deletion apply automatically; disable them explicitly on the child when that is the intended policy. +## 7. Review polymorphic mutation lifecycle + +Variant writes and tracked polymorphic `saveChanges()` now honor lifecycle hooks, +timestamps and base-schema soft deletion for STI and CTI. Previously, explicit +variant updates/deletes and tracked mutations bypassed that pipeline; variant +`delete()` physically removed rows even when soft deletion was configured. + +If physical removal is required, replace `.ofVariant(key).delete()` with +`.ofVariant(key).hardDelete()`. Use `.withDeleted()` to include previously deleted +rows. Restore with `.ofVariant(key).onlyDeleted().restore()`. CTI soft deletion +retains the child row; the base deletion marker controls entity visibility. + +Variant updates now accept base and branch fields together, with correct column +mapping. Remove primary-key, discriminator and CTI join-key changes from update +payloads. Review hooks that may now execute: base hooks precede variant hooks, +insert/update hooks receive the combined payload, and after-insert hooks receive +the completed row. Do not assume hook side effects are undone on transaction rollback. +See the [ORM lifecycle contract](../orm/README.md#variant-deletion-and-lifecycle). + ## Upgrade checklist - Upgrade the fixed Framework package group together; remove `.withRowSchema()` calls. @@ -205,6 +224,7 @@ on the child when that is the intended policy. - Replace raw base queries and shape-changing SQL with explicit output contracts. - Audit API DTOs for exact numbers, dates, SQL nulls and storage/input separation. - Migrate relation customizers and polymorphic branch projections. +- Review polymorphic deletion intent and hooks; use `hardDelete()` for permanent removal. - Verify identity tracking, writes, transaction rollback and concurrency in the app. - Run TypeScript, unit and real PostgreSQL tests; test representative endpoint flows. diff --git a/libs/knex-schema/integration/polymorphic-writes.test.ts b/libs/knex-schema/integration/polymorphic-writes.test.ts new file mode 100644 index 00000000..baf5478f --- /dev/null +++ b/libs/knex-schema/integration/polymorphic-writes.test.ts @@ -0,0 +1,618 @@ +import { randomUUID } from 'node:crypto'; +import { + ConcurrencyError, + createDb, + date, + defineEntity, + generateCreateTable, + number, + object, + string +} from '@cleverbrush/orm'; +import Knex from 'knex'; +import { types as pgTypes } from 'pg'; +import { afterAll, describe, expect, it } from 'vitest'; + +const connection = process.env.QUERY_TEST_DATABASE_URL; +if (!connection) throw new Error('QUERY_TEST_DATABASE_URL is required'); +const knex = Knex({ client: 'pg', connection, pool: { min: 0, max: 8 } }); +const tables: string[] = []; +afterAll(async () => { + for (const table of tables.reverse()) + await knex.schema.dropTableIfExists(table); + await knex.destroy(); +}); + +type Controls = { + events: string[]; + scopeCalls: number; + beforeInsert?: (data: any) => any; + beforeUpdate?: (data: any) => any; + beforeDelete?: (query: any) => Promise; + afterInsert?: (row: any) => void; +}; + +async function fixture( + storage: 'sti' | 'cti', + soft = true, + scopedBody = false +) { + const name = `f07_${randomUUID().replaceAll('-', '')}`; + const childName = `${name}_child`; + const controls: Controls = { events: [], scopeCalls: 0 }; + let base = object({ + id: number() + .hasColumnName('asset_id') + .primaryKey({ autoIncrement: true }), + kind: string().hasColumnName('asset_kind'), + title: string().hasColumnName('title_text'), + tenant: number().defaultTo(1), + version: number().defaultTo(1).rowVersion(), + createdAt: date().hasColumnName('created_on'), + updatedAt: date().hasColumnName('changed_on') + }) + .hasTableName(name) + .hasTimestamps({ createdAt: 'created_on', updatedAt: 'changed_on' }) + .defaultScope(q => { + controls.scopeCalls++; + return q.where(t => t.tenant, 1); + }) + .beforeInsert((data: any) => { + controls.events.push('base:insert'); + return controls.beforeInsert?.(data) ?? data; + }) + .afterInsert((row: any) => { + controls.events.push('base:inserted'); + controls.afterInsert?.(row); + }) + .beforeUpdate((data: any) => { + controls.events.push('base:update'); + return controls.beforeUpdate?.(data) ?? data; + }) + .beforeDelete(async (query: any) => { + controls.events.push('base:delete'); + await controls.beforeDelete?.(query); + }); + if (soft) base = base.softDelete({ column: 'removed_on' }); + let body = object({ + ...(storage === 'cti' + ? { assetId: number().hasColumnName('owner_id').primaryKey() } + : {}), + caption: string().hasColumnName('caption_text'), + document: object({ label: string() }) + .acceptUnknownProps() + .jsonb() + .hasColumnName('document_data'), + bodyCreatedAt: date().hasColumnName('body_created'), + bodyUpdatedAt: date().hasColumnName('body_changed') + }) + .hasTableName(storage === 'cti' ? childName : name) + .hasTimestamps({ createdAt: 'body_created', updatedAt: 'body_changed' }) + .beforeInsert((data: any) => { + controls.events.push('body:insert'); + return data; + }) + .afterInsert(() => { + controls.events.push('body:inserted'); + }) + .beforeUpdate((data: any) => { + controls.events.push('body:update'); + return data; + }) + .beforeDelete(() => { + controls.events.push('body:delete'); + }); + if (scopedBody) + body = body.defaultScope(q => { + controls.scopeCalls++; + return q.where(t => t.caption, 'initial'); + }); + await generateCreateTable(base)(knex); + tables.push(name); + if (storage === 'cti') { + await generateCreateTable(body)(knex); + tables.push(childName); + await knex.schema.alterTable(childName, table => + table.foreign('owner_id').references('asset_id').inTable(name) + ); + } else { + await knex.schema.alterTable(name, table => { + table.text('caption_text').notNullable(); + table.jsonb('document_data').notNullable(); + table.timestamp('body_created').notNullable(); + table.timestamp('body_changed').notNullable(); + }); + } + await knex.schema.alterTable( + storage === 'cti' ? childName : name, + table => { + table.check("caption_text <> 'reject'", [], 'valid_caption'); + } + ); + const entity = + storage === 'cti' + ? defineEntity(base) + .discriminator('kind') + .ctiVariant('photo', defineEntity(body), t => t.assetId!) + : defineEntity(base) + .discriminator('kind') + .stiVariant('photo', body); + const db = createDb(knex, { assets: entity }); + const view = db.assets.ofVariant('photo'); + const insert = (extra: Record = {}) => + view.insert({ + title: 'initial', + caption: 'initial', + document: { label: 'original', extra: [true] }, + ...extra + }); + return { name, childName, entity, view, db, insert, controls }; +} + +for (const storage of ['sti', 'cti'] as const) { + describe(storage, () => { + it('runs logical hooks once and maps base/body updates, JSON and timestamps', async () => { + const f = await fixture(storage); + f.controls.beforeInsert = data => { + expect(data.kind).toBe('photo'); + return { + ...data, + title: 'insert hook', + caption: 'body insert hook' + }; + }; + f.controls.afterInsert = row => + expect(row.caption).toBe('body insert hook'); + const row = await f.insert(); + expect(f.controls.events).toEqual([ + 'base:insert', + 'body:insert', + 'base:inserted', + 'body:inserted' + ]); + expect(row).toMatchObject({ + kind: 'photo', + title: 'insert hook', + caption: 'body insert hook', + document: { extra: [true] } + }); + expect(row.createdAt).toBeInstanceOf(Date); + expect(row.bodyCreatedAt).toBeInstanceOf(Date); + await knex(f.name) + .where('asset_id', row.id) + .update({ changed_on: '2000-01-01' }); + await knex(storage === 'cti' ? f.childName : f.name).update({ + body_changed: '2000-01-01' + }); + f.controls.events.length = 0; + f.controls.beforeUpdate = data => ({ + ...data, + title: 'updated by hook' + }); + await f.view + .where(t => t.id, row.id) + .update({ + caption: 'updated', + document: { label: 'next', custom: { a: 1 } } + }); + expect(f.controls.events).toEqual(['base:update', 'body:update']); + const updated = await f.view.findOrFail(row.id); + expect(updated).toMatchObject({ + title: 'updated by hook', + caption: 'updated', + document: { custom: { a: 1 } } + }); + expect(updated.updatedAt.getTime()).toBeGreaterThan( + new Date('2000-01-01').getTime() + ); + expect(updated.bodyUpdatedAt.getTime()).toBeGreaterThan( + new Date('2000-01-01').getTime() + ); + expect(updated.createdAt).toEqual(row.createdAt); + f.controls.beforeUpdate = undefined; + await f.view + .where(t => t.id, row.id) + .update({ title: 'base only' }); + expect((await f.view.findOrFail(row.id)).title).toBe('base only'); + }); + + it('soft-deletes and restores the base marker without losing body rows', async () => { + const f = await fixture(storage); + const row = await f.insert(); + const before = await f.view.findOrFail(row.id); + f.controls.events.length = 0; + f.controls.beforeDelete = async query => { + expect((await query).map((r: any) => r.id)).toEqual([row.id]); + }; + await f.view.where(t => t.id, row.id).delete(); + expect(f.controls.events).toEqual(['base:delete', 'body:delete']); + expect(await f.view.find(row.id)).toBeUndefined(); + expect( + (await f.view.onlyDeleted().execute()).map(r => r.id) + ).toEqual([row.id]); + expect( + await knex(f.name).where('asset_id', row.id).first() + ).toMatchObject({ removed_on: expect.any(Date) }); + if (storage === 'cti') + expect( + await knex(f.childName).where('owner_id', row.id).first() + ).toBeDefined(); + f.controls.events.length = 0; + expect(await f.view.where(t => t.id, row.id).restore()).toEqual([]); + const restored = await f.view + .onlyDeleted() + .where(t => t.id, row.id) + .restore(); + expect(restored).toEqual([before]); + expect(f.controls.events).toEqual([]); + expect(await f.view.where(t => t.id, row.id).hardDelete()).toBe(1); + expect(f.controls.events).toEqual(['base:delete', 'body:delete']); + expect(await f.view.withDeleted().find(row.id)).toBeUndefined(); + if (storage === 'cti') + expect( + await knex(f.childName).where('owner_id', row.id).first() + ).toBeUndefined(); + }); + + it('honors scopes, explicit predicates, variant filters and pagination', async () => { + const f = await fixture(storage); + const a = await f.insert({ title: 'a' }); + const b = await f.insert({ title: 'b' }); + const c = await f.insert({ title: 'c', tenant: 2 }); + await f.view + .orderBy(t => t.id) + .offset(1) + .limit(1) + .update({ caption: 'page' }); + expect((await f.view.findOrFail(a.id)).caption).toBe('initial'); + expect((await f.view.findOrFail(b.id)).caption).toBe('page'); + expect((await f.view.unscoped().findOrFail(c.id)).caption).toBe( + 'initial' + ); + await knex(f.name) + .where('asset_id', a.id) + .update({ asset_kind: 'another' }); + await f.view.where(t => t.id, a.id).delete(); + expect( + (await knex(f.name).where('asset_id', a.id).first()).removed_on + ).toBeNull(); + await f.view + .unscoped() + .where(t => t.id, c.id) + .delete(); + expect(await f.view.unscoped().find(c.id)).toBeUndefined(); + }); + + it('rolls back insert/update failures and leaves a caller transaction usable', async () => { + const f = await fixture(storage); + const row = await f.insert(); + await expect(f.insert({ caption: 'reject' })).rejects.toThrow(); + expect(await f.view.countValue()).toBe(1); + await knex.transaction(async trx => { + let previousTitle = 'initial'; + for (const bind of [ + 'withTransaction', + 'transacting' + ] as const) { + const scoped = f.view.where(t => t.id, row.id)[bind](trx); + await expect( + scoped.update({ + title: 'must rollback', + caption: 'reject' + }) + ).rejects.toThrow(); + expect((await scoped.findOrFail(row.id)).title).toBe( + previousTitle + ); + await scoped.update({ title: bind }); + previousTitle = bind; + } + }); + expect((await f.view.findOrFail(row.id)).title).toBe('transacting'); + f.controls.beforeDelete = async () => { + throw new Error('hook rejected'); + }; + await expect( + f.view.where(t => t.id, row.id).delete() + ).rejects.toThrow('hook rejected'); + expect(await f.view.find(row.id)).toBeDefined(); + f.controls.afterInsert = () => { + throw new Error('after rejected'); + }; + await expect(f.insert()).rejects.toThrow('after rejected'); + expect(await f.view.countValue()).toBe(1); + }); + + it('physically deletes without soft-delete metadata and rejects restore', async () => { + const f = await fixture(storage, false); + const row = await f.insert(); + await expect(f.view.restore()).rejects.toThrow(/soft delete/); + await f.view.where(t => t.id, row.id).delete(); + expect( + await knex(f.name).where('asset_id', row.id).first() + ).toBeUndefined(); + if (storage === 'cti') + expect( + await knex(f.childName).where('owner_id', row.id).first() + ).toBeUndefined(); + }); + + it('tracks hook results, versions, soft deletion and rollback without advancing snapshots', async () => { + const f = await fixture(storage); + const row = await f.insert(); + const db = createDb(knex, { assets: f.entity }, { tracking: true }); + const current = await db.assets + .ofVariant('photo') + .findOrFail(row.id); + current.caption = 'tracked'; + f.controls.beforeUpdate = data => ({ + ...data, + title: 'tracked hook' + }); + await db.saveChanges(); + expect(current).toMatchObject({ + title: 'tracked hook', + caption: 'tracked', + version: 2 + }); + expect(db.entry(current).state).toBe('Unchanged'); + current.caption = 'reject'; + await expect(db.saveChanges()).rejects.toThrow(); + expect(current.version).toBe(2); + expect(db.entry(current).state).toBe('Modified'); + expect(db.entry(current).originalValues.caption).toBe('tracked'); + current.caption = 'next'; + await knex(f.name) + .where('asset_id', row.id) + .update({ version: 10 }); + await expect(db.saveChanges()).rejects.toBeInstanceOf( + ConcurrencyError + ); + expect(current.version).toBe(2); + await knex(f.name).where('asset_id', row.id).update({ version: 2 }); + db.remove(current); + await db.saveChanges(); + expect(await f.view.find(row.id)).toBeUndefined(); + if (storage === 'cti') + expect( + await knex(f.childName).where('owner_id', row.id).first() + ).toBeDefined(); + expect( + await f.view + .onlyDeleted() + .where(t => t.id, row.id) + .hardDelete() + ).toBe(1); + }); + + it('does not run hooks on empty targets and rejects hook-produced identity changes', async () => { + const f = await fixture(storage); + const row = await f.insert(); + f.controls.events.length = 0; + await f.view.where(t => t.id, -1).update({ caption: 'absent' }); + await f.view.where(t => t.id, -1).delete(); + expect(await f.view.where(t => t.id, -1).hardDelete()).toBe(0); + expect(f.controls.events).toEqual([]); + f.controls.beforeUpdate = data => ({ ...data, id: row.id + 1 }); + await expect( + f.view.where(t => t.id, row.id).update({ title: 'invalid' }) + ).rejects.toThrow(/identity/); + expect((await f.view.findOrFail(row.id)).title).toBe('initial'); + }); + + it('passes exactly the paginated targets to delete hooks', async () => { + const f = await fixture(storage); + const first = await f.insert(); + const second = await f.insert(); + f.controls.beforeDelete = async query => { + expect((await query).map((r: any) => r.id)).toEqual([ + second.id + ]); + }; + await f.view + .orderBy(t => t.id) + .offset(1) + .limit(1) + .delete(); + expect(await f.view.find(first.id)).toBeDefined(); + expect(await f.view.find(second.id)).toBeUndefined(); + }); + + it('does not commit a caller transaction and preserves rows after permanent-delete failure', async () => { + const f = await fixture(storage); + const row = await f.insert(); + const reference = `${f.name}_ref`; + await knex.schema.createTable(reference, table => { + table + .integer('asset_id') + .references('asset_id') + .inTable(f.name); + }); + tables.push(reference); + await knex(reference).insert({ asset_id: row.id }); + await knex.transaction(async trx => { + await expect( + f.view + .withTransaction(trx) + .where(t => t.id, row.id) + .hardDelete() + ).rejects.toThrow(); + expect( + await f.view.withTransaction(trx).find(row.id) + ).toBeDefined(); + if (storage === 'cti') + expect( + await trx(f.childName).where('owner_id', row.id).first() + ).toBeDefined(); + }); + await expect( + knex.transaction(async trx => { + await f.view + .where(t => t.id, row.id) + .transacting(trx) + .delete(); + throw new Error('outer rollback'); + }) + ).rejects.toThrow('outer rollback'); + expect(await f.view.find(row.id)).toBeDefined(); + }); + + it('inserts tracked variants through the same hooks and atomic pipeline', async () => { + const f = await fixture(storage); + const db = createDb(knex, { assets: f.entity }, { tracking: true }); + const row = { + kind: 'photo', + title: 'tracked new', + caption: 'reject', + document: { label: 'new' } + }; + db.attach('assets', row); + await expect(db.saveChanges()).rejects.toThrow(); + expect((row as any).id).toBeUndefined(); + expect(db.entry(row).state).toBe('Added'); + row.caption = 'accepted'; + f.controls.events.length = 0; + await db.saveChanges(); + expect((row as any).id).toEqual(expect.any(Number)); + expect(db.entry(row).state).toBe('Unchanged'); + expect(f.controls.events).toEqual([ + 'base:insert', + 'body:insert', + 'base:inserted', + 'body:inserted' + ]); + expect(await f.view.find((row as any).id)).toMatchObject({ + caption: 'accepted' + }); + }); + + it('rechecks predicates after waiting for a concurrent writer', async () => { + const f = await fixture(storage); + const row = await f.insert(); + f.controls.events.length = 0; + const blocker = await knex.transaction(); + await blocker(f.name) + .where('asset_id', row.id) + .forUpdate() + .select('asset_id'); + const pending = f.view + .where(t => t.id, row.id) + .where(t => t.title, 'initial') + .update({ caption: 'must not write' }) + .then( + () => undefined, + error => error + ); + try { + await expect + .poll( + async () => { + const waiting = await knex('pg_stat_activity') + .where('wait_event_type', 'Lock') + .where('query', 'like', `%${f.name}%`); + return waiting.length; + }, + { timeout: 5000, interval: 20 } + ) + .toBeGreaterThan(0); + await blocker(f.name) + .where('asset_id', row.id) + .update({ title_text: 'concurrent' }); + await blocker.commit(); + expect(await pending).toBeUndefined(); + expect((await f.view.findOrFail(row.id)).caption).toBe( + 'initial' + ); + expect(f.controls.events).toEqual([]); + } finally { + if (!blocker.isCompleted()) await blocker.rollback(); + await pending; + } + }); + + it('keeps native bigint mutation keys exact even with a lossy driver parser', async () => { + const name = `f07_${randomUUID().replaceAll('-', '')}`; + const childName = `${name}_child`; + const base = object({ + id: number().bigint().primaryKey({ autoIncrement: false }), + kind: string() + }) + .hasTableName(name) + .softDelete(); + const body = object({ + ...(storage === 'cti' + ? { + assetId: number() + .bigint() + .hasColumnName('asset_id') + .primaryKey({ autoIncrement: false }) + } + : {}), + caption: string() + }).hasTableName(storage === 'cti' ? childName : name); + await generateCreateTable(base)(knex); + tables.push(name); + if (storage === 'cti') { + await generateCreateTable(body)(knex); + tables.push(childName); + } else await knex.schema.alterTable(name, t => t.text('caption')); + const entity = + storage === 'cti' + ? defineEntity(base) + .discriminator('kind') + .ctiVariant( + 'photo', + defineEntity(body), + t => t.assetId! + ) + : defineEntity(base) + .discriminator('kind') + .stiVariant('photo', body); + const view = createDb(knex, { assets: entity }).assets.ofVariant( + 'photo' + ); + const original = pgTypes.getTypeParser(20, 'text'); + pgTypes.setTypeParser(20, Number); + try { + const id = '9007199254740993'; + await view.insert({ id, caption: 'initial' }); + await view.insert({ + id: '9007199254740992', + caption: 'neighbor' + }); + await view.where(t => t.id, id).update({ caption: 'exact' }); + expect((await view.findOrFail(id)).caption).toBe('exact'); + expect( + (await view.findOrFail('9007199254740992')).caption + ).toBe('neighbor'); + await view.where(t => t.id, id).delete(); + expect(await view.find(id)).toBeUndefined(); + expect( + ( + await view + .onlyDeleted() + .where(t => t.id, id) + .restore() + )[0].id + ).toBe(id); + expect(await view.where(t => t.id, id).hardDelete()).toBe(1); + expect(await view.countValue()).toBe(1); + } finally { + pgTypes.setTypeParser(20, original); + } + }); + + it('captures scopes once and returns writes even after moving out of a child scope', async () => { + const f = await fixture(storage, true, true); + const row = await f.insert(); + const selected = f.view.where(t => t.id, row.id); + const calls = f.controls.scopeCalls; + await selected.update({ caption: 'outside scope' }); + expect(f.controls.scopeCalls).toBe(calls); + const stored = await knex( + storage === 'cti' ? f.childName : f.name + ).first('caption_text'); + expect(stored.caption_text).toBe('outside scope'); + }); + }); +} diff --git a/libs/knex-schema/src/PolymorphicQueryBuilder.ts b/libs/knex-schema/src/PolymorphicQueryBuilder.ts index ad239e97..c81eacf2 100644 --- a/libs/knex-schema/src/PolymorphicQueryBuilder.ts +++ b/libs/knex-schema/src/PolymorphicQueryBuilder.ts @@ -172,6 +172,7 @@ export class PolymorphicQueryBuilder< limit?: number; offset?: number; }; + private readonly deletionColumn: string; private skipDefaults = false; private deleted: 'exclude' | 'include' | 'only' = 'exclude'; private readonly predicateAlias = @@ -202,6 +203,15 @@ export class PolymorphicQueryBuilder< const config = getVariants(source); if (!config) throw new ReadSchemaError('No polymorphic variants are declared'); + this.deletionColumn = privateColumn( + [ + ...Object.keys(source.introspect().properties), + ...Object.values(config.variants).flatMap(v => + Object.keys(v.schema.introspect().properties) + ) + ], + 'read_deleted' + ); this.branches = Object.create(null); for (const key of Object.keys(config.variants)) this.branches[key] = this.branch(key, true); @@ -215,16 +225,17 @@ export class PolymorphicQueryBuilder< common, knex .from(base.clone().as('__read_unknown')) - .select( - Object.fromEntries( + .select({ + ...Object.fromEntries( Object.keys(common.introspect().properties).map(key => [ key, knex.ref( `__read_unknown.${buildColumnMap(source).propToCol.get(key) ?? key}` ) ]) - ) - ) + ), + ...this.deletionSelection('__read_unknown') + }) .where(q => q .whereNotIn(discriminator, Object.keys(config.variants)) @@ -264,6 +275,19 @@ export class PolymorphicQueryBuilder< } } + private deletionSelection(alias: string): Record { + const softDelete = this.source.introspect().extensions?.softDelete as + | { column: string } + | undefined; + return softDelete + ? { + [this.deletionColumn]: this.knex.raw('??', [ + `${alias}.${softDelete.column}` + ]) + } + : {}; + } + private commonSource(): ReadObject { const info = this.source.introspect(); const relationNames = new Set( @@ -350,7 +374,9 @@ export class PolymorphicQueryBuilder< `${baseAlias}.${basePk[0]}` ); } - const columns: Record = Object.create(null); + const columns: Record = { + ...this.deletionSelection(baseAlias) + }; const properties: Record = Object.create(null); for (const [name, schema] of Object.entries(baseProperties)) { if (excluded.has(name)) continue; @@ -712,6 +738,13 @@ export class PolymorphicQueryBuilder< /** @internal Compile one UNION ALL statement; JSON preserves distinct branch shapes. */ compile(correlate?: ReadCorrelation): Knex.QueryBuilder { + return this.compileRows(correlate); + } + + private compileRows( + correlate?: ReadCorrelation, + targetKey?: string + ): Knex.QueryBuilder { const defaults = this.skipDefaults ? undefined : this.defaults; const reserved = Object.values(this.branches).flatMap(branch => Object.keys(branch.rowSchema.introspect().properties) @@ -744,35 +777,35 @@ export class PolymorphicQueryBuilder< const softDelete = this.source.introspect().extensions ?.softDelete as { column: string } | undefined; if (softDelete && this.deleted !== 'include') { - const key = - buildColumnMap(this.source).colToProp.get( - softDelete.column - ) ?? softDelete.column; branch = branch.withPredicate(query => { query[ this.deleted === 'only' ? 'whereNotNull' : 'whereNull' - ](key); + ](this.deletionColumn); }); } + let nativeKey: Knex.Raw | undefined; + const orderColumns: Record = {}; + const compiled = branch.compile((sql, alias, source) => { + correlate?.(sql, alias, source); + if (targetKey) + nativeKey = this.knex.raw('??', [`${alias}.${targetKey}`]); + const columns = buildColumnMap(source).propToCol; + for (const item of order) { + if ('raw' in item) continue; + const { key, hidden } = item; + orderColumns[hidden] = this.knex.raw('cast(?? as text)', [ + `${alias}.${columns.get(key) ?? key}` + ]); + } + if (Object.keys(orderColumns).length) sql.select(orderColumns); + }); + if (targetKey) { + return compiled + .clearSelect() + .select({ __write_pk: nativeKey!, ...orderColumns }); + } return this.knex - .from( - branch - .compile((sql, alias, source) => { - correlate?.(sql, alias, source); - const columns = buildColumnMap(source).propToCol; - for (const item of order) { - if ('raw' in item) continue; - const { key, hidden } = item; - sql.select({ - [hidden]: this.knex.raw( - 'cast(?? as text)', - [`${alias}.${columns.get(key) ?? key}`] - ) - }); - } - }) - .as('__read_branch') - ) + .from(compiled.as('__read_branch')) .select( this.knex.raw('to_jsonb(__read_branch) as __read_poly') ); @@ -784,7 +817,7 @@ export class PolymorphicQueryBuilder< .unionAll(queries, true) .as('__read_variants') ) - .select('__read_poly'); + .select(targetKey ? '__write_pk' : '__read_poly'); for (const item of order) { if ('raw' in item) { query.orderByRaw(item.raw()); @@ -807,7 +840,9 @@ export class PolymorphicQueryBuilder< 'Polymorphic ordering requires a scalar column' ); query.orderByRaw( - `cast(__read_poly ->> ? as ${type}) ${direction}`, + targetKey + ? `cast(?? as ${type}) ${direction}` + : `cast(__read_poly ->> ? as ${type}) ${direction}`, [hidden] ); } @@ -817,6 +852,22 @@ export class PolymorphicQueryBuilder< if (offset !== undefined) query.offset(offset); return query; } + /** @internal Capture native primary keys for a single writable ORM variant. */ + mutationTargets(variantKey: string): Knex.QueryBuilder { + const keys = Object.keys(this.branches); + if (keys.length !== 1 || keys[0] !== variantKey || this.includeUnknown) + throw new ReadSchemaError( + 'Variant writes require exactly their original variant' + ); + this.branches[variantKey].assertWritable(); + const pk = getPrimaryKeyColumns(this.source); + if (pk.propertyKeys.length !== 1) + throw new ReadSchemaError( + 'Variant writes require a single-column primary key' + ); + return this.compileRows(undefined, pk.propertyKeys[0]); + } + /** @internal Decode using exactly the selected branch's schema and codecs. */ decode(row: any, path = 'row'): InferType> { const value = row.__read_poly; diff --git a/libs/knex-schema/src/SchemaQueryBuilder.ts b/libs/knex-schema/src/SchemaQueryBuilder.ts index f29fa4a3..ea979d90 100644 --- a/libs/knex-schema/src/SchemaQueryBuilder.ts +++ b/libs/knex-schema/src/SchemaQueryBuilder.ts @@ -430,6 +430,14 @@ export class SchemaQueryBuilder< return !this.selected && !this.grouped && !this.distinctRows; } + /** @internal Reject writes through projections or loaded relations. */ + assertWritable(): void { + if (!this.returnsEntityRows || this.loaded.length) + throw new ReadSchemaError( + 'Writes require an unprojected table query without relations or aggregation' + ); + } + /** @internal Build an independent statement containing defaults and explicit filters. */ private filtered(): Knex.QueryBuilder { const query = this.base.clone(); @@ -1341,10 +1349,7 @@ export class SchemaQueryBuilder< } private writer(filtered = true): QuerySource> { - if (!this.returnsEntityRows || this.loaded.length) - throw new ReadSchemaError( - 'Writes require an unprojected table query without relations or aggregation' - ); + this.assertWritable(); let base = this.knex(getTableName(this.source)); if (filtered) { const keys = getPrimaryKeyColumns(this.source).columnNames; diff --git a/libs/orm/README.md b/libs/orm/README.md index 52a65622..6a9f6eb0 100644 --- a/libs/orm/README.md +++ b/libs/orm/README.md @@ -320,18 +320,60 @@ await db.activities.ofVariant('assigned').where(t => t.id, 3).delete(); | Method | Description | |--------|-------------| | `.insert(payload)` | Insert a new variant row; discriminator is set automatically | -| `.update(patch)` | Update variant columns for rows matched by the current `WHERE` clause | -| `.delete()` | Delete rows matched by the current `WHERE` clause (CTI: atomic) | +| `.update(patch)` | Atomically update matching base and variant fields; returns `Promise` | +| `.delete()` | Delete matching entities, honoring base soft deletion; returns `Promise` | +| `.restore()` | Clear the base deletion marker and return restored variant rows | +| `.hardDelete()` | Permanently delete matching entities and return their count | | `.find(pk)` | Find a single variant row by PK; `undefined` if not found | | `.findOrFail(pk)` | Like `.find`, but throws `EntityNotFoundError` | | `.findMany([pk…])` | Fetch multiple variant rows by PK in one query | | `.where(col, value)` | Adds a `WHERE` predicate (chainable; returns `VariantDbSet`) | -| `.include(t => t.rel)` | Eager-loads a relation (chainable; returns `VariantDbSet`) | -| `.withTransaction(trx)` | Returns a new `VariantDbSet` bound to an existing transaction | +| `.include(t => t.rel)` | Eager-load a relation; the resulting query is read-only | +| `.withTransaction(trx)` / `.transacting(trx)` | Bind an independent variant view to an existing transaction, preserving filters | -Calling `.insert()` / `.update()` / `.delete()` directly on the polymorphic -base `DbSet` (without `ofVariant`) throws a runtime error — use `ofVariant` -for all writes on polymorphic entities. +Polymorphic root queries are read-only. Use `ofVariant()` for explicit mutations; +tracked polymorphic inserts, updates and removals use the same lifecycle pipeline +when `saveChanges()` runs. Writes require a single-column base primary key and an +unprojected variant view without loaded relations. Primary keys, discriminators +and CTI join keys cannot be changed by an update, including from a hook. + +### Variant deletion and lifecycle + +The **base schema** controls entity soft deletion. With `.softDelete()`, `delete()` +sets its deletion marker and keeps CTI child rows intact. Without that metadata, +`delete()` removes the physical rows. `hardDelete()` always removes CTI children +before their base rows. The child schema's own deletion marker is not changed by +soft deletion or restoration of the entity. + +All mutations retain query predicates, default scopes, ordering/pagination, and +deleted-row visibility. Select hidden rows explicitly when restoring or purging: + +```ts +const assigned = db.activities.ofVariant('assigned'); +await assigned.onlyDeleted().where(t => t.id, activityId).restore(); +await assigned.withDeleted().where(t => t.id, activityId).hardDelete(); +``` + +Base hooks run before variant hooks, in registration order, once per operation. +`beforeInsert` and `beforeUpdate` receive the combined property-name payload before +it is split into storage tables. `afterInsert` receives the complete decoded row. +`beforeDelete` receives a transaction-bound read query restricted to the captured +mutation targets, for both soft and permanent deletion. Query configuration is +immutable; hooks cannot add delete predicates by changing that query. No hooks run +when an update or delete matches no rows. + +Inserts set configured creation/update timestamps on both storage schemas; updates +advance configured update timestamps, including the base timestamp for child-only +changes. Soft deletion and restoration change only the base deletion marker: +restoration runs no lifecycle hook and neither operation advances update timestamps, +matching ordinary query writes. Restoring a schema without base soft deletion fails. + +Target selection and all storage writes execute in one transaction. Existing +transactions use a savepoint, so catching a failed mutation cannot retain half a CTI +write. Hooks run before commit: database changes roll back on failure, but external +side effects performed by hooks cannot be rolled back. Tracked saves preserve +optimistic row-version checks and apply returned values and new snapshots only +after their save transaction succeeds. --- diff --git a/libs/orm/src/change-tracker.ts b/libs/orm/src/change-tracker.ts index f9e003ad..3de537a9 100644 --- a/libs/orm/src/change-tracker.ts +++ b/libs/orm/src/change-tracker.ts @@ -23,7 +23,7 @@ import { } from '@cleverbrush/knex-schema'; import type { Knex } from 'knex'; import { ConcurrencyError, InvariantViolationError } from './errors.js'; -import { insertVariant } from './variant-write.js'; +import { insertVariant, mutateVariant } from './variant-write.js'; // --------------------------------------------------------------------------- // Types @@ -790,29 +790,86 @@ export class ChangeTracker { const pkInfo = getPrimaryKeyColumns(config.schema); const current = entry.entity as Record; + if (entry.variantKey !== undefined) { + const variants = getVariants(config.schema)!; + const patch: Record = {}; + for (const key of new Set([ + ...Object.keys(entry.originalSnapshot), + ...Object.keys(current) + ])) { + if ( + pkInfo.propertyKeys.includes(key) || + key === variants.discriminatorKey + ) + continue; + if ( + !sameColumn( + entry.originalSnapshot, + key, + current[key] + ) + ) + patch[key] = current[key]; + } + const managed: Record = {}; + const rv = entry.rowVersion; + if (rv?.strategy === 'increment') { + const value = + typeof rv.snapshotValue === 'string' + ? (BigInt(rv.snapshotValue) + 1n).toString() + : Number(rv.snapshotValue ?? 0) + 1; + if ( + typeof value === 'number' && + !Number.isSafeInteger(value) + ) + throw new Error( + 'Row-version increment exceeds the safe integer range; use a bigint storage column' + ); + managed[rv.propertyKey] = value; + } else if (rv?.strategy === 'timestamp') + managed[rv.propertyKey] = new Date(); + let selected = (schemaQuery(trx, config.schema) as any) + .selectVariants([entry.variantKey]) + .unscoped() + .withDeleted() + .where( + pkInfo.propertyKeys[0], + current[pkInfo.propertyKeys[0]] + ); + if (rv) + selected = selected.where( + rv.propertyKey, + rv.snapshotValue + ); + const result = await mutateVariant( + trx, + config.schema, + entry.variantKey, + selected, + 'update', + patch, + managed + ); + if (!result.count && rv) + throw new ConcurrencyError( + tableName, + extractPkValues(config.schema, current), + rv.snapshotValue + ); + if (result.rows[0]) + committedValues.set(current, result.rows[0]); + updated++; + continue; + } + // Build the SET clause — only changed columns, excluding PK const pkPropSet = new Set(pkInfo.propertyKeys); const updateData: Record = {}; - const variant = getVariants(config.schema)?.variants[ - entry.variantKey ?? '' - ]; - const variantColumns = variant - ? buildColumnMap(variant.schema).propToCol - : new Map(); - const variantData: Record = {}; const put = (key: string, value: unknown) => { - const property = - config.schema.introspect().properties[key] ?? - variant?.schema.introspect().properties[key]; - value = encodeJsonColumn(property, value); - const variantColumn = - !propToCol.has(key) && variantColumns.get(key); - if (variantColumn && variant?.storage === 'cti') { - if (variantColumn !== variant.foreignKey) - variantData[variantColumn] = value; - } else - updateData[variantColumn || propToCol.get(key) || key] = - value; + updateData[propToCol.get(key) ?? key] = encodeJsonColumn( + config.schema.introspect().properties[key], + value + ); }; for (const propKey of Object.keys(entry.originalSnapshot)) { if (pkPropSet.has(propKey)) continue; @@ -866,11 +923,7 @@ export class ChangeTracker { // 'manual': caller already set the new value in current } - if ( - Object.keys(updateData).length === 0 && - Object.keys(variantData).length === 0 - ) - continue; + if (Object.keys(updateData).length === 0) continue; // Build WHERE clause with PK + optional rowVersion check let qb = trx(tableName); @@ -890,11 +943,7 @@ export class ChangeTracker { qb = qb.andWhere(rvCol, rv.snapshotValue as any) as any; } - const affected = Object.keys(updateData).length - ? await qb.update(updateData) - : (await qb.forUpdate().first()) - ? 1 - : 0; + const affected = await qb.update(updateData); if (affected === 0 && entry.rowVersion) { const tableName2 = config.schema.getExtension?.( 'tableName' @@ -905,17 +954,6 @@ export class ChangeTracker { entry.rowVersion.snapshotValue ); } - if ( - affected && - variant?.storage === 'cti' && - Object.keys(variantData).length - ) - await trx(variant.tableName!) - .where( - variant.foreignKey!, - current[pkInfo.propertyKeys[0]] as any - ) - .update(variantData); updated++; } @@ -934,6 +972,37 @@ export class ChangeTracker { const pkInfo = getPrimaryKeyColumns(config.schema); const current = entry.entity as Record; + if (entry.variantKey !== undefined) { + let selected = (schemaQuery(trx, config.schema) as any) + .selectVariants([entry.variantKey]) + .unscoped() + .withDeleted() + .where( + pkInfo.propertyKeys[0], + current[pkInfo.propertyKeys[0]] + ); + if (entry.rowVersion) + selected = selected.where( + entry.rowVersion.propertyKey, + entry.rowVersion.snapshotValue + ); + const result = await mutateVariant( + trx, + config.schema, + entry.variantKey, + selected, + 'delete' + ); + if (!result.count && entry.rowVersion) + throw new ConcurrencyError( + tableName, + extractPkValues(config.schema, current), + entry.rowVersion.snapshotValue + ); + deleted++; + continue; + } + let qb = trx(tableName); for (let i = 0; i < pkInfo.propertyKeys.length; i++) { const colName = diff --git a/libs/orm/src/dbset.ts b/libs/orm/src/dbset.ts index 1f91e510..92ffe4e8 100644 --- a/libs/orm/src/dbset.ts +++ b/libs/orm/src/dbset.ts @@ -17,6 +17,7 @@ import type { PolymorphicQueryBuilder, PrimaryKeyValueOf, ReadQueryShape, + ReadVariantMetadata, SchemaAwareQuery, SchemaForValue, VariantReadSchemas @@ -44,9 +45,8 @@ import type { } from './result-types.js'; import { saveGraph } from './save-graph.js'; import { - deleteVariant as _deleteVariant, insertVariant as _insertVariant, - updateVariant as _updateVariant + mutateVariant } from './variant-write.js'; type EntityRowSchema = ObjectSchemaBuilder< @@ -120,10 +120,10 @@ export type EntityQuery< TResult, Writable extends boolean = true > = - SchemaAwareQuery> extends PolymorphicQueryBuilder< - any, - any - > + ReadVariantMetadata> extends { + discriminator: string; + variants: Record; + } ? PolymorphicEntityQuery : TableEntityQuery; @@ -338,7 +338,9 @@ export interface DbSetOperations> { * // a.type === 'assigned' and a.assigneeId is available * ``` */ - ofVariant(variantKey: K): VariantDbSet; + ofVariant< + K extends keyof VariantReadSchemas> & string + >(variantKey: K): VariantDbSet; } // --------------------------------------------------------------------------- @@ -390,14 +392,20 @@ export interface VariantDbSet< ): Promise>; /** - * Update variant-specific columns for rows matched by the current - * WHERE clause. For CTI, only the variant table is updated. + * Update matching base and variant columns atomically, honoring hooks + * and configured timestamps. */ update(patch: VariantUpdatePayload): Promise; - /** Delete rows matched by the current WHERE clause (CTI: atomic). */ + /** Delete matching entities, honoring base-schema soft deletion and hooks. */ delete(): Promise; + /** Permanently delete matching entities, including both CTI rows. */ + hardDelete(): Promise; + + /** Restore matching soft-deleted entities; select them with onlyDeleted/withDeleted. */ + restore(): Promise[]>; + /** Return a new variant view bound to `trx`. */ withTransaction(trx: Knex.Transaction): VariantDbSet; } @@ -444,7 +452,9 @@ function wrapQuery, TResult>( if ( prop === 'update' || prop === 'delete' || - prop === 'insert' + prop === 'insert' || + prop === 'restore' || + prop === 'hardDelete' ) { if (getVariants((entity as any).schema)) { return () => { @@ -817,6 +827,9 @@ function wrapVariantQuery< return async ( payload: Record ): Promise => { + ( + sqb as unknown as PolymorphicQueryBuilder + ).mutationTargets(variantKey); const maybeTrx = ( knexInst as unknown as { isTransaction?: boolean } ).isTransaction @@ -840,74 +853,34 @@ function wrapVariantQuery< }; } - // --- Variant-aware update (collect PKs first, then UPDATE) --- - if (prop === 'update') { - return async ( - set: Record - ): Promise => { - const pkInfo = getPrimaryKeyColumns( - entity.schema as Parameters< - typeof getPrimaryKeyColumns - >[0] - ); - const rows = (await ( - proxy as unknown as { - execute: () => Promise< - Array> - >; - } - ).execute()) as Array>; - const pkProp = pkInfo.propertyKeys[0]; - const pkValues = rows - .map(r => r[pkProp]) - .filter(v => v !== undefined); - const maybeTrx = ( - knexInst as unknown as { isTransaction?: boolean } - ).isTransaction - ? (knexInst as unknown as Knex.Transaction) - : undefined; - await _updateVariant( - knexInst, - entity.schema, + if (prop === 'transacting' || prop === 'withTransaction') { + return (trx: Knex.Transaction) => + wrapVariantQuery( + sqb.transacting(trx), + entity, variantKey, - set, - pkValues, - maybeTrx + trx, + onResults ); - }; } - - // --- Variant-aware delete (collect PKs first, then DELETE) --- - if (prop === 'delete') { - return async (): Promise => { - const pkInfo = getPrimaryKeyColumns( - entity.schema as Parameters< - typeof getPrimaryKeyColumns - >[0] - ); - const rows = (await ( - proxy as unknown as { - execute: () => Promise< - Array> - >; - } - ).execute()) as Array>; - const pkProp = pkInfo.propertyKeys[0]; - const pkValues = rows - .map(r => r[pkProp]) - .filter(v => v !== undefined); - const maybeTrx = ( - knexInst as unknown as { isTransaction?: boolean } - ).isTransaction - ? (knexInst as unknown as Knex.Transaction) - : undefined; - await _deleteVariant( + if ( + prop === 'update' || + prop === 'delete' || + prop === 'restore' || + prop === 'hardDelete' + ) { + return async (patch: Record = {}) => { + const result = await mutateVariant( knexInst, entity.schema, variantKey, - pkValues, - maybeTrx + sqb as unknown as PolymorphicQueryBuilder, + prop, + patch ); + if (prop === 'hardDelete') return result.count; + if (prop === 'restore') + return onResults?.(result.rows) ?? result.rows; }; } diff --git a/libs/orm/src/orm.test.ts b/libs/orm/src/orm.test.ts index ae5d3fe1..740c51f8 100644 --- a/libs/orm/src/orm.test.ts +++ b/libs/orm/src/orm.test.ts @@ -974,6 +974,8 @@ describe('VariantDbSet.update', () => { let mock: MockKnex; beforeEach(() => { mock = makeMockKnex(); + (mock.knex.client as any).transacting = true; + stubTransaction(); }); afterEach(async () => { await mock.knex.destroy(); @@ -986,19 +988,10 @@ describe('VariantDbSet.update', () => { } it('STI: emits an UPDATE on the base table filtered by discriminator', async () => { - // First query: execute() to collect PKs → returns matching rows. - mock.responses.push([ - { - __read_poly: { - id: 3, - type: 'assigned', - todoId: 10, - userId: 1, - assigneeId: 4 - } - } - ]); - // Second query: the UPDATE itself. + // Lock matching base keys, then recheck the captured predicates. + mock.responses.push([{ id: 3 }]); + mock.responses.push([{ id: 3 }]); + // Execute the UPDATE after target selection. mock.responses.push([]); const db = createDb(mock.knex, { activities: ActivityEntitySTI }); @@ -1015,17 +1008,10 @@ describe('VariantDbSet.update', () => { it('CTI: emits an UPDATE on the variant table', async () => { stubTransaction(); - // execute() returns matched base-table rows. - mock.responses.push([ - { - __read_poly: { - id: 5, - type: 'assigned', - todoId: 1, - assigneeId: 4 - } - } - ]); + // Lock base and child rows, then confirm matching base keys. + mock.responses.push([{ id: 5 }]); + mock.responses.push([]); + mock.responses.push([{ id: 5 }]); // UPDATE on the variant table. mock.responses.push([]); @@ -1064,6 +1050,8 @@ describe('VariantDbSet.delete', () => { let mock: MockKnex; beforeEach(() => { mock = makeMockKnex(); + (mock.knex.client as any).transacting = true; + stubTransaction(); }); afterEach(async () => { await mock.knex.destroy(); @@ -1077,18 +1065,9 @@ describe('VariantDbSet.delete', () => { it('STI: emits a DELETE on the base table with discriminator filter', async () => { stubTransaction(); - // execute() to collect PKs. - mock.responses.push([ - { - __read_poly: { - id: 2, - type: 'commented', - todoId: 1, - userId: 7, - body: 'hi' - } - } - ]); + // Lock and confirm matching base keys. + mock.responses.push([{ id: 2 }]); + mock.responses.push([{ id: 2 }]); // The DELETE. mock.responses.push([]); @@ -1105,17 +1084,10 @@ describe('VariantDbSet.delete', () => { it('CTI: deletes variant row first then base row', async () => { stubTransaction(); - // execute() → matched base rows. - mock.responses.push([ - { - __read_poly: { - id: 7, - type: 'assigned', - todoId: 3, - assigneeId: 4 - } - } - ]); + // Lock base and child rows, then confirm matching base keys. + mock.responses.push([{ id: 7 }]); + mock.responses.push([]); + mock.responses.push([{ id: 7 }]); // DELETE from variant table. mock.responses.push([]); // DELETE from base table. diff --git a/libs/orm/src/result-types.ts b/libs/orm/src/result-types.ts index 6e69f6ca..da902ef1 100644 --- a/libs/orm/src/result-types.ts +++ b/libs/orm/src/result-types.ts @@ -8,6 +8,7 @@ import type { Entity, EntityRelations, EntitySchema, + PrimaryKeyOf, RelationInfo, SchemaAwareQuery } from '@cleverbrush/knex-schema'; @@ -167,7 +168,7 @@ export type VariantResult< > = ExtractBranch, K>; /** - * Write payload for `DbSet.insertVariant(key, payload)`. + * Write payload for `DbSet.ofVariant(key).insert(payload)`. * * All columns from the matching variant branch are **optional** (so * auto-generated PKs and columns with DB defaults can be omitted), and the @@ -176,7 +177,7 @@ export type VariantResult< * @example * ```ts * // Only valid fields for the 'assigned' variant; type-checked at compile time: - * await db.activities.insertVariant('assigned', { + * await db.activities.ofVariant('assigned').insert({ * todoId: 42, * userId: 7, * assigneeId: 9, @@ -203,17 +204,23 @@ export type VariantInsertPayload< : never; /** - * Write payload for `EntityQuery.updateVariant(key, set)`. + * Write payload for `DbSet.ofVariant(key).update(patch)`. * - * Same shape as {@link VariantInsertPayload} — partial of the variant - * branch with the discriminator excluded. + * Partial variant branch with primary keys and the discriminator excluded. + * The CTI join key is already absent from the public variant row. * * @public */ export type VariantUpdatePayload< TEntity extends Entity, K extends string -> = VariantInsertPayload; +> = Omit< + VariantInsertPayload, + PrimaryKeyOf> extends readonly (infer P extends + string)[] + ? P + : Extract>, string> +>; /** * @internal Re-export ExtractBranch for use in DbSet/EntityQuery types. diff --git a/libs/orm/src/variant-write.test-d.ts b/libs/orm/src/variant-write.test-d.ts new file mode 100644 index 00000000..6fb2c711 --- /dev/null +++ b/libs/orm/src/variant-write.test-d.ts @@ -0,0 +1,62 @@ +import Knex, { type Knex as KnexType } from 'knex'; +import { expectTypeOf, test } from 'vitest'; +import { + createDb, + defineEntity, + number, + object, + string, + type VariantResult +} from './index.js'; + +const base = object({ + id: number().primaryKey(), + kind: string(), + title: string() +}) + .hasTableName('assets') + .softDelete(); +const body = object({ + ownerId: number().hasColumnName('owner_id'), + caption: string() +}).hasTableName('photos'); +const entity = defineEntity(base) + .discriminator('kind') + .ctiVariant('photo', defineEntity(body), t => t.ownerId); +const db = createDb(Knex({ client: 'pg' }), { assets: entity }); +declare const trx: KnexType.Transaction; + +test('variant mutations preserve types through filters and transaction binding', () => { + const view = db.assets + .ofVariant('photo') + .where(t => t.id, 1) + .withTransaction(trx) + .transacting(trx); + expectTypeOf(view.update({ title: 'base', caption: 'body' })).toEqualTypeOf< + Promise + >(); + expectTypeOf(view.delete()).toEqualTypeOf>(); + expectTypeOf(view.hardDelete()).toEqualTypeOf>(); + expectTypeOf(view.onlyDeleted().restore()).toEqualTypeOf< + Promise[]> + >(); + // @ts-expect-error unknown variant + db.assets.ofVariant('missing'); + // @ts-expect-error discriminator is immutable + view.update({ kind: 'photo' }); + // @ts-expect-error primary keys are immutable + view.update({ id: 2 }); + // @ts-expect-error CTI join keys are managed by the framework + view.update({ ownerId: 2 }); + // @ts-expect-error another branch's field is unavailable + view.update({ documentText: 'text' }); + // @ts-expect-error polymorphic root is read-only + db.assets.hardDelete(); + const projected = view.forVariant('photo', q => + q.select(t => ({ kind: t.kind, label: t.caption })) + ); + // @ts-expect-error projected variants cannot be updated + projected.update({ caption: 'text' }); + // @ts-expect-error projected variants cannot be restored + projected.restore(); +}); diff --git a/libs/orm/src/variant-write.test.ts b/libs/orm/src/variant-write.test.ts new file mode 100644 index 00000000..2baa74d0 --- /dev/null +++ b/libs/orm/src/variant-write.test.ts @@ -0,0 +1,86 @@ +import Knex from 'knex'; +import { afterAll, expect, it } from 'vitest'; +import { createDb, defineEntity, number, object, string } from './index.js'; + +const knex = Knex({ client: 'pg' }); +afterAll(() => knex.destroy()); +const owner = object({ + id: number().primaryKey(), + name: string() +}).hasTableName('owners'); +const base = object({ + id: number().primaryKey(), + kind: string(), + ownerId: number(), + owner: owner.optional() +}).hasTableName('assets'); +const entity = defineEntity(base) + .belongsTo( + t => t.owner, + t => t.ownerId, + t => t.id + ) + .discriminator('kind') + .stiVariant('photo', object({ caption: string() })) + .stiVariant('document', object({ text: string() })); +const view = createDb(knex, { assets: entity }).assets.ofVariant('photo'); + +it('rejects identity changes and unknown fields before opening a connection', async () => { + for (const patch of [ + { id: 2 }, + { kind: 'document' }, + { missing: true }, + { owner: { id: 1 } } + ]) { + await expect(view.update(patch as any)).rejects.toThrow( + /identity|Unknown/ + ); + } +}); + +it('rejects writes through a projected, included or expanded variant view', async () => { + const projected = view.forVariant('photo', q => + q.select(t => ({ id: t.id, kind: t.kind })) + ); + const included = view.include(t => t.owner); + expect(() => (view as any).selectVariants(['photo', 'document'])).toThrow( + /declared variants/ + ); + for (const query of [projected, included] as any[]) { + for (const method of [ + 'insert', + 'update', + 'delete', + 'restore', + 'hardDelete' + ]) { + await expect(query[method]({})).rejects.toThrow( + /unprojected|original variant/ + ); + } + } +}); + +it('rejects composite primary keys before executing writes', async () => { + const composite = defineEntity( + object({ a: number(), b: number(), kind: string() }) + .hasTableName('composite') + .hasPrimaryKey(['a', 'b']) + ) + .discriminator('kind') + .stiVariant('photo', object({ caption: string() })); + const query = createDb(knex, { assets: composite }).assets.ofVariant( + 'photo' + ); + for (const method of [ + 'insert', + 'update', + 'delete', + 'restore', + 'hardDelete' + ] as const) { + await expect((query[method] as any)({})).rejects.toThrow( + /single-column/ + ); + } +}); diff --git a/libs/orm/src/variant-write.ts b/libs/orm/src/variant-write.ts index 4240009c..a5b8c452 100644 --- a/libs/orm/src/variant-write.ts +++ b/libs/orm/src/variant-write.ts @@ -1,366 +1,364 @@ -// @cleverbrush/orm — Polymorphic write helpers -// -// Runtime implementations for `insertVariant`, `updateVariant`, -// `deleteVariant`, and `findVariant`. These are invoked from `DbSet` / -// `EntityQuery` and handle the two-table atomicity required for CTI -// (Class Table Inheritance) variants. - +// Polymorphic writes share the ordinary storage codecs and timestamp pipeline. +// Hooks operate on one logical entity before payloads are split across tables. import { buildColumnMap, - encodeJsonColumn, getPrimaryKeyColumns, getVariants, object, + type PolymorphicQueryBuilder, query as schemaQuery } from '@cleverbrush/knex-schema'; -import type { ObjectSchemaBuilder } from '@cleverbrush/schema'; +import type { ObjectSchemaBuilder, SchemaBuilder } from '@cleverbrush/schema'; import type { Knex } from 'knex'; -// --------------------------------------------------------------------------- -// Internal helpers -// --------------------------------------------------------------------------- +type Row = Record; +type Schema = ObjectSchemaBuilder; +type Mutation = 'update' | 'delete' | 'restore' | 'hardDelete'; + +const lifecycle = new Set([ + 'beforeInsert', + 'afterInsert', + 'beforeUpdate', + 'beforeDelete' +]); -/** Keep table metadata and hooks without constructing a polymorphic reader. */ +/** Keep physical-table metadata; lifecycle hooks run once on the logical row. */ function storageSchema( - schema: any, - properties = schema.introspect().properties -): ObjectSchemaBuilder { - let stored: ObjectSchemaBuilder = - object(properties); + schema: Schema, + properties: Record< + string, + SchemaBuilder + > = schema.introspect().properties, + relations: readonly { name: string }[] = [] +): Schema { + const navigation = new Set( + [ + ...((schema.getExtension('relations') as + | { name: string }[] + | undefined) ?? []), + ...relations + ].map(relation => relation.name) + ); + let stored: Schema = object( + Object.fromEntries( + Object.entries(properties).filter(([name]) => !navigation.has(name)) + ) + ); for (const [key, value] of Object.entries( schema.introspect().extensions ?? {} )) { - if (key !== 'variants' && key !== 'polymorphicVariants') - stored = stored.withExtension(key, value) as typeof stored; + if ( + key !== 'variants' && + key !== 'polymorphicVariants' && + key !== 'defaultScope' && + !lifecycle.has(key) + ) + stored = stored.withExtension(key, value) as Schema; } return stored; } -/** - * Resolve the variant config for a schema, throwing if the schema is not - * polymorphic or if the requested variant key is unknown. - * @internal - */ -function requireVariantSpec(schema: any, variantKey: string) { - const raw = getVariants(schema); - if (!raw) { - throw new Error( - `insertVariant / deleteVariant / updateVariant / findVariant: ` + - `entity schema is not polymorphic (no variants declared).` - ); - } - const spec = raw.variants[variantKey]; - if (!spec) { - const known = Object.keys(raw.variants).join(', '); - throw new Error( - `Variant key "${variantKey}" is unknown. Known variants: ${known}.` - ); - } - return { config: raw, spec }; -} - -/** - * Derive the SQL discriminator column name from the schema's propToCol map. - * @internal - */ -function _resolveDiscriminatorColumn( - schema: any, - discriminatorKey: string -): string { - const { propToCol } = buildColumnMap(schema); - return propToCol.get(discriminatorKey) ?? discriminatorKey; +function layout(schema: Schema, key: string) { + const config = getVariants(schema); + const spec = config?.variants[key]; + if (!config || !spec) + throw new Error(`Unknown polymorphic variant: ${key}`); + const pk = getPrimaryKeyColumns(schema); + if (pk.propertyKeys.length !== 1) + throw new Error('Variant writes require a single-column primary key'); + const baseMap = buildColumnMap(storageSchema(schema)); + const body = storageSchema(spec.schema, undefined, spec.relations); + const bodyMap = buildColumnMap(body); + const base = storageSchema( + schema, + spec.storage === 'sti' + ? { + ...schema.introspect().properties, + ...spec.schema.introspect().properties + } + : undefined, + spec.storage === 'sti' + ? [ + ...spec.relations, + ...((spec.schema.getExtension('relations') as + | { name: string }[] + | undefined) ?? []) + ] + : [] + ); + return { + config, + spec, + base, + body, + baseMap, + bodyMap, + pk: pk.propertyKeys[0], + pkColumn: pk.columnNames[0], + foreignKey: bodyMap.colToProp.get(spec.foreignKey ?? ''), + table: schema.getExtension('tableName') as string, + discriminatorColumn: baseMap.propToCol.get(config.discriminatorKey)!, + hooks: (name: string): Function[] => + [schema, spec.schema].flatMap( + s => (s.getExtension(name) as Function[] | undefined) ?? [] + ) + }; } -// --------------------------------------------------------------------------- -// insertVariant -// --------------------------------------------------------------------------- - -/** - * Transactionally insert a polymorphic entity row. - * - * For CTI variants: - * 1. Inserts the base row with `discriminatorColumn = variantKey` and all - * base-table payload columns; captures the auto-generated PK. - * 2. Inserts the variant row with the FK set to the base PK and all - * variant-table payload columns. - * - * For STI variants: - * - Inserts a single row into the base table with the discriminator set. - * - * @internal - */ -export async function insertVariant( - knex: Knex, - schema: any, - variantKey: string, - payload: Record, - trx?: Knex.Transaction -): Promise> { - // Validate before opening a transaction so errors surface immediately. - const { config, spec } = requireVariantSpec(schema, variantKey); - - const run = async (t: Knex): Promise> => { - const discKey = config.discriminatorKey; - const baseTableName = schema.getExtension?.('tableName') as string; - if (!baseTableName) { - throw new Error('insertVariant: base schema has no table name.'); - } - - if (spec.storage === 'sti') { - // Single-table: insert into base table with discriminator column - const row = { ...payload, [discKey]: variantKey }; - const props = { - ...schema.introspect().properties, - ...spec.schema.introspect().properties - }; - const stored = storageSchema(schema, props); - const result = await schemaQuery(t, stored).insert(row as any); - return result as Record; - } - - // CTI: two-table insert - const variantSchema = spec.schema; - const variantTableName = spec.tableName as string; - const fkCol = spec.foreignKey as string; - - // Split payload between base-schema columns and variant-schema columns - const baseIntrospected = (schema as any).introspect?.() as { - properties?: Record; - }; - const variantIntrospected = variantSchema.introspect?.() as { - properties?: Record; - }; - const basePropKeys = new Set( - Object.keys(baseIntrospected?.properties ?? {}) - ); - const variantPropKeys = new Set( - Object.keys(variantIntrospected?.properties ?? {}) - ); - - const basePayload: Record = { [discKey]: variantKey }; - const variantPayload: Record = {}; - - for (const [key, val] of Object.entries(payload)) { - // Discriminator is set automatically; FK is set from base PK - if (key === discKey || key === fkCol) continue; - if (basePropKeys.has(key)) { - basePayload[key] = val; - } else if (variantPropKeys.has(key)) { - variantPayload[key] = val; - } - // Keys that match neither schema are silently dropped - } - - // 1. Decode base RETURNING values before using the PK. In particular, - // bigint IDs must be text-cast in SQL before driver parsers can round them. - const baseRow = (await schemaQuery(t, storageSchema(schema)).insert( - basePayload as any - )) as Record; - - // Resolve the base PK value from the returned row - const pkInfo = getPrimaryKeyColumns(schema); - if (pkInfo.propertyKeys.length === 0) { +function assertPatch(meta: ReturnType, patch: Row): void { + const protectedKeys = new Set([ + meta.pk, + meta.pkColumn, + meta.config.discriminatorKey, + meta.discriminatorColumn, + ...(meta.spec.storage === 'cti' + ? [meta.foreignKey, meta.spec.foreignKey] + : []) + ]); + for (const prop of Object.keys(patch)) { + if (protectedKeys.has(prop)) throw new Error( - `insertVariant: base schema has no primary key declared.` + `Variant updates cannot change identity property "${prop}"` ); - } - if (pkInfo.propertyKeys.length > 1) { - throw new Error( - `insertVariant: composite primary keys are not yet supported for polymorphic CTI inserts.` - ); - } - const pkPropKey = pkInfo.propertyKeys[0]; - const pkColName = pkInfo.columnNames[0]; - const pkValue = baseRow[pkPropKey] ?? baseRow[pkColName]; - if (pkValue === undefined) { - throw new Error( - `insertVariant: could not resolve PK value from base insert result.` - ); - } + if ( + !meta.baseMap.propToCol.has(prop) && + !meta.bodyMap.propToCol.has(prop) + ) + throw new Error(`Unknown variant property: ${prop}`); + } +} - // 2. Insert variant row using raw knex (not SchemaQueryBuilder) so we - // can set the FK column even if it isn't in the schema's column map. - const { propToCol: varPropToCol } = buildColumnMap(variantSchema); - const variantRow: Record = { [fkCol]: pkValue }; - for (const [propKey, val] of Object.entries(variantPayload)) { - const colName = varPropToCol.get(propKey) ?? propKey; - variantRow[colName] = encodeJsonColumn( - variantSchema.introspect().properties[propKey], - val - ); - } - const discColInVariant = varPropToCol.get(discKey); - if (discColInVariant) { - variantRow[discColInVariant] = variantKey; - } - await (t as unknown as Knex)(variantTableName).insert(variantRow); +async function prepare(hooks: Function[], input: Row): Promise { + let data = { ...input }; + for (const hook of hooks) data = (await hook(data)) ?? data; + return data; +} - // Read the completed branch inside the same transaction for one consistent storage representation. - const result = await (schemaQuery(t, schema) as any) - .selectVariants([variantKey]) - .unscoped() - .withDeleted() - .where(pkPropKey, pkValue) - .first(); - if (!result) - throw new Error( - 'insertVariant: inserted row could not be read back' - ); - return result; +function split(meta: ReturnType, data: Row): [Row, Row] { + const base: Row = {}; + const body: Row = {}; + for (const [key, value] of Object.entries(data)) { + if (meta.baseMap.propToCol.has(key)) base[key] = value; + else if (meta.bodyMap.propToCol.has(key)) body[key] = value; + } + return [base, body]; +} + +// An STI body shares the physical base table, but can name additional timestamp +// columns. The ordinary writer handles the base schema's timestamps itself. +function stiTimestamps( + meta: ReturnType, + db: Knex, + data: Row, + insert: boolean +): Row { + const timestamps = meta.body.getExtension('timestamps') as + | { createdAt: string; updatedAt: string } + | undefined; + if (!timestamps) return data; + const columns = buildColumnMap(meta.base).colToProp; + return { + ...data, + ...(insert + ? { + [columns.get(timestamps.createdAt) ?? timestamps.createdAt]: + db.fn.now() + } + : {}), + [columns.get(timestamps.updatedAt) ?? timestamps.updatedAt]: db.fn.now() }; +} - if (trx) return run(trx as unknown as Knex); - return (knex as Knex).transaction(t => run(t as unknown as Knex)); +// RETURNING must not replay default-scope callbacks or hide a row which a +// successful mutation just moved outside a scope. Decode the stored branch. +function storedVariantQuery(db: Knex, schema: Schema, key: string) { + const config = getVariants(schema)!; + const spec = config.variants[key]; + const stored = schema + .withExtension('defaultScope', undefined) + .withExtension('variants', { + ...config, + variants: { + [key]: { + ...spec, + schema: spec.schema + .withExtension('defaultScope', undefined) + .withExtension('softDelete', undefined) + } + } + }); + return (schemaQuery(db, stored) as any).selectVariants([key]).withDeleted(); } -// --------------------------------------------------------------------------- -// updateVariant -// --------------------------------------------------------------------------- +function readRows( + db: Knex, + schema: Schema, + key: string, + ids: readonly unknown[] +): Promise { + const pk = getPrimaryKeyColumns(schema).propertyKeys[0]; + return storedVariantQuery(db, schema, key).whereIn(pk, ids).execute(); +} -/** - * Update variant-specific columns for all rows matching the current query - * whose discriminator equals `variantKey`. - * - * For CTI: runs UPDATE on the variant table, joining on the FK = base PK. - * For STI: runs UPDATE on the base table filtered by discriminator. - * - * The `getPks` callback must return the PK values of all rows the caller's - * WHERE clause matched (resolved before this function is called). - * - * @internal - */ -export async function updateVariant( +/** @internal Atomic insert, including hooks and both CTI storage rows. */ +export async function insertVariant( knex: Knex, - schema: any, + schema: Schema, variantKey: string, - set: Record, - pkValues: readonly unknown[], + payload: Row, trx?: Knex.Transaction -): Promise { - if (pkValues.length === 0) return; - - const db: Knex = (trx as unknown as Knex) ?? knex; - const { config, spec } = requireVariantSpec(schema, variantKey); - const discKey = config.discriminatorKey; - - if (spec.storage === 'sti') { - // STI: update the base table restricted to the discriminator value - const { propToCol } = buildColumnMap(schema); - const discCol = propToCol.get(discKey) ?? discKey; - const pkInfo = getPrimaryKeyColumns(schema); - if (pkInfo.columnNames.length === 0) { - throw new Error( - `updateVariant: base schema has no primary key declared.` - ); - } - if (pkInfo.columnNames.length > 1) { - throw new Error( - `updateVariant: composite primary keys are not supported for STI updates.` - ); - } - const pkColName = pkInfo.columnNames[0]; - const baseTable = schema.getExtension?.('tableName') as string; - - const updateData: Record = {}; - for (const [propKey, val] of Object.entries(set)) { - updateData[propToCol.get(propKey) ?? propKey] = encodeJsonColumn( - spec.schema.introspect().properties[propKey] ?? - schema.introspect().properties[propKey], - val +): Promise { + const meta = layout(schema, variantKey); + // A transaction on an existing transaction creates a savepoint. A caller + // may catch a failed write without retaining a partially inserted entity. + return (trx ?? knex).transaction(async t => { + const data = await prepare(meta.hooks('beforeInsert'), { + ...payload, + [meta.config.discriminatorKey]: variantKey + }); + data[meta.config.discriminatorKey] = variantKey; + const [baseData, bodyData] = split(meta, data); + let result: Row; + if (meta.spec.storage === 'sti') { + result = await schemaQuery(t, meta.base).insert( + stiTimestamps(meta, t, { ...baseData, ...bodyData }, true) ); + } else { + const base = await schemaQuery(t, meta.base).insert(baseData); + // foreignKey is the schema property name, not its SQL column name. + bodyData[meta.foreignKey!] = base[meta.pk]; + if (meta.bodyMap.propToCol.has(meta.config.discriminatorKey)) + bodyData[meta.config.discriminatorKey] = variantKey; + await schemaQuery(t, meta.body).insert(bodyData); + [result] = await readRows(t, schema, variantKey, [base[meta.pk]]); + if (!result) + throw new Error('Inserted variant could not be read back'); } - - await db(baseTable) - .whereIn(pkColName, pkValues as any[]) - .andWhere(discCol, variantKey) - .update(updateData); - return; - } - - // CTI: update the variant table - const variantSchema = spec.schema; - const variantTableName = spec.tableName as string; - const fkCol = spec.foreignKey as string; - const { propToCol: varPropToCol } = buildColumnMap(variantSchema); - - const updateData: Record = {}; - for (const [propKey, val] of Object.entries(set)) { - // Don't allow updating the FK (it's the join column) - if (propKey === fkCol) continue; - updateData[varPropToCol.get(propKey) ?? propKey] = encodeJsonColumn( - variantSchema.introspect().properties[propKey], - val - ); - } - - if (Object.keys(updateData).length === 0) return; - - await db(variantTableName) - .whereIn(fkCol, pkValues as any[]) - .update(updateData); + for (const hook of meta.hooks('afterInsert')) await hook(result); + return result; + }); } -// --------------------------------------------------------------------------- -// deleteVariant -// --------------------------------------------------------------------------- +/** @internal Shared result for explicit writes and optimistic tracked saves. */ +export interface VariantMutationResult { + rows: Row[]; + count: number; +} /** - * Delete rows for a specific variant key, given their base-table PK values. - * - * For CTI: deletes variant rows first (FK-constraint order), then base rows. - * For STI: deletes the single base-table row filtered by discriminator. - * - * @internal + * @internal Apply a single variant mutation within a transaction/savepoint. + * Selection is untracked and locks base rows, then CTI rows, in primary-key order. + * Recheck predicates after waiting for locks before invoking any lifecycle hook. */ -export async function deleteVariant( +export async function mutateVariant( knex: Knex, - schema: any, + schema: Schema, variantKey: string, - pkValues: readonly unknown[], - trx?: Knex.Transaction -): Promise { - if (pkValues.length === 0) return; - - const run = async (t: Knex): Promise => { - const { config, spec } = requireVariantSpec(schema, variantKey); - const discKey = config.discriminatorKey; - const baseTable = schema.getExtension?.('tableName') as string; - const pkInfo = getPrimaryKeyColumns(schema); - if (pkInfo.columnNames.length === 0) { - throw new Error( - `deleteVariant: base schema has no primary key declared.` - ); + query: PolymorphicQueryBuilder, + operation: Mutation, + patch: Row = {}, + managedValues: Row = {} +): Promise { + const meta = layout(schema, variantKey); + query.mutationTargets(variantKey); // Validate before executing any SQL. + if (operation === 'update') assertPatch(meta, patch); + if (operation === 'restore' && !meta.base.getExtension('softDelete')) + throw new Error( + 'Schema does not have soft delete enabled. Use .softDelete() on the base schema.' + ); + return knex.transaction(async t => { + const selected = query.transacting(t); + const targets = () => + t(meta.table) + .whereIn(meta.pkColumn, selected.mutationTargets(variantKey)) + .andWhere(meta.discriminatorColumn, variantKey); + // Cast before pg's parsers, which may otherwise round a bigint key. + const keyColumn = t.raw('cast(?? as text) as ??', [ + meta.pkColumn, + meta.pk + ]); + const locked = await targets() + .orderBy(meta.pkColumn) + .forUpdate() + .select(keyColumn); + let ids = locked.map(row => row[meta.pk]); + if (!ids.length) return { rows: [], count: 0 }; + if (meta.spec.storage === 'cti') { + await t(meta.spec.tableName!) + .whereIn(meta.spec.foreignKey!, ids) + .orderBy(meta.spec.foreignKey!) + .forUpdate() + .select(t.raw('1')); } - if (pkInfo.columnNames.length > 1) { - throw new Error( - `deleteVariant: composite primary keys are not supported for variant deletes.` - ); + const confirmed = await targets() + .whereIn(meta.pkColumn, ids) + .select(keyColumn); + ids = confirmed.map(row => row[meta.pk]); + if (!ids.length) return { rows: [], count: 0 }; + const baseQuery = () => + schemaQuery(t, meta.base) + .unscoped() + .withDeleted() + .whereIn(meta.pk, ids); + const bodyQuery = () => + schemaQuery(t, meta.body) + .unscoped() + .withDeleted() + .whereIn(meta.foreignKey!, ids); + + if (operation === 'update') { + const data = await prepare(meta.hooks('beforeUpdate'), patch); + assertPatch(meta, data); + Object.assign(data, managedValues); // Tracker owns automatic row versions. + const [baseData, bodyData] = split(meta, data); + if (meta.spec.storage === 'sti') { + const combined = stiTimestamps( + meta, + t, + { ...baseData, ...bodyData }, + false + ); + if ( + Object.keys(combined).length || + meta.base.getExtension('timestamps') + ) + await baseQuery().update(combined); + } else { + if ( + Object.keys(baseData).length || + meta.base.getExtension('timestamps') + ) + await baseQuery().update(baseData); + if ( + Object.keys(bodyData).length || + meta.body.getExtension('timestamps') + ) + await bodyQuery().update(bodyData); + } + return { + rows: await readRows(t, schema, variantKey, ids), + count: ids.length + }; } - const pkColName = pkInfo.columnNames[0]; - - if (spec.storage === 'sti') { - const { propToCol } = buildColumnMap(schema); - const discCol = propToCol.get(discKey) ?? discKey; - await (t as unknown as Knex)(baseTable) - .whereIn(pkColName, pkValues as any[]) - .andWhere(discCol, variantKey) - .delete(); - return; + if (operation === 'restore') { + await baseQuery().restore(); + return { + rows: await readRows(t, schema, variantKey, ids), + count: ids.length + }; } - - // CTI: variant rows first, then base rows - const variantTableName = spec.tableName as string; - const fkCol = spec.foreignKey as string; - - await (t as unknown as Knex)(variantTableName) - .whereIn(fkCol, pkValues as any[]) - .delete(); - - await (t as unknown as Knex)(baseTable) - .whereIn(pkColName, pkValues as any[]) - .delete(); - }; - - if (trx) return run(trx as unknown as Knex); - return (knex as Knex).transaction(t => run(t as unknown as Knex)); + // IDs already include scopes and pagination; do not apply an offset twice. + const hookQuery = storedVariantQuery(t, schema, variantKey).whereIn( + meta.pk, + ids + ); + for (const hook of meta.hooks('beforeDelete')) await hook(hookQuery); + if (operation === 'delete' && meta.base.getExtension('softDelete')) { + await baseQuery().delete(); + } else { + if (meta.spec.storage === 'cti') await bodyQuery().hardDelete(); + await baseQuery().hardDelete(); + } + return { rows: [], count: ids.length }; + }); }