cli/notify.py merges the lab's notify-done, notify-ask and ask-form scripts
into one stdlib-only CLI shipped with this repo (the aliases stay as shims),
so a session on a remote worker has it too. Two transports: the hub directly
on the lab, or — on a worker — the sidecar's new /outbox, which the hub's
OutboxMirror polls every 2 s, creating the record through the same functions
/api/notify, /api/ask and /api/forms run and pushing the answer back. The
worker still never calls the hub.
The hub now owns what the CLI used to compute: the default conversation URL
(worker transcripts included), the `notified` meta stamp, and a `--final`
push marks the session finished. A push's --action-cmd button answers on
/api/notify/{id}/action like an ask (HA integration updated). Ids are minted
by the CLI (ntf_/ask_/frm_) so viewer, phone and worker agree. Dropped: the
direct-HA fallback, the DONE button, the Forge cost line, the dashboard
mirror. The viewer's widgets also parse the `notify send|ask|form|wait`
spelling; the worker installer links the skill into ~/.claude/skills.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
703 lines
29 KiB
Python
Executable File
703 lines
29 KiB
Python
Executable File
#!/usr/bin/env python3
|
||
"""notify — the one CLI an agent session uses to reach Gabriel through the
|
||
ai-agent hub: a push notification, a 2-3-button choice, or a structured form.
|
||
|
||
notify send [TITLE] MESSAGE [--type T] [--url U] [--image I] [--say S]
|
||
[--action-cmd "Title::cmd"] [--final] [--no-cost] [--dry-run]
|
||
notify ask "QUESTION" "Opt 1" "Opt 2" ["Opt 3"] [--background|-b]
|
||
[--timeout N] [--type T] [--url U] [--image I] [--dry-run]
|
||
notify form --title T [--description D] [--type T]
|
||
(--fields JSON | --fields-file F | JSON on stdin) [--dry-run]
|
||
notify wait ID [--timeout N] ask_… → chosen label; frm_… → answers JSON
|
||
notify cancel ID frm_… (a pending form)
|
||
notify cost the "<harness> · $x · Nk tokens · 8m" line
|
||
|
||
Two transports, picked by the environment the session was launched with:
|
||
|
||
* **hub** (the lab) — talks to the ai-agent backend at ``AI_AGENT_URL``
|
||
(default ``http://127.0.0.1:8096``), the notification source of truth.
|
||
* **outbox** (a remote worker) — the sidecar that launched this run exported
|
||
``AI_AGENT_OUTBOX_URL``: the request is dropped into *that sidecar's* outbox
|
||
and the hub, which polls every worker it has paired with, picks it up,
|
||
records it, forwards it to the phone and pushes any answer back into the
|
||
outbox. The worker never calls the hub (see sidecar/outbox.py).
|
||
|
||
Ids are minted **here** (``ntf_`` / ``ask_`` / ``frm_`` + 12 hex) so they are
|
||
the same on the worker, in the hub and in the viewer whichever transport
|
||
carried them. ``wait`` and ``cancel`` then work on either side.
|
||
|
||
Deliberately stdlib-only (urllib, argparse): a worker has whatever python3 the
|
||
Mac has and nothing installed for us. The old notify-done / notify-ask /
|
||
ask-form skill scripts are thin shims onto this file.
|
||
|
||
The compat aliases map 1:1: notify-done → send, notify-ask → ask,
|
||
ask-form → form (plus its own wait/cancel).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import json
|
||
import os
|
||
import re
|
||
import subprocess
|
||
import sys
|
||
import time
|
||
import urllib.error
|
||
import urllib.parse
|
||
import urllib.request
|
||
import uuid
|
||
from pathlib import Path
|
||
|
||
HERE = Path(__file__).resolve().parent
|
||
|
||
# One long-poll cycle. The hub (and the outbox) cap the in-request wait at 55 s,
|
||
# so a backgrounded wait is one request a minute that still returns the instant
|
||
# an answer lands.
|
||
POLL_WAIT_S = 55
|
||
|
||
NOTIFY_TYPES = ("deploy", "code", "config", "build", "git", "commit", "service",
|
||
"container", "ai", "mcp", "music", "download", "media",
|
||
"jellyfin", "research", "logs", "metrics", "dns", "auth",
|
||
"security", "mail", "success", "done", "backup", "error",
|
||
"fail", "warn", "warning", "notify", "default", "")
|
||
|
||
# What a token costs, USD per million — the hub's own table (backend/
|
||
# conversations.py PRICES) so the line in the push matches the dashboards.
|
||
# Cache write = 1.25× input (5-minute ephemeral), cache read = 0.1× input.
|
||
PRICES = {"opus": (5.0, 25.0), "sonnet": (3.0, 15.0), "haiku": (1.0, 5.0),
|
||
"fable": (10.0, 50.0)}
|
||
_ID_RE = re.compile(r"^(ntf|ask|frm)_[0-9a-f]{12}$")
|
||
|
||
|
||
def _err(msg: str) -> None:
|
||
print(msg, file=sys.stderr, flush=True)
|
||
|
||
|
||
def _mint(prefix: str) -> str:
|
||
return f"{prefix}_{uuid.uuid4().hex[:12]}"
|
||
|
||
|
||
# ── .env.claude (private defaults, gitignored) ───────────────────────────────
|
||
def load_env_claude(start: Path = HERE) -> None:
|
||
"""Source the first ``.env.claude`` found walking up from ``start`` (bash
|
||
``set -a; . file`` semantics: file values override the environment)."""
|
||
cur = start.resolve()
|
||
while True:
|
||
f = cur / ".env.claude"
|
||
if f.is_file():
|
||
try:
|
||
text = f.read_text(encoding="utf-8", errors="replace")
|
||
except OSError:
|
||
return
|
||
for line in text.splitlines():
|
||
line = line.strip()
|
||
if not line or line.startswith("#"):
|
||
continue
|
||
if line.startswith("export "):
|
||
line = line[7:].lstrip()
|
||
m = re.match(r"^([A-Za-z_][A-Za-z0-9_]*)=(.*)$", line)
|
||
if not m:
|
||
continue
|
||
k, v = m.group(1), m.group(2)
|
||
if len(v) >= 2 and v[0] == v[-1] and v[0] in "'\"":
|
||
v = v[1:-1]
|
||
os.environ[k] = v
|
||
return
|
||
if cur == cur.parent:
|
||
return
|
||
cur = cur.parent
|
||
|
||
|
||
# ── session / transcript ─────────────────────────────────────────────────────
|
||
def config_dirs() -> list[Path]:
|
||
"""Every Claude config dir a transcript may live under, most relevant
|
||
first: $CLAUDE_CONFIG_DIR (the run's own account), ~/.claude, then any
|
||
extra in $CLAUDE_CONFIG_DIRS (colon-separated)."""
|
||
raw = [os.environ.get("CLAUDE_CONFIG_DIR", ""), str(Path.home() / ".claude"),
|
||
str(Path.home() / ".claude-work")]
|
||
raw += [p for p in os.environ.get("CLAUDE_CONFIG_DIRS", "").split(":") if p]
|
||
out: list[Path] = []
|
||
for p in raw:
|
||
if p:
|
||
d = Path(os.path.expanduser(p))
|
||
if d not in out:
|
||
out.append(d)
|
||
return out
|
||
|
||
|
||
def find_transcript(sid: str) -> Path | None:
|
||
for d in config_dirs():
|
||
root = d / "projects"
|
||
if root.is_dir():
|
||
hits = sorted(root.glob(f"*/{sid}.jsonl"))
|
||
if hits:
|
||
return hits[0]
|
||
return None
|
||
|
||
|
||
def newest_transcript_for_cwd() -> Path | None:
|
||
enc = os.getcwd().replace("/", "-")
|
||
best, best_m = None, -1.0
|
||
for d in config_dirs():
|
||
for f in (d / "projects" / enc).glob("*.jsonl"):
|
||
try:
|
||
m = f.stat().st_mtime
|
||
except OSError:
|
||
continue
|
||
if m > best_m:
|
||
best, best_m = f, m
|
||
return best
|
||
|
||
|
||
def resolve_session_id() -> str:
|
||
"""The calling Claude Code session: $CLAUDE_SESSION_ID (the sidecar exports
|
||
it for every run it launches), else an external resolver the lab's shims
|
||
point at ($NOTIFY_SESSION_RESOLVER — conv-meta's PID-ancestry match),
|
||
else the newest transcript written for this cwd."""
|
||
sid = os.environ.get("CLAUDE_SESSION_ID", "").strip()
|
||
if sid:
|
||
return sid
|
||
resolver = os.environ.get("NOTIFY_SESSION_RESOLVER", "")
|
||
if resolver and os.path.isfile(resolver):
|
||
try:
|
||
out = subprocess.run(["bash", resolver], capture_output=True,
|
||
text=True, timeout=30).stdout.strip()
|
||
if out:
|
||
return out
|
||
except (OSError, subprocess.TimeoutExpired):
|
||
pass
|
||
t = newest_transcript_for_cwd()
|
||
return t.name[:-6] if t else ""
|
||
|
||
|
||
# ── cost line ────────────────────────────────────────────────────────────────
|
||
def _rates(model: str) -> tuple[float, float]:
|
||
m = (model or "").lower()
|
||
for key, r in PRICES.items():
|
||
if key in m:
|
||
return r
|
||
if "mythos" in m:
|
||
return PRICES["fable"]
|
||
return PRICES["opus"]
|
||
|
||
|
||
def transcript_cost(path: Path) -> dict:
|
||
"""Sum a transcript's per-turn usage (deduped by requestId) and price it."""
|
||
seen: set[str] = set()
|
||
tokens = 0
|
||
cost = 0.0
|
||
first = last = None
|
||
with open(path, encoding="utf-8", errors="replace") as fh:
|
||
for line in fh:
|
||
try:
|
||
o = json.loads(line)
|
||
except ValueError:
|
||
continue
|
||
if o.get("type") != "assistant":
|
||
continue
|
||
req = o.get("requestId")
|
||
if req:
|
||
if req in seen:
|
||
continue
|
||
seen.add(req)
|
||
msg = o.get("message") or {}
|
||
u = msg.get("usage") or {}
|
||
i = int(u.get("input_tokens") or 0)
|
||
out = int(u.get("output_tokens") or 0)
|
||
cw = int(u.get("cache_creation_input_tokens") or 0)
|
||
cr = int(u.get("cache_read_input_tokens") or 0)
|
||
if not (i or out or cw or cr):
|
||
continue
|
||
pi, po = _rates(msg.get("model") or "")
|
||
cost += (i * pi + out * po + cw * pi * 1.25 + cr * pi * 0.1) / 1e6
|
||
tokens += i + out + cw + cr
|
||
ts = o.get("timestamp")
|
||
if ts:
|
||
first = ts if first is None else min(first, ts)
|
||
last = ts if last is None else max(last, ts)
|
||
duration = ""
|
||
if first and last and last > first:
|
||
from datetime import datetime
|
||
try:
|
||
secs = int((datetime.fromisoformat(last.replace("Z", "+00:00"))
|
||
- datetime.fromisoformat(first.replace("Z", "+00:00")))
|
||
.total_seconds())
|
||
m, s = divmod(secs, 60)
|
||
h, m = divmod(m, 60)
|
||
duration = f"{h}h {m}m {s}s" if h else (f"{m}m {s}s" if m else f"{s}s")
|
||
except ValueError:
|
||
pass
|
||
return {"cost": cost, "tokens": tokens, "duration": duration}
|
||
|
||
|
||
def cost_line(sid: str) -> str:
|
||
"""``Claude Code · $1.57 · 78k tokens · 8m 38s`` for the session, or ''
|
||
when no transcript resolves (never fatal)."""
|
||
try:
|
||
t = find_transcript(sid) if sid else None
|
||
t = t or newest_transcript_for_cwd()
|
||
if not t:
|
||
return ""
|
||
c = transcript_cost(t)
|
||
except Exception:
|
||
return ""
|
||
n = c["tokens"]
|
||
tok = (f"{n / 1e6:.1f}M tokens" if n >= 1_000_000
|
||
else f"{n // 1000}k tokens" if n >= 1000 else f"{n} tokens")
|
||
line = f"Claude Code · ${c['cost']:.2f} · {tok}"
|
||
if c["duration"]:
|
||
line += f" · {c['duration']}"
|
||
return line
|
||
|
||
|
||
# ── transport ────────────────────────────────────────────────────────────────
|
||
class Transport:
|
||
"""``hub`` POSTs straight to the backend; ``outbox`` drops the same payload
|
||
into the local sidecar for the hub to collect."""
|
||
|
||
def __init__(self) -> None:
|
||
self.outbox = os.environ.get("AI_AGENT_OUTBOX_URL", "").rstrip("/")
|
||
self.outbox_token = os.environ.get("AI_AGENT_OUTBOX_TOKEN", "")
|
||
self.hub = (os.environ.get("AI_AGENT_URL") or "http://127.0.0.1:8096").rstrip("/")
|
||
self.mode = "outbox" if self.outbox else "hub"
|
||
|
||
def _call(self, url: str, data: dict | None = None, *, method: str = "POST",
|
||
timeout: float = 15, token: str = "") -> tuple[int, dict]:
|
||
body = json.dumps(data).encode() if data is not None else None
|
||
headers = {"Content-Type": "application/json"}
|
||
if token:
|
||
headers["Authorization"] = f"Bearer {token}"
|
||
req = urllib.request.Request(url, data=body, method=method, headers=headers)
|
||
try:
|
||
with urllib.request.urlopen(req, timeout=timeout) as r:
|
||
return r.status, _json(r.read().decode(errors="replace"))
|
||
except urllib.error.HTTPError as e:
|
||
return e.code, _json(e.read().decode(errors="replace"))
|
||
except (urllib.error.URLError, OSError, ValueError) as e:
|
||
return 0, {"error": str(e)}
|
||
|
||
# create
|
||
def submit(self, kind: str, payload: dict) -> tuple[int, dict]:
|
||
if self.mode == "outbox":
|
||
return self._call(f"{self.outbox}/outbox",
|
||
{"kind": kind, "payload": payload},
|
||
token=self.outbox_token)
|
||
path = {"notify": "/api/notify", "ask": "/api/ask", "form": "/api/forms"}[kind]
|
||
return self._call(f"{self.hub}{path}", payload)
|
||
|
||
# read / long-poll
|
||
def status(self, rid: str, wait_s: int = 0) -> tuple[int, dict]:
|
||
q = f"?waitSecs={wait_s}" if wait_s > 0 else ""
|
||
if self.mode == "outbox":
|
||
return self._call(f"{self.outbox}/outbox/{rid}{q}", method="GET",
|
||
timeout=wait_s + 15, token=self.outbox_token)
|
||
path = {"ntf": "/api/notify/", "ask": "/api/ask/", "frm": "/api/forms/"}[rid[:3]]
|
||
return self._call(f"{self.hub}{path}{rid}{q}", method="GET",
|
||
timeout=wait_s + 15)
|
||
|
||
def cancel(self, rid: str) -> tuple[int, dict]:
|
||
if self.mode == "outbox":
|
||
return self._call(f"{self.outbox}/outbox/{rid}/cancel", {},
|
||
token=self.outbox_token)
|
||
return self._call(f"{self.hub}/api/forms/{rid}/cancel", {})
|
||
|
||
def describe(self) -> str:
|
||
return f"outbox {self.outbox}" if self.mode == "outbox" else f"hub {self.hub}"
|
||
|
||
|
||
def _json(text: str) -> dict:
|
||
try:
|
||
d = json.loads(text or "{}")
|
||
return d if isinstance(d, dict) else {"value": d}
|
||
except ValueError:
|
||
return {"raw": text[:300]}
|
||
|
||
|
||
def _accepted(code: int, resp: dict) -> bool:
|
||
"""A create is accepted when the hub recorded it (partial webhook delivery
|
||
still counts — the record is answerable from the viewer) or the outbox
|
||
queued it."""
|
||
if code == 202 or resp.get("queued"):
|
||
return True
|
||
if code != 200:
|
||
return False
|
||
return bool(resp.get("id")) and (
|
||
resp.get("ok") is True
|
||
or any(d.get("ok") for d in (resp.get("delivered") or []))
|
||
or "delivered" not in resp)
|
||
|
||
|
||
def _delivery_warning(resp: dict) -> str:
|
||
failed = [str(d.get("webhook")) for d in (resp.get("delivered") or [])
|
||
if not d.get("ok")]
|
||
return f"warn: webhook(s) failed: {', '.join(failed)}" if failed else ""
|
||
|
||
|
||
# ── og:image (deploy pushes preview the site's social card) ──────────────────
|
||
def og_image_for(url: str) -> str:
|
||
"""``<origin>/og-image.png`` when it answers 200 — publicly trusted hosts
|
||
only (the phone can't validate the lab's certs)."""
|
||
if os.environ.get("HA_NOTIFY_NO_OG"):
|
||
return ""
|
||
try:
|
||
u = urllib.parse.urlsplit(url.strip())
|
||
except ValueError:
|
||
return ""
|
||
host = u.hostname or ""
|
||
if not u.scheme.startswith("http") or not host or host.endswith(".lab.gabvdl.xyz"):
|
||
return ""
|
||
name = os.environ.get("HA_NOTIFY_OG_NAME", "og-image.png").lstrip("/")
|
||
og = f"{u.scheme}://{u.netloc}/{name}"
|
||
try:
|
||
req = urllib.request.Request(og, method="HEAD")
|
||
with urllib.request.urlopen(req, timeout=5) as r:
|
||
code = r.status
|
||
except urllib.error.HTTPError as e:
|
||
code = e.code
|
||
except (urllib.error.URLError, OSError, ValueError):
|
||
code = 0
|
||
if code == 200:
|
||
_err(f"attaching og:image {og}")
|
||
return og
|
||
_err(f"warn: og:image not attached ({og} → HTTP {code})")
|
||
return ""
|
||
|
||
|
||
# ── wait loops ───────────────────────────────────────────────────────────────
|
||
def wait_for(t: Transport, rid: str, timeout: int, *, quiet: bool = False) -> int:
|
||
"""Long-poll until the record reaches a terminal state.
|
||
|
||
ask_: prints the chosen label (exit 0); expired → 3.
|
||
frm_: prints the answers as one JSON line (0); cancelled → 3; timeout → 4.
|
||
ntf_: prints the tapped action title (0); expired → 3.
|
||
An unknown id is 1. ``timeout`` 0 = no local deadline."""
|
||
deadline = time.time() + timeout if timeout > 0 else 0
|
||
kind = rid[:3]
|
||
while True:
|
||
code, r = t.status(rid, POLL_WAIT_S)
|
||
if code == 404:
|
||
_err(f"error: unknown {rid}")
|
||
return 1
|
||
st = r.get("status") or ""
|
||
if code == 0 or not st:
|
||
_err("warn: backend unreachable (retrying in 10s)")
|
||
time.sleep(10)
|
||
elif kind == "ask" and st == "answered":
|
||
label = (r.get("answer") or "").replace("\n", " ")
|
||
_err(f"chosen: {r.get('answerIndex')} {label}")
|
||
print(label, flush=True)
|
||
return 0
|
||
elif kind == "frm" and st == "submitted":
|
||
print(json.dumps(r.get("answers") or {}, ensure_ascii=False), flush=True)
|
||
return 0
|
||
elif kind == "ntf" and st == "acted":
|
||
print(r.get("actionTitle") or str(r.get("actionIndex")), flush=True)
|
||
return 0
|
||
elif st in ("expired", "cancelled"):
|
||
if not quiet:
|
||
_err(f"{st}: {rid} ended without an answer")
|
||
return 3
|
||
if deadline and time.time() >= deadline:
|
||
if not quiet:
|
||
_err(f"timeout: no answer within {timeout}s ({rid} still pending)")
|
||
return 4 if kind == "frm" else 3
|
||
|
||
|
||
def _spawn_detached(argv: list[str]) -> None:
|
||
try:
|
||
subprocess.Popen(argv, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL,
|
||
stderr=subprocess.DEVNULL, start_new_session=True)
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
# ── commands ─────────────────────────────────────────────────────────────────
|
||
def cmd_send(a: argparse.Namespace) -> int:
|
||
if not a.args:
|
||
_err("usage: notify send [TITLE] MESSAGE [--type T] [--url U] …")
|
||
return 2
|
||
title, message = ("Claude Code", a.args[0]) if len(a.args) == 1 else (a.args[0], a.args[1])
|
||
ntype = (a.type if a.type is not None else os.environ.get("HA_NOTIFY_TYPE", "")).strip()
|
||
if ntype not in NOTIFY_TYPES:
|
||
_err(f"error: unknown --type {ntype!r} (one of: "
|
||
f"{', '.join(t for t in NOTIFY_TYPES if t)})")
|
||
return 2
|
||
url = (a.url if a.url is not None else os.environ.get("HA_NOTIFY_URL", "")).strip()
|
||
image = (a.image if a.image is not None else os.environ.get("HA_NOTIFY_IMAGE", "")).strip()
|
||
spoken = (a.spoken if a.spoken is not None else os.environ.get("HA_NOTIFY_SPOKEN", "")).strip()
|
||
if not image and ntype == "deploy" and url:
|
||
image = og_image_for(url)
|
||
|
||
sid = resolve_session_id()
|
||
if not a.no_cost and not os.environ.get("HA_NOTIFY_NO_COST"):
|
||
line = cost_line(sid)
|
||
if line:
|
||
message = f"{message}\n{line}"
|
||
|
||
nid = _mint("ntf")
|
||
payload: dict = {"id": nid, "title": title, "message": message}
|
||
for k, v in (("type", ntype), ("url", url), ("image", image), ("spoken", spoken),
|
||
("channel", a.channel), ("importance", a.importance),
|
||
("vibrationPattern", a.vibration), ("sessionId", sid)):
|
||
if (v or "").strip():
|
||
payload[k] = v.strip()
|
||
if a.final:
|
||
payload["final"] = True
|
||
action_cmd = ""
|
||
if a.action_cmd:
|
||
if "::" not in a.action_cmd:
|
||
_err("error: --action-cmd must be '<title>::<command>'")
|
||
return 2
|
||
action_title, action_cmd = a.action_cmd.split("::", 1)
|
||
payload["actions"] = [{"action": f"{nid}__0", "title": action_title.strip()}]
|
||
|
||
t = Transport()
|
||
if a.dry_run:
|
||
print(f"DRY RUN — {t.describe()} notify payload:")
|
||
print(json.dumps(payload, ensure_ascii=False))
|
||
return 0
|
||
if not url:
|
||
_err("no --url given: the hub links this conversation in the viewer")
|
||
code, resp = t.submit("notify", payload)
|
||
if not _accepted(code, resp):
|
||
_err(f"error: notification not sent via {t.describe()} "
|
||
f"(HTTP {code}: {json.dumps(resp)[:300]})")
|
||
return 1
|
||
w = _delivery_warning(resp)
|
||
if w:
|
||
_err(w)
|
||
print(f"notification sent via {t.describe()} ({nid})", flush=True)
|
||
if action_cmd:
|
||
# Detached waiter: runs the command when the button is tapped (6 h —
|
||
# a "Merge to main" tap may come much later). Survives this call.
|
||
_spawn_detached([sys.executable, str(Path(__file__).resolve()), "_on-action",
|
||
nid, os.environ.get("HA_ACTION_TIMEOUT", "21600"), action_cmd])
|
||
return 0
|
||
|
||
|
||
def cmd_on_action(a: argparse.Namespace) -> int:
|
||
"""Hidden: wait for a notify's action button, then run its command."""
|
||
rc = wait_for(Transport(), a.id, int(a.timeout), quiet=True)
|
||
if rc != 0:
|
||
return 0
|
||
return subprocess.run(["bash", "-c", a.command]).returncode
|
||
|
||
|
||
def cmd_ask(a: argparse.Namespace) -> int:
|
||
if a.args and a.args[0] == "wait": # `notify-ask wait <id>` compat
|
||
a.id = next((x for x in a.args[1:] if x.startswith("ask_")), "")
|
||
return cmd_wait(a)
|
||
question, options = (a.args[0] if a.args else ""), [o for o in a.args[1:] if o.strip()]
|
||
if not question.strip() or len(options) < 2:
|
||
_err('usage: notify ask "<question>" "Label1" "Label2" [Label3] [--background]')
|
||
return 2
|
||
if len(options) > 3:
|
||
_err("error: Android shows at most 3 action buttons; pass 2 or 3 options")
|
||
return 2
|
||
ntype = (a.type if a.type is not None else os.environ.get("HA_NOTIFY_TYPE", "")).strip()
|
||
if ntype not in NOTIFY_TYPES:
|
||
_err(f"error: unknown --type {ntype!r}")
|
||
return 2
|
||
timeout = int(a.timeout if a.timeout is not None else os.environ.get("HA_ASK_TIMEOUT", "180"))
|
||
sid = resolve_session_id()
|
||
aid = _mint("ask")
|
||
payload: dict = {"id": aid, "question": question.strip(), "options": options,
|
||
"timeoutSecs": timeout}
|
||
url = (a.url if a.url is not None else os.environ.get("HA_NOTIFY_URL", "")).strip()
|
||
image = (a.image if a.image is not None else os.environ.get("HA_NOTIFY_IMAGE", "")).strip()
|
||
for k, v in (("type", ntype), ("url", url), ("image", image), ("sessionId", sid)):
|
||
if v:
|
||
payload[k] = v
|
||
t = Transport()
|
||
if a.dry_run:
|
||
print(f"DRY RUN — {t.describe()} ask payload:")
|
||
print(json.dumps(payload, ensure_ascii=False))
|
||
return 0
|
||
code, resp = t.submit("ask", payload)
|
||
if not _accepted(code, resp):
|
||
_err(f"error: ask not sent via {t.describe()} (HTTP {code}: {json.dumps(resp)[:300]})")
|
||
return 1
|
||
w = _delivery_warning(resp)
|
||
if w:
|
||
_err(w)
|
||
if a.background:
|
||
_err(f"asked via {t.describe()} (not waiting): {question}")
|
||
print(f"ASK_ID={aid}", flush=True)
|
||
_err(f"next: run `notify wait {aid}` under a Monitor (or in a background "
|
||
f"shell) — do NOT block the turn on it")
|
||
return 0
|
||
_err(f"asked via {t.describe()} (waiting up to {timeout}s for a tap): {question}")
|
||
return wait_for(t, aid, timeout + 5 if timeout > 0 else 0)
|
||
|
||
|
||
def cmd_form(a: argparse.Namespace) -> int:
|
||
if a.args and a.args[0] in ("wait", "cancel"): # `ask-form wait|cancel <id>` compat
|
||
a.id = next((x for x in a.args[1:] if x.startswith("frm_")), "")
|
||
return cmd_wait(a) if a.args[0] == "wait" else cmd_cancel(a)
|
||
if not a.title:
|
||
_err("error: --title is required")
|
||
return 2
|
||
raw = a.fields or ""
|
||
if a.fields_file:
|
||
try:
|
||
raw = Path(a.fields_file).read_text()
|
||
except OSError as e:
|
||
_err(f"error: cannot read --fields-file: {e}")
|
||
return 1
|
||
if not raw and not sys.stdin.isatty():
|
||
raw = sys.stdin.read()
|
||
if not raw.strip():
|
||
_err("error: no fields (use --fields / --fields-file / stdin)")
|
||
return 2
|
||
try:
|
||
fields = json.loads(raw)
|
||
except json.JSONDecodeError as e:
|
||
_err(f"error: --fields is not valid JSON: {e}")
|
||
return 1
|
||
if not isinstance(fields, list) or not fields:
|
||
_err("error: --fields must be a non-empty JSON array")
|
||
return 1
|
||
sid = resolve_session_id()
|
||
fid = _mint("frm")
|
||
payload = {"id": fid, "title": a.title, "description": a.description or "",
|
||
"type": a.type or "ai", "sessionId": sid, "fields": fields,
|
||
# The hub pushes the "📋 title" notification itself, deep-linked
|
||
# to this conversation with ?form=<id>.
|
||
"notify": True}
|
||
t = Transport()
|
||
if a.dry_run:
|
||
print(f"DRY RUN — {t.describe()} form payload:")
|
||
print(json.dumps(payload, ensure_ascii=False))
|
||
return 0
|
||
code, resp = t.submit("form", payload)
|
||
if not _accepted(code, resp):
|
||
_err(f"error: form create failed via {t.describe()} "
|
||
f"(HTTP {code}: {json.dumps(resp)[:400]})")
|
||
return 1
|
||
print(f"FORM_ID={fid}", flush=True)
|
||
focus = resp.get("focusUrl") or resp.get("url")
|
||
if focus:
|
||
print(f"focus: {focus}", flush=True)
|
||
_err(f"next: run `notify wait {fid}` under a Monitor (or in a background "
|
||
f"shell) — do NOT block the turn on it")
|
||
return 0
|
||
|
||
|
||
def cmd_wait(a: argparse.Namespace) -> int:
|
||
rid = (a.id or "").strip()
|
||
if not _ID_RE.match(rid):
|
||
_err("usage: notify wait <ask_…|frm_…|ntf_…> [--timeout N]")
|
||
return 2
|
||
return wait_for(Transport(), rid, int(a.timeout or 0))
|
||
|
||
|
||
def cmd_cancel(a: argparse.Namespace) -> int:
|
||
rid = (a.id or "").strip()
|
||
if not rid.startswith("frm_"):
|
||
_err("usage: notify cancel <frm_…>")
|
||
return 2
|
||
code, resp = Transport().cancel(rid)
|
||
if code not in (200, 202):
|
||
_err(f"error: cancel failed (HTTP {code}: {json.dumps(resp)[:200]})")
|
||
return 1
|
||
print(f"cancelled {rid}", flush=True)
|
||
return 0
|
||
|
||
|
||
def cmd_cost(a: argparse.Namespace) -> int:
|
||
line = cost_line(resolve_session_id())
|
||
if not line:
|
||
_err("error: no transcript resolves for this session")
|
||
return 1
|
||
print(line)
|
||
return 0
|
||
|
||
|
||
# ── argv ─────────────────────────────────────────────────────────────────────
|
||
def build_parser() -> argparse.ArgumentParser:
|
||
p = argparse.ArgumentParser(prog="notify", description=__doc__,
|
||
formatter_class=argparse.RawDescriptionHelpFormatter)
|
||
sub = p.add_subparsers(dest="cmd", metavar="COMMAND")
|
||
|
||
s = sub.add_parser("send", help="push a notification (alias: notify-done)")
|
||
s.add_argument("args", nargs="*", metavar="[TITLE] MESSAGE")
|
||
s.add_argument("--url", "-u")
|
||
s.add_argument("--type", "-t")
|
||
s.add_argument("--image", "-i")
|
||
s.add_argument("--spoken", "--say", dest="spoken",
|
||
help="one short French sentence the desk phone reads aloud")
|
||
s.add_argument("--channel")
|
||
s.add_argument("--importance")
|
||
s.add_argument("--vibration", "--vibrate", dest="vibration")
|
||
s.add_argument("--action-cmd", dest="action_cmd",
|
||
help="extra button '<title>::<shell command>' run on tap")
|
||
s.add_argument("--final", action="store_true",
|
||
help="this is the task's closing notification: the hub marks "
|
||
"the session finished")
|
||
s.add_argument("--no-cost", action="store_true")
|
||
s.add_argument("--dry-run", action="store_true")
|
||
s.set_defaults(fn=cmd_send)
|
||
|
||
k = sub.add_parser("ask", help="2-3-button choice (alias: notify-ask)")
|
||
k.add_argument("args", nargs="*", metavar='"QUESTION" OPTION…')
|
||
k.add_argument("--timeout")
|
||
k.add_argument("--type", "-t")
|
||
k.add_argument("--url", "-u")
|
||
k.add_argument("--image", "-i")
|
||
k.add_argument("--target", help=argparse.SUPPRESS) # legacy, ignored
|
||
k.add_argument("--background", "--no-wait", "-b", action="store_true")
|
||
k.add_argument("--dry-run", action="store_true")
|
||
k.set_defaults(fn=cmd_ask, id="")
|
||
|
||
f = sub.add_parser("form", help="structured form in the viewer (alias: ask-form)")
|
||
f.add_argument("args", nargs="*", help=argparse.SUPPRESS)
|
||
f.add_argument("--title", "-T")
|
||
f.add_argument("--description", "-d")
|
||
f.add_argument("--type", "-t")
|
||
f.add_argument("--fields")
|
||
f.add_argument("--fields-file", dest="fields_file")
|
||
f.add_argument("--timeout", help=argparse.SUPPRESS)
|
||
f.add_argument("--dry-run", action="store_true")
|
||
f.set_defaults(fn=cmd_form, id="")
|
||
|
||
w = sub.add_parser("wait", help="collect an ask's / form's answer")
|
||
w.add_argument("id")
|
||
w.add_argument("--timeout", default="0")
|
||
w.set_defaults(fn=cmd_wait)
|
||
|
||
c = sub.add_parser("cancel", help="dismiss a pending form")
|
||
c.add_argument("id")
|
||
c.set_defaults(fn=cmd_cancel)
|
||
|
||
sub.add_parser("cost", help="print this session's cost line").set_defaults(fn=cmd_cost)
|
||
|
||
h = sub.add_parser("_on-action")
|
||
h.add_argument("id")
|
||
h.add_argument("timeout")
|
||
h.add_argument("command")
|
||
h.set_defaults(fn=cmd_on_action)
|
||
return p
|
||
|
||
|
||
def main(argv: list[str] | None = None) -> int:
|
||
load_env_claude()
|
||
argv = list(sys.argv[1:] if argv is None else argv)
|
||
# Compat entry points set NOTIFY_AS=send|ask|form (the shims) so the old
|
||
# argv shapes keep working: `notify-done "msg" --url …`.
|
||
alias = os.environ.get("NOTIFY_AS", "")
|
||
if alias in ("send", "ask", "form") and (not argv or argv[0] not in
|
||
("send", "ask", "form", "wait",
|
||
"cancel", "cost", "_on-action")):
|
||
argv = [alias, *argv]
|
||
p = build_parser()
|
||
a = p.parse_args(argv)
|
||
if not getattr(a, "fn", None):
|
||
p.print_help()
|
||
return 2
|
||
return a.fn(a)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|