Files
juc500/juc500_xfer/protocol.py
T
joyharam ff60cee1cc Bypass pyusb Device I/O on WinUSB with direct libusb bulk transfers
Device.write/read always call get_configuration which WinUSB does not
implement. Route all framed I/O through the backend with IF5 claimed.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-15 11:29:01 +09:00

534 lines
20 KiB
Python

"""Framed packet protocol over the JUC500 vendor bulk bridge.
Frame layout (little-endian):
magic : 4 bytes b"J5FX"
type : u8
flags : u8
seq : u32
length : u32
payload : length bytes
crc32 : u32 zlib.crc32(payload)
"""
from __future__ import annotations
import hashlib
import json
import os
import struct
import threading
import time
import uuid
import zlib
from concurrent.futures import Future
from dataclasses import dataclass
from enum import IntEnum
from pathlib import Path
from typing import Any, Callable, Optional
from .usb_link import Juc500Link, MAX_PACKET
MAGIC = b"J5FX"
HEADER_FMT = "<4sBBII"
HEADER_SIZE = struct.calcsize(HEADER_FMT)
CRC_SIZE = 4
MAX_PAYLOAD = 16 * 1024 * 1024
CHUNK_PAYLOAD = 60 * 1024
class MsgType(IntEnum):
HELLO = 1
HELLO_ACK = 2
PING = 3
PONG = 4
FILE_OFFER = 10
FILE_ACCEPT = 11
FILE_REJECT = 12
FILE_CHUNK = 13
FILE_DONE = 14
FILE_ACK = 15
RPC_REQ = 30
RPC_RESP = 31
ERROR = 255
@dataclass
class Frame:
type: MsgType
flags: int
seq: int
payload: bytes
def pack_frame(msg_type: MsgType, payload: bytes = b"", seq: int = 0, flags: int = 0) -> bytes:
crc = zlib.crc32(payload) & 0xFFFFFFFF
header = struct.pack(HEADER_FMT, MAGIC, int(msg_type), flags, seq & 0xFFFFFFFF, len(payload))
return header + payload + struct.pack("<I", crc)
class FrameReader:
def __init__(self) -> None:
self._buf = bytearray()
def feed(self, data: bytes) -> list[Frame]:
self._buf.extend(data)
frames: list[Frame] = []
while True:
if len(self._buf) < HEADER_SIZE + CRC_SIZE:
break
idx = self._buf.find(MAGIC)
if idx < 0:
self._buf = self._buf[-3:]
break
if idx > 0:
del self._buf[:idx]
if len(self._buf) < HEADER_SIZE + CRC_SIZE:
break
magic, typ, flags, seq, length = struct.unpack_from(HEADER_FMT, self._buf, 0)
if magic != MAGIC or length > MAX_PAYLOAD:
del self._buf[0:1]
continue
total = HEADER_SIZE + length + CRC_SIZE
if len(self._buf) < total:
break
payload = bytes(self._buf[HEADER_SIZE : HEADER_SIZE + length])
(crc,) = struct.unpack_from("<I", self._buf, HEADER_SIZE + length)
del self._buf[:total]
if (zlib.crc32(payload) & 0xFFFFFFFF) != crc:
continue
try:
frames.append(Frame(MsgType(typ), flags, seq, payload))
except ValueError:
continue
return frames
def _default_peer_name() -> str:
host = os.uname().nodename if hasattr(os, "uname") else "host"
try:
user = os.getlogin()
except OSError:
user = os.environ.get("USER") or os.environ.get("USERNAME") or "user"
return f"{host}-{user}"
class Session:
"""Bidirectional session: HELLO + file frames + JSON RPC."""
def __init__(
self,
link: Juc500Link,
peer_name: Optional[str] = None,
on_log: Optional[Callable[[str], None]] = None,
local_root: Optional[Path] = None,
):
self.link = link
self.peer_name = peer_name or _default_peer_name()
self.on_log = on_log or (lambda m: None)
self.local_root = (local_root or Path.home()).resolve()
self._reader = FrameReader()
self._seq = 0
self._lock = threading.Lock()
self._xfer_lock = threading.Lock()
self._stop = threading.Event()
self._rx_thread: Optional[threading.Thread] = None
self._alive_thread: Optional[threading.Thread] = None
self._handlers: dict[MsgType, Callable[[Frame], None]] = {}
self._pending_accept = threading.Event()
self._accept_result: Optional[bool] = None
self._peer_hello: Optional[str] = None
self._file_ack = threading.Event()
self._connected = threading.Event()
self._rpc_waiters: dict[str, Future] = {}
self._rpc_lock = threading.Lock()
# inbound file receive staging (when peer pushes)
self._inbound: dict[str, Any] = {}
self.inbound_dest = self.local_root
self.on_progress: Optional[Callable[[int, int, str], None]] = None
def log(self, msg: str) -> None:
self.on_log(msg)
def start(self) -> None:
self._stop.clear()
self._handlers[MsgType.HELLO] = self._on_hello
self._handlers[MsgType.HELLO_ACK] = self._on_hello_ack
self._handlers[MsgType.PING] = self._on_ping
self._handlers[MsgType.PONG] = lambda _f: None
self._handlers[MsgType.FILE_ACCEPT] = self._on_file_accept
self._handlers[MsgType.FILE_REJECT] = self._on_file_reject
self._handlers[MsgType.FILE_ACK] = self._on_file_ack
self._handlers[MsgType.FILE_OFFER] = self._on_inbound_offer
self._handlers[MsgType.FILE_CHUNK] = self._on_inbound_chunk
self._handlers[MsgType.FILE_DONE] = self._on_inbound_done
self._handlers[MsgType.RPC_REQ] = self._on_rpc_req
self._handlers[MsgType.RPC_RESP] = self._on_rpc_resp
self._rx_thread = threading.Thread(target=self._rx_loop, name="juc500-rx", daemon=True)
self._alive_thread = threading.Thread(target=self._alive_loop, name="juc500-alive", daemon=True)
self._rx_thread.start()
self._alive_thread.start()
self.send(MsgType.HELLO, self.peer_name.encode("utf-8"))
def stop(self) -> None:
self._stop.set()
if self._rx_thread:
self._rx_thread.join(timeout=2)
if self._alive_thread:
self._alive_thread.join(timeout=2)
with self._rpc_lock:
for fut in self._rpc_waiters.values():
if not fut.done():
fut.set_exception(RuntimeError("session stopped"))
self._rpc_waiters.clear()
def wait_peer(self, timeout: float = 30.0) -> str:
deadline = time.monotonic() + timeout
while not self._connected.is_set():
remaining = deadline - time.monotonic()
if remaining <= 0:
raise TimeoutError(
"상대가 JUC500 전송 프로그램으로 연결되지 않았습니다. "
"양쪽에서 GUI/앱을 실행 중인지 확인하세요."
)
try:
self.send(MsgType.HELLO, self.peer_name.encode("utf-8"))
except Exception as exc:
# Peer not reading yet (FIFO full) → USB timeout; keep retrying quietly.
msg = str(exc)
if "timed out" in msg.lower() or "timeout" in msg.lower():
self.log("상대 수신 대기 중… (HELLO)")
else:
self.log(f"HELLO 재전송 실패: {exc}")
self._connected.wait(timeout=min(2.0, remaining))
return self._peer_hello or "peer"
def send(self, msg_type: MsgType, payload: bytes = b"", flags: int = 0) -> int:
with self._lock:
self._seq = (self._seq + 1) & 0xFFFFFFFF
seq = self._seq
frame = pack_frame(msg_type, payload, seq=seq, flags=flags)
return self.link.write_data(frame)
def rpc(self, op: str, timeout: float = 30.0, **kwargs: Any) -> Any:
req_id = uuid.uuid4().hex
body = {"id": req_id, "op": op, **kwargs}
fut: Future = Future()
with self._rpc_lock:
self._rpc_waiters[req_id] = fut
self.send(MsgType.RPC_REQ, json.dumps(body, ensure_ascii=False).encode("utf-8"))
try:
return fut.result(timeout=timeout)
finally:
with self._rpc_lock:
self._rpc_waiters.pop(req_id, None)
def _rx_loop(self) -> None:
while not self._stop.is_set():
chunk = self.link.read_data(MAX_PACKET, timeout_ms=200)
if not chunk:
continue
for frame in self._reader.feed(chunk):
handler = self._handlers.get(frame.type)
if handler:
try:
handler(frame)
except Exception as exc:
self.log(f"핸들러 오류 ({frame.type.name}): {exc}")
def _alive_loop(self) -> None:
while not self._stop.is_set():
try:
self.link.write_ctrl(b"\x00ALIVE" + struct.pack("<d", time.time()), timeout_ms=500)
except Exception:
pass
try:
self.link.read_ctrl(MAX_PACKET, timeout_ms=50)
except Exception:
pass
self._stop.wait(1.0)
if self._connected.is_set():
try:
self.send(MsgType.PING, struct.pack("<d", time.time()))
except Exception:
pass
self._stop.wait(2.0)
def _on_hello(self, frame: Frame) -> None:
name = frame.payload.decode("utf-8", errors="replace")
self._peer_hello = name
self.log(f"상대 HELLO: {name}")
self.send(MsgType.HELLO_ACK, self.peer_name.encode("utf-8"))
self._connected.set()
def _on_hello_ack(self, frame: Frame) -> None:
name = frame.payload.decode("utf-8", errors="replace")
self._peer_hello = name
self.log(f"상대 HELLO_ACK: {name}")
self._connected.set()
def _on_ping(self, frame: Frame) -> None:
self.send(MsgType.PONG, frame.payload)
def _on_file_accept(self, _frame: Frame) -> None:
self._accept_result = True
self._pending_accept.set()
def _on_file_reject(self, _frame: Frame) -> None:
self._accept_result = False
self._pending_accept.set()
def _on_file_ack(self, _frame: Frame) -> None:
self._file_ack.set()
def _safe_under_root(self, rel: str) -> Path:
rel = (rel or "").replace("\\", "/").lstrip("/")
target = (self.local_root / rel).resolve()
try:
target.relative_to(self.local_root)
except ValueError as exc:
raise PermissionError(f"경로가 루트를 벗어남: {rel}") from exc
return target
def list_local(self, rel: str = "") -> dict[str, Any]:
path = self._safe_under_root(rel)
if not path.exists():
raise FileNotFoundError(rel)
if not path.is_dir():
raise NotADirectoryError(rel)
entries = []
for child in sorted(path.iterdir(), key=lambda p: (not p.is_dir(), p.name.lower())):
try:
st = child.stat()
entries.append(
{
"name": child.name,
"is_dir": child.is_dir(),
"size": 0 if child.is_dir() else int(st.st_size),
"mtime": int(st.st_mtime),
}
)
except OSError:
continue
return {"path": rel.replace("\\", "/").strip("/"), "entries": entries}
def mkdir_local(self, rel: str) -> dict[str, Any]:
path = self._safe_under_root(rel)
path.mkdir(parents=True, exist_ok=True)
return {"path": rel.replace("\\", "/").strip("/")}
def _on_rpc_req(self, frame: Frame) -> None:
try:
req = json.loads(frame.payload.decode("utf-8"))
except Exception as exc:
self._rpc_reply({"id": "", "ok": False, "error": f"bad json: {exc}"})
return
req_id = req.get("id", "")
op = req.get("op")
try:
if op == "listdir":
result = self.list_local(req.get("path") or "")
elif op == "mkdir":
result = self.mkdir_local(req.get("path") or "")
elif op == "ping":
result = {"pong": True, "name": self.peer_name}
elif op == "pull":
# Peer asks us to push a file/folder to them.
rel = req.get("path") or ""
# Respond OK first; actual bytes go via FILE_OFFER stream in a worker.
self._rpc_reply({"id": req_id, "ok": True, "result": {"queued": True, "path": rel}})
threading.Thread(
target=self._push_path_for_pull,
args=(rel,),
daemon=True,
name="juc500-pull-push",
).start()
return
else:
raise ValueError(f"unknown op: {op}")
self._rpc_reply({"id": req_id, "ok": True, "result": result})
except Exception as exc:
self._rpc_reply({"id": req_id, "ok": False, "error": str(exc)})
def _rpc_reply(self, body: dict[str, Any]) -> None:
self.send(MsgType.RPC_RESP, json.dumps(body, ensure_ascii=False).encode("utf-8"))
def _on_rpc_resp(self, frame: Frame) -> None:
try:
body = json.loads(frame.payload.decode("utf-8"))
except Exception:
return
req_id = body.get("id")
with self._rpc_lock:
fut = self._rpc_waiters.get(req_id or "")
if not fut or fut.done():
return
if body.get("ok"):
fut.set_result(body.get("result"))
else:
fut.set_exception(RuntimeError(body.get("error") or "RPC failed"))
def _push_path_for_pull(self, rel: str) -> None:
try:
path = self._safe_under_root(rel)
if path.is_file():
self.send_file(path, remote_name=Path(rel).name)
elif path.is_dir():
for file_path in path.rglob("*"):
if file_path.is_file():
remote_name = str(file_path.relative_to(self.local_root)).replace("\\", "/")
# When pulled, receiver picks dest; send basename under relative tree
rel_under = str(file_path.relative_to(path)).replace("\\", "/")
self.send_file(file_path, remote_name=rel_under)
else:
self.log(f"pull 대상 없음: {rel}")
except Exception as exc:
self.log(f"pull push 실패: {exc}")
# --- outbound send ---
def send_file(
self,
path: Path,
remote_name: Optional[str] = None,
progress: Optional[Callable[[int, int], None]] = None,
accept_timeout: float = 120.0,
) -> None:
path = Path(path)
if not path.is_file():
raise FileNotFoundError(path)
with self._xfer_lock:
self._send_file_locked(path, remote_name or path.name, progress, accept_timeout)
def _send_file_locked(
self,
path: Path,
name: str,
progress: Optional[Callable[[int, int], None]],
accept_timeout: float,
) -> None:
size = path.stat().st_size
sha = hashlib.sha256()
with path.open("rb") as f:
while True:
b = f.read(1024 * 1024)
if not b:
break
sha.update(b)
digest = sha.hexdigest()
meta = f"{name}\n{size}\n{digest}".encode("utf-8")
self._pending_accept.clear()
self._accept_result = None
self.log(f"파일 제안: {name} ({size} bytes)")
self.send(MsgType.FILE_OFFER, meta)
if not self._pending_accept.wait(accept_timeout):
raise TimeoutError("상대가 파일 수신을 수락하지 않았습니다.")
if not self._accept_result:
raise RuntimeError("상대가 파일 수신을 거부했습니다.")
sent = 0
with path.open("rb") as f:
while True:
chunk = f.read(CHUNK_PAYLOAD)
if not chunk:
break
self._file_ack.clear()
self.send(MsgType.FILE_CHUNK, chunk)
if not self._file_ack.wait(60.0):
raise TimeoutError("청크 ACK 타임아웃")
sent += len(chunk)
if progress:
progress(sent, size)
if self.on_progress:
self.on_progress(sent, size, name)
self.send(MsgType.FILE_DONE, digest.encode("ascii"))
self.log(f"전송 완료: {name}")
# --- inbound receive (auto-accept into inbound_dest) ---
def _on_inbound_offer(self, frame: Frame) -> None:
# If we're mid outbound wait for ACCEPT, this is our counterpart — ignore auto path
# Actually inbound offers always handled here; outbound waits on _pending_accept which is set by ACCEPT from peer.
try:
name, size_s, digest = frame.payload.decode("utf-8").split("\n", 2)
size = int(size_s)
except Exception:
self.send(MsgType.FILE_REJECT, b"bad offer")
return
# sanitize name: no absolute, no ..
name = name.replace("\\", "/").lstrip("/")
parts = [p for p in name.split("/") if p and p != ".."]
safe = "/".join(parts)
if not safe:
self.send(MsgType.FILE_REJECT, b"bad name")
return
out = (Path(self.inbound_dest) / safe).resolve()
try:
out.relative_to(Path(self.inbound_dest).resolve())
except ValueError:
self.send(MsgType.FILE_REJECT, b"path escape")
return
out.parent.mkdir(parents=True, exist_ok=True)
self._inbound = {
"path": out,
"name": safe,
"size": size,
"digest": digest.strip(),
"fh": out.open("wb"),
"sha": hashlib.sha256(),
"received": 0,
}
self.send(MsgType.FILE_ACCEPT, b"ok")
self.log(f"수신 시작: {safe} ({size} bytes)")
def _on_inbound_chunk(self, frame: Frame) -> None:
st = self._inbound
fh = st.get("fh")
if not fh:
return
fh.write(frame.payload)
st["sha"].update(frame.payload)
st["received"] = st.get("received", 0) + len(frame.payload)
self.send(MsgType.FILE_ACK, b"")
if self.on_progress:
self.on_progress(st["received"], st["size"], st.get("name") or "")
def _on_inbound_done(self, frame: Frame) -> None:
st = self._inbound
fh = st.pop("fh", None)
if fh:
fh.close()
digest = st.get("sha").hexdigest() if st.get("sha") else ""
expected = st.get("digest", "")
remote = frame.payload.decode("ascii", errors="replace")
if digest != expected or (remote and remote != digest):
self.log(f"해시 불일치: {st.get('name')}")
else:
self.log(f"수신 완료: {st.get('path')}")
self._inbound = {}
# compatibility helpers used by CLI
def receive_file(
self,
dest_dir: Path,
allow: Optional[Callable[[str, int], bool]] = None,
progress: Optional[Callable[[int, int], None]] = None,
timeout: float = 600.0,
) -> Path:
"""Wait until an inbound transfer completes into dest_dir."""
self.inbound_dest = Path(dest_dir)
self.inbound_dest.mkdir(parents=True, exist_ok=True)
deadline = time.monotonic() + timeout
last_path = None
while time.monotonic() < deadline:
if self._inbound.get("fh") is None and last_path and not self._inbound:
# completed
if progress and last_path:
pass
return Path(last_path)
if self._inbound.get("path"):
last_path = self._inbound["path"]
if progress and self._inbound.get("size"):
progress(self._inbound.get("received", 0), self._inbound["size"])
time.sleep(0.1)
raise TimeoutError("파일 수신 타임아웃")