← back to Professional Directory

agents/ingest-agent/importers/mbc.js

249 lines

#!/usr/bin/env node
/**
 * California Medical Board (MBC) license importer — Stage 2.
 *
 *  Source: https://www.mbc.ca.gov/Resources/Statistics/Public-Information.aspx
 *    "Profession Specific Public Information File" (free, weekly refresh).
 *    Currently published as Excel (.xlsx) — historically Access (.mdb).
 *
 * This importer reads a CSV that the user has converted from the MBC bulk file
 * to:
 *      data/mbc/mbc_licenses.csv
 * (Use Excel "Save As CSV" or `mdb-export <file.mdb> licenses > mbc_licenses.csv`.)
 *
 * Expected columns (MBC uses several historical schemas — we adapt to whichever
 * is present):
 *      LICENSE TYPE / License Type
 *      LICENSE NUMBER / License #
 *      LICENSE STATUS / Status
 *      ISSUE DATE
 *      EXPIRATION DATE
 *      LAST NAME
 *      FIRST NAME
 *      MIDDLE NAME
 *      SCHOOL NAME / SCHOOL OF GRADUATION
 *      GRADUATION YEAR
 *      ADDRESS LINE 1
 *      CITY
 *      STATE
 *      ZIP
 *
 * Match strategy (priority): license_number → name+last+CA license → fuzzy
 * Updates professionals: license_status, license_issue_date, license_expiration_date,
 * medical_school, graduation_year, license_type. Bumps source_confidence_score.
 */
const fs = require('node:fs');
const path = require('node:path');
const crypto = require('node:crypto');
const { parse } = require('csv-parse');
const { pool, query, withTx } = require('../../shared/db');

const SOURCE_NAME = 'CA Medical Board (MBC) — Access DB';
const SOURCE_URL  = 'https://www.mbc.ca.gov/Resources/Statistics/Public-Information.aspx';
const INPUT_CSV   = process.env.MBC_CSV
  || path.resolve(__dirname, '../../../data/mbc/mbc_licenses.csv');

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

function pick(row, ...keys) {
  for (const k of keys) {
    const v = row[k];
    if (v !== undefined && v !== null && String(v).trim() !== '') return clean(v);
  }
  return null;
}

function toIsoDate(s) {
  if (!s) return null;
  const us = s.match(/^(\d{1,2})\/(\d{1,2})\/(\d{4})$/);
  if (us) return `${us[3]}-${us[1].padStart(2,'0')}-${us[2].padStart(2,'0')}`;
  const iso = s.match(/^(\d{4})-(\d{2})-(\d{2})$/);
  return iso ? s : null;
}

function toYear(s) {
  if (!s) return null;
  const m = String(s).match(/(\d{4})/);
  return m ? Number(m[1]) : null;
}

function sha256(s) { return crypto.createHash('sha256').update(s).digest('hex'); }

async function getSourceId() {
  const r = await query('SELECT id FROM sources WHERE source_name = $1', [SOURCE_NAME]);
  if (r.rowCount === 0) throw new Error(`source not seeded: ${SOURCE_NAME}`);
  return r.rows[0].id;
}

