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>
This commit is contained in:
2026-06-26 12:46:24 +03:00
parent 12b9992079
commit b2874ff580
3 changed files with 21 additions and 12 deletions

View File

@@ -76,7 +76,7 @@
}, },
{ {
"parameters": { "parameters": {
"jsCode": "const WINDOW_MS = 20 * 60 * 1000;\nconst MAX_MESSAGES = 12;\nconst MIN_MESSAGES_TO_FLUSH = 3;\n\nfunction aggregateBuffer(buffer) {\n const lines = buffer.messages.map((m) => {\n const repliedTo = m.replyTo ? ` -> ${m.replyTo}` : '';\n return `[${m.date}] ${m.author}${repliedTo}: ${m.text}`;\n });\n\n return {\n skip: false,\n text: [\n `Контекст диалога (${buffer.messages.length} сообщений, участники: ${buffer.participants.join(', ')})`,\n ...lines,\n ].join('\\n'),\n metadata: {\n ...buffer.lastMeta,\n context_id: buffer.contextId,\n context_message_count: buffer.messages.length,\n context_start_date: new Date(buffer.startTs).toISOString(),\n context_end_date: new Date(buffer.lastTs).toISOString(),\n context_participants: buffer.participants,\n context_mode: 'thread_time_window_20m'\n }\n };\n}\n\nconst staticData = $getWorkflowStaticData('global');\nstaticData.contextBuffers = staticData.contextBuffers || {};\n\nconst meta = $json.metadata || {};\nconst rawText = (meta.original_text || '').trim();\nif (!rawText) {\n return { json: { skip: true, reason: 'empty_original_text' } };\n}\n\nconst nowTs = meta.date ? new Date(meta.date).getTime() : Date.now();\nconst safeTs = Number.isNaN(nowTs) ? Date.now() : nowTs;\nconst chatId = String(meta.chat_id || meta.group_id || 'unknown_chat');\nconst threadId = String(meta.thread_id || meta.reply_to_message_id || meta.message_id || 'single');\nconst contextId = `${chatId}:${threadId}`;\n\nlet buffer = staticData.contextBuffers[contextId];\n\nif (buffer && safeTs - buffer.lastTs > WINDOW_MS) {\n const staleBuffer = buffer;\n buffer = null;\n\n const author = meta.author_label || meta.sender_name || 'Unknown';\n staticData.contextBuffers[contextId] = {\n contextId,\n startTs: safeTs,\n lastTs: safeTs,\n participants: [author],\n lastMeta: meta,\n messages: [\n {\n date: meta.date || new Date(safeTs).toISOString(),\n author,\n replyTo: meta.replied_to_author_label || '',\n text: rawText\n }\n ]\n };\n\n if (staleBuffer.messages.length >= MIN_MESSAGES_TO_FLUSH) {\n return { json: aggregateBuffer(staleBuffer) };\n }\n\n return { json: { skip: true, reason: 'buffer_reset_on_timeout', context_id: contextId } };\n}\n\nif (!buffer) {\n buffer = {\n contextId,\n startTs: safeTs,\n lastTs: safeTs,\n participants: [],\n lastMeta: meta,\n messages: []\n };\n}\n\nconst author = meta.author_label || meta.sender_name || 'Unknown';\nif (!buffer.participants.includes(author)) {\n buffer.participants.push(author);\n}\n\nbuffer.messages.push({\n date: meta.date || new Date(safeTs).toISOString(),\n author,\n replyTo: meta.replied_to_author_label || '',\n text: rawText\n});\nbuffer.lastMeta = meta;\nbuffer.lastTs = safeTs;\n\nconst shouldFlush =\n buffer.messages.length >= MAX_MESSAGES ||\n buffer.lastTs - buffer.startTs >= WINDOW_MS;\n\nif (!shouldFlush) {\n staticData.contextBuffers[contextId] = buffer;\n return {\n json: {\n skip: true,\n reason: 'buffering',\n context_id: contextId,\n buffered_messages: buffer.messages.length\n }\n };\n}\n\ndelete staticData.contextBuffers[contextId];\nreturn { json: aggregateBuffer(buffer) };" "jsCode": "const WINDOW_MS = 5 * 60 * 1000;\nconst MAX_MESSAGES = 3;\nconst MIN_MESSAGES_TO_FLUSH = 1;\n\nfunction aggregateBuffer(buffer) {\n const lines = buffer.messages.map((m) => {\n const repliedTo = m.replyTo ? ` -> ${m.replyTo}` : '';\n return `[${m.date}] ${m.author}${repliedTo}: ${m.text}`;\n });\n\n return {\n skip: false,\n text: [\n `Контекст диалога (${buffer.messages.length} сообщений, участники: ${buffer.participants.join(', ')})`,\n ...lines,\n ].join('\\n'),\n metadata: {\n ...buffer.lastMeta,\n context_id: buffer.contextId,\n context_message_count: buffer.messages.length,\n context_start_date: new Date(buffer.startTs).toISOString(),\n context_end_date: new Date(buffer.lastTs).toISOString(),\n context_participants: buffer.participants,\n context_mode: 'thread_time_window_5m_fast_flush'\n }\n };\n}\n\nconst staticData = $getWorkflowStaticData('global');\nstaticData.contextBuffers = staticData.contextBuffers || {};\n\nconst meta = $json.metadata || {};\nconst rawText = (meta.original_text || '').trim();\nif (!rawText) {\n return { json: { skip: true, reason: 'empty_original_text' } };\n}\n\nconst nowTs = meta.date ? new Date(meta.date).getTime() : Date.now();\nconst safeTs = Number.isNaN(nowTs) ? Date.now() : nowTs;\nconst chatId = String(meta.chat_id || meta.group_id || 'unknown_chat');\nconst threadId = String(meta.thread_id || meta.reply_to_message_id || meta.message_id || 'single');\nconst contextId = `${chatId}:${threadId}`;\n\nlet buffer = staticData.contextBuffers[contextId];\n\nif (buffer && safeTs - buffer.lastTs > WINDOW_MS) {\n const staleBuffer = buffer;\n buffer = null;\n\n const author = meta.author_label || meta.sender_name || 'Unknown';\n staticData.contextBuffers[contextId] = {\n contextId,\n startTs: safeTs,\n lastTs: safeTs,\n participants: [author],\n lastMeta: meta,\n messages: [\n {\n date: meta.date || new Date(safeTs).toISOString(),\n author,\n replyTo: meta.replied_to_author_label || '',\n text: rawText\n }\n ]\n };\n\n if (staleBuffer.messages.length >= MIN_MESSAGES_TO_FLUSH) {\n return { json: aggregateBuffer(staleBuffer) };\n }\n\n return { json: { skip: true, reason: 'buffer_reset_on_timeout', context_id: contextId } };\n}\n\nif (!buffer) {\n buffer = {\n contextId,\n startTs: safeTs,\n lastTs: safeTs,\n participants: [],\n lastMeta: meta,\n messages: []\n };\n}\n\nconst author = meta.author_label || meta.sender_name || 'Unknown';\nif (!buffer.participants.includes(author)) {\n buffer.participants.push(author);\n}\n\nbuffer.messages.push({\n date: meta.date || new Date(safeTs).toISOString(),\n author,\n replyTo: meta.replied_to_author_label || '',\n text: rawText\n});\nbuffer.lastMeta = meta;\nbuffer.lastTs = safeTs;\n\nconst shouldFlush =\n buffer.messages.length >= MAX_MESSAGES ||\n buffer.lastTs - buffer.startTs >= WINDOW_MS;\n\nif (!shouldFlush) {\n staticData.contextBuffers[contextId] = buffer;\n return {\n json: {\n skip: true,\n reason: 'buffering',\n context_id: contextId,\n buffered_messages: buffer.messages.length\n }\n };\n}\n\ndelete staticData.contextBuffers[contextId];\nreturn { json: aggregateBuffer(buffer) };"
}, },
"type": "n8n-nodes-base.code", "type": "n8n-nodes-base.code",
"typeVersion": 2, "typeVersion": 2,

View File

@@ -232,11 +232,11 @@ docker compose logs -f tg-userbot
## Параметры склейки сообщений в n8n ## Параметры склейки сообщений в n8n
В узле `Build Context Block`: В узле `Build Context Block`:
- `WINDOW_MS = 20 минут` - `WINDOW_MS = 5 минут`
- `MAX_MESSAGES = 12` - `MAX_MESSAGES = 3`
- `MIN_MESSAGES_TO_FLUSH = 3` - `MIN_MESSAGES_TO_FLUSH = 1`
Рекомендации: Рекомендации:
- меньше блоки: `MAX_MESSAGES = 8-10`; - если всё ещё есть задержка индексации: `MAX_MESSAGES = 1-2`;
- шире контекст: `WINDOW_MS = 30-40 минут`. - если нужен более широкий контекст: `WINDOW_MS = 10-20 минут`, `MAX_MESSAGES = 5-8`.

View File

@@ -57,13 +57,15 @@ client = TelegramClient(
auto_reconnect=True, auto_reconnect=True,
) )
def build_author(sender) -> dict: def build_author(sender, sender_id=None) -> dict:
effective_id = str(sender_id or 0)
if not sender: if not sender:
unknown_label = f"id:{effective_id}" if effective_id != "0" else "Unknown"
return { return {
"id": "0", "id": effective_id,
"username": "", "username": "",
"display_name": "Unknown", "display_name": "Unknown",
"label": "Unknown", "label": unknown_label,
} }
username = getattr(sender, "username", "") or "" username = getattr(sender, "username", "") or ""
display_name = ( display_name = (
@@ -71,8 +73,11 @@ def build_author(sender) -> dict:
or "Unknown" or "Unknown"
) )
label = f"@{username}" if username else display_name 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 { return {
"id": str(getattr(sender, "id", 0)), "id": sender_obj_id,
"username": username, "username": username,
"display_name": display_name, "display_name": display_name,
"label": label, "label": label,
@@ -105,10 +110,14 @@ def extract_attachment(message) -> dict:
async def send_to_n8n(message): async def send_to_n8n(message):
sender = message.sender or await message.get_sender() sender = message.sender or await message.get_sender()
author = build_author(sender) author = build_author(sender, message.sender_id)
reply_message = await message.get_reply_message() if message.reply_to_msg_id else None 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_sender = await reply_message.get_sender() if reply_message else None
reply_author = build_author(reply_sender) if reply_sender else None reply_author = (
build_author(reply_sender, getattr(reply_message, "sender_id", None))
if reply_message
else None
)
data = { data = {
"text": message.text, "text": message.text,