Compare commits

..

1 Commits

Author SHA1 Message Date
5ef5d5af73 SCRAPE-11: Implement runScrape() main orchestration function
Some checks failed
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 / 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 9m56s
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.
2026-02-06 01:35:15 -07:00
4 changed files with 27 additions and 348 deletions

View File

@ -109,23 +109,11 @@ jobs:
steps: steps:
- name: Send results to n8n webhook - name: Send results to n8n webhook
env:
LINT_STATUS: ${{ needs.lint.result }}
TEST_STATUS: ${{ needs.test.result }}
RAW_LINT_OUTPUT: ${{ needs.lint.outputs.output }}
RAW_TEST_OUTPUT: ${{ needs.test.outputs.output }}
GH_REPO: ${{ github.repository }}
GH_BRANCH: ${{ github.head_ref || github.ref_name }}
GH_SHA: ${{ github.sha }}
GH_COMMIT_MSG: ${{ github.event.head_commit.message || github.event.pull_request.title || 'N/A' }}
GH_ACTOR: ${{ github.actor }}
GH_EVENT: ${{ github.event_name }}
GH_PR_NUMBER: ${{ github.event.pull_request.number || '' }}
GH_RUN_ID: ${{ github.run_id }}
GH_RUN_URL: ${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}
WEBHOOK_URL: ${{ secrets.N8N_WEBHOOK_URL }}
run: | run: |
# Determine overall status # Determine overall status
LINT_STATUS="${{ needs.lint.result }}"
TEST_STATUS="${{ needs.test.result }}"
if [ "$LINT_STATUS" = "success" ] && [ "$TEST_STATUS" = "success" ]; then if [ "$LINT_STATUS" = "success" ] && [ "$TEST_STATUS" = "success" ]; then
OVERALL_STATUS="success" OVERALL_STATUS="success"
else else
@ -133,21 +121,21 @@ jobs:
fi fi
# Truncate outputs if too long (max 10000 chars each) # Truncate outputs if too long (max 10000 chars each)
LINT_OUTPUT=$(echo "$RAW_LINT_OUTPUT" | head -c 10000) LINT_OUTPUT=$(echo '${{ needs.lint.outputs.output }}' | head -c 10000)
TEST_OUTPUT=$(echo "$RAW_TEST_OUTPUT" | head -c 10000) TEST_OUTPUT=$(echo '${{ needs.test.outputs.output }}' | head -c 10000)
# Build JSON payload # Build JSON payload
PAYLOAD=$(jq -n \ PAYLOAD=$(jq -n \
--arg repo "$GH_REPO" \ --arg repo "${{ github.repository }}" \
--arg branch "$GH_BRANCH" \ --arg branch "${{ github.head_ref || github.ref_name }}" \
--arg commit "$GH_SHA" \ --arg commit "${{ github.sha }}" \
--arg commit_short "${GH_SHA:0:7}" \ --arg commit_short "$(echo '${{ github.sha }}' | cut -c1-7)" \
--arg commit_message "$GH_COMMIT_MSG" \ --arg commit_message "${{ github.event.head_commit.message || github.event.pull_request.title || 'N/A' }}" \
--arg author "$GH_ACTOR" \ --arg author "${{ github.actor }}" \
--arg event "$GH_EVENT" \ --arg event "${{ github.event_name }}" \
--arg pr_number "$GH_PR_NUMBER" \ --arg pr_number "${{ github.event.pull_request.number || '' }}" \
--arg run_id "$GH_RUN_ID" \ --arg run_id "${{ github.run_id }}" \
--arg run_url "$GH_RUN_URL" \ --arg run_url "${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}" \
--arg overall_status "$OVERALL_STATUS" \ --arg overall_status "$OVERALL_STATUS" \
--arg lint_status "$LINT_STATUS" \ --arg lint_status "$LINT_STATUS" \
--arg lint_output "$LINT_OUTPUT" \ --arg lint_output "$LINT_OUTPUT" \
@ -187,7 +175,7 @@ jobs:
curl -X POST \ curl -X POST \
-H "Content-Type: application/json" \ -H "Content-Type: application/json" \
-d "$PAYLOAD" \ -d "$PAYLOAD" \
"$WEBHOOK_URL" \ "${{ secrets.N8N_WEBHOOK_URL }}" \
--fail --silent --show-error --fail --silent --show-error
- name: Fail if lint or tests failed - name: Fail if lint or tests failed

View File

@ -1,307 +0,0 @@
/**
* Tests for recordScraperRun() history persistence
*
* Covers:
* - insertOne is called with runData plus recordedAt timestamp
* - All runData fields are preserved (jobId, trigger, status, duration, etc.)
* - Error is caught and logged but NOT thrown (graceful failure)
* - Returns null on error instead of crashing
* - Returns insert result on success
* - recordedAt is a valid Date object
*/
const { MongoClient } = require('mongodb');
const { MongoMemoryServer } = require('mongodb-memory-server');
// Will require after implementation
let recordScraperRun;
let mongoServer;
let client;
let db;
let logger;
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
({ recordScraperRun } = require('../../services/scraperService'));
});
afterAll(async () => {
if (client) await client.close();
if (mongoServer) await mongoServer.stop();
});
beforeEach(async () => {
// Clean all collections before each test
const collections = await db.listCollections().toArray();
for (const col of collections) {
await db.collection(col.name).deleteMany({});
}
// Create fresh mock logger for each test
logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() };
});
describe('recordScraperRun', () => {
const SCRAPER_RUNS_COLLECTION = 'scraper_runs';
// Sample runData matching the structure from runScrape()'s finally block
function createSampleRunData(overrides = {}) {
return {
jobId: 'test-job-001',
trigger: 'manual',
dryRun: false,
status: 'success',
startedAt: '2026-02-06T06:00:00.000Z',
completedAt: '2026-02-06T06:00:12.345Z',
duration: 12345,
unitsProcessed: 50,
pricesInserted: 48,
newUnitsCount: 3,
rentedUnitsCount: 1,
staleUnitsCount: 2,
errors: [],
...overrides
};
}
// ---------------------------------------------------------------
// 1. insertOne is called with runData plus recordedAt timestamp
// ---------------------------------------------------------------
describe('insertOne with runData and recordedAt', () => {
it('should insert a document into the scraper_runs collection', async () => {
const runData = createSampleRunData();
await recordScraperRun(db, runData, logger);
const docs = await db.collection(SCRAPER_RUNS_COLLECTION).find({}).toArray();
expect(docs).toHaveLength(1);
});
it('should include a recordedAt field in the inserted document', async () => {
const runData = createSampleRunData();
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.recordedAt).toBeDefined();
});
});
// ---------------------------------------------------------------
// 2. All runData fields are preserved
// ---------------------------------------------------------------
describe('preserving all runData fields', () => {
it('should preserve jobId, trigger, status, and duration', async () => {
const runData = createSampleRunData({
jobId: 'preserve-test-001',
trigger: 'scheduled',
status: 'failed',
duration: 9999
});
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.jobId).toBe('preserve-test-001');
expect(doc.trigger).toBe('scheduled');
expect(doc.status).toBe('failed');
expect(doc.duration).toBe(9999);
});
it('should preserve unit processing metrics', async () => {
const runData = createSampleRunData({
unitsProcessed: 42,
pricesInserted: 40,
newUnitsCount: 5,
rentedUnitsCount: 2,
staleUnitsCount: 3
});
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.unitsProcessed).toBe(42);
expect(doc.pricesInserted).toBe(40);
expect(doc.newUnitsCount).toBe(5);
expect(doc.rentedUnitsCount).toBe(2);
expect(doc.staleUnitsCount).toBe(3);
});
it('should preserve timing fields (startedAt, completedAt)', async () => {
const runData = createSampleRunData({
startedAt: '2026-02-06T10:00:00.000Z',
completedAt: '2026-02-06T10:00:15.500Z'
});
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.startedAt).toBe('2026-02-06T10:00:00.000Z');
expect(doc.completedAt).toBe('2026-02-06T10:00:15.500Z');
});
it('should preserve dryRun flag', async () => {
const runData = createSampleRunData({ dryRun: true });
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.dryRun).toBe(true);
});
it('should preserve errors array when it has entries', async () => {
const runData = createSampleRunData({
status: 'failed',
errors: ['Connection timeout', 'Retry exhausted']
});
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.errors).toEqual(['Connection timeout', 'Retry exhausted']);
});
it('should preserve empty errors array for successful runs', async () => {
const runData = createSampleRunData({ errors: [] });
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.errors).toEqual([]);
});
});
// ---------------------------------------------------------------
// 3. Error is caught and logged but NOT thrown (graceful failure)
// ---------------------------------------------------------------
describe('graceful error handling', () => {
it('should NOT throw when insertOne fails', async () => {
const runData = createSampleRunData();
// Create a mock db that throws on insertOne
const mockCollection = {
insertOne: jest.fn().mockRejectedValue(new Error('Write concern timeout'))
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
// This should NOT throw
await expect(recordScraperRun(mockDb, runData, logger)).resolves.not.toThrow();
});
it('should log error message when insert fails', async () => {
const runData = createSampleRunData();
const mockCollection = {
insertOne: jest.fn().mockRejectedValue(new Error('Disk full'))
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
await recordScraperRun(mockDb, runData, logger);
expect(logger.error).toHaveBeenCalledWith(
'Failed to record scraper run',
{ errorMessage: 'Disk full' }
);
});
});
// ---------------------------------------------------------------
// 4. Returns null on error instead of crashing
// ---------------------------------------------------------------
describe('return value on error', () => {
it('should return null when insertOne throws', async () => {
const runData = createSampleRunData();
const mockCollection = {
insertOne: jest.fn().mockRejectedValue(new Error('Connection refused'))
};
const mockDb = {
collection: jest.fn().mockReturnValue(mockCollection)
};
const result = await recordScraperRun(mockDb, runData, logger);
expect(result).toBeNull();
});
});
// ---------------------------------------------------------------
// 5. Returns insert result on success
// ---------------------------------------------------------------
describe('return value on success', () => {
it('should return the insertOne result object', async () => {
const runData = createSampleRunData();
const result = await recordScraperRun(db, runData, logger);
expect(result).toBeDefined();
expect(result).not.toBeNull();
});
it('should return result with acknowledged property', async () => {
const runData = createSampleRunData();
const result = await recordScraperRun(db, runData, logger);
expect(result.acknowledged).toBe(true);
});
it('should return result with insertedId', async () => {
const runData = createSampleRunData();
const result = await recordScraperRun(db, runData, logger);
expect(result.insertedId).toBeDefined();
});
});
// ---------------------------------------------------------------
// 6. recordedAt is a valid Date object
// ---------------------------------------------------------------
describe('recordedAt field', () => {
it('should set recordedAt as a Date instance', async () => {
const runData = createSampleRunData();
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.recordedAt).toBeInstanceOf(Date);
});
it('should set recordedAt close to current time', async () => {
const beforeTime = new Date();
const runData = createSampleRunData();
await recordScraperRun(db, runData, logger);
const afterTime = new Date();
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
expect(doc.recordedAt.getTime()).toBeGreaterThanOrEqual(beforeTime.getTime());
expect(doc.recordedAt.getTime()).toBeLessThanOrEqual(afterTime.getTime());
});
it('should not overwrite any existing runData fields with recordedAt', async () => {
const runData = createSampleRunData();
await recordScraperRun(db, runData, logger);
const doc = await db.collection(SCRAPER_RUNS_COLLECTION).findOne({});
// Verify the original fields still exist alongside recordedAt
expect(doc.jobId).toBe(runData.jobId);
expect(doc.status).toBe(runData.status);
expect(doc.recordedAt).toBeInstanceOf(Date);
});
});
});

