← back to Patternbank Archive

src/db.js

66 lines

require('dotenv').config({ path: require('path').join(__dirname, '..', '.env') });
const { Pool } = require('pg');

const pool = process.env.DATABASE_URL
  ? new Pool({ connectionString: process.env.DATABASE_URL, max: 8 })
  : new Pool({
      database: process.env.PG_DATABASE || 'patternbank_archive',
      host: process.env.PG_HOST || '127.0.0.1',
      port: Number(process.env.PG_PORT || 5432),
      user: process.env.PG_USER,
      max: 8,
    });

async function q(sql, params) {
  const c = await pool.connect();
  try { return await c.query(sql, params); } finally { c.release(); }
}

async function startRun(mode, notes) {
  const id = `${mode}-${new Date().toISOString().replace(/[:.]/g,'-')}`;
  await q(`INSERT INTO ingest_runs (run_id, mode, notes) VALUES ($1,$2,$3)`, [id, mode, notes||null]);
  return id;
}

async function finishRun(id, counters) {
  const cols = Object.keys(counters||{});
  if (!cols.length) {
    await q(`UPDATE ingest_runs SET finished_at=NOW() WHERE run_id=$1`, [id]);
    return;
  }
  const set = cols.map((k,i) => `${k}=$${i+2}`).join(', ');
  const vals = cols.map(k => counters[k]);
  await q(`UPDATE ingest_runs SET finished_at=NOW(), ${set} WHERE run_id=$1`, [id, ...vals]);
}

async function enqueue(url, kind, priority=100) {
  await q(
    `INSERT INTO crawl_queue (url, kind, priority) VALUES ($1,$2,$3)
     ON CONFLICT (url) DO NOTHING`,
    [url, kind, priority]
  );
}

async function popQueue(kind, limit=20) {
  const r = await q(
    `SELECT url FROM crawl_queue
     WHERE kind=$1 AND done=FALSE AND attempts<5
     ORDER BY priority ASC, enqueued_at ASC LIMIT $2`,
    [kind, limit]
  );
  return r.rows.map(x => x.url);
}

async function markQueue(url, status, err) {
  await q(
    `UPDATE crawl_queue
     SET attempts=attempts+1, last_attempt_at=NOW(),
         last_status=$2, last_error=$3,
         done=CASE WHEN $2 BETWEEN 200 AND 299 THEN TRUE ELSE done END
     WHERE url=$1`,
    [url, status||null, err||null]
  );
}

module.exports = { pool, q, startRun, finishRun, enqueue, popQueue, markQueue };