From ac4ed0c1090034e44de4f97d0597e3f981801a19 Mon Sep 17 00:00:00 2001 From: Eric Doughty-Papassideris Date: Tue, 26 Mar 2024 12:06:22 +0100 Subject: [PATCH] =?UTF-8?q?=E2=99=BB=EF=B8=8F=20=20Factored=20out=20paging?= =?UTF-8?q?=20repository=20iteration=20(#438)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/cli/cmds/search_cmds/index_all.ts | 80 +++++++++++-------- 1 file changed, 48 insertions(+), 32 deletions(-) diff --git a/tdrive/backend/node/src/cli/cmds/search_cmds/index_all.ts b/tdrive/backend/node/src/cli/cmds/search_cmds/index_all.ts index 740236f8..d7b4e432 100644 --- a/tdrive/backend/node/src/cli/cmds/search_cmds/index_all.ts +++ b/tdrive/backend/node/src/cli/cmds/search_cmds/index_all.ts @@ -6,7 +6,9 @@ import { DatabaseServiceAPI } from "../../../core/platform/services/database/api import { Pagination } from "../../../core/platform/framework/api/crud-service"; import User, { TYPE as UserTYPE } from "../../../services/user/entities/user"; -import Repository from "../../../core/platform/services/database/services/orm/repository/repository"; +import Repository, { + FindFilter, +} from "../../../core/platform/services/database/services/orm/repository/repository"; import { EntityTarget, SearchServiceAPI } from "../../../core/platform/services/search/api"; import CompanyUser, { TYPE as CompanyUserTYPE } from "../../../services/user/entities/company_user"; import runWithPlatform from "../../lib/run_with_platform"; @@ -17,6 +19,24 @@ type Options = { repairEntities: boolean; }; +const waitTimeoutMS = (ms: number) => ms > 0 && new Promise(r => setTimeout(r, ms)); + +async function iterateOverRepoPages( + repository: Repository, + forEachPage: (entities: Entity[]) => Promise, + pageSizeAsStringForReasons: string = "100", + filter: FindFilter = {}, + delayPerPageMS: number = 200, +) { + let page: Pagination = { limitStr: pageSizeAsStringForReasons }; + do { + const list = await repository.find(filter, { pagination: page }, undefined); + await forEachPage(list.getEntities()); + page = list.nextPage as Pagination; + await waitTimeoutMS(delayPerPageMS); + } while (page.page_token); +} + type RepositoryConstructor = (database: DatabaseServiceAPI) => Promise>; const makeRepoConstructor = @@ -50,30 +70,31 @@ class SearchIndexAll { repository: Repository, ): Promise { // Complete user with companies in cache - options.spinner.start("Complete user with companies in cache"); + options.spinner.start("Adding companies to cache of user"); const companiesUsersRepository = await this.database.getRepository( CompanyUserTYPE, CompanyUser, ); - let page: Pagination = { limitStr: "100" }; let count = 0; - do { - const list = await repository.find({}, { pagination: page }, undefined); - - for (const user of list.getEntities()) { - const companies = await companiesUsersRepository.find({ user_id: user.id }, {}, undefined); - - user.cache ||= { companies: [] }; - user.cache.companies = companies.getEntities().map(company => company.group_id); - await repository.save(user, undefined); - } - count += list.getEntities().length; - options.spinner.start(`Completed ${count} users with companies in cache...`); - - page = list.nextPage as Pagination; - await new Promise(r => setTimeout(r, 2000)); - } while (page.page_token); - options.spinner.succeed(`Completed ${count} users with companies in cache`); + await iterateOverRepoPages( + repository, + async entities => { + for (const user of entities) { + const companies = await companiesUsersRepository.find( + { user_id: user.id }, + {}, + undefined, + ); + user.cache ||= { companies: [] }; + user.cache.companies = companies.getEntities().map(company => company.group_id); + await repository.save(user, undefined); + } + count += entities.length; + options.spinner.start(`Adding companies to cache of ${count} users...`); + }, + "2", + ); + options.spinner.succeed(`Added companies to cache of ${count} users`); } private repairEntities(options: Options, repository: Repository): Promise { @@ -93,20 +114,16 @@ class SearchIndexAll { options.spinner.start("Start indexing..."); let count = 0; - // Get all items - let page: Pagination = { limitStr: "100" }; - do { - const list = await repository.find({}, { pagination: page }, undefined); - page = list.nextPage as Pagination; - await this.search.upsert(list.getEntities()); - count += list.getEntities().length; + await iterateOverRepoPages(repository, async entities => { + await this.search.upsert(entities); + count += entities.length; options.spinner.start("Indexed " + count + " items..."); - await new Promise(r => setTimeout(r, 200)); - } while (page.page_token); + }); if (count === 0) options.spinner.warn("Index finieshed; but 0 items included"); else options.spinner.succeed(`${count} items indexed`); - options.spinner.start("Emptying flush (10s)..."); - await new Promise(r => setTimeout(r, 10000)); + const giveFlushAChanceDurationMS = 10000; + options.spinner.start(`Emptying flush (${giveFlushAChanceDurationMS / 1000}s)...`); + await waitTimeoutMS(giveFlushAChanceDurationMS); options.spinner.succeed("Done!"); } } @@ -145,7 +162,6 @@ const command: yargs.CommandModule = { spinner, repairEntities: !!argv.repairEntities, }); - return 0; } catch (err) { spinner.fail(`Error indexing: ${err.stack}`); return 1;