"""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 queue 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 # Stream payload size — keep well under multi-second USB OUT budgets. CHUNK_PAYLOAD = 256 * 1024 # How many FILE_CHUNK frames may be queued on TX at once. TX_WINDOW = 2 # Bulk IN read size (libusb splits to 512-byte packets). RX_READ_SIZE = 64 * 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 # Input share (binary payloads — see pack_input_*) # INPUT_CTRL u8 target: 0=local 1=remote-req 2=remote-ack 3=remote-nack # 4=share-on 5=share-off + u16 w + u16 h (w/h used for share-on screen size) INPUT_CTRL = 40 INPUT_MOVE = 41 # i32 x, i32 y (absolute virtual cursor on peer screen) INPUT_BUTTON = 42 # u8 button(1=left,2=right,3=middle), u8 pressed(0/1) INPUT_WHEEL = 43 # i16 dx, i16 dy INPUT_KEY = 44 # u8 pressed + utf-8 key name (e.g. "a", "space", "ctrl") INPUT_CLIPBOARD = 45 # utf-8 plain text to set on peer clipboard ERROR = 255 # INPUT_CTRL targets INPUT_CTRL_LOCAL = 0 INPUT_CTRL_REMOTE = 1 # request: take control of peer INPUT_CTRL_REMOTE_ACK = 2 # peer accepted; payload = peer screen size INPUT_CTRL_REMOTE_NACK = 3 # peer rejected (share off / no permission) INPUT_CTRL_SHARE_ON = 4 # peer requests both sides enable input share INPUT_CTRL_SHARE_OFF = 5 # peer requests both sides disable input share INPUT_CTRL_PEER_CLICK = 6 # peer → controller: click injected; w=button, h=x def pack_input_ctrl(target: int, width: int = 0, height: int = 0) -> bytes: return struct.pack(" tuple[int, int, int]: if len(payload) < 5: return INPUT_CTRL_LOCAL, 0, 0 target, w, h = struct.unpack_from(" bytes: """Absolute virtual cursor position on peer screen (i32 x, i32 y).""" return struct.pack( " tuple[int, int]: if len(payload) >= 8: return struct.unpack_from("= 4: # Legacy i16 delta packets — treat as no-op absolute return 0, 0 return 0, 0 def pack_input_button( button: int, pressed: bool, x: int = -1, y: int = -1 ) -> bytes: """Legacy 2-byte packets omit coords; extended packs absolute virtual (x,y).""" if x >= 0 and y >= 0: return struct.pack( " tuple[int, bool, int, int]: if len(payload) < 2: return 1, False, -1, -1 button, pressed = struct.unpack_from("= 10: x, y = struct.unpack_from(" bytes: return struct.pack(" tuple[int, int]: if len(payload) < 4: return 0, 0 return struct.unpack_from(" bytes: raw = (name or "?").encode("utf-8")[:48] return bytes([1 if pressed else 0]) + raw def unpack_input_key(payload: bytes) -> tuple[str, bool]: if not payload: return "?", False return payload[1:].decode("utf-8", errors="replace"), bool(payload[0]) @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(" 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(" 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}" # Special browse path: list drives / volumes (Windows C:/D:/…, macOS / + /Volumes) ROOTS_PATH = "__roots__" def list_fs_roots() -> dict[str, Any]: """Return drive / volume entries for the computer roots browser.""" entries: list[dict[str, Any]] = [] if os.name == "nt": import string for letter in string.ascii_uppercase: root = f"{letter}:/" try: if not Path(root).exists(): continue except OSError: continue entries.append( { "name": f"{letter}:", "is_dir": True, "size": 0, "mtime": 0, "nav_path": root, "kind": "drive", } ) else: entries.append( { "name": "Macintosh HD", "is_dir": True, "size": 0, "mtime": 0, "nav_path": "/", "kind": "drive", } ) home = Path.home().resolve() entries.append( { "name": f"홈 ({home.name})", "is_dir": True, "size": 0, "mtime": 0, "nav_path": home.as_posix(), "kind": "home", } ) vol_root = Path("/Volumes") if vol_root.is_dir(): try: children = sorted(vol_root.iterdir(), key=lambda p: p.name.lower()) except OSError: children = [] root_res = Path("/").resolve() for child in children: try: if not child.is_dir(): continue # Skip the boot volume alias under /Volumes if child.resolve() == root_res: continue st = child.stat() entries.append( { "name": child.name, "is_dir": True, "size": 0, "mtime": int(st.st_mtime), "nav_path": child.as_posix(), "kind": "volume", } ) except OSError: continue return {"path": ROOTS_PATH, "entries": entries} 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._tx_thread: Optional[threading.Thread] = None self._tx_q: queue.Queue[tuple[bytes, int, Optional[threading.Event], list, list]] = ( queue.Queue() ) self._handlers: dict[MsgType, Callable[[Frame], None]] = {} self._pending_accept = threading.Event() self._accept_result: Optional[bool] = None self._accept_error: str = "" 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 self._rpc_busy = threading.Event() self._xfer_busy = threading.Event() # Input share callbacks (set by InputShare) self.on_input_ctrl: Optional[Callable[[int, int, int], None]] = None self.on_input_move: Optional[Callable[[int, int], None]] = None self.on_input_button: Optional[Callable[[int, bool], None]] = None self.on_input_wheel: Optional[Callable[[int, int], None]] = None self.on_input_key: Optional[Callable[[str, bool], None]] = None self.on_input_clipboard: Optional[Callable[[str], None]] = None # Coalesce high-rate mouse moves; flushed directly by TX (bypass backlog) self._move_acc_lock = threading.Lock() self._move_acc_x: int | None = None self._move_acc_y: int | None = None self._move_flush_pending = False # Pause keepalives while KM capture / inject is active self._km_busy = threading.Event() # Fired once when HELLO/HELLO_ACK completes the handshake (auto-connect on peer) self.on_connected: Optional[Callable[[str], None]] = None # Peer RPC restart_gui → spawn detached start.sh on this host self.on_restart: Optional[Callable[[], None]] = None # USB bulk Entity not found — reopen IF5 / autolisten self.on_link_broken: Optional[Callable[[], None]] = None self._hello_failures = 0 self._hello_log_at = 0.0 self._hello_out_fails = 0 def _notify_link_broken(self) -> None: cb = self.on_link_broken if cb is None: return try: threading.Thread( target=cb, name="juc500-link-broken", daemon=True ).start() except Exception: pass def _is_link_gone_error(self, exc: BaseException) -> bool: msg = str(exc).lower() errno = getattr(exc, "errno", None) # Pipe/timeout: peer may not be draining yet — keep waiting. if errno in (32, 60) or "pipe error" in msg or "timed out" in msg or "timeout" in msg: return False return ( errno in (2, 19) or "entity not found" in msg or "errno 2" in msg or "errno2" in msg or "not found" in msg or "no device" in msg or "연결 끊" in msg ) 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._handlers[MsgType.INPUT_CTRL] = self._on_input_ctrl self._handlers[MsgType.INPUT_MOVE] = self._on_input_move self._handlers[MsgType.INPUT_BUTTON] = self._on_input_button self._handlers[MsgType.INPUT_WHEEL] = self._on_input_wheel self._handlers[MsgType.INPUT_KEY] = self._on_input_key self._handlers[MsgType.INPUT_CLIPBOARD] = self._on_input_clipboard self._tx_thread = threading.Thread(target=self._tx_loop, name="juc500-tx", daemon=True) 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._tx_thread.start() self._rx_thread.start() self._alive_thread.start() try: wake = getattr(self.link, "wake_bridge", None) if callable(wake): wake() except Exception: pass try: self.send( MsgType.HELLO, self.peer_name.encode("utf-8"), timeout_ms=200, wait=False, ) except Exception: # Peer may not be reading yet; wait_peer / _on_hello will continue. pass def stop(self) -> None: self._stop.set() self._xfer_busy.clear() # Wake TX waiter try: self._tx_q.put_nowait((b"", 0, None, None, None)) except Exception: pass if self._tx_thread: self._tx_thread.join(timeout=2) 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 next_hello = 0.0 self._hello_out_fails = getattr(self, "_hello_out_fails", 0) hello_gap = 0.5 while not self._connected.is_set(): if self._stop.is_set(): raise RuntimeError("session stopped") remaining = deadline - time.monotonic() if remaining <= 0: raise TimeoutError( "상대가 JUC500 전송 프로그램으로 연결되지 않았습니다. " "양쪽에서 GUI/앱을 실행 중인지 확인하세요. " "(Windows: 공식 Wormhole 종료 후 이 앱을 실행하세요.)" ) now = time.monotonic() if now >= next_hello: next_hello = now + hello_gap try: wake = getattr(self.link, "wake_bridge", None) if callable(wake): wake() except Exception: pass self._hello_failures += 1 use_blocking = self._hello_failures % 4 == 0 try: self.send( MsgType.HELLO, self.peer_name.encode("utf-8"), timeout_ms=1200 if use_blocking else 300, wait=use_blocking, ) if use_blocking: self._hello_out_fails = 0 except Exception as exc: msg = str(exc) soft_wait = ( "timed out" in msg.lower() or "timeout" in msg.lower() or "pipe" in msg.lower() or "input/output" in msg.lower() or "errno 5" in msg.lower() or "송신 큐" in msg ) if soft_wait: self._hello_out_fails += 1 if now - self._hello_log_at >= 5.0: self._hello_log_at = now self.log("Mac → Windows HELLO 송신 중… (Win 수신 대기 유지)") elif self._is_link_gone_error(exc): if now - self._hello_log_at >= 2.0: self._hello_log_at = now self.log(f"USB 링크 오류 (HELLO): {exc}") self._hello_failures += 1 if self._hello_failures >= 8: self._hello_failures = 0 self._notify_link_broken() else: if now - self._hello_log_at >= 2.0: self._hello_log_at = now self.log(f"HELLO 재전송 실패: {exc}") else: if now - self._hello_log_at >= 5.0: self._hello_log_at = now self.log("Mac → Windows HELLO 송신 중… (Win 수신 대기 유지)") self._connected.wait(timeout=min(0.25, remaining)) return self._peer_hello or "peer" def poke_hello(self) -> None: """Listen-mode HELLO without USB reopen.""" if self._connected.is_set() or self._stop.is_set(): return try: wake = getattr(self.link, "wake_bridge", None) if callable(wake): wake() except Exception: pass self._hello_failures += 1 payload = self.peer_name.encode("utf-8") try: self.send( MsgType.HELLO, payload, timeout_ms=300, wait=False, ) except Exception: pass if self._hello_failures % 3 != 0: return try: self.send( MsgType.HELLO, payload, timeout_ms=1200, wait=True, ) self._hello_out_fails = 0 except Exception as exc: msg = str(exc).lower() if not ( "timed out" in msg or "timeout" in msg or "pipe" in msg or "input/output" in msg or "errno 5" in msg or "송신 큐" in msg ): return def _tx_wait_budget(self, frame_len: int, timeout_ms: int) -> float: packets = max(1, (frame_len + MAX_PACKET - 1) // MAX_PACKET) return max(60.0, packets * 0.08 + timeout_ms / 1000.0 + 15.0) def _enqueue_send( self, msg_type: MsgType, payload: bytes = b"", flags: int = 0, timeout_ms: int | None = None, ) -> tuple[threading.Event, list[Any], list[Any], int]: """Queue a frame; returns (done, result, error, frame_len) for pipelined waits.""" with self._lock: self._seq = (self._seq + 1) & 0xFFFFFFFF seq = self._seq frame = pack_frame(msg_type, payload, seq=seq, flags=flags) to = 2000 if timeout_ms is None else int(timeout_ms) done = threading.Event() result: list[Any] = [0] error: list[Any] = [None] self._tx_q.put((frame, to, done, result, error)) return done, result, error, len(frame) def _await_tx( self, done: threading.Event, error: list[Any], frame_len: int, timeout_ms: int, ) -> None: if not done.wait(timeout=self._tx_wait_budget(frame_len, timeout_ms)): raise TimeoutError( f"USB 송신 큐 대기 시간 초과 ({frame_len} bytes)" ) if error[0] is not None: raise error[0] def send( self, msg_type: MsgType, payload: bytes = b"", flags: int = 0, timeout_ms: int | None = None, *, wait: bool = True, ) -> int: """Queue a frame on the TX thread (serializes USB OUT).""" # Best-effort frames must not clog the queue on USB OUT stalls. if not wait: with self._lock: self._seq = (self._seq + 1) & 0xFFFFFFFF seq = self._seq frame = pack_frame(msg_type, payload, seq=seq, flags=flags) to = min(150 if timeout_ms is None else int(timeout_ms), 800) self._tx_q.put((frame, to, None, None, None)) return len(frame) to = 2000 if timeout_ms is None else int(timeout_ms) done, result, error, flen = self._enqueue_send( msg_type, payload, flags=flags, timeout_ms=to ) self._await_tx(done, error, flen, to) return int(result[0] or 0) def set_km_busy(self, busy: bool) -> None: if busy: self._km_busy.set() else: self._km_busy.clear() def _tx_loop(self) -> None: while not self._stop.is_set(): # Always flush coalesced mouse first — keeps cursor responsive under USB load. self._flush_input_moves_usb() try: frame, timeout_ms, done, result, error = self._tx_q.get(timeout=0.004) except queue.Empty: continue # Empty frame is a flush barrier from flush_input_moves() if not frame: self._flush_input_moves_usb() if done is not None: done.set() continue # Another move may have arrived while waiting on the queue self._flush_input_moves_usb() try: n = self.link.write_data(frame, timeout_ms=timeout_ms) if result is not None: result[0] = n except Exception as exc: if error is not None: error[0] = exc if done is not None and self._is_link_gone_error(exc): self._notify_link_broken() # Retry only for wait=True (critical) frames. msg = str(exc).lower() if done is not None and ("timed out" in msg or "timeout" in msg): try: self.link.clear_data_halt() except Exception: pass time.sleep(0.08) try: n = self.link.write_data( frame, timeout_ms=max(timeout_ms, 1200), ) if result is not None: result[0] = n if error is not None: error[0] = None except Exception as exc2: if error is not None: error[0] = exc2 finally: if done is not None: done.set() def rpc(self, op: str, timeout: float = 30.0, **kwargs: Any) -> Any: self._rpc_busy.set() last_exc: Exception | None = None # listdir RESP can be large and Win listing can take a few seconds. per_try = max(12.0, timeout / 2) try: for attempt in range(4): # Give RX time to answer peer RPCs before we push ours (half-duplex). # Also wait out a peer that is mid-replying to us / still browsing. settle = 0.4 + 0.5 * attempt time.sleep(settle) try: self.link.drain(0.2) except Exception: pass 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 try: self.send( MsgType.RPC_REQ, json.dumps(body, ensure_ascii=False).encode("utf-8"), timeout_ms=2500, ) return fut.result(timeout=per_try) except Exception as exc: last_exc = exc msg = str(exc).lower() if "timed out" in msg or "timeout" in msg: try: self.link.clear_data_halt() except Exception: pass continue raise finally: with self._rpc_lock: self._rpc_waiters.pop(req_id, None) raise TimeoutError( "상대 폴더 목록을 받지 못했습니다 (USB 타임아웃). " "상대 GUI가 연결 상태인지 확인 후 다시 시도하세요." ) from last_exc finally: with self._rpc_lock: busy = bool(self._rpc_waiters) if not busy: self._rpc_busy.clear() def _rx_loop(self) -> None: while not self._stop.is_set(): try: chunk = self.link.read_data(RX_READ_SIZE, timeout_ms=200) except Exception: if self._stop.is_set(): break continue 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(): # Pre-handshake: keep poking the cable bridge so the peer can detect us. if not self._connected.is_set(): try: wake = getattr(self.link, "wake_bridge", None) if callable(wake): wake() except Exception: pass self._stop.wait(0.5) continue # Pause keepalives during file bulk transfer OR active KM share # (PING steals OUT bandwidth and makes the remote cursor stutter). if ( not self._rpc_busy.is_set() and not self._xfer_busy.is_set() and not self._km_busy.is_set() ): try: self.link.write_ctrl( b"\x00ALIVE" + struct.pack(" bool: """Record peer name and signal first-time connection. Returns True if new.""" first = not self._connected.is_set() self._peer_hello = name if log_msg and first: self.log(log_msg) if not self._connected.is_set(): self._connected.set() if first: cb = self.on_connected if cb is not None: try: threading.Thread( target=cb, args=(name,), name="juc500-on-connected", daemon=True ).start() except Exception: pass return first def _on_hello(self, frame: Frame) -> None: name = frame.payload.decode("utf-8", errors="replace") first = self._mark_connected(name, log_msg=f"상대 HELLO: {name}") if first: self.send( MsgType.HELLO_ACK, self.peer_name.encode("utf-8"), timeout_ms=800, wait=False, ) def _send_hello_ack(self) -> None: try: self.send( MsgType.HELLO_ACK, self.peer_name.encode("utf-8"), timeout_ms=800, ) except Exception as exc: self.log(f"HELLO_ACK 송신 재시도 예정: {exc}") def _on_hello_ack(self, frame: Frame) -> None: name = frame.payload.decode("utf-8", errors="replace") self._mark_connected(name, log_msg=f"상대 HELLO_ACK: {name}") def _on_ping(self, frame: Frame) -> None: # Peer may already show "connected" and only send keepalive PINGs. if not self._connected.is_set(): self._mark_connected( self._peer_hello or "peer", log_msg="상대 PING 수신 — 연결 동기화", ) try: self.send( MsgType.HELLO_ACK, self.peer_name.encode("utf-8"), timeout_ms=800, wait=False, ) except Exception: pass self.send(MsgType.PONG, frame.payload, timeout_ms=800, wait=False) def _safe_send(self, msg_type: MsgType, payload: bytes = b"") -> None: try: self.send(msg_type, payload, timeout_ms=800) except Exception: pass def _on_file_accept(self, _frame: Frame) -> None: self._accept_result = True self._accept_error = "" self._pending_accept.set() def _on_file_reject(self, frame: Frame) -> None: self._accept_result = False self._accept_error = frame.payload.decode("utf-8", errors="replace").strip() self._pending_accept.set() def _on_file_ack(self, _frame: Frame) -> None: self._file_ack.set() def _safe_under_root(self, rel: str) -> Path: """Resolve browse/transfer path. Empty → local_root; absolute paths may leave home.""" raw = (rel or "").replace("\\", "/").strip() if not raw or raw in (".", "./"): return self.local_root.resolve() if raw.startswith("/") or raw.startswith("//") or (len(raw) >= 2 and raw[1] == ":"): return Path(raw).resolve() return (self.local_root / raw.lstrip("/")).resolve() @staticmethod def _path_to_api(path: Path) -> str: s = path.resolve().as_posix() if len(s) == 2 and s[1] == ":": s += "/" return s def _resolve_inbound_path(self, name: str) -> Path: """Map FILE_OFFER name to a local write path without double-home nesting. Senders often pass the remote panel path after ``strip('/')``, e.g. ``Users/macbook/util/a.txt`` while ``inbound_dest`` is ``/Users/macbook``. Absolute POSIX/Windows paths are honored when they are explicit dests. """ raw = (name or "").replace("\\", "/").strip() if not raw or raw in (".", "./"): raise ValueError("empty name") parts = [p for p in raw.split("/") if p and p != ".."] if not parts or ".." in raw.split("/"): raise ValueError("bad name") dest_root = Path(self.inbound_dest).resolve() is_win_abs = len(raw) >= 2 and raw[1] == ":" is_posix_abs = raw.startswith("/") if is_win_abs: if os.name == "nt": out = Path(raw).resolve() else: # Windows path offered to a non-Windows peer — keep basename only. return (dest_root / parts[-1]).resolve() elif is_posix_abs: out = Path(raw).resolve() else: dest_parts = [p for p in dest_root.parts if p not in ("/", "\\")] # Case-insensitive strip of duplicated home prefix (Users/macbook/...). if len(parts) >= len(dest_parts) and [ p.lower() for p in parts[: len(dest_parts)] ] == [p.lower() for p in dest_parts]: parts = parts[len(dest_parts) :] out = dest_root.joinpath(*parts) if parts else dest_root / "download" out = out.resolve() home = Path.home().resolve() for root in (dest_root, home): try: out.relative_to(root) return out except ValueError: continue # Explicit Windows absolute outside home (e.g. C:/Temp/...) if is_win_abs and os.name == "nt": return out return (dest_root / parts[-1]).resolve() def list_local(self, rel: str = "") -> dict[str, Any]: """List a directory for browse RPC. Windows OneDrive/junction entries can block for a long time on ``stat()`` / ``is_dir()``. Bound per-entry latency so peer RPC (Mac→Win listdir) does not time out while Win→Mac still works. """ raw = (rel or "").replace("\\", "/").strip() if raw == ROOTS_PATH: return list_fs_roots() path = self._safe_under_root(rel) if not path.exists(): raise FileNotFoundError(rel) if not path.is_dir(): raise NotADirectoryError(rel) deadline = time.monotonic() + 6.0 max_entries = 1500 entries: list[dict[str, Any]] = [] def _meta(entry: os.DirEntry) -> dict[str, Any] | None: try: is_dir = entry.is_dir(follow_symlinks=False) st = entry.stat(follow_symlinks=False) return { "name": entry.name, "is_dir": is_dir, "size": 0 if is_dir else int(st.st_size), "mtime": int(st.st_mtime), } except OSError: return None try: with os.scandir(path) as it: for entry in it: if len(entries) >= max_entries or time.monotonic() >= deadline: break box: list[dict[str, Any] | None] = [] def _run(e: os.DirEntry = entry, out: list = box) -> None: out.append(_meta(e)) t = threading.Thread(target=_run, daemon=True) t.start() t.join(timeout=0.4) if t.is_alive() or not box or box[0] is None: continue entries.append(box[0]) except OSError as exc: raise FileNotFoundError(f"{rel}: {exc}") from exc entries.sort(key=lambda e: (not e["is_dir"], e["name"].lower())) return {"path": self._path_to_api(path), "entries": entries} def mkdir_local(self, rel: str) -> dict[str, Any]: raw = (rel or "").replace("\\", "/").strip() if raw == ROOTS_PATH: raise ValueError("드라이브 목록에는 폴더를 만들 수 없습니다") path = self._safe_under_root(rel) path.mkdir(parents=True, exist_ok=True) return {"path": self._path_to_api(path)} def _on_rpc_req(self, frame: Frame) -> None: threading.Thread( target=self._handle_rpc_req, args=(frame,), daemon=True, name="juc500-rpc-reply", ).start() def _handle_rpc_req(self, frame: Frame) -> None: # Pause keepalives while we build/send a reply (half-duplex). self._rpc_busy.set() try: 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": try: from .inputshare import INPUTSHARE_BUILD build = INPUTSHARE_BUILD except Exception: build = "unknown" try: from .usb_link import XFER_BUILD xfer_build = XFER_BUILD except Exception: xfer_build = "unknown" result = { "pong": True, "name": self.peer_name, "inputshare_build": build, "xfer_build": xfer_build, } elif op == "pull": rel = req.get("path") or "" 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 elif op == "restart_gui": cb = self.on_restart if cb is not None: threading.Thread( target=cb, name="juc500-restart-gui", daemon=True, ).start() result = {"scheduled": True} 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)}) finally: with self._rpc_lock: busy = bool(self._rpc_waiters) if not busy: self._rpc_busy.clear() def _rpc_reply(self, body: dict[str, Any]) -> None: payload = json.dumps(body, ensure_ascii=False).encode("utf-8") # Scale OUT timeout with frame size (dir listings are often tens of KB). packets = max(1, (len(payload) + HEADER_SIZE + CRC_SIZE + MAX_PACKET - 1) // MAX_PACKET) timeout_ms = min(30000, max(4000, 1500 + packets * 8)) last: Exception | None = None for attempt in range(4): try: self.send(MsgType.RPC_RESP, payload, timeout_ms=timeout_ms) return except Exception as exc: last = exc try: self.link.clear_data_halt() except Exception: pass time.sleep(0.15 * (attempt + 1)) self.log(f"RPC_RESP 송신 실패 ({len(payload)} bytes): {last}") 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")) # --- input share --- def send_input_ctrl(self, target: int, width: int = 0, height: int = 0) -> None: self.flush_input_moves() self.send( MsgType.INPUT_CTRL, pack_input_ctrl(target, width, height), timeout_ms=500, wait=True, ) def send_input_move(self, x: int, y: int) -> None: """Queue latest absolute cursor position for peer (usbLAN-style).""" with self._move_acc_lock: self._move_acc_x = int(x) self._move_acc_y = int(y) def _flush_input_moves_usb(self) -> None: """Write latest absolute INPUT_MOVE (TX thread only — serializes USB OUT).""" with self._move_acc_lock: if self._move_acc_x is None or self._move_acc_y is None: return x, y = self._move_acc_x, self._move_acc_y self._move_acc_x = self._move_acc_y = None with self._lock: self._seq = (self._seq + 1) & 0xFFFFFFFF seq = self._seq frame = pack_frame(MsgType.INPUT_MOVE, pack_input_move(x, y), seq=seq) try: # Short timeout: prefer drop-and-continue over waiting for a stalled OUT self.link.write_data(frame, timeout_ms=25) except Exception: # Keep newest motion only — re-queue lost sample makes lag worse pass def flush_input_moves(self) -> None: """Ask TX thread to drain coalesced moves before click/key/ctrl. Never write USB from caller threads — half-duplex OUT must stay serialized. """ done = threading.Event() try: self._tx_q.put((b"", 0, done, None, None)) except Exception: return done.wait(timeout=0.08) def flush_input_moves_blocking(self, timeout_s: float = 0.15) -> None: """Drain coalesced moves before click/key — longer wait for half-duplex USB.""" done = threading.Event() try: self._tx_q.put((b"", 0, done, None, None)) except Exception: return done.wait(timeout=timeout_s) def send_input_button( self, button: int, pressed: bool, x: int = -1, y: int = -1 ) -> None: """Flush latest cursor position, then send click (reliable USB delivery).""" self.flush_input_moves_blocking(0.12) self.send( MsgType.INPUT_BUTTON, pack_input_button(button, pressed, x, y), timeout_ms=200, wait=True, ) def send_input_wheel(self, dx: int, dy: int) -> None: self.flush_input_moves() self.send(MsgType.INPUT_WHEEL, pack_input_wheel(dx, dy), timeout_ms=25, wait=False) def send_input_key(self, name: str, pressed: bool) -> None: self.flush_input_moves() self.send(MsgType.INPUT_KEY, pack_input_key(name, pressed), timeout_ms=50, wait=False) def send_input_clipboard(self, text: str) -> None: """Push plain-text clipboard to peer so Ctrl/Cmd+V pastes there.""" raw = (text or "").encode("utf-8") if not raw: return # Keep USB frames reasonable; huge pastes fall back to peer's existing clipboard if len(raw) > 200_000: raw = raw[:200_000] self.flush_input_moves() self.send(MsgType.INPUT_CLIPBOARD, raw, timeout_ms=200, wait=False) def _on_input_ctrl(self, frame: Frame) -> None: if self.on_input_ctrl: self.on_input_ctrl(*unpack_input_ctrl(frame.payload)) def _on_input_move(self, frame: Frame) -> None: if self.on_input_move: self.on_input_move(*unpack_input_move(frame.payload)) def _on_input_button(self, frame: Frame) -> None: if self.on_input_button: args = unpack_input_button(frame.payload) try: self.on_input_button(*args) except TypeError: # Legacy handlers: (button, pressed) only self.on_input_button(args[0], args[1]) def _on_input_wheel(self, frame: Frame) -> None: if self.on_input_wheel: self.on_input_wheel(*unpack_input_wheel(frame.payload)) def _on_input_key(self, frame: Frame) -> None: if self.on_input_key: self.on_input_key(*unpack_input_key(frame.payload)) def _on_input_clipboard(self, frame: Frame) -> None: if self.on_input_clipboard: try: text = frame.payload.decode("utf-8", errors="replace") except Exception: return self.on_input_clipboard(text) 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(): 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 # Stream mode: digest "-" → receiver skips per-chunk ACK (USB backpressure only). meta = f"{name}\n{size}\n-".encode("utf-8") self._xfer_busy.set() # Scale USB OUT budget with chunk size (slow bridges need >8s for 256KiB+). chunk_to = max(12000, 3000 + (CHUNK_PAYLOAD // 512) * 10) inflight: list[tuple[threading.Event, list[Any], int]] = [] def drain_inflight(*, all_left: bool = False) -> None: limit = 0 if all_left else TX_WINDOW - 1 while len(inflight) > limit: done, error, flen = inflight.pop(0) self._await_tx(done, error, flen, chunk_to) def offer_and_wait_accept() -> None: last_reject = "" for attempt in range(3): self._pending_accept.clear() self._accept_result = None self._accept_error = "" try: self.link.drain(0.15) except Exception: pass try: self.link.clear_data_halt() except Exception: pass self.log(f"파일 제안: {name} ({size} bytes)" + (f" 재시도 {attempt}" if attempt else "")) self.send(MsgType.FILE_OFFER, meta, timeout_ms=5000) # Give RX a moment after half-duplex OUT before expecting ACCEPT. try: self.link.drain(0.2) except Exception: pass per = max(20.0, accept_timeout / 3) if self._pending_accept.wait(per): if self._accept_result: return last_reject = self._accept_error or "rejected" raise RuntimeError( "상대가 파일 수신을 거부했습니다" + (f": {last_reject}" if last_reject else "") + " — 상대 경로(예: C:/dev) 쓰기 권한·디스크 공간을 확인하세요." ) try: self.link.clear_data_halt() except Exception: pass time.sleep(0.25 * (attempt + 1)) raise TimeoutError( "상대가 파일 수신을 수락하지 않았습니다. " "상대 GUI가 연결 상태인지, 대상 폴더에 쓸 수 있는지 확인하세요." ) try: offer_and_wait_accept() sha = hashlib.sha256() sent = 0 last_log = 0 t0 = time.monotonic() with path.open("rb", buffering=CHUNK_PAYLOAD) as f: while True: chunk = f.read(CHUNK_PAYLOAD) if not chunk: break sha.update(chunk) drain_inflight() done, _result, error, flen = self._enqueue_send( MsgType.FILE_CHUNK, chunk, timeout_ms=chunk_to ) inflight.append((done, error, flen)) sent += len(chunk) if progress: progress(sent, size) if self.on_progress: self.on_progress(sent, size, name) if sent - last_log >= 4 * 1024 * 1024 or sent == size: elapsed = max(time.monotonic() - t0, 1e-3) mbps = (sent / (1024 * 1024)) / elapsed pct = int(100 * sent / size) if size else 100 self.log( f"전송 중: {name} {sent}/{size} ({pct}%) ~{mbps:.1f} MB/s" ) last_log = sent drain_inflight(all_left=True) digest = sha.hexdigest() self._file_ack.clear() self.send(MsgType.FILE_DONE, digest.encode("ascii"), timeout_ms=8000) # Single end-of-transfer ACK (not per-chunk). if not self._file_ack.wait(180.0): self.log(f"전송 완료 ACK 대기 시간 초과 (로컬 송신은 끝남): {name}") elapsed = max(time.monotonic() - t0, 1e-3) mbps = (size / (1024 * 1024)) / elapsed if size else 0.0 self.log(f"전송 완료: {name} ({mbps:.1f} MB/s)") finally: self._xfer_busy.clear() # --- 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", timeout_ms=3000) return try: out = self._resolve_inbound_path(name) safe = self._path_to_api(out) except Exception as exc: self.send( MsgType.FILE_REJECT, f"bad name: {exc}".encode("utf-8")[:200], timeout_ms=3000, ) return try: out.parent.mkdir(parents=True, exist_ok=True) except Exception as exc: self.send( MsgType.FILE_REJECT, f"mkdir: {exc}".encode("utf-8")[:200], timeout_ms=3000, ) return try: fh = out.open("wb", buffering=CHUNK_PAYLOAD) except Exception as exc: self.send( MsgType.FILE_REJECT, str(exc).encode("utf-8")[:200], timeout_ms=3000, ) return dig = digest.strip() # Stream offers ("-") skip per-chunk ACK for throughput. stream = dig in ("-", "pending", "stream") self._xfer_busy.set() self._inbound = { "path": out, "name": safe, "size": size, "digest": dig, "fh": fh, "sha": hashlib.sha256(), "received": 0, "last_log": 0, "stream": stream, "t0": time.monotonic(), } # wait=True so ACCEPT is not dropped on half-duplex USB before peer listens. try: self.send(MsgType.FILE_ACCEPT, b"ok", timeout_ms=5000) except Exception as exc: try: fh.close() except Exception: pass self._inbound = {} self._xfer_busy.clear() self.log(f"FILE_ACCEPT 송신 실패: {exc}") return mode = "stream" if stream else "ack" self.log(f"수신 시작: {safe} ({size} bytes, {mode})") 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) # Legacy (non-stream) peers still need per-chunk ACK. if not st.get("stream"): self.send(MsgType.FILE_ACK, b"", timeout_ms=3000, wait=False) if self.on_progress: self.on_progress(st["received"], st["size"], st.get("name") or "") recv = st["received"] total = int(st.get("size") or 0) last = int(st.get("last_log") or 0) if recv - last >= 8 * 1024 * 1024 or (total and recv >= total): elapsed = max(time.monotonic() - float(st.get("t0") or time.monotonic()), 1e-3) mbps = (recv / (1024 * 1024)) / elapsed pct = int(100 * recv / total) if total else 0 self.log(f"수신 중: {st.get('name')} {recv}/{total} ({pct}%) ~{mbps:.1f} MB/s") st["last_log"] = recv def _on_inbound_done(self, frame: Frame) -> None: st = self._inbound fh = st.pop("fh", None) if fh: try: fh.flush() except Exception: pass fh.close() digest = st.get("sha").hexdigest() if st.get("sha") else "" expected = (st.get("digest") or "").strip() remote = frame.payload.decode("ascii", errors="replace").strip() # OFFER may send "-" when sender streams hash; trust FILE_DONE then. offer_ok = (not expected) or expected in ("-", "pending", "stream") or digest == expected done_ok = (not remote) or digest == remote # End-of-transfer ACK for stream senders (replaces per-chunk ACK). self.send(MsgType.FILE_ACK, b"", timeout_ms=2000, wait=False) if not offer_ok or not done_ok: self.log(f"해시 불일치: {st.get('name')}") else: elapsed = max(time.monotonic() - float(st.get("t0") or time.monotonic()), 1e-3) total = int(st.get("received") or st.get("size") or 0) mbps = (total / (1024 * 1024)) / elapsed if total else 0.0 self.log(f"수신 완료: {st.get('path')} ({mbps:.1f} MB/s)") self._inbound = {} self._xfer_busy.clear() # 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("파일 수신 타임아웃")