catalog.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  1. import { createSessionFormatChain } from './chain.ts'
  2. import { SessionFormatEventCollector } from './context.ts'
  3. import { SessionFormatError, SessionFormatUnsupportedMigrationError } from './error.ts'
  4. import {
  5. inspectSessionFormatVersion,
  6. snapshotSessionFormatHeader,
  7. sessionFormatVersion,
  8. } from './json.ts'
  9. import type {
  10. SessionFormatArtifact,
  11. SessionFormatArtifactDecoder,
  12. SessionFormatCatalog,
  13. SessionFormatCatalogOptions,
  14. SessionFormatCodec,
  15. SessionFormatEvent,
  16. SessionFormatEventRun,
  17. SessionFormatHeaderReadResult,
  18. SessionFormatMigrationContext,
  19. SessionFormatMigrationStream,
  20. SessionFormatRestore,
  21. SessionFormatRestoreOptions,
  22. } from './types.ts'
  23. /**
  24. * Compile a build-static physical codec and adjacent migration catalog.
  25. * @param options - complete codecs, migrations, current version, and restorer.
  26. * @returns immutable physical dispatch and migration operations.
  27. */
  28. export function createSessionFormatCatalog(options: SessionFormatCatalogOptions): SessionFormatCatalog {
  29. const chain = createSessionFormatChain(options)
  30. const codecs = new Map<number, SessionFormatCodec>()
  31. for (const codec of options.codecs) {
  32. const version = sessionFormatVersion(codec.version, 'Session format codec version')
  33. if (codecs.has(version)) throw new SessionFormatError(`Session format codec v${version} is duplicated`)
  34. codecs.set(version, Object.freeze({ ...codec }))
  35. }
  36. for (let version = 0; version <= chain.currentVersion; version += 1) {
  37. if (!codecs.has(version)) throw new SessionFormatError(`Session format codec v${version} is missing`)
  38. }
  39. if (codecs.size !== chain.currentVersion + 1) {
  40. const invalid = [...codecs.keys()].find(version => version > chain.currentVersion) as number
  41. throw new SessionFormatError(`Session format codec v${invalid} is newer than current v${chain.currentVersion}`)
  42. }
  43. function readHeader(headerValue: unknown): SessionFormatHeaderReadResult {
  44. let storedVersion: number | undefined
  45. try {
  46. storedVersion = inspectSessionFormatVersion(headerValue)
  47. } catch (error: unknown) {
  48. return malformed(chain.currentVersion, error)
  49. }
  50. if (storedVersion > chain.currentVersion) {
  51. return Object.freeze({
  52. status: 'unsupported',
  53. storedVersion,
  54. targetVersion: chain.currentVersion,
  55. reason: `stored Session uses newer format v${storedVersion}; this build writes v${chain.currentVersion}`,
  56. })
  57. }
  58. const codec = codecs.get(storedVersion)
  59. /* v8 ignore next -- construction proves every supported version has exactly one codec. */
  60. if (codec === undefined) {
  61. return Object.freeze({
  62. status: 'unsupported',
  63. storedVersion,
  64. targetVersion: chain.currentVersion,
  65. reason: `this build has no Session format codec for v${storedVersion}`,
  66. })
  67. }
  68. try {
  69. const decoded = snapshotSessionFormatHeader(codec.decodeHeader(headerValue), `format v${storedVersion} header`)
  70. const header = chain.migrateHeader(decoded)
  71. return Object.freeze({
  72. status: storedVersion === chain.currentVersion ? 'current' : 'migration-required',
  73. storedVersion,
  74. targetVersion: chain.currentVersion,
  75. header,
  76. })
  77. } catch (error: unknown) {
  78. if (error instanceof SessionFormatUnsupportedMigrationError) {
  79. return Object.freeze({
  80. status: 'unsupported',
  81. storedVersion,
  82. targetVersion: chain.currentVersion,
  83. reason: error.message,
  84. })
  85. }
  86. return malformed(chain.currentVersion, error, storedVersion)
  87. }
  88. }
  89. function artifactCodec(headerValue: unknown): {
  90. readonly storedVersion: number
  91. readonly codec: SessionFormatCodec
  92. } {
  93. const storedVersion = inspectSessionFormatVersion(headerValue)
  94. if (storedVersion > chain.currentVersion) {
  95. throw new SessionFormatUnsupportedMigrationError(
  96. `stored Session uses newer format v${storedVersion}; this build writes v${chain.currentVersion}`,
  97. )
  98. }
  99. const codec = codecs.get(storedVersion)
  100. /* v8 ignore next -- construction proves every supported version has exactly one codec. */
  101. if (codec === undefined) {
  102. throw new SessionFormatUnsupportedMigrationError(`this build has no Session format codec for v${storedVersion}`)
  103. }
  104. return { storedVersion, codec }
  105. }
  106. function encodeCurrentHeader(
  107. header: Parameters<SessionFormatCatalog['encodeCurrentHeader']>[0],
  108. inheritedEventCount: number,
  109. ) {
  110. if (inspectSessionFormatVersion(header) !== chain.currentVersion) {
  111. throw new SessionFormatError(`encodeCurrent requires Session format v${chain.currentVersion}`)
  112. }
  113. const encoded = options.currentEncoder.encodeHeader(header, inheritedEventCount)
  114. if (inspectSessionFormatVersion(encoded) !== chain.currentVersion) {
  115. throw new SessionFormatError('current Session codec returned a non-current header')
  116. }
  117. return encoded
  118. }
  119. function createRestore(
  120. headerValue: unknown,
  121. restoreOptions: SessionFormatRestoreOptions,
  122. ): SessionFormatRestore {
  123. const { storedVersion, codec } = artifactCodec(headerValue)
  124. const decoder = codec.createDecoder(headerValue, restoreOptions.recovery)
  125. const sourceCut = decoder.headerInheritedEventCount
  126. if (storedVersion === chain.currentVersion) {
  127. return new CurrentSessionFormatRestore(
  128. decoder,
  129. sourceCut,
  130. restoreOptions.validation === 'current' ? options.restoreCurrent : identityArtifact,
  131. chain.currentVersion,
  132. )
  133. }
  134. const collector = new SessionFormatEventCollector()
  135. const migration = chain.createStream(
  136. decoder.header,
  137. sourceCut,
  138. collector,
  139. )
  140. return new MigratingSessionFormatRestore(
  141. decoder,
  142. sourceCut,
  143. migration,
  144. collector,
  145. restoreOptions.validation === 'current'
  146. ? options.restoreCurrent
  147. : options.restoreTransformedCurrent,
  148. restoreOptions.validation,
  149. storedVersion,
  150. chain.currentVersion,
  151. )
  152. }
  153. return Object.freeze({
  154. currentVersion: chain.currentVersion,
  155. readHeader,
  156. createRestore,
  157. encodeCurrentHeader,
  158. encodeCurrentEvent: options.currentEncoder.encodeEvent.bind(options.currentEncoder),
  159. })
  160. }
  161. type SessionFormatArtifactRestorer = (artifact: SessionFormatArtifact) => SessionFormatArtifact
  162. class CurrentSessionFormatRestore implements SessionFormatRestore {
  163. readonly header: SessionFormatArtifact['header']
  164. private readonly collector = new SessionFormatEventCollector()
  165. constructor(
  166. private readonly decoder: SessionFormatArtifactDecoder,
  167. private readonly sourceInheritedEventCount: number | undefined,
  168. private readonly restoreArtifact: SessionFormatArtifactRestorer,
  169. private readonly currentVersion: number,
  170. ) {
  171. this.header = decoder.header
  172. }
  173. decodeRow(rowValue: unknown): void {
  174. this.decoder.decodeRow(rowValue, this.collector)
  175. }
  176. finish(): SessionFormatArtifact {
  177. const inheritedEventCount = finishDecoder(
  178. this.decoder,
  179. this.collector,
  180. this.sourceInheritedEventCount,
  181. )
  182. return restoreCurrentVersion(this.restoreArtifact({
  183. header: this.header,
  184. inheritedEventCount,
  185. events: this.collector.values,
  186. }), this.currentVersion)
  187. }
  188. }
  189. class MigratingSessionFormatRestore implements
  190. SessionFormatRestore,
  191. SessionFormatMigrationContext {
  192. readonly header: SessionFormatArtifact['header']
  193. constructor(
  194. private readonly decoder: SessionFormatArtifactDecoder,
  195. private readonly sourceInheritedEventCount: number | undefined,
  196. private readonly migration: SessionFormatMigrationStream,
  197. private readonly collector: SessionFormatEventCollector,
  198. private readonly restoreArtifact: SessionFormatArtifactRestorer,
  199. private readonly validation: SessionFormatRestoreOptions['validation'],
  200. private readonly sourceVersion: number,
  201. private readonly currentVersion: number,
  202. ) {
  203. this.header = migration.header
  204. }
  205. decodeRow(rowValue: unknown): void {
  206. this.decoder.decodeRow(rowValue, this)
  207. }
  208. emitEvent(event: SessionFormatEvent): void {
  209. this.migration.emitEvent(event)
  210. }
  211. emitRun(run: SessionFormatEventRun): void {
  212. this.migration.emitRun(run)
  213. }
  214. finish(): SessionFormatArtifact {
  215. finishDecoder(this.decoder, this, this.sourceInheritedEventCount)
  216. const artifact = {
  217. header: this.header,
  218. inheritedEventCount: this.migration.finish(),
  219. events: this.collector.values,
  220. }
  221. let restored: SessionFormatArtifact
  222. try {
  223. restored = this.restoreArtifact(artifact)
  224. } catch (error: unknown) {
  225. if (this.validation === 'current'
  226. || error instanceof SessionFormatUnsupportedMigrationError) throw error
  227. const detail = error instanceof Error ? error.message : String(error)
  228. throw new SessionFormatUnsupportedMigrationError(
  229. `Session migration from v${this.sourceVersion} to v${this.currentVersion} refuses the transformed artifact: ${detail}`,
  230. { cause: error },
  231. )
  232. }
  233. return restoreCurrentVersion(restored, this.currentVersion)
  234. }
  235. }
  236. function finishDecoder(
  237. decoder: SessionFormatArtifactDecoder,
  238. context: SessionFormatMigrationContext,
  239. sourceInheritedEventCount: number | undefined,
  240. ): number {
  241. const inheritedEventCount = decoder.finish(context)
  242. if (sourceInheritedEventCount !== undefined && inheritedEventCount !== sourceInheritedEventCount) {
  243. throw new SessionFormatError('streaming decoder changed its predeclared inherited cut')
  244. }
  245. return inheritedEventCount
  246. }
  247. function restoreCurrentVersion(
  248. artifact: SessionFormatArtifact,
  249. currentVersion: number,
  250. ): SessionFormatArtifact {
  251. if (artifact.header.version !== currentVersion) {
  252. throw new SessionFormatError(
  253. `current Session restorer returned v${artifact.header.version}; expected v${currentVersion}`,
  254. )
  255. }
  256. return artifact
  257. }
  258. function identityArtifact(artifact: SessionFormatArtifact): SessionFormatArtifact {
  259. return artifact
  260. }
  261. function malformed(targetVersion: number, error: unknown, storedVersion?: number): SessionFormatHeaderReadResult {
  262. return Object.freeze({
  263. status: 'malformed',
  264. ...(storedVersion === undefined ? {} : { storedVersion }),
  265. targetVersion,
  266. reason: error instanceof Error ? error.message : String(error),
  267. })
  268. }