SCRAPE-31: Add scraper indexes and validation collections #30

Merged
stephen merged 2 commits from scraper/indexes into dev 2026-02-07 16:50:51 -07:00
8 changed files with 116 additions and 20 deletions

View File

@ -22,7 +22,7 @@ let mongoServer;
let client; let client;
let db; let db;
const PRICES_COLLECTION = 'unit_prices_migration_test'; const PRICES_COLLECTION = 'unit_prices_scraper';
// Mock logger for capturing log calls // Mock logger for capturing log calls
function createMockLogger() { function createMockLogger() {

View File

@ -21,7 +21,7 @@ let mongoServer;
let client; let client;
let db; let db;
const UNITS_COLLECTION = 'units_migration_test'; const UNITS_COLLECTION = 'units_scraper';
// Mock logger for capturing log calls // Mock logger for capturing log calls
function createMockLogger() { function createMockLogger() {

View File

@ -142,14 +142,14 @@ describe('config/scraper', () => {
describe('COLLECTIONS', () => { describe('COLLECTIONS', () => {
describe('default values', () => { describe('default values', () => {
test('UNITS should default to "units_migration_test"', () => { test('UNITS should default to "units_scraper"', () => {
const config = require('../../config/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'); 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"', () => { test('DAILY_SUMMARIES should default to "daily_summaries"', () => {

View File

@ -118,8 +118,8 @@ jest.mock('../../config/scraper', () => ({
RETRY_CONFIG: { maxRetries: 3, baseDelay: 1000, timeout: 30000 }, RETRY_CONFIG: { maxRetries: 3, baseDelay: 1000, timeout: 30000 },
SHUTDOWN_TIMEOUT: 30000, SHUTDOWN_TIMEOUT: 30000,
COLLECTIONS: { COLLECTIONS: {
UNITS: 'units_migration_test', UNITS: 'units_scraper',
PRICES: 'unit_prices_migration_test', PRICES: 'unit_prices_scraper',
DAILY_SUMMARIES: 'daily_summaries', DAILY_SUMMARIES: 'daily_summaries',
SCRAPER_RUNS: 'scraper_runs' SCRAPER_RUNS: 'scraper_runs'
} }

View File

@ -20,9 +20,10 @@ let mongoServer;
let client; let client;
let db; let db;
// We will require runScrape and recordScraperRun after implementation // We will require runScrape, recordScraperRun, and createScraperIndexes after implementation
let runScrape; let runScrape;
let recordScraperRun; let recordScraperRun;
let createScraperIndexes;
// Store original module references so we can mock individual functions // Store original module references so we can mock individual functions
let scraperService; let scraperService;
@ -36,7 +37,7 @@ beforeAll(async () => {
// Dynamically require to pick up implementation // Dynamically require to pick up implementation
scraperService = require('../../services/scraperService'); scraperService = require('../../services/scraperService');
({ runScrape, recordScraperRun } = scraperService); ({ runScrape, recordScraperRun, createScraperIndexes } = scraperService);
}); });
afterAll(async () => { afterAll(async () => {
@ -219,11 +220,11 @@ describe('runScrape', () => {
expect(result.staleUnitsCount).toBe(0); expect(result.staleUnitsCount).toBe(0);
// Verify no units were written to the units collection // 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); expect(unitsCount).toBe(0);
// Verify no prices were written // 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); expect(pricesCount).toBe(0);
}); });
@ -277,7 +278,7 @@ describe('runScrape', () => {
const faultyDb = { const faultyDb = {
collection: (name) => { collection: (name) => {
const realCollection = db.collection(name); const realCollection = db.collection(name);
if (name === 'units_migration_test') { if (name === 'units_scraper') {
return new Proxy(realCollection, { return new Proxy(realCollection, {
get(target, prop) { get(target, prop) {
if (prop === 'bulkWrite') { if (prop === 'bulkWrite') {
@ -343,7 +344,7 @@ describe('runScrape', () => {
const faultyDb = { const faultyDb = {
collection: (name) => { collection: (name) => {
const realCollection = db.collection(name); const realCollection = db.collection(name);
if (name === 'unit_prices_migration_test') { if (name === 'unit_prices_scraper') {
return new Proxy(realCollection, { return new Proxy(realCollection, {
get(target, prop) { get(target, prop) {
if (prop === 'bulkWrite') { if (prop === 'bulkWrite') {
@ -388,7 +389,7 @@ describe('runScrape', () => {
const yesterdayStr = yesterday.toISOString().split('T')[0]; const yesterdayStr = yesterday.toISOString().split('T')[0];
// Yesterday had units: OLD-A, OLD-B, OLD-C // 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-A', date_checked: yesterdayStr, price: 1000 },
{ unit_code: 'OLD-B', date_checked: yesterdayStr, price: 1100 }, { unit_code: 'OLD-B', date_checked: yesterdayStr, price: 1100 },
{ unit_code: 'OLD-C', date_checked: yesterdayStr, price: 1200 } { unit_code: 'OLD-C', date_checked: yesterdayStr, price: 1200 }
@ -631,7 +632,7 @@ describe('sanitizeError', () => {
const faultyDb = { const faultyDb = {
collection: (name) => { collection: (name) => {
const realCollection = db.collection(name); const realCollection = db.collection(name);
if (name === 'units_migration_test') { if (name === 'units_scraper') {
return new Proxy(realCollection, { return new Proxy(realCollection, {
get(target, prop) { get(target, prop) {
if (prop === 'bulkWrite') { 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');
});
});

View File

@ -56,7 +56,7 @@ beforeEach(async () => {
}); });
describe('upsertUnits', () => { describe('upsertUnits', () => {
const UNITS_COLLECTION = 'units_migration_test'; const UNITS_COLLECTION = 'units_scraper';
// --------------------------------------------------------------- // ---------------------------------------------------------------
// 1. bulkWrite is called with correct updateOne operations // 1. bulkWrite is called with correct updateOne operations

View File

@ -30,10 +30,11 @@ module.exports = {
// Graceful shutdown timeout (how long to wait for running job before force-stopping) // Graceful shutdown timeout (how long to wait for running job before force-stopping)
SHUTDOWN_TIMEOUT: parseInt(process.env.SCRAPER_SHUTDOWN_TIMEOUT, 10) || 30000, 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: { COLLECTIONS: {
UNITS: process.env.SCRAPER_UNITS_COLLECTION || 'units_migration_test', UNITS: process.env.SCRAPER_UNITS_COLLECTION || 'units_scraper',
PRICES: process.env.SCRAPER_PRICES_COLLECTION || 'unit_prices_migration_test', PRICES: process.env.SCRAPER_PRICES_COLLECTION || 'unit_prices_scraper',
DAILY_SUMMARIES: process.env.SCRAPER_SUMMARIES_COLLECTION || 'daily_summaries', DAILY_SUMMARIES: process.env.SCRAPER_SUMMARIES_COLLECTION || 'daily_summaries',
SCRAPER_RUNS: process.env.SCRAPER_RUNS_COLLECTION || 'scraper_runs' SCRAPER_RUNS: process.env.SCRAPER_RUNS_COLLECTION || 'scraper_runs'
} }

View File

@ -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 // Main Orchestration Function
// ============================================================ // ============================================================
@ -871,6 +902,7 @@ module.exports = {
markStaleUnits, markStaleUnits,
updateDailySummary, updateDailySummary,
recordScraperRun, recordScraperRun,
createScraperIndexes,
// Export helpers for testing // Export helpers for testing
getTodayUTC, getTodayUTC,
getYesterdayUnitCodes, getYesterdayUnitCodes,