Merge dev: Node.js scraper migration + CI fix #32
@ -22,7 +22,7 @@ let mongoServer;
|
||||
let client;
|
||||
let db;
|
||||
|
||||
const PRICES_COLLECTION = 'unit_prices_migration_test';
|
||||
const PRICES_COLLECTION = 'unit_prices_scraper';
|
||||
|
||||
// Mock logger for capturing log calls
|
||||
function createMockLogger() {
|
||||
|
||||
@ -21,7 +21,7 @@ let mongoServer;
|
||||
let client;
|
||||
let db;
|
||||
|
||||
const UNITS_COLLECTION = 'units_migration_test';
|
||||
const UNITS_COLLECTION = 'units_scraper';
|
||||
|
||||
// Mock logger for capturing log calls
|
||||
function createMockLogger() {
|
||||
|
||||
@ -142,14 +142,14 @@ describe('config/scraper', () => {
|
||||
|
||||
describe('COLLECTIONS', () => {
|
||||
describe('default values', () => {
|
||||
test('UNITS should default to "units_migration_test"', () => {
|
||||
test('UNITS should default to "units_scraper"', () => {
|
||||
const config = require('../../config/scraper');
|
||||
expect(config.COLLECTIONS.UNITS).toBe('units_migration_test');
|
||||
expect(config.COLLECTIONS.UNITS).toBe('units_scraper');
|
||||
});
|
||||
|
||||
test('PRICES should default to "unit_prices_migration_test"', () => {
|
||||
test('PRICES should default to "unit_prices_scraper"', () => {
|
||||
const config = require('../../config/scraper');
|
||||
expect(config.COLLECTIONS.PRICES).toBe('unit_prices_migration_test');
|
||||
expect(config.COLLECTIONS.PRICES).toBe('unit_prices_scraper');
|
||||
});
|
||||
|
||||
test('DAILY_SUMMARIES should default to "daily_summaries"', () => {
|
||||
|
||||
@ -118,8 +118,8 @@ jest.mock('../../config/scraper', () => ({
|
||||
RETRY_CONFIG: { maxRetries: 3, baseDelay: 1000, timeout: 30000 },
|
||||
SHUTDOWN_TIMEOUT: 30000,
|
||||
COLLECTIONS: {
|
||||
UNITS: 'units_migration_test',
|
||||
PRICES: 'unit_prices_migration_test',
|
||||
UNITS: 'units_scraper',
|
||||
PRICES: 'unit_prices_scraper',
|
||||
DAILY_SUMMARIES: 'daily_summaries',
|
||||
SCRAPER_RUNS: 'scraper_runs'
|
||||
}
|
||||
|
||||
@ -20,9 +20,10 @@ let mongoServer;
|
||||
let client;
|
||||
let db;
|
||||
|
||||
// We will require runScrape and recordScraperRun after implementation
|
||||
// We will require runScrape, recordScraperRun, and createScraperIndexes after implementation
|
||||
let runScrape;
|
||||
let recordScraperRun;
|
||||
let createScraperIndexes;
|
||||
|
||||
// Store original module references so we can mock individual functions
|
||||
let scraperService;
|
||||
@ -36,7 +37,7 @@ beforeAll(async () => {
|
||||
|
||||
// Dynamically require to pick up implementation
|
||||
scraperService = require('../../services/scraperService');
|
||||
({ runScrape, recordScraperRun } = scraperService);
|
||||
({ runScrape, recordScraperRun, createScraperIndexes } = scraperService);
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
@ -219,11 +220,11 @@ describe('runScrape', () => {
|
||||
expect(result.staleUnitsCount).toBe(0);
|
||||
|
||||
// Verify no units were written to the units collection
|
||||
const unitsCount = await db.collection('units_migration_test').countDocuments();
|
||||
const unitsCount = await db.collection('units_scraper').countDocuments();
|
||||
expect(unitsCount).toBe(0);
|
||||
|
||||
// Verify no prices were written
|
||||
const pricesCount = await db.collection('unit_prices_migration_test').countDocuments();
|
||||
const pricesCount = await db.collection('unit_prices_scraper').countDocuments();
|
||||
expect(pricesCount).toBe(0);
|
||||
});
|
||||
|
||||
@ -277,7 +278,7 @@ describe('runScrape', () => {
|
||||
const faultyDb = {
|
||||
collection: (name) => {
|
||||
const realCollection = db.collection(name);
|
||||
if (name === 'units_migration_test') {
|
||||
if (name === 'units_scraper') {
|
||||
return new Proxy(realCollection, {
|
||||
get(target, prop) {
|
||||
if (prop === 'bulkWrite') {
|
||||
@ -343,7 +344,7 @@ describe('runScrape', () => {
|
||||
const faultyDb = {
|
||||
collection: (name) => {
|
||||
const realCollection = db.collection(name);
|
||||
if (name === 'unit_prices_migration_test') {
|
||||
if (name === 'unit_prices_scraper') {
|
||||
return new Proxy(realCollection, {
|
||||
get(target, prop) {
|
||||
if (prop === 'bulkWrite') {
|
||||
@ -388,7 +389,7 @@ describe('runScrape', () => {
|
||||
const yesterdayStr = yesterday.toISOString().split('T')[0];
|
||||
|
||||
// Yesterday had units: OLD-A, OLD-B, OLD-C
|
||||
await db.collection('unit_prices_migration_test').insertMany([
|
||||
await db.collection('unit_prices_scraper').insertMany([
|
||||
{ unit_code: 'OLD-A', date_checked: yesterdayStr, price: 1000 },
|
||||
{ unit_code: 'OLD-B', date_checked: yesterdayStr, price: 1100 },
|
||||
{ unit_code: 'OLD-C', date_checked: yesterdayStr, price: 1200 }
|
||||
@ -631,7 +632,7 @@ describe('sanitizeError', () => {
|
||||
const faultyDb = {
|
||||
collection: (name) => {
|
||||
const realCollection = db.collection(name);
|
||||
if (name === 'units_migration_test') {
|
||||
if (name === 'units_scraper') {
|
||||
return new Proxy(realCollection, {
|
||||
get(target, prop) {
|
||||
if (prop === 'bulkWrite') {
|
||||
@ -674,3 +675,65 @@ describe('sanitizeError', () => {
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
// ============================================================
|
||||
// Test: createScraperIndexes
|
||||
// ============================================================
|
||||
describe('createScraperIndexes', () => {
|
||||
const mockLogger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() };
|
||||
|
||||
beforeEach(() => {
|
||||
mockLogger.info.mockClear();
|
||||
mockLogger.warn.mockClear();
|
||||
mockLogger.error.mockClear();
|
||||
});
|
||||
|
||||
it('should create the status_startedAt compound index', async () => {
|
||||
await createScraperIndexes(db, mockLogger);
|
||||
|
||||
const indexes = await db.collection('scraper_runs').indexes();
|
||||
const statusIndex = indexes.find(idx => idx.name === 'status_startedAt');
|
||||
expect(statusIndex).toBeTruthy();
|
||||
expect(statusIndex.key).toEqual({ status: 1, startedAt: -1 });
|
||||
});
|
||||
|
||||
it('should create the startedAt_desc index', async () => {
|
||||
await createScraperIndexes(db, mockLogger);
|
||||
|
||||
const indexes = await db.collection('scraper_runs').indexes();
|
||||
const startedAtIndex = indexes.find(idx => idx.name === 'startedAt_desc');
|
||||
expect(startedAtIndex).toBeTruthy();
|
||||
expect(startedAtIndex.key).toEqual({ startedAt: -1 });
|
||||
});
|
||||
|
||||
it('should log success after creating indexes', async () => {
|
||||
await createScraperIndexes(db, mockLogger);
|
||||
|
||||
expect(mockLogger.info).toHaveBeenCalledWith('Scraper indexes created successfully');
|
||||
});
|
||||
|
||||
it('should be idempotent (no error on second call)', async () => {
|
||||
await createScraperIndexes(db, mockLogger);
|
||||
|
||||
// Calling again should not throw
|
||||
await expect(createScraperIndexes(db, mockLogger)).resolves.not.toThrow();
|
||||
|
||||
// Indexes should still exist
|
||||
const indexes = await db.collection('scraper_runs').indexes();
|
||||
const statusIndex = indexes.find(idx => idx.name === 'status_startedAt');
|
||||
const startedAtIndex = indexes.find(idx => idx.name === 'startedAt_desc');
|
||||
expect(statusIndex).toBeTruthy();
|
||||
expect(startedAtIndex).toBeTruthy();
|
||||
});
|
||||
|
||||
it('should handle database errors gracefully', async () => {
|
||||
// Pass an object that will throw when collection() is called
|
||||
const faultyDb = {
|
||||
collection: () => {
|
||||
throw new Error('Connection lost');
|
||||
}
|
||||
};
|
||||
|
||||
await expect(createScraperIndexes(faultyDb, mockLogger)).rejects.toThrow('Connection lost');
|
||||
});
|
||||
});
|
||||
|
||||
@ -56,7 +56,7 @@ beforeEach(async () => {
|
||||
});
|
||||
|
||||
describe('upsertUnits', () => {
|
||||
const UNITS_COLLECTION = 'units_migration_test';
|
||||
const UNITS_COLLECTION = 'units_scraper';
|
||||
|
||||
// ---------------------------------------------------------------
|
||||
// 1. bulkWrite is called with correct updateOne operations
|
||||
|
||||
@ -30,10 +30,11 @@ module.exports = {
|
||||
// Graceful shutdown timeout (how long to wait for running job before force-stopping)
|
||||
SHUTDOWN_TIMEOUT: parseInt(process.env.SCRAPER_SHUTDOWN_TIMEOUT, 10) || 30000,
|
||||
|
||||
// MongoDB collection names (environment variable overrides for development isolation)
|
||||
// MongoDB collection names - defaults to validation collections for scraper validation;
|
||||
// switch to production collections via env vars (SCRAPER_UNITS_COLLECTION, SCRAPER_PRICES_COLLECTION) when ready
|
||||
COLLECTIONS: {
|
||||
UNITS: process.env.SCRAPER_UNITS_COLLECTION || 'units_migration_test',
|
||||
PRICES: process.env.SCRAPER_PRICES_COLLECTION || 'unit_prices_migration_test',
|
||||
UNITS: process.env.SCRAPER_UNITS_COLLECTION || 'units_scraper',
|
||||
PRICES: process.env.SCRAPER_PRICES_COLLECTION || 'unit_prices_scraper',
|
||||
DAILY_SUMMARIES: process.env.SCRAPER_SUMMARIES_COLLECTION || 'daily_summaries',
|
||||
SCRAPER_RUNS: process.env.SCRAPER_RUNS_COLLECTION || 'scraper_runs'
|
||||
}
|
||||
|
||||
@ -730,6 +730,37 @@ async function recordScraperRun(db, runData, logger) {
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Index Management
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Create indexes on the scraper_runs collection.
|
||||
* Should be called once during application startup.
|
||||
* Idempotent - safe to call multiple times.
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
async function createScraperIndexes(db, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.SCRAPER_RUNS);
|
||||
|
||||
// Compound index for querying runs by status sorted by most recent
|
||||
await collection.createIndex(
|
||||
{ status: 1, startedAt: -1 },
|
||||
{ name: 'status_startedAt' }
|
||||
);
|
||||
|
||||
// Index for sorting all runs by start time (most recent first)
|
||||
await collection.createIndex(
|
||||
{ startedAt: -1 },
|
||||
{ name: 'startedAt_desc' }
|
||||
);
|
||||
|
||||
logger.info('Scraper indexes created successfully');
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Main Orchestration Function
|
||||
// ============================================================
|
||||
@ -871,6 +902,7 @@ module.exports = {
|
||||
markStaleUnits,
|
||||
updateDailySummary,
|
||||
recordScraperRun,
|
||||
createScraperIndexes,
|
||||
// Export helpers for testing
|
||||
getTodayUTC,
|
||||
getYesterdayUnitCodes,
|
||||
|
||||
Reference in New Issue
Block a user