1506 lines
57 KiB
Python
1506 lines
57 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 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("<BHH", target & 0xFF, width & 0xFFFF, height & 0xFFFF)
|
|
|
|
|
|
def unpack_input_ctrl(payload: bytes) -> tuple[int, int, int]:
|
|
if len(payload) < 5:
|
|
return INPUT_CTRL_LOCAL, 0, 0
|
|
target, w, h = struct.unpack_from("<BHH", payload, 0)
|
|
return int(target), int(w), int(h)
|
|
|
|
|
|
def pack_input_move(x: int, y: int) -> bytes:
|
|
"""Absolute virtual cursor position on peer screen (i32 x, i32 y)."""
|
|
return struct.pack(
|
|
"<ii",
|
|
max(-2_147_483_648, min(2_147_483_647, int(x))),
|
|
max(-2_147_483_648, min(2_147_483_647, int(y))),
|
|
)
|
|
|
|
|
|
def unpack_input_move(payload: bytes) -> tuple[int, int]:
|
|
if len(payload) >= 8:
|
|
return struct.unpack_from("<ii", payload, 0)
|
|
if len(payload) >= 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(
|
|
"<BBii", button & 0xFF, 1 if pressed else 0, int(x), int(y)
|
|
)
|
|
return struct.pack("<BB", button & 0xFF, 1 if pressed else 0)
|
|
|
|
|
|
def unpack_input_button(payload: bytes) -> tuple[int, bool, int, int]:
|
|
if len(payload) < 2:
|
|
return 1, False, -1, -1
|
|
button, pressed = struct.unpack_from("<BB", payload, 0)
|
|
if len(payload) >= 10:
|
|
x, y = struct.unpack_from("<ii", payload, 2)
|
|
return int(button), bool(pressed), int(x), int(y)
|
|
return int(button), bool(pressed), -1, -1
|
|
|
|
|
|
def pack_input_wheel(dx: int, dy: int) -> bytes:
|
|
return struct.pack("<hh", max(-32768, min(32767, dx)), max(-32768, min(32767, dy)))
|
|
|
|
|
|
def unpack_input_wheel(payload: bytes) -> tuple[int, int]:
|
|
if len(payload) < 4:
|
|
return 0, 0
|
|
return struct.unpack_from("<hh", payload, 0)
|
|
|
|
|
|
def pack_input_key(name: str, pressed: bool) -> 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("<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}"
|
|
|
|
|
|
# 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("<d", time.time()),
|
|
timeout_ms=300,
|
|
)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
self.link.read_ctrl(MAX_PACKET, timeout_ms=50)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
self.send(
|
|
MsgType.PING,
|
|
struct.pack("<d", time.time()),
|
|
timeout_ms=400,
|
|
wait=False,
|
|
)
|
|
except Exception:
|
|
pass
|
|
# Poll more often while KM busy so PING resumes quickly after release
|
|
self._stop.wait(0.8 if self._km_busy.is_set() else 3.0)
|
|
|
|
def _mark_connected(self, name: str, *, log_msg: str | None = None) -> 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("파일 수신 타임아웃")
|