{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.11"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Knee MRI ensemble inference — two GPUs, weighted votes, nothing surrendered\n\nInference-only notebook: no training happens here. It reads an attached package of ~20\ntrained DINOv2-small checkpoints (public score 0.891) plus an optional independently\ntrained five-fold bundle, and predicts the 12 findings for every test study.\n\n**Methods used, in one list:**\n\n- **Sequence recovery from DICOM headers** — fat-suppression and T1/T2/PD weighting are\n  re-derived from raw tags (description, scan options, TR/TE), not taken from the CSV.\n- **Slot representation** — each study becomes exactly 6 canonical series slots\n  (plane x contrast), with a presence mask for missing ones.\n- **Laterality normalisation** — every right knee is mapped onto a left-knee convention\n  (mirror on coronal/axial, channel reversal on sagittal).\n- **Geometric slice ordering** — slices sorted by patient-space position, never filename.\n- **Physical-scale crop** — a fixed millimetre field of view before resizing, so anatomy\n  has the same scale everywhere.\n- **One-time uint8 pixel cache** — decode once per pixel-contract group, hold in RAM.\n- **DINOv2 encoder + per-diagnosis slot attention** — each finding attends over the\n  slots where radiologists actually read it.\n- **Fingerprint verification** — every checkpoint proves it computes exactly what it\n  computed when trained before it is allowed to vote.\n- **Overlapping-window TTA** — all 7 consecutive 3-slice windows averaged per member.\n- **Jitter TTA (adaptive)** — an extra augmented view per window, granted only when the\n  time estimate shows slack.\n- **Weighted rank-mean ensembling** — the metric reads only order, so members vote in\n  percentile ranks; the package votes at weight 1.0, the legacy bundle at 0.5/fold.\n- **Dual-GPU work queue, banked submissions, peer retry** — both T4s pull members from a\n  queue, `submission.csv` is rewritten after every member, and a member failing on one\n  GPU retries on the other. A killed run still submits its best partial ensemble.\n\n**Settings: GPU T4 x2, Internet OFF.** Attach: competition data, the member package (the\ndataset with `manifest.json`), optionally the legacy bundle (`rsna_20260807_v1.pt`),\nand the DINOv2 Kaggle models (small, plus base if the manifest lists base members).","metadata":{}},{"cell_type":"markdown","source":"## 1. Configuration\nAll pixel-contract constants: crop size in millimetres, cache resolution, slice band,\nthe 6 slot definitions, the two rule sets (`native`/`legacy`) that pin down slice order,\nlaterality evidence, slot fallback and decode fill, plus regexes for reading sequence\nnames. Ends by listing the available GPUs — the worker roster for the member queue.","metadata":{}},{"cell_type":"code","source":"from __future__ import annotations\n\nimport os\n\nfor _v in (\"OMP_NUM_THREADS\", \"OPENBLAS_NUM_THREADS\", \"MKL_NUM_THREADS\"):\n    os.environ.setdefault(_v, \"4\")\n\nimport gc\nimport hashlib\nimport json\nimport re\nimport time\nimport traceback\nfrom concurrent.futures import ThreadPoolExecutor\nfrom pathlib import Path\n\nimport numpy as np\nimport pandas as pd\nimport pydicom\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\n\n# The label extractor is defined in the cells above when this runs as a notebook. As a\n# plain script it is imported from the package source, so the two paths share one\n# definition rather than keeping a copy each.\n\nT0 = time.time()\nSEED = 2026\nnp.random.seed(SEED)\ntorch.manual_seed(SEED)\n\nTARGETS = [\"ACL\", \"MCL\", \"Medial Meniscus\", \"Lateral Meniscus\", \"Medial OA\",\n           \"Lateral OA\", \"PF OA\", \"Effusion\", \"Synovitis\", \"Baker's\",\n           \"Contusion\", \"Fracture\"]\n\n\n# The centre crop has to be smaller than the smallest field of view in the corpus or it\n# silently does nothing. Measured over every training series, the acquired field of view\n# (Rows x PixelSpacing) has median 160 mm and runs from 70 to 320: a 160 mm crop is\n# larger than the image in 60% of series and is skipped for all of them, which leaves\n# their physical scale unnormalised. 130 mm is below the field of view of 99.6% of\n# series and still contains the joint.\nCROP_MM = 130.0\n\n# Cache resolution. Everything downstream may downsample from this, so it is set by the\n# most demanding configuration rather than by the default one.\nCACHE_IMG = 336\nGROUP = 3                  # slices per encoder input, stacked as the three channels\nN_GROUP_MAX = 1\nCACHE_FRACTION = 0.45      # share of free memory the pixel cache may take\nCACHE_BUDGET_MAX_GB = 24.0 # hard ceiling regardless of what the machine reports\nCACHE_BUDGET_GB = 12.0     # only the fallback, for a machine with no /proc/meminfo\nTEST_SHARE = 0.30          # floor on the test corpus relative to the training one, since\n                           # the visible test split is a stub and the scored one is not\nHDR_THREADS = 16\nPIX_THREADS = 12\nORDER_THREADS = 32         # slice-ordering is latency-bound on the mount, not CPU-bound\n# Ceiling for the ordering pass. It has to be a ceiling because the pass is hundreds of\n# thousands of small reads over a network mount, so its duration is a property of the\n# mount rather than of the work, and varies between runs that do the same reading. It must not be a tight one: giving up leaves\n# those series in file order, which is uncorrelated with anatomy, and that degradation is\n# silent. So the ceiling sits well above what the pass ordinarily needs: its purpose is\n# to stop the pass consuming the whole run on a slow mount, not to trim the ordinary\n# case, and a ceiling tight enough to bind on a normal day would trade a silent\n# degradation for a saving the run does not need.\nORDER_BUDGET_S = 5400\n\n# Resolution is the axis under test. A feature of width d mm survives resampling only if\n# the pixel pitch is at most d/2, and the pitch here is set by the crop above rather than\n# by the acquired field of view: CROP_MM / P. At 224 px that is 0.58 mm, above the 0.5 mm\n# a 1 mm tear needs; at 336 px it is 0.39 mm and clears it. Both configurations read the\n# same cache, so the comparison isolates the resize.\nRUNS = [\n    {\"name\": \"r224\", \"img\": 224},\n    {\"name\": \"r336\", \"img\": 336},\n]\n\nEPOCHS = 10\nBATCH_STUDIES = 8          # a study is a bag of up to N_SLOT slot images\nAUG_ROT_DEG = 8.0          # rigid jitter; see augment() for why neither flip is used\nAUG_SCALE = 0.08\nAUG_SHIFT = 0.05\nAUG_INTENSITY = 0.10\nLAT_MIN_OFFSET_MM = 20.0   # inside this the side is not readable from geometry; see\n                           # side_from_geometry()\nSLICE_BAND = (0.20, 0.80)  # fraction of the ordered stack read_slot samples across\n\n# --- What a slice IS, as opposed to how many of them there are --------------- #\n#\n# A member is a function of the pixels it was fitted on, and img/crop_mm/slices/band do\n# not determine those pixels by themselves. Four further decisions do, none of them\n# visible in any shape:\n#\n#   order          which slice is the next one along the stack\n#   lat            which knees are mirrored, and on what evidence\n#   slot_fallback  whether a T1 slot may be filled from a series that is not T1\n#   decode_fill    what stands in for a slice that would not decode\n#\n# `native` is the reading derived in the sections below. `legacy` is the reading an\n# imported member was fitted under. A member read under the wrong one loads with every\n# shape matching, runs, and writes a plausible submission computed from the wrong image -\n# so the choice travels with the member and is part of the key that decides which members\n# can share a decode. The legacy rules are reproduced rather than corrected: correcting\n# them would hand that member pixels its weights never saw.\nRULES_NATIVE = {\"order\": \"normal\", \"lat\": \"centre\",\n                \"slot_fallback\": False, \"decode_fill\": \"nearest\"}\nRULES_LEGACY = {\"order\": \"dominant_axis\", \"lat\": \"corner_x\",\n                \"slot_fallback\": True, \"decode_fill\": \"zero\"}\nRULES = dict(RULES_NATIVE)\nLEGACY_LAT_OFFSET_MM = 5.0   # the dead zone the legacy laterality rule was fitted with\n\nLR_HEAD = 1e-3\nLR_BACKBONE = 8e-6         # the encoder is adapted, not retrained\nUNFREEZE_LAST = 6          # trainable transformer blocks, from the output end\nWEIGHT_DECAY = 0.02\nEVAL_BATCH = 8\nTIME_BUDGET = 8.0 * 3600\n\n# Six slots: three planes crossed with the acquisition axes. The fat-suppressed\n# fluid-sensitive series exist for nearly every study; the T1 and the non-suppressed\n# fluid-sensitive series are scarcer, which is what the presence mask is for.\nSLOTS_RECOVERED = [\n    (\"SAG_FLUID_FS\", \"Sagittal\", True, True),\n    (\"COR_FLUID_FS\", \"Coronal\", True, True),\n    (\"AX_FLUID_FS\", \"Axial\", True, True),\n    (\"SAG_FLUID_NOFS\", \"Sagittal\", True, False),\n    (\"COR_T1\", \"Coronal\", False, False),\n    (\"SAG_T1\", \"Sagittal\", False, False),\n]\n\n# The alternative: plane x the single axis the delivered flags carry, ignoring the\n# recovered weighting. Kept as\n# a switch so the choice of slot definition can be varied while everything else is held\n# fixed. Under this scheme a `Struct` slot mixes T1 series with non-fat-suppressed PD/T2\n# series, which carry very different tissue contrast.\nSLOTS_PUBLIC = [\n    (\"SAG_FLUID\", \"Sagittal\", None, True),\n    (\"COR_FLUID\", \"Coronal\", None, True),\n    (\"AX_FLUID\", \"Axial\", None, True),\n    (\"SAG_STRUCT\", \"Sagittal\", None, False),\n    (\"COR_STRUCT\", \"Coronal\", None, False),\n    (\"AX_STRUCT\", \"Axial\", None, False),\n]\n\nSLOT_SCHEME = os.environ.get(\"SLOT_SCHEME\", \"recovered\")\nSLOTS = SLOTS_PUBLIC if SLOT_SCHEME == \"public\" else SLOTS_RECOVERED\nN_SLOT = len(SLOTS)\n\n# How many 384-wide parts the per-slot feature is built from. The encoder emits one\n# vector per token; a slot feature is a fixed summary of that grid, and the summary an\n# imported member was fitted with carries a third part.\nPOOL_PARTS = {\"cls_mean\": 2, \"cls_mean_focal\": 3}\n\n# Which slots an imported member's attention is tilted toward, per diagnosis. Indices are\n# into SLOTS. This is a fixed table rather than a learned parameter, so it is part of that\n# member's definition and has to be reproduced exactly for its weights to mean anything.\nSLOT_PRIOR_TABLE = {\n    \"ACL\": (0, 3, 5), \"MCL\": (1, 4),\n    \"Medial Meniscus\": (0, 1, 3, 4), \"Lateral Meniscus\": (0, 1, 3, 4),\n    \"Medial OA\": (1, 4, 5), \"Lateral OA\": (1, 4, 5),\n    \"PF OA\": (0, 2, 5), \"Effusion\": (0, 2), \"Synovitis\": (0, 2),\n    \"Baker's\": (0,), \"Contusion\": (0, 1, 2), \"Fracture\": (0, 1, 2, 4, 5),\n}\nSLOT_PRIOR_STRENGTH = 0.55\n\nFATSAT_OPTS = {\"FS\", \"FATSAT\", \"FAT_SAT\", \"FSAT\"}\n_SEP = re.compile(r\"[_\\-.]\")\n_FATSAT_RX = re.compile(r\"\\bfs\\b|fatsat|fat sat|\\bstir\\b|\\bspair\\b|\\bspir\\b|\\bwe\\b|\"\n                        r\"water excit|\\btirm\\b|\\bsting\\b|\\bfatsup\\b\")\n_T1_RX = re.compile(r\"\\bt1\\b|\\bt1w\\b\")\n_T2_RX = re.compile(r\"\\bt2\\b|\\bt2w\\b\")\n_PD_RX = re.compile(r\"\\bpd\\b|\\bpdw\\b|proton|\\bdp\\b|dens\")\n\n\nimport threading\n\n# ---- dual-GPU roster (the executive change vs the reference notebook) ------------------\n# The reference ran members one after another on a single device and, when its time\n# budget bound, surrendered TTA windows and then whole members (public 0.847 vs the\n# package's 0.891). Members are independent functions of the same cache, so they\n# parallelise across devices with no change to any member's pixels or weights.\nif torch.cuda.is_available():\n    DEVS = [torch.device(f\"cuda:{i}\") for i in range(torch.cuda.device_count())]\nelse:\n    DEVS = [torch.device(\"cpu\")]\nprint(f\"devices: {[str(d) for d in DEVS]}\")\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 2. Paths and memory plan\nFinds the competition mount and the DINOv2 checkpoint directories, then asks the machine\nhow much RAM it will lend and sizes the pixel cache (slices per slot) to fit it.","metadata":{}},{"cell_type":"code","source":"def log(msg):\n    print(f\"[{time.time() - T0:7.1f}s] {msg}\", flush=True)\n\n\ndef find_root():\n    for c in [Path(\"/kaggle/input/competitions/rsna-knee-abnormality-detection\"),\n              Path(\"/kaggle/input/rsna-knee-abnormality-detection\"),\n              Path(\"data\"), Path(\".\")]:\n        if (c / \"test.csv\").is_file() and (c / \"test_series\").is_dir():\n            return c\n    # last resort: two-level scan, because the mount is nested one deeper than usual\n    base = Path(\"/kaggle/input\")\n    if base.is_dir():\n        for depth1 in sorted(p for p in base.iterdir() if p.is_dir()):\n            for cand in [depth1] + sorted(p for p in depth1.iterdir() if p.is_dir()):\n                if (cand / \"test.csv\").is_file():\n                    return cand\n    raise FileNotFoundError(\n        f\"competition mount not found (cwd {Path.cwd()}); expected a directory holding \"\n        f\"test.csv and test_series/\")\n\n\ndef find_dinov2(variant=\"small\"):\n    \"\"\"Locate a mounted DINOv2 checkpoint directory by variant name.\"\"\"\n    base = Path(\"/kaggle/input\")\n    if not base.is_dir():\n        return None\n    hits = []\n    for root, dirs, files in os.walk(base):\n        dirs[:] = [d for d in dirs if d not in (\"train_series\", \"test_series\")]\n        if \"config.json\" in files and \"dinov2\" in root.lower():\n            hits.append(Path(root))\n    for h in hits:\n        if variant in str(h).lower():\n            return h\n    return hits[0] if hits else None\n\n\nROOT = find_root()\nlog(f\"input root: {ROOT}\")\n\n\nIMG = CACHE_IMG            # kept as the name the pixel reader and cache use\n\n\ndef available_gb():\n    \"\"\"Memory this machine will actually lend, read rather than assumed.\n\n    A hardcoded ceiling is a guess about a machine the author is not sitting at, and a\n    guess that is too low costs coverage silently while a guess that is too high ends the\n    run. The machine will say, so it is asked.\n    \"\"\"\n    try:\n        with open(\"/proc/meminfo\") as fh:\n            info = {k.strip(): v for k, v in\n                    (l.split(\":\", 1) for l in fh if \":\" in l)}\n        return int(info[\"MemAvailable\"].split()[0]) / 1024 ** 2\n    except Exception:\n        return CACHE_BUDGET_GB / CACHE_FRACTION      # fall back to the old constant\n\n\ndef plan_cache(n_study, n_test=0):\n    \"\"\"Choose how many slices per slot the memory the machine has will allow.\n\n    The cache is n_study x n_slot x slices x IMG^2 bytes. Coverage is the cheap axis -\n    linear - and resolution the expensive one, so when the budget binds it is the slice\n    count that gives way rather than the pixel grid. Deciding once, from the training\n    corpus size, keeps train and test caches on the same group layout.\n\n    Only a fraction of what is free is taken. The rest is not slack: the encoder, its\n    activations, the pinned batches and the frames all come out of the same pool, and the\n    cache is the one allocation big enough that overshooting it kills the run outright.\n    \"\"\"\n    avail = available_gb()\n    budget = min(avail * CACHE_FRACTION, CACHE_BUDGET_MAX_GB)\n    # Both caches are held at once, and the test half is what the visible run cannot\n    # show: here it is a handful of studies, and at scoring it is the whole hidden set.\n    # Sizing against the training corpus alone therefore passes every run that can be\n    # watched and overruns the one that counts.\n    n_total = n_study + max(n_test, int(TEST_SHARE * n_study))\n    per_slice = n_total * N_SLOT * IMG * IMG\n    afford = int(budget * 1024 ** 3 // max(per_slice, 1))\n    groups = max(1, min(N_GROUP_MAX, afford // GROUP))\n    log(f\"memory: {avail:.1f} GB available, {budget:.1f} GB to the cache; \"\n        f\"sizing for {n_study} train + {n_total - n_study} test studies \"\n        f\"-> {groups} group(s) of {GROUP} = {groups * GROUP} slices per slot\"\n        + (f\" (wanted {N_GROUP_MAX})\" if groups < N_GROUP_MAX else \"\"))\n    return groups\n\n\nN_GROUP = plan_cache(len(pd.read_csv(ROOT / \"train.csv\")),\n                     len(pd.read_csv(ROOT / \"test.csv\")))\nCACHE_SLICES = GROUP * N_GROUP\nlog(f\"cache layout: {N_GROUP} groups x {GROUP} slices = {CACHE_SLICES} per slot\")\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 3. Header pass\nReads one DICOM header per series (no pixels) to recover what each series *is*:\nfat-suppressed or not, T1/T2/PD weighting from description and pulse timing, and which\nknee was scanned — from the laterality tag where present, from geometry where the rule\nset allows it.","metadata":{}},{"cell_type":"code","source":"HDR_TAGS = [\"SeriesDescription\", \"SequenceName\", \"ScanOptions\", \"ScanningSequence\",\n            \"RepetitionTime\", \"EchoTime\", \"Laterality\", \"PixelSpacing\", \"Rows\",\n            \"Columns\", \"RescaleSlope\", \"RescaleIntercept\",\n            # Position and orientation are read from the same header probe() already\n            # opens, so they cost nothing, and they are what recovers the side when the\n            # Laterality tag is absent - which it is for half the studies here.\n            \"ImagePositionPatient\", \"ImageOrientationPatient\"]\n\n\ndef _hdr_vec(s, n):\n    \"\"\"Parse a DICOM multi-value string as stored by probe(): floats joined by `|`.\"\"\"\n    if not isinstance(s, str):\n        return None\n    try:\n        v = [float(x) for x in s.split(\"|\")]\n    except ValueError:\n        return None\n    return np.array(v) if len(v) >= n else None\n\n\ndef side_from_geometry(h):\n    \"\"\"Study -> 'L' / 'R' / None, from where the image sits in the patient.\n\n    `Laterality` (0020,0060) is Type 2C and may legitimately be absent; in this corpus it\n    is missing on exactly half the studies, and the vendors it is missing from are whole\n    vendors rather than scattered series. A study with no tag is not a left knee, but the\n    normalisation upstream treats it as one, so half the corpus was never normalised and\n    the five side-defined targets - the two menisci, the two tibiofemoral compartments\n    and the medial collateral ligament - saw that axis reversed on a large minority of it.\n\n    The patient coordinate system fixes this without the tag: +x is the patient's left, so\n    the centre of a right knee sits at negative x. The centre is used rather than\n    `ImagePositionPatient` itself because that is the corner of the image, which is offset\n    by half a field of view - enough to change the sign on a knee near the midline.\n\n    The median over a study's series is what is thresholded, not a single series: probe()\n    reads one arbitrary slice per series, which on a sagittal stack can sit anywhere\n    across the joint. Studies whose centre falls near the midline are left unresolved\n    rather than guessed - measured against the tagged half, the rule is right 97% of the\n    time overall and no better than chance inside 20 mm.\n    \"\"\"\n    cx = {}\n    for r in h.itertuples(index=False):\n        ipp = _hdr_vec(getattr(r, \"ImagePositionPatient\", None), 3)\n        iop = _hdr_vec(getattr(r, \"ImageOrientationPatient\", None), 6)\n        ps = _hdr_vec(getattr(r, \"PixelSpacing\", None), 2)\n        rows, cols = getattr(r, \"Rows\", None), getattr(r, \"Columns\", None)\n        if ipp is None or iop is None or ps is None or not rows or not cols:\n            continue\n        try:\n            c = ipp[:3] + iop[:3] * ps[1] * float(cols) / 2 + iop[3:6] * ps[0] * float(rows) / 2\n        except (TypeError, ValueError):\n            continue\n        cx.setdefault(r.StudyInstanceUID, []).append(float(c[0]))\n    out = {}\n    for st, xs in cx.items():\n        m = float(np.median(xs))\n        out[st] = None if abs(m) < LAT_MIN_OFFSET_MM else (\"R\" if m < 0 else \"L\")\n    return out\n\n\ndef side_from_corner_x(h):\n    \"\"\"The laterality an imported member was fitted under.\n\n    It thresholds the median raw `ImagePositionPatient` x over a study's series. That is\n    the x of the image *corner*, not of its centre, so it differs from the rule above by\n    up to half a field of view - which is enough to reverse the sign on a knee scanned\n    near the midline. The dead zone is 5 mm rather than 20 mm, so it also commits on\n    studies the rule above leaves unresolved.\n\n    Neither difference changes a shape. Each one decides whether a study is mirrored, and\n    a study mirrored one way at training and the other at inference presents the five\n    side-defined targets with their axis reversed.\n    \"\"\"\n    out = {}\n    for st, g in h.groupby(\"StudyInstanceUID\"):\n        xs = []\n        for r in g.itertuples(index=False):\n            ipp = _hdr_vec(getattr(r, \"ImagePositionPatient\", None), 3)\n            if ipp is not None and np.isfinite(ipp).all():\n                xs.append(float(ipp[0]))\n        if not xs:\n            out[st] = None\n            continue\n        x = float(np.median(xs))\n        # DICOM patient coordinates are LPS: +x is the patient's left.\n        out[st] = None if abs(x) < LEGACY_LAT_OFFSET_MM else (\"R\" if x < 0 else \"L\")\n    return out\n\n\ndef lat_of(h, tag=\"\"):\n    \"\"\"Study -> 'L' / 'R' / None: the tag where it exists, geometry where it does not.\n\n    The tag is present on exactly half the studies here and is sometimes an empty\n    string rather than absent, which is not the same as NaN. Treating the other half\n    as left-sided is what `normalise_laterality` did by omission, so the geometry\n    fallback is not a refinement - it is the difference between normalising half the\n    corpus and normalising all of it.\n    \"\"\"\n    geo = side_from_corner_x(h) if RULES[\"lat\"] == \"corner_x\" else side_from_geometry(h)\n    d, n_tag, n_geo, n_none, n_disagree = {}, 0, 0, 0, 0\n    for st, g in h.groupby(\"StudyInstanceUID\"):\n        v = [str(x).strip().upper() for x in g[\"Laterality\"].dropna()]\n        if RULES[\"lat\"] == \"corner_x\" and \"ImageLaterality\" in g.columns:\n            # The legacy rule reads the second tag too, so a study tagged only there is\n            # resolved from the tag rather than from geometry.\n            v += [str(x).strip().upper() for x in g[\"ImageLaterality\"].dropna()]\n        v = [x[0] for x in v if x and x[0] in (\"L\", \"R\")]\n        side = v[0] if v else None\n        if side is not None:\n            n_tag += 1\n            if geo.get(st) is not None and geo[st] != side:\n                n_disagree += 1\n        else:\n            side = geo.get(st)\n            n_geo += side is not None\n            n_none += side is None\n        d[st] = side\n    log(f\"{tag}laterality: {n_tag} from the tag, {n_geo} from geometry, \"\n        f\"{n_none} unresolved; tag and geometry disagree on {n_disagree} \"\n        f\"({n_disagree / max(n_tag, 1):.1%} of the tagged)\")\n    return d\n\n\n\ndef probe(item):\n    split, study, series, path = item\n    row = {\"split\": split, \"StudyInstanceUID\": study, \"SeriesInstanceUID\": series,\n           \"dir\": path}\n    try:\n        files = sorted(e.name for e in os.scandir(path) if e.name.endswith(\".dcm\"))\n        row[\"files\"] = files\n        row[\"n_slices\"] = len(files)\n        if not files:\n            return row\n        ds = pydicom.dcmread(os.path.join(path, files[len(files) // 2]),\n                             stop_before_pixels=True, force=True)\n        for t in HDR_TAGS:\n            v = getattr(ds, t, None)\n            if v is None:\n                row[t] = None\n            elif isinstance(v, (list, tuple)) or type(v).__name__ == \"MultiValue\":\n                row[t] = \"|\".join(str(x) for x in v)\n            else:\n                row[t] = str(v)\n    except Exception as exc:\n        row[\"err\"] = str(exc)[:120]\n    return row\n\n\ndef walk(split):\n    \"\"\"Every series directory of a split, with one header read per series.\n\n    An absent split returns an empty frame *with the columns annotate expects*. Returning\n    a bare DataFrame looks like the same thing and is not: the next call indexes\n    `SeriesDescription` and raises KeyError, so the branch that exists to survive a\n    missing split is what turns it into a crash.\n    \"\"\"\n    base = ROOT / split\n    items = []\n    if not base.is_dir():\n        return pd.DataFrame(columns=[\"split\", \"StudyInstanceUID\", \"SeriesInstanceUID\",\n                                     \"dir\", \"files\", \"n_slices\"] + HDR_TAGS)\n    for study in os.scandir(base):\n        if study.is_dir():\n            for series in os.scandir(study.path):\n                if series.is_dir():\n                    items.append((split, study.name, series.name, series.path))\n    with ThreadPoolExecutor(max_workers=HDR_THREADS) as pool:\n        rows = list(pool.map(probe, items))\n    return pd.DataFrame(rows)\n\n\ndef annotate(df):\n    \"\"\"Recover fat suppression and pulse-sequence weighting from the header.\"\"\"\n    desc = (df[\"SeriesDescription\"].fillna(\"\") + \" \" + df[\"SequenceName\"].fillna(\"\"))\n    desc = desc.str.lower().str.replace(_SEP, \" \", regex=True)\n\n    opts = df[\"ScanOptions\"].fillna(\"\").str.upper().str.split(\"|\")\n    # GE writes SAT_GEMS for spatial saturation, so ScanOptions must be matched as\n    # exact tokens; a substring test on \"SAT\" fires on non-fat-sat series.\n    opts_fs = opts.apply(lambda ts: any(t.strip() in FATSAT_OPTS for t in ts))\n    df[\"fatsat\"] = desc.str.contains(_FATSAT_RX) | opts_fs\n\n    tr = pd.to_numeric(df[\"RepetitionTime\"], errors=\"coerce\")\n    te = pd.to_numeric(df[\"EchoTime\"], errors=\"coerce\")\n    gre = df[\"ScanningSequence\"].fillna(\"\").str.upper().str.contains(\"GR\")\n    t1, t2, pdw = desc.str.contains(_T1_RX), desc.str.contains(_T2_RX), desc.str.contains(_PD_RX)\n\n    df[\"weight\"] = np.where(t1 & ~t2 & ~pdw, \"T1\",\n                     np.where(t2 & ~pdw, \"T2\",\n                       np.where(pdw, \"PD\",\n                         np.where(gre, \"GRE\",\n                           np.where(tr < 800, \"T1\",\n                             np.where(te > 60, \"T2\",\n                               np.where(tr >= 800, \"PD\", \"UNK\")))))))\n    df[\"fluid\"] = np.isin(df[\"weight\"], [\"PD\", \"T2\"])\n    df[\"px\"] = pd.to_numeric(\n        df[\"PixelSpacing\"].fillna(\"\").str.split(\"|\").str[0].replace(\"\", np.nan),\n        errors=\"coerce\")\n    return df\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 4. Slot selection\nAssigns exactly one series per study to each of the 6 canonical slots, breaking ties\ntoward the thickest stack; under legacy rules a missing T1 slot may fall back to any\nnon-fat-sat series in the plane.","metadata":{}},{"cell_type":"code","source":"def pick_slots(series_df, plane_map):\n    \"\"\"One series per slot per study.\n\n    Ties are broken toward the stack with the most slices: a thicker stack samples the\n    joint more densely, and the three-slice sampler below benefits from the margin.\n    \"\"\"\n    series_df = series_df.copy()\n    series_df[\"plane\"] = series_df[\"SeriesInstanceUID\"].map(plane_map)\n    out = {}\n    for study, g in series_df.groupby(\"StudyInstanceUID\"):\n        chosen = {}\n        for name, plane, fluid, fs in SLOTS:\n            sel = (g[\"plane\"] == plane) & (g[\"fatsat\"] == fs)\n            # fluid=None means \"do not condition on weighting\" - the public scheme,\n            # where the single provided flag stands in for both axes at once.\n            if fluid is not None:\n                sel &= (g[\"fluid\"] == fluid)\n            cand = g[sel]\n            # A slot with no series matching its predicate stays empty, and no substitute\n            # is admitted from a neighbouring predicate. Relaxing the weighting to fill a\n            # T1 slot would draw from the pool `SAG_FLUID_NOFS` selects from, since that\n            # pool is what remains once the weighting is dropped: over the training corpus\n            # it would put one series in two slots for 2383 of 4407 studies and leave 56%\n            # of the T1 slot holding PD or T2. The presence mask would then assert a\n            # sequence that was never acquired, and the per-diagnosis softmax of §6 would\n            # divide its attention across two identical slots, giving one acquisition\n            # about twice the weight it carries in a study that holds both. The mask is\n            # there to say a slot is absent, which is what an absent slot is.\n            if len(cand) == 0 and RULES[\"slot_fallback\"] and fluid is False:\n                # The relaxation the paragraph above rejects, reproduced because an\n                # imported member was fitted with its T1 slots filled this way: over half\n                # of that member's training studies had a T1 slot holding a series that\n                # is not T1. Leaving those slots empty would present it with a presence\n                # mask it never saw.\n                cand = g[(g[\"plane\"] == plane) & (~g[\"fatsat\"])]\n            if len(cand):\n                chosen[name] = cand.sort_values(\"n_slices\", ascending=False).iloc[0]\n        out[study] = chosen\n    return out\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 5. Slice order and pixel decode\nSorts each stack by patient-space position (filename order is anatomically meaningless),\nsamples the configured band of the stack, applies the fixed-millimetre centre crop,\nnormalises intensity to robust 1st–99th percentiles, and resizes to cache resolution.","metadata":{}},{"cell_type":"code","source":"ORDER_TAGS = [(0x0020, 0x0032), (0x0020, 0x0037), (0x0020, 0x0013)]\n\n# Series in which at least one sampled slice would not decode. A list rather than a\n# counter because appending is atomic under the reader threads, and reported rather than\n# swallowed: unreported, a decode failure is indistinguishable from a black knee.\nDECODE_FAILED = []\n\n\ndef cache_tag(rules=None):\n    \"\"\"The name a decoded cache is stored under.\n\n    It has to name everything that decides the pixels, not only their dimensions. Two\n    configurations that agree on resolution, slice count, crop and band but disagree on\n    how a slice is chosen produce different arrays of identical shape - so a tag built\n    from the dimensions alone lets the second attach to the first one's file and train\n    against pixels it never asked for, with nothing anywhere reporting a mismatch.\n\n    A native reading keeps the plain name, so caches decoded before the rules existed\n    stay valid; anything else earns a suffix.\n    \"\"\"\n    r = dict(RULES if rules is None else rules)\n    t = (f\"{CACHE_IMG}px_{CACHE_SLICES}sl_{int(CROP_MM)}mm_\"\n         f\"{SLICE_BAND[0]:.2f}-{SLICE_BAND[1]:.2f}\")\n    if {k: r.get(k, v) for k, v in RULES_NATIVE.items()} != RULES_NATIVE:\n        t += \"_\" + hashlib.md5(json.dumps(r, sort_keys=True).encode()).hexdigest()[:6]\n    return t\n\n\ndef _natural_key(name):\n    return tuple(int(x) if x.isdigit() else x.lower()\n                 for x in re.split(r\"(\\d+)\", str(name)))\n\n\ndef _order_dominant_axis(rec):\n    \"\"\"The slice order an imported member was fitted under.\n\n    It sorts on the raw patient coordinate along whichever axis varies most across the\n    stack, rather than on the projection onto the slice normal. The two differ by a sign,\n    not by a formula: measured over this corpus every sagittal series has a slice normal\n    with n_x in [-1.00, -0.98], so p.n is the negative of the raw x this sorts on and the\n    two stacks come out exactly reversed. Because the band sampler truncates rather than\n    rounds, its nine indices are not symmetric about the middle, so nine slices drawn from\n    a twenty-six slice stack under one order share two with the other.\n\n    Missing geometry falls back to `InstanceNumber` and then to a natural sort of the file\n    name, both at the same 80% threshold the imported pipeline used.\n    \"\"\"\n    files, d = rec[\"files\"], rec[\"dir\"]\n    rows = []\n    for pos, f in enumerate(files):\n        ipp = inst = None\n        try:\n            ds = pydicom.dcmread(os.path.join(d, f), force=True, stop_before_pixels=True,\n                                 specific_tags=[\"ImagePositionPatient\", \"InstanceNumber\"])\n            raw = getattr(ds, \"ImagePositionPatient\", None)\n            if raw is not None and len(raw) >= 3:\n                c = np.asarray(raw[:3], dtype=np.float64)\n                if np.isfinite(c).all():\n                    ipp = c\n            n = getattr(ds, \"InstanceNumber\", None)\n            if n is not None:\n                inst = float(n)\n        except Exception:\n            pass\n        rows.append((f, ipp, inst, pos))\n\n    placed = [r for r in rows if r[1] is not None]\n    need = max(2, int(0.8 * len(rows)))\n    if len(placed) >= need:\n        xyz = np.stack([r[1] for r in placed])\n        axis = int(np.argmax(np.ptp(xyz, axis=0)))\n        spare = float(np.nanmedian(xyz[:, axis]))\n        rows.sort(key=lambda r: (float(r[1][axis]) if r[1] is not None else spare,\n                                 r[2] if r[2] is not None else float(\"inf\"), r[3]))\n    elif sum(r[2] is not None for r in rows) >= need:\n        rows.sort(key=lambda r: (r[2] if r[2] is not None else float(\"inf\"), r[3]))\n    else:\n        rows.sort(key=lambda r: _natural_key(r[0]))\n    return [r[0] for r in rows], True\n\n\ndef order_slices(rec):\n    \"\"\"Return the series' files sorted along the through-plane axis.\n\n    A DICOM file name here is a SOP Instance UID, which is assigned arbitrarily. Sorting\n    by it therefore produces an order uncorrelated with anatomy - measured over one\n    series, Spearman between file-name rank and physical position is 0.009, i.e. none.\n    Anything that assumes the file order means something is then operating on noise: the\n    three channels of a \"2.5D\" input are three unrelated views rather than neighbouring\n    slices, \"the middle of the stack\" is a random subset, and reversing slice order to\n    normalise laterality reverses nothing meaningful.\n\n    The physical order is recoverable exactly. Each slice carries its position in patient\n    coordinates and the in-plane axes; projecting the position onto the slice normal\n    gives a signed through-plane coordinate, monotonic along the stack:\n\n        n = r_x  x  r_y ,      k = p . n\n\n    `InstanceNumber` is the fallback. It usually tracks the projection up to sign, but\n    interleaved and multi-echo acquisitions need not number slices in the order they\n    occupy in space - but the projection is signed in patient\n    coordinates, which is what laterality normalisation needs.\n    \"\"\"\n    if RULES[\"order\"] == \"dominant_axis\":\n        return _order_dominant_axis(rec)\n    files, d = rec[\"files\"], rec[\"dir\"]\n    keyed = []\n    for f in files:\n        k = None\n        try:\n            ds = pydicom.dcmread(os.path.join(d, f), force=True, stop_before_pixels=True,\n                                 specific_tags=ORDER_TAGS)\n            iop = np.asarray(ds.ImageOrientationPatient, dtype=float)\n            ipp = np.asarray(ds.ImagePositionPatient, dtype=float)\n            k = float(np.dot(ipp, np.cross(iop[:3], iop[3:])))\n        except Exception:\n            try:\n                k = float(ds.InstanceNumber)\n            except Exception:\n                k = None\n        keyed.append((k, f))\n    if any(k is None for k, _ in keyed):\n        # A series with no usable geometry keeps its arbitrary order; that is worse than\n        # sorting but better than dropping the series, and it is logged as a count.\n        return files, False\n    return [f for _, f in sorted(keyed, key=lambda t: t[0])], True\n\n\ndef read_slot(rec, n_slice=None, out_size=None):\n    \"\"\"`n_slice` physically spread slices from one series, at `out_size` pixels.\n\n    Returns uint8 [n_slice, out, out] normalised per-series to its 1st-99th\n    percentile. Percentiles rather than min/max because MR intensity has no absolute\n    scale and a single bright vessel would otherwise compress the whole dynamic range.\n\n    Reading is the expensive half of this pipeline, so the caller reads once at the\n    largest configuration it needs and derives the smaller ones from the returned buffer\n    rather than re-reading.\n    \"\"\"\n    n_slice = GROUP if n_slice is None else n_slice\n    out_size = IMG if out_size is None else out_size\n    files, d, px = rec.get(\"ordered\") or rec[\"files\"], rec[\"dir\"], rec[\"px\"]\n    n = len(files)\n    if n == 0:\n        return None\n    # Spread the samples over a central band of the stack: the outermost slices of a knee\n    # series are mostly soft tissue outside the joint. The band is a constant rather than\n    # a literal because how much of the stack is worth reading depends on how many slices\n    # are being taken - at three the middle is all that fits, while at sixteen the ends\n    # are worth having, and a Baker cyst sits at the posteromedial end of a sagittal one.\n    lo, hi = int(SLICE_BAND[0] * (n - 1)), int(SLICE_BAND[1] * (n - 1))\n    idx = np.unique(np.linspace(lo, hi, n_slice).astype(int)) if hi > lo else np.array([n // 2])\n    while len(idx) < n_slice:\n        idx = np.append(idx, idx[-1])\n\n    planes = []\n    for i in idx[:n_slice]:\n        try:\n            ds = pydicom.dcmread(os.path.join(d, files[int(i)]), force=True)\n            a = ds.pixel_array.astype(np.float32)\n            sl = float(getattr(ds, \"RescaleSlope\", 1) or 1)\n            ic = float(getattr(ds, \"RescaleIntercept\", 0) or 0)\n            a = a * sl + ic\n        except Exception:\n            a = None                      # no shape is known here; see below\n        planes.append(a)\n\n    # A slice that would not decode has no shape of its own, and inventing one is how a\n    # single unreadable file erases a whole series: a substitute allocated at the resize\n    # target while the decoded slices are still native makes the shape check below take\n    # the substitute as the authority and zero the good slices with it, leaving a black\n    # slot that the presence mask still reports as acquired.\n    #\n    # A failure is instead filled from the nearest slice that did decode - the same\n    # convention the sampler already uses when the band holds fewer distinct slices than\n    # were asked for - and a series where nothing decodes is reported absent, which the\n    # mask can express, rather than black, which it cannot.\n    got = [k for k, p in enumerate(planes) if p is not None]\n    if RULES[\"decode_fill\"] == \"zero\":\n        # What an imported member was fitted with: a failure becomes a zero plane at the\n        # resize target, which the shape check below then propagates to the whole slot.\n        # It is the behaviour the paragraph above describes and rejects, kept here only\n        # because that member's weights were learned against slots blacked out this way.\n        if not got:\n            DECODE_FAILED.append(rec.get(\"SeriesInstanceUID\", d))\n        planes = [np.zeros((out_size, out_size), np.float32) if p is None else p\n                  for p in planes]\n        got = list(range(len(planes)))\n    if not got:\n        DECODE_FAILED.append(rec.get(\"SeriesInstanceUID\", d))\n        return None\n    if len(got) < len(planes):\n        DECODE_FAILED.append(rec.get(\"SeriesInstanceUID\", d))\n        for k, p in enumerate(planes):\n            if p is None:\n                planes[k] = planes[min(got, key=lambda j: abs(j - k))]\n\n    # Slices of one series can still differ in matrix size - multi-echo and some\n    # reformats do - and those are genuinely not stackable.\n    shp = planes[0].shape\n    planes = [p if p.shape == shp else np.zeros(shp, np.float32) for p in planes]\n    vol = np.stack(planes)\n\n    # constant physical extent, then resize: PixelSpacing varies 3.4x across the corpus\n    if px and np.isfinite(px) and px > 0:\n        want = int(round(CROP_MM / px))\n        h, w = shp\n        if 16 < want < min(h, w):\n            cy, cx = h // 2, w // 2\n            half = want // 2\n            vol = vol[:, max(0, cy - half):cy + half, max(0, cx - half):cx + half]\n\n    lo_v, hi_v = np.percentile(vol, [1, 99])\n    vol = np.clip((vol - lo_v) / max(hi_v - lo_v, 1e-6), 0, 1)\n\n    t = torch.from_numpy(np.ascontiguousarray(vol)).unsqueeze(0)\n    t = F.interpolate(t, size=(out_size, out_size), mode=\"bilinear\", align_corners=False)\n    # uint8, not float32. These buffers queue up between the reader threads and the\n    # encoder, and at this size a float32 slot-series is several megabytes. Intensity is\n    # already normalised into [0, 1] here, so eight bits cost nothing that a bilinear\n    # resize has not already cost, and the queue is a quarter the size.\n    return (t.squeeze(0) * 255).round().clamp(0, 255).to(torch.uint8)\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Laterality normalisation\nMaps every right knee onto a left-knee convention: horizontal mirror for coronal and\naxial views, channel-order reversal for sagittal stacks (which are not mirror images).","metadata":{}},{"cell_type":"code","source":"def normalise_laterality(img, plane, lat):\n    \"\"\"Map every knee onto a left-knee convention.\n\n    Coronal and axial views mirror under a horizontal flip. Sagittal stacks are not\n    mirror images of each other - the slice order runs medial-to-lateral in opposite\n    directions - so the channel order is reversed instead.\n    \"\"\"\n    if lat != \"R\":\n        return img\n    if plane in (\"Coronal\", \"Axial\"):\n        return torch.flip(img, dims=[-1])\n    return torch.flip(img, dims=[0])\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 6. Study cache\nDecodes every (study, slot) once into an in-memory uint8 array of shape\n[studies, slots, slices, px, px] plus a slot-presence mask — read many times, paid for\nonce.","metadata":{}},{"cell_type":"code","source":"# Where the geometric slice order may be remembered between runs. Unset on the platform,\n# because each run gets a fresh machine and there is nothing to remember; set off it,\n# where the same corpus is cached again at every resolution and slice count and the order\n# is a function of neither. It is opt-in so that the scored run's behaviour is decided by\n# the code rather than by whether a file happens to be lying about.\nORDER_CACHE = os.environ.get(\"RSNA_ORDER_CACHE\") or None\n\n\ndef build_cache(slot_map, plane_map, lat_map, tag):\n    \"\"\"Decode every (study, slot) once into an in-memory uint8 array.\n\n    Fine-tuning revisits the same pixels every epoch. Reading them from the mount each\n    time would make the epoch count a function of I/O rather than of learning, so they\n    are decoded once and held as bytes: intensity has already been normalised into\n    [0, 1], and eight bits cost nothing a bilinear resize has not already cost.\n\n    CACHE_SLICES positions are kept per slot, which the training loop reads as N_GROUP\n    groups of GROUP consecutive channels.\n    \"\"\"\n    studies = sorted(slot_map)\n    sidx = {s: i for i, s in enumerate(studies)}\n    cache = np.zeros((len(studies), N_SLOT, CACHE_SLICES, IMG, IMG), np.uint8)\n    mask = np.zeros((len(studies), N_SLOT), np.float32)\n    log(f\"{tag}: cache {cache.shape} = {cache.nbytes / 1024 ** 3:.1f} GB\")\n\n    jobs = [(st, k, plane, slot_map[st][name])\n            for st in studies\n            for k, (name, plane, _, _) in enumerate(SLOTS)\n            if name in slot_map[st]]\n    n_job = len(jobs)\n\n    # Ordering first, and as its own pass. It reads one header per slice of every chosen\n    # series - far more file opens than the decode that follows - and on a network mount\n    # that is latency, not work, so it gets its own wider pool.\n    t_ord = time.time()\n    n_slice_total = sum(len(j[3][\"files\"]) for j in jobs)\n    log(f\"{tag}: ordering {len(jobs)} slot-series ({n_slice_total} slice headers)\")\n    ok = done = 0\n    CHUNK_O = 1024\n\n    # A remembered order, when one is offered. The projection depends on the DICOM\n    # geometry alone, so it is the same at every resolution and every slice count, and\n    # it costs one header read per slice - the largest single cost in this pass. An entry\n    # is validated by the number of files present, so a tree that has changed under it is\n    # recomputed rather than trusted: order is derived data, and a stale entry would be\n    # invisible in the way that matters most.\n    seen = {}\n    if ORDER_CACHE and Path(ORDER_CACHE).is_file():\n        try:\n            import json as _json\n            seen = _json.loads(Path(ORDER_CACHE).read_text())\n        except (OSError, ValueError):\n            seen = {}\n        hit = 0\n        for _, _, _, rec in jobs:\n            e = seen.get(rec[\"SeriesInstanceUID\"])\n            if e and len(e[\"files\"]) == len(rec[\"files\"]):\n                rec[\"ordered\"] = e[\"files\"]\n                ok += int(e[\"good\"])\n                hit += 1\n        jobs = [j for j in jobs if \"ordered\" not in j[3]]\n        log(f\"{tag}: {hit} slot-series ordered from {ORDER_CACHE}, {len(jobs)} to read\")\n\n    with ThreadPoolExecutor(max_workers=ORDER_THREADS) as pool:\n        for c0 in range(0, len(jobs), CHUNK_O):\n            block = jobs[c0:c0 + CHUNK_O]\n            for (_, _, _, rec), (files, good) in zip(\n                    block, pool.map(lambda j: order_slices(j[3]), block)):\n                rec[\"ordered\"] = files\n                ok += int(good)\n                done += 1\n                if ORDER_CACHE:\n                    seen[rec[\"SeriesInstanceUID\"]] = {\"files\": files, \"good\": bool(good)}\n            # The ceiling is whichever comes first: the pass's own budget, or the share\n            # of what is left of the run that it may take. The second is what makes the\n            # first safe to set generously - a mount slow enough to matter cannot spend\n            # the training time, because the budget shrinks as the run does.\n            budget = min(ORDER_BUDGET_S, max(60.0, (TIME_BUDGET - (time.time() - T0)) * 0.35))\n            if time.time() - t_ord > budget:\n                log(f\"{tag}: ordering budget spent at {done}/{len(jobs)}; \"\n                    f\"the rest keep file order\")\n                break\n    if ORDER_CACHE and done:\n        import json as _json\n        _t = Path(ORDER_CACHE).with_suffix(\".tmp\")\n        _t.write_text(_json.dumps(seen))\n        _t.replace(Path(ORDER_CACHE))\n    log(f\"{tag}: ordered {ok}/{n_job} by geometry \"\n        f\"({n_job - ok} kept arbitrary) in {time.time() - t_ord:.0f}s\")\n\n    jobs = [(st, k, plane, slot_map[st][name])\n            for st in studies\n            for k, (name, plane, _, _) in enumerate(SLOTS)\n            if name in slot_map[st]]\n    log(f\"{tag}: decoding {len(jobs)} slot-series\")\n    n_failed_before = len(DECODE_FAILED)\n\n    CHUNK = 512\n    done = 0\n    with ThreadPoolExecutor(max_workers=PIX_THREADS) as pool:\n        for c0 in range(0, len(jobs), CHUNK):\n            block = jobs[c0:c0 + CHUNK]\n            for (st, k, plane, _), img in zip(\n                    block, pool.map(lambda j: read_slot(j[3], CACHE_SLICES, IMG), block)):\n                done += 1\n                if img is None:\n                    continue\n                cache[sidx[st], k] = normalise_laterality(img, plane,\n                                                          lat_map.get(st)).numpy()\n                mask[sidx[st], k] = 1.0\n            if done % 4096 < CHUNK:\n                log(f\"  {tag} {done}/{len(jobs)}\")\n            if time.time() - T0 > TIME_BUDGET:\n                log(f\"  {tag}: time budget reached during decode\")\n                break\n    n_failed = len(DECODE_FAILED) - n_failed_before\n    log(f\"{tag}: {int(mask.sum())}/{len(jobs)} slots filled\"\n        + (f\"; {n_failed} series had a slice that would not decode\" if n_failed else \"\"))\n    gc.collect()\n    return studies, cache, mask\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 7. Model\n`SlotHead`: one learned attention query per diagnosis over the 6 slot embeddings, with\nan optional fixed anatomical prior. `Model`: DINOv2 encoder whose token grid is pooled\ninto CLS + patch-mean (+ focal top-k where the checkpoint used it), then fed to the\nhead. `build_model` loads the encoder variant a checkpoint asks for.","metadata":{}},{"cell_type":"code","source":"class SlotHead(nn.Module):\n    \"\"\"Per-diagnosis attention over the slot embeddings of one study.\n\n    Each finding is read on particular sequences - cruciates sagittally, collateral\n    ligaments and the meniscal body coronally, patellar cartilage axially - so pooling\n    the slots identically would dilute the one that carries the evidence with the rest.\n\n    The aggregation is deliberately this simple. With a study-level label there is no\n    signal telling the model which part of a study matters, so extra attention\n    parameters below the slot level would have nothing to learn from and would spend\n    their capacity fitting noise.\n    \"\"\"\n\n    def __init__(self, dim, n_slot, n_out, hidden=256, p=0.2, prior=False):\n        super().__init__()\n        self.proj = nn.Sequential(nn.LayerNorm(dim), nn.Linear(dim, hidden), nn.GELU())\n        self.slot_emb = nn.Parameter(torch.randn(n_slot, hidden) * 0.02)\n        self.query = nn.Parameter(torch.randn(n_out, hidden) * 0.02)\n        self.drop = nn.Dropout(p)\n        self.out = nn.Linear(hidden, n_out)\n        self.hidden = hidden\n        # An imported member carries a fixed per-(diagnosis, slot) tilt on the attention\n        # logits, set from the anatomy table below rather than learned. It is a buffer, so\n        # it travels in the state dict and must exist for that member to load; exp(0.55)\n        # gives a preferred slot about 1.73x the weight of an unpreferred one, which\n        # biases the softmax without ever excluding a slot.\n        p_ = torch.zeros(n_out, n_slot)\n        if prior and n_slot == len(SLOTS) and n_out == len(TARGETS):\n            for t, slots in SLOT_PRIOR_TABLE.items():\n                if t in TARGETS:\n                    p_[TARGETS.index(t), list(slots)] = SLOT_PRIOR_STRENGTH\n        self.prior = prior\n        if prior:\n            self.register_buffer(\"slot_prior\", p_)\n\n    def forward(self, x, mask):\n        h = self.proj(x) + self.slot_emb\n        att = torch.einsum(\"bsh,oh->bos\", h, self.query) / self.hidden ** 0.5\n        if self.prior:\n            att = att + self.slot_prior.unsqueeze(0)\n        att = att.masked_fill(mask.unsqueeze(1) < 0.5, -1e4).softmax(-1)\n        ctx = self.drop(torch.einsum(\"bos,bsh->boh\", att, h))\n        return (ctx * self.out.weight.unsqueeze(0)).sum(-1) + self.out.bias\n\n\nclass Model(nn.Module):\n    \"\"\"Encoder plus head, trained end to end.\n\n    A study arrives as a bag of slot images. The bag is flattened for the encoder and\n    folded back before the head, so the encoder never sees the study structure and the\n    head never sees pixels.\n    \"\"\"\n\n    def __init__(self, backbone, dim, pool=\"cls_mean\", prior=False):\n        super().__init__()\n        self.backbone = backbone\n        self.pool = pool\n        self.head = SlotHead(dim * POOL_PARTS[pool], N_SLOT, len(TARGETS), prior=prior)\n        self.register_buffer(\"mean\", torch.tensor([0.485, 0.456, 0.406]).view(1, 3, 1, 1))\n        self.register_buffer(\"std\", torch.tensor([0.229, 0.224, 0.225]).view(1, 3, 1, 1))\n\n    def forward(self, imgs, mask, img_size=None):\n        B, S = imgs.shape[:2]\n        x = imgs.reshape(B * S, *imgs.shape[2:]).float().div_(255.0)\n        if img_size is not None and img_size != x.shape[-1]:\n            # The cache is held at the highest resolution any configuration needs; the\n            # rest downsample from it, so every configuration sees the same pixels\n            # through a different sampling grid rather than a different crop.\n            x = F.interpolate(x, size=(img_size, img_size), mode=\"bilinear\",\n                              align_corners=False)\n        x = (x - self.mean) / self.std\n        out = self.backbone(pixel_values=x).last_hidden_state\n        patch = out[:, 1:]\n        parts = [out[:, 0], patch.mean(1)]\n        if self.pool == \"cls_mean_focal\":\n            # The upper tail of each channel over the patch grid, taken per channel\n            # rather than by selecting whole patches: a finding occupies a small part of\n            # the field, so a plain mean over 256 patches dilutes it by two orders of\n            # magnitude, and this keeps the top eighth of each channel's responses.\n            k = max(1, patch.shape[1] // 8)\n            parts.append(patch.topk(k, dim=1).values.mean(1))\n        feat = torch.cat(parts, dim=1).reshape(B, S, -1)\n        return self.head(feat, mask)\n\n\ndef build_model(unfreeze_last, source=None, variant=\"small\", pool=\"cls_mean\",\n                prior=False):\n    \"\"\"Load the encoder and open the last `unfreeze_last` blocks for training.\n\n    The early blocks of a self-supervised transformer are generic edge and texture\n    filters; the late blocks carry semantics. Opening only the late ones is the cautious\n    choice - there may not be enough supervision here to improve the early ones and there\n    is certainly enough to damage them - but how far the line should sit is a question\n    the corpus has to answer rather than the intuition.\n\n    `source` names where the weights come from. Left unset it is the attached model\n    directory, which is the only thing available here. It is a parameter so that a run\n    off the platform builds the same object from the same code rather than from a second\n    definition that has to be kept in step by hand.\n    \"\"\"\n    from transformers import AutoModel\n    p = source if source is not None else find_dinov2(variant)\n    if p is None:\n        raise FileNotFoundError(\"DINOv2 weights not attached\")\n    bb = AutoModel.from_pretrained(str(p))\n    n_layer = len(bb.encoder.layer)\n    for prm in bb.parameters():\n        prm.requires_grad = False\n    for blk in bb.encoder.layer[max(0, n_layer - unfreeze_last):]:\n        for prm in blk.parameters():\n            prm.requires_grad = True\n    for prm in bb.layernorm.parameters():\n        prm.requires_grad = True\n    dim = bb.config.hidden_size\n    trainable = sum(p.numel() for p in bb.parameters() if p.requires_grad)\n    log(f\"backbone: {n_layer} blocks, last {unfreeze_last} trainable \"\n        f\"({trainable / 1e6:.1f}M params), feature dim {dim * POOL_PARTS[pool]}\")\n    return Model(bb, dim, pool=pool, prior=prior)\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Group slicing and jitter augmentation\n`take_group` cuts 3 consecutive cached slices as one pseudo-RGB input. `augment` applies\nthe small rigid jitter (rotation/scale/shift, intensity scale — deliberately no flips)\nthat training used; at inference it powers the optional jitter-TTA views.","metadata":{}},{"cell_type":"code","source":"def take_group(cache_rows, g):\n    \"\"\"Slice GROUP consecutive channels out of the cached slices.\"\"\"\n    return cache_rows[:, :, g * GROUP:(g + 1) * GROUP]\n\n\ndef augment(imgs):\n    \"\"\"A small rigid jitter and an intensity scale, applied to a whole bag at once.\n\n    Neither flip is available here, and for different reasons. A horizontal flip would\n    reintroduce the nuisance axis that the laterality normalisation removed - it would\n    undo, once per batch, what the header pass was run to establish.\n\n    A vertical flip is not a nuisance axis at all. A knee is acquired in a canonical\n    orientation, and no study in this corpus looks like its own vertical mirror. An\n    augmentation is meant to cover directions along which the label does not change; this\n    one moves the input off the distribution the encoder will be asked about, which is a\n    different thing. Where a finding sits in the frame is also information rather than\n    noise - a Baker cyst is identified by lying in the popliteal fossa, not by its\n    appearance alone.\n\n    What is left is jitter that no label depends on: a few degrees of rotation, a few\n    per cent of scale and translation. That still prevents memorising the exact framing,\n    which is what an augmentation is for, while leaving the anatomy where it was.\n    \"\"\"\n    # A bag arrives as [study, slot, GROUP, IMG, IMG]: five axes, not four. The warp is\n    # a 2-D operation, so the two leading axes are folded together and restored after -\n    # every slot image is an independent acquisition and gets its own jitter.\n    lead = imgs.shape[:-3]\n    x = imgs.reshape(-1, *imgs.shape[-3:]).float()\n    n, dev = x.shape[0], x.device\n\n    rot = (torch.rand(n, device=dev) - 0.5) * 2 * (AUG_ROT_DEG * np.pi / 180)\n    # Zoom in only. `border` padding repeats the edge row outward, and the edge of this\n    # crop is where the popliteal fossa sits; zooming out would fabricate tissue exactly\n    # where a Baker cyst is looked for.\n    sc = 1.0 + torch.rand(n, device=dev) * AUG_SCALE\n    tx = (torch.rand(n, device=dev) - 0.5) * 2 * AUG_SHIFT\n    ty = (torch.rand(n, device=dev) - 0.5) * 2 * AUG_SHIFT\n    cos, sin = torch.cos(rot) / sc, torch.sin(rot) / sc\n    theta = torch.zeros(n, 2, 3, device=dev, dtype=torch.float32)\n    theta[:, 0, 0], theta[:, 0, 1], theta[:, 0, 2] = cos, -sin, tx\n    theta[:, 1, 0], theta[:, 1, 1], theta[:, 1, 2] = sin, cos, ty\n    grid = F.affine_grid(theta, x.shape, align_corners=False)\n    x = F.grid_sample(x, grid, mode=\"bilinear\", padding_mode=\"border\", align_corners=False)\n\n    scale = 1.0 + (torch.rand(n, 1, 1, 1, device=dev) - 0.5) * 2 * AUG_INTENSITY\n    x = (x * scale).clamp(0, 255)\n    return x.reshape(*lead, *x.shape[-3:]).to(imgs.dtype)","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 8. Fingerprints, TTA and the dual-GPU member queue\nThe heart of the run. `fingerprint`/`check_fingerprint` verify each checkpoint computes\nwhat it computed when trained. `predict_member` averages a member over its TTA windows\n(plus jittered views when granted). `legacy_group_members` wraps the optional five-fold\nbundle as reduced-weight members under the legacy pixel rules. `infer_from_package`\ngroups members by pixel contract, builds one cache per group, and lets both GPUs pull\nmembers from a shared queue — model loading serialised under a lock (HF loading is not\nthread-safe), inference parallel, failed members retried on the peer GPU, and the\nweighted rank-mean submission rewritten after every banked member.","metadata":{}},{"cell_type":"code","source":"FINGERPRINT_TOL = 2e-3\n\n\ndef fingerprint(model, dev, img_size, n_slot=None, group=None, seed=None):\n    \"\"\"The model's output on a fixed synthetic bag, as a portable identity.\n\n    Weights that are loaded but read through the wrong preprocessing produce predictions,\n    not errors. The submission is well formed, the log says nothing, and the difference is\n    a number no output of the run reveals. Scaling that never happens, or happens twice,\n    is enough on its own and changes no shape anywhere.\n\n    So a set of weights carries the answer it gave to a question with no data in it. The\n    input is generated from a seed rather than read, so it is the same on any machine, and\n    it is pushed through the whole forward path - the byte scaling, the ImageNet\n    normalisation, the resize, the encoder, the slot attention. Any of those differing\n    moves the output by order one. Numerics differing between two GPUs moves it by about\n    1e-5, which is why the tolerance sits between them rather than at zero.\n\n    This checks that the model computes what it computed when it was fitted. It cannot\n    check that the pixels reaching it are the right pixels; `read_slot` and the header\n    pass answer to their own tests.\n    \"\"\"\n    n_slot = N_SLOT if n_slot is None else n_slot\n    group = GROUP if group is None else group\n    seed = SEED if seed is None else seed\n    g = torch.Generator().manual_seed(seed)\n    imgs = torch.randint(0, 256, (2, n_slot, group, img_size, img_size),\n                         generator=g, dtype=torch.uint8).to(dev)\n    mask = torch.ones(2, n_slot, device=dev)\n    mask[1, -1] = 0.0                       # exercise the masked branch of the softmax\n    was_training = model.training\n    model.eval()\n    with torch.no_grad():\n        # float32 throughout: autocast would make the value depend on which device\n        # happened to run it, and the point of the number is that it does not.\n        out = model(imgs, mask, img_size).float().cpu().numpy()\n    if was_training:\n        model.train()\n    return out\n\n\ndef check_fingerprint(model, dev, img_size, expected, tol=FINGERPRINT_TOL, tag=\"\"):\n    \"\"\"Compare against a stored fingerprint; raise when the model is not the same map.\"\"\"\n    got = fingerprint(model, dev, img_size)\n    exp = np.asarray(expected, np.float32)\n    if got.shape != exp.shape:\n        raise WeightsError(f\"{tag}fingerprint shape {got.shape} != stored {exp.shape}: \"\n                           f\"the architecture is not the one these weights were fitted to\")\n    d = float(np.abs(got - exp).max())\n    if d > tol:\n        raise WeightsError(\n            f\"{tag}fingerprint differs by {d:.4g} (tolerance {tol:g}). The weights load \"\n            f\"but do not compute what they computed when fitted - preprocessing, \"\n            f\"resolution or architecture has moved between the two runs.\")\n    log(f\"{tag}fingerprint matches within {d:.2g}\")\n    return d\n\n\nclass WeightsError(RuntimeError):\n    \"\"\"Raised when attached weights cannot be trusted to be the ones that were fitted.\n\n    Deliberately fatal for the same reason as LabelSourceError: a run that predicts from\n    a mismatched model completes, writes a plausible submission, and differs from a\n    correct one only in a number no output of the run reveals.\n    \"\"\"\n\n\ndef find_weights(name=\"manifest.json\"):\n    \"\"\"Locate a mounted weights package, or return None if none is attached.\n\n    Same shape as `find_label_table`: the notebook must keep working for a reader who\n    attaches nothing, so absence is a path rather than an error. What must not be silent\n    is a package that is attached and unusable, and that is what `load_weights` refuses.\n    \"\"\"\n    import json\n    base = Path(\"/kaggle/input\")\n    if not base.is_dir():\n        return None\n    for root, dirs, files in os.walk(base):\n        dirs[:] = [d for d in dirs if d not in (\"train_series\", \"test_series\")]\n        if name not in files:\n            continue\n        # The manifest decides, not the filenames beside it. Testing for a naming\n        # convention makes the search agree with whatever the packager happened to call\n        # its files last, which is a second definition of what a package is.\n        try:\n            man = json.loads((Path(root) / name).read_text())\n        except (OSError, ValueError):\n            continue\n        if isinstance(man.get(\"members\"), list) and man[\"members\"]:\n            missing = [m[\"file\"] for m in man[\"members\"]\n                       if not (Path(root) / m[\"file\"]).is_file()]\n            if missing:\n                raise WeightsError(\n                    f\"{root} holds a manifest listing {len(man['members'])} members but \"\n                    f\"{len(missing)} of their files are absent (first {missing[0]!r})\")\n            return Path(root)\n    return None\n\n\n# How a member is read at inference. Overlapping windows over the slices the cache\n# already holds cost forward passes and no extra decoding, which is the cheap direction\n# to spend; and averaging probabilities rather than logits is an arithmetic mean of risk\n# rather than a geometric mean of odds, which orders studies differently. Both were\n# chosen by measuring them on the folds each member held out rather than by argument.\nTTA_OVERLAP = True\nTTA_POOL = \"prob\"\n\n# Per-target pooling over TTA windows. A fracture, a meniscal tear or a Baker's cyst\n# lives on a few slices: when one window sees it and nine do not, the mean dilutes the\n# detection, so the focal findings take the max (or top-k mean) over windows. Diffuse\n# findings -- the three OA compartments, effusion, synovitis -- spread across the joint\n# and keep the plain mean. Labels not listed keep the original probability mean exactly.\nTTA_TARGET_POOL = {\n    \"Fracture\": \"max\",\n    \"Contusion\": \"max\",\n    \"Medial Meniscus\": \"max\",\n    \"Lateral Meniscus\": \"max\",\n    \"ACL\": \"top3\",\n    \"MCL\": \"top3\",\n    \"Baker's\": \"max\",\n}\n\n\ndef window_starts(n_slice, group, overlap=None):\n    \"\"\"Where each TTA window begins.\"\"\"\n    overlap = TTA_OVERLAP if overlap is None else overlap\n    if overlap and n_slice >= group:\n        return list(range(n_slice - group + 1))\n    return [g * group for g in range(max(n_slice // group, 1))]\n\n\n@torch.no_grad()\ndef predict_member(model, cache, mask, idx, dev, img_size, group=None, pool=None,\n                   starts=None, jitter=False):\n    \"\"\"One member's predictions, pooled over its TTA windows.\n\n    `jitter` adds one extra view per window drawn from the same rigid-jitter family the\n    members were trained under; jittered views are averaged WITHIN their window first,\n    so jitter smooths a window's estimate without blunting the across-window max that\n    the focal findings rely on.\n    \"\"\"\n    group = GROUP if group is None else group\n    pool = TTA_POOL if pool is None else pool\n    starts = window_starts(cache.shape[2], group) if starts is None else list(starts)\n    if not starts:\n        raise ValueError(\"predict_member was given no windows to average over\")\n    target_idx = {t: j for j, t in enumerate(TARGETS)}\n    unknown = set(TTA_TARGET_POOL) - set(target_idx)\n    if unknown:\n        raise ValueError(f\"unknown target(s) in TTA_TARGET_POOL: {unknown}\")\n    model.eval()\n    out = []\n    for b in range(0, len(idx), EVAL_BATCH):\n        sel = idx[b:b + EVAL_BATCH]\n        m = torch.from_numpy(mask[sel]).to(dev)\n        acc, n_pass, win_probs = None, 0, []\n        for st in starts:\n            rows = torch.from_numpy(\n                np.ascontiguousarray(cache[sel, :, st:st + group])).to(dev)\n            views = [rows] + ([augment(rows)] if jitter else [])\n            wacc = None\n            for view in views:\n                with torch.autocast(\"cuda\", enabled=dev.type == \"cuda\"):\n                    z = model(view, m, img_size).float()\n                p = torch.sigmoid(z)\n                v = z if pool == \"logit\" else p\n                acc = v if acc is None else acc + v\n                n_pass += 1\n                wacc = p if wacc is None else wacc + p\n            win_probs.append(wacc / len(views))\n        v = acc / n_pass\n        final = torch.sigmoid(v) if pool == \"logit\" else v\n        if TTA_TARGET_POOL:\n            probs = torch.stack(win_probs, dim=0)      # [window, batch, target]\n            for target, mode in TTA_TARGET_POOL.items():\n                j = target_idx[target]\n                col = probs[:, :, j]\n                if mode == \"mean\":\n                    final[:, j] = col.mean(dim=0)\n                elif mode == \"logit_mean\":\n                    final[:, j] = torch.sigmoid(torch.logit(\n                        col.clamp(1e-6, 1 - 1e-6)).mean(dim=0))\n                elif mode == \"max\":\n                    final[:, j] = col.max(dim=0).values\n                elif mode in (\"top2\", \"top3\"):\n                    k = min(int(mode[3:]), col.shape[0])\n                    final[:, j] = col.topk(k, dim=0).values.mean(dim=0)\n                else:\n                    raise ValueError(f\"unknown TTA pool mode {mode!r}\")\n        out.append(final.cpu().numpy())\n    return np.concatenate(out) if out else np.zeros((0, len(TARGETS)), np.float32)\n\n\n\n# HF from_pretrained mutates process-global state and is not guaranteed thread-safe, so\n# model construction, weight loading and the fingerprint check are serialised; only\n# inference -- the expensive part -- runs on both devices at once.\nBUILD_LOCK = threading.Lock()\nSTATE_LOCK = threading.Lock()\n\n# An independently trained five-fold bundle joins the vote at reduced weight: a second\n# training run (public 0.836 on its own) adds decorrelated errors, which is the only thing an\n# inference-only run can add that the 20-member package does not already have. 0.5 per\n# fold puts the bundle at ~11% of the total vote -- roughly its quality gap.\nLEGACY_BUNDLE_FILE = \"rsna_20260807_v1.pt\"\nLEGACY_WEIGHT = 0.5\n\n\ndef find_legacy_bundle():\n    base = Path(\"/kaggle/input\")\n    if not base.is_dir():\n        return None\n    for root, dirs, files in os.walk(base):\n        dirs[:] = [d for d in dirs if d not in (\"train_series\", \"test_series\")]\n        if LEGACY_BUNDLE_FILE in files:\n            return Path(root) / LEGACY_BUNDLE_FILE\n    return None\n\n\ndef legacy_group_members():\n    \"\"\"The five-fold bundle as extra, lower-weight members under RULES_LEGACY.\n\n    The package's legacy pixel rules exist precisely to reproduce what this bundle was\n    fitted on (dominant-axis slice order, corner-x laterality at 5 mm, T1 slot fallback,\n    zero decode fill, 160 mm crop, central 60% band). The bundle predates fingerprints,\n    which is accepted loudly and priced into its reduced weight; a fold whose state dict\n    does not load, or whose predictions are degenerate, is dropped and costs its own\n    vote only.\n    \"\"\"\n    p = find_legacy_bundle()\n    if p is None:\n        log(\"no legacy bundle attached; blending skipped\")\n        return {}\n    try:\n        b = torch.load(p, map_location=\"cpu\", weights_only=False)\n        folds = b.get(\"fold_states\") or []\n        b_slots = [tuple(s)[0] for s in b.get(\"slots\", SLOTS)]\n        if list(b.get(\"targets\", TARGETS)) != TARGETS or b_slots != [s[0] for s in SLOTS]:\n            log(f\"legacy bundle {p.name}: target/slot contract differs; blending skipped\")\n            return {}\n        gr, n_gr = int(b.get(\"group\", 3)), int(b.get(\"n_group\", 3))\n        variant = str(b.get(\"model_variant\", \"dinov2-small\")).split(\"-\")[-1]\n        key = json.dumps({\"img\": int(b.get(\"img\", 224)), \"group\": gr,\n                          \"slices\": gr * n_gr, \"crop_mm\": 160.0, \"band\": [0.20, 0.80],\n                          \"rules\": RULES_LEGACY, \"slots\": [s[0] for s in SLOTS]},\n                         sort_keys=True)\n        ms = [{\"id\": f\"legacy-f{f.get('fold', k)}\", \"fold\": f.get(\"fold\", k),\n               \"state\": f[\"state_dict\"], \"holdout\": None, \"weight\": LEGACY_WEIGHT,\n               \"pixel_group\": key,\n               \"config\": {\"unfreeze_last\": 6,\n                          \"variant\": \"base\" if variant == \"base\" else \"small\",\n                          \"pool\": \"cls_mean_focal\", \"prior\": True}}\n              for k, f in enumerate(folds)]\n        if ms:\n            log(f\"legacy bundle {p.name}: {len(ms)} fold(s) join at weight \"\n                f\"{LEGACY_WEIGHT} each\")\n        return {key: ms} if ms else {}\n    except Exception as exc:\n        log(f\"legacy bundle unusable ({type(exc).__name__}: {exc}); blending skipped\")\n        return {}\n\n\ndef _run_member(path, m, dev, Cte, Mte, idx, starts, jitter):\n    \"\"\"Load, verify and predict one member on one device. Returns (pred, timings).\"\"\"\n    t0 = time.time()\n    with BUILD_LOCK:\n        if \"state\" in m:\n            state, fp = m[\"state\"], None\n        else:\n            ck = torch.load(Path(path) / m[\"file\"], map_location=\"cpu\",\n                            weights_only=False)\n            state, fp = ck[\"model\"], ck.get(\"fingerprint\")\n        model = build_model(int(m[\"config\"][\"unfreeze_last\"]),\n                            variant=m[\"config\"][\"variant\"],\n                            pool=m[\"config\"].get(\"pool\", \"cls_mean\"),\n                            prior=bool(m[\"config\"].get(\"prior\", False))).to(dev)\n        model.load_state_dict(state)\n        if fp is not None:\n            check_fingerprint(model, dev, IMG, fp, tag=f\"{m['id']}: \")\n        else:\n            log(f\"  {m['id']}: no stored fingerprint (legacy bundle) -- \"\n                f\"accepted at reduced weight\")\n    t_ready = time.time()\n    p = predict_member(model, Cte, Mte, idx, dev, IMG, starts=starts, jitter=jitter)\n    t_done = time.time()\n    del model, state\n    gc.collect()\n    if dev.type == \"cuda\":\n        with torch.cuda.device(dev):\n            torch.cuda.empty_cache()\n    passes = len(starts) * (2 if jitter else 1)\n    return p, (t_ready - t0, (t_done - t_ready) / max(passes, 1))\n\n\ndef _combine(per_member):\n    \"\"\"Weight-aware mean of per-member percentile ranks.\n\n    Rank rather than probability, because the metric reads order and two members\n    calibrated differently would otherwise have unequal say. Package members carry\n    weight 1.0 -- their combiner is the one the 0.891 was measured under -- and legacy\n    folds a reduced weight rather than a holdout-fitted one, which would import a few\n    hundred training studies' sampling noise.\n    \"\"\"\n    all_ids = sorted({s for m in per_member for s in m[\"ids\"]})\n    pos = {s: i for i, s in enumerate(all_ids)}\n    acc = np.zeros((len(all_ids), len(TARGETS)), np.float64)\n    tot = 0.0\n    for m in per_member:\n        w = float(m.get(\"weight\", 1.0))\n        r = pd.DataFrame(m[\"pred\"]).rank(pct=True).to_numpy()\n        acc[[pos[s] for s in m[\"ids\"]]] += w * r\n        tot += w\n    return all_ids, acc / max(tot, 1e-9)\n\n\ndef infer_from_package(path):\n    \"\"\"Predict the test split from an attached package of trained members.\n\n    Identical pixel contract to the reference implementation -- same caches, same\n    windows, same fingerprints, same rank transform -- with executive changes only:\n\n    1. A work queue over the devices: each GPU pops the next member when free (the\n       reference ran one device and surrendered TTA windows, then members: 0.847).\n    2. submission.csv is rewritten after every banked member, so a run killed at any\n       point still submits the best partial ensemble instead of the 0.5 benchmark.\n    3. A member that fails on one device is retried once on the other, then dropped;\n       a dropped member costs one vote, never the run.\n    4. An independently trained legacy bundle joins as five reduced-weight members.\n    5. When the time estimate says the whole remaining ensemble fits with room to\n       spare, each window gains one jittered TTA view (same family as training aug).\n    \"\"\"\n    man = json.loads((Path(path) / \"manifest.json\").read_text())\n    members = man[\"members\"]\n    log(f\"weights package: {len(members)} member(s) from {path}; \"\n        f\"{len(DEVS)} device(s)\")\n\n    test_df = pd.read_csv(ROOT / \"test.csv\")\n    test_series = pd.read_csv(ROOT / \"test_series.csv\")\n    plane_map = dict(zip(test_series[\"SeriesInstanceUID\"],\n                         test_series[\"Anatomical_Plane\"]))\n    hte = annotate(walk(\"test_series\"))\n    log(f\"test header pass: {len(hte)} series\")\n\n    groups = {}\n    for m in members:\n        groups.setdefault(m[\"pixel_group\"], []).append(m)\n    # Legacy votes come last: if time binds after all, the weakest votes are the ones\n    # surrendered, not the package's.\n    groups.update(legacy_group_members())\n\n    per_member = []\n    est = {\"fixed\": None, \"win\": None}\n\n    def bank(m, ids, pred, starts, jitter):\n        if float(np.std(pred)) < 1e-9:\n            log(f\"  {m['id']}: degenerate predictions; not banked\")\n            return\n        with STATE_LOCK:\n            per_member.append({\"id\": m[\"id\"], \"ids\": ids, \"pred\": pred,\n                               \"weight\": m.get(\"weight\", 1.0),\n                               \"holdout\": m.get(\"holdout\")})\n            all_ids, acc = _combine(per_member)\n            write_submission(acc, all_ids, test_df, \"submission.csv\")\n            log(f\"  banked {m['id']} fold {m.get('fold', '?')} \"\n                f\"({len(starts)} window(s){', jitter' if jitter else ''}); \"\n                f\"submission.csv = weighted rank mean of {len(per_member)} member(s)\")\n\n    for gi, (key, gm) in enumerate(groups.items(), 1):\n        cfg = json.loads(key)\n        adopt_config_globals(cfg)\n        log(f\"decode group {gi}/{len(groups)}: {cfg['img']}px x {cfg['slices']} slices, \"\n            f\"crop {cfg['crop_mm']} mm -> {len(gm)} member(s)\")\n        st_te, Cte, Mte = build_cache(pick_slots(hte, plane_map), plane_map,\n                                      lat_of(hte, \"test \"), f\"test g{gi}\")\n        idx = np.arange(len(st_te))\n\n        starts_full = window_starts(Cte.shape[2], GROUP)\n        pending = sorted(gm, key=lambda m: -(m.get(\"holdout\") or 0))\n        left_after = sum(len(g) for j, (_, g) in enumerate(groups.items(), 1) if j > gi)\n\n        def pop_next():\n            \"\"\"Next member, its window count, and whether jitter TTA is affordable.\n\n            Windows are surrendered before members; jitter is granted only when the\n            estimate says the whole remaining ensemble fits at double passes inside\n            60% of the room. The question asked is whether ONE more member fits.\n            \"\"\"\n            with STATE_LOCK:\n                if not pending:\n                    return None, None, False\n                left = TIME_BUDGET - (time.time() - T0)\n                remaining = len(pending) + left_after\n                slots_left = -(-remaining // len(DEVS))       # ceil: concurrent slots\n                starts, jit = starts_full, False\n                if est[\"fixed\"] is not None and est[\"win\"] is not None:\n                    afford = max(left * 0.9, 0.0)\n                    room = afford / max(slots_left, 1)\n                    if est[\"fixed\"] + est[\"win\"] > room:\n                        log(f\"  {left / 60:.0f} min left: surrendering \"\n                            f\"{len(pending)} member(s); not one more fits\")\n                        pending.clear()\n                        return None, None, False\n                    jit = (est[\"fixed\"] + 2 * len(starts_full) * est[\"win\"]\n                           <= room * 0.6)\n                    per_win = est[\"win\"] * (2 if jit else 1)\n                    n_win = (int((room - est[\"fixed\"]) / per_win)\n                             if per_win > 0 else len(starts_full))\n                    n_win = max(1, min(len(starts_full), n_win))\n                    if n_win < len(starts_full):\n                        mid = (len(starts_full) - n_win) // 2\n                        starts = starts_full[mid:mid + n_win]\n                return pending.pop(0), starts, jit\n\n        def worker(dev):\n            others = [d for d in DEVS if d is not dev]\n            while True:\n                m, starts, jit = pop_next()\n                if m is None:\n                    return\n                for attempt, d in enumerate([dev] + others[:1]):\n                    try:\n                        p, (fs, ws) = _run_member(path, m, d, Cte, Mte, idx,\n                                                  starts, jit)\n                        with STATE_LOCK:\n                            est[\"fixed\"], est[\"win\"] = fs, ws\n                        bank(m, st_te, p, starts, jit)\n                        break\n                    except Exception as exc:\n                        log(f\"  MEMBER {m['id']} failed on {d} \"\n                            f\"({type(exc).__name__}: {exc}); \"\n                            + (\"retrying on peer device\" if attempt == 0 and others\n                               else \"dropped -- costs one vote, not the run\"))\n                        if d.type == \"cuda\":\n                            with torch.cuda.device(d):\n                                torch.cuda.empty_cache()\n\n        threads = [threading.Thread(target=worker, args=(d,)) for d in DEVS]\n        for t in threads:\n            t.start()\n        for t in threads:\n            t.join()\n\n        del Cte, Mte\n        gc.collect()\n\n    if not per_member:\n        raise WeightsError(\"no member produced predictions; submission stays at 0.5\")\n\n    all_ids, acc = _combine(per_member)\n    sub = write_submission(acc, all_ids, test_df, \"submission.csv\")\n    log(f\"final submission.csv = weighted rank mean of {len(per_member)} member(s); \"\n        f\"{sub.shape}; nulls {int(sub[TARGETS].isna().sum().sum())}\")\n    return sub\n\n\ndef adopt_config_globals(cfg):\n    \"\"\"Point the pixel path at what one group of members was fitted on.\"\"\"\n    global IMG, CACHE_IMG, GROUP, CACHE_SLICES, N_GROUP, CROP_MM, SLICE_BAND, RULES\n    CACHE_IMG = IMG = int(cfg[\"img\"])\n    GROUP = int(cfg[\"group\"])\n    CACHE_SLICES = int(cfg[\"slices\"])\n    N_GROUP = max(CACHE_SLICES // GROUP, 1)\n    CROP_MM = float(cfg[\"crop_mm\"])\n    SLICE_BAND = tuple(float(x) for x in cfg[\"band\"])\n    # The four decisions that change what a slice is. A member fitted under one reading\n    # and decoded under another gets pixels its weights never saw, with every shape\n    # still agreeing, so an unrecognised name is refused rather than defaulted.\n    rules = cfg.get(\"rules\") or RULES_NATIVE\n    unknown = {k: v for k, v in rules.items()\n               if k not in RULES_NATIVE\n               or v not in (RULES_NATIVE[k], RULES_LEGACY[k])}\n    if unknown:\n        raise WeightsError(f\"the members record pixel rules this pipeline cannot \"\n                           f\"reproduce: {unknown}\")\n    RULES = {**RULES_NATIVE, **rules}\n    if [s[0] for s in SLOTS] != list(cfg[\"slots\"]):\n        raise WeightsError(\n            f\"the members were fitted on slots {cfg['slots']} and this pipeline defines \"\n            f\"{[s[0] for s in SLOTS]}; a weight would be read against the wrong slot\")\n","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 9. Run\nWrites a 0.5 benchmark `submission.csv` first (a run that dies must still have a valid\nfile), locates the weights package, and runs the whole inference. Any error leaves the\nlast banked submission on disk and re-raises loudly.","metadata":{}},{"cell_type":"code","source":"def write_submission(pred, studies, test_df, path):\n    \"\"\"Write one submission file from a prediction matrix.\n\n    Predictions are converted to per-column ranks first: the metric reads only order, so\n    ranks discard nothing, and they make files from different configurations directly\n    comparable and safe to average.\n    \"\"\"\n    sub = pd.DataFrame(pd.DataFrame(pred).rank(pct=True).values, columns=TARGETS)\n    sub.insert(0, \"StudyInstanceUID\", studies)\n    sub = test_df[[\"StudyInstanceUID\"]].merge(sub, on=\"StudyInstanceUID\", how=\"left\")\n    sub[TARGETS] = sub[TARGETS].fillna(0.5)\n    sub.to_csv(path, index=False)\n    return sub\n\n\ndef write_benchmark_submission():\n    \"\"\"Write the 0.5 benchmark file immediately.\n\n    A submission that never writes scores nothing at all, which is strictly worse than\n    scoring badly. The try/except around main() covers exceptions, but a kill for memory\n    is a SIGKILL and never reaches it. So a valid file exists from the first second and\n    is overwritten only once real predictions are ready.\n    \"\"\"\n    t = pd.read_csv(ROOT / \"test.csv\")\n    for c in TARGETS:\n        t[c] = 0.5\n    t.to_csv(\"submission.csv\", index=False)\n\n\n\ndef main():\n    write_benchmark_submission()\n    pkg = find_weights()\n    if pkg is None:\n        raise WeightsError(\n            \"no weights package attached. This notebook is inference-only: attach the \"\n            \"trained member package (the dataset holding manifest.json) and re-run.\")\n    infer_from_package(pkg)\n    log(\"done\")\n\n\ntry:\n    main()\nexcept Exception:\n    traceback.print_exc()\n    print(\"submission.csv holds the last banked state (or the 0.5 benchmark if \"\n          \"nothing was banked)\")\n    raise\n","metadata":{},"outputs":[],"execution_count":null}]}