Files
dlp-io/dlp_io.py
T
p40000043244@byd.com 7b62c23e1e
publish-dlp-io / publish (push) Successful in 1m16s
feat: is_encrypted 加密状态判断与未加密写入通道 (0.2.0)
- is_encrypted(): certutil 原始哈希 + helper 视图对比, Python/pyd/EXE 三环境均可判定
- w/a/x 写入按后缀探测通道, 直写会被加密时自动切换 PowerShell stdin 中转, 保证落盘未加密
- 新增 WriteChannel/DlpWriteError 与 DLP_IO_WRITE_CHANNEL 强制通道开关
- tests: DLP 加解密样本 fixture 与 19 个结构化断言用例
2026-07-31 18:24:22 +08:00

1194 lines
38 KiB
Python

"""io-compatible file reads through an approved Python helper on Windows DLP hosts."""
from __future__ import annotations
import atexit
import builtins
import errno
import hashlib
import io
import json
import operator
import os
import queue
import re
import shutil
import subprocess
import tempfile
import threading
from dataclasses import dataclass
from enum import Enum
from multiprocessing.context import AuthenticationError
from multiprocessing.connection import Client
__version__ = "0.2.0"
__all__ = [
"DlpConfig",
"DlpConfigurationError",
"DlpHelperStartError",
"DlpIoError",
"DlpProtocolError",
"DlpSession",
"DlpSessionBusyError",
"DlpTransportError",
"DlpWriteError",
"WriteChannel",
"configure",
"install_open_patch",
"install_reader_patch",
"is_encrypted",
"open",
"shutdown",
"uninstall_open_patch",
]
# ---------------------------------------------------------------------------
# Errors
class DlpIoError(Exception):
"""Base exception for dlp-io failures."""
class DlpConfigurationError(DlpIoError):
"""The bridge configuration is invalid for the current lifecycle state."""
class DlpHelperStartError(DlpIoError):
"""The approved Python helper could not be started."""
class DlpProtocolError(DlpIoError):
"""The helper and client exchanged an invalid protocol message."""
class DlpTransportError(DlpIoError):
"""The authenticated pipe stream failed before a valid EOF."""
class DlpSessionBusyError(DlpIoError):
"""A session already owns an active file stream."""
class DlpWriteError(DlpIoError):
"""No available write channel can produce an unencrypted file."""
# ---------------------------------------------------------------------------
# Configuration
@dataclass(frozen=True)
class DlpConfig:
python_executable: str | None = None
startup_timeout: float = 60
def python_command(self) -> list[str] | None:
if self.python_executable is None:
return None
return [self.python_executable]
# ---------------------------------------------------------------------------
# Helper protocol
PROTOCOL_VERSION = 1
CONTROL_LIMIT = 64 * 1024
CHUNK_SIZE = 1024 * 1024
HANDSHAKE_PREFIX = "DLP_IO 1 "
class HelperErrorCode(str, Enum):
INVALID_PATH = "invalid_path"
INVALID_OFFSET = "invalid_offset"
FILE_ACCESS = "file_access"
INVALID_REQUEST = "invalid_request"
INVALID_REQUEST_TYPE = "invalid_request_type"
UNSUPPORTED_VERSION = "unsupported_version"
UNKNOWN_OPERATION = "unknown_operation"
INTERNAL = "internal"
@dataclass(frozen=True)
class HelperSourceConfig:
authkey: bytes
protocol_version: int = PROTOCOL_VERSION
control_limit: int = CONTROL_LIMIT
chunk_size: int = CHUNK_SIZE
@classmethod
def from_authkey(cls, authkey: bytes) -> HelperSourceConfig:
if not isinstance(authkey, bytes) or len(authkey) != 32:
raise DlpConfigurationError("authkey must contain exactly 32 bytes")
return cls(authkey=authkey)
_HELPER_TEMPLATE = r'''
import json
import os
import sys
from multiprocessing.connection import AuthenticationError, Listener
VERSION = __PROTOCOL_VERSION__
CONTROL_LIMIT = __CONTROL_LIMIT__
CHUNK_SIZE = __CHUNK_SIZE__
PIPE_PREFIX = r"\\.\pipe\DlpIo_"
AUTHKEY = bytes.fromhex("__AUTHKEY_HEX__")
def send_json(conn, value):
conn.send_bytes(json.dumps(value).encode("utf-8"))
def error(conn, code, message):
send_json(conn, {
"version": VERSION,
"ok": False,
"error_code": code,
"error": message,
})
def handle_read(conn, request):
path = request.get("path")
offset = request.get("offset", 0)
if not isinstance(path, str) or not path:
error(conn, "invalid_path", "path must be a non-empty string")
return
if isinstance(offset, bool) or not isinstance(offset, int) or offset < 0:
error(conn, "invalid_offset", "offset must be a non-negative integer")
return
try:
size = os.path.getsize(path)
stream = open(path, "rb")
except Exception as exc:
error(conn, "file_access", "%s: %s" % (type(exc).__name__, exc))
return
try:
stream.seek(offset)
send_json(conn, {"version": VERSION, "ok": True, "size": size})
while True:
chunk = stream.read(CHUNK_SIZE)
if not chunk:
conn.send_bytes(b"")
return
conn.send_bytes(chunk)
except (BrokenPipeError, EOFError, OSError):
return
finally:
stream.close()
def handle(conn):
try:
request = json.loads(conn.recv_bytes(maxlength=CONTROL_LIMIT).decode("utf-8"))
except Exception as exc:
error(conn, "invalid_request", "invalid request: %s" % exc)
return True
if not isinstance(request, dict):
error(conn, "invalid_request_type", "request must be a JSON object")
return True
if request.get("version") != VERSION:
error(conn, "unsupported_version", "unsupported protocol version")
return True
operation = request.get("op")
if operation == "ping":
send_json(conn, {
"version": VERSION,
"ok": True,
"pid": os.getpid(),
"python": sys.version,
})
return True
if operation == "read":
handle_read(conn, request)
return True
if operation == "quit":
send_json(conn, {"version": VERSION, "ok": True})
return False
error(conn, "unknown_operation", "unknown operation: %r" % operation)
return True
def main():
address = PIPE_PREFIX + os.urandom(16).hex()
listener = Listener(address=address, family="AF_PIPE", authkey=AUTHKEY)
sys.stdout.write("DLP_IO " + str(VERSION) + " " + address + "\n")
sys.stdout.flush()
running = True
while running:
try:
conn = listener.accept()
except AuthenticationError:
continue
try:
running = handle(conn)
except Exception as exc:
try:
error(conn, "internal", "%s: %s" % (type(exc).__name__, exc))
except Exception:
pass
finally:
try:
conn.close()
except Exception:
pass
listener.close()
try:
main()
except Exception as exc:
sys.stderr.write("dlp-io helper fatal: %s: %s\n" % (type(exc).__name__, exc))
sys.stderr.flush()
sys.exit(1)
'''
def render_helper_source(authkey: bytes) -> str:
config = HelperSourceConfig.from_authkey(authkey)
replacements = {
"__AUTHKEY_HEX__": config.authkey.hex(),
"__PROTOCOL_VERSION__": str(config.protocol_version),
"__CONTROL_LIMIT__": str(config.control_limit),
"__CHUNK_SIZE__": str(config.chunk_size),
}
source = _HELPER_TEMPLATE
for marker, value in replacements.items():
source = source.replace(marker, value)
return source
# ---------------------------------------------------------------------------
# Pipe-backed raw stream
class _PipeRawIO(io.RawIOBase):
_DISCARD_LIMIT = 64 * 1024 * 1024
def __init__(self, *, conn, size, path, offset, reopen, release) -> None:
super().__init__()
self._conn = conn
self._size = int(size)
self._path = path
self._pos = int(offset)
self._expected_end = max(self._size, self._pos)
self._reopen = reopen
self._release = release
self._buffer = bytearray()
self._eof = False
self._released = False
def readable(self) -> bool:
return True
def seekable(self) -> bool:
return True
def tell(self) -> int:
if self.closed:
raise ValueError("I/O operation on closed file")
return self._pos
def _receive(self) -> None:
try:
chunk = self._conn.recv_bytes(maxlength=CHUNK_SIZE)
except (EOFError, OSError) as exc:
raise DlpTransportError("helper disconnected before EOF") from exc
if chunk:
received_end = self._pos + len(self._buffer) + len(chunk)
if received_end > self._expected_end:
raise DlpTransportError(
"file size changed during transfer: expected end %d, received at least %d"
% (self._expected_end, received_end)
)
self._buffer.extend(chunk)
return
if self._pos != self._expected_end:
raise DlpTransportError(
"file size changed during transfer: expected end %d, received %d"
% (self._expected_end, self._pos)
)
self._eof = True
def readinto(self, target) -> int:
if self.closed:
raise ValueError("I/O operation on closed file")
if len(target) == 0:
return 0
if not self._buffer:
if self._eof:
return 0
self._receive()
if self._eof:
return 0
count = min(len(target), len(self._buffer))
target[:count] = self._buffer[:count]
del self._buffer[:count]
self._pos += count
return count
def seek(self, offset, whence=io.SEEK_SET) -> int:
if self.closed:
raise ValueError("I/O operation on closed file")
if whence == io.SEEK_SET:
target = offset
elif whence == io.SEEK_CUR:
target = self._pos + offset
elif whence == io.SEEK_END:
target = self._size + offset
else:
raise ValueError("invalid whence: %r" % (whence,))
if target < 0:
raise ValueError("negative seek position: %d" % target)
delta = target - self._pos
if delta == 0:
return self._pos
if 0 < delta <= self._DISCARD_LIMIT:
remaining = delta
while remaining:
scratch = bytearray(min(remaining, CHUNK_SIZE))
count = self.readinto(scratch)
if count == 0:
self._pos = target
self._expected_end = max(self._size, target)
break
remaining -= count
return self._pos
self._conn.close()
try:
conn, size = self._reopen(self._path, target)
except Exception:
self.close()
raise
if int(size) != self._size:
conn.close()
self.close()
raise DlpTransportError(
"file size changed during seek: initial %d, reopened %d"
% (self._size, size)
)
self._conn = conn
self._pos = target
self._expected_end = max(self._size, target)
self._buffer.clear()
self._eof = False
return self._pos
def _release_once(self) -> None:
if self._released:
return
self._released = True
self._release(self)
def close(self) -> None:
if not self.closed:
try:
self._conn.close()
finally:
self._release_once()
super().close()
# ---------------------------------------------------------------------------
# Helper session
_SHUTDOWN_TIMEOUT = 5
def _locate_python() -> list[str] | None:
for variable in ("DLP_IO_PYTHON", "DATAPACKER_DLP_PYTHON"):
value = os.environ.get(variable)
if value:
return [value]
executable = shutil.which("python")
if executable:
return [executable]
if shutil.which("py"):
return ["py", "-3"]
return None
def _helper_process_options(os_name=None) -> dict:
if (os_name or os.name) == "nt":
return {"creationflags": subprocess.CREATE_NO_WINDOW}
return {}
def _stop_process(process) -> None:
if process.poll() is not None:
return
process.terminate()
try:
process.wait(timeout=5)
except subprocess.TimeoutExpired:
process.kill()
try:
process.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
def _close_process_pipes(process) -> None:
for name in ("stdin", "stdout", "stderr"):
stream = getattr(process, name, None)
if stream is None:
continue
try:
stream.close()
except OSError:
pass
class DlpSession:
"""One authenticated, short-lived approved-Python helper session."""
def __init__(self, process, pipe_name: str, authkey: bytes) -> None:
self._proc = process
self.pipe_name = pipe_name
self._authkey = authkey
self._request_lock = threading.Lock()
self._stream_lock = threading.Lock()
self._state_lock = threading.RLock()
self._active_raw = None
self._closed = False
@classmethod
def start(cls, *, python_cmd=None, timeout=60):
if os.name != "nt":
raise DlpConfigurationError("dlp-io bridge reads require Windows AF_PIPE")
if timeout <= 0:
raise DlpConfigurationError("startup timeout must be positive")
command = list(python_cmd) if python_cmd else _locate_python()
if not command:
raise DlpHelperStartError(
"no approved Python found; configure DLP_IO_PYTHON"
)
authkey = os.urandom(32)
source = render_helper_source(authkey)
try:
process = subprocess.Popen(
command + ["-u", "-"],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
**_helper_process_options(),
)
except OSError as exc:
raise DlpHelperStartError(
"failed to start approved Python %r: %s" % (command, exc)
) from exc
try:
process.stdin.write(source.encode("utf-8"))
process.stdin.close()
except (BrokenPipeError, OSError) as exc:
_stop_process(process)
raise DlpHelperStartError("failed to inject helper source: %s" % exc) from exc
result = queue.Queue()
def read_handshake() -> None:
try:
result.put(process.stdout.readline())
except Exception as exc:
result.put(exc)
thread = threading.Thread(target=read_handshake, daemon=True)
thread.start()
try:
line = result.get(timeout=timeout)
except queue.Empty as exc:
_stop_process(process)
raise DlpHelperStartError(
"timed out waiting for helper handshake after %s seconds" % timeout
) from exc
if isinstance(line, Exception):
_stop_process(process)
raise DlpHelperStartError("failed to read helper handshake: %s" % line)
decoded = line.decode("utf-8", "replace").strip()
if not decoded.startswith(HANDSHAKE_PREFIX):
_stop_process(process)
stderr = process.stderr.read().decode("utf-8", "replace").strip()
raise DlpHelperStartError(
"invalid helper handshake %r; stderr: %s" % (decoded, stderr)
)
return cls(process, decoded[len(HANDSHAKE_PREFIX):], authkey)
def _ensure_open(self) -> None:
if self._closed:
raise RuntimeError("dlp-io session is closed")
def _validate_header(self, header) -> dict:
if not isinstance(header, dict):
raise DlpProtocolError("helper response must be a JSON object")
if header.get("version") != PROTOCOL_VERSION:
raise DlpProtocolError("helper response has an unsupported version")
if not isinstance(header.get("ok"), bool):
raise DlpProtocolError("helper response is missing boolean ok")
return header
def _request(self, request):
with self._request_lock:
try:
encoded = json.dumps(request).encode("utf-8")
except (TypeError, ValueError) as exc:
raise DlpProtocolError("helper request is not JSON serializable") from exc
try:
conn = Client(self.pipe_name, family="AF_PIPE", authkey=self._authkey)
except AuthenticationError as exc:
raise DlpTransportError("failed to authenticate helper pipe") from exc
except (EOFError, OSError) as exc:
raise DlpTransportError("failed to connect to helper pipe") from exc
try:
try:
conn.send_bytes(encoded)
except (EOFError, OSError) as exc:
raise DlpTransportError("failed to send helper request") from exc
try:
payload = conn.recv_bytes(maxlength=CONTROL_LIMIT)
except (EOFError, OSError) as exc:
raise DlpTransportError("failed to receive helper response") from exc
try:
header = json.loads(payload.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
raise DlpProtocolError("helper response is not valid UTF-8 JSON") from exc
return conn, self._validate_header(header)
except Exception:
conn.close()
raise
def ping(self) -> dict:
with self._state_lock:
self._ensure_open()
conn, header = self._request({"version": PROTOCOL_VERSION, "op": "ping"})
conn.close()
if not header["ok"]:
raise DlpProtocolError(header.get("error", "ping failed"))
return header
def _open_conn(self, path: str, offset: int = 0):
conn, header = self._request(
{
"version": PROTOCOL_VERSION,
"op": "read",
"path": path,
"offset": offset,
}
)
if not header["ok"]:
conn.close()
raise OSError(header.get("error", "helper read failed"))
size = header.get("size")
if isinstance(size, bool) or not isinstance(size, int) or size < 0:
conn.close()
raise DlpProtocolError("helper returned an invalid file size")
return conn, size
def _release_stream(self, raw) -> None:
with self._state_lock:
if self._active_raw is raw:
self._active_raw = None
self._stream_lock.release()
def _reopen_stream(self, path: str, offset: int):
with self._state_lock:
self._ensure_open()
return self._open_conn(path, offset)
def open_raw(self, path):
with self._state_lock:
self._ensure_open()
if not self._stream_lock.acquire(blocking=False):
raise DlpSessionBusyError("dlp-io session already has an active stream")
normalized = os.fsdecode(os.fspath(path))
try:
conn, size = self._open_conn(normalized)
raw = _PipeRawIO(
conn=conn,
size=size,
path=normalized,
offset=0,
reopen=self._reopen_stream,
release=self._release_stream,
)
with self._state_lock:
if self._closed:
raw.close()
raise RuntimeError("dlp-io session is closed")
self._active_raw = raw
return raw
except Exception:
self._stream_lock.release()
raise
def open(
self,
file,
mode="r",
buffering=-1,
encoding=None,
errors=None,
newline=None,
):
buffering = _validate_path_mode(mode, buffering, encoding, errors, newline)
if "r" not in mode or "+" in mode:
raise io.UnsupportedOperation("DlpSession.open supports read-only modes")
return _wrap_reader(
self.open_raw(file), mode, buffering, encoding, errors, newline
)
def open_read(self, path):
"""Compatibility alias returning a buffered binary reader."""
return self.open(path, "rb")
def close(self) -> None:
with self._state_lock:
if self._closed:
return
self._closed = True
active = self._active_raw
if active is not None:
active.close()
def request_quit() -> None:
try:
conn, header = self._request(
{"version": PROTOCOL_VERSION, "op": "quit"}
)
conn.close()
if not header["ok"]:
raise DlpProtocolError(header.get("error", "helper quit failed"))
except Exception:
pass
quit_thread = threading.Thread(target=request_quit, daemon=True)
quit_thread.start()
quit_thread.join(timeout=_SHUTDOWN_TIMEOUT)
try:
if quit_thread.is_alive():
_stop_process(self._proc)
quit_thread.join(timeout=_SHUTDOWN_TIMEOUT)
else:
try:
self._proc.wait(timeout=_SHUTDOWN_TIMEOUT)
except subprocess.TimeoutExpired:
_stop_process(self._proc)
finally:
_close_process_pipes(self._proc)
def __enter__(self):
return self
def __exit__(self, exc_type, exc_value, traceback):
self.close()
return False
# ---------------------------------------------------------------------------
# Public io-compatible API
_ORIGINAL_OPEN = builtins.open
_DEFAULT_BUFFER_SIZE = 1024 * 1024
_state_lock = threading.RLock()
_config = DlpConfig()
_default_session = None
def _validate_path_mode(mode, buffering, encoding, errors, newline) -> int:
buffering = operator.index(buffering)
if buffering < -1:
raise ValueError("invalid buffering size")
if not isinstance(mode, str):
raise TypeError("open() argument 'mode' must be str, not %s" % type(mode).__name__)
allowed = set("axrwb+t")
invalid = set(mode) - allowed
duplicate = any(mode.count(character) > 1 for character in allowed)
if not mode or invalid or duplicate or sum(c in mode for c in "axrw") != 1:
raise ValueError("invalid mode: %r" % mode)
if "b" in mode and "t" in mode:
raise ValueError("can't have text and binary mode at once")
if "b" in mode:
if encoding is not None:
raise ValueError("binary mode doesn't take an encoding argument")
if errors is not None:
raise ValueError("binary mode doesn't take an errors argument")
if newline is not None:
raise ValueError("binary mode doesn't take a newline argument")
elif buffering == 0:
raise ValueError("can't have unbuffered text I/O")
return buffering
def configure(*, python_executable=None, startup_timeout=60) -> None:
"""Configure the lazy default session before its first bridged read."""
global _config
with _state_lock:
if _default_session is not None:
raise DlpConfigurationError(
"cannot configure dlp_io while the default session is active"
)
if startup_timeout <= 0:
raise DlpConfigurationError("startup_timeout must be positive")
_config = DlpConfig(
python_executable=python_executable,
startup_timeout=startup_timeout,
)
def _get_default_session():
global _default_session
with _state_lock:
if _default_session is None:
_default_session = DlpSession.start(
python_cmd=_config.python_command(),
timeout=_config.startup_timeout,
)
return _default_session
def _wrap_reader(raw, mode, buffering, encoding, errors, newline):
binary = "b" in mode
if binary and buffering == 0:
return raw
if not binary and buffering == 0:
raw.close()
raise ValueError("can't have unbuffered text I/O")
buffer_size = _DEFAULT_BUFFER_SIZE if buffering in (-1, 1) else buffering
if buffer_size <= 0:
raw.close()
raise ValueError("invalid buffering size")
buffered = io.BufferedReader(raw, buffer_size=buffer_size)
if binary:
return buffered
return io.TextIOWrapper(
buffered,
encoding=encoding,
errors=errors,
newline=newline,
line_buffering=buffering == 1,
)
def open(
file,
mode="r",
buffering=-1,
encoding=None,
errors=None,
newline=None,
closefd=True,
opener=None,
):
"""Open a file with io.open-compatible calling conventions."""
is_path = isinstance(file, (str, bytes, os.PathLike))
if not is_path or opener is not None:
return _ORIGINAL_OPEN(
file,
mode,
buffering,
encoding,
errors,
newline,
closefd,
opener,
)
buffering = _validate_path_mode(mode, buffering, encoding, errors, newline)
if "+" in mode:
raise io.UnsupportedOperation(
"dlp_io does not support update mode for path-based files"
)
pure_read = "r" in mode
if pure_read:
if not closefd:
raise ValueError("Cannot use closefd=False with file name")
path = os.fsdecode(os.fspath(file))
raw = _get_default_session().open_raw(path)
return _wrap_reader(raw, mode, buffering, encoding, errors, newline)
return _open_write(file, mode, buffering, encoding, errors, newline, closefd)
def shutdown() -> None:
"""Close and forget the lazy default session. This operation is idempotent."""
global _default_session
with _state_lock:
session = _default_session
_default_session = None
if session is not None:
session.close()
atexit.register(shutdown)
# ---------------------------------------------------------------------------
# Global open() patch
_ORIGINAL_IO_OPEN = io.open
_patch_lock = threading.RLock()
_installed = False
def _reader_wrapper(open_read_func):
def patched_open(
file,
mode="r",
buffering=-1,
encoding=None,
errors=None,
newline=None,
closefd=True,
opener=None,
):
is_path = isinstance(file, (str, bytes, os.PathLike))
if not is_path or opener is not None:
return _ORIGINAL_OPEN(
file,
mode,
buffering,
encoding,
errors,
newline,
closefd,
opener,
)
buffering = _validate_path_mode(mode, buffering, encoding, errors, newline)
if "+" in mode:
raise io.UnsupportedOperation(
"dlp_io does not support update mode for path-based files"
)
pure_read = "r" in mode
if pure_read:
if not closefd:
raise ValueError("Cannot use closefd=False with file name")
reader = open_read_func(os.fsdecode(os.fspath(file)))
if isinstance(reader, io.RawIOBase):
return _wrap_reader(reader, mode, buffering, encoding, errors, newline)
if "b" in mode:
return reader
return io.TextIOWrapper(
reader,
encoding=encoding,
errors=errors,
newline=newline,
line_buffering=buffering == 1,
)
return _open_write(file, mode, buffering, encoding, errors, newline, closefd)
return patched_open
def install_open_patch() -> None:
"""Route global path-based reads through the lazy default dlp-io session."""
global _installed
with _patch_lock:
if _installed:
return
builtins.open = open
io.open = open
_installed = True
def install_reader_patch(open_read_func) -> None:
"""Compatibility hook for an explicitly managed session.open_raw callback."""
global _installed
with _patch_lock:
if _installed:
return
wrapper = _reader_wrapper(open_read_func)
builtins.open = wrapper
io.open = wrapper
_installed = True
def uninstall_open_patch() -> None:
global _installed
with _patch_lock:
if not _installed:
return
builtins.open = _ORIGINAL_OPEN
io.open = _ORIGINAL_IO_OPEN
_installed = False
# ---------------------------------------------------------------------------
# Encryption detection and unencrypted writes
#
# DLP transparent encryption stores ciphertext on disk. Approved processes
# (the whitelisted Python) see plaintext through the filter driver, while
# non-whitelisted tools such as certutil see the raw on-disk bytes. File
# attributes (attrib, ADS, size) are identical between encrypted and plain
# files, so content comparison is the only reliable discriminator.
_DETECT_HEAD_SIZE = 64 * 1024
_RAW_HASH_PATTERN = re.compile(rb"^[ \t]*([0-9a-fA-F]{64})[ \t\r]*$", re.MULTILINE)
def _raw_sha256(path) -> bytes:
"""SHA-256 of the raw on-disk bytes, read by non-whitelisted certutil.
certutil must be spawned through ``cmd /c``: DLP hosts block approved
processes from starting certutil directly, and list-form arguments keep
non-ASCII paths intact through the cmd command line.
"""
try:
completed = subprocess.run(
["cmd", "/c", "certutil", "-hashfile", path, "SHA256"],
capture_output=True,
**_helper_process_options(),
)
except (OSError, subprocess.SubprocessError) as exc:
raise DlpIoError("raw certutil hash unavailable: %s" % exc) from exc
if completed.returncode != 0:
raise DlpIoError(
"certutil -hashfile failed with exit code %d" % completed.returncode
)
match = _RAW_HASH_PATTERN.search(completed.stdout)
if match is None:
raise DlpIoError("certutil output did not contain a SHA-256 digest")
return bytes.fromhex(match.group(1).decode("ascii"))
def _native_sha256(path) -> bytes:
digest = hashlib.sha256()
with _ORIGINAL_OPEN(path, "rb") as stream:
for chunk in iter(lambda: stream.read(CHUNK_SIZE), b""):
digest.update(chunk)
return digest.digest()
def _native_head(path, size) -> bytes:
with _ORIGINAL_OPEN(path, "rb") as stream:
return stream.read(size)
def _helper_head(path, size) -> bytes:
raw = _get_default_session().open_raw(path)
try:
return raw.read(size)
finally:
raw.close()
def is_encrypted(path) -> bool:
"""Return True when the file on disk holds DLP ciphertext.
The raw on-disk digest (certutil) is compared with what the current
process sees. A mismatch means the file is encrypted and the current
process is approved. When both match the result is ambiguous — the file
is either plain, or encrypted while the current process is not approved
(packaged .exe) — so the first bytes seen by the approved helper are
compared with the first bytes seen by the current process.
"""
if os.name != "nt":
return False
normalized = os.fsdecode(os.fspath(path))
native_hash = _native_sha256(normalized)
raw_hash = _raw_sha256(normalized)
if raw_hash != native_hash:
return True
return _helper_head(normalized, _DETECT_HEAD_SIZE) != _native_head(
normalized, _DETECT_HEAD_SIZE
)
class WriteChannel(str, Enum):
"""How path-based writes are routed to keep files unencrypted."""
DIRECT = "direct"
POWERSHELL = "powershell"
_WRITE_CHANNEL_ENV = "DLP_IO_WRITE_CHANNEL"
_WRITE_PROBE_PAYLOAD = b"dlp-io-write-probe\x00\x1a\xff" + bytes(range(256)) * 2
_write_channel_lock = threading.RLock()
_write_channel_cache = {}
def _powershell_script(path, file_mode) -> str:
escaped = path.replace("'", "''")
return (
"$ErrorActionPreference='Stop';"
"$input_stream=[Console]::OpenStandardInput();"
"$file_stream=[IO.File]::Open('" + escaped + "',[IO.FileMode]::" + file_mode + ");"
"$input_stream.CopyTo($file_stream);"
"$file_stream.Close()"
)
class _PowerShellRelayRaw(io.RawIOBase):
"""Binary stream written to disk by a non-whitelisted PowerShell child.
Bytes travel through the stdin pipe (never through a DLP-encrypted temp
file), and the child's file writes are not encrypted because the child
is outside the DLP write-encryption policy.
"""
def __init__(self, path, file_mode) -> None:
super().__init__()
try:
self._proc = subprocess.Popen(
["powershell", "-NoProfile", "-Command", _powershell_script(path, file_mode)],
stdin=subprocess.PIPE,
stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE,
**_helper_process_options(),
)
except OSError as exc:
raise DlpWriteError("failed to start powershell write relay: %s" % exc) from exc
def writable(self) -> bool:
return True
def write(self, data):
if self.closed:
raise ValueError("I/O operation on closed file")
try:
self._proc.stdin.write(data)
except (BrokenPipeError, OSError) as exc:
raise DlpWriteError("powershell write relay rejected data: %s" % exc) from exc
return len(data)
def flush(self) -> None:
if self.closed:
raise ValueError("I/O operation on closed file")
stdin = self._proc.stdin
if stdin is None or stdin.closed:
return
try:
stdin.flush()
except (BrokenPipeError, OSError) as exc:
raise DlpWriteError("powershell write relay rejected data: %s" % exc) from exc
def close(self) -> None:
if self.closed:
return
try:
try:
self._proc.stdin.close()
except (BrokenPipeError, OSError):
pass
returncode = self._proc.wait()
if returncode != 0:
detail = self._proc.stderr.read().decode("utf-8", "replace").strip()
raise DlpWriteError(
"powershell write relay failed with exit code %d: %s"
% (returncode, detail)
)
finally:
super().close()
def _write_probe_path(suffix) -> str:
descriptor, probe = tempfile.mkstemp(prefix="dlp-io-probe-", suffix=suffix or None)
os.close(descriptor)
return probe
def _probe_direct_write(suffix) -> bool:
probe = _write_probe_path(suffix)
try:
with _ORIGINAL_OPEN(probe, "wb") as stream:
stream.write(_WRITE_PROBE_PAYLOAD)
try:
raw_hash = _raw_sha256(probe)
except DlpIoError:
# Without a raw channel there is no DLP filter to bypass.
return True
return raw_hash == hashlib.sha256(_WRITE_PROBE_PAYLOAD).digest()
finally:
os.unlink(probe)
def _probe_powershell_write(suffix) -> bool:
probe = _write_probe_path(suffix)
try:
try:
with _PowerShellRelayRaw(probe, "Create") as stream:
stream.write(_WRITE_PROBE_PAYLOAD)
return _raw_sha256(probe) == hashlib.sha256(_WRITE_PROBE_PAYLOAD).digest()
except (DlpIoError, OSError):
return False
finally:
os.unlink(probe)
def _select_write_channel(suffix) -> WriteChannel:
"""Pick a write channel that lands unencrypted bytes on disk.
Results are cached per file suffix because DLP write policies are
typically scoped by document type. ``DLP_IO_WRITE_CHANNEL`` forces a
channel and skips probing.
"""
override = os.environ.get(_WRITE_CHANNEL_ENV)
if override:
try:
return WriteChannel(override.strip().lower())
except ValueError as exc:
valid = ", ".join(channel.value for channel in WriteChannel)
raise DlpConfigurationError(
"%s must be one of: %s" % (_WRITE_CHANNEL_ENV, valid)
) from exc
key = suffix.lower()
with _write_channel_lock:
cached = _write_channel_cache.get(key)
if cached is not None:
return cached
if _probe_direct_write(suffix):
channel = WriteChannel.DIRECT
elif _probe_powershell_write(suffix):
channel = WriteChannel.POWERSHELL
else:
raise DlpWriteError(
"no write channel can keep %r files unencrypted on this host"
% (suffix or "<none>")
)
_write_channel_cache[key] = channel
return channel
def _wrap_writer(raw, mode, buffering, encoding, errors, newline):
binary = "b" in mode
if buffering == 0:
if not binary:
raw.close()
raise ValueError("can't have unbuffered text I/O")
return raw
buffer_size = _DEFAULT_BUFFER_SIZE if buffering in (-1, 1) else buffering
if buffer_size <= 0:
raw.close()
raise ValueError("invalid buffering size")
buffered = io.BufferedWriter(raw, buffer_size=buffer_size)
if binary:
return buffered
return io.TextIOWrapper(
buffered,
encoding=encoding,
errors=errors,
newline=newline,
line_buffering=buffering == 1,
)
def _open_write(file, mode, buffering, encoding, errors, newline, closefd):
"""Open a path for writing so the file stays unencrypted on disk."""
path = os.fsdecode(os.fspath(file))
if "x" in mode and os.path.exists(path):
raise FileExistsError(errno.EEXIST, "File exists", path)
channel = _select_write_channel(os.path.splitext(path)[1])
if channel is WriteChannel.DIRECT:
return _ORIGINAL_OPEN(file, mode, buffering, encoding, errors, newline, closefd)
if not closefd:
raise ValueError("Cannot use closefd=False with file name")
if "w" in mode:
file_mode = "Create"
elif "a" in mode:
file_mode = "Append"
else:
file_mode = "CreateNew"
raw = _PowerShellRelayRaw(path, file_mode)
return _wrap_writer(raw, mode, buffering, encoding, errors, newline)