Spaces:
Running on Zero
Running on Zero
Download app.py from spacedout-bits/Oracle: direct link, hf CLI and curl.
- Browser
- Download file 16.9 kB
-
https://huggingface.co/spaces/spacedout-bits/Oracle/resolve/main/app.py
- Command line
-
hf download hf://spaces/spacedout-bits/Oracle/app.py
-
curl -L -o app.py https://huggingface.co/spaces/spacedout-bits/Oracle/resolve/main/app.py
16.9 kB
| """Entrypoint. Serves the Telegram webhook alongside a one-page Gradio UI. | |
| Two platform constraints shape the structure here, both learned the hard way: | |
| * **Gradio must own the server.** ZeroGPU (the only free hardware for personal | |
| accounts) instruments Gradio's ``launch()`` path to discover ``@spaces.GPU`` | |
| functions. Serving the same app with ``uvicorn.run()`` bypasses that hook and | |
| the Space dies with "No @spaces.GPU function detected during startup" -- even | |
| when such functions plainly exist at module level. So on a Space we launch | |
| Gradio and graft the FastAPI routes onto ``demo.app``. | |
| * **The webhook answers immediately** and does the real work in a background | |
| task. Reading a statement can take a minute; holding the webhook open that | |
| long makes Telegram time out and redeliver, importing the file twice. | |
| The routes live on an ``APIRouter`` so they can be attached either to our own | |
| FastAPI app (tests, self-hosting via the Dockerfile) or to Gradio's (on a | |
| Space), with the shared state reached through ``request.app.state.rt`` rather | |
| than a module global. | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import logging | |
| import os | |
| import threading | |
| from contextlib import asynccontextmanager | |
| from dataclasses import dataclass | |
| from typing import Optional | |
| import gradio as gr | |
| from fastapi import APIRouter, BackgroundTasks, FastAPI, Header, Request | |
| from fastapi.concurrency import run_in_threadpool | |
| from fastapi.responses import JSONResponse | |
| from finbot import __version__, local_llm | |
| from finbot.advice import Advisor | |
| from finbot.analytics import today_in | |
| from finbot.backend import make_llm | |
| from finbot.config import get_settings | |
| from finbot.dashboard import build_ui | |
| from finbot.features import default_features | |
| from finbot.features.base import Attachment, BotContext, Services | |
| from finbot.ingest import Ingestor | |
| from finbot.limits import RateLimiter | |
| from finbot.registry import FeatureRegistry | |
| from finbot.store import open_store | |
| from finbot.telegram import ( | |
| MAX_DOWNLOAD_BYTES, | |
| ParsedUpdate, | |
| SeenUpdates, | |
| TelegramClient, | |
| parse_update, | |
| redact, | |
| ) | |
| logging.basicConfig( | |
| level=os.environ.get("LOG_LEVEL", "INFO"), | |
| format="%(asctime)s %(levelname)-8s %(name)s: %(message)s", | |
| ) | |
| # httpx logs the full request URL at INFO, and the Telegram Bot API carries the | |
| # bot token in the URL *path* (/bot<token>/getMe). On a public Space the | |
| # container logs are readable by anyone, so leaving this at INFO publishes the | |
| # token. Silence it; our own call sites log what matters, redacted. | |
| for noisy in ("httpx", "httpcore"): | |
| logging.getLogger(noisy).setLevel(logging.WARNING) | |
| log = logging.getLogger("finbot") | |
| WEBHOOK_PATH = "/telegram/webhook" | |
| # ---------------------------------------------------------------- runtime | |
| class Runtime: | |
| """Everything a request needs, built once per process.""" | |
| services: Services | |
| telegram: TelegramClient | |
| seen: SeenUpdates | |
| limiter: RateLimiter | |
| #: Set by build_runtime so the Gradio UI, which is constructed at import time, | |
| #: can reach the services once they exist. | |
| _runtime: Optional["Runtime"] = None | |
| def build_runtime() -> Runtime: | |
| """Synchronous construction, so it works from any startup path.""" | |
| global _runtime | |
| settings = get_settings() | |
| registry = FeatureRegistry() | |
| llm = make_llm(settings) | |
| services = Services( | |
| settings=settings, | |
| store=open_store(settings), | |
| llm=llm, | |
| ingestor=Ingestor(llm, settings.default_currency), | |
| advisor=Advisor(llm), | |
| registry=registry, | |
| ) | |
| registry.register_all(default_features()) | |
| for gap in settings.missing_required(): | |
| log.warning("Not configured: %s", gap) | |
| if settings.open_access: | |
| log.warning( | |
| "OPEN ACCESS: anyone on Telegram may use this bot. Ledgers are " | |
| "isolated per user; uploads are capped at %s per user per day.", | |
| settings.uploads_per_user_per_day or "unlimited", | |
| ) | |
| _runtime = Runtime( | |
| services=services, | |
| telegram=TelegramClient(settings.telegram_bot_token), | |
| seen=SeenUpdates(), | |
| limiter=RateLimiter(settings.uploads_per_user_per_day), | |
| ) | |
| return _runtime | |
| def space_webhook_url() -> str: | |
| """Public URL of this Space, if we are running on one.""" | |
| if host := os.environ.get("SPACE_HOST"): | |
| return f"https://{host.rstrip('/')}{WEBHOOK_PATH}" | |
| return "" | |
| async def bootstrap_telegram(rt: Runtime) -> None: | |
| """Identify the bot, publish its command menu and register the webhook. | |
| Uses a throwaway client so this can run on a different event loop from the | |
| one that will later serve requests -- an httpx client is bound to the loop | |
| that created it, and reusing one across loops fails at the worst time. | |
| """ | |
| settings = rt.services.settings | |
| if not settings.telegram_bot_token: | |
| log.warning("TELEGRAM_BOT_TOKEN unset; running without a bot connection.") | |
| return | |
| client = TelegramClient(settings.telegram_bot_token) | |
| try: | |
| me = await client.get_me() | |
| username = me.get("username", "") | |
| rt.telegram.username = username # so /cmd@Bot mentions resolve | |
| log.info("Connected to Telegram as @%s", username) | |
| await client.set_my_commands(rt.services.registry.telegram_commands()) | |
| if not settings.telegram_webhook_secret: | |
| log.warning("TELEGRAM_WEBHOOK_SECRET unset; webhook not registered.") | |
| elif url := space_webhook_url(): | |
| await client.set_webhook(url, settings.telegram_webhook_secret) | |
| log.info("Webhook registered at %s", url) | |
| else: | |
| log.info("No SPACE_HOST; register manually via POST /admin/setup") | |
| except Exception as exc: | |
| # A Telegram hiccup must not stop the app booting: a live container | |
| # that reports the problem beats a crash loop. Redacted because an | |
| # httpx error carries the token-bearing URL. | |
| log.error( | |
| "Telegram bootstrap failed: %s", redact(exc, settings.telegram_bot_token) | |
| ) | |
| finally: | |
| await client.close() | |
| async def shutdown_runtime(rt: Runtime) -> None: | |
| log.info("Shutting down; flushing ledger to the Hub.") | |
| await rt.telegram.close() | |
| rt.services.llm.close() | |
| await run_in_threadpool(rt.services.store.close) | |
| def get_runtime(request: Request) -> Runtime: | |
| return request.app.state.rt | |
| # ----------------------------------------------------------------- routes | |
| router = APIRouter() | |
| async def health() -> JSONResponse: | |
| """Liveness only. No ledger or configuration detail -- this may be public.""" | |
| return JSONResponse({"status": "ok", "version": __version__}) | |
| async def diagnostics( | |
| request: Request, | |
| x_admin_secret: Optional[str] = Header(default=None), | |
| ) -> JSONResponse: | |
| """Operational detail, gated on the webhook secret.""" | |
| rt = get_runtime(request) | |
| settings = rt.services.settings | |
| if ( | |
| not settings.telegram_webhook_secret | |
| or x_admin_secret != settings.telegram_webhook_secret | |
| ): | |
| return JSONResponse({"ok": False, "error": "forbidden"}, status_code=403) | |
| return JSONResponse( | |
| { | |
| "ok": True, | |
| "version": __version__, | |
| "llm_enabled": rt.services.llm.enabled, | |
| "llm": rt.services.llm.describe(), | |
| "sync_enabled": settings.sync_enabled, | |
| "sync_status": rt.services.store.sync_status, | |
| "rows": rt.services.store.count(), | |
| "missing": settings.missing_required(), | |
| } | |
| ) | |
| async def telegram_webhook( | |
| request: Request, | |
| background: BackgroundTasks, | |
| x_telegram_bot_api_secret_token: Optional[str] = Header(default=None), | |
| ) -> JSONResponse: | |
| rt = get_runtime(request) | |
| settings = rt.services.settings | |
| expected = settings.telegram_webhook_secret | |
| if not expected or x_telegram_bot_api_secret_token != expected: | |
| log.warning("Rejected webhook call with a bad secret token.") | |
| return JSONResponse({"ok": False}, status_code=403) | |
| try: | |
| update = await request.json() | |
| except Exception: | |
| return JSONResponse({"ok": False, "error": "bad json"}, status_code=400) | |
| parsed = parse_update(update, rt.telegram.username) | |
| if parsed is None: | |
| return JSONResponse({"ok": True, "skipped": "unsupported update"}) | |
| if not rt.seen.add_if_new(parsed.update_id): | |
| log.info("Ignoring redelivered update %s", parsed.update_id) | |
| return JSONResponse({"ok": True, "skipped": "duplicate"}) | |
| # Answer now, work later: Telegram redelivers if we hold this open. | |
| background.add_task(handle_update, rt, parsed) | |
| return JSONResponse({"ok": True}) | |
| async def admin_setup( | |
| request: Request, | |
| x_admin_secret: Optional[str] = Header(default=None), | |
| ) -> JSONResponse: | |
| """Register the webhook manually (local dev behind a tunnel, or no SPACE_HOST). | |
| Authorised with the same webhook secret, so there is no extra credential. | |
| """ | |
| rt = get_runtime(request) | |
| settings = rt.services.settings | |
| if ( | |
| not settings.telegram_webhook_secret | |
| or x_admin_secret != settings.telegram_webhook_secret | |
| ): | |
| return JSONResponse({"ok": False, "error": "forbidden"}, status_code=403) | |
| body = {} | |
| try: | |
| body = await request.json() | |
| except Exception: | |
| pass | |
| url = body.get("url") or space_webhook_url() | |
| if not url: | |
| return JSONResponse( | |
| {"ok": False, "error": "no url given and SPACE_HOST is unset"}, | |
| status_code=400, | |
| ) | |
| await rt.telegram.set_webhook(url, settings.telegram_webhook_secret) | |
| await rt.telegram.set_my_commands(rt.services.registry.telegram_commands()) | |
| info = await rt.telegram.get_webhook_info() | |
| return JSONResponse({"ok": True, "webhook": info}) | |
| # ------------------------------------------------------------ update handling | |
| async def handle_update(rt: Runtime, parsed: ParsedUpdate) -> None: | |
| services = rt.services | |
| telegram = rt.telegram | |
| settings = services.settings | |
| if not settings.is_authorised(parsed.user_id, parsed.username): | |
| log.warning( | |
| "Unauthorised access attempt from %s (@%s)", | |
| parsed.user_id, | |
| parsed.username or "no-handle", | |
| ) | |
| who = ( | |
| f"@{parsed.username}" | |
| if parsed.username | |
| else f"id <code>{parsed.user_id}</code> (you have no @username set)" | |
| ) | |
| await telegram.send_message( | |
| parsed.chat_id, | |
| "This is a private bot and you're not on its allowlist.\n\n" | |
| f"If it's yours, add {who} to " | |
| "<code>ALLOWED_TELEGRAM_USERNAMES</code> in the Space secrets " | |
| "and restart it.", | |
| ) | |
| return | |
| # Refresh the handle registry on every message, so a rename is picked up | |
| # without anything else needing to change. | |
| await run_in_threadpool( | |
| services.store.record_user, parsed.user_id, parsed.username, parsed.first_name | |
| ) | |
| try: | |
| attachment = None | |
| if parsed.has_file: | |
| if parsed.file_size > MAX_DOWNLOAD_BYTES: | |
| await telegram.send_message( | |
| parsed.chat_id, | |
| "That file is over Telegram's 20 MB bot limit. " | |
| "Try exporting a narrower date range.", | |
| ) | |
| return | |
| # Reading a file occupies the GPU (or spends remote credits), so | |
| # cap it per user. Typing an expense is free and never metered. | |
| if not rt.limiter.check(parsed.user_id): | |
| await telegram.send_message( | |
| parsed.chat_id, | |
| "You've hit today's limit for reading receipts and " | |
| "statements. It resets in 24 hours — you can still log " | |
| "expenses by typing them.", | |
| ) | |
| return | |
| await telegram.send_typing(parsed.chat_id) | |
| data, file_path = await telegram.download(parsed.file_id) | |
| attachment = Attachment( | |
| data=data, | |
| filename=parsed.file_name or file_path.rsplit("/", 1)[-1], | |
| mime_type=parsed.mime_type, | |
| kind=parsed.file_kind, | |
| caption=parsed.caption, | |
| ) | |
| ctx = BotContext( | |
| user_id=parsed.user_id, | |
| chat_id=parsed.chat_id, | |
| services=services, | |
| today=today_in(settings.timezone), | |
| username=parsed.username, | |
| text=parsed.text or parsed.caption, | |
| command=parsed.command, | |
| args=parsed.args, | |
| attachment=attachment, | |
| ) | |
| if attachment is not None or parsed.command in {"advice", "report"}: | |
| await telegram.send_typing(parsed.chat_id) | |
| reply = await run_in_threadpool(services.registry.dispatch, ctx) | |
| if reply is None: | |
| if not parsed.command and ctx.text.strip(): | |
| await telegram.send_message( | |
| parsed.chat_id, | |
| "I couldn't find an amount in that. Try " | |
| "<code>250 coffee</code>, or /help.", | |
| ) | |
| return | |
| if reply.document is not None: | |
| filename, payload = reply.document | |
| await telegram.send_document( | |
| parsed.chat_id, filename, payload, caption=reply.text | |
| ) | |
| return | |
| await telegram.send_message(parsed.chat_id, reply.text, reply.parse_mode) | |
| except Exception: | |
| log.exception("Failed handling update %s", parsed.update_id) | |
| try: | |
| await telegram.send_message( | |
| parsed.chat_id, | |
| "Something went wrong on my side. The details are in the Space logs.", | |
| ) | |
| except Exception: | |
| pass | |
| # ----------------------------------------------------------------- gradio UI | |
| def _get_services() -> Optional[Services]: | |
| """Late-bound accessor for the dashboard. | |
| The UI object is built at import time, before any runtime exists, so it | |
| reaches for services through this rather than capturing them. | |
| """ | |
| return _runtime.services if _runtime is not None else None | |
| demo = build_ui(_get_services, local_llm.warmup) | |
| # ------------------------------------------------------------------- wiring | |
| async def lifespan(fastapi_app: FastAPI): | |
| """Startup path for our own FastAPI app (tests, and the Dockerfile).""" | |
| rt = build_runtime() | |
| fastapi_app.state.rt = rt | |
| await bootstrap_telegram(rt) | |
| try: | |
| yield | |
| finally: | |
| await shutdown_runtime(rt) | |
| app = FastAPI(title="coco-finbot", version=__version__, lifespan=lifespan) | |
| app.include_router(router) | |
| app = gr.mount_gradio_app(app, demo, path="/") | |
| _launch_lock = threading.Lock() | |
| def serve_via_gradio() -> None: | |
| """Startup path for a Gradio Space, where Gradio must own the server.""" | |
| with _launch_lock: | |
| rt = build_runtime() | |
| # On a Space, let Gradio use its own env-driven defaults (GRADIO_SERVER_PORT | |
| # and the Space root path). Forcing a port breaks the frontend asset URLs | |
| # behind HF's proxy: the CSS/JS 404, no JavaScript runs, and the page never | |
| # leaves its initial server-rendered state. That looks exactly like "the | |
| # dashboard shows nothing" -- demo.load never fires, so the sign-in token in | |
| # the query string is never read. | |
| launch_kwargs: dict = { | |
| "prevent_thread_lock": True, | |
| "show_api": False, | |
| # Spaces sets GRADIO_SSR_MODE=1, which switches the frontend to a | |
| # Node-rendered SvelteKit build served from /_app/immutable/*. Launching | |
| # ourselves never starts that SSR server, so every CSS/JS asset 404s, | |
| # no JavaScript runs, and the page freezes on its server-rendered shell. | |
| # Client-side rendering serves the assets straight out of the installed | |
| # gradio package and needs no Node. | |
| "ssr_mode": False, | |
| } | |
| if port := os.environ.get("APP_PORT"): | |
| # Self-hosting / local dev, where nothing sets GRADIO_SERVER_PORT. | |
| launch_kwargs["server_name"] = "0.0.0.0" | |
| launch_kwargs["server_port"] = int(port) | |
| demo.launch(**launch_kwargs) # returns so we can attach routes | |
| # demo.app only exists once launched. | |
| demo.app.state.rt = rt | |
| demo.app.include_router(router) | |
| log.info("Attached %d API routes to the Gradio app", len(router.routes)) | |
| # Runs on its own loop; bootstrap_telegram uses a throwaway client so this | |
| # does not poison the server loop's connections. | |
| asyncio.run(bootstrap_telegram(rt)) | |
| demo.block_thread() | |
| if __name__ == "__main__": | |
| serve_via_gradio() | |