← back to Omega Watches 2

collectors/runner.mjs

75 lines

/**
 * Collector Runner — Orchestrates collection jobs
 * Usage: node collectors/runner.mjs [source_name] [job_type]
 */
import { pool } from './base-collector.mjs';

const sourceName = process.argv[2];
const jobType = process.argv[3] || 'incremental';

async function run() {
  if (!sourceName) {
    // List available sources
    const { rows } = await pool.query(
      'SELECT name, display_name, source_type, access_method, is_active FROM data_source ORDER BY name'
    );
    console.log('\nAvailable sources:');
    console.log('─'.repeat(80));
    for (const s of rows) {
      const status = s.is_active ? '✓' : '✗';
      console.log(`  ${status} ${s.name.padEnd(20)} ${s.display_name.padEnd(25)} ${s.source_type.padEnd(15)} ${s.access_method}`);
    }
    console.log(`\nUsage: node collectors/runner.mjs <source_name> [job_type]`);
    console.log('Job types: incremental, daily_msrp, backfill, full_crawl\n');
    process.exit(0);
  }

  try {
    // Try both underscore and hyphen directory conventions
    const candidates = [
      `./${sourceName}/collector.mjs`,
      `./${sourceName.replace(/_/g, '-')}/collector.mjs`,
    ];
    let loaded = false;

    for (const modulePath of candidates) {
      try {
        const { default: CollectorClass } = await import(modulePath);
        const collector = new CollectorClass();
        await collector.run();
        loaded = true;
        break;
      } catch (importErr) {
        if (importErr.code !== 'ERR_MODULE_NOT_FOUND') throw importErr;
      }
    }

    if (!loaded) {
      console.log(`No collector module found for "${sourceName}"`);
      console.log(`Expected at: ${candidates.join(' or ')}`);
      console.log(`\nRegistering a placeholder run...`);

      const { rows: sources } = await pool.query(
        'SELECT id FROM data_source WHERE name = $1', [sourceName]
      );
      if (!sources[0]) {
        console.error(`Source "${sourceName}" not found in database`);
        process.exit(1);
      }
      await pool.query(
        `INSERT INTO collector_run (source_id, job_type, status, parser_version, metadata)
         VALUES ($1, $2, 'pending', '1.0.0', '{"note": "No collector module yet"}')`,
        [sources[0].id, jobType]
      );
      console.log('Placeholder run created.');
    }
  } catch (err) {
    console.error('Runner error:', err.message);
    process.exit(1);
  } finally {
    await pool.end();
  }
}

run();