← back to Stayclaim

scripts/ingest-la-additional-permits.ts

222 lines

/**
 * ingest-la-additional-permits.ts
 *
 * Pulls 5 additional LADBS permit datasets that weren't in the main building-permit
 * ingest. Same Socrata pattern; each row creates/updates a listing and inserts
 * a permit row.
 *
 *   - 3f9m-afei  Certificate of Occupancy        (88K)
 *   - 67is-svtd  Mechanical Permits 2020+        (279K)
 *   - 5m3t-xjex  Mechanical Permits 2010-2019    (503K)
 *   - ysqd-apz7  Electrical Permits 2020+        (338K)
 *   - j7mw-thyc  Bureau of Engineering Permits   (626K)
 *
 * Tier: A (LA City government primary record).
 *
 * Usage: npx tsx scripts/ingest-la-additional-permits.ts [--dataset=SLUG]
 */
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 args = Object.fromEntries(
  process.argv.slice(2).map(a => {
    const m = a.match(/^--([^=]+)(?:=(.*))?$/);
    return m ? [m[1], m[2] ?? 'true'] : [a, 'true'];
  })
);
const ONLY_DATASET = args.dataset as string | undefined;

type Dataset = {
  slug: string;
  source: string;
  source_label: string;
  expectedRows: number;
};
// Only datasets sharing the standard LADBS schema (permit_nbr, primary_address, etc.)
// C-of-O (3f9m-afei) and BOE (j7mw-thyc) have different schemas — separate ingesters.
const DATASETS: Dataset[] = [
  { slug: '67is-svtd', source: 'la_city_mechanical', source_label: 'LADBS Mechanical Permits 2020+',     expectedRows: 279087 },
  { slug: '5m3t-xjex', source: 'la_city_mechanical', source_label: 'LADBS Mechanical Permits 2010-2019', expectedRows: 503281 },
  { slug: 'ysqd-apz7', source: 'la_city_electrical', source_label: 'LADBS Electrical Permits 2020+',     expectedRows: 338674 },
];

const PAGE_SIZE = 50000;
const BATCH = 1000;

type SocrataRow = {
  permit_nbr?: string;
  primary_address?: string;
  zip_code?: string;
  cd?: string;
  apn?: string;
  zone?: string;
  permit_group?: string;
  permit_type?: string;
  permit_sub_type?: string;
  use_code?: string;
  use_desc?: string;
  issue_date?: string;
  status_desc?: string;
  status_date?: string;
  valuation?: string;
  square_footage?: string;
  work_desc?: string;
  lat?: string;
  lon?: string;
};

function pn(v?: string): number | null { if (!v) return null; const n = parseFloat(v); return isFinite(n) ? n : null; }
function pi(v?: string): number | null { if (!v) return null; const n = parseInt(v, 10); return isFinite(n) ? n : null; }
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 canonicalize(addr: string): string {
  return addr.toLowerCase().replace(/[^\w\s-]/g, '').replace(/\s+/g, '-').replace(/-+/g, '-').replace(/^-|-$/g, '').slice(0, 90);
}

async function fetchPage(slug: string, offset: number): Promise<SocrataRow[]> {
  const u = `https://data.lacity.org/resource/${slug}.json?$limit=${PAGE_SIZE}&$offset=${offset}&$order=permit_nbr`;
  for (let attempt = 1; attempt <= 5; 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 SocrataRow[];
    } catch (e) {
      if (attempt >= 5) throw e;
      await new Promise(r => setTimeout(r, 1500 * attempt));
    }
  }
  return [];
}

async function upsertListings(rows: SocrataRow[], source: string): Promise<Map<string, string>> {
  const dedup = new Map<string, { addr: string; lat: number | null; lon: number | null; zip: string | null }>();
  const rawToCanonical = new Map<string, string>();
  for (const r of rows) {
    if (!r.primary_address) continue;
    const addr = r.primary_address.trim();
    const canon = canonicalize(addr);
    if (!canon) continue;
    rawToCanonical.set(addr, canon);
    if (!dedup.has(canon)) {
      dedup.set(canon, { addr, lat: pn(r.lat), lon: pn(r.lon), zip: r.zip_code ?? null });
    }
  }
  const canonToLid = new Map<string, string>();
  if (dedup.size === 0) return new Map();
  const arr = Array.from(dedup.entries());
  for (let i = 0; i < arr.length; i += BATCH) {
    const slice = arr.slice(i, i + BATCH);
    const tuples: string[] = [];
    const values: any[] = [];
    slice.forEach(([canonical, r], idx) => {
      const base = idx * 7;
      tuples.push(`($${base+1},$${base+2},$${base+3},$${base+4},$${base+5},$${base+6},$${base+7})`);
      values.push(`${canonical}-ladbs`, `ladbs:${canonical}`, r.addr, r.addr, r.zip, r.lat, r.lon);
    });
    const sql = `
      INSERT INTO listing
        (slug, source, source_id, title, address_line1, city, state, country, postal_code, latitude, longitude, is_public, tier)
      VALUES ${tuples.map((_, k) => {
        const b = k * 7;
        return `($${b+1},'ladbs_permit',$${b+2},$${b+3},$${b+4},'Los Angeles','CA','US',$${b+5},$${b+6},$${b+7},true,'free')`;
      }).join(',')}
      ON CONFLICT (slug) DO UPDATE SET
        latitude    = COALESCE(listing.latitude,    EXCLUDED.latitude),
        longitude   = COALESCE(listing.longitude,   EXCLUDED.longitude),
        postal_code = COALESCE(listing.postal_code, EXCLUDED.postal_code),
        updated_at  = now()
      RETURNING id, source_id`;
    const r = await pool.query<{ id: string; source_id: string }>(sql, values);
    for (const row of r.rows) canonToLid.set(row.source_id.replace(/^ladbs:/, ''), row.id);
  }
  const result = new Map<string, string>();
  for (const [raw, canonical] of rawToCanonical) {
    const lid = canonToLid.get(canonical);
    if (lid) result.set(raw, lid);
  }
  return result;
}

