← back to La Socrata Ingester

src/ingest.js

77 lines

import { upsert, getState, setState, pool } from './db.js';
import { socrataPages } from './adapters/socrata.js';
import { arcgisPages } from './adapters/arcgis.js';
import { sleep } from './adapters/http.js';
import { SOURCES, resolveMap } from './sources.js';

// Tables that carry a dataset_id column (so we stamp it onto each row).
const HAS_DATASET_ID = new Set(['la_building_permits_raw', 'la_code_enforcement_raw']);

const log = (...a) => console.log(new Date().toISOString(), ...a);

// Run one source end-to-end.
// opts: { full, maxRows, pageSize }
export async function ingestSource(name, opts = {}) {
  const src = SOURCES[name];
  if (!src) throw new Error(`unknown source: ${name}`);
  const map = resolveMap(name);
  const datasetLabel = src.datasetId || name;
  const appToken = process.env.SOCRATA_APP_TOKEN || undefined;
  const pageDelay = Number(process.env.INGEST_PAGE_DELAY_MS || 120);

  const prev = await getState(name);
  const since = !opts.full && src.cursorField ? prev?.last_cursor || null : null;

  const pagerOpts = { since, appToken, maxRows: opts.maxRows ?? Infinity, fullScan: !!opts.full };
  if (opts.pageSize) pagerOpts.pageSize = opts.pageSize;
  // --year=YYYY (ArcGIS only): pull exactly one roll year, resumable/idempotent.
  if (opts.year && src.platform === 'arcgis' && src.cursorField) {
    pagerOpts.whereOverride = `${src.cursorField}='${opts.year}'`;
  }
  const pager = src.platform === 'socrata' ? socrataPages(src, pagerOpts) : arcgisPages(src, pagerOpts);

  const mode = opts.year ? `year=${opts.year}` : opts.full ? 'FULL' : since ? `since ${since}` : 'initial';
  log(`▶ ${name} [${src.platform} ${datasetLabel}] ${mode}${opts.maxRows ? ` (max ${opts.maxRows})` : ''}`);

  let total = 0;
  let maxCursor = prev?.last_cursor || null;
  let status = 'ok';
  try {
    for await (const { rows, url } of pager) {
      const mapped = rows.map((r) => {
        const row = { ...map(r), raw: r, source_url: url };
        if (HAS_DATASET_ID.has(src.table)) row.dataset_id = datasetLabel;
        return row;
      });
      // drop rows with a null primary key component (can't upsert those)
      const clean = mapped.filter((row) => src.conflict.every((k) => row[k] != null));
      const skipped = mapped.length - clean.length;
      await upsert(src.table, src.conflict, clean);
      total += clean.length;

      // advance high-water mark
      if (src.cursorField) {
        for (const r of rows) {
          const c = r[src.cursorField];
          if (c != null && (maxCursor == null || String(c) > maxCursor)) maxCursor = String(c);
        }
      }
      log(`  +${clean.length}${skipped ? ` (skipped ${skipped} keyless)` : ''}  total=${total}`);
      if (pageDelay) await sleep(pageDelay);
    }
  } catch (err) {
    status = 'error: ' + err.message;
    log(`✖ ${name} failed: ${err.message}`);
    await setState(name, { dataset_id: datasetLabel, last_cursor: maxCursor, rows_upserted: total, last_status: status });
    throw err;
  }

  await setState(name, { dataset_id: datasetLabel, last_cursor: maxCursor, rows_upserted: total, last_status: status });
  log(`✔ ${name} done — ${total} rows upserted this run${maxCursor ? `, cursor→${maxCursor}` : ''}`);
  return { name, total, maxCursor };
}

export async function closePool() {
  await pool.end();
}