SCRAPE-11: Implement runScrape() main orchestration function
Some checks failed
CI/CD Pipeline - Apartment API / Scan Dependencies (pull_request) Successful in 13s
CI/CD Pipeline - Apartment API / Run Linting (pull_request) Successful in 9m37s
CI/CD Pipeline - Apartment API / Run Tests (pull_request) Successful in 10m0s
CI/CD Pipeline - Apartment API / Send Webhook Notification (pull_request) Failing after 2s
CI/CD Pipeline - Apartment API / Build & Push Image (pull_request) Has been skipped
CI/CD Pipeline - Apartment API / Deploy to Production (pull_request) Has been skipped

Add the top-level runScrape() function that coordinates the full scraper
pipeline: fetch HTML (with retry), parse units, convert data types, upsert
units, insert prices, mark stale units, and update daily summary.

Features:
- dryRun mode skips all database writes while still parsing/validating
- htmlContent parameter allows injecting HTML directly (bypasses fetch)
- Calculates newUnitsCount and rentedUnitsCount by diffing against prior state
- Records every run to scraper_runs history (success or failure)
- Structured logging with jobId correlation throughout the pipeline
- Graceful error handling at each pipeline stage

Includes 15 tests covering full workflow, dry run, error handling,
trigger types, new/rented unit calculation, and empty HTML edge case.
This commit is contained in:
2026-02-06 01:35:15 -07:00
parent d1f717891a
commit 5e1f1f8041
2 changed files with 627 additions and 2 deletions

View File

