← back to Lawyer Directory Builder

src/ingest/wikidata.ts

235 lines

/**
 * Wikidata SPARQL importer — pulls notable law firms with HQ or office in LA County.
 * Free, no auth. Captures the big-name firms that OSM under-represents (Latham,
 * Gibson Dunn, Paul Hastings, Manatt, Lewis Brisbois, Littler Mendelson, etc.).
 *
 * Endpoint: https://query.wikidata.org/sparql
 * License:  CC0 (Wikidata data) — fully reusable.
 */
import 'dotenv/config';
import crypto from 'node:crypto';
import { fetch } from 'undici';
import { pool, query, withTx } from '../db/pool.ts';

const SOURCE_NAME = 'Wikidata (LA County law firms)';
const ENDPOINT = 'https://query.wikidata.org/sparql';
const USER_AGENT = process.env.USER_AGENT
  || 'LawyerDirectoryBuilder/0.1 (research; contact: steveabramsdesigns@gmail.com)';

const SPARQL = `
SELECT DISTINCT ?firm ?firmLabel ?hq ?hqLabel ?street ?coord ?inception ?website ?employees WHERE {
  ?firm wdt:P31/wdt:P279* wd:Q613142 .
  OPTIONAL { ?firm wdt:P159 ?hq . }
  OPTIONAL { ?firm wdt:P969 ?street . }
  OPTIONAL { ?firm wdt:P625 ?coord . }
  OPTIONAL { ?firm wdt:P571 ?inception . }
  OPTIONAL { ?firm wdt:P856 ?website . }
  OPTIONAL { ?firm wdt:P1128 ?employees . }
  {
    ?firm wdt:P159 ?hq2 . ?hq2 wdt:P131* wd:Q104994 .
  } UNION {
    ?firm wdt:P276 ?loc2 . ?loc2 wdt:P131* wd:Q104994 .
  }
  SERVICE wikibase:label { bd:serviceParam wikibase:language "en". }
}
LIMIT 500
`;

async function ensureSource(): Promise<number> {
  await query(`
    INSERT INTO sources (source_name, source_type, base_url, terms_notes, allowed_method, rate_limit_rps)
    VALUES ($1, 'api', 'https://query.wikidata.org/sparql', 'Wikidata SPARQL — CC0. Polite-use; cap ~1 query / 5 s.', 'api', 0.20)
    ON CONFLICT (source_name) DO NOTHING
  `, [SOURCE_NAME]);
  const r = await query<{ id: number }>(`SELECT id FROM sources WHERE source_name = $1`, [SOURCE_NAME]);
  return r.rows[0].id;
}

async function startJob(sourceId: number, label: string) {
  const r = await query<{ id: number }>(`
    INSERT INTO scrape_jobs (source_id, job_label, status, started_at)
    VALUES ($1, $2, 'running', NOW()) RETURNING id
  `, [sourceId, label]);
  return r.rows[0].id;
}

async function finishJob(jobId: number, fields: Record<string, unknown>) {
  const sets: string[] = [];
  const params: unknown[] = [];
  let i = 1;
  for (const [k, v] of Object.entries(fields)) {
    sets.push(`${k} = $${i++}`); params.push(v);
  }
  sets.push(`finished_at = NOW()`);
  params.push(jobId);
  await query(`UPDATE scrape_jobs SET ${sets.join(', ')} WHERE id = $${i}`, params);
}

function clean(s: string | undefined | null) {
  if (!s) return null;
  const t = String(s).trim();
  return t.length === 0 ? null : t;
}

