sp-hydra-veil-cli/cli/observers.py
2026-06-09 09:07:44 -05:00

134 lines
No EOL
5.6 KiB
Python

from core.utils.encrypted_proxy.net import get_proxied_ip
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, progress_bar, label_killswitch
from cli.helpers import sanitize_profile
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"[net] verify_ip 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 {}
label_connected(f"proxy={d.get('proxy_ip','?')} real={d.get('real_ip','?')} port={d.get('socks5_port','?')}")
label_killswitch(killswitch.status())
if _monitor is not None:
_monitor.set_connection_data(
socks5_port=d.get('socks5_port'),
real_ip=d.get('real_ip')
)
if _timer_state is not None:
from cli.ui import session_timer
_timer_state['stop'] = session_timer()
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"real={d.get('ip','?')}")
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('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')))