diff --git a/.env.example b/.env.example index d67e19a..9b26bd8 100644 --- a/.env.example +++ b/.env.example @@ -17,6 +17,10 @@ TELEGRAM_ALLOWED_USERS= # Per-message timeout in seconds (default 300) # CLAUDE_TIMEOUT=300 +# Max concurrent Claude runs across all chats (default 1). Keep at 1 on small boxes +# (a Pi can OOM if several run at once); raise only on a roomy host. +# OGMA_MAX_CONCURRENT=1 + # --- optional: permissions / model --- # Permission posture for tools. Empty = safest (tools needing approval are skipped). # Set to "acceptEdits" to let it edit files, or use OGMA_ALLOWED_TOOLS for a curated set. diff --git a/CHANGELOG.md b/CHANGELOG.md index 6449b8c..c5245f9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,20 @@ 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/). +## [0.5.0] — 2026-06-16 + +### Added +- **Concurrency guard.** Messages are now handled in per-chat worker threads with **one in-flight + message per chat** (a second is dropped with a notice rather than overlapping), plus a global + semaphore that caps concurrent `claude` runs via **`OGMA_MAX_CONCURRENT`** (default 1 — protects + RAM on small boxes like a Pi). A long run in one chat no longer blocks polling or other chats. +- **`bin/setup --check`** — a non-interactive doctor that validates the token (via getMe), allow-list, + model/effort/fallback, the `claude` CLI, the systemd service, and installed skills. Changes nothing. + +### Changed +- The gateway validates the bot token at startup — one clear `token OK` / `TOKEN CHECK FAILED` log + line — instead of emitting a stream of 401s from getUpdates when the token is wrong. + ## [0.4.0] — 2026-06-16 ### Added @@ -92,6 +106,7 @@ First public release. A minimal, self-hosted bridge from Telegram to Claude Code - Single shared brain — multiple allow-listed chats share one persona/workspace/memory. Per-user isolation is planned (see issues). +[0.5.0]: https://github.com/eric-wien/ogma/releases/tag/v0.5.0 [0.4.0]: https://github.com/eric-wien/ogma/releases/tag/v0.4.0 [0.3.0]: https://github.com/eric-wien/ogma/releases/tag/v0.3.0 [0.2.1]: https://github.com/eric-wien/ogma/releases/tag/v0.2.1 diff --git a/README.md b/README.md index 2907d93..02dc7f6 100644 --- a/README.md +++ b/README.md @@ -60,7 +60,9 @@ cd ~/ogma bin/setup # interactive: token, persona, skills, systemd, chat-ID — all guided ``` The script walks you through everything below and never transmits anything off your machine. The -manual steps are documented here too, in case you prefer to do it by hand. +manual steps are documented here too, in case you prefer to do it by hand. To re-check an existing +install without changing anything, run `bin/setup --check` — it validates your token, allow-list, +model/effort/fallback, the `claude` CLI, the service, and installed skills. ### Manual setup @@ -126,6 +128,7 @@ See [`skills/README.md`](skills/README.md) for details and how to write your own | `TELEGRAM_ALLOWED_USERS` | comma-separated allowed chat IDs (**required**) | — | | `CLAUDE_BIN` | path to the `claude` CLI | `~/.local/bin/claude` | | `CLAUDE_TIMEOUT` | per-message timeout (seconds) | `300` | +| `OGMA_MAX_CONCURRENT` | max concurrent Claude runs across chats (raise only on a roomy host) | `1` | | `OGMA_WORKDIR` | Claude's working dir | `./workspace` | | `OGMA_PERMISSION_MODE` | e.g. `acceptEdits`; empty = safest | empty | | `OGMA_ALLOWED_TOOLS` | curated tool allow-list | empty | diff --git a/bin/setup b/bin/setup index 3ca5789..b594b39 100755 --- a/bin/setup +++ b/bin/setup @@ -50,6 +50,43 @@ open(path, "w").write("\n".join(out) + "\n") PY } +# Non-interactive health/config check — `bin/setup --check`. Changes nothing. +do_check() { + local tok allowed cb u problems=0 + step "Ogma config check ${c_dim}($OGMA_DIR)${c_rst}" + tok="$(grep -E '^TELEGRAM_BOT_TOKEN=' "$ENV" 2>/dev/null | cut -d= -f2-)" + allowed="$(grep -E '^TELEGRAM_ALLOWED_USERS=' "$ENV" 2>/dev/null | cut -d= -f2-)" + cb="$(command -v claude || echo "$HOME/.local/bin/claude")" + [ -f "$ENV" ] && ok ".env present" || { err ".env missing — run bin/setup"; problems=1; } + [ -n "$tok" ] && ok "token set" || { err "token MISSING"; problems=1; } + [ -n "$allowed" ] && ok "allowed chats: $allowed" || { err "allowed chats MISSING"; problems=1; } + printf ' model : %s\n' "$(grep -E '^OGMA_MODEL=' "$ENV" 2>/dev/null | cut -d= -f2- | sed 's/^$/(Claude Code default)/')" + printf ' effort : %s\n' "$(grep -E '^OGMA_EFFORT=' "$ENV" 2>/dev/null | cut -d= -f2- | sed 's/^$/(default)/')" + printf ' fallback : %s\n' "$(grep -E '^OGMA_FALLBACK_MODEL=' "$ENV" 2>/dev/null | cut -d= -f2- | sed 's/^$/(none)/')" + [ -x "$cb" ] && ok "claude CLI: $cb" || { err "claude CLI missing ($cb)"; problems=1; } + if [ -n "$tok" ] && command -v curl >/dev/null; then + u="$(curl -sS -m 12 "https://api.telegram.org/bot${tok}/getMe" 2>/dev/null | "$PY" -c 'import sys,json +try: d=json.load(sys.stdin) +except Exception: print("?"); raise SystemExit +print(("OK "+((d.get("result") or {}).get("username") or "?")) if d.get("ok") else "BAD")' 2>/dev/null)" + case "$u" in "OK "*) ok "token valid — bot @${u#OK }";; *) err "token REJECTED by Telegram (getMe failed)"; problems=1;; esac + fi + if command -v systemctl >/dev/null && systemctl --user show-environment >/dev/null 2>&1; then + printf ' service : %s\n' "$(systemctl --user is-active ogma-gateway 2>/dev/null || echo inactive)" + fi + for s in tickets session-search daily-briefing; do + [ -e "$HOME/.claude/skills/$s" ] && ok "skill: $s" || warn "skill not installed: $s" + done + say "" + [ "$problems" = 0 ] && say "${c_grn}All good.${c_rst}" || say "${c_yel}Issues found — see above.${c_rst}" + return "$problems" +} + +case "${1:-}" in + --check|check) do_check; exit $? ;; + -h|--help) echo "Usage: bin/setup [--check]"; exit 0 ;; +esac + say "${c_bold}Ogma setup${c_rst} ${c_dim}($OGMA_DIR)${c_rst}" say "This configures your own copy. Nothing is sent anywhere; secrets you enter stay on this box." @@ -321,4 +358,5 @@ else say "Next: once token + allowed users are set, start it:" say " ${c_dim}systemctl --user enable --now ogma-gateway${c_rst} (or: python3 $OGMA_DIR/gateway.py)" fi -say "\nRead ${c_bold}docs/workflow.md${c_rst} for how Ogma is meant to be used. Enjoy." +say "" +say "Read ${c_bold}docs/workflow.md${c_rst} for how Ogma is meant to be used. Enjoy." diff --git a/gateway.py b/gateway.py index 87b4bae..ffe4494 100755 --- a/gateway.py +++ b/gateway.py @@ -66,6 +66,13 @@ CLAUDE_TIMEOUT = int(os.environ.get("CLAUDE_TIMEOUT", "300")) SESSIONS_FILE = BASE / "sessions.json" ENV_FILE = BASE / ".env" EFFORT_LEVELS = ("low", "medium", "high", "xhigh", "max") +# Max concurrent Claude runs. Default 1 — small boxes (e.g. a Pi) OOM if several run at once. +MAX_CONCURRENT = max(1, int(cfg("MAX_CONCURRENT", "1") or "1")) + +_inflight: set = set() # chat_ids with a message currently being handled +_inflight_lock = threading.Lock() +_sessions_lock = threading.Lock() # guards the shared sessions dict + file +_run_sem = threading.Semaphore(MAX_CONCURRENT) # bounds concurrent `claude` invocations API = f"https://api.telegram.org/bot{TOKEN}" TG_MAX = 4000 # Telegram hard limit is 4096; leave headroom @@ -284,8 +291,9 @@ def handle(chat_id: str, text: str, sessions: dict[str, str]) -> None: send(chat_id, HELP_TEXT) return if cmd == "/new": - sessions.pop(chat_id, None) - save_sessions(sessions) + with _sessions_lock: + sessions.pop(chat_id, None) + save_sessions(sessions) send(chat_id, "🧹 Fresh session.") return if cmd == "/model": @@ -306,15 +314,44 @@ def handle(chat_id: str, text: str, sessions: dict[str, str]) -> None: typer = threading.Thread(target=keep_typing, args=(chat_id, stop), daemon=True) typer.start() try: - reply, sid = ask_claude(text, sessions.get(chat_id)) + with _sessions_lock: + prior = sessions.get(chat_id) + with _run_sem: # bound concurrent claude runs (RAM safety on small boxes) + reply, sid = ask_claude(text, prior) finally: stop.set() - if sid and sid != sessions.get(chat_id): - sessions[chat_id] = sid - save_sessions(sessions) + if sid: + with _sessions_lock: + if sid != sessions.get(chat_id): + sessions[chat_id] = sid + save_sessions(sessions) send(chat_id, reply) +def validate_token() -> tuple[bool, str]: + """Check the bot token via getMe so a bad token is an obvious one-line log, + not an endless stream of 401s from getUpdates.""" + try: + r = tg("getMe", {}, timeout=15) + except Exception as e: # noqa: BLE001 + return (False, str(e)) + if r.get("ok"): + return (True, (r.get("result") or {}).get("username", "?")) + return (False, r.get("description", "not ok")) + + +def _worker(chat_id: str, text: str, sessions: dict[str, str]) -> None: + """Handle one message in its own thread, then release the per-chat slot.""" + try: + handle(chat_id, text, sessions) + except Exception as e: # noqa: BLE001 + log("handler error:", e) + send(chat_id, "⚠️ Something broke handling that. Logged it.") + finally: + with _inflight_lock: + _inflight.discard(chat_id) + + def main() -> None: global EFFORT if not TOKEN: @@ -325,8 +362,14 @@ def main() -> None: if EFFORT and EFFORT not in EFFORT_LEVELS: log(f"ignoring invalid OGMA_EFFORT={EFFORT!r} (use one of {', '.join(EFFORT_LEVELS)})") EFFORT = "" + ok_token, info = validate_token() + if ok_token: + log(f"token OK — bot @{info}") + else: + log(f"⚠️ TOKEN CHECK FAILED ({info}). Fix TELEGRAM_BOT_TOKEN in .env and restart.") sessions = load_sessions() - log(f"Ogma up. workdir={WORKDIR} allowed={ALLOWED or '(none — locked down)'}") + log(f"Ogma up. workdir={WORKDIR} allowed={ALLOWED or '(none — locked down)'} " + f"max_concurrent={MAX_CONCURRENT}") offset = 0 while True: try: @@ -348,11 +391,18 @@ def main() -> None: "add it to TELEGRAM_ALLOWED_USERS to enable access.") continue log(f"[{chat_id}] {text[:80]}") - try: - handle(chat_id, text, sessions) - except Exception as e: # noqa: BLE001 - log("handler error:", e) - send(chat_id, "⚠️ Something broke handling that. Logged it.") + # 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 + # long run in one chat doesn't block polling or other chats. + with _inflight_lock: + busy = chat_id in _inflight + if not busy: + _inflight.add(chat_id) + if busy: + 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, text, sessions), daemon=True).start() if __name__ == "__main__":