← back to Stayclaim

scripts/ingest-la-city-permits.ts

392 lines

/**
 * ingest-la-city-permits.ts
 *
 * Pulls all 1,560,839 LA City Department of Building & Safety building permits
 * from the three Socrata datasets on data.lacity.org:
 *
 *   - pi9x-tg5x  (2020→present, ~387K rows, daily refresh)  ← canonical replacement for stale xnhu-aczu
 *   - dyxf-7hc4  (2010-2019,    ~533K rows, static)
 *   - e67z-kt2n  (<2010,        ~640K rows, static)
 *
 * Verified field shape (consistent across all 3): permit_nbr, primary_address,
 * zip_code, cd, apn, zone, permit_group, permit_type, permit_sub_type,
 * use_code, use_desc, issue_date, status_desc, valuation, lat, lon, plus
 * dataset-specific extras (square_footage on dyxf+e67z, construction on e67z).
 *
 * For each row:
 *   1. Upsert a `listing` row (source='ladbs_permit', source_id=primary_address) —
 *      one listing per distinct LA-City address, geocoded.
 *   2. Insert a `permit` row linked to that listing.
 *
 * Idempotency: ON CONFLICT (source_dataset, permit_number) DO UPDATE for permits;
 * ON CONFLICT (source, source_id) DO UPDATE COALESCE for listings.
 *
 * Batched 1000 rows per upsert; ~50K-row Socrata pages. Resume-safe via
 * --start-page and --dataset flags.
 *
 * Tier: A (LA City government primary record).
 *
 * Usage: npx tsx scripts/ingest-la-city-permits.ts
 *          [--dataset=pi9x-tg5x|dyxf-7hc4|e67z-kt2n|all]   default: all
 *          [--start-page=N]                                 default: 0
 *          [--page-size=50000]                              default: 50000
 *          [--app-token=XXXXXXXX]                           optional Socrata token
 */
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;
const START_PAGE   = args['start-page'] ? parseInt(args['start-page'] as string, 10) : 0;
const PAGE_SIZE    = args['page-size'] ? parseInt(args['page-size'] as string, 10) : 50000;
const APP_TOKEN    = (args['app-token'] as string) || process.env.SOCRATA_APP_TOKEN || undefined;
const BATCH_INSERT_SIZE = 1000;

type Dataset = { slug: string; expectedRows: number };
const DATASETS: Dataset[] = [
  { slug: 'pi9x-tg5x', expectedRows: 387531 },
  { slug: 'dyxf-7hc4', expectedRows: 533366 },
  { slug: 'e67z-kt2n', expectedRows: 639942 },
];

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;
  submitted_date?: string;
  valuation?: string;
  lat?: string;
  lon?: string;
  square_footage?: string;
  work_desc?: string;
  cnc?: string;
  cpa?: string;
};

async function ensureSchema() {
  await pool.query(`
    CREATE TABLE IF NOT EXISTS permit (
      id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
      source TEXT NOT NULL DEFAULT 'ladbs',
      source_dataset TEXT NOT NULL,
      permit_number TEXT NOT NULL,
      listing_id UUID REFERENCES listing(id),
      primary_address TEXT,
      zip_code TEXT,
      council_district TEXT,
      apn TEXT,
      zone TEXT,
      permit_group TEXT,
      permit_type TEXT,
      permit_sub_type TEXT,
      use_code TEXT,
      use_desc TEXT,
      issue_date DATE,
      status_date DATE,
      submitted_date DATE,
      status_desc TEXT,
      valuation NUMERIC(15,2),
      square_footage INT,
      work_desc TEXT,
      neighborhood_council TEXT,
      community_plan_area TEXT,
      latitude NUMERIC(9,6),
      longitude NUMERIC(9,6),
      source_url TEXT,
      source_tier CHAR(1) NOT NULL DEFAULT 'A',
      source_label TEXT NOT NULL DEFAULT 'LA City Department of Building & Safety',
      retrieved_at TIMESTAMPTZ DEFAULT now(),
      UNIQUE (source_dataset, permit_number)
    );
    CREATE INDEX IF NOT EXISTS idx_permit_listing ON permit(listing_id);
    CREATE INDEX IF NOT EXISTS idx_permit_address ON permit(primary_address);
    CREATE INDEX IF NOT EXISTS idx_permit_apn ON permit(apn);
    CREATE INDEX IF NOT EXISTS idx_permit_zip_date ON permit(zip_code, issue_date);
    CREATE INDEX IF NOT EXISTS idx_permit_issue_date ON permit(issue_date);
    CREATE INDEX IF NOT EXISTS idx_permit_council_district ON permit(council_district);
  `);
}

function parseNum(v?: string): number | null {
  if (!v) return null;
  const n = parseFloat(v);
  return isFinite(n) ? n : null;
}

