diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8524003..f4a9ae9 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -18,7 +18,7 @@ jobs: python-version: "3.11" - name: Python syntax (py_compile) - run: python3 -m py_compile gateway.py llm_claude.py transport_telegram.py hooks/*.py bin/news-fetch + run: python3 -m py_compile gateway.py llm_claude.py transport_telegram.py transport_matrix.py hooks/*.py bin/news-fetch - name: Shell syntax (bash -n) run: | diff --git a/CHANGELOG.md b/CHANGELOG.md index ebbda8d..60387d3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,24 @@ All notable changes to this project are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/), and this project aims to follow [Semantic Versioning](https://semver.org/). +## [2.5.0] — 2026-07-10 + +### Added +- **Matrix: inbound images reach the assistant layer (vision).** Images sent to + the bot used to be silently dropped — `updates()` only passed `m.text`. The + matrix transport now yields media messages (`m.image`/`m.file`/`m.video`/ + `m.audio`) with a `media` dict (mxc url, filename, mimetype, size; the + caption arrives as the message text) plus a `download_media()` helper that + speaks the authenticated media endpoint (Matrix 1.11) with a legacy + `/_matrix/media/v3` fallback. The gateway downloads images (capped at 10 MB), + saves them into the LLM workspace (`media/`, pruned after 14 days), and hands + Claude a prompt pointing at the file — Claude Code's Read tool views images + natively, so the existing headless session sees and responds to the picture, + caption included. Non-image files get a polite notice instead; guests and + command-only mode are refused the same way as free text. Encrypted media + (`content.file`) is out of scope, matching this backend's no-E2EE design. + Ticket: 20260709-215539. + ## [2.4.0] — 2026-07-09 ### Added diff --git a/README.md b/README.md index a0d6348..8deac86 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,8 @@ analyze. One always-on Python process (stdlib only) long-polls Telegram — or, with `OGMA_TRANSPORT=matrix`, your own Matrix homeserver, keeping the whole channel on infrastructure -you control. No inbound ports, no pip installs, no API key plumbing — the optional LLM layer +you control. On Matrix, images you send reach the assistant layer too — Claude views and +responds to them. No inbound ports, no pip installs, no API key plumbing — the optional LLM layer reuses your existing Claude Code auth on the box, and without it there is no Claude dependency at all (`OGMA_LLM=off`). diff --git a/gateway.py b/gateway.py index 71c1b2e..707ab1c 100755 --- a/gateway.py +++ b/gateway.py @@ -99,6 +99,15 @@ _confirm_lock = threading.Lock() DOC_THRESHOLD = 8000 # command output longer than this is sent as a file, not chunks +# Inbound media (images sent to the bot). Downloads are capped — this box has +# 2 GB of RAM and Claude's vision path tops out around this size anyway — and +# saved copies are pruned after MEDIA_KEEP_DAYS so the workspace can't grow +# unbounded on the SD card. +MEDIA_MAX_BYTES = 10 * 1024 * 1024 +MEDIA_KEEP_DAYS = 14 +MEDIA_EXT = {"image/jpeg": ".jpg", "image/png": ".png", + "image/gif": ".gif", "image/webp": ".webp"} + def log(*a: object) -> None: print(time.strftime("%Y-%m-%d %H:%M:%S"), *a, flush=True) @@ -450,8 +459,70 @@ def launch_detached(script: str) -> bool: return True +def media_to_prompt(chat_id: str, actor: str, caption: str, guest: bool, + media: dict) -> str | None: + """Turn an inbound media message into a free-text prompt for the LLM layer. + + Sends the refusal/notice itself and returns None when there is nothing to + hand to Claude (guest, no LLM, non-image, download failure). Images are + saved into the LLM workspace and the prompt points Claude at the file — + its Read tool views images natively, so no API plumbing is needed here. + """ + filename, mimetype = media.get("filename", "file"), media.get("mimetype", "") + if guest: + audit("guest-refused", actor, f"media {filename}") + send(chat_id, GUEST_MSG) + return None + if LLM is None: + send(chat_id, NO_LLM_MSG) + return None + if not mimetype.startswith("image/"): + send(chat_id, f"📎 Got {filename} ({mimetype or 'unknown type'}) — " + "I can only look at images for now.") + return None + size = int(media.get("size", 0) or 0) + if size > MEDIA_MAX_BYTES: + send(chat_id, f"⚠️ That image is {size // (1024 * 1024)} MB — " + f"I only take up to {MEDIA_MAX_BYTES // (1024 * 1024)} MB.") + return None + try: + data = TRANSPORT.download_media(media["url"], MEDIA_MAX_BYTES) + except Exception as e: # noqa: BLE001 + log("media download failed:", e) + send(chat_id, "⚠️ Couldn't fetch that image from the homeserver — try resending it.") + return None + media_dir = Path(LLM.workdir) / "media" + safe = re.sub(r"[^A-Za-z0-9._-]+", "_", filename).strip("._")[-60:] or "image" + if not Path(safe).suffix: + safe += MEDIA_EXT.get(mimetype, ".img") + path = media_dir / f"{time.strftime('%Y%m%d-%H%M%S')}-{safe}" + try: + media_dir.mkdir(parents=True, exist_ok=True) + path.write_bytes(data) + cutoff = time.time() - MEDIA_KEEP_DAYS * 86400 + for old in media_dir.iterdir(): + if old.is_file() and old.stat().st_mtime < cutoff: + old.unlink(missing_ok=True) + except OSError as e: + log("media save failed:", e) + send(chat_id, "⚠️ Couldn't store that image on disk.") + return None + audit("media", actor, filename, bytes=len(data), path=str(path)) + prompt = (f"I just sent you an image over chat; it is saved at {path} — " + "view it with the Read tool, then respond to it.") + if caption.strip(): + prompt += f"\nMy message with it: {caption.strip()}" + return prompt + + def handle(chat_id: str, actor: str, text: str, guest: bool, - sent_ts: float = 0.0) -> None: + sent_ts: float = 0.0, media: dict | None = None) -> None: + if media is not None: + prompt = media_to_prompt(chat_id, actor, text, guest, media) + if prompt is None: + return + text = prompt # falls through to the free-text LLM path below + stripped = text.strip() parts = stripped.split(maxsplit=1) cmd = parts[0].lower() if parts else "" @@ -574,13 +645,14 @@ def handle(chat_id: str, actor: str, text: str, guest: bool, send(chat_id, reply) -def _worker(chat_id: str, actor: str, text: str, guest: bool, sent_ts: float) -> None: +def _worker(chat_id: str, actor: str, text: str, guest: bool, sent_ts: float, + media: dict | None = None) -> None: """Handle one message in its own thread, then release the per-chat slot.""" # Log completion + duration: the receipt line alone can't distinguish "still # running", "killed mid-flight", and "reply delayed by a Telegram stall". t0 = time.monotonic() try: - handle(chat_id, actor, text, guest, sent_ts) + handle(chat_id, actor, text, guest, sent_ts, media) except Exception as e: # noqa: BLE001 log("handler error:", e) send(chat_id, "⚠️ Something broke handling that. Logged it.") @@ -639,7 +711,10 @@ def main() -> None: # bot otherwise. Small offsets are clock noise, only log real lag. sent_ts = float(upd.get("date") or 0) lag = (time.time() - sent_ts) if sent_ts else 0.0 - log(f"[{chat_id}] {text[:80]}" + media = upd.get("media") + log(f"[{chat_id}] " + + (f"[{media['msgtype']} {media.get('filename', '')}] " if media else "") + + text[:80] + (f" (sent {lag:.0f}s ago)" if lag > 5 else "")) # Concurrency guard: one in-flight message per chat. Drop a second one # (with a notice) rather than overlapping runs. Handle in a thread so a @@ -652,7 +727,8 @@ def main() -> None: send(chat_id, "⏳ Still working on your previous message — give me a moment, " "then resend if needed.") continue - threading.Thread(target=_worker, args=(chat_id, actor, text, guest, sent_ts), + threading.Thread(target=_worker, + args=(chat_id, actor, text, guest, sent_ts, media), daemon=True).start() diff --git a/transport_matrix.py b/transport_matrix.py index 1abbd5c..d3348ed 100644 --- a/transport_matrix.py +++ b/transport_matrix.py @@ -11,6 +11,13 @@ Implements the same surface as transport_telegram.py: send_document(chat_id, filename, content, caption) -> None long output as a file typing(chat_id) -> None best-effort activity indicator +Beyond that surface, media messages (m.image/m.file/m.video/m.audio) are +yielded with an extra "media" dict ({"url" mxc://, "msgtype", "filename", +"mimetype", "size"}; "text" carries the caption when the client sent one), and +download_media(mxc, max_bytes) fetches the bytes — the gateway uses the pair to +show images to the LLM layer. Transports without these simply never yield +"media", so the gateway needs no capability check. + Speaks the plain Matrix client-server API (long-polling /sync) against your own homeserver over TLS — no third party sees the channel, which is the point of this backend. End-to-end encryption is NOT implemented: E2EE needs olm and a @@ -240,20 +247,73 @@ class MatrixTransport: or ev.get("sender") == self.user_id): continue content = ev.get("content") or {} - if content.get("msgtype") != "m.text": - continue - yield { + msgtype = content.get("msgtype") + upd = { "chat_id": room_id, "from_id": ev.get("sender", ""), "actor": ev.get("sender", ""), "text": str(content.get("body", "")), "date": int(ev.get("origin_server_ts", 0)) // 1000, } + if msgtype == "m.text": + yield upd + continue + if msgtype not in ("m.image", "m.file", "m.video", "m.audio"): + continue + # Media message. Unencrypted events carry a plain mxc:// in + # content.url (encrypted ones use content.file, which this + # backend deliberately doesn't speak — see module docstring). + url = str(content.get("url", "")) + if not url.startswith("mxc://"): + continue + info = content.get("info") or {} + filename = str(content.get("filename", "")) or upd["text"] or "file" + # Per spec, body is the caption only when a separate + # filename field exists and differs — otherwise it is just + # the filename repeated, which would read as a prompt. + if upd["text"] == filename or not content.get("filename"): + upd["text"] = "" + upd["media"] = { + "url": url, + "msgtype": msgtype, + "filename": filename, + "mimetype": str(info.get("mimetype", "")), + "size": int(info.get("size", 0) or 0), + } + yield upd new_since = resp.get("next_batch", "") if new_since and new_since != since: since = new_since self._save_since(since) + def download_media(self, mxc: str, max_bytes: int) -> bytes: + """Fetch an mxc:// object from the homeserver, capped at max_bytes. + + Tries the authenticated media endpoint (Matrix 1.11+, the only one + modern servers still serve) and falls back to the legacy media repo for + older homeservers. Raises on a malformed url, an oversized object, or + when both endpoints fail — the caller owns the user-facing message. + """ + server, _, media_id = mxc.removeprefix("mxc://").partition("/") + if not mxc.startswith("mxc://") or not server or not media_id or "/" in media_id: + raise ValueError(f"not a valid mxc url: {mxc!r}") + last_err: Exception = ValueError("no media endpoint answered") + for path in (f"/_matrix/client/v1/media/download/{_q(server)}/{_q(media_id)}", + f"/_matrix/media/v3/download/{_q(server)}/{_q(media_id)}"): + req = urllib.request.Request( + f"{self.homeserver}{path}", + headers={"Authorization": f"Bearer {self.token}"}) + try: + with urllib.request.urlopen(req, timeout=60) as r: + data = r.read(max_bytes + 1) + except (urllib.error.URLError, OSError) as e: + last_err = e + continue + if len(data) > max_bytes: + raise ValueError(f"media larger than {max_bytes} bytes") + return data + raise last_err + # ----------------------------------------------------------------------- # Sending # -----------------------------------------------------------------------