feat(nbstore): better doc sync logic (#9037)
This commit is contained in:
@@ -10,7 +10,7 @@ export class IndexedDBBlobStorage extends BlobStorage {
|
|||||||
readonly connection = share(new IDBConnection(this.options));
|
readonly connection = share(new IDBConnection(this.options));
|
||||||
|
|
||||||
get db() {
|
get db() {
|
||||||
return this.connection.inner;
|
return this.connection.inner.db;
|
||||||
}
|
}
|
||||||
|
|
||||||
override async get(key: string) {
|
override async get(key: string) {
|
||||||
|
|||||||
@@ -4,8 +4,11 @@ import { Connection } from '../../connection';
|
|||||||
import type { StorageOptions } from '../../storage';
|
import type { StorageOptions } from '../../storage';
|
||||||
import { type DocStorageSchema, migrator } from './schema';
|
import { type DocStorageSchema, migrator } from './schema';
|
||||||
|
|
||||||
export class IDBConnection extends Connection<IDBPDatabase<DocStorageSchema>> {
|
export class IDBConnection extends Connection<{
|
||||||
private readonly dbName = `${this.opts.peer}:${this.opts.type}:${this.opts.id}`;
|
db: IDBPDatabase<DocStorageSchema>;
|
||||||
|
channel: BroadcastChannel;
|
||||||
|
}> {
|
||||||
|
readonly dbName = `${this.opts.peer}:${this.opts.type}:${this.opts.id}`;
|
||||||
|
|
||||||
override get shareId() {
|
override get shareId() {
|
||||||
return `idb(${migrator.version}):${this.dbName}`;
|
return `idb(${migrator.version}):${this.dbName}`;
|
||||||
@@ -16,20 +19,23 @@ export class IDBConnection extends Connection<IDBPDatabase<DocStorageSchema>> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
override async doConnect() {
|
override async doConnect() {
|
||||||
return openDB<DocStorageSchema>(this.dbName, migrator.version, {
|
return {
|
||||||
upgrade: migrator.migrate,
|
db: await openDB<DocStorageSchema>(this.dbName, migrator.version, {
|
||||||
blocking: () => {
|
upgrade: migrator.migrate,
|
||||||
// if, for example, an tab with newer version is opened, this function will be called.
|
blocking: () => {
|
||||||
// we should close current connection to allow the new version to upgrade the db.
|
// if, for example, an tab with newer version is opened, this function will be called.
|
||||||
this.close(
|
// we should close current connection to allow the new version to upgrade the db.
|
||||||
new Error('Blocking a new version. Closing the connection.')
|
this.close(
|
||||||
);
|
new Error('Blocking a new version. Closing the connection.')
|
||||||
},
|
);
|
||||||
blocked: () => {
|
},
|
||||||
// fallback to retry auto retry
|
blocked: () => {
|
||||||
this.setStatus('error', new Error('Blocked by other tabs.'));
|
// fallback to retry auto retry
|
||||||
},
|
this.setStatus('error', new Error('Blocked by other tabs.'));
|
||||||
});
|
},
|
||||||
|
}),
|
||||||
|
channel: new BroadcastChannel(this.dbName),
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
override async doDisconnect() {
|
override async doDisconnect() {
|
||||||
@@ -37,7 +43,8 @@ export class IDBConnection extends Connection<IDBPDatabase<DocStorageSchema>> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private close(error?: Error) {
|
private close(error?: Error) {
|
||||||
this.maybeConnection?.close();
|
this.maybeConnection?.channel.close();
|
||||||
|
this.maybeConnection?.db.close();
|
||||||
this.setStatus('closed', error);
|
this.setStatus('closed', error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,19 +4,37 @@ import {
|
|||||||
type DocClocks,
|
type DocClocks,
|
||||||
type DocRecord,
|
type DocRecord,
|
||||||
DocStorage,
|
DocStorage,
|
||||||
|
type DocStorageOptions,
|
||||||
type DocUpdate,
|
type DocUpdate,
|
||||||
} from '../../storage';
|
} from '../../storage';
|
||||||
import { IDBConnection } from './db';
|
import { IDBConnection } from './db';
|
||||||
|
import { IndexedDBLocker } from './lock';
|
||||||
|
|
||||||
|
interface ChannelMessage {
|
||||||
|
type: 'update';
|
||||||
|
update: DocRecord;
|
||||||
|
origin?: string;
|
||||||
|
}
|
||||||
|
|
||||||
export class IndexedDBDocStorage extends DocStorage {
|
export class IndexedDBDocStorage extends DocStorage {
|
||||||
readonly connection = share(new IDBConnection(this.options));
|
readonly connection = share(new IDBConnection(this.options));
|
||||||
|
|
||||||
get db() {
|
get db() {
|
||||||
return this.connection.inner;
|
return this.connection.inner.db;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
get channel() {
|
||||||
|
return this.connection.inner.channel;
|
||||||
|
}
|
||||||
|
|
||||||
|
override locker = new IndexedDBLocker(this.connection);
|
||||||
|
|
||||||
private _lastTimestamp = new Date(0);
|
private _lastTimestamp = new Date(0);
|
||||||
|
|
||||||
|
constructor(options: DocStorageOptions) {
|
||||||
|
super(options);
|
||||||
|
}
|
||||||
|
|
||||||
private generateTimestamp() {
|
private generateTimestamp() {
|
||||||
const timestamp = new Date();
|
const timestamp = new Date();
|
||||||
if (timestamp.getTime() <= this._lastTimestamp.getTime()) {
|
if (timestamp.getTime() <= this._lastTimestamp.getTime()) {
|
||||||
@@ -47,6 +65,17 @@ export class IndexedDBDocStorage extends DocStorage {
|
|||||||
origin
|
origin
|
||||||
);
|
);
|
||||||
|
|
||||||
|
this.channel.postMessage({
|
||||||
|
type: 'update',
|
||||||
|
update: {
|
||||||
|
docId: update.docId,
|
||||||
|
bin: update.bin,
|
||||||
|
timestamp,
|
||||||
|
editor: update.editor,
|
||||||
|
},
|
||||||
|
origin,
|
||||||
|
} satisfies ChannelMessage);
|
||||||
|
|
||||||
return { docId: update.docId, timestamp };
|
return { docId: update.docId, timestamp };
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -144,4 +173,31 @@ export class IndexedDBDocStorage extends DocStorage {
|
|||||||
trx.commit();
|
trx.commit();
|
||||||
return updates.length;
|
return updates.length;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private docUpdateListener = 0;
|
||||||
|
|
||||||
|
override subscribeDocUpdate(
|
||||||
|
callback: (update: DocRecord, origin?: string) => void
|
||||||
|
): () => void {
|
||||||
|
if (this.docUpdateListener === 0) {
|
||||||
|
this.channel.addEventListener('message', this.handleChannelMessage);
|
||||||
|
}
|
||||||
|
this.docUpdateListener++;
|
||||||
|
|
||||||
|
const dispose = super.subscribeDocUpdate(callback);
|
||||||
|
|
||||||
|
return () => {
|
||||||
|
dispose();
|
||||||
|
this.docUpdateListener--;
|
||||||
|
if (this.docUpdateListener === 0) {
|
||||||
|
this.channel.removeEventListener('message', this.handleChannelMessage);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
handleChannelMessage(event: MessageEvent<ChannelMessage>) {
|
||||||
|
if (event.data.type === 'update') {
|
||||||
|
this.emit('update', event.data.update, event.data.origin);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
83
packages/common/nbstore/src/impls/idb/lock.ts
Normal file
83
packages/common/nbstore/src/impls/idb/lock.ts
Normal file
@@ -0,0 +1,83 @@
|
|||||||
|
import EventEmitter2 from 'eventemitter2';
|
||||||
|
|
||||||
|
import { type Locker } from '../../storage/lock';
|
||||||
|
import type { IDBConnection } from './db';
|
||||||
|
|
||||||
|
interface ChannelMessage {
|
||||||
|
type: 'unlock';
|
||||||
|
key: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class IndexedDBLocker implements Locker {
|
||||||
|
get db() {
|
||||||
|
return this.dbConnection.inner.db;
|
||||||
|
}
|
||||||
|
private readonly eventEmitter = new EventEmitter2();
|
||||||
|
|
||||||
|
get channel() {
|
||||||
|
return this.dbConnection.inner.channel;
|
||||||
|
}
|
||||||
|
|
||||||
|
constructor(private readonly dbConnection: IDBConnection) {}
|
||||||
|
|
||||||
|
async lock(domain: string, resource: string) {
|
||||||
|
const key = `${domain}:${resource}`;
|
||||||
|
|
||||||
|
// eslint-disable-next-line no-constant-condition
|
||||||
|
while (true) {
|
||||||
|
const trx = this.db.transaction('locks', 'readwrite');
|
||||||
|
const record = await trx.store.get(key);
|
||||||
|
const lockTimestamp = record?.lock.getTime();
|
||||||
|
|
||||||
|
if (
|
||||||
|
lockTimestamp &&
|
||||||
|
lockTimestamp > Date.now() - 30000 /* lock timeout 3s */
|
||||||
|
) {
|
||||||
|
trx.commit();
|
||||||
|
|
||||||
|
await new Promise<void>(resolve => {
|
||||||
|
const cleanup = () => {
|
||||||
|
this.channel.removeEventListener('message', channelListener);
|
||||||
|
this.eventEmitter.off('unlock', eventListener);
|
||||||
|
clearTimeout(timer);
|
||||||
|
};
|
||||||
|
const channelListener = (event: MessageEvent<ChannelMessage>) => {
|
||||||
|
if (event.data.type === 'unlock' && event.data.key === key) {
|
||||||
|
cleanup();
|
||||||
|
resolve();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
const eventListener = (unlockKey: string) => {
|
||||||
|
if (unlockKey === key) {
|
||||||
|
cleanup();
|
||||||
|
resolve();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
this.channel.addEventListener('message', channelListener); // add listener
|
||||||
|
this.eventEmitter.on('unlock', eventListener);
|
||||||
|
|
||||||
|
const timer = setTimeout(() => {
|
||||||
|
cleanup();
|
||||||
|
resolve();
|
||||||
|
}, 3000);
|
||||||
|
// timeout to avoid dead lock
|
||||||
|
});
|
||||||
|
continue;
|
||||||
|
} else {
|
||||||
|
await trx.store.put({ key, lock: new Date() });
|
||||||
|
trx.commit();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
[Symbol.asyncDispose]: async () => {
|
||||||
|
const trx = this.db.transaction('locks', 'readwrite');
|
||||||
|
await trx.store.delete(key);
|
||||||
|
trx.commit();
|
||||||
|
this.channel.postMessage({ type: 'unlock', key });
|
||||||
|
this.eventEmitter.emit('unlock', key);
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -101,6 +101,13 @@ export interface DocStorageSchema extends DBSchema {
|
|||||||
peer: string;
|
peer: string;
|
||||||
};
|
};
|
||||||
};
|
};
|
||||||
|
locks: {
|
||||||
|
key: string;
|
||||||
|
value: {
|
||||||
|
key: string;
|
||||||
|
lock: Date;
|
||||||
|
};
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
const migrate: OpenDBCallbacks<DocStorageSchema>['upgrade'] = (
|
const migrate: OpenDBCallbacks<DocStorageSchema>['upgrade'] = (
|
||||||
@@ -162,6 +169,11 @@ const init: Migrate = db => {
|
|||||||
keyPath: 'key',
|
keyPath: 'key',
|
||||||
autoIncrement: false,
|
autoIncrement: false,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
db.createObjectStore('locks', {
|
||||||
|
keyPath: 'key',
|
||||||
|
autoIncrement: false,
|
||||||
|
});
|
||||||
};
|
};
|
||||||
// END REGION
|
// END REGION
|
||||||
|
|
||||||
|
|||||||
@@ -5,9 +5,24 @@ export class IndexedDBSyncStorage extends SyncStorage {
|
|||||||
readonly connection = share(new IDBConnection(this.options));
|
readonly connection = share(new IDBConnection(this.options));
|
||||||
|
|
||||||
get db() {
|
get db() {
|
||||||
return this.connection.inner;
|
return this.connection.inner.db;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override async getPeerRemoteClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null> {
|
||||||
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
|
const record = await trx.store.get([peer, docId]);
|
||||||
|
|
||||||
|
return record
|
||||||
|
? {
|
||||||
|
docId: record.docId,
|
||||||
|
timestamp: record.clock,
|
||||||
|
}
|
||||||
|
: null;
|
||||||
|
}
|
||||||
override async getPeerRemoteClocks(peer: string) {
|
override async getPeerRemoteClocks(peer: string) {
|
||||||
const trx = this.db.transaction('peerClocks', 'readonly');
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
@@ -34,6 +49,21 @@ export class IndexedDBSyncStorage extends SyncStorage {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override async getPeerPulledRemoteClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null> {
|
||||||
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
|
const record = await trx.store.get([peer, docId]);
|
||||||
|
|
||||||
|
return record
|
||||||
|
? {
|
||||||
|
docId: record.docId,
|
||||||
|
timestamp: record.pulledClock,
|
||||||
|
}
|
||||||
|
: null;
|
||||||
|
}
|
||||||
override async getPeerPulledRemoteClocks(peer: string) {
|
override async getPeerPulledRemoteClocks(peer: string) {
|
||||||
const trx = this.db.transaction('peerClocks', 'readonly');
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
@@ -59,6 +89,21 @@ export class IndexedDBSyncStorage extends SyncStorage {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override async getPeerPushedClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null> {
|
||||||
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
|
const record = await trx.store.get([peer, docId]);
|
||||||
|
|
||||||
|
return record
|
||||||
|
? {
|
||||||
|
docId: record.docId,
|
||||||
|
timestamp: record.pushedClock,
|
||||||
|
}
|
||||||
|
: null;
|
||||||
|
}
|
||||||
override async getPeerPushedClocks(peer: string) {
|
override async getPeerPushedClocks(peer: string) {
|
||||||
const trx = this.db.transaction('peerClocks', 'readonly');
|
const trx = this.db.transaction('peerClocks', 'readonly');
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ import EventEmitter2 from 'eventemitter2';
|
|||||||
import { diffUpdate, encodeStateVectorFromUpdate, mergeUpdates } from 'yjs';
|
import { diffUpdate, encodeStateVectorFromUpdate, mergeUpdates } from 'yjs';
|
||||||
|
|
||||||
import { isEmptyUpdate } from '../utils/is-empty-update';
|
import { isEmptyUpdate } from '../utils/is-empty-update';
|
||||||
import type { Lock } from './lock';
|
import type { Locker } from './lock';
|
||||||
import { SingletonLocker } from './lock';
|
import { SingletonLocker } from './lock';
|
||||||
import { Storage, type StorageOptions } from './storage';
|
import { Storage, type StorageOptions } from './storage';
|
||||||
|
|
||||||
@@ -42,7 +42,7 @@ export abstract class DocStorage<
|
|||||||
> extends Storage<Opts> {
|
> extends Storage<Opts> {
|
||||||
private readonly event = new EventEmitter2();
|
private readonly event = new EventEmitter2();
|
||||||
override readonly storageType = 'doc';
|
override readonly storageType = 'doc';
|
||||||
private readonly locker = new SingletonLocker();
|
protected readonly locker: Locker = new SingletonLocker();
|
||||||
|
|
||||||
// REGION: open apis by Op system
|
// REGION: open apis by Op system
|
||||||
/**
|
/**
|
||||||
@@ -243,7 +243,7 @@ export abstract class DocStorage<
|
|||||||
return merge(updates.filter(bin => !isEmptyUpdate(bin)));
|
return merge(updates.filter(bin => !isEmptyUpdate(bin)));
|
||||||
}
|
}
|
||||||
|
|
||||||
protected async lockDocForUpdate(docId: string): Promise<Lock> {
|
protected async lockDocForUpdate(docId: string): Promise<AsyncDisposable> {
|
||||||
return this.locker.lock(`workspace:${this.spaceId}:update`, docId);
|
return this.locker.lock(`workspace:${this.spaceId}:update`, docId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
export interface Locker {
|
export interface Locker {
|
||||||
lock(domain: string, resource: string): Promise<Lock>;
|
lock(domain: string, resource: string): Promise<AsyncDisposable>;
|
||||||
}
|
}
|
||||||
|
|
||||||
export class SingletonLocker implements Locker {
|
export class SingletonLocker implements Locker {
|
||||||
|
|||||||
@@ -8,14 +8,25 @@ export abstract class SyncStorage<
|
|||||||
> extends Storage<Opts> {
|
> extends Storage<Opts> {
|
||||||
override readonly storageType = 'sync';
|
override readonly storageType = 'sync';
|
||||||
|
|
||||||
|
abstract getPeerRemoteClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null>;
|
||||||
abstract getPeerRemoteClocks(peer: string): Promise<DocClocks>;
|
abstract getPeerRemoteClocks(peer: string): Promise<DocClocks>;
|
||||||
abstract setPeerRemoteClock(peer: string, clock: DocClock): Promise<void>;
|
abstract setPeerRemoteClock(peer: string, clock: DocClock): Promise<void>;
|
||||||
|
abstract getPeerPulledRemoteClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null>;
|
||||||
abstract getPeerPulledRemoteClocks(peer: string): Promise<DocClocks>;
|
abstract getPeerPulledRemoteClocks(peer: string): Promise<DocClocks>;
|
||||||
|
|
||||||
abstract setPeerPulledRemoteClock(
|
abstract setPeerPulledRemoteClock(
|
||||||
peer: string,
|
peer: string,
|
||||||
clock: DocClock
|
clock: DocClock
|
||||||
): Promise<void>;
|
): Promise<void>;
|
||||||
|
abstract getPeerPushedClock(
|
||||||
|
peer: string,
|
||||||
|
docId: string
|
||||||
|
): Promise<DocClock | null>;
|
||||||
abstract getPeerPushedClocks(peer: string): Promise<DocClocks>;
|
abstract getPeerPushedClocks(peer: string): Promise<DocClocks>;
|
||||||
abstract setPeerPushedClock(peer: string, clock: DocClock): Promise<void>;
|
abstract setPeerPushedClock(peer: string, clock: DocClock): Promise<void>;
|
||||||
abstract clearClocks(): Promise<void>;
|
abstract clearClocks(): Promise<void>;
|
||||||
|
|||||||
@@ -32,7 +32,7 @@ type Job =
|
|||||||
type: 'save';
|
type: 'save';
|
||||||
docId: string;
|
docId: string;
|
||||||
update?: Uint8Array;
|
update?: Uint8Array;
|
||||||
serverClock: Date;
|
remoteClock: Date;
|
||||||
};
|
};
|
||||||
|
|
||||||
interface Status {
|
interface Status {
|
||||||
@@ -41,8 +41,6 @@ interface Status {
|
|||||||
jobDocQueue: AsyncPriorityQueue;
|
jobDocQueue: AsyncPriorityQueue;
|
||||||
jobMap: Map<string, Job[]>;
|
jobMap: Map<string, Job[]>;
|
||||||
remoteClocks: ClockMap;
|
remoteClocks: ClockMap;
|
||||||
pulledRemoteClocks: ClockMap;
|
|
||||||
pushedClocks: ClockMap;
|
|
||||||
syncing: boolean;
|
syncing: boolean;
|
||||||
retrying: boolean;
|
retrying: boolean;
|
||||||
errorMessage: string | null;
|
errorMessage: string | null;
|
||||||
@@ -81,7 +79,7 @@ export class DocSyncPeer {
|
|||||||
/**
|
/**
|
||||||
* random unique id for recognize self in "update" event
|
* random unique id for recognize self in "update" event
|
||||||
*/
|
*/
|
||||||
private readonly uniqueId = nanoid();
|
private readonly uniqueId = `sync:${this.local.peer}:${this.remote.peer}:${nanoid()}`;
|
||||||
private readonly prioritySettings = new Map<string, number>();
|
private readonly prioritySettings = new Map<string, number>();
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
@@ -97,8 +95,6 @@ export class DocSyncPeer {
|
|||||||
jobDocQueue: new AsyncPriorityQueue(),
|
jobDocQueue: new AsyncPriorityQueue(),
|
||||||
jobMap: new Map(),
|
jobMap: new Map(),
|
||||||
remoteClocks: new ClockMap(new Map()),
|
remoteClocks: new ClockMap(new Map()),
|
||||||
pulledRemoteClocks: new ClockMap(new Map()),
|
|
||||||
pushedClocks: new ClockMap(new Map()),
|
|
||||||
syncing: false,
|
syncing: false,
|
||||||
retrying: false,
|
retrying: false,
|
||||||
errorMessage: null,
|
errorMessage: null,
|
||||||
@@ -107,14 +103,23 @@ export class DocSyncPeer {
|
|||||||
|
|
||||||
private readonly jobs = createJobErrorCatcher({
|
private readonly jobs = createJobErrorCatcher({
|
||||||
connect: async (docId: string, signal?: AbortSignal) => {
|
connect: async (docId: string, signal?: AbortSignal) => {
|
||||||
const pushedClock = this.status.pushedClocks.get(docId);
|
const pushedClock =
|
||||||
|
(await this.syncMetadata.getPeerPushedClock(this.remote.peer, docId))
|
||||||
|
?.timestamp ?? null;
|
||||||
const clock = await this.local.getDocTimestamp(docId);
|
const clock = await this.local.getDocTimestamp(docId);
|
||||||
|
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
if (pushedClock === null || pushedClock !== clock?.timestamp) {
|
if (pushedClock === null || pushedClock !== clock?.timestamp) {
|
||||||
await this.jobs.pullAndPush(docId, signal);
|
await this.jobs.pullAndPush(docId, signal);
|
||||||
} else {
|
} else {
|
||||||
const pulled = this.status.pulledRemoteClocks.get(docId);
|
// no need to push
|
||||||
|
const pulled =
|
||||||
|
(
|
||||||
|
await this.syncMetadata.getPeerPulledRemoteClock(
|
||||||
|
this.remote.peer,
|
||||||
|
docId
|
||||||
|
)
|
||||||
|
)?.timestamp ?? null;
|
||||||
if (pulled === null || pulled !== this.status.remoteClocks.get(docId)) {
|
if (pulled === null || pulled !== this.status.remoteClocks.get(docId)) {
|
||||||
await this.jobs.pull(docId, signal);
|
await this.jobs.pull(docId, signal);
|
||||||
}
|
}
|
||||||
@@ -133,6 +138,7 @@ export class DocSyncPeer {
|
|||||||
(a, b) => (a.getTime() > b.clock.getTime() ? a : b.clock),
|
(a, b) => (a.getTime() > b.clock.getTime() ? a : b.clock),
|
||||||
new Date(0)
|
new Date(0)
|
||||||
);
|
);
|
||||||
|
|
||||||
const merged = await this.mergeUpdates(
|
const merged = await this.mergeUpdates(
|
||||||
jobs.map(j => j.update).filter(update => !isEmptyUpdate(update))
|
jobs.map(j => j.update).filter(update => !isEmptyUpdate(update))
|
||||||
);
|
);
|
||||||
@@ -147,19 +153,22 @@ export class DocSyncPeer {
|
|||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock: timestamp,
|
remoteClock: timestamp,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
await this.actions.updatePushedClock(docId, maxClock);
|
await this.syncMetadata.setPeerPushedClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: maxClock,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
pullAndPush: async (docId: string, signal?: AbortSignal) => {
|
pullAndPush: async (docId: string, signal?: AbortSignal) => {
|
||||||
const docRecord = await this.local.getDoc(docId);
|
const localDocRecord = await this.local.getDoc(docId);
|
||||||
|
|
||||||
const stateVector =
|
const stateVector =
|
||||||
docRecord && !isEmptyUpdate(docRecord.bin)
|
localDocRecord && !isEmptyUpdate(localDocRecord.bin)
|
||||||
? encodeStateVectorFromUpdate(docRecord.bin)
|
? encodeStateVectorFromUpdate(localDocRecord.bin)
|
||||||
: new Uint8Array();
|
: new Uint8Array();
|
||||||
const remoteDocRecord = await this.remote.getDocDiff(docId, stateVector);
|
const remoteDocRecord = await this.remote.getDocDiff(docId, stateVector);
|
||||||
|
|
||||||
@@ -167,12 +176,12 @@ export class DocSyncPeer {
|
|||||||
const {
|
const {
|
||||||
missing: newData,
|
missing: newData,
|
||||||
state: serverStateVector,
|
state: serverStateVector,
|
||||||
timestamp: serverClock,
|
timestamp: remoteClock,
|
||||||
} = remoteDocRecord;
|
} = remoteDocRecord;
|
||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock,
|
remoteClock,
|
||||||
});
|
});
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
const { timestamp: localClock } = await this.local.pushDocUpdate(
|
const { timestamp: localClock } = await this.local.pushDocUpdate(
|
||||||
@@ -183,14 +192,17 @@ export class DocSyncPeer {
|
|||||||
this.uniqueId
|
this.uniqueId
|
||||||
);
|
);
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
await this.actions.updatePulledRemoteClock(docId, serverClock);
|
await this.syncMetadata.setPeerPulledRemoteClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: remoteClock,
|
||||||
|
});
|
||||||
const diff =
|
const diff =
|
||||||
docRecord && serverStateVector && serverStateVector.length > 0
|
localDocRecord && serverStateVector && serverStateVector.length > 0
|
||||||
? diffUpdate(docRecord.bin, serverStateVector)
|
? diffUpdate(localDocRecord.bin, serverStateVector)
|
||||||
: docRecord?.bin;
|
: localDocRecord?.bin;
|
||||||
if (diff && !isEmptyUpdate(diff)) {
|
if (diff && !isEmptyUpdate(diff)) {
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
const { timestamp: serverClock } = await this.remote.pushDocUpdate(
|
const { timestamp: remoteClock } = await this.remote.pushDocUpdate(
|
||||||
{
|
{
|
||||||
bin: diff,
|
bin: diff,
|
||||||
docId,
|
docId,
|
||||||
@@ -200,18 +212,21 @@ export class DocSyncPeer {
|
|||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock,
|
remoteClock,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
await this.actions.updatePushedClock(docId, localClock);
|
await this.syncMetadata.setPeerPushedClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: localClock,
|
||||||
|
});
|
||||||
} else {
|
} else {
|
||||||
if (docRecord) {
|
if (localDocRecord) {
|
||||||
if (!isEmptyUpdate(docRecord.bin)) {
|
if (!isEmptyUpdate(localDocRecord.bin)) {
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
const { timestamp: serverClock } = await this.remote.pushDocUpdate(
|
const { timestamp: remoteClock } = await this.remote.pushDocUpdate(
|
||||||
{
|
{
|
||||||
bin: docRecord.bin,
|
bin: localDocRecord.bin,
|
||||||
docId,
|
docId,
|
||||||
},
|
},
|
||||||
this.uniqueId
|
this.uniqueId
|
||||||
@@ -219,10 +234,13 @@ export class DocSyncPeer {
|
|||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock,
|
remoteClock,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
await this.actions.updatePushedClock(docId, docRecord.timestamp);
|
await this.syncMetadata.setPeerPushedClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: localDocRecord.timestamp,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
@@ -237,7 +255,7 @@ export class DocSyncPeer {
|
|||||||
if (!serverDoc) {
|
if (!serverDoc) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
const { missing: newData, timestamp: serverClock } = serverDoc;
|
const { missing: newData, timestamp: remoteClock } = serverDoc;
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
await this.local.pushDocUpdate(
|
await this.local.pushDocUpdate(
|
||||||
{
|
{
|
||||||
@@ -247,11 +265,14 @@ export class DocSyncPeer {
|
|||||||
this.uniqueId
|
this.uniqueId
|
||||||
);
|
);
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
await this.actions.updatePulledRemoteClock(docId, serverClock);
|
await this.syncMetadata.setPeerPulledRemoteClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: remoteClock,
|
||||||
|
});
|
||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock,
|
remoteClock: remoteClock,
|
||||||
});
|
});
|
||||||
},
|
},
|
||||||
save: async (
|
save: async (
|
||||||
@@ -259,8 +280,8 @@ export class DocSyncPeer {
|
|||||||
jobs: (Job & { type: 'save' })[],
|
jobs: (Job & { type: 'save' })[],
|
||||||
signal?: AbortSignal
|
signal?: AbortSignal
|
||||||
) => {
|
) => {
|
||||||
const serverClock = jobs.reduce(
|
const remoteClock = jobs.reduce(
|
||||||
(a, b) => (a.getTime() > b.serverClock.getTime() ? a : b.serverClock),
|
(a, b) => (a.getTime() > b.remoteClock.getTime() ? a : b.remoteClock),
|
||||||
new Date(0)
|
new Date(0)
|
||||||
);
|
);
|
||||||
if (this.status.connectedDocs.has(docId)) {
|
if (this.status.connectedDocs.has(docId)) {
|
||||||
@@ -282,7 +303,10 @@ export class DocSyncPeer {
|
|||||||
);
|
);
|
||||||
throwIfAborted(signal);
|
throwIfAborted(signal);
|
||||||
|
|
||||||
await this.actions.updatePulledRemoteClock(docId, serverClock);
|
await this.syncMetadata.setPeerPulledRemoteClock(this.remote.peer, {
|
||||||
|
docId,
|
||||||
|
timestamp: remoteClock,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
@@ -298,29 +322,6 @@ export class DocSyncPeer {
|
|||||||
this.statusUpdatedSubject$.next(docId);
|
this.statusUpdatedSubject$.next(docId);
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
updatePushedClock: async (docId: string, pushedClock: Date) => {
|
|
||||||
const updated = this.status.pushedClocks.setIfBigger(docId, pushedClock);
|
|
||||||
if (updated) {
|
|
||||||
await this.syncMetadata.setPeerPushedClock(this.remote.peer, {
|
|
||||||
docId,
|
|
||||||
timestamp: pushedClock,
|
|
||||||
});
|
|
||||||
this.statusUpdatedSubject$.next(docId);
|
|
||||||
}
|
|
||||||
},
|
|
||||||
updatePulledRemoteClock: async (docId: string, pulledClock: Date) => {
|
|
||||||
const updated = this.status.pulledRemoteClocks.setIfBigger(
|
|
||||||
docId,
|
|
||||||
pulledClock
|
|
||||||
);
|
|
||||||
if (updated) {
|
|
||||||
await this.syncMetadata.setPeerPulledRemoteClock(this.remote.peer, {
|
|
||||||
docId,
|
|
||||||
timestamp: pulledClock,
|
|
||||||
});
|
|
||||||
this.statusUpdatedSubject$.next(docId);
|
|
||||||
}
|
|
||||||
},
|
|
||||||
addDoc: (docId: string) => {
|
addDoc: (docId: string) => {
|
||||||
if (!this.status.docs.has(docId)) {
|
if (!this.status.docs.has(docId)) {
|
||||||
this.status.docs.add(docId);
|
this.status.docs.add(docId);
|
||||||
@@ -370,7 +371,7 @@ export class DocSyncPeer {
|
|||||||
this.schedule({
|
this.schedule({
|
||||||
type: 'save',
|
type: 'save',
|
||||||
docId,
|
docId,
|
||||||
serverClock: remoteClock,
|
remoteClock: remoteClock,
|
||||||
update,
|
update,
|
||||||
});
|
});
|
||||||
},
|
},
|
||||||
@@ -396,8 +397,6 @@ export class DocSyncPeer {
|
|||||||
connectedDocs: new Set(),
|
connectedDocs: new Set(),
|
||||||
jobDocQueue: new AsyncPriorityQueue(),
|
jobDocQueue: new AsyncPriorityQueue(),
|
||||||
jobMap: new Map(),
|
jobMap: new Map(),
|
||||||
pulledRemoteClocks: new ClockMap(new Map()),
|
|
||||||
pushedClocks: new ClockMap(new Map()),
|
|
||||||
remoteClocks: new ClockMap(new Map()),
|
remoteClocks: new ClockMap(new Map()),
|
||||||
syncing: false,
|
syncing: false,
|
||||||
// tell ui to show retrying status
|
// tell ui to show retrying status
|
||||||
@@ -478,7 +477,13 @@ export class DocSyncPeer {
|
|||||||
// subscribe local doc updates
|
// subscribe local doc updates
|
||||||
disposes.push(
|
disposes.push(
|
||||||
this.local.subscribeDocUpdate((update, origin) => {
|
this.local.subscribeDocUpdate((update, origin) => {
|
||||||
if (origin === this.uniqueId) {
|
if (
|
||||||
|
origin === this.uniqueId ||
|
||||||
|
origin?.startsWith(
|
||||||
|
`sync:${this.local.peer}:${this.remote.peer}:`
|
||||||
|
// skip if local and remote is same
|
||||||
|
)
|
||||||
|
) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
this.events.localUpdated({
|
this.events.localUpdated({
|
||||||
@@ -517,19 +522,6 @@ export class DocSyncPeer {
|
|||||||
for (const [id, v] of Object.entries(cachedClocks)) {
|
for (const [id, v] of Object.entries(cachedClocks)) {
|
||||||
this.status.remoteClocks.set(id, v);
|
this.status.remoteClocks.set(id, v);
|
||||||
}
|
}
|
||||||
const pulledClocks = await this.syncMetadata.getPeerPulledRemoteClocks(
|
|
||||||
this.remote.peer
|
|
||||||
);
|
|
||||||
for (const [id, v] of Object.entries(pulledClocks)) {
|
|
||||||
this.status.pulledRemoteClocks.set(id, v);
|
|
||||||
}
|
|
||||||
const pushedClocks = await this.syncMetadata.getPeerPushedClocks(
|
|
||||||
this.remote.peer
|
|
||||||
);
|
|
||||||
throwIfAborted(signal);
|
|
||||||
for (const [id, v] of Object.entries(pushedClocks)) {
|
|
||||||
this.status.pushedClocks.set(id, v);
|
|
||||||
}
|
|
||||||
this.statusUpdatedSubject$.next(true);
|
this.statusUpdatedSubject$.next(true);
|
||||||
|
|
||||||
// get new clocks from server
|
// get new clocks from server
|
||||||
|
|||||||
Reference in New Issue
Block a user