import * as plugins from '../ts/plugins.js'; import { ControllerSessionStateModel, controllerSessionStateId, } from '../ts/classes.authmodels.js'; import { ControllerManagedSessionModel, assertControllerFlexSessionGeneration, type IControllerManagedSessionDocument, } from '../ts/classes.managedsessionmodels.js'; import { ControllerSessionIdentityError, ControllerSessionIdentityService, controllerSessionIdentitySnapshotEntryLimit, type IControllerSessionSnapshotEntry, } from '../ts/classes.sessionidentityservice.js'; import { controllerRuntimeIdKey, type TControllerSessionId, } from '../ts_interfaces/index.js'; const base64UrlPattern = /^[A-Za-z0-9_-]+$/; const controllerIdPattern = /^controller:[1-9][0-9]{0,4}$/; const maximumNativeIdBytes = 512; const maximumProviderGenerationBytes = 512; const statePageEntryLimit = 100; export const controllerLegacyManagedSessionMigrationCohortLimit = 500; export type TControllerLegacyManagedSessionMigrationState = 'admitting' | 'completed'; export type TControllerLegacyManagedSessionMigrationOutcome = | 'pending' | 'admitted' | 'excluded'; export interface IControllerLegacyManagedSessionMigrationManifestEntry { stateDocumentId: string; runtimeId: TControllerSessionId; providerSessionGeneration: string; sessionGenerationId?: string; sessionGenerationSequence?: number; outcome: TControllerLegacyManagedSessionMigrationOutcome; } export interface IControllerLegacyManagedSessionMigrationDocument { id: string; controllerId: string; issuerId: string; projectId: string; harnessId: TControllerSessionId['harnessId']; state: TControllerLegacyManagedSessionMigrationState; cohortObservedAt: Date; manifest: IControllerLegacyManagedSessionMigrationManifestEntry[]; updateId: string; completedAt?: Date; } export interface IControllerLegacyManagedSessionMigrationInput { controllerId: string; issuerId: string; projectId: string; harnessId: TControllerSessionId['harnessId']; supervisorGeneration: string; sourceAdmissionFenced: true; snapshotComplete: true; sessions: readonly IControllerSessionSnapshotEntry[]; creationObligationsComplete: true; creationObligationRuntimeIds: readonly TControllerSessionId[]; observedAt: Date; signal: AbortSignal; assertRuntimeAuthorityBeforeCompletion: () => void | Promise; assertRuntimeAuthorityAfterCompletion: () => void | Promise; } export interface IControllerLegacyManagedSessionMigrationResult { status: 'completed'; markerId: string; cohortSize: number; admittedCount: number; excludedCount: number; } export interface IControllerLegacyManagedSessionV26MigrationOptions { issuerId: string; sessionIdentityService: ControllerSessionIdentityService; } interface IResolvedMigrationScope { markerId: string; controllerId: string; issuerId: string; projectId: string; harnessId: TControllerSessionId['harnessId']; } interface IValidatedMigrationInput { snapshotByRuntimeKey: Map; creationObligationRuntimeKeys: Set; observedAt: Date; } type TStoredMarker = NonNullable>>; type TStoredManagedSession = NonNullable>>; const isPlainObject = (valueArg: unknown): valueArg is Record => ( typeof valueArg === 'object' && valueArg !== null && !Array.isArray(valueArg) && ( Object.getPrototypeOf(valueArg) === Object.prototype || Object.getPrototypeOf(valueArg) === null ) ); const hasExactKeys = ( valueArg: Record, requiredKeysArg: readonly string[], optionalKeysArg: readonly string[] = [], ): boolean => { const keys = Object.keys(valueArg); return requiredKeysArg.every((key) => Object.hasOwn(valueArg, key)) && keys.every((key) => requiredKeysArg.includes(key) || optionalKeysArg.includes(key)); }; const isBase64UrlBytes = (valueArg: unknown, byteLengthArg: number): valueArg is string => ( typeof valueArg === 'string' && base64UrlPattern.test(valueArg) && Buffer.from(valueArg, 'base64url').byteLength === byteLengthArg && Buffer.from(valueArg, 'base64url').toString('base64url') === valueArg ); const isValidDate = (valueArg: unknown): valueArg is Date => ( valueArg instanceof Date && Number.isFinite(valueArg.getTime()) && valueArg.getTime() >= 0 ); const isProviderGeneration = (valueArg: unknown): valueArg is string => ( typeof valueArg === 'string' && valueArg.length > 0 && Buffer.byteLength(valueArg, 'utf8') <= maximumProviderGenerationBytes && !/[\u0000-\u001f\u007f]/u.test(valueArg) ); const isRuntimeId = ( valueArg: unknown, harnessIdArg?: TControllerSessionId['harnessId'], ): valueArg is TControllerSessionId => ( isPlainObject(valueArg) && hasExactKeys(valueArg, ['harnessId', 'nativeId']) && (valueArg.harnessId === 'opencode' || valueArg.harnessId === 'flex') && (harnessIdArg === undefined || valueArg.harnessId === harnessIdArg) && typeof valueArg.nativeId === 'string' && valueArg.nativeId.length > 0 && Buffer.byteLength(valueArg.nativeId, 'utf8') <= maximumNativeIdBytes && !/[\u0000-\u001f\u007f]/u.test(valueArg.nativeId) ); const randomId = (bytesArg: number): string => plugins.crypto .randomBytes(bytesArg) .toString('base64url'); const digest = (valueArg: unknown): string => plugins.crypto .createHash('sha256') .update(JSON.stringify(valueArg), 'utf8') .digest('base64url'); const migrationError = ( codeArg: ConstructorParameters[0], messageArg: string, causeArg?: unknown, ): ControllerSessionIdentityError => new ControllerSessionIdentityError( codeArg, messageArg, causeArg === undefined ? undefined : { cause: causeArg }, ); const isAmbiguousWrite = (errorArg: unknown): boolean => ( errorArg instanceof plugins.smartdata.SmartdataExactPersistenceError && errorArg.code === 'ambiguous_write' ); const runtimeIdsEqual = ( leftArg: TControllerSessionId, rightArg: TControllerSessionId, ): boolean => ( leftArg.harnessId === rightArg.harnessId && leftArg.nativeId === rightArg.nativeId ); const datesEqual = (leftArg: Date, rightArg: Date): boolean => ( leftArg.getTime() === rightArg.getTime() ); const generationFactsEqual = ( leftArg: Pick< IControllerLegacyManagedSessionMigrationManifestEntry, 'providerSessionGeneration' | 'sessionGenerationId' | 'sessionGenerationSequence' >, rightArg: Pick< IControllerSessionSnapshotEntry, 'providerSessionGeneration' | 'sessionGenerationId' | 'sessionGenerationSequence' >, ): boolean => ( leftArg.providerSessionGeneration === rightArg.providerSessionGeneration && leftArg.sessionGenerationId === rightArg.sessionGenerationId && leftArg.sessionGenerationSequence === rightArg.sessionGenerationSequence ); const immutableManifestEntriesEqual = ( leftArg: IControllerLegacyManagedSessionMigrationManifestEntry, rightArg: IControllerLegacyManagedSessionMigrationManifestEntry, ): boolean => ( leftArg.stateDocumentId === rightArg.stateDocumentId && runtimeIdsEqual(leftArg.runtimeId, rightArg.runtimeId) && generationFactsEqual(leftArg, rightArg) ); const immutableMarkerCohortEqual = ( leftArg: IControllerLegacyManagedSessionMigrationDocument, rightArg: IControllerLegacyManagedSessionMigrationDocument, ): boolean => ( leftArg.id === rightArg.id && leftArg.controllerId === rightArg.controllerId && leftArg.issuerId === rightArg.issuerId && leftArg.projectId === rightArg.projectId && leftArg.harnessId === rightArg.harnessId && datesEqual(leftArg.cohortObservedAt, rightArg.cohortObservedAt) && leftArg.manifest.length === rightArg.manifest.length && leftArg.manifest.every((entry, index) => immutableManifestEntriesEqual( entry, rightArg.manifest[index], )) ); export const controllerLegacyManagedSessionMigrationDocumentId = ( controllerIdArg: string, issuerIdArg: string, projectIdArg: string, harnessIdArg: TControllerSessionId['harnessId'], ): string => digest([ 'controller-legacy-managed-session-migration-v1', controllerIdArg, issuerIdArg, projectIdArg, harnessIdArg, ]); export const assertControllerLegacyManagedSessionMigrationDocument = ( valueArg: unknown, ): asserts valueArg is IControllerLegacyManagedSessionMigrationDocument => { if ( !isPlainObject(valueArg) || !hasExactKeys( valueArg, [ 'id', 'controllerId', 'issuerId', 'projectId', 'harnessId', 'state', 'cohortObservedAt', 'manifest', 'updateId', ], ['completedAt'], ) || !controllerIdPattern.test(String(valueArg.controllerId)) || Number(String(valueArg.controllerId).slice('controller:'.length)) > 65_535 || !isBase64UrlBytes(valueArg.issuerId, 32) || !isBase64UrlBytes(valueArg.projectId, 16) || (valueArg.harnessId !== 'opencode' && valueArg.harnessId !== 'flex') || valueArg.id !== controllerLegacyManagedSessionMigrationDocumentId( valueArg.controllerId as string, valueArg.issuerId, valueArg.projectId, valueArg.harnessId, ) || (valueArg.state !== 'admitting' && valueArg.state !== 'completed') || !isValidDate(valueArg.cohortObservedAt) || !Array.isArray(valueArg.manifest) || valueArg.manifest.length > controllerLegacyManagedSessionMigrationCohortLimit || !isBase64UrlBytes(valueArg.updateId, 32) ) throw new Error('Invalid controller legacy managed-session migration marker.'); let previousStateDocumentId: string | undefined; const runtimeKeys = new Set(); for (const entry of valueArg.manifest) { if ( !isPlainObject(entry) || !hasExactKeys( entry, ['stateDocumentId', 'runtimeId', 'providerSessionGeneration', 'outcome'], ['sessionGenerationId', 'sessionGenerationSequence'], ) || !isRuntimeId(entry.runtimeId, valueArg.harnessId) || entry.stateDocumentId !== controllerSessionStateId( valueArg.controllerId as string, valueArg.projectId as string, entry.runtimeId, ) || !isProviderGeneration(entry.providerSessionGeneration) || !['pending', 'admitted', 'excluded'].includes(String(entry.outcome)) || ( previousStateDocumentId !== undefined && previousStateDocumentId >= entry.stateDocumentId ) || runtimeKeys.has(controllerRuntimeIdKey(entry.runtimeId)) ) throw new Error('Invalid controller legacy managed-session migration manifest.'); const hasSessionGenerationId = Object.hasOwn(entry, 'sessionGenerationId'); const hasSessionGenerationSequence = Object.hasOwn(entry, 'sessionGenerationSequence'); if (valueArg.harnessId === 'flex') { try { assertControllerFlexSessionGeneration( entry.sessionGenerationId, entry.sessionGenerationSequence, entry.providerSessionGeneration, ); } catch { throw new Error('Invalid controller legacy managed-session Flex generation.'); } } else if (hasSessionGenerationId || hasSessionGenerationSequence) { throw new Error('Invalid controller legacy managed-session OpenCode generation.'); } previousStateDocumentId = entry.stateDocumentId; runtimeKeys.add(controllerRuntimeIdKey(entry.runtimeId)); } const hasCompletedAt = Object.hasOwn(valueArg, 'completedAt'); if ( (valueArg.state === 'admitting' && hasCompletedAt) || (valueArg.state === 'completed' && ( !isValidDate(valueArg.completedAt) || valueArg.completedAt.getTime() < valueArg.cohortObservedAt.getTime() || valueArg.manifest.some((entry) => entry.outcome === 'pending') )) ) throw new Error('Invalid controller legacy managed-session migration state.'); }; @plugins.smartdata.managed({ collectionName: 'agl_controller_legacy_managed_session_migrations', }) @plugins.smartdata.exactPersistence({ assertDocument: assertControllerLegacyManagedSessionMigrationDocument, }) export class ControllerLegacyManagedSessionMigrationModel extends plugins.smartdata.SmartDataDbDoc< ControllerLegacyManagedSessionMigrationModel, IControllerLegacyManagedSessionMigrationDocument > { declare static exact: plugins.smartdata.TExact; @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.svDb() public issuerId!: string; @plugins.smartdata.svDb() public projectId!: string; @plugins.smartdata.svDb() public harnessId!: TControllerSessionId['harnessId']; @plugins.smartdata.svDb() public state!: TControllerLegacyManagedSessionMigrationState; @plugins.smartdata.svDb() public cohortObservedAt!: Date; @plugins.smartdata.svDb() public manifest!: IControllerLegacyManagedSessionMigrationManifestEntry[]; @plugins.smartdata.svDb() public updateId!: string; @plugins.smartdata.svDb() public completedAt?: Date; } export class ControllerLegacyManagedSessionV26Migration { private readonly issuerId: string; private readonly sessionIdentityService: ControllerSessionIdentityService; constructor(optionsArg: IControllerLegacyManagedSessionV26MigrationOptions) { if (!isBase64UrlBytes(optionsArg.issuerId, 32)) { throw new Error('Controller legacy managed-session migration issuer is invalid.'); } if (!(optionsArg.sessionIdentityService instanceof ControllerSessionIdentityService)) { throw new Error('Controller legacy managed-session migration service is invalid.'); } this.issuerId = optionsArg.issuerId; this.sessionIdentityService = optionsArg.sessionIdentityService; } private resolveScope( inputArg: IControllerLegacyManagedSessionMigrationInput, ): IResolvedMigrationScope { if ( !controllerIdPattern.test(inputArg.controllerId) || Number(inputArg.controllerId.slice('controller:'.length)) > 65_535 || inputArg.issuerId !== this.issuerId || !isBase64UrlBytes(inputArg.projectId, 16) || (inputArg.harnessId !== 'opencode' && inputArg.harnessId !== 'flex') || !inputArg.signal || typeof inputArg.signal.throwIfAborted !== 'function' ) { throw migrationError('invalid_input', 'Legacy managed-session migration scope is invalid.'); } return { markerId: controllerLegacyManagedSessionMigrationDocumentId( inputArg.controllerId, inputArg.issuerId, inputArg.projectId, inputArg.harnessId, ), controllerId: inputArg.controllerId, issuerId: inputArg.issuerId, projectId: inputArg.projectId, harnessId: inputArg.harnessId, }; } private validateInput( inputArg: IControllerLegacyManagedSessionMigrationInput, scopeArg: IResolvedMigrationScope, ): IValidatedMigrationInput { if ( !isBase64UrlBytes(inputArg.supervisorGeneration, 32) || inputArg.sourceAdmissionFenced !== true || inputArg.snapshotComplete !== true || inputArg.creationObligationsComplete !== true || !isValidDate(inputArg.observedAt) || typeof inputArg.assertRuntimeAuthorityBeforeCompletion !== 'function' || typeof inputArg.assertRuntimeAuthorityAfterCompletion !== 'function' ) { throw migrationError( 'invalid_input', 'Legacy managed-session migration authority input is invalid.', ); } if ( !Array.isArray(inputArg.sessions) || inputArg.sessions.length > controllerSessionIdentitySnapshotEntryLimit ) { throw migrationError( 'limit_exceeded', 'The complete legacy managed-session provider snapshot exceeds its entry limit.', ); } const snapshotByRuntimeKey = new Map(); for (const entry of inputArg.sessions) { inputArg.signal.throwIfAborted(); if ( !isPlainObject(entry) || !hasExactKeys( entry, ['nativeId', 'providerSessionGeneration'], ['sessionGenerationId', 'sessionGenerationSequence', 'parentNativeId'], ) || !isRuntimeId({ harnessId: scopeArg.harnessId, nativeId: entry.nativeId }) || !isProviderGeneration(entry.providerSessionGeneration) || ( entry.parentNativeId !== undefined && !isRuntimeId({ harnessId: scopeArg.harnessId, nativeId: entry.parentNativeId }) ) ) { throw migrationError( 'invalid_input', 'Legacy managed-session snapshot entries must use the exact canonical shape.', ); } const snapshotEntry = entry as unknown as IControllerSessionSnapshotEntry; const hasSessionGenerationId = Object.hasOwn(snapshotEntry, 'sessionGenerationId'); const hasSessionGenerationSequence = Object.hasOwn( snapshotEntry, 'sessionGenerationSequence', ); if (scopeArg.harnessId === 'flex') { try { assertControllerFlexSessionGeneration( snapshotEntry.sessionGenerationId, snapshotEntry.sessionGenerationSequence, snapshotEntry.providerSessionGeneration, ); } catch (errorArg) { throw migrationError( 'invalid_input', 'Legacy managed-session Flex snapshot generation is invalid.', errorArg, ); } } else if (hasSessionGenerationId || hasSessionGenerationSequence) { throw migrationError( 'invalid_input', 'Legacy managed-session OpenCode snapshots cannot carry Flex generation fields.', ); } const runtimeId: TControllerSessionId = { harnessId: scopeArg.harnessId, nativeId: snapshotEntry.nativeId, }; const runtimeKey = controllerRuntimeIdKey(runtimeId); if (snapshotByRuntimeKey.has(runtimeKey)) { throw migrationError( 'invalid_input', 'The complete legacy managed-session snapshot contains a duplicate runtime ID.', ); } snapshotByRuntimeKey.set(runtimeKey, { nativeId: snapshotEntry.nativeId, providerSessionGeneration: snapshotEntry.providerSessionGeneration, ...(snapshotEntry.sessionGenerationId === undefined ? {} : { sessionGenerationId: snapshotEntry.sessionGenerationId }), ...(snapshotEntry.sessionGenerationSequence === undefined ? {} : { sessionGenerationSequence: snapshotEntry.sessionGenerationSequence }), ...(snapshotEntry.parentNativeId === undefined ? {} : { parentNativeId: snapshotEntry.parentNativeId }), }); } if (!Array.isArray(inputArg.creationObligationRuntimeIds)) { throw migrationError( 'invalid_input', 'Legacy managed-session creation obligations must be a complete runtime-ID list.', ); } if (inputArg.creationObligationRuntimeIds.length > controllerSessionIdentitySnapshotEntryLimit) { throw migrationError( 'limit_exceeded', 'The complete legacy creation-obligation runtime-ID list exceeds its entry limit.', ); } const creationObligationRuntimeKeys = new Set(); for (const runtimeId of inputArg.creationObligationRuntimeIds) { inputArg.signal.throwIfAborted(); if (!isRuntimeId(runtimeId, scopeArg.harnessId)) { throw migrationError( 'invalid_input', 'A legacy managed-session creation obligation runtime ID is invalid.', ); } const runtimeKey = controllerRuntimeIdKey(runtimeId); if (creationObligationRuntimeKeys.has(runtimeKey)) { throw migrationError( 'invalid_input', 'The complete creation-obligation list contains a duplicate runtime ID.', ); } creationObligationRuntimeKeys.add(runtimeKey); } return { snapshotByRuntimeKey, creationObligationRuntimeKeys, observedAt: new Date(inputArg.observedAt), }; } private assertMarkerScope( markerArg: IControllerLegacyManagedSessionMigrationDocument, scopeArg: IResolvedMigrationScope, ): void { if ( markerArg.id !== scopeArg.markerId || markerArg.controllerId !== scopeArg.controllerId || markerArg.issuerId !== scopeArg.issuerId || markerArg.projectId !== scopeArg.projectId || markerArg.harnessId !== scopeArg.harnessId ) { throw migrationError( 'corrupt_state', 'Legacy managed-session migration marker escaped its exact scope.', ); } } private async readMarker( markerIdArg: string, signalArg: AbortSignal, ): Promise { signalArg.throwIfAborted(); const stored = await ControllerLegacyManagedSessionMigrationModel.exact.findStoredOne( { id: markerIdArg }, { signal: signalArg }, ); signalArg.throwIfAborted(); return stored ?? undefined; } private async assertRuntimeAuthorityBeforeCompletion( inputArg: IControllerLegacyManagedSessionMigrationInput, ): Promise { inputArg.signal.throwIfAborted(); await inputArg.assertRuntimeAuthorityBeforeCompletion(); inputArg.signal.throwIfAborted(); } private async buildManifest( scopeArg: IResolvedMigrationScope, validatedArg: IValidatedMigrationInput, signalArg: AbortSignal, ): Promise { const manifest: IControllerLegacyManagedSessionMigrationManifestEntry[] = []; let scannedCount = 0; let lastId: string | undefined; while (true) { signalArg.throwIfAborted(); const pageLimit = Math.min( statePageEntryLimit, controllerLegacyManagedSessionMigrationCohortLimit + 1 - scannedCount, ); const page = await ControllerSessionStateModel.exact.findStored({ filter: { controllerId: scopeArg.controllerId, projectId: scopeArg.projectId, 'sessionId.harnessId': scopeArg.harnessId, ...(lastId === undefined ? {} : { id: { $gt: lastId } }), }, sort: { id: 1 }, limit: pageLimit, signal: signalArg, }); signalArg.throwIfAborted(); if (page.length > pageLimit) { throw migrationError( 'corrupt_state', 'A legacy managed-session state page exceeds its requested limit.', ); } if (scannedCount + page.length > controllerLegacyManagedSessionMigrationCohortLimit) { throw migrationError( 'limit_exceeded', 'The legacy managed-session state scope exceeds 500 entries.', ); } if (page.length === 0) break; for (const stored of page) { signalArg.throwIfAborted(); const state = ControllerSessionStateModel.exact.toPersisted(stored); if ( (lastId !== undefined && state.id <= lastId) || state.controllerId !== scopeArg.controllerId || state.projectId !== scopeArg.projectId || state.sessionId.harnessId !== scopeArg.harnessId ) { throw migrationError( 'corrupt_state', 'Legacy managed-session state pagination escaped its exact scope.', ); } scannedCount += 1; lastId = state.id; if (state.deletedAt !== undefined) continue; const runtimeId: TControllerSessionId = { harnessId: scopeArg.harnessId, nativeId: state.sessionId.nativeId, }; const runtimeKey = controllerRuntimeIdKey(runtimeId); const snapshot = validatedArg.snapshotByRuntimeKey.get(runtimeKey); if (!snapshot || validatedArg.creationObligationRuntimeKeys.has(runtimeKey)) continue; if ( scopeArg.harnessId === 'flex' && state.flexProjectManagement !== undefined && ( state.flexProjectManagement.sessionGenerationId !== snapshot.sessionGenerationId || state.flexProjectManagement.sessionGenerationSequence !== snapshot.sessionGenerationSequence ) ) continue; manifest.push({ stateDocumentId: state.id, runtimeId, providerSessionGeneration: snapshot.providerSessionGeneration, ...(snapshot.sessionGenerationId === undefined ? {} : { sessionGenerationId: snapshot.sessionGenerationId }), ...(snapshot.sessionGenerationSequence === undefined ? {} : { sessionGenerationSequence: snapshot.sessionGenerationSequence }), outcome: 'pending', }); } if (page.length < pageLimit) break; } return manifest; } private async insertMarker( candidateArg: IControllerLegacyManagedSessionMigrationDocument, scopeArg: IResolvedMigrationScope, signalArg: AbortSignal, ): Promise { signalArg.throwIfAborted(); let stored: TStoredMarker | undefined; try { stored = (await ControllerLegacyManagedSessionMigrationModel.exact.insert(candidateArg)) .document; } catch (errorArg) { if (!isAmbiguousWrite(errorArg)) throw errorArg; const reconciled = await this.readMarker(candidateArg.id, signalArg); if (!reconciled) { throw migrationError( 'concurrent_change', 'Legacy managed-session marker insertion has an ambiguous outcome.', errorArg, ); } const body = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(reconciled); if (body.updateId === candidateArg.updateId && !immutableMarkerCohortEqual( body, candidateArg, )) { throw migrationError( 'corrupt_state', 'The reconciled legacy managed-session marker changed its immutable cohort.', ); } stored = reconciled; } signalArg.throwIfAborted(); const body = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); this.assertMarkerScope(body, scopeArg); return stored; } private async transitionMarker( currentArg: TStoredMarker, signalArg: AbortSignal, changeArg: (modelArg: ControllerLegacyManagedSessionMigrationModel, updateIdArg: string) => void, postconditionArg: (documentArg: IControllerLegacyManagedSessionMigrationDocument) => boolean, concurrentMessageArg: string, ): Promise { const current = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(currentArg); const updateId = randomId(32); signalArg.throwIfAborted(); let result: Awaited>; try { result = await ControllerLegacyManagedSessionMigrationModel.exact.transition({ current: currentArg, change: (model) => { changeArg(model, updateId); model.updateId = updateId; }, }); } catch (errorArg) { if (!isAmbiguousWrite(errorArg)) throw errorArg; const reconciled = await this.readMarker(current.id, signalArg); if (reconciled) { const body = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(reconciled); if ( body.updateId === updateId && immutableMarkerCohortEqual(current, body) && postconditionArg(body) ) return reconciled; } throw migrationError( 'concurrent_change', `${concurrentMessageArg} Its write outcome is ambiguous.`, errorArg, ); } signalArg.throwIfAborted(); if (result.status === 'transitioned') { const body = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(result.document); if ( body.updateId !== updateId || !immutableMarkerCohortEqual(current, body) || !postconditionArg(body) ) { throw migrationError( 'corrupt_state', 'Legacy managed-session marker transition committed an invalid postcondition.', ); } return result.document; } const reconciled = await this.readMarker(current.id, signalArg); if (reconciled) { const body = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(reconciled); if (immutableMarkerCohortEqual(current, body) && postconditionArg(body)) return reconciled; } throw migrationError('concurrent_change', concurrentMessageArg); } private async transitionOutcome( currentArg: TStoredMarker, entryArg: IControllerLegacyManagedSessionMigrationManifestEntry, outcomeArg: Exclude, signalArg: AbortSignal, ): Promise { return this.transitionMarker( currentArg, signalArg, (model) => { model.manifest = model.manifest.map((entry) => ( entry.stateDocumentId === entryArg.stateDocumentId && entry.outcome === 'pending' ? { ...entry, outcome: outcomeArg } : { ...entry } )); }, (document) => { const transitionedEntry = document.manifest.find( (entry) => entry.stateDocumentId === entryArg.stateDocumentId, ); return (document.state === 'admitting' || document.state === 'completed') && transitionedEntry?.outcome === outcomeArg; }, `Legacy managed-session outcome for ${entryArg.runtimeId.nativeId} changed concurrently.`, ); } private async readUnretiredRuntimeMembership( scopeArg: IResolvedMigrationScope, runtimeIdArg: TControllerSessionId, signalArg: AbortSignal, ): Promise { signalArg.throwIfAborted(); const stored = await ControllerManagedSessionModel.exact.findStored({ filter: { issuerId: scopeArg.issuerId, projectIdentityId: scopeArg.projectId, harnessId: scopeArg.harnessId, 'runtimeId.harnessId': scopeArg.harnessId, 'runtimeId.nativeId': runtimeIdArg.nativeId, state: { $in: ['active', 'deleting'] }, }, sort: { id: 1 }, limit: 2, signal: signalArg, }); signalArg.throwIfAborted(); if (stored.length > 1) { throw migrationError( 'corrupt_state', 'A legacy runtime has multiple unretired managed-session memberships.', ); } if (stored.length === 0) return undefined; const membership = ControllerManagedSessionModel.exact.toPersisted( stored[0] as TStoredManagedSession, ); if ( membership.issuerId !== scopeArg.issuerId || membership.projectIdentityId !== scopeArg.projectId || membership.harnessId !== scopeArg.harnessId || !runtimeIdsEqual(membership.runtimeId, runtimeIdArg) ) { throw migrationError( 'corrupt_state', 'A managed-session membership escaped its legacy migration scope.', ); } if (membership.state === 'deleting') { throw migrationError('deleting', 'A legacy managed session is being deleted.'); } return membership; } private assertMembershipGeneration( membershipArg: IControllerManagedSessionDocument, entryArg: IControllerLegacyManagedSessionMigrationManifestEntry, ): void { if ( !runtimeIdsEqual(membershipArg.runtimeId, entryArg.runtimeId) || !generationFactsEqual(membershipArg, entryArg) ) { throw migrationError( 'provider_generation_mismatch', 'An active managed-session membership conflicts with the sealed legacy generation.', ); } } private async assertStableManagedMembership( scopeArg: IResolvedMigrationScope, entryArg: IControllerLegacyManagedSessionMigrationManifestEntry, signalArg: AbortSignal, ): Promise { const authority = await this.sessionIdentityService.resolveSessionLayoutAuthority({ projectIdentityId: scopeArg.projectId, runtimeId: entryArg.runtimeId, signal: signalArg, }); signalArg.throwIfAborted(); if (authority !== 'managed') { throw migrationError( 'corrupt_state', 'An active legacy managed-session membership lacks stable layout authority.', ); } } private async admitAndVerify( inputArg: IControllerLegacyManagedSessionMigrationInput, scopeArg: IResolvedMigrationScope, entryArg: IControllerLegacyManagedSessionMigrationManifestEntry, observedAtArg: Date, ): Promise { await this.assertRuntimeAuthorityBeforeCompletion(inputArg); const observation = await this.sessionIdentityService.admitManagedSession({ projectIdentityId: scopeArg.projectId, runtimeId: { ...entryArg.runtimeId }, supervisorGeneration: inputArg.supervisorGeneration, providerSessionGeneration: entryArg.providerSessionGeneration, ...(entryArg.sessionGenerationId === undefined ? {} : { sessionGenerationId: entryArg.sessionGenerationId }), ...(entryArg.sessionGenerationSequence === undefined ? {} : { sessionGenerationSequence: entryArg.sessionGenerationSequence }), observedAt: new Date(observedAtArg), admissionSource: 'legacy-state-migration', signal: inputArg.signal, }); inputArg.signal.throwIfAborted(); if ( observation.identity.issuerId !== scopeArg.issuerId || observation.identity.projectIdentityId !== scopeArg.projectId || !runtimeIdsEqual(observation.binding.runtimeId, entryArg.runtimeId) || observation.binding.supervisorGeneration !== inputArg.supervisorGeneration || observation.binding.providerSessionGeneration !== entryArg.providerSessionGeneration ) { throw migrationError( 'corrupt_state', 'Legacy managed-session admission returned a mismatched identity observation.', ); } const membership = await this.readUnretiredRuntimeMembership( scopeArg, entryArg.runtimeId, inputArg.signal, ); if (!membership) { throw migrationError( 'corrupt_state', 'Legacy managed-session admission did not persist an active membership.', ); } this.assertMembershipGeneration(membership, entryArg); await this.assertStableManagedMembership(scopeArg, entryArg, inputArg.signal); } private async resolvePendingOutcome( inputArg: IControllerLegacyManagedSessionMigrationInput, scopeArg: IResolvedMigrationScope, validatedArg: IValidatedMigrationInput, entryArg: IControllerLegacyManagedSessionMigrationManifestEntry, ): Promise> { const runtimeKey = controllerRuntimeIdKey(entryArg.runtimeId); const currentSnapshot = validatedArg.snapshotByRuntimeKey.get(runtimeKey); let membership = await this.readUnretiredRuntimeMembership( scopeArg, entryArg.runtimeId, inputArg.signal, ); if (membership) { this.assertMembershipGeneration(membership, entryArg); if (currentSnapshot && !generationFactsEqual(entryArg, currentSnapshot)) { throw migrationError( 'provider_generation_mismatch', 'The current provider generation conflicts with an active legacy membership.', ); } await this.assertStableManagedMembership(scopeArg, entryArg, inputArg.signal); if (currentSnapshot) { await this.admitAndVerify( inputArg, scopeArg, entryArg, validatedArg.observedAt, ); } return 'admitted'; } const layoutAuthority = await this.sessionIdentityService.resolveSessionLayoutAuthority({ projectIdentityId: scopeArg.projectId, runtimeId: entryArg.runtimeId, signal: inputArg.signal, }); inputArg.signal.throwIfAborted(); if (layoutAuthority === 'managed') { throw migrationError( 'corrupt_state', 'Stable managed layout authority has no discoverable active membership.', ); } if (!currentSnapshot || !generationFactsEqual(entryArg, currentSnapshot)) return 'excluded'; try { await this.admitAndVerify( inputArg, scopeArg, entryArg, validatedArg.observedAt, ); return 'admitted'; } catch (errorArg) { membership = await this.readUnretiredRuntimeMembership( scopeArg, entryArg.runtimeId, inputArg.signal, ); if (membership) { this.assertMembershipGeneration(membership, entryArg); await this.assertStableManagedMembership(scopeArg, entryArg, inputArg.signal); return 'admitted'; } if ( errorArg instanceof ControllerSessionIdentityError && (errorArg.code === 'tombstoned' || errorArg.code === 'provider_generation_mismatch') ) { const settledAuthority = await this.sessionIdentityService.resolveSessionLayoutAuthority({ projectIdentityId: scopeArg.projectId, runtimeId: entryArg.runtimeId, signal: inputArg.signal, }); inputArg.signal.throwIfAborted(); if (settledAuthority === 'definitively_unmanaged') return 'excluded'; } throw errorArg; } } private resultFromMarker( markerArg: IControllerLegacyManagedSessionMigrationDocument, ): IControllerLegacyManagedSessionMigrationResult { return { status: 'completed', markerId: markerArg.id, cohortSize: markerArg.manifest.length, admittedCount: markerArg.manifest.filter((entry) => entry.outcome === 'admitted').length, excludedCount: markerArg.manifest.filter((entry) => entry.outcome === 'excluded').length, }; } public async run( inputArg: IControllerLegacyManagedSessionMigrationInput, ): Promise { const scope = this.resolveScope(inputArg); inputArg.signal.throwIfAborted(); let stored = await this.readMarker(scope.markerId, inputArg.signal); if (stored) { const existing = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); this.assertMarkerScope(existing, scope); if (existing.state === 'completed') return this.resultFromMarker(existing); } const validated = this.validateInput(inputArg, scope); if (!stored) { const manifest = await this.buildManifest(scope, validated, inputArg.signal); await this.assertRuntimeAuthorityBeforeCompletion(inputArg); stored = await this.insertMarker({ id: scope.markerId, controllerId: scope.controllerId, issuerId: scope.issuerId, projectId: scope.projectId, harnessId: scope.harnessId, state: 'admitting', cohortObservedAt: new Date(validated.observedAt), manifest, updateId: randomId(32), }, scope, inputArg.signal); const inserted = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); if (inserted.state === 'completed') return this.resultFromMarker(inserted); } let marker = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); this.assertMarkerScope(marker, scope); if (validated.observedAt.getTime() < marker.cohortObservedAt.getTime()) { throw migrationError( 'invalid_input', 'The current legacy managed-session snapshot predates the sealed cohort.', ); } const manifestIds = marker.manifest.map((entry) => entry.stateDocumentId); for (const stateDocumentId of manifestIds) { inputArg.signal.throwIfAborted(); marker = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); if (marker.state === 'completed') return this.resultFromMarker(marker); const entry = marker.manifest.find((candidate) => ( candidate.stateDocumentId === stateDocumentId )); if (!entry) { throw migrationError( 'corrupt_state', 'The sealed legacy managed-session manifest changed during admission.', ); } if (entry.outcome !== 'pending') continue; await this.assertRuntimeAuthorityBeforeCompletion(inputArg); const outcome = await this.resolvePendingOutcome(inputArg, scope, validated, entry); await this.assertRuntimeAuthorityBeforeCompletion(inputArg); stored = await this.transitionOutcome(stored, entry, outcome, inputArg.signal); } marker = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); if (marker.state === 'completed') return this.resultFromMarker(marker); if (marker.manifest.some((entry) => entry.outcome === 'pending')) { throw migrationError( 'concurrent_change', 'Legacy managed-session migration still has pending manifest outcomes.', ); } await this.assertRuntimeAuthorityBeforeCompletion(inputArg); stored = await this.transitionMarker( stored, inputArg.signal, (model) => { if (model.state === 'admitting' && model.manifest.every( (entry) => entry.outcome !== 'pending', )) { model.state = 'completed'; model.completedAt = new Date(validated.observedAt); } }, (document) => document.state === 'completed', 'Legacy managed-session migration completion changed concurrently.', ); await inputArg.assertRuntimeAuthorityAfterCompletion(); inputArg.signal.throwIfAborted(); marker = ControllerLegacyManagedSessionMigrationModel.exact.toPersisted(stored); return this.resultFromMarker(marker); } }