Merge dev: Node.js scraper migration + CI fix (#32)
All checks were successful
CI/CD Pipeline - Apartment API / Scan Dependencies (push) Successful in 13s
CI/CD Pipeline - Apartment API / Lint & Test (push) Successful in 44s
CI/CD Pipeline - Apartment API / Send Webhook Notification (push) Successful in 2s
CI/CD Pipeline - Apartment API / Build & Push Image (push) Successful in 1m41s
CI/CD Pipeline - Apartment API / Deploy to Production (push) Successful in 14s
All checks were successful
CI/CD Pipeline - Apartment API / Scan Dependencies (push) Successful in 13s
CI/CD Pipeline - Apartment API / Lint & Test (push) Successful in 44s
CI/CD Pipeline - Apartment API / Send Webhook Notification (push) Successful in 2s
CI/CD Pipeline - Apartment API / Build & Push Image (push) Successful in 1m41s
CI/CD Pipeline - Apartment API / Deploy to Production (push) Successful in 14s
Co-authored-by: Stephen Minakian <stephenminakian@gmail.com> Co-committed-by: Stephen Minakian <stephenminakian@gmail.com>
This commit is contained in:
164
services/scraperLogger.js
Normal file
164
services/scraperLogger.js
Normal file
@ -0,0 +1,164 @@
|
||||
/**
|
||||
* Scraper-specific structured JSON logger
|
||||
* Produces one JSON object per line to stdout
|
||||
*/
|
||||
|
||||
const LEVELS = {
|
||||
info: 'info',
|
||||
warn: 'warn',
|
||||
error: 'error'
|
||||
};
|
||||
|
||||
/**
|
||||
* Create a logger instance scoped to a job ID
|
||||
* @param {string} jobId - Job identifier for correlation
|
||||
* @returns {Object} Logger object with info, warn, error methods
|
||||
*/
|
||||
function createLogger(jobId) {
|
||||
/**
|
||||
* Internal log function
|
||||
* @param {string} level - Log level
|
||||
* @param {string} message - Log message
|
||||
* @param {Object} context - Additional context data
|
||||
*/
|
||||
const log = (level, message, context = {}) => {
|
||||
const entry = {
|
||||
timestamp: new Date().toISOString(),
|
||||
level,
|
||||
message: truncateMessage(message),
|
||||
jobId,
|
||||
context: sanitizeContext(context)
|
||||
};
|
||||
|
||||
// Output as single-line JSON
|
||||
console.log(JSON.stringify(entry));
|
||||
};
|
||||
|
||||
return {
|
||||
info: (message, context) => log(LEVELS.info, message, context),
|
||||
warn: (message, context) => log(LEVELS.warn, message, context),
|
||||
error: (message, context) => log(LEVELS.error, message, context)
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Truncate message to prevent log bloat
|
||||
* @param {string} message - Message to truncate
|
||||
* @param {number} maxLength - Maximum length (default 1000)
|
||||
* @returns {string} Truncated message
|
||||
*/
|
||||
function truncateMessage(message, maxLength = 1000) {
|
||||
if (message === null || message === undefined) {
|
||||
return '';
|
||||
}
|
||||
const str = String(message);
|
||||
if (str.length <= maxLength) {
|
||||
return str;
|
||||
}
|
||||
return str.substring(0, maxLength) + '... [truncated]';
|
||||
}
|
||||
|
||||
/**
|
||||
* Recursively process an object to handle Buffers, circular references, and sensitive data
|
||||
* @param {*} obj - Object to process
|
||||
* @param {WeakSet} seen - Set of seen objects for circular reference detection
|
||||
* @returns {*} Processed value
|
||||
*/
|
||||
function processValue(obj, seen = new WeakSet()) {
|
||||
// Handle null/undefined
|
||||
if (obj === null) {
|
||||
return null;
|
||||
}
|
||||
if (obj === undefined) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Handle Buffer BEFORE checking for object (Buffer is an object)
|
||||
if (Buffer.isBuffer(obj)) {
|
||||
return `[Buffer: ${obj.length} bytes]`;
|
||||
}
|
||||
|
||||
// Handle strings - check for MongoDB connection strings
|
||||
if (typeof obj === 'string') {
|
||||
if (/mongodb(\+srv)?:\/\//.test(obj)) {
|
||||
return redactConnectionString(obj);
|
||||
}
|
||||
return obj;
|
||||
}
|
||||
|
||||
// Handle primitives
|
||||
if (typeof obj !== 'object') {
|
||||
return obj;
|
||||
}
|
||||
|
||||
// Handle circular references
|
||||
if (seen.has(obj)) {
|
||||
return '[Circular]';
|
||||
}
|
||||
seen.add(obj);
|
||||
|
||||
// Handle arrays
|
||||
if (Array.isArray(obj)) {
|
||||
return obj.map(item => processValue(item, seen));
|
||||
}
|
||||
|
||||
// Handle plain objects
|
||||
const result = {};
|
||||
const sensitiveKeys = ['password', 'secret', 'token', 'apikey', 'authorization'];
|
||||
|
||||
for (const key of Object.keys(obj)) {
|
||||
// Check for sensitive keys
|
||||
if (sensitiveKeys.some(k => key.toLowerCase().includes(k))) {
|
||||
result[key] = '[REDACTED]';
|
||||
} else {
|
||||
result[key] = processValue(obj[key], seen);
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sanitize context object for safe logging
|
||||
* - Remove circular references
|
||||
* - Redact sensitive data
|
||||
* - Handle special types (Buffer, undefined)
|
||||
* @param {Object} context - Context object
|
||||
* @returns {Object} Sanitized context
|
||||
*/
|
||||
function sanitizeContext(context) {
|
||||
if (!context || typeof context !== 'object') {
|
||||
return {};
|
||||
}
|
||||
|
||||
try {
|
||||
return processValue(context);
|
||||
} catch (error) {
|
||||
// If sanitization fails, return empty context
|
||||
return { sanitizationError: 'Failed to sanitize context' };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Redact credentials from MongoDB connection string
|
||||
* @param {string} uri - Connection string
|
||||
* @returns {string} Redacted string
|
||||
*/
|
||||
function redactConnectionString(uri) {
|
||||
try {
|
||||
// Match mongodb://user:pass@host or mongodb+srv://user:pass@host
|
||||
return uri.replace(
|
||||
/mongodb(\+srv)?:\/\/([^:]+):([^@]+)@/,
|
||||
'mongodb$1://[user]:[REDACTED]@'
|
||||
);
|
||||
} catch {
|
||||
return '[REDACTED CONNECTION STRING]';
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
createLogger,
|
||||
LEVELS,
|
||||
// Export internal functions for testing
|
||||
redactConnectionString
|
||||
};
|
||||
920
services/scraperService.js
Normal file
920
services/scraperService.js
Normal file
@ -0,0 +1,920 @@
|
||||
/**
|
||||
* Scraper Service
|
||||
*
|
||||
* Core scraper logic including HTTP fetching, HTML parsing,
|
||||
* data transformation, and database operations.
|
||||
*/
|
||||
|
||||
const axios = require('axios');
|
||||
const cheerio = require('cheerio');
|
||||
const crypto = require('crypto');
|
||||
const config = require('../config/scraper');
|
||||
const { createLogger } = require('./scraperLogger');
|
||||
|
||||
/**
|
||||
* Sleep utility for retry delays
|
||||
* @param {number} ms - Milliseconds to sleep
|
||||
* @returns {Promise<void>}
|
||||
*/
|
||||
function sleep(ms) {
|
||||
return new Promise(resolve => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
/**
|
||||
* Determine if an error is retryable
|
||||
* @param {Error} error - Axios error
|
||||
* @returns {boolean} True if should retry
|
||||
*/
|
||||
function isRetryableError(error) {
|
||||
// Network errors (timeout, DNS, connection) don't have a response property
|
||||
if (!error.response) {
|
||||
return true;
|
||||
}
|
||||
|
||||
const status = error.response.status;
|
||||
|
||||
// 5xx server errors are retryable
|
||||
if (status >= 500) {
|
||||
return true;
|
||||
}
|
||||
|
||||
// 429 Too Many Requests - do not retry immediately
|
||||
if (status === 429) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// 4xx client errors - do not retry
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Fetch page HTML with retry logic
|
||||
* @param {string} url - Target URL
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<string>} HTML content
|
||||
* @throws {Error} After all retries exhausted
|
||||
*/
|
||||
async function fetchPage(url, logger) {
|
||||
const { maxRetries, baseDelay, timeout } = config.RETRY_CONFIG;
|
||||
let lastError;
|
||||
|
||||
for (let attempt = 1; attempt <= maxRetries + 1; attempt++) {
|
||||
try {
|
||||
logger.info('Fetching page', { url, attempt });
|
||||
|
||||
const response = await axios.get(url, {
|
||||
timeout,
|
||||
headers: {
|
||||
'User-Agent': config.USER_AGENT
|
||||
},
|
||||
maxRedirects: 5,
|
||||
validateStatus: (status) => status < 400 // Accept 2xx and 3xx
|
||||
});
|
||||
|
||||
logger.info('Page fetched successfully', {
|
||||
status: response.status,
|
||||
contentLength: response.data.length
|
||||
});
|
||||
|
||||
return response.data;
|
||||
|
||||
} catch (error) {
|
||||
lastError = error;
|
||||
|
||||
// Determine if error is retryable
|
||||
const isRetryable = isRetryableError(error);
|
||||
|
||||
logger.warn('Fetch attempt failed', {
|
||||
attempt,
|
||||
errorType: error.name,
|
||||
errorMessage: error.message,
|
||||
statusCode: error.response?.status,
|
||||
isRetryable
|
||||
});
|
||||
|
||||
// Don't retry non-retryable errors (4xx)
|
||||
if (!isRetryable) {
|
||||
throw error;
|
||||
}
|
||||
|
||||
// Don't wait after last attempt
|
||||
if (attempt <= maxRetries) {
|
||||
const delay = baseDelay * Math.pow(2, attempt - 1); // Exponential backoff
|
||||
logger.info('Waiting before retry', { delay });
|
||||
await sleep(delay);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
throw lastError;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse unit data from HTML using cheerio
|
||||
* @param {string} html - HTML content
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Array<Object>} Array of raw unit objects (string values, no type conversion)
|
||||
*/
|
||||
function parseUnits(html, logger) {
|
||||
const $ = cheerio.load(html);
|
||||
const units = [];
|
||||
|
||||
// Find the main container
|
||||
const container = $('section.spaces__tab-unit');
|
||||
|
||||
if (container.length === 0) {
|
||||
logger.error('Container section.spaces__tab-unit not found - possible structure change', {});
|
||||
return [];
|
||||
}
|
||||
|
||||
// Extract each article element
|
||||
container.find('article').each((index, article) => {
|
||||
const $article = $(article);
|
||||
|
||||
const unit = {
|
||||
// Core identifiers
|
||||
id: $article.attr('data-spaces-id'),
|
||||
unit_code: $article.attr('data-spaces-unit'),
|
||||
unit_id: $article.attr('data-spaces-unit-id'),
|
||||
|
||||
// Physical attributes
|
||||
floor: $article.attr('data-spaces-unit-floor'),
|
||||
area: $article.attr('data-spaces-sort-area'),
|
||||
bed_count: $article.attr('data-spaces-sort-bed'),
|
||||
bath_count: $article.attr('data-spaces-sort-bath'),
|
||||
|
||||
// Pricing
|
||||
price: $article.attr('data-spaces-sort-price'),
|
||||
|
||||
// Availability
|
||||
available: $article.attr('data-spaces-available'),
|
||||
unavailable: $article.attr('data-spaces-unavailable'),
|
||||
soonest: $article.attr('data-spaces-soonest'),
|
||||
date_available: $article.attr('data-spaces-sort-date'),
|
||||
|
||||
// Plan information
|
||||
plan_id: $article.attr('data-spaces-plan-id'),
|
||||
plan_name: $article.attr('data-spaces-sort-plan-name'),
|
||||
|
||||
// Property information
|
||||
obj_type: $article.attr('data-spaces-obj'),
|
||||
community: $article.attr('data-spaces-community'),
|
||||
asset: $article.attr('data-spaces-asset'),
|
||||
|
||||
// URLs
|
||||
href: $article.attr('data-spaces-href'),
|
||||
inventory_href: $article.attr('data-spaces-inventory-href'),
|
||||
|
||||
// Specials
|
||||
specials_content: $article.attr('data-spaces-specials-content')
|
||||
};
|
||||
|
||||
// Extract image URL from nested element if present
|
||||
const imgElement = $article.find('img').first();
|
||||
if (imgElement.length > 0) {
|
||||
unit.image_url = imgElement.attr('src') || imgElement.attr('data-src');
|
||||
}
|
||||
|
||||
units.push(unit);
|
||||
});
|
||||
|
||||
logger.info('Units parsed', { count: units.length });
|
||||
|
||||
// Deduplicate by unit_code (in case of duplicate articles)
|
||||
const seen = new Set();
|
||||
const deduplicated = units.filter(unit => {
|
||||
if (!unit.unit_code || seen.has(unit.unit_code)) {
|
||||
return false;
|
||||
}
|
||||
seen.add(unit.unit_code);
|
||||
return true;
|
||||
});
|
||||
|
||||
if (deduplicated.length < units.length) {
|
||||
logger.warn('Duplicate units removed', {
|
||||
original: units.length,
|
||||
deduplicated: deduplicated.length
|
||||
});
|
||||
}
|
||||
|
||||
return deduplicated;
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Data Type Conversion Helpers
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Parse string to integer, return null for invalid
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {number|null} Parsed integer or null
|
||||
*/
|
||||
function parseInteger(value) {
|
||||
if (value === null || value === undefined || value === '') {
|
||||
return null;
|
||||
}
|
||||
const parsed = parseInt(value, 10);
|
||||
return Number.isNaN(parsed) || !Number.isFinite(parsed) ? null : parsed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse string to positive integer, return null for 0 or invalid.
|
||||
* Used for fields like area where 0 means "not available".
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {number|null} Parsed positive integer or null
|
||||
*/
|
||||
function parsePositiveInteger(value) {
|
||||
const parsed = parseInteger(value);
|
||||
return parsed === 0 ? null : parsed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse value that could be integer or string identifier.
|
||||
* Returns integer if cleanly parseable, otherwise trimmed string.
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {number|string|null} Parsed integer, trimmed string, or null
|
||||
*/
|
||||
function parseIntegerOrString(value) {
|
||||
if (value === null || value === undefined || value === '') {
|
||||
return null;
|
||||
}
|
||||
const parsed = parseInt(value, 10);
|
||||
// If it parses cleanly to an integer, return integer
|
||||
if (!Number.isNaN(parsed) && String(parsed) === String(value).trim()) {
|
||||
return parsed;
|
||||
}
|
||||
// Otherwise return as trimmed string
|
||||
return String(value).trim();
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse price, stripping non-numeric characters.
|
||||
* Handles "$1,234" format and "Call for pricing" text.
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {number|null} Parsed price or null
|
||||
*/
|
||||
function parsePrice(value) {
|
||||
if (value === null || value === undefined || value === '') {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Check for "Call for pricing" or similar text
|
||||
if (typeof value === 'string' && /call|contact|inquire/i.test(value)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Strip non-numeric characters except decimal point
|
||||
const cleaned = String(value).replace(/[^0-9.]/g, '');
|
||||
const parsed = parseInt(cleaned, 10);
|
||||
|
||||
if (Number.isNaN(parsed) || !Number.isFinite(parsed)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Negative prices are invalid
|
||||
if (parsed < 0) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return parsed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse boolean from string.
|
||||
* Handles "true"/"false"/"1"/"0" and actual boolean values.
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {boolean|null} Parsed boolean or null
|
||||
*/
|
||||
function parseBoolean(value) {
|
||||
if (value === null || value === undefined || value === '') {
|
||||
return null;
|
||||
}
|
||||
if (typeof value === 'boolean') {
|
||||
return value;
|
||||
}
|
||||
const str = String(value).toLowerCase().trim();
|
||||
if (str === 'true' || str === '1') return true;
|
||||
if (str === 'false' || str === '0') return false;
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse float value.
|
||||
* Named parseFloatValue to avoid shadowing the global parseFloat.
|
||||
* @param {*} value - Value to parse
|
||||
* @returns {number|null} Parsed float or null
|
||||
*/
|
||||
function parseFloatValue(value) {
|
||||
if (value === null || value === undefined || value === '') {
|
||||
return null;
|
||||
}
|
||||
const parsed = Number.parseFloat(value);
|
||||
return Number.isNaN(parsed) || !Number.isFinite(parsed) ? null : parsed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Trim string, return null for empty.
|
||||
* Converts non-string values to string before trimming.
|
||||
* @param {*} value - Value to trim
|
||||
* @returns {string|null} Trimmed string or null
|
||||
*/
|
||||
function trimString(value) {
|
||||
if (value === null || value === undefined) {
|
||||
return null;
|
||||
}
|
||||
const trimmed = String(value).trim();
|
||||
return trimmed === '' ? null : trimmed;
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Main Data Type Conversion Function
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Convert unit data types from strings to proper types.
|
||||
* Applies the appropriate parser to each field based on its expected type.
|
||||
* @param {Object} unit - Raw unit object with string values
|
||||
* @returns {Object} Unit object with converted types
|
||||
*/
|
||||
function convertDataTypes(unit) {
|
||||
return {
|
||||
// Integer fields
|
||||
id: parseInteger(unit.id),
|
||||
unit_id: parseIntegerOrString(unit.unit_id),
|
||||
floor: parseInteger(unit.floor),
|
||||
area: parsePositiveInteger(unit.area),
|
||||
bed_count: parseInteger(unit.bed_count),
|
||||
plan_id: parseInteger(unit.plan_id),
|
||||
asset: parseInteger(unit.asset),
|
||||
date_available: parseInteger(unit.date_available),
|
||||
|
||||
// Float fields
|
||||
bath_count: parseFloatValue(unit.bath_count),
|
||||
|
||||
// Price field (special handling)
|
||||
price: parsePrice(unit.price),
|
||||
|
||||
// Boolean fields
|
||||
available: parseBoolean(unit.available),
|
||||
unavailable: parseBoolean(unit.unavailable),
|
||||
|
||||
// String fields (trim and preserve)
|
||||
unit_code: trimString(unit.unit_code),
|
||||
plan_name: trimString(unit.plan_name),
|
||||
soonest: trimString(unit.soonest),
|
||||
obj_type: trimString(unit.obj_type),
|
||||
community: trimString(unit.community),
|
||||
href: trimString(unit.href),
|
||||
inventory_href: trimString(unit.inventory_href),
|
||||
image_url: trimString(unit.image_url),
|
||||
specials_content: trimString(unit.specials_content)
|
||||
};
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Date Helpers
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Calculate yesterday's date from a given date string.
|
||||
* @param {string} dateStr - Date in YYYY-MM-DD format
|
||||
* @returns {string} Yesterday's date in YYYY-MM-DD format
|
||||
*/
|
||||
function getYesterday(dateStr) {
|
||||
const date = new Date(dateStr + 'T00:00:00Z');
|
||||
date.setUTCDate(date.getUTCDate() - 1);
|
||||
return date.toISOString().split('T')[0];
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 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;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Insert price records for today using bulkWrite with upsert.
|
||||
* Uses unit_code + date_checked as the composite key to prevent
|
||||
* duplicate price records for the same unit on the same day.
|
||||
* Re-running on the same day updates existing records (idempotent).
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Array<Object>} units - Array of unit objects
|
||||
* @param {string} date - Date in YYYY-MM-DD format
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object>} Insert result with insertedCount
|
||||
*/
|
||||
async function insertPrices(db, units, date, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.PRICES);
|
||||
const now = new Date().toISOString();
|
||||
|
||||
// Filter out units without prices
|
||||
const unitsWithPrices = units.filter(u => u.price !== null);
|
||||
|
||||
const priceRecords = unitsWithPrices.map(unit => ({
|
||||
unit_code: unit.unit_code,
|
||||
date_checked: date,
|
||||
price: unit.price,
|
||||
last_updated: now,
|
||||
data_source: 'web_scraper'
|
||||
}));
|
||||
|
||||
if (priceRecords.length === 0) {
|
||||
logger.warn('No price records to insert');
|
||||
return { insertedCount: 0 };
|
||||
}
|
||||
|
||||
// Use updateOne with upsert to handle re-runs on same day
|
||||
const operations = priceRecords.map(record => ({
|
||||
updateOne: {
|
||||
filter: {
|
||||
unit_code: record.unit_code,
|
||||
date_checked: record.date_checked
|
||||
},
|
||||
update: { $set: record },
|
||||
upsert: true
|
||||
}
|
||||
}));
|
||||
|
||||
try {
|
||||
const result = await collection.bulkWrite(operations, { ordered: false });
|
||||
|
||||
logger.info('Prices inserted', {
|
||||
inserted: result.upsertedCount,
|
||||
updated: result.modifiedCount
|
||||
});
|
||||
|
||||
return { insertedCount: result.upsertedCount + result.modifiedCount };
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Failed to insert prices', {
|
||||
errorType: error.name,
|
||||
errorMessage: error.message
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mark units not in current scrape as stale/unavailable.
|
||||
* Uses updateMany to set available: false and marked_stale_date
|
||||
* for all units whose unit_code is NOT in the current scrape
|
||||
* and that are currently available.
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Set<string>} currentUnitCodes - Unit codes from current scrape
|
||||
* @param {string} date - Current date in YYYY-MM-DD format
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object>} Update result
|
||||
*/
|
||||
async function markStaleUnits(db, currentUnitCodes, date, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.UNITS);
|
||||
|
||||
try {
|
||||
const result = await collection.updateMany(
|
||||
{
|
||||
unit_code: { $nin: [...currentUnitCodes] },
|
||||
available: true
|
||||
},
|
||||
{
|
||||
$set: {
|
||||
available: false,
|
||||
marked_stale_date: date
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
logger.info('Stale units marked', { count: result.modifiedCount });
|
||||
|
||||
return result;
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Failed to mark stale units', {
|
||||
errorType: error.name,
|
||||
errorMessage: error.message
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Update daily summary document.
|
||||
* Upserts by date field (YYYY-MM-DD format).
|
||||
* Fetches yesterday's summary for comparison metrics.
|
||||
* Calculates turnover_rate as (rentedUnits / yesterdayTotal * 100).
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Object} summaryData - Summary data
|
||||
* @param {string} summaryData.date - Date in YYYY-MM-DD format
|
||||
* @param {Array<string>} summaryData.newUnits - Unit codes added today
|
||||
* @param {Array<string>} summaryData.rentedUnits - Unit codes removed today
|
||||
* @param {number} summaryData.staleUnitsCount - Count of stale units
|
||||
* @param {number} summaryData.totalAvailable - Total available units today
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object>} MongoDB updateOne result
|
||||
*/
|
||||
async function updateDailySummary(db, summaryData, logger) {
|
||||
const collection = db.collection(config.COLLECTIONS.DAILY_SUMMARIES);
|
||||
const { date, newUnits, rentedUnits, staleUnitsCount, totalAvailable } = summaryData;
|
||||
|
||||
// Calculate yesterday's date for comparison
|
||||
const yesterday = getYesterday(date);
|
||||
|
||||
// Get yesterday's summary for comparison
|
||||
const yesterdaySummary = await collection.findOne({ date: yesterday });
|
||||
const yesterdayTotal = yesterdaySummary?.total_available_today || 0;
|
||||
|
||||
const summary = {
|
||||
date,
|
||||
timestamp: new Date().toISOString(),
|
||||
new_units: newUnits,
|
||||
rented_units: rentedUnits,
|
||||
stale_units: [], // Stale units list is not tracked per PRD
|
||||
new_units_count: newUnits.length,
|
||||
rented_units_count: rentedUnits.length,
|
||||
stale_units_count: staleUnitsCount,
|
||||
net_change: newUnits.length - rentedUnits.length,
|
||||
total_available_today: totalAvailable,
|
||||
total_available_yesterday: yesterdayTotal,
|
||||
turnover_rate: yesterdayTotal > 0
|
||||
? Math.round((rentedUnits.length / yesterdayTotal) * 10000) / 100
|
||||
: 0
|
||||
};
|
||||
|
||||
try {
|
||||
const result = await collection.updateOne(
|
||||
{ date },
|
||||
{ $set: summary },
|
||||
{ upsert: true }
|
||||
);
|
||||
|
||||
logger.info('Daily summary updated', {
|
||||
date,
|
||||
newUnits: newUnits.length,
|
||||
rentedUnits: rentedUnits.length,
|
||||
totalAvailable
|
||||
});
|
||||
|
||||
return result;
|
||||
|
||||
} catch (error) {
|
||||
logger.error('Failed to update daily summary', {
|
||||
errorType: error.name,
|
||||
errorMessage: error.message
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 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));
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Error Sanitization
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Sanitize an error object to remove sensitive information before storage.
|
||||
* Removes file paths, connection strings, and credential patterns while
|
||||
* preserving the error type and a useful general description for debugging.
|
||||
*
|
||||
* @param {Error} error - Error object to sanitize
|
||||
* @returns {Object} Sanitized error with name, message, and optionally stack
|
||||
*/
|
||||
function sanitizeError(error) {
|
||||
const sanitized = {
|
||||
name: error.name || 'Error',
|
||||
message: sanitizeMessage(error.message || ''),
|
||||
};
|
||||
|
||||
if (error.stack) {
|
||||
sanitized.stack = sanitizeMessage(error.stack);
|
||||
}
|
||||
|
||||
return sanitized;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sanitize a string message by removing sensitive patterns.
|
||||
* @param {string} message - Raw error message
|
||||
* @returns {string} Sanitized message
|
||||
*/
|
||||
function sanitizeMessage(message) {
|
||||
let result = message;
|
||||
|
||||
// Redact MongoDB connection strings (mongodb:// and mongodb+srv://)
|
||||
result = result.replace(/mongodb(\+srv)?:\/\/[^\s,;)}\]'"]+/gi, '[REDACTED_CONNECTION_STRING]');
|
||||
|
||||
// Redact credential/secret patterns: KEY=value, password=value, token=value, etc.
|
||||
result = result.replace(/\b(api[_-]?key|secret[_-]?key|secret[_-]?token|token|password|passwd|authorization|credential)\s*=\s*\S+/gi, '$1=[REDACTED]');
|
||||
|
||||
// Remove Unix absolute paths (/home/..., /var/..., /tmp/..., /usr/..., /etc/..., /opt/...)
|
||||
result = result.replace(/\/(?:home|var|tmp|usr|etc|opt)\/[^\s:,;)}\]'"]+/g, '[PATH]');
|
||||
|
||||
// Remove Windows-style absolute paths (C:\..., D:\...)
|
||||
result = result.replace(/[A-Z]:\\[^\s:,;)}\]'"]+/gi, '[PATH]');
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// Scraper Run History
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* Record scraper run to history collection.
|
||||
* Called in the finally block of runScrape() to persist run metadata.
|
||||
* This function must NOT throw errors - history recording should never break the scraper.
|
||||
*
|
||||
* @param {Db} db - MongoDB database instance
|
||||
* @param {Object} runData - Run data to record (jobId, trigger, status, duration, etc.)
|
||||
* @param {Object} logger - Logger instance
|
||||
* @returns {Promise<Object|null>} Insert result, or null on failure
|
||||
*/
|
||||
async function recordScraperRun(db, runData, logger) {
|
||||
try {
|
||||
const collection = db.collection(config.COLLECTIONS.SCRAPER_RUNS);
|
||||
const result = await collection.insertOne({
|
||||
...runData,
|
||||
recordedAt: new Date()
|
||||
});
|
||||
|
||||
return result;
|
||||
|
||||
} catch (error) {
|
||||
// Log but don't throw - recording history should not break scraper
|
||||
logger.error('Failed to record scraper run', { errorMessage: error.message });
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 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
|
||||
// ============================================================
|
||||
|
||||
/**
|
||||
* 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) {
|
||||
const cleanError = sanitizeError(error);
|
||||
logger.error('Scrape failed', {
|
||||
errorType: cleanError.name,
|
||||
errorMessage: cleanError.message
|
||||
});
|
||||
result.status = 'failed';
|
||||
result.errors.push(cleanError.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 = {
|
||||
runScrape,
|
||||
fetchPage,
|
||||
parseUnits,
|
||||
convertDataTypes,
|
||||
upsertUnits,
|
||||
insertPrices,
|
||||
markStaleUnits,
|
||||
updateDailySummary,
|
||||
recordScraperRun,
|
||||
createScraperIndexes,
|
||||
// Export helpers for testing
|
||||
getTodayUTC,
|
||||
getYesterdayUnitCodes,
|
||||
getYesterday,
|
||||
parseInteger,
|
||||
parsePositiveInteger,
|
||||
parseIntegerOrString,
|
||||
parsePrice,
|
||||
parseBoolean,
|
||||
parseFloatValue,
|
||||
trimString,
|
||||
isRetryableError,
|
||||
sleep,
|
||||
sanitizeError
|
||||
};
|
||||
Reference in New Issue
Block a user