function parseInt2(v?: string): number | null {
  if (!v) return null;
  const n = parseInt(v, 10);
  return isFinite(n) ? n : null;
}

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

function buildSourceUrl(slug: string, permitNbr: string): string {
  // Direct Socrata row-filter URL: dataset page, filtered to that one permit.
  // This lets a user click through and see the canonical record.
  return `https://data.lacity.org/resource/${slug}.json?permit_nbr=${encodeURIComponent(permitNbr)}`;
}

/**
 * Fetch one Socrata page. Returns rows or empty array if past the end.
 */
async function fetchPage(slug: string, offset: number, limit: number): Promise<SocrataRow[]> {
  const url = new URL(`https://data.lacity.org/resource/${slug}.json`);
  url.searchParams.set('$limit', String(limit));
  url.searchParams.set('$offset', String(offset));
  url.searchParams.set('$order', 'permit_nbr');
  if (APP_TOKEN) url.searchParams.set('$$app_token', APP_TOKEN);
  for (let attempt = 1; attempt <= 5; attempt++) {
    try {
      const res = await fetch(url, { 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} ${res.statusText}`);
      return (await res.json()) as SocrataRow[];
    } catch (e) {
      if (attempt >= 5) throw e;
      await new Promise(r => setTimeout(r, 1500 * attempt));
    }
  }
  return [];
}

/**
 * Bulk-upsert listings for a batch of (address, lat, lon, zip) tuples.
 * Dedupes by primary_address (the source_id).
 * Returns a map of primary_address → listing_id.
 */
function canonicalize(addr: string): string {
  return addr
    .toLowerCase()
    .replace(/[^\w\s-]/g, '')
    .replace(/\s+/g, '-')
    .replace(/-+/g, '-')
    .replace(/^-|-$/g, '')
    .slice(0, 90);
}

async function upsertListings(rows: SocrataRow[]): Promise<Map<string, string>> {
  // Dedup by canonical (slug-equivalent) form so within-batch collisions are
  // collapsed before the INSERT. Map original raw address → canonical so
  // upsertPermits can look up listing_id by primary_address.
  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();
    if (!addr) continue;
    const canonical = canonicalize(addr);
    if (!canonical) continue;
    rawToCanonical.set(addr, canonical);
    if (dedup.has(canonical)) {
      // merge: prefer non-null lat/lon/zip
      const ex = dedup.get(canonical)!;
      const lat = parseNum(r.lat);
      const lon = parseNum(r.lon);
      if (ex.lat == null && lat != null) ex.lat = lat;
      if (ex.lon == null && lon != null) ex.lon = lon;
      if (!ex.zip && r.zip_code) ex.zip = r.zip_code;
      continue;
    }
    dedup.set(canonical, {
      addr,
      lat: parseNum(r.lat),
      lon: parseNum(r.lon),
      zip: r.zip_code ?? null,
    });
  }
  // canonical → listing.id
  const canonicalToListing = 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_INSERT_SIZE) {
    const slice = arr.slice(i, i + BATCH_INSERT_SIZE);
    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})`);
      const slug = `${canonical}-ladbs`;
      const sourceId = `ladbs:${canonical}`;
      values.push(
        slug,           // 1 slug
        sourceId,       // 2 source_id
        r.addr,         // 3 title
        r.addr,         // 4 address_line1
        r.zip,          // 5 postal_code
        r.lat,          // 6 latitude
        r.lon,          // 7 longitude
      );
    });
    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((t, i) => {
        const base = i * 7;
        return `($${base+1},'ladbs_permit',$${base+2},$${base+3},$${base+4},'Los Angeles','CA','US',$${base+5},$${base+6},$${base+7},true,'free')`;
      }).join(',')}
      ON CONFLICT (source, source_id) 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) {
      const canonical = row.source_id.replace(/^ladbs:/, '');
      canonicalToListing.set(canonical, row.id);
    }
  }
  // Resolve raw → listing via canonical
  const result = new Map<string, string>();
  for (const [raw, canonical] of rawToCanonical) {
    const lid = canonicalToListing.get(canonical);
    if (lid) result.set(raw, lid);
  }
  return result;
}

