← Files PixVerseARCHIVED FILE
pvx/project_context.py
37.1 KB · Oct 2, 2026 · 00:28 UTC
from __future__ import annotations
import json
import shlex
from datetime import datetime
from pathlib import Path
from typing import Any
from .state import list_projects, project_dir, project_handoff, pvx_command, read_jsonl, slugify
MEDIA_SUFFIXES = {
".aac",
".flac",
".gif",
".jpeg",
".jpg",
".m4a",
".mov",
".mp3",
".mp4",
".png",
".wav",
".webm",
".webp",
}
def project_snapshot(slug: str, *, root: Path | None = None, surface: str = "local") -> dict[str, Any]:
path = project_dir(slug, root)
normalized_slug = path.name
if not path.exists():
return {
"slug": normalized_slug,
"title": normalized_slug,
"path": str(path),
"exists": False,
"stage": "missing",
"next_action": {
"label": "Choose or initialize a project",
"command": f"{pvx_command()} project portfolio --format markdown",
"reason": "The requested project workspace does not exist.",
},
}
title = _project_title(path)
updated_at_epoch = _project_updated_at(path)
notebook = read_jsonl(path / "notebook.jsonl")
preferences = read_jsonl(path / "preferences.jsonl")
manifest = read_jsonl(path / "manifest.jsonl")
if surface == "canvas":
return _canvas_project_snapshot(path, root=root, title=title, updated_at_epoch=updated_at_epoch,
notebook=notebook, preferences=preferences)
billing_runs = [row for row in manifest if row.get("event") == "queue.billing"]
latest_run = billing_runs[-1] if billing_runs else {}
totals = project_run_totals(billing_runs)
ledger = totals["assets"]
latest_ledger = latest_run.get("asset_ledger") if isinstance(latest_run.get("asset_ledger"), list) else []
latest_ledger = [item for item in latest_ledger if isinstance(item, dict)]
invoice = latest_run.get("invoice") if isinstance(latest_run.get("invoice"), dict) else {}
timing = latest_run.get("timing") if isinstance(latest_run.get("timing"), dict) else {}
prepared_queues = _queue_plans(path)
queue_specs = [str(item["path"]) for item in prepared_queues]
deliverables = _media_files(path / "deliverables")
local_assets = _media_files(path / "assets")
qa = _latest_qa(path)
latest_run_epoch = _timestamp_epoch(latest_run.get("at"))
deliverable_fresh = not latest_run_epoch or _latest_file_epoch(deliverables) >= latest_run_epoch
if qa:
qa["fresh_for_latest_run"] = not latest_run_epoch or qa["checked_at_epoch"] >= latest_run_epoch
current_qa = qa if qa and qa.get("fresh_for_latest_run") else None
unresolved = _unresolved_tasks(path, manifest)
failed_assets = [item for item in latest_ledger if str(item.get("status") or "") not in {"", "success"}]
successful_latest_assets = [item for item in latest_ledger if item.get("status") == "success"]
successful_assets = [item for item in ledger if item.get("status") == "success"]
stage = _project_stage(
unresolved=unresolved,
failed_assets=failed_assets,
successful_assets=successful_latest_assets,
qa=current_qa,
deliverables=deliverables if deliverable_fresh else [],
queue_specs=queue_specs,
notebook=notebook,
preferences=preferences,
)
highlights = _memory_highlights(notebook, preferences)
significance = _significance_score(
title=title,
slug=normalized_slug,
highlights=highlights,
ledger=ledger,
deliverables=deliverables,
local_assets=local_assets,
qa=current_qa,
queue_specs=queue_specs,
)
result = {
"slug": normalized_slug,
"title": title,
"path": str(path.resolve()),
"exists": True,
"updated_at_epoch": updated_at_epoch,
"stage": stage,
"signals": {
"notebook_entries": len(notebook),
"preferences": len(preferences),
"queue_specs": len(queue_specs),
"billing_runs": len(billing_runs),
"generated_assets": len(ledger),
"successful_assets": len(successful_assets),
"failed_assets": len(failed_assets),
"local_assets": len(local_assets),
"deliverables": len(deliverables),
"latest_run_has_deliverable": bool(deliverables) and deliverable_fresh,
"unresolved_tasks": len(unresolved),
"significance_score": significance,
},
"memory": highlights,
"latest_run": {
"at": latest_run.get("at") or "",
"wall_seconds": timing.get("wall_seconds"),
"generated_assets": invoice.get("generated_assets", len(ledger)),
"total_actual_credits": invoice.get("total_actual_credits"),
"observed_credit_delta": invoice.get("observed_credit_delta"),
}
if latest_run
else None,
"project_totals": {
"billing_runs": len(billing_runs),
"generated_assets": totals["generated_assets"],
"successful_assets": totals["successful_assets"],
"total_attributable_credits": totals["total_attributable_credits"],
"assets_with_unknown_credits": totals["assets_with_unknown_credits"],
},
"assets": [_compact_asset(item) for item in ledger[-12:]],
"deliverables": deliverables[-8:],
"local_assets": local_assets[-8:],
"qa": qa,
"unresolved_tasks": unresolved,
"queue_specs": queue_specs,
"prepared_queues": prepared_queues[-4:],
}
result["next_action"] = _next_action(result)
return result
def _canvas_project_snapshot(path: Path, *, root: Path | None, title: str, updated_at_epoch: int,
notebook: list[dict[str, Any]], preferences: list[dict[str, Any]]) -> dict[str, Any]:
canvas = project_handoff(path.name, root=root, surface="canvas")["canvas"]
bound = bool(canvas.get("bound"))
stage_record = canvas.get("stage_status") or {}
latest = canvas.get("latest_run") or {}
nodes = canvas.get("cloud_nodes") or []
generation = latest.get("generation_status") or "unknown"
slug = shlex.quote(path.name)
command = f"{pvx_command()} project handoff {slug} --surface canvas --format markdown"
stage = "canvas_review"
label, reason = "Review the Canvas and continue cloud nodes", "Cloud outputs do not require local media, QA recovery, or credit reporting."
if not bound:
stage, label, reason = "canvas_binding_missing", "Resolve the intended Canvas project", "No valid Canvas binding exists for this selected target; do not create a replacement automatically."
command = ""
elif latest.get("run_id") and (latest.get("generation_state") != "started" or generation == "unknown"):
stage, label, reason = "canvas_needs_recovery", "Recover the existing Canvas run", "Submission ownership is unresolved; inspect the existing run without resubmitting."
command = f"{pvx_command()} canvas paid reconcile --project {slug} --run-id {shlex.quote(str(latest['run_id']))}"
elif generation in {"running", "blocked"}:
stage, label, reason = "canvas_generating", "Follow the existing Canvas run", "Continue cloud status following without downloads or credits."
command = f"{pvx_command()} canvas paid follow --project {slug} --run-id {shlex.quote(str(latest['run_id']))}"
elif generation == "failed":
stage, label, reason = "canvas_needs_attention", "Review failed Canvas nodes", "Keep existing cloud results and obtain a fresh paid plan before any new generation."
elif (stage_record.get("final_deliverable_status") == "complete"
and _timestamp_epoch(stage_record.get("recorded_at")) >= _timestamp_epoch(latest.get("submitted_at"))):
local_paths = stage_record.get("deliverable_paths") or []
if stage_record.get("delivery_mode") == "local" and (not local_paths or not all(Path(p).is_file() for p in local_paths)):
stage, label, reason = "canvas_local_delivery_incomplete", "Complete requested local delivery", "The user requested local files, but recorded exports are unavailable; do not regenerate."
else:
stage, label, reason = "canvas_delivered", "Continue from the Canvas deliverable", "The recorded stage is complete; cloud delivery does not require a local copy."
return {
"slug": path.name, "title": title, "path": str(path.resolve()), "exists": True,
"surface": "canvas", "stage": stage, "updated_at_epoch": updated_at_epoch,
"canvas": canvas, "memory": _memory_highlights(notebook, preferences),
"assets": nodes, "latest_run": latest or None,
"signals": {"generated_assets": len(nodes), "significance_score": 10 if bound else 0},
"next_action": {"label": label, "reason": reason, "command": command},
}
def project_portfolio(*, root: Path | None = None, limit: int = 20, surface: str = "local") -> dict[str, Any]:
rows = list_projects(root=root, limit=limit)
snapshots = [project_snapshot(str(row["slug"]), root=root, surface=surface) for row in rows]
if surface == "canvas":
snapshots = [item for item in snapshots if item.get("canvas", {}).get("bound")]
for index, snapshot in enumerate(snapshots):
# "Continue" should usually mean recent meaningful work, not the
# historically largest project. Richness filters out empty/demo noise;
# recency then dominates among credible candidates.
recency_score = max(0, 100 - index * 5)
snapshot["resume_score"] = int(snapshot.get("signals", {}).get("significance_score", 0)) + recency_score
meaningful = [snapshot for snapshot in snapshots if snapshot.get("stage") not in {"idle", "missing"}]
suggested = max(meaningful, key=lambda item: int(item.get("resume_score", 0)), default=None)
duplicate_titles = _duplicate_title_groups(snapshots)
return {
"surface": surface,
"projects": snapshots,
"suggested_resume": _compact_project(suggested) if suggested else None,
"possible_duplicate_titles": duplicate_titles,
"selection_rule": "recent activity plus real assets, deliverables, QA, queues, and useful memory",
}
def project_run_totals(runs: list[dict[str, Any]]) -> dict[str, Any]:
"""Summarize unique generated assets without double-counting idempotent reruns."""
unique: dict[str, dict[str, Any]] = {}
anonymous_index = 0
for run in runs:
run_ledger = run.get("asset_ledger") if isinstance(run.get("asset_ledger"), list) else []
for item in run_ledger:
if not isinstance(item, dict):
continue
stable = item.get("task_id") or item.get("url") or item.get("path") or item.get("cover_url")
if stable:
key = str(stable)
else:
anonymous_index += 1
key = f"anonymous:{anonymous_index}:{item.get('id') or ''}"
unique[key] = item
assets = list(unique.values())
known_credits = [
item.get("actual_credits")
for item in assets
if isinstance(item.get("actual_credits"), (int, float))
and not isinstance(item.get("actual_credits"), bool)
]
invoice_fallback_credits: list[float] = []
seen_run_signatures: set[tuple[str, ...]] = set()
for run_index, run in enumerate(runs):
run_ledger = run.get("asset_ledger") if isinstance(run.get("asset_ledger"), list) else []
stable_keys = tuple(
sorted(
str(item.get("task_id") or item.get("url") or item.get("path") or item.get("cover_url"))
for item in run_ledger
if isinstance(item, dict)
and (item.get("task_id") or item.get("url") or item.get("path") or item.get("cover_url"))
)
)
signature = stable_keys or (f"anonymous-run:{run_index}",)
if signature in seen_run_signatures:
continue
seen_run_signatures.add(signature)
unknown_items = [
item
for item in run_ledger
if isinstance(item, dict)
and not (
isinstance(item.get("actual_credits"), (int, float))
and not isinstance(item.get("actual_credits"), bool)
)
]
invoice = run.get("invoice") if isinstance(run.get("invoice"), dict) else {}
invoice_total = invoice.get("total_actual_credits")
if unknown_items and isinstance(invoice_total, (int, float)) and not isinstance(invoice_total, bool):
itemized = sum(
item.get("actual_credits")
for item in run_ledger
if isinstance(item, dict)
and isinstance(item.get("actual_credits"), (int, float))
and not isinstance(item.get("actual_credits"), bool)
)
invoice_fallback_credits.append(max(0, invoice_total - itemized))
total_credits = sum(known_credits) + sum(invoice_fallback_credits)
return {
"assets": assets,
"generated_assets": len(assets),
"successful_assets": sum(item.get("status") == "success" for item in assets),
"total_attributable_credits": total_credits if known_credits or invoice_fallback_credits else None,
"assets_with_unknown_credits": len(assets) - len(known_credits),
}
def resolve_project_resume(
selector: str = "",
*,
root: Path | None = None,
limit: int = 20,
surface: str = "local",
) -> dict[str, Any]:
portfolio = project_portfolio(root=root, limit=limit, surface=surface)
projects = portfolio["projects"]
query = selector.strip()
matched_by = "portfolio_recommendation"
candidates: list[dict[str, Any]] = []
if query:
normalized = slugify(query)
exact = [item for item in projects if item.get("slug") == query or item.get("slug") == normalized]
if exact:
candidates = exact
matched_by = "exact_slug"
else:
needle = query.casefold()
candidates = [
item
for item in projects
if needle in str(item.get("slug") or "").casefold()
or needle in str(item.get("title") or "").casefold()
]
matched_by = "title_or_slug_search"
elif portfolio.get("suggested_resume"):
suggested_slug = portfolio["suggested_resume"]["slug"]
candidates = [item for item in projects if item.get("slug") == suggested_slug]
if not candidates:
return {
"ok": False,
"query": query,
"error": "project_not_found" if query else "project_portfolio_empty",
"message": "No matching meaningful project workspace was found.",
"recent_projects": [_compact_project(item) for item in projects[:5]],
}
selected = max(candidates, key=lambda item: int(item.get("resume_score", 0)))
alternatives = [
_compact_project(item)
for item in sorted(candidates, key=lambda item: int(item.get("resume_score", 0)), reverse=True)
if item.get("slug") != selected.get("slug")
][:4]
return {
"ok": True,
"query": query,
"matched_by": matched_by,
"selection_reason": _selection_reason(selected),
"project": selected,
"alternatives": alternatives,
"possible_duplicate_titles": portfolio.get("possible_duplicate_titles", []),
}
def render_project_resume_markdown(payload: dict[str, Any]) -> str:
if not payload.get("ok"):
lines = ["# PixVerse Project Resume", "", str(payload.get("message") or "Project not found.")]
recent = payload.get("recent_projects") if isinstance(payload.get("recent_projects"), list) else []
if recent:
lines.extend(["", "## Recent projects", ""])
lines.extend(f"- `{item['slug']}` — {item['title']} ({item['stage']})" for item in recent)
return "\n".join(lines)
project = payload["project"]
if project.get("surface") == "canvas":
canvas = project.get("canvas") or {}
stage = canvas.get("stage_status") or {}
action = project.get("next_action") or {}
lines = ["# PixVerse Canvas Resume", "", f"- Project: **{project['title']}** (`{project['slug']}`)",
f"- Canvas: [Open project]({canvas.get('editor_url', '')})",
f"- Stage: `{project['stage']}`; position: `{stage.get('position') or 'not recorded'}`",
f"- Final deliverable: `{stage.get('final_deliverable_status', 'unknown')}`",
f"- Remaining: {', '.join(stage.get('remaining_stages') or []) or 'none recorded'}",
"- Technical QA: `not_checked`; downloads and credits are on demand.", "", "## Cloud nodes", ""]
lines.extend(f"- `{node['node_id']}` — {node.get('title') or node.get('node_type')} (`{node.get('status')}`)"
for node in canvas.get("cloud_nodes", []))
memory = project.get("memory") or {}
if memory.get("entries"):
lines.extend(["", "## Remembered context", ""])
lines.extend(f"- {entry.get('kind', 'note')}: {entry.get('text')}" for entry in memory["entries"])
lines.extend(["", "## Recommended next move", "", f"{action.get('label')}: {action.get('reason')}"])
if action.get("command"):
lines.append(f"`{action['command']}`")
return "\n".join(lines)
signals = project.get("signals") if isinstance(project.get("signals"), dict) else {}
latest_run = project.get("latest_run") if isinstance(project.get("latest_run"), dict) else {}
totals = project.get("project_totals") if isinstance(project.get("project_totals"), dict) else {}
unknown_credit_assets = int(totals.get("assets_with_unknown_credits", 0) or 0)
credit_note = f"; {unknown_credit_assets} asset(s) not itemized" if unknown_credit_assets else ""
project_credits = totals.get("total_attributable_credits")
project_credit_text = project_credits if project_credits is not None else "unknown"
next_action = project.get("next_action") if isinstance(project.get("next_action"), dict) else {}
lines = [
"# PixVerse Project Resume",
"",
f"- Project: **{project.get('title')}** (`{project.get('slug')}`)",
f"- Stage: `{project.get('stage')}`",
f"- Selected because: {payload.get('selection_reason')}",
f"- Project assets: `{signals.get('generated_assets', 0)}` across `{totals.get('billing_runs', 0)}` run(s); deliverables: `{signals.get('deliverables', 0)}`; unresolved: `{signals.get('unresolved_tasks', 0)}`",
f"- Project attributable credits: `{project_credit_text}`{credit_note}; latest run: `{latest_run.get('total_actual_credits', 'unknown')}`",
]
qa = project.get("qa") if isinstance(project.get("qa"), dict) else None
if qa:
qa_state = (
"stale after latest generation"
if qa.get("fresh_for_latest_run") is False
else "passed"
if qa.get("ok")
else "needs attention"
)
lines.append(f"- Latest QA: `{qa_state}` ({qa.get('issues_count', 0)} issue(s))")
memory = project.get("memory") if isinstance(project.get("memory"), dict) else {}
entries = memory.get("entries") if isinstance(memory.get("entries"), list) else []
if entries:
lines.extend(["", "## Remembered context", ""])
lines.extend(f"- **{item.get('kind', 'note')}** — {item.get('text')}" for item in entries)
assets = project.get("assets") if isinstance(project.get("assets"), list) else []
if assets:
lines.extend(["", "## Recent generated assets", "", "| Role | Model | Status | Task | Credits |", "|---|---|---|---|---:|"])
lines.extend(
f"| {_md(item.get('role'))} | {_md(item.get('model'))} | {_md(item.get('status'))} | {_md(item.get('task_id'))} | {_md(item.get('actual_credits'))} |"
for item in assets
)
deliverables = project.get("deliverables") if isinstance(project.get("deliverables"), list) else []
if deliverables:
lines.extend(["", "## Deliverables", ""])
lines.extend(f"- `{item}`" for item in deliverables)
prepared_queues = project.get("prepared_queues") if isinstance(project.get("prepared_queues"), list) else []
if project.get("stage") == "ready_to_generate" and prepared_queues:
prepared = prepared_queues[-1]
route = prepared.get("route") if isinstance(prepared.get("route"), dict) else {}
models = ", ".join(str(model) for model in prepared.get("models", []) if model) or "unspecified"
lines.extend(
[
"",
"## Prepared generation",
"",
f"- Tasks: `{prepared.get('task_count', 0)}`; models: `{models}`",
f"- Route: `{route.get('kind') or 'custom'} / {route.get('control_layer') or 'custom'}`; intent: `{route.get('intent') or 'unspecified'}`",
]
)
if route.get("reason"):
lines.append(f"- Why this route: {route['reason']}")
lines.append(f"- Queue: `{prepared.get('path')}`")
lines.extend(["", "## Recommended next move", "", f"**{next_action.get('label', 'Continue project')}** — {next_action.get('reason', '')}"])
if next_action.get("command"):
lines.extend(["", f"`{next_action['command']}`"])
return "\n".join(lines)
def render_project_portfolio_markdown(payload: dict[str, Any]) -> str:
projects = payload.get("projects") if isinstance(payload.get("projects"), list) else []
if payload.get("surface") == "canvas":
lines = ["# PixVerse Canvas Portfolio", "", "| Project | Stage | Cloud nodes | Next move |", "|---|---|---:|---|"]
for item in projects:
canvas = item.get("canvas") or {}
action = item.get("next_action") or {}
lines.append(
f"| [{_md(item.get('title'))}]({canvas.get('editor_url', '')}) | {_md(item.get('stage'))} | "
f"{canvas.get('cloud_node_count', 0)} | {_md(action.get('label'))} |"
)
lines.extend(["", "Cloud preview is the default. Downloads, local QA, and credits are on demand."])
return "\n".join(lines)
lines = [
"# PixVerse Project Portfolio",
"",
"| Project | Stage | Assets | Deliverables | QA | Credits | Next move |",
"|---|---|---:|---:|---|---:|---|",
]
for item in projects:
signals = item.get("signals") if isinstance(item.get("signals"), dict) else {}
totals = item.get("project_totals") if isinstance(item.get("project_totals"), dict) else {}
qa = item.get("qa") if isinstance(item.get("qa"), dict) else None
next_action = item.get("next_action") if isinstance(item.get("next_action"), dict) else {}
qa_text = (
"stale"
if qa and qa.get("fresh_for_latest_run") is False
else "passed"
if qa and qa.get("ok")
else "issues"
if qa
else "not checked"
)
lines.append(
f"| **{_md(item.get('title'))}** (`{_md(item.get('slug'))}`) | {_md(item.get('stage'))} | "
f"{signals.get('generated_assets', 0)} | {signals.get('deliverables', 0)} | {qa_text} | "
f"{_md(totals.get('total_attributable_credits', ''))} | {_md(next_action.get('label'))} |"
)
suggested = payload.get("suggested_resume") if isinstance(payload.get("suggested_resume"), dict) else None
if suggested:
lines.extend(
[
"",
f"Suggested resume: **{suggested.get('title')}** (`{suggested.get('slug')}`) — stage `{suggested.get('stage')}`.",
]
)
duplicates = payload.get("possible_duplicate_titles") if isinstance(payload.get("possible_duplicate_titles"), list) else []
if duplicates:
lines.extend(["", "Possible duplicate project titles:", ""])
lines.extend(f"- {item['title']}: {', '.join(f'`{slug}`' for slug in item['slugs'])}" for item in duplicates)
return "\n".join(lines)
def _project_title(path: Path) -> str:
project_md = path / "project.md"
if project_md.exists():
first_line = project_md.read_text(encoding="utf-8", errors="ignore").splitlines()[:1]
if first_line and first_line[0].startswith("# "):
return first_line[0][2:].strip() or path.name
return path.name
def _project_updated_at(path: Path) -> int:
return int(
max(
(child.stat().st_mtime for child in path.rglob("*") if child.is_file()),
default=path.stat().st_mtime,
)
)
def _queue_plans(path: Path) -> list[dict[str, Any]]:
plans: list[dict[str, Any]] = []
# Queue files often start at the project root, but users and agents also
# organize them under `briefs/`, `queues/`, or another shallow planning
# folder. Resume should discover that intent instead of treating a prepared
# project as empty. Generated media and QA trees are excluded so provider
# payloads or reports cannot masquerade as executable plans.
ignored_trees = {"assets", "deliverables", "quality"}
candidates = [
candidate
for candidate in path.rglob("*.json")
if not ignored_trees.intersection(candidate.relative_to(path).parts[:-1])
]
for candidate in sorted(candidates):
try:
payload = json.loads(candidate.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
continue
if not isinstance(payload, dict) or not isinstance(payload.get("tasks"), list):
continue
tasks = [item for item in payload["tasks"] if isinstance(item, dict)]
models: list[str] = []
for task in tasks:
model = _command_option(str(task.get("cmd") or task.get("command") or ""), "--model", "-m")
if model and model not in models:
models.append(model)
route = payload.get("route") if isinstance(payload.get("route"), dict) else {}
reasons = route.get("reasons") if isinstance(route.get("reasons"), list) else []
plans.append(
{
"path": str(candidate.resolve()),
"updated_at_epoch": candidate.stat().st_mtime,
"task_count": len(tasks),
"models": models,
"route": {
"kind": route.get("kind") or "",
"intent": route.get("intent") or "",
"control_layer": route.get("control_layer") or "",
"model": route.get("model") or "",
"quality": route.get("quality") or "",
"reason": _compact_text(str(reasons[0])) if reasons else "",
},
}
)
plans.sort(key=lambda item: float(item["updated_at_epoch"]))
return plans
def _command_option(command: str, *names: str) -> str:
try:
argv = shlex.split(command)
except ValueError:
return ""
for index, value in enumerate(argv):
if value in names and index + 1 < len(argv):
return argv[index + 1]
for name in names:
if value.startswith(f"{name}="):
return value.split("=", 1)[1]
return ""
def _media_files(path: Path) -> list[str]:
if not path.exists():
return []
files = [
candidate
for candidate in path.rglob("*")
if candidate.is_file() and candidate.suffix.lower() in MEDIA_SUFFIXES
]
files.sort(key=lambda candidate: candidate.stat().st_mtime)
return [str(candidate.resolve()) for candidate in files]
def _latest_file_epoch(paths: list[str]) -> float:
mtimes: list[float] = []
for value in paths:
try:
mtimes.append(Path(value).stat().st_mtime)
except OSError:
continue
return max(mtimes, default=0.0)
def _latest_qa(path: Path) -> dict[str, Any] | None:
quality = path / "quality"
if not quality.exists():
return None
candidates = [candidate for candidate in quality.glob("*.json") if candidate.is_file()]
candidates.sort(key=lambda candidate: candidate.stat().st_mtime, reverse=True)
for candidate in candidates:
try:
payload = json.loads(candidate.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
continue
if not isinstance(payload, dict):
continue
summary = payload.get("summary") if isinstance(payload.get("summary"), dict) else {}
issues = payload.get("issues") if isinstance(payload.get("issues"), list) else []
issues_count = summary.get("assets_with_issues", len(issues))
ok = payload.get("ok")
if not isinstance(ok, bool):
ok = bool(payload.get("exists", True)) and not issues
return {
"ok": ok,
"issues_count": int(issues_count or 0),
"checked_at": payload.get("checked_at") or "",
"checked_at_epoch": _timestamp_epoch(payload.get("checked_at")) or candidate.stat().st_mtime,
"report_path": str(candidate.resolve()),
}
return None
def _unresolved_tasks(path: Path, manifest: list[dict[str, Any]]) -> list[dict[str, Any]]:
pending_file = path / "pending-reconcile.jsonl"
if pending_file.exists():
rows = read_jsonl(pending_file)
if rows:
return [
{"id": row.get("id") or "", "task_id": row.get("task_id") or "", "reason": row.get("error_class") or "pending_reconcile"}
for row in rows
]
latest_by_id: dict[str, dict[str, Any]] = {}
for row in manifest:
event = str(row.get("event") or "")
task_id = str(row.get("task_id") or "")
if task_id and event.startswith("task."):
latest_by_id[task_id] = row
return [
{"id": row.get("id") or "", "task_id": task_id, "reason": row.get("error_class") or row.get("event")}
for task_id, row in latest_by_id.items()
if row.get("event") in {"task.submitted", "task.progress", "task.unresolved"}
]
def _project_stage(**signals: Any) -> str:
if signals["unresolved"]:
return "needs_reconcile"
if signals["failed_assets"]:
return "needs_attention"
if signals["successful_assets"]:
qa = signals["qa"]
if qa is None:
return "delivered" if signals["deliverables"] else "generated"
if not qa.get("ok"):
return "needs_attention"
return "delivered" if signals["deliverables"] else "qa_passed"
if signals["queue_specs"]:
return "ready_to_generate"
if signals["notebook"] or signals["preferences"]:
return "planning"
return "idle"
def _memory_highlights(notebook: list[dict[str, Any]], preferences: list[dict[str, Any]]) -> dict[str, Any]:
entries: list[dict[str, str]] = []
for record in notebook:
kind = str(record.get("kind") or "note")
if kind == "prompt":
continue
text = str(record.get("text") or "").strip()
if text:
entries.append({"kind": kind, "text": _compact_text(text)})
for record in preferences:
text = str(record.get("preference") or record.get("text") or "").strip()
if text:
entries.append({"kind": "preference", "text": _compact_text(text)})
return {"entries": entries[-8:]}
def _significance_score(**signals: Any) -> int:
score = 0
score += min(20, len(signals["highlights"].get("entries", [])) * 4)
score += min(16, len(signals["queue_specs"]) * 4)
real_assets = [
item
for item in signals["ledger"]
if item.get("task_id") and "example.test" not in str(item.get("url") or item.get("cover_url") or "")
]
score += min(32, len(real_assets) * 8)
score += min(16, len(signals["local_assets"]) * 4)
score += min(32, len(signals["deliverables"]) * 16)
if signals["qa"] and signals["qa"].get("ok"):
score += 12
if str(signals["title"]).casefold() != str(signals["slug"]).replace("-", " ").casefold():
score += 4
return score
def _compact_asset(item: dict[str, Any]) -> dict[str, Any]:
return {
"id": item.get("id") or "",
"role": item.get("role") or item.get("id") or "",
"kind": item.get("kind") or item.get("media_type") or "",
"model": item.get("model") or "",
"status": item.get("status") or "",
"task_id": item.get("task_id") or "",
"local_path": item.get("local_path") or "",
"url_or_path": (
item.get("local_path")
or item.get("url")
or item.get("path")
or item.get("cover_url")
or ""
),
"actual_credits": item.get("actual_credits"),
}
def _next_action(snapshot: dict[str, Any]) -> dict[str, str]:
slug = str(snapshot.get("slug") or "")
stage = snapshot.get("stage")
if stage == "needs_reconcile":
return {
"label": "Reconcile submitted tasks",
"command": f"{pvx_command()} queue reconcile {slug}",
"reason": "At least one billed or submitted task has no terminal result in project memory.",
}
if stage == "needs_attention":
return {
"label": "Inspect the failed or weak asset before spending again",
"command": f"{pvx_command()} project ledger {slug} --format markdown",
"reason": "The latest run or QA contains an issue that should shape any retry.",
}
if stage == "generated":
return {
"label": "Deliver generated assets",
"command": f"{pvx_command()} project ledger {slug} --format markdown",
"reason": "Generation succeeded; deliver the recorded local files or recover missing downloads. QA runs only when requested or required by a specialized workflow.",
}
if stage == "qa_passed":
return {
"label": "Choose or package the final deliverable",
"command": f"{pvx_command()} project ledger {slug} --format markdown",
"reason": "The generated assets passed QA but no local deliverable is recorded yet.",
}
if stage == "delivered":
return {
"label": "Continue from the final deliverable without losing prior choices",
"command": f"{pvx_command()} project search {slug} \"<requested change>\"",
"reason": "A final local deliverable is already present; search memory before revising. QA is optional.",
}
if stage == "ready_to_generate":
spec = snapshot.get("queue_specs", [""])[-1]
return {
"label": "Preflight the prepared queue",
"command": f"{pvx_command()} quote queue {spec} --format markdown",
"reason": "A queue exists but no generated run is recorded.",
}
if stage == "planning":
return {
"label": "Turn remembered decisions into the next queue",
"command": f"{pvx_command()} project summary {slug}",
"reason": "The project has useful memory but no prepared generation queue yet.",
}
return {
"label": "Add a brief or choose another project",
"command": f"{pvx_command()} project portfolio --format markdown",
"reason": "This workspace contains no meaningful creative state yet.",
}
def _duplicate_title_groups(projects: list[dict[str, Any]]) -> list[dict[str, Any]]:
groups: dict[str, list[dict[str, Any]]] = {}
for item in projects:
title_key = slugify(str(item.get("title") or ""), fallback="")
if title_key:
groups.setdefault(title_key, []).append(item)
return [
{"title": items[0].get("title") or key, "slugs": sorted(str(item.get("slug")) for item in items)}
for key, items in groups.items()
if len(items) > 1
]
def _compact_project(item: dict[str, Any] | None) -> dict[str, Any] | None:
if not item:
return None
return {
"slug": item.get("slug"),
"title": item.get("title"),
"stage": item.get("stage"),
"resume_score": item.get("resume_score", 0),
"next_action": item.get("next_action"),
}
def _selection_reason(project: dict[str, Any]) -> str:
signals = project.get("signals") if isinstance(project.get("signals"), dict) else {}
parts = [f"stage {project.get('stage')}"]
for key, label in (("generated_assets", "generated asset(s)"), ("deliverables", "deliverable(s)"), ("unresolved_tasks", "unresolved task(s)")):
value = int(signals.get(key, 0) or 0)
if value:
parts.append(f"{value} {label}")
qa = project.get("qa") if isinstance(project.get("qa"), dict) else None
if qa:
if qa.get("fresh_for_latest_run") is False:
parts.append("QA stale after latest run")
else:
parts.append("passing QA" if qa.get("ok") else "QA needing attention")
return ", ".join(parts)
def _compact_text(value: str, max_length: int = 220) -> str:
rendered = " ".join(value.split())
if len(rendered) <= max_length:
return rendered
return rendered[: max_length - 3].rstrip() + "..."
def _timestamp_epoch(value: Any) -> float:
if not isinstance(value, str) or not value.strip():
return 0.0
try:
return datetime.fromisoformat(value.strip().replace("Z", "+00:00")).timestamp()
except ValueError:
return 0.0
def _md(value: Any) -> str:
if value is None:
return ""
return str(value).replace("|", "\\|").replace("\n", " ")
SHA-256: 9f7a62628190376c412dd3b1a54231c7498e424f662d5f285afa60fc7ab0c2b1