From ec48a60e3edf311e6ec695be3f4f734230f37be7 Mon Sep 17 00:00:00 2001 From: Sebastian Jeong Date: Sun, 16 Aug 2026 17:42:43 +0900 Subject: [PATCH] feat: supervise EXEC with guest liveness instead of a stopwatch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit TCPAGENT는 system()이 도는 동안 통째로 얼어 있어 응답도 진행 보고도 못 한다. 그런데 호스트는 소켓에 30초 고정 타임아웃을 걸고, 만료되면 연결 자체를 버렸다. 그래서 35초짜리 컴파일이 "느린 명령"이 아니라 "죽은 에이전트"로 취급됐다. QEMU는 게스트가 얼어 있어도 계속 돈다. info blockstats의 idle_time_ns로 "작업 중"과 "멈춤"을 구분한다. 실측으로 확인했다: 에이전트가 완전히 벙어리인 동안에도 rd_operations가 7초당 47000씩 증가하고 idle은 0.00s를 유지한다. - EXEC은 짧은 간격으로 깨어나 감시만 하고 소켓은 절대 안 버린다. - --idle-timeout(기본 60s)과 --hard-timeout(기본 900s). 후자는 디스크를 안 쓰는 CPU 바운드 멈춤용 백스톱이다. - 중단은 QEMU 모니터로 Ctrl+C를 주입하고 COMMAND.COM의 "Terminate batch file (Y/N/A)?" 프롬프트에 답한다. - Ctrl+C는 DOS break check에서만 먹는다. FreeDOS 기본값 BREAK=OFF에서 출력을 파일로 돌린 CPU 바운드 자식은 거기 도달 안 할 수 있다. 그래서 중단은 보장이 아니라 요청으로 다루고, 명령이 안 멈춰도 RESULT를 끝까지 수거해 스트림을 깨뜨리지 않는다. - ferro-vm abort 추가. 실행 중에도 응답해야 하므로 파이프 서버를 요청당 스레드로 바꿨다. - 5558 바인딩을 SO_EXCLUSIVEADDRUSE로. Windows의 SO_REUSEADDR는 다른 프로세스가 같은 포트를 잡아 조용히 반쯤 동작하게 만든다. 검증 (QEMU FreeDOS 실측): - 32.4초 명령 정상 완료 (이전에는 30초에 실패) - 실행 중 abort가 0.1초에 응답, exit=95로 종료, 부분 출력 1805B 수거, 연결 유지 - pause처럼 디스크를 안 쓰는 명령을 idle 15s로 검출해 중단 시리얼 시절에 있다가 TCP 전환에서 사라진 TODO 3건을 복구한다. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_012PQm6oAvWX4Lp3iSN5AHGT --- TODO.md | 28 ++++- src/ferrolang_vm/cli.py | 20 +++- src/ferrolang_vm/daemon.py | 228 +++++++++++++++++++++++++++++++------ tools/tcpagent/README.md | 25 +++- 4 files changed, 259 insertions(+), 42 deletions(-) diff --git a/TODO.md b/TODO.md index faeb710..d6fcdd1 100644 --- a/TODO.md +++ b/TODO.md @@ -1,5 +1,27 @@ # TODO -- [x] Add a non-reboot abort path for a hung DOS command: inject `Ctrl+C` through QEMU's monitor and wait for the serial agent to recover. -- [x] Apply a configurable timeout to `dos_exec` and invoke the non-reboot abort path on timeout. -- [x] Expose `dos_abort` for an immediate user-requested command interruption. +## Done + +- [x] Add a non-reboot abort path for a hung DOS command: inject `Ctrl+C` through + QEMU's monitor and wait for the agent to recover. +- [x] Apply a configurable timeout to `exec` and invoke the non-reboot abort path + on timeout. +- [x] Expose `abort` for an immediate user-requested command interruption. + +The three above were first built for the COM1 serial agent, lost in the rewrite +to the resident TCP agent, and rebuilt on the QEMU monitor in `ferro-vm exec`. +The TCP version supervises with guest disk liveness (`info blockstats`) rather +than a fixed stopwatch, so a slow compile is no longer mistaken for a hang. + +## Open + +- [ ] Consider `BREAK=ON` in `C:\FDCONFIG.SYS`. Ctrl+C only takes effect at a DOS + break check, and with the FreeDOS default of `BREAK=OFF` a compute-bound + child whose output is redirected to a file may never reach one, so `abort` + cannot always stop it. `BREAK=ON` checks on every DOS call and makes the + abort reliable, at a small cost to every DOS call. Needs a VM reboot; back + up `FDCONFIG.SYS` first. +- [ ] `TCPAGENT.EXE` connects from a fixed source port (`LOCAL_PORT 2058`). After + the host end closes, a reconnect reuses the same 4-tuple and can flap until + the old state ages out. Observed as a ~10s connect/disconnect cycle after + the daemon is killed mid-connection. diff --git a/src/ferrolang_vm/cli.py b/src/ferrolang_vm/cli.py index 411b144..ff40ddb 100644 --- a/src/ferrolang_vm/cli.py +++ b/src/ferrolang_vm/cli.py @@ -70,6 +70,8 @@ EPILOG = r"""examples: uv run ferro-vm exec 'dir C:\FEC' run a DOS command, print exit code and output uv run ferro-vm put fec/src/check.c 'C:\FEC\SRC\CHECK.C' uv run ferro-vm get 'C:\FEC\TEST.OK' .qemu/TEST.OK + uv run ferro-vm exec --idle-timeout 180 'C:\FEC\BUILD-DOS.BAT' + uv run ferro-vm abort Ctrl+C the command running right now uv run ferro-vm logs follow the structured daemon log The authoritative workspace is C:\FEC inside the VM. Never build on D: (the vvfat @@ -84,6 +86,7 @@ SIMPLE_COMMANDS = { "screenshot": "Capture the VGA console to a PPM/PNG under .qemu/.", "ocr": "Capture the console and print recognized text (RapidOCR).", "logs": "Follow the append-only daemon log. Uses lnav when available.", + "abort": "Interrupt the DOS command currently running (Ctrl+C via QEMU).", } @@ -112,9 +115,19 @@ def main() -> int: execute = commands.add_parser( "exec", help=exec_help, description=exec_help + " Quote the command so the host shell does not eat" - r" backslashes: exec 'wcl386 -q HELLO.C'.") + r" backslashes: exec 'wcl386 -q HELLO.C'." + " A slow command is not a failed one: the wait ends" + " when the guest stops touching its disk, not when a" + " stopwatch expires.") execute.add_argument("command", metavar="DOS_COMMAND", help=r"command line to hand to COMMAND.COM, e.g. 'dir C:\FEC'") + execute.add_argument("--idle-timeout", type=float, default=60, metavar="SECONDS", + help="interrupt once the guest has made no disk access for" + " this long (default: %(default)s)") + execute.add_argument("--hard-timeout", type=float, default=900, metavar="SECONDS", + help="interrupt after this much total time regardless of" + " activity; the backstop for a CPU-bound hang" + " (default: %(default)s)") put_help = "Copy a host file into the VM." put = commands.add_parser("put", help=put_help, description=put_help) @@ -161,7 +174,10 @@ def main() -> int: return 0 raise RuntimeError("TCPAGENT did not become ready") payload: dict[str, object] = {"op": args.op} - if args.op == "exec": payload["command"] = args.command + if args.op == "exec": + payload["command"] = args.command + payload["idle_timeout"] = args.idle_timeout + payload["hard_timeout"] = args.hard_timeout if args.op == "put": payload["source"] = str(args.source.resolve()) payload["destination"] = args.destination diff --git a/src/ferrolang_vm/daemon.py b/src/ferrolang_vm/daemon.py index e35bdf2..cbe9219 100644 --- a/src/ferrolang_vm/daemon.py +++ b/src/ferrolang_vm/daemon.py @@ -8,6 +8,7 @@ from __future__ import annotations import json import logging import os +import re import select import shutil import socket @@ -26,6 +27,22 @@ AGENT_ADDRESS = ("127.0.0.1", 5558) MONITOR_ADDRESS = ("127.0.0.1", 4444) LOG_PATH = QEMU / "ferro-vm.log" +# Short commands answer promptly, so a plain socket timeout is the right guard. +REQUEST_TIMEOUT = 30 +# EXEC is different: TCPAGENT is frozen inside system() for the whole command +# and cannot answer, so silence proves nothing. Wake up often, decide with +# guest liveness instead of a stopwatch, and never discard the connection just +# because a compile is slow. +EXEC_POLL_SECONDS = 2 +DEFAULT_IDLE_TIMEOUT = 60 +DEFAULT_HARD_TIMEOUT = 900 +# Budget for collecting the result after Ctrl+C, before giving up on the stream. +INTERRUPT_GRACE_SECONDS = 15 + + +class ExecInterrupted(RuntimeError): + """The supervisor decided to stop the running DOS command.""" + def configure_logging() -> None: QEMU.mkdir(exist_ok=True) @@ -46,16 +63,37 @@ class Host: self.agent_lock = threading.Lock() self.agent_ready = threading.Event() self.qemu: subprocess.Popen[bytes] | None = None + # QEMU accepts one monitor connection at a time, and the EXEC + # supervisor polls it while other control requests run concurrently. + self.monitor_lock = threading.Lock() + self.abort_requested = threading.Event() + self.exec_active = threading.Event() - def accept_agents(self) -> None: + @staticmethod + def bind_agent_listener() -> socket.socket: + """Bind 5558 exclusively so a second daemon fails loudly. + + SO_REUSEADDR means something different on Windows than on Unix: it lets + another process bind an already-bound port and quietly take over new + connections, so a duplicate daemon would be silently half-working + instead of refusing to start. SO_EXCLUSIVEADDRUSE is the Windows way to + say "only me". + """ server = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + exclusive = getattr(socket, "SO_EXCLUSIVEADDRUSE", None) + if exclusive is not None: + server.setsockopt(socket.SOL_SOCKET, exclusive, 1) + else: + server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) server.bind(AGENT_ADDRESS) server.listen(1) + return server + + def accept_agents(self, server: socket.socket) -> None: log_event(logging.INFO, "agent listener ready", address="127.0.0.1:5558") while True: sock, peer = server.accept() - sock.settimeout(30) + sock.settimeout(REQUEST_TIMEOUT) with self.agent_lock: if self.agent is not None: sock.close() @@ -94,15 +132,62 @@ class Host: log_event(logging.INFO, "agent disconnected", peer=f"{peer[0]}:{peer[1]}") @staticmethod - def _read_line(sock: socket.socket) -> bytes: + def _read_line(sock: socket.socket, supervise=None) -> bytes: + # Partial input is kept across timeouts, so supervise() may fire in the + # middle of a line without losing what has already arrived. out = bytearray() while not out.endswith(b"\n"): - part = sock.recv(1) + try: + part = sock.recv(1) + except TimeoutError: + if supervise is None: + raise + supervise() + continue if not part: raise ConnectionError("TCP agent closed connection") out.extend(part) return bytes(out).rstrip(b"\r\n") + @staticmethod + def _read_exactly(sock: socket.socket, count: int, supervise=None) -> bytes: + chunks: list[bytes] = [] + while count: + try: + chunk = sock.recv(min(65536, count)) + except TimeoutError: + if supervise is None: + raise + supervise() + continue + if not chunk: + raise ConnectionError("TCP agent closed connection") + chunks.append(chunk) + count -= len(chunk) + return b"".join(chunks) + + def guest_idle_seconds(self) -> float | None: + """Seconds since the guest last touched a disk, or None if unknown. + + QEMU keeps counting while TCPAGENT is frozen inside system(), so this + is the one progress signal available during a long DOS command. A + purely CPU-bound command looks idle here, which is what the hard + timeout is for. + """ + try: + text = self.monitor("info blockstats") + except OSError: + return None + idle: float | None = None + for line in text.splitlines(): + operations = re.search(r"rd_operations=(\d+)", line) + elapsed = re.search(r"idle_time_ns=(\d+)", line) + if not operations or not elapsed or int(operations.group(1)) == 0: + continue + seconds = int(elapsed.group(1)) / 1e9 + idle = seconds if idle is None else min(idle, seconds) + return idle + def request(self, command: str, payload: bytes = b"") -> str: with self.agent_lock: if self.agent is None: @@ -125,45 +210,107 @@ class Host: def ping(self) -> dict[str, object]: return {"response": self.request("PING")} - def exec(self, command: str) -> dict[str, object]: - log_event(logging.INFO, "exec start", command=command) + def _read_exec_result(self, sock: socket.socket, supervise=None) -> tuple[int, int, bytes]: + header = self._read_line(sock, supervise).decode("ascii", "replace") + fields = header.split() + if len(fields) == 4 and fields[0] == "RESULT": + code, length, flags = int(fields[1]), int(fields[2]), int(fields[3]) + return code, flags, self._read_exactly(sock, length, supervise) + if fields and fields[0] == "OK": + # Compatibility with an installed pre-RESULT agent. + return int(fields[1]), 1, bytes.fromhex(fields[2]) if len(fields) > 2 else b"" + raise RuntimeError("malformed EXEC response: " + header) + + def exec(self, command: str, idle_timeout: float = DEFAULT_IDLE_TIMEOUT, + hard_timeout: float = DEFAULT_HARD_TIMEOUT) -> dict[str, object]: + log_event(logging.INFO, "exec start", command=command, + idle_timeout=idle_timeout, hard_timeout=hard_timeout) started = time.monotonic() + self.abort_requested.clear() with self.agent_lock: if self.agent is None: raise RuntimeError("TCPAGENT is not connected") sock = self.agent + self.exec_active.set() + reported = started + interrupt_at: float | None = None + interrupted = "" + answered = False + + def supervise() -> None: + """Called every EXEC_POLL_SECONDS while the agent stays silent.""" + nonlocal reported, interrupt_at, interrupted, answered + now = time.monotonic() + idle = self.guest_idle_seconds() + if now - reported >= 15: + reported = now + log_event(logging.INFO, "exec running", elapsed_s=round(now-started, 1), + guest_idle_s=None if idle is None else round(idle, 1)) + if interrupt_at is None: + if self.abort_requested.is_set(): + interrupted = "aborted by request" + elif now - started > hard_timeout: + interrupted = f"hard timeout after {hard_timeout:.0f}s" + elif idle is not None and idle > idle_timeout: + interrupted = f"guest idle {idle:.0f}s exceeds {idle_timeout:.0f}s" + if interrupted: + interrupt_at = now + log_event(logging.WARNING, "exec interrupting", reason=interrupted) + self.monitor("sendkey ctrl-c") + return + waited = now - interrupt_at + if not answered and waited > 4: + # COMMAND.COM asks "Terminate batch file (Y/N/A)?" for .BAT + # targets and sits at that prompt until it is answered. + answered = True + self.monitor("sendkey y") + self.monitor("sendkey ret") + # Ctrl+C only lands at a DOS break check. With BREAK=OFF (the + # FreeDOS default) a compute-bound child whose output we + # redirected to a file may never reach one, so the command runs + # to completion regardless. Keep collecting its result rather + # than abandoning a stream that still owes us one -- give up + # only once the guest has gone quiet too. + if waited > INTERRUPT_GRACE_SECONDS and (idle is None or idle > 5): + raise ExecInterrupted(interrupted + "; command did not stop") + try: encoded = command.encode("ascii", "replace").hex().upper() + sock.settimeout(EXEC_POLL_SECONDS) sock.sendall(f"EXEC {encoded}\n".encode("ascii")) - header = self._read_line(sock).decode("ascii", "replace") - fields = header.split() - if len(fields) == 4 and fields[0] == "RESULT": - code, remaining, flags = int(fields[1]), int(fields[2]), int(fields[3]) - chunks: list[bytes] = [] - while remaining: - chunk = sock.recv(min(65536, remaining)) - if not chunk: - raise ConnectionError("TCP agent closed during EXEC result") - chunks.append(chunk) - remaining -= len(chunk) - raw = b"".join(chunks) - elif fields and fields[0] == "OK": - # Compatibility with an installed pre-RESULT agent. - code = int(fields[1]); flags = 1 - raw = bytes.fromhex(fields[2]) if len(fields) > 2 else b"" - else: - raise RuntimeError("malformed EXEC response: " + header) - except OSError as exc: + code, flags, raw = self._read_exec_result(sock, supervise) + except (OSError, ExecInterrupted) as exc: + # Only now is the stream beyond repair; drop it so the agent + # reconnects with a clean protocol state. if self.agent is sock: self.agent = None self.agent_ready.clear() + sock.close() raise RuntimeError(f"TCPAGENT EXEC failed: {exc}") from exc + finally: + self.exec_active.clear() + self.abort_requested.clear() + try: + sock.settimeout(REQUEST_TIMEOUT) + except OSError: + pass output = raw.decode("cp437", "replace") for line in output.splitlines(): log_event(logging.INFO, "dos output", line=line) log_event(logging.INFO, "exec finish", exit=code, bytes=len(raw), flags=flags, + interrupted=interrupted or None, elapsed_ms=round((time.monotonic()-started)*1000)) - return {"exit": code, "output": output, "bytes": len(raw), "flags": flags} + result = {"exit": code, "output": output, "bytes": len(raw), "flags": flags} + if interrupted: + result["interrupted"] = interrupted + return result + + def abort(self) -> dict[str, object]: + if not self.exec_active.is_set(): + return {"aborted": False, "reason": "no command is running"} + self.abort_requested.set() + log_event(logging.INFO, "abort requested") + return {"aborted": True} def put(self, source: str, destination: str) -> dict[str, object]: data = Path(source).read_bytes() @@ -198,9 +345,8 @@ class Host: log_event(logging.INFO, "get", path=source, bytes=len(data), destination=str(target)) return {"path": source, "destination": str(target), "bytes": len(data)} - @staticmethod - def monitor(command: str) -> str: - with socket.create_connection(MONITOR_ADDRESS, timeout=3) as sock: + def monitor(self, command: str) -> str: + with self.monitor_lock, socket.create_connection(MONITOR_ADDRESS, timeout=3) as sock: sock.settimeout(1) time.sleep(.1) try: @@ -269,7 +415,11 @@ class Host: if op == "start": return self.start() if op == "stop": return self.stop() if op == "ping": return self.ping() - if op == "exec": return self.exec(str(request["command"])) + if op == "abort": return self.abort() + if op == "exec": + return self.exec(str(request["command"]), + float(request.get("idle_timeout", DEFAULT_IDLE_TIMEOUT)), + float(request.get("hard_timeout", DEFAULT_HARD_TIMEOUT))) if op == "put": return self.put(str(request["source"]), str(request["destination"])) if op == "get": return self.get(str(request["source"]), str(request["destination"])) if op == "screenshot": return self.screenshot() @@ -280,8 +430,8 @@ class Host: def serve_pipe(host: Host) -> None: listener = Listener(PIPE, family="AF_PIPE") log_event(logging.INFO, "control pipe ready", pipe=PIPE) - while True: - conn = listener.accept() + + def handle(conn) -> None: try: request = conn.recv() try: @@ -292,13 +442,23 @@ def serve_pipe(host: Host) -> None: finally: conn.close() + while True: + # One thread per request: `abort` has to be answerable while a long + # `exec` is still holding the agent. + threading.Thread(target=handle, args=(listener.accept(),), daemon=True).start() + def main() -> None: if os.name != "nt": raise SystemExit("ferro-vm currently supports Windows only") configure_logging() host = Host() - threading.Thread(target=host.accept_agents, daemon=True).start() + try: + server = host.bind_agent_listener() + except OSError as exc: + log_event(logging.ERROR, "agent listener bind failed", address="127.0.0.1:5558", error=str(exc)) + raise SystemExit(f"another ferro-vm daemon already owns 127.0.0.1:5558 ({exc})") + threading.Thread(target=host.accept_agents, args=(server,), daemon=True).start() serve_pipe(host) diff --git a/tools/tcpagent/README.md b/tools/tcpagent/README.md index 0b2487f..249b8c8 100644 --- a/tools/tcpagent/README.md +++ b/tools/tcpagent/README.md @@ -61,9 +61,28 @@ are not interpreted on this FreeDOS console, so neither of the usual routes works. Elapsed times come from the BIOS tick counter at 18.2065 Hz (~55 ms resolution). -Note that mTCP is not driven while `system()` runs a child, so a DOS command -lasting tens of seconds can drop the TCP connection. The agent logs `link lost` -and reconnects on its own, but the host loses that command's result. +## Long commands + +mTCP is only driven when the agent calls it, and `system()` freezes the agent +for the entire child command. So during a long `EXEC` the DOS side is mute: it +cannot answer, cannot acknowledge, cannot report progress. Silence therefore +proves nothing about whether the command is healthy. + +The host must not read that silence as failure. `ferro-vm exec` waits on +QEMU's own view of the guest instead: `info blockstats` keeps counting while +the agent is frozen, and `idle_time_ns` distinguishes a slow command from a +stuck one. See `--idle-timeout` and `--hard-timeout` in `ferro-vm exec --help`. + +When the host does decide to stop a command it injects Ctrl+C through the QEMU +monitor, then answers COMMAND.COM's `Terminate batch file (Y/N/A)?` prompt. +That is a request, not a guarantee: Ctrl+C only lands at a DOS break check, and +with `BREAK=OFF` (the FreeDOS default in `C:\FDCONFIG.SYS`) a compute-bound +child whose output we redirected to a file may never reach one. The host keeps +collecting the result either way rather than abandoning a stream that still +owes it a `RESULT`. + +Adding `BREAK=ON` to `C:\FDCONFIG.SYS` would make DOS check on every system +call and so make Ctrl+C reliable, at a small cost to every DOS call. ## Rebuilding inside the VM