← back to Stayclaim

scripts/ingest-la-code-enforcement.ts

170 lines

/**
 * ingest-la-code-enforcement.ts
 *
 * LA City code enforcement cases via Socrata.
 *   - u82d-eh7z (Open):    ~28,931 rows
 *   - rken-a55j (Closed): ~810,798 rows
 *
 * Stores as place_event rows with source='la_city_code', source_id=apno.
 * Links to existing listing via address match (from LADBS ingest).
 */
import { Pool } from 'pg';

const pool = new Pool({
  host: process.env.PGHOST ?? '/tmp',
  database: process.env.PGDATABASE ?? 'stayclaim',
  user: process.env.PGUSER ?? process.env.USER,
  password: process.env.PGPASSWORD,
  port: parseInt(process.env.PGPORT ?? '5432', 10),
  max: 8,
});

const PAGE = 50000;
const BATCH = 1000;

type SrcRow = {
  apno?: string;       // case number
  prclid?: string;     // parcel id
  aptype?: string;     // case type
  stat?: string;       // status
  adddttm?: string;    // added date
  housenumberhigh?: string;
  housenumber?: string;
  streetname?: string;
  streetdirection?: string;
  streetsuffix?: string;
  zipcode?: string;
  cd?: string;
};

function canonicalize(addr: string): string {
  return addr.toLowerCase().replace(/[^\w\s-]/g, '').replace(/\s+/g, '-').replace(/-+/g, '-').replace(/^-|-$/g, '').slice(0, 90);
}

function buildAddress(r: SrcRow): string | null {
  const parts = [r.housenumber, r.streetdirection, r.streetname, r.streetsuffix].filter(Boolean);
  const addr = parts.join(' ').trim();
  return addr || null;
}

async function ensureSchema() {
  await pool.query(`
    CREATE TABLE IF NOT EXISTS code_enforcement_case (
      id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
      source TEXT NOT NULL,
      source_dataset TEXT NOT NULL,
      case_number TEXT NOT NULL,
      listing_id UUID REFERENCES listing(id),
      primary_address TEXT,
      apn TEXT,
      zip_code TEXT,
      council_district TEXT,
      case_type TEXT,
      status TEXT,
      opened_at DATE,
      source_url TEXT,
      source_label TEXT NOT NULL DEFAULT 'LA City Department of Building & Safety',
      source_tier CHAR(1) NOT NULL DEFAULT 'A',
      retrieved_at TIMESTAMPTZ DEFAULT now(),
      UNIQUE (source_dataset, case_number)
    );
    CREATE INDEX IF NOT EXISTS idx_cec_listing ON code_enforcement_case(listing_id);
    CREATE INDEX IF NOT EXISTS idx_cec_apn ON code_enforcement_case(apn);
    CREATE INDEX IF NOT EXISTS idx_cec_status ON code_enforcement_case(status);
    CREATE INDEX IF NOT EXISTS idx_cec_opened ON code_enforcement_case(opened_at);
  `);
}

async function fetchPage(slug: string, offset: number): Promise<SrcRow[]> {
  const u = `https://data.lacity.org/resource/${slug}.json?$limit=${PAGE}&$offset=${offset}&$order=apno`;
  for (let attempt = 1; attempt <= 4; attempt++) {
    try {
      const res = await fetch(u, { headers: { 'User-Agent': 'stayclaim-ingest/1.0' } });
      if (res.status === 429 || res.status === 503) { await new Promise(r => setTimeout(r, 2000 * attempt)); continue; }
      if (!res.ok) throw new Error(`HTTP ${res.status}`);
      return await res.json() as SrcRow[];
    } catch (e) {
      if (attempt >= 4) throw e;
      await new Promise(r => setTimeout(r, 1500 * attempt));
    }
  }
  return [];
}

