import * as plugins from './plugins.js'; import { hasFlexExactKeys, isFlexPlainObject, isFlexPublicMessage, isFlexPublicSession, } from './interfaces.flexipc.js'; type TFlexMessage = plugins.flexharness.IFlexMessage; type TFlexSession = plugins.flexharness.IFlexSession; export const flexPublicSessionLimit = 2048; export const flexPublicMessageLimit = 2048; export const flexPublicMessageBytesLimit = 64 * 1024 * 1024; export const flexPublicSingleMessageBytesLimit = 480 * 1024; export const flexPublicPageBytesLimit = 512 * 1024; export const flexPublicPageLimit = 50; export interface IFlexRecordManifest { id: string; sha256: string; bytes: number; } export interface IFlexPublicSourceRevision { incarnationId: string; revision: number; } export interface IFlexPublicProjectionSource extends IFlexPublicSourceRevision { sessionId: string; } export interface IFlexPublicCandidateManifest { candidateId: string; revision: number; sha256: string; bytes: number; sessionRecords: IFlexRecordManifest[]; messageRecords: IFlexRecordManifest[]; sessionsTruncated: boolean; messagesTruncated: boolean; messageBytesTruncated: boolean; visibleDigest: string; scopeSource: IFlexPublicSourceRevision; projectionSources: IFlexPublicProjectionSource[]; } export interface IFlexPublicCandidateDocument extends IFlexPublicCandidateManifest { id: string; controllerId: string; scopeId: string; createdAt: string; } export interface IFlexPublicHeadDocument { id: string; controllerId: string; storageKey: string; revision: number; currentPublicCandidateId: string; publicCandidates: IFlexPublicCandidateManifest[]; } export interface IFlexPublicSessionRecordDocument { id: string; controllerId: string; scopeId: string; candidateId: string; revision: number; session: TFlexSession; createdAt: string; } export interface IFlexPublicMessageRecordDocument { id: string; controllerId: string; scopeId: string; candidateId: string; revision: number; messageIndex: number; message: TFlexMessage; createdAt: string; } export interface IFlexProjectionMessageDescriptor { id: string; index: number; /** Absent only on projection sources created before source-order repair. */ messageIndex?: number; sha256: string; bytes: number; messageCreatedAt: string; sessionId: string; messageId: string; } export interface IFlexProjectionSourceManifest { candidateId: string; /** Absent only on persisted source chains created before source-order repair. */ orderVersion?: 1; sha256: string; bytes: number; messageCount: number; messagesTruncated: boolean; visibleDigest: string; latestMessage?: IFlexProjectionMessageDescriptor; } export interface IFlexProjectionSourceDocument extends IFlexProjectionSourceManifest { id: string; controllerId: string; scopeId: string; sessionId: string; incarnationId: string; revision: number; createdAt: string; } export interface IFlexProjectionMessageRecordDocument extends IFlexProjectionMessageDescriptor { controllerId: string; scopeId: string; candidateId: string; incarnationId: string; revision: number; message: TFlexMessage; previous?: IFlexProjectionMessageDescriptor; createdAt: string; } const digestPattern = /^[a-f0-9]{64}$/u; const candidateIdPattern = /^[A-Za-z0-9_-]{22}$/u; const isBoundedString = (valueArg: unknown, maximumBytesArg: number): valueArg is string => typeof valueArg === 'string' && valueArg.length > 0 && Buffer.byteLength(valueArg, 'utf8') <= maximumBytesArg; const isNonNegativeInteger = (valueArg: unknown): valueArg is number => Number.isSafeInteger(valueArg) && (valueArg as number) >= 0; const isIsoDate = (valueArg: unknown): valueArg is string => typeof valueArg === 'string' && Number.isFinite(Date.parse(valueArg)) && new Date(valueArg).toISOString() === valueArg; export const isFlexProjectionMessageDescriptor = ( valueArg: unknown, ): valueArg is IFlexProjectionMessageDescriptor => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, [ 'id', 'index', 'sha256', 'bytes', 'messageCreatedAt', 'sessionId', 'messageId', ], ['messageIndex']) && isBoundedString(valueArg.id, 1024) && isNonNegativeInteger(valueArg.index) && (valueArg.messageIndex === undefined || isNonNegativeInteger(valueArg.messageIndex)) && typeof valueArg.sha256 === 'string' && digestPattern.test(valueArg.sha256) && isNonNegativeInteger(valueArg.bytes) && valueArg.bytes <= flexPublicSingleMessageBytesLimit + 4096 && isIsoDate(valueArg.messageCreatedAt) && isBoundedString(valueArg.sessionId, 512) && isBoundedString(valueArg.messageId, 512); export const isFlexProjectionSourceManifest = ( valueArg: unknown, ): valueArg is IFlexProjectionSourceManifest => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, [ 'candidateId', 'sha256', 'bytes', 'messageCount', 'messagesTruncated', 'visibleDigest', ], ['latestMessage', 'orderVersion']) && typeof valueArg.candidateId === 'string' && candidateIdPattern.test(valueArg.candidateId) && (valueArg.orderVersion === undefined || valueArg.orderVersion === 1) && typeof valueArg.sha256 === 'string' && digestPattern.test(valueArg.sha256) && isNonNegativeInteger(valueArg.bytes) && valueArg.bytes <= 4096 && isNonNegativeInteger(valueArg.messageCount) && typeof valueArg.messagesTruncated === 'boolean' && typeof valueArg.visibleDigest === 'string' && digestPattern.test(valueArg.visibleDigest) && (valueArg.latestMessage === undefined ? valueArg.messageCount === 0 : isFlexProjectionMessageDescriptor(valueArg.latestMessage) && valueArg.messageCount === valueArg.latestMessage.index + 1) && ( valueArg.orderVersion !== 1 || valueArg.latestMessage === undefined || ( isFlexProjectionMessageDescriptor(valueArg.latestMessage) && valueArg.latestMessage.messageIndex !== undefined ) ); export const isFlexRecordManifest = (valueArg: unknown): valueArg is IFlexRecordManifest => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, ['id', 'sha256', 'bytes']) && isBoundedString(valueArg.id, 1024) && typeof valueArg.sha256 === 'string' && digestPattern.test(valueArg.sha256) && isNonNegativeInteger(valueArg.bytes); const isFlexPublicSourceRevision = (valueArg: unknown): valueArg is IFlexPublicSourceRevision => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, ['incarnationId', 'revision']) && typeof valueArg.incarnationId === 'string' && candidateIdPattern.test(valueArg.incarnationId) && isNonNegativeInteger(valueArg.revision) && valueArg.revision >= 1; const isFlexPublicProjectionSource = ( valueArg: unknown, ): valueArg is IFlexPublicProjectionSource => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, ['sessionId', 'incarnationId', 'revision']) && isBoundedString(valueArg.sessionId, 512) && typeof valueArg.incarnationId === 'string' && candidateIdPattern.test(valueArg.incarnationId) && isNonNegativeInteger(valueArg.revision) && valueArg.revision >= 1; export const isFlexPublicCandidateManifest = ( valueArg: unknown, ): valueArg is IFlexPublicCandidateManifest => isFlexPlainObject(valueArg) && hasFlexExactKeys(valueArg, [ 'candidateId', 'revision', 'sha256', 'bytes', 'sessionRecords', 'messageRecords', 'sessionsTruncated', 'messagesTruncated', 'messageBytesTruncated', 'visibleDigest', 'scopeSource', 'projectionSources', ]) && typeof valueArg.candidateId === 'string' && candidateIdPattern.test(valueArg.candidateId) && isNonNegativeInteger(valueArg.revision) && typeof valueArg.sha256 === 'string' && digestPattern.test(valueArg.sha256) && isNonNegativeInteger(valueArg.bytes) && Array.isArray(valueArg.sessionRecords) && valueArg.sessionRecords.length <= flexPublicSessionLimit && valueArg.sessionRecords.every(isFlexRecordManifest) && Array.isArray(valueArg.messageRecords) && valueArg.messageRecords.length <= flexPublicMessageLimit && valueArg.messageRecords.every(isFlexRecordManifest) && typeof valueArg.sessionsTruncated === 'boolean' && typeof valueArg.messagesTruncated === 'boolean' && typeof valueArg.messageBytesTruncated === 'boolean' && typeof valueArg.visibleDigest === 'string' && digestPattern.test(valueArg.visibleDigest) && isFlexPublicSourceRevision(valueArg.scopeSource) && Array.isArray(valueArg.projectionSources) && valueArg.projectionSources.length <= flexPublicSessionLimit && valueArg.projectionSources.every(isFlexPublicProjectionSource) && new Set(valueArg.projectionSources.map((source) => source.sessionId)).size === valueArg.projectionSources.length; export function assertFlexPublicCandidateDocument( valueArg: unknown, ): asserts valueArg is IFlexPublicCandidateDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'scopeId', 'candidateId', 'revision', 'sha256', 'bytes', 'sessionRecords', 'messageRecords', 'sessionsTruncated', 'messagesTruncated', 'messageBytesTruncated', 'visibleDigest', 'scopeSource', 'projectionSources', 'createdAt', ]) || !isBoundedString(valueArg.id, 1024) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.scopeId, 512) || !isIsoDate(valueArg.createdAt) || !isFlexPublicCandidateManifest({ candidateId: valueArg.candidateId, revision: valueArg.revision, sha256: valueArg.sha256, bytes: valueArg.bytes, sessionRecords: valueArg.sessionRecords, messageRecords: valueArg.messageRecords, sessionsTruncated: valueArg.sessionsTruncated, messagesTruncated: valueArg.messagesTruncated, messageBytesTruncated: valueArg.messageBytesTruncated, visibleDigest: valueArg.visibleDigest, scopeSource: valueArg.scopeSource, projectionSources: valueArg.projectionSources, }) ) throw new Error('Invalid Flex public candidate document.'); } export function assertFlexPublicHeadDocument( valueArg: unknown, ): asserts valueArg is IFlexPublicHeadDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'storageKey', 'revision', 'currentPublicCandidateId', 'publicCandidates', ]) || !isBoundedString(valueArg.id, 1024) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.storageKey, 512) || !isNonNegativeInteger(valueArg.revision) || valueArg.revision < 1 || typeof valueArg.currentPublicCandidateId !== 'string' || !candidateIdPattern.test(valueArg.currentPublicCandidateId) || !Array.isArray(valueArg.publicCandidates) || valueArg.publicCandidates.length < 1 || valueArg.publicCandidates.length > 4 || !valueArg.publicCandidates.every(isFlexPublicCandidateManifest) || valueArg.publicCandidates[0]?.candidateId !== valueArg.currentPublicCandidateId || valueArg.publicCandidates[0]?.revision !== valueArg.revision || valueArg.publicCandidates.some((candidate, index, candidates) => ( index > 0 && candidate.revision >= candidates[index - 1]!.revision )) ) throw new Error('Invalid Flex public head document.'); } export function assertFlexPublicSessionRecordDocument( valueArg: unknown, ): asserts valueArg is IFlexPublicSessionRecordDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'scopeId', 'candidateId', 'revision', 'session', 'createdAt', ]) || !isBoundedString(valueArg.id, 1024) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.scopeId, 512) || typeof valueArg.candidateId !== 'string' || !candidateIdPattern.test(valueArg.candidateId) || !isNonNegativeInteger(valueArg.revision) || !isIsoDate(valueArg.createdAt) || !isFlexPublicSession(valueArg.session) || valueArg.session.scopeId !== valueArg.scopeId ) throw new Error('Invalid Flex public session record.'); } export function assertFlexPublicMessageRecordDocument( valueArg: unknown, ): asserts valueArg is IFlexPublicMessageRecordDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'scopeId', 'candidateId', 'revision', 'messageIndex', 'message', 'createdAt', ]) || !isBoundedString(valueArg.id, 1024) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.scopeId, 512) || typeof valueArg.candidateId !== 'string' || !candidateIdPattern.test(valueArg.candidateId) || !isNonNegativeInteger(valueArg.revision) || !isNonNegativeInteger(valueArg.messageIndex) || !isIsoDate(valueArg.createdAt) || !isFlexPublicMessage(valueArg.message) ) throw new Error('Invalid Flex public message record.'); } export function assertFlexProjectionSourceDocument( valueArg: unknown, ): asserts valueArg is IFlexProjectionSourceDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'scopeId', 'sessionId', 'incarnationId', 'revision', 'candidateId', 'sha256', 'bytes', 'messageCount', 'messagesTruncated', 'visibleDigest', 'createdAt', ], ['latestMessage', 'orderVersion']) || !isBoundedString(valueArg.id, 1024) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.scopeId, 512) || !isBoundedString(valueArg.sessionId, 512) || typeof valueArg.incarnationId !== 'string' || !candidateIdPattern.test(valueArg.incarnationId) || !isNonNegativeInteger(valueArg.revision) || valueArg.revision < 1 || !isIsoDate(valueArg.createdAt) || !isFlexProjectionSourceManifest({ candidateId: valueArg.candidateId, ...(valueArg.orderVersion === undefined ? {} : { orderVersion: valueArg.orderVersion }), sha256: valueArg.sha256, bytes: valueArg.bytes, messageCount: valueArg.messageCount, messagesTruncated: valueArg.messagesTruncated, visibleDigest: valueArg.visibleDigest, ...(valueArg.latestMessage === undefined ? {} : { latestMessage: valueArg.latestMessage }), }) ) throw new Error('Invalid Flex projection source document.'); } export function assertFlexProjectionMessageRecordDocument( valueArg: unknown, ): asserts valueArg is IFlexProjectionMessageRecordDocument { if ( !isFlexPlainObject(valueArg) || !hasFlexExactKeys(valueArg, [ 'id', 'controllerId', 'scopeId', 'candidateId', 'sessionId', 'incarnationId', 'revision', 'index', 'sha256', 'bytes', 'messageCreatedAt', 'messageId', 'message', 'createdAt', ], ['messageIndex', 'previous']) || !isBoundedString(valueArg.controllerId, 512) || !isBoundedString(valueArg.scopeId, 512) || typeof valueArg.candidateId !== 'string' || !candidateIdPattern.test(valueArg.candidateId) || !isBoundedString(valueArg.sessionId, 512) || typeof valueArg.incarnationId !== 'string' || !candidateIdPattern.test(valueArg.incarnationId) || !isNonNegativeInteger(valueArg.revision) || valueArg.revision < 1 || !isIsoDate(valueArg.createdAt) || !isFlexProjectionMessageDescriptor({ id: valueArg.id, index: valueArg.index, ...(valueArg.messageIndex === undefined ? {} : { messageIndex: valueArg.messageIndex }), sha256: valueArg.sha256, bytes: valueArg.bytes, messageCreatedAt: valueArg.messageCreatedAt, sessionId: valueArg.sessionId, messageId: valueArg.messageId, }) || !isFlexPublicMessage(valueArg.message) || valueArg.message.sessionId !== valueArg.sessionId || valueArg.message.messageId !== valueArg.messageId || valueArg.message.createdAt !== valueArg.messageCreatedAt || (valueArg.previous !== undefined && !isFlexProjectionMessageDescriptor(valueArg.previous)) ) throw new Error('Invalid Flex projection message record.'); } @plugins.smartdata.compoundIndex({ name: 'flex_public_candidate_orphans', key: { controllerId: 1, createdAt: 1, _id: 1 }, }) @plugins.smartdata.managed({ collectionName: 'flex_public_candidates' }) export class FlexPublicCandidateModel extends plugins.smartdata.SmartDataDbDoc< FlexPublicCandidateModel, IFlexPublicCandidateDocument > { @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.index() public scopeId!: string; @plugins.smartdata.index() public candidateId!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public sha256!: string; @plugins.smartdata.svDb() public bytes!: number; @plugins.smartdata.svDb() public sessionRecords!: IFlexRecordManifest[]; @plugins.smartdata.svDb() public messageRecords!: IFlexRecordManifest[]; @plugins.smartdata.svDb() public sessionsTruncated!: boolean; @plugins.smartdata.svDb() public messagesTruncated!: boolean; @plugins.smartdata.svDb() public messageBytesTruncated!: boolean; @plugins.smartdata.svDb() public visibleDigest!: string; @plugins.smartdata.svDb() public scopeSource!: IFlexPublicSourceRevision; @plugins.smartdata.svDb() public projectionSources!: IFlexPublicProjectionSource[]; @plugins.smartdata.index() public createdAt!: string; } @plugins.smartdata.compoundIndex({ name: 'flex_public_head_candidate_owner', key: { controllerId: 1, 'publicCandidates.candidateId': 1 }, }) @plugins.smartdata.managed({ collectionName: 'flex_public_heads' }) @plugins.smartdata.exactPersistence({ assertDocument: assertFlexPublicHeadDocument }) export class FlexPublicHeadModel extends plugins.smartdata.SmartDataDbDoc< FlexPublicHeadModel, IFlexPublicHeadDocument > { declare static exact: plugins.smartdata.TExact; @plugins.smartdata.unI() public id!: string; @plugins.smartdata.index() public controllerId!: string; @plugins.smartdata.index() public storageKey!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public currentPublicCandidateId!: string; @plugins.smartdata.svDb() public publicCandidates!: IFlexPublicCandidateManifest[]; } @plugins.smartdata.managed({ collectionName: 'flex_public_sessions' }) export class FlexPublicSessionRecordModel extends plugins.smartdata.SmartDataDbDoc< FlexPublicSessionRecordModel, IFlexPublicSessionRecordDocument > { @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.index() public scopeId!: string; @plugins.smartdata.index() public candidateId!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public session!: TFlexSession; @plugins.smartdata.index() public createdAt!: string; } @plugins.smartdata.managed({ collectionName: 'flex_public_messages' }) export class FlexPublicMessageRecordModel extends plugins.smartdata.SmartDataDbDoc< FlexPublicMessageRecordModel, IFlexPublicMessageRecordDocument > { @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.index() public scopeId!: string; @plugins.smartdata.index() public candidateId!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public messageIndex!: number; @plugins.smartdata.svDb() public message!: TFlexMessage; @plugins.smartdata.index() public createdAt!: string; } @plugins.smartdata.compoundIndex({ name: 'flex_projection_source_orphans', key: { controllerId: 1, createdAt: 1, _id: 1 }, }) @plugins.smartdata.managed({ collectionName: 'flex_projection_sources' }) export class FlexProjectionSourceModel extends plugins.smartdata.SmartDataDbDoc< FlexProjectionSourceModel, IFlexProjectionSourceDocument > { @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.index() public scopeId!: string; @plugins.smartdata.index() public sessionId!: string; @plugins.smartdata.index() public candidateId!: string; @plugins.smartdata.svDb() public incarnationId!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public orderVersion?: 1; @plugins.smartdata.svDb() public sha256!: string; @plugins.smartdata.svDb() public bytes!: number; @plugins.smartdata.svDb() public messageCount!: number; @plugins.smartdata.svDb() public messagesTruncated!: boolean; @plugins.smartdata.svDb() public visibleDigest!: string; @plugins.smartdata.svDb() public latestMessage?: IFlexProjectionMessageDescriptor; @plugins.smartdata.index() public createdAt!: string; } @plugins.smartdata.managed({ collectionName: 'flex_projection_messages' }) export class FlexProjectionMessageRecordModel extends plugins.smartdata.SmartDataDbDoc< FlexProjectionMessageRecordModel, IFlexProjectionMessageRecordDocument > { @plugins.smartdata.unI() public id!: string; @plugins.smartdata.svDb() public controllerId!: string; @plugins.smartdata.index() public scopeId!: string; @plugins.smartdata.index() public candidateId!: string; @plugins.smartdata.index() public sessionId!: string; @plugins.smartdata.svDb() public incarnationId!: string; @plugins.smartdata.svDb() public revision!: number; @plugins.smartdata.svDb() public index!: number; @plugins.smartdata.svDb() public messageIndex?: number; @plugins.smartdata.svDb() public sha256!: string; @plugins.smartdata.svDb() public bytes!: number; @plugins.smartdata.svDb() public messageCreatedAt!: string; @plugins.smartdata.svDb() public messageId!: string; @plugins.smartdata.svDb() public message!: TFlexMessage; @plugins.smartdata.svDb() public previous?: IFlexProjectionMessageDescriptor; @plugins.smartdata.index() public createdAt!: string; } export const flexHeadId = (controllerIdArg: string, scopeIdArg: string): string => `flex-head:${plugins.crypto.createHash('sha256') .update(controllerIdArg, 'utf8') .update('\0', 'utf8') .update(scopeIdArg, 'utf8') .digest('hex')}`; export const flexPublicHeadId = (controllerIdArg: string, scopeIdArg: string): string => `flex-public-head:${plugins.crypto.createHash('sha256') .update(controllerIdArg, 'utf8') .update('\0', 'utf8') .update(scopeIdArg, 'utf8') .digest('hex')}`; export const registerFlexPublicModels = ( managerArg: { db: plugins.smartdata.SmartdataDb }, ): void => { plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexPublicCandidateModel); plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexPublicHeadModel); plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexPublicSessionRecordModel); plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexPublicMessageRecordModel); plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexProjectionSourceModel); plugins.smartdata.setDefaultManagerForDoc(managerArg, FlexProjectionMessageRecordModel); }; export const initFlexPublicModels = async (): Promise => { await FlexPublicCandidateModel.init(); await FlexPublicHeadModel.init(); await FlexPublicSessionRecordModel.init(); await FlexPublicMessageRecordModel.init(); await FlexProjectionSourceModel.init(); await FlexProjectionMessageRecordModel.init(); };