{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[],"dockerImageVersionId":28755,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Microsoft Malware Classification 2015 — Phase 1\n\n**Author:** Huynh Bao Khang, Nguyen Phuoc Khang, Pham Minh Dat\n\n**Objective:** Convert the raw `.bytes` + `.asm` files from the competition archive into a clean, checkpointed set of tabular features and fixed-size grayscale images, ready for the Phase 2 hybrid model.\n","metadata":{}},{"cell_type":"markdown","source":"### Configuration\nSet the dataset paths, batch size, worker count, and output image size used throughout the pipeline.","metadata":{}},{"cell_type":"code","source":"MODE = \"train\"\n\n# --- Directories ---\nARCHIVE_PATH = \"/kaggle/input/competitions/malware-classification/train.7z\"\nLABEL_PATH = \"/kaggle/input/competitions/malware-classification/trainLabels.csv\"\n\nOUTPUT_DIR = \"/kaggle/working/malware_phase1_output\"\nIMAGE_DIR_NAME = \"train_malware_images\"\n\n# --- Batching and Parallelism ---\nBATCH_SIZE = 64\nMAX_WORKERS = 4\n\n# --- Image Output ---\nIMAGE_SIZE = 224\n\n# --- Packaging ---\nZIP_IMAGES = True\nREMOVE_IMAGE_FOLDER_AFTER_ZIP = True\n\n# --- Debug ---\nDEBUG_MAX_FILES = None","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Imports and Config\nImport libraries, set up input/output directories, and load the target file IDs from the label file.","metadata":{}},{"cell_type":"code","source":"import os\nimport re\nimport gc\nimport json\nimport math\nimport time\nimport shutil\nimport hashlib\nimport subprocess\nfrom pathlib import Path\nfrom collections import Counter\nfrom concurrent.futures import ProcessPoolExecutor, as_completed\n\nimport numpy as np\nimport pandas as pd\nfrom PIL import Image, ImageOps\n\n# --- Thread Limits ---\nos.environ[\"OMP_NUM_THREADS\"] = \"1\"\nos.environ[\"MKL_NUM_THREADS\"] = \"1\"\n\n# --- Output Paths ---\nOUTPUT_PATH = Path(OUTPUT_DIR)\nIMAGE_DIR = OUTPUT_PATH / IMAGE_DIR_NAME\nBATCH_CSV_DIR = OUTPUT_PATH / \"tabular_batches\"\nFAILED_LOG_PATH = OUTPUT_PATH / \"failed_ids.csv\"\nMANIFEST_PATH = OUTPUT_PATH / \"phase1_manifest.json\"\n\n# --- Temporary Directory ---\n# Prefer RAM disk; fall back to /kaggle/working if it's not usable.\nshm = Path(\"/dev/shm\")\nTEMP_ROOT = shm / \"malware_phase1_batch\" if shm.exists() and os.access(shm, os.W_OK) else Path(\"/kaggle/working/malware_phase1_batch\")\n\n# --- Create Directories ---\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 ---\nassert Path(ARCHIVE_PATH).exists(), f\"Archive not found: {ARCHIVE_PATH}\"\nassert Path(LABEL_PATH).exists(), f\"Labels not found: {LABEL_PATH}\"\n\n# --- Load Target IDs ---\nlabels_df = pd.read_csv(LABEL_PATH, dtype={\"Id\": str})\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\"Output image size: {IMAGE_SIZE} x {IMAGE_SIZE}\")\nprint(f\"Temporary directory: {TEMP_ROOT}\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Feature Schema\nDefine the byte, opcode, and PE-section vocabularies used for feature extraction, along with a few small helper functions.","metadata":{}},{"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\n# --- Lookup Tables and Patterns ---\nHEX_LOOKUP = {f\"{i:02X}\": i for i in range(256)}\nOPCODE_PATTERN = re.compile(r\"\\b(\" + \"|\".join(map(re.escape, OPCODES)) + r\")\\b\", flags=re.IGNORECASE)\nKEYWORD_PATTERN = re.compile(r\"\\b(\" + \"|\".join(map(re.escape, ASM_KEYWORDS)) + r\")\\b\", flags=re.IGNORECASE)\n\n# --- Helper Functions ---\ndef get_image_width(raw_byte_count: int) -> int:\n    \"\"\"Choose image width based on raw byte count, not the size of the `.bytes` text file.\"\"\"\n    size_kb = raw_byte_count / 1024.0\n    if size_kb < 10: return 64\n    if size_kb < 30: return 128\n    if size_kb < 60: return 256\n    if size_kb < 100: return 384\n    if size_kb < 200: return 512\n    if size_kb < 500: return 768\n    if size_kb < 1000: return 1024\n    return 2048\n\ndef entropy_from_counts(counts: np.ndarray) -> float:\n    total = int(counts.sum())\n    if total == 0:\n        return 0.0\n    p = counts[counts > 0].astype(np.float64) / total\n    return float(-(p * np.log2(p)).sum())\n\ndef safe_remove(path: Path) -> None:\n    try:\n        path.unlink(missing_ok=True)\n    except Exception:\n        pass","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Process a Single File\nParse one sample's `.bytes` and `.asm` files into an image plus a row of tabular features.","metadata":{}},{"cell_type":"code","source":"\n# --- 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 file line by line so the entire text file isn't held in RAM.\n    with byte_path.open(\"r\", encoding = \"latin1\", errors = \"ignore\") as f:\n        for line in f:\n            for token in line.split():\n                token = token.upper()\n                if token == \"??\":\n                    unknown_count += 1\n                    # Image still needs a pixel; use 0 but keep unknown_ratio separate.\n                    values.append(0)\n                else:\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(\"File .bytes does not contain valid bytes\")\n\n    array = np.frombuffer(values, dtype = np.uint8)\n    width = get_image_width(raw_count)\n    height = int(math.ceil(raw_count / width))\n    padded_length = height * width\n    if padded_length > raw_count:\n        array = np.pad(array, (0, padded_length - raw_count), constant_values = 0)\n\n    raw_image = array.reshape(height, width)\n    image = Image.fromarray(raw_image, mode = \"L\")\n    # Resize the whole image to a fixed size; avoid cropping and losing the end of the file.\n    image = image.resize((IMAGE_SIZE, IMAGE_SIZE), Image.Resampling.BILINEAR)\n    image.save(image_path, format = \"PNG\", optimize = False)\n\n    # --- Byte-Level Statistics ---\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\": float(array[:raw_count].mean()),\n        \"byte_std\": float(array[:raw_count].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    # --- Per-Value Counts and Frequencies ---\n    # Keep both raw counts and normalized frequencies: counts capture the size signal,\n    # while frequencies describe the byte distribution independently of file size.\n    frequencies = counts.astype(np.float64) / max(known_count, 1)\n    for i in range(256):\n        features[f\"byte_count_{i:02X}\"] = int(counts[i])\n        features[f\"byte_freq_{i:02X}\"] = float(frequencies[i])\n\n    return features\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 f:\n        for line in f:\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 = 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# --- Combine Both Parsers for One Sample ---\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, \"has_bytes\": int(byte_path.exists()), \"has_asm\": int(asm_path.exists())}\n    try:\n        if not byte_path.exists():\n            raise FileNotFoundError(f\"Lack of {fid}.bytes\")\n        if not asm_path.exists():\n            raise FileNotFoundError(f\"Lack of {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","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Extract and Process a Batch\nExtract a batch of raw files from the archive and process them in parallel.","metadata":{}},{"cell_type":"code","source":"# --- Reset Batch Directory ---\ndef 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# --- Extract Batch from Archive ---\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 f:\n        for fid in id_list:\n            f.write(f\"{MODE}/{fid}.bytes\\n\")\n            f.write(f\"{MODE}/{fid}.asm\\n\")\n\n    command = [\n        \"7z\", \"e\", ARCHIVE_PATH, f\"@{list_path}\",\n        f\"-o{batch_dir}\", \"-y\", \"-bd\"\n    ]\n    completed = subprocess.run(command, capture_output=True, text=True)\n    safe_remove(list_path)\n\n    if completed.returncode != 0:\n        raise RuntimeError(\n            \"7z extraction failed\\n\"\n            + completed.stderr[-1500:]\n            + completed.stdout[-1500:]\n        )\n    return batch_dir\n\n# --- Process Batch in Parallel ---\ndef process_batch(id_list, batch_number: int):\n    batch_dir = extract_batch(id_list, batch_number)\n    successes, failures = [], []\n\n    with ProcessPoolExecutor(max_workers = MAX_WORKERS) as executor:\n        futures = {\n            executor.submit(process_single_file, fid, str(batch_dir), str(IMAGE_DIR)): fid\n            for fid in id_list\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\n    shutil.rmtree(batch_dir, ignore_errors=True)\n    return successes, failures","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Checkpointing and Full Run\nRun all batches with checkpointing, so an interrupted run can resume without reprocessing completed files.","metadata":{}},{"cell_type":"code","source":"# --- Checkpoint Helpers ---\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\n# --- Resume from Existing Checkpoints ---\nexisting_batch_files = sorted(BATCH_CSV_DIR.glob(\"train_tabular_features_batch_*.csv\"), key=numeric_batch_id)\nprocessed_ids = set()\nexisting_batch_numbers = []\n\nfor path in existing_batch_files:\n    try:\n        checkpoint = pd.read_csv(path, usecols=[\"Id\"], dtype={\"Id\": str})\n        processed_ids.update(checkpoint[\"Id\"].astype(str))\n        existing_batch_numbers.append(numeric_batch_id(path))\n    except Exception as exc:\n        print(f\"Bỏ qua checkpoint hỏng {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\"Đã có checkpoint : {len(processed_ids):,} file\")\nprint(f\"Còn lại           : {len(remaining_ids):,} file\")\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)} file\", flush=True)\n\n    try:\n        successes, failures = process_batch(batch_ids, batch_number)\n    except Exception as exc:\n        print(f\"  Lỗi toàn batch: {exc}\")\n        all_failures.extend({\"Id\": fid, \"error\": f\"BatchError: {exc}\"} for fid in batch_ids)\n        continue\n\n    if successes:\n        batch_df = pd.DataFrame(successes).sort_values(\"Id\").reset_index(drop=True)\n        output_path = BATCH_CSV_DIR / f\"train_tabular_features_batch_{batch_number:04d}.csv\"\n        batch_df.to_csv(output_path, index=False)\n        print(f\"  Thành công: {len(successes):3d} | lưu {output_path.name}\")\n    else:\n        print(\"  Không có file nào xử lý thành công\")\n\n    if failures:\n        print(f\"  Thất bại : {len(failures):3d} | các file này sẽ chưa được checkpoint\")\n        all_failures.extend(failures)\n\n    elapsed = (time.time() - start_time) / 60\n    completed_now = min(offset + len(batch_ids), len(remaining_ids))\n    print(f\"  Tiến độ  : {completed_now:,}/{len(remaining_ids):,} file mới | {elapsed:.1f} phút\")\n    gc.collect()\n\n# --- Save Failure Log ---\nif all_failures:\n    pd.DataFrame(all_failures).drop_duplicates(\"Id\", keep=\"last\").to_csv(FAILED_LOG_PATH, index=False)\n    print(f\"\\nĐã lưu danh sách lỗi: {FAILED_LOG_PATH}\")\nelif FAILED_LOG_PATH.exists():\n    FAILED_LOG_PATH.unlink()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Validate and Package Output\nCombine, validate, and package the final tabular features and image archive for Phase 2.","metadata":{}},{"cell_type":"code","source":"# --- Combine Batch Results ---\nbatch_files = sorted(BATCH_CSV_DIR.glob(\"train_tabular_features_batch_*.csv\"), key = numeric_batch_id)\nassert batch_files, \"Have no batch CSV files.\"\n\nall_features = pd.concat(\n    [pd.read_csv(path, dtype={\"Id\": str}) for path in batch_files],\n    ignore_index=True,\n)\nall_features = all_features.drop_duplicates(\"Id\", keep=\"last\").sort_values(\"Id\").reset_index(drop = True)\n\n# --- Validate Completeness ---\nexpected_ids = set(target_ids)\nfeature_ids = set(all_features[\"Id\"])\nmissing_feature_ids = sorted(expected_ids - feature_ids)\nduplicate_count = int(all_features[\"Id\"].duplicated().sum())\nmissing_image_ids = [fid for fid in all_features[\"Id\"] if not (IMAGE_DIR / f\"{fid}.png\").exists()]\n\nprint(f\"Feature row count: {len(all_features):,}\")\nprint(f\"Feature column count: {all_features.shape[1] - 1:,}\")\nprint(f\"Duplicate Id: {duplicate_count}\")\nprint(f\"Missing features: {len(missing_feature_ids):,}\")\nprint(f\"Missing images: {len(missing_image_ids):,}\")\n\nif missing_feature_ids:\n    print(\"Some missing IDs:\", missing_feature_ids[:10])\nif missing_image_ids:\n    print(\"Some missing images:\", missing_image_ids[:10])\n\n# --- Save Combined Outputs ---\n# Save the combined file for fast loading in Phase 2; batch CSVs are kept for checkpointing.\ncombined_csv_path = OUTPUT_PATH / \"train_tabular_features.csv\"\nall_features.to_csv(combined_csv_path, index=False)\n\ncombined_parquet_path = OUTPUT_PATH / \"train_tabular_features.parquet\"\nparquet_saved = False\ntry:\n    all_features.to_parquet(combined_parquet_path, index=False)\n    parquet_saved = True\nexcept Exception as exc:\n    print(f\"Could not save Parquet, Phrase 2 will use CSV instead: {exc}\")\n\n# --- Build Preview Grid ---\nsample_ids = all_features[\"Id\"].head(25).tolist()\ncanvas = Image.new(\"L\", (IMAGE_SIZE * 5, IMAGE_SIZE * 5), color=0)\nfor idx, fid in enumerate(sample_ids):\n    path = IMAGE_DIR / f\"{fid}.png\"\n    if path.exists():\n        with Image.open(path) as img:\n            canvas.paste(img.convert(\"L\"), ((idx % 5) * IMAGE_SIZE, (idx // 5) * IMAGE_SIZE))\ncanvas.save(OUTPUT_PATH / \"sample_malware_images.png\")\n\nimage_zip_path = OUTPUT_PATH / f\"{IMAGE_DIR_NAME}.zip\"\nif ZIP_IMAGES:\n    print(\"Compressing the image folder...\")\n    if image_zip_path.exists():\n        image_zip_path.unlink()\n    shutil.make_archive(str(image_zip_path.with_suffix(\"\")), \"zip\", root_dir=IMAGE_DIR)\n    print(f\"Đã tạo: {image_zip_path}\")\n    if REMOVE_IMAGE_FOLDER_AFTER_ZIP:\n        shutil.rmtree(IMAGE_DIR, ignore_errors=True)\n\n# --- Write Manifest ---\nfeature_columns = [c for c in all_features.columns if c != \"Id\"]\nmanifest = {\n    \"phase\": 1,\n    \"version\": \"2.0\",\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    \"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}\nMANIFEST_PATH.write_text(json.dumps(manifest, ensure_ascii=False, indent=2), encoding=\"utf-8\")\n\n# --- Summary ---\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}\")","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}