async function processSlug(slug: string, dataset: string) {
  console.log(`\n=== ${dataset} (${slug}) ===`);
  let offset = 0;
  let total = 0;
  const t0 = Date.now();
  while (true) {
    const rows = await fetchPage(slug, offset);
    if (!rows.length) break;

    // bulk insert; dedupe by case_number within batch first, then resolve listing_id.
    for (let i = 0; i < rows.length; i += BATCH) {
      const slice = rows.slice(i, i + BATCH);
      const seen = new Set<string>();
      const dedup: SrcRow[] = [];
      for (const r of slice) {
        if (!r.apno) continue;
        if (seen.has(r.apno)) continue;
        seen.add(r.apno);
        dedup.push(r);
      }
      const tuples: string[] = [];
      const values: any[] = [];
      dedup.forEach((r, idx) => {
        const addr = buildAddress(r);
        const base = idx * 9;
        tuples.push(`($${base+1},$${base+2},$${base+3},$${base+4},$${base+5},$${base+6},$${base+7},$${base+8},$${base+9})`);
        values.push(
          slug,                                                      // 1 source_dataset
          r.apno,                                                    // 2 case_number
          addr,                                                      // 3 primary_address
          r.prclid ?? null,                                          // 4 apn
          r.zipcode ?? null,                                         // 5 zip_code
          r.cd ?? null,                                              // 6 council_district
          r.aptype ?? null,                                          // 7 case_type
          r.stat ?? null,                                            // 8 status
          r.adddttm ? new Date(r.adddttm).toISOString().slice(0, 10) : null, // 9 opened_at
        );
      });
      if (tuples.length === 0) continue;
      await pool.query(
        `INSERT INTO code_enforcement_case
           (source, source_dataset, case_number, primary_address, apn, zip_code, council_district, case_type, status, opened_at, listing_id, source_url)
         SELECT 'la_city', t.dataset, t.cn, t.addr, t.apn, t.zip, t.cd, t.ct, t.st, t.opened, l.id,
                'https://data.lacity.org/resource/' || t.dataset || '.json?apno=' || t.cn
         FROM (VALUES ${tuples.map((tp, k) => {
           const b = k * 9;
           return `($${b+1},$${b+2},$${b+3},$${b+4},$${b+5},$${b+6},$${b+7},$${b+8},$${b+9}::date)`;
         }).join(',')}) AS t(dataset, cn, addr, apn, zip, cd, ct, st, opened)
         LEFT JOIN listing l ON l.source='ladbs_permit' AND l.source_id = 'ladbs:' || lower(regexp_replace(regexp_replace(coalesce(t.addr,''), '[^\\w\\s-]', '', 'g'), '\\s+', '-', 'g'))
         ON CONFLICT (source_dataset, case_number) DO UPDATE SET
           status = COALESCE(EXCLUDED.status, code_enforcement_case.status),
           listing_id = COALESCE(code_enforcement_case.listing_id, EXCLUDED.listing_id),
           retrieved_at = now()`,
        values
      );
    }
    total += rows.length;
    const dt = (Date.now() - t0) / 1000;
    console.log(`  ${slug}: ${total} (${(total/dt).toFixed(0)}/s)`);
    if (rows.length < PAGE) break;
    offset += PAGE;
  }
}

async function main() {
  await ensureSchema();
  await processSlug('u82d-eh7z', 'open');
  await processSlug('rken-a55j', 'closed');
  const { rows } = await pool.query<{ ds: string; n: string; matched: string }>(`
    SELECT source_dataset as ds, count(*)::text as n, count(*) FILTER (WHERE listing_id IS NOT NULL)::text as matched
    FROM code_enforcement_case GROUP BY source_dataset`);
  for (const r of rows) console.log(`  ${r.ds}: ${parseInt(r.n).toLocaleString()} cases (${parseInt(r.matched).toLocaleString()} matched to listing)`);
  await pool.end();
}

main().catch(e => { console.error('FATAL', e); process.exit(1); });