function normAddress(s: string | null) {
  if (!s) return null;
  return s.toLowerCase().replace(/[.,]/g, ' ').replace(/\b(suite|ste|unit|apt|#)\b/g, '').replace(/\s+/g, ' ').trim();
}

function parseCoord(p: string | undefined) {
  if (!p) return { lat: null, lng: null };
  const m = p.match(/Point\(\s*(-?\d+\.?\d*)\s+(-?\d+\.?\d*)\s*\)/);
  if (!m) return { lat: null, lng: null };
  return { lat: parseFloat(m[2]), lng: parseFloat(m[1]) };
}

interface SparqlBinding {
  firm?: { value: string };
  firmLabel?: { value: string };
  hqLabel?: { value: string };
  street?: { value: string };
  coord?: { value: string };
  inception?: { value: string };
  website?: { value: string };
  employees?: { value: string };
}

async function fetchSparql(): Promise<SparqlBinding[]> {
  const url = `${ENDPOINT}?query=${encodeURIComponent(SPARQL)}`;
  const r = await fetch(url, {
    headers: {
      'User-Agent': USER_AGENT,
      Accept: 'application/sparql-results+json',
    },
    signal: AbortSignal.timeout(60000),
  });
  if (!r.ok) {
    const t = await r.text();
    throw new Error(`Wikidata ${r.status}: ${t.slice(0, 300)}`);
  }
  const j = await r.json() as { results: { bindings: SparqlBinding[] } };
  return j.results.bindings;
}

async function upsertFirm(b: SparqlBinding, sourceId: number) {
  const name = clean(b.firmLabel?.value);
  if (!name || /^Q\d+$/.test(name)) return null;          // skip Q-IDs (no English label)

  const wdQid = b.firm?.value?.split('/').pop() || null;
  const hq = clean(b.hqLabel?.value);
  const street = clean(b.street?.value);
  const website = clean(b.website?.value);
  const { lat, lng } = parseCoord(b.coord?.value);
  const employees = b.employees?.value ? parseInt(b.employees.value, 10) : null;

  const fullAddress = street && hq ? `${street}, ${hq}, CA` : (street || hq || null);
  const sourceUrl = b.firm?.value || 'https://www.wikidata.org/';

  return await withTx(async (client) => {
    // Try to merge with existing org by website (BigLaw firms usually have unique websites).
    let orgId: number | null = null;
    if (website) {
      const r = await client.query<{ id: number }>(
        `SELECT id FROM organizations WHERE LOWER(website) = LOWER($1) LIMIT 1`, [website]);
      if (r.rowCount && r.rowCount > 0) orgId = r.rows[0].id;
    }
    if (!orgId) {
      // Match by exact name (case-insensitive)
      const r = await client.query<{ id: number }>(
        `SELECT id FROM organizations WHERE LOWER(name) = LOWER($1) AND county = 'Los Angeles' LIMIT 1`, [name]);
      if (r.rowCount && r.rowCount > 0) orgId = r.rows[0].id;
    }

    const sizeBand = employees && employees >= 500 ? 'biglaw'
      : employees && employees >= 100 ? 'large'
      : employees && employees >= 25 ? 'medium'
      : null;

    if (orgId) {
      await client.query(`
        UPDATE organizations
        SET name = $2,
            website = COALESCE(website, $3),
            address = COALESCE(address, $4),
            address_norm = COALESCE(address_norm, $5),
            city = COALESCE(city, $6),
            neighborhood = COALESCE(neighborhood, $6),
            lat = COALESCE(lat, $7),
            lng = COALESCE(lng, $8),
            firm_size_band = COALESCE($9, firm_size_band),
            attorney_count = GREATEST(attorney_count, COALESCE($10, 0)),
            source_url = COALESCE(source_url, $11),
            updated_at = NOW()
        WHERE id = $1
      `, [orgId, name, website, fullAddress, normAddress(fullAddress), hq, lat, lng, sizeBand, employees, sourceUrl]);
    } else {
      const r = await client.query<{ id: number }>(`
        INSERT INTO organizations (
          name, type, address, address_norm, city, neighborhood, state, county,
          lat, lng, geocoded_at, website, firm_size_band, attorney_count, source_url
        ) VALUES (
          $1,'law_firm',$2,$3,$4,$4,'CA','Los Angeles',
          $5::double precision, $6::double precision,
          CASE WHEN $5::double precision IS NOT NULL THEN NOW() END,
          $7, $8, COALESCE($9::int, 0), $10
        )
        RETURNING id
      `, [name, fullAddress, normAddress(fullAddress), hq, lat, lng, website, sizeBand, employees, sourceUrl]);
      orgId = r.rows[0].id;
    }

    // Provenance
    const rawJson = JSON.stringify(b);
    const hash = crypto.createHash('sha256').update(rawJson + '|wikidata|' + (wdQid || name) + '|' + orgId).digest('hex');
    await client.query(`
      INSERT INTO raw_records (source_id, source_url, entity_type, entity_id, raw_json, fetched_at, hash)
      VALUES ($1,$2,'organization',$3,$4::jsonb, NOW(),$5)
      ON CONFLICT (source_id, hash) DO NOTHING
    `, [sourceId, sourceUrl, orgId, rawJson, hash]);

    return orgId;
  });
}

async function main() {
  console.log('[wikidata] querying SPARQL for law firms with LA County HQ/office…');
  const sourceId = await ensureSource();
  const jobId = await startJob(sourceId, 'wikidata:la-county-firms');

  let bindings: SparqlBinding[] = [];
  try {
    bindings = await fetchSparql();
  } catch (e) {
    await finishJob(jobId, { status: 'failed', error_message: (e as Error).message });
    throw e;
  }
  console.log(`[wikidata] received ${bindings.length} bindings`);

  let inserted = 0, skipped = 0;
  for (const b of bindings) {
    try {
      const id = await upsertFirm(b, sourceId);
      if (id) inserted++; else skipped++;
    } catch (e) {
      console.error(`[wikidata] upsert err: ${(e as Error).message}`);
      skipped++;
    }
  }

  await finishJob(jobId, {
    status: 'completed',
    records_found: bindings.length,
    records_inserted: inserted,
    records_skipped: skipped,
  });

  console.log(`[wikidata] done. seen=${bindings.length} kept=${inserted} skipped=${skipped}`);
  await pool.end();
}

main().catch(async (err) => {
  console.error('[wikidata] fatal:', err);
  try { await pool.end(); } catch {}
  process.exit(1);
});