SCRAPE-6: Implement upsertUnits() bulk operation #11

Merged
stephen merged 1 commits from scraper/upsert-units into dev 2026-02-05 22:10:54 -07:00
2 changed files with 560 additions and 0 deletions
Showing only changes of commit 31c3c0c8ee - Show all commits

View File

@ -0,0 +1,498 @@
/**
* Tests for upsertUnits() bulk operation
*
* Covers:
* - bulkWrite is called with correct updateOne operations
* - upsert: true is set for all operations
* - $set updates last_scraped and data_source fields
* - $setOnInsert sets first_seen for new units
* - Empty units array logs warning and returns early
* - ordered: false is set for parallel execution
* - Error handling when bulkWrite fails
* - Logging of matched, modified, and upserted counts
*/
const { MongoClient } = require('mongodb');
const { MongoMemoryServer } = require('mongodb-memory-server');
// We will require upsertUnits after implementation
let upsertUnits;
let mongoServer;
let client;
let db;
// Mock logger for capturing log calls
function createMockLogger() {
return {
info: jest.fn(),
warn: jest.fn(),
error: jest.fn()
};
}
beforeAll(async () => {
mongoServer = await MongoMemoryServer.create();
const uri = mongoServer.getUri();
client = new MongoClient(uri);
await client.connect();
db = client.db('test_apartments');
// Dynamically require to pick up implementation
({ upsertUnits } = require('../../services/scraperService'));
});
afterAll(async () => {
if (client) await client.close();
if (mongoServer) await mongoServer.stop();
});
beforeEach(async () => {
// Clean the units collection before each test
const collections = await db.listCollections().toArray();
for (const col of collections) {
await db.collection(col.name).deleteMany({});
}
});
describe('upsertUnits', () => {
const UNITS_COLLECTION = 'units_migration_test';
// ---------------------------------------------------------------
// 1. bulkWrite is called with correct updateOne operations
// ---------------------------------------------------------------
describe('bulkWrite operations', () => {
it('should create updateOne operations for each unit', async () => {
const logger = createMockLogger();
const units = [
{ unit_code: 'A101', price: 1200, bed_count: 1, area: 500 },
{ unit_code: 'B202', price: 1500, bed_count: 2, area: 750 }
];
const result = await upsertUnits(db, units, logger);
// Both units should be upserted (new inserts)
expect(result.upsertedCount).toBe(2);
// Verify documents exist in the collection
const docs = await db.collection(UNITS_COLLECTION).find({}).toArray();
expect(docs).toHaveLength(2);
const unitCodes = docs.map(d => d.unit_code).sort();
expect(unitCodes).toEqual(['A101', 'B202']);
});
it('should filter by unit_code as the unique key', async () => {
const logger = createMockLogger();
// Insert a unit first
const units = [{ unit_code: 'A101', price: 1200, bed_count: 1 }];
await upsertUnits(db, units, logger);
// Upsert again with updated price
const updatedUnits = [{ unit_code: 'A101', price: 1400, bed_count: 1 }];
const result = await upsertUnits(db, updatedUnits, logger);
// Should match and modify, not insert new
expect(result.matchedCount).toBe(1);
expect(result.upsertedCount).toBe(0);
// Verify only one document exists
const docs = await db.collection(UNITS_COLLECTION).find({}).toArray();
expect(docs).toHaveLength(1);
expect(docs[0].price).toBe(1400);
});
});
// ---------------------------------------------------------------
// 2. upsert: true is set for all operations
// ---------------------------------------------------------------
describe('upsert behavior', () => {
it('should insert new units when they do not exist (upsert: true)', async () => {
const logger = createMockLogger();
const units = [
{ unit_code: 'NEW-001', price: 1100 },
{ unit_code: 'NEW-002', price: 1200 },
{ unit_code: 'NEW-003', price: 1300 }
];
const result = await upsertUnits(db, units, logger);
expect(result.upsertedCount).toBe(3);
// Verify all three are in the database
const count = await db.collection(UNITS_COLLECTION).countDocuments();
expect(count).toBe(3);
});
it('should update existing units without creating duplicates', async () => {
const logger = createMockLogger();
// First insert
await upsertUnits(db, [{ unit_code: 'X100', price: 900 }], logger);
// Second upsert with same unit_code
await upsertUnits(db, [{ unit_code: 'X100', price: 950 }], logger);
const count = await db.collection(UNITS_COLLECTION).countDocuments();
expect(count).toBe(1);
});
});
// ---------------------------------------------------------------
// 3. $set updates last_scraped and data_source fields
// ---------------------------------------------------------------
describe('$set fields', () => {
it('should set last_scraped timestamp on every upsert', async () => {
const logger = createMockLogger();
const beforeTime = new Date().toISOString();
const units = [{ unit_code: 'T100', price: 1000 }];
await upsertUnits(db, units, logger);
const afterTime = new Date().toISOString();
const doc = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'T100' });
expect(doc.last_scraped).toBeDefined();
expect(typeof doc.last_scraped).toBe('string');
// The last_scraped should be between before and after timestamps
expect(doc.last_scraped >= beforeTime).toBe(true);
expect(doc.last_scraped <= afterTime).toBe(true);
});
it('should set data_source to "web_scraper"', async () => {
const logger = createMockLogger();
const units = [{ unit_code: 'T200', price: 1100 }];
await upsertUnits(db, units, logger);
const doc = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'T200' });
expect(doc.data_source).toBe('web_scraper');
});
it('should update last_scraped on subsequent upserts', async () => {
const logger = createMockLogger();
// First insert
await upsertUnits(db, [{ unit_code: 'T300', price: 800 }], logger);
const doc1 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'T300' });
const firstScraped = doc1.last_scraped;
// Small delay to ensure different timestamp
await new Promise(resolve => setTimeout(resolve, 10));
// Second upsert
await upsertUnits(db, [{ unit_code: 'T300', price: 850 }], logger);
const doc2 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'T300' });
const secondScraped = doc2.last_scraped;
expect(secondScraped > firstScraped).toBe(true);
});
it('should spread all unit fields into the $set operation', async () => {
const logger = createMockLogger();
const units = [{
unit_code: 'T400',
price: 1500,
bed_count: 2,
bath_count: 1.5,
area: 900,
floor: 3,
plan_name: 'Luxury',
community: 'Tower A'
}];
await upsertUnits(db, units, logger);
const doc = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'T400' });
expect(doc.price).toBe(1500);
expect(doc.bed_count).toBe(2);
expect(doc.bath_count).toBe(1.5);
expect(doc.area).toBe(900);
expect(doc.floor).toBe(3);
expect(doc.plan_name).toBe('Luxury');
expect(doc.community).toBe('Tower A');
});
});
// ---------------------------------------------------------------
// 4. $setOnInsert sets first_seen for new units
// ---------------------------------------------------------------
describe('$setOnInsert for first_seen', () => {
it('should set first_seen on initial insert', async () => {
const logger = createMockLogger();
const beforeTime = new Date().toISOString();
await upsertUnits(db, [{ unit_code: 'F100', price: 1000 }], logger);
const doc = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'F100' });
expect(doc.first_seen).toBeDefined();
expect(typeof doc.first_seen).toBe('string');
expect(doc.first_seen >= beforeTime).toBe(true);
});
it('should NOT update first_seen on subsequent upserts', async () => {
const logger = createMockLogger();
// First insert
await upsertUnits(db, [{ unit_code: 'F200', price: 1000 }], logger);
const doc1 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'F200' });
const originalFirstSeen = doc1.first_seen;
// Small delay to ensure different timestamp
await new Promise(resolve => setTimeout(resolve, 10));
// Second upsert (update)
await upsertUnits(db, [{ unit_code: 'F200', price: 1100 }], logger);
const doc2 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'F200' });
// first_seen should remain unchanged
expect(doc2.first_seen).toBe(originalFirstSeen);
});
it('should set different first_seen for different units inserted at different times', async () => {
const logger = createMockLogger();
await upsertUnits(db, [{ unit_code: 'F300', price: 1000 }], logger);
const doc1 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'F300' });
await new Promise(resolve => setTimeout(resolve, 10));
await upsertUnits(db, [{ unit_code: 'F400', price: 1200 }], logger);
const doc2 = await db.collection(UNITS_COLLECTION).findOne({ unit_code: 'F400' });
// Both should have first_seen, but different values
expect(doc1.first_seen).toBeDefined();
expect(doc2.first_seen).toBeDefined();
expect(doc2.first_seen > doc1.first_seen).toBe(true);
});
});
// ---------------------------------------------------------------
// 5. Empty units array logs warning and returns early
// ---------------------------------------------------------------
describe('empty units array', () => {
it('should log a warning when units array is empty', async () => {
const logger = createMockLogger();
await upsertUnits(db, [], logger);
expect(logger.warn).toHaveBeenCalledWith('No units to upsert');
});
it('should return early with zero counts for empty array', async () => {
const logger = createMockLogger();
const result = await upsertUnits(db, [], logger);
expect(result.modifiedCount).toBe(0);
expect(result.upsertedCount).toBe(0);
});
it('should not perform any database writes for empty array', async () => {
const logger = createMockLogger();
await upsertUnits(db, [], logger);
const count = await db.collection(UNITS_COLLECTION).countDocuments();
expect(count).toBe(0);
});
});
// ---------------------------------------------------------------
// 6. ordered: false is set for parallel execution
// ---------------------------------------------------------------
describe('ordered: false for parallel execution', () => {
it('should successfully process all units even if one has a conflict', async () => {
const logger = createMockLogger();
// Insert multiple units - ordered: false means all operations execute
// even if some fail. We verify by checking all valid units are present.
const units = [
{ unit_code: 'P100', price: 1000 },
{ unit_code: 'P200', price: 1100 },
{ unit_code: 'P300', price: 1200 },
{ unit_code: 'P400', price: 1300 },
{ unit_code: 'P500', price: 1400 }
];
const result = await upsertUnits(db, units, logger);
expect(result.upsertedCount).toBe(5);
// All 5 units should be in the database
const count = await db.collection(UNITS_COLLECTION).countDocuments();
expect(count).toBe(5);
});
it('should handle a mix of new inserts and existing updates', async () => {
const logger = createMockLogger();
// First, insert some units
await upsertUnits(db, [
{ unit_code: 'M100', price: 900 },
{ unit_code: 'M200', price: 1000 }
], logger);
// Now upsert a mix of existing and new units
const mixedUnits = [
{ unit_code: 'M100', price: 950 }, // existing - update
{ unit_code: 'M200', price: 1050 }, // existing - update
{ unit_code: 'M300', price: 1100 }, // new - insert
];
const result = await upsertUnits(db, mixedUnits, logger);
expect(result.matchedCount).toBe(2);
expect(result.upsertedCount).toBe(1);
// Total should be 3 documents
const count = await db.collection(UNITS_COLLECTION).countDocuments();
expect(count).toBe(3);
});
});
// ---------------------------------------------------------------
// 7. Error handling when bulkWrite fails
// ---------------------------------------------------------------
describe('error handling', () => {
it('should throw error when bulkWrite fails', async () => {
const logger = createMockLogger();
const units = [{ unit_code: 'E100', price: 1000 }];
// Create a mock db that throws on bulkWrite
const mockCollection = {
bulkWrite: jest.fn().mockRejectedValue(new Error('Connection lost'))
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
await expect(upsertUnits(mockDb, units, logger)).rejects.toThrow('Connection lost');
});
it('should log error details when bulkWrite fails', async () => {
const logger = createMockLogger();
const units = [{ unit_code: 'E200', price: 1000 }];
const mockCollection = {
bulkWrite: jest.fn().mockRejectedValue(new Error('Write concern timeout'))
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
try {
await upsertUnits(mockDb, units, logger);
} catch (e) {
// Expected to throw
}
expect(logger.error).toHaveBeenCalledWith(
'Failed to upsert units',
expect.objectContaining({
errorType: 'Error',
errorMessage: 'Write concern timeout'
})
);
});
it('should re-throw the original error', async () => {
const logger = createMockLogger();
const units = [{ unit_code: 'E300', price: 1000 }];
const originalError = new TypeError('Invalid operation');
const mockCollection = {
bulkWrite: jest.fn().mockRejectedValue(originalError)
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
await expect(upsertUnits(mockDb, units, logger)).rejects.toBe(originalError);
});
});
// ---------------------------------------------------------------
// 8. Logging of matched, modified, and upserted counts
// ---------------------------------------------------------------
describe('logging of result counts', () => {
it('should log matched, modified, and upserted counts on success', async () => {
const logger = createMockLogger();
const units = [
{ unit_code: 'L100', price: 1000 },
{ unit_code: 'L200', price: 1100 }
];
await upsertUnits(db, units, logger);
expect(logger.info).toHaveBeenCalledWith(
'Units upserted',
expect.objectContaining({
matched: expect.any(Number),
modified: expect.any(Number),
upserted: expect.any(Number)
})
);
});
it('should log correct counts for new inserts', async () => {
const logger = createMockLogger();
const units = [
{ unit_code: 'L300', price: 1000 },
{ unit_code: 'L400', price: 1100 }
];
await upsertUnits(db, units, logger);
// For new inserts: matched=0, modified=0, upserted=2
expect(logger.info).toHaveBeenCalledWith(
'Units upserted',
expect.objectContaining({
matched: 0,
modified: 0,
upserted: 2
})
);
});
it('should log correct counts for updates to existing units', async () => {
const logger = createMockLogger();
// First insert
await upsertUnits(db, [{ unit_code: 'L500', price: 1000 }], logger);
// Reset mock to capture only the second call
logger.info.mockClear();
// Update existing
await upsertUnits(db, [{ unit_code: 'L500', price: 1100 }], logger);
expect(logger.info).toHaveBeenCalledWith(
'Units upserted',
expect.objectContaining({
matched: 1,
modified: 1,
upserted: 0
})
);
});
});
// ---------------------------------------------------------------
// 9. Return value
// ---------------------------------------------------------------
describe('return value', () => {
it('should return the bulkWrite result object', async () => {
const logger = createMockLogger();
const units = [{ unit_code: 'R100', price: 1000 }];
const result = await upsertUnits(db, units, logger);
// bulkWrite result should have these standard properties
expect(result).toHaveProperty('matchedCount');
expect(result).toHaveProperty('modifiedCount');
expect(result).toHaveProperty('upsertedCount');
});
});
});

View File

@ -369,10 +369,72 @@ function convertDataTypes(unit) {
}; };
} }
// ============================================================
// Database Operations
// ============================================================
/**
* Upsert unit records to database using bulkWrite.
* Each unit is matched by unit_code as the unique key.
* Sets last_scraped and data_source on every update.
* Sets first_seen only on initial insert via $setOnInsert.
*
* @param {Db} db - MongoDB database instance
* @param {Array<Object>} units - Array of unit objects
* @param {Object} logger - Logger instance
* @returns {Promise<Object>} Bulk write result
*/
async function upsertUnits(db, units, logger) {
const collection = db.collection(config.COLLECTIONS.UNITS);
const now = new Date().toISOString();
const operations = units.map(unit => ({
updateOne: {
filter: { unit_code: unit.unit_code },
update: {
$set: {
...unit,
last_scraped: now,
data_source: 'web_scraper'
},
$setOnInsert: {
first_seen: now
}
},
upsert: true
}
}));
if (operations.length === 0) {
logger.warn('No units to upsert');
return { modifiedCount: 0, upsertedCount: 0 };
}
try {
const result = await collection.bulkWrite(operations, { ordered: false });
logger.info('Units upserted', {
matched: result.matchedCount,
modified: result.modifiedCount,
upserted: result.upsertedCount
});
return result;
} catch (error) {
logger.error('Failed to upsert units', {
errorType: error.name,
errorMessage: error.message
});
throw error;
}
}
module.exports = { module.exports = {
fetchPage, fetchPage,
parseUnits, parseUnits,
convertDataTypes, convertDataTypes,
upsertUnits,
// Export helpers for testing // Export helpers for testing
parseInteger, parseInteger,
parsePositiveInteger, parsePositiveInteger,