#!/usr/bin/python3
"""Installed only inside an isolated disposable Linux extractor VM.
Never run on PVE, the control CT, or the Windows source. No repair or write mount.
"""
import json
import os
import re
import shutil
import stat
import subprocess
import sys
from pathlib import Path


def command(argv, timeout=30):
    return subprocess.run(argv, check=True, text=True, capture_output=True, timeout=timeout).stdout


def atomic_json(path, value):
    temp = path.with_suffix('.tmp')
    with temp.open('x', encoding='utf-8') as handle:
        json.dump(value, handle, ensure_ascii=True)
        handle.flush()
        os.fsync(handle.fileno())
    os.replace(temp, path)


def exact_directory(root, parts):
    current = root
    for part in parts:
        current = current / part
        if current.is_symlink() or not current.is_dir():
            raise RuntimeError('expected evidence directory missing or symlink')
    return current


def main():
    if os.geteuid() != 0 or len(sys.argv) != 2:
        raise RuntimeError('fixed extractor requires root QGA control and one manifest')
    manifest_path = Path(sys.argv[1])
    if manifest_path.parent != Path('/var/lib/otche/control') or not re.fullmatch(r'[0-9a-f-]{36}\.json', manifest_path.name):
        raise RuntimeError('manifest path outside fixed control directory')
    if manifest_path.stat().st_size > 40960:
        raise RuntimeError('manifest too large')
    spec = json.loads(manifest_path.read_text(encoding='utf-8'))
    cid = spec['command_id']
    if not re.fullmatch(r'[0-9a-f]{8}(-[0-9a-f]{4}){3}-[0-9a-f]{12}', cid) or manifest_path.stem != cid:
        raise RuntimeError('invalid command identity')
    maximum = int(spec['max_bytes'])
    if not 1 <= maximum <= 1073741824 or not re.fullmatch(r'otche[0-9a-f]{16}', spec['serial']):
        raise RuntimeError('invalid collection limits or disk identity')
    target = Path('/var/lib/otche/export') / cid
    target.mkdir(mode=0o700, parents=False, exist_ok=False)  # durable acceptance; never repeat
    report = {'command_id': cid, 'complete': False, 'errors': [], 'files': []}
    atomic_json(target / 'receipt.json', {'command_id': cid, 'state': 'accepted'})
    mount = Path('/var/lib/otche/mount') / cid
    mount.mkdir(mode=0o700, parents=False, exist_ok=False)
    mounted = False
    try:
        listing = json.loads(command(['lsblk', '-J', '-b', '-o', 'PATH,TYPE,SERIAL,SIZE,RO,FSTYPE']))
        matching = [x for x in listing['blockdevices'] if x.get('type') == 'disk' and str(x.get('serial', '')).strip() == spec['serial']]
        if len(matching) != 1:
            raise RuntimeError('allocated evidence serial is not unique')
        disk = matching[0]
        if int(disk['size']) != int(spec['disk_bytes']) or str(disk['ro']).lower() not in ('1', 'true'):
            raise RuntimeError('evidence size or hypervisor read-only property mismatch')
        if command(['blockdev', '--getro', disk['path']]).strip() != '1':
            raise RuntimeError('kernel does not report block device read-only')
        partitions = [disk] + disk.get('children', [])
        found = False
        copied = 0
        for partition in partitions:
            if str(partition.get('fstype', '')).lower() != 'ntfs':
                continue
            if str(partition.get('ro')).lower() not in ('1', 'true'):
                raise RuntimeError('partition is not read-only')
            try:
                command(['ntfs-3g', '-o', 'ro,norecover,noexec,nodev,nosuid', partition['path'], str(mount)])
                mounted = True
                try:
                    folder = exact_directory(mount, ['ProgramData', 'Otche', 'results', cid])
                except RuntimeError:
                    command(['umount', str(mount)])
                    mounted = False
                    continue
                found = True
                candidates = []
                with os.scandir(folder) as files:
                    for entry in files:
                        if entry.is_file(follow_symlinks=False):
                            candidates.append(Path(entry.path))
                if spec.get('include_crash_dump'):
                    windows = exact_directory(mount, ['Windows'])
                    dump = windows / 'MEMORY.DMP'
                    if dump.is_file() and not dump.is_symlink():
                        candidates.append(dump)
                    else:
                        report['errors'].append('Requested system crash dump unavailable')
                if len(candidates) > 64:
                    raise RuntimeError('artifact count exceeds bound')
                for index, source in enumerate(sorted(candidates)):
                    name = re.sub(r'[^A-Za-z0-9_.-]', '_', source.name)[:120]
                    output_name = f'{index:03d}-{name}'
                    fd = os.open(source, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK)
                    try:
                        size = os.fstat(fd).st_size
                        if not stat.S_ISREG(os.fstat(fd).st_mode):
                            raise RuntimeError('evidence is not a regular file')
                        if copied + size > maximum:
                            report['errors'].append('Artifact byte budget exceeded: ' + name)
                            continue
                        with os.fdopen(fd, 'rb', closefd=False) as src, (target / output_name).open('xb') as dst:
                            remaining = size
                            while remaining:
                                data = src.read(min(262144, remaining))
                                if not data:
                                    raise RuntimeError('evidence file short read')
                                dst.write(data)
                                remaining -= len(data)
                            dst.flush()
                            os.fsync(dst.fileno())
                        copied += size
                        report['files'].append({'name': output_name, 'size': size})
                    finally:
                        os.close(fd)
                command(['umount', str(mount)])
                mounted = False
                break
            finally:
                if mounted:
                    command(['umount', str(mount)])
                    mounted = False
        if not found:
            report['errors'].append('Command directory unavailable: dirty/encrypted/unsupported filesystem or no flushed evidence')
        report['complete'] = found and not report['errors']
    except Exception as error:
        report['errors'].append(type(error).__name__ + ': ' + str(error)[:500])
    finally:
        if mounted:
            try:
                command(['umount', str(mount)])
            except Exception:
                report['errors'].append('Unmount incomplete; preserve stopped extractor')
                report['complete'] = False
        atomic_json(target / 'manifest.json', report)
    print(json.dumps({'command_id': cid, 'complete': report['complete']}))


if __name__ == '__main__':
    main()
