Rewrite analysis pipeline, projects, email, and LCSC client.

Native job_workspace / datasheet_extract / review_session stay the live
BOM-extract-graph-review path. Project writes follow owner_user_id.
This commit is contained in:
2026-09-20 19:10:28 +02:00
parent 6a4903fa3a
commit 3d7fa21eb6
5 changed files with 3021 additions and 0 deletions
+316
View File
@@ -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",
)
File diff suppressed because it is too large Load Diff
+950
View File
@@ -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))
@@ -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
+36
View File
@@ -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")