async function upsertPermits(slug: string, rows: SocrataRow[], listingMap: Map<string, string>): Promise<number> {
  const valid = rows.filter(r => r.permit_nbr && r.permit_nbr.trim());
  if (valid.length === 0) return 0;
  let inserted = 0;
  for (let i = 0; i < valid.length; i += BATCH_INSERT_SIZE) {
    const slice = valid.slice(i, i + BATCH_INSERT_SIZE);
    const tuples: string[] = [];
    const values: any[] = [];
    slice.forEach((r, idx) => {
      const base = idx * 21;
      tuples.push(`(${Array.from({length: 21}, (_, k) => `$${base+k+1}`).join(',')})`);
      values.push(
        slug,                                                   // 1 source_dataset
        r.permit_nbr!.trim(),                                   // 2 permit_number
        r.primary_address ? listingMap.get(r.primary_address.trim()) ?? null : null, // 3 listing_id
        r.primary_address ?? null,                              // 4 primary_address
        r.zip_code ?? null,                                     // 5 zip_code
        r.cd ?? null,                                           // 6 council_district
        r.apn ?? null,                                          // 7 apn
        r.zone ?? null,                                         // 8 zone
        r.permit_group ?? null,                                 // 9 permit_group
        r.permit_type ?? null,                                  // 10 permit_type
        r.permit_sub_type ?? null,                              // 11 permit_sub_type
        r.use_code ?? null,                                     // 12 use_code
        r.use_desc ?? null,                                     // 13 use_desc
        parseDate(r.issue_date),                                // 14 issue_date
        parseDate(r.status_date),                               // 15 status_date
        parseDate(r.submitted_date),                            // 16 submitted_date
        r.status_desc ?? null,                                  // 17 status_desc
        parseNum(r.valuation),                                  // 18 valuation
        parseInt2(r.square_footage),                            // 19 square_footage
        r.work_desc ?? null,                                    // 20 work_desc
        buildSourceUrl(slug, r.permit_nbr!.trim()),             // 21 source_url
      );
    });
    const sql = `
      INSERT INTO permit
        (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, submitted_date, status_desc, valuation, square_footage,
         work_desc, source_url)
      VALUES ${tuples.join(',')}
      ON CONFLICT (source_dataset, permit_number) DO UPDATE SET
        listing_id     = COALESCE(EXCLUDED.listing_id, permit.listing_id),
        status_date    = COALESCE(EXCLUDED.status_date, permit.status_date),
        status_desc    = COALESCE(EXCLUDED.status_desc, permit.status_desc),
        retrieved_at   = now()
    `;
    const r = await pool.query(sql, values);
    inserted += r.rowCount ?? 0;
  }
  return inserted;
}

async function ingestDataset(ds: Dataset) {
  console.log(`\n=== ${ds.slug} (~${ds.expectedRows.toLocaleString()} rows) ===`);
  const t0 = Date.now();
  let offset = START_PAGE * PAGE_SIZE;
  let pulled = 0, listingsTouched = 0, permitsInserted = 0;
  while (true) {
    const tFetch = Date.now();
    const rows = await fetchPage(ds.slug, offset, PAGE_SIZE);
    if (rows.length === 0) break;
    const listingMap = await upsertListings(rows);
    const inserted = await upsertPermits(ds.slug, rows, listingMap);
    pulled += rows.length;
    listingsTouched += listingMap.size;
    permitsInserted += inserted;
    const dt = (Date.now() - t0) / 1000;
    const dtFetch = (Date.now() - tFetch) / 1000;
    console.log(`  page off=${offset} rows=${rows.length} listings=${listingMap.size} permits=${inserted} (page=${dtFetch.toFixed(1)}s, total ${pulled.toLocaleString()}/${ds.expectedRows.toLocaleString()} @ ${(pulled/dt).toFixed(0)}/s)`);
    if (rows.length < PAGE_SIZE) break;
    offset += PAGE_SIZE;
  }
  const dt = (Date.now() - t0) / 1000;
  console.log(`  ✓ ${ds.slug}: pulled=${pulled.toLocaleString()} listings=${listingsTouched.toLocaleString()} permits=${permitsInserted.toLocaleString()} elapsed=${(dt/60).toFixed(1)}m`);
}

async function main() {
  await ensureSchema();
  const datasets = ONLY_DATASET && ONLY_DATASET !== 'all'
    ? DATASETS.filter(d => d.slug === ONLY_DATASET)
    : DATASETS;
  if (datasets.length === 0) {
    console.error(`Unknown dataset: ${ONLY_DATASET}`);
    process.exit(1);
  }
  const t0 = Date.now();
  for (const ds of datasets) {
    await ingestDataset(ds);
  }
  const dt = (Date.now() - t0) / 1000;
  console.log(`\n=== ALL DONE in ${(dt/60).toFixed(1)}m ===`);
  const { rows } = await pool.query<{ ds: string; n: string }>(
    `SELECT source_dataset as ds, count(*)::text as n FROM permit GROUP BY source_dataset ORDER BY source_dataset`
  );
  for (const r of rows) console.log(`  ${r.ds}: ${parseInt(r.n).toLocaleString()} permits`);
  const { rows: lrows } = await pool.query<{ n: string }>(
    `SELECT count(*)::text as n FROM listing WHERE source = 'ladbs_permit'`
  );
  console.log(`  listings: ${parseInt(lrows[0].n).toLocaleString()}`);
  await pool.end();
}

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