mirror of
https://github.com/MCKero6423/uv-k5-v3-emulator.git
synced 2026-10-02 11:07:31 +00:00
Add a QMP client module for the web UI
Wraps the connect/negotiate/command dance with a lock, since only one client can hold the QMP socket at a time. Events interleave with replies, so command() reads until it sees the matching return rather than trusting the first message. Connection failure carries the reason and the likely cause: key.py already holding the socket.
This commit is contained in:
1 parent
b246132a35
commit
59ca82403f
3 files changed
+164
No files matched your search
@@ -10,3 +10,6 @@ build/
|
||||
|
||||
# Screenshots used in the README are committed on purpose.
|
||||
!docs/screenshots/*.png
|
||||
|
||||
# Plan documents from the plan skill: local working notes.
|
||||
.hermes/
|
||||
@@ -0,0 +1,97 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Unit tests for the QMP client. Uses a fake server, so no emulator needed."""
|
||||
import json
|
||||
import os
|
||||
import socket
|
||||
import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
|
||||
from uvk5_qmp import QmpClient
|
||||
|
||||
|
||||
def fake_server(path, script):
|
||||
"""Minimal QMP server: greets, then replies to each command from `script`."""
|
||||
srv = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
srv.bind(path)
|
||||
srv.listen(1)
|
||||
|
||||
def run():
|
||||
conn, _ = srv.accept()
|
||||
conn.sendall(json.dumps({"QMP": {"version": {}}}).encode() + b"\n")
|
||||
buf = b""
|
||||
for reply in script:
|
||||
while b"\n" not in buf:
|
||||
chunk = conn.recv(4096)
|
||||
if not chunk:
|
||||
return
|
||||
buf += chunk
|
||||
_, buf = buf.split(b"\n", 1)
|
||||
conn.sendall(json.dumps(reply).encode() + b"\n")
|
||||
conn.close()
|
||||
srv.close()
|
||||
|
||||
threading.Thread(target=run, daemon=True).start()
|
||||
return srv
|
||||
|
||||
|
||||
class TestQmpClient(unittest.TestCase):
|
||||
def test_negotiates_and_returns_command_result(self):
|
||||
path = os.path.join(tempfile.mkdtemp(), "qmp.sock")
|
||||
# reply 1 = qmp_capabilities, reply 2 = our command
|
||||
fake_server(path, [{"return": {}}, {"return": {"status": "running"}}])
|
||||
|
||||
client = QmpClient(path)
|
||||
self.addCleanup(client.close)
|
||||
self.assertEqual(client.command("query-status"), {"status": "running"})
|
||||
|
||||
def test_raises_on_qmp_error(self):
|
||||
path = os.path.join(tempfile.mkdtemp(), "qmp.sock")
|
||||
fake_server(path, [{"return": {}},
|
||||
{"error": {"class": "GenericError", "desc": "nope"}}])
|
||||
client = QmpClient(path)
|
||||
self.addCleanup(client.close)
|
||||
with self.assertRaises(RuntimeError) as ctx:
|
||||
client.command("bogus")
|
||||
self.assertIn("nope", str(ctx.exception))
|
||||
|
||||
def test_skips_interleaved_events(self):
|
||||
"""Events arrive unsolicited and must not be mistaken for a result.
|
||||
|
||||
The fake server answers one command with an event followed by the real
|
||||
return, so a client that stopped at the first message would hand back the
|
||||
event instead.
|
||||
"""
|
||||
path = os.path.join(tempfile.mkdtemp(), "qmp.sock")
|
||||
srv = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
srv.bind(path)
|
||||
srv.listen(1)
|
||||
|
||||
def run():
|
||||
conn, _ = srv.accept()
|
||||
conn.sendall(json.dumps({"QMP": {"version": {}}}).encode() + b"\n")
|
||||
buf = b""
|
||||
|
||||
def next_command():
|
||||
nonlocal buf
|
||||
while b"\n" not in buf:
|
||||
buf += conn.recv(4096)
|
||||
line, buf = buf.split(b"\n", 1)
|
||||
return json.loads(line)
|
||||
|
||||
next_command() # qmp_capabilities
|
||||
conn.sendall(json.dumps({"return": {}}).encode() + b"\n")
|
||||
next_command() # query-status
|
||||
# An event first, then the actual return.
|
||||
conn.sendall(json.dumps({"event": "RESUME"}).encode() + b"\n")
|
||||
conn.sendall(json.dumps({"return": {"status": "running"}}).encode() + b"\n")
|
||||
|
||||
threading.Thread(target=run, daemon=True).start()
|
||||
|
||||
client = QmpClient(path)
|
||||
self.addCleanup(client.close)
|
||||
self.assertEqual(client.command("query-status"), {"status": "running"})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,64 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Minimal QMP client for the UV-K5 emulator.
|
||||
|
||||
Deliberately not shared with tools/key.py: that script is standalone on purpose,
|
||||
so it runs with nothing but the stdlib and no sys.path juggling.
|
||||
|
||||
Only one client can hold the QMP socket at a time -- the emulator is started with
|
||||
server=on,wait=off, which accepts a single connection. A long-lived server owns
|
||||
it and callers serialise through the lock.
|
||||
"""
|
||||
import json
|
||||
import socket
|
||||
import threading
|
||||
|
||||
|
||||
class QmpClient:
|
||||
def __init__(self, path: str, timeout: float = 5.0):
|
||||
self._lock = threading.Lock()
|
||||
self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
self._sock.settimeout(timeout)
|
||||
try:
|
||||
self._sock.connect(path)
|
||||
except OSError as exc:
|
||||
raise RuntimeError(
|
||||
f"cannot reach the emulator at {path}: {exc}\n"
|
||||
"Start it with tools/run.sh first. Note the QMP socket takes a "
|
||||
"single client, so tools/key.py cannot be connected at the same "
|
||||
"time."
|
||||
) from exc
|
||||
self._buf = b""
|
||||
self._read_json() # greeting
|
||||
self.command("qmp_capabilities")
|
||||
|
||||
def _read_json(self) -> dict:
|
||||
while b"\n" not in self._buf:
|
||||
chunk = self._sock.recv(65536)
|
||||
if not chunk:
|
||||
raise RuntimeError("QMP connection closed")
|
||||
self._buf += chunk
|
||||
line, self._buf = self._buf.split(b"\n", 1)
|
||||
return json.loads(line)
|
||||
|
||||
def command(self, name: str, **args):
|
||||
payload = {"execute": name}
|
||||
if args:
|
||||
payload["arguments"] = args
|
||||
with self._lock:
|
||||
self._sock.sendall(json.dumps(payload).encode() + b"\n")
|
||||
while True:
|
||||
msg = self._read_json()
|
||||
if "error" in msg:
|
||||
raise RuntimeError(
|
||||
f"QMP {name} failed: "
|
||||
f"{msg['error'].get('desc', msg['error'])}")
|
||||
if "return" in msg:
|
||||
return msg["return"]
|
||||
# Events (STOP, RESET, RESUME, ...) interleave with replies;
|
||||
# keep reading until the matching return arrives.
|
||||
|
||||
def close(self):
|
||||
try:
|
||||
self._sock.close()
|
||||
except OSError:
|
||||
pass
|
||||
Reference in new issue
Block a user