{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"codemirror_mode":{"name":"ipython","version":3},"file_extension":".py","mimetype":"text/x-python","name":"python","nbconvert_exporter":"python","pygments_lexer":"ipython3","version":"3.12.13"},"kaggle":{"accelerator":"none","dataSources":[],"dockerImageVersionId":28755,"isGpuEnabled":false,"isInternetEnabled":false,"language":"python","sourceType":"notebook"},"papermill":{"default_parameters":{},"duration":938.690005,"end_time":"2026-07-21T05:09:31.616356+00:00","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2026-07-21T04:53:52.926351+00:00","version":"2.7.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"id":"872bfb25","cell_type":"markdown","source":"# Microsoft Malware Classification 2015 — Phase 1\n## Data Conversion and Static Feature Extraction\n\n**Authors:** Huỳnh Bảo Khang, Nguyễn Phước Khang, Phạm Minh Đạt\n\n**Objective:** Convert the raw `.bytes` and `.asm` files into a validated,\ncheckpointed dataset containing fixed-size grayscale images and tabular\nstatic features for Phase 2.\n\nThe pipeline is designed for long-running execution on Kaggle:\n- batch extraction and parallel processing;\n- atomic checkpoints and automatic retry;\n- optional restoration from a previous Phase 1 output;\n- safe resume only when both a feature row and its image are valid;\n- automatic restoration from an existing image ZIP;\n- zero-variance feature removal;\n- strict integrity validation before packaging.\n","metadata":{"papermill":{"duration":0.005053,"end_time":"2026-07-21T04:53:55.963135+00:00","exception":false,"start_time":"2026-07-21T04:53:55.958082+00:00","status":"completed"},"tags":[]}},{"id":"4ebeb810","cell_type":"markdown","source":"### Configuration\nSet the dataset paths, batch size, worker count, and output image size used throughout the pipeline.","metadata":{"papermill":{"duration":0.003512,"end_time":"2026-07-21T04:53:55.970509+00:00","exception":false,"start_time":"2026-07-21T04:53:55.966997+00:00","status":"completed"},"tags":[]}},{"id":"167e18f3","cell_type":"code","source":"MODE = \"train\"\n\n# --- Input ---\nARCHIVE_PATH = \"/kaggle/input/competitions/malware-classification/train.7z\"\nLABEL_PATH = \"/kaggle/input/competitions/malware-classification/trainLabels.csv\"\n\n# --- Output ---\nOUTPUT_DIR = \"/kaggle/working/malware_phase1_output\"\nIMAGE_DIR_NAME = \"train_malware_images\"\n\n# --- Batching and Parallelism ---\nBATCH_SIZE = 64\nMAX_WORKERS = 4\nMAX_BATCH_ATTEMPTS = 3\nRETRY_DELAY_SECONDS = 3\n\n# --- Image Output ---\nIMAGE_SIZE = 224\n\n# --- Packaging and Resume ---\nZIP_IMAGES = True\nREMOVE_IMAGE_FOLDER_AFTER_ZIP = True\nRESTORE_IMAGES_FROM_ZIP_ON_RESUME = True\nFAIL_ON_INCOMPLETE_OUTPUT = True\n\n# Search Kaggle inputs for one previous `malware_phase1_output`.\n# If no checkpoint is found, the pipeline starts from the beginning.\nRESTORE_CHECKPOINT_FROM_INPUT = True\n\n# Optional substring for selecting one checkpoint input when several\n# Phase 1 outputs are attached. Leave as None for automatic selection.\nCHECKPOINT_INPUT_HINT = None\n\n# --- Debug ---\n# Set an integer for a quick test; None processes the full training set.\nDEBUG_MAX_FILES = None\n","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","execution":{"iopub.execute_input":"2026-07-21T04:53:55.979513Z","iopub.status.busy":"2026-07-21T04:53:55.979161Z","iopub.status.idle":"2026-07-21T04:53:55.989153Z","shell.execute_reply":"2026-07-21T04:53:55.987987Z"},"papermill":{"duration":0.017324,"end_time":"2026-07-21T04:53:55.991479+00:00","exception":false,"start_time":"2026-07-21T04:53:55.974155+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"ad55198f-d9b6-4c5c-8906-dff4393ec4ca","cell_type":"markdown","source":"### Interruption-Safe Checkpoint Strategy\n\nProcessing the complete dataset may exceed the maximum duration of a single\nKaggle session. The pipeline therefore saves tabular features and generated\nimages after every batch.\n\nWhen one previous Phase 1 output is attached as a Kaggle input, its valid\ncheckpoints are restored automatically and only incomplete samples are\nprocessed. If no previous output is attached, the pipeline starts from the\nbeginning. Final artifacts are accepted only after all expected samples pass\nstrict integrity validation.\n","metadata":{}},{"id":"80e2daa5","cell_type":"code","source":"# ============================================================\n# OPTIONAL CHECKPOINT RESTORATION\n# ============================================================\n\nfrom pathlib import Path\nimport shutil\n\nRESUME_TARGET = Path(OUTPUT_DIR)\n\nif not RESTORE_CHECKPOINT_FROM_INPUT:\n    print(\"Checkpoint restoration is disabled. Phase 1 will start from scratch.\")\nelse:\n    resume_candidates = [\n        path\n        for path in Path(\"/kaggle/input\").rglob(\"malware_phase1_output\")\n        if (path / \"tabular_batches\").exists()\n    ]\n\n    if CHECKPOINT_INPUT_HINT:\n        resume_candidates = [\n            path\n            for path in resume_candidates\n            if CHECKPOINT_INPUT_HINT.lower() in str(path).lower()\n        ]\n\n    print(f\"Checkpoint sources found: {len(resume_candidates)}\")\n    for path in resume_candidates:\n        print(\" -\", path)\n\n    if len(resume_candidates) == 0:\n        print(\n            \"No previous Phase 1 checkpoint was found. \"\n            \"The pipeline will start from the beginning.\"\n        )\n\n    elif len(resume_candidates) == 1:\n        resume_source = resume_candidates[0]\n\n        print(\"\\nRestoring previous Phase 1 output:\")\n        print(f\"Source : {resume_source}\")\n        print(f\"Target : {RESUME_TARGET}\")\n\n        shutil.copytree(\n            resume_source,\n            RESUME_TARGET,\n            dirs_exist_ok=True,\n        )\n\n        checkpoint_dir = RESUME_TARGET / \"tabular_batches\"\n        image_dir = RESUME_TARGET / IMAGE_DIR_NAME\n        image_zip_path = RESUME_TARGET / f\"{IMAGE_DIR_NAME}.zip\"\n\n        checkpoint_count = len(\n            list(checkpoint_dir.glob(\"train_tabular_features_batch_*.csv\"))\n        )\n        image_count = len(list(image_dir.glob(\"*.png\")))\n\n        print(f\"Restored checkpoint files : {checkpoint_count}\")\n        print(f\"Restored PNG images       : {image_count}\")\n        print(f\"Existing image ZIP        : {image_zip_path.exists()}\")\n\n        if checkpoint_count == 0:\n            raise RuntimeError(\n                \"The selected checkpoint input contains no batch CSV files.\"\n            )\n\n        if image_count == 0 and not image_zip_path.exists():\n            raise RuntimeError(\n                \"The selected checkpoint input contains neither PNG images \"\n                \"nor a restorable image ZIP.\"\n            )\n\n        print(\"Previous Phase 1 checkpoint restored successfully.\")\n\n    else:\n        raise RuntimeError(\n            \"Multiple Phase 1 checkpoint inputs were found. \"\n            \"Remove unnecessary inputs or set CHECKPOINT_INPUT_HINT \"\n            \"to select exactly one source.\"\n        )\n","metadata":{"execution":{"iopub.execute_input":"2026-07-21T04:53:56.001538Z","iopub.status.busy":"2026-07-21T04:53:56.001178Z","iopub.status.idle":"2026-07-21T04:54:37.707588Z","shell.execute_reply":"2026-07-21T04:54:37.706514Z"},"papermill":{"duration":41.716452,"end_time":"2026-07-21T04:54:37.712056+00:00","exception":false,"start_time":"2026-07-21T04:53:55.995604+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"83002421","cell_type":"markdown","source":"### Imports and Runtime Setup\nImport dependencies, validate inputs, restore images from an existing ZIP\nwhen necessary, and initialize all runtime paths.\n","metadata":{"papermill":{"duration":0.003556,"end_time":"2026-07-21T04:54:37.719484+00:00","exception":false,"start_time":"2026-07-21T04:54:37.715928+00:00","status":"completed"},"tags":[]}},{"id":"49ff16c7","cell_type":"code","source":"import os\n\n# Limit numerical-library threads before importing NumPy/Pandas. Each worker handles one file.\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# --- Output Paths ---\nOUTPUT_PATH = Path(OUTPUT_DIR)\nIMAGE_DIR = OUTPUT_PATH / IMAGE_DIR_NAME\nIMAGE_ZIP_PATH = OUTPUT_PATH / f\"{IMAGE_DIR_NAME}.zip\"\nBATCH_CSV_DIR = OUTPUT_PATH / \"tabular_batches\"\nFAILED_LOG_PATH = OUTPUT_PATH / \"failed_ids.csv\"\nMANIFEST_PATH = OUTPUT_PATH / \"phase1_manifest.json\"\n\n# Prefer RAM disk; fall back to Kaggle working storage.\nshm = Path(\"/dev/shm\")\nTEMP_ROOT = (\n    shm / \"malware_phase1_batch\"\n    if shm.exists() and os.access(shm, os.W_OK)\n    else Path(\"/kaggle/working/malware_phase1_batch\")\n)\n\nOUTPUT_PATH.mkdir(parents=True, exist_ok=True)\nIMAGE_DIR.mkdir(parents=True, exist_ok=True)\nBATCH_CSV_DIR.mkdir(parents=True, exist_ok=True)\nTEMP_ROOT.mkdir(parents=True, exist_ok=True)\n\n# --- Validate Inputs and Runtime ---\nif not Path(ARCHIVE_PATH).exists():\n    raise FileNotFoundError(f\"Archive not found: {ARCHIVE_PATH}\")\nif not Path(LABEL_PATH).exists():\n    raise FileNotFoundError(f\"Labels not found: {LABEL_PATH}\")\nif shutil.which(\"7z\") is None:\n    raise RuntimeError(\"The `7z` executable is not available in this Kaggle runtime.\")\n\n# Restore images when a previous completed run removed the image folder after zipping.\nif (\n    RESTORE_IMAGES_FROM_ZIP_ON_RESUME\n    and IMAGE_ZIP_PATH.exists()\n    and not any(IMAGE_DIR.glob(\"*.png\"))\n):\n    print(f\"Restoring images from existing archive: {IMAGE_ZIP_PATH.name}\")\n    with zipfile.ZipFile(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(IMAGE_DIR)\n\n# --- Load Target IDs ---\nlabels_df = pd.read_csv(LABEL_PATH, dtype={\"Id\": str})\nif \"Id\" not in labels_df.columns:\n    raise ValueError(\"The label file must contain an `Id` column.\")\nif labels_df[\"Id\"].duplicated().any():\n    duplicated = labels_df.loc[labels_df[\"Id\"].duplicated(), \"Id\"].head(10).tolist()\n    raise ValueError(f\"Duplicate IDs in label file, examples: {duplicated}\")\n\ntarget_ids = labels_df[\"Id\"].astype(str).tolist()\nif DEBUG_MAX_FILES is not None:\n    target_ids = target_ids[: int(DEBUG_MAX_FILES)]\n\nMAX_WORKERS = max(1, min(int(MAX_WORKERS), os.cpu_count() or 2))\n\nprint(f\"Total files to process : {len(target_ids):,}\")\nprint(f\"Batch size             : {BATCH_SIZE}\")\nprint(f\"Workers                : {MAX_WORKERS}\")\nprint(f\"Maximum batch attempts : {MAX_BATCH_ATTEMPTS}\")\nprint(f\"Output image size      : {IMAGE_SIZE} x {IMAGE_SIZE}\")\nprint(f\"Temporary directory    : {TEMP_ROOT}\")\n","metadata":{"execution":{"iopub.execute_input":"2026-07-21T04:54:37.729295Z","iopub.status.busy":"2026-07-21T04:54:37.728955Z","iopub.status.idle":"2026-07-21T04:54:38.762495Z","shell.execute_reply":"2026-07-21T04:54:38.761273Z"},"papermill":{"duration":1.040875,"end_time":"2026-07-21T04:54:38.764510+00:00","exception":false,"start_time":"2026-07-21T04:54:37.723635+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"7f7c9407","cell_type":"markdown","source":"### Feature Schema and Utility Functions\nDefine the byte/opcode vocabularies and utility functions used for entropy, atomic writes, image validation, and checkpoint ordering.\n","metadata":{"papermill":{"duration":0.004182,"end_time":"2026-07-21T04:54:38.772527+00:00","exception":false,"start_time":"2026-07-21T04:54:38.768345+00:00","status":"completed"},"tags":[]}},{"id":"b0e5e7db","cell_type":"code","source":"# --- Vocabulary ---\nSEGMENTS = [\".text\", \".data\", \".bss\", \".rdata\", \".edata\", \".idata\", \".rsrc\", \".tls\"]\n\nOPCODES = [\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\nASM_KEYWORDS = [\"db\", \"dw\", \"dd\", \"dq\", \"proc\", \"endp\", \"align\", \"extrn\"]\n\nHEX_LOOKUP = {f\"{i:02X}\": i for i in range(256)}\nOPCODE_PATTERN = re.compile(\n    r\"\\b(\" + \"|\".join(map(re.escape, OPCODES)) + r\")\\b\",\n    flags=re.IGNORECASE,\n)\nKEYWORD_PATTERN = re.compile(\n    r\"\\b(\" + \"|\".join(map(re.escape, ASM_KEYWORDS)) + r\")\\b\",\n    flags=re.IGNORECASE,\n)\n\n\ndef 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\n\ndef 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\n\ndef safe_remove(path: Path) -> None:\n    try:\n        path.unlink(missing_ok=True)\n    except Exception:\n        pass\n\n\ndef 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\n\ndef 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\n\ndef 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\n\ndef 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":{"execution":{"iopub.execute_input":"2026-07-21T04:54:38.781915Z","iopub.status.busy":"2026-07-21T04:54:38.781544Z","iopub.status.idle":"2026-07-21T04:54:38.796908Z","shell.execute_reply":"2026-07-21T04:54:38.795884Z"},"papermill":{"duration":0.022592,"end_time":"2026-07-21T04:54:38.798957+00:00","exception":false,"start_time":"2026-07-21T04:54:38.776365+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"8e6f4328","cell_type":"markdown","source":"### Process One Malware Sample\nStream one `.bytes` file and one `.asm` file, create the grayscale image atomically, and return one tabular feature row. Missing-file flags are intentionally not used as model features because every valid sample must contain both files.\n","metadata":{"papermill":{"duration":0.003937,"end_time":"2026-07-21T04:54:38.806836+00:00","exception":false,"start_time":"2026-07-21T04:54:38.802899+00:00","status":"completed"},"tags":[]}},{"id":"0acb3584","cell_type":"code","source":"# --- Parse .bytes File ---\ndef 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    # Stream the text file so the entire source is never held in memory.\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                    # Unknown bytes remain 0 in the image, while their ratio is retained separately.\n                    values.append(0)\n                    continue\n\n                value = 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 = 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)  # PIL infers grayscale mode; avoids deprecated `mode=`.\n    image = image.resize((IMAGE_SIZE, IMAGE_SIZE), Image.Resampling.BILINEAR)\n\n    temporary_image_path = image_path.with_suffix(\".tmp.png\")\n    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        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\": 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 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 SEGMENTS:\n                if segment + \":\" in lower:\n                    segment_counts[segment] += 1\n\n            opcode_counts.update(OPCODE_PATTERN.findall(lower))\n            keyword_counts.update(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 SEGMENTS:\n        name = segment.replace(\".\", \"\")\n        features[f\"seg_count_{name}\"] = int(segment_counts[segment])\n\n    for opcode in 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 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 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(parse_bytes_file(byte_path, image_path))\n        result.update(parse_asm_file(asm_path))\n        result[\"_error\"] = None\n    except Exception as exc:\n        safe_remove(image_path)\n        result[\"_error\"] = f\"{type(exc).__name__}: {exc}\"\n    finally:\n        safe_remove(byte_path)\n        safe_remove(asm_path)\n\n    return result\n","metadata":{"execution":{"iopub.execute_input":"2026-07-21T04:54:38.817454Z","iopub.status.busy":"2026-07-21T04:54:38.817141Z","iopub.status.idle":"2026-07-21T04:54:38.840117Z","shell.execute_reply":"2026-07-21T04:54:38.839213Z"},"papermill":{"duration":0.030306,"end_time":"2026-07-21T04:54:38.842256+00:00","exception":false,"start_time":"2026-07-21T04:54:38.811950+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"4c32f16b","cell_type":"markdown","source":"### Extract and Process a Batch\nExtract only the requested files from `train.7z`, process samples in parallel, and always remove temporary raw files after each attempt.\n","metadata":{"papermill":{"duration":0.003781,"end_time":"2026-07-21T04:54:38.849891+00:00","exception":false,"start_time":"2026-07-21T04:54:38.846110+00:00","status":"completed"},"tags":[]}},{"id":"efe2a93e","cell_type":"code","source":"def 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 extract_batch(id_list, batch_number: int) -> Path:\n    batch_dir = TEMP_ROOT / f\"batch_{batch_number:04d}\"\n    reset_batch_directory(batch_dir)\n\n    list_path = TEMP_ROOT / f\"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\"{MODE}/{fid}.bytes\\n\")\n            list_file.write(f\"{MODE}/{fid}.asm\\n\")\n\n    command = [\n        \"7z\",\n        \"e\",\n        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    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 process_batch(id_list, batch_number: int):\n    batch_dir = extract_batch(id_list, batch_number)\n    successes, failures = [], []\n\n    try:\n        with ProcessPoolExecutor(max_workers=MAX_WORKERS) as executor:\n            futures = {\n                executor.submit(\n                    process_single_file,\n                    fid,\n                    str(batch_dir),\n                    str(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":{"execution":{"iopub.execute_input":"2026-07-21T04:54:38.859665Z","iopub.status.busy":"2026-07-21T04:54:38.859350Z","iopub.status.idle":"2026-07-21T04:54:38.874698Z","shell.execute_reply":"2026-07-21T04:54:38.873854Z"},"papermill":{"duration":0.023184,"end_time":"2026-07-21T04:54:38.877246+00:00","exception":false,"start_time":"2026-07-21T04:54:38.854062+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"7ca3d012","cell_type":"markdown","source":"### Checkpointing and Full Run\nResume only samples whose checkpoint row **and** PNG both exist. Each batch is retried automatically, and successful checkpoint files are written atomically.\n","metadata":{"papermill":{"duration":0.009086,"end_time":"2026-07-21T04:54:38.896750+00:00","exception":false,"start_time":"2026-07-21T04:54:38.887664+00:00","status":"completed"},"tags":[]}},{"id":"ae6bf224","cell_type":"code","source":"# --- Resume from Existing Checkpoints ---\nexisting_batch_files = sorted(\n    BATCH_CSV_DIR.glob(\"train_tabular_features_batch_*.csv\"),\n    key=numeric_batch_id,\n)\n\nprocessed_ids = set()\nexisting_batch_numbers = []\nstale_checkpoint_ids = set()\ncorrupted_checkpoints = []\n\nfor path in existing_batch_files:\n    existing_batch_numbers.append(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 target_ids and image_is_valid(IMAGE_DIR / f\"{fid}.png\"):\n                processed_ids.add(fid)\n            elif fid in target_ids:\n                stale_checkpoint_ids.add(fid)\n    except Exception as exc:\n        corrupted_checkpoints.append({\"file\": path.name, \"error\": str(exc)})\n        print(f\"Ignoring unreadable checkpoint {path.name}: {exc}\")\n\nremaining_ids = [fid for fid in target_ids if fid not in processed_ids]\nnext_batch_number = max(existing_batch_numbers, default=0) + 1\n\nprint(f\"Valid checkpointed samples : {len(processed_ids):,}\")\nprint(f\"Stale checkpoint rows      : {len(stale_checkpoint_ids):,}\")\nprint(f\"Remaining samples          : {len(remaining_ids):,}\")\n\n# --- Main Processing Loop ---\nall_failures = []\nstart_time = time.time()\n\nfor offset in range(0, len(remaining_ids), BATCH_SIZE):\n    batch_ids = remaining_ids[offset : offset + BATCH_SIZE]\n    batch_number = next_batch_number + offset // 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, MAX_BATCH_ATTEMPTS + 1):\n        if not pending_ids:\n            break\n\n        if attempt > 1:\n            print(f\"  Retry {attempt}/{MAX_BATCH_ATTEMPTS}: {len(pending_ids)} samples\")\n            time.sleep(RETRY_DELAY_SECONDS)\n\n        try:\n            successes, failures = 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            BATCH_CSV_DIR / f\"train_tabular_features_batch_{batch_number:04d}.csv\"\n        )\n        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        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() - start_time) / 60\n    completed_now = min(offset + len(batch_ids), len(remaining_ids))\n    print(\n        f\"  Progress   : {completed_now:,}/{len(remaining_ids):,} new samples | \"\n        f\"{elapsed_minutes:.1f} minutes\"\n    )\n    gc.collect()\n\n# --- Save Current Failure Log ---\nif all_failures:\n    failure_df = pd.DataFrame(all_failures).drop_duplicates(\"Id\", keep=\"last\")\n    atomic_write_csv(failure_df, FAILED_LOG_PATH)\n    print(f\"\\nFailure log saved: {FAILED_LOG_PATH}\")\nelif FAILED_LOG_PATH.exists():\n    FAILED_LOG_PATH.unlink()\n","metadata":{"execution":{"iopub.execute_input":"2026-07-21T04:54:38.916379Z","iopub.status.busy":"2026-07-21T04:54:38.915696Z","iopub.status.idle":"2026-07-21T05:08:58.866374Z","shell.execute_reply":"2026-07-21T05:08:58.865001Z"},"papermill":{"duration":859.967909,"end_time":"2026-07-21T05:08:58.872155+00:00","exception":false,"start_time":"2026-07-21T04:54:38.904246+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"14505c8a","cell_type":"markdown","source":"### Validate and Package Output\nMerge checkpoints, remove zero-variance features, require complete feature/image coverage, save tabular files, validate the image ZIP, and write a complete Phase 1 manifest.\n","metadata":{"papermill":{"duration":0.004107,"end_time":"2026-07-21T05:08:58.880536+00:00","exception":false,"start_time":"2026-07-21T05:08:58.876429+00:00","status":"completed"},"tags":[]}},{"id":"5fad051d","cell_type":"code","source":"# --- Combine Batch Results ---\nbatch_files = sorted(\n    BATCH_CSV_DIR.glob(\"train_tabular_features_batch_*.csv\"),\n    key=numeric_batch_id,\n)\nif not batch_files:\n    raise RuntimeError(\"No valid batch checkpoint CSV files were produced.\")\n\nframes = []\nfor path in 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\"] = numeric_batch_id(path)\n    frames.append(frame)\n\nraw_features = pd.concat(frames, ignore_index=True, sort=False)\nraw_features[\"Id\"] = raw_features[\"Id\"].astype(str)\n\nexpected_ids = set(target_ids)\nunexpected_feature_ids = sorted(set(raw_features[\"Id\"]) - expected_ids)\nraw_features = raw_features.loc[raw_features[\"Id\"].isin(expected_ids)].copy()\n\nduplicate_rows_resolved = int(raw_features.duplicated(\"Id\", keep=False).sum())\nall_features = (\n    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.\nlegacy_constant_columns = [\n    column for column in [\"has_bytes\", \"has_asm\"] if column in all_features.columns\n]\nif legacy_constant_columns:\n    all_features = all_features.drop(columns=legacy_constant_columns)\n\ncandidate_feature_columns = [column for column in all_features.columns if column != \"Id\"]\nconstant_feature_columns = [\n    column\n    for column in candidate_feature_columns\n    if all_features[column].nunique(dropna=False) <= 1\n]\nif constant_feature_columns:\n    all_features = all_features.drop(columns=constant_feature_columns)\n\norder = {fid: index for index, fid in enumerate(target_ids)}\nall_features[\"_target_order\"] = all_features[\"Id\"].map(order)\nall_features = (\n    all_features.sort_values(\"_target_order\")\n    .drop(columns=[\"_target_order\"])\n    .reset_index(drop=True)\n)\n\n# --- Strict Completeness and Integrity Validation ---\nfeature_ids = set(all_features[\"Id\"])\nmissing_feature_ids = sorted(expected_ids - feature_ids)\nmissing_image_ids = [\n    fid for fid in target_ids if not image_is_valid(IMAGE_DIR / f\"{fid}.png\")\n]\nfinal_duplicate_count = int(all_features[\"Id\"].duplicated().sum())\nfeature_columns = [column for column in all_features.columns if column != \"Id\"]\n\nnon_numeric_columns = [\n    column\n    for column in feature_columns\n    if not pd.api.types.is_numeric_dtype(all_features[column])\n]\ninvalid_value_count = 0\nif not non_numeric_columns and feature_columns:\n    numeric_values = all_features[feature_columns].to_numpy(dtype=np.float64, copy=False)\n    invalid_value_count = int((~np.isfinite(numeric_values)).sum())\n\nprint(f\"Feature row count         : {len(all_features):,}\")\nprint(f\"Feature column count      : {len(feature_columns):,}\")\nprint(f\"Final duplicate IDs       : {final_duplicate_count}\")\nprint(f\"Resolved duplicate rows   : {duplicate_rows_resolved}\")\nprint(f\"Missing feature rows      : {len(missing_feature_ids):,}\")\nprint(f\"Missing/empty images      : {len(missing_image_ids):,}\")\nprint(f\"Non-numeric columns       : {len(non_numeric_columns):,}\")\nprint(f\"NaN/Inf values            : {invalid_value_count:,}\")\nprint(f\"Removed constant features : {len(legacy_constant_columns) + len(constant_feature_columns):,}\")\n\nvalidation_errors = []\nif missing_feature_ids:\n    validation_errors.append(f\"missing feature IDs, examples: {missing_feature_ids[:10]}\")\nif missing_image_ids:\n    validation_errors.append(f\"missing image IDs, examples: {missing_image_ids[:10]}\")\nif final_duplicate_count:\n    validation_errors.append(f\"{final_duplicate_count} duplicate IDs remain\")\nif non_numeric_columns:\n    validation_errors.append(f\"non-numeric columns: {non_numeric_columns[:10]}\")\nif invalid_value_count:\n    validation_errors.append(f\"{invalid_value_count} NaN/Inf values\")\n\nif validation_errors and FAIL_ON_INCOMPLETE_OUTPUT:\n    raise RuntimeError(\"Phase 1 validation failed: \" + \"; \".join(validation_errors))\n\n# --- Save Combined Tabular Outputs Atomically ---\ncombined_csv_path = OUTPUT_PATH / \"train_tabular_features.csv\"\natomic_write_csv(all_features, combined_csv_path)\n\ncombined_parquet_path = OUTPUT_PATH / \"train_tabular_features.parquet\"\nparquet_saved = False\ntry:\n    temporary_parquet = combined_parquet_path.with_suffix(\".parquet.tmp\")\n    all_features.to_parquet(temporary_parquet, index=False)\n    os.replace(temporary_parquet, combined_parquet_path)\n    parquet_saved = True\nexcept Exception as exc:\n    safe_remove(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 ---\nsample_ids = all_features[\"Id\"].head(25).tolist()\ncanvas = Image.new(\"L\", (IMAGE_SIZE * 5, IMAGE_SIZE * 5), color=0)\nfor index, fid in enumerate(sample_ids):\n    path = IMAGE_DIR / f\"{fid}.png\"\n    with Image.open(path) as image:\n        canvas.paste(\n            image.convert(\"L\"),\n            ((index % 5) * IMAGE_SIZE, (index // 5) * IMAGE_SIZE),\n        )\ncanvas.save(OUTPUT_PATH / \"sample_malware_images.png\")\n\n# --- Create and Validate the Image Archive Atomically ---\nif ZIP_IMAGES:\n    print(\"Compressing the validated image folder...\")\n    temporary_base = OUTPUT_PATH / f\".{IMAGE_DIR_NAME}_temporary\"\n    temporary_zip = temporary_base.with_suffix(\".zip\")\n    safe_remove(temporary_zip)\n\n    created_archive = Path(\n        shutil.make_archive(str(temporary_base), \"zip\", root_dir=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        safe_remove(created_archive)\n        raise RuntimeError(f\"Image ZIP validation failed at member: {bad_member}\")\n    if png_count != len(target_ids):\n        safe_remove(created_archive)\n        raise RuntimeError(\n            f\"Image ZIP contains {png_count:,} PNG files; expected {len(target_ids):,}.\"\n        )\n\n    os.replace(created_archive, IMAGE_ZIP_PATH)\n    print(f\"Image archive created: {IMAGE_ZIP_PATH}\")\n\n# --- Write Complete Manifest Atomically ---\nremoved_constant_features = legacy_constant_columns + constant_feature_columns\nmanifest = {\n    \"phase\": 1,\n    \"version\": \"3.0\",\n    \"status\": \"complete\" if not validation_errors else \"incomplete\",\n    \"completed_at_utc\": datetime.now(timezone.utc).isoformat(),\n    \"mode\": MODE,\n    \"sample_count\": int(len(all_features)),\n    \"target_count\": int(len(target_ids)),\n    \"feature_count\": int(len(feature_columns)),\n    \"feature_columns\": feature_columns,\n    \"constant_features_removed\": removed_constant_features,\n    \"image_conversion\": {\n        \"method\": \"adaptive_width_grayscale\",\n        \"unknown_byte_pixel_value\": 0,\n        \"resize\": \"bilinear\",\n        \"image_size\": int(IMAGE_SIZE),\n    },\n    \"image_size\": int(IMAGE_SIZE),\n    \"image_folder_name\": IMAGE_DIR_NAME,\n    \"image_zip_name\": IMAGE_ZIP_PATH.name if ZIP_IMAGES else None,\n    \"combined_csv\": combined_csv_path.name,\n    \"combined_parquet\": combined_parquet_path.name if parquet_saved else None,\n    \"missing_feature_count\": int(len(missing_feature_ids)),\n    \"missing_image_count\": int(len(missing_image_ids)),\n    \"unexpected_feature_id_count\": int(len(unexpected_feature_ids)),\n    \"duplicate_rows_resolved\": int(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}\natomic_write_text(\n    json.dumps(manifest, ensure_ascii=False, indent=2),\n    MANIFEST_PATH,\n)\n\n# Remove the image folder only after the ZIP and manifest are successfully validated.\nif ZIP_IMAGES and REMOVE_IMAGE_FOLDER_AFTER_ZIP:\n    shutil.rmtree(IMAGE_DIR, ignore_errors=True)\n\nprint(\"\\n\" + \"=\" * 72)\nprint(\"PHASE 1 COMPLETE\")\nprint(f\"Output directory : {OUTPUT_PATH}\")\nprint(f\"Manifest         : {MANIFEST_PATH.name}\")\nprint(f\"Tabular CSV      : {combined_csv_path.name}\")\nif parquet_saved:\n    print(f\"Tabular Parquet  : {combined_parquet_path.name}\")\nif ZIP_IMAGES:\n    print(f\"Image ZIP        : {IMAGE_ZIP_PATH.name}\")\nprint(f\"Samples          : {len(all_features):,}\")\nprint(f\"Features         : {len(feature_columns):,}\")\nprint(\"=\" * 72)\n","metadata":{"execution":{"iopub.execute_input":"2026-07-21T05:08:58.891364Z","iopub.status.busy":"2026-07-21T05:08:58.890635Z","iopub.status.idle":"2026-07-21T05:09:31.088878Z","shell.execute_reply":"2026-07-21T05:09:31.087745Z"},"papermill":{"duration":32.206883,"end_time":"2026-07-21T05:09:31.091525+00:00","exception":false,"start_time":"2026-07-21T05:08:58.884642+00:00","status":"completed"},"tags":[]},"outputs":[],"execution_count":null},{"id":"145d8417-5a59-4140-8ea3-63fe3a1e9fd3","cell_type":"markdown","source":"### Final Validated Artifact\n\nThe recorded final execution completed all **10,868 samples** with\n**590 non-constant numeric features**, zero duplicate IDs, zero missing\nfeature rows, zero missing images, and zero invalid numeric values.\nThe image archive and manifest were validated before the pipeline reported\n`PHASE 1 COMPLETE`.\n","metadata":{}}]}