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>
217 lines
9.3 KiB
Python
217 lines
9.3 KiB
Python
"""Creating and recording pushes and asks, and the flat notification feed."""
|
|
|
|
import datetime
|
|
import os
|
|
import pathlib
|
|
import re
|
|
import uuid as uuidlib
|
|
|
|
from fastapi import HTTPException
|
|
|
|
from conversations.cards import _conv_meta, _conv_title
|
|
from conversations.catalog import _conv_ids_by_session
|
|
from conversations.tree import _conv_id
|
|
from core.config import SOURCE_DIRS
|
|
from core.state import AppState
|
|
from notifications import store as notify_mod
|
|
|
|
# ── flat notification feed (server-side mirror of the frontend's buildFeed) ──
|
|
# The PWA flattens /api/notifications + /api/notify-log into one list of feed
|
|
# items (lib/notifications.ts). The phone can't run that, so we do the same
|
|
# flattening here — same stable ids — for the unread endpoint and the read
|
|
# ledger. Keep the two in sync: id scheme is `<convId>#<at>#<index>` for
|
|
# conv-linked pushes and `log#<id>` for hub-only entries.
|
|
_DEDUP_WINDOW_MS = 120_000
|
|
|
|
|
|
def _parse_iso_ms(at: str | None) -> float | None:
|
|
if not at:
|
|
return None
|
|
try:
|
|
s = at.replace("Z", "+00:00")
|
|
return datetime.datetime.fromisoformat(s).timestamp() * 1000.0
|
|
except (ValueError, AttributeError):
|
|
return None
|
|
|
|
|
|
def _flat_notifications(deps: AppState) -> list[dict]:
|
|
"""Every notification as a flat, de-duplicated, newest-first feed item —
|
|
the exact set the PWA shows, so "unread" means the same on the phone."""
|
|
feed: list[dict] = []
|
|
for path, s in deps.store.all_summaries():
|
|
if not s.get("messages") or s.get("isSidechain"):
|
|
continue
|
|
sid = s.get("sessionId") or ""
|
|
notes = deps.meta_store.peek(sid).get("notifications") or []
|
|
if not notes:
|
|
continue
|
|
m = _conv_meta(deps, s)
|
|
cid = _conv_id(path)
|
|
title = _conv_title(deps, s)
|
|
last_seen: dict[str, float] = {}
|
|
for i, n in enumerate(notes):
|
|
key = (f"{n.get('title') or ''} {n.get('body') or ''} "
|
|
f"{n.get('type') or ''} {n.get('url') or ''}")
|
|
t = _parse_iso_ms(n.get("at"))
|
|
prev = last_seen.get(key)
|
|
if prev is not None and (t is None or abs(t - prev) <= _DEDUP_WINDOW_MS):
|
|
continue # duplicate push within the conversation
|
|
if t is not None:
|
|
last_seen[key] = t
|
|
feed.append({
|
|
"id": f"{cid}#{n.get('at') or ''}#{i}",
|
|
"convId": cid, "convTitle": title,
|
|
"title": n.get("title") or "", "body": n.get("body") or "",
|
|
"type": n.get("type") or "", "url": n.get("url") or "",
|
|
"image": n.get("image") or "",
|
|
"at": n.get("at"),
|
|
"projects": m.get("projects") or [],
|
|
"services": m.get("services") or [],
|
|
"kind": "notify",
|
|
"spoken": n.get("spoken") or "",
|
|
})
|
|
# Merge the hub log: a log entry that matches a conv push (title+body within
|
|
# the dedup window) is a duplicate and skipped; hub-only entries are added.
|
|
for e in deps.notify_store.log(limit=500):
|
|
t = _parse_iso_ms(e.get("at"))
|
|
dup = False
|
|
for n in feed:
|
|
if n["title"] == (e.get("title") or "") and \
|
|
n["body"] == (e.get("body") or ""):
|
|
nt = _parse_iso_ms(n.get("at"))
|
|
if t is None or nt is None or abs(nt - t) <= _DEDUP_WINDOW_MS:
|
|
dup = True
|
|
break
|
|
if dup:
|
|
continue
|
|
feed.append({
|
|
"id": f"log#{e.get('id')}",
|
|
"convId": "", "convTitle": "",
|
|
"title": e.get("title") or "", "body": e.get("body") or "",
|
|
"type": e.get("type") or "", "url": e.get("url") or "",
|
|
"image": e.get("image") or "",
|
|
"at": e.get("at"),
|
|
"projects": [], "services": [],
|
|
"kind": e.get("kind") or "notify",
|
|
})
|
|
feed.sort(key=lambda n: n.get("at") or "", reverse=True)
|
|
return feed
|
|
|
|
|
|
# Where the viewer lives, for the links a push carries when the caller gave
|
|
# no --url: the conversation that sent it.
|
|
VIEW_URL = os.environ.get("AI_AGENT_VIEW_URL",
|
|
"https://ai-agent.lab.gabvdl.xyz").rstrip("/")
|
|
_CLIENT_ID_RE = re.compile(r"^(ntf|ask|frm)_[0-9a-f]{12}$")
|
|
|
|
|
|
def _client_id(given: str | None, prefix: str) -> str:
|
|
"""A caller-minted id (the notify CLI's), validated, else a fresh one."""
|
|
given = (given or "").strip()
|
|
if given:
|
|
if not _CLIENT_ID_RE.match(given) or not given.startswith(prefix + "_"):
|
|
raise HTTPException(400, f"bad id {given!r} (want {prefix}_<12 hex>)")
|
|
return given
|
|
return f"{prefix}_{uuidlib.uuid4().hex[:12]}"
|
|
|
|
|
|
def _conversation_url(deps: AppState, session_id: str) -> str:
|
|
"""The viewer page of the session that is notifying — the default link
|
|
of every push, ask and form. Resolved **here** rather than in the CLI
|
|
because only the hub knows where a worker's transcript was mirrored to:
|
|
the indexed summary first, else the transcript on disk under any source
|
|
dir (a session too young to be indexed), else the conversations list."""
|
|
sid = (session_id or "").strip()
|
|
if not sid:
|
|
return f"{VIEW_URL}/conversations"
|
|
cid = _conv_ids_by_session(deps, ).get(sid)
|
|
if not cid:
|
|
for root in SOURCE_DIRS:
|
|
for hit in pathlib.Path(root).glob(f"*/{sid}.jsonl"):
|
|
cid = f"{hit.parent.name}/{hit.name}"
|
|
break
|
|
if cid:
|
|
break
|
|
return f"{VIEW_URL}/conversation/{cid}" if cid else f"{VIEW_URL}/conversations"
|
|
|
|
|
|
def _stamp_notified(deps: AppState, session_id: str, note: dict, final: bool) -> None:
|
|
"""What ``conv-meta notify`` used to do from the skill: record the push on
|
|
the conversation's metadata (the viewer's notification feed + the
|
|
``notified`` lifecycle check) and, for the closing push, mark it
|
|
finished. Best-effort; never fails the push."""
|
|
sid = (session_id or "").strip()
|
|
if not sid:
|
|
return
|
|
patch: dict = {"notification": note}
|
|
if final:
|
|
patch["state"] = "finished"
|
|
try:
|
|
deps.meta_store.update(sid, patch)
|
|
deps.hub.publish({"type": "meta", "id": sid})
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def _create_notify(deps: AppState, d: dict) -> dict:
|
|
"""Record a push and forward it to the subscribed webhooks. ``d`` is a
|
|
NotifyBody dump — the route and the outbox mirror both come through
|
|
here."""
|
|
nid = _client_id(d.get("id"), "ntf")
|
|
sid = (d.get("sessionId") or "").strip()
|
|
if not (d.get("url") or "").strip():
|
|
d["url"] = _conversation_url(deps, sid)
|
|
actions = [a for a in (d.get("actions") or [])
|
|
if isinstance(a, dict) and a.get("action") and a.get("title")]
|
|
payload = {"event": "notify", "id": nid,
|
|
**{k: v for k, v in d.items()
|
|
if k not in ("sessionId", "id", "final") and v not in (None, "", [])}}
|
|
targets = deps.notify_store.targets("notify")
|
|
delivered = notify_mod.forward(targets, payload)
|
|
entry = {
|
|
"id": nid, "kind": "notify", "title": d.get("title") or "",
|
|
"body": d.get("message") or "", "type": d.get("type") or "",
|
|
"url": d.get("url") or "", "image": d.get("image") or "",
|
|
"at": notify_mod._now_iso(), "sessionId": sid,
|
|
"delivered": delivered,
|
|
# Kept (not just forwarded) so the UI can replay what the phone said.
|
|
"spoken": d.get("spoken") or "",
|
|
}
|
|
if actions:
|
|
entry.update({"actions": actions, "status": "sent",
|
|
"actionIndex": None, "actionTitle": None, "actedAt": None})
|
|
deps.notify_store.watch_actions(nid)
|
|
deps.notify_store.record(entry)
|
|
_stamp_notified(deps, sid, {"title": entry["title"], "body": entry["body"],
|
|
"type": entry["type"], "url": entry["url"],
|
|
"at": entry["at"]}, bool(d.get("final")))
|
|
deps.hub.publish({"type": "notification", "id": nid})
|
|
return {"ok": bool(delivered) and all(x["ok"] for x in delivered),
|
|
"id": nid, "delivered": delivered, "url": entry["url"]}
|
|
|
|
|
|
def _create_ask(deps: AppState, d: dict) -> dict:
|
|
"""Create an ask, forward it, return its id (AskBody dump in)."""
|
|
options = [o.strip() for o in (d.get("options") or []) if o and o.strip()]
|
|
if not (d.get("question") or "").strip() or not 2 <= len(options) <= 3:
|
|
raise HTTPException(400, "need a question and 2-3 options")
|
|
nid = _client_id(d.get("id"), "ask")
|
|
sid = (d.get("sessionId") or "").strip()
|
|
url = (d.get("url") or "").strip() or _conversation_url(deps, sid)
|
|
ask = deps.notify_store.create_ask(
|
|
d["question"].strip(), options, type_=d.get("type") or "",
|
|
url=url, image=d.get("image") or "", session_id=sid,
|
|
timeout_s=d.get("timeoutSecs") or notify_mod.DEFAULT_ASK_TIMEOUT_S,
|
|
nid=nid)
|
|
payload = {"event": "ask", "id": ask["id"], "question": ask["body"],
|
|
"options": options, "url": url,
|
|
**{k: v for k, v in d.items()
|
|
if k in ("type", "image", "icon", "color")
|
|
and v not in (None, "")}}
|
|
targets = deps.notify_store.targets("ask")
|
|
delivered = notify_mod.forward(targets, payload)
|
|
deps.notify_store.set_delivered(ask["id"], delivered)
|
|
deps.hub.publish({"type": "notification", "id": ask["id"]})
|
|
return {"ok": bool(delivered) and all(x["ok"] for x in delivered),
|
|
"id": ask["id"], "delivered": delivered, "url": url}
|