Files
DBBot-telegram/script.py
bilal b2874ff580 Reduce ingestion buffering and improve author metadata fallback.
Flush context blocks faster to avoid delayed indexing and use sender_id-based labels when Telegram sender objects are unavailable in historical messages.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-06-26 12:46:24 +03:00

208 lines
7.2 KiB
Python

import asyncio
import os
from pathlib import Path
import requests
from telethon import TelegramClient, events
from telethon.sessions import StringSession
API_ID = int(os.getenv("TG_API_ID", "0"))
API_HASH = os.getenv("TG_API_HASH", "")
GROUP_ID = int(os.getenv("TG_GROUP_ID", "0"))
N8N_WEBHOOK_URL = os.getenv("N8N_WEBHOOK_URL", "http://n8n:5678/webhook/tg-inbox")
HISTORY_CHECK_FILE = os.getenv("HISTORY_CHECK_FILE", "/app/session/history_done.flag")
SESSION_FILE_PATH = os.getenv("SESSION_FILE_PATH", "/app/session/session.session")
SESSION_STRING = os.getenv("TG_SESSION_STRING", "")
TG_USE_IPV6 = os.getenv("TG_USE_IPV6", "false").lower() in {"1", "true", "yes"}
TG_CONNECT_TIMEOUT = int(os.getenv("TG_CONNECT_TIMEOUT", "20"))
TG_REQUEST_RETRIES = int(os.getenv("TG_REQUEST_RETRIES", "10"))
TG_CONNECTION_RETRIES = int(os.getenv("TG_CONNECTION_RETRIES", "20"))
TG_RETRY_DELAY = int(os.getenv("TG_RETRY_DELAY", "5"))
TG_RECONNECT_WAIT = int(os.getenv("TG_RECONNECT_WAIT", "15"))
def load_session_string() -> str:
if SESSION_STRING:
return SESSION_STRING.strip()
session_path = Path(SESSION_FILE_PATH)
if not session_path.exists():
raise FileNotFoundError(
f"Session file not found: {SESSION_FILE_PATH}. "
"Provide TG_SESSION_STRING or mount session file."
)
return session_path.read_text(encoding="utf-8").strip()
def validate_config() -> None:
if API_ID <= 0:
raise ValueError("TG_API_ID is required and must be > 0")
if not API_HASH:
raise ValueError("TG_API_HASH is required")
if GROUP_ID == 0:
raise ValueError("TG_GROUP_ID is required")
if not N8N_WEBHOOK_URL:
raise ValueError("N8N_WEBHOOK_URL is required")
validate_config()
client = TelegramClient(
StringSession(load_session_string()),
API_ID,
API_HASH,
use_ipv6=TG_USE_IPV6,
timeout=TG_CONNECT_TIMEOUT,
request_retries=TG_REQUEST_RETRIES,
connection_retries=TG_CONNECTION_RETRIES,
retry_delay=TG_RETRY_DELAY,
auto_reconnect=True,
)
def build_author(sender, sender_id=None) -> dict:
effective_id = str(sender_id or 0)
if not sender:
unknown_label = f"id:{effective_id}" if effective_id != "0" else "Unknown"
return {
"id": effective_id,
"username": "",
"display_name": "Unknown",
"label": unknown_label,
}
username = getattr(sender, "username", "") or ""
display_name = (
f"{getattr(sender, 'first_name', '')} {getattr(sender, 'last_name', '')}".strip()
or "Unknown"
)
label = f"@{username}" if username else display_name
sender_obj_id = str(getattr(sender, "id", None) or sender_id or 0)
if label == "Unknown" and sender_obj_id != "0":
label = f"id:{sender_obj_id}"
return {
"id": sender_obj_id,
"username": username,
"display_name": display_name,
"label": label,
}
def extract_entities(message) -> list:
entities = []
for entity in (message.entities or []):
entities.append(
{
"type": entity.__class__.__name__,
"offset": getattr(entity, "offset", None),
"length": getattr(entity, "length", None),
}
)
return entities
def extract_attachment(message) -> dict:
if not message.media:
return {"has_media": False}
return {
"has_media": True,
"media_type": message.media.__class__.__name__,
"file_name": getattr(message.file, "name", None) if message.file else None,
"mime_type": getattr(message.file, "mime_type", None) if message.file else None,
}
async def send_to_n8n(message):
sender = message.sender or await message.get_sender()
author = build_author(sender, message.sender_id)
reply_message = await message.get_reply_message() if message.reply_to_msg_id else None
reply_sender = await reply_message.get_sender() if reply_message else None
reply_author = (
build_author(reply_sender, getattr(reply_message, "sender_id", None))
if reply_message
else None
)
data = {
"text": message.text,
"metadata": {
"date": str(message.date),
"message_id": message.id,
"chat_id": str(message.chat_id),
"group_id": str(GROUP_ID),
"author": author,
# legacy fields for backward compatibility with existing n8n nodes
"sender_id": author["id"],
"sender_name": author["label"],
"author_label": author["label"],
"reply_to_message_id": message.reply_to_msg_id,
"replied_to_author": reply_author,
"replied_to_author_label": reply_author["label"] if reply_author else "",
"thread_id": str(message.reply_to.reply_to_top_id) if getattr(message, "reply_to", None) and getattr(message.reply_to, "reply_to_top_id", None) else "",
"edit_date": str(message.edit_date) if message.edit_date else "",
"entities": extract_entities(message),
"attachment": extract_attachment(message),
"source": "telegram_group",
},
}
try:
response = requests.post(N8N_WEBHOOK_URL, json=data, timeout=10)
return response.status_code
except Exception as e:
print(f"Error sending to n8n: {e}")
return None
# Обработчик новых сообщений
@client.on(events.NewMessage(chats=GROUP_ID))
async def handler(event):
if event.message.text:
await send_to_n8n(event.message)
async def main():
print(
"Connecting to Telegram with settings:",
{
"use_ipv6": TG_USE_IPV6,
"connect_timeout": TG_CONNECT_TIMEOUT,
"request_retries": TG_REQUEST_RETRIES,
"connection_retries": TG_CONNECTION_RETRIES,
"retry_delay": TG_RETRY_DELAY,
},
)
while True:
try:
await client.start()
break
except Exception as e:
print(f"Telegram connect failed: {e}. Retrying in {TG_RECONNECT_WAIT}s...")
await asyncio.sleep(TG_RECONNECT_WAIT)
print("UserBot started...")
if not os.path.exists(HISTORY_CHECK_FILE):
print("Starting massive history sync...")
count = 0
async for message in client.iter_messages(GROUP_ID, reverse=True):
if message.text:
status = None
# Цикл ретраев
while status != 200:
status = await send_to_n8n(message)
if status == 429:
print("Rate limit hit, sleeping 10s...")
await asyncio.sleep(10)
elif status != 200:
print(f"Error {status}, retrying in 2s...")
await asyncio.sleep(2)
count += 1
if count % 20 == 0:
print(f"Processed {count} messages...")
await asyncio.sleep(0.3)
Path(HISTORY_CHECK_FILE).write_text("done", encoding="utf-8")
print("History sync complete.")
await client.run_until_disconnected()
if __name__ == '__main__':
asyncio.run(main())