refactor: remote data provider

This commit is contained in:
DarkSky
2022-08-10 22:10:34 +08:00
parent a5a8b32a25
commit 89191290e4
5 changed files with 137 additions and 111 deletions

View File

@@ -134,7 +134,7 @@ export class IndexedDBProvider extends Observable<string> {
} }
/** /**
* Destroys this instance and removes all data from SQLite. * Destroys this instance and removes all data from indexeddb.
* *
* @return {Promise<void>} * @return {Promise<void>}
*/ */

View File

@@ -41,6 +41,7 @@ const initSQLiteInstance = async () => {
_sqliteProcessing = true; _sqliteProcessing = true;
_sqliteInstance = await sqlite({ _sqliteInstance = await sqlite({
locateFile: () => locateFile: () =>
// @ts-ignore
new URL('sql.js/dist/sql-wasm.wasm', import.meta.url).href, new URL('sql.js/dist/sql-wasm.wasm', import.meta.url).href,
}); });
_sqliteProcessing = false; _sqliteProcessing = false;

View File

@@ -18,11 +18,7 @@ import {
snapshot, snapshot,
} from 'yjs'; } from 'yjs';
import { import { IndexedDBProvider } from '@toeverything/datasource/jwt-rpc';
IndexedDBProvider,
SQLiteProvider,
WebsocketProvider,
} from '@toeverything/datasource/jwt-rpc';
import { import {
AsyncDatabaseAdapter, AsyncDatabaseAdapter,
@@ -31,7 +27,7 @@ import {
Connectivity, Connectivity,
HistoryManager, HistoryManager,
} from '../../adapter'; } from '../../adapter';
import { BucketBackend, BlockItem, BlockTypes } from '../../types'; import { BlockItem, BlockTypes } from '../../types';
import { getLogger, sha3, sleep } from '../../utils'; import { getLogger, sha3, sleep } from '../../utils';
import { YjsRemoteBinaries } from './binary'; import { YjsRemoteBinaries } from './binary';
@@ -43,51 +39,26 @@ import {
} from './operation'; } from './operation';
import { EmitEvents, Suspend } from './listener'; import { EmitEvents, Suspend } from './listener';
import { YjsHistoryManager } from './history'; import { YjsHistoryManager } from './history';
import { YjsProvider } from './provider';
declare const JWT_DEV: boolean; declare const JWT_DEV: boolean;
const logger = getLogger('BlockDB:yjs'); const logger = getLogger('BlockDB:yjs');
type ConnectivityListener = (
workspace: string,
connectivity: Connectivity
) => void;
type YjsProviders = { type YjsProviders = {
awareness: Awareness; awareness: Awareness;
idb: IndexedDBProvider; idb: IndexedDBProvider;
binariesIdb: IndexedDBProvider; binariesIdb: IndexedDBProvider;
fstore?: SQLiteProvider;
ws?: WebsocketProvider;
backend: string;
gatekeeper: GateKeeper; gatekeeper: GateKeeper;
connListener: { listeners?: ConnectivityListener };
userId: string; userId: string;
remoteToken?: string; // remote storage token remoteToken?: string; // remote storage token
}; };
const _yjsDatabaseInstance = new Map<string, YjsProviders>(); const _yjsDatabaseInstance = new Map<string, YjsProviders>();
async function _initWebsocketProvider(
url: string,
room: string,
doc: Doc,
token?: string,
params?: YjsInitOptions['params']
): Promise<[Awareness, WebsocketProvider | undefined]> {
const awareness = new Awareness(doc);
if (token) {
const ws = new WebsocketProvider(token, url, room, doc, {
awareness,
params,
}) as any; // TODO: type is erased after cascading references
// Wait for ws synchronization to complete, otherwise the data will be modified in reverse, which can be optimized later
return 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
ws.once('synced', () => resolve([awareness, ws]));
ws.once('lost-connection', () => resolve([awareness, ws]));
ws.once('connection-error', () => reject());
});
} else {
return [awareness, undefined];
}
}
const _asyncInitLoading = new Set<string>(); const _asyncInitLoading = new Set<string>();
const _waitLoading = async (workspace: string) => { const _waitLoading = async (workspace: string) => {
while (_asyncInitLoading.has(workspace)) { while (_asyncInitLoading.has(workspace)) {
@@ -96,14 +67,11 @@ const _waitLoading = async (workspace: string) => {
}; };
async function _initYjsDatabase( async function _initYjsDatabase(
backend: string,
workspace: string, workspace: string,
options: { options: {
params: YjsInitOptions['params'];
userId: string; userId: string;
token?: string; token?: string;
importData?: Uint8Array; provider?: Record<string, YjsProvider>;
exportData?: (binary: Uint8Array) => void;
} }
): Promise<YjsProviders> { ): Promise<YjsProviders> {
if (_asyncInitLoading.has(workspace)) { if (_asyncInitLoading.has(workspace)) {
@@ -119,28 +87,10 @@ async function _initYjsDatabase(
} }
// if (instance) return instance; // if (instance) return instance;
_asyncInitLoading.add(workspace); _asyncInitLoading.add(workspace);
const { params, userId, token: remoteToken } = options; const { userId, token } = options;
const doc = new Doc({ autoLoad: true, shouldLoad: true }); const doc = new Doc({ autoLoad: true, shouldLoad: true });
const idb = await new IndexedDBProvider(workspace, doc).whenSynced;
const idbp = new IndexedDBProvider(workspace, doc).whenSynced;
const fs = new SQLiteProvider(workspace, doc, options.importData);
if (options.exportData) fs.registerExporter(options.exportData);
const wsp = _initWebsocketProvider(
backend,
workspace,
doc,
remoteToken,
params
);
const [idb, [awareness, ws], fstore] = await Promise.all([
idbp,
wsp,
fs.whenSynced,
]);
const binaries = new Doc({ autoLoad: true, shouldLoad: true }); const binaries = new Doc({ autoLoad: true, shouldLoad: true });
const binariesIdb = await new IndexedDBProvider( const binariesIdb = await new IndexedDBProvider(
@@ -148,6 +98,8 @@ async function _initYjsDatabase(
binaries binaries
).whenSynced; ).whenSynced;
const awareness = new Awareness(doc);
const gateKeeperData = doc.getMap<YMap<string>>('gatekeeper'); const gateKeeperData = doc.getMap<YMap<string>>('gatekeeper');
const gatekeeper = new GateKeeper( const gatekeeper = new GateKeeper(
@@ -157,44 +109,45 @@ async function _initYjsDatabase(
gateKeeperData.get('common') || gateKeeperData.set('common', new YMap()) gateKeeperData.get('common') || gateKeeperData.set('common', new YMap())
); );
_yjsDatabaseInstance.set(workspace, { const connListener: { listeners?: ConnectivityListener } = {};
if (options.provider) {
const emitState = (c: Connectivity) =>
connListener.listeners?.(workspace, c);
await Promise.all(
Object.entries(options.provider).map(async ([, p]) =>
p({ awareness, doc, token, workspace, emitState })
)
);
}
const newInstance = {
awareness, awareness,
idb, idb,
binariesIdb, binariesIdb,
fstore,
ws,
backend,
gatekeeper, gatekeeper,
connListener,
userId, userId,
remoteToken, remoteToken: token,
}); };
_yjsDatabaseInstance.set(workspace, newInstance);
_asyncInitLoading.delete(workspace); _asyncInitLoading.delete(workspace);
return { return newInstance;
awareness,
idb,
binariesIdb,
fstore,
ws,
backend,
gatekeeper,
userId,
remoteToken,
};
} }
export type { YjsBlockInstance } from './block'; export type { YjsBlockInstance } from './block';
export type { YjsContentOperation } from './operation'; export type { YjsContentOperation } from './operation';
export type YjsInitOptions = { export type YjsInitOptions = {
backend: typeof BucketBackend[keyof typeof BucketBackend];
params?: Record<string, string>;
userId?: string; userId?: string;
token?: string; token?: string;
importData?: Uint8Array; provider?: Record<string, YjsProvider>;
exportData?: (binary: Uint8Array) => void;
}; };
export { getYjsProviders } from './provider';
export type { YjsProviderOptions } from './provider';
export class YjsAdapter implements AsyncDatabaseAdapter<YjsContentOperation> { export class YjsAdapter implements AsyncDatabaseAdapter<YjsContentOperation> {
private readonly _provider: YjsProviders; private readonly _provider: YjsProviders;
private readonly _doc: Doc; // doc instance private readonly _doc: Doc; // doc instance
@@ -217,20 +170,11 @@ export class YjsAdapter implements AsyncDatabaseAdapter<YjsContentOperation> {
workspace: string, workspace: string,
options: YjsInitOptions options: YjsInitOptions
): Promise<YjsAdapter> { ): Promise<YjsAdapter> {
const { const { userId = 'default', token, provider } = options;
backend, const providers = await _initYjsDatabase(workspace, {
params = {},
userId = 'default',
token,
importData,
exportData,
} = options;
const providers = await _initYjsDatabase(backend, workspace, {
params,
userId, userId,
token, token,
importData, provider,
exportData,
}); });
return new YjsAdapter(providers); return new YjsAdapter(providers);
} }
@@ -255,18 +199,14 @@ export class YjsAdapter implements AsyncDatabaseAdapter<YjsContentOperation> {
this._listener = new Map(); this._listener = new Map();
const ws = providers.ws as any; providers.connListener.listeners = (
if (ws) { workspace: string,
const workspace = providers.idb.name; connectivity: Connectivity
const emitState = (connectivity: Connectivity) => { ) => {
this._listener.get('connectivity')?.( this._listener.get('connectivity')?.(
new Map([[workspace, connectivity]]) new Map([[workspace, connectivity]])
); );
}; };
ws.on('synced', () => emitState('connected'));
ws.on('lost-connection', () => emitState('retry'));
ws.on('connection-error', () => emitState('retry'));
}
const debounced_editing_notifier = debounce( const debounced_editing_notifier = debounce(
() => { () => {

View File

@@ -0,0 +1,75 @@
import { Doc } from 'yjs';
import { Awareness } from 'y-protocols/awareness.js';
import {
SQLiteProvider,
WebsocketProvider,
} from '@toeverything/datasource/jwt-rpc';
import { Connectivity } from '../../adapter';
import { BucketBackend } from '../../types';
type YjsDefaultInstances = {
awareness: Awareness;
doc: Doc;
token?: string;
workspace: string;
emitState: (connectivity: Connectivity) => void;
};
export type YjsProvider = (instances: YjsDefaultInstances) => Promise<void>;
export type YjsProviderOptions = {
backend: typeof BucketBackend[keyof typeof BucketBackend];
params?: Record<string, string>;
importData?: Uint8Array;
exportData?: (binary: Uint8Array) => void;
};
export const getYjsProviders = (
options: YjsProviderOptions
): Record<string, YjsProvider> => {
return {
sqlite: async (instances: YjsDefaultInstances) => {
const fs = new SQLiteProvider(
instances.workspace,
instances.doc,
options.importData
);
if (options.exportData) fs.registerExporter(options.exportData);
await fs.whenSynced;
},
ws: async (instances: YjsDefaultInstances) => {
if (instances.token) {
const ws = new WebsocketProvider(
instances.token,
options.backend,
instances.workspace,
instances.doc,
{
awareness: instances.awareness,
params: options.params,
}
) as any; // TODO: type is erased after cascading references
// Wait for ws synchronization to complete, otherwise the data will be modified in reverse, which can be optimized later
return new Promise<void>((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
ws.once('synced', () => resolve());
ws.once('lost-connection', () => resolve());
ws.once('connection-error', () => reject());
ws.on('synced', () => instances.emitState('connected'));
ws.on('lost-connection', () =>
instances.emitState('retry')
);
ws.on('connection-error', () =>
instances.emitState('retry')
);
});
} else {
return;
}
},
};
};

View File

@@ -15,7 +15,11 @@ import {
ContentTypes, ContentTypes,
Connectivity, Connectivity,
} from './adapter'; } from './adapter';
import { YjsBlockInstance } from './adapter/yjs'; import {
getYjsProviders,
YjsBlockInstance,
YjsProviderOptions,
} from './adapter/yjs';
import { import {
BaseBlock, BaseBlock,
BlockIndexer, BlockIndexer,
@@ -27,11 +31,11 @@ import {
BlockTypes, BlockTypes,
BlockTypeKeys, BlockTypeKeys,
BlockFlavors, BlockFlavors,
BucketBackend,
UUID, UUID,
BlockFlavorKeys, BlockFlavorKeys,
BlockItem, BlockItem,
ExcludeFunction, ExcludeFunction,
BucketBackend,
} from './types'; } from './types';
import { BlockEventBus, genUUID, getLogger } from './utils'; import { BlockEventBus, genUUID, getLogger } from './utils';
@@ -588,10 +592,16 @@ export class BlockClient<
public static async init( public static async init(
workspace: string, workspace: string,
options: Partial<YjsInitOptions & BlockClientOptions> = {} options: Partial<
YjsInitOptions & YjsProviderOptions & BlockClientOptions
> = {}
): Promise<BlockClientInstance> { ): Promise<BlockClientInstance> {
const instance = await YjsAdapter.init(workspace, { const instance = await YjsAdapter.init(workspace, {
backend: BucketBackend.YjsWebSocketAffine, provider: getYjsProviders({
backend: BucketBackend.YjsWebSocketAffine,
exportData: console.log.bind(console),
...options,
}),
...options, ...options,
}); });
return new BlockClient(instance, workspace, options); return new BlockClient(instance, workspace, options);