loggorilla implementation

This commit is contained in:
Dita Aji Pratama 2026-09-26 19:21:36 +07:00
parent 19447fdef7
commit d84328d7a1
3 changed files with 57 additions and 49 deletions

13
scripts/loggorilla.py Normal file
View File

@ -0,0 +1,13 @@
import datetime
def prcss(loc, msg):
print(f"[loggorilla][{datetime.datetime.now()}][\033[32mprcss\033[39m][\033[95m{loc}\033[39m] {msg}", flush=True)
def accss(loc, msg):
print(f"[loggorilla][{datetime.datetime.now()}][\033[36maccss\033[39m][\033[95m{loc}\033[39m] {msg}", flush=True)
def fyinf(loc, msg):
print(f"[loggorilla][{datetime.datetime.now()}][\033[93mfyinf\033[39m][\033[95m{loc}\033[39m] {msg}", flush=True)
def error(loc, msg):
print(f"[loggorilla][{datetime.datetime.now()}][\033[31merror\033[39m][\033[95m{loc}\033[39m] {msg}", flush=True)

View File

@ -10,6 +10,7 @@ from services.session_manager import SessionManager
from lib.agent_loop import run_agent_loop from lib.agent_loop import run_agent_loop
from lib import personality from lib import personality
from lib import ragroleplay from lib import ragroleplay
from scripts import loggorilla
from tools.roleplayer import _name_mentioned from tools.roleplayer import _name_mentioned
@ -47,7 +48,7 @@ class TelegramClient:
self._app = Application.builder().token(token).build() self._app = Application.builder().token(token).build()
def start(self): def start(self):
print(f'[{_ts()}] Starting Telegram service...', flush=True) loggorilla.prcss('telegram', 'Starting Telegram service...')
asyncio.run(self._async_run()) asyncio.run(self._async_run())
def stop(self): def stop(self):
@ -63,14 +64,14 @@ class TelegramClient:
bot_user = await self._app.bot.get_me() bot_user = await self._app.bot.get_me()
self._bot_username = bot_user.username or "" self._bot_username = bot_user.username or ""
print(f'[{_ts()}] Telegram bot: @{self._bot_username}', flush=True) loggorilla.prcss('telegram', f'Telegram bot: @{self._bot_username}')
self._register_handlers() self._register_handlers()
await self._app.initialize() await self._app.initialize()
await self._app.start() await self._app.start()
await self._app.updater.start_polling() await self._app.updater.start_polling()
print(f'[{_ts()}] Telegram bot is polling', flush=True) loggorilla.prcss('telegram', 'Telegram bot is polling')
try: try:
await self._stopped.wait() await self._stopped.wait()
@ -80,15 +81,15 @@ class TelegramClient:
await self._app.updater.stop() await self._app.updater.stop()
await self._app.stop() await self._app.stop()
await self._app.shutdown() await self._app.shutdown()
print(f'[{_ts()}] Telegram service stopped', flush=True) loggorilla.prcss('telegram', 'Telegram service stopped')
async def _error_handler(self, update, context): async def _error_handler(self, update, context):
from telegram.error import NetworkError from telegram.error import NetworkError
cause = context.error cause = context.error
if isinstance(cause, NetworkError) and "ReadError" in str(cause): if isinstance(cause, NetworkError) and "ReadError" in str(cause):
print(f'[{_ts()}] [Network Warning] Telegram ReadError occurred. Retrying...', flush=True) loggorilla.fyinf('telegram', 'Network Warning: Telegram ReadError occurred. Retrying...')
else: else:
print(f'[{_ts()}] Exception while handling an update: {cause}', flush=True) loggorilla.error('telegram', f'Exception while handling an update: {cause}')
# Untuk error selain ReadError, tetap cetak full traceback jika perlu via logging # Untuk error selain ReadError, tetap cetak full traceback jika perlu via logging
# atau biarkan default handler mengurusnya. # atau biarkan default handler mengurusnya.
@ -131,7 +132,7 @@ class TelegramClient:
msg_id = update.message.message_id msg_id = update.message.message_id
user = update.effective_user user = update.effective_user
tg_username = user.username if user else None tg_username = user.username if user else None
print(f'[{_ts()}] Telegram DM from {chat_id}: {text[:60]}', flush=True) loggorilla.accss('telegram', f'Telegram DM from {chat_id}: {text}')
threading.Thread( threading.Thread(
target=self._process_message, target=self._process_message,
args=(chat_id, text, 'private', '', msg_id, tg_username), args=(chat_id, text, 'private', '', msg_id, tg_username),
@ -172,7 +173,7 @@ class TelegramClient:
name_mentioned = _name_mentioned(personality.PERSONALITY.codename, text) name_mentioned = _name_mentioned(personality.PERSONALITY.codename, text)
if not replied_to_bot and not has_mention and not name_mentioned: if not replied_to_bot and not has_mention and not name_mentioned:
print(f'[{_ts()}] Telegram Group [{chat_id}] NO-REPLY: {text[:60]}', flush=True) loggorilla.fyinf('telegram', f'Telegram Group [{chat_id}] NO-REPLY: {text}')
return return
if has_mention and self._bot_username: if has_mention and self._bot_username:
@ -183,7 +184,7 @@ class TelegramClient:
msg_id = update.message.message_id msg_id = update.message.message_id
sender = update.effective_user sender = update.effective_user
sender_name = sender.full_name or sender.username or str(sender.id) if sender else str(chat_id) sender_name = sender.full_name or sender.username or str(sender.id) if sender else str(chat_id)
print(f'[{_ts()}] Telegram Group [{chat_id}] <{sender_name}>: {text[:60]}', flush=True) loggorilla.accss('telegram', f'Telegram Group [{chat_id}] <{sender_name}>: {text}')
threading.Thread( threading.Thread(
target=self._process_message, target=self._process_message,
args=(chat_id, text, 'group', sender_name, msg_id), args=(chat_id, text, 'group', sender_name, msg_id),
@ -202,7 +203,7 @@ class TelegramClient:
if body in (':new', '/new'): if body in (':new', '/new'):
self._session_mgr.reset(str(chat_id)) self._session_mgr.reset(str(chat_id))
print(f'[{_ts()}] Session reset for {chat_id}', flush=True) loggorilla.prcss('telegram', f'Session reset for {chat_id}')
self._schedule_send(chat_id, 'Memulai sesi baru. Ada yang bisa dibantu?', reply_to_msg_id) self._schedule_send(chat_id, 'Memulai sesi baru. Ada yang bisa dibantu?', reply_to_msg_id)
return return
@ -305,17 +306,17 @@ class TelegramClient:
recent_history=recent_history, recent_history=recent_history,
my_name=my_name, my_name=my_name,
): ):
print(f'[{_ts()}] need_response=True → sending response', flush=True) loggorilla.fyinf('telegram', 'need_response=True → sending response')
self._schedule_send(chat_id, final_content, reply_to_msg_id) self._schedule_send(chat_id, final_content, reply_to_msg_id)
else: else:
print(f'[{_ts()}] need_response=False → staying silent', flush=True) loggorilla.fyinf('telegram', 'need_response=False → staying silent')
else: else:
from tools.roleplayer import _name_mentioned from tools.roleplayer import _name_mentioned
if _name_mentioned(my_name, body): if _name_mentioned(my_name, body):
print(f'[{_ts()}] Name mentioned → sending response', flush=True) loggorilla.fyinf('telegram', 'Name mentioned → sending response')
self._schedule_send(chat_id, final_content, reply_to_msg_id) self._schedule_send(chat_id, final_content, reply_to_msg_id)
else: else:
print(f'[{_ts()}] Name not mentioned → staying silent', flush=True) loggorilla.fyinf('telegram', 'Name not mentioned → staying silent')
else: else:
self._schedule_send(chat_id, final_content, reply_to_msg_id) self._schedule_send(chat_id, final_content, reply_to_msg_id)
else: else:
@ -328,7 +329,7 @@ class TelegramClient:
# Natural close via tool end_session: berlaku untuk private & group # Natural close via tool end_session: berlaku untuk private & group
# roleplay. Konten penutup sudah dikirim, sekarang tutup sesi. # roleplay. Konten penutup sudah dikirim, sekarang tutup sesi.
if should_close and is_roleplay: if should_close and is_roleplay:
print(f'[{_ts()}] Natural close triggered for {chat_id}', flush=True) loggorilla.prcss('telegram', f'Natural close triggered for {chat_id}')
self._session_mgr.reset(str(chat_id)) self._session_mgr.reset(str(chat_id))
else: else:
timeout = 86400 if chat_type == 'private' else 300 timeout = 86400 if chat_type == 'private' else 300
@ -338,7 +339,7 @@ class TelegramClient:
if self._loop and not self._loop.is_closed(): if self._loop and not self._loop.is_closed():
char_count = len(text) if text else 0 char_count = len(text) if text else 0
sleep_delay = max(1.0, min(char_count / config.TYPING_SPEED, config.TYPING_MAX)) sleep_delay = max(1.0, min(char_count / config.TYPING_SPEED, config.TYPING_MAX))
print(f'[{_ts()}] Typing delay: {sleep_delay:.1f}s ({char_count} chars)', flush=True) loggorilla.fyinf('telegram', f'Typing delay: {sleep_delay:.1f}s ({char_count} chars)')
time.sleep(sleep_delay) time.sleep(sleep_delay)
asyncio.run_coroutine_threadsafe( asyncio.run_coroutine_threadsafe(
@ -350,10 +351,10 @@ class TelegramClient:
self._loop self._loop
) )
else: else:
print(f'[{_ts()}] WARNING: cannot send to {chat_id} — loop unavailable', flush=True) loggorilla.error('telegram', f'WARNING: cannot send to {chat_id} — loop unavailable')
def _timeout_session(self, chat_id): def _timeout_session(self, chat_id):
print(f'[{_ts()}] Session timeout: {chat_id}', flush=True) loggorilla.fyinf('telegram', f'Session timeout: {chat_id}')
body = '*session end*' if 'roleplayer' in self._skill else 'Sesi ditutup. Sampai jumpa' body = '*session end*' if 'roleplayer' in self._skill else 'Sesi ditutup. Sampai jumpa'
self._schedule_send(chat_id, body) self._schedule_send(chat_id, body)
self._session_mgr.reset(str(chat_id)) self._session_mgr.reset(str(chat_id))

View File

@ -1,5 +1,5 @@
import asyncio import asyncio
import logging from scripts import loggorilla
import random import random
import signal import signal
import threading import threading
@ -350,14 +350,14 @@ class XMPPClient(ClientXMPP):
self._schedule_muc_rejoin(room) self._schedule_muc_rejoin(room)
async def _on_connected(self, event): async def _on_connected(self, event):
print(f'[{_ts()}] XMPP connected', flush=True) loggorilla.prcss('xmpp', 'XMPP connected')
_dbg(f'connected state: authenticated={self.authenticated}, bound={self.bound}, ' _dbg(f'connected state: authenticated={self.authenticated}, bound={self.bound}, '
f'sessionstarted={self.sessionstarted}, ' f'sessionstarted={self.sessionstarted}, '
f'_session_started={getattr(self, "_session_started", None)}') f'_session_started={getattr(self, "_session_started", None)}')
_dbg(f'connected via: {self.transport}') _dbg(f'connected via: {self.transport}')
async def _on_disconnected(self, event): async def _on_disconnected(self, event):
print(f'[{_ts()}] XMPP disconnected', flush=True) loggorilla.prcss('xmpp', 'XMPP disconnected')
# Anti-ban: cancel all pending rejoin tasks on disconnect # Anti-ban: cancel all pending rejoin tasks on disconnect
for room, task in list(self._muc_rejoin_tasks.items()): for room, task in list(self._muc_rejoin_tasks.items()):
if not task.done(): if not task.done():
@ -398,7 +398,7 @@ class XMPPClient(ClientXMPP):
f'_session_started={getattr(self, "_session_started", None)}') f'_session_started={getattr(self, "_session_started", None)}')
_dbg(f'boundjid: full={self.boundjid.full}, bare={self.boundjid.bare}, ' _dbg(f'boundjid: full={self.boundjid.full}, bare={self.boundjid.bare}, '
f'resource={self.boundjid.resource}, host={self.boundjid.host}') f'resource={self.boundjid.resource}, host={self.boundjid.host}')
print(f'[{_ts()}] XMPP online as {self.boundjid.full}', flush=True) loggorilla.prcss('xmpp', f'XMPP online as {self.boundjid.full}')
# Anti-ban: set disco identity SEKARANG (setelah bind, boundjid ber-resource) # Anti-ban: set disco identity SEKARANG (setelah bind, boundjid ber-resource)
# supaya identity masuk ke caps ver (presence) DAN respon disco#info server. # supaya identity masuk ke caps ver (presence) DAN respon disco#info server.
@ -418,15 +418,15 @@ class XMPPClient(ClientXMPP):
_dbg(f'disco identity check failed: {e}') _dbg(f'disco identity check failed: {e}')
# Anti-ban: send presence + request roster (wajib sebelum kirim pesan DM) # Anti-ban: send presence + request roster (wajib sebelum kirim pesan DM)
print(f'[{_ts()}] Sending initial presence...', flush=True) loggorilla.prcss('xmpp', 'Sending initial presence...')
self.send_presence() self.send_presence()
print(f'[{_ts()}] Requesting roster...', flush=True) loggorilla.prcss('xmpp', 'Requesting roster...')
self.get_roster() self.get_roster()
# Anti-ban: delay sebelum join MUC pertama agar startup tidak terlihat bot # Anti-ban: delay sebelum join MUC pertama agar startup tidak terlihat bot
if self._muc_rooms: if self._muc_rooms:
pre_delay = random.uniform(5.0, 15.0) pre_delay = random.uniform(5.0, 15.0)
print(f'[{_ts()}] MUC pre-join delay {pre_delay:.1f}s (anti-ban)...', flush=True) loggorilla.prcss('xmpp', f'MUC pre-join delay {pre_delay:.1f}s (anti-ban)...')
await asyncio.sleep(pre_delay) await asyncio.sleep(pre_delay)
# Anti-ban: delay sebelum join MUC untuk menghindari koneksi yang terlalu agresif # Anti-ban: delay sebelum join MUC untuk menghindari koneksi yang terlalu agresif
@ -434,7 +434,7 @@ class XMPPClient(ClientXMPP):
# Delay antar room join (3-8 detik per room) # Delay antar room join (3-8 detik per room)
if i > 0: if i > 0:
join_delay = random.uniform(3.0, 8.0) join_delay = random.uniform(3.0, 8.0)
print(f'[{_ts()}] MUC [{room}] Waiting {join_delay:.1f}s before join...', flush=True) loggorilla.prcss('xmpp', f'MUC [{room}] Waiting {join_delay:.1f}s before join...')
await asyncio.sleep(join_delay) await asyncio.sleep(join_delay)
# Anti-ban: retry join dengan incremental delay & nick fallback # Anti-ban: retry join dengan incremental delay & nick fallback
@ -443,26 +443,26 @@ class XMPPClient(ClientXMPP):
nick = self._get_muc_nick(room) nick = self._get_muc_nick(room)
try: try:
await self.plugin['xep_0045'].join_muc_wait(room, nick, maxstanzas=0) await self.plugin['xep_0045'].join_muc_wait(room, nick, maxstanzas=0)
print(f'[{_ts()}] Joined MUC room: {room} as {nick}', flush=True) loggorilla.prcss('xmpp', f'Joined MUC room: {room} as {nick}')
self._muc_last_join[room] = datetime.now() self._muc_last_join[room] = datetime.now()
self._muc_rejoin_attempts.pop(room, None) self._muc_rejoin_attempts.pop(room, None)
self._muc_rejoin_attempts.pop("_nick_" + room, None) self._muc_rejoin_attempts.pop("_nick_" + room, None)
success = True success = True
break break
except Exception as e: except Exception as e:
print(f'[{_ts()}] MUC join attempt #{attempt} failed ({room}): {e}', flush=True) loggorilla.error('xmpp', f'MUC join attempt #{attempt} failed ({room}): {e}')
# Anti-ban: handle 409 Conflict - coba nick alternatif # Anti-ban: handle 409 Conflict - coba nick alternatif
if '409' in str(e) or 'conflict' in str(e).lower(): if '409' in str(e) or 'conflict' in str(e).lower():
nick_attempts = self._muc_rejoin_attempts.get("_nick_" + room, 0) nick_attempts = self._muc_rejoin_attempts.get("_nick_" + room, 0)
if nick_attempts < MUC_NICK_SUFFIX_MAX: if nick_attempts < MUC_NICK_SUFFIX_MAX:
nick_attempts += 1 nick_attempts += 1
self._muc_rejoin_attempts["_nick_" + room] = nick_attempts self._muc_rejoin_attempts["_nick_" + room] = nick_attempts
print(f'[{_ts()}] MUC [{room}] Nick conflict, switching to: {self._get_muc_nick(room)}', flush=True) loggorilla.fyinf('xmpp', f'MUC [{room}] Nick conflict, switching to: {self._get_muc_nick(room)}')
# Retry segera dengan nick baru (jangan wait) # Retry segera dengan nick baru (jangan wait)
continue continue
else: else:
# Anti-ban: semua nick alternatif habis # Anti-ban: semua nick alternatif habis
print(f'[{_ts()}] MUC [{room}] All nick variations exhausted', flush=True) loggorilla.fyinf('xmpp', f'MUC [{room}] All nick variations exhausted')
break break
elif attempt < 3: elif attempt < 3:
# Anti-ban: error biasa, wait before retry (5s, 10s, 15s) # Anti-ban: error biasa, wait before retry (5s, 10s, 15s)
@ -481,18 +481,14 @@ class XMPPClient(ClientXMPP):
body = msg['body'].strip() body = msg['body'].strip()
if not body: if not body:
return return
print(f'[{_ts()}] DM from {jid}: {body[:60]}', flush=True) loggorilla.accss('xmpp', f'DM from {jid}: {body}')
# Anti-ban: proses langsung di event loop, bukan thread baru # Anti-ban: proses langsung di event loop, bukan thread baru
# Ini memastikan msg.send() dipanggil dari thread yang benar # Ini memastikan msg.send() dipanggil dari thread yang benar
asyncio.create_task(self._process_dm_async(jid, body)) asyncio.create_task(self._process_dm_async(jid, body))
def _on_message_error(self, msg): def _on_message_error(self, msg):
"""Anti-ban: handle error message dari server.""" """Anti-ban: handle error message dari server."""
print(f'[{_ts()}] MESSAGE ERROR from server:', flush=True) loggorilla.error('xmpp', f'MESSAGE ERROR from server:\n Type: {msg.get("type", "unknown")}\n From: {msg.get("from", "unknown")}\n To: {msg.get("to", "unknown")}\n Error: {msg.get("error", "unknown")}')
print(f'[{_ts()}] Type: {msg.get("type", "unknown")}', flush=True)
print(f'[{_ts()}] From: {msg.get("from", "unknown")}', flush=True)
print(f'[{_ts()}] To: {msg.get("to", "unknown")}', flush=True)
print(f'[{_ts()}] Error: {msg.get("error", "unknown")}', flush=True)
print(f'[{_ts()}] Full stanza: {msg}', flush=True) print(f'[{_ts()}] Full stanza: {msg}', flush=True)
# Catat error untuk adaptive rate limiting # Catat error untuk adaptive rate limiting
@ -500,13 +496,11 @@ class XMPPClient(ClientXMPP):
def _on_stream_error(self, error): def _on_stream_error(self, error):
"""Anti-ban: handle stream error dari server.""" """Anti-ban: handle stream error dari server."""
print(f'[{_ts()}] STREAM ERROR from server:', flush=True) loggorilla.error('xmpp', f'STREAM ERROR from server:\n Condition: {error.get("condition", "unknown")}\n Text: {error.get("text", "unknown")}')
print(f'[{_ts()}] Condition: {error.get("condition", "unknown")}', flush=True)
print(f'[{_ts()}] Text: {error.get("text", "unknown")}', flush=True)
print(f'[{_ts()}] Full error: {error}', flush=True) print(f'[{_ts()}] Full error: {error}', flush=True)
async def _on_roster_subscription_request(self, presence): async def _on_roster_subscription_request(self, presence):
print(f'[{_ts()}] Subscription request from {presence["from"]}', flush=True) loggorilla.fyinf('xmpp', f'Subscription request from {presence["from"]}')
# Anti-ban: TIDAK merespons subscribe otomatis (auto_subscribe=False). # Anti-ban: TIDAK merespons subscribe otomatis (auto_subscribe=False).
# Balasan 'subscribe' balik adalah red flag (contact farming) di conversations.im. # Balasan 'subscribe' balik adalah red flag (contact farming) di conversations.im.
_dbg(f'NOT auto-responding subscribe from {presence["from"]} ' _dbg(f'NOT auto-responding subscribe from {presence["from"]} '
@ -526,7 +520,7 @@ class XMPPClient(ClientXMPP):
body = msg['body'].strip() body = msg['body'].strip()
if not body: if not body:
return return
print(f'[{_ts()}] MUC [{room}] <{nick}>: {body[:60]}', flush=True) loggorilla.accss('xmpp', f'MUC [{room}] <{nick}>: {body}')
# Anti-ban: proses langsung di event loop, bukan thread baru # Anti-ban: proses langsung di event loop, bukan thread baru
asyncio.create_task(self._process_muc_async(room, nick, body)) asyncio.create_task(self._process_muc_async(room, nick, body))
@ -674,7 +668,7 @@ class XMPPClient(ClientXMPP):
if body == ':new': if body == ':new':
self._session_mgr.reset(room) self._session_mgr.reset(room)
print(f'[{_ts()}] Session reset for MUC room {room}', flush=True) loggorilla.prcss('xmpp', f'Session reset for MUC room {room}')
self._schedule_send(room, 'Memulai sesi baru. Ada yang bisa di bantu?', mtype='groupchat') self._schedule_send(room, 'Memulai sesi baru. Ada yang bisa di bantu?', mtype='groupchat')
return return
@ -745,17 +739,17 @@ class XMPPClient(ClientXMPP):
# END-SESSION: sinyal model bahwa sesi selesai. # END-SESSION: sinyal model bahwa sesi selesai.
# Jangan kirim apa pun ke room, tutup sesi langsung. # Jangan kirim apa pun ke room, tutup sesi langsung.
if sentinel == "END-SESSION": if sentinel == "END-SESSION":
print(f'[{_ts()}] MUC: Agent decided END-SESSION → closing session', flush=True) loggorilla.prcss('xmpp', 'MUC: Agent decided END-SESSION → closing session')
_pop_ghost("END-SESSION") _pop_ghost("END-SESSION")
self._session_mgr.reset(room) self._session_mgr.reset(room)
return return
# NO-REPLY: agent memilih diam → jangan kirim apa pun. # NO-REPLY: agent memilih diam → jangan kirim apa pun.
if sentinel == "NO-REPLY": if sentinel == "NO-REPLY":
print(f'[{_ts()}] MUC: Agent decided NO-REPLY → staying silent', flush=True) loggorilla.fyinf('xmpp', 'MUC: Agent decided NO-REPLY → staying silent')
_pop_ghost("NO-REPLY") _pop_ghost("NO-REPLY")
else: else:
print(f'[{_ts()}] MUC: sending response', flush=True) loggorilla.fyinf('xmpp', 'MUC: sending response')
self._schedule_send(room, final_content, 'groupchat') self._schedule_send(room, final_content, 'groupchat')
_flush_pending() _flush_pending()
else: else:
@ -773,7 +767,7 @@ class XMPPClient(ClientXMPP):
# Natural close via tool end_session: kirim sisa buffer (farewell), # Natural close via tool end_session: kirim sisa buffer (farewell),
# lalu tutup sesi. # lalu tutup sesi.
if _is_roleplay and should_close: if _is_roleplay and should_close:
print(f'[{_ts()}] MUC: natural close triggered for {room}', flush=True) loggorilla.prcss('xmpp', f'MUC: natural close triggered for {room}')
_flush_pending() _flush_pending()
self._session_mgr.reset(room) self._session_mgr.reset(room)
return return
@ -890,7 +884,7 @@ class XMPPClient(ClientXMPP):
self._session_mgr.reset(session_id) self._session_mgr.reset(session_id)
def start(self): def start(self):
print(f'[{_ts()}] Starting XMPP service...', flush=True) loggorilla.prcss('xmpp', 'Starting XMPP service...')
_setup_debug_logging() _setup_debug_logging()
asyncio.run(self._run()) asyncio.run(self._run())
@ -910,14 +904,14 @@ class XMPPClient(ClientXMPP):
except (NotImplementedError, RuntimeError): except (NotImplementedError, RuntimeError):
pass pass
print(f'[{_ts()}] Connecting to server (jid={self.jid})...', flush=True) loggorilla.prcss('xmpp', f'Connecting to server (jid={self.jid})...')
try: try:
await self.connect() await self.connect()
except Exception as e: except Exception as e:
print(f'[{_ts()}] CONNECT ERROR: {e}', flush=True) print(f'[{_ts()}] CONNECT ERROR: {e}', flush=True)
print(traceback.format_exc(), flush=True) print(traceback.format_exc(), flush=True)
raise raise
print(f'[{_ts()}] Connected, waiting for stream events...', flush=True) loggorilla.prcss('xmpp', 'Connected, waiting for stream events...')
try: try:
await self._stopped.wait() await self._stopped.wait()
except (asyncio.CancelledError, KeyboardInterrupt): except (asyncio.CancelledError, KeyboardInterrupt):