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

263 lines
7.8 KiB
Python

from __future__ import annotations
import hashlib
import json
import sqlite3
import uuid
from dataclasses import dataclass
from typing import Any
from cracklab.server.repositories.audit import AuditRepository
from cracklab.server.repositories.handshakes import HandshakesRepository
from cracklab.server.repositories.jobs import JobsRepository
from cracklab.server.repositories.methods import MethodsRepository
from cracklab.server.repositories.reservations import ReservationsRepository
from cracklab.server.services.method_definition import (
MethodDefinitionError,
validate_method_definition,
)
class JobIssuanceError(ValueError):
pass
@dataclass(frozen=True)
class IssuedJob:
job_id: str
method_id: str
method_version: int
status: str
created_at: str
hash_file_name: str
hash_file_sha256: str
hash_count: int
handshake_count: int
step_count: int
hash_content: str
handshake_ids: tuple[int, ...]
steps: tuple[dict[str, Any], ...]
class JobIssuanceService:
def __init__(self, conn: sqlite3.Connection) -> None:
self._conn = conn
self._methods = MethodsRepository(conn)
self._handshakes = HandshakesRepository(conn)
self._jobs = JobsRepository(conn)
self._reservations = ReservationsRepository(conn)
self._audit = AuditRepository(conn)
def issue_job(
self,
method_id: str,
method_version: int,
limit: int,
client_id: str | None = None,
) -> IssuedJob | None:
if limit <= 0:
raise ValueError("limit must be greater than zero")
method_version_row = self._methods.get_method_version(
method_id,
method_version,
)
if method_version_row is None:
raise JobIssuanceError(
f"Method version not found: {method_id} v{method_version}"
)
definition = self._decode_method_definition(
method_version_row["definition_json"],
)
steps = self._build_steps(
method_id,
method_version,
definition,
)
self._conn.execute("BEGIN IMMEDIATE")
try:
handshake_rows = self._handshakes.list_eligible(
method_id,
method_version,
limit,
)
if not handshake_rows:
self._conn.rollback()
return None
job_id = f"JOB-{uuid.uuid4().hex[:12].upper()}"
hash_file_name = f"{job_id}.22000"
handshake_ids = tuple(
int(row["handshake_id"])
for row in handshake_rows
)
unique_hashes: list[str] = []
seen_hashes: set[str] = set()
for row in handshake_rows:
hash22000 = str(row["hash22000"]).strip()
if hash22000 not in seen_hashes:
seen_hashes.add(hash22000)
unique_hashes.append(hash22000)
hash_content = "\n".join(unique_hashes) + "\n"
hash_file_sha256 = hashlib.sha256(
hash_content.encode("utf-8"),
).hexdigest()
created_at = self._jobs.create_job(
job_id=job_id,
method_id=method_id,
method_version=method_version,
status="QUEUED",
hash_file_name=hash_file_name,
hash_file_sha256=hash_file_sha256,
hash_count=len(unique_hashes),
step_count=len(steps),
client_id=client_id,
)
for row in handshake_rows:
self._jobs.add_job_handshake(
job_id=job_id,
handshake_id=int(row["handshake_id"]),
hash22000=str(row["hash22000"]).strip(),
access_point_id=(
int(row["access_point_id"])
if row["access_point_id"] is not None
else None
),
)
self._reservations.create_reservation(
reservation_id=uuid.uuid4().hex,
handshake_id=int(row["handshake_id"]),
method_id=method_id,
method_version=method_version,
job_id=job_id,
client_id=client_id,
)
for step_no, step in enumerate(steps, start=1):
self._jobs.create_job_step(
job_id=job_id,
step_no=step_no,
step_id=step["step_id"],
status="PENDING",
session_name=self._session_name(
job_id,
step["step_id"],
),
definition=step,
)
self._audit.create_event(
event_type="JOB_ISSUED",
job_id=job_id,
client_id=client_id,
method_id=method_id,
method_version=method_version,
details={
"handshake_count": len(handshake_rows),
"hash_count": len(unique_hashes),
"step_count": len(steps),
"hash_file_name": hash_file_name,
},
)
self._conn.commit()
except Exception:
self._conn.rollback()
raise
return IssuedJob(
job_id=job_id,
method_id=method_id,
method_version=method_version,
status="QUEUED",
created_at=created_at,
hash_file_name=hash_file_name,
hash_file_sha256=hash_file_sha256,
hash_count=len(unique_hashes),
handshake_count=len(handshake_rows),
step_count=len(steps),
hash_content=hash_content,
handshake_ids=handshake_ids,
steps=tuple(steps),
)
@staticmethod
def _decode_method_definition(
definition_json: str,
) -> dict[str, Any]:
try:
definition = json.loads(definition_json)
except json.JSONDecodeError as exc:
raise JobIssuanceError(
"Method version contains invalid definition JSON"
) from exc
try:
return validate_method_definition(definition)
except MethodDefinitionError as exc:
raise JobIssuanceError(str(exc)) from exc
@staticmethod
def _build_steps(
method_id: str,
method_version: int,
definition: dict[str, Any],
) -> list[dict[str, Any]]:
del method_id
del method_version
source_steps = definition["steps"]
result: list[dict[str, Any]] = []
seen_step_ids: set[str] = set()
for step in source_steps:
if not isinstance(step, dict):
raise JobIssuanceError(
"Each method step must be a JSON object"
)
step_id = step.get("step_id")
if not isinstance(step_id, str) or not step_id:
raise JobIssuanceError(
"Each method step must have a non-empty step_id"
)
if step_id in seen_step_ids:
raise JobIssuanceError(
f"Duplicate step_id: {step_id}"
)
seen_step_ids.add(step_id)
attack_mode = step.get("attack_mode")
if not isinstance(attack_mode, int) or attack_mode < 0:
raise JobIssuanceError(
f"Invalid attack_mode for step: {step_id}"
)
result.append(dict(step))
return result
@staticmethod
def _session_name(
job_id: str,
step_id: str,
) -> str:
return f"cl_{job_id[:12]}_{step_id}"