|
|
|
@@ -71,7 +71,6 @@ _retry_timer: retry.RetryThread | None = None
|
|
|
|
|
_destination: RNS.Destination | None = None
|
|
|
|
|
_loop: asyncio.AbstractEventLoop | None = None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _check_finished(timeout: float = 0):
|
|
|
|
|
return _finished is not None and await process.event_wait(_finished, timeout=timeout)
|
|
|
|
|
|
|
|
|
@@ -92,13 +91,11 @@ async def _spin_tty(until=None, msg=None, timeout=None):
|
|
|
|
|
print(("\b\b"+syms[i]+" "), end="")
|
|
|
|
|
sys.stdout.flush()
|
|
|
|
|
i = (i+1)%len(syms)
|
|
|
|
|
|
|
|
|
|
print("\r"+" "*len(msg)+" \r", end="")
|
|
|
|
|
|
|
|
|
|
if timeout != None and time.time() > timeout: return False
|
|
|
|
|
else: return True
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _spin_pipe(until: callable = None, msg=None, timeout: float | None = None) -> bool:
|
|
|
|
|
if timeout is not None: timeout += time.time()
|
|
|
|
|
|
|
|
|
@@ -136,14 +133,11 @@ def compute_target_rns_loglevel(verbosity: int, quietness: int, base_level: int
|
|
|
|
|
if target < RNS.LOG_NONE: target = RNS.LOG_NONE
|
|
|
|
|
if target > RNS.LOG_EXTREME: target = RNS.LOG_EXTREME
|
|
|
|
|
return target
|
|
|
|
|
|
|
|
|
|
except Exception: return base_level
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class RemoteExecutionError(Exception):
|
|
|
|
|
def __init__(self, msg): self.msg = msg
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _initiate_link(configdir, rnsconfigdir, identitypath=None, logfile=None, verbosity=0, quietness=0, noid=False, destination=None,
|
|
|
|
|
timeout=RNS.Transport.PATH_REQUEST_TIMEOUT):
|
|
|
|
|
global _identity, _reticulum, _link, _destination, _remote_exec_grace
|
|
|
|
@@ -153,13 +147,12 @@ async def _initiate_link(configdir, rnsconfigdir, identitypath=None, logfile=Non
|
|
|
|
|
raise RemoteExecutionError(
|
|
|
|
|
"Allowed destination length is invalid, must be {hex} hexadecimal characters ({byte} bytes).".format(
|
|
|
|
|
hex=dest_len, byte=dest_len // 2))
|
|
|
|
|
try:
|
|
|
|
|
destination_hash = bytes.fromhex(destination)
|
|
|
|
|
except Exception as e:
|
|
|
|
|
raise RemoteExecutionError("Invalid destination entered. Check your input.")
|
|
|
|
|
|
|
|
|
|
try: destination_hash = bytes.fromhex(destination)
|
|
|
|
|
except Exception as e: raise RemoteExecutionError("Invalid destination entered. Check your input.")
|
|
|
|
|
|
|
|
|
|
if _reticulum is None:
|
|
|
|
|
targetloglevel = compute_target_rns_loglevel(verbosity, quietness, RNS.LOG_ERROR)
|
|
|
|
|
targetloglevel = compute_target_rns_loglevel(verbosity, quietness, RNS.LOG_INFO)
|
|
|
|
|
RNS.logfile = logfile
|
|
|
|
|
_reticulum = RNS.Reticulum(configdir=rnsconfigdir, loglevel=targetloglevel, logdest=RNS.LOG_FILE)
|
|
|
|
|
|
|
|
|
@@ -168,19 +161,14 @@ async def _initiate_link(configdir, rnsconfigdir, identitypath=None, logfile=Non
|
|
|
|
|
|
|
|
|
|
if not RNS.Transport.has_path(destination_hash):
|
|
|
|
|
RNS.Transport.request_path(destination_hash)
|
|
|
|
|
RNS.log(f"Requesting path...", RNS.LOG_INFO)
|
|
|
|
|
if not await _spin(until=lambda: RNS.Transport.has_path(destination_hash), msg="Requesting path...",
|
|
|
|
|
RNS.log(f"Requesting path to {RNS.prettyhexrep(destination_hash)}", RNS.LOG_INFO)
|
|
|
|
|
if not await _spin(until=lambda: RNS.Transport.has_path(destination_hash), msg=f"Requesting path to {RNS.prettyhexrep(destination_hash)}",
|
|
|
|
|
timeout=timeout, quiet=quietness > 0):
|
|
|
|
|
raise RemoteExecutionError("Path not found")
|
|
|
|
|
|
|
|
|
|
if _destination is None:
|
|
|
|
|
listener_identity = RNS.Identity.recall(destination_hash)
|
|
|
|
|
_destination = RNS.Destination(
|
|
|
|
|
listener_identity,
|
|
|
|
|
RNS.Destination.OUT,
|
|
|
|
|
RNS.Destination.SINGLE,
|
|
|
|
|
rnsh.APP_NAME
|
|
|
|
|
)
|
|
|
|
|
_destination = RNS.Destination(listener_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, rnsh.APP_NAME)
|
|
|
|
|
|
|
|
|
|
if _link is None or _link.status == RNS.Link.PENDING:
|
|
|
|
|
RNS.log("No link", RNS.LOG_DEBUG)
|
|
|
|
@@ -189,34 +177,30 @@ async def _initiate_link(configdir, rnsconfigdir, identitypath=None, logfile=Non
|
|
|
|
|
|
|
|
|
|
_link.set_link_closed_callback(_client_link_closed)
|
|
|
|
|
|
|
|
|
|
RNS.log(f"Establishing link...", RNS.LOG_VERBOSE)
|
|
|
|
|
if not await _spin(until=lambda: _link.status == RNS.Link.ACTIVE, msg="Establishing link...",
|
|
|
|
|
RNS.log(f"Establishing link with {RNS.prettyhexrep(destination_hash)}", RNS.LOG_INFO)
|
|
|
|
|
if not await _spin(until=lambda: _link.status == RNS.Link.ACTIVE, msg=f"Establishing link with {RNS.prettyhexrep(destination_hash)}",
|
|
|
|
|
timeout=timeout, quiet=quietness > 0):
|
|
|
|
|
raise RemoteExecutionError("Could not establish link with " + RNS.prettyhexrep(destination_hash))
|
|
|
|
|
|
|
|
|
|
RNS.log("Have link", RNS.LOG_DEBUG)
|
|
|
|
|
RNS.log(f"Link established with {RNS.prettyhexrep(destination_hash)}", RNS.LOG_INFO)
|
|
|
|
|
if not noid and not _link.did_identify:
|
|
|
|
|
# Delay a tiny bit to allow listener to fully enter WAIT_IDENT state
|
|
|
|
|
await asyncio.sleep(min(1, _link.rtt * 1.1 + 0.05))
|
|
|
|
|
_link.identify(_identity)
|
|
|
|
|
_link.did_identify = True
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _handle_error(errmsg: RNS.MessageBase):
|
|
|
|
|
if isinstance(errmsg, protocol.ErrorMessage):
|
|
|
|
|
with contextlib.suppress(Exception):
|
|
|
|
|
if _link and _link.status == RNS.Link.ACTIVE:
|
|
|
|
|
_link.teardown()
|
|
|
|
|
if _link and _link.status == RNS.Link.ACTIVE: _link.teardown()
|
|
|
|
|
await asyncio.sleep(0.1)
|
|
|
|
|
raise RemoteExecutionError(f"Remote error: {errmsg.msg}")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:str, verbosity: int, quietness: int, noid: bool, destination: str,
|
|
|
|
|
timeout: float, command: [str] | None = None):
|
|
|
|
|
|
|
|
|
|
global _finished, _link
|
|
|
|
|
if timeout is None:
|
|
|
|
|
timeout = RNS.Transport.PATH_REQUEST_TIMEOUT
|
|
|
|
|
if timeout is None: timeout = RNS.Transport.PATH_REQUEST_TIMEOUT
|
|
|
|
|
with process.TTYRestorer(sys.stdin.fileno()) as ttyRestorer:
|
|
|
|
|
loop = asyncio.get_running_loop()
|
|
|
|
|
state = InitiatorState.IS_INITIAL
|
|
|
|
@@ -233,8 +217,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
destination=destination,
|
|
|
|
|
timeout=timeout)
|
|
|
|
|
|
|
|
|
|
if not _link or _link.status not in [RNS.Link.ACTIVE, RNS.Link.PENDING]:
|
|
|
|
|
return 255
|
|
|
|
|
if not _link or _link.status not in [RNS.Link.ACTIVE, RNS.Link.PENDING]: return 255
|
|
|
|
|
|
|
|
|
|
state = InitiatorState.IS_LINKED
|
|
|
|
|
outlet = session.RNSOutlet(_link)
|
|
|
|
@@ -245,9 +228,8 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
try:
|
|
|
|
|
vm = _pq.get(timeout=max(outlet.rtt * 20, 5))
|
|
|
|
|
await _handle_error(vm)
|
|
|
|
|
if not isinstance(vm, protocol.VersionInfoMessage):
|
|
|
|
|
raise Exception("Invalid message received")
|
|
|
|
|
RNS.log(f"Server version info: sw {vm.sw_version} prot {vm.protocol_version}", RNS.LOG_DEBUG)
|
|
|
|
|
if not isinstance(vm, protocol.VersionInfoMessage): raise Exception("Invalid message received")
|
|
|
|
|
RNS.log(f"Connected server version info: sw {vm.sw_version}, proto {vm.protocol_version}", RNS.LOG_INFO)
|
|
|
|
|
state = InitiatorState.IS_RUNNING
|
|
|
|
|
except queue.Empty:
|
|
|
|
|
print("Protocol error")
|
|
|
|
@@ -323,8 +305,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
pre_esc = False
|
|
|
|
|
data.append(b)
|
|
|
|
|
|
|
|
|
|
if not line_mode:
|
|
|
|
|
data_buffer.extend(data)
|
|
|
|
|
if not line_mode: data_buffer.extend(data)
|
|
|
|
|
else:
|
|
|
|
|
line_buffer.extend(data)
|
|
|
|
|
if line_flush:
|
|
|
|
@@ -338,8 +319,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
blind_write_count += len(data)
|
|
|
|
|
|
|
|
|
|
except EOFError:
|
|
|
|
|
if os.isatty(0):
|
|
|
|
|
data_buffer.extend(process.CTRL_D)
|
|
|
|
|
if os.isatty(0): data_buffer.extend(process.CTRL_D)
|
|
|
|
|
stdin_eof = True
|
|
|
|
|
process.tty_unset_reader_callbacks(sys.stdin.fileno())
|
|
|
|
|
|
|
|
|
@@ -414,8 +394,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
_link.teardown()
|
|
|
|
|
return 200
|
|
|
|
|
|
|
|
|
|
except queue.Empty:
|
|
|
|
|
processed = False
|
|
|
|
|
except queue.Empty: processed = False
|
|
|
|
|
|
|
|
|
|
if channel.is_ready_to_send():
|
|
|
|
|
def compress_adaptive(buf: bytes):
|
|
|
|
@@ -437,8 +416,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
if compressed_length < max_data_len and compressed_length < chunk_segment_length:
|
|
|
|
|
comp_success = True
|
|
|
|
|
break
|
|
|
|
|
else:
|
|
|
|
|
comp_try += 1
|
|
|
|
|
else: comp_try += 1
|
|
|
|
|
|
|
|
|
|
if comp_success:
|
|
|
|
|
diff = max_data_len - len(compressed_chunk)
|
|
|
|
@@ -467,6 +445,7 @@ async def initiate(configdir: str, rnsconfigdir:str, identitypath: str, logfile:
|
|
|
|
|
r, c, h, v = process.tty_get_winsize(0)
|
|
|
|
|
channel.send(protocol.WindowSizeMessage(r, c, h, v))
|
|
|
|
|
processed = True
|
|
|
|
|
|
|
|
|
|
except RemoteExecutionError as e:
|
|
|
|
|
print(e.msg)
|
|
|
|
|
return 255
|
|
|
|
|