// // Copyright 2022 DXOS.org // import { runInContextAsync, synchronized } from '@dxos/async'; import { Context } from '@dxos/context'; import { PublicKey } from '@dxos/keys'; import { log } from '@dxos/log'; import { type TypedMessage } from '@dxos/protocols/proto'; import { type Credential, MembershipPolicy, SpaceMember } from '@dxos/protocols/proto/dxos/halo/credentials'; import { type DelegateSpaceInvitation } from '@dxos/protocols/proto/dxos/halo/invitations'; import { type AsyncCallback, Callback, ComplexMap, ComplexSet } from '@dxos/util'; import { getCredentialAssertion, verifyCredential } from '../credentials'; import { type CredentialProcessor } from '../processor/credential-processor'; import { type FeedInfo, FeedStateMachine } from './feed-state-machine'; import { InvitationStateMachine } from './invitation-state-machine'; import { type MemberInfo, MemberStateMachine } from './member-state-machine'; export interface SpaceState { readonly members: ReadonlyMap; readonly membershipChainHeads: PublicKey[]; readonly feeds: ReadonlyMap; readonly credentials: Credential[]; readonly genesisCredential: Credential | undefined; readonly tags: string[]; readonly membershipPolicy: MembershipPolicy; readonly creator: MemberInfo | undefined; readonly invitations: ReadonlyMap; addCredentialProcessor(processor: CredentialProcessor): Promise; removeCredentialProcessor(processor: CredentialProcessor): Promise; getCredentialsOfType(type: TypedMessage['@type']): Credential[]; getMemberRole(memberKey: PublicKey): SpaceMember.Role; hasMembershipManagementPermission(memberKey: PublicKey): boolean; } export type ProcessOptions = { sourceFeed: PublicKey; skipVerification?: boolean; }; export type CredentialEntry = { credential: Credential; sourceFeed: PublicKey; revoked: boolean; }; /** * Validates and processes credentials for a single space. * Keeps a list of members and feeds. * Keeps and in-memory index of credentials and allows to query them. */ export class SpaceStateMachine implements SpaceState { private readonly _members: MemberStateMachine; private readonly _feeds: FeedStateMachine; private readonly _invitations = new InvitationStateMachine(); private readonly _credentials: CredentialEntry[] = []; private readonly _credentialsById = new ComplexMap(PublicKey.hash); private readonly _processedCredentials = new ComplexSet(PublicKey.hash); private _genesisCredential: Credential | undefined; private _tags: string[] = []; private _membershipPolicy: MembershipPolicy = MembershipPolicy.INVITE; private _credentialProcessors: CredentialConsumer[] = []; readonly onCredentialProcessed = new Callback>(); readonly onMemberRoleChanged: Callback>; readonly onFeedAdmitted: Callback>; readonly onDelegatedInvitation = this._invitations.onDelegatedInvitation; readonly onDelegatedInvitationRemoved = this._invitations.onDelegatedInvitationRemoved; constructor(private readonly _spaceKey: PublicKey) { this._members = new MemberStateMachine(this._spaceKey); this._feeds = new FeedStateMachine(this._spaceKey); this.onMemberRoleChanged = this._members.onMemberRoleChanged; this.onFeedAdmitted = this._feeds.onFeedAdmitted; } get creator(): MemberInfo | undefined { return this._members.creator; } get members(): ReadonlyMap { return this._members.members; } get membershipChainHeads(): PublicKey[] { return this._members.membershipChainHeads; } get feeds(): ReadonlyMap { return this._feeds.feeds; } get credentials(): Credential[] { return this._credentials.map((entry) => entry.credential); } get credentialEntries(): CredentialEntry[] { return this._credentials; } get genesisCredential(): Credential | undefined { return this._genesisCredential; } get tags(): string[] { return this._tags; } get membershipPolicy(): MembershipPolicy { return this._membershipPolicy; } get invitations(): ReadonlyMap { return this._invitations.invitations; } async addCredentialProcessor(processor: CredentialProcessor): Promise { if (this._credentialProcessors.find((p) => p.processor === processor)) { throw new Error('Credential processor already added.'); } const consumer = new CredentialConsumer( processor, async () => { for (const credential of this.credentials) { await consumer._process(credential); } // NOTE: It is important to set this flag after immediately after processing existing credentials. // Otherwise, we might miss some credentials. // Having an `await` statement between the end of the loop and setting the flag would cause a race condition. consumer._isReadyForLiveCredentials = true; }, async () => { this._credentialProcessors = this._credentialProcessors.filter((p) => p !== consumer); }, ); this._credentialProcessors.push(consumer); await consumer.open(); } async removeCredentialProcessor(processor: CredentialProcessor): Promise { const consumer = this._credentialProcessors.find((p) => p.processor === processor); await consumer?.close(); } getCredentialsOfType(type: TypedMessage['@type']): Credential[] { return this.credentials.filter((credential) => getCredentialAssertion(credential)['@type'] === type); } /** * @param credential Message to process. * @param fromFeed Key of the feed where this credential is recorded. */ @synchronized async process(credential: Credential, { sourceFeed, skipVerification }: ProcessOptions): Promise { if (credential.id) { if (this._processedCredentials.has(credential.id)) { return true; } this._processedCredentials.add(credential.id); } if (!skipVerification) { const result = await verifyCredential(credential); if (result.kind !== 'pass') { log.warn(`Invalid credential: ${result.errors.join(', ')}`); return false; } } const assertion = getCredentialAssertion(credential); switch (assertion['@type']) { case 'dxos.halo.credentials.SpaceGenesis': { if (this._genesisCredential) { log.warn('Space already has a genesis credential.'); return false; } if (!credential.issuer.equals(this._spaceKey)) { log.warn('Space genesis credential must be issued by space.'); return false; } if (!credential.subject.id.equals(this._spaceKey)) { log.warn('Space genesis credential must be issued to space.'); return false; } this._genesisCredential = credential; this._tags = assertion.tags ?? []; this._membershipPolicy = assertion.membershipPolicy ?? MembershipPolicy.INVITE; break; } case 'dxos.halo.credentials.SpaceMember': { if (!assertion.spaceKey.equals(this._spaceKey)) { break; // Ignore credentials for other spaces. } if (!this._genesisCredential) { log.warn('Space must have a genesis credential before adding members.'); return false; } if (!this._canInviteNewMembers(credential.issuer)) { log.warn(`Space member is not authorized to invite new members: ${credential.issuer}`); return false; } await this._members.process(credential); await this._invitations.process(credential); break; } case 'dxos.halo.credentials.MemberProfile': { if (!this._genesisCredential) { log.warn('Space must have a genesis credential before adding members.'); return false; } await this._members.process(credential); break; } case 'dxos.halo.credentials.AdmittedFeed': { if (!this._genesisCredential) { log.warn('Space must have a genesis credential before admitting feeds.'); return false; } // We don't do any validation on feed admission since we would perform the same validation on the credentials inside . await this._feeds.process(credential, sourceFeed); break; } case 'dxos.halo.invitations.CancelDelegatedInvitation': case 'dxos.halo.invitations.DelegateSpaceInvitation': { if (!this._canInviteNewMembers(credential.issuer)) { log.warn(`Invalid invitation, space member is not authorized to invite new members: ${credential.issuer}`); return false; } await this._invitations.process(credential); break; } } const newEntry: CredentialEntry = { credential, sourceFeed, revoked: false }; this._credentials.push(newEntry); // TODO(dmaretskyi): Invariant on every credential having an id? if (credential.id) { this._credentialsById.set(credential.id, newEntry); } for (const processor of this._credentialProcessors) { if (processor._isReadyForLiveCredentials) { await processor._process(credential); } } await this.onCredentialProcessed.callIfSet(credential); return true; } public getMemberRole(memberKey: PublicKey): SpaceMember.Role { return this._members.getRole(memberKey); } public hasMembershipManagementPermission(memberKey: PublicKey): boolean { return this._canInviteNewMembers(memberKey); } private _canInviteNewMembers(key: PublicKey): boolean { if (this._membershipPolicy === MembershipPolicy.LOCKED) { // When locked, only the space key can add the initial owner during genesis. // Once a member exists, no new members can be added. return key.equals(this._spaceKey) && this._members.members.size === 0; } return ( key.equals(this._spaceKey) || this._members.getRole(key) === SpaceMember.Role.ADMIN || this._members.getRole(key) === SpaceMember.Role.OWNER ); } } // TODO(dmaretskyi): Simplify. class CredentialConsumer { private _ctx = new Context(); /** * @internal * Processor is ready to process live credentials. * NOTE: Setting this flag before all existing credentials are processed will cause them to be processed out of order. * Set externally. */ _isReadyForLiveCredentials = false; constructor( public readonly processor: T, private readonly _onOpen: () => Promise, private readonly _onClose: () => Promise, ) {} /** * @internal */ async _process(credential: Credential): Promise { await runInContextAsync(this._ctx, async () => { await this.processor.processCredential(credential); }); } async open(): Promise { if (this._ctx.disposed) { throw new Error('CredentialProcessor is disposed'); } await this._onOpen(); } async close(): Promise { await this._ctx.dispose(); await this._onClose(); } }