diff --git a/packages/app/src/pages/testing-only.tsx b/packages/app/src/pages/testing-only.tsx new file mode 100644 index 000000000..a050496bb --- /dev/null +++ b/packages/app/src/pages/testing-only.tsx @@ -0,0 +1,18 @@ +import { useEffect } from 'react'; +import { getDataCenter } from '@affine/datacenter'; +/** + * testing only when development + */ + +const Page = () => { + useEffect(() => { + getDataCenter().then(dc => { + // @ts-expect-error global variable + window.dc = dc; + }); + }, []); + + return
Testing only
; +}; + +export default Page; diff --git a/packages/data-center/src/datacenter.bk.ts b/packages/data-center/src/datacenter.bk.ts deleted file mode 100644 index dd01e7785..000000000 --- a/packages/data-center/src/datacenter.bk.ts +++ /dev/null @@ -1,234 +0,0 @@ -import assert from 'assert'; -import { BlockSchema } from '@blocksuite/blocks/models'; -import { Workspace, Signal } from '@blocksuite/store'; - -import { getLogger } from './index.js'; -import { getApis, Apis } from './apis/index.js'; -import { AffineProvider, BaseProvider } from './provider/index.js'; -import { LocalProvider } from './provider/index.js'; -import { getKVConfigure } from './store.js'; -import { Workspaces } from './workspaces'; - -// load workspace's config -type LoadConfig = { - // use witch provider load data - providerId?: string; - // provider config - config?: Record; -}; - -export type DataCenterSignals = DataCenter['signals']; -type WorkspaceItem = { - // provider id - provider: string; - // data exists locally - locally: boolean; -}; -type WorkspaceLoadEvent = WorkspaceItem & { - workspace: string; -}; - -export class DataCenter { - private readonly _apis: Apis; - private readonly _providers = new Map(); - private readonly _workspaces = new Map>(); - private readonly _config; - private readonly _logger; - public readonly workspacesList = new Workspaces(this); - - readonly signals = { - listAdd: new Signal(), - listRemove: new Signal(), - }; - - static async init(debug: boolean): Promise { - const dc = new DataCenter(debug); - dc.addProvider(AffineProvider); - dc.addProvider(LocalProvider); - - return dc; - } - - private constructor(debug: boolean) { - this._apis = getApis(); - this._config = getKVConfigure('sys'); - this._logger = getLogger('dc'); - this._logger.enabled = debug; - - this.signals.listAdd.on(e => { - this._config.set(`list:${e.workspace}`, { - provider: e.provider, - locally: e.locally, - }); - }); - this.signals.listRemove.on(workspace => { - this._config.delete(`list:${workspace}`); - }); - this.workspacesList.init(); - } - - get apis(): Readonly { - return this._apis; - } - - private addProvider(provider: typeof BaseProvider) { - this._providers.set(provider.id, provider); - } - - private async _getProvider( - id: string, - providerId = 'local' - ): Promise { - const providerKey = `${id}:provider`; - if (this._providers.has(providerId)) { - await this._config.set(providerKey, providerId); - return providerId; - } else { - const providerValue = await this._config.get(providerKey); - if (providerValue) return providerValue; - } - throw Error(`Provider ${providerId} not found`); - } - - private async _getWorkspace( - id: string, - params: LoadConfig - ): Promise { - this._logger(`Init workspace ${id} with ${params.providerId}`); - - const providerId = await this._getProvider(id, params.providerId); - - // init workspace & register block schema - const workspace = new Workspace({ room: id }).register(BlockSchema); - - const Provider = this._providers.get(providerId); - assert(Provider); - - // initial configurator - const config = getKVConfigure(`workspace:${id}`); - // set workspace configs - const values = Object.entries(params.config || {}); - if (values.length) await config.setMany(values); - - // init data by provider - const provider = new Provider(); - await provider.init({ - apis: this._apis, - config, - debug: this._logger.enabled, - logger: this._logger.extend(`${Provider.id}:${id}`), - signals: this.signals, - workspace, - }); - await provider.initData(); - this._logger(`Workspace ${id} loaded`); - - return provider; - } - - async auth(providerId: string, globalConfig?: Record) { - const Provider = this._providers.get(providerId); - if (Provider) { - // initial configurator - const config = getKVConfigure(`provider:${providerId}`); - // set workspace configs - const values = Object.entries(globalConfig || {}); - if (values.length) await config.setMany(values); - - const logger = this._logger.extend(`auth:${providerId}`); - logger.enabled = this._logger.enabled; - await Provider.auth(config, logger, this.signals); - } - } - - /** - * load workspace data to memory - * @param workspaceId workspace id - * @param config.providerId provider id - * @param config.config provider config - * @returns Workspace instance - */ - async load( - workspaceId: string, - params: LoadConfig = {} - ): Promise { - if (workspaceId) { - if (!this._workspaces.has(workspaceId)) { - this._workspaces.set( - workspaceId, - this._getWorkspace(workspaceId, params) - ); - } - const workspace = this._workspaces.get(workspaceId); - assert(workspace); - return workspace.then(w => w.workspace); - } - return null; - } - - /** - * destroy workspace's instance in memory - * @param workspaceId workspace id - */ - async destroy(workspaceId: string) { - const provider = await this._workspaces.get(workspaceId); - if (provider) { - this._workspaces.delete(workspaceId); - await provider.destroy(); - } - } - - /** - * reload new workspace instance to memory to refresh config - * @param workspaceId workspace id - * @param config.providerId provider id - * @param config.config provider config - * @returns Workspace instance - */ - async reload( - workspaceId: string, - config: LoadConfig = {} - ): Promise { - await this.destroy(workspaceId); - return this.load(workspaceId, config); - } - - /** - * get workspace list,return a map of workspace id and data state - * data state is also map, the key is the provider id, and the data exists locally when the value is true, otherwise it does not exist - */ - async list(): Promise>> { - const listString = 'list:'; - const entries: [string, WorkspaceItem][] = await this._config.entries(); - return entries.reduce((acc, [k, i]) => { - if (k.startsWith(listString)) { - const key = k.slice(listString.length); - acc[key] = acc[key] || {}; - acc[key][i.provider] = i.locally; - } - return acc; - }, {} as Record>); - } - - /** - * delete local workspace's data - * @param workspaceId workspace id - */ - async delete(workspaceId: string) { - await this._config.delete(`${workspaceId}:provider`); - const provider = await this._workspaces.get(workspaceId); - if (provider) { - this._workspaces.delete(workspaceId); - // clear workspace data implement by provider - await provider.clear(); - } - } - - /** - * clear all local workspace's data - */ - async clear() { - const workspaces = await this.list(); - await Promise.all(Object.keys(workspaces).map(id => this.delete(id))); - } -} diff --git a/packages/data-center/src/datacenter.ts b/packages/data-center/src/datacenter.ts index d782c98bc..28de94043 100644 --- a/packages/data-center/src/datacenter.ts +++ b/packages/data-center/src/datacenter.ts @@ -5,7 +5,7 @@ import { LocalProvider } from './provider/local/local'; import { AffineProvider } from './provider'; import { Workspace as WS, WorkspaceMeta } from 'src/types'; import assert from 'assert'; -import { getLogger } from 'src'; +import { getLogger } from './logger'; import { BlockSchema } from '@blocksuite/blocks/models'; import { applyUpdate, encodeStateAsUpdate } from 'yjs'; @@ -231,7 +231,7 @@ export class DataCenter { const currentProvider = this.providerMap.get(workspaceInfo.provider); assert(currentProvider, 'Provider not found'); const newProvider = this.providerMap.get(provider); - assert(newProvider, 'AffineProvider is not registered'); + assert(newProvider, `${provider} provider is not registered`); const newWorkspace = await newProvider.createWorkspace({ name: workspaceInfo.name, avatar: workspaceInfo.avatar, diff --git a/packages/data-center/src/index.ts b/packages/data-center/src/index.ts index fb7e2f044..6fbb0e45d 100644 --- a/packages/data-center/src/index.ts +++ b/packages/data-center/src/index.ts @@ -1,5 +1,4 @@ -import debug from 'debug'; -import { DataCenter } from './dataCenter'; +import { DataCenter } from './datacenter'; const _initializeDataCenter = () => { let _dataCenterInstance: Promise; @@ -26,11 +25,6 @@ const _initializeDataCenter = () => { export const getDataCenter = _initializeDataCenter(); -export function getLogger(namespace: string) { - const logger = debug(namespace); - logger.log = console.log.bind(console); - return logger; -} - -export type { AccessTokenMessage } from './apis'; +export type { AccessTokenMessage } from './provider/affine/apis'; export type { Workspace } from './types'; +export { getLogger } from './logger'; diff --git a/packages/data-center/src/logger.ts b/packages/data-center/src/logger.ts new file mode 100644 index 000000000..83d68b932 --- /dev/null +++ b/packages/data-center/src/logger.ts @@ -0,0 +1,7 @@ +import debug from 'debug'; + +export function getLogger(namespace: string) { + const logger = debug(namespace); + logger.log = console.log.bind(console); + return logger; +} diff --git a/packages/data-center/src/provider/affine/affine.ts b/packages/data-center/src/provider/affine/affine.ts index b534565b8..c3ec22ca9 100644 --- a/packages/data-center/src/provider/affine/affine.ts +++ b/packages/data-center/src/provider/affine/affine.ts @@ -8,16 +8,16 @@ import { inviteMember, removeMember, createWorkspace, -} from 'src/apis/workspace'; +} from './apis/workspace'; import { BaseProvider } from '../base'; -import { User, Workspace as WS, WorkspaceMeta } from 'src/types'; +import { User, Workspace as WS, WorkspaceMeta } from '../../types'; import { Workspace } from '@blocksuite/store'; import { BlockSchema } from '@blocksuite/blocks/models'; import { applyUpdate } from 'yjs'; -import { token, Callback } from 'src/apis'; +import { token, Callback } from './apis'; import { varStorage as storage } from 'lib0/storage'; import assert from 'assert'; -import { getAuthorizer } from 'src/apis/token'; +import { getAuthorizer } from './apis/token'; import { WebsocketProvider } from './sync'; import { IndexedDBProvider } from '../indexeddb'; diff --git a/packages/data-center/src/apis/business.ts b/packages/data-center/src/provider/affine/apis/business.ts similarity index 99% rename from packages/data-center/src/apis/business.ts rename to packages/data-center/src/provider/affine/apis/business.ts index 2612fedd9..a43cac241 100644 --- a/packages/data-center/src/apis/business.ts +++ b/packages/data-center/src/provider/affine/apis/business.ts @@ -1,7 +1,7 @@ import { uuidv4 } from '@blocksuite/store'; import { getDataCenter } from 'src'; import { DataCenter } from 'src/datacenter'; -import { Workspace } from '../types'; +// import { Workspace } from '../types'; // export class Business { // private _dc: DataCenter | undefined; diff --git a/packages/data-center/src/apis/index.ts b/packages/data-center/src/provider/affine/apis/index.ts similarity index 100% rename from packages/data-center/src/apis/index.ts rename to packages/data-center/src/provider/affine/apis/index.ts diff --git a/packages/data-center/src/apis/request.ts b/packages/data-center/src/provider/affine/apis/request.ts similarity index 100% rename from packages/data-center/src/apis/request.ts rename to packages/data-center/src/provider/affine/apis/request.ts diff --git a/packages/data-center/src/apis/token.ts b/packages/data-center/src/provider/affine/apis/token.ts similarity index 98% rename from packages/data-center/src/apis/token.ts rename to packages/data-center/src/provider/affine/apis/token.ts index 3bd1dd65a..d7a2c34f8 100644 --- a/packages/data-center/src/apis/token.ts +++ b/packages/data-center/src/provider/affine/apis/token.ts @@ -2,8 +2,8 @@ import { initializeApp } from 'firebase/app'; import { getAuth, GoogleAuthProvider, signInWithPopup } from 'firebase/auth'; import type { User } from 'firebase/auth'; -import { getLogger } from '../index.js'; -import { bareClient } from './request.js'; +import { getLogger } from '../../../logger'; +import { bareClient } from './request'; export interface AccessTokenMessage { create_at: number; diff --git a/packages/data-center/src/apis/user.ts b/packages/data-center/src/provider/affine/apis/user.ts similarity index 100% rename from packages/data-center/src/apis/user.ts rename to packages/data-center/src/provider/affine/apis/user.ts diff --git a/packages/data-center/src/apis/workspace.ts b/packages/data-center/src/provider/affine/apis/workspace.ts similarity index 100% rename from packages/data-center/src/apis/workspace.ts rename to packages/data-center/src/provider/affine/apis/workspace.ts diff --git a/packages/data-center/src/provider_bk/affine/index.ts b/packages/data-center/src/provider_bk/affine/index.ts deleted file mode 100644 index aad66af8d..000000000 --- a/packages/data-center/src/provider_bk/affine/index.ts +++ /dev/null @@ -1,175 +0,0 @@ -import assert from 'assert'; -import { applyUpdate, Doc } from 'yjs'; - -import type { - ConfigStore, - DataCenterSignals, - InitialParams, - Logger, -} from '../index.js'; -import { token, Callback, getApis } from '../../apis/index.js'; -import { LocalProvider } from '../local/index.js'; - -import { WebsocketProvider } from './sync.js'; -import { IndexedDBProvider } from '../local/indexeddb.js'; - -export class AffineProvider extends LocalProvider { - static id = 'affine'; - private _onTokenRefresh?: Callback = undefined; - private _ws?: WebsocketProvider; - - constructor() { - super(); - } - - async init(params: InitialParams) { - super.init(params); - - this._onTokenRefresh = () => { - if (token.refresh) { - this._config.set('token', token.refresh); - } - }; - assert(this._onTokenRefresh); - - token.onChange(this._onTokenRefresh); - - // initial login token - if (token.isExpired) { - try { - const refreshToken = await this._config.get('token'); - await token.refreshToken(refreshToken); - - if (token.refresh) { - this._config.set('token', token.refresh); - } - - assert(token.isLogin); - } catch (_) { - this._logger('Authorization failed, fallback to local mode'); - } - } else { - this._config.set('token', token.refresh); - } - } - - async destroy() { - if (this._onTokenRefresh) { - token.offChange(this._onTokenRefresh); - } - this._ws?.disconnect(); - } - - async initData() { - const databases = await indexedDB.databases(); - await super.initData( - // set locally to true if exists a same name db - databases - .map(db => db.name) - .filter(v => v) - .includes(this._workspace.room) - ); - - const workspace = this._workspace; - const doc = workspace.doc; - - this._logger(`Login: ${token.isLogin}`); - - if (workspace.room && token.isLogin) { - try { - // init data from cloud - await AffineProvider._initCloudDoc( - workspace.room, - doc, - this._logger, - this._signals - ); - - // Wait for ws synchronization to complete, otherwise the data will be modified in reverse, which can be optimized later - this._ws = new WebsocketProvider('/', workspace.room, doc); - await new Promise((resolve, reject) => { - // TODO: synced will also be triggered on reconnection after losing sync - // There needs to be an event mechanism to emit the synchronization state to the upper layer - assert(this._ws); - this._ws.once('synced', () => resolve()); - this._ws.once('lost-connection', () => resolve()); - this._ws.once('connection-error', () => reject()); - }); - this._signals.listAdd.emit({ - workspace: workspace.room, - provider: this.id, - locally: true, - }); - } catch (e) { - this._logger('Failed to init cloud workspace', e); - } - } - - // if after update, the space:meta is empty - // then we need to get map with doc - // just a workaround for yjs - doc.getMap('space:meta'); - } - - private static async _initCloudDoc( - workspace: string, - doc: Doc, - logger: Logger, - signals: DataCenterSignals - ) { - const apis = getApis(); - logger(`Loading ${workspace}...`); - const updates = await apis.downloadWorkspace(workspace); - if (updates) { - await new Promise(resolve => { - doc.once('update', resolve); - applyUpdate(doc, new Uint8Array(updates)); - }); - logger(`Loaded: ${workspace}`); - - // only add to list as online workspace - signals.listAdd.emit({ - workspace, - provider: this.id, - // at this time we always download full workspace - // but after we support sub doc, we can only download metadata - locally: false, - }); - } - } - - static async auth( - config: Readonly>, - logger: Logger, - signals: DataCenterSignals - ) { - const refreshToken = await config.get('token'); - if (refreshToken) { - await token.refreshToken(refreshToken); - if (token.isLogin && !token.isExpired) { - logger('check login success'); - // login success - return; - } - } - - logger('start login'); - // login with google - const apis = getApis(); - assert(apis.signInWithGoogle); - const user = await apis.signInWithGoogle(); - assert(user); - logger(`login success: ${user.name}`); - - // TODO: refresh local workspace data - const workspaces = await apis.getWorkspaces(); - await Promise.all( - workspaces.map(async ({ id }) => { - const doc = new Doc(); - const idb = new IndexedDBProvider(id, doc); - await idb.whenSynced; - await this._initCloudDoc(id, doc, logger, signals); - }) - ); - } -} diff --git a/packages/data-center/src/provider_bk/affine/sync.js b/packages/data-center/src/provider_bk/affine/sync.js deleted file mode 100644 index 09004a51d..000000000 --- a/packages/data-center/src/provider_bk/affine/sync.js +++ /dev/null @@ -1,508 +0,0 @@ -/* eslint-disable no-undef */ -/** - * @module provider/websocket - */ - -/* eslint-env browser */ - -// import * as Y from 'yjs'; // eslint-disable-line -import * as bc from 'lib0/broadcastchannel'; -import * as time from 'lib0/time'; -import * as encoding from 'lib0/encoding'; -import * as decoding from 'lib0/decoding'; -import * as syncProtocol from 'y-protocols/sync'; -import * as authProtocol from 'y-protocols/auth'; -import * as awarenessProtocol from 'y-protocols/awareness'; -import { Observable } from 'lib0/observable'; -import * as math from 'lib0/math'; -import * as url from 'lib0/url'; - -export const messageSync = 0; -export const messageQueryAwareness = 3; -export const messageAwareness = 1; -export const messageAuth = 2; - -/** - * encoder, decoder, provider, emitSynced, messageType - * @type {Array} - */ -const messageHandlers = []; - -messageHandlers[messageSync] = ( - encoder, - decoder, - provider, - emitSynced, - _messageType -) => { - encoding.writeVarUint(encoder, messageSync); - const syncMessageType = syncProtocol.readSyncMessage( - decoder, - encoder, - provider.doc, - provider - ); - if ( - emitSynced && - syncMessageType === syncProtocol.messageYjsSyncStep2 && - !provider.synced - ) { - provider.synced = true; - } -}; - -messageHandlers[messageQueryAwareness] = ( - encoder, - _decoder, - provider, - _emitSynced, - _messageType -) => { - encoding.writeVarUint(encoder, messageAwareness); - encoding.writeVarUint8Array( - encoder, - awarenessProtocol.encodeAwarenessUpdate( - provider.awareness, - Array.from(provider.awareness.getStates().keys()) - ) - ); -}; - -messageHandlers[messageAwareness] = ( - _encoder, - decoder, - provider, - _emitSynced, - _messageType -) => { - awarenessProtocol.applyAwarenessUpdate( - provider.awareness, - decoding.readVarUint8Array(decoder), - provider - ); -}; - -messageHandlers[messageAuth] = ( - _encoder, - decoder, - provider, - _emitSynced, - _messageType -) => { - authProtocol.readAuthMessage(decoder, provider.doc, (_ydoc, reason) => - permissionDeniedHandler(provider, reason) - ); -}; - -// @todo - this should depend on awareness.outdatedTime -const messageReconnectTimeout = 30000; - -/** - * @param {WebsocketProvider} provider - * @param {string} reason - */ -const permissionDeniedHandler = (provider, reason) => - console.warn(`Permission denied to access ${provider.url}.\n${reason}`); - -/** - * @param {WebsocketProvider} provider - * @param {Uint8Array} buf - * @param {boolean} emitSynced - * @return {encoding.Encoder} - */ -const readMessage = (provider, buf, emitSynced) => { - const decoder = decoding.createDecoder(buf); - const encoder = encoding.createEncoder(); - const messageType = decoding.readVarUint(decoder); - const messageHandler = provider.messageHandlers[messageType]; - if (/** @type {any} */ (messageHandler)) { - messageHandler(encoder, decoder, provider, emitSynced, messageType); - } else { - console.error('Unable to compute message'); - } - return encoder; -}; - -/** - * @param {WebsocketProvider} provider - */ -const setupWS = provider => { - if (provider.shouldConnect && provider.ws === null) { - const websocket = new provider._WS(provider.url); - websocket.binaryType = 'arraybuffer'; - provider.ws = websocket; - provider.wsconnecting = true; - provider.wsconnected = false; - provider.synced = false; - - websocket.onmessage = event => { - provider.wsLastMessageReceived = time.getUnixTime(); - const encoder = readMessage(provider, new Uint8Array(event.data), true); - if (encoding.length(encoder) > 1) { - websocket.send(encoding.toUint8Array(encoder)); - } - }; - websocket.onerror = event => { - provider.emit('connection-error', [event, provider]); - }; - websocket.onclose = event => { - provider.emit('connection-close', [event, provider]); - provider.ws = null; - provider.wsconnecting = false; - if (provider.wsconnected) { - provider.wsconnected = false; - provider.synced = false; - // update awareness (all users except local left) - awarenessProtocol.removeAwarenessStates( - provider.awareness, - Array.from(provider.awareness.getStates().keys()).filter( - client => client !== provider.doc.clientID - ), - provider - ); - provider.emit('status', [ - { - status: 'disconnected', - }, - ]); - } else { - provider.wsUnsuccessfulReconnects++; - } - // Start with no reconnect timeout and increase timeout by - // using exponential backoff starting with 100ms - setTimeout( - setupWS, - math.min( - math.pow(2, provider.wsUnsuccessfulReconnects) * 100, - provider.maxBackoffTime - ), - provider - ); - }; - websocket.onopen = () => { - provider.wsLastMessageReceived = time.getUnixTime(); - provider.wsconnecting = false; - provider.wsconnected = true; - provider.wsUnsuccessfulReconnects = 0; - provider.emit('status', [ - { - status: 'connected', - }, - ]); - // always send sync step 1 when connected - const encoder = encoding.createEncoder(); - encoding.writeVarUint(encoder, messageSync); - syncProtocol.writeSyncStep1(encoder, provider.doc); - websocket.send(encoding.toUint8Array(encoder)); - // broadcast local awareness state - if (provider.awareness.getLocalState() !== null) { - const encoderAwarenessState = encoding.createEncoder(); - encoding.writeVarUint(encoderAwarenessState, messageAwareness); - encoding.writeVarUint8Array( - encoderAwarenessState, - awarenessProtocol.encodeAwarenessUpdate(provider.awareness, [ - provider.doc.clientID, - ]) - ); - websocket.send(encoding.toUint8Array(encoderAwarenessState)); - } - }; - - provider.emit('status', [ - { - status: 'connecting', - }, - ]); - } -}; - -/** - * @param {WebsocketProvider} provider - * @param {ArrayBuffer} buf - */ -const broadcastMessage = (provider, buf) => { - if (provider.wsconnected) { - /** @type {WebSocket} */ (provider.ws).send(buf); - } - if (provider.bcconnected) { - bc.publish(provider.bcChannel, buf, provider); - } -}; - -/** - * Websocket Provider for Yjs. Creates a websocket connection to sync the shared document. - * The document name is attached to the provided url. I.e. the following example - * creates a websocket connection to http://localhost:1234/my-document-name - * - * @example - * import * as Y from 'yjs' - * import { WebsocketProvider } from 'y-websocket' - * const doc = new Y.Doc() - * const provider = new WebsocketProvider('http://localhost:1234', 'my-document-name', doc) - * - * @extends {Observable} - */ -export class WebsocketProvider extends Observable { - /** - * @param {string} serverUrl - * @param {string} roomname - * @param {Y.Doc} doc - * @param {object} [opts] - * @param {boolean} [opts.connect] - * @param {awarenessProtocol.Awareness} [opts.awareness] - * @param {Object} [opts.params] - * @param {typeof WebSocket} [opts.WebSocketPolyfill] Optionall provide a WebSocket polyfill - * @param {number} [opts.resyncInterval] Request server state every `resyncInterval` milliseconds - * @param {number} [opts.maxBackoffTime] Maximum amount of time to wait before trying to reconnect (we try to reconnect using exponential backoff) - * @param {boolean} [opts.disableBc] Disable cross-tab BroadcastChannel communication - */ - constructor( - serverUrl, - roomname, - doc, - { - connect = true, - awareness = new awarenessProtocol.Awareness(doc), - params = {}, - WebSocketPolyfill = WebSocket, - resyncInterval = -1, - maxBackoffTime = 2500, - disableBc = false, - } = {} - ) { - super(); - // ensure that url is always ends with / - while (serverUrl[serverUrl.length - 1] === '/') { - serverUrl = serverUrl.slice(0, serverUrl.length - 1); - } - const encodedParams = url.encodeQueryParams(params); - this.maxBackoffTime = maxBackoffTime; - this.bcChannel = serverUrl + '/' + roomname; - this.url = - serverUrl + - '/' + - roomname + - (encodedParams.length === 0 ? '' : '?' + encodedParams); - this.roomname = roomname; - this.doc = doc; - this._WS = WebSocketPolyfill; - this.awareness = awareness; - this.wsconnected = false; - this.wsconnecting = false; - this.bcconnected = false; - this.disableBc = disableBc; - this.wsUnsuccessfulReconnects = 0; - this.messageHandlers = messageHandlers.slice(); - /** - * @type {boolean} - */ - this._synced = false; - /** - * @type {WebSocket?} - */ - this.ws = null; - this.wsLastMessageReceived = 0; - /** - * Whether to connect to other peers or not - * @type {boolean} - */ - this.shouldConnect = connect; - - /** - * @type {number} - */ - this._resyncInterval = 0; - if (resyncInterval > 0) { - this._resyncInterval = /** @type {any} */ ( - setInterval(() => { - if (this.ws && this.ws.readyState === WebSocket.OPEN) { - // resend sync step 1 - const encoder = encoding.createEncoder(); - encoding.writeVarUint(encoder, messageSync); - syncProtocol.writeSyncStep1(encoder, doc); - this.ws.send(encoding.toUint8Array(encoder)); - } - }, resyncInterval) - ); - } - - /** - * @param {ArrayBuffer} data - * @param {any} origin - */ - this._bcSubscriber = (data, origin) => { - if (origin !== this) { - const encoder = readMessage(this, new Uint8Array(data), false); - if (encoding.length(encoder) > 1) { - bc.publish(this.bcChannel, encoding.toUint8Array(encoder), this); - } - } - }; - /** - * Listens to Yjs updates and sends them to remote peers (ws and broadcastchannel) - * @param {Uint8Array} update - * @param {any} origin - */ - this._updateHandler = (update, origin) => { - if (origin !== this) { - const encoder = encoding.createEncoder(); - encoding.writeVarUint(encoder, messageSync); - syncProtocol.writeUpdate(encoder, update); - broadcastMessage(this, encoding.toUint8Array(encoder)); - } - }; - this.doc.on('update', this._updateHandler); - /** - * @param {any} changed - * @param {any} _origin - */ - this._awarenessUpdateHandler = ({ added, updated, removed }, _origin) => { - const changedClients = added.concat(updated).concat(removed); - const encoder = encoding.createEncoder(); - encoding.writeVarUint(encoder, messageAwareness); - encoding.writeVarUint8Array( - encoder, - awarenessProtocol.encodeAwarenessUpdate(awareness, changedClients) - ); - broadcastMessage(this, encoding.toUint8Array(encoder)); - }; - this._unloadHandler = () => { - awarenessProtocol.removeAwarenessStates( - this.awareness, - [doc.clientID], - 'window unload' - ); - }; - if (typeof window !== 'undefined') { - window.addEventListener('unload', this._unloadHandler); - } else if (typeof process !== 'undefined') { - process.on('exit', this._unloadHandler); - } - awareness.on('update', this._awarenessUpdateHandler); - this._checkInterval = /** @type {any} */ ( - setInterval(() => { - if ( - this.wsconnected && - messageReconnectTimeout < - time.getUnixTime() - this.wsLastMessageReceived - ) { - // no message received in a long time - not even your own awareness - // updates (which are updated every 15 seconds) - /** @type {WebSocket} */ (this.ws).close(); - } - }, messageReconnectTimeout / 10) - ); - if (connect) { - this.connect(); - } - } - - /** - * @type {boolean} - */ - get synced() { - return this._synced; - } - - set synced(state) { - if (this._synced !== state) { - this._synced = state; - this.emit('synced', [state]); - this.emit('sync', [state]); - } - } - - destroy() { - if (this._resyncInterval !== 0) { - clearInterval(this._resyncInterval); - } - clearInterval(this._checkInterval); - this.disconnect(); - if (typeof window !== 'undefined') { - window.removeEventListener('unload', this._unloadHandler); - } else if (typeof process !== 'undefined') { - process.off('exit', this._unloadHandler); - } - this.awareness.off('update', this._awarenessUpdateHandler); - this.doc.off('update', this._updateHandler); - super.destroy(); - } - - connectBc() { - if (this.disableBc) { - return; - } - if (!this.bcconnected) { - bc.subscribe(this.bcChannel, this._bcSubscriber); - this.bcconnected = true; - } - // send sync step1 to bc - // write sync step 1 - const encoderSync = encoding.createEncoder(); - encoding.writeVarUint(encoderSync, messageSync); - syncProtocol.writeSyncStep1(encoderSync, this.doc); - bc.publish(this.bcChannel, encoding.toUint8Array(encoderSync), this); - // broadcast local state - const encoderState = encoding.createEncoder(); - encoding.writeVarUint(encoderState, messageSync); - syncProtocol.writeSyncStep2(encoderState, this.doc); - bc.publish(this.bcChannel, encoding.toUint8Array(encoderState), this); - // write queryAwareness - const encoderAwarenessQuery = encoding.createEncoder(); - encoding.writeVarUint(encoderAwarenessQuery, messageQueryAwareness); - bc.publish( - this.bcChannel, - encoding.toUint8Array(encoderAwarenessQuery), - this - ); - // broadcast local awareness state - const encoderAwarenessState = encoding.createEncoder(); - encoding.writeVarUint(encoderAwarenessState, messageAwareness); - encoding.writeVarUint8Array( - encoderAwarenessState, - awarenessProtocol.encodeAwarenessUpdate(this.awareness, [ - this.doc.clientID, - ]) - ); - bc.publish( - this.bcChannel, - encoding.toUint8Array(encoderAwarenessState), - this - ); - } - - disconnectBc() { - // broadcast message with local awareness state set to null (indicating disconnect) - const encoder = encoding.createEncoder(); - encoding.writeVarUint(encoder, messageAwareness); - encoding.writeVarUint8Array( - encoder, - awarenessProtocol.encodeAwarenessUpdate( - this.awareness, - [this.doc.clientID], - new Map() - ) - ); - broadcastMessage(this, encoding.toUint8Array(encoder)); - if (this.bcconnected) { - bc.unsubscribe(this.bcChannel, this._bcSubscriber); - this.bcconnected = false; - } - } - - disconnect() { - this.shouldConnect = false; - this.disconnectBc(); - if (this.ws !== null) { - this.ws.close(); - } - } - - connect() { - this.shouldConnect = true; - if (!this.wsconnected && this.ws === null) { - setupWS(this); - this.connectBc(); - } - } -} diff --git a/packages/data-center/src/provider_bk/base.ts b/packages/data-center/src/provider_bk/base.ts deleted file mode 100644 index bd8d8a4e0..000000000 --- a/packages/data-center/src/provider_bk/base.ts +++ /dev/null @@ -1,79 +0,0 @@ -/* eslint-disable @typescript-eslint/no-unused-vars */ -import type { Workspace } from '@blocksuite/store'; - -import type { - Apis, - DataCenterSignals, - Logger, - InitialParams, - ConfigStore, -} from './index'; - -export class BaseProvider { - static id = 'base'; - protected _apis!: Readonly; - protected _config!: Readonly; - protected _logger!: Logger; - protected _signals!: DataCenterSignals; - protected _workspace!: Workspace; - - constructor() { - // Nothing to do here - } - - get id(): string { - return (this.constructor as any).id; - } - - async init(params: InitialParams) { - this._apis = params.apis; - this._config = params.config; - this._logger = params.logger; - this._signals = params.signals; - this._workspace = params.workspace; - this._logger.enabled = params.debug; - } - - async clear() { - await this.destroy(); - await this._config.clear(); - } - - async destroy() { - // Nothing to do here - } - - async initData() { - throw Error('Not implemented: initData'); - } - - // should return a blob url - async getBlob(_id: string): Promise { - throw Error('Not implemented: getBlob'); - } - - // should return a blob unique id - async setBlob(_blob: Blob): Promise { - throw Error('Not implemented: setBlob'); - } - - get workspace() { - return this._workspace; - } - - static async auth( - _config: Readonly, - logger: Logger, - _signals: DataCenterSignals - ) { - logger("This provider doesn't require authentication"); - } - - // get workspace list,return a map of workspace id and boolean - // if value is true, it exists locally, otherwise it does not exist locally - static async list( - _config: Readonly - ): Promise | undefined> { - throw Error('Not implemented: list'); - } -} diff --git a/packages/data-center/src/provider_bk/baseProvider.ts b/packages/data-center/src/provider_bk/baseProvider.ts deleted file mode 100644 index e69de29bb..000000000 diff --git a/packages/data-center/src/provider_bk/index.ts b/packages/data-center/src/provider_bk/index.ts deleted file mode 100644 index e80bdf675..000000000 --- a/packages/data-center/src/provider_bk/index.ts +++ /dev/null @@ -1,22 +0,0 @@ -import type { Workspace } from '@blocksuite/store'; - -import type { Apis } from '../apis'; -// import type { DataCenterSignals } from '../datacenter'; -import type { getLogger } from '../index'; -import type { ConfigStore } from '../store'; - -export type Logger = ReturnType; - -// export type InitialParams = { -// apis: Apis; -// config: Readonly; -// debug: boolean; -// logger: Logger; -// signals: DataCenterSignals; -// workspace: Workspace; -// }; - -// export type { Apis, ConfigStore, DataCenterSignals, Workspace }; -export type { BaseProvider } from './base.js'; -export { AffineProvider } from './affine/index.js'; -export { LocalProvider } from './local/index.js'; diff --git a/packages/data-center/src/provider_bk/local/index.ts b/packages/data-center/src/provider_bk/local/index.ts deleted file mode 100644 index 24bf9584a..000000000 --- a/packages/data-center/src/provider_bk/local/index.ts +++ /dev/null @@ -1,73 +0,0 @@ -import type { BlobStorage } from '@blocksuite/store'; -import assert from 'assert'; - -import type { ConfigStore, InitialParams } from '../index.js'; -import { BaseProvider } from '../base.js'; -import { IndexedDBProvider } from './indexeddb.js'; - -export class LocalProvider extends BaseProvider { - static id = 'local'; - private _blobs!: BlobStorage; - private _idb?: IndexedDBProvider = undefined; - - constructor() { - super(); - } - async init(params: InitialParams) { - super.init(params); - - const blobs = await this._workspace.blobs; - assert(blobs); - this._blobs = blobs; - } - - async initData(locally = true) { - assert(this._workspace.room); - this._logger('Loading local data'); - this._idb = new IndexedDBProvider( - this._workspace.room, - this._workspace.doc - ); - - await this._idb.whenSynced; - this._logger('Local data loaded'); - - this._signals.listAdd.emit({ - workspace: this._workspace.room, - provider: this.id, - locally, - }); - } - - async clear() { - assert(this._workspace.room); - await super.clear(); - await this._blobs.clear(); - await this._idb?.clearData(); - this._signals.listRemove.emit(this._workspace.room); - } - - async destroy(): Promise { - super.destroy(); - await this._idb?.destroy(); - } - - async getBlob(id: string): Promise { - return this._blobs.get(id); - } - - async setBlob(blob: Blob): Promise { - return this._blobs.set(blob); - } - - static async list( - config: Readonly> - ): Promise | undefined> { - const entries = await config.entries(); - return new Map( - entries - .filter(([key]) => key.startsWith('list:')) - .map(([key, value]) => [key.slice(5), value]) - ); - } -} diff --git a/packages/data-center/src/provider_bk/local/indexeddb.ts b/packages/data-center/src/provider_bk/local/indexeddb.ts deleted file mode 100644 index db7fab29d..000000000 --- a/packages/data-center/src/provider_bk/local/indexeddb.ts +++ /dev/null @@ -1,199 +0,0 @@ -/* eslint-disable @typescript-eslint/no-explicit-any */ -import * as idb from 'lib0/indexeddb.js'; -import { Observable } from 'lib0/observable.js'; -import type { Doc } from 'yjs'; -import { applyUpdate, encodeStateAsUpdate, transact } from 'yjs'; - -const customStoreName = 'custom'; -const updatesStoreName = 'updates'; - -const PREFERRED_TRIM_SIZE = 500; - -const fetchUpdates = async (provider: IndexedDBProvider) => { - const [updatesStore] = idb.transact(provider.db as IDBDatabase, [ - updatesStoreName, - ]); // , 'readonly') - if (updatesStore) { - const updates = await idb.getAll( - updatesStore, - idb.createIDBKeyRangeLowerBound(provider._dbref, false) - ); - transact( - provider.doc, - () => { - updates.forEach(val => applyUpdate(provider.doc, val)); - }, - provider, - false - ); - const lastKey = await idb.getLastKey(updatesStore); - provider._dbref = lastKey + 1; - const cnt = await idb.count(updatesStore); - provider._dbsize = cnt; - } - return updatesStore; -}; - -const storeState = (provider: IndexedDBProvider, forceStore = true) => - fetchUpdates(provider).then(updatesStore => { - if ( - updatesStore && - (forceStore || provider._dbsize >= PREFERRED_TRIM_SIZE) - ) { - idb - .addAutoKey(updatesStore, encodeStateAsUpdate(provider.doc)) - .then(() => - idb.del( - updatesStore, - idb.createIDBKeyRangeUpperBound(provider._dbref, true) - ) - ) - .then(() => - idb.count(updatesStore).then(cnt => { - provider._dbsize = cnt; - }) - ); - } - }); - -export class IndexedDBProvider extends Observable { - doc: Doc; - name: string; - _dbref: number; - _dbsize: number; - private _destroyed: boolean; - whenSynced: Promise; - db: IDBDatabase | null; - private _db: Promise; - private _storeTimeout: number; - private _storeTimeoutId: NodeJS.Timeout | null; - private _storeUpdate: (update: Uint8Array, origin: any) => void; - - constructor(name: string, doc: Doc) { - super(); - this.doc = doc; - this.name = name; - this._dbref = 0; - this._dbsize = 0; - this._destroyed = false; - this.db = null; - this._db = idb.openDB(name, db => - idb.createStores(db, [['updates', { autoIncrement: true }], ['custom']]) - ); - - this.whenSynced = this._db.then(async db => { - this.db = db; - const currState = encodeStateAsUpdate(doc); - const updatesStore = await fetchUpdates(this); - if (updatesStore) { - await idb.addAutoKey(updatesStore, currState); - } - if (this._destroyed) { - return this; - } - this.emit('synced', [this]); - return this; - }); - - // Timeout in ms untill data is merged and persisted in idb. - this._storeTimeout = 1000; - - this._storeTimeoutId = null; - - this._storeUpdate = (update: Uint8Array, origin: any) => { - if (this.db && origin !== this) { - const [updatesStore] = idb.transact( - /** @type {IDBDatabase} */ this.db, - [updatesStoreName] - ); - if (updatesStore) { - idb.addAutoKey(updatesStore, update); - } - if (++this._dbsize >= PREFERRED_TRIM_SIZE) { - // debounce store call - if (this._storeTimeoutId !== null) { - clearTimeout(this._storeTimeoutId); - } - this._storeTimeoutId = setTimeout(() => { - storeState(this, false); - this._storeTimeoutId = null; - }, this._storeTimeout); - } - } - }; - doc.on('update', this._storeUpdate); - this.destroy = this.destroy.bind(this); - doc.on('destroy', this.destroy); - } - - override destroy() { - if (this._storeTimeoutId) { - clearTimeout(this._storeTimeoutId); - } - this.doc.off('update', this._storeUpdate); - this.doc.off('destroy', this.destroy); - this._destroyed = true; - return this._db.then(db => { - db.close(); - }); - } - - /** - * Destroys this instance and removes all data from indexeddb. - * - * @return {Promise} - */ - async clearData(): Promise { - return this.destroy().then(() => { - idb.deleteDB(this.name); - }); - } - - /** - * @param {String | number | ArrayBuffer | Date} key - * @return {Promise} - */ - async get( - key: string | number | ArrayBuffer | Date - ): Promise { - return this._db.then(db => { - const [custom] = idb.transact(db, [customStoreName], 'readonly'); - if (custom) { - return idb.get(custom, key); - } - return undefined; - }); - } - - /** - * @param {String | number | ArrayBuffer | Date} key - * @param {String | number | ArrayBuffer | Date} value - * @return {Promise} - */ - async set( - key: string | number | ArrayBuffer | Date, - value: string | number | ArrayBuffer | Date - ): Promise { - return this._db.then(db => { - const [custom] = idb.transact(db, [customStoreName]); - if (custom) { - return idb.put(custom, value, key); - } - return undefined; - }); - } - - /** - * @param {String | number | ArrayBuffer | Date} key - * @return {Promise} - */ - async del(key: string | number | ArrayBuffer | Date): Promise { - return this._db.then(db => { - const [custom] = idb.transact(db, [customStoreName]); - if (custom) { - return idb.del(custom, key); - } - return undefined; - }); - } -} diff --git a/packages/data-center/src/workspaces/workspaces.ts b/packages/data-center/src/workspaces/workspaces.ts index f6abafb87..b8ee7f5df 100644 --- a/packages/data-center/src/workspaces/workspaces.ts +++ b/packages/data-center/src/workspaces/workspaces.ts @@ -1,8 +1,8 @@ -import { Workspace as WS } from 'src/types'; +import { Workspace as WS } from '../types'; import { Observable } from 'lib0/observable'; import { uuidv4 } from '@blocksuite/store'; -import { DataCenter } from 'src/dataCenterNew'; +import { DataCenter } from '../datacenter'; export class Workspaces extends Observable { private _workspaces: WS[];