1
import { SessionFormatError, SessionFormatUnsupportedMigrationError } from './error.ts'2
import {3
snapshotSessionFormatHeader,4
sessionFormatCount,5
sessionFormatVersion,6
} from './json.ts'7
import type {8
SessionFormatChain,9
SessionFormatChainOptions,10
SessionFormatEvent,11
SessionFormatEventRun,12
SessionFormatHeader,13
SessionFormatMigration,14
SessionFormatMigrationContext,15
SessionFormatMigrationStage,16
SessionFormatMigrationStream,17
} from './types.ts'19
/**20
* Validate and freeze one adjacent migration declaration.21
* @param migration - named exact adjacent conversion.22
* @returns immutable validated declaration.23
*/24
export function defineSessionFormatMigration(migration: SessionFormatMigration): SessionFormatMigration {25
if (typeof migration.name !== 'string' || migration.name.length === 0) {26
throw new SessionFormatError('Session migration name must be a non-empty string')27
}28
const from = sessionFormatVersion(migration.fromVersion, `${migration.name} fromVersion`)29
const to = sessionFormatVersion(migration.toVersion, `${migration.name} toVersion`)30
if (to !== from + 1) {31
throw new SessionFormatError(`${migration.name} must declare adjacent v${from}->v${from + 1}`)32
}33
return Object.freeze({ ...migration })34
}36
/**37
* Compile a unique, complete adjacent migration chain.38
* @param options - current version, adjacent declarations, and current restorer.39
* @returns immutable planner and streaming migration compiler.40
*/41
export function createSessionFormatChain(options: SessionFormatChainOptions): SessionFormatChain {42
return new CompiledSessionFormatChain(options)43
}45
class CompiledSessionFormatChain implements SessionFormatChain {46
readonly currentVersion: number47
private readonly migrations: readonly SessionFormatMigration[]48
private readonly restoreCurrentHeader: SessionFormatChainOptions['restoreCurrentHeader']50
constructor(options: SessionFormatChainOptions) {51
this.currentVersion = sessionFormatVersion(options.currentVersion, 'current Session format version')52
this.restoreCurrentHeader = options.restoreCurrentHeader53
const byFrom = new Map<number, SessionFormatMigration>()54
const names = new Set<string>()55
for (const candidate of options.migrations) {56
const migration = defineSessionFormatMigration(candidate)57
if (byFrom.has(migration.fromVersion)) {58
throw new SessionFormatError(`Session migration v${migration.fromVersion}->v${migration.toVersion} is duplicated`)59
}60
if (names.has(migration.name)) throw new SessionFormatError(`Session migration name ${JSON.stringify(migration.name)} is duplicated`)61
byFrom.set(migration.fromVersion, migration)62
names.add(migration.name)63
}64
const ordered: SessionFormatMigration[] = []65
for (let version = 0; version < this.currentVersion; version += 1) {66
const migration = byFrom.get(version)67
if (migration === undefined) {68
throw new SessionFormatUnsupportedMigrationError(`Session migration v${version}->v${version + 1} is missing`)69
}70
ordered.push(migration)71
}72
if (byFrom.size !== ordered.length) {73
const invalid = [...byFrom.keys()].find(version => version >= this.currentVersion) as number74
throw new SessionFormatError(`Session migration from v${invalid} does not lead to current v${this.currentVersion}`)75
}76
this.migrations = Object.freeze(ordered)77
}79
private plan(fromVersion: number): readonly SessionFormatMigration[] {80
const from = sessionFormatVersion(fromVersion, 'stored Session format version')81
if (from > this.currentVersion) {82
throw new SessionFormatUnsupportedMigrationError(83
`stored Session uses newer format v${from}; this build writes v${this.currentVersion}`,84
)85
}86
return Object.freeze(this.migrations.slice(from))87
}89
createStream(90
sourceHeader: SessionFormatHeader,91
sourceCut: number | undefined,92
output: SessionFormatMigrationContext,93
): SessionFormatMigrationStream {94
let header = sourceHeader95
const validatedSourceCut = sourceCut === undefined96
? undefined97
: sessionFormatCount(sourceCut, 'Session inherited event count')98
let inheritedEventCount = validatedSourceCut99
const stages: Array<{100
readonly migration: SessionFormatMigration101
readonly stage: SessionFormatMigrationStage102
}> = []103
const plan = this.plan(header.version)104
for (const [index, migration] of plan.entries()) {105
const targetHeader = this.advanceHeader(migration, header)106
let stage: SessionFormatMigrationStage107
try {108
stage = migration.createStage({109
sourceHeader: header,110
targetHeader,111
sourceInheritedEventCount: inheritedEventCount,112
sourceKind: index === 0 ? 'decoded' : 'transformed',113
})114
} catch (error: unknown) {115
throwUnsupportedRefusal(migration, error)116
}117
header = targetHeader118
stages.push({ migration, stage })119
inheritedEventCount = stage.headerInheritedEventCount120
}121
return new CompiledSessionFormatMigrationStream(122
header,123
validatedSourceCut,124
stages,125
output,126
)127
}129
migrateHeader(source: SessionFormatHeader): SessionFormatHeader {130
let current = snapshotSessionFormatHeader(source, 'stored Session header')131
for (const migration of this.plan(current.version)) {132
current = this.advanceHeader(migration, current)133
}134
current = snapshotSessionFormatHeader(this.restoreCurrentHeader(current), 'current Session header restoration')135
if (current.version !== this.currentVersion) {136
throw new SessionFormatError(137
`current Session header restorer returned v${current.version}; expected v${this.currentVersion}`,138
)139
}140
return current141
}143
private advanceHeader(144
migration: SessionFormatMigration,145
source: SessionFormatHeader,146
): SessionFormatHeader {147
let target: SessionFormatHeader148
try {149
target = migration.migrateHeader(snapshotSessionFormatHeader(source, `${migration.name} header input`))150
} catch (error: unknown) {151
throwUnsupportedRefusal(migration, error, 'Session header')152
}153
const current = snapshotSessionFormatHeader(target, `${migration.name} header output`)154
if (current.version !== migration.toVersion) {155
throw new SessionFormatError(`${migration.name} header returned v${current.version}; expected v${migration.toVersion}`)156
}157
try {158
migration.validateTargetHeader(current)159
} catch (error: unknown) {160
throwUnsupportedRefusal(migration, error, 'Session header')161
}162
return current163
}164
}166
interface CompiledMigrationStage {167
readonly migration: SessionFormatMigration168
readonly stage: SessionFormatMigrationStage169
}171
class ChainedMigrationContext implements SessionFormatMigrationContext {172
constructor(173
private readonly entry: CompiledMigrationStage,174
private readonly output: SessionFormatMigrationContext,175
) {}177
emitEvent(event: SessionFormatEvent): void {178
try {179
this.entry.stage.transformEvent(event, this.output)180
} catch (error: unknown) {181
throwUnsupportedRefusal(this.entry.migration, error)182
}183
}185
emitRun(run: SessionFormatEventRun): void {186
try {187
this.entry.stage.transformRun(run, this.output)188
} catch (error: unknown) {189
throwUnsupportedRefusal(this.entry.migration, error)190
}191
}193
finish(): number {194
let targetCut: number195
try {196
targetCut = this.entry.stage.finish(this.output)197
} catch (error: unknown) {198
throwUnsupportedRefusal(this.entry.migration, error)199
}200
if (this.entry.stage.headerInheritedEventCount !== undefined201
&& this.entry.stage.headerInheritedEventCount !== targetCut) {202
throw new SessionFormatError(`${this.entry.migration.name} changed its predeclared inherited cut`)203
}204
return targetCut205
}206
}208
class CompiledSessionFormatMigrationStream implements SessionFormatMigrationStream {209
private readonly input: SessionFormatMigrationContext210
private readonly stages: readonly ChainedMigrationContext[]212
constructor(213
readonly header: SessionFormatHeader,214
private readonly sourceInheritedEventCount: number | undefined,215
entries: readonly CompiledMigrationStage[],216
output: SessionFormatMigrationContext,217
) {218
const stages = new Array<ChainedMigrationContext>(entries.length)219
let downstream = output220
for (const [offset, entry] of entries.toReversed().entries()) {221
const context = new ChainedMigrationContext(entry, downstream)222
stages[entries.length - offset - 1] = context223
downstream = context224
}225
this.input = downstream226
this.stages = stages227
}229
emitEvent(event: SessionFormatEvent): void {230
this.input.emitEvent(event)231
}233
emitRun(run: SessionFormatEventRun): void {234
this.input.emitRun(run)235
}237
finish(): number {238
let inheritedEventCount = this.sourceInheritedEventCount239
for (const stage of this.stages) inheritedEventCount = stage.finish()240
return sessionFormatCount(inheritedEventCount, 'finished Session inherited event count')241
}242
}244
function throwUnsupportedRefusal(245
migration: SessionFormatMigration,246
error: unknown,247
subject = 'Session',248
): never {249
if (error instanceof SessionFormatUnsupportedMigrationError) throw error250
const detail = error instanceof Error ? error.message : String(error)251
throw new SessionFormatUnsupportedMigrationError(252
`${migration.name} refuses this format v${migration.fromVersion} ${subject}: ${detail}`,253
{ cause: error },254
)255
}