diff --git a/periscope/src/backend/services/email.py b/periscope/src/backend/services/email.py new file mode 100644 index 0000000..391a202 --- /dev/null +++ b/periscope/src/backend/services/email.py @@ -0,0 +1,316 @@ +"""Outbound mail via Gmail API. No-op when email settings are empty. + +Clerk lookup is leftover-router compatible. Hosted auth is optional; self-host +with JWT never calls Clerk. +""" + +from __future__ import annotations + +import asyncio +import base64 +import html +import logging +from email.mime.multipart import MIMEMultipart +from email.mime.text import MIMEText + +import httpx + +from backend.config import settings + +log = logging.getLogger(__name__) + + +async def _resolve_clerk_user(user_id: str) -> dict | None: + if not getattr(settings, "use_auth", False): + return None + secret = getattr(settings, "clerk_secret_key", "") or "" + if not secret: + return None + try: + async with httpx.AsyncClient(timeout=10) as client: + resp = await client.get( + f"https://api.clerk.com/v1/users/{user_id}", + headers={"Authorization": f"Bearer {secret}"}, + ) + if resp.status_code == 200: + return resp.json() + except Exception: + log.warning("Clerk lookup failed for %s", user_id) + return None + + +def _addr(user: dict | None) -> tuple[str, str | None]: + if not user: + return "there", None + first = user.get("first_name") or "" + last = user.get("last_name") or "" + name = f"{first} {last}".strip() or "there" + emails = user.get("email_addresses") or [] + to = emails[0].get("email_address") if emails else None + return name, to + + +def _build_gmail_service(): + try: + import google.auth + import google.auth.transport.requests + from google.auth import iam + from google.oauth2 import service_account + from googleapiclient.discovery import build + except ImportError: + log.warning("Gmail client libraries missing") + return None + scopes = ["https://www.googleapis.com/auth/gmail.send"] + try: + creds, _ = google.auth.default() + if hasattr(creds, "_signer"): + delegated = creds.with_subject(settings.email_sender) + return build("gmail", "v1", credentials=delegated, cache_discovery=False) + req = google.auth.transport.requests.Request() + creds.refresh(req) + sa = getattr(creds, "service_account_email", None) + if not sa: + return None + signer = iam.Signer(req, creds, sa) + sa_creds = service_account.Credentials( + signer, sa, token_uri="https://oauth2.googleapis.com/token", + scopes=scopes, subject=settings.email_sender, + ) + return build("gmail", "v1", credentials=sa_creds, cache_discovery=False) + except Exception: + log.exception("Gmail service build failed") + return None + + +def _encode_message(msg: MIMEMultipart) -> dict: + raw = base64.urlsafe_b64encode(msg.as_bytes()).decode() + return {"raw": raw} + + +async def _send_raw(to_email: str, msg: MIMEMultipart, label: str) -> None: + def _send() -> None: + svc = _build_gmail_service() + if svc is None: + log.warning("%s skipped: no Gmail service", label) + return + svc.users().messages().send(userId="me", body=_encode_message(msg)).execute() + + try: + await asyncio.to_thread(_send) + except Exception: + log.exception("%s failed to=%s", label, to_email) + + +def _plain_msg(to_email: str, subject: str, body: str) -> MIMEMultipart: + msg = MIMEMultipart("alternative") + msg["From"] = f"Periscope <{settings.email_sender}>" + msg["To"] = to_email + msg["Subject"] = subject + msg.attach(MIMEText(body, "plain")) + return msg + + +async def send_topup_failed_email(user_id: str, *, amount_usd: float, reason: str) -> None: + if not settings.use_email: + return + name, to = _addr(await _resolve_clerk_user(user_id)) + if not to: + return + body = ( + f"Hi {name},\n\nAuto top-up of ${amount_usd:.2f} failed: {reason}\n" + f"Top up: {settings.email_frontend_url}/credits\n" + ) + await _send_raw(to, _plain_msg(to, "Periscope: auto top-up failed", body), "topup-failed") + + +async def send_low_balance_email(user_id: str, *, balance: float, threshold: float) -> None: + if not settings.use_email: + return + name, to = _addr(await _resolve_clerk_user(user_id)) + if not to: + return + body = ( + f"Hi {name},\n\nBalance is {balance:.2f} credits (threshold {threshold:.2f}).\n" + f"Top up: {settings.email_frontend_url}/credits\n" + ) + await _send_raw(to, _plain_msg(to, "Periscope: low credit balance", body), "low-balance") + + +async def send_pipeline_paused_email( + user_id: str, + project_name: str, + project_id: str, + *, + last_completed: str | None, + stage: str | None, + balance: float, + credits_needed_low: float, +) -> None: + if not settings.use_email: + return + name, to = _addr(await _resolve_clerk_user(user_id)) + if not to: + return + url = f"{settings.email_frontend_url}/project/{project_id}" + body = ( + f"Hi {name},\n\n\"{project_name}\" paused for credits.\n" + f"Last: {last_completed or '—'} stage={stage or '—'}\n" + f"Balance {balance:.2f}; about {credits_needed_low:.2f} more needed.\n{url}\n" + ) + await _send_raw(to, _plain_msg(to, f"Periscope paused: {project_name}", body), "paused") + + +async def send_report_ready_email( + user_id: str, + project_name: str, + project_id: str, + summary: dict[str, int], + total_cost_usd: float | None = None, +) -> None: + if not settings.use_email: + return + name, to = _addr(await _resolve_clerk_user(user_id)) + if not to: + return + bits = ", ".join(f"{k}={v}" for k, v in (summary or {}).items()) or "no summary" + cost = f"\nCost: ${total_cost_usd:.4f}" if total_cost_usd is not None else "" + url = f"{settings.email_frontend_url}/project/{project_id}" + body = f"Hi {name},\n\nReport ready for \"{project_name}\".\n{bits}{cost}\n{url}\n" + await _send_raw(to, _plain_msg(to, f"Periscope report: {project_name}", body), "report-ready") + + +async def send_test_email(to_email: str) -> dict: + result = {"ok": False, "step": "", "error": ""} + if not settings.use_email: + result["step"] = "config" + result["error"] = ( + f"use_email=False (email_sender={settings.email_sender!r}, " + f"email_frontend_url={getattr(settings, 'email_frontend_url', '')!r})" + ) + return result + result["step"] = "build_service" + try: + svc = _build_gmail_service() + except Exception as exc: + result["error"] = str(exc) + return result + if svc is None: + result["error"] = "Gmail service unavailable" + return result + result["step"] = "send" + try: + msg = _plain_msg(to_email, "Periscope test email", "This is a Periscope mail test.") + await asyncio.to_thread( + lambda: svc.users().messages().send( + userId="me", body=_encode_message(msg), + ).execute() + ) + result["ok"] = True + return result + except Exception as exc: + result["error"] = str(exc) + return result + + +async def send_pipeline_started_email( + user_id: str, + project_name: str, + project_id: str, + num_components: int, + num_nets: int, + num_ics: int, + num_passives: int, + num_simple: int, +) -> None: + if not settings.use_email or not getattr(settings, "email_admin_notify", ""): + return + name, email = _addr(await _resolve_clerk_user(user_id)) + to = settings.email_admin_notify + url = f"{settings.email_frontend_url}/project/{project_id}" + body = ( + f"Pipeline started for \"{project_name}\"\n" + f"Created by: {name} ({email or user_id})\n" + f"Components: {num_components} ({num_ics} ICs, {num_passives} passives, " + f"{num_simple} discrete)\nNets: {num_nets}\n{url}\n" + ) + await _send_raw( + to, + _plain_msg(to, f"Pipeline started: {project_name} ({num_components} components)", body), + "pipeline-started", + ) + + +async def send_feedback_received_email( + ticket_id: str, + user_id: str, + feedback_type: str, + message: str, + *, + submitter_name: str | None = None, + submitter_email: str | None = None, + project_name: str | None = None, + project_id: str | None = None, + finding_designator: str | None = None, + finding_mpn: str | None = None, + finding_status: str | None = None, + finding_text: str | None = None, +) -> None: + if not settings.use_email or not getattr(settings, "email_admin_notify", ""): + return + name = (submitter_name or "").strip() + email = (submitter_email or "").strip() + if not name or not email: + n, e = _addr(await _resolve_clerk_user(user_id)) + name = name or n + email = email or (e or user_id) + to = settings.email_admin_notify + ctx = project_name or "general" + extra = "" + if finding_designator: + extra = f"\nFinding: {finding_designator} {finding_mpn or ''} {finding_status or ''}\n{finding_text or ''}" + pid = f"\nProject: {project_id}" if project_id else "" + body = ( + f"Ticket {ticket_id} ({feedback_type}) from {name} <{email}>\n" + f"Context: {ctx}{pid}{extra}\n\n{message}\n" + ) + await _send_raw( + to, + _plain_msg(to, f"Feedback {ticket_id}: {html.escape(ctx)}", body), + "feedback-in", + ) + + +async def send_feedback_reply_email( + user_id: str, + reply_text: str, + original_message: str, + *, + recipient_name: str | None = None, + recipient_email: str | None = None, + project_name: str | None = None, + finding_designator: str | None = None, + finding_mpn: str | None = None, +) -> None: + if not settings.use_email: + return + name = (recipient_name or "").strip() + to = (recipient_email or "").strip() + if not name or not to: + n, e = _addr(await _resolve_clerk_user(user_id)) + name = name or n + to = to or (e or "") + if not to: + log.warning("feedback reply: no address for %s", user_id) + return + first = name.split()[0] if name else "there" + extra = f"\n{finding_designator} {finding_mpn or ''}" if finding_designator else "" + proj = f"\nProject: {project_name}" if project_name else "" + body = ( + f"Hi {first},\n\nThe Periscope team replied to your feedback.{proj}{extra}\n\n" + f"— Reply —\n{reply_text}\n\n— Original —\n{original_message}\n" + ) + await _send_raw( + to, + _plain_msg(to, "The Periscope team replied to your feedback", body), + "feedback-reply", + ) diff --git a/periscope/src/backend/services/pipeline.py b/periscope/src/backend/services/pipeline.py new file mode 100644 index 0000000..c7b7a2e --- /dev/null +++ b/periscope/src/backend/services/pipeline.py @@ -0,0 +1,1519 @@ +"""Analysis job stages: BOM → extract → graph → review. + +Workspace and SSE broker live in ``job_workspace``. Extraction is +``datasheet_extract``. Review is ``validate_design_async`` (review_session). +This module must not import PinScope ``validate.py`` for datasheet loading. +""" + +from __future__ import annotations + +import asyncio +import json +import logging +import re +import subprocess +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any, Awaitable, Callable + +from backend.config import settings +from backend.periscopex.bom_summary import build_bom_summary +from backend.periscopex.constraints_lookup import build_constraints_map, load_datasheets +from backend.periscopex.derating import build_derating_table +from backend.periscopex.graph import build_graph +from backend.periscopex.models import ComponentType +from backend.periscopex.parsers import parse_bom, parse_netlist_any +from backend.periscopex.resolve_passives import SkippedItem, load_patterns, resolve_mpn +from backend.periscopex.taxonomy import SIMPLE_TYPES, type_for_ref +from backend.periscopex.utils import natural_sort_key, safe_mpn +from backend.services import admin_settings as settings_svc +from backend.services import datasheet_extract as extraction +from backend.services import job_workspace as _job_ws +from backend.services import projects as proj_svc +from backend.services.api_logs import ApiLogger, total_cost +from backend.services.billing_hook import InsufficientCredits, get_billing +from backend.services.cost_estimator import estimate_stage_cost_usd +from backend.services.datasheet_store import compute_md5_from_path, store_datasheet, store_datasheet_bytes +from backend.services.storage import StorageBackend +from backend.services.validation import validate_design_async + +log = logging.getLogger(__name__) + +EventBroker = _job_ws.EventBroker +PipelineWorkspace = _job_ws.PipelineWorkspace +broker = _job_ws.broker + +_CANCEL_POLL_S = 3.0 +_GIT_COMMIT: str | None = None + + +def set_broker(b) -> None: + _job_ws.set_broker(b) + globals()["broker"] = b + + +def _ic_descriptions(extracted_dir: Path) -> dict[str, str]: + out: dict[str, str] = {} + try: + for mpn, c in load_datasheets(extracted_dir).items(): + desc = c.package_info.description if c.package_info else None + if desc: + out[mpn] = desc + except Exception: + log.exception("IC descriptions failed for %s", extracted_dir) + return out + + +def _git_commit() -> str: + global _GIT_COMMIT + if _GIT_COMMIT is None: + try: + _GIT_COMMIT = subprocess.run( + ["git", "rev-parse", "--short", "HEAD"], + capture_output=True, text=True, timeout=5, + cwd=Path(__file__).resolve().parent, + ).stdout.strip() or "unknown" + except Exception: + log.exception("git SHA for traces unavailable") + _GIT_COMMIT = "unknown" + return _GIT_COMMIT + + +class CancelRequested(Exception): + """User cancel flag observed in the worker gate.""" + + +@dataclass +class PipelineContext: + storage: StorageBackend + user_id: str + project_id: str + ws: PipelineWorkspace + api_logger: ApiLogger + meta: Any + min_ver: str + skipped: list[SkippedItem] = field(default_factory=list) + ic_mpns: dict[str, list[str]] = field(default_factory=dict) + passive_mpns: dict[str, list[str]] = field(default_factory=dict) + passive_values: dict[str, str] = field(default_factory=dict) + simple_mpns: dict[str, list[str]] = field(default_factory=dict) + simple_mpn_types: dict[str, str] = field(default_factory=dict) + datasheet_urls: dict[str, str] = field(default_factory=dict) + lcsc_data: dict[str, dict] = field(default_factory=dict) + ref_col: str = "Reference" + mpn_col: str = "Manufacturer Part Number" + patterns: list = field(default_factory=list) + graph: Any | None = None + report: Any | None = None + paused: bool = False + pause_stage: str | None = None + pause_unit_id: str | None = None + pause_last_completed: str | None = None + credits_spent: float = 0.0 + completed_review_refs: set[str] = field(default_factory=set) + all_review_refs: list[str] = field(default_factory=list) + free: bool = False + resume: bool = False + + +def _cancel_gate_check(ctx: PipelineContext) -> None: + import time as _time + + last = getattr(ctx, "_last_cancel_poll", 0.0) + now = _time.monotonic() + if now - last < _CANCEL_POLL_S: + return + ctx._last_cancel_poll = now # type: ignore[attr-defined] + try: + meta = proj_svc.get_project(ctx.storage, ctx.user_id, ctx.project_id) + except Exception: + return + if meta is not None and meta.cancel_requested: + raise CancelRequested(f"cancel requested for {ctx.project_id}") + + +def _check_credit_gate( + ctx: PipelineContext, stage: str, unit_id: str, estimated_cost_usd: float, +) -> bool: + if ctx.paused: + return False + if ctx.free: + return True + billing = get_billing() + need = billing.credits_for_api_cost(estimated_cost_usd) + if need <= 0: + return True + if billing.get_balance(ctx.storage, ctx.user_id) < need: + ctx.paused = True + ctx.pause_stage = stage + ctx.pause_unit_id = unit_id + return False + return True + + +def _charge_for_logs(ctx: PipelineContext, before_count: int) -> None: + _cancel_gate_check(ctx) + new_entries = ctx.api_logger.entries[before_count:] + total_credits = sum(float(e.get("credits_charged") or 0) for e in new_entries) + if total_credits <= 0: + return + amount = round(total_credits, 4) + unit_id = new_entries[-1].get("identifier") if new_entries else None + stage = new_entries[-1].get("stage") if new_entries else None + billing = get_billing() + try: + billing.charge( + ctx.storage, ctx.user_id, amount, + reason="pipeline_charge", + run_id=ctx.project_id, + unit_id=f"{stage}:{unit_id}" if stage else None, + allow_overdraft=True, + ) + ctx.credits_spent += amount + broker.publish( + ctx.project_id, "credits_update", + { + "credits_spent": round(ctx.credits_spent, 4), + "balance_after": round(billing.get_balance(ctx.storage, ctx.user_id), 4), + "delta": round(amount, 4), + "stage": stage, + "unit_id": unit_id, + }, + ) + except InsufficientCredits: + pass + try: + async def _notify() -> None: + failure = await billing.maybe_auto_topup(ctx.storage, ctx.user_id) + if failure: + broker.publish(ctx.project_id, "auto_topup_failed", failure) + + asyncio.create_task(_notify()) + except Exception: + pass + + +def _charge_private_logger(ctx: PipelineContext, private: ApiLogger) -> None: + before = len(ctx.api_logger.entries) + ctx.api_logger.entries.extend(private.entries) + _charge_for_logs(ctx, before) + + +async def _paused_stage_publish(ctx: PipelineContext, stage: str, reason: str) -> None: + broker.publish( + ctx.project_id, "step_update", + {"stage": stage, "status": "paused", "detail": reason}, + ) + + +async def _resolve_lcsc_codes(ctx: PipelineContext | None, bom: dict[str, dict]) -> None: + """Fill empty MPN from the LCSC column only. No-op when catalogue is off.""" + if not settings.use_purple_parts: + return + from backend.services.purple_parts import is_lcsc_code, lookup_lcsc_batch + + todo: list[tuple[str, str]] = [] + for ref, info in bom.items(): + mpn = (info.get("mpn") or "").strip() + lcsc = (info.get("lcsc") or "").strip() + if not mpn and is_lcsc_code(lcsc): + todo.append((ref, lcsc)) + if not todo: + return + unique = sorted({c for _, c in todo}) + if ctx is not None: + broker.publish( + ctx.project_id, "step_update", + {"stage": "bom_parse", "status": "running", + "detail": f"Resolving {len(unique)} LCSC code(s) via purple-parts"}, + ) + resolved = await lookup_lcsc_batch(unique) + hits = 0 + for ref, code in todo: + part = resolved.get(code) + if part and part.get("mpn"): + mpn = part["mpn"] + bom[ref]["mpn"] = mpn + if ctx is not None: + ctx.lcsc_data.setdefault(mpn, part) + hits += 1 + log.info("LCSC backstop: %d/%d codes, %d refs", hits, len(unique), len(todo)) + + +async def _stage_bom_parse(ctx: PipelineContext) -> None: + broker.publish(ctx.project_id, "step_update", {"stage": "bom_parse", "status": "running"}) + col_map = ctx.meta.bom_columns or {} + ctx.ref_col = col_map.get("reference", "Reference") + ctx.mpn_col = col_map.get("mpn", "Manufacturer Part Number") + bom = parse_bom( + str(ctx.ws.local_path("uploads/bom.csv")), + reference_col=ctx.ref_col, mpn_col=ctx.mpn_col, + ) + await _resolve_lcsc_codes(ctx, bom) + for ref, info in sorted(bom.items()): + mpn = info.get("mpn") + url = (info.get("datasheet_url") or "").strip() + if mpn and url and mpn not in ctx.datasheet_urls: + ctx.datasheet_urls[mpn] = url + if not mpn: + continue + typ = type_for_ref(ref) + if typ == "ic": + ctx.ic_mpns.setdefault(mpn, []).append(ref) + elif typ == "passive": + ctx.passive_mpns.setdefault(mpn, []).append(ref) + val = (info.get("value") or "").strip() + if val and not ctx.passive_values.get(mpn): + ctx.passive_values[mpn] = val + elif typ and typ in SIMPLE_TYPES: + ctx.simple_mpns.setdefault(mpn, []).append(ref) + ctx.simple_mpn_types[mpn] = typ + proj_svc.update_project( + ctx.storage, ctx.user_id, ctx.project_id, + component_mpns={ + "ic": list(ctx.ic_mpns.keys()), + "passive": list(ctx.passive_mpns.keys()), + "simple": list(ctx.simple_mpns.keys()), + }, + ) + _, nets, _ = parse_netlist_any( + str(ctx.ws.netlist_local_path()), + known_refs=set(bom.keys()), + include_subdesigns=( + set(ctx.meta.netlist_subdesigns) + if ctx.meta.netlist_subdesigns is not None else None + ), + ) + broker.publish( + ctx.project_id, "step_update", + {"stage": "bom_parse", "status": "complete", + "detail": f"{len(bom)} refs, {len(ctx.ic_mpns)} ICs, " + f"{len(ctx.simple_mpns)} discrete/simple, {len(ctx.passive_mpns)} passives"}, + ) + from backend.services.email import send_pipeline_started_email + try: + await send_pipeline_started_email( + user_id=ctx.user_id, + project_name=ctx.meta.name, + project_id=ctx.project_id, + num_components=len(bom), + num_nets=len(nets), + num_ics=len(ctx.ic_mpns), + num_passives=len(ctx.passive_mpns), + num_simple=len(ctx.simple_mpns), + ) + except Exception: + pass + + +def _lcsc_id_for_mpn(ctx: PipelineContext, mpn: str) -> str | None: + payload = ctx.lcsc_data.get(mpn) or {} + code = payload.get("lcsc") or payload.get("lcsc_id") + if isinstance(code, str) and code.strip(): + return code.strip() + mapping = getattr(ctx.meta, "lcsc_to_mpn", None) or {} + for lcsc, resolved in mapping.items(): + if resolved == mpn: + return lcsc + return None + + +async def _ensure_local_datasheet( + ctx: PipelineContext, mpn: str, pdf_path: Path, *, stage: str = "ic_extraction", +) -> bool: + if pdf_path.is_file(): + return True + from backend.services.datasheet_finder import find_datasheet, find_local_pdf, mpn_query_variants + + alt = find_local_pdf(pdf_path.parent, mpn) + if alt is not None and alt.is_file(): + if alt.resolve() != pdf_path.resolve(): + pdf_path.parent.mkdir(parents=True, exist_ok=True) + pdf_path.write_bytes(alt.read_bytes()) + return True + for name in mpn_query_variants(mpn) or [mpn]: + lib = proj_svc.library_has_datasheet(ctx.storage, name) + if lib: + ctx.storage.download_to_local(lib, pdf_path) + return True + broker.publish( + ctx.project_id, "step_update", + {"stage": stage, "substep": mpn, "status": "running", "detail": "finding datasheet"}, + ) + hit = await find_datasheet( + mpn, lcsc_id=_lcsc_id_for_mpn(ctx, mpn), url_hint=ctx.datasheet_urls.get(mpn), + ) + if not hit.ok or not hit.pdf_bytes: + return False + pdf_path.parent.mkdir(parents=True, exist_ok=True) + pdf_path.write_bytes(hit.pdf_bytes) + try: + proj_svc.save_datasheet(ctx.storage, ctx.user_id, ctx.project_id, mpn, hit.pdf_bytes) + except Exception: + log.exception("project datasheet persist failed for %s", mpn) + try: + store_datasheet_bytes(ctx.storage, hit.pdf_bytes, mpn, extra_mpns=hit.alias_mpns) + except Exception: + log.exception("library datasheet persist failed for %s", mpn) + return True + + +def _prior_extract_ready(ctx: PipelineContext) -> bool: + graph = ctx.ws.local_path("design_graph.json") + extracted = ctx.ws.local_path("extracted") + return graph.is_file() and extracted.is_dir() and any(extracted.glob("*.json")) + + +async def _stage_ic_extraction(ctx: PipelineContext) -> None: + if ctx.resume and _prior_extract_ready(ctx): + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "status": "complete", "detail": "reusing previous extraction"}, + ) + return + extracted_dir = ctx.ws.local_path("extracted") + from backend.periscopex.layout_rules import needs_layout_rules_refresh + + layout_ver = settings.get_default_model_version() + cache: dict[str, tuple] = {} + n_new = 0 + for mpn in ctx.ic_mpns: + safe = safe_mpn(mpn) + json_path = extracted_dir / f"{safe}.json" + if json_path.is_file(): + existing = json.loads(json_path.read_text()) + if existing.get("pintable"): + ws_ver = existing.get("model_version", "0.0.0") + if ( + not settings_svc.version_is_stale(ws_ver, ctx.min_ver) + and not needs_layout_rules_refresh(existing, min_scan_version=layout_ver) + ): + cache[mpn] = ("workspace",) + continue + lib_key = proj_svc.library_has_extraction(ctx.storage, mpn, min_version=ctx.min_ver) + if lib_key: + try: + lib_payload = ctx.storage.read_json(lib_key) + except Exception: + lib_payload = {} + if not needs_layout_rules_refresh( + lib_payload if isinstance(lib_payload, dict) else {}, + min_scan_version=layout_ver, + ): + cache[mpn] = ("library", lib_key) + continue + n_new += 1 + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "status": "running", "total_new": n_new}, + ) + pending: list[tuple[str, str, Path, Path]] = [] + for mpn, _refs in ctx.ic_mpns.items(): + safe = safe_mpn(mpn) + json_path = extracted_dir / f"{safe}.json" + hit = cache.get(mpn) + if hit: + if hit[0] == "library": + ctx.storage.download_to_local(hit[1], json_path) + detail = "already extracted" if hit[0] == "workspace" else "from library" + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "substep": mpn, "status": "complete", "detail": detail}, + ) + continue + pdf_path = ctx.ws.local_path(f"uploads/datasheets/{safe}.pdf") + if not await _ensure_local_datasheet(ctx, mpn, pdf_path): + ctx.skipped.append(SkippedItem(mpn, "ic_extraction", "No datasheet found")) + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "substep": mpn, "status": "failed", "error": "No datasheet found"}, + ) + continue + pending.append((mpn, safe, json_path, pdf_path)) + + sem = asyncio.Semaphore(settings.ic_concurrency) + + async def _one(mpn: str, safe: str, json_path: Path, pdf_path: Path) -> None: + async with sem: + if ctx.paused: + return + if not _check_credit_gate(ctx, "ic_extraction", mpn, estimate_stage_cost_usd("ic_extraction")): + await _paused_stage_publish(ctx, "ic_extraction", "out of credits") + return + private = ApiLogger(free=ctx.api_logger.free) + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "substep": mpn, + "status": "running", "detail": "extracting pintable"}, + ) + await extraction.extract_pintable( + mpn, str(pdf_path), extracted_dir, + taxonomy_dir=ctx.ws.taxonomy_dir, api_logger=private, + ) + extracted_key = f"{ctx.ws.prefix}/extracted/{safe}.json" + ctx.storage.upload_from_local(json_path, extracted_key) + try: + from backend.periscopex.library_gate import should_promote_extraction + + payload = json.loads(json_path.read_text(encoding="utf-8")) + ok, reason = should_promote_extraction(payload) + if ok: + proj_svc.save_to_library(ctx.storage, extracted_key, "extracted", f"{safe}.json") + else: + log.warning("not promoting %s: %s", mpn, reason) + except Exception: + log.exception("library gate failed for %s — promoting", mpn) + proj_svc.save_to_library(ctx.storage, extracted_key, "extracted", f"{safe}.json") + store_datasheet(ctx.storage, pdf_path, mpn) + _charge_private_logger(ctx, private) + ctx.pause_last_completed = f"Extracted {mpn}" + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "substep": mpn, "status": "complete"}, + ) + except CancelRequested: + if private.entries: + ctx.api_logger.entries.extend(private.entries) + raise + except Exception as exc: + if private.entries: + ctx.api_logger.entries.extend(private.entries) + ctx.skipped.append(SkippedItem(mpn, "ic_extraction", str(exc))) + broker.publish( + ctx.project_id, "step_update", + {"stage": "ic_extraction", "substep": mpn, "status": "failed", "error": str(exc)}, + ) + + results = await asyncio.gather( + *(_one(*row) for row in pending), return_exceptions=True, + ) + for r in results: + if isinstance(r, (asyncio.CancelledError, CancelRequested)): + raise r + broker.publish(ctx.project_id, "step_update", {"stage": "ic_extraction", "status": "complete"}) + + +async def _stage_simple_extraction(ctx: PipelineContext) -> None: + if not ctx.simple_mpns: + return + if ctx.resume and _prior_extract_ready(ctx): + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "status": "complete", "detail": "reusing previous extraction"}, + ) + return + models_dir = ctx.ws.local_path("models") + cache: dict[str, tuple] = {} + n_new = 0 + for mpn in ctx.simple_mpns: + model_path = models_dir / f"{safe_mpn(mpn)}.json" + if model_path.is_file(): + cache[mpn] = ("workspace",) + continue + lib = proj_svc.library_has_model(ctx.storage, mpn) + if lib: + cache[mpn] = ("library", lib) + continue + n_new += 1 + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "status": "running", "total_new": n_new}, + ) + for mpn, _refs in ctx.simple_mpns.items(): + safe = safe_mpn(mpn) + model_path = models_dir / f"{safe}.json" + hit = cache.get(mpn) + if hit: + if hit[0] == "library": + ctx.storage.download_to_local(hit[1], model_path) + detail = "already extracted" if hit[0] == "workspace" else "from library" + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, "status": "complete", "detail": detail}, + ) + continue + pdf_path = ctx.ws.local_path(f"uploads/datasheets/{safe}.pdf") + if not await _ensure_local_datasheet(ctx, mpn, pdf_path, stage="simple_extraction"): + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "complete", "detail": "no datasheet (optional)"}, + ) + continue + if not _check_credit_gate(ctx, "simple_extraction", mpn, estimate_stage_cost_usd("simple_extraction")): + await _paused_stage_publish(ctx, "simple_extraction", "out of credits") + return + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "running", "detail": "extracting specs"}, + ) + before = len(ctx.api_logger.entries) + await extraction.extract_specs( + mpn, str(pdf_path), ctx.simple_mpn_types[mpn], models_dir, + taxonomy_dir=ctx.ws.taxonomy_dir, api_logger=ctx.api_logger, + ) + model_key = f"{ctx.ws.prefix}/models/{safe}.json" + ctx.storage.upload_from_local(model_path, model_key) + proj_svc.save_to_library(ctx.storage, model_key, "models", f"{safe}.json") + store_datasheet(ctx.storage, pdf_path, mpn) + _charge_for_logs(ctx, before) + ctx.pause_last_completed = f"Extracted {mpn}" + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, "status": "complete"}, + ) + except Exception as exc: + ctx.skipped.append(SkippedItem(mpn, "simple_extraction", str(exc))) + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, "status": "failed", "error": str(exc)}, + ) + if settings.use_digikey: + missing = [ + m for m in ctx.simple_mpns + if m not in cache and not (models_dir / f"{safe_mpn(m)}.json").is_file() + ] + if missing: + from backend.services.digikey import fetch_params + + for mpn in missing: + safe = safe_mpn(mpn) + model_path = models_dir / f"{safe}.json" + try: + lib = proj_svc.library_has_model(ctx.storage, mpn) + if lib: + ctx.storage.download_to_local(lib, model_path) + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "complete", "detail": "specs from library"}, + ) + continue + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "running", "detail": "auto-resolving via DigiKey"}, + ) + result = await fetch_params(mpn) + if not result.ok or not result.params: + raise RuntimeError(result.error or "No DigiKey parameters") + model = await extraction.auto_resolve_specs( + mpn=mpn, + digikey_params=result.params.parameters, + digikey_category=result.params.category, + digikey_description=result.params.description, + component_type=ctx.simple_mpn_types[mpn], + taxonomy_dir=ctx.ws.taxonomy_dir, + api_logger=ctx.api_logger, + ) + model_path.write_text(model.model_dump_json(indent=2) + "\n") + model_key = f"{ctx.ws.prefix}/models/{safe}.json" + ctx.storage.upload_from_local(model_path, model_key) + proj_svc.save_to_library(ctx.storage, model_key, "models", f"{safe}.json") + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "complete", "detail": "auto-resolved via DigiKey"}, + ) + except Exception as exc: + ctx.skipped.append(SkippedItem(mpn, "simple_digikey_resolve", str(exc))) + broker.publish( + ctx.project_id, "step_update", + {"stage": "simple_extraction", "substep": mpn, + "status": "failed", "error": str(exc)}, + ) + broker.publish(ctx.project_id, "step_update", {"stage": "simple_extraction", "status": "complete"}) + + +async def _catalog_resolve_unresolved_passives( + ctx: PipelineContext, still_unresolved: list[str], models_dir: Path, +) -> None: + if not still_unresolved: + return + from backend.services.digikey import fetch_params + + need = [m for m in still_unresolved if m not in ctx.lcsc_data] + if need and settings.use_purple_parts: + try: + from backend.services.purple_parts import lookup_mpn_batch + + parts = await lookup_mpn_batch(need) + for m, part in parts.items(): + if part and part.get("description"): + ctx.lcsc_data.setdefault(m, part) + except Exception: + log.warning("LCSC by-mpn backstop failed", exc_info=True) + + sem = asyncio.Semaphore(settings.ic_concurrency) + + async def _one(mpn: str) -> None: + async with sem: + if ctx.paused: + return + safe = safe_mpn(mpn) + model_path = models_dir / f"{safe}.json" + lib = proj_svc.library_has_passive_model(ctx.storage, mpn) + if lib: + ctx.storage.download_to_local(lib, model_path) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "complete", "detail": "specs from library"}, + ) + return + private = ApiLogger(free=ctx.api_logger.free) + model = None + via: str | None = None + first_error: str | None = None + llm_args: tuple[list[dict[str, str]], str, str, str] | None = None + try: + from backend.services.passive_from_distributor import ( + lcsc_payload_args, + specs_from_distributor, + specs_from_lcsc_payload, + ) + + async def _gate() -> bool: + if not _check_credit_gate( + ctx, "passive_extraction", mpn, + estimate_stage_cost_usd("digikey_resolve"), + ): + await _paused_stage_publish(ctx, "passive_extraction", "out of credits") + return False + return True + + lcsc = ctx.lcsc_data.get(mpn) + if lcsc and lcsc.get("description"): + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", "detail": "resolving from LCSC catalog"}, + ) + model = specs_from_lcsc_payload(mpn, lcsc) + if model is not None: + via = "lcsc" + else: + params, category, description = lcsc_payload_args(lcsc) + llm_args = (params, category, description, "lcsc") + if model is None and settings.use_digikey: + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", "detail": "resolving from DigiKey catalog"}, + ) + result = await fetch_params(mpn) + if result.ok and result.params: + model = specs_from_distributor( + mpn=mpn, + params=result.params.parameters, + category=result.params.category, + description=result.params.description, + ) + if model is not None: + via = "digikey" + else: + llm_args = ( + result.params.parameters, + result.params.category, + result.params.description, + "digikey", + ) + else: + first_error = result.error or "no DigiKey parameters" + except Exception as exc: + first_error = str(exc) + if model is None: + from backend.services.passive_from_mpn import specs_from_mpn + + model = specs_from_mpn(mpn) + if model is not None: + via = "mpn" + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", "detail": "decoded from MPN"}, + ) + if model is None and llm_args is not None: + if not await _gate(): + return + params, category, description, src = llm_args + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", "detail": f"auto-resolving via {src} (LLM)"}, + ) + model = await extraction.auto_resolve_specs( + mpn=mpn, digikey_params=params, + digikey_category=category, digikey_description=description, + component_type="passive", taxonomy_dir=ctx.ws.taxonomy_dir, + api_logger=private, + ) + if model is not None: + via = src + except Exception as exc: + first_error = str(exc) + if model is None: + from backend.services.passive_from_value import ( + is_placeholder_value, + specs_from_bom_value, + ) + + bom_value = ctx.passive_values.get(mpn, "").strip() + refs = ctx.passive_mpns.get(mpn, []) + pref = re.match(r"^[A-Za-z]+", refs[0]) if refs else None + prefix = pref.group(0).upper() if pref else "" + if bom_value and prefix in {"C", "R", "L", "FB"}: + model = specs_from_bom_value(mpn, bom_value, prefix) + if model is not None: + via = "value" + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", + "detail": f"parsed BOM value {bom_value!r}"}, + ) + elif not is_placeholder_value(bom_value): + if not await _gate(): + return + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "running", + "detail": f"resolving from BOM value {bom_value!r}"}, + ) + model = await extraction.resolve_from_value( + mpn=mpn, value=bom_value, ref_prefix=prefix, + component_type="passive", + taxonomy_dir=ctx.ws.taxonomy_dir, + api_logger=private, + ) + via = "value" + except Exception as exc: + first_error = first_error or str(exc) + if model is None: + err = first_error or "no LCSC/DigiKey hit and no usable BOM value" + if private.entries: + ctx.api_logger.entries.extend(private.entries) + ctx.skipped.append(SkippedItem(mpn, "passive_resolve", err)) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "failed", "error": err}, + ) + return + model_path.write_text(model.model_dump_json(indent=2) + "\n") + model_key = f"{ctx.ws.prefix}/models/{safe}.json" + ctx.storage.upload_from_local(model_path, model_key) + if via in ("digikey", "lcsc", "mpn"): + proj_svc.save_to_library(ctx.storage, model_key, "passives", f"{safe}.json") + _charge_private_logger(ctx, private) + ctx.pause_last_completed = f"Resolved {mpn}" + detail = { + "lcsc": "auto-resolved via LCSC", + "digikey": "auto-resolved via DigiKey", + "mpn": "decoded from MPN (saved to library)", + "value": "resolved from BOM value (not saved to library)", + }.get(via, "resolved") + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "complete", "detail": detail}, + ) + except CancelRequested: + if private.entries: + ctx.api_logger.entries.extend(private.entries) + raise + except Exception as exc: + if private.entries: + ctx.api_logger.entries.extend(private.entries) + ctx.skipped.append(SkippedItem(mpn, "passive_resolve", str(exc))) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "failed", "error": str(exc)}, + ) + + results = await asyncio.gather(*(_one(m) for m in still_unresolved), return_exceptions=True) + for r in results: + if isinstance(r, (asyncio.CancelledError, CancelRequested)): + raise r + + +async def _stage_passive_extraction(ctx: PipelineContext) -> None: + broker.publish(ctx.project_id, "step_update", {"stage": "passive_extraction", "status": "running"}) + patterns_dir = ctx.ws.local_path("patterns") + models_dir = ctx.ws.local_path("models") + for lib_key in proj_svc.list_library_patterns(ctx.storage): + dest = patterns_dir / lib_key.rsplit("/", 1)[-1] + if not dest.exists(): + ctx.storage.download_to_local(lib_key, dest) + ctx.patterns = load_patterns(str(patterns_dir)) if patterns_dir.is_dir() else [] + if ctx.resume: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "status": "complete", "detail": "reusing previous extraction"}, + ) + return + unresolved: dict[str, list[str]] = {} + for mpn, refs in ctx.passive_mpns.items(): + if resolve_mpn(mpn, ctx.patterns) is not None: + continue + safe = safe_mpn(mpn) + if (models_dir / f"{safe}.json").is_file(): + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "complete", "detail": "specs already resolved"}, + ) + continue + lib = proj_svc.library_has_passive_model(ctx.storage, mpn) + if lib: + dest = models_dir / f"{safe}.json" + if not dest.is_file(): + ctx.storage.download_to_local(lib, dest) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": mpn, + "status": "complete", "detail": "specs from library"}, + ) + continue + unresolved[mpn] = refs + if not unresolved: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "status": "complete", + "detail": "all passives already resolved"}, + ) + return + await _catalog_resolve_unresolved_passives(ctx, list(unresolved), models_dir) + unresolved = { + m: r for m, r in unresolved.items() + if not (models_dir / f"{safe_mpn(m)}.json").is_file() + } + if not unresolved: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "status": "complete", + "detail": "all passives resolved from catalog"}, + ) + return + ds_dir = ctx.ws.local_path("uploads/datasheets") + seen_lib: set[str] = set() + seen_hash: set[str] = set() + pdfs: list[Path] = [] + for mpn in unresolved: + pdf = ds_dir / f"{safe_mpn(mpn)}.pdf" + if not pdf.is_file(): + lib = proj_svc.library_has_datasheet(ctx.storage, mpn, patterns=ctx.patterns) + if lib: + if lib in seen_lib: + continue + seen_lib.add(lib) + ctx.storage.download_to_local(lib, pdf) + if pdf.is_file(): + h = compute_md5_from_path(pdf) + if h in seen_hash: + continue + seen_hash.add(h) + pdfs.append(pdf) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "status": "running", "total_new": len(pdfs)}, + ) + extracted_ds = {getattr(p, "datasheet_key", None) or "" for p in ctx.patterns} + extracted_ds.discard("") + for pdf_path in pdfs: + if not unresolved: + break + if pdf_path.stem not in {safe_mpn(m) for m in unresolved}: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "complete", "detail": "resolved by pattern"}, + ) + continue + blob = f"library/datasheets/blobs/{compute_md5_from_path(pdf_path)}.pdf" + if blob in extracted_ds: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "complete", "detail": "pattern already extracted from this PDF"}, + ) + continue + trigger = next((m for m in unresolved if safe_mpn(m) == pdf_path.stem), None) + if not _check_credit_gate( + ctx, "passive_extraction", pdf_path.stem, + estimate_stage_cost_usd("passive_pattern"), + ): + await _paused_stage_publish(ctx, "passive_extraction", "out of credits") + return + try: + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "running", "detail": "extracting pattern"}, + ) + before = len(ctx.api_logger.entries) + try: + out = await asyncio.wait_for( + extraction.extract_pattern( + str(pdf_path), list(unresolved.keys()), patterns_dir, + trigger_mpn=trigger, taxonomy_dir=ctx.ws.taxonomy_dir, + api_logger=ctx.api_logger, + ), + timeout=90, + ) + except asyncio.TimeoutError: + ctx.skipped.append(SkippedItem( + pdf_path.stem, "passive_extraction", "pattern extraction timed out (90s)", + )) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "failed", "error": "pattern extraction timed out"}, + ) + continue + if out: + blob_k = store_datasheet(ctx.storage, pdf_path, out.stem) + extracted_ds.add(blob_k) + data = json.loads(out.read_text()) + data["datasheet_key"] = blob_k + out.write_text(json.dumps(data, indent=2) + "\n") + rel = out.relative_to(ctx.ws.local_dir) + pattern_key = f"{ctx.ws.prefix}/{rel}" + ctx.storage.upload_from_local(out, pattern_key) + proj_svc.save_to_library(ctx.storage, pattern_key, "patterns", out.name) + ctx.patterns = load_patterns(str(patterns_dir)) + prev = len(unresolved) + unresolved = {m: r for m, r in unresolved.items() if resolve_mpn(m, ctx.patterns) is None} + _charge_for_logs(ctx, before) + ctx.pause_last_completed = f"Pattern from {pdf_path.stem}" + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "complete", + "detail": ( + f"pattern extracted, resolved {prev - len(unresolved)} MPNs" + if out else "pattern failed, MPN falls to DigiKey" + )}, + ) + except Exception as exc: + ctx.skipped.append(SkippedItem(pdf_path.stem, "passive_extraction", str(exc))) + broker.publish( + ctx.project_id, "step_update", + {"stage": "passive_extraction", "substep": pdf_path.stem, + "status": "failed", "error": str(exc)}, + ) + leftover = [m for m in unresolved if not (models_dir / f"{safe_mpn(m)}.json").is_file()] + await _catalog_resolve_unresolved_passives(ctx, leftover, models_dir) + broker.publish(ctx.project_id, "step_update", {"stage": "passive_extraction", "status": "complete"}) + + +def _write_layout_graph(ws: PipelineWorkspace, project_id: str) -> None: + pcb = ws.local_path("uploads/pcb.kicad_pcb") + if not pcb.is_file(): + return + try: + from backend.periscopex.parsers_kicad_pcb import parse_kicad_pcb + + layout = parse_kicad_pcb(pcb) + ws.local_path("layout_graph.json").write_text(layout.model_dump_json(indent=2) + "\n") + log.info("layout_graph footprints=%s nets=%s %s", len(layout.footprints), len(layout.nets), project_id) + except Exception: + log.exception("kicad_pcb parse failed — continuing without layout") + + +def _write_functional_groups(ws: PipelineWorkspace, graph) -> None: + try: + from backend.periscopex.functional_groups import build_functional_groups + + extracted = ws.local_path("extracted") + cmap = {} + if extracted.is_dir(): + cmap = build_constraints_map(load_datasheets(extracted)) + payload = build_functional_groups(graph, cmap).model_dump_json(indent=2) + "\n" + for name in ("functional_groups.json", "placement_plan.json"): + ws.local_path(name).write_text(payload) + ws._upload_file(name) + except Exception: + log.exception("functional_groups write failed") + + +def _write_impedance_nets(ws: PipelineWorkspace, graph) -> None: + path = ws.local_path("layout_graph.json") + if not path.is_file(): + return + try: + from backend.periscopex.impedance_traces import analyze_where_needed + from backend.periscopex.models import LayoutGraph + + layout = LayoutGraph.model_validate_json(path.read_text()) + report = analyze_where_needed(layout, graph) + ws.local_path("impedance_nets.json").write_text(json.dumps(report, indent=2) + "\n") + log.info("impedance_nets n=%s skipped=%s", len(report.get("nets") or []), report.get("skipped")) + except Exception: + log.exception("impedance analysis failed") + + +async def _stage_graph_build(ctx: PipelineContext) -> None: + broker.publish(ctx.project_id, "step_update", {"stage": "graph_build", "status": "running"}) + ctx.graph = build_graph( + str(ctx.ws.netlist_local_path()), + str(ctx.ws.local_path("uploads/bom.csv")), + str(ctx.ws.local_path("extracted")), + str(ctx.ws.local_path("patterns")), + str(ctx.ws.local_path("models")), + reference_col=ctx.ref_col, + mpn_col=ctx.mpn_col, + skipped=ctx.skipped, + include_subdesigns=( + set(ctx.meta.netlist_subdesigns) + if ctx.meta.netlist_subdesigns is not None else None + ), + pcb_path=ctx.ws.local_path("uploads/pcb.kicad_pcb"), + ) + ctx.ws.local_path("design_graph.json").write_text(ctx.graph.model_dump_json(indent=2) + "\n") + _write_layout_graph(ctx.ws, ctx.project_id) + _write_impedance_nets(ctx.ws, ctx.graph) + _write_functional_groups(ctx.ws, ctx.graph) + broker.publish( + ctx.project_id, "step_update", + {"stage": "graph_build", "status": "complete", + "detail": f"{len(ctx.graph.components)} components, {len(ctx.graph.nets)} nets"}, + ) + + +async def _stage_validation(ctx: PipelineContext) -> None: + ds_dir = ctx.ws.local_path("uploads/datasheets") + extracted_dir = ctx.ws.local_path("extracted") + graph_path = ctx.ws.local_path("design_graph.json") + ds_mpns: set[str] = set() + if ds_dir.is_dir(): + ds_mpns = {p.stem for p in ds_dir.glob("*.pdf")} + for comp in ctx.graph.components.values(): + if comp.mpn and comp.mpn not in ds_mpns: + if proj_svc.library_has_datasheet(ctx.storage, comp.mpn, patterns=ctx.patterns): + ds_mpns.add(comp.mpn) + rows = build_bom_summary(ctx.graph, datasheet_mpns=ds_mpns, descriptions=_ic_descriptions(extracted_dir)) + ctx.ws.local_path("bom_summary.json").write_text(json.dumps(rows, indent=2) + "\n") + ctx.ws.local_path("derating.json").write_text(json.dumps(build_derating_table(ctx.graph), indent=2) + "\n") + + mpns: list[str] = [] + seen: set[str] = set() + for mpn in list(ctx.ic_mpns): + if mpn not in seen: + seen.add(mpn) + mpns.append(mpn) + for comp in ctx.graph.components.values(): + if comp.component_type != ComponentType.IC: + continue + extra = (comp.mpn or "").strip() + if extra and extra not in seen: + seen.add(extra) + mpns.append(extra) + + async def _place(mpn: str) -> None: + pdf_path = ds_dir / f"{safe_mpn(mpn)}.pdf" + if pdf_path.is_file(): + return + if ctx.resume: + from backend.services.datasheet_finder import find_local_pdf, mpn_query_variants + + alt = find_local_pdf(ds_dir, mpn) + if alt is not None and alt.is_file() and alt.resolve() != pdf_path.resolve(): + pdf_path.write_bytes(alt.read_bytes()) + return + for name in mpn_query_variants(mpn) or [mpn]: + lib = proj_svc.library_has_datasheet(ctx.storage, name) + if lib: + ctx.storage.download_to_local(lib, pdf_path) + return + return + try: + await asyncio.wait_for( + _ensure_local_datasheet(ctx, mpn, pdf_path, stage="review"), + timeout=45, + ) + except TimeoutError: + log.warning("datasheet lookup timed out for %s", mpn) + except Exception: + log.exception("datasheet lookup failed for %s", mpn) + + await asyncio.gather(*(_place(m) for m in mpns)) + from backend.services.datasheet_finder import find_local_pdf + + planned: list[str] = [] + for ref, comp in ctx.graph.components.items(): + if comp.component_type != ComponentType.IC: + continue + mpn = (comp.mpn or "").strip() or (comp.value or "").strip() + if mpn and find_local_pdf(ds_dir, mpn) is not None: + planned.append(ref) + ctx.all_review_refs = sorted(planned, key=natural_sort_key) + broker.publish(ctx.project_id, "step_update", {"stage": "validation", "status": "running"}) + report_path = ctx.ws.local_path("report.json") + + async def on_progress(ref: str, turn: int, tool: str, detail: str): + if tool == "error": + broker.publish( + ctx.project_id, "step_update", + {"stage": "validation", "substep": ref, "status": "failed", "detail": detail}, + ) + return + done = tool in ("submit_review", "skipped") + broker.publish( + ctx.project_id, "step_update", + {"stage": "validation", "substep": ref, + "status": "complete" if done else "running", + "detail": detail if done else tool}, + ) + + async def on_ic_error(ref: str, exc: BaseException) -> None: + ctx.skipped.append(SkippedItem(ref, "validation", f"{type(exc).__name__}: {exc}")) + + async def before_ic(ref: str) -> bool: + cost = estimate_stage_cost_usd("review") + if settings.normalize_findings_enabled: + cost += estimate_stage_cost_usd("normalize") + if not _check_credit_gate(ctx, "validation", ref, cost): + await _paused_stage_publish(ctx, "validation", "out of credits") + return False + return True + + async def on_ic_done(ref: str, result: Any, private: ApiLogger | None = None) -> None: + if private is not None: + _charge_private_logger(ctx, private) + ctx.completed_review_refs.add(ref) + ctx.pause_last_completed = f"Reviewed {ref}" + + async def on_dedupe_done(private: ApiLogger | None = None) -> None: + if private is not None: + _charge_private_logger(ctx, private) + + skip_refs = set(ctx.completed_review_refs) + fp_path = ctx.ws.local_path("review_fingerprints.json") + current_fp: dict[str, str] = {} + try: + from backend.periscopex.review_fingerprint import graph_ic_fingerprints, skip_unchanged_ics + + cmap = build_constraints_map(load_datasheets(extracted_dir)) + current_fp = graph_ic_fingerprints(ctx.graph, cmap) + previous: dict[str, str] = {} + if fp_path.is_file(): + try: + previous = json.loads(fp_path.read_text()) + except json.JSONDecodeError: + previous = {} + if previous: + skip_refs = skip_unchanged_ics(skip_refs, previous, current_fp) + for ref in sorted(skip_refs): + broker.publish( + ctx.project_id, "step_update", + {"stage": "validation", "substep": ref, "status": "complete", + "detail": "unchanged since last review"}, + ) + except Exception: + log.exception("review fingerprints failed — reviewing kept refs") + + ctx.report = await validate_design_async( + str(graph_path), str(report_path), str(extracted_dir), + pdf_dir=str(ds_dir), + on_progress=on_progress, + api_logger=ctx.api_logger, + storage=ctx.storage, + skip_refs=skip_refs, + before_ic=before_ic, + on_ic_done=on_ic_done, + on_ic_error=on_ic_error, + on_dedupe_done=on_dedupe_done, + project_prefix=proj_svc.project_prefix(ctx.user_id, ctx.project_id), + run_meta={"git_commit": _git_commit()}, + ) + if current_fp: + fp_path.write_text(json.dumps(current_fp, indent=2) + "\n") + if ctx.paused: + return + broker.publish(ctx.project_id, "step_update", {"stage": "validation", "status": "complete"}) + + +@dataclass +class StageSpec: + stage_id: str + title: str + fn: Callable[[PipelineContext], Awaitable[None]] + + +PIPELINE_STAGES: list[StageSpec] = [ + StageSpec("bom_parse", "Parse BOM", _stage_bom_parse), + StageSpec("ic_extraction", "IC Datasheet Extraction", _stage_ic_extraction), + StageSpec("simple_extraction", "Component Specs Extraction", _stage_simple_extraction), + StageSpec("passive_extraction", "Passive Pattern Extraction", _stage_passive_extraction), + StageSpec("graph_build", "Build Design Graph", _stage_graph_build), + StageSpec("validation", "Review Design", _stage_validation), +] + + +async def run_pipeline( + storage: StorageBackend, user_id: str, project_id: str, + *, + resume: bool = False, + free: bool = False, +) -> None: + ctx: PipelineContext | None = None + api_logger: ApiLogger | None = None + try: + meta = proj_svc.get_project(storage, user_id, project_id) + if not meta: + raise ValueError(f"Project {project_id} not found") + if not proj_svc.get_bom_key(storage, user_id, project_id) or not proj_svc.get_netlist_key(storage, user_id, project_id): + proj_svc.update_project( + storage, user_id, project_id, status="error", + pipeline_state={"error": "Missing BOM or netlist"}, + ) + broker.publish(project_id, "pipeline_error", {"error": "Missing BOM or netlist"}) + return + api_logger = ApiLogger(free=free) + try: + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_QUEUED, proj_svc.STATUS_RUNNING}, + to_status=proj_svc.STATUS_RUNNING, + pause_checkpoint=None, pause_reason=None, cancel_requested=False, + ) + except proj_svc.StatusConflict: + log.warning("worker booted into non-queued project %s", project_id) + return + async with PipelineWorkspace(storage, user_id, project_id) as ws: + ctx = PipelineContext( + storage=storage, user_id=user_id, project_id=project_id, ws=ws, + api_logger=api_logger, meta=meta, + min_ver=settings_svc.get_min_model_version(storage), + free=free, resume=resume, + ) + if resume and meta.completed_review_refs: + ctx.completed_review_refs = set(meta.completed_review_refs) + ctx.credits_spent = float(meta.credits_spent or 0) + for spec in PIPELINE_STAGES: + await spec.fn(ctx) + try: + api_logger.flush(storage, user_id, project_id) + except Exception: + log.exception("api_logs flush failed") + if ctx.paused: + break + blob = api_logger.to_jsonl() + if blob: + ws.local_path("api_logs.jsonl").write_text(blob) + + skipped = [s.to_dict() for s in ctx.skipped] if ctx.skipped else None + project_cost = float(meta.total_cost_usd or 0) if ctx.free else ( + total_cost(api_logger.entries) + float(meta.total_cost_usd or 0) + ) + if ctx.paused: + pending = [r for r in ctx.all_review_refs if r not in ctx.completed_review_refs] + checkpoint = { + "paused_at": ctx.pause_unit_id, + "paused_stage": ctx.pause_stage, + "last_completed_label": ctx.pause_last_completed, + "completed_review_refs": sorted(ctx.completed_review_refs, key=natural_sort_key), + "pending_review_refs": pending, + } + proj_svc.update_project( + storage, user_id, project_id, + status="paused_insufficient_credits", + skipped_components=skipped or None, + total_cost_usd=project_cost, + credits_spent=ctx.credits_spent, + pause_checkpoint=checkpoint, + pause_reason="insufficient_credits", + completed_review_refs=sorted(ctx.completed_review_refs, key=natural_sort_key), + ) + broker.publish( + project_id, "pipeline_paused", + {"reason": "insufficient_credits", "last_completed": ctx.pause_last_completed, + "stage": ctx.pause_stage, "unit_id": ctx.pause_unit_id, + "completed_review_refs": sorted(ctx.completed_review_refs, key=natural_sort_key), + "pending_review_refs": pending}, + ) + from backend.services.cost_estimator import estimate_pipeline_cost + from backend.services.email import send_pipeline_paused_email + + try: + balance = get_billing().get_balance(storage, user_id) + needed = 0.0 + try: + remaining = estimate_pipeline_cost(storage, user_id, project_id) + needed = max(0.0, remaining.credits_low - max(0.0, balance)) + except Exception: + pass + await send_pipeline_paused_email( + user_id=user_id, project_name=meta.name, project_id=project_id, + last_completed=ctx.pause_last_completed, stage=ctx.pause_stage, + balance=balance, credits_needed_low=needed, + ) + except Exception: + pass + return + summary = ctx.report.summary if ctx.report else {} + proj_svc.update_project( + storage, user_id, project_id, + status="complete", summary=summary, + skipped_components=skipped or None, + total_cost_usd=project_cost, credits_spent=ctx.credits_spent, + pause_checkpoint=None, pause_reason=None, + completed_review_refs=sorted(ctx.completed_review_refs), + ) + broker.publish(project_id, "pipeline_complete", {"summary": summary, "skipped": skipped or []}) + from backend.services.email import send_report_ready_email + + try: + await send_report_ready_email( + user_id=user_id, project_name=meta.name, project_id=project_id, + summary=summary, total_cost_usd=project_cost, + ) + except Exception: + pass + except (asyncio.CancelledError, CancelRequested): + try: + extra: dict = { + "pipeline_state": {"error": "Pipeline cancelled by user"}, + "cancel_requested": False, + } + if ctx is not None: + extra["completed_review_refs"] = sorted(ctx.completed_review_refs, key=natural_sort_key) + extra["skipped_components"] = [s.to_dict() for s in ctx.skipped] or None + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED}, + to_status=proj_svc.STATUS_CANCELLED, + **extra, + ) + except proj_svc.StatusConflict: + pass + broker.publish(project_id, "pipeline_cancelled", {"error": "Pipeline cancelled by user"}) + try: + if api_logger is not None: + api_logger.flush(storage, user_id, project_id) + except Exception: + pass + except Exception as exc: + log.exception("pipeline crashed for %s", project_id) + try: + extra = {"pipeline_state": {"error": str(exc)}, "cancel_requested": False} + if ctx is not None: + extra["completed_review_refs"] = sorted(ctx.completed_review_refs, key=natural_sort_key) + extra["skipped_components"] = [s.to_dict() for s in ctx.skipped] or None + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED}, + to_status=proj_svc.STATUS_ERROR, + **extra, + ) + except proj_svc.StatusConflict: + pass + broker.publish(project_id, "pipeline_error", {"error": str(exc)}) + try: + if api_logger is not None: + api_logger.flush(storage, user_id, project_id) + except Exception: + pass + + +async def run_regen_pipeline( + storage: StorageBackend, user_id: str, project_id: str, stages: list[str], +) -> None: + try: + meta = proj_svc.get_project(storage, user_id, project_id) + if not meta: + raise ValueError(f"Project {project_id} not found") + api_logger = ApiLogger(free=True) + try: + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_QUEUED, proj_svc.STATUS_RUNNING}, + to_status=proj_svc.STATUS_RUNNING, cancel_requested=False, + ) + except proj_svc.StatusConflict: + log.warning("regen worker booted into non-queued project %s", project_id) + return + async with PipelineWorkspace(storage, user_id, project_id) as ws: + col = meta.bom_columns or {} + ref_col = col.get("reference", "Reference") + mpn_col = col.get("mpn", "Manufacturer Part Number") + broker.publish(project_id, "step_update", {"stage": "graph_build", "status": "running"}) + graph = build_graph( + str(ws.netlist_local_path()), + str(ws.local_path("uploads/bom.csv")), + str(ws.local_path("extracted")), + str(ws.local_path("patterns")), + str(ws.local_path("models")), + reference_col=ref_col, mpn_col=mpn_col, + include_subdesigns=( + set(meta.netlist_subdesigns) + if meta.netlist_subdesigns is not None else None + ), + pcb_path=ws.local_path("uploads/pcb.kicad_pcb"), + ) + ws.local_path("design_graph.json").write_text(graph.model_dump_json(indent=2) + "\n") + _write_layout_graph(ws, project_id) + _write_impedance_nets(ws, graph) + _write_functional_groups(ws, graph) + broker.publish( + project_id, "step_update", + {"stage": "graph_build", "status": "complete", + "detail": f"{len(graph.components)} components, {len(graph.nets)} nets"}, + ) + patterns_dir = ws.local_path("patterns") + patterns = load_patterns(str(patterns_dir)) if patterns_dir.is_dir() else [] + ds_dir = ws.local_path("uploads/datasheets") + ds_mpns: set[str] = {p.stem for p in ds_dir.glob("*.pdf")} if ds_dir.is_dir() else set() + for comp in graph.components.values(): + if comp.mpn and comp.mpn not in ds_mpns: + if proj_svc.library_has_datasheet(storage, comp.mpn, patterns=patterns): + ds_mpns.add(comp.mpn) + rows = build_bom_summary( + graph, datasheet_mpns=ds_mpns, + descriptions=_ic_descriptions(ws.local_path("extracted")), + ) + ws.local_path("bom_summary.json").write_text(json.dumps(rows, indent=2) + "\n") + if "derating" in stages: + ws.local_path("derating.json").write_text( + json.dumps(build_derating_table(graph), indent=2) + "\n", + ) + blob = api_logger.to_jsonl() + if blob: + ws.local_path("api_logs.jsonl").write_text(blob) + proj_svc.update_project(storage, user_id, project_id, status="complete") + broker.publish( + project_id, "pipeline_complete", + {"summary": meta.summary or {}, "regen_stages": stages}, + ) + except (asyncio.CancelledError, CancelRequested): + try: + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED}, + to_status=proj_svc.STATUS_CANCELLED, + pipeline_state={"error": "Regen cancelled"}, cancel_requested=False, + ) + except proj_svc.StatusConflict: + pass + broker.publish(project_id, "pipeline_cancelled", {"error": "Regen cancelled"}) + except Exception as exc: + log.exception("regen crashed for %s", project_id) + try: + proj_svc.transition_status( + storage, user_id, project_id, + from_status={proj_svc.STATUS_RUNNING, proj_svc.STATUS_QUEUED}, + to_status=proj_svc.STATUS_ERROR, + pipeline_state={"error": str(exc)}, cancel_requested=False, + ) + except proj_svc.StatusConflict: + pass + broker.publish(project_id, "pipeline_error", {"error": str(exc)}) diff --git a/periscope/src/backend/services/projects.py b/periscope/src/backend/services/projects.py new file mode 100644 index 0000000..225eb85 --- /dev/null +++ b/periscope/src/backend/services/projects.py @@ -0,0 +1,950 @@ +"""Project metadata and library keys on StorageBackend. + +Prefix is ``users/{owner}/projects/{id}/``. Writes follow the owner used to +read the object, not a stale ``user_id`` field (Emmaforo JWT vs local JSON). +""" + +from __future__ import annotations + +import logging +import uuid +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +from pydantic import AliasChoices, BaseModel, Field + +from backend.periscopex.utils import natural_sort_key, safe_mpn +from backend.services.storage import StaleGeneration, StorageBackend + +log = logging.getLogger(__name__) + + +class ProjectNotFound(Exception): + def __init__(self, project_id: str): + self.project_id = project_id + super().__init__(f"Project {project_id} not found") + + +STATUS_DRAFT = "draft" +STATUS_QUEUED = "queued" +STATUS_RUNNING = "running" +STATUS_COMPLETE = "complete" +STATUS_ERROR = "error" +STATUS_CANCELLED = "cancelled" +STATUS_PAUSED = "paused_insufficient_credits" + +TERMINAL_STATUSES = frozenset({ + STATUS_COMPLETE, STATUS_ERROR, STATUS_CANCELLED, STATUS_PAUSED, +}) + + +class StatusConflict(Exception): + """Status precondition failed, or generation retries exhausted.""" + + +class ProjectMeta(BaseModel): + id: str + name: str + user_id: str = "" + status: str = "draft" + created: str = "" + updated: str = "" + has_bom: bool = False + has_netlist: bool = False + has_pcb: bool = False + netlist_format: str | None = None + netlist_subdesigns: list[str] | None = None + datasheet_count: int = 0 + summary: dict[str, int] | None = None + component_mpns: dict[str, list[str]] | None = None + bom_columns: dict[str, str] | None = None + lcsc_to_mpn: dict[str, str] | None = None + lcsc_payloads: dict[str, dict] | None = None + skipped_components: list[dict[str, str]] | None = None + pipeline_state: dict[str, Any] | None = None + total_cost_usd: float | None = None + collaborators: list[str] = [] + tier: str = "demo" + credits_spent: float = 0.0 + estimate: dict[str, Any] | None = None + pause_checkpoint: dict[str, Any] | None = None + pause_reason: str | None = None + completed_review_refs: list[str] = [] + periscope_version: str | None = Field( + default=None, + validation_alias=AliasChoices("periscope_version", "pinscope_version"), + ) + execution_name: str | None = None + queued_at: str | None = None + cancel_requested: bool = False + placement_status: str = "draft" + placement_state: dict[str, Any] | None = None + placement_execution_name: str | None = None + placement_cancel_requested: bool = False + pcb_status: str = "draft" + pcb_state: dict[str, Any] | None = None + pcb_execution_name: str | None = None + pcb_cancel_requested: bool = False + + +def _project_prefix(user_id: str, project_id: str) -> str: + return f"users/{user_id}/projects/{project_id}" + + +def project_prefix(user_id: str, project_id: str) -> str: + return _project_prefix(user_id, project_id) + + +def _meta_key(user_id: str, project_id: str) -> str: + return f"{_project_prefix(user_id, project_id)}/project.json" + + +def _now() -> str: + return datetime.now(timezone.utc).isoformat() + + +def _read_meta(storage: StorageBackend, user_id: str, project_id: str) -> ProjectMeta: + return ProjectMeta.model_validate(storage.read_json(_meta_key(user_id, project_id))) + + +def _read_meta_with_generation( + storage: StorageBackend, user_id: str, project_id: str, +) -> tuple[ProjectMeta, int]: + data, gen = storage.read_json_with_generation(_meta_key(user_id, project_id)) + return ProjectMeta.model_validate(data), gen + + +def _write_meta( + storage: StorageBackend, + meta: ProjectMeta, + *, + owner_user_id: str | None = None, +) -> None: + uid = owner_user_id or meta.user_id + meta.updated = _now() + storage.write_json(_meta_key(uid, meta.id), meta.model_dump()) + + +def get_project( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta | None: + key = _meta_key(user_id, project_id) + if not storage.exists(key): + return None + return _read_meta(storage, user_id, project_id) + + +def completed_review_refs_for_retry( + storage: StorageBackend, user_id: str, project_id: str, +) -> list[str]: + meta = get_project(storage, user_id, project_id) + if not meta: + return [] + failed: set[str] = set() + for item in meta.skipped_components or []: + stage = item.get("stage") or "" + ident = (item.get("identifier") or "").strip() + if ident and stage in ("validation", "review"): + failed.add(ident) + report_key = f"{_project_prefix(user_id, project_id)}/report.json" + if storage.exists(report_key): + try: + report = storage.read_json(report_key) + except Exception: + report = {} + for ref in (report.get("review_errors") or {}): + if ref: + failed.add(str(ref)) + kept = [r for r in (meta.completed_review_refs or []) if r and r not in failed] + return sorted(kept, key=natural_sort_key) + + +def transition_status( + storage: StorageBackend, + user_id: str, + project_id: str, + *, + from_status: str | set[str] | frozenset[str], + to_status: str, + **fields: Any, +) -> ProjectMeta: + allowed = frozenset({from_status}) if isinstance(from_status, str) else frozenset(from_status) + last: Exception | None = None + for _ in range(5): + meta, gen = _read_meta_with_generation(storage, user_id, project_id) + if meta.status not in allowed: + raise StatusConflict( + f"project {project_id} is in status {meta.status!r}; " + f"expected one of {sorted(allowed)} for transition to {to_status!r}" + ) + meta.status = to_status + for k, v in fields.items(): + setattr(meta, k, v) + meta.updated = _now() + try: + storage.write_json_if_match( + _meta_key(user_id, project_id), meta.model_dump(), gen, + ) + return meta + except StaleGeneration as exc: + last = exc + continue + raise StatusConflict( + f"project {project_id}: lost optimistic-concurrency race after retries" + ) from last + + +def request_cancel( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta: + return update_project(storage, user_id, project_id, cancel_requested=True) + + +def mark_stale_running( + storage: StorageBackend, user_id: str, project_id: str, error: str, +) -> ProjectMeta | None: + try: + return transition_status( + storage, user_id, project_id, + from_status={STATUS_RUNNING, STATUS_QUEUED}, + to_status=STATUS_ERROR, + pipeline_state={"error": error}, + cancel_requested=False, + ) + except StatusConflict: + return None + + +def _last_event(storage: StorageBackend, prefix: str) -> dict | None: + events = f"{prefix}/events/" + try: + keys = sorted( + k for k in storage.list_prefix(events) + if k.endswith(".json") and "/events/" in k + ) + if keys: + return storage.read_json(keys[-1]) + except Exception: + return None + return None + + +def heal_if_pipeline_finished( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta | None: + meta = get_project(storage, user_id, project_id) + if meta is None or meta.status not in (STATUS_RUNNING, STATUS_QUEUED): + return None + prefix = _project_prefix(user_id, project_id) + last = _last_event(storage, prefix) + if (last or {}).get("event") == "pipeline_complete": + summary = (last.get("data") or {}).get("summary") + try: + return transition_status( + storage, user_id, project_id, + from_status={STATUS_RUNNING, STATUS_QUEUED}, + to_status=STATUS_COMPLETE, + summary=summary if isinstance(summary, dict) else meta.summary, + cancel_requested=False, + pipeline_state=None, + ) + except StatusConflict: + return None + from backend.services import job_runner + + exec_name = meta.execution_name or f"local/projects/{project_id}" + try: + state = job_runner.get_execution_state(exec_name) + except Exception: + state = "unknown" + if state in ("pending", "running"): + return None + if storage.exists(f"{prefix}/report.json"): + try: + return transition_status( + storage, user_id, project_id, + from_status={STATUS_RUNNING, STATUS_QUEUED}, + to_status=STATUS_COMPLETE, + cancel_requested=False, + pipeline_state=None, + ) + except StatusConflict: + return None + return mark_stale_running( + storage, user_id, project_id, + f"Analysis worker terminated ({state})", + ) + + +def heal_if_placement_stuck( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta | None: + meta = get_project(storage, user_id, project_id) + if meta is None: + return None + pst = meta.placement_status or "draft" + if pst not in ("queued", "running"): + return None + prefix = _project_prefix(user_id, project_id) + last = _last_event(storage, prefix) + if (last or {}).get("event") == "placement_complete": + data = (last or {}).get("data") or {} + return update_project( + storage, user_id, project_id, + placement_status="complete", + placement_cancel_requested=False, + placement_state={"domains": data.get("domains"), "groups": data.get("groups")}, + ) + from backend.services import job_runner + + exec_name = meta.placement_execution_name or f"local/placement/{project_id}" + try: + state = job_runner.get_execution_state(exec_name) + except Exception: + state = "unknown" + if state in ("pending", "running"): + return None + has_plan = ( + storage.exists(f"{prefix}/placement_plan.json") + or storage.exists(f"{prefix}/functional_groups.json") + ) + if has_plan: + return update_project( + storage, user_id, project_id, + placement_status="complete", + placement_cancel_requested=False, + placement_state=meta.placement_state, + ) + return update_project( + storage, user_id, project_id, + placement_status="error", + placement_cancel_requested=False, + placement_state={"error": f"Placement worker terminated ({state})"}, + ) + + +def heal_if_pcb_stuck( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta | None: + meta = get_project(storage, user_id, project_id) + if meta is None: + return None + pst = meta.pcb_status or "draft" + if pst not in ("queued", "running"): + return None + prefix = _project_prefix(user_id, project_id) + last = _last_event(storage, prefix) + if (last or {}).get("event") == "pcb_complete": + data = (last or {}).get("data") or {} + return update_project( + storage, user_id, project_id, + pcb_status="complete", + pcb_cancel_requested=False, + pcb_state={ + "findings": data.get("findings"), + "domains": data.get("domains"), + "groups": data.get("groups"), + }, + ) + from backend.services import job_runner + + exec_name = meta.pcb_execution_name or f"local/pcb/{project_id}" + try: + state = job_runner.get_execution_state(exec_name) + except Exception: + state = "unknown" + if state in ("pending", "running"): + return None + if storage.exists(f"{prefix}/pcb_report.json"): + return update_project( + storage, user_id, project_id, + pcb_status="complete", + pcb_cancel_requested=False, + pcb_state=meta.pcb_state, + ) + return update_project( + storage, user_id, project_id, + pcb_status="error", + pcb_cancel_requested=False, + pcb_state={"error": f"PCB worker terminated ({state})"}, + ) + + +def create_project(storage: StorageBackend, user_id: str, name: str) -> ProjectMeta: + meta = ProjectMeta( + id=uuid.uuid4().hex[:12], + name=name, + user_id=user_id, + created=_now(), + ) + _write_meta(storage, meta) + return meta + + +def list_projects(storage: StorageBackend, user_id: str) -> list[ProjectMeta]: + prefix = f"users/{user_id}/projects/" + out: list[ProjectMeta] = [] + for entry in storage.list_prefix(prefix): + key = f"{entry}/project.json" if not entry.endswith("/project.json") else entry + if storage.exists(key): + out.append(ProjectMeta.model_validate(storage.read_json(key))) + return out + + +def update_project( + storage: StorageBackend, user_id: str, project_id: str, **fields: Any, +) -> ProjectMeta: + if not storage.exists(_meta_key(user_id, project_id)): + raise ProjectNotFound(project_id) + meta = _read_meta(storage, user_id, project_id) + for k, v in fields.items(): + setattr(meta, k, v) + _write_meta(storage, meta, owner_user_id=user_id) + return meta + + +def _shared_ref_key(user_id: str, project_id: str) -> str: + return f"users/{user_id}/shared/{project_id}.json" + + +def delete_project(storage: StorageBackend, user_id: str, project_id: str) -> bool: + if not storage.exists(_meta_key(user_id, project_id)): + return False + meta = _read_meta(storage, user_id, project_id) + for collab in meta.collaborators: + ref = _shared_ref_key(collab, project_id) + if storage.exists(ref): + storage.delete_key(ref) + storage.delete_prefix(_project_prefix(user_id, project_id)) + return True + + +_DERIVED = ( + "design_graph.json", + "bom_summary.json", + "derating.json", + "report.json", + "periscope-findings.json", + "review_fingerprints.json", + "api_logs.jsonl", + "graph_voltage_updates.json", +) + + +def clear_project_extractions( + storage: StorageBackend, user_id: str, project_id: str, +) -> None: + prefix = _project_prefix(user_id, project_id) + for sub in ("extracted", "patterns", "models"): + storage.delete_prefix(f"{prefix}/{sub}") + for name in _DERIVED: + key = f"{prefix}/{name}" + if storage.exists(key): + storage.delete_key(key) + update_project( + storage, user_id, project_id, + summary=None, + skipped_components=None, + pipeline_state=None, + pause_checkpoint=None, + pause_reason=None, + completed_review_refs=[], + ) + + +def reopen_project( + storage: StorageBackend, user_id: str, project_id: str, +) -> ProjectMeta: + prefix = _project_prefix(user_id, project_id) + for name in _DERIVED: + key = f"{prefix}/{name}" + if storage.exists(key): + storage.delete_key(key) + return update_project( + storage, user_id, project_id, + status="draft", + summary=None, + skipped_components=None, + pipeline_state=None, + pause_checkpoint=None, + pause_reason=None, + completed_review_refs=[], + ) + + +def list_project_datasheets( + storage: StorageBackend, user_id: str, project_id: str, +) -> list[str]: + prefix = f"{_project_prefix(user_id, project_id)}/uploads/datasheets/" + return [ + key.rsplit("/", 1)[-1][:-4] + for key in storage.list_prefix(prefix) + if key.endswith(".pdf") + ] + + +def resolve_project_access( + storage: StorageBackend, caller_user_id: str, project_id: str, +) -> tuple[str, ProjectMeta] | None: + meta = get_project(storage, caller_user_id, project_id) + if meta is not None: + return caller_user_id, meta + ref_key = _shared_ref_key(caller_user_id, project_id) + if not storage.exists(ref_key): + return None + ref = storage.read_json(ref_key) + owner = ref.get("owner_user_id") + if not owner: + return None + meta = get_project(storage, owner, project_id) + if meta is None: + return None + if caller_user_id not in meta.collaborators: + storage.delete_key(ref_key) + return None + return owner, meta + + +def find_project_any_user( + storage: StorageBackend, project_id: str, +) -> tuple[str, ProjectMeta] | None: + seen: set[str] = set() + for entry in storage.list_prefix("users/"): + parts = entry.split("/") + if len(parts) < 2: + continue + uid = parts[1] + if uid in seen: + continue + seen.add(uid) + meta = get_project(storage, uid, project_id) + if meta is not None: + return uid, meta + return None + + +def add_collaborator( + storage: StorageBackend, owner_user_id: str, project_id: str, collaborator_user_id: str, +) -> ProjectMeta: + meta = _read_meta(storage, owner_user_id, project_id) + if collaborator_user_id not in meta.collaborators: + meta.collaborators.append(collaborator_user_id) + _write_meta(storage, meta, owner_user_id=owner_user_id) + storage.write_json( + _shared_ref_key(collaborator_user_id, project_id), + {"owner_user_id": owner_user_id}, + ) + return meta + + +def remove_collaborator( + storage: StorageBackend, owner_user_id: str, project_id: str, collaborator_user_id: str, +) -> ProjectMeta: + meta = _read_meta(storage, owner_user_id, project_id) + meta.collaborators = [c for c in meta.collaborators if c != collaborator_user_id] + _write_meta(storage, meta, owner_user_id=owner_user_id) + ref = _shared_ref_key(collaborator_user_id, project_id) + if storage.exists(ref): + storage.delete_key(ref) + return meta + + +def transfer_ownership( + storage: StorageBackend, + current_owner_user_id: str, + project_id: str, + new_owner_user_id: str, +) -> ProjectMeta: + meta = _read_meta(storage, current_owner_user_id, project_id) + if new_owner_user_id == current_owner_user_id: + raise ValueError("target user is already the owner") + if new_owner_user_id not in meta.collaborators: + raise ValueError("target user must currently be a collaborator") + collabs = [c for c in meta.collaborators if c != new_owner_user_id] + if current_owner_user_id not in collabs: + collabs.append(current_owner_user_id) + meta.user_id = new_owner_user_id + meta.collaborators = collabs + meta.updated = _now() + old_p = _project_prefix(current_owner_user_id, project_id) + new_p = _project_prefix(new_owner_user_id, project_id) + for old_key in storage.list_recursive(old_p): + rel = old_key[len(old_p):].lstrip("/") + storage.copy_object(old_key, f"{new_p}/{rel}") + storage.write_json(_meta_key(new_owner_user_id, project_id), meta.model_dump()) + storage.delete_prefix(old_p) + new_ref = _shared_ref_key(new_owner_user_id, project_id) + if storage.exists(new_ref): + storage.delete_key(new_ref) + storage.write_json( + _shared_ref_key(current_owner_user_id, project_id), + {"owner_user_id": new_owner_user_id}, + ) + for cid in collabs: + if cid == current_owner_user_id: + continue + storage.write_json( + _shared_ref_key(cid, project_id), + {"owner_user_id": new_owner_user_id}, + ) + return meta + + +def list_shared_projects(storage: StorageBackend, user_id: str) -> list[ProjectMeta]: + prefix = f"users/{user_id}/shared/" + out: list[ProjectMeta] = [] + for entry in storage.list_prefix(prefix): + if not entry.endswith(".json"): + continue + try: + ref = storage.read_json(entry) + owner = ref.get("owner_user_id") + if not owner: + continue + pid = entry.rsplit("/", 1)[-1].replace(".json", "") + meta = get_project(storage, owner, pid) + if meta and user_id in meta.collaborators: + out.append(meta) + except Exception: + continue + return out + + +def save_bom( + storage: StorageBackend, user_id: str, project_id: str, data: bytes, +) -> str: + key = f"{_project_prefix(user_id, project_id)}/uploads/bom.csv" + storage.write_bytes(key, data) + update_project(storage, user_id, project_id, has_bom=True) + return key + + +_NETLIST_EXT = { + "pads": "asc", + "edif": "edn", + "kicad_xml": "xml", + "kicad_sexp": "kicad_net", + "kicad_sch": "kicad_sch", +} + + +def _netlist_key(user_id: str, project_id: str, fmt: str) -> str: + return f"{_project_prefix(user_id, project_id)}/uploads/netlist.{_NETLIST_EXT.get(fmt, 'asc')}" + + +def save_netlist( + storage: StorageBackend, + user_id: str, + project_id: str, + data: bytes, + *, + fmt: str = "pads", +) -> str: + key = _netlist_key(user_id, project_id, fmt) + storage.write_bytes(key, data) + for other in _NETLIST_EXT: + if other == fmt: + continue + other_key = _netlist_key(user_id, project_id, other) + if storage.exists(other_key): + storage.delete_key(other_key) + update_project( + storage, user_id, project_id, + has_netlist=True, netlist_format=fmt, netlist_subdesigns=None, + ) + return key + + +def clear_companion_sheets( + storage: StorageBackend, user_id: str, project_id: str, +) -> None: + prefix = f"{_project_prefix(user_id, project_id)}/uploads/" + for key in storage.list_recursive(prefix): + rel = key[len(prefix):] + if rel.endswith(".kicad_sch") and rel != "netlist.kicad_sch": + storage.delete_key(key) + + +def save_companion_sheets( + storage: StorageBackend, + user_id: str, + project_id: str, + root: Path, + extras: list[Path], +) -> None: + clear_companion_sheets(storage, user_id, project_id) + parent = root.parent + prefix = f"{_project_prefix(user_id, project_id)}/uploads/" + for extra in extras: + rel = extra.relative_to(parent).as_posix() + if rel == "netlist.kicad_sch": + continue + storage.write_bytes(prefix + rel, extra.read_bytes()) + + +def save_pcb( + storage: StorageBackend, user_id: str, project_id: str, data: bytes, +) -> str: + key = f"{_project_prefix(user_id, project_id)}/uploads/pcb.kicad_pcb" + storage.write_bytes(key, data) + update_project(storage, user_id, project_id, has_pcb=True) + return key + + +def remember_datasheet( + storage: StorageBackend, mpn: str, data: bytes, extra_mpns: list[str] | None = None, +) -> None: + try: + from backend.services.datasheet_store import store_datasheet_bytes + + store_datasheet_bytes(storage, data, mpn, extra_mpns=extra_mpns) + except Exception: + log.exception("library datasheet store failed for %s", mpn) + + +def save_datasheet( + storage: StorageBackend, user_id: str, project_id: str, mpn: str, data: bytes, +) -> str: + safe = safe_mpn(mpn) + key = f"{_project_prefix(user_id, project_id)}/uploads/datasheets/{safe}.pdf" + storage.write_bytes(key, data) + remember_datasheet(storage, mpn, data) + ds = f"{_project_prefix(user_id, project_id)}/uploads/datasheets/" + count = sum(1 for k in storage.list_prefix(ds) if k.endswith(".pdf")) + update_project(storage, user_id, project_id, datasheet_count=count) + return key + + +def get_bom_key( + storage: StorageBackend, user_id: str, project_id: str, +) -> str | None: + key = f"{_project_prefix(user_id, project_id)}/uploads/bom.csv" + return key if storage.exists(key) else None + + +def get_netlist_key( + storage: StorageBackend, user_id: str, project_id: str, +) -> str | None: + for fmt in _NETLIST_EXT: + key = _netlist_key(user_id, project_id, fmt) + if storage.exists(key): + return key + return None + + +def get_datasheet_key( + storage: StorageBackend, user_id: str, project_id: str, mpn: str, +) -> str | None: + key = f"{_project_prefix(user_id, project_id)}/uploads/datasheets/{safe_mpn(mpn)}.pdf" + return key if storage.exists(key) else None + + +def library_has_extraction( + storage: StorageBackend, mpn: str, min_version: str | None = None, +) -> str | None: + key = f"library/extracted/{safe_mpn(mpn)}.json" + if not storage.exists(key): + return None + data = storage.read_json(key) + if not data.get("pintable"): + return None + if min_version: + from backend.services.admin_settings import version_is_stale + + if version_is_stale(data.get("model_version", "0.0.0"), min_version): + return None + return key + + +def library_has_datasheet( + storage: StorageBackend, mpn: str, patterns: list | None = None, +) -> str | None: + from backend.services.datasheet_store import resolve_datasheet + + hit = resolve_datasheet(storage, mpn) + if hit: + return hit + legacy = f"library/datasheets/{safe_mpn(mpn)}.pdf" + if storage.exists(legacy): + return legacy + if patterns: + from backend.periscopex.resolve_passives import resolve_mpn + + match = resolve_mpn(mpn, patterns) + if match is not None: + ds = match[0].datasheet_key + if ds and storage.exists(ds): + return ds + return None + + +def library_has_model(storage: StorageBackend, mpn: str) -> str | None: + key = f"library/models/{safe_mpn(mpn)}.json" + return key if storage.exists(key) else None + + +def library_has_passive_model(storage: StorageBackend, mpn: str) -> str | None: + safe = safe_mpn(mpn) + key = f"library/passives/{safe}.json" + if storage.exists(key): + return key + legacy = f"library/models/{safe}.json" + return legacy if storage.exists(legacy) else None + + +def save_to_library( + storage: StorageBackend, src_key: str, category: str, filename: str, +) -> str: + dst = f"library/{category}/{filename}" + storage.copy_object(src_key, dst) + return dst + + +def _specs_param_count(specs: dict) -> int: + if not isinstance(specs, dict): + return 0 + values = specs.get("values") + if isinstance(values, dict): + return sum(1 for v in values.values() if v not in (None, "", [])) + skip = {"specs_type", "component_subtype"} + return sum(1 for k, v in specs.items() if k not in skip and v not in (None, "", [])) + + +def _catalog_model_row(data: dict, key: str, *, row_type: str) -> dict: + mpn = data.get("mpn", "") or key.rsplit("/", 1)[-1].replace(".json", "") + specs = data.get("specs", {}) or {} + return { + "mpn": mpn, + "type": row_type, + "specs_type": specs.get("specs_type", ""), + "subtype": specs.get("component_subtype", ""), + "param_count": _specs_param_count(specs), + } + + +def list_library_catalog(storage: StorageBackend) -> dict: + from backend.services.datasheet_store import REF_PREFIX, resolve_datasheet + + ics: list[dict] = [] + seen_ic: set[str] = set() + for key in storage.list_prefix("library/extracted/"): + if not key.endswith(".json"): + continue + try: + data = storage.read_json(key) + mpn = data.get("mpn") or key.rsplit("/", 1)[-1].replace(".json", "") + if mpn in seen_ic: + continue + seen_ic.add(mpn) + ics.append({ + "mpn": mpn, + "type": "ic", + "subtype": data.get("component_subtype", ""), + "pin_count": len(data.get("pintable", [])), + "has_ratings": bool(data.get("absolute_maximum_ratings")), + "has_datasheet": bool(resolve_datasheet(storage, mpn)), + }) + except Exception: + continue + + passives: list[dict] = [] + seen_p: set[str] = set() + for key in storage.list_prefix("library/patterns/"): + if not key.endswith(".json"): + continue + try: + data = storage.read_json(key) + name = data.get("name") or key.rsplit("/", 1)[-1].replace(".json", "") + if name in seen_p: + continue + seen_p.add(name) + passives.append({ + "mpn": name, + "type": "passive", + "subtype": data.get("component_type", ""), + "description": data.get("description", ""), + "regex": data.get("regex", ""), + }) + except Exception: + continue + + simple_models: list[dict] = [] + passive_parts: list[dict] = [] + seen_m: set[str] = set() + for prefix, row_type, dest in ( + ("library/passives/", "passive_part", passive_parts), + ("library/models/", "simple", simple_models), + ): + for key in storage.list_prefix(prefix): + if not key.endswith(".json"): + continue + try: + data = storage.read_json(key) + row = _catalog_model_row(data, key, row_type=row_type) + if row["mpn"] in seen_m: + continue + seen_m.add(row["mpn"]) + row["has_datasheet"] = bool(resolve_datasheet(storage, row["mpn"])) + dest.append(row) + except Exception: + continue + + datasheets: list[dict] = [] + seen_d: set[str] = set() + for key in storage.list_prefix(REF_PREFIX): + if not key.endswith(".json"): + continue + try: + ref = storage.read_json(key) + mpn = ref.get("mpn") or key.rsplit("/", 1)[-1].replace(".json", "") + if mpn in seen_d: + continue + seen_d.add(mpn) + datasheets.append({ + "mpn": mpn, + "hash": ref.get("hash"), + "has_extraction": mpn in seen_ic, + "has_model": mpn in seen_m, + }) + except Exception: + continue + + ics.sort(key=lambda r: r["mpn"].lower()) + passives.sort(key=lambda r: r["mpn"].lower()) + passive_parts.sort(key=lambda r: r["mpn"].lower()) + simple_models.sort(key=lambda r: r["mpn"].lower()) + datasheets.sort(key=lambda r: r["mpn"].lower()) + return { + "ics": ics, + "passives": passives, + "passive_parts": passive_parts, + "simple": simple_models, + "datasheets": datasheets, + } + + +def list_library_patterns(storage: StorageBackend) -> list[str]: + return [k for k in storage.list_prefix("library/patterns/") if k.endswith(".json")] + + +def load_library_patterns(storage: StorageBackend): + from backend.periscopex.resolve_passives import load_patterns + from backend.services.storage import LocalStorageBackend + + if isinstance(storage, LocalStorageBackend): + d = storage._path("library/patterns") + if not d.is_dir(): + return [] + return load_patterns(str(d)) + import tempfile + + keys = list_library_patterns(storage) + if not keys: + return [] + with tempfile.TemporaryDirectory() as tmp: + dest = Path(tmp) / "patterns" + dest.mkdir() + for key in keys: + storage.download_to_local(key, dest / key.rsplit("/", 1)[-1]) + return load_patterns(str(dest)) diff --git a/periscope/src/backend/services/purple_parts.py b/periscope/src/backend/services/purple_parts.py new file mode 100644 index 0000000..6b908e6 --- /dev/null +++ b/periscope/src/backend/services/purple_parts.py @@ -0,0 +1,200 @@ +"""LCSC catalogue client (purple-parts HTTP API). + +Exact-MPN only on reverse lookup. Pipeline backstop never treats an LCSC-shaped +value sitting in the MPN slot as a code — that is upload-time column rewrite. +""" + +from __future__ import annotations + +import asyncio +import csv +import io +import logging +import re +import time +from typing import Optional + +import httpx + +from backend.config import settings + +log = logging.getLogger(__name__) + +_LCSC_RE = re.compile(r"^C\d+$", re.IGNORECASE) +_TOKEN_TTL_S = 50 * 60 +_token: dict[str, float | str] = {"value": "", "exp": 0.0} +_token_lock = asyncio.Lock() +_BATCH_SIZE = 400 + + +def is_lcsc_code(value: str | None) -> bool: + if not value: + return False + return bool(_LCSC_RE.match(value.strip())) + + +async def _get_identity_token() -> str | None: + now = time.time() + cached = str(_token.get("value") or "") + if cached and float(_token.get("exp") or 0) > now: + return cached + async with _token_lock: + cached = str(_token.get("value") or "") + if cached and float(_token.get("exp") or 0) > now: + return cached + try: + from google.auth.transport.requests import Request + from google.oauth2 import id_token as gid + except ImportError: + log.warning("google-auth missing; LCSC catalogue disabled") + return None + loop = asyncio.get_running_loop() + try: + tok = await loop.run_in_executor( + None, + lambda: gid.fetch_id_token(Request(), settings.purple_parts_url), + ) + except Exception as exc: + log.debug("LCSC identity token skipped: %s", exc) + return None + _token["value"] = tok + _token["exp"] = now + _TOKEN_TTL_S + return tok + + +def detect_lcsc_column(csv_bytes: bytes, mpn_col: str) -> bool: + """True only when every non-empty cell in ``mpn_col`` is ``^C\\d+$``.""" + text = csv_bytes.decode("utf-8", errors="replace") + reader = csv.DictReader(io.StringIO(text)) + if not reader.fieldnames or mpn_col not in reader.fieldnames: + return False + saw = False + for row in reader: + cell = (row.get(mpn_col) or "").strip() + if not cell: + continue + if not is_lcsc_code(cell): + return False + saw = True + return saw + + +def _headers(token: str) -> dict[str, str]: + return { + "Authorization": f"Bearer {token}", + "X-API-Key": settings.purple_parts_api_key, + "Content-Type": "application/json", + } + + +async def lookup_lcsc_batch(lcsc_codes: list[str]) -> dict[str, Optional[dict]]: + if not settings.use_purple_parts: + return {c: None for c in lcsc_codes} + codes = [c.strip() for c in lcsc_codes if c and c.strip()] + if not codes: + return {} + token = await _get_identity_token() + if token is None: + return {c: None for c in codes} + out: dict[str, Optional[dict]] = {c: None for c in codes} + base = settings.purple_parts_url.rstrip("/") + async with httpx.AsyncClient(timeout=15) as client: + for i in range(0, len(codes), _BATCH_SIZE): + chunk = codes[i:i + _BATCH_SIZE] + try: + resp = await client.post( + f"{base}/v1/parts/by-lcsc/batch", + headers=_headers(token), + json={"ids": chunk}, + ) + resp.raise_for_status() + body = resp.json() + except Exception as exc: + log.warning("LCSC by-id batch failed: %s", exc) + continue + for code, part in (body.get("results") or {}).items(): + out[code] = part + return out + + +def _norm_mpn(value: str | None) -> str: + return "".join((value or "").split()).upper() + + +def _pick_exact(query: str, candidates: list[dict]) -> Optional[dict]: + want = _norm_mpn(query) + for part in candidates: + if part and _norm_mpn(part.get("mpn")) == want: + return part + return None + + +async def lookup_mpn_batch(mpns: list[str]) -> dict[str, Optional[dict]]: + if not settings.use_purple_parts: + return {m: None for m in mpns} + names = list(dict.fromkeys(m.strip() for m in mpns if m and m.strip())) + if not names: + return {} + token = await _get_identity_token() + if token is None: + return {m: None for m in names} + out: dict[str, Optional[dict]] = {m: None for m in names} + base = settings.purple_parts_url.rstrip("/") + async with httpx.AsyncClient(timeout=15) as client: + for i in range(0, len(names), _BATCH_SIZE): + chunk = names[i:i + _BATCH_SIZE] + try: + resp = await client.post( + f"{base}/v1/parts/by-mpn/batch", + headers=_headers(token), + json={"mpns": chunk}, + ) + resp.raise_for_status() + body = resp.json() + except Exception as exc: + log.warning("LCSC by-mpn batch failed: %s", exc) + continue + for mpn, cands in (body.get("results") or {}).items(): + out[mpn] = _pick_exact(mpn, cands or []) + return out + + +async def resolve_lcsc_column_bytes( + csv_bytes: bytes, + *, + mpn_col: str = "Manufacturer Part Number", +) -> tuple[bytes, int, dict[str, str], dict[str, dict]]: + if not settings.use_purple_parts: + return csv_bytes, 0, {}, {} + text = csv_bytes.decode("utf-8", errors="replace") + reader = csv.DictReader(io.StringIO(text)) + fields = list(reader.fieldnames or []) + rows = list(reader) + if not rows or mpn_col not in fields: + return csv_bytes, 0, {}, {} + jobs: list[tuple[int, str]] = [] + for i, row in enumerate(rows): + code = (row.get(mpn_col) or "").strip() + if is_lcsc_code(code): + jobs.append((i, code)) + if not jobs: + return csv_bytes, 0, {}, {} + unique = sorted({c for _, c in jobs}) + found = await lookup_lcsc_batch(unique) + n = 0 + mapping: dict[str, str] = {} + payloads: dict[str, dict] = {} + for i, code in jobs: + part = found.get(code) + if part and part.get("mpn"): + rows[i][mpn_col] = part["mpn"] + mapping[code] = part["mpn"] + payloads[code] = dict(part) + n += 1 + if n == 0: + return csv_bytes, 0, {}, {} + buf = io.StringIO() + writer = csv.DictWriter(buf, fieldnames=fields, extrasaction="ignore") + writer.writeheader() + writer.writerows(rows) + return buf.getvalue().encode("utf-8"), n, mapping, payloads diff --git a/tests/test_periscope_pipeline_rewrite.py b/tests/test_periscope_pipeline_rewrite.py new file mode 100644 index 0000000..ec6ae59 --- /dev/null +++ b/tests/test_periscope_pipeline_rewrite.py @@ -0,0 +1,36 @@ +"""Prove analysis pipeline src is native, not PinScope pipeline.""" + +from __future__ import annotations + +import ast +from pathlib import Path + + +def test_src_pipeline_does_not_import_pinscope_validate(): + from backend.services import pipeline as pipe + + tree = ast.parse(Path(pipe.__file__).read_text()) + imported = { + node.module + for node in ast.walk(tree) + if isinstance(node, ast.ImportFrom) and node.module + } + from_names = { + (node.module, alias.name) + for node in ast.walk(tree) + if isinstance(node, ast.ImportFrom) and node.module + for alias in node.names + } + assert "backend.periscopex.validate" not in imported + assert ("backend.services", "datasheet_extract") in from_names + assert ("backend.services", "job_workspace") in from_names + assert ("backend.services.validation", "validate_design_async") in from_names + assert Path(pipe.__file__).resolve().parts[-4:-1] == ("src", "backend", "services") + + +def test_src_projects_write_meta_uses_owner_user_id(): + from backend.services import projects as proj + + src = Path(proj.__file__).read_text() + assert "owner_user_id" in src + assert Path(proj.__file__).resolve().parts[-4:-1] == ("src", "backend", "services")