main.py (5146 lines, 106 routes) becomes an assembly only: one package per domain — core, files, settings, dashboards, conversations, diff, notifications, forms, runs, accounts, workers, models, cron, agents, projects, services, goals, memories, plans, templates — each exposing an APIRouter; the flat domain modules move into their package behind a barrel that keeps the old `import conversations` / `import projects` spellings. The shared singletons (store, meta_store, hub, indexer, …) are built once by core.state.build_state() and attached to app.state.ai; routes take them as the `deps: State` dependency and helpers as an explicit `deps: AppState`. conversations/pricing.py carries the per-model rates out of the parser. Verified: route table and OpenAPI byte-identical; 90 read endpoints golden-diffed against the monolith on a copy of the live data (identical); write routes smoke-tested; 66 backend tests pass. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
383 lines
16 KiB
Python
383 lines
16 KiB
Python
"""
|
|
Forms — rich structured questions the agent asks the user, inline in the
|
|
conversation.
|
|
|
|
The ``ask-form`` skill POSTs ``/api/forms`` with a field spec (text inputs,
|
|
select, multi-select, radio, checkboxes, sliders, file uploads, …), then sends
|
|
a normal push notification linking to the conversation with ``?form=<id>`` so
|
|
the PWA opens the form focused in a full-screen modal. The conversation thread
|
|
also renders the form inline on the Bash tool card that created it; submitting
|
|
POSTs the answers back and the skill's ``wait`` command (which the agent runs
|
|
under a Monitor, not a blocking foreground call) exits with them.
|
|
|
|
This replaces Claude Code's interactive AskUserQuestion tool for spawned
|
|
sessions (the sidecar disallows it) — unlike a 2-3-button phone ask, a form
|
|
can carry many typed fields and both harnesses (claude and pi) can drive it,
|
|
since it's just HTTP against the backend.
|
|
|
|
Store shape (single JSON file, atomic rewrite, same pattern as ``notify.py``):
|
|
|
|
{ "forms": [ {"id", "title", "description", "type", "url", "sessionId",
|
|
"at", "fields": [...], "status": "pending"|"submitted"|
|
|
"cancelled", "answers", "submittedAt", "cancelledAt"}, … ] }
|
|
|
|
The log is capped (newest kept). Waiting for a submit is in-memory (one
|
|
``threading.Event`` per pending form) — a backend restart drops the wait; the
|
|
skill's poll loop just re-polls.
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import pathlib
|
|
import re
|
|
import threading
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
_MAX_LOG = 500
|
|
_ID_RE = re.compile(r"[^a-z0-9_-]+")
|
|
|
|
# Field types the UI knows how to render. An unknown type is rejected at
|
|
# create time — better a loud 400 for the skill than a card the user can't
|
|
# answer.
|
|
FIELD_TYPES = (
|
|
"text", "textarea", "number", "select", "multiselect", "radio",
|
|
"checkbox", "slider", "file", "date", "markdown", "html",
|
|
)
|
|
_OPTION_TYPES = ("select", "multiselect", "radio")
|
|
# The "Other…" choice an option field can carry (``other``: on by default for a
|
|
# required one — a mandatory question whose options miss the real answer must
|
|
# still be answerable): the viewer adds it to the list with a free-text input,
|
|
# and the typed text comes back as the value. ``OTHER`` is the draft sentinel
|
|
# for "Other picked, nothing typed yet" — never a valid answer.
|
|
OTHER = "__other__"
|
|
ASSET_RE = re.compile(r"\{\{\s*asset:\s*([^}\s]+)\s*\}\}")
|
|
# Display-only blocks: rendered in the form's flow, never answered — skipped by
|
|
# answer validation and recaps. ``markdown`` carries its source in ``text``;
|
|
# ``html`` carries raw markup in ``html`` (a visualisation, a schema, an image
|
|
# or video, …) that the viewer renders either sanitized inline or — when
|
|
# ``sandbox`` is set, or the markup needs a document of its own (a <script>,
|
|
# a <style> sheet, a <link>, an <iframe>) — inside a sandboxed iframe where
|
|
# scripts may run and styles stay scoped. A block's content is capped so one form can't
|
|
# bloat the single-file store.
|
|
DISPLAY_TYPES = ("markdown", "html")
|
|
_MAX_BLOCK_CHARS = 512 * 1024
|
|
_NEEDS_DOC_RE = re.compile(r"<(script|style|link|iframe)\b", re.I)
|
|
|
|
|
|
def _now_iso() -> str:
|
|
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
|
|
|
|
def _slug(s: str) -> str:
|
|
return _ID_RE.sub("-", (s or "").strip().lower()).strip("-")
|
|
|
|
|
|
def _num(v, fallback=None):
|
|
try:
|
|
f = float(v)
|
|
except (TypeError, ValueError):
|
|
return fallback
|
|
return int(f) if f.is_integer() else f
|
|
|
|
|
|
def valid_fields(raw: list) -> list[dict]:
|
|
"""Validate + normalize a field spec list; raises ValueError on bad input."""
|
|
if not isinstance(raw, list) or not raw:
|
|
raise ValueError("fields must be a non-empty list")
|
|
out: list[dict] = []
|
|
seen: set[str] = set()
|
|
for i, f in enumerate(raw):
|
|
if not isinstance(f, dict):
|
|
raise ValueError(f"field #{i}: must be an object")
|
|
ftype = str(f.get("type") or "text").strip().lower()
|
|
if ftype not in FIELD_TYPES:
|
|
raise ValueError(
|
|
f"field #{i}: unknown type {ftype!r} (one of {', '.join(FIELD_TYPES)})")
|
|
label = str(f.get("label") or "").strip()
|
|
key = _slug(str(f.get("key") or "")) or _slug(label)
|
|
if ftype in DISPLAY_TYPES:
|
|
src = "html" if ftype == "html" else "text"
|
|
# An html block also accepts ``text`` so the two blocks are spelled
|
|
# alike; ``html`` wins when both are given.
|
|
text = str(f.get(src) or f.get("text") or "").strip()
|
|
if not text:
|
|
raise ValueError(f"field #{i}: {ftype} needs a non-empty \"{src}\"")
|
|
if len(text) > _MAX_BLOCK_CHARS:
|
|
raise ValueError(
|
|
f"field #{i}: {ftype} block exceeds {_MAX_BLOCK_CHARS // 1024} KB")
|
|
key = key or f"{ftype}-{i + 1}"
|
|
if not key:
|
|
raise ValueError(f"field #{i}: needs a key or a label")
|
|
base, n = key, 2
|
|
while key in seen:
|
|
key = f"{base}-{n}"
|
|
n += 1
|
|
seen.add(key)
|
|
field: dict = {
|
|
"key": key,
|
|
"label": label or key,
|
|
"type": ftype,
|
|
"required": bool(f.get("required", False)),
|
|
}
|
|
if ftype in DISPLAY_TYPES:
|
|
# No label fallback to the key: a block's heading is optional.
|
|
field["label"] = label
|
|
field["required"] = False
|
|
if ftype == "html":
|
|
field["html"] = text
|
|
# Scripts and sheets only work in the iframe, so their presence
|
|
# implies it; recorded so the viewer needn't re-scan.
|
|
field["sandbox"] = bool(f.get("sandbox")) or bool(_NEEDS_DOC_RE.search(text))
|
|
h = _num(f.get("height"))
|
|
if h is not None and h > 0:
|
|
field["height"] = int(h)
|
|
else:
|
|
field["text"] = text
|
|
out.append(field)
|
|
continue
|
|
for opt_key in ("placeholder", "help", "accept"):
|
|
v = str(f.get(opt_key) or "").strip()
|
|
if v:
|
|
field[opt_key] = v
|
|
if ftype in _OPTION_TYPES:
|
|
opts = []
|
|
for o in (f.get("options") or []):
|
|
if isinstance(o, dict):
|
|
val = str(o.get("value") if o.get("value") is not None
|
|
else o.get("label") or "").strip()
|
|
lab = str(o.get("label") or val).strip()
|
|
else:
|
|
val = lab = str(o).strip()
|
|
if val:
|
|
opts.append({"value": val, "label": lab or val})
|
|
if len(opts) < 2:
|
|
raise ValueError(f"field {key!r}: {ftype} needs >= 2 options")
|
|
field["options"] = opts
|
|
other = f.get("other")
|
|
field["other"] = field["required"] if other is None else bool(other)
|
|
if ftype in ("number", "slider"):
|
|
for k in ("min", "max", "step"):
|
|
v = _num(f.get(k))
|
|
if v is not None:
|
|
field[k] = v
|
|
if ftype == "slider":
|
|
field.setdefault("min", 0)
|
|
field.setdefault("max", 100)
|
|
field.setdefault("step", 1)
|
|
if ftype == "file" and f.get("multiple"):
|
|
field["multiple"] = True
|
|
if f.get("default") is not None:
|
|
field["default"] = f.get("default")
|
|
out.append(field)
|
|
return out
|
|
|
|
|
|
def apply_assets(value, urls: dict[str, str]):
|
|
"""Replace ``{{asset:<name>}}`` placeholders in every string of a field
|
|
spec (recursively) with the URL the asset was stored under. A placeholder
|
|
naming an unknown asset is left as is."""
|
|
if isinstance(value, str):
|
|
return ASSET_RE.sub(lambda m: urls.get(m.group(1), m.group(0)), value)
|
|
if isinstance(value, list):
|
|
return [apply_assets(x, urls) for x in value]
|
|
if isinstance(value, dict):
|
|
return {k: apply_assets(x, urls) for k, x in value.items()}
|
|
return value
|
|
|
|
|
|
def check_answers(fields: list[dict], answers: dict) -> dict:
|
|
"""Validate submitted answers against the spec; raises ValueError.
|
|
|
|
Returns the answers reduced to known keys — lenient on shapes (the UI is
|
|
the trusted producer) but strict on required fields and option membership,
|
|
so the agent never reads back a value the form couldn't have produced."""
|
|
if not isinstance(answers, dict):
|
|
raise ValueError("answers must be an object")
|
|
out: dict = {}
|
|
for f in fields:
|
|
key, ftype = f["key"], f["type"]
|
|
if ftype in DISPLAY_TYPES:
|
|
continue
|
|
v = answers.get(key)
|
|
# "Other" ticked but nothing typed is no answer at all.
|
|
empty = v is None or v == "" or v == [] or v == OTHER
|
|
if empty:
|
|
if f.get("required") and ftype != "checkbox":
|
|
raise ValueError(f"field {f['label']!r} is required")
|
|
if ftype == "checkbox":
|
|
out[key] = bool(v)
|
|
continue
|
|
if ftype in ("select", "radio"):
|
|
allowed = {o["value"] for o in f.get("options") or []}
|
|
if str(v) not in allowed and not f.get("other"):
|
|
raise ValueError(f"field {f['label']!r}: {v!r} not an option")
|
|
out[key] = str(v)
|
|
elif ftype == "multiselect":
|
|
allowed = {o["value"] for o in f.get("options") or []}
|
|
vals = [str(x) for x in (v if isinstance(v, list) else [v])]
|
|
vals = [x for x in vals if x != OTHER]
|
|
bad = [x for x in vals if x not in allowed]
|
|
if bad and not f.get("other"):
|
|
raise ValueError(f"field {f['label']!r}: {bad} not options")
|
|
if not vals:
|
|
if f.get("required"):
|
|
raise ValueError(f"field {f['label']!r} is required")
|
|
continue
|
|
out[key] = vals
|
|
elif ftype == "checkbox":
|
|
out[key] = bool(v)
|
|
elif ftype in ("number", "slider"):
|
|
n = _num(v)
|
|
if n is None:
|
|
raise ValueError(f"field {f['label']!r}: not a number")
|
|
out[key] = n
|
|
elif ftype == "file":
|
|
files = v if isinstance(v, list) else [v]
|
|
keep = []
|
|
for x in files:
|
|
if isinstance(x, dict) and (x.get("repoPath") or x.get("url")):
|
|
keep.append({k: x[k] for k in
|
|
("name", "size", "contentType", "repoPath",
|
|
"url") if k in x})
|
|
out[key] = keep
|
|
else:
|
|
out[key] = str(v)
|
|
return out
|
|
|
|
|
|
def answer_lines(fields: list[dict], answers: dict) -> list[str]:
|
|
"""Checked answers → ``- **Label** — value`` markdown bullets.
|
|
|
|
Mirrors the frontend's ``answerLines`` (PlanQuestionsForm.tsx) so a plan's
|
|
``## Answers`` section reads the same as the composer prompt the answers
|
|
also feed. Unanswered non-checkbox fields are skipped."""
|
|
out: list[str] = []
|
|
for f in fields:
|
|
if f["type"] in DISPLAY_TYPES:
|
|
continue
|
|
v = answers.get(f["key"])
|
|
if f["type"] != "checkbox" and (v is None or v == "" or v == []):
|
|
continue
|
|
if f["type"] == "checkbox":
|
|
text = "yes" if v else "no"
|
|
elif f["type"] == "file":
|
|
files = v if isinstance(v, list) else []
|
|
text = ", ".join(x.get("repoPath") or x.get("name") or ""
|
|
for x in files if isinstance(x, dict)).strip(", ")
|
|
else:
|
|
by_val = {o["value"]: o["label"] for o in f.get("options") or []}
|
|
if isinstance(v, list):
|
|
text = ", ".join(by_val.get(str(x), str(x)) for x in v)
|
|
else:
|
|
text = by_val.get(str(v), str(v))
|
|
if text:
|
|
# A multiline textarea answer stays one markdown bullet.
|
|
out.append(f"- **{f['label']}** — " + text.replace("\n", "\n "))
|
|
return out
|
|
|
|
|
|
class FormStore:
|
|
"""Form definitions + answers in one JSON sidecar."""
|
|
|
|
def __init__(self, path: str):
|
|
self.path = pathlib.Path(path)
|
|
self._lock = threading.Lock()
|
|
self._waiters: dict[str, threading.Event] = {}
|
|
self._data: dict = {"forms": []}
|
|
self._load()
|
|
|
|
def _load(self) -> None:
|
|
try:
|
|
data = json.loads(self.path.read_text(encoding="utf-8"))
|
|
self._data = {"forms": list(data.get("forms") or [])}
|
|
except (OSError, ValueError):
|
|
self._data = {"forms": []}
|
|
|
|
def _save(self) -> None:
|
|
self.path.parent.mkdir(parents=True, exist_ok=True)
|
|
tmp = self.path.with_suffix(".json.tmp")
|
|
tmp.write_text(json.dumps(self._data, indent=2), encoding="utf-8")
|
|
os.replace(tmp, self.path)
|
|
|
|
def create(self, title: str, fields: list[dict], *, description: str = "",
|
|
type_: str = "", url: str = "", session_id: str = "",
|
|
fid: str = "") -> dict:
|
|
"""``fid``: the notify CLI mints ids itself (``frm_`` + 12 hex) so a
|
|
form created through a worker's outbox keeps the id the CLI printed."""
|
|
entry = {
|
|
"id": fid or "frm_" + uuid.uuid4().hex[:12],
|
|
"title": title, "description": description or "",
|
|
"type": type_ or "", "url": url or "",
|
|
"sessionId": session_id or "", "at": _now_iso(),
|
|
"fields": fields, "status": "pending",
|
|
"answers": None, "submittedAt": None, "cancelledAt": None,
|
|
}
|
|
with self._lock:
|
|
self._waiters[entry["id"]] = threading.Event()
|
|
self._data["forms"].append(entry)
|
|
if len(self._data["forms"]) > _MAX_LOG:
|
|
self._data["forms"] = self._data["forms"][-_MAX_LOG:]
|
|
self._save()
|
|
return json.loads(json.dumps(entry))
|
|
|
|
def get(self, fid: str) -> dict | None:
|
|
with self._lock:
|
|
for f in self._data["forms"]:
|
|
if f.get("id") == fid:
|
|
return json.loads(json.dumps(f))
|
|
return None
|
|
|
|
def list(self, *, session_id: str | None = None,
|
|
limit: int = 100) -> list[dict]:
|
|
with self._lock:
|
|
forms = [f for f in self._data["forms"]
|
|
if not session_id or f.get("sessionId") == session_id]
|
|
out = forms[-max(1, limit):]
|
|
out.reverse() # newest first
|
|
return json.loads(json.dumps(out))
|
|
|
|
def _finish(self, fid: str, patch: dict) -> dict | None:
|
|
"""Apply a terminal patch; idempotent — a form that already reached a
|
|
terminal state is returned untouched (a double-submit or a cancel
|
|
racing a submit can't clobber the recorded answers)."""
|
|
with self._lock:
|
|
cur = next((f for f in self._data["forms"]
|
|
if f.get("id") == fid), None)
|
|
if not cur:
|
|
return None
|
|
if cur.get("status") == "pending":
|
|
cur.update(patch)
|
|
self._save()
|
|
ev = self._waiters.pop(fid, None)
|
|
out = json.loads(json.dumps(cur))
|
|
if ev:
|
|
ev.set()
|
|
return out
|
|
|
|
def submit(self, fid: str, answers: dict) -> dict | None:
|
|
"""Record the answers; raises ValueError on invalid input."""
|
|
cur = self.get(fid)
|
|
if not cur:
|
|
return None
|
|
checked = check_answers(cur.get("fields") or [], answers)
|
|
return self._finish(fid, {"status": "submitted", "answers": checked,
|
|
"submittedAt": _now_iso()})
|
|
|
|
def cancel(self, fid: str) -> dict | None:
|
|
return self._finish(fid, {"status": "cancelled",
|
|
"cancelledAt": _now_iso()})
|
|
|
|
def wait(self, fid: str, wait_s: float) -> dict | None:
|
|
"""Block up to ``wait_s`` for a submit/cancel; returns the form.
|
|
|
|
The waiter event is (re-)armed lazily, so waits survive a backend
|
|
restart (which empties ``_waiters``) — the poll loop just re-arms."""
|
|
cur = self.get(fid)
|
|
if not cur or cur.get("status") != "pending" or wait_s <= 0:
|
|
return cur
|
|
with self._lock:
|
|
ev = self._waiters.setdefault(fid, threading.Event())
|
|
ev.wait(min(wait_s, 55.0))
|
|
return self.get(fid)
|