263 lines
7.8 KiB
Python
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}"
|