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 101 additions and 20 deletions
Showing only changes of commit 2e3ef0580c - Show all commits

View File

@ -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() {

View File

@ -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() {

View File

@ -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"', () => {

View File

@ -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'
}

View File

@ -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,51 @@ describe('sanitizeError', () => {
});
});
});
// ============================================================
// Test: createScraperIndexes
// ============================================================
describe('createScraperIndexes', () => {
it('should create the status_startedAt compound index', async () => {
await createScraperIndexes(db);
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);
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 be idempotent (no error on second call)', async () => {
await createScraperIndexes(db);
// Calling again should not throw
await expect(createScraperIndexes(db)).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)).rejects.toThrow('Connection lost');
});
});

View File

@ -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

View File

@ -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'
}

View File

@ -730,6 +730,36 @@ 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
* @returns {Promise<void>}
*/
async function createScraperIndexes(db) {
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' }
);
console.log('Scraper indexes created successfully');
}
// ============================================================
// Main Orchestration Function
// ============================================================
@ -871,6 +901,7 @@ module.exports = {
markStaleUnits,
updateDailySummary,
recordScraperRun,
createScraperIndexes,
// Export helpers for testing
getTodayUTC,
getYesterdayUnitCodes,