Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGLOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
# Changelog

## 1.8.2 (2026-07-24)

- Fix bug: rollback progress correctly logged

## 1.8.1 (2026-07-21)

- Allow update functions to return the NO_UPDATE symbol
Expand Down
42 changes: 41 additions & 1 deletion __tests__/MongoBulkDataMigration.rollback.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
const COLLECTION = 'testCollection';
const ROLLBACK_COLLECTION = `_rollback_${COLLECTION}_scriptId`;
const SCRIPT_ID = 'scriptId';
const ROLLBACK_BULK_LOG = 'Documents rollback is successful';
const ROLLBACK_SUMMARY_LOG = 'Documents rollbacked';

type DmDemoCollection = {
_id: ObjectId;
Expand All @@ -30,11 +32,11 @@
collectionName: string;
db: Db;
id: string;
query: any;

Check warning on line 35 in __tests__/MongoBulkDataMigration.rollback.test.ts

View workflow job for this annotation

GitHub Actions / lint

Unexpected any. Specify a different type
projection: any;

Check warning on line 36 in __tests__/MongoBulkDataMigration.rollback.test.ts

View workflow job for this annotation

GitHub Actions / lint

Unexpected any. Specify a different type
logger: LoggerInterface;
};
let loggerMock: LoggerInterface;
let loggerMock: jest.Mocked<LoggerInterface>;

beforeEach(async () => {
db = global.db;
Expand Down Expand Up @@ -141,7 +143,7 @@
const dataMigration = new MongoBulkDataMigration<DmDemoCollection>({
...DM_DEFAULT_SETUP,
projection: { a: 1, b: 1 },
update: ({ a, b }) => ({ $set: { a: a + b } }),

Check warning on line 146 in __tests__/MongoBulkDataMigration.rollback.test.ts

View workflow job for this annotation

GitHub Actions / lint

Missing return type on function
options: {
projectionBackupFilter: ['a'],
},
Expand Down Expand Up @@ -206,7 +208,7 @@
const insertedDocuments = await collection.find().toArray();
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
update: (doc: any) => ({ $set: { a: doc.a.b } }),

Check warning on line 211 in __tests__/MongoBulkDataMigration.rollback.test.ts

View workflow job for this annotation

GitHub Actions / lint

Missing return type on function
});

await dataMigration.update();
Expand Down Expand Up @@ -1104,5 +1106,43 @@
expect(restoredDocuments).toEqual([document]);
});
});

describe('bulk splitting', () => {
it('should perform updates in batches', async () => {
await collection.insertMany(
Array.from({ length: 100 }, (_, i) => ({ value: i + 1 })),
);
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
update: { $set: { value: 1000 } },
options: {
maxBulkSize: 30,
maxConcurrentUpdateCalls: 1000,
},
});

await dataMigration.update();
await dataMigration.rollback();

expect(extractLogsPayload(ROLLBACK_BULK_LOG)).toEqual([
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 10, nModified: 10, ok: 1 },
]);
expect(extractLogsPayload(ROLLBACK_SUMMARY_LOG)).toEqual([
{ treatedDocumentsCount: 30, totalEntries: 100, progress: '30.00' },
{ treatedDocumentsCount: 60, totalEntries: 100, progress: '60.00' },
{ treatedDocumentsCount: 90, totalEntries: 100, progress: '90.00' },
{ treatedDocumentsCount: 100, totalEntries: 100, progress: '100.00' },
]);
});
});
});

function extractLogsPayload(logText: string) {
return loggerMock.info.mock.calls
.filter(([_, msg]) => msg === logText)
.map(([payload]) => payload);
}
});
108 changes: 88 additions & 20 deletions __tests__/MongoBulkDataMigration.update.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,13 @@ import { Collection, Db, Document, ObjectId, UpdateFilter } from 'mongodb';
import { MongoBulkDataMigration, DELETE_OPERATION, FETCH_ALL } from '../src';
import { INITIAL_BULK_INFOS } from '../src/lib/AbstractBulkOperationResults';
import { LoggerInterface } from '../src/types';
import { NO_UPDATE } from '../src/MongoBulkDataMigration';
import { DONT_COUNT_LOG, NO_UPDATE } from '../src/MongoBulkDataMigration';

