{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[],"dockerImageVersionId":28755,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os\nfrom kaggle_secrets import UserSecretsClient\n\nsecrets = UserSecretsClient()\nos.environ[\"KAGGLE_API_TOKEN\"] = secrets.get_secret(\"KAGGLE_API_TOKEN\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:21:00.118281Z","iopub.execute_input":"2026-09-06T20:21:00.119003Z","iopub.status.idle":"2026-09-06T20:21:00.168031Z","shell.execute_reply.started":"2026-09-06T20:21:00.118967Z","shell.execute_reply":"2026-09-06T20:21:00.167203Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!pip install -U kaggle -q\n!kaggle --version","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:21:18.822277Z","iopub.execute_input":"2026-09-06T20:21:18.82257Z","iopub.status.idle":"2026-09-06T20:21:26.804508Z","shell.execute_reply.started":"2026-09-06T20:21:18.822546Z","shell.execute_reply":"2026-09-06T20:21:26.803577Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nRole 2 real preprocessing pipeline.\n\nLocal test:  python3 preprocess_pipeline.py --input test_dicoms/train_series --limit 2\nKaggle:      python3 preprocess_pipeline.py --input /kaggle/input/<comp>/train_series --limit 5\n\n--limit caps how many StudyInstanceUID folders get processed, so you can\nsanity-check on a handful of studies before running the full 4,400+ set.\n\"\"\"\n\nimport os\nimport argparse\nimport numpy as np\nimport pydicom\nimport h5py\nimport pandas as pd\n\nCANONICAL_LATERALITY = \"L\"  # everything gets flipped to match this side\n\n\ndef load_series(series_dir):\n    \"\"\"Load all .dcm files in a series folder, sorted by slice position.\"\"\"\n    files = [os.path.join(series_dir, f) for f in os.listdir(series_dir)\n              if f.endswith(\".dcm\")]\n    datasets = [pydicom.dcmread(f) for f in files]\n\n    # sort by ImagePositionPatient z-coordinate when available, else InstanceNumber\n    def sort_key(ds):\n        if hasattr(ds, \"ImagePositionPatient\"):\n            return float(ds.ImagePositionPatient[2])\n        return int(getattr(ds, \"InstanceNumber\", 0))\n\n    datasets.sort(key=sort_key)\n    return datasets\n\n\ndef rescale(ds):\n    \"\"\"Apply RescaleSlope/Intercept to get true intensity values.\"\"\"\n    pixels = ds.pixel_array.astype(np.float32)\n    slope = float(getattr(ds, \"RescaleSlope\", 1.0))\n    intercept = float(getattr(ds, \"RescaleIntercept\", 0.0))\n    return pixels * slope + intercept\n\n\ndef normalize(volume, low_pct=0.5, high_pct=99.5):\n    \"\"\"Percentile clip then rescale to [0, 1].\"\"\"\n    lo, hi = np.percentile(volume, [low_pct, high_pct])\n    volume = np.clip(volume, lo, hi)\n    if hi > lo:\n        volume = (volume - lo) / (hi - lo)\n    return volume.astype(np.float32)\n\n\ndef get_plane(ds):\n    \"\"\"Infer imaging plane from SeriesDescription, fallback to orientation.\n    Real descriptions use abbreviations (sag/cor/tra), not full words.\"\"\"\n    desc = str(getattr(ds, \"SeriesDescription\", \"\")).lower()\n    plane_keywords = {\n        \"sagittal\": [\"sag\"],\n        \"coronal\": [\"cor\"],\n        \"axial\": [\"tra\", \"ax\", \"transverse\", \"axial\"],\n    }\n    for plane, keywords in plane_keywords.items():\n        if any(kw in desc for kw in keywords):\n            return plane\n    return \"unknown\"\n\n\ndef check_laterality(ds):\n    \"\"\"Return (laterality, needs_flip).\"\"\"\n    lat = str(getattr(ds, \"Laterality\", \"\")).upper() or None\n    if lat is None:\n        return None, False\n    needs_flip = lat != CANONICAL_LATERALITY\n    return lat, needs_flip\n\n\ndef process_series(series_dir):\n    datasets = load_series(series_dir)\n    if not datasets:\n        return None\n\n    volume = np.stack([rescale(ds) for ds in datasets], axis=0)\n    volume = normalize(volume)\n\n    first = datasets[0]\n    plane = get_plane(first)\n    laterality, needs_flip = check_laterality(first)\n    if needs_flip:\n        volume = np.flip(volume, axis=-1)  # mirror left-right axis\n\n    return {\n        \"array\": volume,\n        \"plane\": plane,\n        \"sequence_type\": str(getattr(first, \"SeriesDescription\", \"unknown\")),\n        \"n_slices\": volume.shape[0],\n        \"laterality_flipped\": needs_flip,\n        \"array_shape\": str(tuple(volume.shape)),\n    }\n\n\ndef main(input_root, limit, out_h5, out_csv):\n    study_ids = sorted(os.listdir(input_root))[:limit]\n    manifest_rows = []\n\n    with h5py.File(out_h5, \"w\") as h5f:\n        for study_uid in study_ids:\n            study_dir = os.path.join(input_root, study_uid)\n            if not os.path.isdir(study_dir):\n                continue\n            study_grp = h5f.create_group(study_uid)\n\n            for series_uid in sorted(os.listdir(study_dir)):\n                series_dir = os.path.join(study_dir, series_uid)\n                if not os.path.isdir(series_dir):\n                    continue\n\n                result = process_series(series_dir)\n                if result is None:\n                    continue\n\n                series_grp = study_grp.create_group(series_uid)\n                series_grp.create_dataset(\"slices\", data=result[\"array\"],\n                                           compression=\"gzip\")\n\n                manifest_rows.append({\n                    \"StudyInstanceUID\": study_uid,\n                    \"SeriesInstanceUID\": series_uid,\n                    \"plane\": result[\"plane\"],\n                    \"sequence_type\": result[\"sequence_type\"],\n                    \"n_slices\": result[\"n_slices\"],\n                    \"laterality_flipped\": result[\"laterality_flipped\"],\n                    \"array_shape\": result[\"array_shape\"],\n                })\n\n    pd.DataFrame(manifest_rows).to_csv(out_csv, index=False)\n    print(f\"Processed {len(study_ids)} studies -> {len(manifest_rows)} series\")\n    print(f\"Wrote {out_h5} and {out_csv}\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T18:48:21.059103Z","iopub.execute_input":"2026-09-06T18:48:21.059417Z","iopub.status.idle":"2026-09-06T18:48:22.150689Z","shell.execute_reply.started":"2026-09-06T18:48:21.059388Z","shell.execute_reply":"2026-09-06T18:48:22.149676Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"main(\n    input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n    limit=5,\n    out_h5=\"preprocessed_test.h5\",\n    out_csv=\"series_manifest_test.csv\"\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T18:48:22.152351Z","iopub.execute_input":"2026-09-06T18:48:22.152823Z","iopub.status.idle":"2026-09-06T18:49:00.010351Z","shell.execute_reply.started":"2026-09-06T18:48:22.152795Z","shell.execute_reply":"2026-09-06T18:49:00.009412Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import pandas as pd\ndf = pd.read_csv(\"series_manifest_test.csv\")\nprint(df.to_string())\n\nimport h5py\nimport numpy as np\nimport matplotlib.pyplot as plt\n\nwith h5py.File(\"preprocessed_test.h5\", \"r\") as f:\n    study_ids = list(f.keys())\n    print(f\"Studies in file: {len(study_ids)}\")\n    \n    # pick the first study, first series\n    study = study_ids[0]\n    series_ids = list(f[study].keys())\n    series = series_ids[0]\n    \n    arr = f[study][series][\"slices\"][:]\n    print(f\"Shape: {arr.shape}\")\n    print(f\"Value range: min={arr.min():.3f}, max={arr.max():.3f}\")\n    print(f\"Mean: {arr.mean():.3f}\")\n\n    # look at a middle slice\n    mid = arr.shape[0] // 2\n    plt.figure(figsize=(6,6))\n    plt.imshow(arr[mid], cmap=\"gray\")\n    plt.title(f\"{study[:20]}... / slice {mid}\")\n    plt.axis(\"off\")\n    plt.show()\n\nimport h5py\nimport numpy as np\nimport matplotlib.pyplot as plt\n\nwith h5py.File(\"preprocessed_test.h5\", \"r\") as f:\n    study_ids = list(f.keys())\n    print(f\"Studies in file: {len(study_ids)}\")\n    \n    study = study_ids[0]\n    series_ids = list(f[study].keys())\n    series = series_ids[0]\n    \n    arr = f[study][series][\"slices\"][:]\n    print(f\"Shape: {arr.shape}\")\n    print(f\"Value range: min={arr.min():.3f}, max={arr.max():.3f}\")\n\n    mid = arr.shape[0] // 2\n    plt.figure(figsize=(6,6))\n    plt.imshow(arr[mid], cmap=\"gray\")\n    plt.axis(\"off\")\n    plt.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T18:52:35.03037Z","iopub.execute_input":"2026-09-06T18:52:35.030809Z","iopub.status.idle":"2026-09-06T18:52:35.605887Z","shell.execute_reply.started":"2026-09-06T18:52:35.030768Z","shell.execute_reply":"2026-09-06T18:52:35.605095Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import h5py\nimport numpy as np\nimport matplotlib.pyplot as plt\n\nwith h5py.File(\"preprocessed_test.h5\", \"r\") as f:\n    study_ids = list(f.keys())\n    print(f\"Studies in file: {len(study_ids)}\")\n    \n    study = study_ids[0]\n    series_ids = list(f[study].keys())\n    series = series_ids[0]\n    \n    arr = f[study][series][\"slices\"][:]\n    print(f\"Shape: {arr.shape}\")\n    print(f\"Value range: min={arr.min():.3f}, max={arr.max():.3f}\")\n\n    mid = arr.shape[0] // 2\n    plt.figure(figsize=(6,6))\n    plt.imshow(arr[mid], cmap=\"gray\")\n    plt.axis(\"off\")\n    plt.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T18:52:50.116529Z","iopub.execute_input":"2026-09-06T18:52:50.117382Z","iopub.status.idle":"2026-09-06T18:52:50.387813Z","shell.execute_reply.started":"2026-09-06T18:52:50.117347Z","shell.execute_reply":"2026-09-06T18:52:50.387071Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nSchema for Role 2's output artifact: preprocessed.h5 + series_manifest.csv\n\nSource of truth: ROLES.md § Role 2 — Output Deliverable.\n\npreprocessed.h5 layout (validated structurally, since HDF5 isn't tabular):\n    /{StudyInstanceUID}/{SeriesInstanceUID} -> float32 array, shape (n_slices, H, W)\n\nseries_manifest.csv is the tabular sidecar describing that structure and IS\nvalidated row-by-row via SeriesManifestRow below.\n\"\"\"\nfrom __future__ import annotations\n\nfrom pydantic import BaseModel, Field, field_validator\n\nVALID_PLANES = {\"sagittal\", \"coronal\", \"axial\", \"unknown\"}\n# TODO: confirm the real protocol/sequence vocabulary against the competition's\n# DICOM metadata; T1/T2/PD are placeholders per ROLES.md's competition context.\nVALID_SEQUENCES = {\"T1\", \"T2\", \"PD\"}\n\n\nclass SeriesManifestRow(BaseModel):\n    \"\"\"One row of series_manifest.csv — one series within one study.\"\"\"\n\n    StudyInstanceUID: str = Field(..., min_length=1)\n    SeriesInstanceUID: str = Field(..., min_length=1)\n    plane: str\n    sequence_type: str\n    n_slices: int = Field(..., gt=0)\n    laterality_flipped: bool\n    array_shape: str  # serialized tuple e.g. \"(24, 320, 320)\" — str for CSV portability\n\n    @field_validator(\"plane\")\n    @classmethod\n    def plane_known(cls, v: str) -> str:\n        if v not in VALID_PLANES:\n            raise ValueError(f\"Unknown plane '{v}', expected one of {sorted(VALID_PLANES)}\")\n        return v\n\n\ndef validate_series_manifest_dataframe(df) -> list[SeriesManifestRow]:\n    \"\"\"Validate series_manifest.csv row-by-row. Used by the Role 2 contract test.\"\"\"\n    required = {\n        \"StudyInstanceUID\", \"SeriesInstanceUID\", \"plane\",\n        \"sequence_type\", \"n_slices\", \"laterality_flipped\", \"array_shape\",\n    }\n    missing = required - set(df.columns)\n    if missing:\n        raise ValueError(f\"series_manifest.csv missing required columns: {sorted(missing)}\")\n    return [SeriesManifestRow(**row) for row in df.to_dict(orient=\"records\")]\n\n\ndef validate_preprocessed_h5_structure(h5file, manifest_rows: list[SeriesManifestRow]) -> None:\n    \"\"\"\n    Structural check for preprocessed.h5: every (StudyInstanceUID, SeriesInstanceUID)\n    listed in series_manifest.csv must exist as a 3D dataset in the HDF5 file.\n    Used by the Role 2 contract test.\n    \"\"\"\n    for row in manifest_rows:\n        path = f\"{row.StudyInstanceUID}/{row.SeriesInstanceUID}\"\n        if path not in h5file:\n            raise ValueError(f\"preprocessed.h5 missing dataset for manifest row: {path}\")\n        arr = h5file[path]\n        if arr.ndim != 3:\n            raise ValueError(\n                f\"{path}: expected 3D array (n_slices, H, W), got shape {arr.shape}\"\n            )","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T19:37:32.335251Z","iopub.execute_input":"2026-09-06T19:37:32.335552Z","iopub.status.idle":"2026-09-06T19:37:32.347579Z","shell.execute_reply.started":"2026-09-06T19:37:32.335528Z","shell.execute_reply":"2026-09-06T19:37:32.346967Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def validate_preprocessed_h5_structure_legacy(h5file, manifest_rows):\n    for row in manifest_rows:\n        path = f\"{row.StudyInstanceUID}/{row.SeriesInstanceUID}/slices\"\n        if path not in h5file:\n            raise ValueError(f\"missing dataset: {path}\")\n        arr = h5file[path]\n        if arr.ndim != 3:\n            raise ValueError(f\"{path}: expected 3D, got {arr.shape}\")\n\nwith h5py.File(\"preprocessed_test.h5\", \"r\") as h5f:\n    validate_preprocessed_h5_structure_legacy(h5f, rows)\nprint(\"H5 structure OK (legacy nested layout)\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T19:43:23.063404Z","iopub.execute_input":"2026-09-06T19:43:23.064359Z","iopub.status.idle":"2026-09-06T19:43:23.083747Z","shell.execute_reply.started":"2026-09-06T19:43:23.064321Z","shell.execute_reply":"2026-09-06T19:43:23.082896Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# Test the pipline","metadata":{}},{"cell_type":"code","source":"import os\nfrom kaggle_secrets import UserSecretsClient\n\nsecrets = UserSecretsClient()\nos.environ[\"KAGGLE_API_TOKEN\"] = secrets.get_secret(\"KAGGLE_API_TOKEN\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:24:04.717546Z","iopub.execute_input":"2026-09-06T20:24:04.718411Z","iopub.status.idle":"2026-09-06T20:24:04.776361Z","shell.execute_reply.started":"2026-09-06T20:24:04.718369Z","shell.execute_reply":"2026-09-06T20:24:04.77564Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!pip install -U kaggle -q\n!kaggle --version","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:24:10.182023Z","iopub.execute_input":"2026-09-06T20:24:10.182305Z","iopub.status.idle":"2026-09-06T20:24:14.575135Z","shell.execute_reply.started":"2026-09-06T20:24:10.182281Z","shell.execute_reply":"2026-09-06T20:24:14.574337Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!kaggle datasets list -s knee","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:24:17.286752Z","iopub.execute_input":"2026-09-06T20:24:17.287632Z","iopub.status.idle":"2026-09-06T20:24:18.120897Z","shell.execute_reply.started":"2026-09-06T20:24:17.287588Z","shell.execute_reply":"2026-09-06T20:24:18.119907Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nRole 2 batch pipeline with resume support.\n\nSplits studies into batches. Each batch gets saved as its own h5 file\n(preprocessed_batch_0000.h5, preprocessed_batch_0001.h5, ...) plus rows\nappended to one running series_manifest.csv.\n\nA completed_studies.txt log tracks which StudyInstanceUIDs are fully done.\nIf the notebook disconnects/restarts, just re-run this same cell -- it\nreads the log, skips finished studies, and continues from where it stopped.\n\nUsage in a notebook cell:\n    run_batched(\n        input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n        batch_size=200,\n        out_dir=\"preprocessed_batches\",\n    )\n\"\"\"\n\nimport os\nimport numpy as np\nimport pydicom\nimport h5py\nimport pandas as pd\n\nCANONICAL_LATERALITY = \"L\"\nCOMPLETED_LOG = \"completed_studies.txt\"\nMANIFEST_CSV = \"series_manifest.csv\"\nBATCH_COUNTER_FILE = \"next_batch_number.txt\"\n\n\n# ---- same per-series functions as before ----\n\ndef load_series(series_dir):\n    files = [os.path.join(series_dir, f) for f in os.listdir(series_dir)\n              if f.endswith(\".dcm\")]\n    datasets = [pydicom.dcmread(f) for f in files]\n\n    def sort_key(ds):\n        if hasattr(ds, \"ImagePositionPatient\"):\n            return float(ds.ImagePositionPatient[2])\n        return int(getattr(ds, \"InstanceNumber\", 0))\n\n    datasets.sort(key=sort_key)\n    return datasets\n\n\ndef rescale(ds):\n    pixels = ds.pixel_array.astype(np.float32)\n    slope = float(getattr(ds, \"RescaleSlope\", 1.0))\n    intercept = float(getattr(ds, \"RescaleIntercept\", 0.0))\n    return pixels * slope + intercept\n\n\ndef normalize(volume, low_pct=0.5, high_pct=99.5):\n    lo, hi = np.percentile(volume, [low_pct, high_pct])\n    volume = np.clip(volume, lo, hi)\n    if hi > lo:\n        volume = (volume - lo) / (hi - lo)\n    return volume.astype(np.float32)\n\n\ndef get_plane(ds):\n    desc = str(getattr(ds, \"SeriesDescription\", \"\")).lower()\n    plane_keywords = {\n        \"sagittal\": [\"sag\"],\n        \"coronal\": [\"cor\"],\n        \"axial\": [\"tra\", \"ax\", \"transverse\", \"axial\"],\n    }\n    for plane, keywords in plane_keywords.items():\n        if any(kw in desc for kw in keywords):\n            return plane\n    return \"unknown\"\n\n\ndef check_laterality(ds):\n    lat = str(getattr(ds, \"Laterality\", \"\")).upper() or None\n    if lat is None:\n        return None, False\n    needs_flip = lat != CANONICAL_LATERALITY\n    return lat, needs_flip\n\n\ndef process_series(series_dir):\n    datasets = load_series(series_dir)\n    if not datasets:\n        return None\n\n    volume = np.stack([rescale(ds) for ds in datasets], axis=0)\n    volume = normalize(volume)\n\n    first = datasets[0]\n    plane = get_plane(first)\n    laterality, needs_flip = check_laterality(first)\n    if needs_flip:\n        volume = np.flip(volume, axis=-1)\n\n    return {\n        \"array\": volume,\n        \"plane\": plane,\n        \"sequence_type\": str(getattr(first, \"SeriesDescription\", \"unknown\")),\n        \"n_slices\": volume.shape[0],\n        \"laterality_flipped\": needs_flip,\n        \"array_shape\": str(tuple(volume.shape)),\n    }\n\n\n# ---- batching + resume logic ----\n\ndef get_next_batch_number(out_dir):\n    \"\"\"Reads a persistent counter file so batch numbers never repeat, even\n    if this run is stopped/restarted after batch files were already deleted\n    (e.g. by the API upload's cleanup step).\"\"\"\n    path = os.path.join(out_dir, BATCH_COUNTER_FILE)\n    if not os.path.exists(path):\n        return 0\n    with open(path) as f:\n        return int(f.read().strip())\n\n\ndef save_next_batch_number(out_dir, n):\n    path = os.path.join(out_dir, BATCH_COUNTER_FILE)\n    with open(path, \"w\") as f:\n        f.write(str(n))\n    path = os.path.join(out_dir, COMPLETED_LOG)\n    if not os.path.exists(path):\n        return set()\n    with open(path) as f:\n        return set(line.strip() for line in f if line.strip())\n\n\ndef mark_completed(out_dir, study_uid):\n    path = os.path.join(out_dir, COMPLETED_LOG)\n    with open(path, \"a\") as f:\n        f.write(study_uid + \"\\n\")\n\n\ndef append_manifest(out_dir, rows):\n    if not rows:\n        return\n    path = os.path.join(out_dir, MANIFEST_CSV)\n    df = pd.DataFrame(rows)\n    write_header = not os.path.exists(path)\n    df.to_csv(path, mode=\"a\", header=write_header, index=False)\n\n\ndef process_one_study(h5f, study_dir, study_uid):\n    \"\"\"Returns list of manifest row dicts for this study.\n\n    Writes flat: /{StudyInstanceUID}/{SeriesInstanceUID} -> array directly.\n    This matches schemas/preprocessed_schema.py's expected layout (no extra\n    \"slices\" nesting) -- required for the team's contract validation to pass.\n    \"\"\"\n    manifest_rows = []\n\n    for series_uid in sorted(os.listdir(study_dir)):\n        series_dir = os.path.join(study_dir, series_uid)\n        if not os.path.isdir(series_dir):\n            continue\n        try:\n            result = process_series(series_dir)\n        except Exception as e:\n            print(f\"  [WARN] {study_uid}/{series_uid} failed: {e}\")\n            continue\n        if result is None:\n            continue\n\n        h5f.create_dataset(\n            f\"{study_uid}/{series_uid}\",\n            data=result[\"array\"],\n            compression=\"gzip\",\n        )\n\n        manifest_rows.append({\n            \"StudyInstanceUID\": study_uid,\n            \"SeriesInstanceUID\": series_uid,\n            \"plane\": result[\"plane\"],\n            \"sequence_type\": result[\"sequence_type\"],\n            \"n_slices\": result[\"n_slices\"],\n            \"laterality_flipped\": result[\"laterality_flipped\"],\n            \"array_shape\": result[\"array_shape\"],\n        })\n\n    return manifest_rows\n\n\ndef clear_batch_file(batch_path, auto_clear=True):\n    \"\"\"\n    Pause after a batch finishes so you can download/offload it, then delete\n    it locally to keep /kaggle/working under the storage quota.\n\n    auto_clear=True: waits for you to type 'y' after you've downloaded the\n    file, then deletes it. Nothing is deleted until you confirm.\n    auto_clear=False: skips this step entirely (files just pile up -- only\n    safe for small test runs).\n    \"\"\"\n    if not auto_clear:\n        return\n\n    size_mb = os.path.getsize(batch_path) / (1024 * 1024)\n    print(f\"\\n>>> Batch file ready: {batch_path} ({size_mb:.1f} MB)\")\n    print(\">>> Download it now (or push it to wherever your team stores batches).\")\n    confirm = input(\">>> Type 'y' once it's saved elsewhere, to delete it and continue: \")\n    if confirm.strip().lower() == \"y\":\n        os.remove(batch_path)\n        print(f\">>> Deleted {batch_path}, freed {size_mb:.1f} MB. Continuing...\\n\")\n    else:\n        print(\">>> Not deleted. Re-run this cell later to retry clearing it.\\n\")\n\n\ndef upload_batch_to_kaggle_dataset(batch_path, dataset_slug_prefix, kaggle_username):\n    \"\"\"\n    Uploads a single batch h5 file as its own Kaggle Dataset via the Kaggle\n    API, then deletes the local copy once the upload succeeds. Fully\n    automatic -- no manual download/confirm needed.\n\n    Requires: kaggle.json set up at /root/.kaggle/kaggle.json (see setup\n    instructions -- Kaggle Secrets + UserSecretsClient).\n\n    Each batch becomes its own dataset, e.g.:\n        yourusername/role2-preprocessed-batch-0000\n        yourusername/role2-preprocessed-batch-0001\n        ...\n    Your team can later combine these, or reference them individually.\n    \"\"\"\n    import subprocess\n    import json as _json\n\n    folder = os.path.dirname(batch_path) or \".\"\n    batch_name = os.path.splitext(os.path.basename(batch_path))[0]\n    slug = f\"{dataset_slug_prefix}-{batch_name}\".replace(\"_\", \"-\").lower()\n\n    upload_dir = os.path.join(folder, f\"_upload_{batch_name}\")\n    os.makedirs(upload_dir, exist_ok=True)\n    dest_path = os.path.join(upload_dir, os.path.basename(batch_path))\n    os.rename(batch_path, dest_path)\n\n    metadata = {\n        \"title\": slug,\n        \"id\": f\"{kaggle_username}/{slug}\",\n        \"licenses\": [{\"name\": \"CC0-1.0\"}],\n    }\n    with open(os.path.join(upload_dir, \"dataset-metadata.json\"), \"w\") as f:\n        _json.dump(metadata, f)\n\n    print(f\"\\n>>> Uploading {batch_name} as Kaggle Dataset '{slug}'...\")\n    result = subprocess.run(\n        [\"kaggle\", \"datasets\", \"create\", \"-p\", upload_dir, \"--dir-mode\", \"zip\"],\n        capture_output=True, text=True,\n    )\n    print(result.stdout)\n\n    # Kaggle's CLI sometimes returns exit code 0 even when it printed an\n    # error (e.g. duplicate title) -- so check the actual output text too,\n    # not just returncode, or a failed upload gets silently treated as done.\n    output_text = (result.stdout + result.stderr).lower()\n    failure_markers = [\"error\", \"traceback\", \"failed\"]\n    looks_failed = result.returncode != 0 or any(m in output_text for m in failure_markers)\n\n    if looks_failed:\n        print(result.stderr)\n        print(f\">>> UPLOAD FAILED for {slug} -- local file kept at {dest_path}, not deleted.\")\n        return False\n\n    os.remove(dest_path)\n    os.remove(os.path.join(upload_dir, \"dataset-metadata.json\"))\n    os.rmdir(upload_dir)\n    print(f\">>> Uploaded and cleared: {slug}\\n\")\n    return True\n\n\ndef run_batched(input_root, batch_size=200, out_dir=\"preprocessed_batches\",\n                 clear_mode=\"manual\", dataset_slug_prefix=None, kaggle_username=None):\n    \"\"\"\n    clear_mode:\n        \"manual\" -- pause after each batch, you download it yourself, confirm, it deletes\n        \"api\"    -- automatically upload each batch to a Kaggle Dataset, then delete\n                    (requires dataset_slug_prefix + kaggle_username, and kaggle.json set up)\n        \"none\"   -- don't clear anything (files pile up, only for small tests)\n    \"\"\"\n    os.makedirs(out_dir, exist_ok=True)\n\n    all_study_ids = sorted(os.listdir(input_root))\n    completed = load_completed(out_dir)\n    remaining = [s for s in all_study_ids if s not in completed]\n\n    print(f\"Total studies: {len(all_study_ids)}\")\n    print(f\"Already done (from previous run): {len(completed)}\")\n    print(f\"Remaining: {len(remaining)}\")\n\n    if not remaining:\n        print(\"Nothing left to process.\")\n        return\n\n    # persistent counter -- survives across restarts even after batch files\n    # get deleted post-upload, so batch numbers never collide/repeat\n    next_batch_num = get_next_batch_number(out_dir)\n\n    for i in range(0, len(remaining), batch_size):\n        chunk = remaining[i:i + batch_size]\n        batch_path = os.path.join(out_dir, f\"preprocessed_batch_{next_batch_num:04d}.h5\")\n        print(f\"\\n--- Batch {next_batch_num}: {len(chunk)} studies -> {batch_path} ---\")\n\n        with h5py.File(batch_path, \"w\") as h5f:\n            for study_uid in chunk:\n                study_dir = os.path.join(input_root, study_uid)\n                if not os.path.isdir(study_dir):\n                    continue\n\n                rows = process_one_study(h5f, study_dir, study_uid)\n                append_manifest(out_dir, rows)\n                mark_completed(out_dir, study_uid)\n\n        next_batch_num += 1\n        save_next_batch_number(out_dir, next_batch_num)\n        print(f\"Batch {next_batch_num - 1} done, progress saved.\")\n\n        if clear_mode == \"manual\":\n            clear_batch_file(batch_path, auto_clear=True)\n        elif clear_mode == \"api\":\n            if not dataset_slug_prefix or not kaggle_username:\n                raise ValueError(\"clear_mode='api' requires dataset_slug_prefix and kaggle_username\")\n            upload_batch_to_kaggle_dataset(batch_path, dataset_slug_prefix, kaggle_username)\n        # clear_mode == \"none\": do nothing, file stays\n\n    print(\"\\nAll remaining studies processed.\")\n\n\n# Example (run this in your notebook):\n# run_batched(\n#     input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n#     batch_size=200,\n#     out_dir=\"preprocessed_batches\",\n# )","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:57:56.57864Z","iopub.execute_input":"2026-09-06T20:57:56.579004Z","iopub.status.idle":"2026-09-06T20:57:56.611903Z","shell.execute_reply.started":"2026-09-06T20:57:56.578976Z","shell.execute_reply":"2026-09-06T20:57:56.610919Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"with open(\"preprocessed_batches/completed_studies.txt\") as f:\n    completed = [line.strip() for line in f if line.strip()]\n\n# remove the first 2 (the ones from the lost batch 0)\nstill_completed = completed[2:]\n\nwith open(\"preprocessed_batches/completed_studies.txt\", \"w\") as f:\n    f.write(\"\\n\".join(still_completed) + \"\\n\")\n\nprint(f\"Removed 2 studies from completed log. {len(still_completed)} remain marked done.\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:58:15.532413Z","iopub.execute_input":"2026-09-06T20:58:15.533097Z","iopub.status.idle":"2026-09-06T20:58:15.541367Z","shell.execute_reply.started":"2026-09-06T20:58:15.533066Z","shell.execute_reply":"2026-09-06T20:58:15.54041Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"run_batched(\n    input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n    batch_size=2,\n    out_dir=\"preprocessed_batches\",\n    clear_mode=\"api\",\n    dataset_slug_prefix=\"role2-preprocessed-test\",\n    kaggle_username=\"yahyaismail12\",\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:56:01.881209Z","iopub.execute_input":"2026-09-06T20:56:01.881489Z","iopub.status.idle":"2026-09-06T20:56:04.124055Z","shell.execute_reply.started":"2026-09-06T20:56:01.881465Z","shell.execute_reply":"2026-09-06T20:56:04.122786Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!kaggle datasets list -m","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T21:00:21.06804Z","iopub.execute_input":"2026-09-06T21:00:21.068423Z","iopub.status.idle":"2026-09-06T21:00:21.921885Z","shell.execute_reply.started":"2026-09-06T21:00:21.068393Z","shell.execute_reply":"2026-09-06T21:00:21.920843Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nos.makedirs(\"preprocessed_batches\", exist_ok=True)\nwith open(\"preprocessed_batches/next_batch_number.txt\", \"w\") as f:\n    f.write(\"15\")  # adjust this number based on what the dataset list actually shows\nprint(\"Counter initialized.\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T20:59:55.038669Z","iopub.execute_input":"2026-09-06T20:59:55.039591Z","iopub.status.idle":"2026-09-06T20:59:55.047358Z","shell.execute_reply.started":"2026-09-06T20:59:55.03955Z","shell.execute_reply":"2026-09-06T20:59:55.046513Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"from kaggle.api.kaggle_api_extended import KaggleApi\nimport re\n\napi = KaggleApi()\napi.authenticate()\n\nall_datasets = []\npage = 1\nwhile True:\n    batch = api.dataset_list(mine=True, page=page)\n    if not batch:\n        break\n    all_datasets.extend(batch)\n    page += 1\n\nbatch_nums = []\nfor d in all_datasets:\n    m = re.search(r\"role2-preprocessed-test-preprocessed-batch-(\\d+)\", d.ref)\n    if m:\n        batch_nums.append(int(m.group(1)))\n\nbatch_nums.sort()\nprint(f\"Total batch datasets found: {len(batch_nums)}\")\nprint(f\"Range: {min(batch_nums)} to {max(batch_nums)}\")\n\nexpected = set(range(min(batch_nums), max(batch_nums) + 1))\nmissing = sorted(expected - set(batch_nums))\nprint(f\"Missing batch numbers (gaps): {missing}\")\n\nduplicates = [n for n in set(batch_nums) if batch_nums.count(n) > 1]\nprint(f\"Duplicate batch numbers: {duplicates}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T21:01:14.158222Z","iopub.execute_input":"2026-09-06T21:01:14.158697Z","iopub.status.idle":"2026-09-06T21:01:14.838103Z","shell.execute_reply.started":"2026-09-06T21:01:14.158652Z","shell.execute_reply":"2026-09-06T21:01:14.837124Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nos.makedirs(\"preprocessed_batches\", exist_ok=True)\nwith open(\"preprocessed_batches/next_batch_number.txt\", \"w\") as f:\n    f.write(\"56\")\nprint(\"Counter set to 56.\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T21:03:45.240034Z","iopub.execute_input":"2026-09-06T21:03:45.240812Z","iopub.status.idle":"2026-09-06T21:03:45.248648Z","shell.execute_reply.started":"2026-09-06T21:03:45.240772Z","shell.execute_reply":"2026-09-06T21:03:45.247921Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\nrun_batched(\n    input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n    batch_size=2,\n    out_dir=\"preprocessed_batches\",\n    clear_mode=\"api\",\n    dataset_slug_prefix=\"role2-preprocessed-test\",\n    kaggle_username=\"yahyaismail12\",\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-06T21:03:51.545765Z","iopub.execute_input":"2026-09-06T21:03:51.54611Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"from kaggle.api.kaggle_api_extended import KaggleApi\nimport re\n\napi = KaggleApi()\napi.authenticate()\n\nall_datasets = []\npage = 1\nwhile True:\n    batch = api.dataset_list(mine=True, page=page)\n    if not batch:\n        break\n    all_datasets.extend(batch)\n    page += 1\n\nbatch_nums = []\nfor d in all_datasets:\n    m = re.search(r\"role2-preprocessed-test-preprocessed-batch-(\\d+)\", d.ref)\n    if m:\n        batch_nums.append(int(m.group(1)))\n\nbatch_nums.sort()\nprint(f\"Total: {len(batch_nums)}, range: {min(batch_nums)}-{max(batch_nums)}\")\nexpected = set(range(min(batch_nums), max(batch_nums)+1))\nprint(\"Missing:\", sorted(expected - set(batch_nums)))\nprint(\"Duplicates:\", [n for n in set(batch_nums) if batch_nums.count(n) > 1])","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T00:00:10.045533Z","iopub.execute_input":"2026-09-07T00:00:10.045902Z","iopub.status.idle":"2026-09-07T00:00:14.847925Z","shell.execute_reply.started":"2026-09-07T00:00:10.045868Z","shell.execute_reply":"2026-09-07T00:00:14.847048Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nBackfills the dataset_ref column into series_manifest.csv for batches that\nwere already uploaded before dataset_ref tracking existed.\n\nWorks because batch numbering is deterministic: completed_studies.txt lines\nare in the exact order studies were processed, and batches were sequential\nchunks of `batch_size` studies each, named batch_0000, batch_0001, ...\n\nRun this in your Kaggle notebook, in the same out_dir as your real run.\n\"\"\"\n\nimport os\nimport pandas as pd\n\nOUT_DIR = \"preprocessed_batches\"\nBATCH_SIZE = 2  # must match whatever batch_size you actually used\nDATASET_SLUG_PREFIX = \"role2-preprocessed-test\"\nKAGGLE_USERNAME = \"yahyaismail12\"\n\n\ndef backfill_dataset_ref():\n    with open(os.path.join(OUT_DIR, \"completed_studies.txt\")) as f:\n        completed_in_order = [line.strip() for line in f if line.strip()]\n\n    # study_uid -> batch number, based on position in the completed log\n    study_to_batch = {}\n    for i, study_uid in enumerate(completed_in_order):\n        batch_num = i // BATCH_SIZE\n        study_to_batch[study_uid] = batch_num\n\n    manifest_path = os.path.join(OUT_DIR, \"series_manifest.csv\")\n    df = pd.read_csv(manifest_path)\n\n    def make_ref(study_uid):\n        batch_num = study_to_batch.get(study_uid)\n        if batch_num is None:\n            return None  # study not in completed log -- shouldn't happen, flag it\n        batch_name = f\"preprocessed_batch_{batch_num:04d}\"\n        slug = f\"{DATASET_SLUG_PREFIX}-{batch_name}\".replace(\"_\", \"-\").lower()\n        return f\"{KAGGLE_USERNAME}/{slug}\"\n\n    df[\"dataset_ref\"] = df[\"StudyInstanceUID\"].apply(make_ref)\n\n    missing = df[\"dataset_ref\"].isna().sum()\n    if missing > 0:\n        print(f\"WARNING: {missing} manifest rows could not be matched to a batch. \"\n              f\"Investigate before trusting this file.\")\n    else:\n        print(f\"All {len(df)} manifest rows successfully matched to a dataset_ref.\")\n\n    df.to_csv(manifest_path, index=False)\n    print(f\"Backfilled and saved: {manifest_path}\")\n    print(df[[\"StudyInstanceUID\", \"dataset_ref\"]].drop_duplicates().tail(10))\n\n\nif __name__ == \"__main__\":\n    backfill_dataset_ref()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:12:48.717712Z","iopub.execute_input":"2026-09-07T01:12:48.718063Z","iopub.status.idle":"2026-09-07T01:12:48.824977Z","shell.execute_reply.started":"2026-09-07T01:12:48.718033Z","shell.execute_reply":"2026-09-07T01:12:48.823963Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nprint(os.listdir(\".\"))\nprint(os.listdir(\"preprocessed_batches\") if os.path.exists(\"preprocessed_batches\") else \"preprocessed_batches folder doesn't exist\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:09:44.824314Z","iopub.execute_input":"2026-09-07T01:09:44.824702Z","iopub.status.idle":"2026-09-07T01:09:44.83124Z","shell.execute_reply.started":"2026-09-07T01:09:44.824671Z","shell.execute_reply":"2026-09-07T01:09:44.830044Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nRecovery script for a wiped Kaggle working directory.\n\nRebuilds completed_studies.txt and series_manifest.csv WITHOUT redoing any\nactual pixel processing or re-uploading anything. Your 449 batches are\nalready safely on Kaggle -- this just reconstructs the local bookkeeping\nthat got lost when the session reset.\n\nHow: reads DICOM headers only (fast, no pixel decoding) for the studies we\nknow are already done (the first N studies in sorted order, where\nN = num_batches_done * batch_size), derives the same plane/laterality/shape\nmetadata the real pipeline would have produced, and assigns dataset_ref\nusing the same deterministic batch-number mapping as before.\n\"\"\"\n\nimport os\nimport pydicom\nimport pandas as pd\n\nINPUT_ROOT = \"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\"\nOUT_DIR = \"preprocessed_batches\"\nNUM_BATCHES_DONE = 449\nBATCH_SIZE = 2\nDATASET_SLUG_PREFIX = \"role2-preprocessed-test\"\nKAGGLE_USERNAME = \"yahyaismail12\"\n\nCANONICAL_LATERALITY = \"L\"\n\n\ndef get_plane(desc):\n    desc = str(desc).lower()\n    plane_keywords = {\n        \"sagittal\": [\"sag\"],\n        \"coronal\": [\"cor\"],\n        \"axial\": [\"tra\", \"ax\", \"transverse\", \"axial\"],\n    }\n    for plane, keywords in plane_keywords.items():\n        if any(kw in desc for kw in keywords):\n            return plane\n    return \"unknown\"\n\n\ndef reconstruct_series_row(series_dir, study_uid, series_uid):\n    \"\"\"Reads only DICOM headers (stop_before_pixels) -- fast, no pixel decode.\"\"\"\n    files = [os.path.join(series_dir, f) for f in os.listdir(series_dir) if f.endswith(\".dcm\")]\n    if not files:\n        return None\n\n    first = pydicom.dcmread(files[0], stop_before_pixels=True)\n    n_slices = len(files)\n    rows = int(getattr(first, \"Rows\", 0))\n    cols = int(getattr(first, \"Columns\", 0))\n\n    plane = get_plane(getattr(first, \"SeriesDescription\", \"\"))\n    lat = str(getattr(first, \"Laterality\", \"\")).upper() or None\n    laterality_flipped = (lat is not None) and (lat != CANONICAL_LATERALITY)\n\n    return {\n        \"StudyInstanceUID\": study_uid,\n        \"SeriesInstanceUID\": series_uid,\n        \"plane\": plane,\n        \"sequence_type\": str(getattr(first, \"SeriesDescription\", \"unknown\")),\n        \"n_slices\": n_slices,\n        \"laterality_flipped\": laterality_flipped,\n        \"array_shape\": str((n_slices, rows, cols)),\n    }\n\n\ndef recover():\n    os.makedirs(OUT_DIR, exist_ok=True)\n\n    all_study_ids = sorted(os.listdir(INPUT_ROOT))\n    n_done = NUM_BATCHES_DONE * BATCH_SIZE\n    done_study_ids = all_study_ids[:n_done]\n\n    print(f\"Reconstructing metadata for {len(done_study_ids)} already-processed studies...\")\n\n    manifest_rows = []\n    for idx, study_uid in enumerate(done_study_ids):\n        study_dir = os.path.join(INPUT_ROOT, study_uid)\n        batch_num = idx // BATCH_SIZE\n        batch_name = f\"preprocessed_batch_{batch_num:04d}\"\n        slug = f\"{DATASET_SLUG_PREFIX}-{batch_name}\".replace(\"_\", \"-\").lower()\n        dataset_ref = f\"{KAGGLE_USERNAME}/{slug}\"\n\n        for series_uid in sorted(os.listdir(study_dir)):\n            series_dir = os.path.join(study_dir, series_uid)\n            if not os.path.isdir(series_dir):\n                continue\n            row = reconstruct_series_row(series_dir, study_uid, series_uid)\n            if row is None:\n                continue\n            row[\"dataset_ref\"] = dataset_ref\n            manifest_rows.append(row)\n\n        if (idx + 1) % 100 == 0:\n            print(f\"  ...{idx + 1}/{len(done_study_ids)} studies rebuilt\")\n\n    # write completed_studies.txt\n    with open(os.path.join(OUT_DIR, \"completed_studies.txt\"), \"w\") as f:\n        f.write(\"\\n\".join(done_study_ids) + \"\\n\")\n\n    # write series_manifest.csv\n    pd.DataFrame(manifest_rows).to_csv(os.path.join(OUT_DIR, \"series_manifest.csv\"), index=False)\n\n    # write the persistent batch counter so future runs continue correctly\n    with open(os.path.join(OUT_DIR, \"next_batch_number.txt\"), \"w\") as f:\n        f.write(str(NUM_BATCHES_DONE))\n\n    print(f\"\\nRecovered {len(done_study_ids)} studies, {len(manifest_rows)} series.\")\n    print(f\"Wrote completed_studies.txt, series_manifest.csv, next_batch_number.txt to {OUT_DIR}/\")\n\n\nif __name__ == \"__main__\":\n    recover()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:11:04.458171Z","iopub.execute_input":"2026-09-07T01:11:04.458529Z","iopub.status.idle":"2026-09-07T01:12:25.679552Z","shell.execute_reply.started":"2026-09-07T01:11:04.458491Z","shell.execute_reply":"2026-09-07T01:12:25.678668Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import kaggle\n\ntest_ref = \"yahyaismail12/role2-preprocessed-test-preprocessed-batch-0448\"  # or any ref from your manifest\nkaggle.api.dataset_download_files(test_ref, path=\"spot_check/\", unzip=True)\n\nimport h5py, os\nh5_path = [f for f in os.listdir(\"spot_check\") if f.endswith(\".h5\")][0]\nwith h5py.File(f\"spot_check/{h5_path}\", \"r\") as f:\n    print(list(f.keys()))  # should show the actual StudyInstanceUID(s) you expect","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:13:45.814379Z","iopub.execute_input":"2026-09-07T01:13:45.814828Z","iopub.status.idle":"2026-09-07T01:13:57.931828Z","shell.execute_reply.started":"2026-09-07T01:13:45.81479Z","shell.execute_reply":"2026-09-07T01:13:57.930604Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\"\"\"\nRole 2 batch pipeline with resume support.\n\nSplits studies into batches. Each batch gets saved as its own h5 file\n(preprocessed_batch_0000.h5, preprocessed_batch_0001.h5, ...) plus rows\nappended to one running series_manifest.csv.\n\nA completed_studies.txt log tracks which StudyInstanceUIDs are fully done.\nIf the notebook disconnects/restarts, just re-run this same cell -- it\nreads the log, skips finished studies, and continues from where it stopped.\n\nUsage in a notebook cell:\n    run_batched(\n        input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n        batch_size=200,\n        out_dir=\"preprocessed_batches\",\n    )\n\"\"\"\n\nimport os\nimport numpy as np\nimport pydicom\nimport h5py\nimport pandas as pd\n\nCANONICAL_LATERALITY = \"L\"\nCOMPLETED_LOG = \"completed_studies.txt\"\nMANIFEST_CSV = \"series_manifest.csv\"\nBATCH_COUNTER_FILE = \"next_batch_number.txt\"\n\n\n# ---- same per-series functions as before ----\n\ndef load_series(series_dir):\n    files = [os.path.join(series_dir, f) for f in os.listdir(series_dir)\n              if f.endswith(\".dcm\")]\n    datasets = [pydicom.dcmread(f) for f in files]\n\n    def sort_key(ds):\n        if hasattr(ds, \"ImagePositionPatient\"):\n            return float(ds.ImagePositionPatient[2])\n        return int(getattr(ds, \"InstanceNumber\", 0))\n\n    datasets.sort(key=sort_key)\n    return datasets\n\n\ndef rescale(ds):\n    pixels = ds.pixel_array.astype(np.float32)\n    slope = float(getattr(ds, \"RescaleSlope\", 1.0))\n    intercept = float(getattr(ds, \"RescaleIntercept\", 0.0))\n    return pixels * slope + intercept\n\n\ndef normalize(volume, low_pct=0.5, high_pct=99.5):\n    lo, hi = np.percentile(volume, [low_pct, high_pct])\n    volume = np.clip(volume, lo, hi)\n    if hi > lo:\n        volume = (volume - lo) / (hi - lo)\n    return volume.astype(np.float32)\n\n\ndef get_plane(ds):\n    desc = str(getattr(ds, \"SeriesDescription\", \"\")).lower()\n    plane_keywords = {\n        \"sagittal\": [\"sag\"],\n        \"coronal\": [\"cor\"],\n        \"axial\": [\"tra\", \"ax\", \"transverse\", \"axial\"],\n    }\n    for plane, keywords in plane_keywords.items():\n        if any(kw in desc for kw in keywords):\n            return plane\n    return \"unknown\"\n\n\ndef check_laterality(ds):\n    lat = str(getattr(ds, \"Laterality\", \"\")).upper() or None\n    if lat is None:\n        return None, False\n    needs_flip = lat != CANONICAL_LATERALITY\n    return lat, needs_flip\n\n\ndef process_series(series_dir):\n    datasets = load_series(series_dir)\n    if not datasets:\n        return None\n\n    volume = np.stack([rescale(ds) for ds in datasets], axis=0)\n    volume = normalize(volume)\n\n    first = datasets[0]\n    plane = get_plane(first)\n    laterality, needs_flip = check_laterality(first)\n    if needs_flip:\n        volume = np.flip(volume, axis=-1)\n\n    return {\n        \"array\": volume,\n        \"plane\": plane,\n        \"sequence_type\": str(getattr(first, \"SeriesDescription\", \"unknown\")),\n        \"n_slices\": volume.shape[0],\n        \"laterality_flipped\": needs_flip,\n        \"array_shape\": str(tuple(volume.shape)),\n    }\n\n\n# ---- batching + resume logic ----\n\ndef get_next_batch_number(out_dir):\n    \"\"\"Reads a persistent counter file so batch numbers never repeat, even\n    if this run is stopped/restarted after batch files were already deleted\n    (e.g. by the API upload's cleanup step).\"\"\"\n    path = os.path.join(out_dir, BATCH_COUNTER_FILE)\n    if not os.path.exists(path):\n        return 0\n    with open(path) as f:\n        return int(f.read().strip())\n\n\ndef save_next_batch_number(out_dir, n):\n    path = os.path.join(out_dir, BATCH_COUNTER_FILE)\n    with open(path, \"w\") as f:\n        f.write(str(n))\n    path = os.path.join(out_dir, COMPLETED_LOG)\n    if not os.path.exists(path):\n        return set()\n    with open(path) as f:\n        return set(line.strip() for line in f if line.strip())\n\n\ndef load_completed(out_dir):\n    path = os.path.join(out_dir, COMPLETED_LOG)\n    if not os.path.exists(path):\n        return set()\n    with open(path) as f:\n        return set(line.strip() for line in f if line.strip())\n\n\ndef mark_completed(out_dir, study_uid):\n    path = os.path.join(out_dir, COMPLETED_LOG)\n    with open(path, \"a\") as f:\n        f.write(study_uid + \"\\n\")\n\n\ndef append_manifest(out_dir, rows):\n    if not rows:\n        return\n    path = os.path.join(out_dir, MANIFEST_CSV)\n    df = pd.DataFrame(rows)\n    write_header = not os.path.exists(path)\n    df.to_csv(path, mode=\"a\", header=write_header, index=False)\n\n\ndef process_one_study(h5f, study_dir, study_uid):\n    \"\"\"Returns list of manifest row dicts for this study.\n\n    Writes flat: /{StudyInstanceUID}/{SeriesInstanceUID} -> array directly.\n    This matches schemas/preprocessed_schema.py's expected layout (no extra\n    \"slices\" nesting) -- required for the team's contract validation to pass.\n    \"\"\"\n    manifest_rows = []\n\n    for series_uid in sorted(os.listdir(study_dir)):\n        series_dir = os.path.join(study_dir, series_uid)\n        if not os.path.isdir(series_dir):\n            continue\n        try:\n            result = process_series(series_dir)\n        except Exception as e:\n            print(f\"  [WARN] {study_uid}/{series_uid} failed: {e}\")\n            continue\n        if result is None:\n            continue\n\n        h5f.create_dataset(\n            f\"{study_uid}/{series_uid}\",\n            data=result[\"array\"],\n            compression=\"gzip\",\n        )\n\n        manifest_rows.append({\n            \"StudyInstanceUID\": study_uid,\n            \"SeriesInstanceUID\": series_uid,\n            \"plane\": result[\"plane\"],\n            \"sequence_type\": result[\"sequence_type\"],\n            \"n_slices\": result[\"n_slices\"],\n            \"laterality_flipped\": result[\"laterality_flipped\"],\n            \"array_shape\": result[\"array_shape\"],\n        })\n\n    return manifest_rows\n\n\ndef clear_batch_file(batch_path, auto_clear=True):\n    \"\"\"\n    Pause after a batch finishes so you can download/offload it, then delete\n    it locally to keep /kaggle/working under the storage quota.\n    Returns True once you've confirmed it's safely saved elsewhere.\n    \"\"\"\n    if not auto_clear:\n        return True\n\n    size_mb = os.path.getsize(batch_path) / (1024 * 1024)\n    print(f\"\\n>>> Batch file ready: {batch_path} ({size_mb:.1f} MB)\")\n    print(\">>> Download it now (or push it to wherever your team stores batches).\")\n    confirm = input(\">>> Type 'y' once it's saved elsewhere, to delete it and continue: \")\n    if confirm.strip().lower() == \"y\":\n        os.remove(batch_path)\n        print(f\">>> Deleted {batch_path}, freed {size_mb:.1f} MB. Continuing...\\n\")\n        return True\n    else:\n        print(\">>> Not deleted. Re-run this cell later to retry clearing it.\\n\")\n        return False\n\n\ndef upload_batch_to_kaggle_dataset(batch_path, dataset_slug_prefix, kaggle_username):\n    \"\"\"\n    Uploads a single batch h5 file as its own Kaggle Dataset via the Kaggle\n    API, then deletes the local copy once the upload succeeds. Fully\n    automatic -- no manual download/confirm needed.\n\n    Requires: kaggle.json set up at /root/.kaggle/kaggle.json (see setup\n    instructions -- Kaggle Secrets + UserSecretsClient).\n\n    Each batch becomes its own dataset, e.g.:\n        yourusername/role2-preprocessed-batch-0000\n        yourusername/role2-preprocessed-batch-0001\n        ...\n    Your team can later combine these, or reference them individually.\n    \"\"\"\n    import subprocess\n    import json as _json\n\n    folder = os.path.dirname(batch_path) or \".\"\n    batch_name = os.path.splitext(os.path.basename(batch_path))[0]\n    slug = f\"{dataset_slug_prefix}-{batch_name}\".replace(\"_\", \"-\").lower()\n\n    upload_dir = os.path.join(folder, f\"_upload_{batch_name}\")\n    os.makedirs(upload_dir, exist_ok=True)\n    dest_path = os.path.join(upload_dir, os.path.basename(batch_path))\n    os.rename(batch_path, dest_path)\n\n    metadata = {\n        \"title\": slug,\n        \"id\": f\"{kaggle_username}/{slug}\",\n        \"licenses\": [{\"name\": \"CC0-1.0\"}],\n    }\n    with open(os.path.join(upload_dir, \"dataset-metadata.json\"), \"w\") as f:\n        _json.dump(metadata, f)\n\n    print(f\"\\n>>> Uploading {batch_name} as Kaggle Dataset '{slug}'...\")\n    result = subprocess.run(\n        [\"kaggle\", \"datasets\", \"create\", \"-p\", upload_dir, \"--dir-mode\", \"zip\"],\n        capture_output=True, text=True,\n    )\n    print(result.stdout)\n\n    # Kaggle's CLI sometimes returns exit code 0 even when it printed an\n    # error (e.g. duplicate title) -- so check the actual output text too,\n    # not just returncode, or a failed upload gets silently treated as done.\n    output_text = (result.stdout + result.stderr).lower()\n    failure_markers = [\"error\", \"traceback\", \"failed\"]\n    looks_failed = result.returncode != 0 or any(m in output_text for m in failure_markers)\n\n    if looks_failed:\n        print(result.stderr)\n        print(f\">>> UPLOAD FAILED for {slug} -- local file kept at {dest_path}, not deleted.\")\n        return False, None\n\n    os.remove(dest_path)\n    os.remove(os.path.join(upload_dir, \"dataset-metadata.json\"))\n    os.rmdir(upload_dir)\n    dataset_ref = f\"{kaggle_username}/{slug}\"\n    print(f\">>> Uploaded and cleared: {dataset_ref}\\n\")\n    return True, dataset_ref\n\n\ndef backup_bookkeeping(out_dir, dataset_slug, kaggle_username):\n    \"\"\"\n    Uploads completed_studies.txt + series_manifest.csv + next_batch_number.txt\n    to a small, persistent Kaggle Dataset. If the notebook session resets and\n    wipes /kaggle/working again, you just re-download this one small dataset\n    instead of re-running the DICOM-header recovery script.\n\n    Safe to call repeatedly -- creates the dataset once, then versions\n    (overwrites) it on every subsequent call with the latest state.\n    \"\"\"\n    import subprocess\n    import json as _json\n    import shutil\n\n    backup_dir = os.path.join(out_dir, \"_bookkeeping_backup\")\n    os.makedirs(backup_dir, exist_ok=True)\n\n    for fname in [\"completed_studies.txt\", \"series_manifest.csv\", \"next_batch_number.txt\"]:\n        src = os.path.join(out_dir, fname)\n        if os.path.exists(src):\n            shutil.copy(src, os.path.join(backup_dir, fname))\n\n    marker_path = os.path.join(out_dir, \"_backup_dataset_created.flag\")\n    already_created = os.path.exists(marker_path)\n\n    if not already_created:\n        metadata = {\n            \"title\": dataset_slug,\n            \"id\": f\"{kaggle_username}/{dataset_slug}\",\n            \"licenses\": [{\"name\": \"CC0-1.0\"}],\n        }\n        with open(os.path.join(backup_dir, \"dataset-metadata.json\"), \"w\") as f:\n            _json.dump(metadata, f)\n        cmd = [\"kaggle\", \"datasets\", \"create\", \"-p\", backup_dir, \"--dir-mode\", \"zip\"]\n    else:\n        cmd = [\"kaggle\", \"datasets\", \"version\", \"-p\", backup_dir, \"-m\", \"auto backup\", \"--dir-mode\", \"zip\"]\n\n    print(f\"\\n>>> Backing up bookkeeping files to '{dataset_slug}'...\")\n    result = subprocess.run(cmd, capture_output=True, text=True)\n    print(result.stdout)\n\n    output_text = (result.stdout + result.stderr).lower()\n    failed = result.returncode != 0 or any(m in output_text for m in [\"error\", \"traceback\", \"failed\"])\n\n    if failed:\n        print(result.stderr)\n        print(\">>> Bookkeeping backup FAILED (not fatal -- will retry at the next backup interval)\\n\")\n        return False\n\n    if not already_created:\n        with open(marker_path, \"w\") as f:\n            f.write(\"done\")\n\n    print(\">>> Bookkeeping backup uploaded successfully.\\n\")\n    return True\n\n\ndef run_batched(input_root, batch_size=200, out_dir=\"preprocessed_batches\",\n                 clear_mode=\"manual\", dataset_slug_prefix=None, kaggle_username=None,\n                 backup_every_n_batches=None, backup_dataset_slug=None):\n    \"\"\"\n    clear_mode:\n        \"manual\" -- pause after each batch, you download it yourself, confirm, it deletes\n        \"api\"    -- automatically upload each batch to a Kaggle Dataset, then delete\n                    (requires dataset_slug_prefix + kaggle_username, and kaggle.json set up)\n        \"none\"   -- don't clear anything (files pile up, only for small tests)\n\n    backup_every_n_batches: if set (e.g. 10), backs up completed_studies.txt +\n        series_manifest.csv + next_batch_number.txt to a small Kaggle Dataset\n        every N batches, so a session reset never requires full DICOM-header\n        recovery again -- just re-download this one small backup dataset.\n        Requires backup_dataset_slug + kaggle_username.\n\n    Studies are only marked \"completed\" and their manifest rows only written\n    AFTER their batch's upload is confirmed successful. If upload fails, the\n    batch's local h5 file is kept and those studies stay unmarked, so the\n    next run retries them -- nothing is silently lost or falsely marked done.\n\n    When clear_mode=\"api\", each manifest row gets a \"dataset_ref\" column\n    recording exactly which Kaggle Dataset holds that series' data. This is\n    the single series_manifest.csv you hand to Role 3 -- they look up\n    dataset_ref per study/series and pull only what they need via the\n    Kaggle API, instead of anyone attaching hundreds of datasets by hand.\n    \"\"\"\n    os.makedirs(out_dir, exist_ok=True)\n\n    all_study_ids = sorted(os.listdir(input_root))\n    completed = load_completed(out_dir)\n    remaining = [s for s in all_study_ids if s not in completed]\n\n    print(f\"Total studies: {len(all_study_ids)}\")\n    print(f\"Already done (from previous run): {len(completed)}\")\n    print(f\"Remaining: {len(remaining)}\")\n\n    if not remaining:\n        print(\"Nothing left to process.\")\n        return\n\n    next_batch_num = get_next_batch_number(out_dir)\n\n    for i in range(0, len(remaining), batch_size):\n        chunk = remaining[i:i + batch_size]\n        batch_path = os.path.join(out_dir, f\"preprocessed_batch_{next_batch_num:04d}.h5\")\n        print(f\"\\n--- Batch {next_batch_num}: {len(chunk)} studies -> {batch_path} ---\")\n\n        # buffer rows in memory -- NOT written to disk / marked done yet\n        pending_rows = []\n        pending_study_ids = []\n\n        with h5py.File(batch_path, \"w\") as h5f:\n            for study_uid in chunk:\n                study_dir = os.path.join(input_root, study_uid)\n                if not os.path.isdir(study_dir):\n                    continue\n\n                rows = process_one_study(h5f, study_dir, study_uid)\n                pending_rows.extend(rows)\n                pending_study_ids.append(study_uid)\n\n        print(f\"Batch {next_batch_num} processed (not yet marked done -- pending upload).\")\n\n        if clear_mode == \"manual\":\n            uploaded_ok = clear_batch_file(batch_path, auto_clear=True)\n            dataset_ref = None\n        elif clear_mode == \"api\":\n            if not dataset_slug_prefix or not kaggle_username:\n                raise ValueError(\"clear_mode='api' requires dataset_slug_prefix and kaggle_username\")\n            uploaded_ok, dataset_ref = upload_batch_to_kaggle_dataset(\n                batch_path, dataset_slug_prefix, kaggle_username\n            )\n        else:\n            uploaded_ok, dataset_ref = True, None  # clear_mode == \"none\": treat as done, file stays\n\n        if uploaded_ok:\n            if dataset_ref:\n                for row in pending_rows:\n                    row[\"dataset_ref\"] = dataset_ref\n            append_manifest(out_dir, pending_rows)\n            for study_uid in pending_study_ids:\n                mark_completed(out_dir, study_uid)\n            next_batch_num += 1\n            save_next_batch_number(out_dir, next_batch_num)\n            print(f\"Batch {next_batch_num - 1} confirmed uploaded -- studies marked done, manifest updated.\")\n\n            if backup_every_n_batches and next_batch_num % backup_every_n_batches == 0:\n                if not backup_dataset_slug or not kaggle_username:\n                    print(\">>> Skipping backup: backup_dataset_slug/kaggle_username not set.\")\n                else:\n                    backup_bookkeeping(out_dir, backup_dataset_slug, kaggle_username)\n        else:\n            print(f\"Batch {next_batch_num} upload FAILED -- studies NOT marked done, \"\n                  f\"manifest NOT updated, local file kept at {batch_path}. \"\n                  f\"Re-run this cell to retry this batch.\")\n            return  # stop here so failures don't cascade silently\n\n    print(\"\\nAll remaining studies processed.\")\n    if backup_every_n_batches and backup_dataset_slug and kaggle_username:\n        backup_bookkeeping(out_dir, backup_dataset_slug, kaggle_username)\n\n\n# Example (run this in your notebook):\n# run_batched(\n#     input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n#     batch_size=200,\n#     out_dir=\"preprocessed_batches\",\n# )","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:40:57.167424Z","iopub.execute_input":"2026-09-07T01:40:57.168111Z","iopub.status.idle":"2026-09-07T01:40:57.21968Z","shell.execute_reply.started":"2026-09-07T01:40:57.168059Z","shell.execute_reply":"2026-09-07T01:40:57.218733Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!pip install -U kaggle -q","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T01:40:29.953742Z","iopub.execute_input":"2026-09-07T01:40:29.955405Z","iopub.status.idle":"2026-09-07T01:40:37.572627Z","shell.execute_reply.started":"2026-09-07T01:40:29.955331Z","shell.execute_reply":"2026-09-07T01:40:37.571186Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"run_batched(\n    input_root=\"/kaggle/input/competitions/rsna-knee-abnormality-detection/train_series\",\n    batch_size=10,\n    out_dir=\"preprocessed_batches\",\n    clear_mode=\"api\",\n    dataset_slug_prefix=\"role2-preprocessed-test\",\n    kaggle_username=\"yahyaismail12\",\n    backup_every_n_batches=5,\n    backup_dataset_slug=\"role2-bookkeeping-backup\",\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-09-07T02:26:37.249525Z","iopub.execute_input":"2026-09-07T02:26:37.249866Z","iopub.status.idle":"2026-09-07T10:09:43.655201Z","shell.execute_reply.started":"2026-09-07T02:26:37.249836Z","shell.execute_reply":"2026-09-07T10:09:43.653717Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}