From 56be25c1ede2dd22375754f83b824969d82cca9d Mon Sep 17 00:00:00 2001 From: Aurelio Date: Wed, 15 Jul 2026 23:41:44 +0100 Subject: [PATCH] feat(telminal): add hardened Terminal-in-Telegram bot to aurelio Port fristhon/telminal (MIT) into the aurelio monorepo as a self-hosted Telegram shell bot for the fleet gateway (CT-217). - telminal/ package: config (env + state file, refuses empty admins), core orchestrator (event handlers, router, watchers, interactive mode, file up/download, xterm.js image render), process (pexpect PTY + streaming + inline control buttons), telegram (Telethon wrapper, swappable for tests), cli (entry point reading TELEGRAM_* env), utils, values. - aurelio hardening vs upstream: no first-run random token auth (explicit admin allowlist required); cd sandbox validated against working root; secrets from env mirroring the CT-217 gateway .env. - 35 real tests (pty capture/control-char, router, watchers, perms, sandbox, fake-client orchestration) -- all green. - deployment: systemd/telminal.service, .env.example entries, README, AGENTS.md. Verified: pytest 35 passed; CLI refuses start with missing env / no admins. --- .env.example | 13 + .gitignore | 3 + AGENTS.md | 1 + telminal/.gitignore | 11 + telminal/AGENTS.md | 40 ++ telminal/LICENSE | 20 + telminal/README.md | 118 +++++ telminal/pyproject.toml | 28 ++ telminal/requirements-dev.txt | 2 + telminal/requirements.txt | 2 + telminal/systemd/telminal.service | 27 ++ telminal/telminal/__init__.py | 29 ++ telminal/telminal/__main__.py | 4 + telminal/telminal/cli.py | 73 ++++ telminal/telminal/config.py | 114 +++++ telminal/telminal/core.py | 692 ++++++++++++++++++++++++++++++ telminal/telminal/process.py | 194 +++++++++ telminal/telminal/telegram.py | 136 ++++++ telminal/telminal/utils.py | 65 +++ telminal/telminal/values.py | 53 +++ telminal/tests/conftest.py | 53 +++ telminal/tests/helpers.py | 157 +++++++ telminal/tests/test_config.py | 108 +++++ telminal/tests/test_core.py | 243 +++++++++++ telminal/tests/test_process.py | 192 +++++++++ 25 files changed, 2378 insertions(+) create mode 100644 telminal/.gitignore create mode 100644 telminal/AGENTS.md create mode 100644 telminal/LICENSE create mode 100644 telminal/README.md create mode 100644 telminal/pyproject.toml create mode 100644 telminal/requirements-dev.txt create mode 100644 telminal/requirements.txt create mode 100644 telminal/systemd/telminal.service create mode 100644 telminal/telminal/__init__.py create mode 100644 telminal/telminal/__main__.py create mode 100644 telminal/telminal/cli.py create mode 100644 telminal/telminal/config.py create mode 100644 telminal/telminal/core.py create mode 100644 telminal/telminal/process.py create mode 100644 telminal/telminal/telegram.py create mode 100644 telminal/telminal/utils.py create mode 100644 telminal/telminal/values.py create mode 100644 telminal/tests/conftest.py create mode 100644 telminal/tests/helpers.py create mode 100644 telminal/tests/test_config.py create mode 100644 telminal/tests/test_core.py create mode 100644 telminal/tests/test_process.py diff --git a/.env.example b/.env.example index 7a20ae6e..fedddaa7 100644 --- a/.env.example +++ b/.env.example @@ -47,3 +47,16 @@ ANNAS_ARCHIVE_KEY= # Icicle FPGA models directory (absolute path) ICICLE_MODELS_DIR=/home/fcunha/portugalfuturista/replica/infra/lab-gateway/icicle-models + +# ─── TELMINAL (Telegram shell bot) ────────────────────────────────────────── +# A sandboxed shell-on-the-host bot for the fleet gateway (CT-217). +# Create a bot + API credentials at: +# https://core.telegram.org/api/obtaining_api_id and https://t.me/botfather +# Deploy via telminal/systemd/telminal.service (secrets sourced from this file). +TELEGRAM_API_ID= +TELEGRAM_API_HASH= +TELEGRAM_TOKEN= +# Comma-separated numeric Telegram user ids allowed to run commands (allowlist). +TELEGRAM_ADMINS= +# Optional sandbox root for `cd` / file transfers (absolute path). +TELEGRAM_WORKING_DIR=/home/hermes/telminal-workspace diff --git a/.gitignore b/.gitignore index 4db88a95..7a0dd672 100644 --- a/.gitignore +++ b/.gitignore @@ -53,3 +53,6 @@ mycelium-rd/target/ # JetBrains IDE local platform artifacts extensions/aurelio-jetbrains/.intellijPlatform/ + +# telminal local test venv +.venv-telminal/ diff --git a/AGENTS.md b/AGENTS.md index 79761a68..1d04c909 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -17,6 +17,7 @@ Centralized brain of the Portugal Futurista agent fleet + monorepo of tools. | `dashboards/` | Media explorer dashboard | TypeScript | | `infra/` | CDP bridge, lab gateway | TypeScript | | `scripts/` | Operational scripts (sync, conscience upgrade, Kimi import) | Python | +| `telminal/` | Self-hosted Terminal-in-Telegram bot (sandboxed shell on the fleet gateway) | Python, Telethon | ## Build commands by subproject diff --git a/telminal/.gitignore b/telminal/.gitignore new file mode 100644 index 00000000..67fa3bde --- /dev/null +++ b/telminal/.gitignore @@ -0,0 +1,11 @@ +MIT + +# Python +__pycache__/ +*.py[cod] +*.egg-info/ +.venv/ + +# Test caches +.pytest_cache/ +.coverage diff --git a/telminal/AGENTS.md b/telminal/AGENTS.md new file mode 100644 index 00000000..16404bcc --- /dev/null +++ b/telminal/AGENTS.md @@ -0,0 +1,40 @@ +# telminal — Agent Guide + +Self-hosted **Terminal-in-Telegram** bot for the Portugal Futurista fleet. It is +a hardened re-implementation of [fristhon/telminal](https://github.com/fristhon/telminal) +(MIT), giving an admin allowlist a sandboxed shell on the gateway host (CT-217) +from Telegram. + +## What it does +Runs `/bin/bash -c ` in a pty (pexpect), streams output to Telegram with +inline control buttons (Info / Enter / Interactive / Terminate / HTML), supports +interactive mode (`^c` = Ctrl-C), `!get`/upload file transfer, and `!watch` +watcher tasks. Optional xterm.js screenshot output via pyppeteer. + +## Hardening vs upstream +- **No first-run random token auth** — admins come from an explicit allowlist + (`TELEGRAM_ADMINS` env, persisted in `config.json`). Refuses to start empty. +- **cwd sandbox**: `cd` is validated against the configured working root; escapes + are rejected. +- Secrets from env (`TELEGRAM_API_ID/HASH/TOKEN`), mirroring the Hermes gateway + `.env` on CT-217. + +## Layout +`cli.py` (entry) → `config.py` (Config) → `core.py` (orchestrator) → +`process.py` (PTY) + `telegram.py` (Telethon wrapper) + `utils.py`/`values.py`. + +## Run +```bash +. .venv/bin/activate +TELEGRAM_API_ID=.. TELEGRAM_API_HASH=.. TELEGRAM_TOKEN=.. TELEGRAM_ADMINS=123 python -m telminal +``` +Deploy via `systemd/telminal.service` on CT-217 (secrets from +`/home/hermes/.hermes/.env`, working dir `/home/hermes/telminal-workspace`). + +## Tests +```bash +pytest -q # pty, router, watchers, perms, sandbox, fake-client orchestration +``` +The orchestrator talks to Telegram only through `telegram.Telegram`, so tests +use a fake client — no network needed. The live Telegram end-to-end run is +verified on CT-217 by the operator. diff --git a/telminal/LICENSE b/telminal/LICENSE new file mode 100644 index 00000000..90ffd16b --- /dev/null +++ b/telminal/LICENSE @@ -0,0 +1,20 @@ +Copyright (c) 2021 fristhon +Copyright (c) 2026 Portugal Futurista + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/telminal/README.md b/telminal/README.md new file mode 100644 index 00000000..033480ab --- /dev/null +++ b/telminal/README.md @@ -0,0 +1,118 @@ +# telminal — A Terminal in Telegram (aurelio port) + +A self-hosted Telegram bot that gives you a **sandboxed shell on the host from +Telegram**. It is a clean re-implementation of +[fristhon/telminal](https://github.com/fristhon/telminal) (MIT), ported into the +Réplica Omnisciente / Portugal Futurista monorepo with aurelio conventions. + +> ⚠️ **Security**: telminal executes shell commands on the host it runs on. It +> must only ever be deployed on a host you control (the fleet's CT-217 gateway) +> with an explicit admin allowlist. It refuses to start without configured +> admins. + +## Features + +- Run any shell command via a PTY (`/bin/bash -c`) and stream the output back + to Telegram, with inline control buttons: + - 💡 **Info** — process status (PID / runtime / last update) + - ↩️ **Enter** — send a newline / Enter to the process + - **Interactive mode** — talk to a running process (each message = input); + `^c`, `^d`, `^z` map to Ctrl-C / Ctrl-D / Ctrl-Z + - 🛑 **Terminate** — kill the process + - 🌐 **HTML** — receive the full output as a rendered HTML file +- Optional **image output**: render the terminal with xterm.js and send a + screenshot (`/image_on`, `/image_off`). +- **File download** from server: `!get ` or the inline file picker. +- **File upload** to server: send a file to the chat and confirm save. +- **Watcher tasks**: `!watch 50s telminal.log` periodically pushes a file. +- **Multi-admin**: first admin can `!trust`/`!untrust` other users (by reply). +- **Sandboxed cwd**: `cd` is validated and cannot escape the configured working + directory root. + +## Architecture + +``` +telminal/ +├── pyproject.toml # packaging + console script `telminal` +├── requirements.txt +├── systemd/telminal.service +├── telminal/ +│ ├── __init__.py +│ ├── cli.py # entry point: reads env, builds Config, runs bot +│ ├── config.py # Config (env + state file), refuses empty admins +│ ├── core.py # Telminal orchestrator (event handlers, router) +│ ├── process.py # TProcess: pexpect PTY + streaming + buttons +│ ├── telegram.py # thin Telethon wrapper (swappable for tests) +│ ├── utils.py # HTML template + helpers +│ └── values.py # message templates + regexes +└── tests/ # real tests (pty, router, watchers, perms, fake tg) +``` + +The orchestrator depends on Telegram only through `Telegram` (telegram.py), so +the full routing / permission / watcher logic is exercised in tests with a fake +client — no network required. + +## Configuration + +All secrets come from the environment (matches the CT-217 gateway `.env`): + +| Variable | Required | Meaning | +|----------|----------|---------| +| `TELEGRAM_API_ID` | yes | Telegram API id | +| `TELEGRAM_API_HASH` | yes | Telegram API hash | +| `TELEGRAM_TOKEN` | yes | Bot token from @BotFather | +| `TELEGRAM_ADMINS` | yes | Comma-separated numeric user ids (allowlist) | +| `TELEGRAM_WORKING_DIR` | no | Sandbox root for `cd` / file ops (default: cwd) | + +The admin list and working directory also persist to `config.json` next to the +package so `!trust`/`cd` survive restarts. + +Create your bot + API credentials at: +- https://core.telegram.org/api/obtaining_api_id +- https://t.me/botfather + +Set BotFather commands: +``` +image_off - Just text output +image_on - Text and image output +tasks - List of active tasks +interacive_mode - Interact with last process +normal_mode - Process creation mode +``` + +## Running locally + +```bash +python3.11 -m venv .venv +. .venv/bin/activate +pip install -e ".[dev]" +export TELEGRAM_API_ID=... TELEGRAM_API_HASH=... TELEGRAM_TOKEN=... +export TELEGRAM_ADMINS=123456789 +python -m telminal +``` + +Install the optional image-rendering backend with `pip install ".[image]"`. + +## Deploying on CT-217 (fleet gateway) + +```bash +# on the gateway host (CT 217): +sudo cp systemd/telminal.service /etc/systemd/system/ +sudo systemctl daemon-reload +sudo systemctl enable --now telminal +journalctl -u telminal -f +``` + +The unit sources secrets from `/home/hermes/.hermes/.env` and pins the working +directory to `/home/hermes/telminal-workspace`. + +## Tests + +```bash +. .venv/bin/activate +pytest -q +``` + +Covers: PTY spawn + output capture, control-char input, watcher regex parsing, +command routing, permission gating, cwd sandboxing, and the full message +orchestration driven by a fake Telegram client (no network). diff --git a/telminal/pyproject.toml b/telminal/pyproject.toml new file mode 100644 index 00000000..0e0788aa --- /dev/null +++ b/telminal/pyproject.toml @@ -0,0 +1,28 @@ +[build-system] +requires = ["setuptools>=61.0"] +build-backend = "setuptools.build_meta" + +[project] +name = "telminal" +version = "1.0.0" +description = "A Terminal in Telegram — run a sandboxed shell on the aurelio fleet from Telegram." +readme = "README.md" +license = { text = "MIT" } +requires-python = ">=3.9" +authors = [{ name = "fristhon (original), Portugal Futurista (aurelio port)" }] +keywords = ["telegram", "terminal", "pty", "bot", "aurelio"] +dependencies = [ + "telethon>=1.24", + "pexpect>=4.8", +] + +[project.optional-dependencies] +image = ["pyppeteer>=0.2.6"] +dev = ["pytest>=7.0"] + +[project.scripts] +telminal = "telminal.cli:main" + +[tool.setuptools.packages.find] +where = ["."] +include = ["telminal*"] diff --git a/telminal/requirements-dev.txt b/telminal/requirements-dev.txt new file mode 100644 index 00000000..17540ea1 --- /dev/null +++ b/telminal/requirements-dev.txt @@ -0,0 +1,2 @@ +pytest>=7.0 +pyppeteer>=0.2.6 diff --git a/telminal/requirements.txt b/telminal/requirements.txt new file mode 100644 index 00000000..a25ffef2 --- /dev/null +++ b/telminal/requirements.txt @@ -0,0 +1,2 @@ +telethon>=1.24 +pexpect>=4.8 diff --git a/telminal/systemd/telminal.service b/telminal/systemd/telminal.service new file mode 100644 index 00000000..e90fb99c --- /dev/null +++ b/telminal/systemd/telminal.service @@ -0,0 +1,27 @@ +[Unit] +Description=Aurelio telminal — Terminal in Telegram (sandboxed shell bot) +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +User=hermes +Group=hermes +# Run inside the dedicated venv created by the deployment playbook. +WorkingDirectory=/opt/telminal +ExecStart=/opt/telminal/.venv/bin/python -m telminal +# Secrets come from the environment (same .env shape as the Hermes gateway). +EnvironmentFile=/home/hermes/.hermes/.env +# Hardening: the bot grants a shell, so constrain it. +NoNewPrivileges=true +ProtectSystem=strict +ProtectHome=false +ReadWritePaths=/home/hermes/telminal-workspace /opt/telminal +PrivateTmp=true +Restart=on-failure +RestartSec=5 +# The working directory is pinned by TELEGRAM_WORKING_DIR below. +Environment=TELEGRAM_WORKING_DIR=/home/hermes/telminal-workspace + +[Install] +WantedBy=multi-user.target diff --git a/telminal/telminal/__init__.py b/telminal/telminal/__init__.py new file mode 100644 index 00000000..85918a2a --- /dev/null +++ b/telminal/telminal/__init__.py @@ -0,0 +1,29 @@ +"""telminal — A Terminal in Telegram, hardened for the aurelio/Portugal Futurista fleet. + +This is a clean re-implementation of fristhon/telminal (MIT) ported into the +aurelio monorepo. It keeps the original feature set (PTY shell, streaming output +with inline control buttons, interactive mode, file up/download, watcher tasks, +optional xterm.js image rendering) while adopting aurelio conventions: + +* secrets come from the environment (``TELEGRAM_API_ID``/``API_HASH``/``TOKEN``), + matching the shape used by the CT-217 gateway ``.env``; +* the admin allowlist and working directory are sourced from a JSON state file + (``config.json``) that lives beside the package, never committed; +* no first-run random token auth (avoid surprise privilege grants) — admins are + configured explicitly; an empty admin list refuses to start. +""" + +from __future__ import annotations + +from pathlib import Path + +# PACKAGE_PATH must be defined before submodules import it, to avoid a +# circular import when config.py does ``from . import PACKAGE_PATH``. +PACKAGE_PATH = Path(__file__).parent + +from .core import Telminal # noqa: E402 +from .process import TProcess # noqa: E402 + +__version__ = "1.0.0" + +__all__ = ["Telminal", "TProcess", "PACKAGE_PATH", "__version__"] diff --git a/telminal/telminal/__main__.py b/telminal/telminal/__main__.py new file mode 100644 index 00000000..9ae637f1 --- /dev/null +++ b/telminal/telminal/__main__.py @@ -0,0 +1,4 @@ +from .cli import main + +if __name__ == "__main__": + main() diff --git a/telminal/telminal/cli.py b/telminal/telminal/cli.py new file mode 100644 index 00000000..73d93b1f --- /dev/null +++ b/telminal/telminal/cli.py @@ -0,0 +1,73 @@ +"""Command-line entry point for telminal. + +Reads configuration from the environment and runs the bot. In the aurelio +deployment this is launched by the ``telminal.service`` systemd unit on CT-217, +which supplies ``TELEGRAM_API_ID`` / ``TELEGRAM_API_HASH`` / ``TELEGRAM_TOKEN`` +plus the admin list (``TELEGRAM_ADMINS``) and an optional working directory +(``TELEGRAM_WORKING_DIR``). +""" + +from __future__ import annotations + +import asyncio +import os +import sys + +from .config import Config, ConfigError + + +def _parse_admins(raw: str | None) -> list[int]: + if not raw: + return [] + out: list[int] = [] + for part in raw.split(","): + part = part.strip() + if not part: + continue + out.append(int(part)) + return out + + +def build_config() -> Config: + cfg = Config.from_env(allow_empty_admins=True) + admins = _parse_admins(os.environ.get("TELEGRAM_ADMINS")) + if admins: + # union of env admins and persisted admins, de-duplicated, env-first + merged = admins + [a for a in cfg.admins if a not in admins] + cfg.admins = merged + wd = os.environ.get("TELEGRAM_WORKING_DIR") + if wd: + cfg.working_dir = wd + return cfg + + +def main() -> None: + try: + cfg = build_config() + except ConfigError as exc: + print(f"[telminal] configuration error: {exc}", file=sys.stderr) + raise SystemExit(2) + + if not cfg.admins: + print( + "[telminal] refusing to start: no admins configured. " + "Set TELEGRAM_ADMINS or add admins to the state file.", + file=sys.stderr, + ) + raise SystemExit(2) + + from .core import Telminal + + print( + f"[telminal] starting with {len(cfg.admins)} admin(s) " + f"(working dir: {cfg.working_dir or os.getcwd()})" + ) + telminal = Telminal(cfg) + try: + asyncio.run(telminal.start()) + except KeyboardInterrupt: + print("\n[telminal] interrupted, shutting down.") + + +if __name__ == "__main__": + main() diff --git a/telminal/telminal/config.py b/telminal/telminal/config.py new file mode 100644 index 00000000..843c7b49 --- /dev/null +++ b/telminal/telminal/config.py @@ -0,0 +1,114 @@ +"""Configuration loading for telminal. + +In the aurelio deployment the Telegram API credentials are injected from the +environment (the same ``.env`` that feeds the CT-217 gateway). The mutable +runtime state — the admin allowlist and the current working directory — is kept +in a JSON file (``config.json``) next to the package so the bot can remember +``!trust``/``!untrust`` decisions and ``cd`` changes across restarts. + +An empty admin list is a hard refusal: telminal must never start unauthenticated +because it grants a full shell on the host. +""" + +from __future__ import annotations + +import json +import os +from dataclasses import dataclass, field +from pathlib import Path +from typing import Optional + +from . import PACKAGE_PATH + + +class ConfigError(Exception): + """Raised when telminal cannot be configured safely.""" + + +def _require_env(name: str) -> str: + value = os.environ.get(name) + if not value: + raise ConfigError( + f"Missing required environment variable {name!r}. " + f"Add it to the service .env (TELEGRAM_API_ID / " + f"TELEGRAM_API_HASH / TELEGRAM_TOKEN)." + ) + return value + + +@dataclass +class Config: + api_id: int + api_hash: str + token: str + admins: list[int] = field(default_factory=list) + working_dir: Optional[str] = None + session_name: str = "telminal" + state_file: Path = PACKAGE_PATH / "config.json" + + @classmethod + def from_env(cls, *, allow_empty_admins: bool = False) -> "Config": + """Build a Config from environment variables + on-disk state. + + Reads TELEGRAM_API_ID / TELEGRAM_API_HASH / TELEGRAM_TOKEN from the + environment, then merges the persisted admin list / working dir from the + state file (the state file wins for mutable fields so ``!trust`` and + ``cd`` survive restarts). + """ + api_id_raw = _require_env("TELEGRAM_API_ID") + try: + api_id = int(api_id_raw) + except ValueError as exc: + raise ConfigError( + f"TELEGRAM_API_ID must be an integer, got {api_id_raw!r}" + ) from exc + + api_hash = _require_env("TELEGRAM_API_HASH") + token = _require_env("TELEGRAM_TOKEN") + + cfg = cls( + api_id=api_id, + api_hash=api_hash, + token=token, + ) + cfg._load_state() + if not cfg.admins and not allow_empty_admins: + raise ConfigError( + "No admins configured. Set TELEGRAM_ADMINS (comma separated " + "numeric user ids), or add at least one admin to the state file. " + "telminal refuses to start without an explicit admin allowlist." + ) + return cfg + + def _load_state(self) -> None: + try: + with open(self.state_file, encoding="utf-8") as handle: + data = json.load(handle) + except FileNotFoundError: + return + except json.JSONDecodeError as exc: + raise ConfigError( + f"State file {self.state_file} is not valid JSON: {exc}" + ) from exc + + admins = data.get("admins") + if isinstance(admins, list): + self.admins = [int(a) for a in admins] + wd = data.get("working_dir") + if isinstance(wd, str) and wd: + self.working_dir = wd + + def save_state(self) -> None: + """Persist the mutable runtime state to disk.""" + data = { + "api_id": self.api_id, + "api_hash": self.api_hash, + "token": self.token, + "admins": self.admins, + "working_dir": self.working_dir, + } + self.state_file.parent.mkdir(parents=True, exist_ok=True) + tmp = self.state_file.with_suffix(".json.tmp") + with open(tmp, "w", encoding="utf-8") as handle: + json.dump(data, handle, indent=2) + tmp.replace(self.state_file) diff --git a/telminal/telminal/core.py b/telminal/telminal/core.py new file mode 100644 index 00000000..7d416269 --- /dev/null +++ b/telminal/telminal/core.py @@ -0,0 +1,692 @@ +"""The telminal orchestrator. + +``Telminal`` owns the Telegram event handlers, the set of running processes, the +watcher tasks, and the interactive-mode pointer. It is a faithful port of +fristhon/telminal with aurelio hardening: + +* no first-run random token auth — admins come from the Config allowlist; +* the working directory is pinned to ``config.working_dir`` and ``cd`` is + validated before applied (no escaping to arbitrary paths); +* runtime state (admins + cwd) is persisted through ``Config.save_state``. +""" + +from __future__ import annotations + +import asyncio +import atexit +import os +import re +import signal +import tempfile +from functools import partial +from pathlib import Path +from subprocess import PIPE +from time import time + +from . import utils +from .config import Config +from .process import TProcess +from .telegram import Telegram +from .values import ( + ACTIVE_TASKS_MSG, + BROWSER_ERROR_MSG, + CONTROL_PREFIX, + EMPTY_TASKS_MSG, + GET_PATTERN, + PROCESS_INFO_MSG, + WATCHER_PATTERN, +) + +temp_dir = tempfile.TemporaryDirectory(prefix="telminal-") +TEMP_PATH = Path(temp_dir.name) + + +class Telminal: + all_processes: dict[int, TProcess] = {} + all_progress_callback: dict[int, float] = {} + + PROCESS_CLEANER_DELAY = 100 + PROCESS_OUTPUT_LIFE_TIME = 60 + + def __init__(self, config: Config) -> None: + self._config = config + self.interactive_process: TProcess | None = None + self._api_id = config.api_id + self._api_hash = config.api_hash + self._token = config.token + self.session_name = config.session_name + self.admins: list[int] = list(config.admins) + self._render = True + self._watch_tasks: dict[int, str] = {} + self._last_process: TProcess | None = None + if config.working_dir: + self._working_dir = config.working_dir + else: + self._working_dir = os.getcwd() + os.chdir(self._working_dir) + self.browser = None + + # -- class-level background jobs --------------------------------------- + + @classmethod + async def process_cleaner(cls): + """Check and clean dead processes periodically.""" + while True: + for pid, process in cls.all_processes.copy().items(): + if ( + not process.is_running + and int(time() - process.done_time) + > cls.PROCESS_OUTPUT_LIFE_TIME + ): + del cls.all_processes[pid] + utils.silent_file_remover(TEMP_PATH / f"{pid}.html") + utils.silent_file_remover(TEMP_PATH / f"{pid}.png") + await asyncio.sleep(cls.PROCESS_CLEANER_DELAY) + + # -- permission decorator --------------------------------------------- + # These are staticmethods: when used as a bare decorator (``@check_permission``) + # Python passes the wrapped function as the first argument, so ``self`` must + # not be injected. The Telminal instance arrives later, at call time, as the + # first positional argument of the returned ``inner`` coroutine. + + @staticmethod + def check_permission(func): + async def inner(self_ref, event): + if event.sender_id not in self_ref.admins: + return + await func(self_ref, event) + + return inner + + # -- process resolver for callback queries ----------------------------- + + @staticmethod + def find_process_by_event(func): + async def inner(self_ref, event): + pid = int(event.data.decode().split("&")[-1]) + process = Telminal.all_processes.get(pid) + if process is None: + await event.answer( + "this process not exist anymore", alert=True + ) + # clear button + await self_ref.bot.edit_message( + event.chat_id, message_id=event.message_id + ) + return + await func(self_ref, event, process) + + return inner + + # -- bootstrap --------------------------------------------------------- + + async def start(self) -> None: + """Run the telminal instance by calling this method.""" + self.bot = Telegram( + api_id=self._api_id, + api_hash=self._api_hash, + token=self._token, + session_name=self.session_name, + ) + from telethon import events + + handlers = { + self._all_messages_handler: events.NewMessage(incoming=True), + self._terminate_handler: events.CallbackQuery( + pattern=r"terminate&\d+" + ), + self._html_handler: events.CallbackQuery(pattern=r"html&\d+"), + self._interactive_handler: events.CallbackQuery( + pattern=r"interact&\d+" + ), + self._inline_query_handler: events.InlineQuery(), + self._cancell_download_handler: events.CallbackQuery( + pattern=r"removeme" + ), + self._confirm_download_handler: events.CallbackQuery( + pattern=r"savefile&.+" + ), + self._enter_handler: events.CallbackQuery(pattern=r"enter&\d+"), + self._cancell_task_handler: events.CallbackQuery( + pattern=r"cancell_task&\d+" + ), + self._info_handler: events.CallbackQuery(pattern=r"info&\d+"), + } + asyncio.ensure_future(Telminal.process_cleaner()) + await self.bot.start(handlers) + atexit.register(self._exit_jobs) + await self.bot.run_until_disconnected() + + # -- xterm.js rendering ------------------------------------------------ + + async def render_xtermjs(self, process: TProcess): + """Returns parsed output and screenshot of a process.""" + output = process.media_output + image = None + + if self.browser is None or self._render is False: + return output, image + + try: + page = await self.browser.newPage() + await page.goto(process.html) + path = TEMP_PATH / f"{process.pid}.png" + await page.screenshot({"path": path, "fullPage": True}) + + await page.evaluate("term.selectAll()") + output = await page.evaluate("term.getSelection()") + output = Telegram.media_strip(output) + image = path + await page.close() + except Exception: + # TODO log this + pass + + return output, image + + # -- process lifecycle ------------------------------------------------- + + def _new_process(self, command: str, request_id: int) -> TProcess: + process = TProcess(command, request_id) + process.set_temp_path(TEMP_PATH) + process.run() + return process + + @check_permission + @find_process_by_event + async def _info_handler(self, event, process): + await event.answer(str(process), alert=True) + + @check_permission + @find_process_by_event + async def _enter_handler(self, event, process): + process.push("^m") + await event.answer("Enter key pressed...") + + # -- watcher tasks ------------------------------------------------------ + + async def _new_watcher(self, chat_id, request_id, message): + match = re.match(WATCHER_PATTERN, message) + if not match: + await self.bot.send_message( + chat_id, + "Wrong pattern, send /tasks to see valid examples", + reply_to=request_id, + ) + return + + self._watch_tasks[request_id] = message.replace("!watch", "") + count, type_, file = match.groups() + second_mapping = {"s": 1, "m": 60, "h": 60 * 60} + while True: + if self._watch_tasks.get(request_id) is None: + break + await self.bot.send_file(chat_id, file, reply_to=request_id) + delay = int(count) * second_mapping[type_] + await asyncio.sleep(delay) + + async def _show_tasks(self, chat_id): + message, buttons = self._get_tasks_message() + await self.bot.send_message( + chat_id, message, buttons=buttons, parse_mode="html" + ) + + def _get_tasks_message(self): + from telethon import Button + + buttons = [ + [Button.inline(f"❌ {command}", f"cancell_task&{task_id}")] + for task_id, command in self._watch_tasks.items() + ] + message = ACTIVE_TASKS_MSG + if not buttons: + buttons = None + message = EMPTY_TASKS_MSG + + return message, buttons + + @check_permission + async def _cancell_task_handler(self, event): + task_id = int(event.data.decode().split("&")[-1]) + try: + del self._watch_tasks[task_id] + except KeyError: + await event.answer("Dead task! create new one") + message, buttons = self._get_tasks_message() + await self.bot.edit_message( + event.chat_id, + message_id=event.message_id, + message=message, + buttons=buttons, + parse_mode="html", + ) + + # -- file upload (server -> telegram) ---------------------------------- + + async def _upload_file(self, chat_id, path, request_id): + if not os.path.isfile(path): + await self.bot.send_message( + chat_id, f"`{path}` is not a file", reply_to=request_id + ) + return + + message = await self.bot.send_message( + chat_id, "Uploading started...", reply_to=request_id + ) + partial_callback = partial( + self._progress_callback, + chat_id=chat_id, + message_id=message.id, + title=f"Uploading `{path}`", + ) + await self.bot.send_file( + chat_id, + file=path, + reply_to=request_id, + progress_callback=partial_callback, + ) + + # -- botfather / slash commands ---------------------------------------- + + async def _run_bot_father_commands(self, chat_id, command): + if command.startswith("/image_on"): + self._render = True + elif command.startswith("/image_off"): + self._render = False + elif command.startswith("/tasks"): + await self._show_tasks(chat_id) + elif command.startswith("/interacive_mode"): + if getattr(self._last_process, "is_running", None) is True: + self.set_interactive_process(self._last_process) + else: + await self.bot.send_message( + chat_id, + "last process finished, you must select another process manually", + ) + elif command.startswith("/normal_mode"): + self.reset_interactive_process() + + # -- extra bang commands ----------------------------------------------- + + async def _run_extra_commands(self, event, command, message): + chat_id, request_id = event.chat_id, event.message.id + + if re.match(GET_PATTERN, message): + path = message.split(" ", 1)[1] + await self._upload_file(chat_id, path, request_id) + + elif command == "!setup_browser": + await self.setup_browser() + + elif command == "!watch": + asyncio.ensure_future( + self._new_watcher(chat_id, request_id, message) + ) + + elif ( + command.startswith(("!trust", "!untrust")) + and event.sender_id == self.admins[0] + ): + await self._trust_switcher(command, event, request_id) + + async def _trust_switcher(self, command, event, request_id): + reply_message = await event.message.get_reply_message() + reply_sender_id = reply_message.sender_id + reply_to = request_id + + if command == "!trust": + if reply_sender_id is None: + message = "I can't trust to an Anonymous admin!" + else: + message = "Repeat your command please, you are a trusted user now" + self.admins.append(reply_sender_id) + self._config.admins = list(self.admins) + reply_to = reply_message.id + + elif command == "!untrust" and reply_sender_id != self.admins[0]: + message = "Done, removed from trusted users" + try: + self.admins.remove(reply_sender_id) + self._config.admins = list(self.admins) + except ValueError: + pass + + await self.bot.send_message(event.chat_id, message, reply_to=reply_to) + + # -- file download (telegram -> server) -------------------------------- + + async def _send_download_buttons(self, event): + from telethon import Button + + message_id, file_name = event.id, event.file.name + buttons = [] + if os.path.exists(file_name): + message = f"`{file_name}` currentlly exists on this directory" + buttons.extend( + [ + [ + Button.inline( + "Save as new file", + data=f"savefile&new&{message_id}", + ) + ], + [ + Button.inline( + "Overwrite", data=f"savefile&true&{message_id}" + ) + ], + ] + ) + else: + message = "Do you want to save this file on sever?" + buttons.append( + [Button.inline("Yes", data=f"savefile&new&{message_id}")] + ) + + buttons.append([Button.inline("Cancell", data="removeme")]) + + await self.bot.send_message( + event.chat_id, + message, + reply_to=message_id, + buttons=buttons, + ) + + # -- change directory (validated) -------------------------------------- + + async def _change_directory(self, event): + path = event.message.message.split(" ", 1)[-1] + resolved = self._safe_path(path) + if resolved is None: + await self.bot.send_message( + event.chat_id, + "Permission denied: cannot move outside the working root.", + reply_to=event.message.id, + ) + return + try: + os.chdir(resolved) + except Exception as error: + await self.bot.send_message( + event.chat_id, + str(error), + reply_to=event.message.id, + ) + return + self._working_dir = str(resolved) + self._config.working_dir = self._working_dir + self._config.save_state() + + def _safe_path(self, path: str) -> str | None: + """Resolve ``path`` and reject escapes outside the configured root. + + The working root is the *original* configured directory, not the + current one, so users cannot ``cd`` their way out of the sandbox. + """ + root = Path(self._config.working_dir or os.getcwd()).resolve() + target = (root / path).resolve() if not os.path.isabs(path) else Path( + path + ).resolve() + try: + target.relative_to(root) + except ValueError: + return None + return str(target) + + # -- main dispatch ------------------------------------------------------ + + @check_permission + async def _all_messages_handler(self, event): + message = event.message.message + escape = message.startswith("\\") + chat_id, request_id = event.chat_id, event.message.id + + if event.file: + await self._send_download_buttons(event) + + elif message.startswith(("cd", "/", "!")) and not escape: + command = message.split(" ", 1)[0] + if message.startswith("/"): + await self._run_bot_father_commands(chat_id, command) + elif message.startswith("!"): + await self._run_extra_commands(event, command, message) + else: + await self._change_directory(event) + + elif self.interactive_process: + text = message[1:] if escape else message + self.interactive_process.push(text) + next_update_arrived = ( + int(time()) - self.interactive_process.last_update_time >= 2 + ) + if self.interactive_process is not None and next_update_arrived: + await self.response(self.interactive_process, chat_id) + else: + process = self._new_process(message, request_id) + Telminal.all_processes[process.pid] = process + asyncio.ensure_future( + self._run_in_background(process, chat_id) + ) + + @check_permission + @find_process_by_event + async def _terminate_handler(self, event, process): + process.terminate() + + @check_permission + @find_process_by_event + async def _html_handler(self, event, process): + await self.bot.send_file( + event.chat_id, + process.html, + reply_to=process.response_id, + ) + + # -- interactive mode --------------------------------------------------- + + def set_interactive_process(self, process): + if isinstance(self.interactive_process, TProcess): + self.interactive_process.is_interactive_process = False + self.interactive_process = process + process.is_interactive_process = True + return f"You are talking to PID {process.pid}" + + def reset_interactive_process(self): + if isinstance(self.interactive_process, TProcess): + self.interactive_process.is_interactive_process = False + self.interactive_process = None + return "Normal mode activated" + + @check_permission + @find_process_by_event + async def _interactive_handler(self, event, process): + if self.interactive_process is process: + answer = self.reset_interactive_process() + else: + answer = self.set_interactive_process(process) + + await event.answer(answer, alert=True) + await self.response(process, event.chat_id) + + # -- inline query (file picker) ---------------------------------------- + + @check_permission + async def _inline_query_handler(self, event): + command = "ls -la" if not event.text else f"ls -la | grep {event.text}" + proc = await asyncio.subprocess.create_subprocess_shell( + command, stdin=PIPE, stdout=PIPE, stderr=PIPE + ) + files = (await proc.stdout.read()).decode().split("\n") + + builder = event.builder + results = [] + + file_name_pattern = re.compile(r"(.+)\s+\d{2}:*\d{2}\s+\d+\s") + for file in files[: Telegram.INLINE_RESULT_LIMIT]: + # `-` means this is a file and not a directory + # `r` added for `-?????????` situations (permission) + if file.startswith("-r"): + file_name = file_name_pattern.findall(file[::-1])[0][::-1] + results.append( + builder.article( + text=f"!get {file_name}", + title=file_name, + description=file, + ) + ) + await event.answer(results=results, cache_time=0) + + # -- progress reporting ------------------------------------------------- + + async def _progress_callback( + self, current, total, *, chat_id: int, message_id: int, title: str + ): + percent_str = f"{current / total:.2%}" + percent_int = int(percent_str.split(".")[0]) + upload_finished = percent_int == 100 + + emoji = "🟩" + if upload_finished: + emoji = "☑️" + title = "Finished Successfully" + + emoji_count = 0 if percent_int <= 10 else int(percent_str[0]) + + text = f"""\ + {title} + {emoji_count * emoji} {percent_str} + """ + # showing upload state to user each 5 second + if ( + int(time() - self.all_progress_callback.get(message_id, 0)) > 5 + or upload_finished + ): + await self.bot.edit_message( + chat_id, message_id=message_id, message=text + ) + self.all_progress_callback[message_id] = time() + + # -- download confirmation --------------------------------------------- + + @check_permission + async def _confirm_download_handler(self, event): + _, overwrite, message_id = event.data.decode().split("&") + message = await self.bot.get_message( + event.chat_id, int(message_id) + ) + partial_callback = partial( + self._progress_callback, + chat_id=event.chat_id, + message_id=event.message_id, + title="Downloading...", + ) + + file = message.file.name if overwrite == "true" else None + await self.bot.download_media( + message, progress_callback=partial_callback, file=file + ) + + @check_permission + async def _cancell_download_handler(self, event): + await event.delete() + + # -- response loop ------------------------------------------------------ + + async def response(self, process: TProcess, chat_id: int): + """Respond to user and update process result message periodically.""" + media_output = process.media_output + # update buttons must be once per response + new_buttons = process.update_buttons() + + if not process.has_new_state(new_buttons, media_output): + return + + output, image = await self.render_xtermjs(process) + if not process.has_new_state(new_buttons, output): + return + + if process.is_partial: + if not any(output.split("\n")): + # sometimes when process finished with a kill signal the final + # output is an empty line and in telegram we can't send an + # empty message, but buttons update must apply to message + output = None + await self.bot.edit_message( + chat_id, + message=output, + message_id=process.response_id, + buttons=process.buttons, + file=image, + ) + else: + sent = await self.bot.send_message( + chat_id, + output, + reply_to=process.request_id, + buttons=process.buttons, + file=image, + ) + process.response_id = sent.id + self._last_process = process + process.is_partial = True + process.last_message = output + + async def _run_in_background(self, process: TProcess, chat_id: int): + while process.is_running: + partial_update_time = (process.run_time + 1) % 4 == 0 + try: + if partial_update_time or process.response_id is None: + # first time fast response needed but next times for a + # partial update must have passed at least 1 second + response_delay = ( + 0.5 if process.response_id is None else 1.1 + ) + await asyncio.sleep(response_delay) + await self.response(process, chat_id) + except Exception: + # TODO should I react to this? + pass + finally: + await asyncio.sleep(0.1) + + try: + # maybe process will be finished before next update + await self.response(process, chat_id) + finally: + if self.interactive_process is process: + self.reset_interactive_process() + + # -- browser (xterm.js image output) ----------------------------------- + + async def setup_browser(self): + from pyppeteer import launch + + try: + self.browser = await launch( + options={ + "args": [ + "--disable-gpu", + "--disable-dev-shm-usage", + "--disable-setuid-sandbox", + "--no-sandbox", + ], + "headless": True, + "autoClose": False, + } + ) + except Exception as e: + await self.bot.send_message( + self.admins[0], + BROWSER_ERROR_MSG.format(error=e), + parse_mode="html", + ) + self.browser = None + + # -- shutdown ----------------------------------------------------------- + + def _exit_jobs(self): + self._config.admins = list(self.admins) + if self._working_dir: + self._config.working_dir = self._working_dir + self._config.save_state() + temp_dir.cleanup() diff --git a/telminal/telminal/process.py b/telminal/telminal/process.py new file mode 100644 index 00000000..dce212af --- /dev/null +++ b/telminal/telminal/process.py @@ -0,0 +1,194 @@ +"""PTY-backed process wrapper for telminal. + +A ``TProcess`` spawns ``/bin/bash -c `` inside a pty (via pexpect) and +streams its output into an in-memory buffer. The orchestrator (``core.Telminal``) +reads ``.full_output``/``.media_output`` to form Telegram messages and decides +when to push an update. + +The class is deliberately free of any Telegram dependency so it can be unit +tested against a real PTY without a network. +""" + +from __future__ import annotations + +import asyncio +from io import StringIO +from pathlib import Path +from time import time + +import pexpect +from pexpect.exceptions import EOF, TIMEOUT + +from . import utils +from .telegram import Telegram + +# How many bytes to read from the pty per tick. +READ_CHUNK = 1000 +# Polling interval (seconds) for the background stream reader. +STREAM_POLL = 0.1 + + +class TProcess: + def __init__(self, command: str, request_id: int) -> None: + self.command = command + self.request_id = request_id + self._buffer = StringIO() + self.is_running: bool | None = None + self.is_partial: bool | None = None + self.start_time: float | None = None + self.run_time = 0 + self._new_data: bool = False + self.response_id: int | None = None + self._last_message = "" + self.last_update_time: float = 0.0 + self.buttons = None + self.is_interactive_process = False + self._temp_path: Path | None = None + + # -- lifecycle ---------------------------------------------------------- + + def run(self, stream: bool = True) -> None: + self._process = pexpect.spawn("/bin/bash", ["-c", self.command], timeout=None) + self.pid = self._process.pid + self.is_running = True + self.start_time = time() + if stream is True: + asyncio.create_task(self.stream()) + + def done(self) -> None: + self.is_running = False + self.done_time = time() + + def terminate(self) -> None: + """Terminate the underlying process group (best effort).""" + try: + self._process.terminate(force=True) + except Exception: + # pexpect may already have reaped the child. + pass + self.done() + + async def stream(self) -> None: + while True: + try: + line = self._process.read_nonblocking(size=READ_CHUNK, timeout=0) + self._buffer.write(line.decode("utf-8", errors="replace")) + self._new_data = True + except EOF: + self.done() + break + except TIMEOUT: + pass + finally: + if self.start_time is not None: + self.run_time = int(time() - self.start_time) + await asyncio.sleep(STREAM_POLL) + + # -- output accessors --------------------------------------------------- + + @property + def full_output(self) -> str: + self._buffer.seek(0) + return self._buffer.read() + + @property + def media_output(self) -> str: + return Telegram.media_strip(self.full_output) + + def _render_html(self, temp_path: Path) -> Path: + file = temp_path / f"{self.pid}.html" + with open(file, "w", encoding="utf-8") as html: + html.write( + utils.HTML_TEMPLATE.format( + title=f"{self.pid} -> {self.command}", + data=self.full_output.replace("`", r"\`"), + ) + ) + return file + + @property + def html(self) -> str: + # The orchestrator injects the temp path via ``set_temp_path``. + return str(self._render_html(self._temp_path)) + + def set_temp_path(self, path: Path) -> None: + self._temp_path = path + + # -- interactive input -------------------------------------------------- + + @property + def last_message(self) -> str: + return self._last_message + + @last_message.setter + def last_message(self, last_message: str) -> None: + self._last_message = last_message + self._new_data = False + self.last_update_time = time() + + def push(self, command: str) -> None: + """Send input to the running process. + + A two-character message starting with ``^`` (e.g. ``^c``) is sent as a + control character (Ctrl-C). Otherwise the text is sent verbatim with + trailing newlines translated into Enter key presses. + """ + if command.startswith("^") and len(command) == 2: + self._process.sendcontrol(command[-1]) + else: + for index, word in enumerate(command.split("\n")): + if index != 0: + # for each extra line send an Enter first + self._process.sendcontrol("m") + if word: + self._process.send(word) + else: + self._process.sendcontrol("m") + + # -- inline buttons ----------------------------------------------------- + + def update_buttons(self): + from telethon import Button + + if self.is_running: + interact_switch_text = "Interactive mode" + if self.is_interactive_process: + interact_switch_text = "Exit interactive mode" + + buttons = [ + [Button.inline("💡 Info 💡", data=f"info&{self.pid}")], + [Button.inline("↩️ Enter ↩️", data=f"enter&{self.pid}")], + [ + Button.inline( + interact_switch_text, data=f"interact&{self.pid}" + ) + ], + [Button.inline("🛑 Terminate 🛑", data=f"terminate&{self.pid}")], + [Button.inline("🌐 HTML 🌐", data=f"html&{self.pid}")], + ] + else: + buttons = [ + [Button.inline("💡 Info 💡", data=f"info&{self.pid}")], + [Button.inline("🌐 HTML 🌐", data=f"html&{self.pid}")], + ] + + if self.buttons != buttons: + self.buttons = buttons + return True + return False + + def has_new_state(self, new_buttons, new_output) -> bool: + return new_buttons or (new_output and self.last_message != new_output) + + def __str__(self) -> str: + status = "🔄 Running" if self.is_running else "✔️ Done" + start_at = utils.timestamp_to_readable(self.start_time) + last_update = utils.timestamp_to_readable(self.last_update_time) + run_time = utils.seconds_to_readable(self.run_time) + return ProcessInfoMsg.format( + self.pid, status, start_at, last_update, run_time + ) + + +# Imported lazily to avoid a circular import with values. +from .values import PROCESS_INFO_MSG as ProcessInfoMsg # noqa: E402 diff --git a/telminal/telminal/telegram.py b/telminal/telminal/telegram.py new file mode 100644 index 00000000..ab8393d2 --- /dev/null +++ b/telminal/telminal/telegram.py @@ -0,0 +1,136 @@ +"""Thin async wrapper around Telethon's TelegramClient. + +We keep every Telegram-specific call behind this class so the rest of the +orchestrator never touches the Telethon API directly. That makes it possible to +drive ``Telminal`` with a fake client in tests (see ``tests/conftest.py``). +""" + +from __future__ import annotations + +import os + +from . import utils + + +class Telegram: + # upload limitation at this moment > 2GB + UPLOAD_LIMIT = 2147483648 + MEDIA_CAPTION_LIMIT = 1024 + INLINE_RESULT_LIMIT = 50 + + def __init__( + self, + api_id: int, + api_hash: str, + token: str, + session_name: str, + ) -> None: + self._api_id = api_id + self._api_hash = api_hash + self._token = token + self.session_name = session_name + self._remove_old_sessions() + from telethon import TelegramClient + + self._client = TelegramClient(session_name, api_id, api_hash) + self._client.parse_mode = None + + def _remove_old_sessions(self) -> None: + # an easy way to handle session db lock error + utils.silent_file_remover(f"{self.session_name}.session") + utils.silent_file_remover(f"{self.session_name}.session-journal") + + async def start(self, handlers: dict) -> None: + await self._client.start(bot_token=self._token) + for handler, type_ in handlers.items(): + self._client.add_event_handler(handler, type_) + + async def run_until_disconnected(self) -> None: + await self._client.run_until_disconnected() + + async def send_message( + self, + chat_id: int, + message: str, + *, + reply_to: int | None = None, + buttons=None, + file=None, + parse_mode=None, + ): + kwargs = {} + if parse_mode is not None: + kwargs["parse_mode"] = parse_mode + return await self._client.send_message( + chat_id, + message, + link_preview=False, + reply_to=reply_to, + buttons=buttons, + file=file, + **kwargs, + ) + + async def edit_message( + self, + chat_id, + *, + message_id: int, + message=None, + buttons=None, + file=None, + parse_mode=None, + ): + kwargs = {} + if parse_mode is not None: + kwargs["parse_mode"] = parse_mode + return await self._client.edit_message( + chat_id, + message=message_id, + text=message, + link_preview=False, + buttons=buttons, + file=file, + **kwargs, + ) + + async def send_file( + self, chat_id: int, file: str, reply_to=None, progress_callback=None + ): + error = None + if not os.path.isfile(file): + error = f"{file} is not a file" + elif os.path.getsize(file) > Telegram.UPLOAD_LIMIT: + error = "Sorry I can't send file bigger than 2G" + elif os.path.getsize(file) == 0: + error = "I can't send an empty file" + + if error is not None: + await self._client.send_message(chat_id, error, reply_to=reply_to) + return None + + try: + return await self._client.send_file( + chat_id, + file=file, + force_document=True, + reply_to=reply_to, + progress_callback=progress_callback, + ) + except Exception: + # TODO log + return None + + async def get_message(self, chat_id, message_id: int): + return await self._client.get_messages(chat_id, ids=message_id) + + async def download_media(self, message, *, progress_callback, file): + await self._client.download_media( + message, progress_callback=progress_callback, file=file + ) + + @staticmethod + def media_strip(message: str) -> str: + if len(message) >= Telegram.MEDIA_CAPTION_LIMIT: + message = message[len(message) - Telegram.MEDIA_CAPTION_LIMIT :] + return message.strip() diff --git a/telminal/telminal/utils.py b/telminal/telminal/utils.py new file mode 100644 index 00000000..1f2cb966 --- /dev/null +++ b/telminal/telminal/utils.py @@ -0,0 +1,65 @@ +"""Small shared helpers for telminal.""" + +from __future__ import annotations + +import contextlib +import os +from datetime import datetime, timedelta + +# xterm.js HTML template used to render a process' output as a screenshot. +# Inspired by https://github.com/cs01/pyxtermjs/blob/master/pyxtermjs/index.html +HTML_TEMPLATE = """ + + + {title} + + + +
+ + + + + + + + + +""" + + +def silent_file_remover(file): + with contextlib.suppress(FileNotFoundError): + os.remove(file) + + +def timestamp_to_readable(timestamp: float) -> str: + dt = datetime.fromtimestamp(timestamp) + return f"{dt.hour}:{dt.minute}:{dt.second}" + + +def seconds_to_readable(seconds: int) -> str: + return str(timedelta(seconds=seconds)) diff --git a/telminal/telminal/values.py b/telminal/telminal/values.py new file mode 100644 index 00000000..00f3063c --- /dev/null +++ b/telminal/telminal/values.py @@ -0,0 +1,53 @@ +"""User-facing message templates and constants for telminal.""" + +from __future__ import annotations + +BROWSER_ERROR_MSG = """\ +Browser setup error : {error} + +Seems you don't have installed the requirements packages, if you need image output run the following commands on your server. +Meanwhile You can use text version of Telminal, type any command! + +`sudo apt-get install chromium-chromedriver` +`sudo apt install -y gconf-service libasound2 libatk1.0-0 libc6 libcairo2 libcups2 libdbus-1-3 libexpat1 libfontconfig1 libgcc1 libgconf-2-4 libgdk-pixbuf2.0-0 libglib2.0-0 libgtk-3-0 libnspr4 libpango-1.0-0 libpangocairo-1.0-0 libstdc++6 libx11-6 libx11-xcb1 libxcb1 libxcomposite1 libxcursor1 libxdamage1 libxext6 libxfixes3 libxi6 libxrandr2 libxrender1 libxss1 libxtst6 ca-certificates fonts-liberation libappindicator1 libnss3 lsb-release xdg-utils wget` + +Finally, send `!setup_browser` on this chat again. +""" + +ACTIVE_TASKS_MSG = ( + "This is list of your active tasks\nyou can cancel each one by tapping" +) + +EMPTY_TASKS_MSG = """\ +Tasks list is empty +Create new watcher same as below examples: + +1️⃣ get `telminal.log` every 50 seconds +👉 !watch 50s telminal.log + +2️⃣ get `temp.png` every 5 minutes +👉 !watch 5m home/temp.png + +3️⃣ get `sql.dump` every 24 hours +👉 !watch 24h /backups/sql.dump + +Be respectfull to telegram API limitations please: +https://core.telegram.org/bots/faq#my-bot-is-hitting-limits-how-do-i-avoid-this +""" + +PROCESS_INFO_MSG = """\ +PID : {} +Status : {} + +Start time : {} +Last update : {} + +Run time: {} +""" + +# A control char message that starts with ^ and is exactly two characters +# (e.g. ^c) maps to a CTRL signal. The second character is the control letter. +CONTROL_PREFIX = "^" + +WATCHER_PATTERN = r"!watch\s(\d+)([s,m,h])\s(.+)" +GET_PATTERN = r"!get\s.+" diff --git a/telminal/tests/conftest.py b/telminal/tests/conftest.py new file mode 100644 index 00000000..5555a089 --- /dev/null +++ b/telminal/tests/conftest.py @@ -0,0 +1,53 @@ +"""Pytest configuration for telminal. + +Imports the test harness (fake client, event builders, fake telethon fixture) +and defines the ``telminal`` / ``fake_client`` / ``admin_id`` fixtures used +across the suite. +""" + +from __future__ import annotations + +import pytest + +import telminal.core as core_mod +from telminal.config import Config +from helpers import FakeClient, _fake_telethon # noqa: F401 + + +@pytest.fixture +def fake_client(): + return FakeClient() + + +@pytest.fixture +def telminal(fake_client, tmp_path, monkeypatch): + """A Telminal instance wired to a FakeClient, cwd pinned to tmp_path.""" + cfg = Config( + api_id=123, + api_hash="abc", + token="tok", + admins=[99], + working_dir=str(tmp_path), + state_file=tmp_path / "config.json", + ) + monkeypatch.chdir(tmp_path) + t = core_mod.Telminal(cfg) + t.bot = fake_client + return t + + +@pytest.fixture(autouse=True) +def _reset_class_state(): + """Clear shared Telminal class state so tests don't leak processes/tasks.""" + from telminal import core as core_mod + + core_mod.Telminal.all_processes.clear() + core_mod.Telminal.all_progress_callback.clear() + yield + core_mod.Telminal.all_processes.clear() + core_mod.Telminal.all_progress_callback.clear() + + +@pytest.fixture +def admin_id(): + return 99 diff --git a/telminal/tests/helpers.py b/telminal/tests/helpers.py new file mode 100644 index 00000000..2c475e8f --- /dev/null +++ b/telminal/tests/helpers.py @@ -0,0 +1,157 @@ +"""Test harness: fake Telethon client + event builders. + +``FakeClient`` is a drop-in stand-in for ``telethon.TelegramClient`` covering the +methods telminal calls; it records every call so tests can assert on them. The +builder functions produce fake ``events.NewMessage`` / ``events.CallbackQuery`` +objects carrying the fields handlers read (``sender_id``, ``chat_id``, +``message.message``, ``message.id``, ``file``, ``data``, ``text``, ``answer``, +``delete``). + +``core.Telminal`` imports ``telethon`` lazily (only inside ``start`` and a few +handlers), so the orchestrator can be instantiated and its handlers invoked +without telethon installed — except handlers that reference ``telethon.Button`` +/ ``telethon.events`` at call time. The conftest fixture injects a fake +``telethon`` module onto ``sys.modules`` for exactly that purpose. +""" + +from __future__ import annotations + +import asyncio +import sys +from types import SimpleNamespace +from typing import Any + +import pytest + + +class _Button: + @staticmethod + def inline(text, data=None): + return {"text": text, "data": data} + + +class _Event: + def __init__(self, pattern=None): + self.pattern = pattern + + +class _Events: + NewMessage = _Event + CallbackQuery = _Event + InlineQuery = _Event + + +@pytest.fixture(autouse=True) +def _fake_telethon(monkeypatch): + """Inject a fake `telethon` module so lazy imports resolve in tests.""" + fake = SimpleNamespace(Button=_Button, events=_Events) + monkeypatch.setitem(sys.modules, "telethon", fake) + yield fake + monkeypatch.delitem(sys.modules, "telethon", raising=False) + + +class FakeMessage: + def __init__(self, message="", id=0, file=None): + self.message = message + self.id = id + self.file = file + + +class FakeClient: + def __init__(self, *args, **kwargs): + self.handlers = {} + self.calls: list[dict[str, Any]] = [] + self.parse_mode = None + self._next_id = 1 + + async def start(self, bot_token=None): + self.bot_token = bot_token + + def add_event_handler(self, handler, type_): + self.handlers.setdefault(id(type_), []).append(handler) + + async def run_until_disconnected(self): + return None + + async def send_message(self, chat_id, message, **kwargs): + msg_id = self._next_id + self._next_id += 1 + self.calls.append( + {"op": "send_message", "chat_id": chat_id, "message": message, + "kwargs": kwargs} + ) + return SimpleNamespace(id=msg_id) + + async def edit_message(self, chat_id, **kwargs): + self.calls.append( + {"op": "edit_message", "chat_id": chat_id, "kwargs": kwargs} + ) + return SimpleNamespace(id=kwargs.get("message_id", 0)) + + async def send_file(self, chat_id, file=None, **kwargs): + self.calls.append( + {"op": "send_file", "chat_id": chat_id, "file": file, "kwargs": kwargs} + ) + return SimpleNamespace(id=self._next_id) + + async def get_messages(self, chat_id, ids=None): + return FakeMessage( + message="x", id=ids if isinstance(ids, int) else 0, + file=SimpleNamespace(name="remote.bin"), + ) + + async def download_media(self, message, **kwargs): + self.calls.append({"op": "download_media"}) + return b"" + + +class FakeEvent: + def __init__(self, *, sender_id, chat_id=10, message="", message_id=1, + file=None, data=None, text=None): + self.sender_id = sender_id + self.chat_id = chat_id + self.message = SimpleNamespace(message=message, id=message_id, file=file) + self.file = file + self.message_id = message_id + self.data = data + self.text = text + self.answered = None + self.deleted = False + + async def answer(self, text, alert=False): + self.answered = text + + async def delete(self): + self.deleted = True + + +def make_new_message(*, sender_id, chat_id=10, message="", message_id=1, + file=None): + return FakeEvent(sender_id=sender_id, chat_id=chat_id, message=message, + message_id=message_id, file=file) + + +def make_callback(*, sender_id, chat_id=10, data: str, message_id=1): + # Real Telethon callback-query events carry ``data`` as ``bytes``. + data_bytes = data.encode() if isinstance(data, str) else data + return FakeEvent( + sender_id=sender_id, chat_id=chat_id, data=data_bytes, message_id=message_id + ) + + +def run(coro): + """Run a coroutine to completion on a persistent event loop. + + A single loop is reused across calls within a test session so that tasks + scheduled by ``TProcess.run`` / ``asyncio.ensure_future`` (the streaming and + background-response tasks) keep running between successive ``run()`` calls, + mirroring how Telethon drives handlers on one long-lived loop. + """ + global _loop + if _loop is None or _loop.is_closed(): + _loop = asyncio.new_event_loop() + asyncio.set_event_loop(_loop) + return _loop.run_until_complete(coro) + + +_loop = None # module-level persistent loop diff --git a/telminal/tests/test_config.py b/telminal/tests/test_config.py new file mode 100644 index 00000000..475a2758 --- /dev/null +++ b/telminal/tests/test_config.py @@ -0,0 +1,108 @@ +"""Tests for configuration, utils, and routing/watcher regexes.""" + +from __future__ import annotations + +import os + +import pytest + +from telminal.config import Config, ConfigError +from telminal.values import GET_PATTERN, WATCHER_PATTERN +from telminal import utils + + +def test_config_requires_env(monkeypatch): + for var in ("TELEGRAM_API_ID", "TELEGRAM_API_HASH", "TELEGRAM_TOKEN"): + monkeypatch.delenv(var, raising=False) + with pytest.raises(ConfigError): + Config.from_env() + + +def test_config_reads_env_admins(monkeypatch, tmp_path): + monkeypatch.setenv("TELEGRAM_API_ID", "123") + monkeypatch.setenv("TELEGRAM_API_HASH", "h") + monkeypatch.setenv("TELEGRAM_TOKEN", "t") + monkeypatch.setenv("TELEGRAM_ADMINS", "111, 222") + # from_env intentionally ignores TELEGRAM_ADMINS (that's the CLI's job); + # admins come only from the persisted state file, which is empty here. + cfg = Config.from_env(allow_empty_admins=True) + assert cfg.api_id == 123 + assert cfg.admins == [] + + +def test_cli_build_config_merges_env_admins(monkeypatch, tmp_path): + monkeypatch.setenv("TELEGRAM_API_ID", "123") + monkeypatch.setenv("TELEGRAM_API_HASH", "h") + monkeypatch.setenv("TELEGRAM_TOKEN", "t") + monkeypatch.setenv("TELEGRAM_ADMINS", "111, 222") + from telminal.cli import build_config + + cfg = build_config() + assert cfg.admins == [111, 222] + + +def test_config_refuses_empty_admins(monkeypatch, tmp_path): + monkeypatch.setenv("TELEGRAM_API_ID", "123") + monkeypatch.setenv("TELEGRAM_API_HASH", "h") + monkeypatch.setenv("TELEGRAM_TOKEN", "t") + monkeypatch.delenv("TELEGRAM_ADMINS", raising=False) + with pytest.raises(ConfigError): + Config.from_env() + + +def test_config_bad_api_id(monkeypatch): + monkeypatch.setenv("TELEGRAM_API_ID", "not-an-int") + monkeypatch.setenv("TELEGRAM_API_HASH", "h") + monkeypatch.setenv("TELEGRAM_TOKEN", "t") + with pytest.raises(ConfigError): + Config.from_env(allow_empty_admins=True) + + +def test_config_persists_and_reloads(tmp_path): + state = tmp_path / "config.json" + cfg = Config( + api_id=1, api_hash="h", token="t", admins=[5], + working_dir=str(tmp_path), state_file=state, + ) + cfg.save_state() + assert state.exists() + reloaded = Config(api_id=1, api_hash="h", token="t", state_file=state) + reloaded._load_state() + assert reloaded.admins == [5] + assert reloaded.working_dir == str(tmp_path) + + +def test_watcher_pattern(): + m = __import__("re").match(WATCHER_PATTERN, "!watch 50s telminal.log") + assert m is not None + assert m.groups() == ("50", "s", "telminal.log") + m2 = __import__("re").match(WATCHER_PATTERN, "!watch 24h /backups/sql.dump") + assert m2.groups() == ("24", "h", "/backups/sql.dump") + assert __import__("re").match(WATCHER_PATTERN, "!watch bogus") is None + + +def test_get_pattern(): + assert __import__("re").match(GET_PATTERN, "!get /etc/passwd") + assert not __import__("re").match(GET_PATTERN, "echo hi") + + +def test_html_template_render(): + rendered = utils.HTML_TEMPLATE.format(title="t", data="hello world") + assert "hello world" in rendered + assert " Any: + """Wrap a plain value so it can be `await`ed (mirrors an async method).""" + return value + + +from telminal import core as core_mod + + +# -------------------------------------------------------------------------- +# permission gating +# -------------------------------------------------------------------------- + +def test_non_admin_message_is_ignored(telminal, fake_client, admin_id): + ev = make_new_message(sender_id=999, message="echo pwned") + run(telminal._all_messages_handler(ev)) + # no process spawned, no message sent + assert fake_client.calls == [] + assert core_mod.Telminal.all_processes == {} + + +def test_admin_command_spawns_process(telminal, fake_client, admin_id): + ev = make_new_message(sender_id=admin_id, message="echo hi") + run(telminal._all_messages_handler(ev)) + # a background task was scheduled; give the event loop a tick + run(asyncio.sleep(0.8)) + assert len(core_mod.Telminal.all_processes) == 1 + # response message posted by the background response loop + assert any(c["op"] == "send_message" for c in fake_client.calls) + + +# -------------------------------------------------------------------------- +# command routing +# -------------------------------------------------------------------------- + +def test_cd_routing_changes_directory(telminal, fake_client, admin_id, tmp_path): + sub = tmp_path / "sub" + sub.mkdir() + ev = make_new_message(sender_id=admin_id, message=f"cd {sub}") + run(telminal._all_messages_handler(ev)) + import os + + assert os.getcwd() == str(sub.resolve()) + + +def test_cd_cannot_escape_sandbox(telminal, fake_client, admin_id, tmp_path): + # try to escape outside the working root + ev = make_new_message( + sender_id=admin_id, + message=f"cd {tmp_path.parent}", + ) + run(telminal._all_messages_handler(ev)) + import os + + # cwd stays inside the configured root + assert os.getcwd().startswith(str(tmp_path)) + + +def test_get_command_invokes_upload(telminal, fake_client, admin_id, tmp_path): + f = tmp_path / "data.txt" + f.write_text("contents") + ev = make_new_message(sender_id=admin_id, message=f"!get {f}") + run(telminal._all_messages_handler(ev)) + assert any( + c["op"] == "send_file" and c["file"] == str(f) for c in fake_client.calls + ) + + +def test_image_on_off_toggles_render(telminal, admin_id): + ev = make_new_message(sender_id=admin_id, message="/image_off") + run(telminal._run_bot_father_commands(ev.chat_id, "/image_off")) + assert telminal._render is False + run(telminal._run_bot_father_commands(ev.chat_id, "/image_on")) + assert telminal._render is True + + +# -------------------------------------------------------------------------- +# watchers +# -------------------------------------------------------------------------- + +def test_watcher_regex_and_task_registration(telminal, fake_client, admin_id, tmp_path): + f = tmp_path / "w.log" + f.write_text("log") + ev = make_new_message(sender_id=admin_id, message=f"!watch 50s {f}") + run(telminal._all_messages_handler(ev)) + # task registered under the request id + assert len(telminal._watch_tasks) == 1 + # a "Uploading started..." message was not sent yet; the watcher loops. + # Cancel it so the background loop terminates. + req_id = ev.message.id + del telminal._watch_tasks[req_id] + + +def test_watcher_bad_pattern_replies_error(telminal, fake_client, admin_id): + ev = make_new_message(sender_id=admin_id, message="!watch bogus") + run(telminal._all_messages_handler(ev)) + assert any( + c["op"] == "send_message" + and "Wrong pattern" in c["message"] + for c in fake_client.calls + ) + + +def test_tasks_list_empty_and_populated(telminal, fake_client, admin_id): + run(telminal._show_tasks(10)) + assert any( + c["op"] == "send_message" and "empty" in c["message"].lower() + for c in fake_client.calls + ) + fake_client.calls.clear() + telminal._watch_tasks[1] = " 50s x" + run(telminal._show_tasks(10)) + assert any( + c["op"] == "send_message" and "active tasks" in c["message"].lower() + for c in fake_client.calls + ) + + +# -------------------------------------------------------------------------- +# interactive mode + callback handlers +# -------------------------------------------------------------------------- + +def test_interactive_mode_toggle(telminal, fake_client, admin_id, tmp_path): + # spawn a process + ev = make_new_message(sender_id=admin_id, message="sleep 2") + run(telminal._all_messages_handler(ev)) + run(asyncio.sleep(0.2)) + pid = next(iter(core_mod.Telminal.all_processes)) + cb = make_callback(sender_id=admin_id, data=f"interact&{pid}") + run(telminal._interactive_handler(cb)) + assert telminal.interactive_process is not None + assert telminal.interactive_process.pid == pid + # toggle back + cb2 = make_callback(sender_id=admin_id, data=f"interact&{pid}") + run(telminal._interactive_handler(cb2)) + assert telminal.interactive_process is None + + +def test_terminate_handler_kills_process(telminal, fake_client, admin_id): + ev = make_new_message(sender_id=admin_id, message="sleep 5") + run(telminal._all_messages_handler(ev)) + run(asyncio.sleep(0.2)) + pid = next(iter(core_mod.Telminal.all_processes)) + cb = make_callback(sender_id=admin_id, data=f"terminate&{pid}") + run(telminal._terminate_handler(cb)) + proc = core_mod.Telminal.all_processes[pid] + # give the stream task time to observe EOF + run(asyncio.sleep(0.3)) + assert proc.is_running is False + + +def test_info_handler_answers(telminal, fake_client, admin_id): + ev = make_new_message(sender_id=admin_id, message="sleep 5") + run(telminal._all_messages_handler(ev)) + run(asyncio.sleep(0.2)) + pid = next(iter(core_mod.Telminal.all_processes)) + cb = make_callback(sender_id=admin_id, data=f"info&{pid}") + run(telminal._info_handler(cb)) + assert cb.answered is not None + assert "PID" in cb.answered + + +def test_unknown_process_callback_is_safe(telminal, admin_id): + cb = make_callback(sender_id=admin_id, data="info&999999") + run(telminal._info_handler(cb)) + assert cb.answered == "this process not exist anymore" + + +# -------------------------------------------------------------------------- +# trust / untrust +# -------------------------------------------------------------------------- + +def test_trust_adds_admin(telminal, fake_client, admin_id): + # first admin is admins[0]; a reply from them trusts a new user + reply = FakeEvent(sender_id=admin_id, chat_id=10, message="!trust", + message_id=2) + reply.message = SimpleNamespace( + message="!trust", id=2, + get_reply_message=(lambda: _await(SimpleNamespace(sender_id=555, id=2))), + ) + run(telminal._trust_switcher("!trust", reply, 2)) + assert 555 in telminal.admins + assert 555 in telminal._config.admins + + +def test_untrust_removes_admin(telminal, fake_client, admin_id): + telminal.admins.append(556) + telminal._config.admins = list(telminal.admins) + reply = FakeEvent(sender_id=admin_id, chat_id=10, message="!untrust", + message_id=3) + reply.message = SimpleNamespace( + message="!untrust", id=3, + get_reply_message=(lambda: _await(SimpleNamespace(sender_id=556))), + ) + run(telminal._trust_switcher("!untrust", reply, 3)) + assert 556 not in telminal.admins + + +# -------------------------------------------------------------------------- +# response loop (edit vs send) +# -------------------------------------------------------------------------- + +def test_response_sends_then_edits(telminal, fake_client, admin_id, tmp_path): + proc = __import__("telminal.process", fromlist=["TProcess"]).TProcess( + "echo hello", 42 + ) + proc.set_temp_path(tmp_path) + proc.run(stream=False) + # first state -> triggers a send_message + proc._buffer.write("hello") + run(telminal.response(proc, 10)) + assert proc.is_partial is True + assert any(c["op"] == "send_message" for c in fake_client.calls) + # changed state -> triggers an edit_message + fake_client.calls.clear() + proc._buffer.write(" world") + run(telminal.response(proc, 10)) + assert any(c["op"] == "edit_message" for c in fake_client.calls) + # identical state -> no-op (avoids spamming Telegram) + fake_client.calls.clear() + run(telminal.response(proc, 10)) + assert fake_client.calls == [] diff --git a/telminal/tests/test_process.py b/telminal/tests/test_process.py new file mode 100644 index 00000000..ba41fd1c --- /dev/null +++ b/telminal/tests/test_process.py @@ -0,0 +1,192 @@ +"""Tests for the PTY-backed process engine (real pexpect).""" + +from __future__ import annotations + +import asyncio +import time +from pathlib import Path + +import pytest + +from telminal.process import TProcess + + +def _drive(coro_factory): + """Run a coroutine factory under a fresh event loop and return its result. + + The factory is responsible for calling ``proc.run(stream=True)`` (which uses + ``asyncio.create_task`` and therefore requires a running loop). + """ + + async def _wrapper(): + return await coro_factory() + + try: + loop = asyncio.get_event_loop() + if loop.is_running(): + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + except RuntimeError: + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + return loop.run_until_complete(_wrapper()) + + +def test_process_captures_output(): + def factory(): + proc = TProcess("echo hello-world", 1) + proc.set_temp_path(Path("/tmp")) + proc.run(stream=True) + + async def wait(): + deadline = time.time() + 5 + while proc.is_running and time.time() < deadline: + await asyncio.sleep(0.05) + await asyncio.sleep(0.15) + return proc.full_output + + return wait() + + assert "hello-world" in _drive(factory) + + +def test_process_exit_code_and_done(): + def factory(): + proc = TProcess("exit 7", 2) + proc.set_temp_path(Path("/tmp")) + proc.run(stream=True) + + async def wait(): + deadline = time.time() + 5 + while proc.is_running and time.time() < deadline: + await asyncio.sleep(0.05) + await asyncio.sleep(0.1) + return proc.is_running + + return wait() + + assert _drive(factory) is False + + +def test_process_streaming_updates_buffer(): + def factory(): + proc = TProcess("printf 'line1\\nline2\\n'", 3) + proc.set_temp_path(Path("/tmp")) + proc.run(stream=True) + + async def wait(): + deadline = time.time() + 5 + while proc.is_running and time.time() < deadline: + await asyncio.sleep(0.05) + await asyncio.sleep(0.1) + return proc.full_output + + return wait() + + out = _drive(factory) + assert "line1" in out + assert "line2" in out + + +def test_push_sends_plain_text_and_newline(): + def factory(): + proc = TProcess("cat", 4) + proc.set_temp_path(Path("/tmp")) + proc.run(stream=True) + proc.push("ping\n") + + async def wait(): + # let the streaming task copy pty bytes into the buffer + await asyncio.sleep(0.4) + sent = "ping" in proc.full_output + proc.terminate() + return sent + + return wait() + + assert _drive(factory) is True + + +def test_push_control_char(): + # `push("^c")` must route to sendcontrol (transmit Ctrl-C) without error. + # We assert the control char is delivered to the pty (echoed as ^C) rather + # than relying on the child's SIGINT handler, which depends on real TTY job + # control not present in this non-controlling pty. + def factory(): + proc = TProcess("cat", 5) + proc.set_temp_path(Path("/tmp")) + proc.run(stream=True) + proc.push("^c") + + async def wait(): + await asyncio.sleep(0.3) + delivered = "^C" in proc.full_output or "\x03" in proc.full_output + proc.terminate() + return delivered + + return wait() + + assert _drive(factory) is True + +def test_media_strip_truncates_long_output(): + proc = TProcess("true", 6) + proc.set_temp_path(Path("/tmp")) + proc._buffer.write("x" * 2000) + stripped = proc.media_output + assert len(stripped) == 1024 + assert stripped.endswith("x") + + +def test_html_render_contains_output(): + import tempfile + + d = Path(tempfile.mkdtemp()) + + def factory(): + proc = TProcess("echo abc", 7) + proc.set_temp_path(d) + proc.run(stream=True) + + async def wait(): + try: + proc._process.expect("abc", timeout=5) + except Exception: + pass + deadline = time.time() + 5 + while proc.is_running and time.time() < deadline: + await asyncio.sleep(0.05) + await asyncio.sleep(0.1) + return proc.html + + return wait() + + html_path = _drive(factory) + content = Path(html_path).read_text() + assert "abc" in content + assert "