← back to Stayclaim

scripts/ingest-la-cofo.ts

145 lines

/**
 * ingest-la-cofo.ts
 *
 * LA City Certificate of Occupancy (3f9m-afei, ~88K rows).
 * Fields: cofo_number, address_start, street_direction, street_name, street_suffix,
 *         issue_date, latitude_longitude, valuation, work_description, etc.
 *
 * Tier: A (LADBS primary record).
 */
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: 6,
});

const SLUG = '3f9m-afei';
const PAGE = 50000;
const BATCH = 1000;

type Row = {
  cofo_number?: string;
  pcis_permit?: string;
  address_start?: string;
  address_end?: string;
  street_direction?: string;
  street_name?: string;
  street_suffix?: string;
  zip_code?: string;
  zone?: string;
  permit_type?: string;
  permit_sub_type?: string;
  permit_category?: string;
  issue_date?: string;
  cofo_issue_date?: string;
  latest_status?: string;
  status_date?: string;
  valuation?: string;
  work_description?: string;
  assessor_parcel?: string;
  latitude_longitude?: { latitude?: string; longitude?: string };
};

function pd(v?: string): string | null { if (!v) return null; const m = v.match(/^(\d{4}-\d{2}-\d{2})/); return m ? m[1] : null; }
function pn(v?: string): number | null { if (!v) return null; const n = parseFloat(v); return isFinite(n) ? n : null; }

function buildAddress(r: Row): string | null {
  const parts = [r.address_start, r.street_direction, r.street_name, r.street_suffix].filter(Boolean);
  const a = parts.join(' ').trim();
  return a || null;
}

async function fetchPage(offset: number): Promise<Row[]> {
  const u = `https://data.lacity.org/resource/${SLUG}.json?$limit=${PAGE}&$offset=${offset}&$order=cofo_number`;
  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 Row[];
    } catch (e) {
      if (attempt >= 4) throw e;
      await new Promise(r => setTimeout(r, 1500 * attempt));
    }
  }
  return [];
}

async function processBatch(rows: Row[]) {
  const seen = new Set<string>();
  const valid = rows.filter(r => r.cofo_number && !seen.has(r.cofo_number) && (seen.add(r.cofo_number), true));
  if (valid.length === 0) return 0;
  const tuples: string[] = [];
  const values: any[] = [];
  valid.forEach((r, idx) => {
    const addr = buildAddress(r);
    const lat = r.latitude_longitude?.latitude ? parseFloat(r.latitude_longitude.latitude) : null;
    const lon = r.latitude_longitude?.longitude ? parseFloat(r.latitude_longitude.longitude) : null;
    const base = idx * 14;
    tuples.push(`(${Array.from({length: 14}, (_, k) => `$${base+k+1}`).join(',')})`);
    values.push(
      r.cofo_number!.trim(),                          // 1 permit_number
      addr,                                           // 2 primary_address
      r.zip_code ?? null,                             // 3 zip_code
      r.assessor_parcel ?? null,                      // 4 apn
      r.zone ?? null,                                 // 5 zone
      r.permit_type ?? null,                          // 6 permit_type
      r.permit_sub_type ?? null,                      // 7 permit_sub_type
      r.permit_category ?? null,                      // 8 permit_group
      pd(r.cofo_issue_date ?? r.issue_date),          // 9 issue_date
      r.latest_status ?? null,                        // 10 status_desc
      pn(r.valuation),                                // 11 valuation
      r.work_description ?? null,                     // 12 work_desc
      lat,                                            // 13 latitude
      lon,                                            // 14 longitude
    );
  });
  await pool.query(
    `INSERT INTO permit
       (source, source_dataset, permit_number, primary_address, zip_code, apn, zone,
        permit_type, permit_sub_type, permit_group,
        issue_date, status_desc, valuation, work_desc, latitude, longitude,
        listing_id, source_url, source_label)
     SELECT 'la_city_cofo', '${SLUG}', t.pn, t.addr, t.zip, t.apn, t.zone,
            t.pt, t.pst, t.pg, t.iss, t.status, t.val, t.wd, t.lat, t.lon, l.id,
            'https://data.lacity.org/resource/${SLUG}.json?cofo_number=' || t.pn,
            'LADBS Certificate of Occupancy'
     FROM (VALUES ${tuples.map((_, k) => {
       const b = k * 14;
       return `($${b+1},$${b+2},$${b+3},$${b+4},$${b+5},$${b+6},$${b+7},$${b+8},$${b+9}::date,$${b+10},$${b+11}::numeric,$${b+12},$${b+13}::numeric,$${b+14}::numeric)`;
     }).join(',')}) AS t(pn, addr, zip, apn, zone, pt, pst, pg, iss, status, val, wd, lat, lon)
     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, permit_number) DO UPDATE SET
       status_desc = COALESCE(EXCLUDED.status_desc, permit.status_desc),
       listing_id = COALESCE(permit.listing_id, EXCLUDED.listing_id),
       retrieved_at = now()`,
    values
  );
  return valid.length;
}

async function main() {
  let offset = 0, total = 0;
  const t0 = Date.now();
  while (true) {
    const rows = await fetchPage(offset);
    if (rows.length === 0) break;
    for (let i = 0; i < rows.length; i += BATCH) await processBatch(rows.slice(i, i + BATCH));
    total += rows.length;
    const dt = (Date.now() - t0) / 1000;
    console.log(`  cofo: ${total.toLocaleString()} (${(total/dt).toFixed(0)}/s)`);
    if (rows.length < PAGE) break;
    offset += PAGE;
  }
  console.log(`✓ Certificate of Occupancy: ${total.toLocaleString()} total`);
  await pool.end();
}

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