feat: sync to keck
This commit is contained in:
@@ -41,7 +41,7 @@ async function _getCurrentToken() {
|
|||||||
|
|
||||||
const _enabled = {
|
const _enabled = {
|
||||||
demo: [],
|
demo: [],
|
||||||
AFFiNE: process.env['NX_KECK'] ? ['idb', 'ws'] : ['idb'],
|
AFFiNE: process.env['NX_KECK'] ? ['idb', 'keck'] : ['idb'],
|
||||||
} as any;
|
} as any;
|
||||||
|
|
||||||
async function _getBlockDatabase(
|
async function _getBlockDatabase(
|
||||||
|
|||||||
@@ -4,7 +4,8 @@ import * as authProtocol from 'y-protocols/auth';
|
|||||||
import * as awarenessProtocol from 'y-protocols/awareness';
|
import * as awarenessProtocol from 'y-protocols/awareness';
|
||||||
import * as syncProtocol from 'y-protocols/sync';
|
import * as syncProtocol from 'y-protocols/sync';
|
||||||
|
|
||||||
import { WebsocketProvider } from './provider';
|
import { KeckProvider } from './keckprovider';
|
||||||
|
import { WebsocketProvider } from './wsprovider';
|
||||||
|
|
||||||
const permissionDeniedHandler = (provider: WebsocketProvider, reason: string) =>
|
const permissionDeniedHandler = (provider: WebsocketProvider, reason: string) =>
|
||||||
console.warn(`Permission denied to access ${provider.url}.\n${reason}`);
|
console.warn(`Permission denied to access ${provider.url}.\n${reason}`);
|
||||||
@@ -19,7 +20,7 @@ export enum Message {
|
|||||||
export type MessageCallback = (
|
export type MessageCallback = (
|
||||||
encoder: encoding.Encoder,
|
encoder: encoding.Encoder,
|
||||||
decoder: decoding.Decoder,
|
decoder: decoding.Decoder,
|
||||||
provider: WebsocketProvider,
|
provider: WebsocketProvider | KeckProvider,
|
||||||
emitSynced: boolean,
|
emitSynced: boolean,
|
||||||
messageType: number
|
messageType: number
|
||||||
) => void;
|
) => void;
|
||||||
@@ -48,14 +49,16 @@ export const handler: Record<Message, MessageCallback> = {
|
|||||||
emitSynced,
|
emitSynced,
|
||||||
messageType
|
messageType
|
||||||
) => {
|
) => {
|
||||||
encoding.writeVarUint(encoder, Message.queryAwareness);
|
if (provider instanceof WebsocketProvider) {
|
||||||
encoding.writeVarUint8Array(
|
encoding.writeVarUint(encoder, Message.queryAwareness);
|
||||||
encoder,
|
encoding.writeVarUint8Array(
|
||||||
awarenessProtocol.encodeAwarenessUpdate(
|
encoder,
|
||||||
provider.awareness,
|
awarenessProtocol.encodeAwarenessUpdate(
|
||||||
Array.from(provider.awareness.getStates().keys())
|
provider.awareness,
|
||||||
)
|
Array.from(provider.awareness.getStates().keys())
|
||||||
);
|
)
|
||||||
|
);
|
||||||
|
}
|
||||||
},
|
},
|
||||||
|
|
||||||
[Message.awareness]: (
|
[Message.awareness]: (
|
||||||
@@ -65,11 +68,13 @@ export const handler: Record<Message, MessageCallback> = {
|
|||||||
emitSynced,
|
emitSynced,
|
||||||
messageType
|
messageType
|
||||||
) => {
|
) => {
|
||||||
awarenessProtocol.applyAwarenessUpdate(
|
if (provider instanceof WebsocketProvider) {
|
||||||
provider.awareness,
|
awarenessProtocol.applyAwarenessUpdate(
|
||||||
decoding.readVarUint8Array(decoder),
|
provider.awareness,
|
||||||
provider
|
decoding.readVarUint8Array(decoder),
|
||||||
);
|
provider
|
||||||
|
);
|
||||||
|
}
|
||||||
},
|
},
|
||||||
|
|
||||||
[Message.auth]: (encoder, decoder, provider, emitSynced, messageType) => {
|
[Message.auth]: (encoder, decoder, provider, emitSynced, messageType) => {
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
export { IndexedDBProvider } from './indexeddb';
|
export { IndexedDBProvider } from './indexeddb';
|
||||||
export { WebsocketProvider } from './provider';
|
export { KeckProvider } from './keckprovider';
|
||||||
export { SQLiteProvider } from './sqlite';
|
export { SQLiteProvider } from './sqlite';
|
||||||
|
export { WebsocketProvider } from './wsprovider';
|
||||||
|
|||||||
@@ -120,13 +120,13 @@ async function _initYjsDatabase(
|
|||||||
[name]: p,
|
[name]: p,
|
||||||
};
|
};
|
||||||
}),
|
}),
|
||||||
// p({
|
p({
|
||||||
// awareness,
|
awareness,
|
||||||
// doc: binaries,
|
doc: binaries,
|
||||||
// token,
|
token,
|
||||||
// workspace: `${workspace}_binaries`,
|
workspace: `${workspace}_binaries`,
|
||||||
// emitState,
|
emitState,
|
||||||
// }).then(p => ({ [`${name}_binaries`]: p })),
|
}).then(p => ({ [`${name}_binaries`]: p })),
|
||||||
])
|
])
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ import { Doc } from 'yjs';
|
|||||||
|
|
||||||
import {
|
import {
|
||||||
IndexedDBProvider,
|
IndexedDBProvider,
|
||||||
|
KeckProvider,
|
||||||
SQLiteProvider,
|
SQLiteProvider,
|
||||||
WebsocketProvider,
|
WebsocketProvider,
|
||||||
} from '@toeverything/datasource/jwt-rpc';
|
} from '@toeverything/datasource/jwt-rpc';
|
||||||
@@ -22,7 +23,7 @@ export type YjsProvider = (
|
|||||||
instances: YjsDefaultInstances
|
instances: YjsDefaultInstances
|
||||||
) => Promise<unknown | undefined>;
|
) => Promise<unknown | undefined>;
|
||||||
|
|
||||||
type ProviderType = 'idb' | 'sqlite' | 'ws';
|
type ProviderType = 'idb' | 'sqlite' | 'ws' | 'keck';
|
||||||
|
|
||||||
export type YjsProviderOptions = {
|
export type YjsProviderOptions = {
|
||||||
enabled: ProviderType[];
|
enabled: ProviderType[];
|
||||||
@@ -100,9 +101,9 @@ export const getYjsProviders = (
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
keck: async (instances: YjsDefaultInstances) => {
|
keck: async (instances: YjsDefaultInstances) => {
|
||||||
if (options.enabled.includes('ws')) {
|
if (options.enabled.includes('keck')) {
|
||||||
if (instances.token) {
|
if (instances.token) {
|
||||||
const ws = new WebsocketProvider(
|
const ws = new KeckProvider(
|
||||||
instances.token,
|
instances.token,
|
||||||
options.backend,
|
options.backend,
|
||||||
instances.workspace,
|
instances.workspace,
|
||||||
|
|||||||
Reference in New Issue
Block a user