View File

@ -110,8 +110,7 @@ describe('recordScraperRun', () => {
errors: [] errors: []
}; };
const mockLogger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() }; const result = await recordScraperRun(db, runData);
const result = await recordScraperRun(db, runData, mockLogger);
expect(result).toBeTruthy(); expect(result).toBeTruthy();
expect(result.insertedId).toBeTruthy(); expect(result.insertedId).toBeTruthy();
@ -128,10 +127,8 @@ describe('recordScraperRun', () => {
const runData = { jobId: 'test-fail', status: 'success' }; const runData = { jobId: 'test-fail', status: 'success' };
// Should not throw // Should not throw
const mockLogger = { info: jest.fn(), warn: jest.fn(), error: jest.fn() }; const result = await recordScraperRun(null, runData);
const result = await recordScraperRun(null, runData, mockLogger);
expect(result).toBeNull(); expect(result).toBeNull();
expect(mockLogger.error).toHaveBeenCalled();
}); });
}); });

View File

@ -657,17 +657,18 @@ async function getYesterdayUnitCodes(db, today) {
/** /**
* Record scraper run to history collection. * Record scraper run to history collection.
* Called in the finally block of runScrape() to persist run metadata. * This function intentionally catches errors and returns null
* This function must NOT throw errors - history recording should never break the scraper. * rather than throwing, because recording history should not
* break the main scraper workflow.
* *
* @param {Db} db - MongoDB database instance * @param {Db} db - MongoDB database instance
* @param {Object} runData - Run data to record (jobId, trigger, status, duration, etc.) * @param {Object} runData - Run data to record
* @param {Object} logger - Logger instance * @returns {Promise<Object|null>} Insert result or null on error
* @returns {Promise<Object|null>} Insert result, or null on failure
*/ */
async function recordScraperRun(db, runData, logger) { async function recordScraperRun(db, runData) {
try { try {
const collection = db.collection(config.COLLECTIONS.SCRAPER_RUNS); 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()
@ -677,7 +678,7 @@ async function recordScraperRun(db, runData, logger) {
} catch (error) { } catch (error) {
// Log but don't throw - recording history should not break scraper // Log but don't throw - recording history should not break scraper
logger.error('Failed to record scraper run', { errorMessage: error.message }); console.error('Failed to record scraper run:', error.message);
return null; return null;
} }
} }
@ -799,7 +800,7 @@ async function runScrape(db, options = {}) {
result.duration = Date.now() - startTime; result.duration = Date.now() - startTime;
// Record run to history (always runs, even on failure) // Record run to history (always runs, even on failure)
await recordScraperRun(db, result, logger); await recordScraperRun(db, result);
logger.info('Scrape completed', { logger.info('Scrape completed', {
status: result.status, status: result.status,