@ -0,0 +1,457 @@
/**
* Tests for runScrape() orchestration function
*
* Covers:
* - Full workflow completes with mocked dependencies
* - Returns correct result structure with jobId, status, metrics
* - Handles dryRun option (skip DB writes)
* - Handles htmlContent option (use provided HTML)
* - Records scraper run to history on success
* - Records scraper run to history on failure
* - Catches and logs errors from fetchPage
* - Catches and logs errors from database operations
* - Calculates newUnitsCount and rentedUnitsCount correctly
*/
const { MongoClient } = require('mongodb');
const { MongoMemoryServer } = require('mongodb-memory-server');
let mongoServer;
let client;
let db;
// We will require runScrape and recordScraperRun after implementation
let runScrape;
let recordScraperRun;
// Store original module references so we can mock individual functions
let scraperService;
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
scraperService = require('../../services/scraperService');
({ runScrape, recordScraperRun } = scraperService);
});
afterAll(async () => {
if (client) await client.close();
if (mongoServer) await mongoServer.stop();
});
beforeEach(async () => {
// Clean all relevant collections before each test
const collections = await db.listCollections().toArray();
for (const col of collections) {
await db.collection(col.name).deleteMany({});
}
// Reset all mocks
jest.restoreAllMocks();
});
// ============================================================
// Helper: Create sample HTML with units
// ============================================================
function createSampleHtml(unitCodes) {
const articles = unitCodes.map(code => `
<article
data-spaces-id="100"
data-spaces-unit="${code}"
data-spaces-unit-id="200"
data-spaces-unit-floor="5"
data-spaces-sort-area="750"
data-spaces-sort-bed="1"
data-spaces-sort-bath="1"
data-spaces-sort-price="1500"
data-spaces-available="true"
data-spaces-unavailable="false"
data-spaces-soonest="Now"
data-spaces-sort-date="1700000000"
data-spaces-plan-id="10"
data-spaces-sort-plan-name="Studio"
data-spaces-obj="unit"
data-spaces-community="TestCommunity"
data-spaces-asset="1"
data-spaces-href="/unit/${code}"
data-spaces-inventory-href="/inventory/${code}"
>
<img src="http://example.com/${code}.jpg" />
</article>
`).join('\n');
return `
<html><body>
<section class="spaces__tab-unit">
${articles}
</section>
</body></html>
`;
}
// ============================================================
// Test: recordScraperRun
// ============================================================
describe('recordScraperRun', () => {
it('should insert a run record into scraper_runs collection', async () => {
const runData = {
jobId: 'test-job-001',
trigger: 'manual',
status: 'success',
startedAt: new Date().toISOString(),
completedAt: new Date().toISOString(),
duration: 1234,
unitsProcessed: 10,
pricesInserted: 8,
errors: []
};
const mockLogger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() };
const result = await recordScraperRun(db, runData, mockLogger);
expect(result).toBeTruthy();
expect(result.insertedId).toBeTruthy();
// Verify it was inserted
const saved = await db.collection('scraper_runs').findOne({ jobId: 'test-job-001' });
expect(saved).toBeTruthy();
expect(saved.status).toBe('success');
expect(saved.recordedAt).toBeInstanceOf(Date);
});
it('should not throw when insert fails', async () => {
// Pass null db to cause an error
const runData = { jobId: 'test-fail', status: 'success' };
// Should not throw
const mockLogger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() };
const result = await recordScraperRun(null, runData, mockLogger);
expect(result).toBeNull();
expect(mockLogger.error).toHaveBeenCalled();
});
});
// ============================================================
// Test: runScrape - Full workflow
// ============================================================
describe('runScrape', () => {
it('should complete full workflow with mocked dependencies', async () => {
const html = createSampleHtml(['UNIT-A', 'UNIT-B']);
const result = await runScrape(db, {
trigger: 'manual',
jobId: 'test-full-workflow',
htmlContent: html
});
expect(result).toBeTruthy();
expect(result.jobId).toBe('test-full-workflow');
expect(result.status).toBe('success');
expect(result.trigger).toBe('manual');
expect(result.unitsProcessed).toBe(2);
expect(result.errors).toEqual([]);
});
it('should return correct result structure with jobId, status, metrics', async () => {
const html = createSampleHtml(['UNIT-X']);
const result = await runScrape(db, {
jobId: 'test-structure',
htmlContent: html
});
// Verify all expected fields exist
expect(result).toHaveProperty('jobId', 'test-structure');
expect(result).toHaveProperty('trigger', 'manual');
expect(result).toHaveProperty('dryRun', false);
expect(result).toHaveProperty('status', 'success');
expect(result).toHaveProperty('startedAt');
expect(result).toHaveProperty('completedAt');
expect(result).toHaveProperty('duration');
expect(result).toHaveProperty('unitsProcessed');
expect(result).toHaveProperty('pricesInserted');
expect(result).toHaveProperty('newUnitsCount');
expect(result).toHaveProperty('rentedUnitsCount');
expect(result).toHaveProperty('staleUnitsCount');
expect(result).toHaveProperty('errors');
// Verify types
expect(typeof result.duration).toBe('number');
expect(result.duration).toBeGreaterThanOrEqual(0);
expect(typeof result.startedAt).toBe('string');
expect(typeof result.completedAt).toBe('string');
expect(Array.isArray(result.errors)).toBe(true);
});
it('should generate a jobId when not provided', async () => {
const html = createSampleHtml(['UNIT-GEN']);
const result = await runScrape(db, {
htmlContent: html
});
expect(result.jobId).toBeTruthy();
expect(typeof result.jobId).toBe('string');
// UUID format: xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx
expect(result.jobId).toMatch(/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i);
});
// ============================================================
// Test: dryRun option
// ============================================================
it('should skip DB writes when dryRun is true', async () => {
const html = createSampleHtml(['DRY-A', 'DRY-B']);
const result = await runScrape(db, {
jobId: 'test-dry-run',
htmlContent: html,
dryRun: true
});
expect(result.status).toBe('success');
expect(result.dryRun).toBe(true);
expect(result.unitsProcessed).toBe(2);
expect(result.pricesInserted).toBe(0);
expect(result.staleUnitsCount).toBe(0);
// Verify no units were written to the units collection
const unitsCount = await db.collection('units_migration_test').countDocuments();
expect(unitsCount).toBe(0);
// Verify no prices were written
const pricesCount = await db.collection('unit_prices_migration_test').countDocuments();
expect(pricesCount).toBe(0);
});
// ============================================================
// Test: htmlContent option
// ============================================================
it('should use provided HTML instead of fetching when htmlContent is given', async () => {
const html = createSampleHtml(['HTML-1', 'HTML-2', 'HTML-3']);
// Mock axios.get to track if HTTP fetch is attempted
const axios = require('axios');
const axiosSpy = jest.spyOn(axios, 'get');
const result = await runScrape(db, {
jobId: 'test-html-content',
htmlContent: html
});
expect(result.status).toBe('success');
expect(result.unitsProcessed).toBe(3);
// axios.get should NOT have been called since we provided htmlContent
expect(axiosSpy).not.toHaveBeenCalled();
});
// ============================================================
// Test: Records scraper run on success
// ============================================================
it('should record scraper run to history on success', async () => {
const html = createSampleHtml(['REC-A']);
await runScrape(db, {
jobId: 'test-record-success',
htmlContent: html
});
const runRecord = await db.collection('scraper_runs').findOne({ jobId: 'test-record-success' });
expect(runRecord).toBeTruthy();
expect(runRecord.status).toBe('success');
expect(runRecord.recordedAt).toBeInstanceOf(Date);
expect(runRecord.unitsProcessed).toBe(1);
});
// ============================================================
// Test: Records scraper run on failure
// ============================================================
it('should record scraper run to history on failure', async () => {
const html = createSampleHtml(['FAIL-A']);
// Create a db proxy that throws on bulkWrite (used by upsertUnits)
// but allows other operations (like scraper_runs insertOne) to pass through
const faultyDb = {
collection: (name) => {
const realCollection = db.collection(name);
if (name === 'units_migration_test') {
return new Proxy(realCollection, {
get(target, prop) {
if (prop === 'bulkWrite') {
return async () => {
throw new Error('Database connection lost');
};
}
const value = target[prop];
if (typeof value === 'function') {
return value.bind(target);
}
return value;
}
});
}
return realCollection;
}
};
const result = await runScrape(faultyDb, {
jobId: 'test-record-failure',
htmlContent: html
});
expect(result.status).toBe('failed');
expect(result.errors).toContain('Database connection lost');
// The run should still be recorded (via the real db passed through for scraper_runs)
const runRecord = await db.collection('scraper_runs').findOne({ jobId: 'test-record-failure' });
expect(runRecord).toBeTruthy();
expect(runRecord.status).toBe('failed');
expect(runRecord.errors).toContain('Database connection lost');
});
// ============================================================
// Test: Catches errors from fetchPage
// ============================================================
it('should catch and handle errors from fetchPage', async () => {
// Mock axios.get to throw a network error (fetchPage uses axios internally)
const axios = require('axios');
jest.spyOn(axios, 'get').mockRejectedValue(
new Error('Network timeout')
);
const result = await runScrape(db, {
jobId: 'test-fetch-error'
// No htmlContent, so it will call fetchPage which uses axios
});
expect(result.status).toBe('failed');
expect(result.errors).toContain('Network timeout');
expect(result.unitsProcessed).toBe(0);
});
// ============================================================
// Test: Catches errors from database operations
// ============================================================
it('should catch and handle errors from database operations', async () => {
const html = createSampleHtml(['DB-ERR']);
// Create a db proxy that throws on the prices collection bulkWrite
// Use Proxy to properly delegate all methods to the real collection
const faultyDb = {
collection: (name) => {
const realCollection = db.collection(name);
if (name === 'unit_prices_migration_test') {
return new Proxy(realCollection, {
get(target, prop) {
if (prop === 'bulkWrite') {
return async () => {
throw new Error('Write concern timeout');
};
}
const value = target[prop];
if (typeof value === 'function') {
return value.bind(target);
}
return value;
}
});
}
return realCollection;
}
};
const result = await runScrape(faultyDb, {
jobId: 'test-db-error',
htmlContent: html
});
expect(result.status).toBe('failed');
expect(result.errors).toContain('Write concern timeout');
// Run should still be recorded despite error
const runRecord = await db.collection('scraper_runs').findOne({ jobId: 'test-db-error' });
expect(runRecord).toBeTruthy();
expect(runRecord.status).toBe('failed');
});
// ============================================================
// Test: Calculates newUnitsCount and rentedUnitsCount correctly
// ============================================================
it('should calculate newUnitsCount and rentedUnitsCount correctly', async () => {
// First, simulate yesterday's data by inserting price records for yesterday
const today = new Date();
const yesterday = new Date(today);
yesterday.setUTCDate(yesterday.getUTCDate() - 1);
const yesterdayStr = yesterday.toISOString().split('T')[0];
// Yesterday had units: OLD-A, OLD-B, OLD-C
await db.collection('unit_prices_migration_test').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 }
]);
// Today's scrape has: OLD-A, OLD-B, NEW-D (OLD-C is gone / rented)
const html = createSampleHtml(['OLD-A', 'OLD-B', 'NEW-D']);
const result = await runScrape(db, {
jobId: 'test-calc-changes',
htmlContent: html
});
expect(result.status).toBe('success');
// NEW-D is new (not in yesterday's data)
expect(result.newUnitsCount).toBe(1);
// OLD-C was in yesterday's data but not in today's scrape
expect(result.rentedUnitsCount).toBe(1);
expect(result.unitsProcessed).toBe(3);
});
// ============================================================
// Test: Default trigger is 'manual'
// ============================================================
it('should default trigger to manual', async () => {
const html = createSampleHtml(['DEF-A']);
const result = await runScrape(db, {
jobId: 'test-default-trigger',
htmlContent: html
});
expect(result.trigger).toBe('manual');
});
// ============================================================
// Test: Supports scheduled trigger
// ============================================================
it('should support scheduled trigger', async () => {
const html = createSampleHtml(['SCHED-A']);
const result = await runScrape(db, {
jobId: 'test-scheduled-trigger',
trigger: 'scheduled',
htmlContent: html
});
expect(result.trigger).toBe('scheduled');
});
// ============================================================
// Test: Handles empty HTML (no units found)
// ============================================================
it('should handle HTML with no units gracefully', async () => {
const emptyHtml = '<html><body><section class="spaces__tab-unit"></section></body></html>';
const result = await runScrape(db, {
jobId: 'test-empty-html',
htmlContent: emptyHtml
});
expect(result.status).toBe('success');
expect(result.unitsProcessed).toBe(0);
expect(result.errors).toContain('No units found in HTML');
});
});

