Source document

pipeline_philly/discover.py

Served verbatim from the project repository. Internal working document conventions apply: documents may reference file paths, branch names, and findings-ledger anchors from the repo.

"""First-fetch / discovery for Philly sources.

Reads docs/cities/philly/sources.yaml, fetches each source whose URL has been
resolved past TBD, records provenance (sha256 + retrieved_at + bytes) to
data/cities/philly/sources/<source_id>/<date>/, and emits a per-source
findings markdown documenting the sheets / columns / sample rows it found.

Decision-neutral: runs as a standalone script that produces artifacts;
doesn't touch the NYC DB, NYC code, or any Postgres. Mirrors
verify/download.py and verify/connectors/source_file.py in spirit.

Usage:
  python -m pipeline_philly.discover --source pde_pssa
  python -m pipeline_philly.discover --all
"""
from __future__ import annotations

import argparse
import datetime as dt
import hashlib
import json
import pathlib
import sys
from typing import Any

import requests
import yaml

REPO_ROOT = pathlib.Path(__file__).resolve().parent.parent
SOURCES_YAML = REPO_ROOT / "docs" / "cities" / "philly" / "sources.yaml"
DATA_DIR = REPO_ROOT / "data" / "cities" / "philly" / "sources"
FINDINGS_DIR = REPO_ROOT / "pipeline_philly" / "findings"


def load_sources() -> list[dict[str, Any]]:
    with open(SOURCES_YAML) as f:
        return yaml.safe_load(f)


def is_resolved(src: dict[str, Any]) -> bool:
    """A source is fetchable once its url isn't TBD."""
    url = src.get("url")
    return bool(url and url != "TBD")


def fetch(src: dict[str, Any]) -> pathlib.Path:
    """Download a source's file, write provenance, return the local path."""
    sid = src["source_id"]
    url = src["url"]
    today = dt.date.today().isoformat()
    out_dir = DATA_DIR / sid / today
    out_dir.mkdir(parents=True, exist_ok=True)

    # Filename heuristic: last URL segment, or sid + extension
    fname = url.rsplit("/", 1)[-1].split("?")[0]
    if not fname or "." not in fname:
        ext = src.get("format", "bin")
        fname = f"{sid}.{ext}"
    path = out_dir / fname

    print(f"[fetch] {sid} -> {url}", file=sys.stderr)
    r = requests.get(url, timeout=120, stream=True)
    r.raise_for_status()
    h = hashlib.sha256()
    n_bytes = 0
    with open(path, "wb") as f:
        for chunk in r.iter_content(chunk_size=1 << 16):
            if not chunk:
                continue
            f.write(chunk)
            h.update(chunk)
            n_bytes += len(chunk)
    sha = h.hexdigest()
    print(f"[fetch] wrote {path} ({n_bytes:,} bytes, sha256={sha[:12]}...)", file=sys.stderr)

    provenance = {
        "source_id": sid,
        "url": url,
        "retrieved_at": dt.datetime.utcnow().isoformat() + "Z",
        "sha256": sha,
        "bytes": n_bytes,
        "publisher": src.get("publisher"),
        "format": src.get("format"),
    }
    (out_dir / "provenance.json").write_text(json.dumps(provenance, indent=2))
    return path


def inspect_xlsx(path: pathlib.Path) -> dict[str, Any]:
    """List sheet names + column names + sample rows for each sheet."""
    import pandas as pd

    out: dict[str, Any] = {"format": "xlsx", "sheets": {}}
    xf = pd.ExcelFile(path, engine="openpyxl")
    out["sheet_names"] = list(xf.sheet_names)
    for sheet in xf.sheet_names[:8]:  # first 8 sheets; many files have a "Notes" sheet
        try:
            df = pd.read_excel(xf, sheet_name=sheet, nrows=5)
        except Exception as e:
            out["sheets"][sheet] = {"error": str(e)}
            continue
        out["sheets"][sheet] = {
            "n_columns": len(df.columns),
            "columns": [str(c) for c in df.columns],
            "sample_row": (
                {str(k): (None if pd.isna(v) else str(v)) for k, v in df.iloc[0].items()}
                if len(df) > 0 else None
            ),
            "n_rows_in_sample": len(df),
        }
    return out


def inspect_zip(path: pathlib.Path) -> dict[str, Any]:
    import zipfile

    out: dict[str, Any] = {"format": "zip", "entries": []}
    with zipfile.ZipFile(path) as z:
        for info in z.infolist()[:30]:
            out["entries"].append({
                "name": info.filename,
                "bytes": info.file_size,
                "compressed": info.compress_size,
            })
        out["entry_count"] = len(z.namelist())
    return out


