← back to Costa Rica

scripts/ingest/_lib.js

112 lines

'use strict';

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

const pool = new Pool({ connectionString: process.env.DATABASE_URL });

const UA = process.env.SCRAPER_UA || 'Mozilla/5.0 (Macintosh; Intel Mac OS X 14_5) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.5 Safari/605.1.15 (CR-Directory bot; contact info@agentabrams.com)';

async function fetchText(url, opts = {}) {
  const res = await fetch(url, {
    headers: { 'User-Agent': UA, 'Accept': 'text/html,application/xhtml+xml,application/json;q=0.9,*/*;q=0.8', 'Accept-Language': 'es-CR,es;q=0.9,en;q=0.7', ...(opts.headers||{}) },
    redirect: 'follow',
    signal: opts.signal,
  });
  if (!res.ok) throw new Error(`HTTP ${res.status} on ${url}`);
  return await res.text();
}

async function fetchBuffer(url, opts = {}) {
  const res = await fetch(url, { headers: { 'User-Agent': UA, ...(opts.headers||{}) }, redirect: 'follow', signal: opts.signal });
  if (!res.ok) throw new Error(`HTTP ${res.status} on ${url}`);
  return Buffer.from(await res.arrayBuffer());
}

async function fetchJson(url, opts = {}) {
  const res = await fetch(url, { headers: { 'User-Agent': UA, 'Accept': 'application/json', ...(opts.headers||{}) }, redirect: 'follow', signal: opts.signal });
  if (!res.ok) throw new Error(`HTTP ${res.status} on ${url}`);
  return await res.json();
}

const slug = (s, maxLen = 96) => String(s||'')
  .normalize('NFD').replace(/[̀-ͯ]/g,'')
  .toLowerCase()
  .replace(/&/g,' and ')
  .replace(/[^a-z0-9]+/g,'-')
  .replace(/^-+|-+$/g,'')
  .slice(0, maxLen);

const sleep = ms => new Promise(r => setTimeout(r, ms));

async function regionMap() {
  const { rows } = await pool.query('SELECT id, slug, name, lat, lng FROM regions');
  const bySlug = new Map(rows.map(r => [r.slug, r]));
  const byName = new Map(rows.map(r => [r.name.toLowerCase(), r]));
  return { rows, bySlug, byName, fallback: bySlug.get('cr-other') };
}

function resolveRegion(rmap, hint) {
  if (!hint) return rmap.fallback;
  const h = String(hint).toLowerCase().trim();
  if (rmap.bySlug.has(h)) return rmap.bySlug.get(h);
  if (rmap.byName.has(h)) return rmap.byName.get(h);
  for (const r of rmap.rows) {
    if (h.includes(r.name.toLowerCase()) || h.includes(r.slug.replace(/-/g,' '))) return r;
  }
  return rmap.fallback;
}

async function ensureRegion(rmap, name, province) {
  if (!name) return rmap.fallback;
  const cleanName = String(name).trim();
  const cleanProv = String(province||'').trim() || 'CR';
  const key = cleanName.toLowerCase();
  if (rmap.byName.has(key)) return rmap.byName.get(key);
  const s = slug(cleanName);
  if (rmap.bySlug.has(s)) return rmap.bySlug.get(s);
  // INSERT new cantón
  const { rows } = await pool.query(
    `INSERT INTO regions (slug, name, province, region_type) VALUES ($1, $2, $3, 'canton')
     ON CONFLICT (slug) DO UPDATE SET province = EXCLUDED.province
     RETURNING id, slug, name, province, lat, lng`,
    [s, cleanName, cleanProv]
  );
  const created = rows[0];
  rmap.bySlug.set(created.slug, created);
  rmap.byName.set(created.name.toLowerCase(), created);
  rmap.rows.push(created);
  return created;
}

async function upsertPlace(p) {
  const cols = ['slug','name','category','vertical','region_id','description','address','phone','email','website','price_range','rating','image_url','tags','lat','lng','cedula_juridica','source','source_url'];
  const vals = cols.map(c => p[c] ?? null);
  const placeholders = cols.map((_, i) => `$${i+1}`).join(',');
  const updates = cols.filter(c => c!=='slug').map(c => `${c} = COALESCE(EXCLUDED.${c}, places.${c})`).join(', ');
  const sql = `
    INSERT INTO places (${cols.join(',')}) VALUES (${placeholders})
    ON CONFLICT (slug) DO UPDATE SET ${updates}, updated_at = NOW()
    RETURNING (xmax = 0) AS inserted, id
  `;
  const { rows } = await pool.query(sql, vals);
  return rows[0]; // { inserted: bool, id }
}

async function startRun(source, notes='') {
  const { rows } = await pool.query(
    `INSERT INTO ingest_runs (source, status, notes) VALUES ($1, 'running', $2) RETURNING id`,
    [source, notes]
  );
  return rows[0].id;
}

async function finishRun(id, { rows_in=0, rows_added=0, rows_updated=0, status='ok', notes='' } = {}) {
  await pool.query(
    `UPDATE ingest_runs SET finished_at=NOW(), rows_in=$2, rows_added=$3, rows_updated=$4, status=$5, notes=COALESCE($6, notes) WHERE id=$1`,
    [id, rows_in, rows_added, rows_updated, status, notes]
  );
}

module.exports = { pool, UA, fetchText, fetchJson, fetchBuffer, slug, sleep, regionMap, resolveRegion, ensureRegion, upsertPlace, startRun, finishRun };