View File

@ -7,7 +7,9 @@
const axios = require('axios'); const axios = require('axios');
const cheerio = require('cheerio'); const cheerio = require('cheerio');
const crypto = require('crypto');
const config = require('../config/scraper'); const config = require('../config/scraper');
const { createLogger } = require('./scraperLogger');
/** /**
* Sleep utility for retry delays * Sleep utility for retry delays
@ -619,6 +621,40 @@ async function updateDailySummary(db, summaryData, logger) {
} }
} }
// ============================================================
// Additional Helpers for runScrape Orchestration
// ============================================================
/**
* Get today's date in YYYY-MM-DD format (UTC).
* @returns {string} Today's date string
*/
function getTodayUTC() {
return new Date().toISOString().split('T')[0];
}
/**
* Get the set of unit codes that had price records yesterday.
* Used to calculate new and rented units by comparison.
* @param {Db} db - MongoDB database instance
* @param {string} today - Today's date in YYYY-MM-DD format
* @returns {Promise<Set<string>>} Set of unit codes from yesterday
*/
async function getYesterdayUnitCodes(db, today) {
const yesterday = getYesterday(today);
const collection = db.collection(config.COLLECTIONS.PRICES);
const yesterdayRecords = await collection
.find({ date_checked: yesterday }, { projection: { unit_code: 1 } })
.toArray();
return new Set(yesterdayRecords.map(r => r.unit_code));
}
// ============================================================
// Scraper Run History
// ============================================================
/** /**
* Record scraper run to history collection. * Record scraper run to history collection.
* Called in the finally block of runScrape() to persist run metadata. * Called in the finally block of runScrape() to persist run metadata.
@ -630,9 +666,8 @@ async function updateDailySummary(db, summaryData, logger) {
* @returns {Promise<Object|null>} Insert result, or null on failure * @returns {Promise<Object|null>} Insert result, or null on failure
*/ */
async function recordScraperRun(db, runData, logger) { async function recordScraperRun(db, runData, logger) {
const collection = db.collection(config.COLLECTIONS.SCRAPER_RUNS);
try { try {
const collection = db.collection(config.COLLECTIONS.SCRAPER_RUNS);
const result = await collection.insertOne({ const result = await collection.insertOne({
...runData, ...runData,
recordedAt: new Date() recordedAt: new Date()
@ -647,7 +682,138 @@ async function recordScraperRun(db, runData, logger) {
} }
} }
// ============================================================
// Main Orchestration Function
// ============================================================
/**
* Execute a complete scrape operation.
* Orchestrates the full workflow: fetch -> parse -> convert -> DB ops.
*
* @param {Db} db - MongoDB database instance
* @param {Object} options - Scrape options
* @param {string} [options.trigger='manual'] - Trigger type ('scheduled' | 'manual')
* @param {string} [options.jobId] - Optional job ID (generated if not provided)
* @param {boolean} [options.dryRun=false] - Skip database writes for safe testing
* @param {string} [options.htmlContent] - Use provided HTML instead of fetching
* @returns {Promise<Object>} Scrape result with status and metrics
*/
async function runScrape(db, options = {}) {
const jobId = options.jobId || crypto.randomUUID();
const trigger = options.trigger || 'manual';
const dryRun = options.dryRun || false;
const htmlContent = options.htmlContent || null;
const logger = createLogger(jobId);
const startTime = Date.now();
let result = {
jobId,
trigger,
dryRun,
status: 'running',
startedAt: new Date().toISOString(),
completedAt: null,
duration: null,
unitsProcessed: 0,
pricesInserted: 0,
newUnitsCount: 0,
rentedUnitsCount: 0,
staleUnitsCount: 0,
errors: []
};
try {
logger.info('Scrape started', { trigger, dryRun, usingProvidedHtml: !!htmlContent });
// Step 1: Fetch HTML (or use provided content for testing)
const html = htmlContent || await fetchPage(config.TARGET_URL, logger);
// Step 2: Parse units
const rawUnits = parseUnits(html, logger);
if (rawUnits.length === 0) {
logger.warn('No units found in HTML - possible structure change');
result.errors.push('No units found in HTML');
}
// Step 3: Convert data types
const units = rawUnits.map(unit => convertDataTypes(unit));
result.unitsProcessed = units.length;
// Step 4: Database operations
const today = getTodayUTC();
// Get yesterday's unit codes for comparison
const yesterdayUnits = await getYesterdayUnitCodes(db, today);
// Determine new and rented units
const currentUnitCodes = new Set(units.map(u => u.unit_code));
const newUnits = units.filter(u => !yesterdayUnits.has(u.unit_code));
const rentedUnits = [...yesterdayUnits].filter(code => !currentUnitCodes.has(code));
result.newUnitsCount = newUnits.length;
result.rentedUnitsCount = rentedUnits.length;
// Database operations (skip if dryRun)
if (dryRun) {
logger.info('Dry run mode - skipping database writes', {
wouldUpsert: units.length,
wouldInsertPrices: units.filter(u => u.price !== null).length
});
result.pricesInserted = 0;
result.staleUnitsCount = 0;
} else {
// Upsert units
await upsertUnits(db, units, logger);
// Insert prices
const pricesResult = await insertPrices(db, units, today, logger);
result.pricesInserted = pricesResult.insertedCount;
// Mark stale units
const staleResult = await markStaleUnits(db, currentUnitCodes, today, logger);
result.staleUnitsCount = staleResult.modifiedCount;
// Update daily summary
await updateDailySummary(db, {
date: today,
newUnits: newUnits.map(u => u.unit_code),
rentedUnits,
staleUnitsCount: result.staleUnitsCount,
totalAvailable: units.length
}, logger);
}
result.status = 'success';
} catch (error) {
logger.error('Scrape failed', {
errorType: error.name,
errorMessage: error.message
});
result.status = 'failed';
result.errors.push(error.message);
} finally {
result.completedAt = new Date().toISOString();
result.duration = Date.now() - startTime;
// Record run to history (always runs, even on failure)
await recordScraperRun(db, result, logger);
logger.info('Scrape completed', {
status: result.status,
duration: result.duration,
unitsProcessed: result.unitsProcessed,
pricesInserted: result.pricesInserted
});
}
return result;
}
module.exports = { module.exports = {
runScrape,
fetchPage, fetchPage,
parseUnits, parseUnits,
convertDataTypes, convertDataTypes,
@ -657,6 +823,8 @@ module.exports = {
updateDailySummary, updateDailySummary,
recordScraperRun, recordScraperRun,
// Export helpers for testing // Export helpers for testing
getTodayUTC,
getYesterdayUnitCodes,
getYesterday, getYesterday,
parseInteger, parseInteger,
parsePositiveInteger, parsePositiveInteger,