perf: use lazy load provider for IDB and SQLITE (#3351)

This commit is contained in:
Peng Xiao
2023-07-26 00:56:48 +08:00
committed by GitHub
parent e3f66d7e22
commit 20ee9d485d
25 changed files with 481 additions and 758 deletions
@@ -45,7 +45,7 @@ const ImagePreviewModal = lazy(() =>
);
const BlockSuiteEditorImpl = (props: EditorProps): ReactElement => {
const { onLoad, page, mode, style, onInit } = props;
const { onLoad, page, mode, style } = props;
if (!page.loaded) {
use(page.waitForLoaded());
}
@@ -66,14 +66,9 @@ const BlockSuiteEditorImpl = (props: EditorProps): ReactElement => {
editor.mode = mode;
}
useEffect(() => {
if (editor.page !== page) {
editor.page = page;
if (page.root === null) {
onInit(page, editor);
}
}
}, [editor, page, onInit]);
if (editor.page !== page) {
editor.page = page;
}
useEffect(() => {
if (editor.page && onLoad) {
@@ -11,7 +11,7 @@ import {
type TitleCellProps = {
icon: JSX.Element;
text: string;
desc?: string;
desc?: React.ReactNode;
suffix?: JSX.Element;
/**
* Customize the children of the cell
@@ -15,7 +15,7 @@ export type ListData = {
pageId: string;
icon: JSX.Element;
title: string;
preview?: string;
preview?: React.ReactNode;
tags: Tag[];
favorite: boolean;
createDate: Date;
@@ -34,7 +34,7 @@ export type TrashListData = {
pageId: string;
icon: JSX.Element;
title: string;
preview?: string;
preview?: React.ReactNode;
createDate: Date;
// TODO remove optional after assert that trashDate is always set
trashDate?: Date;
+1 -1
View File
@@ -13,7 +13,7 @@ export type ContentProps = {
lineHeight?: CSSProperties['lineHeight'];
ellipsis?: boolean;
lineNum?: number;
children: string;
children: React.ReactNode;
};
export const Content = styled('div', {
shouldForwardProp: prop => {
@@ -26,6 +26,7 @@ export function useBlockSuitePagePreview(page: Page): Atom<string> {
const disposable = page.slots.yUpdated.on(() => {
set(getPagePreviewText(page));
});
set(getPagePreviewText(page));
return () => {
disposable.dispose();
};
@@ -2,6 +2,7 @@ import { assertExists, DisposableGroup } from '@blocksuite/global/utils';
import type { Page, Workspace } from '@blocksuite/store';
import type { Atom } from 'jotai';
import { atom, useAtomValue } from 'jotai';
import { useEffect } from 'react';
const weakMap = new WeakMap<Workspace, Map<string, Atom<Page | null>>>();
@@ -51,5 +52,13 @@ export function useBlockSuiteWorkspacePage(
): Page | null {
const pageAtom = getAtom(blockSuiteWorkspace, pageId);
assertExists(pageAtom);
return useAtomValue(pageAtom);
const page = useAtomValue(pageAtom);
useEffect(() => {
if (!page?.loaded) {
page?.waitForLoaded().catch(console.error);
}
}, [page]);
return page;
}
+1 -7
View File
@@ -2,7 +2,6 @@ import { DebugLogger } from '@affine/debug';
import type { LocalWorkspace, WorkspaceCRUD } from '@affine/env/workspace';
import { WorkspaceFlavour } from '@affine/env/workspace';
import { nanoid, Workspace as BlockSuiteWorkspace } from '@blocksuite/store';
import { createIndexedDBProvider } from '@toeverything/y-indexeddb';
import { createJSONStorage } from 'jotai/utils';
import { z } from 'zod';
@@ -75,12 +74,7 @@ export const CRUD: WorkspaceCRUD<WorkspaceFlavour.LOCAL> = {
}
});
});
const persistence = createIndexedDBProvider(blockSuiteWorkspace.doc);
persistence.connect();
await persistence.whenSynced.then(() => {
persistence.disconnect();
});
// todo: do we need to persist doc to persistence datasource?
saveWorkspaceToLocalStorage(id);
return id;
},
@@ -68,7 +68,15 @@ describe('download provider', () => {
) as LocalIndexedDBDownloadProvider;
provider.sync();
await provider.whenReady;
expect(workspace.doc.toJSON()).toEqual(prev);
expect(workspace.doc.toJSON()).toEqual({
...prev,
// download provider only download the root doc
spaces: {
'space:page0': {
blocks: {},
},
},
});
}
});
});
@@ -2,9 +2,11 @@ import type {
SQLiteDBDownloadProvider,
SQLiteProvider,
} from '@affine/env/workspace';
import { getDoc } from '@affine/y-provider';
import { __unstableSchemas, AffineSchemas } from '@blocksuite/blocks/models';
import type { Y as YType } from '@blocksuite/store';
import { uuidv4, Workspace } from '@blocksuite/store';
import { setTimeout } from 'timers/promises';
import { beforeEach, describe, expect, test, vi } from 'vitest';
import {
@@ -30,18 +32,26 @@ const mockedAddBlob = vi.fn();
vi.stubGlobal('window', {
apis: {
db: {
getDocAsUpdates: async () => {
return Y.encodeStateAsUpdate(offlineYdoc);
getDocAsUpdates: async (workspaceId, guid) => {
const subdoc = guid ? getDoc(offlineYdoc, guid) : offlineYdoc;
if (!subdoc) {
return false;
}
return Y.encodeStateAsUpdate(subdoc);
},
applyDocUpdate: async (id: string, update: Uint8Array) => {
Y.applyUpdate(offlineYdoc, update, 'sqlite');
applyDocUpdate: async (id, update, subdocId) => {
const subdoc = subdocId ? getDoc(offlineYdoc, subdocId) : offlineYdoc;
if (!subdoc) {
return;
}
Y.applyUpdate(subdoc, update, 'sqlite');
},
getBlobKeys: async () => {
// todo: may need to hack the way to get hash keys of blobs
return [];
},
addBlob: mockedAddBlob,
} satisfies Partial<NonNullable<typeof window.apis>['db']>,
} as Partial<NonNullable<typeof window.apis>['db']>,
},
events: {
db: {
@@ -53,7 +63,7 @@ vi.stubGlobal('window', {
};
},
},
} satisfies Partial<NonNullable<typeof window.events>>,
} as Partial<NonNullable<typeof window.events>>,
});
vi.stubGlobal('environment', {
@@ -84,48 +94,25 @@ beforeEach(() => {
describe('SQLite download provider', () => {
test('sync updates', async () => {
// on connect, the updates from sqlite should be sync'ed to the existing ydoc
// and ydoc should be sync'ed back to sqlite
// Workspace.Y.applyUpdate(workspace.doc);
workspace.doc.getText('text').insert(0, 'mem-hello');
expect(offlineYdoc.getText('text').toString()).toBe('sqlite-hello');
downloadProvider.sync();
await downloadProvider.whenReady;
// depending on the nature of the sync, the data can be sync'ed in either direction
const options = ['mem-hellosqlite-hello', 'sqlite-hellomem-hello'];
const options = ['sqlite-hellomem-hello', 'mem-hellosqlite-hello'];
const synced = options.filter(
o => o === offlineYdoc.getText('text').toString()
o => o === workspace.doc.getText('text').toString()
);
expect(synced.length).toBe(1);
expect(workspace.doc.getText('text').toString()).toBe(synced[0]);
// workspace.doc.getText('text').insert(0, 'world');
// // check if the data are sync'ed
// expect(offlineYdoc.getText('text').toString()).toBe('world' + synced[0]);
});
test.fails('blobs will be synced to sqlite on connect', async () => {
// mock bs.list
const bin = new Uint8Array([1, 2, 3]);
const blob = new Blob([bin]);
workspace.blobs.list = vi.fn(async () => ['blob1']);
workspace.blobs.get = vi.fn(async () => {
return blob;
});
downloadProvider.sync();
await downloadProvider.whenReady;
await new Promise(resolve => setTimeout(resolve, 100));
expect(mockedAddBlob).toBeCalledWith(id, 'blob1', bin);
});
test('on db update', async () => {
// there is no updates from sqlite for now
test.skip('on db update', async () => {
provider.connect();
await setTimeout(200);
offlineYdoc.getText('text').insert(0, 'sqlite-world');
// @ts-expect-error
+4 -24
View File
@@ -10,7 +10,6 @@ import { createBroadcastChannelProvider } from '@blocksuite/store/providers/broa
import {
createIndexedDBProvider as create,
downloadBinary,
EarlyDisconnectError,
} from '@toeverything/y-indexeddb';
import type { Doc } from 'yjs';
@@ -40,17 +39,6 @@ const createIndexedDBBackgroundProvider: DocProviderCreator = (
connect: () => {
logger.info('connect indexeddb provider', id);
indexeddbProvider.connect();
indexeddbProvider.whenSynced
.then(() => {
connected = true;
})
.catch(error => {
connected = false;
if (error instanceof EarlyDisconnectError) {
return;
}
throw error;
});
},
disconnect: () => {
assertExists(indexeddbProvider);
@@ -61,7 +49,6 @@ const createIndexedDBBackgroundProvider: DocProviderCreator = (
};
};
const cache: WeakMap<Doc, Uint8Array> = new WeakMap();
const indexedDBDownloadOrigin = 'indexeddb-download-provider';
const createIndexedDBDownloadProvider: DocProviderCreator = (
@@ -74,18 +61,11 @@ const createIndexedDBDownloadProvider: DocProviderCreator = (
_resolve = resolve;
_reject = reject;
});
async function downloadBinaryRecursively(doc: Doc) {
if (cache.has(doc)) {
const binary = cache.get(doc) as Uint8Array;
async function downloadAndApply(doc: Doc) {
const binary = await downloadBinary(doc.guid);
if (binary) {
Y.applyUpdate(doc, binary, indexedDBDownloadOrigin);
} else {
const binary = await downloadBinary(doc.guid);
if (binary) {
Y.applyUpdate(doc, binary, indexedDBDownloadOrigin);
cache.set(doc, binary);
}
}
await Promise.all([...doc.subdocs].map(downloadBinaryRecursively));
}
return {
flavour: 'local-indexeddb',
@@ -98,7 +78,7 @@ const createIndexedDBDownloadProvider: DocProviderCreator = (
},
sync: () => {
logger.info('sync indexeddb provider', id);
downloadBinaryRecursively(doc).then(_resolve).catch(_reject);
downloadAndApply(doc).then(_resolve).catch(_reject);
},
};
};
@@ -2,8 +2,10 @@ import type {
SQLiteDBDownloadProvider,
SQLiteProvider,
} from '@affine/env/workspace';
import { getDoc } from '@affine/y-provider';
import { assertExists } from '@blocksuite/global/utils';
import {
createLazyProvider,
type DatasourceDocAdapter,
} from '@affine/y-provider';
import type { DocProviderCreator } from '@blocksuite/store';
import { Workspace as BlockSuiteWorkspace } from '@blocksuite/store';
import type { Doc } from 'yjs';
@@ -14,32 +16,26 @@ const Y = BlockSuiteWorkspace.Y;
const sqliteOrigin = Symbol('sqlite-provider-origin');
type SubDocsEvent = {
added: Set<Doc>;
removed: Set<Doc>;
loaded: Set<Doc>;
};
// workaround: there maybe new updates before SQLite is connected
// we need to exchange them with the SQLite db
// will be removed later when we have lazy load doc provider
const syncDiff = async (rootDoc: Doc, subdocId?: string) => {
try {
const workspaceId = rootDoc.guid;
const doc = subdocId ? getDoc(rootDoc, subdocId) : rootDoc;
if (!doc) {
logger.error('doc not found', workspaceId, subdocId);
return;
}
const update = await window.apis?.db.getDocAsUpdates(workspaceId, subdocId);
const diff = Y.encodeStateAsUpdate(
doc,
Y.encodeStateVectorFromUpdate(update)
);
await window.apis.db.applyDocUpdate(workspaceId, diff, subdocId);
} catch (err) {
logger.error('failed to sync diff', err);
const createDatasource = (workspaceId: string): DatasourceDocAdapter => {
if (!window.apis?.db) {
throw new Error('sqlite datasource is not available');
}
return {
queryDocState: async guid => {
return window.apis.db.getDocAsUpdates(
workspaceId,
workspaceId === guid ? undefined : guid
);
},
sendDocUpdate: async (guid, update) => {
return window.apis.db.applyDocUpdate(
guid,
update,
workspaceId === guid ? undefined : guid
);
},
};
};
/**
@@ -49,126 +45,27 @@ export const createSQLiteProvider: DocProviderCreator = (
id,
rootDoc
): SQLiteProvider => {
const { apis, events } = window;
// make sure it is being used in Electron with APIs
assertExists(apis);
assertExists(events);
const updateHandlerMap = new WeakMap<
Doc,
(update: Uint8Array, origin: unknown) => void
>();
const subDocsHandlerMap = new WeakMap<Doc, (event: SubDocsEvent) => void>();
const createOrHandleUpdate = (doc: Doc) => {
if (updateHandlerMap.has(doc)) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
return updateHandlerMap.get(doc)!;
}
function handleUpdate(update: Uint8Array, origin: unknown) {
if (origin === sqliteOrigin) {
return;
}
const subdocId = doc.guid === id ? undefined : doc.guid;
apis.db.applyDocUpdate(id, update, subdocId).catch(err => {
logger.error(err);
});
}
updateHandlerMap.set(doc, handleUpdate);
return handleUpdate;
};
const createOrGetHandleSubDocs = (doc: Doc) => {
if (subDocsHandlerMap.has(doc)) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
return subDocsHandlerMap.get(doc)!;
}
function handleSubdocs(event: SubDocsEvent) {
event.removed.forEach(doc => {
untrackDoc(doc);
});
event.loaded.forEach(doc => {
trackDoc(doc);
});
}
subDocsHandlerMap.set(doc, handleSubdocs);
return handleSubdocs;
};
function trackDoc(doc: Doc) {
syncDiff(rootDoc, rootDoc !== doc ? doc.guid : undefined).catch(
logger.error
);
doc.on('update', createOrHandleUpdate(doc));
doc.on('subdocs', createOrGetHandleSubDocs(doc));
doc.subdocs.forEach(doc => {
trackDoc(doc);
});
}
function untrackDoc(doc: Doc) {
doc.subdocs.forEach(doc => {
untrackDoc(doc);
});
doc.off('update', createOrHandleUpdate(doc));
doc.off('subdocs', createOrGetHandleSubDocs(doc));
}
let unsubscribe = () => {};
let datasource: ReturnType<typeof createDatasource> | null = null;
let provider: ReturnType<typeof createLazyProvider> | null = null;
let connected = false;
const connect = () => {
if (connected) {
return;
}
logger.info('connecting sqlite provider', id);
trackDoc(rootDoc);
unsubscribe = events.db.onExternalUpdate(
({
update,
workspaceId,
docId,
}: {
workspaceId: string;
update: Uint8Array;
docId?: string;
}) => {
if (workspaceId === id) {
if (docId) {
for (const doc of rootDoc.subdocs) {
if (doc.guid === docId) {
Y.applyUpdate(doc, update, sqliteOrigin);
return;
}
}
} else {
Y.applyUpdate(rootDoc, update, sqliteOrigin);
}
}
}
);
connected = true;
logger.info('connecting sqlite done', id);
};
const cleanup = () => {
logger.info('disconnecting sqlite provider', id);
unsubscribe();
untrackDoc(rootDoc);
connected = false;
};
return {
flavour: 'sqlite',
passive: true,
get connected(): boolean {
connect: () => {
datasource = createDatasource(id);
provider = createLazyProvider(rootDoc, datasource);
provider.connect();
connected = true;
},
disconnect: () => {
provider?.disconnect();
datasource = null;
provider = null;
connected = false;
},
get connected() {
return connected;
},
cleanup,
connect,
disconnect: cleanup,
};
};
@@ -180,7 +77,6 @@ export const createSQLiteDBDownloadProvider: DocProviderCreator = (
rootDoc
): SQLiteDBDownloadProvider => {
const { apis } = window;
let disconnected = false;
let _resolve: () => void;
let _reject: (error: unknown) => void;
@@ -194,33 +90,13 @@ export const createSQLiteDBDownloadProvider: DocProviderCreator = (
const subdocId = doc.guid === id ? undefined : doc.guid;
const updates = await apis.db.getDocAsUpdates(id, subdocId);
if (disconnected) {
return false;
}
if (updates) {
Y.applyUpdate(doc, updates, sqliteOrigin);
}
const mergedUpdates = Y.encodeStateAsUpdate(
doc,
Y.encodeStateVectorFromUpdate(updates)
);
// also apply updates to sqlite
await apis.db.applyDocUpdate(id, mergedUpdates, subdocId);
return true;
}
async function syncAllUpdates(doc: Doc) {
if (await syncUpdates(doc)) {
// load all subdocs
const subdocs = Array.from(doc.subdocs);
await Promise.all(subdocs.map(syncAllUpdates));
}
}
return {
flavour: 'sqlite-download',
active: true,
@@ -228,12 +104,12 @@ export const createSQLiteDBDownloadProvider: DocProviderCreator = (
return promise;
},
cleanup: () => {
disconnected = true;
// todo
},
sync: async () => {
logger.info('connect sqlite download provider', id);
try {
await syncAllUpdates(rootDoc);
await syncUpdates(rootDoc);
_resolve();
} catch (error) {
_reject(error);
+1
View File
@@ -0,0 +1 @@
This benchmark is outdated because our new API for IndexedDB has no direct parity with the official provider.
+1
View File
@@ -36,6 +36,7 @@
"idb": "^7.1.1"
},
"devDependencies": {
"@affine/y-provider": "workspace:*",
"@blocksuite/blocks": "0.0.0-20230721134812-6e0e3bef-nightly",
"@blocksuite/store": "0.0.0-20230721134812-6e0e3bef-nightly",
"vite": "^4.4.6",
@@ -3,18 +3,18 @@
*/
import 'fake-indexeddb/auto';
import { setTimeout } from 'node:timers/promises';
import { __unstableSchemas, AffineSchemas } from '@blocksuite/blocks/models';
import { assertExists } from '@blocksuite/global/utils';
import type { Page } from '@blocksuite/store';
import { uuidv4, Workspace } from '@blocksuite/store';
import { openDB } from 'idb';
import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest';
import { IndexeddbPersistence } from 'y-indexeddb';
import { applyUpdate, Doc, encodeStateAsUpdate } from 'yjs';
import type { WorkspacePersist } from '../index';
import {
CleanupWhenConnectingError,
createIndexedDBProvider,
dbVersion,
DEFAULT_DB_NAME,
@@ -43,7 +43,7 @@ function initEmptyPage(page: Page) {
async function getUpdates(id: string): Promise<Uint8Array[]> {
const db = await openDB(rootDBName, dbVersion);
const store = await db
const store = db
.transaction('workspace', 'readonly')
.objectStore('workspace');
const data = (await store.get(id)) as WorkspacePersist | undefined;
@@ -74,7 +74,10 @@ describe('indexeddb provider', () => {
test('connect', async () => {
const provider = createIndexedDBProvider(workspace.doc);
provider.connect();
await provider.whenSynced;
// todo: has a better way to know when data is synced
await setTimeout(200);
const db = await openDB(rootDBName, dbVersion);
{
const store = db
@@ -96,9 +99,9 @@ describe('indexeddb provider', () => {
const frameId = page.addBlock('affine:note', {}, pageBlockId);
page.addBlock('affine:paragraph', {}, frameId);
}
await new Promise(resolve => setTimeout(resolve, 1000));
await setTimeout(200);
{
const store = await db
const store = db
.transaction('workspace', 'readonly')
.objectStore('workspace');
const data = (await store.get(id)) as WorkspacePersist | undefined;
@@ -130,22 +133,11 @@ describe('indexeddb provider', () => {
}
});
test('disconnect suddenly', async () => {
const provider = createIndexedDBProvider(workspace.doc, rootDBName);
const fn = vi.fn();
provider.connect();
provider.disconnect();
expect(fn).toBeCalledTimes(0);
await provider.whenSynced.catch(fn);
expect(fn).toBeCalledTimes(1);
});
test('connect and disconnect', async () => {
const provider = createIndexedDBProvider(workspace.doc, rootDBName);
provider.connect();
expect(provider.connected).toBe(true);
const p1 = provider.whenSynced;
await p1;
await setTimeout(200);
const snapshot = encodeStateAsUpdate(workspace.doc);
provider.disconnect();
expect(provider.connected).toBe(false);
@@ -164,8 +156,7 @@ describe('indexeddb provider', () => {
expect(provider.connected).toBe(false);
provider.connect();
expect(provider.connected).toBe(true);
const p2 = provider.whenSynced;
await p2;
await setTimeout(200);
{
const updates = await getUpdates(workspace.id);
expect(updates).not.toEqual([]);
@@ -173,13 +164,12 @@ describe('indexeddb provider', () => {
expect(provider.connected).toBe(true);
provider.disconnect();
expect(provider.connected).toBe(false);
expect(p1).not.toBe(p2);
});
test('cleanup', async () => {
const provider = createIndexedDBProvider(workspace.doc);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
const db = await openDB(rootDBName, dbVersion);
{
@@ -190,8 +180,8 @@ describe('indexeddb provider', () => {
expect(keys).contain(workspace.id);
}
provider.disconnect();
await provider.cleanup();
provider.disconnect();
{
const store = db
@@ -202,17 +192,6 @@ describe('indexeddb provider', () => {
}
});
test('cleanup when connecting', async () => {
const provider = createIndexedDBProvider(workspace.doc);
provider.connect();
await expect(() => provider.cleanup()).rejects.toThrowError(
CleanupWhenConnectingError
);
await provider.whenSynced;
provider.disconnect();
await provider.cleanup();
});
test('merge', async () => {
setMergeCount(5);
const provider = createIndexedDBProvider(workspace.doc, rootDBName);
@@ -226,7 +205,7 @@ describe('indexeddb provider', () => {
page.addBlock('affine:paragraph', {}, frameId);
}
}
await provider.whenSynced;
await setTimeout(200);
{
const updates = await getUpdates(id);
expect(updates.length).lessThanOrEqual(5);
@@ -242,14 +221,12 @@ describe('indexeddb provider', () => {
{
const provider = createIndexedDBProvider(doc, rootDBName);
provider.connect();
await provider.whenSynced;
provider.disconnect();
}
{
const newDoc = new Workspace.Y.Doc();
const provider = createIndexedDBProvider(newDoc, rootDBName);
provider.connect();
await provider.whenSynced;
provider.disconnect();
newDoc.getMap('map').forEach((value, key) => {
expect(value).toBe(parseInt(key));
@@ -257,42 +234,6 @@ describe('indexeddb provider', () => {
}
});
test('migration', async () => {
{
const yDoc = new Doc();
yDoc.getMap().set('foo', 'bar');
const persistence = new IndexeddbPersistence('test', yDoc);
await persistence.whenSynced;
await persistence.destroy();
}
{
const yDoc = new Doc({
guid: 'test',
});
const provider = createIndexedDBProvider(yDoc);
provider.connect();
await provider.whenSynced;
await new Promise(resolve => setTimeout(resolve, 0));
expect(yDoc.getMap().get('foo')).toBe('bar');
}
localStorage.clear();
{
indexedDB.databases = vi.fn(async () => {
throw new Error('not supported');
});
await expect(indexedDB.databases).rejects.toThrow('not supported');
const yDoc = new Doc({
guid: 'test',
});
expect(indexedDB.databases).toBeCalledTimes(1);
const provider = createIndexedDBProvider(yDoc);
provider.connect();
await provider.whenSynced;
expect(indexedDB.databases).toBeCalledTimes(2);
expect(yDoc.getMap().get('foo')).toBe('bar');
}
});
test('beforeunload', async () => {
const oldAddEventListener = window.addEventListener;
window.addEventListener = vi.fn((event: string, fn, options) => {
@@ -311,7 +252,8 @@ describe('indexeddb provider', () => {
const map = doc.getMap('map');
map.set('1', 1);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
expect(window.addEventListener).toBeCalledTimes(1);
expect(window.removeEventListener).toBeCalledTimes(1);
@@ -396,7 +338,7 @@ describe('subDoc', () => {
map.set('2', 'test');
const provider = createIndexedDBProvider(doc);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
provider.disconnect();
json1 = doc.toJSON();
}
@@ -406,7 +348,7 @@ describe('subDoc', () => {
});
const provider = createIndexedDBProvider(doc);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
const map = doc.getMap();
const subDoc = map.get('1') as Doc;
subDoc.load();
@@ -431,7 +373,7 @@ describe('subDoc', () => {
});
await page1.waitForLoaded();
const { paragraphBlockId: paragraphBlockIdPage2 } = initEmptyPage(page1);
await new Promise(resolve => setTimeout(resolve, 1000));
await setTimeout(200);
provider.disconnect();
{
const newWorkspace = new Workspace({
@@ -441,15 +383,17 @@ describe('subDoc', () => {
newWorkspace.register(AffineSchemas).register(__unstableSchemas);
const provider = createIndexedDBProvider(newWorkspace.doc, rootDBName);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
const page0 = newWorkspace.getPage('page0') as Page;
await page0.waitForLoaded();
await setTimeout(200);
{
const block = page0.getBlockById(paragraphBlockIdPage1);
assertExists(block);
}
const page1 = newWorkspace.getPage('page1') as Page;
await page1.waitForLoaded();
await setTimeout(200);
{
const block = page1.getBlockById(paragraphBlockIdPage2);
assertExists(block);
@@ -465,7 +409,7 @@ describe('utils', () => {
initEmptyPage(page);
const provider = createIndexedDBProvider(workspace.doc, rootDBName);
provider.connect();
await provider.whenSynced;
await setTimeout(200);
provider.disconnect();
const update = (await downloadBinary(
workspace.id,
@@ -478,16 +422,12 @@ describe('utils', () => {
});
newWorkspace.register(AffineSchemas).register(__unstableSchemas);
applyUpdate(newWorkspace.doc, update);
await new Promise<void>(resolve =>
setTimeout(() => {
expect(workspace.doc.toJSON()['meta']).toEqual(
newWorkspace.doc.toJSON()['meta']
);
expect(Object.keys(workspace.doc.toJSON()['spaces'])).toEqual(
Object.keys(newWorkspace.doc.toJSON()['spaces'])
);
resolve();
}, 0)
await setTimeout();
expect(workspace.doc.toJSON()['meta']).toEqual(
newWorkspace.doc.toJSON()['meta']
);
expect(Object.keys(workspace.doc.toJSON()['spaces'])).toEqual(
Object.keys(newWorkspace.doc.toJSON()['spaces'])
);
});
+2 -262
View File
@@ -1,26 +1,17 @@
import { openDB } from 'idb';
import {
applyUpdate,
diffUpdate,
Doc,
encodeStateAsUpdate,
encodeStateVector,
UndoManager,
} from 'yjs';
import type {
BlockSuiteBinaryDB,
IndexedDBProvider,
WorkspaceMilestone,
} from './shared';
import type { BlockSuiteBinaryDB, WorkspaceMilestone } from './shared';
import { dbVersion, DEFAULT_DB_NAME, upgradeDB } from './shared';
import { tryMigrate } from './utils';
const indexeddbOrigin = 'indexeddb-provider-origin';
const snapshotOrigin = 'snapshot-origin';
let mergeCount = 500;
/**
* @internal
*/
@@ -40,10 +31,6 @@ export const writeOperation = async (op: Promise<unknown>) => {
});
};
export function setMergeCount(count: number) {
mergeCount = count;
}
export function revertUpdate(
doc: Doc,
snapshotUpdate: Uint8Array,
@@ -142,257 +129,10 @@ export const getMilestones = async (
return milestone.milestone;
};
type SubDocsEvent = {
added: Set<Doc>;
removed: Set<Doc>;
loaded: Set<Doc>;
};
/**
* We use `doc.guid` as the unique key, please make sure it not changes.
*/
export const createIndexedDBProvider = (
doc: Doc,
dbName: string = DEFAULT_DB_NAME,
/**
* In the future, migrate will be removed and there will be a separate function
*/
migrate = true
): IndexedDBProvider => {
let resolve: () => void;
let reject: (reason?: unknown) => void;
let early = true;
let connected = false;
const dbPromise = openDB<BlockSuiteBinaryDB>(dbName, dbVersion, {
upgrade: upgradeDB,
});
const updateHandlerMap = new WeakMap<
Doc,
(update: Uint8Array, origin: unknown) => void
>();
const destroyHandlerMap = new WeakMap<Doc, () => void>();
const subDocsHandlerMap = new WeakMap<Doc, (event: SubDocsEvent) => void>();
const createOrGetHandleUpdate = (id: string, doc: Doc) => {
if (updateHandlerMap.has(doc)) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
return updateHandlerMap.get(doc)!;
}
const fn = async function handleUpdate(
update: Uint8Array,
origin: unknown
) {
const db = await dbPromise;
if (!connected) {
return;
}
if (origin === indexeddbOrigin) {
return;
}
const store = db
.transaction('workspace', 'readwrite')
.objectStore('workspace');
let data = await store.get(id);
if (!data) {
data = {
id,
updates: [],
};
}
data.updates.push({
timestamp: Date.now(),
update,
});
if (data.updates.length > mergeCount) {
const updates = data.updates.map(({ update }) => update);
const doc = new Doc();
doc.transact(() => {
updates.forEach(update => {
applyUpdate(doc, update, indexeddbOrigin);
});
}, indexeddbOrigin);
const update = encodeStateAsUpdate(doc);
data = {
id,
updates: [
{
timestamp: Date.now(),
update,
},
],
};
await writeOperation(store.put(data));
} else {
await writeOperation(store.put(data));
}
};
updateHandlerMap.set(doc, fn);
return fn;
};
/* deepscan-disable UNUSED_PARAM */
const createOrGetHandleDestroy = (_: string, doc: Doc) => {
if (destroyHandlerMap.has(doc)) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
return destroyHandlerMap.get(doc)!;
}
const fn = async function handleDestroy() {
unTrackDoc(doc.guid, doc);
};
destroyHandlerMap.set(doc, fn);
return fn;
};
/* deepscan-disable UNUSED_PARAM */
const createOrGetHandleSubDocs = (_: string, doc: Doc) => {
if (subDocsHandlerMap.has(doc)) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
return subDocsHandlerMap.get(doc)!;
}
const fn = async function handleSubDocs(event: SubDocsEvent) {
event.removed.forEach(doc => {
unTrackDoc(doc.guid, doc);
});
event.loaded.forEach(doc => {
trackDoc(doc.guid, doc);
});
};
subDocsHandlerMap.set(doc, fn);
return fn;
};
function trackDoc(id: string, doc: Doc) {
doc.on('update', createOrGetHandleUpdate(id, doc));
doc.on('destroy', createOrGetHandleDestroy(id, doc));
doc.on('subdocs', createOrGetHandleSubDocs(id, doc));
doc.subdocs.forEach(doc => {
trackDoc(doc.guid, doc);
});
}
function unTrackDoc(id: string, doc: Doc) {
doc.subdocs.forEach(doc => {
unTrackDoc(doc.guid, doc);
});
doc.off('update', createOrGetHandleUpdate(id, doc));
doc.off('destroy', createOrGetHandleDestroy(id, doc));
doc.off('subdocs', createOrGetHandleSubDocs(id, doc));
}
async function saveDocOperation(id: string, doc: Doc) {
const db = await dbPromise;
const store = db
.transaction('workspace', 'readwrite')
.objectStore('workspace');
const data = await store.get(id);
if (!connected) {
return;
}
if (!data) {
await writeOperation(
db.put('workspace', {
id,
updates: [
{
timestamp: Date.now(),
update: encodeStateAsUpdate(doc),
},
],
})
);
} else {
const updates = data.updates.map(({ update }) => update);
const fakeDoc = new Doc();
fakeDoc.transact(() => {
updates.forEach(update => {
applyUpdate(fakeDoc, update, indexeddbOrigin);
});
}, indexeddbOrigin);
const newUpdate = diffUpdate(
encodeStateAsUpdate(doc),
encodeStateAsUpdate(fakeDoc)
);
await writeOperation(
store.put({
...data,
updates: [
...data.updates,
{
timestamp: Date.now(),
update: newUpdate,
},
],
})
);
doc.transact(() => {
updates.forEach(update => {
applyUpdate(doc, update, indexeddbOrigin);
});
}, indexeddbOrigin);
}
}
const apis = {
connect: async () => {
if (connected) return;
apis.whenSynced = new Promise<void>((_resolve, _reject) => {
early = true;
resolve = _resolve;
reject = _reject;
});
connected = true;
trackDoc(doc.guid, doc);
// only the runs `await` below, otherwise the logic is incorrect
const db = await dbPromise;
if (migrate) {
// Tips:
// this is only backward compatible with the yjs official version of y-indexeddb
await tryMigrate(db, doc.guid, dbName);
}
if (!connected) {
return;
}
// recursively save all docs into indexeddb
const docs: [string, Doc][] = [];
docs.push([doc.guid, doc]);
while (docs.length > 0) {
const [id, doc] = docs.pop() as [string, Doc];
await saveDocOperation(id, doc);
doc.subdocs.forEach(doc => {
docs.push([doc.guid, doc]);
});
}
early = false;
resolve();
},
disconnect() {
connected = false;
if (early) {
reject(new EarlyDisconnectError());
}
unTrackDoc(doc.guid, doc);
},
async cleanup() {
if (connected) {
throw new CleanupWhenConnectingError();
}
await (await dbPromise).delete('workspace', doc.guid);
},
whenSynced: Promise.resolve(),
get connected() {
return connected;
},
};
return apis;
};
export * from './provider';
export * from './shared';
export * from './utils';
+133
View File
@@ -0,0 +1,133 @@
import {
createLazyProvider,
type DatasourceDocAdapter,
writeOperation,
} from '@affine/y-provider';
import { openDB } from 'idb';
import type { Doc } from 'yjs';
import { diffUpdate, mergeUpdates } from 'yjs';
import {
type BlockSuiteBinaryDB,
dbVersion,
DEFAULT_DB_NAME,
type IndexedDBProvider,
type UpdateMessage,
upgradeDB,
} from './shared';
let mergeCount = 500;
export function setMergeCount(count: number) {
mergeCount = count;
}
const createDatasource = ({
dbName,
mergeCount,
}: {
dbName: string;
mergeCount?: number;
}) => {
const dbPromise = openDB<BlockSuiteBinaryDB>(dbName, dbVersion, {
upgrade: upgradeDB,
});
const adapter = {
queryDocState: async (guid, options) => {
try {
const db = await dbPromise;
const store = db
.transaction('workspace', 'readonly')
.objectStore('workspace');
const data = await store.get(guid);
if (!data) {
return false;
}
const { updates } = data;
const update = mergeUpdates(updates.map(({ update }) => update));
const diff = options?.stateVector
? diffUpdate(update, options?.stateVector)
: update;
return diff;
} catch (err: any) {
if (!err.message?.includes('The database connection is closing.')) {
throw err;
}
return false;
}
},
sendDocUpdate: async (guid, update) => {
try {
const db = await dbPromise;
const store = db
.transaction('workspace', 'readwrite')
.objectStore('workspace');
// TODO: maybe we do not need to get data every time
const { updates } = (await store.get(guid)) ?? { updates: [] };
let rows: UpdateMessage[] = [
...updates,
{ timestamp: Date.now(), update },
];
if (mergeCount && rows.length >= mergeCount) {
const merged = mergeUpdates(rows.map(({ update }) => update));
rows = [{ timestamp: Date.now(), update: merged }];
}
await writeOperation(
store.put({
id: guid,
updates: rows,
})
);
} catch (err: any) {
if (!err.message?.includes('The database connection is closing.')) {
throw err;
}
}
},
} satisfies DatasourceDocAdapter;
return {
...adapter,
disconnect: () => {
dbPromise.then(db => db.close()).catch(console.error);
},
cleanup: async () => {
const db = await dbPromise;
await db.clear('workspace');
},
};
};
export const createIndexedDBProvider = (
doc: Doc,
dbName: string = DEFAULT_DB_NAME
): IndexedDBProvider => {
let datasource: ReturnType<typeof createDatasource> | null = null;
let provider: ReturnType<typeof createLazyProvider> | null = null;
return {
connect: () => {
datasource = createDatasource({ dbName, mergeCount });
provider = createLazyProvider(doc, datasource);
provider.connect();
},
disconnect: () => {
datasource?.disconnect();
provider?.disconnect();
datasource = null;
provider = null;
},
cleanup: async () => {
await datasource?.cleanup();
},
get connected() {
return provider?.connected || false;
},
};
};
-1
View File
@@ -12,7 +12,6 @@ export interface IndexedDBProvider {
connect: () => void;
disconnect: () => void;
cleanup: () => Promise<void>;
whenSynced: Promise<void>;
readonly connected: boolean;
}
+3
View File
@@ -9,6 +9,9 @@
"references": [
{
"path": "./tsconfig.node.json"
},
{
"path": "../y-provider"
}
]
}
-4
View File
@@ -3,11 +3,7 @@
"type": "module",
"version": "0.7.0-canary.51",
"description": "Yjs provider utilities for AFFiNE",
"exports": {
".": "./src/index.ts"
},
"main": "./src/index.ts",
"module": "./src/index.ts",
"devDependencies": {
"@blocksuite/store": "0.0.0-20230721134812-6e0e3bef-nightly"
},
+17 -17
View File
@@ -3,6 +3,7 @@ import {
applyUpdate,
type Doc,
encodeStateAsUpdate,
encodeStateVector,
encodeStateVectorFromUpdate,
} from 'yjs';
@@ -33,30 +34,30 @@ export const createLazyProvider = (
let connected = false;
const pendingMap = new Map<string, Uint8Array[]>(); // guid -> pending-updates
const disposableMap = new Map<string, Set<() => void>>();
const connectedDocs = new Set();
const connectedDocs = new Set<string>();
let datasourceUnsub: (() => void) | undefined;
async function syncDoc(doc: Doc) {
const guid = doc.guid;
// perf: optimize me
const currentUpdate = encodeStateAsUpdate(doc);
const remoteUpdate = await datasource.queryDocState(guid, {
stateVector: encodeStateVectorFromUpdate(currentUpdate),
stateVector: encodeStateVector(doc),
});
const updates = [currentUpdate];
pendingMap.set(guid, []);
if (remoteUpdate) {
applyUpdate(doc, remoteUpdate, selfUpdateOrigin);
const newUpdate = encodeStateAsUpdate(
doc,
encodeStateVectorFromUpdate(remoteUpdate)
);
updates.push(newUpdate);
await datasource.sendDocUpdate(guid, newUpdate);
}
const sv = remoteUpdate
? encodeStateVectorFromUpdate(remoteUpdate)
: undefined;
// perf: optimize me
// it is possible the doc is only in memory but not yet in the datasource
// we need to send the whole update to the datasource
await datasource.sendDocUpdate(guid, encodeStateAsUpdate(doc, sv));
}
/**
@@ -73,10 +74,7 @@ export const createLazyProvider = (
datasource.sendDocUpdate(doc.guid, update).catch(console.error);
};
const subdocLoadHandler = (event: {
loaded: Set<Doc>;
removed: Set<Doc>;
}) => {
const subdocsHandler = (event: { loaded: Set<Doc>; removed: Set<Doc> }) => {
event.loaded.forEach(subdoc => {
connectDoc(subdoc).catch(console.error);
});
@@ -86,11 +84,11 @@ export const createLazyProvider = (
};
doc.on('update', updateHandler);
doc.on('subdocs', subdocLoadHandler);
doc.on('subdocs', subdocsHandler);
// todo: handle destroy?
disposables.add(() => {
doc.off('update', updateHandler);
doc.off('subdocs', subdocLoadHandler);
doc.off('subdocs', subdocsHandler);
});
}
@@ -127,6 +125,7 @@ export const createLazyProvider = (
connectedDocs.add(doc.guid);
setupDocListener(doc);
await syncDoc(doc);
await Promise.all(
[...doc.subdocs]
.filter(subdoc => subdoc.shouldLoad)
@@ -150,6 +149,7 @@ export const createLazyProvider = (
disposables.forEach(dispose => dispose());
});
disposableMap.clear();
connectedDocs.clear();
}
/**
+16
View File
@@ -12,3 +12,19 @@ export function getDoc(doc: Doc, guid: string): Doc | undefined {
}
return undefined;
}
const saveAlert = (event: BeforeUnloadEvent) => {
event.preventDefault();
return (event.returnValue =
'Data is not saved. Are you sure you want to leave?');
};
export const writeOperation = async (op: Promise<unknown>) => {
window.addEventListener('beforeunload', saveAlert, {
capture: true,
});
await op;
window.removeEventListener('beforeunload', saveAlert, {
capture: true,
});
};