refactor(nbstore): improve doc state management (#12359)

Move the `waitForSynced` method from `frontend` to `nbstore worker` to make the wait more reliable

<!-- This is an auto-generated comment: release notes by coderabbit.ai -->
## Summary by CodeRabbit

- **New Features**
  - Added explicit tracking of document updating state to indicate when data is being applied or saved.
  - Introduced new methods to wait for update and synchronization completion with abort support.

- **Improvements**
  - Applied throttling with leading and trailing emissions to state observables for smoother UI updates.
  - Refined synchronization waiting logic for clearer separation between update completion and sync completion.
  - Removed throttling in workspace selector component for more immediate state feedback.
  - Updated import and clipper services to use the new synchronization waiting methods.
  - Simplified asynchronous waiting logic in indexer synchronization methods.

- **Bug Fixes**
  - Enhanced accuracy and reliability of document update and sync status indicators.

- **Tests**
  - Increased wait timeout in avatar selection test to improve stability.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
EYHN
2025-05-20 07:33:20 +00:00
parent 151f499154
commit 59ef4b227b
10 changed files with 157 additions and 167 deletions

View File

@@ -1,13 +1,16 @@
import { groupBy } from 'lodash-es'; import { groupBy } from 'lodash-es';
import { nanoid } from 'nanoid'; import { nanoid } from 'nanoid';
import type { Subscription } from 'rxjs';
import { import {
combineLatest, combineLatest,
filter,
first,
lastValueFrom,
map, map,
Observable, Observable,
ReplaySubject, ReplaySubject,
share, share,
Subject, Subject,
throttleTime,
} from 'rxjs'; } from 'rxjs';
import { import {
applyUpdate, applyUpdate,
@@ -22,6 +25,7 @@ import type { DocRecord, DocStorage } from '../storage';
import type { DocSync } from '../sync/doc'; import type { DocSync } from '../sync/doc';
import { AsyncPriorityQueue } from '../utils/async-priority-queue'; import { AsyncPriorityQueue } from '../utils/async-priority-queue';
import { isEmptyUpdate } from '../utils/is-empty-update'; import { isEmptyUpdate } from '../utils/is-empty-update';
import { takeUntilAbort } from '../utils/take-until-abort';
import { MANUALLY_STOP, throwIfAborted } from '../utils/throw-if-aborted'; import { MANUALLY_STOP, throwIfAborted } from '../utils/throw-if-aborted';
const NBSTORE_ORIGIN = 'nbstore-frontend'; const NBSTORE_ORIGIN = 'nbstore-frontend';
@@ -86,6 +90,10 @@ export type DocFrontendState = {
* number of docs that have been loaded to yjs doc instance * number of docs that have been loaded to yjs doc instance
*/ */
loaded: number; loaded: number;
/**
* some data is being applied to yjs doc instance, or some data is being saved to local doc storage
*/
updating: boolean;
/** /**
* number of docs that are syncing with remote peers * number of docs that are syncing with remote peers
*/ */
@@ -128,7 +136,7 @@ export class DocFrontend {
readonly options: DocFrontendOptions = {} readonly options: DocFrontendOptions = {}
) {} ) {}
docState$(docId: string): Observable<DocFrontendDocState> { private _docState$(docId: string): Observable<DocFrontendDocState> {
const frontendState$ = new Observable<{ const frontendState$ = new Observable<{
ready: boolean; ready: boolean;
loaded: boolean; loaded: boolean;
@@ -160,24 +168,38 @@ export class DocFrontend {
); );
} }
state$ = combineLatest([ docState$(docId: string): Observable<DocFrontendDocState> {
new Observable<{ total: number; loaded: number }>(subscriber => { return this._docState$(docId).pipe(
const next = () => { throttleTime(1000, undefined, {
subscriber.next({ trailing: true,
total: this.status.docs.size, leading: true,
loaded: this.status.connectedDocs.size, })
}); );
}; }
next();
return this.statusUpdatedSubject$.subscribe(() => { private readonly _state$ = combineLatest([
new Observable<{ total: number; loaded: number; updating: boolean }>(
subscriber => {
const next = () => {
subscriber.next({
total: this.status.docs.size,
loaded: this.status.connectedDocs.size,
updating:
this.status.jobMap.size > 0 || this.status.currentJob !== null,
});
};
next(); next();
}); return this.statusUpdatedSubject$.subscribe(() => {
}), next();
});
}
),
this.sync.state$, this.sync.state$,
]).pipe( ]).pipe(
map(([frontend, sync]) => ({ map(([frontend, sync]) => ({
total: sync.total ?? frontend.total, total: sync.total ?? frontend.total,
loaded: frontend.loaded, loaded: frontend.loaded,
updating: frontend.updating,
syncing: sync.syncing, syncing: sync.syncing,
synced: sync.synced, synced: sync.synced,
syncRetrying: sync.retrying, syncRetrying: sync.retrying,
@@ -188,6 +210,13 @@ export class DocFrontend {
}) })
) satisfies Observable<DocFrontendState>; ) satisfies Observable<DocFrontendState>;
state$ = this._state$.pipe(
throttleTime(1000, undefined, {
leading: true,
trailing: true,
})
);
start() { start() {
if (this.abort.signal.aborted) { if (this.abort.signal.aborted) {
throw new Error('doc frontend can only start once'); throw new Error('doc frontend can only start once');
@@ -463,96 +492,43 @@ ${changedList}
return merge(updates.filter(bin => !isEmptyUpdate(bin))); return merge(updates.filter(bin => !isEmptyUpdate(bin)));
} }
async waitForSynced(abort?: AbortSignal) { async waitForUpdated(docId?: string, abort?: AbortSignal) {
let sub: Subscription | undefined = undefined; const source$: Observable<DocFrontendDocState | DocFrontendState> = docId
return Promise.race([ ? this._docState$(docId)
new Promise<void>(resolve => { : this._state$;
sub = this.state$?.subscribe(status => { await lastValueFrom(
if (status.synced) { source$.pipe(
resolve(); filter(status => !status.updating),
} takeUntilAbort(abort),
}); first()
}), )
new Promise<void>((_, reject) => { );
if (abort?.aborted) { return;
reject(abort?.reason);
}
abort?.addEventListener('abort', () => {
reject(abort.reason);
});
}),
]).finally(() => {
sub?.unsubscribe();
});
} }
async waitForDocLoaded(docId: string, abort?: AbortSignal) { async waitForDocLoaded(docId: string, abort?: AbortSignal) {
let sub: Subscription | undefined = undefined; await lastValueFrom(
return Promise.race([ this._docState$(docId).pipe(
new Promise<void>(resolve => { filter(state => state.loaded),
sub = this.docState$(docId).subscribe(state => { takeUntilAbort(abort),
if (state.loaded) { first()
resolve(); )
} );
});
}),
new Promise<void>((_, reject) => {
if (abort?.aborted) {
reject(abort?.reason);
}
abort?.addEventListener('abort', () => {
reject(abort.reason);
});
}),
]).finally(() => {
sub?.unsubscribe();
});
} }
async waitForDocSynced(docId: string, abort?: AbortSignal) { async waitForSynced(docId?: string, abort?: AbortSignal) {
let sub: Subscription | undefined = undefined; await this.waitForUpdated(docId, abort);
return Promise.race([ await this.sync.waitForSynced(docId, abort);
new Promise<void>(resolve => {
sub = this.docState$(docId).subscribe(state => {
if (state.synced && !state.updating) {
resolve();
}
});
}),
new Promise<void>((_, reject) => {
if (abort?.aborted) {
reject(abort?.reason);
}
abort?.addEventListener('abort', () => {
reject(abort.reason);
});
}),
]).finally(() => {
sub?.unsubscribe();
});
} }
async waitForDocReady(docId: string, abort?: AbortSignal) { async waitForDocReady(docId: string, abort?: AbortSignal) {
let sub: Subscription | undefined = undefined; await lastValueFrom(
return Promise.race([ this._docState$(docId).pipe(
new Promise<void>(resolve => { filter(state => state.ready),
sub = this.docState$(docId).subscribe(state => { takeUntilAbort(abort),
if (state.ready) { first()
resolve(); )
} );
});
}),
new Promise<void>((_, reject) => {
if (abort?.aborted) {
reject(abort?.reason);
}
abort?.addEventListener('abort', () => {
reject(abort.reason);
});
}),
]).finally(() => {
sub?.unsubscribe();
});
} }
async resetSync() { async resetSync() {

View File

@@ -1,9 +1,20 @@
import type { Observable } from 'rxjs'; import type { Observable } from 'rxjs';
import { combineLatest, map, of, ReplaySubject, share } from 'rxjs'; import {
combineLatest,
filter,
first,
lastValueFrom,
map,
of,
ReplaySubject,
share,
throttleTime,
} from 'rxjs';
import type { DocStorage, DocSyncStorage } from '../../storage'; import type { DocStorage, DocSyncStorage } from '../../storage';
import { DummyDocStorage } from '../../storage/dummy/doc'; import { DummyDocStorage } from '../../storage/dummy/doc';
import { DummyDocSyncStorage } from '../../storage/dummy/doc-sync'; import { DummyDocSyncStorage } from '../../storage/dummy/doc-sync';
import { takeUntilAbort } from '../../utils/take-until-abort';
import { MANUALLY_STOP } from '../../utils/throw-if-aborted'; import { MANUALLY_STOP } from '../../utils/throw-if-aborted';
import type { PeerStorageOptions } from '../types'; import type { PeerStorageOptions } from '../types';
import { DocSyncPeer } from './peer'; import { DocSyncPeer } from './peer';
@@ -26,6 +37,7 @@ export interface DocSyncDocState {
export interface DocSync { export interface DocSync {
readonly state$: Observable<DocSyncState>; readonly state$: Observable<DocSyncState>;
docState$(docId: string): Observable<DocSyncDocState>; docState$(docId: string): Observable<DocSyncDocState>;
waitForSynced(docId?: string, abort?: AbortSignal): Promise<void>;
addPriority(id: string, priority: number): () => void; addPriority(id: string, priority: number): () => void;
resetSync(): Promise<void>; resetSync(): Promise<void>;
} }
@@ -39,7 +51,9 @@ export class DocSyncImpl implements DocSync {
); );
private abort: AbortController | null = null; private abort: AbortController | null = null;
state$ = combineLatest(this.peers.map(peer => peer.peerState$)).pipe( private readonly _state$ = combineLatest(
this.peers.map(peer => peer.peerState$)
).pipe(
map(allPeers => map(allPeers =>
allPeers.length === 0 allPeers.length === 0
? { ? {
@@ -66,6 +80,14 @@ export class DocSyncImpl implements DocSync {
}) })
) as Observable<DocSyncState>; ) as Observable<DocSyncState>;
state$ = this._state$.pipe(
// throttle the state to 1 second to avoid spamming the UI
throttleTime(1000, undefined, {
leading: true,
trailing: true,
})
);
constructor( constructor(
readonly storages: PeerStorageOptions<DocStorage>, readonly storages: PeerStorageOptions<DocStorage>,
readonly sync: DocSyncStorage readonly sync: DocSyncStorage
@@ -84,7 +106,7 @@ export class DocSyncImpl implements DocSync {
); );
} }
docState$(docId: string): Observable<DocSyncDocState> { private _docState$(docId: string): Observable<DocSyncDocState> {
if (this.peers.length === 0) { if (this.peers.length === 0) {
return of({ return of({
errorMessage: null, errorMessage: null,
@@ -106,6 +128,29 @@ export class DocSyncImpl implements DocSync {
); );
} }
docState$(docId: string): Observable<DocSyncDocState> {
return this._docState$(docId).pipe(
// throttle the state to 1 second to avoid spamming the UI
throttleTime(1000, undefined, {
leading: true,
trailing: true,
})
);
}
async waitForSynced(docId?: string, abort?: AbortSignal): Promise<void> {
const source$: Observable<DocSyncDocState | DocSyncState> = docId
? this._docState$(docId)
: this._state$;
await lastValueFrom(
source$.pipe(
filter(state => state.synced),
takeUntilAbort(abort),
first()
)
);
}
start() { start() {
if (this.abort) { if (this.abort) {
this.abort.abort(MANUALLY_STOP); this.abort.abort(MANUALLY_STOP);

View File

@@ -2,6 +2,7 @@ import { readAllDocsFromRootDoc } from '@affine/reader';
import { import {
filter, filter,
first, first,
lastValueFrom,
Observable, Observable,
ReplaySubject, ReplaySubject,
share, share,
@@ -71,60 +72,37 @@ export class IndexerSyncImpl implements IndexerSync {
state$ = this.status.state$.pipe( state$ = this.status.state$.pipe(
// throttle the state to 1 second to avoid spamming the UI // throttle the state to 1 second to avoid spamming the UI
throttleTime(1000) throttleTime(1000, undefined, {
leading: true,
trailing: true,
})
); );
docState$(docId: string) { docState$(docId: string) {
return this.status.docState$(docId).pipe( return this.status.docState$(docId).pipe(
// throttle the state to 1 second to avoid spamming the UI // throttle the state to 1 second to avoid spamming the UI
throttleTime(1000) throttleTime(1000, undefined, { leading: true, trailing: true })
); );
} }
waitForCompleted(signal?: AbortSignal) { async waitForCompleted(signal?: AbortSignal) {
return new Promise<void>((resolve, reject) => { await lastValueFrom(
this.status.state$ this.status.state$.pipe(
.pipe( filter(state => state.completed),
filter(state => state.completed), takeUntilAbort(signal),
takeUntilAbort(signal), first()
first() )
)
.subscribe({
next: () => {
resolve();
},
error: err => {
reject(err);
},
});
});
}
waitForDocCompleted(docId: string, signal?: AbortSignal) {
return new Promise<void>((resolve, reject) => {
this.status
.docState$(docId)
.pipe(
filter(state => state.completed),
takeUntilAbort(signal),
first()
)
.subscribe({
next: () => {
resolve();
},
error: err => {
reject(err);
},
});
});
}
readonly interval = () =>
new Promise<void>(resolve =>
requestIdleCallback(() => resolve(), {
timeout: 200,
})
); );
}
async waitForDocCompleted(docId: string, signal?: AbortSignal) {
await lastValueFrom(
this.status.docState$(docId).pipe(
filter(state => state.completed),
takeUntilAbort(signal),
first()
)
);
}
constructor( constructor(
readonly doc: DocStorage, readonly doc: DocStorage,

View File

@@ -242,6 +242,10 @@ class WorkerDocSync implements DocSync {
return this.client.ob$('docSync.docState', docId); return this.client.ob$('docSync.docState', docId);
} }
async waitForSynced(docId?: string, abort?: AbortSignal): Promise<void> {
await this.client.call('docSync.waitForSynced', docId ?? null, abort);
}
addPriority(docId: string, priority: number) { addPriority(docId: string, priority: number) {
const subscription = this.client const subscription = this.client
.ob$('docSync.addPriority', { docId, priority }) .ob$('docSync.addPriority', { docId, priority })

View File

@@ -212,6 +212,8 @@ class StoreConsumer {
const undo = this.docSync.addPriority(docId, priority); const undo = this.docSync.addPriority(docId, priority);
return () => undo(); return () => undo();
}), }),
'docSync.waitForSynced': (docId, ctx) =>
this.docSync.waitForSynced(docId ?? undefined, ctx.signal),
'docSync.resetSync': () => this.docSync.resetSync(), 'docSync.resetSync': () => this.docSync.resetSync(),
'blobSync.state': () => this.blobSync.state$, 'blobSync.state': () => this.blobSync.state$,
'blobSync.blobState': blobId => this.blobSync.blobState$(blobId), 'blobSync.blobState': blobId => this.blobSync.blobState$(blobId),

View File

@@ -102,6 +102,7 @@ interface GroupedWorkerOps {
docSync: { docSync: {
state: [void, DocSyncState]; state: [void, DocSyncState];
docState: [string, DocSyncDocState]; docState: [string, DocSyncDocState];
waitForSynced: [string | null, void];
addPriority: [{ docId: string; priority: number }, boolean]; addPriority: [{ docId: string; priority: number }, boolean];
resetSync: [void, void]; resetSync: [void, void];
}; };

View File

@@ -88,7 +88,7 @@ const useSyncEngineSyncProgress = (meta: WorkspaceMetadata) => {
const engineState = useLiveData( const engineState = useLiveData(
useMemo(() => { useMemo(() => {
return workspace return workspace
? LiveData.from(workspace.engine.doc.state$, null).throttleTime(500) ? LiveData.from(workspace.engine.doc.state$, null)
: null; : null;
}, [workspace]) }, [workspace])
); );

View File

@@ -8,15 +8,7 @@ import {
onStart, onStart,
} from '@toeverything/infra'; } from '@toeverything/infra';
import { fileTypeFromBuffer } from 'file-type'; import { fileTypeFromBuffer } from 'file-type';
import { import { switchMap, tap } from 'rxjs';
filter,
firstValueFrom,
fromEvent,
map,
switchMap,
takeUntil,
tap,
} from 'rxjs';
import type { DocsSearchService } from '../../docs-search'; import type { DocsSearchService } from '../../docs-search';
import type { WorkspaceService } from '../../workspace'; import type { WorkspaceService } from '../../workspace';
@@ -85,15 +77,7 @@ export class UnusedBlobs extends Entity {
async getUnusedBlobs(abortSignal?: AbortSignal) { async getUnusedBlobs(abortSignal?: AbortSignal) {
// Wait for both sync and indexing to complete // Wait for both sync and indexing to complete
const ready$ = this.workspaceService.workspace.engine.doc.state$ await this.workspaceService.workspace.engine.doc.waitForSynced();
.pipe(filter(state => state.syncing === 0 && !state.syncRetrying))
.pipe(map(() => true));
await firstValueFrom(
abortSignal
? ready$.pipe(takeUntil(fromEvent(abortSignal, 'abort')))
: ready$
);
await this.docsSearchService.indexer.waitForCompleted(abortSignal); await this.docsSearchService.indexer.waitForCompleted(abortSignal);

View File

@@ -44,8 +44,8 @@ export class ImportClipperService extends Service {
docsService.list.setPrimaryMode(docId, 'page'); docsService.list.setPrimaryMode(docId, 'page');
workspace.engine.doc.addPriority(workspace.id, 100); workspace.engine.doc.addPriority(workspace.id, 100);
workspace.engine.doc.addPriority(docId, 100); workspace.engine.doc.addPriority(docId, 100);
await workspace.engine.doc.waitForDocSynced(workspace.id); await workspace.engine.doc.waitForSynced(workspace.id);
await workspace.engine.doc.waitForDocSynced(docId); await workspace.engine.doc.waitForSynced(docId);
disposeWorkspace(); disposeWorkspace();
return docId; return docId;
} else { } else {

View File

@@ -33,7 +33,7 @@ test('should create a page with a local first avatar and remove it', async ({
.nth(0) .nth(0)
.getByTestId('workspace-avatar') .getByTestId('workspace-avatar')
.click(); .click();
await page.waitForTimeout(1000); await page.waitForTimeout(2000);
await page.getByTestId('workspace-name').click(); await page.getByTestId('workspace-name').click();
await page await page
.getByTestId('workspace-card') .getByTestId('workspace-card')