diff --git a/dashboard/plugin_api.py b/dashboard/plugin_api.py index 9964f50..ee932ec 100644 --- a/dashboard/plugin_api.py +++ b/dashboard/plugin_api.py @@ -1,8 +1,10 @@ import asyncio import subprocess import re +import os import shutil import logging +import time import threading from typing import Optional from fastapi import APIRouter @@ -11,8 +13,9 @@ from fastapi.responses import JSONResponse logger = logging.getLogger(__name__) router = APIRouter() -_sshx_process = None _sshx_link = None +_sshx_pid = None +_sshx_log = "/tmp/sshx_link.log" LINK_RE = re.compile(r"https://sshx\.io/s/[A-Za-z0-9_-]+") @@ -31,92 +34,90 @@ def _ensure_sshx() -> bool: return False -def _start_reader(proc, future): - global _sshx_link - - def _read(): +def _kill_sshx(): + global _sshx_pid, _sshx_link + if _sshx_pid: try: - found = False - for line in proc.stdout: - line = line.strip() - if not found: - m = LINK_RE.search(line) - if m: - _sshx_link = m.group(0) - found = True - if not future.done(): - future.set_result(_sshx_link) - if not found and not future.done(): - future.set_exception(RuntimeError("sshx exited without link")) - except Exception as e: - if not future.done(): - future.set_exception(e) + os.kill(_sshx_pid, 15) + except ProcessLookupError: + pass + _sshx_pid = None + _sshx_link = None + if os.path.exists(_sshx_log): + os.remove(_sshx_log) - t = threading.Thread(target=_read, daemon=True) - t.start() - return t + +def _is_alive() -> bool: + global _sshx_pid + if not _sshx_pid: + return False + try: + os.kill(_sshx_pid, 0) + return True + except ProcessLookupError: + return False + + +def _start_sshx_background(): + global _sshx_pid, _sshx_link + _kill_sshx() + _sshx_link = None + + subprocess.Popen( + ["nohup", "sshx", "--quiet"], + stdout=open(_sshx_log, "w"), + stderr=subprocess.DEVNULL, + stdin=subprocess.DEVNULL, + start_new_session=True, + ) + + for _ in range(30): + time.sleep(1) + if os.path.exists(_sshx_log) and os.path.getsize(_sshx_log) > 0: + with open(_sshx_log) as f: + content = f.read().strip() + m = LINK_RE.search(content) + if m: + _sshx_link = m.group(0) + try: + with open("/tmp/sshx_pid") as pf: + _sshx_pid = int(pf.read().strip()) + except Exception: + pass + return True + return False @router.get("/start") async def start_sshx(): - global _sshx_process, _sshx_link + global _sshx_link, _sshx_pid - if _sshx_process and _sshx_process.poll() is None and _sshx_link: + if _is_alive() and _sshx_link: return {"status": "running", "link": _sshx_link} if not _ensure_sshx(): return JSONResponse(status_code=500, content={"status": "error", "message": "Failed to install sshx"}) - if _sshx_process and _sshx_process.poll() is None: - _sshx_process.terminate() - try: - _sshx_process.wait(timeout=5) - except subprocess.TimeoutExpired: - _sshx_process.kill() - - _sshx_process = subprocess.Popen( - ["sshx", "--quiet"], - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - text=True, - bufsize=1, - ) - - _sshx_link = None loop = asyncio.get_event_loop() - link_future = loop.create_future() - _start_reader(_sshx_process, link_future) + result = await loop.run_in_executor(None, _start_sshx_background) - try: - link = await asyncio.wait_for(link_future, timeout=30) - return {"status": "running", "link": link} - except asyncio.TimeoutError: - return JSONResponse(status_code=504, content={"status": "error", "message": "sshx did not return a link in time"}) - except Exception as e: - return JSONResponse(status_code=500, content={"status": "error", "message": str(e)}) + if result: + return {"status": "running", "link": _sshx_link} + return JSONResponse(status_code=504, content={"status": "error", "message": "sshx did not return a link in time"}) @router.get("/status") async def sshx_status(): - global _sshx_process, _sshx_link - if _sshx_process and _sshx_process.poll() is None and _sshx_link: + global _sshx_link + if _is_alive() and _sshx_link: return {"status": "running", "link": _sshx_link} return {"status": "stopped", "link": None} @router.post("/stop") async def stop_sshx(): - global _sshx_process, _sshx_link - if _sshx_process and _sshx_process.poll() is None: - _sshx_process.terminate() - try: - _sshx_process.wait(timeout=5) - except subprocess.TimeoutExpired: - _sshx_process.kill() - _sshx_process = None - _sshx_link = None - return {"status": "stopped"} - return {"status": "already_stopped"} + _kill_sshx() + return {"status": "stopped"} @router.post("/restart-dashboard") @@ -125,6 +126,7 @@ async def restart_dashboard(): subprocess.Popen( ["nohup", "bash", "-c", "sleep 1 && sudo systemctl restart hermes-dashboard"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, + start_new_session=True, ) return {"status": "restarting"} except Exception as e: