From bd967ce5e49a06d433980dc1b64aebf215fdf8e8 Mon Sep 17 00:00:00 2001 From: "Claude (thor)" Date: Mon, 10 Aug 2026 18:31:09 -0400 Subject: [PATCH] refactor(server): collapse duplicated sync stream handlers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The sync service hand-rolled the same delete/upsert/backfill/create stream loops across ~27 per-entity handlers. Extract four typed helpers on SyncService — streamDeletes, streamUpserts (overloaded, optional map), streamBackfill, and streamCreates — and rewrite each handler to call them. Behavior is unchanged: the per-parent backfill bookkeeping (isEntityBackfillComplete / getStartId / completion acks / checkpoint seeding) and the create-stream's one-time ack are preserved exactly, and the OCR delete keeps its whole-row payload. Verified by the sync medium test suite (26 files, 145 tests) plus tsc and eslint. Co-Authored-By: Claude Opus 4.8 --- server/src/services/sync.service.ts | 924 +++++++++++++--------------- 1 file changed, 413 insertions(+), 511 deletions(-) diff --git a/server/src/services/sync.service.ts b/server/src/services/sync.service.ts index e3842b1503..de6bd427e4 100644 --- a/server/src/services/sync.service.ts +++ b/server/src/services/sync.service.ts @@ -21,6 +21,8 @@ import { hexOrBufferToBase64 } from 'src/utils/bytes'; import { fromAck, serialize, SerializeOptions, toAck } from 'src/utils/sync'; type CheckpointMap = Partial>; +type SyncPartnerParent = { sharedById: string; createId: string }; +type SyncAlbumParent = { id: string; createId: string }; type AssetLike = Omit & { checksum: Buffer; thumbhash: Buffer | null; @@ -238,40 +240,145 @@ export class SyncService extends BaseService { return DateTime.fromMillis(milliseconds) < DateTime.now().minus(MAX_DURATION); } - private async syncAuthUsersV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const upsertType = SyncEntityType.AuthUserV1; - const upserts = this.syncRepository.authUser.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, profileImagePath, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data: { ...data, hasProfileImage: !!profileImagePath } }); + private async streamDeletes( + response: Writable, + type: T, + deletes: AsyncIterableIterator<{ id: string } & SyncItem[T]>, + ) { + for await (const { id, ...data } of deletes) { + send(response, { type, ids: [id], data: data as unknown as SyncItem[T] }); } } + private async streamUpserts( + response: Writable, + type: T, + upserts: AsyncIterableIterator<{ updateId: string } & SyncItem[T]>, + ): Promise; + private async streamUpserts( + response: Writable, + type: T, + upserts: AsyncIterableIterator, + map: (row: Omit) => SyncItem[T] | Promise, + ): Promise; + private async streamUpserts( + response: Writable, + type: T, + upserts: AsyncIterableIterator, + map?: (row: Omit) => SyncItem[T] | Promise, + ): Promise { + for await (const { updateId, ...row } of upserts) { + const data = map ? await map(row) : (row as unknown as SyncItem[T]); + send(response, { type, ids: [updateId], data }); + } + } + + private async streamBackfill< + T extends keyof SyncItem, + Parent extends { createId: string }, + Row extends { updateId: string }, + >( + response: Writable, + args: { + sessionId: string; + backfillType: T; + backfillCheckpoint: SyncAck | undefined; + endCheckpoint: SyncAck | undefined; + getParents: (afterCreateId: string | undefined) => Promise; + getBackfill: ( + parent: Parent, + range: { afterUpdateId: string | undefined; beforeUpdateId: string }, + ) => AsyncIterableIterator; + map?: (row: Omit, 'updateId'>) => SyncItem[T]; + }, + ) { + const { sessionId, backfillType, backfillCheckpoint, endCheckpoint, getParents, getBackfill, map } = args; + const parents = await getParents(backfillCheckpoint?.updateId); + + if (endCheckpoint) { + const beforeUpdateId = endCheckpoint.updateId; + + for (const parent of parents) { + const createId = parent.createId; + if (isEntityBackfillComplete(createId, backfillCheckpoint)) { + continue; + } + + const afterUpdateId = getStartId(createId, backfillCheckpoint); + const backfill = getBackfill(parent, { afterUpdateId, beforeUpdateId }); + + for await (const { updateId, ...row } of backfill) { + const data = map ? map(row) : (row as unknown as SyncItem[T]); + send(response, { type: backfillType, ids: [createId, updateId], data }); + } + + sendEntityBackfillCompleteAck(response, backfillType, createId); + } + } else if (parents.length > 0) { + await this.upsertBackfillCheckpoint({ + type: backfillType, + sessionId, + createId: parents.at(-1)!.createId, + }); + } + } + + private async streamCreates( + response: Writable, + args: { + createType: T; + ackType: SyncEntityType; + nowId: string; + creates: AsyncIterableIterator; + map?: (row: Omit, 'updateId'>) => SyncItem[T]; + }, + ) { + const { createType, ackType, nowId, creates, map } = args; + let isFirst = true; + for await (const { updateId, ...row } of creates) { + if (isFirst) { + send(response, { type: SyncEntityType.SyncAckV1, data: {}, ackType, ids: [nowId] }); + isFirst = false; + } + const data = map ? map(row) : (row as unknown as SyncItem[T]); + send(response, { type: createType, ids: [updateId], data }); + } + } + + private async syncAuthUsersV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { + await this.streamUpserts( + response, + SyncEntityType.AuthUserV1, + this.syncRepository.authUser.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AuthUserV1] }), + ({ profileImagePath, ...data }) => ({ ...data, hasProfileImage: !!profileImagePath }), + ); + } + private async syncUsersV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.UserDeleteV1; - const deletes = this.syncRepository.user.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.UserV1; - const upserts = this.syncRepository.user.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, profileImagePath, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data: { ...data, hasProfileImage: !!profileImagePath } }); - } + await this.streamDeletes( + response, + SyncEntityType.UserDeleteV1, + this.syncRepository.user.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.UserDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.UserV1, + this.syncRepository.user.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.UserV1] }), + ({ profileImagePath, ...data }) => ({ ...data, hasProfileImage: !!profileImagePath }), + ); } private async syncPartnersV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.PartnerDeleteV1; - const deletes = this.syncRepository.partner.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.PartnerV1; - const upserts = this.syncRepository.partner.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.PartnerDeleteV1, + this.syncRepository.partner.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.PartnerDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.PartnerV1, + this.syncRepository.partner.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.PartnerV1] }), + ); } private syncAssetsV1(): Promise { @@ -279,17 +386,17 @@ export class SyncService extends BaseService { } private async syncAssetsV2(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.AssetDeleteV1; - const deletes = this.syncRepository.asset.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.AssetV2; - const upserts = this.syncRepository.asset.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data: mapSyncAssetV2(data) }); - } + await this.streamDeletes( + response, + SyncEntityType.AssetDeleteV1, + this.syncRepository.asset.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AssetDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.AssetV2, + this.syncRepository.asset.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AssetV2] }), + mapSyncAssetV2, + ); } private syncPartnerAssetsV1(): Promise { @@ -304,80 +411,56 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const deleteType = SyncEntityType.PartnerAssetDeleteV1; - const deletes = this.syncRepository.partnerAsset.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } + await this.streamDeletes( + response, + SyncEntityType.PartnerAssetDeleteV1, + this.syncRepository.partnerAsset.getDeletes({ + ...options, + ack: checkpointMap[SyncEntityType.PartnerAssetDeleteV1], + }), + ); - const backfillType = SyncEntityType.PartnerAssetBackfillV2; - const backfillCheckpoint = checkpointMap[backfillType]; - const partners = await this.syncRepository.partner.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, - }); - const upsertType = SyncEntityType.PartnerAssetV2; - const upsertCheckpoint = checkpointMap[upsertType]; - if (upsertCheckpoint) { - const endId = upsertCheckpoint.updateId; - - for (const partner of partners) { - const createId = partner.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.partnerAsset.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - partner.sharedById, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { - type: backfillType, - ids: [createId, updateId], - data: mapSyncAssetV2(data), - }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (partners.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, + await this.streamBackfill( + response, + { sessionId, - createId: partners.at(-1)!.createId, - }); - } + backfillType: SyncEntityType.PartnerAssetBackfillV2, + backfillCheckpoint: checkpointMap[SyncEntityType.PartnerAssetBackfillV2], + endCheckpoint: checkpointMap[SyncEntityType.PartnerAssetV2], + getParents: (afterCreateId) => this.syncRepository.partner.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (partner, range) => + this.syncRepository.partnerAsset.getBackfill({ ...options, ...range }, partner.sharedById), + map: mapSyncAssetV2, + }, + ); - const upserts = this.syncRepository.partnerAsset.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data: mapSyncAssetV2(data) }); - } + await this.streamUpserts( + response, + SyncEntityType.PartnerAssetV2, + this.syncRepository.partnerAsset.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.PartnerAssetV2] }), + mapSyncAssetV2, + ); } private async syncAssetExifsV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const upsertType = SyncEntityType.AssetExifV1; - const upserts = this.syncRepository.assetExif.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamUpserts( + response, + SyncEntityType.AssetExifV1, + this.syncRepository.assetExif.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AssetExifV1] }), + ); } private async syncAssetEditsV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.AssetEditDeleteV1; - const deletes = this.syncRepository.assetEdit.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - const upsertType = SyncEntityType.AssetEditV1; - const upserts = this.syncRepository.assetEdit.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.AssetEditDeleteV1, + this.syncRepository.assetEdit.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AssetEditDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.AssetEditV1, + this.syncRepository.assetEdit.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AssetEditV1] }), + ); } private async syncPartnerAssetExifsV1( @@ -386,83 +469,57 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const backfillType = SyncEntityType.PartnerAssetExifBackfillV1; - const backfillCheckpoint = checkpointMap[backfillType]; - const partners = await this.syncRepository.partner.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, + await this.streamBackfill(response, { + sessionId, + backfillType: SyncEntityType.PartnerAssetExifBackfillV1, + backfillCheckpoint: checkpointMap[SyncEntityType.PartnerAssetExifBackfillV1], + endCheckpoint: checkpointMap[SyncEntityType.PartnerAssetExifV1], + getParents: (afterCreateId) => this.syncRepository.partner.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (partner: SyncPartnerParent, range) => + this.syncRepository.partnerAssetExif.getBackfill({ ...options, ...range }, partner.sharedById), }); - const upsertType = SyncEntityType.PartnerAssetExifV1; - const upsertCheckpoint = checkpointMap[upsertType]; - if (upsertCheckpoint) { - const endId = upsertCheckpoint.updateId; - - for (const partner of partners) { - const createId = partner.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.partnerAssetExif.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - partner.sharedById, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { type: backfillType, ids: [partner.createId, updateId], data }); - } - - sendEntityBackfillCompleteAck(response, backfillType, partner.createId); - } - } else if (partners.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, - sessionId, - createId: partners.at(-1)!.createId, - }); - } - - const upserts = this.syncRepository.partnerAssetExif.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamUpserts( + response, + SyncEntityType.PartnerAssetExifV1, + this.syncRepository.partnerAssetExif.getUpserts({ + ...options, + ack: checkpointMap[SyncEntityType.PartnerAssetExifV1], + }), + ); } private async syncAlbumsV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.AlbumDeleteV1; - const deletes = this.syncRepository.album.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.AlbumV1; - const upserts = this.syncRepository.album.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - const albumUsers = await this.syncRepository.album.getAlbumUsers(data.id); - send(response, { - type: upsertType, - ids: [updateId], + await this.streamDeletes( + response, + SyncEntityType.AlbumDeleteV1, + this.syncRepository.album.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AlbumDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.AlbumV1, + this.syncRepository.album.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AlbumV1] }), + async (data) => { + const albumUsers = await this.syncRepository.album.getAlbumUsers(data.id); // TODO: return null instead of '' in v4 - data: syncAlbumV2ToV1({ ...data, description: data.description ?? '' }, albumUsers), - }); - } + return syncAlbumV2ToV1({ ...data, description: data.description ?? '' }, albumUsers); + }, + ); } private async syncAlbumsV2(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.AlbumDeleteV1; - const deletes = this.syncRepository.album.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.AlbumV2; - const upserts = this.syncRepository.album.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { + await this.streamDeletes( + response, + SyncEntityType.AlbumDeleteV1, + this.syncRepository.album.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AlbumDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.AlbumV2, + this.syncRepository.album.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AlbumV2] }), // TODO: return null instead of '' in v4 - send(response, { type: upsertType, ids: [updateId], data: { ...data, description: data.description ?? '' } }); - } + (data) => ({ ...data, description: data.description ?? '' }), + ); } private async syncAlbumUsersV1( @@ -471,53 +528,26 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const deleteType = SyncEntityType.AlbumUserDeleteV1; - const deletes = this.syncRepository.albumUser.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } + await this.streamDeletes( + response, + SyncEntityType.AlbumUserDeleteV1, + this.syncRepository.albumUser.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AlbumUserDeleteV1] }), + ); - const backfillType = SyncEntityType.AlbumUserBackfillV1; - const backfillCheckpoint = checkpointMap[backfillType]; - const albums = await this.syncRepository.album.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, + await this.streamBackfill(response, { + sessionId, + backfillType: SyncEntityType.AlbumUserBackfillV1, + backfillCheckpoint: checkpointMap[SyncEntityType.AlbumUserBackfillV1], + endCheckpoint: checkpointMap[SyncEntityType.AlbumUserV1], + getParents: (afterCreateId) => this.syncRepository.album.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (album: SyncAlbumParent, range) => this.syncRepository.albumUser.getBackfill({ ...options, ...range }, album.id), }); - const upsertType = SyncEntityType.AlbumUserV1; - const upsertCheckpoint = checkpointMap[upsertType]; - if (upsertCheckpoint) { - const endId = upsertCheckpoint.updateId; - for (const album of albums) { - const createId = album.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.albumUser.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - album.id, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { type: backfillType, ids: [createId, updateId], data }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (albums.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, - sessionId, - createId: albums.at(-1)!.createId, - }); - } - - const upserts = this.syncRepository.albumUser.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamUpserts( + response, + SyncEntityType.AlbumUserV1, + this.syncRepository.albumUser.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AlbumUserV1] }), + ); } private syncAlbumAssetsV1(): Promise { @@ -532,70 +562,39 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const backfillType = SyncEntityType.AlbumAssetBackfillV2; - const backfillCheckpoint = checkpointMap[backfillType]; - const albums = await this.syncRepository.album.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, - }); const updateType = SyncEntityType.AlbumAssetUpdateV2; - const createType = SyncEntityType.AlbumAssetCreateV2; - const updateCheckpoint = checkpointMap[updateType]; - const createCheckpoint = checkpointMap[createType]; - if (createCheckpoint) { - const endId = createCheckpoint.updateId; + const createCheckpoint = checkpointMap[SyncEntityType.AlbumAssetCreateV2]; - for (const album of albums) { - const createId = album.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.albumAsset.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - album.id, - options.userId, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { type: backfillType, ids: [createId, updateId], data: mapSyncAssetV2(data) }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (albums.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, + await this.streamBackfill( + response, + { sessionId, - createId: albums.at(-1)!.createId, - }); - } + backfillType: SyncEntityType.AlbumAssetBackfillV2, + backfillCheckpoint: checkpointMap[SyncEntityType.AlbumAssetBackfillV2], + endCheckpoint: createCheckpoint, + getParents: (afterCreateId) => this.syncRepository.album.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (album, range) => + this.syncRepository.albumAsset.getBackfill({ ...options, ...range }, album.id, options.userId), + map: mapSyncAssetV2, + }, + ); if (createCheckpoint) { - const updates = this.syncRepository.albumAsset.getUpdates( - { ...options, ack: updateCheckpoint }, - createCheckpoint, + await this.streamUpserts( + response, + updateType, + this.syncRepository.albumAsset.getUpdates({ ...options, ack: checkpointMap[updateType] }, createCheckpoint), + mapSyncAssetV2, ); - for await (const { updateId, ...data } of updates) { - send(response, { type: updateType, ids: [updateId], data: mapSyncAssetV2(data) }); - } } - const creates = this.syncRepository.albumAsset.getCreates({ ...options, ack: createCheckpoint }); - let isFirst = true; - for await (const { updateId, ...data } of creates) { - if (isFirst) { - send(response, { - type: SyncEntityType.SyncAckV1, - data: {}, - ackType: SyncEntityType.AlbumAssetUpdateV2, - ids: [options.nowId], - }); - isFirst = false; - } - send(response, { type: createType, ids: [updateId], data: mapSyncAssetV2(data) }); - } + await this.streamCreates(response, { + createType: SyncEntityType.AlbumAssetCreateV2, + ackType: updateType, + nowId: options.nowId, + creates: this.syncRepository.albumAsset.getCreates({ ...options, ack: createCheckpoint }), + map: mapSyncAssetV2, + }); } private async syncAlbumAssetExifsV1( @@ -604,69 +603,32 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const backfillType = SyncEntityType.AlbumAssetExifBackfillV1; - const backfillCheckpoint = checkpointMap[backfillType]; - const albums = await this.syncRepository.album.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, - }); const updateType = SyncEntityType.AlbumAssetExifUpdateV1; - const createType = SyncEntityType.AlbumAssetExifCreateV1; - const upsertCheckpoint = checkpointMap[updateType]; - const createCheckpoint = checkpointMap[createType]; - if (createCheckpoint) { - const endId = createCheckpoint.updateId; + const createCheckpoint = checkpointMap[SyncEntityType.AlbumAssetExifCreateV1]; - for (const album of albums) { - const createId = album.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.albumAssetExif.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - album.id, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { type: backfillType, ids: [createId, updateId], data }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (albums.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, - sessionId, - createId: albums.at(-1)!.createId, - }); - } + await this.streamBackfill(response, { + sessionId, + backfillType: SyncEntityType.AlbumAssetExifBackfillV1, + backfillCheckpoint: checkpointMap[SyncEntityType.AlbumAssetExifBackfillV1], + endCheckpoint: createCheckpoint, + getParents: (afterCreateId) => this.syncRepository.album.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (album: SyncAlbumParent, range) => this.syncRepository.albumAssetExif.getBackfill({ ...options, ...range }, album.id), + }); if (createCheckpoint) { - const updates = this.syncRepository.albumAssetExif.getUpdates( - { ...options, ack: upsertCheckpoint }, - createCheckpoint, + await this.streamUpserts( + response, + updateType, + this.syncRepository.albumAssetExif.getUpdates({ ...options, ack: checkpointMap[updateType] }, createCheckpoint), ); - for await (const { updateId, ...data } of updates) { - send(response, { type: updateType, ids: [updateId], data }); - } } - const creates = this.syncRepository.albumAssetExif.getCreates({ ...options, ack: createCheckpoint }); - let isFirst = true; - for await (const { updateId, ...data } of creates) { - if (isFirst) { - send(response, { - type: SyncEntityType.SyncAckV1, - data: {}, - ackType: SyncEntityType.AlbumAssetExifUpdateV1, - ids: [options.nowId], - }); - isFirst = false; - } - send(response, { type: createType, ids: [updateId], data }); - } + await this.streamCreates(response, { + createType: SyncEntityType.AlbumAssetExifCreateV1, + ackType: updateType, + nowId: options.nowId, + creates: this.syncRepository.albumAssetExif.getCreates({ ...options, ack: createCheckpoint }), + }); } private async syncAlbumToAssetsV1( @@ -675,95 +637,71 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const deleteType = SyncEntityType.AlbumToAssetDeleteV1; - const deletes = this.syncRepository.albumToAsset.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } + await this.streamDeletes( + response, + SyncEntityType.AlbumToAssetDeleteV1, + this.syncRepository.albumToAsset.getDeletes({ + ...options, + ack: checkpointMap[SyncEntityType.AlbumToAssetDeleteV1], + }), + ); - const backfillType = SyncEntityType.AlbumToAssetBackfillV1; - const backfillCheckpoint = checkpointMap[backfillType]; - const albums = await this.syncRepository.album.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, + await this.streamBackfill(response, { + sessionId, + backfillType: SyncEntityType.AlbumToAssetBackfillV1, + backfillCheckpoint: checkpointMap[SyncEntityType.AlbumToAssetBackfillV1], + endCheckpoint: checkpointMap[SyncEntityType.AlbumToAssetV1], + getParents: (afterCreateId) => this.syncRepository.album.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (album: SyncAlbumParent, range) => this.syncRepository.albumToAsset.getBackfill({ ...options, ...range }, album.id), }); - const upsertType = SyncEntityType.AlbumToAssetV1; - const upsertCheckpoint = checkpointMap[upsertType]; - if (upsertCheckpoint) { - const endId = upsertCheckpoint.updateId; - for (const album of albums) { - const createId = album.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.albumToAsset.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - album.id, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { type: backfillType, ids: [createId, updateId], data }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (albums.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, - sessionId, - createId: albums.at(-1)!.createId, - }); - } - - const upserts = this.syncRepository.albumToAsset.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamUpserts( + response, + SyncEntityType.AlbumToAssetV1, + this.syncRepository.albumToAsset.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AlbumToAssetV1] }), + ); } private async syncMemoriesV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.MemoryDeleteV1; - const deletes = this.syncRepository.memory.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.MemoryV1; - const upserts = this.syncRepository.memory.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.MemoryDeleteV1, + this.syncRepository.memory.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.MemoryDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.MemoryV1, + this.syncRepository.memory.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.MemoryV1] }), + ); } private async syncMemoryAssetsV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.MemoryToAssetDeleteV1; - const deletes = this.syncRepository.memoryToAsset.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.MemoryToAssetV1; - const upserts = this.syncRepository.memoryToAsset.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.MemoryToAssetDeleteV1, + this.syncRepository.memoryToAsset.getDeletes({ + ...options, + ack: checkpointMap[SyncEntityType.MemoryToAssetDeleteV1], + }), + ); + await this.streamUpserts( + response, + SyncEntityType.MemoryToAssetV1, + this.syncRepository.memoryToAsset.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.MemoryToAssetV1] }), + ); } private async syncStackV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.StackDeleteV1; - const deletes = this.syncRepository.stack.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.StackV1; - const upserts = this.syncRepository.stack.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.StackDeleteV1, + this.syncRepository.stack.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.StackDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.StackV1, + this.syncRepository.stack.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.StackV1] }), + ); } private async syncPartnerStackV1( @@ -772,71 +710,43 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, sessionId: string, ) { - const deleteType = SyncEntityType.PartnerStackDeleteV1; - const deletes = this.syncRepository.partnerStack.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } + await this.streamDeletes( + response, + SyncEntityType.PartnerStackDeleteV1, + this.syncRepository.partnerStack.getDeletes({ + ...options, + ack: checkpointMap[SyncEntityType.PartnerStackDeleteV1], + }), + ); - const backfillType = SyncEntityType.PartnerStackBackfillV1; - const backfillCheckpoint = checkpointMap[backfillType]; - const partners = await this.syncRepository.partner.getCreatedAfter({ - ...options, - afterCreateId: backfillCheckpoint?.updateId, + await this.streamBackfill(response, { + sessionId, + backfillType: SyncEntityType.PartnerStackBackfillV1, + backfillCheckpoint: checkpointMap[SyncEntityType.PartnerStackBackfillV1], + endCheckpoint: checkpointMap[SyncEntityType.PartnerStackV1], + getParents: (afterCreateId) => this.syncRepository.partner.getCreatedAfter({ ...options, afterCreateId }), + getBackfill: (partner: SyncPartnerParent, range) => + this.syncRepository.partnerStack.getBackfill({ ...options, ...range }, partner.sharedById), }); - const upsertType = SyncEntityType.PartnerStackV1; - const upsertCheckpoint = checkpointMap[upsertType]; - if (upsertCheckpoint) { - const endId = upsertCheckpoint.updateId; - for (const partner of partners) { - const createId = partner.createId; - if (isEntityBackfillComplete(createId, backfillCheckpoint)) { - continue; - } - - const startId = getStartId(createId, backfillCheckpoint); - const backfill = this.syncRepository.partnerStack.getBackfill( - { ...options, afterUpdateId: startId, beforeUpdateId: endId }, - partner.sharedById, - ); - - for await (const { updateId, ...data } of backfill) { - send(response, { - type: backfillType, - ids: [createId, updateId], - data, - }); - } - - sendEntityBackfillCompleteAck(response, backfillType, createId); - } - } else if (partners.length > 0) { - await this.upsertBackfillCheckpoint({ - type: backfillType, - sessionId, - createId: partners.at(-1)!.createId, - }); - } - - const upserts = this.syncRepository.partnerStack.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamUpserts( + response, + SyncEntityType.PartnerStackV1, + this.syncRepository.partnerStack.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.PartnerStackV1] }), + ); } private async syncPeopleV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.PersonDeleteV1; - const deletes = this.syncRepository.person.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.PersonV1; - const upserts = this.syncRepository.person.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.PersonDeleteV1, + this.syncRepository.person.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.PersonDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.PersonV1, + this.syncRepository.person.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.PersonV1] }), + ); } private syncAssetFacesV1(): Promise { @@ -846,33 +756,32 @@ export class SyncService extends BaseService { } private async syncAssetFacesV2(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.AssetFaceDeleteV1; - const deletes = this.syncRepository.assetFace.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.AssetFaceV2; - const upserts = this.syncRepository.assetFace.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.AssetFaceDeleteV1, + this.syncRepository.assetFace.getDeletes({ ...options, ack: checkpointMap[SyncEntityType.AssetFaceDeleteV1] }), + ); + await this.streamUpserts( + response, + SyncEntityType.AssetFaceV2, + this.syncRepository.assetFace.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AssetFaceV2] }), + ); } private async syncUserMetadataV1(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) { - const deleteType = SyncEntityType.UserMetadataDeleteV1; - const deletes = this.syncRepository.userMetadata.getDeletes({ ...options, ack: checkpointMap[deleteType] }); - - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.UserMetadataV1; - const upserts = this.syncRepository.userMetadata.getUpserts({ ...options, ack: checkpointMap[upsertType] }); - - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } + await this.streamDeletes( + response, + SyncEntityType.UserMetadataDeleteV1, + this.syncRepository.userMetadata.getDeletes({ + ...options, + ack: checkpointMap[SyncEntityType.UserMetadataDeleteV1], + }), + ); + await this.streamUpserts( + response, + SyncEntityType.UserMetadataV1, + this.syncRepository.userMetadata.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.UserMetadataV1] }), + ); } private async syncAssetMetadataV1( @@ -881,25 +790,22 @@ export class SyncService extends BaseService { checkpointMap: CheckpointMap, auth: AuthDto, ) { - const deleteType = SyncEntityType.AssetMetadataDeleteV1; - const deletes = this.syncRepository.assetMetadata.getDeletes( - { ...options, ack: checkpointMap[deleteType] }, - auth.user.id, + await this.streamDeletes( + response, + SyncEntityType.AssetMetadataDeleteV1, + this.syncRepository.assetMetadata.getDeletes( + { ...options, ack: checkpointMap[SyncEntityType.AssetMetadataDeleteV1] }, + auth.user.id, + ), ); - - for await (const { id, ...data } of deletes) { - send(response, { type: deleteType, ids: [id], data }); - } - - const upsertType = SyncEntityType.AssetMetadataV1; - const upserts = this.syncRepository.assetMetadata.getUpserts( - { ...options, ack: checkpointMap[upsertType] }, - auth.user.id, + await this.streamUpserts( + response, + SyncEntityType.AssetMetadataV1, + this.syncRepository.assetMetadata.getUpserts( + { ...options, ack: checkpointMap[SyncEntityType.AssetMetadataV1] }, + auth.user.id, + ), ); - - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } } private async syncAssetOcrV1( @@ -918,15 +824,11 @@ export class SyncService extends BaseService { send(response, { type: deleteType, ids: [row.id], data: row }); } - const upsertType = SyncEntityType.AssetOcrV1; - const upserts = this.syncRepository.assetOcr.getUpserts( - { ...options, ack: checkpointMap[upsertType] }, - auth.user.id, + await this.streamUpserts( + response, + SyncEntityType.AssetOcrV1, + this.syncRepository.assetOcr.getUpserts({ ...options, ack: checkpointMap[SyncEntityType.AssetOcrV1] }, auth.user.id), ); - - for await (const { updateId, ...data } of upserts) { - send(response, { type: upsertType, ids: [updateId], data }); - } } private async upsertBackfillCheckpoint(item: { type: SyncEntityType; sessionId: string; createId: string }) {