Files
2026-10-06 21:08:12 +03:00

1066 lines
25 KiB
Python

from __future__ import annotations
import hashlib
import json
import re
import shutil
import tempfile
import zipfile
from dataclasses import dataclass
from pathlib import Path
from .state import initialize_state
JOB_ID_RE = re.compile(r"^JOB-[0-9A-F]{12}$")
class JobPackageError(RuntimeError):
pass
@dataclass(frozen=True)
class ImportedJob:
job_id: str
job_path: Path
package_path: Path
idempotent: bool
@dataclass(frozen=True)
class _ValidatedPackage:
job: dict
hash_content: bytes
job_bytes: bytes
hash_path: str
class JobPackageService:
def __init__(self, paths) -> None:
self._paths = paths
def import_package(
self,
package_path: str | Path,
) -> ImportedJob:
package_path = Path(package_path).resolve()
if not package_path.is_file():
raise JobPackageError(
f"Job package not found: {package_path}"
)
validated = self._read_and_validate(
package_path
)
job_id = validated.job["job_id"]
job_path = (
self._paths.jobs
/ package_path.stem
)
if job_path.exists():
idempotent = self._validate_existing_job(
job_path,
validated,
)
pending_path = (
self._paths.pending
/ package_path.name
)
self._move_package(
package_path,
pending_path,
)
return ImportedJob(
job_id=job_id,
job_path=job_path,
package_path=pending_path,
idempotent=idempotent,
)
self._install_new_job(
job_path,
validated,
)
pending_path = (
self._paths.pending
/ package_path.name
)
self._move_package(
package_path,
pending_path,
)
return ImportedJob(
job_id=job_id,
job_path=job_path,
package_path=pending_path,
idempotent=False,
)
def _read_and_validate(
self,
package_path: Path,
) -> _ValidatedPackage:
try:
with zipfile.ZipFile(
package_path,
"r",
) as archive:
self._validate_zip_members(
archive
)
names = set(
archive.namelist()
)
job_bytes = archive.read(
"job.json"
)
hash_name = next(
name
for name in names
if name.startswith("hashes/")
)
hash_bytes = archive.read(
hash_name
)
manifest_bytes = archive.read(
"manifest.json"
)
except (
zipfile.BadZipFile,
KeyError,
StopIteration,
OSError,
) as exc:
raise JobPackageError(
f"Invalid job package: {package_path}"
) from exc
job = self._load_json_object(
job_bytes,
"job.json",
)
manifest = self._load_json_object(
manifest_bytes,
"manifest.json",
)
self._validate_job(
job,
hash_name,
hash_bytes,
)
self._validate_manifest(
manifest,
job,
package_path.name,
job_bytes,
hash_bytes,
hash_name,
)
return _ValidatedPackage(
job=job,
hash_content=hash_bytes,
job_bytes=job_bytes,
hash_path=hash_name,
)
@staticmethod
def _validate_zip_members(
archive: zipfile.ZipFile,
) -> None:
infos = archive.infolist()
names = [info.filename for info in infos]
if len(names) != len(set(names)):
raise JobPackageError(
"Job package contains duplicate ZIP members"
)
expected_fixed = {
"job.json",
"manifest.json",
}
hash_members = [
name
for name in names
if name.startswith("hashes/")
]
if len(names) != 3:
raise JobPackageError(
"Job package must contain exactly three files"
)
if not expected_fixed.issubset(
set(names)
):
raise JobPackageError(
"Job package is missing job.json "
"or manifest.json"
)
if len(hash_members) != 1:
raise JobPackageError(
"Job package must contain exactly one hash file"
)
expected = expected_fixed | set(
hash_members
)
if set(names) != expected:
raise JobPackageError(
"Job package contains unexpected files"
)
for info in infos:
name = info.filename
if not name:
raise JobPackageError(
"Job package contains an empty member name"
)
if name.startswith("/"):
raise JobPackageError(
f"Absolute ZIP member path: {name}"
)
parts = Path(name).parts
if ".." in parts:
raise JobPackageError(
f"Path traversal in ZIP member: {name}"
)
if "\\" in name:
raise JobPackageError(
f"Invalid ZIP member path: {name}"
)
if name.endswith("/"):
raise JobPackageError(
f"Directories are not allowed: {name}"
)
if info.is_dir():
raise JobPackageError(
f"Directory member is not allowed: {name}"
)
unix_mode = (
info.external_attr >> 16
) & 0xFFFF
if unix_mode:
file_type = unix_mode & 0o170000
if file_type == 0o120000:
raise JobPackageError(
f"Symlink member is not allowed: {name}"
)
@staticmethod
def _load_json_object(
content: bytes,
name: str,
) -> dict:
try:
value = json.loads(
content.decode("utf-8")
)
except (
UnicodeDecodeError,
json.JSONDecodeError,
) as exc:
raise JobPackageError(
f"Invalid JSON: {name}"
) from exc
if not isinstance(value, dict):
raise JobPackageError(
f"{name} must contain a JSON object"
)
return value
@staticmethod
def _require_string(
data: dict,
name: str,
) -> str:
value = data.get(name)
if not isinstance(value, str):
raise JobPackageError(
f"{name} must be a string"
)
if not value:
raise JobPackageError(
f"{name} must not be empty"
)
return value
@staticmethod
def _require_int(
data: dict,
name: str,
) -> int:
value = data.get(name)
if (
isinstance(value, bool)
or not isinstance(value, int)
):
raise JobPackageError(
f"{name} must be an integer"
)
return value
def _validate_job(
self,
job: dict,
hash_path: str,
hash_bytes: bytes,
) -> None:
required = (
"job_id",
"created_at",
"method_id",
"method_version",
"hash_mode",
"hash_file_name",
"hash_file_sha256",
"hash_count",
"handshake_count",
"step_count",
"handshakes",
"steps",
)
for name in required:
if name not in job:
raise JobPackageError(
f"job.json is missing field: {name}"
)
job_id = self._require_string(
job,
"job_id",
)
if not JOB_ID_RE.fullmatch(job_id):
raise JobPackageError(
f"Invalid Job ID: {job_id}"
)
self._require_string(
job,
"created_at",
)
self._require_string(
job,
"method_id",
)
method_version = self._require_int(
job,
"method_version",
)
hash_mode = self._require_int(
job,
"hash_mode",
)
hash_count = self._require_int(
job,
"hash_count",
)
handshake_count = self._require_int(
job,
"handshake_count",
)
step_count = self._require_int(
job,
"step_count",
)
if method_version < 1:
raise JobPackageError(
"method_version must be positive"
)
if hash_mode < 0:
raise JobPackageError(
"hash_mode must not be negative"
)
if hash_count <= 0:
raise JobPackageError(
"hash_count must be greater than zero"
)
if handshake_count <= 0:
raise JobPackageError(
"handshake_count must be greater than zero"
)
if step_count <= 0:
raise JobPackageError(
"step_count must be greater than zero"
)
hash_file_name = self._require_string(
job,
"hash_file_name",
)
expected_hash_name = (
f"{job_id}.22000"
)
if hash_file_name != expected_hash_name:
raise JobPackageError(
"hash_file_name does not match "
"the Job ID"
)
expected_zip_path = (
f"hashes/{hash_file_name}"
)
if hash_path != expected_zip_path:
raise JobPackageError(
"Hash file path does not match "
"job.json"
)
hash_sha256 = self._require_string(
job,
"hash_file_sha256",
)
if not re.fullmatch(
r"[0-9a-fA-F]{64}",
hash_sha256,
):
raise JobPackageError(
"hash_file_sha256 must be a "
"64-character hexadecimal SHA-256"
)
calculated_hash_sha256 = (
hashlib.sha256(
hash_bytes
).hexdigest()
)
if (
calculated_hash_sha256
!= hash_sha256.lower()
):
raise JobPackageError(
"Hash file SHA-256 mismatch"
)
lines = hash_bytes.decode(
"utf-8"
).splitlines()
if len(lines) != hash_count:
raise JobPackageError(
"hash_count does not match "
"the hash file"
)
if any(
not line.strip()
for line in lines
):
raise JobPackageError(
"Hash file contains an empty line"
)
if len(set(lines)) != hash_count:
raise JobPackageError(
"Hash file contains duplicate hashes"
)
handshakes = job["handshakes"]
if not isinstance(
handshakes,
list,
):
raise JobPackageError(
"handshakes must be a list"
)
if len(handshakes) != handshake_count:
raise JobPackageError(
"handshake_count does not match "
"handshakes"
)
handshake_ids: set[int] = set()
for item in handshakes:
if not isinstance(
item,
dict,
):
raise JobPackageError(
"Each handshake must be an object"
)
handshake_id = item.get(
"handshake_id"
)
if (
isinstance(handshake_id, bool)
or not isinstance(
handshake_id,
int,
)
):
raise JobPackageError(
"handshake_id must be an integer"
)
if handshake_id <= 0:
raise JobPackageError(
"handshake_id must be positive"
)
if handshake_id in handshake_ids:
raise JobPackageError(
"Duplicate handshake_id"
)
handshake_ids.add(
handshake_id
)
hash22000 = item.get(
"hash22000"
)
if (
not isinstance(
hash22000,
str,
)
or not hash22000
):
raise JobPackageError(
"handshake hash22000 must be "
"a non-empty string"
)
if item.get(
"access_point_id"
) is not None:
access_point_id = item[
"access_point_id"
]
if (
isinstance(
access_point_id,
bool,
)
or not isinstance(
access_point_id,
int,
)
or access_point_id <= 0
):
raise JobPackageError(
"access_point_id must be "
"a positive integer or null"
)
steps = job["steps"]
if not isinstance(
steps,
list,
):
raise JobPackageError(
"steps must be a list"
)
if len(steps) != step_count:
raise JobPackageError(
"step_count does not match steps"
)
step_numbers: set[int] = set()
step_ids: set[str] = set()
for item in steps:
if not isinstance(
item,
dict,
):
raise JobPackageError(
"Each step must be an object"
)
step_no = item.get(
"step_no"
)
if (
isinstance(step_no, bool)
or not isinstance(
step_no,
int,
)
):
raise JobPackageError(
"step_no must be an integer"
)
if step_no <= 0:
raise JobPackageError(
"step_no must be positive"
)
if step_no in step_numbers:
raise JobPackageError(
"Duplicate step_no"
)
step_numbers.add(
step_no
)
step_id = item.get(
"step_id"
)
if (
not isinstance(
step_id,
str,
)
or not step_id
):
raise JobPackageError(
"step_id must be a "
"non-empty string"
)
if step_id in step_ids:
raise JobPackageError(
"Duplicate step_id"
)
step_ids.add(
step_id
)
session_name = item.get(
"session_name"
)
if (
not isinstance(
session_name,
str,
)
or not session_name
):
raise JobPackageError(
"session_name must be a "
"non-empty string"
)
if "definition" not in item:
raise JobPackageError(
"Step is missing definition"
)
if not isinstance(
item["definition"],
dict,
):
raise JobPackageError(
"Step definition must be "
"an object"
)
def _validate_manifest(
self,
manifest: dict,
job: dict,
zip_filename: str,
job_bytes: bytes,
hash_bytes: bytes,
hash_path: str,
) -> None:
if manifest.get(
"format_version"
) != 1:
raise JobPackageError(
"Unsupported manifest format_version"
)
if manifest.get(
"job_id"
) != job["job_id"]:
raise JobPackageError(
"Manifest Job ID mismatch"
)
if manifest.get(
"created_at"
) != job["created_at"]:
raise JobPackageError(
"Manifest created_at mismatch"
)
if manifest.get(
"zip_filename"
) != zip_filename:
raise JobPackageError(
"Manifest zip_filename mismatch"
)
files = manifest.get(
"files"
)
if not isinstance(
files,
list,
):
raise JobPackageError(
"Manifest files must be a list"
)
if len(files) != 2:
raise JobPackageError(
"Manifest must contain exactly two files"
)
expected = {
"job.json": job_bytes,
hash_path: hash_bytes,
}
seen: set[str] = set()
for item in files:
if not isinstance(
item,
dict,
):
raise JobPackageError(
"Manifest file entry must be an object"
)
path = item.get(
"path"
)
if (
not isinstance(
path,
str,
)
or not path
):
raise JobPackageError(
"Manifest file path must be a "
"non-empty string"
)
if path in seen:
raise JobPackageError(
"Duplicate manifest file path"
)
seen.add(path)
if path not in expected:
raise JobPackageError(
f"Unexpected manifest file: {path}"
)
size = item.get(
"size"
)
if (
isinstance(size, bool)
or not isinstance(size, int)
or size < 0
):
raise JobPackageError(
f"Invalid manifest size: {path}"
)
sha256 = item.get(
"sha256"
)
if (
not isinstance(
sha256,
str,
)
or not re.fullmatch(
r"[0-9a-fA-F]{64}",
sha256,
)
):
raise JobPackageError(
f"Invalid manifest SHA-256: {path}"
)
content = expected[path]
if size != len(content):
raise JobPackageError(
f"Manifest size mismatch: {path}"
)
calculated = hashlib.sha256(
content
).hexdigest()
if calculated != sha256.lower():
raise JobPackageError(
f"Manifest SHA-256 mismatch: {path}"
)
if seen != set(expected):
raise JobPackageError(
"Manifest does not describe "
"exactly the package files"
)
def _validate_existing_job(
self,
job_path: Path,
package: _ValidatedPackage,
) -> bool:
if not job_path.is_dir():
raise JobPackageError(
f"Existing Job path is not a directory: "
f"{job_path}"
)
job_json_path = (
job_path / "job.json"
)
hash_path = (
job_path
/ "hashes"
/ package.job["hash_file_name"]
)
if not job_json_path.is_file():
raise JobPackageError(
f"Existing Job is incomplete: {job_path}"
)
if not hash_path.is_file():
raise JobPackageError(
f"Existing Job hash file is missing: "
f"{job_path}"
)
try:
existing_job_bytes = (
job_json_path.read_bytes()
)
existing_hash_bytes = (
hash_path.read_bytes()
)
except OSError as exc:
raise JobPackageError(
f"Failed to read existing Job: "
f"{job_path}"
) from exc
if existing_job_bytes != package.job_bytes:
raise JobPackageError(
f"Existing Job snapshot conflicts: "
f"{package.job['job_id']}"
)
if existing_hash_bytes != package.hash_content:
raise JobPackageError(
f"Existing Job hash snapshot conflicts: "
f"{package.job['job_id']}"
)
return True
def _install_new_job(
self,
job_path: Path,
package: _ValidatedPackage,
) -> None:
self._paths.jobs.mkdir(
parents=True,
exist_ok=True,
)
temporary_path = Path(
tempfile.mkdtemp(
prefix=f".{package.job['job_id']}.",
dir=self._paths.jobs,
)
)
try:
(temporary_path / "hashes").mkdir()
(temporary_path / "state").mkdir()
(temporary_path / "logs").mkdir()
(temporary_path / "results").mkdir()
(temporary_path / "job.json").write_bytes(
package.job_bytes
)
(
temporary_path
/ "hashes"
/ package.job["hash_file_name"]
).write_bytes(
package.hash_content
)
initialize_state(
temporary_path,
package.job,
)
temporary_path.replace(
job_path
)
except Exception:
shutil.rmtree(
temporary_path,
ignore_errors=True,
)
raise
def archive_job(
self,
job_name: str,
) -> Path:
package_name = f"{job_name}.zip"
pending_path = (
self._paths.pending
/ package_name
)
archive_path = (
self._paths.archive
/ package_name
)
if not pending_path.exists():
if archive_path.is_file():
return archive_path
raise JobPackageError(
f"Pending Job package not found: "
f"{pending_path}"
)
if not pending_path.is_file():
raise JobPackageError(
f"Pending Job package is not a file: "
f"{pending_path}"
)
self._move_package(
pending_path,
archive_path,
)
return archive_path
def _move_package(
self,
package_path: Path,
destination: Path,
) -> None:
destination.parent.mkdir(
parents=True,
exist_ok=True,
)
if destination.exists():
if not destination.is_file():
raise JobPackageError(
f"Package destination is not a file: "
f"{destination}"
)
if self._sha256_file(
destination
) != self._sha256_file(
package_path
):
raise JobPackageError(
f"Package destination already exists "
f"with different content: "
f"{destination.name}"
)
package_path.unlink()
return
package_path.replace(
destination
)
@staticmethod
def _sha256_file(
path: Path,
) -> str:
digest = hashlib.sha256()
with path.open(
"rb"
) as fp:
while True:
chunk = fp.read(
1024 * 1024
)
if not chunk:
break
digest.update(chunk)
return digest.hexdigest()