const COLLECTION = 'testCollection';
const SCRIPT_ID = 'scriptId';
const END_OF_BULK_LOG = 'Documents migration is successful';
const END_OF_ROLLBACK_BULK_LOG = 'Documents backup is successful';
const BULK_SUMMARY_LOG = 'Documents migrated';

describe('MongoBulkDataMigration', () => {
let db: Db;
Expand Down Expand Up @@ -226,7 +229,6 @@ describe('MongoBulkDataMigration', () => {
});

describe('Bulk splitting', () => {
const END_OF_BULK_LOG = 'Documents migration is successful';
let update: (arg: { value: number }) => UpdateFilter<{ value: number }>;
beforeEach(async () => {
await collection.insertMany(
Expand All @@ -238,6 +240,39 @@ describe('MongoBulkDataMigration', () => {
});

it('should perform updates in batches', async () => {
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
options: {
maxBulkSize: 30,
maxConcurrentUpdateCalls: 1000,
rollbackable: true,
},
update,
});

await dataMigration.update();

expect(extractLogsPayload(END_OF_BULK_LOG)).toEqual([
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 10, nModified: 10, ok: 1 },
]);
expect(extractLogsPayload(END_OF_ROLLBACK_BULK_LOG)).toEqual([
{ nUpserted: 30, ok: 1 },
{ nUpserted: 30, ok: 1 },
{ nUpserted: 30, ok: 1 },
{ nUpserted: 10, ok: 1 },
]);
expect(extractLogsPayload(BULK_SUMMARY_LOG)).toEqual([
{ treatedDocumentsCount: 30, totalEntries: 100, progress: '30.00' },
{ treatedDocumentsCount: 60, totalEntries: 100, progress: '60.00' },
{ treatedDocumentsCount: 90, totalEntries: 100, progress: '90.00' },
{ treatedDocumentsCount: 100, totalEntries: 100, progress: '100.00' },
]);
});

it('should perform updates in batches, even when not rollbackable', async () => {
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
options: {
Expand All @@ -250,29 +285,56 @@ describe('MongoBulkDataMigration', () => {

await dataMigration.update();

const batchSizes = loggerMock.info.mock.calls
.filter(([_, msg]) => msg === END_OF_BULK_LOG)
.map(([{ nModified }]) => nModified);
expect(batchSizes).toEqual([30, 30, 30, 10]);
expect(extractLogsPayload(END_OF_BULK_LOG)).toEqual([
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 30, nModified: 30, ok: 1 },
{ nMatched: 10, nModified: 10, ok: 1 },
]);
expect(extractLogsPayload(END_OF_ROLLBACK_BULK_LOG)).toEqual([]);
expect(extractLogsPayload(BULK_SUMMARY_LOG)).toEqual([
{ treatedDocumentsCount: 30, totalEntries: 100, progress: '30.00' },
{ treatedDocumentsCount: 60, totalEntries: 100, progress: '60.00' },
{ treatedDocumentsCount: 90, totalEntries: 100, progress: '90.00' },
{ treatedDocumentsCount: 100, totalEntries: 100, progress: '100.00' },
]);
});

it('should perform updates in batches, event when not rollbackable', async () => {
it('should not log progress when dontCount option is ON', async () => {
const dataMigration = new MongoBulkDataMigration({
...DM_DEFAULT_SETUP,
options: {
maxBulkSize: 30,
maxConcurrentUpdateCalls: 1000,
rollbackable: true,
dontCount: true,
},
update,
});

await dataMigration.update();

const batchSizes = loggerMock.info.mock.calls
.filter(([_, msg]) => msg === END_OF_BULK_LOG)
.map(([{ nModified }]) => nModified);
expect(batchSizes).toEqual([30, 30, 30, 10]);
expect(extractLogsPayload(BULK_SUMMARY_LOG)).toEqual([
{
treatedDocumentsCount: 30,
totalEntries: DONT_COUNT_LOG,
progress: DONT_COUNT_LOG,
},
{
treatedDocumentsCount: 60,
totalEntries: DONT_COUNT_LOG,
progress: DONT_COUNT_LOG,
},
{
treatedDocumentsCount: 90,
totalEntries: DONT_COUNT_LOG,
progress: DONT_COUNT_LOG,
},
{
treatedDocumentsCount: 100,
totalEntries: DONT_COUNT_LOG,
progress: DONT_COUNT_LOG,
},
]);
});
});

Expand Down Expand Up @@ -503,10 +565,12 @@ describe('MongoBulkDataMigration', () => {

await dataMigration.update();

const batchSizes = loggerMock.info.mock.calls
.filter(([_, msg]) => msg === 'Documents migration is successful')
.map(([{ nMatched, nModified }]) => ({ nMatched, nModified }));
expect(batchSizes).toEqual([{ nMatched: 2, nModified: 2 }]);
expect(extractLogsPayload(END_OF_BULK_LOG)).toEqual([
{ nMatched: 2, nModified: 2, ok: 1 },
]);
expect(extractLogsPayload(END_OF_ROLLBACK_BULK_LOG)).toEqual([
{ nUpserted: 2, ok: 1 },
]);
});

it('should not send any batch update if all documents are ignored', async () => {
Expand All @@ -518,10 +582,8 @@ describe('MongoBulkDataMigration', () => {

await dataMigration.update();

const batchSizes = loggerMock.info.mock.calls.filter(
([_, msg]) => msg === 'Documents migration is successful',
);
expect(batchSizes).toEqual([]);
expect(extractLogsPayload(END_OF_BULK_LOG)).toEqual([]);
expect(extractLogsPayload(END_OF_ROLLBACK_BULK_LOG)).toEqual([]);
});
});
});
Expand Down Expand Up @@ -693,4 +755,10 @@ describe('MongoBulkDataMigration', () => {
expect(updatedDocuments).toEqual([{ key: 1 }, { key: 3 }]);
});
});

function extractLogsPayload(logText: string) {
return loggerMock.info.mock.calls
.filter(([_, msg]) => msg === logText)
.map(([payload]) => payload);
}
});
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@360-l/mongo-bulk-data-migration",
"version": "1.8.1",
"version": "1.8.2",
"description": "MongoDB bulk data migration for node scripts",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
Expand Down
9 changes: 4 additions & 5 deletions src/MongoBulkDataMigration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ export const DELETE_COLLECTION = Symbol();
export const FETCH_ALL = Symbol();
/** Special return value for update function, when we don't need to update a given document */
export const NO_UPDATE = Symbol();
export const DONT_COUNT_LOG = 'N/A (dontCount option ON)';

