| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289 |
- import { createSessionFormatChain } from './chain.ts'
- import { SessionFormatEventCollector } from './context.ts'
- import { SessionFormatError, SessionFormatUnsupportedMigrationError } from './error.ts'
- import {
- inspectSessionFormatVersion,
- snapshotSessionFormatHeader,
- sessionFormatVersion,
- } from './json.ts'
- import type {
- SessionFormatArtifact,
- SessionFormatArtifactDecoder,
- SessionFormatCatalog,
- SessionFormatCatalogOptions,
- SessionFormatCodec,
- SessionFormatEvent,
- SessionFormatEventRun,
- SessionFormatHeaderReadResult,
- SessionFormatMigrationContext,
- SessionFormatMigrationStream,
- SessionFormatRestore,
- SessionFormatRestoreOptions,
- } from './types.ts'
- /**
- * Compile a build-static physical codec and adjacent migration catalog.
- * @param options - complete codecs, migrations, current version, and restorer.
- * @returns immutable physical dispatch and migration operations.
- */
- export function createSessionFormatCatalog(options: SessionFormatCatalogOptions): SessionFormatCatalog {
- const chain = createSessionFormatChain(options)
- const codecs = new Map<number, SessionFormatCodec>()
- for (const codec of options.codecs) {
- const version = sessionFormatVersion(codec.version, 'Session format codec version')
- if (codecs.has(version)) throw new SessionFormatError(`Session format codec v${version} is duplicated`)
- codecs.set(version, Object.freeze({ ...codec }))
- }
- for (let version = 0; version <= chain.currentVersion; version += 1) {
- if (!codecs.has(version)) throw new SessionFormatError(`Session format codec v${version} is missing`)
- }
- if (codecs.size !== chain.currentVersion + 1) {
- const invalid = [...codecs.keys()].find(version => version > chain.currentVersion) as number
- throw new SessionFormatError(`Session format codec v${invalid} is newer than current v${chain.currentVersion}`)
- }
- function readHeader(headerValue: unknown): SessionFormatHeaderReadResult {
- let storedVersion: number | undefined
- try {
- storedVersion = inspectSessionFormatVersion(headerValue)
- } catch (error: unknown) {
- return malformed(chain.currentVersion, error)
- }
- if (storedVersion > chain.currentVersion) {
- return Object.freeze({
- status: 'unsupported',
- storedVersion,
- targetVersion: chain.currentVersion,
- reason: `stored Session uses newer format v${storedVersion}; this build writes v${chain.currentVersion}`,
- })
- }
- const codec = codecs.get(storedVersion)
- /* v8 ignore next -- construction proves every supported version has exactly one codec. */
- if (codec === undefined) {
- return Object.freeze({
- status: 'unsupported',
- storedVersion,
- targetVersion: chain.currentVersion,
- reason: `this build has no Session format codec for v${storedVersion}`,
- })
- }
- try {
- const decoded = snapshotSessionFormatHeader(codec.decodeHeader(headerValue), `format v${storedVersion} header`)
- const header = chain.migrateHeader(decoded)
- return Object.freeze({
- status: storedVersion === chain.currentVersion ? 'current' : 'migration-required',
- storedVersion,
- targetVersion: chain.currentVersion,
- header,
- })
- } catch (error: unknown) {
- if (error instanceof SessionFormatUnsupportedMigrationError) {
- return Object.freeze({
- status: 'unsupported',
- storedVersion,
- targetVersion: chain.currentVersion,
- reason: error.message,
- })
- }
- return malformed(chain.currentVersion, error, storedVersion)
- }
- }
- function artifactCodec(headerValue: unknown): {
- readonly storedVersion: number
- readonly codec: SessionFormatCodec
- } {
- const storedVersion = inspectSessionFormatVersion(headerValue)
- if (storedVersion > chain.currentVersion) {
- throw new SessionFormatUnsupportedMigrationError(
- `stored Session uses newer format v${storedVersion}; this build writes v${chain.currentVersion}`,
- )
- }
- const codec = codecs.get(storedVersion)
- /* v8 ignore next -- construction proves every supported version has exactly one codec. */
- if (codec === undefined) {
- throw new SessionFormatUnsupportedMigrationError(`this build has no Session format codec for v${storedVersion}`)
- }
- return { storedVersion, codec }
- }
- function encodeCurrentHeader(
- header: Parameters<SessionFormatCatalog['encodeCurrentHeader']>[0],
- inheritedEventCount: number,
- ) {
- if (inspectSessionFormatVersion(header) !== chain.currentVersion) {
- throw new SessionFormatError(`encodeCurrent requires Session format v${chain.currentVersion}`)
- }
- const encoded = options.currentEncoder.encodeHeader(header, inheritedEventCount)
- if (inspectSessionFormatVersion(encoded) !== chain.currentVersion) {
- throw new SessionFormatError('current Session codec returned a non-current header')
- }
- return encoded
- }
- function createRestore(
- headerValue: unknown,
- restoreOptions: SessionFormatRestoreOptions,
- ): SessionFormatRestore {
- const { storedVersion, codec } = artifactCodec(headerValue)
- const decoder = codec.createDecoder(headerValue, restoreOptions.recovery)
- const sourceCut = decoder.headerInheritedEventCount
- if (storedVersion === chain.currentVersion) {
- return new CurrentSessionFormatRestore(
- decoder,
- sourceCut,
- restoreOptions.validation === 'current' ? options.restoreCurrent : identityArtifact,
- chain.currentVersion,
- )
- }
- const collector = new SessionFormatEventCollector()
- const migration = chain.createStream(
- decoder.header,
- sourceCut,
- collector,
- )
- return new MigratingSessionFormatRestore(
- decoder,
- sourceCut,
- migration,
- collector,
- restoreOptions.validation === 'current'
- ? options.restoreCurrent
- : options.restoreTransformedCurrent,
- restoreOptions.validation,
- storedVersion,
- chain.currentVersion,
- )
- }
- return Object.freeze({
- currentVersion: chain.currentVersion,
- readHeader,
- createRestore,
- encodeCurrentHeader,
- encodeCurrentEvent: options.currentEncoder.encodeEvent.bind(options.currentEncoder),
- })
- }
- type SessionFormatArtifactRestorer = (artifact: SessionFormatArtifact) => SessionFormatArtifact
- class CurrentSessionFormatRestore implements SessionFormatRestore {
- readonly header: SessionFormatArtifact['header']
- private readonly collector = new SessionFormatEventCollector()
- constructor(
- private readonly decoder: SessionFormatArtifactDecoder,
- private readonly sourceInheritedEventCount: number | undefined,
- private readonly restoreArtifact: SessionFormatArtifactRestorer,
- private readonly currentVersion: number,
- ) {
- this.header = decoder.header
- }
- decodeRow(rowValue: unknown): void {
- this.decoder.decodeRow(rowValue, this.collector)
- }
- finish(): SessionFormatArtifact {
- const inheritedEventCount = finishDecoder(
- this.decoder,
- this.collector,
- this.sourceInheritedEventCount,
- )
- return restoreCurrentVersion(this.restoreArtifact({
- header: this.header,
- inheritedEventCount,
- events: this.collector.values,
- }), this.currentVersion)
- }
- }
- class MigratingSessionFormatRestore implements
- SessionFormatRestore,
- SessionFormatMigrationContext {
- readonly header: SessionFormatArtifact['header']
- constructor(
- private readonly decoder: SessionFormatArtifactDecoder,
- private readonly sourceInheritedEventCount: number | undefined,
- private readonly migration: SessionFormatMigrationStream,
- private readonly collector: SessionFormatEventCollector,
- private readonly restoreArtifact: SessionFormatArtifactRestorer,
- private readonly validation: SessionFormatRestoreOptions['validation'],
- private readonly sourceVersion: number,
- private readonly currentVersion: number,
- ) {
- this.header = migration.header
- }
- decodeRow(rowValue: unknown): void {
- this.decoder.decodeRow(rowValue, this)
- }
- emitEvent(event: SessionFormatEvent): void {
- this.migration.emitEvent(event)
- }
- emitRun(run: SessionFormatEventRun): void {
- this.migration.emitRun(run)
- }
- finish(): SessionFormatArtifact {
- finishDecoder(this.decoder, this, this.sourceInheritedEventCount)
- const artifact = {
- header: this.header,
- inheritedEventCount: this.migration.finish(),
- events: this.collector.values,
- }
- let restored: SessionFormatArtifact
- try {
- restored = this.restoreArtifact(artifact)
- } catch (error: unknown) {
- if (this.validation === 'current'
- || error instanceof SessionFormatUnsupportedMigrationError) throw error
- const detail = error instanceof Error ? error.message : String(error)
- throw new SessionFormatUnsupportedMigrationError(
- `Session migration from v${this.sourceVersion} to v${this.currentVersion} refuses the transformed artifact: ${detail}`,
- { cause: error },
- )
- }
- return restoreCurrentVersion(restored, this.currentVersion)
- }
- }
- function finishDecoder(
- decoder: SessionFormatArtifactDecoder,
- context: SessionFormatMigrationContext,
- sourceInheritedEventCount: number | undefined,
- ): number {
- const inheritedEventCount = decoder.finish(context)
- if (sourceInheritedEventCount !== undefined && inheritedEventCount !== sourceInheritedEventCount) {
- throw new SessionFormatError('streaming decoder changed its predeclared inherited cut')
- }
- return inheritedEventCount
- }
- function restoreCurrentVersion(
- artifact: SessionFormatArtifact,
- currentVersion: number,
- ): SessionFormatArtifact {
- if (artifact.header.version !== currentVersion) {
- throw new SessionFormatError(
- `current Session restorer returned v${artifact.header.version}; expected v${currentVersion}`,
- )
- }
- return artifact
- }
- function identityArtifact(artifact: SessionFormatArtifact): SessionFormatArtifact {
- return artifact
- }
- function malformed(targetVersion: number, error: unknown, storedVersion?: number): SessionFormatHeaderReadResult {
- return Object.freeze({
- status: 'malformed',
- ...(storedVersion === undefined ? {} : { storedVersion }),
- targetVersion,
- reason: error instanceof Error ? error.message : String(error),
- })
- }
|