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

794 lines
19 KiB
Python

from __future__ import annotations
import json
import logging
import traceback
from pathlib import Path
from typing import Any
from .config import ClientConfig, load_config
from .hashcat import (
HashcatError,
run_show,
run_step,
)
from .job_package import (
ImportedJob,
JobPackageError,
JobPackageService,
)
from .client_logging import (
configure_job_logging,
configure_logging,
)
from .paths import ClientPaths
from .result_package import (
ResultPackageError,
create_result_package,
)
from .state import (
JOB_STATUS_COMPLETED,
JOB_STATUS_FAILED,
JOB_STATUS_PARTIAL,
JOB_STATUS_RUNNING,
STEP_STATUS_COMPLETED,
STEP_STATUS_FAILED,
STEP_STATUS_PARTIAL,
STEP_STATUS_RUNNING,
StateError,
get_step_state,
initialize_state,
load_state,
update_job_state,
update_step_state,
)
class ClientError(RuntimeError):
pass
class CrackLabClient:
def __init__(
self,
script_dir: Path,
) -> None:
self.script_dir = script_dir.resolve()
self.config: ClientConfig = load_config(
self.script_dir
)
self.config.potfile_path.parent.mkdir(
parents=True,
exist_ok=True,
)
self.paths = ClientPaths.from_root(
self.config.client_root
)
self.paths.initialize()
self.logger = configure_logging(
self.paths.logs / "client.log"
)
self.package_service = JobPackageService(
self.paths
)
self.logger.info(
"CrackLab client initialized"
)
self.logger.info(
"Client root: %s",
self.paths.root,
)
self.logger.info(
"Hashcat: %s",
self.config.hashcat_exe,
)
self.logger.info(
"Potfile: %s",
self.config.potfile_path,
)
self.logger.info(
"poll_seconds=%s max_parallel_jobs=%s",
self.config.poll_seconds,
self.config.max_parallel_jobs,
)
def import_jobs(self) -> list[ImportedJob]:
imported: list[ImportedJob] = []
packages = sorted(
self.paths.inbox.glob(
"JOB-*.zip"
)
)
if not packages:
self.logger.info(
"Inbox is empty."
)
return imported
for package in packages:
try:
self.logger.info(
"Importing Job package: %s",
package.name,
)
item = (
self.package_service
.import_package(package)
)
imported.append(item)
if item.idempotent:
self.logger.info(
"Job package already imported: "
"%s",
item.job_id,
)
else:
self.logger.info(
"Job imported: %s",
item.job_id,
)
except JobPackageError as exc:
self.logger.error(
"Job package rejected: %s: %s",
package.name,
exc,
)
except Exception:
self.logger.exception(
"Unexpected error importing %s",
package.name,
)
return imported
def _load_job(
self,
job_dir: Path,
) -> dict[str, Any]:
path = job_dir / "job.json"
try:
data = json.loads(
path.read_text(
encoding="utf-8"
)
)
except (OSError, json.JSONDecodeError) as exc:
raise ClientError(
f"Failed to read {path}: {exc}"
) from exc
if not isinstance(data, dict):
raise ClientError(
f"Invalid Job JSON: {path}"
)
return data
def _ensure_state(
self,
job_dir: Path,
job: dict[str, Any],
) -> dict[str, Any]:
path = job_dir / "state" / "client_state.json"
if not path.is_file():
self.logger.info(
"%s: creating client state",
job["job_id"],
)
return initialize_state(
job_dir,
job,
)
try:
state = load_state(
job_dir
)
except StateError:
self.logger.warning(
"%s: invalid state; rebuilding state",
job["job_id"],
)
return initialize_state(
job_dir,
job,
)
if state.get("job_id") != job["job_id"]:
raise ClientError(
"State job_id does not match job.json."
)
expected_steps = {
str(step["step_no"])
for step in job["steps"]
}
actual_steps = state.get(
"steps",
{},
)
if not isinstance(actual_steps, dict):
self.logger.warning(
"%s: state has invalid steps; rebuilding",
job["job_id"],
)
return initialize_state(
job_dir,
job,
)
if set(actual_steps) != expected_steps:
self.logger.warning(
"%s: state steps do not match "
"job.json; rebuilding",
job["job_id"],
)
return initialize_state(
job_dir,
job,
)
return state
def _step_result(
self,
job_dir: Path,
step: dict[str, Any],
) -> dict[str, Any]:
state = get_step_state(
job_dir,
step["step_no"],
)
return {
"step_no": step["step_no"],
"step_id": step["step_id"],
"session_name": step["session_name"],
"status": state["status"],
"started_at": state.get(
"started_at"
),
"completed_at": state.get(
"completed_at"
),
"exit_code": state.get(
"exit_code"
),
"restore_seen": bool(
state.get("restore_seen", False)
),
"error": state.get("error"),
}
def _run_step(
self,
job_dir: Path,
job: dict[str, Any],
step: dict[str, Any],
job_logger: logging.Logger,
) -> bool:
step_no = step["step_no"]
state = get_step_state(
job_dir,
step_no,
)
if (
state.get("status")
== STEP_STATUS_COMPLETED
):
job_logger.info(
"Step %s already COMPLETED; skipping.",
step_no,
)
return True
restore_file = (
job_dir
/ "state"
/ f"step-{step_no:03d}.restore"
)
restore_seen = restore_file.is_file()
mode = (
"RESTORE"
if restore_seen
else "START"
)
job_logger.info(
"Step %s (%s) starting: %s "
"session=%s",
step_no,
step["step_id"],
mode,
step["session_name"],
)
if restore_seen:
job_logger.warning(
"Step %s has restore file; "
"resuming existing Hashcat session.",
step_no,
)
update_step_state(
job_dir,
step_no,
STEP_STATUS_RUNNING,
restore_seen=restore_seen,
started=True,
)
log_path = (
job_dir
/ "logs"
/ f"step-{step_no:03d}.log"
)
try:
result = run_step(
self.config,
job_dir,
job,
step,
log_path,
restore=restore_seen,
)
except HashcatError as exc:
current_restore = (
restore_file.is_file()
)
update_step_state(
job_dir,
step_no,
(
STEP_STATUS_PARTIAL
if current_restore
else STEP_STATUS_FAILED
),
exit_code=exc.exit_code,
restore_seen=current_restore,
error=str(exc),
completed=True,
)
if current_restore:
job_logger.warning(
"Step %s interrupted/partial: %s",
step_no,
exc,
)
else:
job_logger.error(
"Step %s failed: %s",
step_no,
exc,
)
return False
except Exception as exc:
update_step_state(
job_dir,
step_no,
STEP_STATUS_FAILED,
restore_seen=restore_file.is_file(),
error=str(exc),
completed=True,
)
job_logger.exception(
"Unexpected Step %s error.",
step_no,
)
return False
if result.restore_seen:
update_step_state(
job_dir,
step_no,
STEP_STATUS_PARTIAL,
exit_code=result.exit_code,
restore_seen=True,
error=(
"Hashcat finished while "
"restore file still exists."
),
completed=True,
)
job_logger.warning(
"Step %s remains PARTIAL: "
"restore file exists after Hashcat "
"exit=%s.",
step_no,
result.exit_code,
)
return False
if result.exit_code not in (0, 1):
update_step_state(
job_dir,
step_no,
STEP_STATUS_FAILED,
exit_code=result.exit_code,
restore_seen=False,
error=(
f"Hashcat exited with "
f"code {result.exit_code}."
),
completed=True,
)
job_logger.error(
"Step %s failed: Hashcat exit=%s.",
step_no,
result.exit_code,
)
return False
marker = (
job_dir
/ "state"
/ f"step-{step_no:03d}.completed"
)
marker.write_text(
"COMPLETED\n",
encoding="utf-8",
)
update_step_state(
job_dir,
step_no,
STEP_STATUS_COMPLETED,
exit_code=result.exit_code,
restore_seen=False,
error=None,
completed=True,
)
job_logger.info(
"Step %s COMPLETED.",
step_no,
)
return True
def run_job(
self,
job_dir: Path,
) -> bool:
job = self._load_job(
job_dir
)
job_id = job["job_id"]
job_logger = configure_job_logging(
job_dir / "logs" / "client.log",
job_id,
)
job_logger.info(
"Job processing started."
)
job_logger.info(
"method=%s version=%s hashes=%s "
"handshakes=%s steps=%s",
job["method_id"],
job["method_version"],
job["hash_count"],
job["handshake_count"],
job["step_count"],
)
try:
state = self._ensure_state(
job_dir,
job,
)
if (
state.get("status")
== JOB_STATUS_COMPLETED
):
try:
archived_package = (
self.package_service.archive_job(
job_dir.name
)
)
except JobPackageError as exc:
job_logger.error(
"Job package archival failed: %s",
exc,
)
return False
except Exception:
job_logger.exception(
"Unexpected Job package archival failure."
)
return False
job_logger.info(
"Job already COMPLETED; "
"Job package finalized: %s",
archived_package,
)
return True
update_job_state(
job_dir,
JOB_STATUS_RUNNING,
)
all_completed = True
for step in job["steps"]:
if not self._run_step(
job_dir,
job,
step,
job_logger,
):
all_completed = False
break
if not all_completed:
update_job_state(
job_dir,
JOB_STATUS_PARTIAL,
)
job_logger.warning(
"Job remains PARTIAL."
)
return False
job_logger.info(
"All steps completed; running Hashcat --show."
)
show_path = (
job_dir
/ "results"
/ "hashcat-show.txt"
)
show_log = (
job_dir
/ "logs"
/ "show.log"
)
try:
show_output = run_show(
self.config,
job_dir,
job,
show_path,
show_log,
)
job_logger.info(
"Hashcat --show completed: "
"%s bytes, %s lines.",
show_path.stat().st_size,
len(show_output.splitlines()),
)
except HashcatError as exc:
update_job_state(
job_dir,
JOB_STATUS_FAILED,
error=str(exc),
)
job_logger.error(
"Hashcat --show failed: %s",
exc,
)
return False
except Exception as exc:
update_job_state(
job_dir,
JOB_STATUS_FAILED,
error=str(exc),
)
job_logger.exception(
"Unexpected --show error."
)
return False
step_states = [
self._step_result(
job_dir,
step,
)
for step in job["steps"]
]
try:
result_package = (
create_result_package(
job_dir,
job,
step_states,
self.paths.outbox,
logger=job_logger,
)
)
except ResultPackageError as exc:
update_job_state(
job_dir,
JOB_STATUS_FAILED,
error=str(exc),
)
job_logger.error(
"Result package creation failed: %s",
exc,
)
return False
except Exception as exc:
update_job_state(
job_dir,
JOB_STATUS_FAILED,
error=str(exc),
)
job_logger.exception(
"Unexpected result package error."
)
return False
update_job_state(
job_dir,
JOB_STATUS_COMPLETED,
)
try:
archived_package = (
self.package_service.archive_job(
job_dir.name
)
)
except JobPackageError as exc:
job_logger.error(
"Job package archival failed: %s",
exc,
)
return False
except Exception:
job_logger.exception(
"Unexpected Job package archival failure."
)
return False
job_logger.info(
"Job COMPLETED."
)
job_logger.info(
"Result package: %s",
result_package,
)
job_logger.info(
"Job package archived: %s",
archived_package,
)
return True
except Exception as exc:
update_job_state(
job_dir,
JOB_STATUS_FAILED,
error=str(exc),
)
job_logger.error(
"Job failed: %s",
exc,
)
job_logger.debug(
"%s",
traceback.format_exc(),
)
return False
def run_once(self) -> list[str]:
imported = self.import_jobs()
for item in imported:
self.logger.info(
"Queue item: %s idempotent=%s",
item.job_id,
item.idempotent,
)
completed: list[str] = []
for job_dir in sorted(
self.paths.jobs.glob("JOB-*")
):
if not job_dir.is_dir():
continue
try:
if self.run_job(
job_dir
):
completed.append(
job_dir.name
)
except Exception:
self.logger.exception(
"Unhandled error processing %s",
job_dir.name,
)
return completed
def run_forever(self) -> None:
self.logger.info(
"Client loop started."
)
while True:
try:
self.run_once()
except KeyboardInterrupt:
self.logger.info(
"Client stopped by user."
)
return
except Exception:
self.logger.exception(
"Unhandled client loop error."
)
import time
time.sleep(
self.config.poll_seconds
)