async function startJob(sourceId, label) {
  const r = await query(
    `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, fields) {
  const sets = []; const params = []; 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);
}

async function findProfessional(client, mbc) {
  // Priority 1: license_number direct match (already populated by NPI's CA license).
  if (mbc.license_number) {
    const r = await client.query(
      `SELECT id FROM professionals WHERE license_number = $1 LIMIT 2`,
      [mbc.license_number]
    );
    if (r.rowCount === 1) return r.rows[0].id;
  }
  // Priority 2: name + license_number fuzzy.
  if (mbc.last && mbc.first && mbc.license_number) {
    const r = await client.query(
      `SELECT id FROM professionals
        WHERE LOWER(last_name)  = LOWER($1)
          AND LOWER(first_name) = LOWER($2)
          AND (license_number = $3 OR license_number IS NULL)
        LIMIT 2`,
      [mbc.last, mbc.first, mbc.license_number]
    );
    if (r.rowCount === 1) return r.rows[0].id;
  }
  // Priority 3: exact name match if uniquely identifying in DB.
  if (mbc.last && mbc.first) {
    const r = await client.query(
      `SELECT id FROM professionals
        WHERE LOWER(last_name)  = LOWER($1)
          AND LOWER(first_name) = LOWER($2)
        LIMIT 2`,
      [mbc.last, mbc.first]
    );
    if (r.rowCount === 1) return r.rows[0].id;
  }
  return null;
}

async function applyMbc(client, professionalId, mbc, sourceId) {
  await client.query(`
    UPDATE professionals SET
      license_number          = COALESCE($2, license_number),
      license_type            = COALESCE($3, license_type),
      license_status          = COALESCE($4, license_status),
      license_issue_date      = COALESCE($5, license_issue_date),
      license_expiration_date = COALESCE($6, license_expiration_date),
      medical_school          = COALESCE($7, medical_school),
      graduation_year         = COALESCE($8, graduation_year),
      source_confidence_score = GREATEST(COALESCE(source_confidence_score,0), 0.85),
      updated_at              = NOW()
    WHERE id = $1
  `, [
    professionalId,
    mbc.license_number,
    mbc.license_type,
    mbc.license_status,
    mbc.issue_date,
    mbc.expiration_date,
    mbc.school,
    mbc.graduation_year,
  ]);

  // Provenance.
  const json = JSON.stringify(mbc);
  const hash = sha256(json + '|professional|' + professionalId);
  await client.query(`
    INSERT INTO raw_records (source_id, source_url, entity_type, entity_id, raw_json, hash)
    VALUES ($1,$2,'professional',$3,$4::jsonb,$5)
    ON CONFLICT (hash) DO NOTHING
  `, [sourceId, SOURCE_URL, professionalId, json, hash]);
}

function rowToMbc(row) {
  return {
    license_type:    pick(row, 'LICENSE TYPE', 'License Type'),
    license_number:  pick(row, 'LICENSE NUMBER', 'License #', 'License Number'),
    license_status:  pick(row, 'LICENSE STATUS', 'Status', 'License Status'),
    issue_date:      toIsoDate(pick(row, 'ISSUE DATE', 'Issue Date')),
    expiration_date: toIsoDate(pick(row, 'EXPIRATION DATE', 'Expiration Date')),
    last:            pick(row, 'LAST NAME', 'Last Name'),
    first:           pick(row, 'FIRST NAME', 'First Name'),
    middle:          pick(row, 'MIDDLE NAME', 'Middle Name'),
    school:          pick(row, 'SCHOOL NAME', 'SCHOOL OF GRADUATION', 'School Name', 'School'),
    graduation_year: toYear(pick(row, 'GRADUATION YEAR', 'Graduation Year', 'GRAD YEAR')),
    address1:        pick(row, 'ADDRESS LINE 1', 'Address Line 1', 'ADDRESS'),
    city:            pick(row, 'CITY', 'City'),
    state:           pick(row, 'STATE', 'State'),
    zip:             pick(row, 'ZIP', 'Zip', 'ZIP CODE'),
  };
}

async function main() {
  if (!fs.existsSync(INPUT_CSV)) {
    console.error(`[mbc] input missing: ${INPUT_CSV}
Steps to provide it (free, no key, no money):
  1. Visit ${SOURCE_URL}
  2. Download the "Profession Specific Public Information" file for the
     Medical Board (License Type: M.D., D.O., or full set).
  3. If it's .xlsx → open in Excel/Numbers, Save As CSV → ${INPUT_CSV}
     If it's .mdb  → 'brew install mdb-tools' then
         mdb-export your_file.mdb LICENSES > ${INPUT_CSV}`);
    process.exit(2);
  }

  const sourceId = await getSourceId();
  const jobId = await startJob(sourceId, 'mbc:bulk');

  let scanned = 0, matched = 0, updated = 0, unmatched = 0;

  const parser = fs.createReadStream(INPUT_CSV).pipe(parse({
    columns: true, skip_empty_lines: true, bom: true, relax_column_count: true,
  }));

  for await (const row of parser) {
    scanned++;
    const mbc = rowToMbc(row);
    if (!mbc.license_number && !(mbc.last && mbc.first)) continue;

    try {
      await withTx(async (client) => {
        const id = await findProfessional(client, mbc);
        if (!id) { unmatched++; return; }
        matched++;
        await applyMbc(client, id, mbc, sourceId);
        updated++;
      });
    } catch (e) {
      console.error(`[mbc] row error: ${e.message}`);
    }

    if (scanned % 5000 === 0) {
      console.log(`[mbc] scanned=${scanned} matched=${matched} unmatched=${unmatched}`);
    }
  }

  await finishJob(jobId, {
    status: 'done',
    records_found: scanned,
    records_inserted: 0,
    records_updated: updated,
    records_skipped: unmatched,
  });

  console.log(`[mbc] done. scanned=${scanned} matched=${matched} updated=${updated} unmatched=${unmatched}`);
  await pool.end();
}

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