import * as plugins from '../ts/plugins.js'; import { assertControllerSessionStateDocument, assertControllerSettingsDocument, ControllerSessionStateModel, ControllerSettingsModel, } from '../ts/classes.authmodels.js'; import { AuthError } from '../ts/interfaces.auth.js'; import type { IControllerSessionStateDocument, IControllerSettingsDocument, } from '../ts/interfaces.projects.js'; type TObject = Record; export interface IProtocolV27DocumentMigrationResult { document: TObject; migrated: boolean; } const migrationPageLimit = 128; const isPlainObject = (valueArg: unknown): valueArg is TObject => ( typeof valueArg === 'object' && valueArg !== null && !Array.isArray(valueArg) && ( Object.getPrototypeOf(valueArg) === Object.prototype || Object.getPrototypeOf(valueArg) === null ) ); const persistedBodyFromRawDocument = (valueArg: TObject): TObject => { const body = { ...valueArg }; delete body._id; delete body._smartdataRevision; return body; }; const withoutProviderConnectionId = (modelArg: TObject): TObject => { const model = { ...modelArg }; delete model.providerConnectionId; return model; }; export const migrateProtocolV27SettingsDocument = ( valueArg: unknown, ): IProtocolV27DocumentMigrationResult => { if (!isPlainObject(valueArg)) { throw new Error('Invalid controller settings model migration document.'); } if (!Array.isArray(valueArg.defaultModels)) { return { document: valueArg, migrated: false }; } let migrated = false; const defaultModels = valueArg.defaultModels.map((modelArg) => { if (!isPlainObject(modelArg) || !Object.hasOwn(modelArg, 'providerConnectionId')) { return modelArg; } migrated = true; return withoutProviderConnectionId(modelArg); }); return { document: migrated ? { ...valueArg, defaultModels } : valueArg, migrated, }; }; export const migrateProtocolV27SessionStateDocument = ( valueArg: unknown, ): IProtocolV27DocumentMigrationResult => { if (!isPlainObject(valueArg)) { throw new Error('Invalid controller session model migration document.'); } if ( !isPlainObject(valueArg.modelChoice) || !Object.hasOwn(valueArg.modelChoice, 'providerConnectionId') ) { return { document: valueArg, migrated: false }; } const nestedProviderConnectionId = valueArg.modelChoice.providerConnectionId; if ( typeof nestedProviderConnectionId !== 'string' || nestedProviderConnectionId.length === 0 || Buffer.byteLength(nestedProviderConnectionId, 'utf8') > 512 ) { throw new Error('Invalid persisted Flex provider connection selection.'); } if ( valueArg.modelChoice.harnessId === 'flex' && valueArg.providerConnectionId !== undefined && valueArg.providerConnectionId !== nestedProviderConnectionId ) { throw new Error('Conflicting persisted Flex provider connection selections.'); } const modelChoice = withoutProviderConnectionId(valueArg.modelChoice); return { document: { ...valueArg, modelChoice, ...(valueArg.modelChoice.harnessId === 'flex' ? { providerConnectionId: nestedProviderConnectionId } : {}), }, migrated: true, }; }; export class ProtocolV27AccountIndependentModelMigration { constructor(private readonly database: plugins.smartdata.SmartdataDb | undefined) {} /** Runs before exact-persistence reads and is safe to repeat after partial completion. */ public async run(): Promise { if (!this.database) { throw new AuthError('not_initialized', 'The authentication store database is unavailable.'); } const settingsCollection = ControllerSettingsModel.collection.mongoDbCollection; let lastSettingsId: plugins.smartdata.TStoredDocument['_id'] | undefined; while (true) { const cursor = settingsCollection.find( lastSettingsId ? { _id: { $gt: lastSettingsId } } : {}, ).sort({ _id: 1 }).limit(migrationPageLimit); let settingsDocuments: Awaited>; try { settingsDocuments = await cursor.toArray(); } finally { await cursor.close(); } for (const rawDocument of settingsDocuments) { const migration = migrateProtocolV27SettingsDocument( persistedBodyFromRawDocument(rawDocument as unknown as TObject), ); if (!migration.migrated) continue; assertControllerSettingsDocument(migration.document); const selector = rawDocument._smartdataRevision === undefined ? { _id: rawDocument._id, _smartdataRevision: { $exists: false } } : { _id: rawDocument._id, _smartdataRevision: rawDocument._smartdataRevision }; const replaced = await settingsCollection.findOneAndReplace( selector, { _id: rawDocument._id, ...migration.document, _smartdataRevision: plugins.crypto.randomUUID(), }, { returnDocument: 'after', includeResultMetadata: false, upsert: false }, ); if (!replaced) { const concurrent = await settingsCollection.findOne({ _id: rawDocument._id }); if (!concurrent) { throw new AuthError( 'concurrent_change', 'The settings model migration changed concurrently.', ); } const remaining = migrateProtocolV27SettingsDocument( persistedBodyFromRawDocument(concurrent as unknown as TObject), ); if (remaining.migrated) { throw new AuthError( 'concurrent_change', 'The settings model migration changed concurrently.', ); } assertControllerSettingsDocument(remaining.document); } } const lastDocument = settingsDocuments.at(-1); if (!lastDocument || settingsDocuments.length < migrationPageLimit) break; lastSettingsId = lastDocument._id; } const sessionStateCollection = ControllerSessionStateModel.collection.mongoDbCollection; let lastSessionStateId: | plugins.smartdata.TStoredDocument['_id'] | undefined; while (true) { const cursor = sessionStateCollection.find( lastSessionStateId ? { _id: { $gt: lastSessionStateId } } : {}, ).sort({ _id: 1 }).limit(migrationPageLimit); let sessionStateDocuments: Awaited>; try { sessionStateDocuments = await cursor.toArray(); } finally { await cursor.close(); } for (const rawDocument of sessionStateDocuments) { const migration = migrateProtocolV27SessionStateDocument( persistedBodyFromRawDocument(rawDocument as unknown as TObject), ); if (!migration.migrated) continue; assertControllerSessionStateDocument(migration.document); const selector = rawDocument._smartdataRevision === undefined ? { _id: rawDocument._id, _smartdataRevision: { $exists: false } } : { _id: rawDocument._id, _smartdataRevision: rawDocument._smartdataRevision }; const replaced = await sessionStateCollection.findOneAndReplace( selector, { _id: rawDocument._id, ...migration.document, _smartdataRevision: plugins.crypto.randomUUID(), }, { returnDocument: 'after', includeResultMetadata: false, upsert: false }, ); if (!replaced) { const concurrent = await sessionStateCollection.findOne({ _id: rawDocument._id }); if (!concurrent) { throw new AuthError( 'concurrent_change', 'The session model migration changed concurrently.', ); } const remaining = migrateProtocolV27SessionStateDocument( persistedBodyFromRawDocument(concurrent as unknown as TObject), ); if (remaining.migrated) { throw new AuthError( 'concurrent_change', 'The session model migration changed concurrently.', ); } assertControllerSessionStateDocument(remaining.document); } } const lastDocument = sessionStateDocuments.at(-1); if (!lastDocument || sessionStateDocuments.length < migrationPageLimit) break; lastSessionStateId = lastDocument._id; } } }