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}"