const defaultLogger = {
info: (...args: unknown[]) => {
Expand Down Expand Up @@ -140,9 +141,7 @@ export default class MongoBulkDataMigration<
rollbackCollection,
);
const formattedTotalEntries =
totalEntries === NO_COUNT_AVAILABLE
? 'N/A (dontCount option ON)'
: totalEntries;
totalEntries === NO_COUNT_AVAILABLE ? DONT_COUNT_LOG : totalEntries;
this.logger.info(
{
collectionName: this.collectionName,
Expand Down Expand Up @@ -190,7 +189,7 @@ export default class MongoBulkDataMigration<
totalEntries: formattedTotalEntries,
progress:
totalEntries === NO_COUNT_AVAILABLE
? 'N/A (dontCount option ON)'
? DONT_COUNT_LOG
: ((treatedDocumentsCount / totalEntries) * 100).toFixed(2),
},
'Documents migrated',
Expand Down Expand Up @@ -384,9 +383,9 @@ export default class MongoBulkDataMigration<
rollbackDocument =
(await cursor.next()) as unknown as RollbackDocument | null;
if (!rollbackDocument || bulkRollback.size >= this.options.maxBulkSize) {
treatedDocumentsCount += bulkRollback.size;
await bulkRollback.execute();

treatedDocumentsCount += bulkRollback.size;
Comment thread
alexisvapillon marked this conversation as resolved.
this.logger.info(
{
treatedDocumentsCount,
Expand Down
Loading