{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","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"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Phase 1 — Malware Feature Extraction & Image Conversion","metadata":{}},{"cell_type":"code","source":"# ============================================================\n# ============================================================\n\n# --- Mode ---\nSUB_MODE = \"train\"\n\n# --- Input ---\n# Kaggle default:\nSUB_ARCHIVE_PATH = \"/kaggle/input/competitions/malware-classification/train.7z\"\nSUB_LABEL_PATH   = \"/kaggle/input/competitions/malware-classification/trainLabels.csv\"\n\n# Local Windows (uncomment + edit if running locally):\n# import platform\n# if platform.system() == 'Windows':\n#     SUB_ARCHIVE_PATH = r\"C:\\path\\to\\train.7z\"\n#     SUB_LABEL_PATH   = r\"C:\\path\\to\\trainLabels.csv\"\n\n# --- Output (isolated namespace) ---\nSUB_OUTPUT_DIR       = \"/kaggle/working/data\"\nSUB_IMAGE_DIR_NAME  = \"images\"\n\n# --- Batching and Parallelism ---\nSUB_BATCH_SIZE          = 64\nSUB_MAX_WORKERS         = 4\nSUB_MAX_BATCH_ATTEMPTS  = 3\nSUB_RETRY_DELAY_SECONDS = 3\n\n# --- Image Output ---\nSUB_IMAGE_SIZE = 224\n\n# --- Packaging and Resume ---\nSUB_ZIP_IMAGES                    = True\nSUB_REMOVE_IMAGE_FOLDER_AFTER_ZIP = True\nSUB_RESTORE_IMAGES_FROM_ZIP_ON_RESUME = True\nSUB_FAIL_ON_INCOMPLETE_OUTPUT    = True\n\n# --- Debug ---\n# Set an integer for a quick test; None processes the full training set.\nSUB_DEBUG_MAX_FILES = None\n\nprint(\"Sub-pipeline configuration loaded.\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:40.136236Z","iopub.execute_input":"2026-07-24T13:57:40.136629Z","iopub.status.idle":"2026-07-24T13:57:40.148912Z","shell.execute_reply.started":"2026-07-24T13:57:40.136599Z","shell.execute_reply":"2026-07-24T13:57:40.147768Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Imports and Runtime Setup\n\nImport libraries, validate input paths, restore images from ZIP on resume, and set up output directories.","metadata":{}},{"cell_type":"code","source":"import os\n\n# Limit numerical-library threads before importing NumPy/Pandas.\nos.environ.setdefault(\"OMP_NUM_THREADS\", \"1\")\nos.environ.setdefault(\"MKL_NUM_THREADS\", \"1\")\n\nimport gc\nimport json\nimport math\nimport platform\nimport re\nimport shutil\nimport subprocess\nimport time\nimport zipfile\nfrom collections import Counter\nfrom concurrent.futures import ProcessPoolExecutor, as_completed\nfrom datetime import datetime, timezone\nfrom pathlib import Path\n\nimport numpy as np\nimport pandas as pd\nfrom PIL import Image\n\n# --- Derived Output Paths ---\nSUB_OUTPUT_PATH = Path(SUB_OUTPUT_DIR)\nSUB_IMAGE_DIR  = SUB_OUTPUT_PATH / SUB_IMAGE_DIR_NAME\nSUB_IMAGE_ZIP_PATH  = SUB_OUTPUT_PATH / f\"{SUB_IMAGE_DIR_NAME}.zip\"\nSUB_BATCH_CSV_DIR   = SUB_OUTPUT_PATH / \"table_data\"\nSUB_FAILED_LOG_PATH = SUB_OUTPUT_PATH / \"failed_ids.csv\"\nSUB_MANIFEST_PATH   = SUB_OUTPUT_PATH / \"manifest.json\"\n\n# Prefer RAM disk; fall back to working storage.\nSHM = Path(\"/dev/shm\")\nSUB_TEMP_ROOT = (\n    SHM / \"batch_temp\"\n    if SHM.exists() and os.access(SHM, os.W_OK)\n    else SUB_OUTPUT_PATH / \"batch_temp\"\n)\n\nSUB_OUTPUT_PATH.mkdir(parents=True, exist_ok=True)\nSUB_IMAGE_DIR.mkdir(parents=True, exist_ok=True)\nSUB_BATCH_CSV_DIR.mkdir(parents=True, exist_ok=True)\nSUB_TEMP_ROOT.mkdir(parents=True, exist_ok=True)\n\n# --- Validate Inputs ---\nif not Path(SUB_ARCHIVE_PATH).exists():\n    raise FileNotFoundError(f\"Archive not found: {SUB_ARCHIVE_PATH}\")\nif not Path(SUB_LABEL_PATH).exists():\n    raise FileNotFoundError(f\"Labels not found: {SUB_LABEL_PATH}\")\nif shutil.which(\"7z\") is None:\n    raise RuntimeError(\"The `7z` executable is not available in this runtime.\")\n\n# --- Restore images from ZIP when resuming ---\nif (\n    SUB_RESTORE_IMAGES_FROM_ZIP_ON_RESUME\n    and SUB_IMAGE_ZIP_PATH.exists()\n    and not any(SUB_IMAGE_DIR.glob(\"*.png\"))\n):\n    print(f\"Restoring images from existing archive: {SUB_IMAGE_ZIP_PATH.name}\")\n    with zipfile.ZipFile(SUB_IMAGE_ZIP_PATH, \"r\") as archive:\n        bad_member = archive.testzip()\n        if bad_member is not None:\n            raise RuntimeError(f\"Corrupted image ZIP member: {bad_member}\")\n        archive.extractall(SUB_IMAGE_DIR)\n\n# --- Load Target IDs ---\nsub_labels_df = pd.read_csv(SUB_LABEL_PATH, dtype={\"Id\": str})\nif \"Id\" not in sub_labels_df.columns:\n    raise ValueError(\"The label file must contain an `Id` column.\")\nif sub_labels_df[\"Id\"].duplicated().any():\n    duplicated = sub_labels_df.loc[sub_labels_df[\"Id\"].duplicated(), \"Id\"].head(10).tolist()\n    raise ValueError(f\"Duplicate IDs in label file, examples: {duplicated}\")\n\nsub_target_ids = sub_labels_df[\"Id\"].astype(str).tolist()\nif SUB_DEBUG_MAX_FILES is not None:\n    sub_target_ids = sub_target_ids[: int(SUB_DEBUG_MAX_FILES)]\n\nSUB_MAX_WORKERS = max(1, min(int(SUB_MAX_WORKERS), os.cpu_count() or 2))\n\nprint(f\"Total files to process : {len(sub_target_ids):,}\")\nprint(f\"Batch size             : {SUB_BATCH_SIZE}\")\nprint(f\"Workers                : {SUB_MAX_WORKERS}\")\nprint(f\"Maximum batch attempts : {SUB_MAX_BATCH_ATTEMPTS}\")\nprint(f\"Output image size      : {SUB_IMAGE_SIZE} x {SUB_IMAGE_SIZE}\")\nprint(f\"Temporary directory    : {SUB_TEMP_ROOT}\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:40.150168Z","iopub.execute_input":"2026-07-24T13:57:40.150493Z","iopub.status.idle":"2026-07-24T13:57:41.847914Z","shell.execute_reply.started":"2026-07-24T13:57:40.150466Z","shell.execute_reply":"2026-07-24T13:57:41.846689Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Feature Schema and Utility Functions\n\nDefine byte/opcode vocabularies and helper functions for entropy calculation, atomic file writes, and image validation.","metadata":{}},{"cell_type":"code","source":"# --- Vocabulary ---\nSUB_SEGMENTS = [\".text\", \".data\", \".bss\", \".rdata\", \".edata\", \".idata\", \".rsrc\", \".tls\"]\n\nSUB_OPCODES = [\n    \"mov\", \"push\", \"pop\", \"jmp\", \"call\", \"ret\", \"cmp\", \"test\",\n    \"add\", \"sub\", \"inc\", \"dec\", \"xor\", \"jz\", \"jnz\", \"lea\",\n    \"and\", \"or\", \"shr\", \"shl\", \"nop\", \"imul\", \"idiv\", \"int\",\n]\n\nSUB_ASM_KEYWORDS = [\"db\", \"dw\", \"dd\", \"dq\", \"proc\", \"endp\", \"align\", \"extrn\"]\n\nSUB_HEX_LOOKUP = {f\"{i:02X}\": i for i in range(256)}\nSUB_OPCODE_PATTERN = re.compile(\n    r\"\\b(\" + \"|\".join(map(re.escape, SUB_OPCODES)) + r\")\\b\",\n    flags=re.IGNORECASE,\n)\nSUB_KEYWORD_PATTERN = re.compile(\n    r\"\\b(\" + \"|\".join(map(re.escape, SUB_ASM_KEYWORDS)) + r\")\\b\",\n    flags=re.IGNORECASE,\n)\n\ndef sub_get_image_width(raw_byte_count: int) -> int:\n    \"\"\"Choose the adaptive image width from the decoded byte count.\"\"\"\n    size_kb = raw_byte_count / 1024.0\n    if size_kb < 10:\n        return 64\n    if size_kb < 30:\n        return 128\n    if size_kb < 60:\n        return 256\n    if size_kb < 100:\n        return 384\n    if size_kb < 200:\n        return 512\n    if size_kb < 500:\n        return 768\n    if size_kb < 1000:\n        return 1024\n    return 2048\n\ndef sub_entropy_from_counts(counts: np.ndarray) -> float:\n    total = int(counts.sum())\n    if total == 0:\n        return 0.0\n    probabilities = counts[counts > 0].astype(np.float64) / total\n    return float(-(probabilities * np.log2(probabilities)).sum())\n\ndef sub_safe_remove(path: Path) -> None:\n    try:\n        path.unlink(missing_ok=True)\n    except Exception:\n        pass\n\ndef sub_image_is_valid(path: Path) -> bool:\n    \"\"\"A completed image is an existing non-empty PNG produced by an atomic rename.\"\"\"\n    try:\n        return path.is_file() and path.stat().st_size > 0\n    except OSError:\n        return False\n\ndef sub_atomic_write_csv(dataframe: pd.DataFrame, destination: Path) -> None:\n    destination.parent.mkdir(parents=True, exist_ok=True)\n    temporary = destination.with_suffix(destination.suffix + \".tmp\")\n    dataframe.to_csv(temporary, index=False)\n    os.replace(temporary, destination)\n\ndef sub_atomic_write_text(text: str, destination: Path) -> None:\n    destination.parent.mkdir(parents=True, exist_ok=True)\n    temporary = destination.with_suffix(destination.suffix + \".tmp\")\n    temporary.write_text(text, encoding=\"utf-8\")\n    os.replace(temporary, destination)\n\ndef sub_numeric_batch_id(path: Path) -> int:\n    match = re.search(r\"batch_(\\d+)\\.csv$\", path.name)\n    return int(match.group(1)) if match else -1\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:41.84995Z","iopub.execute_input":"2026-07-24T13:57:41.850392Z","iopub.status.idle":"2026-07-24T13:57:41.869535Z","shell.execute_reply.started":"2026-07-24T13:57:41.850361Z","shell.execute_reply":"2026-07-24T13:57:41.86785Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Process One Malware Sample\n\nParse `.bytes` files into grayscale images and extract static features from `.asm` files using opcode/segment counting.","metadata":{}},{"cell_type":"code","source":"# --- Parse .bytes File ---\ndef sub_parse_bytes_file(byte_path: Path, image_path: Path) -> dict:\n    counts = np.zeros(256, dtype=np.int64)\n    values = bytearray()\n    unknown_count = 0\n\n    with byte_path.open(\"r\", encoding=\"latin1\", errors=\"ignore\") as file:\n        for line in file:\n            for token in line.split():\n                token = token.upper()\n                if token == \"??\":\n                    unknown_count += 1\n                    values.append(0)\n                    continue\n\n                value = SUB_HEX_LOOKUP.get(token)\n                if value is not None:\n                    counts[value] += 1\n                    values.append(value)\n\n    raw_count = len(values)\n    known_count = int(counts.sum())\n    if raw_count == 0:\n        raise ValueError(\"The .bytes file does not contain valid byte tokens.\")\n\n    raw_array = np.frombuffer(values, dtype=np.uint8)\n    byte_mean = float(raw_array.mean())\n    byte_std = float(raw_array.std())\n\n    width = sub_get_image_width(raw_count)\n    height = int(math.ceil(raw_count / width))\n    padded_length = height * width\n    image_array = raw_array\n    if padded_length > raw_count:\n        image_array = np.pad(\n            raw_array,\n            (0, padded_length - raw_count),\n            mode=\"constant\",\n            constant_values=0,\n        )\n\n    raw_image = image_array.reshape(height, width)\n    image = Image.fromarray(raw_image)\n    image = image.resize((SUB_IMAGE_SIZE, SUB_IMAGE_SIZE), Image.Resampling.BILINEAR)\n\n    temporary_image_path = image_path.with_suffix(\".tmp.png\")\n    sub_safe_remove(temporary_image_path)\n    try:\n        image.save(temporary_image_path, format=\"PNG\", optimize=False)\n        os.replace(temporary_image_path, image_path)\n    finally:\n        sub_safe_remove(temporary_image_path)\n\n    features = {\n        \"byte_total\": raw_count,\n        \"byte_known\": known_count,\n        \"byte_unknown\": unknown_count,\n        \"byte_unknown_ratio\": unknown_count / max(raw_count, 1),\n        \"byte_entropy\": sub_entropy_from_counts(counts),\n        \"byte_mean\": byte_mean,\n        \"byte_std\": byte_std,\n        \"byte_zero_ratio\": float(counts[0] / max(known_count, 1)),\n        \"byte_ff_ratio\": float(counts[255] / max(known_count, 1)),\n    }\n\n    frequencies = counts.astype(np.float64) / max(known_count, 1)\n    for value in range(256):\n        features[f\"byte_count_{value:02X}\"] = int(counts[value])\n        features[f\"byte_freq_{value:02X}\"] = float(frequencies[value])\n\n    return features\n\n\n# --- Parse .asm File ---\ndef sub_parse_asm_file(asm_path: Path) -> dict:\n    segment_counts = Counter()\n    opcode_counts = Counter()\n    keyword_counts = Counter()\n    line_count = 0\n    nonempty_line_count = 0\n    asm_chars = 0\n\n    with asm_path.open(\"r\", encoding=\"latin1\", errors=\"ignore\") as file:\n        for line in file:\n            lower = line.lower()\n            line_count += 1\n            asm_chars += len(line)\n            if lower.strip():\n                nonempty_line_count += 1\n\n            for segment in SUB_SEGMENTS:\n                if segment + \":\" in lower:\n                    segment_counts[segment] += 1\n\n            opcode_counts.update(SUB_OPCODE_PATTERN.findall(lower))\n            keyword_counts.update(SUB_KEYWORD_PATTERN.findall(lower))\n\n    total_opcodes = int(sum(opcode_counts.values()))\n    features = {\n        \"asm_file_size\": int(asm_path.stat().st_size),\n        \"asm_char_count\": asm_chars,\n        \"asm_line_count\": line_count,\n        \"asm_nonempty_line_count\": nonempty_line_count,\n        \"asm_total_opcodes\": total_opcodes,\n    }\n\n    for segment in SUB_SEGMENTS:\n        name = segment.replace(\".\", \"\")\n        features[f\"seg_count_{name}\"] = int(segment_counts[segment])\n\n    for opcode in SUB_OPCODES:\n        count = int(opcode_counts[opcode])\n        features[f\"op_count_{opcode}\"] = count\n        features[f\"op_freq_{opcode}\"] = count / max(total_opcodes, 1)\n\n    for keyword in SUB_ASM_KEYWORDS:\n        features[f\"asm_keyword_{keyword}\"] = int(keyword_counts[keyword])\n\n    return features\n\n\n# --- Process One Pair of Files ---\ndef sub_process_single_file(fid: str, batch_dir_str: str, image_dir_str: str) -> dict:\n    batch_dir = Path(batch_dir_str)\n    image_dir = Path(image_dir_str)\n    byte_path = batch_dir / f\"{fid}.bytes\"\n    asm_path = batch_dir / f\"{fid}.asm\"\n    image_path = image_dir / f\"{fid}.png\"\n\n    result = {\"Id\": fid}\n    try:\n        if not byte_path.exists():\n            raise FileNotFoundError(f\"Missing {fid}.bytes\")\n        if not asm_path.exists():\n            raise FileNotFoundError(f\"Missing {fid}.asm\")\n\n        result.update(sub_parse_bytes_file(byte_path, image_path))\n        result.update(sub_parse_asm_file(asm_path))\n        result[\"_error\"] = None\n    except Exception as exc:\n        sub_safe_remove(image_path)\n        result[\"_error\"] = f\"{type(exc).__name__}: {exc}\"\n    finally:\n        sub_safe_remove(byte_path)\n        sub_safe_remove(asm_path)\n\n    return result\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:41.871155Z","iopub.execute_input":"2026-07-24T13:57:41.871593Z","iopub.status.idle":"2026-07-24T13:57:41.913271Z","shell.execute_reply.started":"2026-07-24T13:57:41.871549Z","shell.execute_reply":"2026-07-24T13:57:41.911847Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Extract and Process a Batch\n\nExtract files from the archive using `7z`, then process them in parallel with `ProcessPoolExecutor`. Temporary raw files are removed after each batch.","metadata":{}},{"cell_type":"code","source":"def sub_reset_batch_directory(batch_dir: Path) -> None:\n    if batch_dir.exists():\n        shutil.rmtree(batch_dir, ignore_errors=True)\n    batch_dir.mkdir(parents=True, exist_ok=True)\n\n\ndef sub_extract_batch(id_list, batch_number: int) -> Path:\n    batch_dir = SUB_TEMP_ROOT / f\"sub_batch_{batch_number:04d}\"\n    sub_reset_batch_directory(batch_dir)\n\n    list_path = SUB_TEMP_ROOT / f\"sub_batch_{batch_number:04d}.txt\"\n    with list_path.open(\"w\", encoding=\"utf-8\") as list_file:\n        for fid in id_list:\n            list_file.write(f\"{SUB_MODE}/{fid}.bytes\\n\")\n            list_file.write(f\"{SUB_MODE}/{fid}.asm\\n\")\n\n    command = [\n        \"7z\",\n        \"e\",\n        SUB_ARCHIVE_PATH,\n        f\"@{list_path}\",\n        f\"-o{batch_dir}\",\n        \"-y\",\n        \"-bd\",\n        \"-bb0\",\n    ]\n    completed = subprocess.run(command, capture_output=True, text=True)\n    sub_safe_remove(list_path)\n\n    if completed.returncode != 0:\n        shutil.rmtree(batch_dir, ignore_errors=True)\n        details = (completed.stderr + \"\\n\" + completed.stdout)[-2000:]\n        raise RuntimeError(f\"7z extraction failed (code {completed.returncode}):\\n{details}\")\n\n    return batch_dir\n\n\ndef sub_process_batch(id_list, batch_number: int):\n    batch_dir = sub_extract_batch(id_list, batch_number)\n    successes, failures = [], []\n\n    try:\n        with ProcessPoolExecutor(max_workers=SUB_MAX_WORKERS) as executor:\n            futures = {\n                executor.submit(\n                    sub_process_single_file,\n                    fid,\n                    str(batch_dir),\n                    str(SUB_IMAGE_DIR),\n                ): fid\n                for fid in id_list\n            }\n\n            for future in as_completed(futures):\n                fid = futures[future]\n                try:\n                    result = future.result()\n                except Exception as exc:\n                    result = {\"Id\": fid, \"_error\": f\"WorkerError: {exc}\"}\n\n                if result.get(\"_error\"):\n                    failures.append({\"Id\": fid, \"error\": result[\"_error\"]})\n                else:\n                    result.pop(\"_error\", None)\n                    successes.append(result)\n    finally:\n        shutil.rmtree(batch_dir, ignore_errors=True)\n\n    return successes, failures\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:41.914721Z","iopub.execute_input":"2026-07-24T13:57:41.915277Z","iopub.status.idle":"2026-07-24T13:57:41.948595Z","shell.execute_reply.started":"2026-07-24T13:57:41.915232Z","shell.execute_reply":"2026-07-24T13:57:41.947022Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Checkpointing and Full Run\n\nResume from existing checkpoints. Each batch is retried automatically on failure, and successful results are written atomically to disk.","metadata":{}},{"cell_type":"code","source":"# --- Resume from Existing Checkpoints ---\nsub_existing_batch_files = sorted(\n    SUB_BATCH_CSV_DIR.glob(f\"batch_*.csv\"),\n    key=sub_numeric_batch_id,\n)\n\nsub_processed_ids = set()\nsub_existing_batch_numbers = []\nsub_stale_checkpoint_ids = set()\nsub_corrupted_checkpoints = []\n\nfor path in sub_existing_batch_files:\n    sub_existing_batch_numbers.append(sub_numeric_batch_id(path))\n    try:\n        checkpoint = pd.read_csv(path, usecols=[\"Id\"], dtype={\"Id\": str})\n        for fid in checkpoint[\"Id\"].astype(str):\n            if fid in sub_target_ids and sub_image_is_valid(SUB_IMAGE_DIR / f\"{fid}.png\"):\n                sub_processed_ids.add(fid)\n            elif fid in sub_target_ids:\n                sub_stale_checkpoint_ids.add(fid)\n    except Exception as exc:\n        sub_corrupted_checkpoints.append({\"file\": path.name, \"error\": str(exc)})\n        print(f\"Ignoring unreadable checkpoint {path.name}: {exc}\")\n\nsub_remaining_ids = [fid for fid in sub_target_ids if fid not in sub_processed_ids]\nsub_next_batch_number = max(sub_existing_batch_numbers, default=0) + 1\n\nprint(f\"Valid checkpointed samples : {len(sub_processed_ids):,}\")\nprint(f\"Stale checkpoint rows      : {len(sub_stale_checkpoint_ids):,}\")\nprint(f\"Remaining samples          : {len(sub_remaining_ids):,}\")\n\n# --- Main Processing Loop ---\nsub_all_failures = []\nsub_start_time = time.time()\n\nfor offset in range(0, len(sub_remaining_ids), SUB_BATCH_SIZE):\n    batch_ids = sub_remaining_ids[offset : offset + SUB_BATCH_SIZE]\n    batch_number = sub_next_batch_number + offset // SUB_BATCH_SIZE\n    print(f\"\\n[Batch {batch_number:04d}] {len(batch_ids)} samples\", flush=True)\n\n    successful_rows = {}\n    pending_ids = list(batch_ids)\n    last_errors = {fid: \"Not processed\" for fid in pending_ids}\n\n    for attempt in range(1, SUB_MAX_BATCH_ATTEMPTS + 1):\n        if not pending_ids:\n            break\n\n        if attempt > 1:\n            print(f\"  Retry {attempt}/{SUB_MAX_BATCH_ATTEMPTS}: {len(pending_ids)} samples\")\n            time.sleep(SUB_RETRY_DELAY_SECONDS)\n\n        try:\n            successes, failures = sub_process_batch(pending_ids, batch_number)\n        except Exception as exc:\n            successes = []\n            failures = [\n                {\"Id\": fid, \"error\": f\"BatchError: {type(exc).__name__}: {exc}\"}\n                for fid in pending_ids\n            ]\n\n        for row in successes:\n            successful_rows[row[\"Id\"]] = row\n\n        failed_map = {item[\"Id\"]: item[\"error\"] for item in failures}\n        for fid, error in failed_map.items():\n            last_errors[fid] = error\n\n        pending_ids = [fid for fid in pending_ids if fid not in successful_rows]\n\n    if successful_rows:\n        batch_df = (\n            pd.DataFrame(successful_rows.values())\n            .sort_values(\"Id\")\n            .reset_index(drop=True)\n        )\n        checkpoint_path = (\n            SUB_BATCH_CSV_DIR / f\"batch_{batch_number:04d}.csv\"\n        )\n        sub_atomic_write_csv(batch_df, checkpoint_path)\n        print(f\"  Successful : {len(batch_df):3d} | saved {checkpoint_path.name}\")\n    else:\n        print(\"  Successful :   0\")\n\n    if pending_ids:\n        failures = [{\"Id\": fid, \"error\": last_errors[fid]} for fid in pending_ids]\n        sub_all_failures.extend(failures)\n        print(f\"  Failed     : {len(pending_ids):3d} after all retries\")\n    else:\n        print(\"  Failed     :   0\")\n\n    elapsed_minutes = (time.time() - sub_start_time) / 60\n    completed_now = min(offset + len(batch_ids), len(sub_remaining_ids))\n    print(\n        f\"  Progress   : {completed_now:,}/{len(sub_remaining_ids):,} new samples | \"\n        f\"{elapsed_minutes:.1f} minutes\"\n    )\n    gc.collect()\n\n# --- Save Current Failure Log ---\nif sub_all_failures:\n    failure_df = pd.DataFrame(sub_all_failures).drop_duplicates(\"Id\", keep=\"last\")\n    sub_atomic_write_csv(failure_df, SUB_FAILED_LOG_PATH)\n    print(f\"\\nFailure log saved: {SUB_FAILED_LOG_PATH}\")\nelif SUB_FAILED_LOG_PATH.exists():\n    SUB_FAILED_LOG_PATH.unlink()\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-07-24T13:57:41.950315Z","iopub.execute_input":"2026-07-24T13:57:41.950737Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Validate and Package Output\n\nMerge all batch checkpoints, remove zero-variance features, validate completeness, save tabular outputs, compress images to ZIP, and write the Phase 1 manifest.","metadata":{}},{"cell_type":"code","source":"# --- Combine Batch Results ---\nsub_batch_files = sorted(\n    SUB_BATCH_CSV_DIR.glob(\"batch_*.csv\"),\n    key=sub_numeric_batch_id,\n)\nif not sub_batch_files:\n    raise RuntimeError(\"No valid batch checkpoint CSV files were produced.\")\n\nsub_frames = []\nfor path in sub_batch_files:\n    frame = pd.read_csv(path, dtype={\"Id\": str})\n    if \"Id\" not in frame.columns:\n        raise ValueError(f\"Checkpoint has no `Id` column: {path.name}\")\n    frame[\"_checkpoint_order\"] = sub_numeric_batch_id(path)\n    sub_frames.append(frame)\n\nsub_raw_features = pd.concat(sub_frames, ignore_index=True, sort=False)\nsub_raw_features[\"Id\"] = sub_raw_features[\"Id\"].astype(str)\n\nsub_expected_ids = set(sub_target_ids)\nsub_unexpected_feature_ids = sorted(set(sub_raw_features[\"Id\"]) - sub_expected_ids)\nsub_raw_features = sub_raw_features.loc[sub_raw_features[\"Id\"].isin(sub_expected_ids)].copy()\n\nsub_duplicate_rows_resolved = int(sub_raw_features.duplicated(\"Id\", keep=False).sum())\nsub_all_features = (\n    sub_raw_features.sort_values([\"Id\", \"_checkpoint_order\"])\n    .drop_duplicates(\"Id\", keep=\"last\")\n    .drop(columns=[\"_checkpoint_order\"])\n)\n\n# Remove legacy presence flags and any other zero-variance feature columns.\nsub_legacy_constant_columns = [\n    column for column in [\"has_bytes\", \"has_asm\"] if column in sub_all_features.columns\n]\nif sub_legacy_constant_columns:\n    sub_all_features = sub_all_features.drop(columns=sub_legacy_constant_columns)\n\nsub_candidate_feature_columns = [column for column in sub_all_features.columns if column != \"Id\"]\nsub_constant_feature_columns = [\n    column\n    for column in sub_candidate_feature_columns\n    if sub_all_features[column].nunique(dropna=False) <= 1\n]\nif sub_constant_feature_columns:\n    sub_all_features = sub_all_features.drop(columns=sub_constant_feature_columns)\n\nsub_order = {fid: index for index, fid in enumerate(sub_target_ids)}\nsub_all_features[\"_target_order\"] = sub_all_features[\"Id\"].map(sub_order)\nsub_all_features = (\n    sub_all_features.sort_values(\"_target_order\")\n    .drop(columns=[\"_target_order\"])\n    .reset_index(drop=True)\n)\n\n# --- Strict Completeness and Integrity Validation ---\nsub_feature_ids = set(sub_all_features[\"Id\"])\nsub_missing_feature_ids = sorted(sub_expected_ids - sub_feature_ids)\nsub_missing_image_ids = [\n    fid for fid in sub_target_ids if not sub_image_is_valid(SUB_IMAGE_DIR / f\"{fid}.png\")\n]\nsub_final_duplicate_count = int(sub_all_features[\"Id\"].duplicated().sum())\nsub_feature_columns = [column for column in sub_all_features.columns if column != \"Id\"]\n\nsub_non_numeric_columns = [\n    column\n    for column in sub_feature_columns\n    if not pd.api.types.is_numeric_dtype(sub_all_features[column])\n]\nsub_invalid_value_count = 0\nif not sub_non_numeric_columns and sub_feature_columns:\n    sub_numeric_values = sub_all_features[sub_feature_columns].to_numpy(dtype=np.float64, copy=False)\n    sub_invalid_value_count = int((~np.isfinite(sub_numeric_values)).sum())\n\nprint(f\"Feature row count         : {len(sub_all_features):,}\")\nprint(f\"Feature column count      : {len(sub_feature_columns):,}\")\nprint(f\"Final duplicate IDs       : {sub_final_duplicate_count}\")\nprint(f\"Resolved duplicate rows   : {sub_duplicate_rows_resolved}\")\nprint(f\"Missing feature rows      : {len(sub_missing_feature_ids):,}\")\nprint(f\"Missing/empty images      : {len(sub_missing_image_ids):,}\")\nprint(f\"Non-numeric columns       : {len(sub_non_numeric_columns):,}\")\nprint(f\"NaN/Inf values            : {sub_invalid_value_count:,}\")\nprint(f\"Removed constant features : {len(sub_legacy_constant_columns) + len(sub_constant_feature_columns):,}\")\n\nsub_validation_errors = []\nif sub_missing_feature_ids:\n    sub_validation_errors.append(f\"missing feature IDs, examples: {sub_missing_feature_ids[:10]}\")\nif sub_missing_image_ids:\n    sub_validation_errors.append(f\"missing image IDs, examples: {sub_missing_image_ids[:10]}\")\nif sub_final_duplicate_count:\n    sub_validation_errors.append(f\"{sub_final_duplicate_count} duplicate IDs remain\")\nif sub_non_numeric_columns:\n    sub_validation_errors.append(f\"non-numeric columns: {sub_non_numeric_columns[:10]}\")\nif sub_invalid_value_count:\n    sub_validation_errors.append(f\"{sub_invalid_value_count} NaN/Inf values\")\n\nif sub_validation_errors and SUB_FAIL_ON_INCOMPLETE_OUTPUT:\n    raise RuntimeError(\"Sub-pipeline Phase 1 validation failed: \" + \"; \".join(sub_validation_errors))\n\n# --- Save Combined Tabular Outputs Atomically ---\nsub_combined_csv_path = SUB_OUTPUT_PATH / \"all_features.csv\"\nsub_atomic_write_csv(sub_all_features, sub_combined_csv_path)\n\nsub_combined_parquet_path = SUB_OUTPUT_PATH / \"all_features.parquet\"\nsub_parquet_saved = False\ntry:\n    temporary_parquet = sub_combined_parquet_path.with_suffix(\".parquet.tmp\")\n    sub_all_features.to_parquet(temporary_parquet, index=False)\n    os.replace(temporary_parquet, sub_combined_parquet_path)\n    sub_parquet_saved = True\nexcept Exception as exc:\n    sub_safe_remove(sub_combined_parquet_path.with_suffix(\".parquet.tmp\"))\n    print(f\"Parquet was not saved; Phase 2 will use CSV: {exc}\")\n\n# --- Build a Small Preview Grid ---\nsub_sample_ids = sub_all_features[\"Id\"].head(25).tolist()\ncanvas = Image.new(\"L\", (SUB_IMAGE_SIZE * 5, SUB_IMAGE_SIZE * 5), color=0)\nfor index, fid in enumerate(sub_sample_ids):\n    path = SUB_IMAGE_DIR / f\"{fid}.png\"\n    with Image.open(path) as image:\n        canvas.paste(\n            image.convert(\"L\"),\n            ((index % 5) * SUB_IMAGE_SIZE, (index // 5) * SUB_IMAGE_SIZE),\n        )\ncanvas.save(SUB_OUTPUT_PATH / \"sample_images.png\")\n\n# --- Create and Validate the Image Archive Atomically ---\nif SUB_ZIP_IMAGES:\n    print(\"Compressing the validated image folder...\")\n    temporary_base = SUB_OUTPUT_PATH / f\".{SUB_IMAGE_DIR_NAME}_temporary\"\n    temporary_zip = temporary_base.with_suffix(\".zip\")\n    sub_safe_remove(temporary_zip)\n\n    created_archive = Path(\n        shutil.make_archive(str(temporary_base), \"zip\", root_dir=SUB_IMAGE_DIR)\n    )\n    with zipfile.ZipFile(created_archive, \"r\") as archive:\n        bad_member = archive.testzip()\n        png_count = sum(1 for name in archive.namelist() if name.lower().endswith(\".png\"))\n\n    if bad_member is not None:\n        sub_safe_remove(created_archive)\n        raise RuntimeError(f\"Image ZIP validation failed at member: {bad_member}\")\n    if png_count != len(sub_target_ids):\n        sub_safe_remove(created_archive)\n        raise RuntimeError(\n            f\"Image ZIP contains {png_count:,} PNG files; expected {len(sub_target_ids):,}.\"\n        )\n\n    os.replace(created_archive, SUB_IMAGE_ZIP_PATH)\n    print(f\"Image archive created: {SUB_IMAGE_ZIP_PATH}\")\n\n# --- Write Complete Manifest Atomically ---\nsub_removed_constant_features = sub_legacy_constant_columns + sub_constant_feature_columns\nsub_manifest = {\n    \"phase\": 1,\n    \"version\": \"3.0\",\n    \"status\": \"complete\" if not sub_validation_errors else \"incomplete\",\n    \"completed_at_utc\": datetime.now(timezone.utc).isoformat(),\n    \"mode\": SUB_MODE,\n    \"sample_count\": int(len(sub_all_features)),\n    \"target_count\": int(len(sub_target_ids)),\n    \"feature_count\": int(len(sub_feature_columns)),\n    \"feature_columns\": sub_feature_columns,\n    \"constant_features_removed\": sub_removed_constant_features,\n    \"image_conversion\": {\n        \"method\": \"adaptive_width_grayscale\",\n        \"unknown_byte_pixel_value\": 0,\n        \"resize\": \"bilinear\",\n        \"image_size\": int(SUB_IMAGE_SIZE),\n    },\n    \"image_size\": int(SUB_IMAGE_SIZE),\n    \"image_folder_name\": SUB_IMAGE_DIR_NAME,\n    \"image_zip_name\": SUB_IMAGE_ZIP_PATH.name if SUB_ZIP_IMAGES else None,\n    \"combined_csv\": sub_combined_csv_path.name,\n    \"combined_parquet\": sub_combined_parquet_path.name if sub_parquet_saved else None,\n    \"missing_feature_count\": int(len(sub_missing_feature_ids)),\n    \"missing_image_count\": int(len(sub_missing_image_ids)),\n    \"unexpected_feature_id_count\": int(len(sub_unexpected_feature_ids)),\n    \"duplicate_rows_resolved\": int(sub_duplicate_rows_resolved),\n    \"runtime\": {\n        \"python\": platform.python_version(),\n        \"numpy\": np.__version__,\n        \"pandas\": pd.__version__,\n        \"pillow\": Image.__version__ if hasattr(Image, \"__version__\") else None,\n    },\n}\nsub_atomic_write_text(\n    json.dumps(sub_manifest, ensure_ascii=False, indent=2),\n    SUB_MANIFEST_PATH,\n)\n\n# Remove the image folder only after the ZIP and manifest are successfully validated.\nif SUB_ZIP_IMAGES and SUB_REMOVE_IMAGE_FOLDER_AFTER_ZIP:\n    shutil.rmtree(SUB_IMAGE_DIR, ignore_errors=True)\n\nprint(\"\\n\" + \"=\" * 72)\nprint(\"SUB-PIPELINE PHASE 1 COMPLETE\")\nprint(f\"Output directory : {SUB_OUTPUT_PATH}\")\nprint(f\"Manifest         : {SUB_MANIFEST_PATH.name}\")\nprint(f\"Tabular CSV      : {sub_combined_csv_path.name}\")\nif sub_parquet_saved:\n    print(f\"Tabular Parquet  : {sub_combined_parquet_path.name}\")\nif SUB_ZIP_IMAGES:\n    print(f\"Image ZIP        : {SUB_IMAGE_ZIP_PATH.name}\")\nprint(f\"Samples          : {len(sub_all_features):,}\")\nprint(f\"Features         : {len(sub_feature_columns):,}\")\nprint(\"=\" * 72)\n","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}