Files
juc500/dist-offline/juc500-offline-update-20260720-144424/juc500_xfer/protocol.py
T
2026-07-21 21:04:25 +09:00

1538 lines
60 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
# Peer RPC set_cable_model → switch local cable selection + reconnect
self.on_cable_switch: Optional[Callable[[str], 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
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 (Mac JUC700 IF0) — 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()
# Mac↔Win JUC700: poke bridge before first HELLO so peer IN can drain.
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
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()
# Connector: fire-and-forget HELLO — blocking OUT stalls the TX queue
# ("송신 큐 대기 시간 초과") and Access-denied reopen storms on IF0.
if now >= next_hello:
next_hello = now + 0.6
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 as exc:
msg = str(exc)
now = time.monotonic()
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_failures += 1
if self._hello_failures % 8 == 0:
try:
rot = getattr(self.link, "rotate_juc700_if0_eps", None)
if callable(rot) and rot():
lay = getattr(self.link, "layout", None)
if lay is not None:
self.log(
f"JUC700 IF0 EP 전환 → "
f"{lay.ep_data_out:#x}/{lay.ep_data_in:#x}"
)
except Exception:
pass
if now - self._hello_log_at >= 4.0:
self._hello_log_at = now
self.log(
"상대 수신 대기 중… (HELLO) — 상대에서도 이 앱 실행"
)
elif self._is_link_gone_error(exc):
self._hello_failures += 1
if now - self._hello_log_at >= 2.0:
self._hello_log_at = now
self.log(f"USB 링크 오류 (HELLO): {exc} — 재오픈")
if self._hello_failures >= 3:
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:
# Periodic EP rotate while waiting (OUT may soft-fail silently).
self._hello_failures += 1
if self._hello_failures % 10 == 0:
try:
rot = getattr(self.link, "rotate_juc700_if0_eps", None)
if callable(rot) and rot():
lay = getattr(self.link, "layout", None)
if lay is not None and now - self._hello_log_at >= 4.0:
self._hello_log_at = now
self.log(
f"JUC700 IF0 EP 전환 → "
f"{lay.ep_data_out:#x}/{lay.ep_data_in:#x}"
)
except Exception:
pass
if now - self._hello_log_at >= 4.0:
self._hello_log_at = now
self.log(
"상대 수신 대기 중… (HELLO) — 상대에서도 이 앱 실행"
)
# Every ~8s: confirm OUT actually lands (drain vs black-hole).
if self._hello_failures % 12 == 6:
try:
self.send(
MsgType.HELLO,
self.peer_name.encode("utf-8"),
timeout_ms=500,
wait=True,
)
if now - self._hello_log_at >= 3.0:
self._hello_log_at = now
self.log(
"HELLO USB 송신은 됨 — 상대 J5FX 응답 없음. "
"Windows: Wormhole 종료, Zadig MI_05→WinUSB, "
"빌드 smartkmlink2, 케이블 모델 SmartKMLink 확인"
)
except Exception:
pass
self._connected.wait(timeout=min(0.25, remaining))
return self._peer_hello or "peer"
def poke_hello(self) -> None:
"""Non-blocking HELLO for listen mode (keep IF0 claimed; avoid reopen races)."""
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
try:
self.send(
MsgType.HELLO,
self.peer_name.encode("utf-8"),
timeout_ms=200,
wait=False,
)
except Exception:
pass
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():
# Always use larger URBs — RPC_RESP (dir listing) is often multi-KB.
chunk = self.link.read_data(RX_READ_SIZE, 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():
# Pre-handshake: keep poking the cable bridge so the peer can detect
# Mac-initiated connect (JUC700 IF0↔IF5) once it opens USB.
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}
elif op == "set_cable_model":
model = str(req.get("model") or "").strip().upper()
if model not in ("JUC500", "JUC700"):
raise ValueError(f"unsupported model: {model}")
cb = self.on_cable_switch
if cb is not None:
threading.Thread(
target=cb,
args=(model,),
name="juc500-cable-switch",
daemon=True,
).start()
result = {"scheduled": True, "model": model}
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("파일 수신 타임아웃")