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>
237 lines
9.0 KiB
Python
237 lines
9.0 KiB
Python
"""Notification feed, pushes (/api/notify) and asks (/api/ask)."""
|
|
|
|
|
|
from fastapi import APIRouter, HTTPException
|
|
from fastapi.responses import Response
|
|
from pydantic import BaseModel
|
|
|
|
import schemas
|
|
from conversations.cards import _conv_meta, _conv_title
|
|
from conversations.tree import _conv_id
|
|
from core.http import _json_in_worker, _r
|
|
from core.state import State
|
|
from notifications import audio as notify_audio_mod
|
|
from notifications.service import (
|
|
_create_ask,
|
|
_create_notify,
|
|
_flat_notifications,
|
|
)
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
class NotifSeenBody(BaseModel):
|
|
# Notification ids to mark read/seed; `all` marks/seeds the whole feed.
|
|
ids: list[str] | None = None
|
|
all: bool = False
|
|
|
|
|
|
@router.get("/api/notifications", responses=_r(schemas.NotificationsResponse))
|
|
@_json_in_worker
|
|
def notifications_feed(deps: State):
|
|
"""Feed source for the notification history: every conversation that
|
|
recorded a push, with just the fields the feed needs. Replaces walking the
|
|
whole conversation list client-side."""
|
|
out = []
|
|
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)
|
|
out.append({"id": _conv_id(path), "title": _conv_title(deps, s),
|
|
"endedAt": s.get("endedAt"),
|
|
"projects": m.get("projects") or [],
|
|
"services": m.get("services") or [],
|
|
"notifications": notes})
|
|
out.sort(key=lambda c: c.get("endedAt") or "", reverse=True)
|
|
return {"conversations": out, "count": len(out)}
|
|
|
|
|
|
@router.get("/api/notifications/unread",
|
|
responses=_r(schemas.UnreadNotificationsResponse))
|
|
def notifications_unread(deps: State, limit: int = 20):
|
|
"""Unread notifications (not yet opened in the UI), newest first.
|
|
|
|
Pure read: fetching this NEVER marks anything read — the phone can list and
|
|
read these aloud without "consuming" them; only the UI (or an explicit
|
|
``POST /api/notifications/seen``) advances the read ledger. This is the
|
|
endpoint the desk phone hits with its read-only API key.
|
|
"""
|
|
feed = _flat_notifications(deps, )
|
|
unread_ids = set(deps.notif_read_store.unread([n["id"] for n in feed]))
|
|
items = [n for n in feed if n["id"] in unread_ids]
|
|
total = len(items)
|
|
if limit and limit > 0:
|
|
items = items[:limit]
|
|
return {"notifications": items, "count": len(items), "total": total}
|
|
|
|
|
|
class NotifAudioBody(BaseModel):
|
|
"""What a notification sounded like: its type (→ jingle) and spoken line."""
|
|
type: str | None = None
|
|
spoken: str | None = None
|
|
# Read instead of `spoken` when the push carried none (most of the log —
|
|
# --spoken only became mandatory recently).
|
|
title: str | None = None
|
|
message: str | None = None
|
|
|
|
|
|
@router.post("/api/notifications/audio")
|
|
def notification_audio(body: NotifAudioBody):
|
|
"""Render a notification's sound as a WAV so the UI can replay it.
|
|
|
|
The desk phone plays a per-type jazz jingle and then reads the French
|
|
`spoken` line; the notification card's play button asks for those same
|
|
bytes. Rendering happens in the phone service (one synthesizer, one voice —
|
|
see notify_audio.py). A push with no spoken line still gets a voice: its own
|
|
title + message, minus the bits written only for a screen.
|
|
"""
|
|
spoken = (body.spoken or "").strip() or notify_audio_mod.fallback_speech(
|
|
body.title or "", body.message or "")
|
|
try:
|
|
wav = notify_audio_mod.render(spoken, body.type or "")
|
|
except notify_audio_mod.RenderError as e:
|
|
raise HTTPException(502, str(e))
|
|
return Response(content=wav, media_type="audio/wav",
|
|
headers={"Cache-Control": "private, max-age=86400"})
|
|
|
|
|
|
@router.post("/api/notifications/seen", responses=_r(schemas.OkResponse))
|
|
def notifications_mark_seen(deps: State, body: NotifSeenBody):
|
|
"""Mark notification ids as read (opened/dismissed in the UI). Mutating —
|
|
refused to a read-only API-key caller by the gate, so the phone can't use
|
|
it. Passing no ids with ``all=true`` marks the whole current feed read."""
|
|
ids = list(body.ids or [])
|
|
if body.all:
|
|
ids += [n["id"] for n in _flat_notifications(deps, )]
|
|
count = deps.notif_read_store.mark_seen(ids)
|
|
return {"ok": True, "seen": count}
|
|
|
|
|
|
@router.post("/api/notifications/seed", responses=_r(schemas.OkResponse))
|
|
def notifications_seed(deps: State, body: NotifSeenBody):
|
|
"""Adopt the given ids (or the whole current feed) as the read baseline —
|
|
the first-sync op so history isn't reported as unread."""
|
|
ids = list(body.ids or [])
|
|
if body.all or not ids:
|
|
ids += [n["id"] for n in _flat_notifications(deps, )]
|
|
count = deps.notif_read_store.seed(ids)
|
|
return {"ok": True, "seen": count}
|
|
|
|
|
|
class NotifyBody(BaseModel):
|
|
title: str
|
|
message: str
|
|
type: str | None = None
|
|
url: str | None = None
|
|
image: str | None = None
|
|
icon: str | None = None
|
|
color: str | None = None
|
|
channel: str | None = None
|
|
importance: str | None = None
|
|
vibrationPattern: str | None = None
|
|
tag: str | None = None
|
|
persistent: bool | None = None
|
|
actions: list[dict] | None = None
|
|
# A phone-only line to read aloud (the phone webhook's TTS prefers it over
|
|
# title/message). Forwarded verbatim to webhooks; ignored by the HA screen
|
|
# notification. See services/phone/bridge/app.py.
|
|
spoken: str | None = None
|
|
sessionId: str | None = None
|
|
# The notify CLI mints the id itself (``ntf_`` + 12 hex) so a push that
|
|
# travelled through a worker's outbox keeps one id end to end.
|
|
id: str | None = None
|
|
# The task's closing push (``notify send --final``): the session is
|
|
# stamped ``finished`` — the same thing the terminal's DONE line means.
|
|
final: bool | None = None
|
|
|
|
|
|
class AskBody(BaseModel):
|
|
question: str
|
|
options: list[str]
|
|
type: str | None = None
|
|
url: str | None = None
|
|
image: str | None = None
|
|
icon: str | None = None
|
|
color: str | None = None
|
|
timeoutSecs: int | None = None
|
|
sessionId: str | None = None
|
|
id: str | None = None
|
|
|
|
|
|
class NotifyActionBody(BaseModel):
|
|
index: int | None = None
|
|
title: str | None = None
|
|
|
|
|
|
class AskAnswerBody(BaseModel):
|
|
index: int | None = None
|
|
label: str | None = None
|
|
|
|
|
|
@router.post("/api/notify", responses=_r(schemas.NotifyResult))
|
|
def notify_push(deps: State, body: NotifyBody):
|
|
"""Record a push notification and forward it to the subscribed webhooks."""
|
|
return _create_notify(deps, body.model_dump())
|
|
|
|
|
|
@router.get("/api/notify/{notify_id}")
|
|
def notify_get(deps: State, notify_id: str, waitSecs: float = 0):
|
|
"""A recorded push; ``waitSecs`` long-polls for one of its action buttons
|
|
(``notify send --action-cmd``) to be tapped."""
|
|
n = (deps.notify_store.wait_for_action(notify_id, waitSecs) if waitSecs > 0
|
|
else deps.notify_store.get(notify_id))
|
|
if not n or n.get("kind") != "notify":
|
|
raise HTTPException(404, "unknown notification")
|
|
return n
|
|
|
|
|
|
@router.post("/api/notify/{notify_id}/action")
|
|
def notify_action(deps: State, notify_id: str, body: NotifyActionBody):
|
|
"""Callback for the HA integration (and the viewer): which extra button of
|
|
a push was tapped."""
|
|
n = deps.notify_store.record_action(notify_id, body.index, body.title)
|
|
if not n:
|
|
raise HTTPException(404, "unknown notification")
|
|
deps.hub.publish({"type": "notification", "id": notify_id})
|
|
return n
|
|
|
|
|
|
@router.post("/api/ask", responses=_r(schemas.NotifyResult))
|
|
def ask_push(deps: State, body: AskBody):
|
|
"""Create an ask (choice question), forward it to the webhooks, return its
|
|
id. The HA integration renders the tappable notification and POSTs the
|
|
answer back to /api/ask/{id}/answer; callers long-poll GET /api/ask/{id}."""
|
|
return _create_ask(deps, body.model_dump())
|
|
|
|
|
|
@router.get("/api/ask/{ask_id}", responses=_r(schemas.AskRecord))
|
|
def ask_get(deps: State, ask_id: str, waitSecs: float = 0):
|
|
"""The ask's current state; ``waitSecs`` long-polls for the answer."""
|
|
ask = (deps.notify_store.wait_for_answer(ask_id, waitSecs) if waitSecs > 0
|
|
else deps.notify_store.ask_status(ask_id))
|
|
if not ask:
|
|
raise HTTPException(404, "unknown ask")
|
|
return ask
|
|
|
|
|
|
@router.post("/api/ask/{ask_id}/answer", responses=_r(schemas.AskRecord))
|
|
def ask_answer(deps: State, ask_id: str, body: AskAnswerBody):
|
|
"""Callback for the HA integration: record which button was tapped."""
|
|
ask = deps.notify_store.answer_ask(ask_id, body.index, body.label)
|
|
if not ask:
|
|
raise HTTPException(404, "unknown ask")
|
|
deps.hub.publish({"type": "notification", "id": ask_id})
|
|
return ask
|
|
|
|
|
|
@router.get("/api/notify-log", responses=_r(schemas.NotifyLogResponse))
|
|
@_json_in_worker
|
|
def notify_log(deps: State, limit: int = 200):
|
|
"""The first-class notification/ask log, newest first."""
|
|
entries = deps.notify_store.log(limit=limit)
|
|
return {"notifications": entries, "count": len(entries)}
|