async function upsertPermits(ds: Dataset, rows: SocrataRow[], listingMap: Map<string, string>): Promise<number> {
  const seen = new Set<string>();
  const valid = rows.filter(r => r.permit_nbr && r.permit_nbr.trim() && !seen.has(r.permit_nbr.trim()) && (seen.add(r.permit_nbr.trim()), true));
  if (valid.length === 0) return 0;
  let inserted = 0;
  for (let i = 0; i < valid.length; i += BATCH) {
    const slice = valid.slice(i, i + BATCH);
    const tuples: string[] = [];
    const values: any[] = [];
    slice.forEach((r, idx) => {
      const base = idx * 18;
      tuples.push(`(${Array.from({length: 18}, (_, k) => `$${base+k+1}`).join(',')})`);
      values.push(
        ds.source, ds.slug, r.permit_nbr!.trim(),
        r.primary_address ? listingMap.get(r.primary_address.trim()) ?? null : null,
        r.primary_address ?? null, r.zip_code ?? null, r.cd ?? null, r.apn ?? null, r.zone ?? null,
        r.permit_group ?? null, r.permit_type ?? null, r.permit_sub_type ?? null,
        r.use_code ?? null, r.use_desc ?? null, pd(r.issue_date), pd(r.status_date),
        r.status_desc ?? null, pn(r.valuation),
      );
    });
    await pool.query(
      `INSERT INTO permit
         (source, source_dataset, permit_number, listing_id, primary_address, zip_code, council_district,
          apn, zone, permit_group, permit_type, permit_sub_type, use_code, use_desc,
          issue_date, status_date, status_desc, valuation)
       VALUES ${tuples.join(',')}
       ON CONFLICT (source_dataset, permit_number) DO UPDATE SET
         status_desc = COALESCE(EXCLUDED.status_desc, permit.status_desc),
         retrieved_at = now()`,
      values
    );
    inserted += slice.length;
  }
  // Set source_url + source_label after insert (single update, indexed by source_dataset)
  return inserted;
}

async function ingestDataset(ds: Dataset) {
  console.log(`\n=== ${ds.slug} ${ds.source_label} (${ds.expectedRows.toLocaleString()}) ===`);
  const t0 = Date.now();
  let offset = 0, pulled = 0, permitsIns = 0;
  while (true) {
    const rows = await fetchPage(ds.slug, offset);
    if (rows.length === 0) break;
    const listingMap = await upsertListings(rows, ds.source);
    const ins = await upsertPermits(ds, rows, listingMap);
    pulled += rows.length;
    permitsIns += ins;
    const dt = (Date.now() - t0) / 1000;
    console.log(`  ${ds.slug}: ${pulled.toLocaleString()}/${ds.expectedRows.toLocaleString()} (${(pulled/dt).toFixed(0)}/s)`);
    if (rows.length < PAGE_SIZE) break;
    offset += PAGE_SIZE;
  }
  // Set source_url + label for this dataset
  await pool.query(
    `UPDATE permit SET source_url = $1, source_label = $2 WHERE source_dataset = $3 AND source_url IS NULL`,
    [`https://data.lacity.org/resource/${ds.slug}.json?permit_nbr=`, ds.source_label, ds.slug]
  );
  // Source URL with permit number suffix needs per-row interpolation
  await pool.query(
    `UPDATE permit SET source_url = source_url || permit_number WHERE source_dataset = $1 AND source_url LIKE '%permit_nbr=' AND source_url NOT LIKE '%permit_nbr=' || permit_number`,
    [ds.slug]
  );
  console.log(`  ✓ ${ds.slug}: ${pulled.toLocaleString()} pulled, ${permitsIns.toLocaleString()} permits inserted`);
}

async function main() {
  const datasets = ONLY_DATASET ? DATASETS.filter(d => d.slug === ONLY_DATASET) : DATASETS;
  for (const ds of datasets) await ingestDataset(ds);
  await pool.end();
}

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