import type { Db } from "../libs/db.js"; import type { ChatMessage } from "../types.js"; interface MessageRow { chat_id: number; message_id: number; user_id: number | null; date: number; author: string; text: string; reply_to_message_id: number | null; } const toMessage = (row: MessageRow): ChatMessage => ({ chatId: row.chat_id, messageId: row.message_id, userId: row.user_id, date: row.date, author: row.author, text: row.text, replyToMessageId: row.reply_to_message_id, }); export const createMessageRepository = (db: Db) => { const statements = { upsert: db.prepare(` INSERT INTO messages (chat_id, message_id, user_id, date, author, text, reply_to_message_id) VALUES (@chatId, @messageId, @userId, @date, @author, @text, @replyToMessageId) ON CONFLICT(chat_id, message_id) DO UPDATE SET text = excluded.text `), insertIfMissing: db.prepare(` INSERT OR IGNORE INTO messages (chat_id, message_id, user_id, date, author, text, reply_to_message_id) VALUES (@chatId, @messageId, @userId, @date, @author, @text, @replyToMessageId) `), pruneChat: db.prepare(` DELETE FROM messages WHERE chat_id = ? AND message_id <= ( SELECT message_id FROM messages WHERE chat_id = ? ORDER BY message_id DESC LIMIT 1 OFFSET ? ) `), latest: db.prepare("SELECT * FROM messages WHERE chat_id = ? ORDER BY message_id DESC LIMIT ?"), latestByUser: db.prepare( "SELECT * FROM messages WHERE chat_id = ? AND user_id = ? ORDER BY message_id DESC LIMIT ?", ), findById: db.prepare("SELECT * FROM messages WHERE chat_id = ? AND message_id = ?"), statsByChat: db.prepare("SELECT COUNT(*) AS count, MIN(date) AS oldest FROM messages WHERE chat_id = ?"), }; const toChronological = (rows: unknown[]): ChatMessage[] => (rows as MessageRow[]).reverse().map(toMessage); /** Сохраняет сообщение (при повторе — обновляет текст) и оставляет в чате не больше `keep` последних. */ const upsertAndPrune = db.transaction((message: ChatMessage, keep: number) => { statements.upsert.run(message); statements.pruneChat.run(message.chatId, message.chatId, keep); }); /** * Массовая загрузка (импорт): уже сохранённые сообщения не трогаем — живые данные точнее экспорта. * Возвращает, сколько сообщений добавлено; затем историю чата обрезаем до `keep` последних. */ const insertManyAndPrune = db.transaction((chatId: number, messages: ChatMessage[], keep: number): number => { const inserted = messages.reduce((sum, message) => sum + statements.insertIfMissing.run(message).changes, 0); statements.pruneChat.run(chatId, chatId, keep); return inserted; }); /** Последние `limit` сообщений чата в хронологическом порядке. */ const latest = (chatId: number, limit: number): ChatMessage[] => toChronological(statements.latest.all(chatId, limit)); /** Последние `limit` сообщений конкретного участника в хронологическом порядке. */ const latestByUser = (chatId: number, userId: number, limit: number): ChatMessage[] => toChronological(statements.latestByUser.all(chatId, userId, limit)); const findById = (chatId: number, messageId: number): ChatMessage | undefined => { const row = statements.findById.get(chatId, messageId) as MessageRow | undefined; return row && toMessage(row); }; /** Сколько сообщений чата сохранено и дата самого старого (null — если пусто). */ const statsByChat = (chatId: number): { count: number; oldest: number | null } => statements.statsByChat.get(chatId) as { count: number; oldest: number | null }; return { upsertAndPrune, insertManyAndPrune, latest, latestByUser, findById, statsByChat }; }; export type MessageRepository = ReturnType;