fix(server): avoid snapshot write conflict (#5174)
This commit is contained in:
@@ -11,6 +11,7 @@ import { chunk } from 'lodash-es';
|
|||||||
import { defer, retry } from 'rxjs';
|
import { defer, retry } from 'rxjs';
|
||||||
import {
|
import {
|
||||||
applyUpdate,
|
applyUpdate,
|
||||||
|
decodeStateVector,
|
||||||
Doc,
|
Doc,
|
||||||
encodeStateAsUpdate,
|
encodeStateAsUpdate,
|
||||||
encodeStateVector,
|
encodeStateVector,
|
||||||
@@ -40,6 +41,36 @@ function compare(yBinary: Buffer, jwstBinary: Buffer, strict = false): boolean {
|
|||||||
return compare(yBinary, yBinary2, true);
|
return compare(yBinary, yBinary2, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Detect whether rhs state is newer than lhs state.
|
||||||
|
*
|
||||||
|
* How could we tell a state is newer:
|
||||||
|
*
|
||||||
|
* i. if the state vector size is larger, it's newer
|
||||||
|
* ii. if the state vector size is same, compare each client's state
|
||||||
|
*/
|
||||||
|
function isStateNewer(lhs: Buffer, rhs: Buffer): boolean {
|
||||||
|
const lhsVector = decodeStateVector(lhs);
|
||||||
|
const rhsVector = decodeStateVector(rhs);
|
||||||
|
|
||||||
|
if (lhsVector.size < rhsVector.size) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
for (const [client, state] of lhsVector) {
|
||||||
|
const rstate = rhsVector.get(client);
|
||||||
|
if (!rstate) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (state < rstate) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
function isEmptyBuffer(buf: Buffer): boolean {
|
function isEmptyBuffer(buf: Buffer): boolean {
|
||||||
return (
|
return (
|
||||||
buf.length === 0 ||
|
buf.length === 0 ||
|
||||||
@@ -374,23 +405,17 @@ export class DocManager implements OnModuleInit, OnModuleDestroy {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const { id, workspaceId } = candidate;
|
const { id, workspaceId } = candidate;
|
||||||
// acquire lock
|
|
||||||
const ok = await this.lockUpdatesForAutoSquash(workspaceId, id);
|
|
||||||
|
|
||||||
if (!ok) {
|
await this.lockUpdatesForAutoSquash(workspaceId, id, async () => {
|
||||||
return;
|
try {
|
||||||
}
|
await this._get(workspaceId, id);
|
||||||
|
} catch (e) {
|
||||||
try {
|
this.logger.error(
|
||||||
await this._get(workspaceId, id);
|
`Failed to apply updates for workspace: ${workspaceId}, guid: ${id}`
|
||||||
} catch (e) {
|
);
|
||||||
this.logger.error(
|
this.logger.error(e);
|
||||||
`Failed to apply updates for workspace: ${workspaceId}, guid: ${id}`
|
}
|
||||||
);
|
});
|
||||||
this.logger.error(e);
|
|
||||||
} finally {
|
|
||||||
await this.unlockUpdatesForAutoSquash(workspaceId, id);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private async getAutoSquashCandidate() {
|
private async getAutoSquashCandidate() {
|
||||||
@@ -414,34 +439,67 @@ export class DocManager implements OnModuleInit, OnModuleDestroy {
|
|||||||
doc: Doc,
|
doc: Doc,
|
||||||
initialSeq?: number
|
initialSeq?: number
|
||||||
) {
|
) {
|
||||||
const blob = Buffer.from(encodeStateAsUpdate(doc));
|
return this.lockSnapshotForUpsert(workspaceId, guid, async () => {
|
||||||
const state = Buffer.from(encodeStateVector(doc));
|
const blob = Buffer.from(encodeStateAsUpdate(doc));
|
||||||
|
|
||||||
if (isEmptyBuffer(blob)) {
|
if (isEmptyBuffer(blob)) {
|
||||||
return;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
await this.db.snapshot.upsert({
|
const state = Buffer.from(encodeStateVector(doc));
|
||||||
select: {
|
|
||||||
seq: true,
|
return await this.db.$transaction(async db => {
|
||||||
},
|
const snapshot = await db.snapshot.findUnique({
|
||||||
where: {
|
where: {
|
||||||
id_workspaceId: {
|
id_workspaceId: {
|
||||||
id: guid,
|
id: guid,
|
||||||
workspaceId,
|
workspaceId,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
create: {
|
});
|
||||||
id: guid,
|
|
||||||
workspaceId,
|
// update
|
||||||
blob,
|
if (snapshot) {
|
||||||
state,
|
// only update if state is newer
|
||||||
seq: initialSeq,
|
if (isStateNewer(snapshot.state ?? Buffer.from([0]), state)) {
|
||||||
},
|
await db.snapshot.update({
|
||||||
update: {
|
select: {
|
||||||
blob,
|
seq: true,
|
||||||
state,
|
},
|
||||||
},
|
where: {
|
||||||
|
id_workspaceId: {
|
||||||
|
workspaceId,
|
||||||
|
id: guid,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
data: {
|
||||||
|
blob,
|
||||||
|
state,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
return true;
|
||||||
|
} else {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// create
|
||||||
|
await db.snapshot.create({
|
||||||
|
select: {
|
||||||
|
seq: true,
|
||||||
|
},
|
||||||
|
data: {
|
||||||
|
id: guid,
|
||||||
|
workspaceId,
|
||||||
|
blob,
|
||||||
|
state,
|
||||||
|
seq: initialSeq,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
});
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -484,21 +542,26 @@ export class DocManager implements OnModuleInit, OnModuleDestroy {
|
|||||||
this.event.emit('doc:manager:snapshot:beforeUpdate', snapshot);
|
this.event.emit('doc:manager:snapshot:beforeUpdate', snapshot);
|
||||||
}
|
}
|
||||||
|
|
||||||
await this.upsert(workspaceId, id, doc, last.seq);
|
const done = await this.upsert(workspaceId, id, doc, last.seq);
|
||||||
this.logger.debug(
|
|
||||||
`Squashed ${updates.length} updates for ${id} in workspace ${workspaceId}`
|
if (done) {
|
||||||
);
|
this.logger.debug(
|
||||||
await this.db.update.deleteMany({
|
`Squashed ${updates.length} updates for ${id} in workspace ${workspaceId}`
|
||||||
where: {
|
);
|
||||||
id,
|
|
||||||
workspaceId,
|
await this.db.update.deleteMany({
|
||||||
seq: {
|
where: {
|
||||||
in: updates.map(u => u.seq),
|
id,
|
||||||
},
|
workspaceId,
|
||||||
},
|
seq: {
|
||||||
});
|
in: updates.map(u => u.seq),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
await this.updateCachedUpdatesCount(workspaceId, id, -updates.length);
|
||||||
|
}
|
||||||
|
|
||||||
await this.updateCachedUpdatesCount(workspaceId, id, -updates.length);
|
|
||||||
return doc;
|
return doc;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -581,22 +644,44 @@ export class DocManager implements OnModuleInit, OnModuleDestroy {
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
private async lockUpdatesForAutoSquash(workspaceId: string, guid: string) {
|
private async doWithLock<T>(lock: string, job: () => Promise<T>) {
|
||||||
return this.cache.setnx(
|
const acquired = await this.cache.setnx(lock, 1, {
|
||||||
|
ttl: 60 * 1000,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!acquired) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
return await job();
|
||||||
|
} finally {
|
||||||
|
await this.cache.delete(lock).catch(e => {
|
||||||
|
// safe, the lock will be expired when ttl ends
|
||||||
|
this.logger.error(`Failed to release lock ${lock}`, e);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async lockUpdatesForAutoSquash<T>(
|
||||||
|
workspaceId: string,
|
||||||
|
guid: string,
|
||||||
|
job: () => Promise<T>
|
||||||
|
) {
|
||||||
|
return this.doWithLock(
|
||||||
`doc:manager:updates-lock:${workspaceId}::${guid}`,
|
`doc:manager:updates-lock:${workspaceId}::${guid}`,
|
||||||
1,
|
job
|
||||||
{
|
|
||||||
ttl: 60 * 1000,
|
|
||||||
}
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async unlockUpdatesForAutoSquash(workspaceId: string, guid: string) {
|
async lockSnapshotForUpsert<T>(
|
||||||
return this.cache
|
workspaceId: string,
|
||||||
.delete(`doc:manager:updates-lock:${workspaceId}::${guid}`)
|
guid: string,
|
||||||
.catch(e => {
|
job: () => Promise<T>
|
||||||
// safe, the lock will be expired when ttl ends
|
) {
|
||||||
this.logger.error('Failed to release updates lock', e);
|
return this.doWithLock(
|
||||||
});
|
`doc:manager:snapshot-lock:${workspaceId}::${guid}`,
|
||||||
|
job
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,7 +6,12 @@ import { Test, TestingModule } from '@nestjs/testing';
|
|||||||
import test from 'ava';
|
import test from 'ava';
|
||||||
import { register } from 'prom-client';
|
import { register } from 'prom-client';
|
||||||
import * as Sinon from 'sinon';
|
import * as Sinon from 'sinon';
|
||||||
import { Doc as YDoc, encodeStateAsUpdate } from 'yjs';
|
import {
|
||||||
|
applyUpdate,
|
||||||
|
decodeStateVector,
|
||||||
|
Doc as YDoc,
|
||||||
|
encodeStateAsUpdate,
|
||||||
|
} from 'yjs';
|
||||||
|
|
||||||
import { CacheModule } from '../src/cache';
|
import { CacheModule } from '../src/cache';
|
||||||
import { Config, ConfigModule } from '../src/config';
|
import { Config, ConfigModule } from '../src/config';
|
||||||
@@ -283,3 +288,73 @@ test('should throw if meet max retry times', async t => {
|
|||||||
);
|
);
|
||||||
t.is(stub.callCount, 5);
|
t.is(stub.callCount, 5);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('should not update snapshot if state is outdated', async t => {
|
||||||
|
const db = m.get(PrismaService);
|
||||||
|
const manager = m.get(DocManager);
|
||||||
|
|
||||||
|
await db.snapshot.create({
|
||||||
|
data: {
|
||||||
|
id: '2',
|
||||||
|
workspaceId: '2',
|
||||||
|
blob: Buffer.from([0, 0]),
|
||||||
|
seq: 1,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
const doc = new YDoc();
|
||||||
|
const text = doc.getText('content');
|
||||||
|
const updates: Buffer[] = [];
|
||||||
|
|
||||||
|
doc.on('update', update => {
|
||||||
|
updates.push(Buffer.from(update));
|
||||||
|
});
|
||||||
|
|
||||||
|
text.insert(0, 'hello');
|
||||||
|
text.insert(5, 'world');
|
||||||
|
text.insert(5, ' ');
|
||||||
|
|
||||||
|
await Promise.all(updates.map(update => manager.push('2', '2', update)));
|
||||||
|
|
||||||
|
const updateWith3Records = await manager.getUpdates('2', '2');
|
||||||
|
text.insert(11, '!');
|
||||||
|
await manager.push('2', '2', updates[3]);
|
||||||
|
const updateWith4Records = await manager.getUpdates('2', '2');
|
||||||
|
|
||||||
|
// Simulation:
|
||||||
|
// Node A get 3 updates and squash them at time 1, will finish at time 10
|
||||||
|
// Node B get 4 updates and squash them at time 3, will finish at time 8
|
||||||
|
// Node B finish the squash first, and update the snapshot
|
||||||
|
// Node A finish the squash later, and update the snapshot to an outdated state
|
||||||
|
// Time: ---------------------->
|
||||||
|
// A: ^get ^upsert
|
||||||
|
// B: ^get ^upsert
|
||||||
|
//
|
||||||
|
// We should avoid such situation
|
||||||
|
// @ts-expect-error private
|
||||||
|
await manager.squash(updateWith4Records, null);
|
||||||
|
// @ts-expect-error private
|
||||||
|
await manager.squash(updateWith3Records, null);
|
||||||
|
|
||||||
|
const result = await db.snapshot.findUnique({
|
||||||
|
where: {
|
||||||
|
id_workspaceId: {
|
||||||
|
id: '2',
|
||||||
|
workspaceId: '2',
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!result) {
|
||||||
|
t.fail('snapshot not found');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const state = decodeStateVector(result.state!);
|
||||||
|
t.is(state.get(doc.clientID), 12);
|
||||||
|
|
||||||
|
const d = new YDoc();
|
||||||
|
applyUpdate(d, result.blob!);
|
||||||
|
|
||||||
|
const dtext = d.getText('content');
|
||||||
|
t.is(dtext.toString(), 'hello world!');
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user