diff --git a/src/core/storage/sql/SqlKVRepository.ts b/src/core/storage/sql/SqlKVRepository.ts index c378f98c..02a0b313 100644 --- a/src/core/storage/sql/SqlKVRepository.ts +++ b/src/core/storage/sql/SqlKVRepository.ts @@ -79,7 +79,7 @@ export class SqlKVRepository { } } - private async upsertBatch(entities: Entity[]): Promise { + protected async upsertBatch(entities: Entity[]): Promise { const all = entities.map((entity) => this.dump(entity)); // make it unique by .id const data = lodash.uniqBy(all, (d: any) => d.id); diff --git a/src/core/storage/sqlite3/Sqlite3KVRepository.ts b/src/core/storage/sqlite3/Sqlite3KVRepository.ts index 27150a18..f06667d3 100644 --- a/src/core/storage/sqlite3/Sqlite3KVRepository.ts +++ b/src/core/storage/sqlite3/Sqlite3KVRepository.ts @@ -1,6 +1,7 @@ import { SqlKVRepository } from '@waha/core/storage/sql/SqlKVRepository'; import { Sqlite3Engine } from '@waha/core/storage/sqlite3/Sqlite3Engine'; import { Sqlite3JsonQuery } from '@waha/core/storage/sqlite3/Sqlite3JsonQuery'; +import { sleep } from '@waha/utils/promiseTimeout'; import { Database } from 'better-sqlite3'; import Knex from 'knex'; @@ -18,4 +19,12 @@ export class Sqlite3KVRepository extends SqlKVRepository { super(engine, knex); this.db = db; } + + protected async upsertBatch(entities: Entity[]): Promise { + await super.upsertBatch(entities); + // Give some time to the Node.js loop because we're using sync better-sqlite + if (entities.length >= this.UPSERT_BATCH_SIZE) { + await sleep(1); + } + } }