← 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 };