返回源码地图

packages/session/session-format/src/chain.ts

main snapshot · da00f7f5358f · 正文引用章节 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1import { SessionFormatError, SessionFormatUnsupportedMigrationError } from './error.ts'
2import {
3 snapshotSessionFormatHeader,
4 sessionFormatCount,
5 sessionFormatVersion,
6} from './json.ts'
7import type {
8 SessionFormatChain,
9 SessionFormatChainOptions,
10 SessionFormatEvent,
11 SessionFormatEventRun,
12 SessionFormatHeader,
13 SessionFormatMigration,
14 SessionFormatMigrationContext,
15 SessionFormatMigrationStage,
16 SessionFormatMigrationStream,
17} from './types.ts'
18
19/**
20 * Validate and freeze one adjacent migration declaration.
21 * @param migration - named exact adjacent conversion.
22 * @returns immutable validated declaration.
23 */
24export 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}
35
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 */
41export function createSessionFormatChain(options: SessionFormatChainOptions): SessionFormatChain {
42 return new CompiledSessionFormatChain(options)
43}
44
45class CompiledSessionFormatChain implements SessionFormatChain {
46 readonly currentVersion: number
47 private readonly migrations: readonly SessionFormatMigration[]
48 private readonly restoreCurrentHeader: SessionFormatChainOptions['restoreCurrentHeader']
49
50 constructor(options: SessionFormatChainOptions) {
51 this.currentVersion = sessionFormatVersion(options.currentVersion, 'current Session format version')
52 this.restoreCurrentHeader = options.restoreCurrentHeader
53 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 number
74 throw new SessionFormatError(`Session migration from v${invalid} does not lead to current v${this.currentVersion}`)
75 }
76 this.migrations = Object.freeze(ordered)
77 }
78
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 }
88
89 createStream(
90 sourceHeader: SessionFormatHeader,
91 sourceCut: number | undefined,
92 output: SessionFormatMigrationContext,
93 ): SessionFormatMigrationStream {
94 let header = sourceHeader
95 const validatedSourceCut = sourceCut === undefined
96 ? undefined
97 : sessionFormatCount(sourceCut, 'Session inherited event count')
98 let inheritedEventCount = validatedSourceCut
99 const stages: Array<{
100 readonly migration: SessionFormatMigration
101 readonly stage: SessionFormatMigrationStage
102 }> = []
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: SessionFormatMigrationStage
107 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 = targetHeader
118 stages.push({ migration, stage })
119 inheritedEventCount = stage.headerInheritedEventCount
120 }
121 return new CompiledSessionFormatMigrationStream(
122 header,
123 validatedSourceCut,
124 stages,
125 output,
126 )
127 }
128
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 current
141 }
142
143 private advanceHeader(
144 migration: SessionFormatMigration,
145 source: SessionFormatHeader,
146 ): SessionFormatHeader {
147 let target: SessionFormatHeader
148 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 current
163 }
164}
165
166interface CompiledMigrationStage {
167 readonly migration: SessionFormatMigration
168 readonly stage: SessionFormatMigrationStage
169}
170
171class ChainedMigrationContext implements SessionFormatMigrationContext {
172 constructor(
173 private readonly entry: CompiledMigrationStage,
174 private readonly output: SessionFormatMigrationContext,
175 ) {}
176
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 }
184
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 }
192
193 finish(): number {
194 let targetCut: number
195 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 !== undefined
201 && this.entry.stage.headerInheritedEventCount !== targetCut) {
202 throw new SessionFormatError(`${this.entry.migration.name} changed its predeclared inherited cut`)
203 }
204 return targetCut
205 }
206}
207
208class CompiledSessionFormatMigrationStream implements SessionFormatMigrationStream {
209 private readonly input: SessionFormatMigrationContext
210 private readonly stages: readonly ChainedMigrationContext[]
211
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 = output
220 for (const [offset, entry] of entries.toReversed().entries()) {
221 const context = new ChainedMigrationContext(entry, downstream)
222 stages[entries.length - offset - 1] = context
223 downstream = context
224 }
225 this.input = downstream
226 this.stages = stages
227 }
228
229 emitEvent(event: SessionFormatEvent): void {
230 this.input.emitEvent(event)
231 }
232
233 emitRun(run: SessionFormatEventRun): void {
234 this.input.emitRun(run)
235 }
236
237 finish(): number {
238 let inheritedEventCount = this.sourceInheritedEventCount
239 for (const stage of this.stages) inheritedEventCount = stage.finish()
240 return sessionFormatCount(inheritedEventCount, 'finished Session inherited event count')
241 }
242}
243
244function throwUnsupportedRefusal(
245 migration: SessionFormatMigration,
246 error: unknown,
247 subject = 'Session',
248): never {
249 if (error instanceof SessionFormatUnsupportedMigrationError) throw error
250 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}