from core.utils.encrypted_proxy import killswitch from core.controllers.ApplicationController import ApplicationController from core.observers.ApplicationVersionObserver import ApplicationVersionObserver from core.observers.ClientObserver import ClientObserver from core.observers.ConnectionObserver import ConnectionObserver from core.observers.InvoiceObserver import InvoiceObserver from core.observers.ProfileObserver import ProfileObserver from cli.ui import ( Spinner, label_ok, label_error, label_warn, label_connected, label_disconnect, label_killswitch, label_dns, label_ipv6, label_ipv4, progress_bar, ) from cli.helpers import sanitize_profile from cli.killswitch_monitor import drop_state, tunnel_state import pprint application_version_observer = ApplicationVersionObserver() client_observer = ClientObserver() connection_observer = ConnectionObserver() invoice_observer = InvoiceObserver() profile_observer = ProfileObserver() _spinner = None _monitor = None _timer_state = None application_version_observer.subscribe('downloading', lambda event: print(f'Downloading {ApplicationController.get(event.subject.application_code).name}, version {event.subject.version_number}...')) application_version_observer.subscribe('download_progressing',lambda event: print(f'Current progress: {event.meta.get("progress"):.2f}%', flush=True, end='\r')) application_version_observer.subscribe('downloaded', lambda event: print('\n')) client_observer.subscribe('synchronizing', lambda event: print('Synchronizing...\n')) client_observer.subscribe('updating', lambda event: print('Updating client...')) client_observer.subscribe('update_progressing', lambda event: print(f'Current progress: {event.meta.get("progress"):.2f}%', flush=True, end='\r')) client_observer.subscribe('updated', lambda event: print('\n')) def _on_connecting(event): global _spinner, _timer_state if _timer_state and _timer_state['stop']: _timer_state['stop'].set() _timer_state['stop'] = None if _timer_state and _timer_state['reconnecting'] is None and _timer_state.get('_retrying'): from cli.ui import connecting_spinner _timer_state['reconnecting'] = connecting_spinner() attempt = event.subject.get("attempt_count", "?") total = event.subject.get("maximum_number_of_attempts", "?") _spinner = Spinner(f"Connecting... attempt {attempt}/{total}") _spinner.start() def _on_connected(event): global _spinner, _monitor, _timer_state if _spinner: _spinner.stop() _spinner = None if _timer_state and _timer_state.get('reconnecting'): _timer_state['reconnecting'].set() _timer_state['reconnecting'] = None if _timer_state: _timer_state['_retrying'] = False d = event.subject or {} tunnel_if = d.get('tunnel_if', '?') server_ip = d.get('server_ip') socks5_port = d.get('socks5_port', '?') label_connected(f"tunnel={tunnel_if} port={socks5_port}") label_killswitch(killswitch.status()) label_dns( enabled=d.get('dns_enabled', True), active=d.get('dns_active', False), ) label_ipv6(blocked=True) if server_ip: label_ipv4(server_ip=server_ip, reachable=True) if _monitor is not None: _monitor.set_connection_data(socks5_port=d.get('socks5_port')) if _timer_state is not None: from cli.ui import session_timer _timer_state['stop'] = session_timer(drop_state=drop_state, tunnel_state=tunnel_state) def _on_connected_token(event): global _monitor if _monitor is not None: token = (event.subject or {}).get('session_token') if token: _monitor.set_session_token(token) def _on_disconnected(event): global _spinner, _monitor, _timer_state if _spinner: _spinner.stop() _spinner = None if _timer_state and _timer_state['stop']: _timer_state['stop'].set() _timer_state['stop'] = None if _monitor is not None: _monitor.stop() _monitor = None d = event.subject or {} label_disconnect(f"tunnel={d.get('tunnel_if','?')}") def _on_error(event): global _spinner, _monitor, _timer_state if _spinner: _spinner.stop() _spinner = None if _timer_state and _timer_state['stop']: _timer_state['stop'].set() _timer_state['stop'] = None if _monitor is not None: _monitor.stop() _monitor = None label_error(str(event.subject or "Unknown error")) def _on_tor_bootstrapping(event): global _spinner _spinner = Spinner("Starting Tor...") _spinner.start() def _on_tor_bootstrap_progressing(event): pct = int(event.meta.get("progress", 0)) progress_bar(pct, 100, "Tor bootstrap") def _on_tor_bootstrapped(event): global _spinner if _spinner: _spinner.stop() _spinner = None print() label_ok("Tor ready") connection_observer.subscribe('connecting', _on_connecting) connection_observer.subscribe('connected', _on_connected) connection_observer.subscribe('connected_token', _on_connected_token) connection_observer.subscribe('disconnected', _on_disconnected) connection_observer.subscribe('error', _on_error) connection_observer.subscribe('tor_bootstrapping', _on_tor_bootstrapping) connection_observer.subscribe('tor_bootstrap_progressing', _on_tor_bootstrap_progressing) connection_observer.subscribe('tor_bootstrapped', _on_tor_bootstrapped) invoice_observer.subscribe('retrieved', lambda event: print(f'\n{pprint.pp(event.subject)}\n')) invoice_observer.subscribe('processing', lambda event: print('A payment has been detected and is being verified...\n')) invoice_observer.subscribe('settled', lambda event: print('The payment has been successfully verified.\n')) profile_observer.subscribe('created', lambda event: pprint.pp((sanitize_profile(event.subject), 'Created'))) profile_observer.subscribe('destroyed', lambda event: pprint.pp((sanitize_profile(event.subject), 'Destroyed'))) profile_observer.subscribe('disabled', lambda event: pprint.pp((sanitize_profile(event.subject), 'Disabled')) if event.meta.get('explicitly') else None) profile_observer.subscribe('enabled', lambda event: pprint.pp((sanitize_profile(event.subject), 'Enabled')))