def inspect_csv(path: pathlib.Path) -> dict[str, Any]:
    import pandas as pd

    out: dict[str, Any] = {"format": "csv"}
    df = pd.read_csv(path, nrows=5, low_memory=False)
    # Count total rows efficiently (avoid loading the whole file)
    with open(path, "rb") as f:
        n_rows = sum(1 for _ in f) - 1  # minus header
    out["n_rows"] = n_rows
    out["n_columns"] = len(df.columns)
    out["columns"] = [str(c) for c in df.columns]
    out["sample_row"] = (
        {str(k): (None if pd.isna(v) else str(v)) for k, v in df.iloc[0].items()}
        if len(df) > 0 else None
    )
    return out


def inspect(path: pathlib.Path, fmt: str | None) -> dict[str, Any]:
    if fmt == "xlsx" or path.suffix.lower() == ".xlsx":
        return inspect_xlsx(path)
    if fmt == "zip" or path.suffix.lower() == ".zip":
        return inspect_zip(path)
    if fmt == "csv" or path.suffix.lower() == ".csv":
        return inspect_csv(path)
    return {"format": fmt or path.suffix, "note": "no inspector for format"}


def render_findings(src: dict[str, Any], path: pathlib.Path, inspection: dict[str, Any]) -> str:
    sid = src["source_id"]
    out: list[str] = []
    out.append(f"# {sid} — first-fetch findings\n")
    out.append(f"- publisher: {src.get('publisher')}")
    out.append(f"- url: {src.get('url')}")
    out.append(f"- format: {src.get('format')}")
    out.append(f"- local path: `{path.relative_to(REPO_ROOT)}`")
    out.append("")
    if inspection.get("format") == "xlsx":
        out.append(f"## Sheets ({len(inspection.get('sheet_names', []))})")
        for name in inspection.get("sheet_names", []):
            out.append(f"- `{name}`")
        out.append("")
        for name, sheet in inspection.get("sheets", {}).items():
            out.append(f"### `{name}` columns ({sheet.get('n_columns')})")
            for c in sheet.get("columns", []):
                out.append(f"- `{c}`")
            if sheet.get("sample_row"):
                out.append("\nSample first row:\n")
                out.append("```json")
                out.append(json.dumps(sheet["sample_row"], indent=2))
                out.append("```")
            out.append("")
    elif inspection.get("format") == "zip":
        out.append(f"## Archive contents ({inspection.get('entry_count')} entries)")
        for entry in inspection.get("entries", []):
            out.append(f"- `{entry['name']}` ({entry['bytes']:,} bytes)")
    elif inspection.get("format") == "csv":
        out.append(f"## CSV — {inspection.get('n_rows', '?'):,} rows × {inspection.get('n_columns', '?')} columns")
        out.append("")
        for c in inspection.get("columns", []):
            out.append(f"- `{c}`")
        if inspection.get("sample_row"):
            out.append("\nSample first row:\n")
            out.append("```json")
            out.append(json.dumps(inspection["sample_row"], indent=2))
            out.append("```")
    return "\n".join(out)


def discover_one(src: dict[str, Any]) -> None:
    sid = src["source_id"]
    if not is_resolved(src):
        print(f"[skip] {sid}: url is TBD", file=sys.stderr)
        return
    try:
        path = fetch(src)
        inspection = inspect(path, src.get("format"))
        findings_md = render_findings(src, path, inspection)
        FINDINGS_DIR.mkdir(parents=True, exist_ok=True)
        (FINDINGS_DIR / f"{sid}.md").write_text(findings_md)
        print(f"[findings] -> pipeline_philly/findings/{sid}.md", file=sys.stderr)
    except Exception as e:
        print(f"[error] {sid}: {type(e).__name__}: {e}", file=sys.stderr)


def main() -> None:
    p = argparse.ArgumentParser()
    p.add_argument("--source", help="single source_id to fetch", default=None)
    p.add_argument("--all", action="store_true", help="fetch every resolved source")
    args = p.parse_args()

    sources = load_sources()
    if args.source:
        match = [s for s in sources if s["source_id"] == args.source]
        if not match:
            print(f"no such source: {args.source}", file=sys.stderr)
            sys.exit(1)
        for s in match:
            discover_one(s)
    elif args.all:
        for s in sources:
            discover_one(s)
    else:
        # default: list sources and their resolution state
        for s in sources:
            mark = "✓" if is_resolved(s) else "·"
            print(f"  {mark} {s['source_id']:38s} {s.get('publisher','-')}")


if __name__ == "__main__":
    main()