Initial project setup for Telegram-to-Qdrant knowledge base.

Add unified n8n workflow, userbot ingestion script, and detailed README for running n8n with a separate Dockerized Telethon collector.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-06-23 20:20:56 +03:00
commit a8d634f422
13 changed files with 1768 additions and 0 deletions

164
script.py Normal file
View File

@@ -0,0 +1,164 @@
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", "")
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)
def build_author(sender) -> dict:
if not sender:
return {
"id": "0",
"username": "",
"display_name": "Unknown",
"label": "Unknown",
}
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
return {
"id": str(getattr(sender, "id", 0)),
"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)
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) if reply_sender 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():
await client.start()
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())