Skip to content

Commit 56b67e3

Browse files
committed
fix(sql): make migrated created_at NOT NULL
migrateScheduleDates() adds each epoch column as nullable, fills it, then renames it over the legacy column. created_at therefore ended up nullable, while both a new 0.8 table and the 0.7 one declare it NOT NULL. On PostgreSQL and MySQL, the migration now sets the epoch column NOT NULL once it is filled, before the legacy columns are dropped. MySQL commits each schema change, so a failure there leaves a table the next run resumes instead of one that already looks migrated. SQLite cannot add NOT NULL to an existing column, so its epoch column is created NOT NULL with a default of 0, which every row then overwrites. Removing that default would need a table rebuild, which Kysely does not provide, so both schema services keep the same result. A row without created_at, possible only in a table changed by hand, now fails the migration with an error that names the schedule, before the first schema change. It used to fail with a raw database error and, on MySQL, left the epoch columns behind.
1 parent 197c298 commit 56b67e3

5 files changed

Lines changed: 121 additions & 19 deletions

File tree

‎src/services/knex_queue_schema.ts‎

Lines changed: 24 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import type { Knex } from 'knex'
22
import {
33
assertTimeZone,
44
epochColumn,
5-
legacyScheduleDateToEpoch,
5+
legacyScheduleRowToEpochs,
66
processTimeZone,
77
scheduleDatesMigrationState,
88
type ScheduleDateColumn,
@@ -235,12 +235,7 @@ export class KnexQueueSchemaService {
235235
const timeZone = dialect === 'sqlite' ? 'UTC' : writerTimeZone
236236
const updates = rows.map((row) => ({
237237
id: row.id,
238-
values: Object.fromEntries(
239-
state.legacy.map((name) => [
240-
epochColumn(name),
241-
legacyScheduleDateToEpoch(row[name], timeZone),
242-
])
243-
),
238+
values: legacyScheduleRowToEpochs(row, state.legacy, timeZone),
244239
}))
245240

246241
await this.#dropNextRunIndexes(trx, dialect, tableName)
@@ -250,21 +245,41 @@ export class KnexQueueSchemaService {
250245
)
251246
if (missingEpochColumns.length > 0) {
252247
await trx.schema.alterTable(tableName, (table) => {
253-
for (const name of missingEpochColumns) table.bigint(epochColumn(name)).nullable()
248+
for (const name of missingEpochColumns) {
249+
// SQLite cannot add NOT NULL to an existing column, only to a new
250+
// one with a default. Every row gets its value below.
251+
if (name === 'created_at' && dialect === 'sqlite') {
252+
table.bigint(epochColumn(name)).notNullable().defaultTo(0)
253+
} else {
254+
table.bigint(epochColumn(name)).nullable()
255+
}
256+
}
254257
})
255258
}
256259

257260
for (const { id, values } of updates) {
258261
await trx(tableName).where('id', id).update(values)
259262
}
260263

264+
// created_at is required, as in a new table. Before the legacy columns
265+
// are dropped, so a failure on MySQL leaves a table the next run resumes.
266+
const epochColumns = [...new Set([...state.legacy, ...state.withEpochColumn])]
267+
if (epochColumns.includes('created_at') && dialect !== 'sqlite') {
268+
await trx.schema.alterTable(tableName, (table) => {
269+
if (dialect === 'pg') {
270+
table.dropNullable(epochColumn('created_at'))
271+
} else {
272+
table.bigint(epochColumn('created_at')).notNullable().alter()
273+
}
274+
})
275+
}
276+
261277
if (state.legacy.length > 0) {
262278
await trx.schema.alterTable(tableName, (table) => {
263279
table.dropColumns(...state.legacy)
264280
})
265281
}
266282

267-
const epochColumns = [...new Set([...state.legacy, ...state.withEpochColumn])]
268283
await trx.schema.alterTable(tableName, (table) => {
269284
for (const name of epochColumns) table.renameColumn(epochColumn(name), name)
270285
})

‎src/services/kysely_queue_schema.ts‎

Lines changed: 27 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { sql, type AlterTableColumnAlteringBuilder, type Kysely, type Transactio
22
import {
33
assertTimeZone,
44
epochColumn,
5-
legacyScheduleDateToEpoch,
5+
legacyScheduleRowToEpochs,
66
processTimeZone,
77
scheduleDatesMigrationState,
88
type ScheduleDateColumn,
@@ -231,19 +231,23 @@ export class KyselyQueueSchemaService<DB> {
231231
const timeZone = this.#dialect === 'sqlite' ? 'UTC' : writerTimeZone
232232
const updates = rows.map((row) => ({
233233
id: row.id,
234-
values: Object.fromEntries(
235-
state.legacy.map((name) => [
236-
epochColumn(name),
237-
legacyScheduleDateToEpoch(row[name], timeZone),
238-
])
239-
),
234+
values: legacyScheduleRowToEpochs(row, state.legacy, timeZone),
240235
}))
241236

242237
await this.#dropNextRunIndexes(trx, tableName)
243238

244239
for (const name of state.legacy) {
245240
if (state.withEpochColumn.includes(name)) continue
246-
await trx.schema.alterTable(tableName).addColumn(epochColumn(name), 'bigint').execute()
241+
await trx.schema
242+
.alterTable(tableName)
243+
.addColumn(epochColumn(name), 'bigint', (column) =>
244+
// SQLite cannot add NOT NULL to an existing column, only to a new
245+
// one with a default. Every row gets its value below.
246+
name === 'created_at' && this.#dialect === 'sqlite'
247+
? column.notNull().defaultTo(0)
248+
: column
249+
)
250+
.execute()
247251
}
248252

249253
for (const { id, values } of updates) {
@@ -254,11 +258,25 @@ export class KyselyQueueSchemaService<DB> {
254258
.execute()
255259
}
256260

261+
// created_at is required, as in a new table. Before the legacy columns
262+
// are dropped, so a failure on MySQL leaves a table the next run resumes.
263+
const epochColumns = new Set([...state.legacy, ...state.withEpochColumn])
264+
if (epochColumns.has('created_at') && this.#dialect !== 'sqlite') {
265+
const alterTable = trx.schema.alterTable(tableName)
266+
await (
267+
this.#dialect === 'postgres'
268+
? alterTable.alterColumn(epochColumn('created_at'), (column) => column.setNotNull())
269+
: alterTable.modifyColumn(epochColumn('created_at'), 'bigint', (column) =>
270+
column.notNull()
271+
)
272+
).execute()
273+
}
274+
257275
for (const name of state.legacy) {
258276
await trx.schema.alterTable(tableName).dropColumn(name).execute()
259277
}
260278

261-
for (const name of new Set([...state.legacy, ...state.withEpochColumn])) {
279+
for (const name of epochColumns) {
262280
await trx.schema.alterTable(tableName).renameColumn(epochColumn(name), name).execute()
263281
}
264282

‎src/services/schedule_dates.ts‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,34 @@ export function legacyScheduleDateToEpoch(
110110
return epoch
111111
}
112112

113+
/**
114+
* The epoch values of the legacy date columns of one schedule row, keyed by
115+
* epoch column. created_at is required once migrated, so a row without one
116+
* fails the migration here, before the first schema change.
117+
*/
118+
export function legacyScheduleRowToEpochs(
119+
row: Record<string, string | number | null>,
120+
columns: readonly ScheduleDateColumn[],
121+
wallClockTimeZone: string
122+
): Record<string, number | null> {
123+
const values: Record<string, number | null> = {}
124+
125+
for (const name of columns) {
126+
const epoch = legacyScheduleDateToEpoch(row[name], wallClockTimeZone)
127+
128+
if (name === 'created_at' && epoch === null) {
129+
throw new Error(
130+
`Cannot migrate schedule "${row.id}": its created_at is empty, and the migrated column ` +
131+
`is required. Set a date, then run the migration again.`
132+
)
133+
}
134+
135+
values[epochColumn(name)] = epoch
136+
}
137+
138+
return values
139+
}
140+
113141
type WallClockParts = [
114142
year: number,
115143
month: number,

‎tests/schedule_dates.spec.ts‎

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,8 @@
11
import { test } from '@japa/runner'
2-
import { legacyScheduleDateToEpoch } from '../src/services/schedule_dates.js'
2+
import {
3+
legacyScheduleDateToEpoch,
4+
legacyScheduleRowToEpochs,
5+
} from '../src/services/schedule_dates.js'
36

47
const toIso = (value: number | null) => (value === null ? null : new Date(value).toISOString())
58

@@ -54,3 +57,27 @@ test.group('legacyScheduleDateToEpoch', () => {
5457
assert.throws(() => legacyScheduleDateToEpoch('not a date', 'UTC'), /Cannot convert/)
5558
})
5659
})
60+
61+
test.group('legacyScheduleRowToEpochs', () => {
62+
test('converts the legacy columns of a row into its epoch columns', ({ assert }) => {
63+
const row = { id: 'a', next_run_at: '2026-07-01 12:00:00', last_run_at: null, created_at: 0 }
64+
65+
assert.deepEqual(
66+
legacyScheduleRowToEpochs(row, ['next_run_at', 'last_run_at', 'created_at'], 'Europe/Paris'),
67+
{
68+
next_run_at__epoch: Date.parse('2026-07-01T10:00:00.000Z'),
69+
last_run_at__epoch: null,
70+
created_at__epoch: 0,
71+
}
72+
)
73+
})
74+
75+
test('rejects a row without created_at, which the migrated table requires', ({ assert }) => {
76+
const row = { id: 'no-created-at', next_run_at: null, created_at: null }
77+
78+
assert.throws(
79+
() => legacyScheduleRowToEpochs(row, ['next_run_at', 'created_at'], 'UTC'),
80+
/Cannot migrate schedule "no-created-at": its created_at is empty/
81+
)
82+
})
83+
})

‎tests/schedule_dates_migration.spec.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,6 +220,7 @@ for (const dialect of ['sqlite', 'postgres', 'mysql'] as const) {
220220

221221
const column = await connection(TABLE).columnInfo('next_run_at')
222222
assert.match(column.type, /int/i)
223+
assert.isFalse((await connection(TABLE).columnInfo('created_at')).nullable)
223224
assert.equal((await connection(TABLE).where('id', 'legacy').first()).tenant, 'acme')
224225

225226
await assertMigratedSchedules(assert, adapter(), createdAt)
@@ -249,6 +250,10 @@ for (const dialect of ['sqlite', 'postgres', 'mysql'] as const) {
249250
timezone: WRITER_TIME_ZONE,
250251
})
251252

253+
// SQLite migrates in one transaction, so it never resumes from nullable epoch columns.
254+
if (dialect !== 'sqlite') {
255+
assert.isFalse((await connection(TABLE).columnInfo('created_at')).nullable)
256+
}
252257
await assertMigratedSchedules(assert, adapter(), createdAt)
253258
})
254259

@@ -487,6 +492,10 @@ for (const dialect of ['sqlite', 'postgres', 'mysql'] as const) {
487492
const table = (await connection.introspection.getTables()).find(({ name }) => name === TABLE)
488493
return new Map(table!.columns.map(({ name, dataType }) => [name, dataType]))
489494
}
495+
const isNullable = async (column: string) => {
496+
const table = (await connection.introspection.getTables()).find(({ name }) => name === TABLE)
497+
return table!.columns.find(({ name }) => name === column)!.isNullable
498+
}
490499

491500
test('converts the 0.7 dates to epoch milliseconds', async ({ assert }) => {
492501
const createdAt = await insertLegacySchedule()
@@ -495,6 +504,7 @@ for (const dialect of ['sqlite', 'postgres', 'mysql'] as const) {
495504
await schema.migrateScheduleDates(TABLE, { timezone: WRITER_TIME_ZONE })
496505

497506
assert.match((await columnTypes()).get('next_run_at')!, /int/i)
507+
assert.isFalse(await isNullable('created_at'))
498508
await assertMigratedSchedules(assert, adapter(), createdAt)
499509
})
500510

@@ -522,6 +532,10 @@ for (const dialect of ['sqlite', 'postgres', 'mysql'] as const) {
522532
timezone: WRITER_TIME_ZONE,
523533
})
524534

535+
// SQLite migrates in one transaction, so it never resumes from nullable epoch columns.
536+
if (dialect !== 'sqlite') {
537+
assert.isFalse(await isNullable('created_at'))
538+
}
525539
await assertMigratedSchedules(assert, adapter(), createdAt)
526540
})
527541

0 commit comments

Comments
 (0)