import { find, getSerialisableError, getSyncChannelName, ReplicationProtocol, tryCatchV2, } from "prostgles-types"; import type { AddSyncParams, BasicCallback, PubSubManager, SyncParams } from "./PubSubManager"; import { parseCondition } from "./PubSubManagerUtils"; /** * Returns a sync channel * A sync channel is unique per socket for each filter */ export async function addSync( this: PubSubManager, syncParams: AddSyncParams, ): Promise<{ channelName: string }> { const sid = this.dboBuilder.prostgles.authHandler.getSIDNoError({ socket: syncParams.socket }); const res = await tryCatchV2(async () => { const { socket, rawSelect, table_info, table_rules, filter, params, condition, initialData } = syncParams; const conditionParsed = parseCondition(condition); const { name: table_name } = table_info; const channelName = getSyncChannelName({ tableName: table_name, filter, select: rawSelect, }); this.upsertSocket(socket); const syncConfig = this.dboBuilder.prostgles.tableConfigurator?.getTableSyncConfig(table_name); if (!syncConfig) { throw `Sync not configured for table ${table_name}`; } const upsertSync = () => { /* Only a sync per socket per table/condition/select allowed */ const existing = find(this.syncs, { socket_id: socket.id, channel_name: channelName }); if (existing) { console.warn("addSync: Client tried to create a duplicate sync", existing.channel_name); return existing; } /* Server will: 1. Ask for last_synced emit(onSyncRequest) 2. Ask for data >= server_synced emit(onPullRequest) -> Upsert that data 2. Push data >= last_synced emit(data.data) Client will: 1. Send last_synced on(onSyncRequest) 2. Send data >= server_synced on(onPullRequest) 3. Send data on CRUD emit(data.data | data.deleted) 4. Upsert data.data | deleted on(data.data | data.deleted) */ const handlers = ReplicationProtocol.getServerHandlers(channelName, socket, { ClientSyncRequest: async (data) => { const err = await this.syncData(newSync, data, "client") .then(() => undefined) .catch((err) => { console.error("Error syncing data with client: ", err); return { success: false, err: getSerialisableError(err) } as const; }); return err ?? ({ success: true } as const); }, }); const newSync = { channel_name: channelName, table_name, filter, condition: conditionParsed, sid, table_rules, ...syncConfig, socket_id: socket.id, last_synced: 0, //initialData.isSynced ? Date.now() : 0, lr: undefined, // initialData.data.at(-1), table_info, is_syncing: false, wal: undefined, socket, params, handlers, }; const unsyncChn = channelName + "unsync"; socket.removeAllListeners(unsyncChn); socket.once(unsyncChn, (_data: any, cb: BasicCallback) => { void this._log({ type: "sync", command: "unsync", socketId: socket.id, tableName: table_name, condition, channelName, sid, connectedSocketIds: this.connectedSocketIds, duration: -1, syncParams: newSync, }); socket.removeAllListeners(channelName); socket.removeAllListeners(unsyncChn); this.syncs = this.syncs.filter((s) => { const isMatch = s.socket_id && s.socket_id === socket.id && s.channel_name === channelName; return !isMatch; }); cb(null, { res: "ok" }); }); return newSync; }; const newSync = upsertSync(); await this.addTrigger( { table_name, condition: conditionParsed, tracked_columns: undefined }, undefined, socket, ); this.syncs.push(newSync); return { channelName, newSync }; }); await this._log({ type: "sync", command: "addSync", tableName: syncParams.table_info.name, condition: syncParams.condition, socketId: syncParams.socket.id, connectedSocketIds: this.connectedSocketIds, duration: res.duration, error: res.error, sid, channelName: res.data?.channelName || "", syncParams: res.data?.newSync ?? ({} as SyncParams), }); if (res.hasError) throw res.error; return res.data; }