Spaces:
Running on CPU Upgrade
Running on CPU Upgrade
| from __future__ import annotations | |
| import base64 | |
| import textwrap | |
| def _encode(script: str) -> str: | |
| return base64.b64encode(script.encode("utf-8")).decode("ascii") | |
| UNIVERSAL_MODEL_CARD_WORKER_SCRIPT = r''' | |
| import hashlib | |
| import hmac | |
| import hashlib | |
| import hmac | |
| import json | |
| import os | |
| import re | |
| import shutil | |
| import subprocess | |
| import sys | |
| import time | |
| from datetime import datetime, timezone | |
| from pathlib import Path | |
| from textwrap import dedent | |
| TARGET_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{1,95}/[A-Za-z0-9][A-Za-z0-9._-]{1,95}$") | |
| GIST_URL = "https://gist.github.com/gary149/2aba2962375fa9ca56bb9ef53f00b73d" | |
| AUTO_SPACE_HARDWARE_CHOICES = {"zero-a10g", "cpu-basic", "t4-small", "t4-medium", "a10g-small", "a10g-large", "l4x1", "l40sx1"} | |
| DEFAULT_PREFERRED_SPACE_HARDWARE = "zero-a10g" | |
| DEFAULT_FALLBACK_SPACE_HARDWARE = "a10g-large" | |
| DEFAULT_MODEL_ID = "sshleifer/tiny-gpt2" | |
| # Internal agent/recovery files may be needed inside the transient Pi | |
| # workspace, but they should not be published to the generated Space or shown | |
| # twice as generated application files. Canonical copies live in run_dir/logs, | |
| # run_dir/repair, and run_dir/traces. | |
| INTERNAL_WORKSPACE_ARTIFACT_NAMES = { | |
| "GOAL.md", | |
| "INCIDENT_BRIEF.md", | |
| "PI_DIAGNOSIS_GOAL.md", | |
| "DEPENDENCY_ERROR_BRIEF.md", | |
| "REPAIR_DECISION.json", | |
| "REPAIR_GOAL.md", | |
| "REPAIR_BRIEF.md", | |
| "REPAIR_PLAN.md", | |
| "REPAIR_SUMMARY.md", | |
| "PI_SUMMARY.md", | |
| "TECHNICAL_BLOCKERS.json", | |
| } | |
| def internal_workspace_upload_ignore_patterns() -> list[str]: | |
| patterns = [".git/*", "node_modules/*", "__pycache__/*", "*.pyc"] | |
| for name in sorted(INTERNAL_WORKSPACE_ARTIFACT_NAMES): | |
| patterns.extend([name, f"**/{name}"]) | |
| return patterns | |
| def internal_workspace_copy_ignore(_dir: str, names: list[str]) -> set[str]: | |
| ignored = {".git", "node_modules", "__pycache__"} | |
| ignored.update(name for name in names if name in INTERNAL_WORKSPACE_ARTIFACT_NAMES or name.endswith(".pyc")) | |
| return ignored | |
| def is_publishable_workspace_file(path: Path) -> bool: | |
| if path.name in INTERNAL_WORKSPACE_ARTIFACT_NAMES: | |
| return False | |
| if any(part in {".git", "node_modules", "__pycache__"} for part in path.parts): | |
| return False | |
| if path.suffix == ".pyc": | |
| return False | |
| return True | |
| def now(): | |
| return datetime.now(timezone.utc).isoformat() | |
| def write_json(path: Path, payload: dict): | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| path.write_text(json.dumps(payload, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| def append_event(path: Path, step: str, status: str, message: str, data: dict | None = None): | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| event = {"ts": now(), "step": step, "status": status, "message": message, "data": data or {}} | |
| line = json.dumps(event, ensure_ascii=False) | |
| with path.open("a", encoding="utf-8") as f: | |
| f.write(line + "\n") | |
| print(line, flush=True) | |
| def _artifact_entry(run_dir: Path, rel_path: str, *, kind: str = "file") -> dict: | |
| path = run_dir / rel_path | |
| if kind == "folder": | |
| present = path.exists() and any(child.is_file() for child in path.rglob("*")) | |
| else: | |
| present = path.exists() and path.is_file() | |
| payload = {"path": rel_path, "kind": kind, "present": bool(present)} | |
| if present and kind == "file": | |
| try: | |
| payload["size"] = path.stat().st_size | |
| except Exception: | |
| pass | |
| return payload | |
| def write_artifact_manifest(run_dir: Path, *, events_path: Path | None = None, reason: str = "snapshot") -> dict: | |
| """Snapshot the artifacts that the worker actually wrote to the mounted bucket. | |
| The web UI should prefer this manifest over guessed paths. It is written | |
| during failures as well as success so a run that created a Space and then | |
| crashed can still expose traces, logs and reports in the dock. | |
| """ | |
| run_id = os.environ.get("RUN_ID", "") | |
| bucket_source = os.environ.get("BUCKET_SOURCE", "") | |
| artifacts = [ | |
| _artifact_entry(run_dir, "events.jsonl"), | |
| _artifact_entry(run_dir, "state.json"), | |
| _artifact_entry(run_dir, "report.md"), | |
| _artifact_entry(run_dir, "model_analysis.json"), | |
| _artifact_entry(run_dir, "hardware_strategy.json"), | |
| _artifact_entry(run_dir, "hardware_attempts.json"), | |
| _artifact_entry(run_dir, "inference_gate.json"), | |
| _artifact_entry(run_dir, "space_runtime.json"), | |
| _artifact_entry(run_dir, "tests/generation_smoke.json"), | |
| _artifact_entry(run_dir, "tests/api_schema.json"), | |
| _artifact_entry(run_dir, "generated/TECHNICAL_BLOCKERS.json"), | |
| _artifact_entry(run_dir, "repair/INCIDENT_BRIEF.md"), | |
| _artifact_entry(run_dir, "repair/DEPENDENCY_ERROR_BRIEF.md"), | |
| _artifact_entry(run_dir, "repair/PI_DIAGNOSIS_GOAL.md"), | |
| _artifact_entry(run_dir, "repair/REPAIR_DECISION.json"), | |
| _artifact_entry(run_dir, "repair/REPAIR_BRIEF.md"), | |
| _artifact_entry(run_dir, "repair/REPAIR_PLAN.md"), | |
| _artifact_entry(run_dir, "repair/REPAIR_SUMMARY.md"), | |
| _artifact_entry(run_dir, "repair/BLOCKAGE.json"), | |
| _artifact_entry(run_dir, "logs/pi_live_output.txt"), | |
| _artifact_entry(run_dir, "logs/pi_output.txt"), | |
| _artifact_entry(run_dir, "logs/pi_diagnosis_output.txt"), | |
| _artifact_entry(run_dir, "logs/pi_repair_output.txt"), | |
| _artifact_entry(run_dir, "logs/space_logs_build.txt"), | |
| _artifact_entry(run_dir, "logs/space_logs_runtime.txt"), | |
| _artifact_entry(run_dir, "traces/raw", kind="folder"), | |
| _artifact_entry(run_dir, "traces/redacted", kind="folder"), | |
| _artifact_entry(run_dir, "generated", kind="folder"), | |
| _artifact_entry(run_dir, "repair", kind="folder"), | |
| _artifact_entry(run_dir, "logs", kind="folder"), | |
| ] | |
| payload = { | |
| "run_id": run_id, | |
| "bucket_source": bucket_source, | |
| "updated_at": now(), | |
| "reason": reason, | |
| "artifacts": artifacts, | |
| "present_paths": [a["path"] for a in artifacts if a.get("present")], | |
| } | |
| write_json(run_dir / "artifact_manifest.json", payload) | |
| if os.environ.get("ASF_EVAL_ENABLED", "").strip().lower() in {"1", "true", "yes", "on"}: | |
| try: | |
| publish_eval_record(run_dir, phase=f"manifest:{reason}", events_path=events_path) | |
| except Exception: | |
| pass | |
| if events_path is not None: | |
| try: | |
| append_event(events_path, "artifact_manifest", "success", "Updated run artifact manifest", {"reason": reason, "present_count": len(payload["present_paths"])}) | |
| except Exception: | |
| pass | |
| return payload | |
| def update_state(run_dir: Path, patch: dict) -> dict: | |
| path = run_dir / "state.json" | |
| current = load_json_if_exists(path) if path.exists() else {} | |
| if not isinstance(current, dict): | |
| current = {} | |
| merged = {**current, **patch, "updated_at": now()} | |
| write_json(path, merged) | |
| return merged | |
| def redact_text(text: str | None) -> str: | |
| if not text: | |
| return "" | |
| value = text | |
| for secret_name in ["HF_TOKEN", "HUGGING_FACE_HUB_TOKEN"]: | |
| secret = os.environ.get(secret_name) | |
| if secret: | |
| value = value.replace(secret, "[REDACTED]") | |
| value = re.sub(r"Bearer\s+[A-Za-z0-9_\-.=]+", "Bearer [REDACTED]", value) | |
| value = re.sub(r"hf_[A-Za-z0-9_\-]{10,}", "hf_[REDACTED]", value) | |
| return value | |
| def ensure_hf_token_context(run_dir: Path, events_path: Path | None = None) -> dict: | |
| """Ensure Pi and HF tooling see the same private HF token safely. | |
| HF Jobs receive HF_TOKEN as a secret. Pi, huggingface_hub, gradio_client, | |
| and the hf CLI can use either HF_TOKEN or HUGGING_FACE_HUB_TOKEN depending | |
| on the tool/version, so expose both aliases inside the isolated Job process. | |
| Never write the token value to artifacts. | |
| """ | |
| token = os.environ.get("HF_TOKEN") or os.environ.get("HUGGING_FACE_HUB_TOKEN") or "" | |
| if token: | |
| os.environ.setdefault("HF_TOKEN", token) | |
| os.environ.setdefault("HUGGING_FACE_HUB_TOKEN", token) | |
| os.environ.setdefault("HF_HUB_ENABLE_HF_TRANSFER", "1") | |
| payload = { | |
| "hf_token_present": bool(token), | |
| "hf_token_length": len(token) if token else 0, | |
| "hugging_face_hub_token_alias_present": bool(os.environ.get("HUGGING_FACE_HUB_TOKEN")), | |
| "hf_username": os.environ.get("HF_USERNAME") or "", | |
| "bucket_source": os.environ.get("BUCKET_SOURCE") or "", | |
| "target_space_id": os.environ.get("TARGET_SPACE_ID") or "", | |
| "token_value": "[REDACTED]" if token else "", | |
| } | |
| run_dir.mkdir(parents=True, exist_ok=True) | |
| write_json(run_dir / "token_context.json", payload) | |
| if events_path: | |
| append_event(events_path, "token_context", "success" if token else "warning", "Verified HF token context for Pi/HF operations", {k: v for k, v in payload.items() if k != "token_value"}) | |
| return payload | |
| def pi_tooling_context_note() -> str: | |
| return """HF tooling context: | |
| - You are running inside an HF Job with an HF_TOKEN secret configured. | |
| - The Factory exposes both HF_TOKEN and HUGGING_FACE_HUB_TOKEN aliases for huggingface_hub, gradio_client, hf CLI, and Pi provider access. | |
| - You may use the token through normal tools when needed to inspect the model/Space or validate private resources. | |
| - Never print, write, echo, or commit token values. Do not include Authorization headers in artifacts. | |
| - If a Hub operation fails with 401/403/quota/billing, report it as evidence instead of retrying blindly. | |
| """ | |
| def _agent_trace_json_safe(payload): | |
| try: | |
| return json.loads(redact_text(json.dumps(payload, ensure_ascii=False))) | |
| except Exception: | |
| return {"value": redact_text(str(payload))[:8000]} | |
| def write_agent_trace_record( | |
| run_dir: Path, | |
| *, | |
| phase: str, | |
| event: str, | |
| status: str = "info", | |
| message: str = "", | |
| data: dict | None = None, | |
| text: str | None = None, | |
| artifacts: list[str] | None = None, | |
| ): | |
| """Append one run-level agent trace record to RAW and redacted journals. | |
| Pi can be invoked several times during a run: initial build, blockage | |
| diagnosis, and optional repair patch. The individual Pi session files and | |
| stdout logs remain useful, but this journal is the canonical narrative tying | |
| those phases together into one auditable agent run. | |
| """ | |
| try: | |
| payload = { | |
| "ts": now(), | |
| "phase": phase, | |
| "event": event, | |
| "status": status, | |
| "message": message, | |
| "data": data or {}, | |
| "artifacts": artifacts or [], | |
| } | |
| if text is not None: | |
| payload["text"] = text | |
| raw_dir = run_dir / "traces" / "raw" | |
| redacted_dir = run_dir / "traces" / "redacted" | |
| raw_dir.mkdir(parents=True, exist_ok=True) | |
| redacted_dir.mkdir(parents=True, exist_ok=True) | |
| raw_line = json.dumps(payload, ensure_ascii=False) | |
| redacted_payload = _agent_trace_json_safe(payload) | |
| redacted_line = json.dumps(redacted_payload, ensure_ascii=False) | |
| with (raw_dir / "agent_trace.jsonl").open("a", encoding="utf-8") as f: | |
| f.write(raw_line + "\n") | |
| with (redacted_dir / "agent_trace.jsonl").open("a", encoding="utf-8") as f: | |
| f.write(redacted_line + "\n") | |
| except Exception: | |
| # Trace journaling must never break the build/repair path. | |
| pass | |
| def append_agent_trace_artifact(run_dir: Path, *, phase: str, event: str, artifact: str, text: str | None = None, status: str = "success", data: dict | None = None): | |
| tail = None if text is None else text[-20000:] | |
| write_agent_trace_record(run_dir, phase=phase, event=event, status=status, message=f"Agent phase wrote {artifact}", data=data or {}, text=tail, artifacts=[artifact]) | |
| def safe_details(details: dict | None) -> dict: | |
| if not details: | |
| return {} | |
| try: | |
| return json.loads(redact_text(json.dumps(details, ensure_ascii=False))) | |
| except Exception: | |
| return {"redacted_details": redact_text(str(details))[-4000:]} | |
| def exception_payload(exc: Exception) -> dict: | |
| """Extract the most useful redacted details from HF Hub / HTTP exceptions.""" | |
| payload = { | |
| "error": redact_text(str(exc))[:4000], | |
| "exception_type": type(exc).__name__, | |
| } | |
| response = getattr(exc, "response", None) | |
| if response is not None: | |
| status_code = getattr(response, "status_code", None) | |
| if status_code is not None: | |
| payload["status_code"] = status_code | |
| try: | |
| text_body = response.text | |
| except Exception: | |
| text_body = "" | |
| if text_body: | |
| payload["response_text"] = redact_text(text_body)[:4000] | |
| payload["hf_error_detail"] = redact_text(text_body)[:4000] | |
| try: | |
| json_body = response.json() | |
| except Exception: | |
| json_body = None | |
| if json_body is not None: | |
| payload["response_json"] = safe_details(json_body) | |
| if isinstance(json_body, dict): | |
| for key in ("error", "message", "detail", "details"): | |
| value = json_body.get(key) | |
| if value: | |
| payload["hf_error_message"] = redact_text(value if isinstance(value, str) else json.dumps(value, ensure_ascii=False))[:4000] | |
| break | |
| for attr in ("server_message", "request_id"): | |
| value = getattr(exc, attr, None) | |
| if value: | |
| payload[attr] = redact_text(str(value))[:1000] | |
| return payload | |
| def eval_enabled() -> bool: | |
| return os.environ.get("ASF_EVAL_RECORD_ENABLED", "").strip().lower() in {"1", "true", "yes", "on"} | |
| def eval_safe_hash(value: str | None, prefix: str = "h") -> str: | |
| raw = (value or "").strip() | |
| salt = os.environ.get("ASF_EVAL_SALT", "").strip() | |
| if not raw: | |
| return "" | |
| if salt: | |
| digest = hmac.new(salt.encode("utf-8"), raw.encode("utf-8"), hashlib.sha256).hexdigest()[:24] | |
| else: | |
| digest = hashlib.sha256(raw.encode("utf-8")).hexdigest()[:24] | |
| return f"{prefix}_{digest}" | |
| def eval_load_json(path: Path) -> dict: | |
| if not path.exists(): | |
| return {} | |
| try: | |
| return json.loads(path.read_text(encoding="utf-8", errors="replace")) | |
| except Exception as exc: | |
| return {"_parse_error": str(exc)} | |
| def eval_load_events(path: Path) -> list[dict]: | |
| if not path.exists(): | |
| return [] | |
| events = [] | |
| for line in path.read_text(encoding="utf-8", errors="replace").splitlines(): | |
| try: | |
| evt = json.loads(line) | |
| if isinstance(evt, dict): | |
| events.append(evt) | |
| except Exception: | |
| pass | |
| return events | |
| def eval_parse_ts(value: str | None) -> float | None: | |
| if not value: | |
| return None | |
| try: | |
| return datetime.fromisoformat(str(value).replace("Z", "+00:00")).timestamp() | |
| except Exception: | |
| return None | |
| def eval_compact_timeline(events: list[dict]) -> list[dict]: | |
| stages: dict[str, dict] = {} | |
| order: list[str] = [] | |
| for evt in events: | |
| step = str(evt.get("step") or "unknown")[:80] | |
| status = str(evt.get("status") or "")[:80] | |
| ts = str(evt.get("ts") or "") | |
| tsv = eval_parse_ts(ts) | |
| if step not in stages: | |
| stages[step] = {"stage": step, "first_ts": ts, "last_ts": ts, "status": status, "event_count": 0} | |
| order.append(step) | |
| item = stages[step] | |
| item["last_ts"] = ts or item.get("last_ts") | |
| item["status"] = status or item.get("status") | |
| item["event_count"] = int(item.get("event_count") or 0) + 1 | |
| if tsv is not None: | |
| if "_first" not in item: | |
| item["_first"] = tsv | |
| item["_last"] = tsv | |
| timeline = [] | |
| for step in order: | |
| item = dict(stages[step]) | |
| first = item.pop("_first", None) | |
| last = item.pop("_last", None) | |
| if first is not None and last is not None and last >= first: | |
| item["duration_ms"] = int((last - first) * 1000) | |
| timeline.append(item) | |
| return timeline[-80:] | |
| def eval_redacted_tail(path: Path, limit: int = 2400) -> str: | |
| if os.environ.get("ASF_EVAL_INCLUDE_REDACTED_TAILS", "").strip().lower() not in {"1", "true", "yes", "on"}: | |
| return "" | |
| try: | |
| return redact_text(path.read_text(encoding="utf-8", errors="replace"))[-limit:] | |
| except Exception: | |
| return "" | |
| def publish_eval_record(run_dir: Path, *, phase: str, events_path: Path | None = None) -> dict: | |
| """Write an anonymized cross-run evaluation record to the optional eval bucket. | |
| This intentionally stores structured metrics and hashed identifiers, not raw | |
| generated code, prompts, tokens, user bucket paths, or target Space IDs. | |
| """ | |
| if not eval_enabled(): | |
| return {"enabled": False} | |
| # Jobs only write the anonymized eval record back into the user's run bucket. | |
| # The ASF backend later publishes this record to the operator eval bucket | |
| # mounted on the ASF Space. User Jobs must never mount or write the private | |
| # operator archive directly. | |
| run_id = os.environ.get("RUN_ID", run_dir.name) | |
| events = eval_load_events(events_path or (run_dir / "events.jsonl")) | |
| state = eval_load_json(run_dir / "state.json") | |
| analysis = eval_load_json(run_dir / "model_analysis.json") | |
| hardware_strategy = eval_load_json(run_dir / "hardware_strategy.json") | |
| hardware_attempts = eval_load_json(run_dir / "hardware_attempts.json") | |
| inference_gate = eval_load_json(run_dir / "inference_gate.json") | |
| generation_smoke = eval_load_json(run_dir / "tests" / "generation_smoke.json") | |
| space_runtime = eval_load_json(run_dir / "space_runtime.json") | |
| manifest = eval_load_json(run_dir / "artifact_manifest.json") | |
| repair_decision = eval_load_json(run_dir / "repair" / "REPAIR_DECISION.json") | |
| blockage = eval_load_json(run_dir / "repair" / "BLOCKAGE.json") | |
| pi_resolution = state.get("pi_model_resolution") if isinstance(state.get("pi_model_resolution"), dict) else {} | |
| status = str(state.get("status") or inference_gate.get("status") or phase or "unknown") | |
| failure_message = str(state.get("message") or (state.get("details") or {}).get("error") or "")[:1200] | |
| started_at = events[0].get("ts") if events else (state.get("created_at") or os.environ.get("LAUNCHED_AT") or "") | |
| finished_at = state.get("updated_at") or (events[-1].get("ts") if events else now()) | |
| model_id = os.environ.get("MODEL_ID") or state.get("model_id") or analysis.get("model_id") or "" | |
| target_space = state.get("target_space") or os.environ.get("TARGET_SPACE_ID") or "" | |
| include_model_id = os.environ.get("ASF_EVAL_INCLUDE_MODEL_ID", "").strip().lower() in {"1", "true", "yes", "on"} | |
| anon_run_id = eval_safe_hash(run_id, "run") | |
| record = { | |
| "schema_version": "1.0", | |
| "app_version": os.environ.get("ASF_VERSION", "unknown"), | |
| "phase": phase, | |
| "anonymous_run_id": anon_run_id, | |
| "run_id_hash": eval_safe_hash(run_id, "run"), | |
| "anonymous_user_hash": eval_safe_hash(os.environ.get("HF_USERNAME"), "user"), | |
| "salt_configured": bool(os.environ.get("ASF_EVAL_SALT")), | |
| "run_kind": state.get("kind") or ("validate_existing_space" if os.environ.get("API_NAME") else "universal_model_card_builder"), | |
| "started_at": started_at, | |
| "finished_at": finished_at, | |
| "input_model": { | |
| "model_id_redacted": not include_model_id, | |
| "model_id": model_id if include_model_id else "", | |
| "model_hash": eval_safe_hash(model_id, "model"), | |
| "pipeline_tag": analysis.get("pipeline_tag"), | |
| "library_name": analysis.get("library_name"), | |
| "default_model_target": bool(analysis.get("default_model_target")), | |
| }, | |
| "config": { | |
| "implementation_mode": os.environ.get("IMPLEMENTATION_MODE", ""), | |
| "expected_output_type": os.environ.get("EXPECTED_OUTPUT_TYPE", ""), | |
| "pi_model_requested": os.environ.get("PI_MODEL") or state.get("pi_model") or "", | |
| "pi_model_effective": pi_resolution.get("effective_model") or pi_resolution.get("observed_model") or "", | |
| "try_zero_gpu_first": os.environ.get("TRY_ZERO_GPU_FIRST", ""), | |
| "allow_fixed_gpu_fallback": os.environ.get("ALLOW_FIXED_GPU_FALLBACK", ""), | |
| "preferred_gpu": os.environ.get("PREFERRED_SPACE_HARDWARE", ""), | |
| "fallback_gpu": os.environ.get("FALLBACK_SPACE_HARDWARE", ""), | |
| }, | |
| "outcome": { | |
| "status": status, | |
| "message_redacted": redact_text(failure_message)[:1200], | |
| "target_space_hash": eval_safe_hash(target_space, "space"), | |
| "space_created": bool(target_space), | |
| "space_runtime_stage": space_runtime.get("stage"), | |
| "space_runtime_hardware": space_runtime.get("hardware") or space_runtime.get("requested_hardware"), | |
| }, | |
| "timeline": eval_compact_timeline(events), | |
| "hardware": { | |
| "selected": hardware_strategy.get("selected_hardware") or state.get("selected_hardware") or "", | |
| "strategy": hardware_strategy.get("strategy") or hardware_attempts.get("strategy") or "", | |
| "attempts": hardware_strategy.get("attempts") or hardware_attempts.get("attempts") or state.get("hardware_attempts") or [], | |
| "single_space_per_run": True, | |
| }, | |
| "pi": { | |
| "requested_model": os.environ.get("PI_MODEL") or state.get("pi_model") or "", | |
| "effective_model": pi_resolution.get("effective_model") or pi_resolution.get("observed_model") or "", | |
| "model_mismatch": bool(pi_resolution.get("mismatch")), | |
| "build_log_tail_redacted": eval_redacted_tail(run_dir / "logs" / "pi_output.txt"), | |
| "diagnosis_log_tail_redacted": eval_redacted_tail(run_dir / "logs" / "pi_diagnosis_output.txt"), | |
| }, | |
| "recovery": { | |
| "entered": any(str(e.get("step", "")).startswith(("repair", "failure_diagnosis", "pi_diagnosis", "technical_blocker")) for e in events), | |
| "decision": repair_decision.get("action") or repair_decision.get("decision") or "", | |
| "confidence": repair_decision.get("confidence") or "", | |
| "blockage": bool(blockage), | |
| "blockage_category": blockage.get("category") or blockage.get("reason") or "", | |
| }, | |
| "validation": { | |
| "generation_smoke_status": generation_smoke.get("status") or "", | |
| "latency_seconds": generation_smoke.get("latency_seconds") or generation_smoke.get("observed_latency_seconds"), | |
| "expected_output_type": generation_smoke.get("expected_output_type") or os.environ.get("EXPECTED_OUTPUT_TYPE", ""), | |
| "inference_gate_status": inference_gate.get("status") or "", | |
| "manual_hardware_required": bool(inference_gate.get("manual_hardware_required")), | |
| }, | |
| "artifacts": { | |
| "present_count": len(manifest.get("present_paths") or []), | |
| "present_paths": sorted(str(p) for p in (manifest.get("present_paths") or []))[:120], | |
| }, | |
| "privacy": { | |
| "tokens_stored": False, | |
| "raw_prompts_stored": False, | |
| "generated_code_stored": False, | |
| "target_space_redacted": True, | |
| "user_bucket_redacted": True, | |
| "redacted_tails_enabled": os.environ.get("ASF_EVAL_INCLUDE_REDACTED_TAILS", "").strip().lower() in {"1", "true", "yes", "on"}, | |
| }, | |
| } | |
| day = (finished_at or now())[:10] | |
| try: | |
| yyyy, mm, dd = day.split("-") | |
| except Exception: | |
| yyyy, mm, dd = now()[:10].split("-") | |
| dest = run_dir / "eval_record.json" | |
| write_json(dest, record) | |
| events_dest = run_dir / "events_compact.jsonl" | |
| events_dest.write_text("".join(json.dumps(item, ensure_ascii=False) + "\n" for item in record["timeline"]), encoding="utf-8") | |
| return {"enabled": True, "written": True, "path": str(dest), "publish_mode": "backend"} | |
| def fail(run_dir: Path, events_path: Path, message: str, details: dict | None = None, status: str = "failed"): | |
| safe = safe_details(details) | |
| append_event(events_path, "failure", "failed", message, safe) | |
| existing_state = load_json_if_exists(run_dir / "state.json") if (run_dir / "state.json").exists() else {} | |
| if not isinstance(existing_state, dict): | |
| existing_state = {} | |
| failure_state = { | |
| **existing_state, | |
| "run_id": os.environ.get("RUN_ID"), | |
| "kind": existing_state.get("kind") or "universal_model_card_builder", | |
| "status": status, | |
| "message": message, | |
| "updated_at": now(), | |
| "details": safe, | |
| } | |
| # Preserve the target Space when the worker fails after repository creation. | |
| target_space = failure_state.get("target_space") or os.environ.get("TARGET_SPACE_ID") or "" | |
| if target_space: | |
| failure_state["target_space"] = target_space | |
| failure_state["target_space_url"] = f"https://huggingface.co/spaces/{target_space}" | |
| write_json(run_dir / "state.json", failure_state) | |
| report = f"""# Agentic Space Factory — model Article Reproduction Report | |
| Status: **{status}** | |
| {message} | |
| Target Space: {failure_state.get('target_space_url') or failure_state.get('target_space') or 'not created / unknown'} | |
| ```json | |
| {json.dumps(safe, indent=2, ensure_ascii=False)} | |
| ``` | |
| """ | |
| (run_dir / "report.md").write_text(report, encoding="utf-8") | |
| write_artifact_manifest(run_dir, events_path=events_path, reason="failure") | |
| try: | |
| publish_eval_record(run_dir, phase="failure", events_path=events_path) | |
| except Exception: | |
| pass | |
| raise SystemExit(1) | |
| def run_cmd(cmd: list[str], *, cwd: Path | None = None, env: dict | None = None, timeout: int = 600): | |
| result = subprocess.run(cmd, cwd=str(cwd) if cwd else None, env=env, text=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, timeout=timeout) | |
| return result.returncode, redact_text(result.stdout) | |
| def sync_pi_traces(run_dir: Path, *, emit_event: bool = False, events_path: Path | None = None): | |
| """Best-effort incremental copy of Pi session JSONL traces into the run bucket.""" | |
| traces_dir = Path.home() / ".pi" / "agent" / "sessions" | |
| raw_dir = run_dir / "traces" / "raw" | |
| redacted_dir = run_dir / "traces" / "redacted" | |
| raw_dir.mkdir(parents=True, exist_ok=True) | |
| redacted_dir.mkdir(parents=True, exist_ok=True) | |
| count = 0 | |
| if traces_dir.exists(): | |
| for path in traces_dir.rglob("*.jsonl"): | |
| rel = path.relative_to(traces_dir) | |
| target_raw = raw_dir / rel | |
| target_raw.parent.mkdir(parents=True, exist_ok=True) | |
| text = path.read_text(encoding="utf-8", errors="ignore") | |
| target_raw.write_text(text, encoding="utf-8") | |
| target_redacted = redacted_dir / rel | |
| target_redacted.parent.mkdir(parents=True, exist_ok=True) | |
| target_redacted.write_text(redact_text(text), encoding="utf-8") | |
| count += 1 | |
| if emit_event and events_path: | |
| append_event(events_path, "pi_traces", "running", "Synchronized Pi traces while Pi is running", {"count": count}) | |
| return count | |
| def run_cmd_streaming( | |
| cmd: list[str], | |
| *, | |
| cwd: Path | None = None, | |
| env: dict | None = None, | |
| timeout: int = 600, | |
| live_log_path: Path | None = None, | |
| events_path: Path | None = None, | |
| run_dir: Path | None = None, | |
| step: str = "pi_run", | |
| trace_phase: str | None = None, | |
| ): | |
| """Run a command while incrementally writing stdout to the bucket. | |
| subprocess.run() only returns output at the end. Pi can work for many minutes, | |
| so stream stdout to logs/pi_live_output.txt and periodically copy Pi session | |
| JSONL traces into traces/redacted while the process is still alive. | |
| """ | |
| started = time.monotonic() | |
| trace_phase = trace_phase or step | |
| last_event = 0.0 | |
| last_trace_sync = 0.0 | |
| lines: list[str] = [] | |
| if live_log_path: | |
| live_log_path.parent.mkdir(parents=True, exist_ok=True) | |
| live_log_path.write_text("", encoding="utf-8") | |
| if run_dir: | |
| write_agent_trace_record(run_dir, phase=trace_phase, event="command_started", status="started", message="Pi command started", data={"step": step, "cwd": str(cwd) if cwd else ""}) | |
| proc = subprocess.Popen( | |
| cmd, | |
| cwd=str(cwd) if cwd else None, | |
| env=env, | |
| text=True, | |
| stdout=subprocess.PIPE, | |
| stderr=subprocess.STDOUT, | |
| bufsize=1, | |
| ) | |
| try: | |
| assert proc.stdout is not None | |
| while True: | |
| if timeout and time.monotonic() - started > timeout: | |
| proc.kill() | |
| message = f"Command timed out after {timeout}s" | |
| lines.append(message) | |
| if live_log_path: | |
| with live_log_path.open("a", encoding="utf-8") as f: | |
| f.write(message + "\n") | |
| return 124, redact_text("\n".join(lines)) | |
| line = proc.stdout.readline() | |
| if line: | |
| clean = redact_text(line.rstrip("\n")) | |
| lines.append(clean) | |
| if live_log_path: | |
| with live_log_path.open("a", encoding="utf-8") as f: | |
| f.write(clean + "\n") | |
| now_m = time.monotonic() | |
| if events_path and now_m - last_event > 8: | |
| output_tail = "\n".join(lines[-20:])[-2000:] | |
| append_event(events_path, step, "running", "Pi is still working", {"output_tail": output_tail}) | |
| if run_dir: | |
| write_agent_trace_record(run_dir, phase=trace_phase, event="stdout_tail", status="running", message="Pi is still working", text=output_tail, data={"step": step}) | |
| last_event = now_m | |
| if run_dir and now_m - last_trace_sync > 15: | |
| sync_pi_traces(run_dir, emit_event=False) | |
| last_trace_sync = now_m | |
| continue | |
| if proc.poll() is not None: | |
| remainder = proc.stdout.read() if proc.stdout else "" | |
| if remainder: | |
| for raw in remainder.splitlines(): | |
| clean = redact_text(raw) | |
| lines.append(clean) | |
| if live_log_path: | |
| with live_log_path.open("a", encoding="utf-8") as f: | |
| f.write(clean + "\n") | |
| break | |
| time.sleep(0.2) | |
| finally: | |
| if run_dir: | |
| sync_pi_traces(run_dir, emit_event=False) | |
| output = redact_text("\n".join(lines)) | |
| if run_dir: | |
| write_agent_trace_record(run_dir, phase=trace_phase, event="command_finished", status="success" if int(proc.returncode or 0) == 0 else "failed", message="Pi command finished", text=output[-20000:], data={"step": step, "returncode": int(proc.returncode or 0)}, artifacts=[str(live_log_path.relative_to(run_dir)) if live_log_path and str(live_log_path).startswith(str(run_dir)) else ""]) | |
| return int(proc.returncode or 0), output | |
| def install_python_deps(events_path: Path): | |
| append_event(events_path, "dependencies", "started", "Installing Python worker dependencies") | |
| code, out = run_cmd([sys.executable, "-m", "pip", "install", "-q", "--upgrade", "huggingface_hub>=1.0.0", "gradio_client>=2.0.0", "requests>=2.31.0"], timeout=600) | |
| if code != 0: | |
| append_event(events_path, "dependencies", "failed", "Python dependency installation failed", {"output_tail": out[-4000:]}) | |
| raise RuntimeError(out) | |
| append_event(events_path, "dependencies", "success", "Python worker dependencies installed") | |
| def ensure_node(events_path: Path): | |
| node = shutil.which("node") | |
| npm = shutil.which("npm") | |
| if node and npm: | |
| _, node_v = run_cmd([node, "--version"], timeout=30) | |
| _, npm_v = run_cmd([npm, "--version"], timeout=30) | |
| append_event(events_path, "node", "success", "Node/npm already available", {"node": node_v.strip(), "npm": npm_v.strip()}) | |
| return | |
| append_event(events_path, "node", "started", "Installing nodejs/npm through apt-get") | |
| code, out = run_cmd(["bash", "-lc", "apt-get update -qq && apt-get install -y -qq nodejs npm"], timeout=600) | |
| if code != 0: | |
| append_event(events_path, "node", "failed", "Could not install nodejs/npm", {"output_tail": out[-4000:]}) | |
| raise RuntimeError(out) | |
| append_event(events_path, "node", "success", "Installed nodejs/npm") | |
| def install_pi(events_path: Path): | |
| ensure_node(events_path) | |
| append_event(events_path, "pi_install", "started", "Installing Pi coding agent from npm") | |
| code, out = run_cmd(["npm", "install", "-g", "@mariozechner/pi-coding-agent"], timeout=900) | |
| if code != 0: | |
| append_event(events_path, "pi_install", "failed", "Pi npm installation failed", {"output_tail": out[-4000:]}) | |
| raise RuntimeError(out) | |
| code, version = run_cmd(["pi", "--version"], timeout=60) | |
| append_event(events_path, "pi_install", "success", "Pi installed", {"version_output": version.strip()[-300:]}) | |
| def configure_pi(events_path: Path, model: str): | |
| pi_dir = Path.home() / ".pi" / "agent" | |
| pi_dir.mkdir(parents=True, exist_ok=True) | |
| token = os.environ.get("HF_TOKEN") or os.environ.get("HUGGING_FACE_HUB_TOKEN") or "" | |
| if token: | |
| os.environ.setdefault("HF_TOKEN", token) | |
| os.environ.setdefault("HUGGING_FACE_HUB_TOKEN", token) | |
| (pi_dir / "auth.json").write_text(json.dumps({"huggingface": {"type": "api_key", "key": token}}, indent=2), encoding="utf-8") | |
| (pi_dir / "settings.json").write_text(json.dumps({"model": model, "provider": "huggingface", "autoRun": True, "autoApply": True}, indent=2), encoding="utf-8") | |
| append_event(events_path, "pi_config", "success" if token else "warning", "Configured Pi with Hugging Face provider token context", {"model": model, "hf_token_present": bool(token), "token_value": "[REDACTED]" if token else ""}) | |
| def normalize_pi_model_name(value: str | None) -> str: | |
| raw = (value or "").strip().lower() | |
| # Pi/provider traces may report either a full Hub id (moonshotai/Kimi-K2-...) | |
| # or only the served model name (Kimi-K2-...). Compare on the model leaf so | |
| # provider fallbacks are detected without false positives from missing owners. | |
| if "/" in raw: | |
| raw = raw.rsplit("/", 1)[-1] | |
| return re.sub(r"[^a-z0-9]+", "", raw) | |
| def extract_pi_models_from_text(text: str) -> list[str]: | |
| """Best-effort extraction of assistant/model names from Pi stdout/session traces.""" | |
| if not text: | |
| return [] | |
| patterns = [ | |
| r"\bQwen/[A-Za-z0-9_.-]+", | |
| r"\bQwen(?:2(?:\.5)?|3)?[-_/A-Za-z0-9.]*Coder[-_/A-Za-z0-9.]*", | |
| r"\bmoonshotai/[A-Za-z0-9_.-]+", | |
| r"\bKimi[-_/A-Za-z0-9.]+", | |
| r"\bzai-org/[A-Za-z0-9_.-]+", | |
| r"\bGLM[-_/A-Za-z0-9.]+", | |
| r"\bdeepseek-ai/[A-Za-z0-9_.-]+", | |
| r"\bDeepSeek[-_/A-Za-z0-9.]+", | |
| r"\bClaude[-_/A-Za-z0-9.]+", | |
| r"\bGPT[-_/A-Za-z0-9.]+", | |
| r"\b[Mm]odel(?:\\s+used|\\s+selected|\\s*:)?\\s*[:=]?\\s*([A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+)", | |
| ] | |
| found: list[str] = [] | |
| for pattern in patterns: | |
| for match in re.finditer(pattern, text, flags=re.IGNORECASE): | |
| value = match.group(1) if match.groups() else match.group(0) | |
| value = value.strip().strip('",` ') | |
| if value and value not in found: | |
| found.append(value) | |
| return found[:12] | |
| def detect_pi_model_resolution(requested_model: str, run_dir: Path, pi_output: str) -> dict: | |
| settings_path = Path.home() / ".pi" / "agent" / "settings.json" | |
| configured_model = requested_model | |
| try: | |
| settings = json.loads(settings_path.read_text(encoding="utf-8")) | |
| configured_model = settings.get("model") or requested_model | |
| except Exception: | |
| settings = {"model": requested_model, "read_error": "could_not_read_settings"} | |
| corpus = [pi_output or ""] | |
| for trace in (run_dir / "traces" / "redacted").rglob("*.jsonl") if (run_dir / "traces" / "redacted").exists() else []: | |
| try: | |
| corpus.append(trace.read_text(encoding="utf-8", errors="ignore")[:200000]) | |
| except Exception: | |
| pass | |
| observed = extract_pi_models_from_text("\n".join(corpus)) | |
| normalized_requested = normalize_pi_model_name(requested_model) | |
| normalized_configured = normalize_pi_model_name(configured_model) | |
| normalized_observed = [normalize_pi_model_name(x) for x in observed] | |
| effective_model = "" | |
| for raw, normalized in zip(observed, normalized_observed): | |
| if normalized and normalized != normalized_requested and normalized != normalized_configured: | |
| effective_model = raw | |
| break | |
| if not effective_model and observed: | |
| effective_model = observed[0] | |
| mismatch = bool(effective_model and normalize_pi_model_name(effective_model) not in {normalized_requested, normalized_configured}) | |
| payload = { | |
| "requested_model": requested_model, | |
| "configured_model": configured_model, | |
| "observed_models": observed, | |
| "effective_model": effective_model or configured_model, | |
| "provider": "huggingface", | |
| "mismatch": mismatch, | |
| "source": "published_pi_traces_after_sync", | |
| "settings": settings, | |
| } | |
| write_json(run_dir / "pi_model_resolution.json", payload) | |
| return payload | |
| def emit_pi_model_resolution(events_path: Path, resolution: dict): | |
| if resolution.get("mismatch"): | |
| append_event( | |
| events_path, | |
| "pi_model_resolution", | |
| "warning", | |
| "Pi assistant model appears to differ from the requested model", | |
| { | |
| "requested_model": resolution.get("requested_model"), | |
| "configured_model": resolution.get("configured_model"), | |
| "effective_model": resolution.get("effective_model"), | |
| "observed_models": resolution.get("observed_models", [])[:8], | |
| "provider": resolution.get("provider"), | |
| }, | |
| ) | |
| else: | |
| append_event( | |
| events_path, | |
| "pi_model_resolution", | |
| "success", | |
| "Pi assistant model matched the requested configuration", | |
| { | |
| "requested_model": resolution.get("requested_model"), | |
| "configured_model": resolution.get("configured_model"), | |
| "effective_model": resolution.get("effective_model"), | |
| "provider": resolution.get("provider"), | |
| }, | |
| ) | |
| def collect_pi_traces(run_dir: Path, events_path: Path): | |
| count = sync_pi_traces(run_dir, emit_event=False) | |
| write_agent_trace_record(run_dir, phase="initial_build", event="pi_sessions_collected", status="success", message="Synchronized Pi session traces into the run-level trace folder", data={"count": count}, artifacts=["traces/raw/agent_trace.jsonl", "traces/redacted/agent_trace.jsonl"]) | |
| append_event(events_path, "traces", "success", "Collected Pi traces", {"count": count, "agent_trace": "traces/redacted/agent_trace.jsonl"}) | |
| write_artifact_manifest(run_dir, reason="pi_traces_collected") | |
| return count | |
| def sanitize_model_id(model_id: str) -> str: | |
| model_id = (model_id or DEFAULT_MODEL_ID).strip().replace("https://huggingface.co/", "") | |
| model_id = model_id.split("?", 1)[0].strip("/") | |
| if not re.match(r"^[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+$", model_id): | |
| raise ValueError("MODEL_ID must look like owner/model-name") | |
| return model_id | |
| def make_gradio_client(target_space_id: str, token: str): | |
| import inspect | |
| from gradio_client import Client | |
| params = inspect.signature(Client).parameters | |
| if "token" in params: | |
| return Client(target_space_id, token=token) | |
| if "hf_token" in params: | |
| return Client(target_space_id, hf_token=token) | |
| if "api_key" in params: | |
| return Client(target_space_id, api_key=token) | |
| if "headers" in params: | |
| return Client(target_space_id, headers={"Authorization": f"Bearer {token}"}) | |
| return Client(target_space_id) | |
| def api_names_from_schema(schema) -> list[str]: | |
| names: list[str] = [] | |
| if isinstance(schema, dict): | |
| endpoints = schema.get("named_endpoints") or schema.get("endpoints") or {} | |
| if isinstance(endpoints, dict): | |
| for key, value in endpoints.items(): | |
| if isinstance(key, str) and key.startswith("/"): | |
| names.append(key) | |
| if isinstance(value, dict): | |
| api_name = value.get("api_name") | |
| if isinstance(api_name, str) and api_name.startswith("/"): | |
| names.append(api_name) | |
| if isinstance(schema.get("dependencies"), list): | |
| for dep in schema["dependencies"]: | |
| if isinstance(dep, dict): | |
| api_name = dep.get("api_name") | |
| if isinstance(api_name, str): | |
| names.append(api_name if api_name.startswith("/") else f"/{api_name}") | |
| return list(dict.fromkeys(names)) | |
| def normalize_api_name(name: str | None) -> str: | |
| value = (name or "").strip() | |
| if not value: | |
| return "/generate" | |
| return value if value.startswith("/") else "/" + value | |
| def default_smoke_args(expected_output_type: str) -> list: | |
| expected = (expected_output_type or "any").lower() | |
| if expected == "image": | |
| return ["a cinematic robot cat astronaut, detailed, studio lighting"] | |
| if expected == "video": | |
| return ["a short cinematic shot of a robot cat astronaut walking on the moon"] | |
| if expected == "audio": | |
| return ["A calm voice saying hello from Agentic Space Factory."] | |
| return ["Explain the concept of agentic AI in one paragraph."] | |
| def endpoint_schema_for_api(schema, api_name: str): | |
| target = normalize_api_name(api_name) | |
| alternatives = {target, target.lstrip("/")} | |
| found = None | |
| def walk(obj): | |
| nonlocal found | |
| if found is not None: | |
| return | |
| if isinstance(obj, dict): | |
| for key, value in obj.items(): | |
| if isinstance(key, str) and key in alternatives and isinstance(value, (dict, list)): | |
| found = value | |
| return | |
| value = obj.get("api_name") or obj.get("apiName") | |
| if isinstance(value, str) and normalize_api_name(value) == target: | |
| found = obj | |
| return | |
| for value in obj.values(): | |
| walk(value) | |
| elif isinstance(obj, list): | |
| for item in obj: | |
| walk(item) | |
| walk(schema) | |
| return found | |
| def endpoint_parameter_names(endpoint) -> list[str]: | |
| names: list[str] = [] | |
| def add(value): | |
| if not isinstance(value, str): | |
| return | |
| cleaned = value.strip() | |
| if cleaned and cleaned not in names: | |
| names.append(cleaned) | |
| def from_parameter(param): | |
| if isinstance(param, dict): | |
| for key in ["parameter_name", "parameterName", "name", "label"]: | |
| add(param.get(key)) | |
| component = param.get("component") or param.get("component_type") | |
| if isinstance(component, dict): | |
| for key in ["label", "name"]: | |
| add(component.get(key)) | |
| elif isinstance(param, str): | |
| add(param) | |
| def walk(obj): | |
| if isinstance(obj, dict): | |
| params = obj.get("parameters") or obj.get("inputs") | |
| if isinstance(params, list): | |
| for param in params: | |
| from_parameter(param) | |
| for value in obj.values(): | |
| if isinstance(value, (dict, list)): | |
| walk(value) | |
| elif isinstance(obj, list): | |
| for item in obj: | |
| walk(item) | |
| walk(endpoint) | |
| return names | |
| def result_contains_expected_output(result, expected: str): | |
| expected = (expected or "any").lower() | |
| info = {"result_type": type(result).__name__, "result_repr": repr(result)[:2000]} | |
| paths = [] | |
| type_hints = [] | |
| text_values = [] | |
| def visit(obj): | |
| type_hints.append(type(obj).__name__.lower()) | |
| if isinstance(obj, (str, Path)): | |
| value = str(obj) | |
| if value: | |
| paths.append(value) | |
| text_values.append(value) | |
| elif isinstance(obj, dict): | |
| for key, value in obj.items(): | |
| key_l = str(key).lower() | |
| if key_l in {"path", "name", "url", "file", "filepath"}: | |
| visit(value) | |
| elif key_l in {"mime_type", "mime", "type", "format"} and value: | |
| type_hints.append(str(value).lower()) | |
| elif key_l in {"image", "video", "audio", "text"}: | |
| type_hints.append(key_l) | |
| visit(value) | |
| elif isinstance(value, (dict, list, tuple)): | |
| visit(value) | |
| elif isinstance(obj, (list, tuple)): | |
| for item in obj: | |
| visit(item) | |
| else: | |
| # PIL Images, numpy arrays returned through Gradio objects, etc. | |
| module = getattr(type(obj), "__module__", "") | |
| if module: | |
| type_hints.append(module.lower()) | |
| visit(result) | |
| lower_paths = [str(p).lower() for p in paths] | |
| hints = " ".join(type_hints + lower_paths) | |
| info["detected_paths"] = paths[:20] | |
| info["type_hints"] = type_hints[:30] | |
| if expected == "any": | |
| return result is not None, info | |
| image_ext = (".png", ".jpg", ".jpeg", ".webp", ".gif", ".bmp", ".tiff") | |
| video_ext = (".mp4", ".mov", ".webm", ".avi", ".mkv") | |
| audio_ext = (".wav", ".mp3", ".flac", ".ogg", ".m4a") | |
| if expected == "text": | |
| if isinstance(result, str) and bool(result.strip()): | |
| return True, info | |
| return any(isinstance(x, str) and x.strip() and not Path(x).suffix.lower() in image_ext + video_ext + audio_ext for x in text_values), info | |
| if expected == "image": | |
| return any(p.endswith(image_ext) for p in lower_paths) or any(h in hints for h in ["image", "pil.", "png", "jpeg", "jpg", "webp"]), info | |
| if expected == "video": | |
| return any(p.endswith(video_ext) for p in lower_paths) or any(h in hints for h in ["video", "mp4", "webm", "moviepy"]), info | |
| if expected == "audio": | |
| return any(p.endswith(audio_ext) for p in lower_paths) or any(h in hints for h in ["audio", "wav", "soundfile"]), info | |
| return result is not None, info | |
| def copy_result_artifacts(result, run_dir: Path): | |
| artifacts = run_dir / "artifacts" | |
| artifacts.mkdir(parents=True, exist_ok=True) | |
| copied = [] | |
| def maybe_copy(obj): | |
| if isinstance(obj, (str, Path)): | |
| path = Path(str(obj)) | |
| if path.exists() and path.is_file(): | |
| target = artifacts / path.name | |
| try: | |
| shutil.copy2(path, target) | |
| copied.append(str(target)) | |
| except Exception: | |
| pass | |
| elif isinstance(obj, dict): | |
| for key in ["path", "name"]: | |
| if key in obj: | |
| maybe_copy(obj[key]) | |
| for value in obj.values(): | |
| if isinstance(value, (dict, list, tuple)): | |
| maybe_copy(value) | |
| elif isinstance(obj, (list, tuple)): | |
| for item in obj: | |
| maybe_copy(item) | |
| maybe_copy(result) | |
| return copied | |
| def measured_zero_gpu_recommendation(latency_seconds: float | None): | |
| if latency_seconds is None: | |
| return { | |
| "observed_latency_seconds": None, | |
| "recommended_zero_gpu_duration_seconds": None, | |
| "recommendation_source": "not_measured", | |
| "recommendation_confidence": "none", | |
| } | |
| recommended = int(max(30, min(300, latency_seconds * 2 + 15))) | |
| return { | |
| "observed_latency_seconds": round(latency_seconds, 3), | |
| "recommended_zero_gpu_duration_seconds": recommended, | |
| "recommendation_source": "live_gradio_predict", | |
| "recommendation_confidence": "measured", | |
| "measurement_note": "Measured from a live gradio_client.predict call, including Gradio/API/network/result serialization overhead.", | |
| } | |
| def run_generation_smoke(target_space_id: str, token: str, run_dir: Path, events_path: Path, expected_output_type: str): | |
| api_name = "/generate" | |
| test_args = default_smoke_args(expected_output_type) | |
| test_kwargs = {} | |
| append_event(events_path, "generation_smoke", "started", "Calling live generation endpoint for ZeroGPU timing", {"api_name": api_name, "expected_output_type": expected_output_type}) | |
| client = make_gradio_client(target_space_id, token) | |
| schema = client.view_api(return_format="dict") | |
| discovered = api_names_from_schema(schema) | |
| if api_name not in discovered and discovered: | |
| non_health = [name for name in discovered if name != "/health"] | |
| if non_health: | |
| api_name = non_health[0] | |
| endpoint = endpoint_schema_for_api(schema, api_name) | |
| endpoint_parameters = endpoint_parameter_names(endpoint) | |
| write_json(run_dir / "tests" / "generation_api_schema.json", {"schema": schema, "api_names": discovered, "selected_api_name": api_name, "endpoint_parameters": endpoint_parameters}) | |
| started = time.time() | |
| result = client.predict(*test_args, api_name=api_name, **test_kwargs) | |
| latency = time.time() - started | |
| ok, info = result_contains_expected_output(result, expected_output_type) | |
| copied = copy_result_artifacts(result, run_dir) | |
| recommendation = measured_zero_gpu_recommendation(latency) | |
| payload = { | |
| "status": "success" if ok else "failed", | |
| "target_space": target_space_id, | |
| "api_name": api_name, | |
| "discovered_api_names": discovered, | |
| "endpoint_parameters": endpoint_parameters, | |
| "test_args": test_args, | |
| "test_kwargs": test_kwargs, | |
| "expected_output_type": expected_output_type, | |
| "latency_seconds": round(latency, 3), | |
| "result_info": info, | |
| "copied_artifacts": copied, | |
| "validated_at": now(), | |
| **recommendation, | |
| } | |
| write_json(run_dir / "tests" / "generation_smoke.json", payload) | |
| if ok: | |
| append_event(events_path, "generation_smoke", "success", "Live generation smoke test passed and ZeroGPU timing was measured", {"latency_seconds": payload["latency_seconds"], "recommended_zero_gpu_duration_seconds": payload["recommended_zero_gpu_duration_seconds"], "recommended_zerogpu_duration_seconds": payload["recommended_zero_gpu_duration_seconds"], "api_name": api_name}) | |
| else: | |
| append_event(events_path, "generation_smoke", "failed", "Live generation returned an unexpected output type", payload) | |
| return payload | |
| def space_subdomain_url(target_space_id: str) -> str: | |
| owner, name = target_space_id.split("/", 1) | |
| # This matches the common Spaces app URL pattern. Keep conservative: our | |
| # generated slugs are ASCII and hyphen-friendly. | |
| return f"https://{owner}-{name}.hf.space".replace("_", "-").lower() | |
| def runtime_to_dict(runtime) -> dict: | |
| payload = {} | |
| for attr in ["stage", "hardware", "requested_hardware", "sleep_time", "storage", "gc_timeout"]: | |
| value = getattr(runtime, attr, None) | |
| payload[attr] = getattr(value, "value", value) | |
| return {k: str(v) if v is not None else None for k, v in payload.items()} | |
| def write_space_runtime(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, attempt: int | None = None) -> dict: | |
| try: | |
| runtime = api.get_space_runtime(repo_id=target_space_id, token=token) | |
| payload = runtime_to_dict(runtime) | |
| payload["attempt"] = attempt | |
| write_json(run_dir / "space_runtime.json", payload) | |
| return payload | |
| except Exception as exc: | |
| payload = {"error": str(exc)[:2000], "attempt": attempt} | |
| write_json(run_dir / "space_runtime.json", payload) | |
| append_event(events_path, "space_runtime", "warning", "Could not fetch Space runtime", payload) | |
| return payload | |
| RESTART_GUARDRAIL_STATE = { | |
| "restart": {"count": 0, "last_ts": 0.0}, | |
| "factory_reboot": {"count": 0, "last_ts": 0.0}, | |
| "same_code_upload": {"count": 0, "last_ts": 0.0}, | |
| } | |
| RESTART_GUARDRAIL_LIMITS = { | |
| "restart": 1, | |
| "factory_reboot": 1, | |
| "same_code_upload": 1, | |
| } | |
| RESTART_GUARDRAIL_COOLDOWN_SECONDS = int(os.environ.get("SPACE_FACTORY_RESTART_COOLDOWN_SECONDS", "240")) | |
| def runtime_stage_value(runtime_payload: dict | None) -> str: | |
| return str((runtime_payload or {}).get("stage") or (runtime_payload or {}).get("runtime_stage") or "").lower() | |
| def restart_guardrail_allows(action: str, events_path: Path, *, reason: str = "") -> bool: | |
| state = RESTART_GUARDRAIL_STATE.setdefault(action, {"count": 0, "last_ts": 0.0}) | |
| limit = RESTART_GUARDRAIL_LIMITS.get(action, 1) | |
| now = time.time() | |
| if state["count"] >= limit: | |
| append_event(events_path, "factory_rebuild", "warning", f"Blocked {action}: restart budget exhausted", {"action": action, "limit": limit, "reason": reason}) | |
| return False | |
| if state["last_ts"] and now - state["last_ts"] < RESTART_GUARDRAIL_COOLDOWN_SECONDS: | |
| append_event(events_path, "factory_rebuild", "warning", f"Blocked {action}: restart cooldown active", {"action": action, "cooldown_seconds": RESTART_GUARDRAIL_COOLDOWN_SECONDS, "elapsed_seconds": round(now - state["last_ts"], 1), "reason": reason}) | |
| return False | |
| state["count"] += 1 | |
| state["last_ts"] = now | |
| return True | |
| def safe_restart_space(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, *, factory_reboot: bool = False, reason: str = "", require_logs_checked: bool = True) -> bool: | |
| """Restart a Space with strict guardrails. | |
| Hugging Face exposes HfApi.restart_space(repo_id=..., factory_reboot=True) | |
| for a factory rebuild and restart_space(..., factory_reboot=False) for a | |
| normal restart. The Factory only calls it through this wrapper so we do not | |
| create restart loops or mask actionable build errors. | |
| """ | |
| action = "factory_reboot" if factory_reboot else "restart" | |
| if require_logs_checked: | |
| logs_dir = run_dir / "logs" | |
| if not ((logs_dir / "space_logs_build.txt").exists() or (logs_dir / "space_logs_runtime.txt").exists()): | |
| append_event(events_path, "factory_rebuild", "warning", f"Skipped {action}: logs have not been collected yet", {"reason": reason}) | |
| return False | |
| runtime_before = write_space_runtime(api, target_space_id, token, run_dir, events_path) | |
| stage = runtime_stage_value(runtime_before) | |
| if stage in {"building", "buildqueued", "build_queued", "starting", "pending"}: | |
| append_event(events_path, "factory_rebuild", "warning", f"Skipped {action}: Space is already busy", {"runtime": runtime_before, "reason": reason}) | |
| return False | |
| if not restart_guardrail_allows(action, events_path, reason=reason): | |
| return False | |
| try: | |
| append_event(events_path, "factory_rebuild", "started", f"Requesting {'factory reboot' if factory_reboot else 'Space restart'} via Hugging Face API", {"runtime_before": runtime_before, "reason": reason}) | |
| result = api.restart_space(repo_id=target_space_id, token=token, factory_reboot=factory_reboot) | |
| payload = {"factory_reboot": factory_reboot, "reason": reason, "result": str(result)[:1000]} | |
| write_json(run_dir / ("space_factory_reboot.json" if factory_reboot else "space_restart.json"), payload) | |
| append_event(events_path, "factory_rebuild", "success", f"{'Factory reboot' if factory_reboot else 'Space restart'} requested", payload) | |
| return True | |
| except TypeError: | |
| # Older huggingface_hub versions may not expose factory_reboot. Do not | |
| # loop; fall back to a plain restart only for non-factory requests. | |
| if factory_reboot: | |
| append_event(events_path, "factory_rebuild", "warning", "Installed huggingface_hub does not support restart_space(factory_reboot=True)", {"reason": reason}) | |
| return False | |
| try: | |
| result = api.restart_space(repo_id=target_space_id, token=token) | |
| append_event(events_path, "factory_rebuild", "success", "Space restart requested", {"reason": reason, "result": str(result)[:1000]}) | |
| return True | |
| except Exception as exc: | |
| append_event(events_path, "factory_rebuild", "failed", "Space restart API call failed", {"reason": reason, **exception_payload(exc)}) | |
| return False | |
| except Exception as exc: | |
| append_event(events_path, "factory_rebuild", "failed", f"{'Factory reboot' if factory_reboot else 'Space restart'} API call failed", {"reason": reason, **exception_payload(exc)}) | |
| return False | |
| def safe_same_code_reupload(api, workspace: Path, target_space_id: str, token: str, run_dir: Path, events_path: Path, *, reason: str = "") -> bool: | |
| """Re-upload the current workspace once, with restart-loop guardrails.""" | |
| if not restart_guardrail_allows("same_code_upload", events_path, reason=reason): | |
| return False | |
| runtime_before = write_space_runtime(api, target_space_id, token, run_dir, events_path) | |
| stage = runtime_stage_value(runtime_before) | |
| if stage in {"building", "buildqueued", "build_queued", "starting", "pending"}: | |
| append_event(events_path, "factory_rebuild", "warning", "Skipped same-code upload: Space is already building/starting", {"runtime": runtime_before, "reason": reason}) | |
| return False | |
| upload_workspace(api, workspace, target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "factory_rebuild", "success", "Same-code workspace uploaded under restart guardrails", {"reason": reason, "runtime_before": runtime_before}) | |
| return True | |
| def collect_space_logs(target_space_id: str, token: str, run_dir: Path, events_path: Path): | |
| # Import locally because main() imports HfApi inside its own scope after | |
| # installing dependencies. Recovery/log collection can run from helper | |
| # functions where that local name is otherwise unavailable. | |
| from huggingface_hub import HfApi | |
| logs_dir = run_dir / "logs" | |
| logs_dir.mkdir(parents=True, exist_ok=True) | |
| written = [] | |
| api = HfApi(token=token) | |
| def write_log(filename: str, text: str, source: str): | |
| (logs_dir / filename).write_text(text or "", encoding="utf-8") | |
| written.append({"file": filename, "source": source, "returncode": 0, "tail": (text or "")[-1000:]}) | |
| try: | |
| # Prefer the Python API. Recent hf CLIs do not expose the old Space logs CLI command, | |
| # which produced misleading "invalid choice: spaces" logs. | |
| runtime_logs = api.get_space_runtime(repo_id=target_space_id, token=token) | |
| write_log("space_runtime_state.json", json.dumps(runtime_to_dict(runtime_logs), indent=2, ensure_ascii=False), "get_space_runtime") | |
| except Exception as exc: | |
| written.append({"file": "space_runtime_state.json", "source": "get_space_runtime", "error": str(exc)[:1000]}) | |
| for filename, factory in { | |
| "space_logs_runtime.txt": lambda: api.get_space_logs(repo_id=target_space_id, token=token, build=False), | |
| "space_logs_build.txt": lambda: api.get_space_logs(repo_id=target_space_id, token=token, build=True), | |
| }.items(): | |
| try: | |
| out = factory() | |
| if isinstance(out, (list, tuple)): | |
| out = "\n".join(str(x) for x in out) | |
| write_log(filename, str(out or ""), "HfApi.get_space_logs") | |
| except Exception as exc: | |
| # Keep a useful artifact instead of calling unsupported hf CLI commands. | |
| msg = f"Could not collect {filename} via HfApi.get_space_logs: {type(exc).__name__}: {str(exc)[:1500]}" | |
| (logs_dir / filename).write_text(msg, encoding="utf-8") | |
| written.append({"file": filename, "source": "HfApi.get_space_logs", "error": msg[-1000:]}) | |
| append_event(events_path, "space_logs", "success", "Collected best-effort Space logs", {"files": written}) | |
| write_artifact_manifest(run_dir, reason="space_logs_collected") | |
| return written | |
| def validate_http_health(target_space_id: str, token: str, run_dir: Path, events_path: Path, attempt: int): | |
| import requests | |
| base_url = space_subdomain_url(target_space_id) | |
| url = base_url.rstrip("/") + "/health" | |
| headers = {"Authorization": f"Bearer {token}", "Accept": "application/json,text/plain,*/*"} | |
| response = requests.get(url, headers=headers, timeout=20) | |
| payload = { | |
| "status": "success" if response.ok else "failed", | |
| "attempt": attempt, | |
| "url": url, | |
| "status_code": response.status_code, | |
| "content_type": response.headers.get("content-type"), | |
| "text": response.text[:2000], | |
| } | |
| if response.ok: | |
| try: | |
| payload["json"] = response.json() | |
| except Exception: | |
| pass | |
| write_json(run_dir / "tests" / "http_health.json", payload) | |
| write_json(run_dir / "tests" / "test_result.json", payload | {"validator": "http_get_health"}) | |
| append_event(events_path, "api_validation", "success", "HTTP /health validation passed", {"attempt": attempt, "url": url, "status_code": response.status_code}) | |
| return payload | {"validator": "http_get_health"} | |
| raise RuntimeError(f"HTTP /health returned {response.status_code}: {response.text[:500]}") | |
| def validate_gradio_api(target_space_id: str, token: str, run_dir: Path, events_path: Path, attempt: int): | |
| client = make_gradio_client(target_space_id, token) | |
| schema = client.view_api(return_format="dict") | |
| write_json(run_dir / "tests" / "api_schema.json", schema if isinstance(schema, dict) else {"schema": str(schema)}) | |
| discovered = api_names_from_schema(schema) | |
| candidates = [] | |
| for name in ["/health", "/predict", "/greet"] + discovered: | |
| if name not in candidates: | |
| candidates.append(name) | |
| errors = [] | |
| for api_name in candidates: | |
| try: | |
| if api_name == "/greet": | |
| result = client.predict("Agentic Space Factory", api_name=api_name) | |
| else: | |
| result = client.predict(api_name=api_name) | |
| payload = {"status": "success", "attempt": attempt, "api_name": api_name, "discovered_api_names": discovered, "result_repr": repr(result)[:2000], "validator": "gradio_client"} | |
| write_json(run_dir / "tests" / "test_result.json", payload) | |
| append_event(events_path, "api_validation", "success", "Gradio API validation passed", {"attempt": attempt, "api_name": api_name, "discovered_api_names": discovered}) | |
| return payload | |
| except Exception as exc: | |
| errors.append({"api_name": api_name, "error": str(exc)[:1000]}) | |
| raise RuntimeError("; ".join(f"{e['api_name']}: {e['error']}" for e in errors[:5]) or "No callable API endpoints found") | |
| def latest_collected_build_log(run_dir: Path) -> str: | |
| path = run_dir / "logs" / "space_logs_build.txt" | |
| return path.read_text(encoding="utf-8", errors="ignore") if path.exists() else "" | |
| def runtime_stage_is_build_error(stage: str) -> bool: | |
| value = (stage or "").upper() | |
| return "BUILD_ERROR" in value or value in {"BUILDING_ERROR", "ERROR", "FAILED", "BUILD_FAILED"} | |
| def runtime_stage_is_busy(stage: str) -> bool: | |
| value = (stage or "").upper() | |
| return any(marker in value for marker in ["BUILDING", "BUILDQUEUED", "BUILD_QUEUED", "STARTING", "PENDING"]) | |
| def build_log_has_terminal_error(text: str) -> bool: | |
| low = (text or "").lower() | |
| terminal_markers = [ | |
| "--> error:", | |
| "error: resolutionimpossible", | |
| "resolutionimpossible", | |
| "no matching distribution found", | |
| "could not find a version that satisfies the requirement", | |
| "did not complete successfully", | |
| "process \"/bin/sh -c pip install", | |
| "cannot install", | |
| "conflicting dependencies", | |
| "metadata-generation-failed", | |
| "exit code: 1", | |
| ] | |
| return any(marker in low for marker in terminal_markers) | |
| def raise_if_collected_build_error(target_space_id: str, token: str, run_dir: Path, events_path: Path, *, attempt: int, reason: str = "", runtime_payload: dict | None = None): | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| build_log = latest_collected_build_log(run_dir) | |
| stage = str((runtime_payload or {}).get("stage") or "") | |
| if runtime_stage_is_build_error(stage) or build_log_has_terminal_error(build_log): | |
| issue = extract_pip_dependency_issue(build_log) | |
| payload = { | |
| "attempt": attempt, | |
| "reason": reason, | |
| "runtime": runtime_payload or {}, | |
| "dependency_issue": issue, | |
| "tail": build_log[-4000:], | |
| } | |
| append_event(events_path, "api_validation", "failed", "Space is in build error; stopping wait and entering recovery", payload) | |
| write_json(run_dir / "build_error_observation.json", payload) | |
| raise RuntimeError(f"Space build failed according to runtime/logs. {json.dumps(issue, ensure_ascii=False)[:1000]} See build_error_observation.json and logs/space_logs_build.txt") | |
| def validate_live_api(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, timeout_s: int = 900): | |
| append_event(events_path, "api_validation", "started", "Waiting for live HTTP /health or Gradio API to become available") | |
| deadline = time.time() + timeout_s | |
| attempt = 0 | |
| last_error = None | |
| runtime_error_count = 0 | |
| while time.time() < deadline: | |
| attempt += 1 | |
| runtime_payload = write_space_runtime(api, target_space_id, token, run_dir, events_path, attempt) | |
| stage = str(runtime_payload.get("stage") or "").upper() | |
| # Never keep waiting when the Hub already reports a BUILD_ERROR. This | |
| # is the critical observer loop: manual refresh should not be required. | |
| if runtime_stage_is_build_error(stage): | |
| raise_if_collected_build_error(target_space_id, token, run_dir, events_path, attempt=attempt, reason=f"runtime_stage={stage}", runtime_payload=runtime_payload) | |
| # While the Space is building/starting, periodically inspect build logs. | |
| # A terminal pip/build error in logs is enough to enter recovery even if | |
| # runtime state has not propagated yet. | |
| if attempt == 1 or attempt % 2 == 0 or any(marker in stage for marker in ["BUILD", "ERROR", "RUNTIME"]): | |
| raise_if_collected_build_error(target_space_id, token, run_dir, events_path, attempt=attempt, reason=f"runtime_stage={stage}", runtime_payload=runtime_payload) | |
| if "RUNTIME_ERROR" in stage: | |
| runtime_error_count += 1 | |
| last_error = f"Space runtime stage is {stage}" | |
| if runtime_error_count >= 2: | |
| raise RuntimeError(f"Space is in RUNTIME_ERROR. See logs/space_logs_runtime.txt and logs/space_logs_build.txt. Last runtime: {runtime_payload}") | |
| try: | |
| return validate_http_health(target_space_id, token, run_dir, events_path, attempt) | |
| except Exception as exc: | |
| last_error = f"HTTP /health failed: {exc}" | |
| try: | |
| return validate_gradio_api(target_space_id, token, run_dir, events_path, attempt) | |
| except Exception as exc: | |
| last_error = (last_error or "") + f"; Gradio API failed: {exc}" | |
| append_event(events_path, "api_validation", "waiting", "Live health/API not ready yet", {"attempt": attempt, "runtime": runtime_payload, "error": last_error[-1500:] if last_error else None}) | |
| time.sleep(30) | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| raise RuntimeError(f"Live health/API validation did not pass before timeout: {last_error}") | |
| def is_auth_or_billing_like_error(error: str | None) -> bool: | |
| value = error or "" | |
| markers = [ | |
| "401", | |
| "402", | |
| "403", | |
| "Invalid username or password", | |
| "Unauthorized", | |
| "Repository Not Found", | |
| "payment", | |
| "billing", | |
| "quota", | |
| "grant", | |
| ] | |
| return any(marker.lower() in value.lower() for marker in markers) | |
| def request_hardware(api, target_space_id: str, hardware: str, token: str, events_path: Path, step: str, retries: int = 2): | |
| """Best-effort hardware request after Space creation. | |
| V23 tries hardware at create_repo time first. This function remains as a | |
| fallback for cases where a Space was created on CPU and the Hub later | |
| accepts a hardware switch. Auth/billing/quota errors are not retried. | |
| """ | |
| if not hardware: | |
| return {"phase": "post_create_request", "requested": False, "hardware": hardware, "ok": False, "error": "empty hardware"} | |
| last_error = None | |
| for attempt in range(1, retries + 1): | |
| try: | |
| runtime = api.request_space_hardware(repo_id=target_space_id, hardware=hardware, token=token) | |
| payload = { | |
| "phase": "post_create_request", | |
| "requested": True, | |
| "hardware": hardware, | |
| "ok": True, | |
| "attempt": attempt, | |
| "runtime_stage": getattr(getattr(runtime, "stage", None), "value", str(getattr(runtime, "stage", None))), | |
| "requested_hardware": getattr(runtime, "requested_hardware", None), | |
| "hardware_current": getattr(runtime, "hardware", None), | |
| } | |
| append_event(events_path, step, "success", f"Requested Space hardware {hardware}", payload) | |
| return payload | |
| except Exception as exc: | |
| exc_payload = exception_payload(exc) | |
| last_error = str(exc_payload.get("hf_error_message") or exc_payload.get("hf_error_detail") or exc_payload.get("error") or "")[:4000] | |
| auth_like = is_auth_or_billing_like_error(last_error) | |
| payload = {"phase": "post_create_request", "attempt": attempt, "hardware": hardware, "error": last_error, "manual_action_required": auth_like, **exc_payload} | |
| append_event(events_path, step, "failed" if auth_like or attempt == retries else "waiting", f"Could not request Space hardware {hardware}", payload) | |
| if auth_like: | |
| return {"phase": "post_create_request", "requested": True, "hardware": hardware, "ok": False, "attempts": attempt, "error": last_error, "manual_action_required": True, **exc_payload} | |
| if attempt < retries: | |
| time.sleep(8 * attempt) | |
| return {"phase": "post_create_request", "requested": True, "hardware": hardware, "ok": False, "attempts": retries, "error": last_error, "manual_action_required": False} | |
| def normalize_auto_space_hardware(value: str | None, default: str) -> str: | |
| candidate = (value or default).strip() | |
| return candidate if candidate in AUTO_SPACE_HARDWARE_CHOICES else default | |
| def build_hardware_sequence(preferred_hardware: str, fallback_hardware: str, allow_fixed_gpu_fallback: bool, try_zero_gpu_first: bool = True) -> list[str]: | |
| sequence = [] | |
| preferred_hardware = normalize_auto_space_hardware(preferred_hardware, DEFAULT_PREFERRED_SPACE_HARDWARE) | |
| fallback_hardware = normalize_auto_space_hardware(fallback_hardware, DEFAULT_FALLBACK_SPACE_HARDWARE) | |
| planned = [] | |
| if try_zero_gpu_first: | |
| planned.append("zero-a10g") | |
| planned.append(preferred_hardware) | |
| if allow_fixed_gpu_fallback: | |
| planned.append(fallback_hardware) | |
| for hw in planned: | |
| value = (hw or "").strip() | |
| if value and value not in sequence: | |
| sequence.append(value) | |
| return sequence | |
| def is_repo_already_exists_error(error: str) -> bool: | |
| text = str(error or "").lower() | |
| return ("already" in text and ("created this space repo" in text or "already exists" in text or "repo already" in text or "name already" in text)) | |
| def collision_safe_target_space_id(target_space_id: str, attempt: int) -> str: | |
| owner, slug = target_space_id.split("/", 1) | |
| base = re.sub(r"[^A-Za-z0-9._-]+", "-", slug).strip(".-_") or "space-factory" | |
| suffix = f"-r{int(time.time()) % 100000}-{attempt}" | |
| max_slug_len = min(95, 96 - len(owner) - 1) | |
| return f"{owner}/{base[:max_slug_len - len(suffix)].rstrip('.-_')}{suffix}" | |
| def create_space_with_hardware_strategy(api, target_space_id: str, token: str, preferred_hardware: str, fallback_hardware: str, allow_fixed_gpu_fallback: bool, events_path: Path, allow_auto_rename: bool = False, try_zero_gpu_first: bool = True): | |
| """Create a private Space and request hardware as early as possible. | |
| HF supports `space_hardware` directly on create_repo. This is the cleanest | |
| moment to request hardware because the Space does not need a second restart. | |
| If OAuth/billing/quota prevents automatic hardware selection, fall back to | |
| a normal private CPU Space and mark manual hardware as required. | |
| When the target Space name was generated automatically, a stale duplicate | |
| run id or a double-submit should not kill the build. In that one case, retry | |
| creation with a fresh suffix and return the final target Space id so the rest | |
| of the run, report and UI stay coherent. User-provided names still fail | |
| loudly if they already exist. | |
| """ | |
| sequence = build_hardware_sequence(preferred_hardware, fallback_hardware, allow_fixed_gpu_fallback, try_zero_gpu_first) | |
| attempts = [] | |
| current_target_space_id = target_space_id | |
| rename_attempt = 0 | |
| def maybe_rename_after_collision(error: str, phase: str, hardware: str | None = None) -> bool: | |
| nonlocal current_target_space_id, rename_attempt | |
| if not (allow_auto_rename and is_repo_already_exists_error(error) and rename_attempt < 5): | |
| return False | |
| previous = current_target_space_id | |
| rename_attempt += 1 | |
| current_target_space_id = collision_safe_target_space_id(previous, rename_attempt) | |
| append_event( | |
| events_path, | |
| "target_space_collision", | |
| "warning", | |
| "Auto-generated target Space already existed; retrying with a fresh generated name", | |
| {"phase": phase, "hardware": hardware, "previous_target_space": previous, "target_space": current_target_space_id}, | |
| ) | |
| return True | |
| def previous_create_attempt_for_current_target() -> bool: | |
| return any( | |
| str(item.get("target_space") or "") == current_target_space_id | |
| and str(item.get("phase") or "").startswith("create_repo_") | |
| for item in attempts | |
| if isinstance(item, dict) | |
| ) | |
| def adopt_existing_after_ambiguous_create(error: str, phase: str, hardware: str | None = None): | |
| """Keep one Space per run when a previous create may have succeeded. | |
| HF Hub calls can occasionally create the repo but still return an error | |
| to the client, or a hardware-at-creation ZeroGPU attempt can leave a CPU | |
| Space behind. In that situation the next hardware fallback sees | |
| "already exists" for the same target. That must not trigger an auto-rename, | |
| otherwise one run could create a second Space. | |
| """ | |
| if not (is_repo_already_exists_error(error) and previous_create_attempt_for_current_target()): | |
| return None | |
| payload = { | |
| "phase": "create_repo_existing_after_ambiguous_create", | |
| "hardware": hardware or "cpu-basic", | |
| "ok": True, | |
| "target_space": current_target_space_id, | |
| "source_error": error, | |
| "single_space_per_run": True, | |
| } | |
| append_event( | |
| events_path, | |
| "create_space", | |
| "warning", | |
| "Target Space already exists after an earlier create attempt; reusing it instead of generating another Space", | |
| payload, | |
| ) | |
| if hardware: | |
| hw_payload = request_hardware(api, current_target_space_id, hardware, token, events_path, "hardware_post_create_reuse") | |
| attempts.append({**payload, "post_create_hardware_request": hw_payload}) | |
| selected = hardware if hw_payload.get("ok") else "default-cpu-or-existing" | |
| manual = bool(hw_payload.get("manual_action_required") or not hw_payload.get("ok")) | |
| return {"created": True, "target_space_id": current_target_space_id, "selected_hardware": selected, "requested_sequence": sequence, "attempts": attempts, "manual_action_required": manual} | |
| attempts.append(payload) | |
| return {"created": True, "target_space_id": current_target_space_id, "selected_hardware": "default-cpu-or-existing", "requested_sequence": sequence, "attempts": attempts, "manual_action_required": True} | |
| for hardware in sequence: | |
| while True: | |
| try: | |
| append_event(events_path, "create_space_hardware", "started", f"Creating private Space with requested hardware {hardware}", {"target_space": current_target_space_id, "hardware": hardware}) | |
| api.create_repo( | |
| repo_id=current_target_space_id, | |
| repo_type="space", | |
| space_sdk="gradio", | |
| private=True, | |
| exist_ok=False, | |
| space_hardware=hardware, | |
| token=token, | |
| ) | |
| payload = {"phase": "create_repo_space_hardware", "hardware": hardware, "ok": True, "target_space": current_target_space_id} | |
| append_event(events_path, "create_space", "success", f"Private target Space created with requested hardware {hardware}", payload) | |
| return {"created": True, "target_space_id": current_target_space_id, "selected_hardware": hardware, "requested_sequence": sequence, "attempts": attempts + [payload], "manual_action_required": False} | |
| except Exception as exc: | |
| exc_payload = exception_payload(exc) | |
| error = str(exc_payload.get("hf_error_message") or exc_payload.get("hf_error_detail") or exc_payload.get("error") or "")[:4000] | |
| reused = adopt_existing_after_ambiguous_create(error, "create_repo_space_hardware", hardware) | |
| if reused: | |
| return reused | |
| if maybe_rename_after_collision(error, "create_repo_space_hardware", hardware): | |
| continue | |
| manual = is_auth_or_billing_like_error(error) | |
| payload = {"phase": "create_repo_space_hardware", "hardware": hardware, "ok": False, "error": error, "manual_action_required": manual, "target_space": current_target_space_id, **exc_payload} | |
| attempts.append(payload) | |
| append_event(events_path, "create_space_hardware", "failed", f"Could not create Space with requested hardware {hardware}", payload) | |
| break | |
| # Continue through the sequence: ZeroGPU quota/auth can fail while a fixed GPU | |
| # may still be worth trying. If fixed GPU also fails, we'll create CPU below. | |
| append_event(events_path, "create_space", "started", "Creating private target Space on default CPU after hardware-at-creation attempts failed", {"target_space": current_target_space_id}) | |
| while True: | |
| try: | |
| api.create_repo(repo_id=current_target_space_id, repo_type="space", space_sdk="gradio", private=True, exist_ok=False, token=token) | |
| cpu_payload = {"phase": "create_repo_default_cpu", "hardware": "cpu-basic", "ok": True, "target_space": current_target_space_id, "manual_action_required": True} | |
| append_event(events_path, "create_space", "success", "Private target Space created on default CPU; manual hardware selection may be required", cpu_payload) | |
| return {"created": True, "target_space_id": current_target_space_id, "selected_hardware": "default-cpu-or-existing", "requested_sequence": sequence, "attempts": attempts + [cpu_payload], "manual_action_required": True} | |
| except Exception as exc: | |
| exc_payload = exception_payload(exc) | |
| error = str(exc_payload.get("hf_error_message") or exc_payload.get("hf_error_detail") or exc_payload.get("error") or "")[:4000] | |
| reused = adopt_existing_after_ambiguous_create(error, "create_repo_default_cpu") | |
| if reused: | |
| return reused | |
| if maybe_rename_after_collision(error, "create_repo_default_cpu"): | |
| continue | |
| raise | |
| def create_initial_workspace(workspace: Path, model_id: str, target_space_id: str, preferred_hardware: str, fallback_hardware: str, allow_fallback: bool, implementation_mode: str, model_analysis: dict | None = None): | |
| workspace.mkdir(parents=True, exist_ok=True) | |
| model_analysis = model_analysis or {} | |
| pipeline_tag = model_analysis.get("pipeline_tag") | |
| library_name = model_analysis.get("library_name") | |
| tags = model_analysis.get("tags", [])[:40] | |
| siblings = model_analysis.get("siblings", [])[:60] | |
| app_py = f"""import gradio as gr | |
| from huggingface_hub import model_info, list_repo_files | |
| MODEL_ID = {model_id!r} | |
| TARGET_SPACE_ID = {target_space_id!r} | |
| def health(): | |
| return {{ | |
| "status": "booted", | |
| "model_id": MODEL_ID, | |
| "target_space_id": TARGET_SPACE_ID, | |
| "stage": "initial-scaffold", | |
| "note": "Pi should replace this scaffold with a model-specific demo while preserving a cheap health endpoint.", | |
| }} | |
| def placeholder(*args): | |
| return "Initial scaffold. Pi should replace this with a model-specific inference path, or write TECHNICAL_BLOCKERS.json." | |
| with gr.Blocks(title="Generated Model Space — Agentic Space Factory") as demo: | |
| gr.Markdown("# Generated Model Space — Agentic Space Factory") | |
| gr.Markdown(f"Private generated Space for `{{MODEL_ID}}`.") | |
| gr.JSON(label="Health", value=health(), every=None) | |
| gr.Button("Health check").click(fn=health, inputs=None, outputs=gr.JSON(), api_name="health") | |
| gr.Textbox(label="Input", value="Hello from Agentic Space Factory").submit(fn=placeholder, inputs=None, outputs=gr.Textbox(), api_name="predict") | |
| gr.Button("Run placeholder").click(fn=placeholder, inputs=None, outputs=gr.Textbox(), api_name="predict") | |
| if __name__ == "__main__": | |
| demo.launch() | |
| """ | |
| (workspace / "app.py").write_text(app_py, encoding="utf-8") | |
| req = """gradio>=6.0.0 | |
| huggingface_hub>=0.34.0,<2.0.0 | |
| spaces | |
| transformers>=4.45.0,<6.0.0 | |
| diffusers | |
| accelerate | |
| safetensors | |
| torch | |
| kernels | |
| pillow | |
| numpy | |
| requests | |
| """ | |
| (workspace / "requirements.txt").write_text(req, encoding="utf-8") | |
| readme = f"""--- | |
| title: Generated Model Space | |
| sdk: gradio | |
| app_file: app.py | |
| python_version: "3.10" | |
| suggested_hardware: {preferred_hardware or fallback_hardware or "cpu-basic"} | |
| short_description: "Generated model demo" | |
| --- | |
| # Generated Model Space — Agentic Space Factory | |
| Private generated Space for `{model_id}`. | |
| This Space is created by Agentic Space Factory. It should remain private until manually reviewed. | |
| """ | |
| (workspace / "README.md").write_text(readme, encoding="utf-8") | |
| analysis_json = json.dumps({"pipeline_tag": pipeline_tag, "library_name": library_name, "tags": tags, "siblings": siblings}, indent=2, ensure_ascii=False) | |
| goal = f"""You are Pi running inside a Hugging Face Job for Agentic Space Factory. | |
| Goal: build the best possible private Hugging Face Space demo for an arbitrary model card. | |
| MODEL_ID: {model_id} | |
| TARGET_SPACE_ID: {target_space_id} | |
| IMPLEMENTATION_MODE: {implementation_mode} | |
| MODEL_METADATA: | |
| ```json | |
| {analysis_json} | |
| ``` | |
| First read and follow the operational rules from this gist: | |
| {GIST_URL} | |
| {pi_tooling_context_note()} | |
| Non-negotiable safety and product constraints: | |
| - The target Space must remain private. | |
| - Do not delete any user resources. | |
| - Do not print secrets or tokens. | |
| - Work only inside the current workspace. | |
| - The wrapper will create the private Space, request allowed hardware best-effort, upload files, and validate the live app. Do not create/delete repos yourself in this builder worker. | |
| - Preserve a cheap health endpoint named `health` with `api_name="health"`. It must not load weights, run GPU work, or download large files. | |
| - Do not pin huggingface_hub below 1.0. Use huggingface_hub>=0.34.0,<2.0.0 unless the model card requires a narrower compatible range. If transformers>=5 is used, keep huggingface_hub compatible with it, for example huggingface_hub>=1.5.0,<2.0.0. | |
| - README.md frontmatter must remain valid; if it uses short_description, it must be 60 characters or fewer. | |
| Implementation contract: | |
| - If IMPLEMENTATION_MODE is `full-inference-gated`, you are not allowed to silently replace generation with a placeholder and call it success. | |
| - Try to implement the closest real inference path for the model card using evidence from README, model metadata, config files, and repo files. | |
| - You may choose an appropriate Gradio UI for the task: text, image, audio, video, multimodal, embeddings, classification, etc. | |
| - If the model is standard and feasible, implement a real generate/predict function and expose it as a Gradio endpoint. | |
| - If the model requires GPU, add ZeroGPU-compatible `@spaces.GPU(...)` only around the inference function. Do not decorate health. Do not assume dedicated A100/H200 hardware is available automatically; if such hardware is needed, document it as a manual requirement. | |
| - If the model requires special dependencies, include them only when needed and document risks. | |
| - Investigate compatibility fallbacks before declaring a blocker: PyTorch SDPA, xformers, HF Kernels where relevant, CPU/offload/lazy loading, smaller resolution/steps, safe smoke-test inputs. | |
| - If real inference is impossible or unsafe in a Space, write TECHNICAL_BLOCKERS.json with concrete evidence for every blocker. | |
| Deliverables: | |
| - app.py must boot on Hugging Face Spaces. | |
| - app.py must expose health/api_name="health". | |
| - If real generation is implemented, generate/predict must attempt a real model call, not only return a textual diagnostic. | |
| - If real generation is not implemented, write TECHNICAL_BLOCKERS.json with: full_inference_implemented=false, blockers[], evidence[], minimum_runtime, and suggested_next_step. | |
| - Write INFERENCE_CONTRACT.json with: full_inference_implemented, health_endpoint, primary_api_name, expected_output_type, validation_level, requires_gpu, estimated_vram, and blockers_count. | |
| - README.md must explain the runtime strategy, task, limitations, and how to test. | |
| - Write a concise PI_SUMMARY.md with what you changed and whether full inference is implemented. | |
| """ | |
| (workspace / "GOAL.md").write_text(goal, encoding="utf-8") | |
| return ["app.py", "requirements.txt", "README.md", "GOAL.md"] | |
| def sanitize_readme_metadata(workspace: Path, events_path: Path): | |
| """Ensure the generated Space README always keeps valid HF metadata. | |
| Pi is allowed to rewrite README.md, but Spaces require a valid YAML | |
| frontmatter block for sdk/app_file. If Pi drops it, the Space can fail with | |
| a Hub configuration error before the app even builds. Treat this as a | |
| factory invariant, not a Pi preference. | |
| """ | |
| readme_path = workspace / "README.md" | |
| if not readme_path.exists(): | |
| readme_path.write_text("# Generated Model Space\n", encoding="utf-8") | |
| text = readme_path.read_text(encoding="utf-8", errors="ignore") | |
| body = text | |
| metadata = {} | |
| changed = False | |
| if text.startswith("---"): | |
| parts = text.split("---", 2) | |
| if len(parts) >= 3: | |
| _, frontmatter, body = parts | |
| for line in frontmatter.splitlines(): | |
| if ":" not in line: | |
| continue | |
| key, value = line.split(":", 1) | |
| metadata[key.strip()] = value.strip() | |
| else: | |
| body = text | |
| changed = True | |
| else: | |
| changed = True | |
| required = { | |
| "title": metadata.get("title") or "Generated Model Space", | |
| "sdk": "gradio", | |
| "app_file": "app.py", | |
| "python_version": metadata.get("python_version") or "3.10", | |
| } | |
| suggested = metadata.get("suggested_hardware") | |
| if suggested: | |
| required["suggested_hardware"] = suggested | |
| short = metadata.get("short_description") or "Generated model demo" | |
| if len(short.strip('"\'')) > 60: | |
| short = "Generated model demo" | |
| required["short_description"] = short | |
| ordered_keys = ["title", "sdk", "app_file", "python_version", "suggested_hardware", "short_description"] | |
| lines = [] | |
| for key in ordered_keys: | |
| if key in required and required[key]: | |
| value = str(required[key]).strip() | |
| if key in {"title", "short_description"} and not (value.startswith('"') or value.startswith("'")): | |
| value = json.dumps(value, ensure_ascii=False) | |
| lines.append(f"{key}: {value}") | |
| normalized_body = body.lstrip("\n") or "# Generated Model Space\n" | |
| new_text = "---\n" + "\n".join(lines) + "\n---\n\n" + normalized_body | |
| if new_text != text: | |
| readme_path.write_text(new_text, encoding="utf-8") | |
| append_event(events_path, "metadata_sanitize", "success", "Ensured README Space metadata", {"metadata_keys": [k for k in ordered_keys if k in required]}) | |
| def normalize_requirements_for_modern_hub(workspace: Path, events_path: Path): | |
| """Normalize only broad known-dangerous base dependencies before upload. | |
| Do not try to solve every dependency conflict here. Pi is responsible for | |
| reading concrete build logs and patching requirements when a Space build | |
| fails. The Factory only prevents obviously unsafe broad ranges that would | |
| make the initial build non-deterministic, then observes the Space status and | |
| hands real build errors back to Pi with evidence. | |
| """ | |
| req_path = workspace / "requirements.txt" | |
| if not req_path.exists(): | |
| return | |
| raw = req_path.read_text(encoding="utf-8", errors="ignore") | |
| lines = [line.rstrip() for line in raw.splitlines()] | |
| prefix_lines: list[str] = [] | |
| package_lines: list[str] = [] | |
| for line in lines: | |
| stripped = line.strip() | |
| if not stripped: | |
| continue | |
| if stripped.startswith("--") or stripped.startswith("-f "): | |
| prefix_lines.append(line) | |
| else: | |
| package_lines.append(line) | |
| # Minimal base pins only. Do not globally pin diffusers here: newer model | |
| # cards may legitimately need a recent Diffusers release. If diffusers causes | |
| # a pip conflict, Pi must repair from the concrete build log. | |
| policy: dict[str, str] = { | |
| "huggingface-hub": "huggingface_hub>=0.34.0,<2.0.0", | |
| "transformers": "transformers>=4.51.0,<5.0.0", | |
| } | |
| aliases = { | |
| "huggingface_hub": "huggingface-hub", | |
| "huggingface-hub": "huggingface-hub", | |
| "transformers": "transformers", | |
| } | |
| seen_policy: set[str] = set() | |
| filtered: list[str] = [] | |
| changed = False | |
| for line in package_lines: | |
| stripped = line.strip() | |
| if stripped.startswith("#") or "://" in stripped or stripped.startswith((".", "/")): | |
| filtered.append(line) | |
| continue | |
| name = re.split(r"[<>=!~;\[]", stripped, 1)[0].strip().lower().replace("_", "-") | |
| canonical = aliases.get(name) | |
| if canonical in policy: | |
| if stripped != policy[canonical]: | |
| changed = True | |
| if canonical not in seen_policy: | |
| filtered.append(policy[canonical]) | |
| seen_policy.add(canonical) | |
| else: | |
| changed = True | |
| continue | |
| filtered.append(line) | |
| stable_policy_lines: list[str] = [] | |
| for canonical in ["huggingface-hub", "transformers"]: | |
| if canonical not in seen_policy: | |
| stable_policy_lines.append(policy[canonical]) | |
| changed = True | |
| new_lines = prefix_lines + stable_policy_lines + filtered | |
| new = "\n".join(line for line in new_lines if line.strip()) + "\n" | |
| if new != raw: | |
| changed = True | |
| if changed: | |
| req_path.write_text(new, encoding="utf-8") | |
| append_event( | |
| events_path, | |
| "requirements_sanitize", | |
| "success", | |
| "Normalized broad base dependencies; concrete pip conflicts remain Pi repair work", | |
| { | |
| "huggingface_hub": policy["huggingface-hub"], | |
| "transformers": policy["transformers"], | |
| "reason": "Avoid uncontrolled Transformers 5.x while preserving model-specific dependency choices for Pi to repair from build logs.", | |
| }, | |
| ) | |
| def useful_log_signals(text: str) -> list[str]: | |
| low = (text or "").lower() | |
| signals = [] | |
| patterns = { | |
| "python_exception": ["traceback", "modulenotfounderror", "importerror", "runtimeerror", "valueerror", "typeerror"], | |
| "dependency_resolution": ["resolutionimpossible", "no matching distribution", "could not find a version", "pip subprocess", "metadata-generation-failed"], | |
| "gradio_api": ["api_name", "endpoint", "unexpected keyword", "too many arguments", "not enough arguments", "no api found"], | |
| "model_loading": ["from_pretrained", "model_index", "safetensors", "does not appear to have", "failed to load", "pipeline"], | |
| "hardware_memory": ["cuda out of memory", "outofmemoryerror", "oom", "not enough memory"], | |
| "auth_access": ["401", "403", "gated", "unauthorized", "forbidden", "repository not found"], | |
| "zerogpu": ["zerogpu", "spaces.gpu", "duration"], | |
| } | |
| for name, needles in patterns.items(): | |
| if any(n in low for n in needles): | |
| signals.append(name) | |
| return signals | |
| def log_quality(build_log: str = "", runtime_log: str = "", failure_reason: str = "") -> str: | |
| combined = f"{build_log}\n{runtime_log}\n{failure_reason}".strip() | |
| low = combined.lower() | |
| if not combined or len(combined) < 80: | |
| return "empty" | |
| no_reason_markers = [ | |
| "no logs", "logs are empty", "no reason", "unknown error", "runtime error", "build error", "error: none", "status=error", | |
| ] | |
| if any(m in low for m in no_reason_markers) and not useful_log_signals(combined): | |
| return "no_reason" | |
| if useful_log_signals(combined): | |
| return "useful" | |
| return "partial" | |
| def classify_repair_failure(failure_reason: str, build_log: str = "", runtime_log: str = "") -> dict: | |
| """Classify a failed live validation before asking Pi to choose an action. | |
| This classifier is deliberately conservative. It gives Pi and the Factory | |
| evidence quality and a suggested class, but the diagnosis step must still | |
| make a bounded decision before any patch is allowed. | |
| """ | |
| text = f"{failure_reason}\n{build_log[-6000:]}\n{runtime_log[-6000:]}".lower() | |
| quality = log_quality(build_log, runtime_log, failure_reason) | |
| phase = "unknown" | |
| if "build" in text or "pip" in text or "container" in text: | |
| phase = "space_build" | |
| if "runtime" in text or "health" in text or "client.predict" in text or "api" in text: | |
| phase = "space_runtime" | |
| if "generation" in text or "smoke" in text or "expected output" in text: | |
| phase = "api_validation" | |
| checks = [ | |
| ("cuda_oom", ["cuda out of memory", "outofmemoryerror", "oom", "not enough memory"], "Prefer hardware/manual action or memory reduction; do not fake inference."), | |
| ("dependency_error", ["resolutionimpossible", "could not find a version", "no matching distribution", "dependency conflict", "pip subprocess", "metadata-generation-failed"], "Fix requirements with compatible, modern pins."), | |
| ("import_error", ["modulenotfounderror", "importerror", "cannot import name"], "Fix imports or requirements without hiding the error."), | |
| ("gradio_api_mismatch", ["no api found", "api_name", "endpoint", "unexpected keyword", "too many arguments", "not enough arguments"], "Preserve the expected Gradio endpoint and align input/output schema."), | |
| ("model_loading_error", ["from_pretrained", "safetensors", "model_index", "pipeline", "failed to load", "does not appear to have"], "Re-read the model card usage and fix the pipeline loading path."), | |
| ("hf_auth_error", ["401", "403", "gated", "unauthorized", "forbidden", "repository not found"], "Do not bypass auth; report gated/private access if needed."), | |
| ("space_boot_timeout", ["timed out", "timeout", "space did not become ready", "sleeping"], "Inspect logs first; if logs are empty, wait/rebuild same code before patching."), | |
| ("wrong_output_type", ["wrong output", "expected output", "invalid output", "not an image", "not text"], "Return the expected output type from real inference."), | |
| ("zero_gpu_duration_error", ["duration", "spaces.gpu", "zerogpu", "zero gpu"], "Tune @spaces.GPU(duration=...) from observed smoke-test latency when available."), | |
| ] | |
| category = "unknown_runtime_error" | |
| recommendation = "Use the Space logs, Gradio API schema and smoke error to choose the smallest truthful action." | |
| for code, needles, rec in checks: | |
| if any(n in text for n in needles): | |
| category = code | |
| recommendation = rec | |
| break | |
| if quality in {"empty", "no_reason"} and category == "unknown_runtime_error": | |
| category = "hf_runtime_flake" if phase != "space_build" else "hf_build_flake" | |
| recommendation = "Logs are not actionable. Prefer wait/inspect or a same-code factory rebuild before any code patch." | |
| dependency_issue = extract_pip_dependency_issue(f"{failure_reason}\n{build_log}\n{runtime_log}") | |
| if dependency_issue: | |
| category = "dependency_error" | |
| phase = "space_build" | |
| quality = "useful" | |
| recommendation = "Patch requirements.txt or dependency pins from the first pip resolver error; do not rebuild same code first." | |
| return { | |
| "category": category, | |
| "failure_phase": phase, | |
| "logs_quality": quality, | |
| "signals": useful_log_signals(text), | |
| "recommendation": recommendation, | |
| "dependency_issue": dependency_issue, | |
| "gist_alignment": "read first actionable pip error, patch minimally, rebuild, validate live" if dependency_issue else "standard blockage diagnosis", | |
| } | |
| def workspace_file_inventory(workspace: Path, max_files: int = 80) -> list[str]: | |
| files = [] | |
| for path in sorted(workspace.rglob("*")): | |
| if not path.is_file(): | |
| continue | |
| if any(part in {".git", "node_modules", "__pycache__"} for part in path.parts): | |
| continue | |
| try: | |
| files.append(str(path.relative_to(workspace))) | |
| except Exception: | |
| files.append(path.name) | |
| if len(files) >= max_files: | |
| break | |
| return files | |
| def write_incident_brief(workspace: Path, run_dir: Path, *, target_space_id: str, model_id: str, pi_model: str, failure_reason: str, build_log: str, runtime_log: str, classification: dict, implementation_mode: str, expected_output_type: str, iteration: int, budgets: dict | None = None) -> str: | |
| repair_dir = run_dir / "repair" | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| inventory = workspace_file_inventory(workspace) | |
| budgets = budgets or {} | |
| brief = f"""# Agentic Space Factory incident brief | |
| This is a blockage diagnosis pass for an existing generated Hugging Face Space. | |
| Pi must use the HF Spaces gist method: read logs first, act once, choose the cheapest useful iteration rung, and verify on the live Space. | |
| ## Target | |
| - Model ID: `{model_id}` | |
| - Target Space: `{target_space_id}` | |
| - Pi model: `{pi_model}` | |
| - Implementation mode: `{implementation_mode}` | |
| - Expected output type: `{expected_output_type}` | |
| - Diagnosis iteration: `{iteration}` | |
| ## Current classifier hints | |
| ```json | |
| {json.dumps(classification, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Dependency build error policy | |
| If `dependency_issue` is non-empty, this is an actionable pip build failure. | |
| Prefer `patch_code` with a minimal `requirements.txt` change. Do not choose | |
| `factory_rebuild_same_code` until a dependency patch has been attempted. | |
| ## HF tooling/token context | |
| ```json | |
| {json.dumps(load_json_if_exists(run_dir / "token_context.json"), indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Remaining action budget | |
| ```json | |
| {json.dumps(budgets, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Observed failure | |
| ```text | |
| {failure_reason[:5000]} | |
| ``` | |
| ## Build log tail | |
| ```text | |
| {build_log[-12000:]} | |
| ``` | |
| ## Runtime log tail | |
| ```text | |
| {runtime_log[-12000:]} | |
| ``` | |
| ## Workspace files | |
| ```text | |
| {chr(10).join(inventory)} | |
| ``` | |
| ## Allowed decisions | |
| You must choose exactly one action: | |
| - `wait_for_logs`: logs are empty/late; wait and collect again. | |
| - `inspect_more_logs`: enough uncertainty remains; collect logs/API status again without changing code. | |
| - `factory_rebuild_same_code`: likely HF build/runtime flake; re-upload the same workspace to force a rebuild, then revalidate. | |
| - `patch_code`: logs are actionable and point to a code/dependency/API/model-loading issue. | |
| - `request_manual_hardware`: issue is hardware/quota/restricted GPU/memory and should not be patched away. | |
| - `declare_technical_blocker`: no safe automated action remains. | |
| ## Hard rules | |
| - Do not choose `patch_code` when logs are empty or no-reason unless you cite concrete evidence from files/API schema. | |
| - Do not fake inference or replace model output with placeholders. | |
| - Do not patch hardware/quota/authorization failures as code bugs. | |
| - If logs are actionable, identify the first real error, not only the last line. | |
| - If using `factory_rebuild_same_code`, do not change files. | |
| ## Required output | |
| Write `REPAIR_DECISION.json` in the workspace root with this schema: | |
| ```json | |
| {{ | |
| "action": "wait_for_logs|inspect_more_logs|factory_rebuild_same_code|patch_code|request_manual_hardware|declare_technical_blocker", | |
| "confidence": "low|medium|high", | |
| "reason": "short explanation", | |
| "evidence": ["specific evidence"], | |
| "patch_allowed": false, | |
| "requires_manual_hardware": false | |
| }} | |
| ``` | |
| """ | |
| (repair_dir / "INCIDENT_BRIEF.md").write_text(brief, encoding="utf-8") | |
| (workspace / "INCIDENT_BRIEF.md").write_text(brief, encoding="utf-8") | |
| write_json(repair_dir / "classification.json", classification) | |
| return brief | |
| def extract_json_object(text: str) -> dict: | |
| if not text: | |
| return {} | |
| candidates = [] | |
| fenced = re.findall(r"```(?:json)?\s*(\{.*?\})\s*```", text, flags=re.S) | |
| candidates.extend(fenced) | |
| first = text.find("{") | |
| last = text.rfind("}") | |
| if first != -1 and last != -1 and last > first: | |
| candidates.append(text[first:last + 1]) | |
| for candidate in candidates: | |
| try: | |
| value = json.loads(candidate) | |
| if isinstance(value, dict): | |
| return value | |
| except Exception: | |
| continue | |
| return {} | |
| def extract_pip_dependency_issue(text: str = "") -> dict: | |
| """Extract the first actionable pip dependency resolver/build error. | |
| The point is not to solve every dependency automatically. It is to make | |
| dependency build failures first-class evidence so Pi gets a precise repair | |
| brief and the Factory does not choose wait/rebuild for a deterministic | |
| requirements problem. | |
| """ | |
| raw = text or "" | |
| low = raw.lower() | |
| issue: dict = {} | |
| patterns = [ | |
| ("no_matching_distribution", r"Could not find a version that satisfies the requirement\s+([^\s]+)(?:\s+\(from[^\n]*\))?"), | |
| ("no_matching_distribution", r"No matching distribution found for\s+([^\s]+)"), | |
| ("resolution_impossible", r"ResolutionImpossible[^\n]*"), | |
| ("metadata_generation_failed", r"metadata-generation-failed[^\n]*"), | |
| ("subprocess_exited", r"pip subprocess.*?did not run successfully[^\n]*"), | |
| ] | |
| for kind, pattern in patterns: | |
| match = re.search(pattern, raw, flags=re.I | re.S) | |
| if match: | |
| requirement = match.group(1).strip(" ,.;:") if match.groups() else "" | |
| package = re.split(r"[<>=!~;\[]", requirement, 1)[0].strip().lower().replace("_", "-") | |
| issue = { | |
| "kind": kind, | |
| "requirement": requirement, | |
| "package": package, | |
| "line": match.group(0)[:1000], | |
| } | |
| break | |
| if not issue and any(marker in low for marker in ["could not find a version", "no matching distribution", "resolutionimpossible", "metadata-generation-failed"]): | |
| issue = {"kind": "dependency_resolution", "requirement": "", "package": "", "line": raw[-2000:]} | |
| if issue: | |
| issue["category"] = "dependency_error" | |
| issue["actionable"] = True | |
| issue["recommended_action"] = "patch_code" | |
| issue["recommended_patch_scope"] = "requirements.txt or dependency pins only, unless the first error proves otherwise" | |
| return issue | |
| def write_dependency_error_brief(workspace: Path, run_dir: Path, events_path: Path, issue: dict, build_log: str = "", runtime_log: str = "", failure_reason: str = "") -> str: | |
| repair_dir = run_dir / "repair" | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| brief = f"""# Dependency build error brief | |
| The Space build failed during dependency installation. | |
| This brief does **not** replace the HF Spaces gist. It is a focused evidence | |
| appendix for applying the gist method: read logs first, find the first | |
| actionable build error, make one minimal requirements patch, rebuild, then | |
| validate live. | |
| ## First dependency error | |
| ```json | |
| {json.dumps(issue, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Gist-aligned repair guidance | |
| - Read the first real pip resolver error, not only the last log line. | |
| - Use the cheapest useful iteration rung from the gist. | |
| - For actionable pip errors, the cheapest useful action is normally a minimal | |
| `requirements.txt` / dependency pin patch, followed by rebuild and live | |
| validation. | |
| - Do not use `factory_rebuild_same_code` before a dependency patch when the | |
| build log already contains a concrete pip resolver error. | |
| - Do not change inference behavior to work around dependency resolution. | |
| - Keep the patch minimal and explain why the chosen versions are compatible. | |
| ## Build log tail | |
| ```text | |
| {build_log[-12000:]} | |
| ``` | |
| ## Runtime log tail | |
| ```text | |
| {runtime_log[-4000:]} | |
| ``` | |
| ## Failure reason | |
| ```text | |
| {failure_reason[:4000]} | |
| ``` | |
| """ | |
| (repair_dir / "DEPENDENCY_ERROR_BRIEF.md").write_text(brief, encoding="utf-8") | |
| (workspace / "DEPENDENCY_ERROR_BRIEF.md").write_text(brief, encoding="utf-8") | |
| append_event(events_path, "failure_diagnosis", "warning", "Detected pip dependency build error; applying gist log-first method", {"dependency_issue": issue, "artifact": "repair/DEPENDENCY_ERROR_BRIEF.md"}) | |
| write_agent_trace_record(run_dir, phase="diagnosis", event="dependency_error_detected", status="warning", message="Detected pip dependency build error; applying gist log-first method", data={"dependency_issue": issue}, artifacts=["repair/DEPENDENCY_ERROR_BRIEF.md"]) | |
| return brief | |
| def known_dependency_guardrail_needed(workspace: Path, failure_reason: str = "", build_log: str = "", runtime_log: str = "") -> bool: | |
| text = f"{failure_reason}\n{build_log}\n{runtime_log}" | |
| low = text.lower() | |
| issue = extract_pip_dependency_issue(text) | |
| # Keep automatic deterministic fixes extremely narrow. General pip conflicts | |
| # must be handed to Pi with DEPENDENCY_ERROR_BRIEF.md. | |
| if issue and issue.get("package") == "safetensors" and "transformers-5" in low and "safetensors>=0.8.0" in low: | |
| return True | |
| return False | |
| def apply_dependency_guardrail_repair(workspace: Path, run_dir: Path, events_path: Path, failure_reason: str = "", build_log: str = "", runtime_log: str = "") -> bool: | |
| """Apply deterministic requirements repair for known pip resolver failures. | |
| This handles build failures before Pi can produce a useful runtime patch. | |
| """ | |
| if not known_dependency_guardrail_needed(workspace, failure_reason, build_log, runtime_log): | |
| return False | |
| before = (workspace / "requirements.txt").read_text(encoding="utf-8", errors="ignore") if (workspace / "requirements.txt").exists() else "" | |
| normalize_requirements_for_modern_hub(workspace, events_path) | |
| after = (workspace / "requirements.txt").read_text(encoding="utf-8", errors="ignore") if (workspace / "requirements.txt").exists() else "" | |
| repair_dir = run_dir / "repair" | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| summary = """# Deterministic dependency repair | |
| The Space build failed during `pip install` because the resolver selected a known Transformers 5.x path requiring unavailable `safetensors>=0.8.0`. | |
| The Factory applied only the narrow known trap guardrail: | |
| - `transformers>=4.51.0,<5.0.0` | |
| - `huggingface_hub>=0.34.0,<2.0.0` | |
| Other dependency conflicts are passed to Pi through `DEPENDENCY_ERROR_BRIEF.md`. | |
| No inference code was changed. | |
| """ | |
| (repair_dir / "REPAIR_SUMMARY.md").write_text(summary, encoding="utf-8") | |
| (repair_dir / "REPAIR_PLAN.md").write_text("# Deterministic dependency repair plan\n\nNormalize requirements and re-upload the same workspace for a clean Space rebuild.\n", encoding="utf-8") | |
| write_agent_trace_record( | |
| run_dir, | |
| phase="repair_patch", | |
| event="dependency_guardrail_repair", | |
| status="success", | |
| message="Factory applied deterministic requirements repair for safetensors>=0.8.0 resolver failure", | |
| data={"before": before[-2000:], "after": after[-2000:]}, | |
| artifacts=["repair/REPAIR_SUMMARY.md", "repair/REPAIR_PLAN.md"], | |
| ) | |
| append_event( | |
| events_path, | |
| "repair_patch", | |
| "success", | |
| "Applied deterministic dependency guardrail for pip resolver build error", | |
| { | |
| "category": "dependency_error", | |
| "transformers": "transformers>=4.51.0,<5.0.0", | |
| "note": "General diffusers/safetensors conflicts are Pi repair work from DEPENDENCY_ERROR_BRIEF.md.", | |
| }, | |
| ) | |
| return before != after | |
| def fallback_decision_from_classifier(classification: dict, budgets: dict | None = None) -> dict: | |
| budgets = budgets or {} | |
| category = classification.get("category") or "unknown_runtime_error" | |
| quality = classification.get("logs_quality") or "partial" | |
| if category in {"hf_auth_error"}: | |
| action = "declare_technical_blocker" | |
| elif category == "cuda_oom": | |
| action = "request_manual_hardware" | |
| elif quality in {"empty", "no_reason"}: | |
| if budgets.get("wait_for_logs", 0) > 0: | |
| action = "wait_for_logs" | |
| elif budgets.get("factory_rebuild_same_code", 0) > 0: | |
| action = "factory_rebuild_same_code" | |
| else: | |
| action = "declare_technical_blocker" | |
| elif category in {"dependency_error", "import_error", "gradio_api_mismatch", "model_loading_error", "wrong_output_type", "zero_gpu_duration_error", "space_boot_timeout"}: | |
| action = "patch_code" if budgets.get("patch_code", 0) > 0 else "declare_technical_blocker" | |
| elif budgets.get("factory_rebuild_same_code", 0) > 0: | |
| action = "factory_rebuild_same_code" | |
| elif budgets.get("patch_code", 0) > 0 and quality == "useful": | |
| action = "patch_code" | |
| else: | |
| action = "declare_technical_blocker" | |
| return { | |
| "action": action, | |
| "confidence": "low", | |
| "reason": f"Factory fallback decision from classifier category={category}, logs_quality={quality}.", | |
| "evidence": classification.get("signals") or [], | |
| "patch_allowed": action == "patch_code", | |
| "requires_manual_hardware": action == "request_manual_hardware", | |
| "source": "factory_fallback", | |
| } | |
| def normalize_repair_decision(decision: dict, classification: dict, budgets: dict | None = None) -> dict: | |
| budgets = budgets or {} | |
| allowed = {"wait_for_logs", "inspect_more_logs", "factory_rebuild_same_code", "patch_code", "request_manual_hardware", "declare_technical_blocker"} | |
| action = str(decision.get("action") or "").strip().lower().replace("-", "_") | |
| if action not in allowed: | |
| decision = fallback_decision_from_classifier(classification, budgets) | |
| action = decision["action"] | |
| quality = classification.get("logs_quality") or "partial" | |
| category = classification.get("category") or "unknown_runtime_error" | |
| overrides = [] | |
| if action == "patch_code" and quality in {"empty", "no_reason"} and category in {"hf_runtime_flake", "hf_build_flake", "unknown_runtime_error"}: | |
| overrides.append("patch_code_blocked_without_actionable_logs") | |
| if budgets.get("wait_for_logs", 0) > 0: | |
| action = "wait_for_logs" | |
| elif budgets.get("factory_rebuild_same_code", 0) > 0: | |
| action = "factory_rebuild_same_code" | |
| else: | |
| action = "declare_technical_blocker" | |
| if category == "dependency_error" and quality == "useful" and action in {"wait_for_logs", "inspect_more_logs", "factory_rebuild_same_code"}: | |
| overrides.append("gist_log_first_dependency_error_needs_minimal_requirements_patch") | |
| # The gist says read the first actionable error and act once surgically. | |
| # For a concrete pip resolver error, a same-code rebuild/wait is not the | |
| # cheapest useful rung. Ask Pi to patch requirements minimally instead. | |
| action = "patch_code" if budgets.get("patch_code", 0) > 0 else "declare_technical_blocker" | |
| if action in budgets and budgets.get(action, 0) <= 0: | |
| overrides.append(f"budget_exhausted:{action}") | |
| action = "declare_technical_blocker" | |
| if action == "patch_code" and category in {"hf_auth_error"}: | |
| overrides.append("auth_failure_not_patchable") | |
| action = "declare_technical_blocker" | |
| if action == "patch_code" and category == "cuda_oom": | |
| # The agent may sometimes reduce memory pressure, but default to manual | |
| # hardware unless the logs are useful and it explicitly gave evidence. | |
| evidence = " ".join(str(x).lower() for x in decision.get("evidence") or []) | |
| if "dtype" not in evidence and "offload" not in evidence and "steps" not in evidence: | |
| overrides.append("oom_needs_hardware_or_explicit_memory_plan") | |
| action = "request_manual_hardware" | |
| normalized = { | |
| "action": action, | |
| "confidence": str(decision.get("confidence") or "low"), | |
| "reason": str(decision.get("reason") or "No reason provided.")[:2000], | |
| "evidence": list(decision.get("evidence") or [])[:20], | |
| "patch_allowed": action == "patch_code", | |
| "requires_manual_hardware": action == "request_manual_hardware", | |
| "factory_overrides": overrides, | |
| "classification": classification, | |
| } | |
| return normalized | |
| def diagnose_failure_with_pi(workspace: Path, run_dir: Path, events_path: Path, pi_model: str, target_space_id: str, model_id: str, failure_reason: str, implementation_mode: str, expected_output_type: str, iteration: int, budgets: dict) -> dict: | |
| logs_dir = run_dir / "logs" | |
| build_log = (logs_dir / "space_logs_build.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_build.txt").exists() else "" | |
| runtime_log = (logs_dir / "space_logs_runtime.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_runtime.txt").exists() else "" | |
| repair_dir = run_dir / "repair" | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| classification = classify_repair_failure(failure_reason, build_log, runtime_log) | |
| dependency_issue = classification.get("dependency_issue") or {} | |
| if dependency_issue: | |
| write_dependency_error_brief(workspace, run_dir, events_path, dependency_issue, build_log, runtime_log, failure_reason) | |
| append_event(events_path, "failure_diagnosis", "started", "Preparing Pi blockage diagnosis from Space status and logs", {"iteration": iteration, **classification}) | |
| brief = write_incident_brief( | |
| workspace, run_dir, | |
| target_space_id=target_space_id, model_id=model_id, pi_model=pi_model, | |
| failure_reason=failure_reason, build_log=build_log, runtime_log=runtime_log, | |
| classification=classification, implementation_mode=implementation_mode, expected_output_type=expected_output_type, | |
| iteration=iteration, budgets=budgets, | |
| ) | |
| goal = f"""You are Pi in BLOCKAGE DIAGNOSIS MODE for Agentic Space Factory. | |
| First read `INCIDENT_BRIEF.md` and the HF Spaces gist operational rules: {GIST_URL} | |
| {pi_tooling_context_note()} | |
| You are not allowed to edit code in this diagnosis step. Your task is to decide the next action for the Factory. | |
| Use the gist method: read logs first, identify the first actionable error, use the cheapest useful iteration rung, and require a live Gradio/API validation before success. | |
| Write `REPAIR_DECISION.json` exactly as requested in INCIDENT_BRIEF.md. Do not patch files during this step. | |
| """ | |
| (workspace / "PI_DIAGNOSIS_GOAL.md").write_text(goal, encoding="utf-8") | |
| (repair_dir / "PI_DIAGNOSIS_GOAL.md").write_text(goal, encoding="utf-8") | |
| append_event(events_path, "pi_diagnosis", "started", "Running Pi diagnosis decision before any repair action", {"model": pi_model, "iteration": iteration}) | |
| write_agent_trace_record(run_dir, phase="diagnosis", event="command_started", status="started", message="Pi blockage diagnosis started", data={"model": pi_model, "iteration": iteration}, artifacts=["repair/INCIDENT_BRIEF.md", "repair/PI_DIAGNOSIS_GOAL.md"]) | |
| code, out = run_cmd(["pi", "-p", goal], cwd=workspace, timeout=900) | |
| logs_dir.mkdir(parents=True, exist_ok=True) | |
| (logs_dir / "pi_diagnosis_output.txt").write_text(out, encoding="utf-8") | |
| append_agent_trace_artifact(run_dir, phase="diagnosis", event="diagnosis_output", artifact="logs/pi_diagnosis_output.txt", text=out, status="success" if code == 0 else "failed", data={"returncode": code}) | |
| raw_decision = {} | |
| decision_path = workspace / "REPAIR_DECISION.json" | |
| if decision_path.exists(): | |
| raw_decision = load_json_if_exists(decision_path) | |
| if not raw_decision: | |
| raw_decision = extract_json_object(out) | |
| if code != 0: | |
| append_event(events_path, "pi_diagnosis", "warning", "Pi diagnosis returned non-zero; using conservative factory fallback", {"returncode": code, "output_tail": out[-3000:]}) | |
| if not raw_decision: | |
| raw_decision = fallback_decision_from_classifier(classification, budgets) | |
| decision = normalize_repair_decision(raw_decision, classification, budgets) | |
| write_json(repair_dir / "REPAIR_DECISION.json", decision) | |
| (workspace / "REPAIR_DECISION.json").write_text(json.dumps(decision, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| write_agent_trace_record(run_dir, phase="diagnosis", event="repair_decision", status="success", message=f"Pi diagnosis selected action: {decision.get('action')}", data={"decision": decision}, artifacts=["repair/REPAIR_DECISION.json"]) | |
| append_event(events_path, "repair_decision", "success", f"Pi diagnosis selected action: {decision.get('action')}", {"decision": decision}) | |
| return decision | |
| def write_blockage_artifact(workspace: Path, run_dir: Path, events_path: Path, decision: dict, failure_reason: str, *, status: str = "technical_blocker"): | |
| blocker = { | |
| "full_inference_implemented": False, | |
| "source": "agentic_space_factory_blockage_protocol", | |
| "status": status, | |
| "decision": decision, | |
| "failure_reason": failure_reason[:4000], | |
| "blockers": [ | |
| { | |
| "type": decision.get("action") or status, | |
| "claim": decision.get("reason") or "Automated recovery could not safely continue.", | |
| "evidence": decision.get("evidence") or [], | |
| "severity": "blocking", | |
| } | |
| ], | |
| "suggested_next_step": "Inspect repair/INCIDENT_BRIEF.md and repair/REPAIR_DECISION.json, then retry manually or adjust hardware/settings if appropriate.", | |
| } | |
| (workspace / "TECHNICAL_BLOCKERS.json").write_text(json.dumps(blocker, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| gen_dir = run_dir / "generated" | |
| gen_dir.mkdir(parents=True, exist_ok=True) | |
| (gen_dir / "TECHNICAL_BLOCKERS.json").write_text(json.dumps(blocker, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| write_json(run_dir / "repair" / "BLOCKAGE.json", blocker) | |
| append_event(events_path, "technical_blocker", "failed", "Automated recovery stopped with an auditable blocker", {"decision": decision}) | |
| def write_repair_brief(workspace: Path, run_dir: Path, *, target_space_id: str, model_id: str, pi_model: str, failure_reason: str, build_log: str, runtime_log: str, classification: dict, implementation_mode: str, expected_output_type: str, decision: dict | None = None) -> str: | |
| repair_dir = run_dir / "repair" | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| inventory = workspace_file_inventory(workspace) | |
| decision = decision or {} | |
| brief = f"""# Agentic Space Factory repair brief | |
| ## Scope | |
| This is a structured patch pass for an existing generated Hugging Face Space. | |
| A separate Pi diagnosis already decided that `patch_code` is justified. Do not rebuild from scratch unless the current architecture is impossible to boot. | |
| ## Diagnosis decision | |
| ```json | |
| {json.dumps(decision, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Original target | |
| - Model ID: `{model_id}` | |
| - Target Space: `{target_space_id}` | |
| - Pi model: `{pi_model}` | |
| - Implementation mode: `{implementation_mode}` | |
| - Expected output type: `{expected_output_type}` | |
| ## Failure classification | |
| - Category: `{classification.get('category', 'unknown_runtime_error')}` | |
| - Logs quality: `{classification.get('logs_quality', 'partial')}` | |
| - Failure phase: `{classification.get('failure_phase', 'unknown')}` | |
| - Recommended strategy: {classification.get('recommendation', '')} | |
| ## HF tooling/token context | |
| ```json | |
| {json.dumps(load_json_if_exists(run_dir / "token_context.json"), indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Observed failure | |
| ```text | |
| {failure_reason[:5000]} | |
| ``` | |
| ## Build log tail | |
| ```text | |
| {build_log[-12000:]} | |
| ``` | |
| ## Runtime log tail | |
| ```text | |
| {runtime_log[-12000:]} | |
| ``` | |
| ## Workspace files | |
| ```text | |
| {chr(10).join(inventory)} | |
| ``` | |
| ## Repair constraints | |
| - Patch the current workspace; do not change the product goal. | |
| - Make the smallest patch that fixes the classified failure. | |
| - Preserve real inference for Strict inference mode. | |
| - Do not replace inference with static placeholders, fake images, canned text, or swallowed exceptions. | |
| - Do not change the model ID unless you write a technical blocker explaining why. | |
| - Preserve or restore a cheap `health` endpoint. | |
| - Preserve the expected Gradio API endpoint when possible. | |
| - Keep README metadata valid and short_description <= 60 chars. | |
| - Keep Hugging Face Hub requirements modern; do not pin below 1.0 unless unavoidable and explained. | |
| ## Required repair artifacts | |
| Before modifying files, write `REPAIR_PLAN.md` with: | |
| - root cause | |
| - strategy | |
| - files to change | |
| - risk level | |
| - whether hardware/manual action is required | |
| After modifying files, write `REPAIR_SUMMARY.md` with: | |
| - files changed | |
| - why the patch is minimal | |
| - how it preserves real inference | |
| - how to validate it | |
| """ | |
| (repair_dir / "REPAIR_BRIEF.md").write_text(brief, encoding="utf-8") | |
| (workspace / "REPAIR_BRIEF.md").write_text(brief, encoding="utf-8") | |
| write_json(repair_dir / "classification.json", classification) | |
| return brief | |
| def sanity_check_repair_workspace(workspace: Path, implementation_mode: str) -> tuple[bool, str]: | |
| app_path = workspace / "app.py" | |
| req_path = workspace / "requirements.txt" | |
| readme_path = workspace / "README.md" | |
| for required in [app_path, req_path, readme_path]: | |
| if not required.exists() or not required.read_text(encoding="utf-8", errors="ignore").strip(): | |
| return False, f"Repair removed or emptied required file: {required.name}" | |
| if implementation_mode == "full-inference-gated": | |
| app_text = app_path.read_text(encoding="utf-8", errors="ignore").lower() | |
| # Only block obvious fake-return implementations. Do not fail a valid | |
| # targeted runtime patch just because a comment/README still contains | |
| # words like "placeholder". The previous broad string scan blocked the | |
| # exact kind of repair Pi should be allowed to do. | |
| hard_fake_patterns = [ | |
| r"return\s+[\"'](?:placeholder|dummy image|fake inference|not implemented|canned response)", | |
| r"def\s+placeholder\s*\(", | |
| r"gr\.button\([^\n]*placeholder", | |
| ] | |
| if any(re.search(pattern, app_text) for pattern in hard_fake_patterns): | |
| return False, "Strict inference repair appears to introduce executable placeholder/fake inference markers." | |
| real_inference_markers = ["from_pretrained", "diffusionpipeline", "zimagepipeline", "pipe(", ".to(\"cuda\")", "@spaces.gpu"] | |
| if not any(marker in app_text for marker in real_inference_markers): | |
| return False, "Strict inference repair no longer shows a real model inference path." | |
| return True, "Repair workspace sanity checks passed." | |
| def repair_workspace_with_pi(workspace: Path, run_dir: Path, events_path: Path, pi_model: str, target_space_id: str, model_id: str, failure_reason: str, implementation_mode: str = "full-inference-gated", expected_output_type: str = "any", decision: dict | None = None): | |
| """Structured one-shot patch pass, only after Pi diagnosis chooses patch_code.""" | |
| logs_dir = run_dir / "logs" | |
| build_log = (logs_dir / "space_logs_build.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_build.txt").exists() else "" | |
| runtime_log = (logs_dir / "space_logs_runtime.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_runtime.txt").exists() else "" | |
| repair_dir = run_dir / "repair" | |
| before_dir = repair_dir / "before" | |
| after_dir = repair_dir / "after" | |
| if before_dir.exists(): | |
| shutil.rmtree(before_dir) | |
| repair_dir.mkdir(parents=True, exist_ok=True) | |
| shutil.copytree(workspace, before_dir, ignore=shutil.ignore_patterns(".git", "node_modules", "__pycache__", "*.pyc")) | |
| classification = (decision or {}).get("classification") or classify_repair_failure(failure_reason, build_log, runtime_log) | |
| append_event(events_path, "repair_diagnosis", "success", "Patch repair is allowed by Pi diagnosis", {"category": classification.get("category"), "decision": decision or {}}) | |
| brief = write_repair_brief( | |
| workspace, run_dir, | |
| target_space_id=target_space_id, model_id=model_id, pi_model=pi_model, | |
| failure_reason=failure_reason, build_log=build_log, runtime_log=runtime_log, | |
| classification=classification, implementation_mode=implementation_mode, expected_output_type=expected_output_type, | |
| decision=decision, | |
| ) | |
| append_event(events_path, "repair_brief", "success", "Repair brief generated for minimal patch", {"category": classification.get("category"), "brief": "repair/REPAIR_BRIEF.md"}) | |
| goal = f"""You are Pi in STRUCTURED REPAIR MODE / STRUCTURED PATCH REPAIR MODE for Agentic Space Factory. | |
| First read `REPAIR_BRIEF.md`, `INCIDENT_BRIEF.md`, `DEPENDENCY_ERROR_BRIEF.md` if present, and the HF Spaces gist operational rules: {GIST_URL} | |
| {pi_tooling_context_note()} | |
| You are continuing the same build run, not starting a separate project. | |
| This patch is allowed only because the diagnosis decision selected `patch_code`. | |
| If `DEPENDENCY_ERROR_BRIEF.md` exists, treat it as evidence for the gist method: identify the first pip error, patch dependency pins minimally, and do not modify inference code unless the dependency fix alone cannot address that first error. | |
| Use the available HF token context to inspect private Hub resources when needed, but never print or persist token values. | |
| Critical method: | |
| 1. Diagnose the root cause from the brief and current files. | |
| 2. Write `REPAIR_PLAN.md` before editing files. | |
| 3. Apply the smallest patch possible. | |
| 4. Write `REPAIR_SUMMARY.md` after editing. | |
| 5. Preserve real inference. Never fake outputs just to satisfy validation. | |
| Failure category: {classification.get('category')} | |
| Recommended strategy: {classification.get('recommendation')} | |
| Hard constraints: | |
| - Do not rebuild from scratch unless you clearly justify it in REPAIR_PLAN.md. | |
| - Do not change MODEL_ID `{model_id}` unless you write a technical blocker. | |
| - Do not remove the expected Gradio/API contract. | |
| - For Strict inference, do not replace inference with placeholders, static sample files, canned text, or broad try/except blocks that hide failures. | |
| - Keep or restore a cheap health endpoint with api_name="health". | |
| - Keep README metadata valid. | |
| - Do not publish anything. Work only in the current workspace. | |
| Required deliverables: REPAIR_PLAN.md, patched files, REPAIR_SUMMARY.md. | |
| """ | |
| (workspace / "REPAIR_GOAL.md").write_text(goal, encoding="utf-8") | |
| (repair_dir / "REPAIR_GOAL.md").write_text(goal, encoding="utf-8") | |
| append_event(events_path, "repair_plan", "started", "Running Pi minimal patch repair", {"model": pi_model, "category": classification.get("category")}) | |
| write_agent_trace_record(run_dir, phase="repair_patch", event="command_started", status="started", message="Pi structured patch repair started", data={"model": pi_model, "category": classification.get("category"), "decision": decision or {}}, artifacts=["repair/REPAIR_BRIEF.md", "repair/REPAIR_GOAL.md", "repair/REPAIR_DECISION.json"]) | |
| code, out = run_cmd(["pi", "-p", goal], cwd=workspace, timeout=1500) | |
| logs_dir.mkdir(parents=True, exist_ok=True) | |
| (logs_dir / "pi_repair_output.txt").write_text(out, encoding="utf-8") | |
| append_agent_trace_artifact(run_dir, phase="repair_patch", event="repair_output", artifact="logs/pi_repair_output.txt", text=out, status="success" if code == 0 else "failed", data={"returncode": code}) | |
| if code != 0: | |
| append_event(events_path, "repair_patch", "failed", "Pi patch repair returned a non-zero exit code", {"returncode": code, "output_tail": out[-3000:]}) | |
| return False | |
| plan_path = workspace / "REPAIR_PLAN.md" | |
| summary_path = workspace / "REPAIR_SUMMARY.md" | |
| if not plan_path.exists(): | |
| plan_path.write_text(f"# Repair plan\n\nPi did not create a separate plan file. Classified category: `{classification.get('category')}`.\n\nStrategy: {classification.get('recommendation')}\n", encoding="utf-8") | |
| if not summary_path.exists(): | |
| summary_path.write_text("# Repair summary\n\nPi completed a structured repair pass but did not write a separate summary file. See logs/pi_repair_output.txt.\n", encoding="utf-8") | |
| shutil.copy2(plan_path, repair_dir / "REPAIR_PLAN.md") | |
| shutil.copy2(summary_path, repair_dir / "REPAIR_SUMMARY.md") | |
| write_agent_trace_record(run_dir, phase="repair_patch", event="repair_plan_summary", status="success", message="Repair plan and summary artifacts are available", data={}, artifacts=["repair/REPAIR_PLAN.md", "repair/REPAIR_SUMMARY.md"]) | |
| append_event(events_path, "repair_plan", "success", "Repair plan and summary artifacts are available", {"plan": "repair/REPAIR_PLAN.md", "summary": "repair/REPAIR_SUMMARY.md"}) | |
| normalize_requirements_for_modern_hub(workspace, events_path) | |
| ok, sanity_message = sanity_check_repair_workspace(workspace, implementation_mode) | |
| if not ok: | |
| append_event(events_path, "repair_patch", "failed", sanity_message, {"category": classification.get("category")}) | |
| return False | |
| append_event(events_path, "repair_patch", "success", sanity_message, {"category": classification.get("category")}) | |
| if after_dir.exists(): | |
| shutil.rmtree(after_dir) | |
| shutil.copytree(workspace, after_dir, ignore=shutil.ignore_patterns(".git", "node_modules", "__pycache__", "*.pyc")) | |
| append_event(events_path, "repair", "success", "Structured repair patch completed; ready to re-upload and revalidate", {"output_tail": out[-3000:], "category": classification.get("category")}) | |
| return True | |
| def recover_after_live_validation_failure(api, workspace: Path, run_dir: Path, events_path: Path, *, pi_model: str, target_space_id: str, model_id: str, token: str, failure_reason: str, implementation_mode: str, expected_output_type: str): | |
| """Let Pi diagnose the blockage, then execute one bounded Factory action at a time. | |
| This is the core blockage protocol. Pi decides among allowed actions, but the | |
| Factory validates budgets and refuses unsafe patches without actionable logs. | |
| """ | |
| budgets = {"wait_for_logs": 1, "inspect_more_logs": 1, "factory_rebuild_same_code": 1, "patch_code": 1} | |
| current_error = failure_reason | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| logs_dir = run_dir / "logs" | |
| build_log = (logs_dir / "space_logs_build.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_build.txt").exists() else "" | |
| runtime_log = (logs_dir / "space_logs_runtime.txt").read_text(encoding="utf-8", errors="ignore") if (logs_dir / "space_logs_runtime.txt").exists() else "" | |
| dependency_issue = extract_pip_dependency_issue(f"{current_error}\n{build_log}\n{runtime_log}") | |
| if dependency_issue: | |
| write_dependency_error_brief(workspace, run_dir, events_path, dependency_issue, build_log, runtime_log, current_error) | |
| if apply_dependency_guardrail_repair(workspace, run_dir, events_path, current_error, build_log, runtime_log): | |
| append_event(events_path, "factory_rebuild", "started", "Re-uploading dependency-guardrailed workspace after pip resolver build error") | |
| if not safe_same_code_reupload(api, workspace, target_space_id, token, run_dir, events_path, reason="dependency_guardrail_repair"): | |
| raise RuntimeError("Dependency guardrail rebuild skipped by restart guardrails") | |
| append_event(events_path, "factory_rebuild", "success", "Dependency-guardrailed workspace uploaded; revalidating live Space") | |
| append_event(events_path, "repair_validation", "started", "Revalidating after deterministic dependency repair") | |
| try: | |
| validation = validate_live_api(api, target_space_id, token, run_dir, events_path, timeout_s=1200) | |
| append_event(events_path, "repair_validation", "success", "Deterministic dependency repair resolved the build blockage") | |
| return validation | |
| except Exception as exc: | |
| current_error = f"{current_error}\n\nDependency guardrail rebuild did not resolve validation: {str(exc)[:4000]}" | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "repair_validation", "failed", "Dependency guardrail rebuild did not resolve validation; falling back to Pi diagnosis", {"error": str(exc)[:4000]}) | |
| for iteration in range(1, 5): | |
| decision = diagnose_failure_with_pi( | |
| workspace, run_dir, events_path, pi_model, target_space_id, model_id, | |
| current_error, implementation_mode, expected_output_type, iteration, budgets, | |
| ) | |
| action = decision.get("action") | |
| if action in {"wait_for_logs", "inspect_more_logs"}: | |
| if budgets.get(action, 0) <= 0: | |
| decision["factory_overrides"] = list(decision.get("factory_overrides") or []) + [f"budget_exhausted:{action}"] | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error) | |
| raise RuntimeError(f"Automated recovery stopped after {action} budget was exhausted") | |
| budgets[action] -= 1 | |
| wait_seconds = int(os.environ.get("SPACE_FACTORY_FAILURE_LOG_WAIT_SECONDS", "45")) | |
| append_event(events_path, "wait_for_logs", "started", "Waiting for delayed HF Space logs before changing code", {"seconds": wait_seconds, "decision": decision}) | |
| time.sleep(max(1, wait_seconds)) | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "wait_for_logs", "success", "Collected Space logs after wait", {"remaining_budget": budgets}) | |
| current_error = f"{current_error}\n\nAfter wait_for_logs/inspect_more_logs, validation is still considered failed; re-diagnose with refreshed logs." | |
| continue | |
| if action == "factory_rebuild_same_code": | |
| if budgets.get(action, 0) <= 0: | |
| decision["factory_overrides"] = list(decision.get("factory_overrides") or []) + ["budget_exhausted:factory_rebuild_same_code"] | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error) | |
| raise RuntimeError("Automated recovery stopped after factory rebuild budget was exhausted") | |
| budgets[action] -= 1 | |
| append_event(events_path, "factory_rebuild", "started", "Re-uploading the same workspace to force a same-code Space rebuild", {"decision": decision}) | |
| if not safe_same_code_reupload(api, workspace, target_space_id, token, run_dir, events_path, reason="pi_decision_factory_rebuild_same_code"): | |
| decision["factory_overrides"] = list(decision.get("factory_overrides") or []) + ["same_code_upload_blocked_by_restart_guardrail"] | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error) | |
| raise RuntimeError("Same-code rebuild skipped by restart guardrails") | |
| append_event(events_path, "factory_rebuild", "success", "Same-code workspace re-uploaded; revalidating live Space") | |
| append_event(events_path, "repair_validation", "started", "Revalidating after same-code factory rebuild") | |
| try: | |
| validation = validate_live_api(api, target_space_id, token, run_dir, events_path, timeout_s=1200) | |
| append_event(events_path, "repair_validation", "success", "Same-code factory rebuild resolved the blockage") | |
| return validation | |
| except Exception as exc: | |
| current_error = f"{current_error}\n\nSame-code factory rebuild did not resolve validation: {str(exc)[:4000]}" | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "repair_validation", "failed", "Same-code factory rebuild did not resolve validation; re-diagnosing", {"error": str(exc)[:4000]}) | |
| continue | |
| if action == "patch_code": | |
| if budgets.get(action, 0) <= 0: | |
| decision["factory_overrides"] = list(decision.get("factory_overrides") or []) + ["budget_exhausted:patch_code"] | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error) | |
| raise RuntimeError("Automated recovery stopped after patch budget was exhausted") | |
| budgets[action] -= 1 | |
| append_event(events_path, "repair", "started", "Pi diagnosis allows a minimal code patch", {"decision": decision}) | |
| repaired = repair_workspace_with_pi(workspace, run_dir, events_path, pi_model, target_space_id, model_id, current_error, implementation_mode, expected_output_type, decision=decision) | |
| if not repaired: | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error, status="failed_after_repair") | |
| append_event(events_path, "failure", "failed", "Structured patch repair failed before redeploy", {"decision": decision}) | |
| raise RuntimeError("Structured patch repair failed before redeploy") | |
| append_event(events_path, "repair_upload", "started", "Uploading repaired workspace") | |
| upload_workspace(api, workspace, target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "repair_upload", "success", "Repaired workspace uploaded") | |
| append_event(events_path, "repair_validation", "started", "Revalidating repaired Space") | |
| try: | |
| validation = validate_live_api(api, target_space_id, token, run_dir, events_path, timeout_s=1200) | |
| append_event(events_path, "repair_validation", "success", "Repaired Space passed live API validation") | |
| return validation | |
| except Exception as exc: | |
| current_error = f"{current_error}\n\nPatch repair did not resolve validation: {str(exc)[:4000]}" | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "repair_validation", "failed", "Repair attempted, but validation still failed", {"error": str(exc)[:4000]}) | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error, status="failed_after_repair") | |
| append_event(events_path, "failure", "failed", "Run failed after structured repair attempt", {"repair_error": str(exc)[:4000]}) | |
| raise | |
| if action == "request_manual_hardware": | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error, status="manual_hardware_required") | |
| append_event(events_path, "manual_hardware_required", "failed", "Pi diagnosis requested manual hardware or quota action", {"decision": decision}) | |
| raise RuntimeError("Pi diagnosis requested manual hardware or quota action") | |
| write_blockage_artifact(workspace, run_dir, events_path, decision, current_error) | |
| append_event(events_path, "failure", "failed", "Pi diagnosis declared a technical blocker", {"decision": decision}) | |
| raise RuntimeError("Pi diagnosis declared a technical blocker") | |
| final_decision = {"action": "declare_technical_blocker", "reason": "Blockage protocol exhausted all bounded iterations.", "evidence": [], "patch_allowed": False} | |
| write_blockage_artifact(workspace, run_dir, events_path, final_decision, current_error) | |
| append_event(events_path, "failure", "failed", "Blockage protocol exhausted all bounded iterations") | |
| raise RuntimeError("Blockage protocol exhausted all bounded iterations") | |
| def upload_workspace(api, workspace: Path, target_space_id: str, token: str, run_dir: Path, events_path: Path): | |
| sanitize_readme_metadata(workspace, events_path) | |
| normalize_requirements_for_modern_hub(workspace, events_path) | |
| append_event(events_path, "upload_files", "started", "Uploading generated universal model-card workspace recursively") | |
| gen_dir = run_dir / "generated" | |
| if gen_dir.exists(): | |
| shutil.rmtree(gen_dir) | |
| shutil.copytree(workspace, gen_dir, ignore=internal_workspace_copy_ignore) | |
| for filename in ["app.py", "README.md", "requirements.txt"]: | |
| if not (workspace / filename).exists(): | |
| raise RuntimeError(f"Missing required generated file: {filename}") | |
| api.upload_folder( | |
| folder_path=str(workspace), | |
| repo_id=target_space_id, | |
| repo_type="space", | |
| token=token, | |
| ignore_patterns=internal_workspace_upload_ignore_patterns(), | |
| ) | |
| uploaded_files = sorted(str(p.relative_to(workspace)) for p in workspace.rglob("*") if p.is_file() and is_publishable_workspace_file(p.relative_to(workspace))) | |
| append_event(events_path, "upload_files", "success", "Uploaded generated workspace folder", {"file_count": len(uploaded_files), "files_sample": uploaded_files[:50]}) | |
| def load_json_if_exists(path: Path) -> dict: | |
| if not path.exists(): | |
| return {} | |
| try: | |
| return json.loads(path.read_text(encoding="utf-8", errors="replace")) | |
| except Exception as exc: | |
| return {"parse_error": str(exc), "raw_tail": path.read_text(encoding="utf-8", errors="replace")[-2000:]} | |
| def infer_generation_gate(workspace: Path, implementation_mode: str, validation: dict, generation_smoke: dict | None, run_dir: Path, events_path: Path) -> dict: | |
| """Classify the run separately from process success. | |
| /health passing means the Space boots. It does not mean the generated Space | |
| performs model inference. In full-inference-gated mode we require either | |
| an actual implementation signal or a machine-readable blocker report. | |
| """ | |
| app_text = (workspace / "app.py").read_text(encoding="utf-8", errors="ignore") if (workspace / "app.py").exists() else "" | |
| summary_text = (workspace / "PI_SUMMARY.md").read_text(encoding="utf-8", errors="ignore") if (workspace / "PI_SUMMARY.md").exists() else "" | |
| req_text = (workspace / "requirements.txt").read_text(encoding="utf-8", errors="ignore") if (workspace / "requirements.txt").exists() else "" | |
| blockers_path = workspace / "TECHNICAL_BLOCKERS.json" | |
| blockers = load_json_if_exists(blockers_path) | |
| combined = (app_text + "\n" + summary_text).lower() | |
| blocked_markers = [ | |
| "full generation is not implemented", | |
| "full generation is intentionally not wired", | |
| "full inference is blocked", | |
| "returns a detailed diagnostic", | |
| "diagnostic report instead", | |
| "placeholder generator", | |
| "placeholder generation", | |
| "info-only", | |
| "not implemented", | |
| "cannot run in this environment", | |
| "out of scope", | |
| ] | |
| blocker_detected = bool(blockers) or any(m in combined for m in blocked_markers) | |
| smoke_ok = isinstance(generation_smoke, dict) and generation_smoke.get("status") == "success" | |
| recommendation = generation_smoke if isinstance(generation_smoke, dict) else measured_zero_gpu_recommendation(None) | |
| implementation_signals = { | |
| "has_spaces_gpu": "@spaces.GPU" in app_text, | |
| "has_torch": "torch" in req_text or "import torch" in app_text, | |
| "has_diffusers": "diffusers" in req_text or "diffusers" in app_text, | |
| "has_video_output_hint": any(x in app_text.lower() for x in ["gr.video", "video", ".mp4", "ffmpeg"]), | |
| "health_passed": validation.get("method") in {"http_health", "gradio"}, | |
| "generation_smoke_passed": smoke_ok, | |
| "zero_gpu_duration_measured": recommendation.get("recommendation_confidence") == "measured", | |
| } | |
| if blocker_detected: | |
| status = "technical_blocker" | |
| message = "Space boots, but full model inference was not implemented. See TECHNICAL_BLOCKERS.json / PI_SUMMARY.md." | |
| elif implementation_mode in {"full-inference-gated", "full-inference-attempt"} and smoke_ok: | |
| status = "full_inference_success" | |
| message = "Space boots and a live generation smoke test passed. ZeroGPU duration recommendation was measured from real inference." | |
| elif implementation_mode in {"full-inference-gated", "full-inference-attempt"}: | |
| status = "full_inference_candidate_health_passed" | |
| message = "Space boots, but live generation smoke test did not produce a verified output. ZeroGPU duration recommendation was not measured." | |
| else: | |
| status = "health_only" | |
| message = "Safe scaffold health validation passed. Full inference was not requested." | |
| if blocker_detected and not blockers: | |
| blockers = { | |
| "full_inference_implemented": False, | |
| "source": "worker_heuristic_from_PI_SUMMARY_or_app.py", | |
| "blockers": [ | |
| { | |
| "type": "agent_declared_or_detected_blocker", | |
| "claim": "Pi-generated artifacts state that full inference is blocked/not implemented or generation returns diagnostics/placeholders.", | |
| "evidence": "See PI_SUMMARY.md and app.py in generated artifacts.", | |
| "severity": "blocking", | |
| } | |
| ], | |
| "required_investigations_for_next_run": [ | |
| "Check whether PyTorch SDPA can replace flash-attn calls.", | |
| "Check whether HF Kernels flash-attn2/3/4 can replace required flash-attn APIs.", | |
| "Verify whether 2-GPU context parallelism is strictly required or can be reduced to a single-GPU smoke test.", | |
| ], | |
| } | |
| (workspace / "TECHNICAL_BLOCKERS.json").write_text(json.dumps(blockers, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| (run_dir / "generated" / "TECHNICAL_BLOCKERS.json").write_text(json.dumps(blockers, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| gate = { | |
| "status": status, | |
| "message": message, | |
| "implementation_mode": implementation_mode, | |
| "blocker_detected": blocker_detected, | |
| "implementation_signals": implementation_signals, | |
| "validation_method": validation.get("method"), | |
| "generation_smoke": generation_smoke, | |
| "zero_gpu_duration_recommendation": { | |
| "observed_latency_seconds": recommendation.get("observed_latency_seconds"), | |
| "recommended_zero_gpu_duration_seconds": recommendation.get("recommended_zero_gpu_duration_seconds"), | |
| "recommendation_source": recommendation.get("recommendation_source"), | |
| "recommendation_confidence": recommendation.get("recommendation_confidence"), | |
| "measurement_note": recommendation.get("measurement_note"), | |
| }, | |
| "blockers": blockers, | |
| } | |
| write_json(run_dir / "inference_gate.json", gate) | |
| append_event(events_path, "inference_gate", status, message, gate) | |
| return gate | |
| def main(): | |
| run_id = os.environ["RUN_ID"] | |
| hf_username = os.environ.get("HF_USERNAME", "unknown") | |
| bucket_source = os.environ.get("BUCKET_SOURCE", "unknown") | |
| output_root = Path(os.environ.get("OUTPUT_ROOT", "/output")) | |
| target_space_id = os.environ.get("TARGET_SPACE_ID", "") | |
| target_space_auto_generated = os.environ.get("TARGET_SPACE_AUTO_GENERATED", "false").lower() in {"1", "true", "yes", "on"} | |
| model_id = sanitize_model_id(os.environ.get("MODEL_ID", DEFAULT_MODEL_ID)) | |
| pi_model = os.environ.get("PI_MODEL", "Qwen/Qwen3-Coder-Next") | |
| preferred_hardware = normalize_auto_space_hardware(os.environ.get("PREFERRED_SPACE_HARDWARE"), DEFAULT_PREFERRED_SPACE_HARDWARE) | |
| fallback_hardware = normalize_auto_space_hardware(os.environ.get("FALLBACK_SPACE_HARDWARE"), DEFAULT_FALLBACK_SPACE_HARDWARE) | |
| try_zero_gpu_first = os.environ.get("TRY_ZERO_GPU_FIRST", "true").lower() in {"1", "true", "yes", "on"} | |
| allow_fixed_gpu_fallback = os.environ.get("ALLOW_FIXED_GPU_FALLBACK", "true").lower() in {"1", "true", "yes", "on"} | |
| implementation_mode = os.environ.get("IMPLEMENTATION_MODE", "full-inference-attempt") | |
| expected_output_type = os.environ.get("EXPECTED_OUTPUT_TYPE", "any") | |
| token = os.environ.get("HF_TOKEN") | |
| run_dir = output_root / "runs" / run_id | |
| events_path = run_dir / "events.jsonl" | |
| state_path = run_dir / "state.json" | |
| workspace = Path("/tmp/universal_workspace") | |
| ensure_hf_token_context(run_dir, events_path) | |
| pi_model_resolution = {"requested_model": pi_model, "configured_model": pi_model, "effective_model": pi_model, "provider": "huggingface", "mismatch": False} | |
| append_event(events_path, "bootstrap", "started", "Universal model-card builder worker started", {"model_id": model_id, "target_space_id": target_space_id}) | |
| write_json(state_path, {"run_id": run_id, "kind": "universal_model_card_builder", "status": "running", "message": "Attempting Universal model-card builder Space creation", "model_id": model_id, "pi_model": pi_model, "pi_model_resolution": pi_model_resolution, "target_space": target_space_id, "created_by": hf_username, "bucket_source": bucket_source, "created_at": now(), "updated_at": now()}) | |
| if not token: | |
| fail(run_dir, events_path, "HF_TOKEN is missing from Job secrets") | |
| if not TARGET_RE.match(target_space_id): | |
| fail(run_dir, events_path, "Invalid TARGET_SPACE_ID", {"target_space_id": target_space_id}) | |
| try: | |
| install_python_deps(events_path) | |
| from huggingface_hub import HfApi | |
| api = HfApi(token=token) | |
| whoami = api.whoami(token=token) | |
| append_event(events_path, "auth", "success", "Authenticated inside Job", {"whoami_name": whoami.get("name")}) | |
| append_event(events_path, "model_analysis", "started", "Fetching model metadata", {"model_id": model_id}) | |
| info = api.model_info(model_id, token=token, files_metadata=True) | |
| siblings = [getattr(s, "rfilename", "") for s in (info.siblings or [])] | |
| analysis = {"model_id": model_id, "pipeline_tag": getattr(info, "pipeline_tag", None), "library_name": getattr(info, "library_name", None), "tags": list(getattr(info, "tags", []) or [])[:100], "siblings": siblings[:160], "default_model_target": model_id == DEFAULT_MODEL_ID, "preferred_hardware": preferred_hardware, "fallback_hardware": fallback_hardware, "allow_fixed_gpu_fallback": allow_fixed_gpu_fallback, "implementation_mode": implementation_mode} | |
| write_json(run_dir / "model_analysis.json", analysis) | |
| append_event(events_path, "model_analysis", "success", "Model metadata fetched", {"pipeline_tag": analysis["pipeline_tag"], "library_name": analysis["library_name"]}) | |
| create_initial_workspace(workspace, model_id, target_space_id, preferred_hardware, fallback_hardware, allow_fixed_gpu_fallback, implementation_mode, analysis) | |
| append_event(events_path, "workspace", "success", "Prepared universal model-card workspace", {"files": sorted(p.name for p in workspace.iterdir())}) | |
| install_pi(events_path) | |
| configure_pi(events_path, pi_model) | |
| append_event(events_path, "pi_run", "started", "Running Pi on universal model-card workspace", {"model": pi_model}) | |
| (run_dir / "logs").mkdir(parents=True, exist_ok=True) | |
| code, pi_out = run_cmd_streaming( | |
| ["pi", "-p", (workspace / "GOAL.md").read_text(encoding="utf-8")], | |
| cwd=workspace, | |
| timeout=2400, | |
| live_log_path=run_dir / "logs" / "pi_live_output.txt", | |
| events_path=events_path, | |
| run_dir=run_dir, | |
| step="pi_run", | |
| trace_phase="initial_build", | |
| ) | |
| (run_dir / "logs" / "pi_output.txt").write_text(pi_out, encoding="utf-8") | |
| if code != 0: | |
| append_event(events_path, "pi_run", "failed", "Pi returned a non-zero exit code", {"returncode": code, "output_tail": pi_out[-4000:]}) | |
| collect_pi_traces(run_dir, events_path) | |
| fail(run_dir, events_path, "Pi failed before Space upload", {"returncode": code, "output_tail": pi_out[-4000:]}) | |
| append_event(events_path, "pi_run", "success", "Pi completed universal model-card workspace pass", {"output_tail": pi_out[-2000:]}) | |
| # Pi sessions are the source of truth for the assistant/model actually used. | |
| # Sync them first, then resolve requested/configured/observed model identity. | |
| collect_pi_traces(run_dir, events_path) | |
| pi_model_resolution = detect_pi_model_resolution(pi_model, run_dir, pi_out) | |
| emit_pi_model_resolution(events_path, pi_model_resolution) | |
| if not (workspace / "PI_SUMMARY.md").exists(): | |
| (workspace / "PI_SUMMARY.md").write_text("# Pi Summary\n\nPi did not create a PI_SUMMARY.md. See logs/pi_output.txt.\n", encoding="utf-8") | |
| app_text = (workspace / "app.py").read_text(encoding="utf-8", errors="ignore") | |
| if "/health" not in app_text and "api_name=\"health\"" not in app_text and "api_name='health'" not in app_text: | |
| append_event(events_path, "pi_verification", "failed", "app.py does not appear to expose /health; injecting safe health endpoint is not implemented") | |
| fail(run_dir, events_path, "Pi output did not preserve a /health endpoint") | |
| append_event(events_path, "pi_verification", "success", "Pi output preserved health validation endpoint") | |
| append_event(events_path, "hardware_strategy", "started", "Creating Space with hardware-at-creation strategy", {"preferred_hardware": preferred_hardware, "fallback_hardware": fallback_hardware, "allow_fixed_gpu_fallback": allow_fixed_gpu_fallback}) | |
| hardware_strategy = create_space_with_hardware_strategy( | |
| api, | |
| target_space_id, | |
| token, | |
| preferred_hardware, | |
| fallback_hardware, | |
| allow_fixed_gpu_fallback, | |
| events_path, | |
| allow_auto_rename=target_space_auto_generated, | |
| try_zero_gpu_first=try_zero_gpu_first, | |
| ) | |
| target_space_id = hardware_strategy.get("target_space_id") or target_space_id | |
| selected_hardware = hardware_strategy.get("selected_hardware") or "default-cpu-or-existing" | |
| hardware_attempts = list(hardware_strategy.get("attempts") or []) | |
| requested_hardware_sequence = list(hardware_strategy.get("requested_sequence") or []) | |
| update_state(run_dir, { | |
| "run_id": run_id, | |
| "kind": "universal_model_card_builder", | |
| "status": "running", | |
| "message": "Private Space created; uploading generated workspace", | |
| "model_id": model_id, | |
| "pi_model": pi_model, | |
| "target_space": target_space_id, | |
| "target_space_url": f"https://huggingface.co/spaces/{target_space_id}", | |
| "selected_hardware": selected_hardware, | |
| "hardware_attempts": hardware_attempts, | |
| "requested_hardware_sequence": requested_hardware_sequence, | |
| "created_by": hf_username, | |
| "bucket_source": bucket_source, | |
| }) | |
| write_artifact_manifest(run_dir, reason="space_created") | |
| # Upload after create. If create_repo(space_hardware=...) succeeded, the build | |
| # starts directly on the requested hardware. If it fell back to CPU, the run | |
| # remains valid but will be marked manual_hardware_required when inference | |
| # signals indicate GPU is needed. | |
| upload_workspace(api, workspace, target_space_id, token, run_dir, events_path) | |
| write_artifact_manifest(run_dir, reason="workspace_uploaded") | |
| if selected_hardware == "default-cpu-or-existing": | |
| append_event(events_path, "hardware", "warning", "Automatic hardware-at-creation failed; Space is on default CPU unless user changes it manually", {"attempts": hardware_attempts}) | |
| write_json(run_dir / "hardware_attempts.json", {"selected_hardware": selected_hardware, "requested_sequence": requested_hardware_sequence, "attempts": hardware_attempts, "strategy": "create_repo_space_hardware_first", "try_zero_gpu_first": try_zero_gpu_first}) | |
| write_json(run_dir / "hardware_strategy.json", {"selected_hardware": selected_hardware, "requested_sequence": requested_hardware_sequence, "attempts": hardware_attempts, "manual_action_required": selected_hardware == "default-cpu-or-existing", "strategy": "create_repo_space_hardware_first", "try_zero_gpu_first": try_zero_gpu_first}) | |
| try: | |
| validation = validate_live_api(api, target_space_id, token, run_dir, events_path, timeout_s=1200) | |
| except Exception as validation_error: | |
| append_event(events_path, "failure_detected", "warning", "Initial live validation failed; entering Pi blockage diagnosis protocol", {"error": str(validation_error)[:2000]}) | |
| validation = recover_after_live_validation_failure( | |
| api, workspace, run_dir, events_path, | |
| pi_model=pi_model, target_space_id=target_space_id, model_id=model_id, token=token, | |
| failure_reason=str(validation_error), implementation_mode=implementation_mode, expected_output_type=expected_output_type, | |
| ) | |
| generation_smoke = None | |
| if implementation_mode in {"full-inference-gated", "full-inference-attempt"}: | |
| try: | |
| generation_smoke = run_generation_smoke(target_space_id, token, run_dir, events_path, expected_output_type) | |
| except Exception as smoke_error: | |
| generation_smoke = { | |
| "status": "failed", | |
| "target_space": target_space_id, | |
| "expected_output_type": expected_output_type, | |
| "error": str(smoke_error)[:4000], | |
| **measured_zero_gpu_recommendation(None), | |
| } | |
| write_json(run_dir / "tests" / "generation_smoke.json", generation_smoke) | |
| append_event(events_path, "generation_smoke", "failed", "Live generation smoke test failed; ZeroGPU duration was not measured", generation_smoke) | |
| else: | |
| generation_smoke = measured_zero_gpu_recommendation(None) | {"status": "skipped", "expected_output_type": expected_output_type} | |
| write_json(run_dir / "tests" / "generation_smoke.json", generation_smoke) | |
| inference_gate = infer_generation_gate(workspace, implementation_mode, validation, generation_smoke, run_dir, events_path) | |
| # If the generated app looks like real GPU inference but automatic | |
| # hardware requests failed, classify the run honestly as needing manual | |
| # hardware instead of pretending CPU/default hardware is enough. the existing-Space validation workflow | |
| # can then smoke-test generation after the user sets a GPU manually. | |
| manual_hw_required = selected_hardware == "default-cpu-or-existing" and inference_gate.get("status") not in {"technical_blocker", "health_only"} and ( | |
| inference_gate.get("implementation_signals", {}).get("has_spaces_gpu") | |
| or inference_gate.get("implementation_signals", {}).get("has_torch") | |
| or any((a.get("manual_action_required") for a in hardware_attempts if isinstance(a, dict))) | |
| ) | |
| if manual_hw_required: | |
| inference_gate = dict(inference_gate) | |
| inference_gate["status"] = "manual_hardware_required" | |
| inference_gate["message"] = "Space was generated and boots, but automatic ZeroGPU/fixed-GPU assignment failed. Set hardware manually, then run the existing-Space validation workflow." | |
| inference_gate["manual_hardware_required"] = True | |
| inference_gate["hardware_attempts"] = hardware_attempts | |
| write_json(run_dir / "inference_gate.json", inference_gate) | |
| append_event(events_path, "inference_gate", "manual_hardware_required", inference_gate["message"], inference_gate) | |
| collect_pi_traces(run_dir, events_path) | |
| final_state = { | |
| "run_id": run_id, | |
| "kind": "universal_model_card_builder", | |
| "status": inference_gate["status"], | |
| "message": inference_gate["message"], | |
| "model_id": model_id, | |
| "pi_model": pi_model, | |
| "target_space": target_space_id, | |
| "target_space_url": f"https://huggingface.co/spaces/{target_space_id}", | |
| "selected_hardware": selected_hardware, | |
| "hardware_attempts": hardware_attempts, | |
| "pi_model_resolution": pi_model_resolution, | |
| "validation": validation, | |
| "generation_smoke": generation_smoke, | |
| "inference_gate": inference_gate, | |
| "updated_at": now(), | |
| "created_by": hf_username, | |
| "bucket_source": bucket_source, | |
| } | |
| write_json(state_path, final_state) | |
| report = f"""# Agentic Space Factory — Universal Model-Card Builder Report | |
| Run ID: `{run_id}` | |
| Status: **{inference_gate['status']}** | |
| {inference_gate['message']} | |
| Target Space: https://huggingface.co/spaces/{target_space_id} | |
| Model: `{model_id}` | |
| ## Pi assistant model | |
| ```json | |
| {json.dumps(pi_model_resolution, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Hardware | |
| Selected/requested hardware: `{selected_hardware}` | |
| Hardware changes are best-effort with OAuth. If requests fail with 401/auth/billing errors, set the Space hardware manually and rerun validation. | |
| ```json | |
| {json.dumps(hardware_attempts, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Health validation | |
| The wrapper validated the live Space using HTTP `/health` first, with Gradio Client as fallback. This only proves bootability. | |
| ```json | |
| {json.dumps(validation, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Live generation smoke / ZeroGPU duration | |
| ```json | |
| {json.dumps(generation_smoke, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Full-inference gate | |
| ```json | |
| {json.dumps(inference_gate, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Pi summary | |
| {(workspace / 'PI_SUMMARY.md').read_text(encoding='utf-8', errors='ignore') if (workspace / 'PI_SUMMARY.md').exists() else 'No PI_SUMMARY.md was produced.'} | |
| ## Safety | |
| - The target Space was created private. | |
| - No public publication was attempted. | |
| - Raw traces should remain private; redacted traces are stored separately. | |
| - If fallback fixed GPU was used or selected manually, review billing/hardware settings manually after the run. | |
| """ | |
| (run_dir / "report.md").write_text(report, encoding="utf-8") | |
| append_event(events_path, "report_write", "success", "Wrote report.md") | |
| write_artifact_manifest(run_dir, events_path=events_path, reason="final") | |
| try: | |
| publish_eval_record(run_dir, phase="final", events_path=events_path) | |
| append_event(events_path, "anonymous_eval", "success", "Published anonymized evaluation record", {"enabled": True}) | |
| except Exception as eval_exc: | |
| append_event(events_path, "anonymous_eval", "warning", "Could not publish anonymized evaluation record", {"error": str(eval_exc)[:1000]}) | |
| append_event(events_path, "done", inference_gate["status"], "Universal model-card builder completed", {"target_space": target_space_id, "selected_hardware": selected_hardware, "gate_status": inference_gate["status"]}) | |
| except SystemExit: | |
| raise | |
| except Exception as exc: | |
| try: | |
| collect_pi_traces(run_dir, events_path) | |
| except Exception: | |
| pass | |
| fail(run_dir, events_path, "Universal model-card builder worker failed", {"error": str(exc)}) | |
| if __name__ == "__main__": | |
| main() | |
| ''' | |
| VALIDATE_EXISTING_SPACE_WORKER_SCRIPT = r''' | |
| import json | |
| import os | |
| import re | |
| import shutil | |
| import subprocess | |
| import sys | |
| import time | |
| from datetime import datetime, timezone | |
| from pathlib import Path | |
| TARGET_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{1,95}/[A-Za-z0-9][A-Za-z0-9._-]{1,95}$") | |
| def now(): | |
| return datetime.now(timezone.utc).isoformat() | |
| def write_json(path: Path, payload: dict): | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| path.write_text(json.dumps(payload, indent=2, ensure_ascii=False) + "\n", encoding="utf-8") | |
| def append_event(path: Path, step: str, status: str, message: str, data: dict | None = None): | |
| path.parent.mkdir(parents=True, exist_ok=True) | |
| event = {"ts": now(), "step": step, "status": status, "message": message, "data": data or {}} | |
| line = json.dumps(event, ensure_ascii=False) | |
| with path.open("a", encoding="utf-8") as f: | |
| f.write(line + "\n") | |
| print(line, flush=True) | |
| def redact_text(text: str | None) -> str: | |
| if not text: | |
| return "" | |
| value = text | |
| for secret_name in ["HF_TOKEN", "HUGGING_FACE_HUB_TOKEN"]: | |
| secret = os.environ.get(secret_name) | |
| if secret: | |
| value = value.replace(secret, "[REDACTED]") | |
| value = re.sub(r"Bearer\s+[A-Za-z0-9_\-.=]+", "Bearer [REDACTED]", value) | |
| value = re.sub(r"hf_[A-Za-z0-9_\-]{10,}", "hf_[REDACTED]", value) | |
| return value | |
| def eval_enabled() -> bool: | |
| return os.environ.get("ASF_EVAL_RECORD_ENABLED", "").strip().lower() in {"1", "true", "yes", "on"} | |
| def eval_safe_hash(value: str | None, prefix: str = "h") -> str: | |
| raw = (value or "").strip() | |
| salt = os.environ.get("ASF_EVAL_SALT", "").strip() | |
| if not raw: | |
| return "" | |
| if salt: | |
| digest = hmac.new(salt.encode("utf-8"), raw.encode("utf-8"), hashlib.sha256).hexdigest()[:24] | |
| else: | |
| digest = hashlib.sha256(raw.encode("utf-8")).hexdigest()[:24] | |
| return f"{prefix}_{digest}" | |
| def eval_load_json(path: Path) -> dict: | |
| if not path.exists(): | |
| return {} | |
| try: | |
| return json.loads(path.read_text(encoding="utf-8", errors="replace")) | |
| except Exception as exc: | |
| return {"_parse_error": str(exc)} | |
| def eval_load_events(path: Path) -> list[dict]: | |
| if not path.exists(): | |
| return [] | |
| events = [] | |
| for line in path.read_text(encoding="utf-8", errors="replace").splitlines(): | |
| try: | |
| evt = json.loads(line) | |
| if isinstance(evt, dict): | |
| events.append(evt) | |
| except Exception: | |
| pass | |
| return events | |
| def eval_parse_ts(value: str | None) -> float | None: | |
| if not value: | |
| return None | |
| try: | |
| return datetime.fromisoformat(str(value).replace("Z", "+00:00")).timestamp() | |
| except Exception: | |
| return None | |
| def eval_compact_timeline(events: list[dict]) -> list[dict]: | |
| stages: dict[str, dict] = {} | |
| order: list[str] = [] | |
| for evt in events: | |
| step = str(evt.get("step") or "unknown")[:80] | |
| status = str(evt.get("status") or "")[:80] | |
| ts = str(evt.get("ts") or "") | |
| tsv = eval_parse_ts(ts) | |
| if step not in stages: | |
| stages[step] = {"stage": step, "first_ts": ts, "last_ts": ts, "status": status, "event_count": 0} | |
| order.append(step) | |
| item = stages[step] | |
| item["last_ts"] = ts or item.get("last_ts") | |
| item["status"] = status or item.get("status") | |
| item["event_count"] = int(item.get("event_count") or 0) + 1 | |
| if tsv is not None: | |
| if "_first" not in item: | |
| item["_first"] = tsv | |
| item["_last"] = tsv | |
| timeline = [] | |
| for step in order: | |
| item = dict(stages[step]) | |
| first = item.pop("_first", None) | |
| last = item.pop("_last", None) | |
| if first is not None and last is not None and last >= first: | |
| item["duration_ms"] = int((last - first) * 1000) | |
| timeline.append(item) | |
| return timeline[-80:] | |
| def eval_redacted_tail(path: Path, limit: int = 2400) -> str: | |
| if os.environ.get("ASF_EVAL_INCLUDE_REDACTED_TAILS", "").strip().lower() not in {"1", "true", "yes", "on"}: | |
| return "" | |
| try: | |
| return redact_text(path.read_text(encoding="utf-8", errors="replace"))[-limit:] | |
| except Exception: | |
| return "" | |
| def publish_eval_record(run_dir: Path, *, phase: str, events_path: Path | None = None) -> dict: | |
| """Write an anonymized cross-run evaluation record to the optional eval bucket. | |
| This intentionally stores structured metrics and hashed identifiers, not raw | |
| generated code, prompts, tokens, user bucket paths, or target Space IDs. | |
| """ | |
| if not eval_enabled(): | |
| return {"enabled": False} | |
| # Jobs only write the anonymized eval record back into the user's run bucket. | |
| # The ASF backend later publishes this record to the operator eval bucket | |
| # mounted on the ASF Space. User Jobs must never mount or write the private | |
| # operator archive directly. | |
| run_id = os.environ.get("RUN_ID", run_dir.name) | |
| events = eval_load_events(events_path or (run_dir / "events.jsonl")) | |
| state = eval_load_json(run_dir / "state.json") | |
| analysis = eval_load_json(run_dir / "model_analysis.json") | |
| hardware_strategy = eval_load_json(run_dir / "hardware_strategy.json") | |
| hardware_attempts = eval_load_json(run_dir / "hardware_attempts.json") | |
| inference_gate = eval_load_json(run_dir / "inference_gate.json") | |
| generation_smoke = eval_load_json(run_dir / "tests" / "generation_smoke.json") | |
| space_runtime = eval_load_json(run_dir / "space_runtime.json") | |
| manifest = eval_load_json(run_dir / "artifact_manifest.json") | |
| repair_decision = eval_load_json(run_dir / "repair" / "REPAIR_DECISION.json") | |
| blockage = eval_load_json(run_dir / "repair" / "BLOCKAGE.json") | |
| pi_resolution = state.get("pi_model_resolution") if isinstance(state.get("pi_model_resolution"), dict) else {} | |
| status = str(state.get("status") or inference_gate.get("status") or phase or "unknown") | |
| failure_message = str(state.get("message") or (state.get("details") or {}).get("error") or "")[:1200] | |
| started_at = events[0].get("ts") if events else (state.get("created_at") or os.environ.get("LAUNCHED_AT") or "") | |
| finished_at = state.get("updated_at") or (events[-1].get("ts") if events else now()) | |
| model_id = os.environ.get("MODEL_ID") or state.get("model_id") or analysis.get("model_id") or "" | |
| target_space = state.get("target_space") or os.environ.get("TARGET_SPACE_ID") or "" | |
| include_model_id = os.environ.get("ASF_EVAL_INCLUDE_MODEL_ID", "").strip().lower() in {"1", "true", "yes", "on"} | |
| anon_run_id = eval_safe_hash(run_id, "run") | |
| record = { | |
| "schema_version": "1.0", | |
| "app_version": os.environ.get("ASF_VERSION", "unknown"), | |
| "phase": phase, | |
| "anonymous_run_id": anon_run_id, | |
| "run_id_hash": eval_safe_hash(run_id, "run"), | |
| "anonymous_user_hash": eval_safe_hash(os.environ.get("HF_USERNAME"), "user"), | |
| "salt_configured": bool(os.environ.get("ASF_EVAL_SALT")), | |
| "run_kind": state.get("kind") or ("validate_existing_space" if os.environ.get("API_NAME") else "validate_existing_space"), | |
| "started_at": started_at, | |
| "finished_at": finished_at, | |
| "input_model": { | |
| "model_id_redacted": not include_model_id, | |
| "model_id": model_id if include_model_id else "", | |
| "model_hash": eval_safe_hash(model_id, "model"), | |
| "pipeline_tag": analysis.get("pipeline_tag"), | |
| "library_name": analysis.get("library_name"), | |
| "default_model_target": bool(analysis.get("default_model_target")), | |
| }, | |
| "config": { | |
| "implementation_mode": os.environ.get("IMPLEMENTATION_MODE", ""), | |
| "expected_output_type": os.environ.get("EXPECTED_OUTPUT_TYPE", ""), | |
| "pi_model_requested": os.environ.get("PI_MODEL") or state.get("pi_model") or "", | |
| "pi_model_effective": pi_resolution.get("effective_model") or pi_resolution.get("observed_model") or "", | |
| "try_zero_gpu_first": os.environ.get("TRY_ZERO_GPU_FIRST", ""), | |
| "allow_fixed_gpu_fallback": os.environ.get("ALLOW_FIXED_GPU_FALLBACK", ""), | |
| "preferred_gpu": os.environ.get("PREFERRED_SPACE_HARDWARE", ""), | |
| "fallback_gpu": os.environ.get("FALLBACK_SPACE_HARDWARE", ""), | |
| }, | |
| "outcome": { | |
| "status": status, | |
| "message_redacted": redact_text(failure_message)[:1200], | |
| "target_space_hash": eval_safe_hash(target_space, "space"), | |
| "space_created": bool(target_space), | |
| "space_runtime_stage": space_runtime.get("stage"), | |
| "space_runtime_hardware": space_runtime.get("hardware") or space_runtime.get("requested_hardware"), | |
| }, | |
| "timeline": eval_compact_timeline(events), | |
| "hardware": { | |
| "selected": hardware_strategy.get("selected_hardware") or state.get("selected_hardware") or "", | |
| "strategy": hardware_strategy.get("strategy") or hardware_attempts.get("strategy") or "", | |
| "attempts": hardware_strategy.get("attempts") or hardware_attempts.get("attempts") or state.get("hardware_attempts") or [], | |
| "single_space_per_run": True, | |
| }, | |
| "pi": { | |
| "requested_model": os.environ.get("PI_MODEL") or state.get("pi_model") or "", | |
| "effective_model": pi_resolution.get("effective_model") or pi_resolution.get("observed_model") or "", | |
| "model_mismatch": bool(pi_resolution.get("mismatch")), | |
| "build_log_tail_redacted": eval_redacted_tail(run_dir / "logs" / "pi_output.txt"), | |
| "diagnosis_log_tail_redacted": eval_redacted_tail(run_dir / "logs" / "pi_diagnosis_output.txt"), | |
| }, | |
| "recovery": { | |
| "entered": any(str(e.get("step", "")).startswith(("repair", "failure_diagnosis", "pi_diagnosis", "technical_blocker")) for e in events), | |
| "decision": repair_decision.get("action") or repair_decision.get("decision") or "", | |
| "confidence": repair_decision.get("confidence") or "", | |
| "blockage": bool(blockage), | |
| "blockage_category": blockage.get("category") or blockage.get("reason") or "", | |
| }, | |
| "validation": { | |
| "generation_smoke_status": generation_smoke.get("status") or "", | |
| "latency_seconds": generation_smoke.get("latency_seconds") or generation_smoke.get("observed_latency_seconds"), | |
| "expected_output_type": generation_smoke.get("expected_output_type") or os.environ.get("EXPECTED_OUTPUT_TYPE", ""), | |
| "inference_gate_status": inference_gate.get("status") or "", | |
| "manual_hardware_required": bool(inference_gate.get("manual_hardware_required")), | |
| }, | |
| "artifacts": { | |
| "present_count": len(manifest.get("present_paths") or []), | |
| "present_paths": sorted(str(p) for p in (manifest.get("present_paths") or []))[:120], | |
| }, | |
| "privacy": { | |
| "tokens_stored": False, | |
| "raw_prompts_stored": False, | |
| "generated_code_stored": False, | |
| "target_space_redacted": True, | |
| "user_bucket_redacted": True, | |
| "redacted_tails_enabled": os.environ.get("ASF_EVAL_INCLUDE_REDACTED_TAILS", "").strip().lower() in {"1", "true", "yes", "on"}, | |
| }, | |
| } | |
| day = (finished_at or now())[:10] | |
| try: | |
| yyyy, mm, dd = day.split("-") | |
| except Exception: | |
| yyyy, mm, dd = now()[:10].split("-") | |
| dest = run_dir / "eval_record.json" | |
| write_json(dest, record) | |
| events_dest = run_dir / "events_compact.jsonl" | |
| events_dest.write_text("".join(json.dumps(item, ensure_ascii=False) + "\n" for item in record["timeline"]), encoding="utf-8") | |
| return {"enabled": True, "written": True, "path": str(dest), "publish_mode": "backend"} | |
| def ensure_hf_token_context(run_dir: Path, events_path: Path | None = None) -> dict: | |
| """Ensure Pi and HF tooling see the same private HF token safely. | |
| HF Jobs receive HF_TOKEN as a secret. Pi, huggingface_hub, gradio_client, | |
| and the hf CLI can use either HF_TOKEN or HUGGING_FACE_HUB_TOKEN depending | |
| on the tool/version, so expose both aliases inside the isolated Job process. | |
| Never write the token value to artifacts. | |
| """ | |
| token = os.environ.get("HF_TOKEN") or os.environ.get("HUGGING_FACE_HUB_TOKEN") or "" | |
| if token: | |
| os.environ.setdefault("HF_TOKEN", token) | |
| os.environ.setdefault("HUGGING_FACE_HUB_TOKEN", token) | |
| os.environ.setdefault("HF_HUB_ENABLE_HF_TRANSFER", "1") | |
| payload = { | |
| "hf_token_present": bool(token), | |
| "hf_token_length": len(token) if token else 0, | |
| "hugging_face_hub_token_alias_present": bool(os.environ.get("HUGGING_FACE_HUB_TOKEN")), | |
| "hf_username": os.environ.get("HF_USERNAME") or "", | |
| "bucket_source": os.environ.get("BUCKET_SOURCE") or "", | |
| "target_space_id": os.environ.get("TARGET_SPACE_ID") or "", | |
| "token_value": "[REDACTED]" if token else "", | |
| } | |
| run_dir.mkdir(parents=True, exist_ok=True) | |
| write_json(run_dir / "token_context.json", payload) | |
| if events_path: | |
| append_event(events_path, "token_context", "success" if token else "warning", "Verified HF token context for Pi/HF operations", {k: v for k, v in payload.items() if k != "token_value"}) | |
| return payload | |
| def pi_tooling_context_note() -> str: | |
| return """HF tooling context: | |
| - You are running inside an HF Job with an HF_TOKEN secret configured. | |
| - The Factory exposes both HF_TOKEN and HUGGING_FACE_HUB_TOKEN aliases for huggingface_hub, gradio_client, hf CLI, and Pi provider access. | |
| - You may use the token through normal tools when needed to inspect the model/Space or validate private resources. | |
| - Never print, write, echo, or commit token values. Do not include Authorization headers in artifacts. | |
| - If a Hub operation fails with 401/403/quota/billing, report it as evidence instead of retrying blindly. | |
| """ | |
| def run_cmd(cmd: list[str], *, env: dict | None = None, timeout: int = 120): | |
| result = subprocess.run(cmd, env=env, text=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, timeout=timeout) | |
| return result.returncode, redact_text(result.stdout) | |
| def install_deps(events_path: Path): | |
| append_event(events_path, "dependencies", "started", "Installing validation dependencies") | |
| code, out = run_cmd([sys.executable, "-m", "pip", "install", "-q", "--upgrade", "huggingface_hub>=1.0.0", "gradio_client>=2.0.0", "requests>=2.31.0"], timeout=600) | |
| if code != 0: | |
| append_event(events_path, "dependencies", "failed", "Dependency installation failed", {"output_tail": out[-4000:]}) | |
| raise RuntimeError(out) | |
| append_event(events_path, "dependencies", "success", "Validation dependencies installed") | |
| def make_gradio_client(target_space_id: str, token: str): | |
| import inspect | |
| from gradio_client import Client | |
| params = inspect.signature(Client).parameters | |
| if "token" in params: | |
| return Client(target_space_id, token=token) | |
| if "hf_token" in params: | |
| return Client(target_space_id, hf_token=token) | |
| if "api_key" in params: | |
| return Client(target_space_id, api_key=token) | |
| if "headers" in params: | |
| return Client(target_space_id, headers={"Authorization": f"Bearer {token}"}) | |
| return Client(target_space_id) | |
| def api_names_from_schema(schema) -> list[str]: | |
| names: list[str] = [] | |
| def add(name): | |
| if not isinstance(name, str) or not name: | |
| return | |
| if not name.startswith("/"): | |
| name = "/" + name | |
| if name not in names: | |
| names.append(name) | |
| def walk(obj): | |
| if isinstance(obj, dict): | |
| for k, v in obj.items(): | |
| if k in {"api_name", "apiName"}: | |
| add(v) | |
| if isinstance(k, str) and k.startswith("/"): | |
| add(k) | |
| walk(v) | |
| elif isinstance(obj, list): | |
| for item in obj: | |
| walk(item) | |
| walk(schema) | |
| return names | |
| def normalize_api_name(name: str | None) -> str: | |
| value = (name or "").strip() | |
| if not value: | |
| return "/generate" | |
| return value if value.startswith("/") else "/" + value | |
| def endpoint_schema_for_api(schema, api_name: str): | |
| target = normalize_api_name(api_name) | |
| alternatives = {target, target.lstrip("/")} | |
| found = None | |
| def walk(obj): | |
| nonlocal found | |
| if found is not None: | |
| return | |
| if isinstance(obj, dict): | |
| for key, value in obj.items(): | |
| if isinstance(key, str) and key in alternatives and isinstance(value, (dict, list)): | |
| found = value | |
| return | |
| value = obj.get("api_name") or obj.get("apiName") | |
| if isinstance(value, str) and normalize_api_name(value) == target: | |
| found = obj | |
| return | |
| for value in obj.values(): | |
| walk(value) | |
| elif isinstance(obj, list): | |
| for item in obj: | |
| walk(item) | |
| walk(schema) | |
| return found | |
| def endpoint_parameter_names(endpoint) -> list[str]: | |
| names: list[str] = [] | |
| def add(value): | |
| if not isinstance(value, str): | |
| return | |
| cleaned = value.strip() | |
| if cleaned and cleaned not in names: | |
| names.append(cleaned) | |
| def from_parameter(param): | |
| if isinstance(param, dict): | |
| for key in ["parameter_name", "parameterName", "name", "label"]: | |
| add(param.get(key)) | |
| component = param.get("component") or param.get("component_type") | |
| if isinstance(component, dict): | |
| for key in ["label", "name"]: | |
| add(component.get(key)) | |
| elif isinstance(param, str): | |
| add(param) | |
| def walk(obj): | |
| if isinstance(obj, dict): | |
| params = obj.get("parameters") or obj.get("inputs") | |
| if isinstance(params, list): | |
| for param in params: | |
| from_parameter(param) | |
| for value in obj.values(): | |
| if isinstance(value, (dict, list)): | |
| walk(value) | |
| elif isinstance(obj, list): | |
| for item in obj: | |
| walk(item) | |
| walk(endpoint) | |
| return names | |
| def sanitize_kwargs_for_schema(api_name: str, schema, kwargs: dict): | |
| endpoint = endpoint_schema_for_api(schema, api_name) | |
| parameter_names = endpoint_parameter_names(endpoint) | |
| if not parameter_names: | |
| return kwargs, {}, [] | |
| allowed = set(parameter_names) | |
| sanitized = {key: value for key, value in kwargs.items() if key in allowed} | |
| dropped = {key: value for key, value in kwargs.items() if key not in allowed} | |
| return sanitized, dropped, parameter_names | |
| def runtime_to_dict(runtime) -> dict: | |
| payload = {} | |
| for attr in ["stage", "hardware", "requested_hardware", "sleep_time", "storage", "gc_timeout"]: | |
| value = getattr(runtime, attr, None) | |
| payload[attr] = getattr(value, "value", value) | |
| return {k: str(v) if v is not None else None for k, v in payload.items()} | |
| def write_space_runtime(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, attempt: int | None = None) -> dict: | |
| try: | |
| runtime = api.get_space_runtime(repo_id=target_space_id, token=token) | |
| payload = runtime_to_dict(runtime) | |
| payload["attempt"] = attempt | |
| write_json(run_dir / "space_runtime.json", payload) | |
| return payload | |
| except Exception as exc: | |
| payload = {"error": str(exc)[:2000], "attempt": attempt} | |
| write_json(run_dir / "space_runtime.json", payload) | |
| append_event(events_path, "space_runtime", "warning", "Could not fetch Space runtime", payload) | |
| return payload | |
| RESTART_GUARDRAIL_STATE = { | |
| "restart": {"count": 0, "last_ts": 0.0}, | |
| "factory_reboot": {"count": 0, "last_ts": 0.0}, | |
| "same_code_upload": {"count": 0, "last_ts": 0.0}, | |
| } | |
| RESTART_GUARDRAIL_LIMITS = { | |
| "restart": 1, | |
| "factory_reboot": 1, | |
| "same_code_upload": 1, | |
| } | |
| RESTART_GUARDRAIL_COOLDOWN_SECONDS = int(os.environ.get("SPACE_FACTORY_RESTART_COOLDOWN_SECONDS", "240")) | |
| def runtime_stage_value(runtime_payload: dict | None) -> str: | |
| return str((runtime_payload or {}).get("stage") or (runtime_payload or {}).get("runtime_stage") or "").lower() | |
| def restart_guardrail_allows(action: str, events_path: Path, *, reason: str = "") -> bool: | |
| state = RESTART_GUARDRAIL_STATE.setdefault(action, {"count": 0, "last_ts": 0.0}) | |
| limit = RESTART_GUARDRAIL_LIMITS.get(action, 1) | |
| now = time.time() | |
| if state["count"] >= limit: | |
| append_event(events_path, "factory_rebuild", "warning", f"Blocked {action}: restart budget exhausted", {"action": action, "limit": limit, "reason": reason}) | |
| return False | |
| if state["last_ts"] and now - state["last_ts"] < RESTART_GUARDRAIL_COOLDOWN_SECONDS: | |
| append_event(events_path, "factory_rebuild", "warning", f"Blocked {action}: restart cooldown active", {"action": action, "cooldown_seconds": RESTART_GUARDRAIL_COOLDOWN_SECONDS, "elapsed_seconds": round(now - state["last_ts"], 1), "reason": reason}) | |
| return False | |
| state["count"] += 1 | |
| state["last_ts"] = now | |
| return True | |
| def safe_restart_space(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, *, factory_reboot: bool = False, reason: str = "", require_logs_checked: bool = True) -> bool: | |
| """Restart a Space with strict guardrails. | |
| Hugging Face exposes HfApi.restart_space(repo_id=..., factory_reboot=True) | |
| for a factory rebuild and restart_space(..., factory_reboot=False) for a | |
| normal restart. The Factory only calls it through this wrapper so we do not | |
| create restart loops or mask actionable build errors. | |
| """ | |
| action = "factory_reboot" if factory_reboot else "restart" | |
| if require_logs_checked: | |
| logs_dir = run_dir / "logs" | |
| if not ((logs_dir / "space_logs_build.txt").exists() or (logs_dir / "space_logs_runtime.txt").exists()): | |
| append_event(events_path, "factory_rebuild", "warning", f"Skipped {action}: logs have not been collected yet", {"reason": reason}) | |
| return False | |
| runtime_before = write_space_runtime(api, target_space_id, token, run_dir, events_path) | |
| stage = runtime_stage_value(runtime_before) | |
| if stage in {"building", "buildqueued", "build_queued", "starting", "pending"}: | |
| append_event(events_path, "factory_rebuild", "warning", f"Skipped {action}: Space is already busy", {"runtime": runtime_before, "reason": reason}) | |
| return False | |
| if not restart_guardrail_allows(action, events_path, reason=reason): | |
| return False | |
| try: | |
| append_event(events_path, "factory_rebuild", "started", f"Requesting {'factory reboot' if factory_reboot else 'Space restart'} via Hugging Face API", {"runtime_before": runtime_before, "reason": reason}) | |
| result = api.restart_space(repo_id=target_space_id, token=token, factory_reboot=factory_reboot) | |
| payload = {"factory_reboot": factory_reboot, "reason": reason, "result": str(result)[:1000]} | |
| write_json(run_dir / ("space_factory_reboot.json" if factory_reboot else "space_restart.json"), payload) | |
| append_event(events_path, "factory_rebuild", "success", f"{'Factory reboot' if factory_reboot else 'Space restart'} requested", payload) | |
| return True | |
| except TypeError: | |
| # Older huggingface_hub versions may not expose factory_reboot. Do not | |
| # loop; fall back to a plain restart only for non-factory requests. | |
| if factory_reboot: | |
| append_event(events_path, "factory_rebuild", "warning", "Installed huggingface_hub does not support restart_space(factory_reboot=True)", {"reason": reason}) | |
| return False | |
| try: | |
| result = api.restart_space(repo_id=target_space_id, token=token) | |
| append_event(events_path, "factory_rebuild", "success", "Space restart requested", {"reason": reason, "result": str(result)[:1000]}) | |
| return True | |
| except Exception as exc: | |
| append_event(events_path, "factory_rebuild", "failed", "Space restart API call failed", {"reason": reason, **exception_payload(exc)}) | |
| return False | |
| except Exception as exc: | |
| append_event(events_path, "factory_rebuild", "failed", f"{'Factory reboot' if factory_reboot else 'Space restart'} API call failed", {"reason": reason, **exception_payload(exc)}) | |
| return False | |
| def safe_same_code_reupload(api, workspace: Path, target_space_id: str, token: str, run_dir: Path, events_path: Path, *, reason: str = "") -> bool: | |
| """Re-upload the current workspace once, with restart-loop guardrails.""" | |
| if not restart_guardrail_allows("same_code_upload", events_path, reason=reason): | |
| return False | |
| runtime_before = write_space_runtime(api, target_space_id, token, run_dir, events_path) | |
| stage = runtime_stage_value(runtime_before) | |
| if stage in {"building", "buildqueued", "build_queued", "starting", "pending"}: | |
| append_event(events_path, "factory_rebuild", "warning", "Skipped same-code upload: Space is already building/starting", {"runtime": runtime_before, "reason": reason}) | |
| return False | |
| upload_workspace(api, workspace, target_space_id, token, run_dir, events_path) | |
| append_event(events_path, "factory_rebuild", "success", "Same-code workspace uploaded under restart guardrails", {"reason": reason, "runtime_before": runtime_before}) | |
| return True | |
| def collect_space_logs(target_space_id: str, token: str, run_dir: Path, events_path: Path): | |
| # Import locally because main() imports HfApi inside its own scope after | |
| # installing dependencies. Recovery/log collection can run from helper | |
| # functions where that local name is otherwise unavailable. | |
| from huggingface_hub import HfApi | |
| logs_dir = run_dir / "logs" | |
| logs_dir.mkdir(parents=True, exist_ok=True) | |
| written = [] | |
| api = HfApi(token=token) | |
| def write_log(filename: str, text: str, source: str): | |
| (logs_dir / filename).write_text(text or "", encoding="utf-8") | |
| written.append({"file": filename, "source": source, "returncode": 0, "tail": (text or "")[-1000:]}) | |
| try: | |
| # Prefer the Python API. Recent hf CLIs do not expose the old Space logs CLI command, | |
| # which produced misleading "invalid choice: spaces" logs. | |
| runtime_logs = api.get_space_runtime(repo_id=target_space_id, token=token) | |
| write_log("space_runtime_state.json", json.dumps(runtime_to_dict(runtime_logs), indent=2, ensure_ascii=False), "get_space_runtime") | |
| except Exception as exc: | |
| written.append({"file": "space_runtime_state.json", "source": "get_space_runtime", "error": str(exc)[:1000]}) | |
| for filename, factory in { | |
| "space_logs_runtime.txt": lambda: api.get_space_logs(repo_id=target_space_id, token=token, build=False), | |
| "space_logs_build.txt": lambda: api.get_space_logs(repo_id=target_space_id, token=token, build=True), | |
| }.items(): | |
| try: | |
| out = factory() | |
| if isinstance(out, (list, tuple)): | |
| out = "\n".join(str(x) for x in out) | |
| write_log(filename, str(out or ""), "HfApi.get_space_logs") | |
| except Exception as exc: | |
| # Keep a useful artifact instead of calling unsupported hf CLI commands. | |
| msg = f"Could not collect {filename} via HfApi.get_space_logs: {type(exc).__name__}: {str(exc)[:1500]}" | |
| (logs_dir / filename).write_text(msg, encoding="utf-8") | |
| written.append({"file": filename, "source": "HfApi.get_space_logs", "error": msg[-1000:]}) | |
| append_event(events_path, "space_logs", "success", "Collected best-effort Space logs", {"files": written}) | |
| write_artifact_manifest(run_dir, reason="space_logs_collected") | |
| return written | |
| def space_subdomain_url(target_space_id: str) -> str: | |
| owner, name = target_space_id.split("/", 1) | |
| return f"https://{owner}-{name}.hf.space".replace("_", "-").lower() | |
| def validate_http_health(target_space_id: str, token: str, run_dir: Path, attempt: int): | |
| import requests | |
| url = space_subdomain_url(target_space_id).rstrip("/") + "/health" | |
| headers = {"Authorization": f"Bearer {token}", "Accept": "application/json,text/plain,*/*"} | |
| response = requests.get(url, headers=headers, timeout=20) | |
| payload = { | |
| "status": "success" if response.ok else "failed", | |
| "attempt": attempt, | |
| "url": url, | |
| "status_code": response.status_code, | |
| "content_type": response.headers.get("content-type"), | |
| "text": response.text[:2000], | |
| } | |
| if response.ok: | |
| try: | |
| payload["json"] = response.json() | |
| except Exception: | |
| pass | |
| write_json(run_dir / "tests" / "http_health.json", payload) | |
| return payload | |
| raise RuntimeError(f"HTTP /health returned {response.status_code}: {response.text[:500]}") | |
| def wait_until_live(api, target_space_id: str, token: str, run_dir: Path, events_path: Path, timeout_s: int = 1800): | |
| append_event(events_path, "live_wait", "started", "Waiting for existing Space to become live") | |
| deadline = time.time() + timeout_s | |
| attempt = 0 | |
| last_error = None | |
| while time.time() < deadline: | |
| attempt += 1 | |
| runtime_payload = write_space_runtime(api, target_space_id, token, run_dir, events_path, attempt) | |
| stage = str(runtime_payload.get("stage") or "").upper() | |
| if "RUNTIME_ERROR" in stage: | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| last_error = f"Space is in RUNTIME_ERROR: {runtime_payload}" | |
| append_event(events_path, "live_wait", "waiting", "Space is in runtime error; still waiting in case hardware was changed manually", {"attempt": attempt, "runtime": runtime_payload}) | |
| time.sleep(30) | |
| continue | |
| try: | |
| health = validate_http_health(target_space_id, token, run_dir, attempt) | |
| append_event(events_path, "live_wait", "success", "HTTP /health is live", {"attempt": attempt}) | |
| return {"validator": "http_health", "health": health, "runtime": runtime_payload} | |
| except Exception as http_exc: | |
| last_error = f"HTTP health failed: {http_exc}" | |
| try: | |
| client = make_gradio_client(target_space_id, token) | |
| schema = client.view_api(return_format="dict") | |
| names = api_names_from_schema(schema) | |
| write_json(run_dir / "tests" / "api_schema.json", {"schema": schema, "api_names": names}) | |
| if names: | |
| append_event(events_path, "live_wait", "success", "Gradio API schema is live", {"attempt": attempt, "api_names": names}) | |
| return {"validator": "gradio_schema", "api_names": names, "runtime": runtime_payload} | |
| except Exception as gr_exc: | |
| last_error = (last_error or "") + f"; Gradio schema failed: {gr_exc}" | |
| append_event(events_path, "live_wait", "waiting", "Space not live yet", {"attempt": attempt, "runtime": runtime_payload, "error": last_error[-1500:] if last_error else None}) | |
| time.sleep(30) | |
| collect_space_logs(target_space_id, token, run_dir, events_path) | |
| raise RuntimeError(f"Space did not become live before timeout: {last_error}") | |
| def parse_json_env(name: str, default): | |
| value = os.environ.get(name) | |
| if not value: | |
| return default | |
| try: | |
| return json.loads(value) | |
| except Exception as exc: | |
| raise ValueError(f"Invalid JSON for {name}: {exc}") | |
| def result_contains_expected_output(result, expected_output_type: str) -> tuple[bool, dict]: | |
| expected = (expected_output_type or "any").lower().strip() | |
| info = {"expected_output_type": expected, "result_type": type(result).__name__, "result_repr": repr(result)[:2000]} | |
| paths = [] | |
| type_hints = [] | |
| text_values = [] | |
| def visit(obj): | |
| type_hints.append(type(obj).__name__.lower()) | |
| if isinstance(obj, (str, Path)): | |
| value = str(obj) | |
| if value: | |
| paths.append(value) | |
| text_values.append(value) | |
| elif isinstance(obj, dict): | |
| for key, value in obj.items(): | |
| key_l = str(key).lower() | |
| if key_l in {"path", "name", "url", "file", "filepath"}: | |
| visit(value) | |
| elif key_l in {"mime_type", "mime", "type", "format"} and value: | |
| type_hints.append(str(value).lower()) | |
| elif key_l in {"image", "video", "audio", "text"}: | |
| type_hints.append(key_l) | |
| visit(value) | |
| elif isinstance(value, (dict, list, tuple)): | |
| visit(value) | |
| elif isinstance(obj, (list, tuple)): | |
| for item in obj: | |
| visit(item) | |
| else: | |
| module = getattr(type(obj), "__module__", "") | |
| if module: | |
| type_hints.append(module.lower()) | |
| visit(result) | |
| lower_paths = [str(p).lower() for p in paths] | |
| hints = " ".join(type_hints + lower_paths) | |
| info["detected_paths"] = paths[:20] | |
| info["type_hints"] = type_hints[:30] | |
| if expected == "any": | |
| return result is not None, info | |
| image_ext = (".png", ".jpg", ".jpeg", ".webp", ".gif", ".bmp", ".tiff") | |
| video_ext = (".mp4", ".mov", ".webm", ".avi", ".mkv") | |
| audio_ext = (".wav", ".mp3", ".flac", ".ogg", ".m4a") | |
| if expected == "text": | |
| if isinstance(result, str) and bool(result.strip()): | |
| return True, info | |
| return any(isinstance(x, str) and x.strip() and Path(x).suffix.lower() not in image_ext + video_ext + audio_ext for x in text_values), info | |
| if expected == "image": | |
| return any(p.endswith(image_ext) for p in lower_paths) or any(h in hints for h in ["image", "pil.", "png", "jpeg", "jpg", "webp", "gradio.data_classes.filedata"]), info | |
| if expected == "video": | |
| return any(p.endswith(video_ext) for p in lower_paths) or any(h in hints for h in ["video", "mp4", "webm", "moviepy"]), info | |
| if expected == "audio": | |
| return any(p.endswith(audio_ext) for p in lower_paths) or any(h in hints for h in ["audio", "wav", "soundfile"]), info | |
| return result is not None, info | |
| def copy_result_artifacts(result, run_dir: Path): | |
| artifacts = run_dir / "artifacts" | |
| artifacts.mkdir(parents=True, exist_ok=True) | |
| copied = [] | |
| def maybe_copy(obj): | |
| if isinstance(obj, (str, Path)): | |
| path = Path(str(obj)) | |
| if path.exists() and path.is_file(): | |
| target = artifacts / path.name | |
| try: | |
| shutil.copy2(path, target) | |
| copied.append(str(target)) | |
| except Exception: | |
| pass | |
| elif isinstance(obj, dict): | |
| for key in ["path", "name"]: | |
| if key in obj: | |
| maybe_copy(obj[key]) | |
| for value in obj.values(): | |
| if isinstance(value, (dict, list, tuple)): | |
| maybe_copy(value) | |
| elif isinstance(obj, (list, tuple)): | |
| for item in obj: | |
| maybe_copy(item) | |
| maybe_copy(result) | |
| return copied | |
| def smoke_generate(target_space_id: str, token: str, run_dir: Path, events_path: Path): | |
| api_name = normalize_api_name(os.environ.get("API_NAME") or "/generate") | |
| expected_output_type = (os.environ.get("EXPECTED_OUTPUT_TYPE") or "any").strip() | |
| test_args = parse_json_env("TEST_ARGS_JSON", ["a cinematic robot cat astronaut, detailed, studio lighting"]) | |
| test_kwargs = parse_json_env("TEST_KWARGS_JSON", {}) | |
| if not isinstance(test_args, list): | |
| raise ValueError("TEST_ARGS_JSON must be a JSON list") | |
| if not isinstance(test_kwargs, dict): | |
| raise ValueError("TEST_KWARGS_JSON must be a JSON object") | |
| append_event(events_path, "generation_smoke", "started", "Calling live generation endpoint", {"api_name": api_name, "expected_output_type": expected_output_type}) | |
| client = make_gradio_client(target_space_id, token) | |
| schema = client.view_api(return_format="dict") | |
| discovered = api_names_from_schema(schema) | |
| safe_kwargs, dropped_kwargs, endpoint_parameters = sanitize_kwargs_for_schema(api_name, schema, test_kwargs) | |
| write_json(run_dir / "tests" / "api_schema.json", {"schema": schema, "api_names": discovered, "selected_api_name": api_name, "endpoint_parameters": endpoint_parameters}) | |
| if dropped_kwargs: | |
| append_event( | |
| events_path, | |
| "generation_smoke", | |
| "warning", | |
| "Ignored unsupported validation kwargs for this Gradio endpoint", | |
| {"api_name": api_name, "ignored_kwargs": sorted(dropped_kwargs), "endpoint_parameters": endpoint_parameters}, | |
| ) | |
| started = time.time() | |
| try: | |
| result = client.predict(*test_args, api_name=api_name, **safe_kwargs) | |
| except Exception as exc: | |
| message = str(exc) | |
| if safe_kwargs and "not a valid key-word argument" in message: | |
| append_event( | |
| events_path, | |
| "generation_smoke", | |
| "warning", | |
| "Retrying validation without keyword arguments after Gradio rejected kwargs", | |
| {"api_name": api_name, "error": message[:1500], "dropped_kwargs": sorted(safe_kwargs)}, | |
| ) | |
| safe_kwargs = {} | |
| result = client.predict(*test_args, api_name=api_name) | |
| else: | |
| raise | |
| latency = time.time() - started | |
| ok, info = result_contains_expected_output(result, expected_output_type) | |
| copied = copy_result_artifacts(result, run_dir) | |
| payload = { | |
| "status": "success" if ok else "failed", | |
| "target_space": target_space_id, | |
| "api_name": api_name, | |
| "discovered_api_names": discovered, | |
| "test_args": test_args, | |
| "test_kwargs": safe_kwargs, | |
| "original_test_kwargs": test_kwargs, | |
| "ignored_test_kwargs": dropped_kwargs, | |
| "endpoint_parameters": endpoint_parameters, | |
| "expected_output_type": expected_output_type, | |
| "latency_seconds": round(latency, 3), | |
| "observed_latency_seconds": round(latency, 3), | |
| "result_info": info, | |
| "copied_artifacts": copied, | |
| "recommended_zero_gpu_duration_seconds": int(max(30, min(300, latency * 2 + 15))), | |
| "recommendation_source": "live_gradio_predict", | |
| "recommendation_confidence": "measured", | |
| "measurement_note": "Measured from a live gradio_client.predict call, including Gradio/API/network/result serialization overhead.", | |
| "validated_at": now(), | |
| } | |
| write_json(run_dir / "tests" / "generation_smoke.json", payload) | |
| write_json(run_dir / "tests" / "test_result.json", payload) | |
| if ok: | |
| append_event(events_path, "generation_smoke", "success", "Live generation smoke test passed", {"latency_seconds": payload["latency_seconds"], "recommended_zero_gpu_duration_seconds": payload["recommended_zero_gpu_duration_seconds"], "recommended_zerogpu_duration_seconds": payload["recommended_zero_gpu_duration_seconds"], "copied_artifacts": copied[:5]}) | |
| return payload | |
| append_event(events_path, "generation_smoke", "failed", "Live generation returned an unexpected output type", payload) | |
| raise RuntimeError("Generation smoke test failed: unexpected output type") | |
| def main(): | |
| run_id = os.environ["RUN_ID"] | |
| username = os.environ.get("HF_USERNAME", "unknown") | |
| output_root = Path(os.environ.get("OUTPUT_ROOT", "/output")) | |
| target_space_id = os.environ["TARGET_SPACE_ID"].strip() | |
| token = os.environ.get("HF_TOKEN") | |
| run_dir = output_root / "runs" / run_id | |
| events_path = run_dir / "events.jsonl" | |
| state_path = run_dir / "state.json" | |
| append_event(events_path, "bootstrap", "started", "Existing Space validation worker started", {"target_space_id": target_space_id}) | |
| write_json(state_path, {"run_id": run_id, "kind": "validate_existing_space", "status": "running", "target_space": target_space_id, "created_by": username, "updated_at": now()}) | |
| if not token: | |
| raise RuntimeError("HF_TOKEN is missing") | |
| if not TARGET_RE.match(target_space_id): | |
| raise ValueError("TARGET_SPACE_ID must look like owner/space-name") | |
| try: | |
| install_deps(events_path) | |
| from huggingface_hub import HfApi | |
| api = HfApi(token=token) | |
| whoami = api.whoami(token=token) | |
| append_event(events_path, "auth", "success", "Authenticated inside validation Job", {"whoami_name": whoami.get("name")}) | |
| live = wait_until_live(api, target_space_id, token, run_dir, events_path, timeout_s=int(os.environ.get("LIVE_TIMEOUT_SECONDS", "1800"))) | |
| smoke = smoke_generate(target_space_id, token, run_dir, events_path) | |
| final_state = { | |
| "run_id": run_id, | |
| "kind": "validate_existing_space", | |
| "status": "full_inference_success", | |
| "message": "Existing Space passed live health/schema validation and generation smoke test.", | |
| "target_space": target_space_id, | |
| "target_space_url": f"https://huggingface.co/spaces/{target_space_id}", | |
| "live_validation": live, | |
| "generation_smoke": smoke, | |
| "updated_at": now(), | |
| } | |
| write_json(state_path, final_state) | |
| report = f"""# Agentic Space Factory — Existing Space Validation Report | |
| Status: **full_inference_success** | |
| Target Space: [`{target_space_id}`](https://huggingface.co/spaces/{target_space_id}) | |
| ## Generation smoke test | |
| ```json | |
| {json.dumps(smoke, indent=2, ensure_ascii=False)} | |
| ``` | |
| ## Notes | |
| - This validation is intended for Spaces whose hardware was set manually after generation. | |
| - Latency is measured from the live Gradio endpoint call. | |
| - The recommended ZeroGPU duration is a rough estimate from this live run, not a guarantee. | |
| """ | |
| (run_dir / "report.md").write_text(report, encoding="utf-8") | |
| append_event(events_path, "report_write", "success", "Wrote report.md") | |
| try: | |
| publish_eval_record(run_dir, phase="final", events_path=events_path) | |
| append_event(events_path, "anonymous_eval", "success", "Published anonymized validation evaluation record", {"enabled": True}) | |
| except Exception as eval_exc: | |
| append_event(events_path, "anonymous_eval", "warning", "Could not publish anonymized validation evaluation record", {"error": str(eval_exc)[:1000]}) | |
| append_event(events_path, "done", "full_inference_success", "Existing Space validation completed", {"latency_seconds": smoke.get("latency_seconds")}) | |
| except Exception as exc: | |
| collect_space_logs(target_space_id, token or "", run_dir, events_path) | |
| details = {"error": str(exc)[:4000]} | |
| write_json(state_path, {"run_id": run_id, "kind": "validate_existing_space", "status": "failed", "target_space": target_space_id, "details": details, "updated_at": now()}) | |
| (run_dir / "report.md").write_text(f"# Existing Space Validation Failed\n\n```json\n{json.dumps(details, indent=2, ensure_ascii=False)}\n```\n", encoding="utf-8") | |
| append_event(events_path, "failure", "failed", "Existing Space validation failed", details) | |
| try: | |
| publish_eval_record(run_dir, phase="failure", events_path=events_path) | |
| append_event(events_path, "anonymous_eval", "success", "Published anonymized failed validation evaluation record", {"enabled": True}) | |
| except Exception as eval_exc: | |
| append_event(events_path, "anonymous_eval", "warning", "Could not publish anonymized failed validation evaluation record", {"error": str(eval_exc)[:1000]}) | |
| raise SystemExit(1) | |
| if __name__ == "__main__": | |
| main() | |
| ''' | |
| def universal_model_card_worker_script() -> str: | |
| """Return the universal model-card builder worker script source.""" | |
| return UNIVERSAL_MODEL_CARD_WORKER_SCRIPT | |
| def validate_existing_space_worker_script() -> str: | |
| """Return the existing-Space validation worker script source.""" | |
| return VALIDATE_EXISTING_SPACE_WORKER_SCRIPT | |
| def encoded_universal_model_card_worker_script() -> str: | |
| """Return the base64-encoded universal model-card builder worker script. | |
| Kept for backwards compatibility in tests/imports. New Jobs read the worker | |
| from WORKER_SCRIPT_PATH to avoid oversized env/argv payloads. | |
| """ | |
| return _encode(UNIVERSAL_MODEL_CARD_WORKER_SCRIPT) | |
| def encoded_validate_existing_space_worker_script() -> str: | |
| """Return the base64-encoded existing-Space validation worker script. | |
| Kept for backwards compatibility in tests/imports. New Jobs read the worker | |
| from WORKER_SCRIPT_PATH to avoid oversized env/argv payloads. | |
| """ | |
| return _encode(VALIDATE_EXISTING_SPACE_WORKER_SCRIPT) | |
| def python_decode_and_run_command() -> list[str]: | |
| """Small command list for `run_job`. | |
| The worker script is stored in the run bucket and exposed through the Job | |
| volume as WORKER_SCRIPT_PATH. This avoids putting the large worker source in | |
| argv/env, which can make HF Jobs fail before Python starts with | |
| `argument list too long`. WORKER_SCRIPT_B64 remains as a tiny compatibility | |
| fallback for older tests/manual launch paths. | |
| """ | |
| runner = textwrap.dedent( | |
| """ | |
| import base64, os, pathlib, subprocess, sys | |
| script_path = os.environ.get('WORKER_SCRIPT_PATH') | |
| if script_path: | |
| path = pathlib.Path(script_path) | |
| if not path.exists(): | |
| raise SystemExit(f'Worker script not found at WORKER_SCRIPT_PATH={script_path}') | |
| else: | |
| script_b64 = os.environ.get('WORKER_SCRIPT_B64') | |
| if not script_b64: | |
| raise SystemExit('Missing WORKER_SCRIPT_PATH') | |
| path = pathlib.Path('/tmp/space_factory_worker.py') | |
| path.write_text(base64.b64decode(script_b64).decode('utf-8'), encoding='utf-8') | |
| raise SystemExit(subprocess.call([sys.executable, str(path)])) | |
| """ | |
| ).strip() | |
| return ["python", "-c", runner] | |