fix(nbstore): adjust doc sync logic (#10342)

This commit is contained in:
EYHN
2025-02-21 08:17:07 +00:00
parent f3218ab3bc
commit 244d683d83
2 changed files with 47 additions and 27 deletions

View File

@@ -399,23 +399,31 @@ export class DocFrontend {
this.statusUpdatedSubject$.next(job.docId); this.statusUpdatedSubject$.next(job.docId);
} }
/**
* skip listen doc update when apply update
*/
private skipDocUpdate = false;
applyUpdate(docId: string, update: Uint8Array) { applyUpdate(docId: string, update: Uint8Array) {
const doc = this.status.docs.get(docId); const doc = this.status.docs.get(docId);
if (doc && !isEmptyUpdate(update)) { if (doc && !isEmptyUpdate(update)) {
try { try {
this.skipDocUpdate = true;
applyUpdate(doc, update, NBSTORE_ORIGIN); applyUpdate(doc, update, NBSTORE_ORIGIN);
} catch (err) { } catch (err) {
console.error('failed to apply update yjs doc', err); console.error('failed to apply update yjs doc', err);
} finally {
this.skipDocUpdate = false;
} }
} }
} }
private readonly handleDocUpdate = ( private readonly handleDocUpdate = (
update: Uint8Array, update: Uint8Array,
origin: any, _origin: any,
doc: YDoc doc: YDoc
) => { ) => {
if (origin === NBSTORE_ORIGIN) { if (this.skipDocUpdate) {
return; return;
} }
if (!this.status.docs.has(doc.guid)) { if (!this.status.docs.has(doc.guid)) {

View File

@@ -17,7 +17,7 @@ type Job =
| { | {
type: 'push'; type: 'push';
docId: string; docId: string;
update: Uint8Array; update?: Uint8Array;
clock: Date; clock: Date;
} }
| { | {
@@ -278,7 +278,9 @@ export class DocSyncPeer {
); );
const merged = await this.mergeUpdates( const merged = await this.mergeUpdates(
jobs.map(j => j.update).filter(update => !isEmptyUpdate(update)) jobs
.map(j => j.update ?? new Uint8Array())
.filter(update => !isEmptyUpdate(update))
); );
if (!isEmptyUpdate(merged)) { if (!isEmptyUpdate(merged)) {
const { timestamp } = await this.remote.pushDocUpdate( const { timestamp } = await this.remote.pushDocUpdate(
@@ -316,11 +318,6 @@ export class DocSyncPeer {
state: serverStateVector, state: serverStateVector,
timestamp: remoteClock, timestamp: remoteClock,
} = remoteDocRecord; } = remoteDocRecord;
this.schedule({
type: 'save',
docId,
remoteClock,
});
throwIfAborted(signal); throwIfAborted(signal);
const { timestamp: localClock } = await this.local.pushDocUpdate( const { timestamp: localClock } = await this.local.pushDocUpdate(
{ {
@@ -359,9 +356,10 @@ export class DocSyncPeer {
}); });
} }
throwIfAborted(signal); throwIfAborted(signal);
await this.syncMetadata.setPeerPushedClock(this.peerId, { this.schedule({
type: 'push',
docId, docId,
timestamp: localClock, clock: localClock,
}); });
} else { } else {
if (localDocRecord) { if (localDocRecord) {
@@ -380,6 +378,11 @@ export class DocSyncPeer {
remoteClock, remoteClock,
}); });
} }
this.schedule({
type: 'push',
docId,
clock: localDocRecord.timestamp,
});
await this.syncMetadata.setPeerPushedClock(this.peerId, { await this.syncMetadata.setPeerPushedClock(this.peerId, {
docId, docId,
timestamp: localDocRecord.timestamp, timestamp: localDocRecord.timestamp,
@@ -400,7 +403,7 @@ export class DocSyncPeer {
} }
const { missing: newData, timestamp: remoteClock } = serverDoc; const { missing: newData, timestamp: remoteClock } = serverDoc;
throwIfAborted(signal); throwIfAborted(signal);
await this.local.pushDocUpdate( const { timestamp } = await this.local.pushDocUpdate(
{ {
docId, docId,
bin: newData, bin: newData,
@@ -413,9 +416,9 @@ export class DocSyncPeer {
timestamp: remoteClock, timestamp: remoteClock,
}); });
this.schedule({ this.schedule({
type: 'save', type: 'push',
docId, docId,
remoteClock: remoteClock, clock: timestamp,
}); });
}, },
save: async ( save: async (
@@ -438,13 +441,20 @@ export class DocSyncPeer {
throwIfAborted(signal); throwIfAborted(signal);
if (!isEmptyUpdate(update)) { if (!isEmptyUpdate(update)) {
await this.local.pushDocUpdate( const { timestamp } = await this.local.pushDocUpdate(
{ {
docId, docId,
bin: update, bin: update,
}, },
this.uniqueId this.uniqueId
); );
// schedule push job to mark the timestamp as pushed timestamp
this.schedule({
type: 'push',
docId,
clock: timestamp,
});
} }
throwIfAborted(signal); throwIfAborted(signal);
@@ -457,15 +467,9 @@ export class DocSyncPeer {
}); });
private readonly actions = { private readonly actions = {
updateRemoteClock: async (docId: string, remoteClock: Date) => { updateRemoteClock: (docId: string, remoteClock: Date) => {
const updated = this.status.remoteClocks.setIfBigger(docId, remoteClock); this.status.remoteClocks.setIfBigger(docId, remoteClock);
if (updated) {
await this.syncMetadata.setPeerRemoteClock(this.peerId, {
docId,
timestamp: remoteClock,
});
this.statusUpdatedSubject$.next(docId); this.statusUpdatedSubject$.next(docId);
}
}, },
addDoc: (docId: string) => { addDoc: (docId: string) => {
if (!this.status.docs.has(docId)) { if (!this.status.docs.has(docId)) {
@@ -511,6 +515,7 @@ export class DocSyncPeer {
}) => { }) => {
// try add doc for new doc // try add doc for new doc
this.actions.addDoc(docId); this.actions.addDoc(docId);
this.actions.updateRemoteClock(docId, remoteClock);
// schedule push job // schedule push job
this.schedule({ this.schedule({
@@ -684,7 +689,14 @@ export class DocSyncPeer {
const maxClockValue = this.status.remoteClocks.max; const maxClockValue = this.status.remoteClocks.max;
const newClocks = await this.remote.getDocTimestamps(maxClockValue); const newClocks = await this.remote.getDocTimestamps(maxClockValue);
for (const [id, v] of Object.entries(newClocks)) { for (const [id, v] of Object.entries(newClocks)) {
await this.actions.updateRemoteClock(id, v); this.actions.updateRemoteClock(id, v);
}
for (const [id, v] of Object.entries(newClocks)) {
await this.syncMetadata.setPeerRemoteClock(this.peerId, {
docId: id,
timestamp: v,
});
} }
// add all docs from remote // add all docs from remote
@@ -778,9 +790,9 @@ export class DocSyncPeer {
}; };
} }
protected mergeUpdates(updates: Uint8Array[]) { protected mergeUpdates = (updates: Uint8Array[]) => {
const merge = this.options?.mergeUpdates ?? mergeUpdates; const merge = this.options?.mergeUpdates ?? mergeUpdates;
return merge(updates.filter(bin => !isEmptyUpdate(bin))); return merge(updates.filter(bin => !isEmptyUpdate(bin)));
} };
} }