diff --git a/server/src/queries/integrity.repository.sql b/server/src/queries/integrity.repository.sql index d9940970d9..937d4454a2 100644 --- a/server/src/queries/integrity.repository.sql +++ b/server/src/queries/integrity.repository.sql @@ -130,10 +130,9 @@ from where "asset"."deletedAt" is null and "asset"."isExternal" = false - and "integrity_report"."createdAt" >= $2 - and "integrity_report"."createdAt" <= $3 + and "asset"."createdAt" >= $2 order by - "integrity_report"."createdAt" asc + "asset"."createdAt" asc -- IntegrityRepository.streamIntegrityReports select diff --git a/server/src/repositories/integrity.repository.ts b/server/src/repositories/integrity.repository.ts index 12f0b98634..19435650df 100644 --- a/server/src/repositories/integrity.repository.ts +++ b/server/src/repositories/integrity.repository.ts @@ -160,8 +160,8 @@ export class IntegrityRepository { >; } - @GenerateSql({ params: [DummyValue.DATE, DummyValue.DATE], stream: true }) - streamAssetChecksums(startMarker?: Date, endMarker?: Date) { + @GenerateSql({ params: [DummyValue.DATE], stream: true }) + streamAssetChecksums(startMarker?: Date) { return this.db .selectFrom('asset') .where('asset.deletedAt', 'is', null) @@ -178,9 +178,8 @@ export class IntegrityRepository { 'integrity_report.id as reportId', ]) .where('asset.isExternal', '=', sql.lit(false)) - .$if(startMarker !== undefined, (qb) => qb.where('integrity_report.createdAt', '>=', startMarker!)) - .$if(endMarker !== undefined, (qb) => qb.where('integrity_report.createdAt', '<=', endMarker!)) - .orderBy('integrity_report.createdAt', 'asc') + .$if(startMarker !== undefined, (qb) => qb.where('asset.createdAt', '>=', startMarker!)) + .orderBy('asset.createdAt', 'asc') .stream(); } diff --git a/server/src/services/integrity.service.ts b/server/src/services/integrity.service.ts index a2c603fde1..b3aa3fcd3f 100644 --- a/server/src/services/integrity.service.ts +++ b/server/src/services/integrity.service.ts @@ -532,8 +532,7 @@ export class IntegrityService extends BaseService { const { count } = await this.integrityRepository.getAssetCount(); const checkpoint = await this.systemMetadataRepository.get(SystemMetadataKey.IntegrityChecksumCheckpoint); - let startMarker: Date | undefined = checkpoint?.date ? new Date(checkpoint.date) : undefined; - let endMarker: Date | undefined; + const startMarker = checkpoint?.date ? new Date(checkpoint.date) : undefined; const printStats = () => { const averageTime = ((Date.now() - startedAt) / processed).toFixed(2); @@ -546,31 +545,25 @@ export class IntegrityService extends BaseService { let lastCreatedAt: Date | undefined; - finishEarly: do { - this.logger.log( - `Processing assets in range [${startMarker?.toISOString() ?? 'beginning'}, ${endMarker?.toISOString() ?? 'end'}]`, - ); + this.logger.log(`Processing assets from ${startMarker?.toISOString() ?? 'beginning'}`); - const assets = this.integrityRepository.streamAssetChecksums(startMarker, endMarker); - endMarker = startMarker; - startMarker = undefined; + const assets = this.integrityRepository.streamAssetChecksums(startMarker); - for await (const { originalPath, checksum, createdAt, assetId, reportId } of assets) { - await this.checkAssetChecksum(originalPath, checksum, assetId, reportId); + for await (const { originalPath, checksum, createdAt, assetId, reportId } of assets) { + await this.checkAssetChecksum(originalPath, checksum, assetId, reportId); - processed++; + processed++; - if (processed % 100 === 0) { - printStats(); - } - - if (Date.now() > startedAt + timeLimit || processed > count * percentageLimit) { - this.logger.log('Reached stop criteria.'); - lastCreatedAt = createdAt; - break finishEarly; - } + if (processed % 100 === 0) { + printStats(); } - } while (endMarker); + + if (Date.now() > startedAt + timeLimit || processed > count * percentageLimit) { + this.logger.log('Reached stop criteria.'); + lastCreatedAt = createdAt; + break; + } + } await this.systemMetadataRepository.set(SystemMetadataKey.IntegrityChecksumCheckpoint, { date: lastCreatedAt?.toISOString(), diff --git a/server/test/medium/specs/services/integrity.service.spec.ts b/server/test/medium/specs/services/integrity.service.spec.ts index db0f6f3ff2..f5de014867 100644 --- a/server/test/medium/specs/services/integrity.service.spec.ts +++ b/server/test/medium/specs/services/integrity.service.spec.ts @@ -1,9 +1,10 @@ import { Kysely } from 'kysely'; +import { DateTime } from 'luxon'; import { createHash, randomUUID } from 'node:crypto'; import { Readable } from 'node:stream'; import { text } from 'node:stream/consumers'; import { StorageCore } from 'src/cores/storage.core'; -import { AssetFileType, IntegrityReport, JobName, JobStatus } from 'src/enum'; +import { AssetFileType, IntegrityReport, JobName, JobStatus, SystemMetadataKey } from 'src/enum'; import { AssetRepository } from 'src/repositories/asset.repository'; import { ConfigRepository } from 'src/repositories/config.repository'; import { EventRepository } from 'src/repositories/event.repository'; @@ -702,6 +703,30 @@ describe(IntegrityService.name, () => { ctx.get(IntegrityRepository).getIntegrityReport({ limit: 100 }, IntegrityReport.ChecksumFail), ).resolves.toEqual({ items: [], nextCursor: undefined }); }); + + it('should continue from checkpoint', async () => { + const { sut, ctx } = setup(); + const job = ctx.getMock(JobRepository); + job.queue.mockResolvedValue(); + + const { user } = await ctx.newUser(); + + await ctx.newAsset({ + ownerId: user.id, + createdAt: DateTime.now().minus({ days: 1 }).toISO(), + originalPath: '/foo/bar', + }); + await ctx.newAsset({ ownerId: user.id, originalPath: '/foo/baz' }); + + await ctx + .get(SystemMetadataRepository) + .set(SystemMetadataKey.IntegrityChecksumCheckpoint, { date: DateTime.now().minus({ minutes: 5 }).toISO() }); + + await sut.handleChecksumFiles({ refreshOnly: false }); + + expect(ctx.getMock(StorageRepository).createPlainReadStream).not.toHaveBeenCalledWith('/foo/bar'); + expect(ctx.getMock(StorageRepository).createPlainReadStream).toHaveBeenCalledWith('/foo/baz'); + }); }); describe('handleChecksumRefresh', () => {