

























@@ -92,6 +92,10 @@ function resolveReplyMessage(msg: Message): Message | undefined {
9292return msg.reply_to_message ?? externalReply;
9393}
949495+function resolveEmbeddedReplyMessage(msg: Message): Message | undefined {
96+return msg.reply_to_message;
97+}
98+9599function resolveMessageBody(msg: Message): string | undefined {
96100const text = getTelegramTextParts(msg).text.trim();
97101if (text) {
@@ -139,6 +143,52 @@ function normalizeMessageNode(
139143};
140144}
141145146+function normalizeRequiredMessageNode(
147+msg: Message,
148+params: { threadId?: number },
149+): TelegramCachedMessageNode {
150+const node = normalizeMessageNode(msg, params);
151+if (!node) {
152+throw new Error("Telegram message cache node missing message id");
153+}
154+return node;
155+}
156+157+function resolveMessageThreadId(msg: Message): number | undefined {
158+const threadId = (msg as { message_thread_id?: unknown }).message_thread_id;
159+return typeof threadId === "number" && Number.isFinite(threadId)
160+ ? Math.trunc(threadId)
161+ : undefined;
162+}
163+164+function normalizeMessageNodes(
165+msg: Message,
166+params: { threadId?: number },
167+): TelegramCachedMessageNode[] {
168+const nodes: TelegramCachedMessageNode[] = [];
169+const visited = new Set<string>();
170+const nodeThreadId = (node: TelegramCachedMessageNode) => {
171+const threadId = Number(node.threadId);
172+return Number.isFinite(threadId) ? threadId : undefined;
173+};
174+const visit = (message: Message, inheritedThreadId?: number) => {
175+const node = normalizeMessageNode(message, {
176+threadId: resolveMessageThreadId(message) ?? inheritedThreadId,
177+});
178+if (!node?.messageId || visited.has(node.messageId)) {
179+return;
180+}
181+visited.add(node.messageId);
182+const replyMessage = resolveEmbeddedReplyMessage(message);
183+if (replyMessage?.message_id != null) {
184+visit(replyMessage, nodeThreadId(node) ?? inheritedThreadId);
185+}
186+nodes.push(node);
187+};
188+visit(msg, params.threadId);
189+return nodes;
190+}
191+142192function isRecord(value: unknown): value is Record<string, unknown> {
143193return typeof value === "object" && value !== null && !Array.isArray(value);
144194}
@@ -162,23 +212,27 @@ function isTelegramSourceMessage(value: unknown): value is Message {
162212);
163213}
164214165-function parsePersistedNode(value: unknown): TelegramCachedMessageNode | null {
166-if (!isRecord(value) || !isTelegramSourceMessage(value.sourceMessage)) {
167-return null;
168-}
169-const threadId = Number(readOptionalString(value, "threadId"));
170-return normalizeMessageNode(value.sourceMessage, Number.isFinite(threadId) ? { threadId } : {});
171-}
172-173-function parsePersistedEntry(value: unknown): {
215+function parsePersistedEntry(value: unknown): Array<{
174216key: string;
175217node: TelegramCachedMessageNode;
176-} | null {
218+}> {
177219if (!isRecord(value) || !isString(value.key)) {
178-return null;
220+return [];
221+}
222+const separatorIndex = value.key.lastIndexOf(":");
223+if (
224+separatorIndex === -1 ||
225+!isRecord(value.node) ||
226+!isTelegramSourceMessage(value.node.sourceMessage)
227+) {
228+return [];
179229}
180-const node = parsePersistedNode(value.node);
181-return node ? { key: value.key, node } : null;
230+const keyPrefix = value.key.slice(0, separatorIndex + 1);
231+const threadId = Number(readOptionalString(value.node, "threadId"));
232+return normalizeMessageNodes(
233+value.node.sourceMessage,
234+Number.isFinite(threadId) ? { threadId } : {},
235+).map((node) => ({ key: `${keyPrefix}${node.messageId}`, node }));
182236}
183237184238function findJsonArrayEnd(text: string): number {
@@ -270,6 +324,44 @@ function trimMessages(messages: Map<string, TelegramCachedMessageNode>, maxMessa
270324}
271325}
272326327+function mergeTelegramSourceMessage(existing: Message, incoming: Message): Message {
328+const existingReply = resolveEmbeddedReplyMessage(existing);
329+const incomingReply = resolveEmbeddedReplyMessage(incoming);
330+const merged = { ...existing, ...incoming };
331+if (existingReply?.message_id != null && incomingReply?.message_id === existingReply.message_id) {
332+return {
333+ ...merged,
334+reply_to_message: mergeTelegramSourceMessage(existingReply, incomingReply),
335+};
336+}
337+return merged;
338+}
339+340+function mergeCachedMessageNode(
341+existing: TelegramCachedMessageNode,
342+incoming: TelegramCachedMessageNode,
343+): TelegramCachedMessageNode {
344+const threadId = Number(incoming.threadId ?? existing.threadId);
345+return normalizeRequiredMessageNode(
346+mergeTelegramSourceMessage(existing.sourceMessage, incoming.sourceMessage),
347+{
348+ ...(Number.isFinite(threadId) ? { threadId } : {}),
349+},
350+);
351+}
352+353+function upsertCachedMessageNode(params: {
354+messages: Map<string, TelegramCachedMessageNode>;
355+key: string;
356+node: TelegramCachedMessageNode;
357+}): TelegramCachedMessageNode {
358+const existing = params.messages.get(params.key);
359+const node = existing ? mergeCachedMessageNode(existing, params.node) : params.node;
360+params.messages.delete(params.key);
361+params.messages.set(params.key, node);
362+return node;
363+}
364+273365function readPersistedMessages(filePath: string, maxMessages: number): PersistedMessageReadResult {
274366const messages = new Map<string, TelegramCachedMessageNode>();
275367let persistedEntryCount = 0;
@@ -281,14 +373,11 @@ function readPersistedMessages(filePath: string, maxMessages: number): Persisted
281373const persisted = readPersistedEntryValues(fs.readFileSync(filePath, "utf-8"));
282374needsRewrite = persisted.needsRewrite;
283375for (const value of persisted.values) {
284-const entry = parsePersistedEntry(value);
285-if (!entry) {
286-continue;
376+for (const entry of parsePersistedEntry(value)) {
377+persistedEntryCount++;
378+upsertCachedMessageNode({ messages, key: entry.key, node: entry.node });
379+trimMessages(messages, maxMessages);
287380}
288-persistedEntryCount++;
289-messages.delete(entry.key);
290-messages.set(entry.key, entry.node);
291-trimMessages(messages, maxMessages);
292381}
293382} catch (error) {
294383logVerbose(`telegram: failed to read message cache: ${String(error)}`);
@@ -426,30 +515,36 @@ export function createTelegramMessageCache(params?: {
426515427516return {
428517record: ({ accountId, chatId, msg, threadId }) => {
429-const entry = normalizeMessageNode(msg, { threadId });
430-if (!entry?.messageId) {
518+const entries = normalizeMessageNodes(msg, { threadId });
519+const entry = entries.at(-1);
520+if (!entry) {
431521return null;
432522}
433-const key = telegramMessageCacheKey({ accountId, chatId, messageId: entry.messageId });
434-messages.delete(key);
435-messages.set(key, entry);
436-trimMessages(messages, maxMessages);
437-try {
438-bucket.persistedEntryCount += appendPersistedMessage({
439- key,
440-node: entry,
441-persistedPath: params?.persistedPath,
442-});
443-if (bucket.persistedEntryCount > maxMessages * COMPACT_THRESHOLD_RATIO) {
444-bucket.persistedEntryCount = replacePersistedMessages({
445- messages,
523+let recordedEntry: TelegramCachedMessageNode | null = null;
524+for (const node of entries) {
525+const key = telegramMessageCacheKey({ accountId, chatId, messageId: node.messageId });
526+const cachedNode = upsertCachedMessageNode({ messages, key, node });
527+if (node.messageId === entry.messageId) {
528+recordedEntry = cachedNode;
529+}
530+trimMessages(messages, maxMessages);
531+try {
532+bucket.persistedEntryCount += appendPersistedMessage({
533+ key,
534+node: cachedNode,
446535persistedPath: params?.persistedPath,
447536});
537+if (bucket.persistedEntryCount > maxMessages * COMPACT_THRESHOLD_RATIO) {
538+bucket.persistedEntryCount = replacePersistedMessages({
539+ messages,
540+persistedPath: params?.persistedPath,
541+});
542+}
543+} catch (error) {
544+logVerbose(`telegram: failed to persist message cache: ${String(error)}`);
448545}
449-} catch (error) {
450-logVerbose(`telegram: failed to persist message cache: ${String(error)}`);
451546}
452-return entry;
547+return recordedEntry ?? entry;
453548},
454549 get,
455550recentBefore: ({ accountId, chatId, messageId, threadId, limit }) => {
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。