chore: improve event flow (#14266)

This commit is contained in:
DarkSky
2026-01-16 16:07:27 +08:00
committed by GitHub
parent d4581b839a
commit 924d58603f
43 changed files with 2306 additions and 567 deletions
+35 -1
View File
@@ -28,6 +28,12 @@ import type { AwarenessSync } from '../sync/awareness';
import type { BlobSync } from '../sync/blob';
import type { DocSync } from '../sync/doc';
import type { IndexerPreferOptions, IndexerSync } from '../sync/indexer';
import type {
TelemetryAck,
TelemetryContext,
TelemetryEvent,
TelemetryQueueState,
} from '../telemetry/types';
import type { StoreInitOptions, WorkerManagerOps, WorkerOps } from './ops';
export type { StoreInitOptions as WorkerInitOptions } from './ops';
@@ -41,7 +47,11 @@ export class StoreManagerClient {
}
>();
constructor(private readonly client: OpClient<WorkerManagerOps>) {}
constructor(private readonly client: OpClient<WorkerManagerOps>) {
this.telemetry = new TelemetryClient(this.client);
}
readonly telemetry: TelemetryClient;
open(key: string, options: StoreInitOptions) {
const { port1, port2 } = new MessageChannel();
@@ -104,6 +114,30 @@ export class StoreManagerClient {
}
}
class TelemetryClient {
constructor(private readonly client: OpClient<WorkerManagerOps>) {}
setContext(context: TelemetryContext): Promise<void> {
return this.client.call('telemetry.setContext', context);
}
track(event: TelemetryEvent): Promise<{ queued: boolean }> {
return this.client.call('telemetry.track', event);
}
pageview(event: TelemetryEvent): Promise<{ queued: boolean }> {
return this.client.call('telemetry.pageview', event);
}
flush(): Promise<TelemetryAck> {
return this.client.call('telemetry.flush');
}
getQueueState(): Promise<TelemetryQueueState> {
return this.client.call('telemetry.getQueueState');
}
}
export class StoreClient {
constructor(private readonly client: OpClient<WorkerOps>) {
this.docStorage = new WorkerDocStorage(this.client);
@@ -6,6 +6,7 @@ import { SpaceStorage } from '../storage';
import type { AwarenessRecord } from '../storage/awareness';
import { Sync } from '../sync';
import type { PeerStorageOptions } from '../sync/types';
import { TelemetryManager } from '../telemetry/manager';
import { MANUALLY_STOP } from '../utils/throw-if-aborted';
import type { StoreInitOptions, WorkerManagerOps, WorkerOps } from './ops';
@@ -338,6 +339,7 @@ export class StoreManagerConsumer {
string,
{ store: StoreConsumer; refCount: number }
>();
private readonly telemetry = new TelemetryManager();
constructor(
private readonly availableStorageImplementations: StorageConstructor[]
@@ -386,6 +388,11 @@ export class StoreManagerConsumer {
workerDisposer();
this.storeDisposers.delete(key);
},
'telemetry.setContext': context => this.telemetry.setContext(context),
'telemetry.track': event => this.telemetry.track(event),
'telemetry.pageview': event => this.telemetry.pageview(event),
'telemetry.flush': () => this.telemetry.flush(),
'telemetry.getQueueState': () => this.telemetry.getQueueState(),
});
}
}
+11
View File
@@ -16,6 +16,12 @@ import type { AwarenessRecord } from '../storage/awareness';
import type { BlobSyncBlobState, BlobSyncState } from '../sync/blob';
import type { DocSyncDocState, DocSyncState } from '../sync/doc';
import type { IndexerDocSyncState, IndexerSyncState } from '../sync/indexer';
import type {
TelemetryAck,
TelemetryContext,
TelemetryEvent,
TelemetryQueueState,
} from '../telemetry/types';
type StorageInitOptions = Values<{
[key in keyof AvailableStorageImplementations]: {
@@ -178,4 +184,9 @@ export type WorkerManagerOps = {
string,
];
close: [string, void];
'telemetry.setContext': [TelemetryContext, void];
'telemetry.track': [TelemetryEvent, { queued: boolean }];
'telemetry.pageview': [TelemetryEvent, { queued: boolean }];
'telemetry.flush': [void, TelemetryAck];
'telemetry.getQueueState': [void, TelemetryQueueState];
};