# retoor <retoor@molodetz.nl>
from __future__ import annotations
import asyncio
import hashlib
import logging
from pathlib import Path
from typing import AsyncIterator
from devplacepy.services.jobs.isslop.acquisition.domcapture import DomSnapshot
from devplacepy.services.jobs.isslop.acquisition.git import CloneFailedError, RepositoryTooLargeError, clone_repository
from devplacepy.services.jobs.isslop.acquisition.source import KIND_GIT, is_private_host, resolve_source
from devplacepy.services.jobs.isslop.acquisition.website import crawl_website
from devplacepy.services.jobs.isslop.acquisition.workspace import content_hash, remove_workspace, reset_workspace
from devplacepy.services.jobs.isslop.agent.classifier import AiVerdict, classify_file, select_samples
from devplacepy.services.jobs.isslop.agent.llm import LlmClient
from devplacepy.services.jobs.isslop.agent.reporter import generate_report
from devplacepy.services.jobs.isslop.agent.vision import (
ImageVerdict,
classify_image,
collect_images,
image_summary,
make_thumbnail,
persist_screenshot,
)
from devplacepy.services.jobs.isslop.analysis.domsignals.aggregate import aggregate_dom_evidence
from devplacepy.services.jobs.isslop.analysis.domsignals.base import DomPageContext
from devplacepy.services.jobs.isslop.analysis.domsignals.builders import detected_builder as detected_builder_for
from devplacepy.services.jobs.isslop.analysis.domsignals.registry import run_dom_checks
from devplacepy.services.jobs.isslop.analysis.engine import build_inventory, compute_repo_baselines, load_context, score_file
from devplacepy.services.jobs.isslop.analysis.scoring import (
FileScore,
adjust_for_dom_signals,
adjust_for_images,
adjust_for_template,
aggregate,
categorize,
criticality_for,
repo_scores_to_dict,
)
from devplacepy.services.jobs.isslop.analysis.signals import SEVERITY_STRONG
from devplacepy.services.jobs.isslop.analysis.templates import TemplateEvidence, detect_template
from devplacepy.services.jobs.isslop.config import AI_EXCERPT_CHARS, IMAGE_CONCURRENCY, WorkerSettings
from devplacepy.services.jobs.isslop.events import (
KIND_AI,
KIND_DOM,
KIND_DONE,
KIND_ERROR,
KIND_FILE,
KIND_LOG,
KIND_PROGRESS,
KIND_REPORT,
KIND_SCORE,
KIND_SIGNAL,
KIND_STAGE,
STAGE_ACQUIRE,
STAGE_AI,
STAGE_DOM,
STAGE_IMAGES,
STAGE_INVENTORY,
STAGE_REPORT,
STAGE_RESOLVE,
STAGE_SCORE,
STAGE_STATIC,
KIND_IMAGE,
WorkerEvent,
)
logger = logging.getLogger(__name__)
PROGRESS_EVERY_FILES: int = 5
SOURCE_CAP_FILES: int = 60
SOURCE_CAP_BYTES: int = 200000
STATIC_BLEND_WEIGHT: float = 0.6
AI_BLEND_WEIGHT: float = 0.4
AI_REVIEW_CONCURRENCY: int = 4
def _persist_source(media_dir: Path | None, relative: str, text: str) -> str | None:
if media_dir is None:
return None
name = f"s{hashlib.sha1(relative.encode('utf-8')).hexdigest()[:16]}.txt"
try:
media_dir.mkdir(parents=True, exist_ok=True)
(media_dir / name).write_text(text[:SOURCE_CAP_BYTES], encoding="utf-8")
except OSError as error:
logger.warning("Source snippet persist failed for %s: %s", relative, error)
return None
return name
def _stage(stage: str, message: str) -> WorkerEvent:
return WorkerEvent(kind=KIND_STAGE, message=message, data={"stage": stage})
def _blend(static_value: float, ai_value: float) -> float:
return round(STATIC_BLEND_WEIGHT * static_value + AI_BLEND_WEIGHT * ai_value, 1)
def apply_ai_verdicts(scores: list[FileScore], verdicts: list[AiVerdict]) -> list[FileScore]:
by_path = {verdict.path: verdict for verdict in verdicts}
adjusted: list[FileScore] = []
for score in scores:
verdict = by_path.get(score.relative)
if verdict is None:
adjusted.append(score)
continue
origin = _blend(score.origin_score, verdict.origin_score)
quality = _blend(score.quality_deficit, verdict.quality_deficit)
adjusted.append(
FileScore(
relative=score.relative,
language=score.language,
sloc=score.sloc,
origin_score=origin,
quality_deficit=quality,
category=categorize(origin, quality),
criticality=criticality_for(score.relative),
signals=score.signals,
)
)
return adjusted
async def _cleanup_event(workspace: Path, stage: str) -> WorkerEvent:
removed = await asyncio.to_thread(remove_workspace, workspace)
message = "Workspace deleted, only the persisted report remains" if removed else "Workspace already absent"
return WorkerEvent(KIND_LOG, message, {"stage": stage, "workspace_removed": removed})
async def run_pipeline(source_url: str, workspace: Path, config: WorkerSettings) -> AsyncIterator[WorkerEvent]:
yield _stage(STAGE_RESOLVE, f"Resolving source type for {source_url}")
if not config.allow_private_hosts and is_private_host(source_url):
yield WorkerEvent(KIND_ERROR, "Refusing to analyze private or loopback hosts", {"stage": STAGE_RESOLVE})
return
resolution = await resolve_source(source_url)
yield WorkerEvent(
KIND_LOG,
f"Source classified as {resolution.kind}",
{"stage": STAGE_RESOLVE, "kind_detected": resolution.kind, "url": resolution.url},
)
yield _stage(STAGE_ACQUIRE, f"Acquiring source into workspace ({resolution.kind})")
reset_workspace(workspace)
dom_snapshots: list[DomSnapshot] = []
try:
if resolution.kind == KIND_GIT:
async for progress in clone_repository(resolution.url, workspace):
yield WorkerEvent(KIND_LOG, progress, {"stage": STAGE_ACQUIRE})
else:
async for progress in crawl_website(
resolution.url, workspace, dom_sink=dom_snapshots, allow_private=config.allow_private_hosts
):
yield WorkerEvent(KIND_LOG, progress, {"stage": STAGE_ACQUIRE})
except (RepositoryTooLargeError, CloneFailedError, RuntimeError) as error:
yield await _cleanup_event(workspace, STAGE_ACQUIRE)
yield WorkerEvent(KIND_ERROR, str(error), {"stage": STAGE_ACQUIRE})
return
digest = content_hash(workspace)
yield WorkerEvent(KIND_LOG, f"Workspace content hash {digest[:16]}", {"stage": STAGE_ACQUIRE, "content_hash": digest})
yield _stage(STAGE_INVENTORY, "Building file inventory and applying exclusion rules")
inventory = build_inventory(workspace)
yield WorkerEvent(
KIND_LOG,
f"{len(inventory.analyzable)} files analyzable, {len(inventory.excluded)} excluded "
f"(vendored, generated, binary, lockfiles, oversize)",
{
"stage": STAGE_INVENTORY,
"analyzable": len(inventory.analyzable),
"excluded": len(inventory.excluded),
"excluded_samples": [
{"path": path, "reason": reason} for path, reason in inventory.excluded[:15]
],
},
)
if not inventory.analyzable:
yield await _cleanup_event(workspace, STAGE_INVENTORY)
yield WorkerEvent(KIND_ERROR, "No analyzable source files found in this source", {"stage": STAGE_INVENTORY})
return
template = detect_template(workspace) if config.template_detection else TemplateEvidence()
if template.markers:
yield WorkerEvent(
KIND_SIGNAL,
f"Starter template provenance detected (score {template.score}/100): "
+ "; ".join(template.markers[:4]),
{"stage": STAGE_INVENTORY, "template_score": template.score, "markers": template.markers},
)
yield _stage(STAGE_STATIC, f"Running static multi-signal analysis on {len(inventory.analyzable)} files")
contexts = []
for entry in inventory.analyzable:
context = load_context(entry, inventory.repo)
if context is not None:
contexts.append(context)
else:
yield WorkerEvent(
KIND_LOG,
f"Skipped {entry.relative}: unreadable content",
{"stage": STAGE_STATIC, "path": entry.relative, "reason": "unreadable"},
)
if not contexts:
yield await _cleanup_event(workspace, STAGE_STATIC)
yield WorkerEvent(
KIND_ERROR,
"Insufficient analyzable content: every candidate file is minified, generated or unreadable. "
"Abstaining rather than guessing, per methodology.",
{"stage": STAGE_STATIC, "candidates": len(inventory.analyzable)},
)
return
baselines = compute_repo_baselines(contexts)
yield WorkerEvent(
KIND_LOG,
f"Repo baselines computed over {baselines.file_count} files "
f"(mean indent deviation {baselines.indent_variance_mean:.2f}, mean comment ratio {baselines.comment_ratio_mean:.2f})",
{"stage": STAGE_STATIC},
)
scores: list[FileScore] = []
excerpts: dict[str, str] = {}
source_media_dir = Path(config.media_dir) if config.media_dir else None
sources_persisted = 0
for index, context in enumerate(contexts, start=1):
score = score_file(context, baselines)
scores.append(score)
excerpts[score.relative] = context.text[:AI_EXCERPT_CHARS]
source_name = None
if score.signals and sources_persisted < SOURCE_CAP_FILES:
source_name = await asyncio.to_thread(
_persist_source, source_media_dir, score.relative, context.text
)
if source_name:
sources_persisted += 1
mode = "fingerprint scan" if context.fingerprint_only else "full analysis"
yield WorkerEvent(
KIND_FILE,
f"Checked {score.relative} ({mode}): origin {score.origin_score}, quality deficit {score.quality_deficit}, {score.category}",
{
"stage": STAGE_STATIC,
"path": score.relative,
"language": score.language,
"sloc": score.sloc,
"origin_score": score.origin_score,
"quality_deficit": score.quality_deficit,
"category": score.category,
"fingerprint_only": context.fingerprint_only,
"signals": [signal.to_dict() for signal in score.signals[:20]],
"source": source_name,
},
)
for signal in score.signals:
if signal.severity == SEVERITY_STRONG:
yield WorkerEvent(
KIND_SIGNAL,
f"Strong signal in {score.relative}:{signal.line} - {signal.title}",
{"stage": STAGE_STATIC, "path": score.relative, "signal": signal.to_dict()},
)
if index % PROGRESS_EVERY_FILES == 0 or index == len(contexts):
yield WorkerEvent(
KIND_PROGRESS,
f"Static analysis {index}/{len(contexts)} files",
{"stage": STAGE_STATIC, "current": index, "total": len(contexts)},
)
dom_pages = [
DomPageContext(
url=snapshot.render.final_url,
dom=snapshot.dom,
console_warnings=snapshot.console_warnings,
console_errors=snapshot.console_errors,
response_headers=snapshot.response_headers,
resource_hosts=snapshot.resource_hosts,
screenshot_bytes=snapshot.screenshot_bytes,
)
for snapshot in dom_snapshots
]
dom_evidence = aggregate_dom_evidence(dom_pages)
dom_media_dir = Path(config.media_dir) if config.media_dir else None
if dom_pages:
yield _stage(STAGE_DOM, f"Analyzing rendered DOM for AI-builder and AI-slop tells on {len(dom_pages)} page(s)")
for page in dom_pages:
screenshot_name = None
if page.screenshot_bytes is not None and dom_media_dir is not None:
screenshot_name = await asyncio.to_thread(
persist_screenshot, page.screenshot_bytes, dom_media_dir, page.url
)
page_signals = run_dom_checks(page)
page_builder, page_builder_confidence = detected_builder_for([page])
yield WorkerEvent(
KIND_DOM,
f"DOM analysis of {page.url}: {len(page_signals)} signals, "
f"builder {page_builder or 'none'}",
{
"stage": STAGE_DOM,
"url": page.url,
"detected_builder": page_builder,
"builder_confidence": page_builder_confidence,
"signal_count": len(page_signals),
"signals": [signal.to_dict() for signal in page_signals[:20]],
"screenshot": screenshot_name,
},
)
yield _stage(STAGE_AI, "Selecting representative files for AI review")
llm = LlmClient(config)
samples = select_samples(scores) if llm.review_available else []
if not llm.review_available:
yield WorkerEvent(
KIND_LOG,
"AI review disabled or unavailable; static signals remain authoritative",
{"stage": STAGE_AI},
)
yield WorkerEvent(
KIND_LOG,
f"AI reviewing {len(samples)} representative files, up to {AI_REVIEW_CONCURRENCY} concurrently "
f"(deterministic selection, temperature 0)",
{"stage": STAGE_AI, "samples": [sample.relative for sample in samples]},
)
verdicts: list[AiVerdict] = []
review_semaphore = asyncio.Semaphore(AI_REVIEW_CONCURRENCY)
async def review(sample: FileScore) -> tuple[FileScore, AiVerdict | None]:
async with review_semaphore:
return sample, await classify_file(llm, sample, excerpts.get(sample.relative, ""))
review_tasks = [asyncio.create_task(review(sample)) for sample in samples]
for index, finished in enumerate(asyncio.as_completed(review_tasks), start=1):
sample, verdict = await finished
if verdict is None:
yield WorkerEvent(
KIND_LOG,
f"AI review unavailable for {sample.relative}, static signals remain authoritative",
{"stage": STAGE_AI, "path": sample.relative},
)
else:
verdicts.append(verdict)
yield WorkerEvent(
KIND_AI,
f"AI verdict {sample.relative}: {verdict.category} ({verdict.ai_probability:.0f}% AI)",
{
"stage": STAGE_AI,
"path": verdict.path,
"category": verdict.category,
"origin_score": verdict.origin_score,
"quality_deficit": verdict.quality_deficit,
"ai_probability": verdict.ai_probability,
"reasoning": verdict.reasoning,
"notable_signals": verdict.notable_signals,
},
)
yield WorkerEvent(
KIND_PROGRESS,
f"AI review {index}/{len(samples)} files",
{"stage": STAGE_AI, "current": index, "total": len(samples)},
)
yield _stage(STAGE_IMAGES, "Reviewing images for AI-generated content")
image_verdicts: list[ImageVerdict] = []
image_stats: dict[str, object] = {"count": 0, "mean_ai_probability": 0.0, "grade": "n/a", "ai_generated_count": 0}
if not llm.vision_available:
yield WorkerEvent(KIND_LOG, "Vision backend unavailable; skipping image analysis", {"stage": STAGE_IMAGES})
else:
images = await asyncio.to_thread(collect_images, workspace, config.image_max_count)
if not images:
yield WorkerEvent(KIND_LOG, "No qualifying images found to review", {"stage": STAGE_IMAGES})
else:
yield WorkerEvent(
KIND_LOG,
f"Reviewing {len(images)} images with the vision model, up to {IMAGE_CONCURRENCY} concurrently "
f"(deterministic sample, capped at {config.image_max_count})",
{"stage": STAGE_IMAGES, "images": [str(path.relative_to(workspace)) for path in images]},
)
media_dir = Path(config.media_dir) if config.media_dir else None
thumbs: dict[str, str | None] = {}
for image_path in images:
relative = str(image_path.relative_to(workspace))
thumb = None
if media_dir is not None:
thumb = await asyncio.to_thread(make_thumbnail, image_path, media_dir, relative)
thumbs[relative] = thumb
yield WorkerEvent(
KIND_LOG,
f"Reviewing image {relative}",
{"stage": STAGE_IMAGES, "path": relative, "thumb": thumb},
)
image_semaphore = asyncio.Semaphore(IMAGE_CONCURRENCY)
async def review_image(image_path: Path) -> ImageVerdict | None:
async with image_semaphore:
return await classify_image(llm, workspace, image_path)
image_tasks = [asyncio.create_task(review_image(image_path)) for image_path in images]
for index, finished in enumerate(asyncio.as_completed(image_tasks), start=1):
verdict = await finished
if verdict is None:
continue
image_verdicts.append(verdict)
yield WorkerEvent(
KIND_IMAGE,
f"Image {verdict.relative}: grade {verdict.grade}, {verdict.verdict} "
f"({verdict.ai_probability:.0f}% AI) - {verdict.image_kind}",
{
"stage": STAGE_IMAGES,
"path": verdict.relative,
"ai_probability": verdict.ai_probability,
"grade": verdict.grade,
"verdict": verdict.verdict,
"image_kind": verdict.image_kind,
"tells": verdict.tells,
"description": verdict.description,
"thumb": thumbs.get(verdict.relative),
},
)
yield WorkerEvent(
KIND_PROGRESS,
f"Image review {index}/{len(images)}",
{"stage": STAGE_IMAGES, "current": index, "total": len(images)},
)
image_stats = image_summary(image_verdicts)
yield WorkerEvent(
KIND_PROGRESS,
f"Image review complete: {image_stats['count']} images, "
f"{image_stats['ai_generated_count']} look AI-generated, grade {image_stats['grade']}",
{"stage": STAGE_IMAGES, **image_stats},
)
yield _stage(STAGE_SCORE, "Aggregating per-file scores into repository verdict")
static_scores = aggregate(scores)
blended = bool(verdicts)
adjusted = apply_ai_verdicts(scores, verdicts) if blended else scores
final_scores = aggregate(adjusted)
image_influenced = False
if image_verdicts:
final_scores = adjust_for_images(final_scores, float(image_stats["mean_ai_probability"]))
image_influenced = True
final_scores = adjust_for_dom_signals(final_scores, dom_evidence)
final_scores = adjust_for_template(final_scores, template.score)
score_payload = {
"stage": STAGE_SCORE,
"static": repo_scores_to_dict(static_scores),
"blended_with_ai": blended,
"template_score": template.score,
"template_markers": template.markers,
**repo_scores_to_dict(final_scores),
"files_total": len(inventory.analyzable) + len(inventory.excluded),
"files_analyzed": len(scores),
"content_hash": digest,
"source_kind": resolution.kind,
"image_review": image_stats,
"image_influenced": image_influenced,
"detected_builder": dom_evidence.detected_builder,
"dom_slop_score": dom_evidence.score,
}
yield WorkerEvent(
KIND_SCORE,
f"Final verdict: grade {final_scores.grade}, {final_scores.category}, "
f"{final_scores.human_percent}% human / {final_scores.ai_percent}% AI, confidence {final_scores.confidence}",
score_payload,
)
yield _stage(STAGE_REPORT, "Generating final report with findings")
markdown, model_used = await generate_report(
llm,
source_url,
resolution.kind,
final_scores,
adjusted,
verdicts,
len(inventory.excluded),
image_verdicts,
image_stats,
template,
dom_evidence,
)
await llm.aclose()
yield WorkerEvent(
KIND_REPORT,
f"Report generated via {model_used} ({len(markdown)} chars)",
{"stage": STAGE_REPORT, "markdown": markdown, "model_used": model_used},
)
yield await _cleanup_event(workspace, STAGE_REPORT)
yield WorkerEvent(KIND_DONE, "Analysis complete", score_payload)