// Philly loader common module. Parallel to scripts/loaders/_lib.ts (NYC) // but writes to the philly_* tables. Decisions: docs/cities/philly/00_overview.md. // // Provenance records (sha256 + bytes) come from the discovered file in // data/cities/philly/sources///provenance.json — written by // pipeline_philly/discover.py — so loaders can re-use them rather than // re-hashing on every load. import "dotenv/config"; import { config as loadEnv } from "dotenv"; import path from "node:path"; import crypto from "node:crypto"; import fs from "node:fs"; import { execSync } from "node:child_process"; loadEnv({ path: path.join(process.cwd(), ".env.local") }); // Same Postgres binding as NYC loaders — DIRECT_URL bypasses RLS. const direct = process.env.DIRECT_URL ?? process.env.POSTGRES_URL_NON_POOLING; if (direct) { process.env.DATABASE_URL = direct; } import { prisma } from "@/prisma/client"; export { prisma }; export function gitSha(): string | null { try { return execSync("git rev-parse HEAD", { stdio: ["ignore", "pipe", "ignore"] }) .toString() .trim(); } catch { return null; } } export type Provenance = { source_id: string; url: string; retrieved_at: string; sha256: string; bytes: number; publisher?: string; format?: string; }; /** * Read the latest provenance.json for a source_id from * data/cities/philly/sources///provenance.json * Returns null if no fetch has run yet. */ export function readLatestProvenance(sourceId: string): Provenance | null { const dir = path.join(process.cwd(), "data", "cities", "philly", "sources", sourceId); if (!fs.existsSync(dir)) return null; const dated = fs .readdirSync(dir) .filter((n) => /^\d{4}-\d{2}-\d{2}$/.test(n)) .sort() .reverse(); for (const d of dated) { const file = path.join(dir, d, "provenance.json"); if (fs.existsSync(file)) { return JSON.parse(fs.readFileSync(file, "utf-8")) as Provenance; } } return null; } /** * Path to the latest downloaded source file (the xlsx/csv/zip living alongside * provenance.json). The download step keeps the filename as the URL's last * segment; the loader needs to know that path. */ export function readLatestSourcePath(sourceId: string): string | null { const dir = path.join(process.cwd(), "data", "cities", "philly", "sources", sourceId); if (!fs.existsSync(dir)) return null; const dated = fs .readdirSync(dir) .filter((n) => /^\d{4}-\d{2}-\d{2}$/.test(n)) .sort() .reverse(); for (const d of dated) { const here = path.join(dir, d); const entries = fs.readdirSync(here).filter((n) => n !== "provenance.json"); if (entries.length > 0) return path.join(here, entries[0]); } return null; } export async function startDataLoad(opts: { sourceId: string; // matches docs/cities/philly/sources.yaml source_id sourceUrl?: string; yearCovered?: string; scriptPath: string; sha256?: string; bytes?: number; notes?: string; }) { return prisma.phillyDataLoad.create({ data: { id: crypto.randomUUID(), sourceId: opts.sourceId, sourceUrl: opts.sourceUrl, sha256: opts.sha256, bytes: opts.bytes, yearCovered: opts.yearCovered, scriptPath: opts.scriptPath, gitSha: gitSha(), notes: opts.notes, }, }); } export async function finishDataLoad( id: string, counts: { inserted: number; updated: number; suppressed: number } ) { await prisma.phillyDataLoad.update({ where: { id }, data: { rowsInserted: counts.inserted, rowsUpdated: counts.updated, rowsSuppressed: counts.suppressed, }, }); } /** * Run an async fn against batches of `items` of size `batchSize`, with bounded * concurrency. Identical to the NYC `batched` helper. */ export async function batched( items: T[], batchSize: number, fn: (batch: T[], batchIndex: number) => Promise, concurrency = 4, ): Promise { const batches: T[][] = []; for (let i = 0; i < items.length; i += batchSize) { batches.push(items.slice(i, i + batchSize)); } let next = 0; async function worker() { while (next < batches.length) { const idx = next++; await fn(batches[idx], idx); } } await Promise.all(Array.from({ length: concurrency }, worker)); } /** * Year semantics for state (Future Ready PA Index) indicators. * * PDE's Future Ready glossary defines Regular/Persistent Attendance as "a * lagging indicator indicating data is from the year prior to the reporting * year": the attendance element in the 2024-25 workbook describes 2023-24 * attendance. Every served row is keyed by the year the students were * actually in school (the OBSERVATION year), so this metric is stored under * report year − 1. Other indicators keep the workbook's report year; add a * key here only with documentation from PDE (never shift a whole workbook). * (Corrected 2026-09-07 after an external audit: the site previously stored * the attendance element under the report year, which made the state's and * district's figures disagree by ~11 points in 2020-21/2021-22 — an artifact * of comparing different school years, not a change in the state's rules.) */ export const STATE_METRIC_YEAR_OFFSET: Record = { attendance_persistence_rate: -1, }; /** "2024-25" + (−1) → "2023-24". */ export function shiftSchoolYear(year: string, offset: number): string { const m = /^(\d{4})-(\d{2})$/.exec(year); if (!m) throw new Error(`bad school year: ${year}`); const start = Number(m[1]) + offset; return `${start}-${String((start + 1) % 100).padStart(2, "0")}`; } /** The school year a state indicator's value describes, given the workbook's report year. */ export function observationYear(metricKey: string, reportYear: string): string { return shiftSchoolYear(reportYear, STATE_METRIC_YEAR_OFFSET[metricKey] ?? 0); } /** * Map the state's (AUN, PA/Schl code) key to a ULCS code. Two SDP schools * share a state code with their "Continuation Academy" program (Olney HS * 8508, Stetson MS 8507): the state publishes ONE row per code, which * belongs to the main school. Resolve duplicates deterministically — prefer * the school whose name does not contain "Continuation" — and log every * ambiguity so it is never silent. (Before this helper, the last directory * row iterated won, which is order-dependent.) */ export function buildPdeToUlcs( schools: { ulcsCode: string; aunCode: string; paCode: string | null; name?: string | null }[], ): Map { const byKey = new Map(); for (const s of schools) { if (!s.paCode || s.paCode === "N/A") continue; const k = `${s.aunCode}-${s.paCode}`; byKey.set(k, [...(byKey.get(k) ?? []), s]); } const out = new Map(); for (const [k, list] of byKey) { if (list.length === 1) { out.set(k, list[0].ulcsCode); continue; } const preferred = list.find((s) => !/continuation/i.test(s.name ?? "")) ?? list[0]; out.set(k, preferred.ulcsCode); console.warn( `[map] state code ${k} shared by ${list.map((s) => `${s.ulcsCode} ${s.name ?? ""}`).join(" | ")} → assigned to ${preferred.ulcsCode}`, ); } return out; }