← Files NGS Analysis WorkbenchARCHIVED FILE

mcp/ngs_workbench_daemon/remote_operations.py

6.75 KB · Sep 30, 2026 · 23:20 UTC

↓ Download file

"""Execute a fixed, bounded workflow operation on an approved Linux SSH host."""

from __future__ import annotations

import hashlib
import json
import os
import signal
import stat
import subprocess
import sys
from pathlib import Path
from typing import Any, BinaryIO

LOG_TAIL_BYTES = 12_000


def read_tail(handle: BinaryIO, limit: int = LOG_TAIL_BYTES) -> str:
    """Shared local/SSH reader, included in this standalone remote program."""
    handle.seek(0, os.SEEK_END)
    handle.seek(max(0, handle.tell() - limit))
    return handle.read(limit).decode("utf-8", errors="replace")


def _path(root: Path, value: str) -> Path:
    candidate = Path(value)
    if (
        not candidate.is_absolute()
        or (candidate != root and root not in candidate.parents)
        or candidate.resolve(strict=False) != candidate
    ):
        raise ValueError("approved remote path escapes its run or contains a symlink")
    return candidate


def _digest(candidate: Path) -> tuple[int, str]:
    value = hashlib.sha256()
    with candidate.open("rb") as handle:
        for chunk in iter(lambda: handle.read(1024 * 1024), b""):
            value.update(chunk)
    return candidate.stat().st_size, "sha256:" + value.hexdigest()


def execute(request: dict[str, Any], controller: str) -> Any:
    root = Path(request["run_dir"])
    if not root.is_absolute() or root.resolve(strict=False) != root:
        raise ValueError("approved remote run path is unsafe")

    operation = request["operation"]
    if operation == "file_select":
        for pattern in request["patterns"]:
            query = Path(pattern)
            directory = _path(root, str(root / query.parent))
            candidates = [_path(root, str(path)) for path in directory.glob(query.name)]
            files = [path for path in candidates if path.is_file()]
            if files:
                return str(max(files, key=lambda path: (path.stat().st_mtime_ns, path.name)))
        return None

    if operation in {"file_tail", "file_stream"}:
        path = _path(root, str(root / request["relative_path"]))
        with os.fdopen(os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK), "rb") as handle:
            metadata = os.fstat(handle.fileno())
            if not stat.S_ISREG(metadata.st_mode):
                raise ValueError("remote evidence is not a regular file")
            if operation == "file_tail":
                return read_tail(handle)
            remaining = metadata.st_size
            while remaining:
                chunk = handle.read(min(remaining, 64 * 1024))
                if not chunk:
                    raise ValueError("remote evidence was truncated during observation")
                sys.stdout.buffer.write(chunk)
                remaining -= len(chunk)
            return None

    if operation == "prepare":
        if root.exists():
            raise ValueError("approved remote run directory already exists")
        root.mkdir(parents=True, mode=0o700)
        for value in request["destinations"]:
            _path(root, value).parent.mkdir(parents=True, exist_ok=True)
        for item in request["generated"]:
            destination = _path(root, item["destination"])
            destination.parent.mkdir(parents=True, exist_ok=True)
            destination.write_text(item["content"], encoding="utf-8")
        return {"ok": True}

    if operation == "verify":
        for item in request["files"]:
            candidate = _path(root, item["destination"])
            if candidate.is_symlink() or _digest(candidate) != (item["bytes"], item["sha256"]):
                raise ValueError("remote staged file differs from approved bytes")
        return {"ok": True}

    if operation == "launch":
        receipt = _path(root, str(root / "workflow" / "controller.json"))
        if receipt.is_file():
            identity = json.loads(receipt.read_text())
            execute({**request, "operation": "status", "identity": identity}, controller)
            return identity
        log = _path(root, str(root / request.get("log_path", "logs/nextflow.log")))
        log.parent.mkdir(parents=True, exist_ok=True)
        terminal = _path(root, str(root / "workflow" / "controller.exit"))
        payload = {
            "argv": request["argv"],
            "run_dir": str(root),
            "log": str(log),
            "exit": str(terminal),
        }
        process = subprocess.Popen(
            [sys.executable, "-c", controller, json.dumps(payload)],
            stdin=subprocess.DEVNULL,
            stdout=subprocess.DEVNULL,
            stderr=subprocess.DEVNULL,
            start_new_session=True,
        )
        fields = Path(f"/proc/{process.pid}/stat").read_text().split()
        result = {"pid": process.pid, "start_time": fields[21], "run_dir": str(root)}
        receipt.write_text(json.dumps(result))
        return result

    if operation == "identity":
        result = json.loads(_path(root, str(root / "workflow" / "controller.json")).read_text())
        if result["pid"] != request.get("pid", result["pid"]) or result["run_dir"] != str(root):
            raise ValueError("remote controller does not match its durable run identity")
        return result

    if operation in {"status", "stop"}:
        identity = request["identity"]
        receipt = json.loads(_path(root, str(root / "workflow" / "controller.json")).read_text())
        if identity != receipt:
            raise ValueError("remote controller identity does not match its durable receipt")
        terminal = _path(root, str(root / "workflow" / "controller.exit"))
        if terminal.is_file() and not terminal.is_symlink():
            return {"running": False, "return_code": int(terminal.read_text().strip())}

        process = Path(f"/proc/{identity['pid']}")
        try:
            fields = (process / "stat").read_text().split()
            command = (process / "cmdline").read_bytes()
            running = (
                fields[2] != "Z"
                and fields[21] == identity["start_time"]
                and json.loads(command.split(b"\0")[-2])["run_dir"] == str(root)
            )
        except (OSError, IndexError, KeyError, TypeError, ValueError):
            running = False
        if operation == "stop" and running:
            os.killpg(identity["pid"], signal.SIGTERM)
        elif operation == "stop":
            raise ValueError("refusing to signal an unverified remote controller")
        return {"running": running, "return_code": None}

    raise ValueError("unsupported remote workflow operation")


if __name__ == "__main__" and len(sys.argv) > 1:
    request = json.loads(sys.argv[1])
    controller = Path(__file__).with_name("remote_controller.py").read_text(encoding="utf-8")
    result = execute(request, controller)
    if request["operation"] != "file_stream":
        sys.stdout.write(json.dumps(result, separators=(",", ":")) + "\n")

SHA-256: a4ba212180447040884007cb4debb38ebae3327102ea67251df6be82a285aad9