import type { ChatMessage, ChatRequest, Conversation, MessageStatus, ModelId, } from "@/lib/chat/types"; import type { Queryable } from "@/lib/db/types"; type ConversationRow = { id: string; title: string; model_id: ModelId; updated_at: Date | string; }; type MessageRow = { id: string; conversation_id: string; role: ChatMessage["role"]; content: string; status: MessageStatus; created_at: Date | string; }; type AttachmentRow = { id: string; message_id: string; original_name: string; content_type: string; byte_size: number }; type GatewayMessageRow = Pick; type GatewayAttachmentRow = { message_id: string; gate_file_id: string }; export async function listConversations( database: Queryable, userId: string, ): Promise { await database.query( `UPDATE chat_messages AS message SET status = 'stopped' FROM conversations AS conversation WHERE message.conversation_id = conversation.id AND conversation.user_id = $1 AND message.status = 'streaming'`, [userId], ); const conversations = await database.query( `SELECT id, title, model_id, updated_at FROM conversations WHERE user_id = $1 ORDER BY updated_at DESC`, [userId], ); const messages = await database.query( `SELECT message.id, message.conversation_id, message.role, message.content, message.status, message.created_at FROM chat_messages AS message JOIN conversations AS conversation ON conversation.id = message.conversation_id WHERE conversation.user_id = $1 ORDER BY message.created_at ASC`, [userId], ); const attachments = await database.query( `SELECT attachment.id, attachment.message_id, attachment.original_name, attachment.content_type, attachment.byte_size FROM chat_attachments AS attachment JOIN chat_messages AS message ON message.id = attachment.message_id JOIN conversations AS conversation ON conversation.id = message.conversation_id WHERE conversation.user_id = $1`, [userId], ); const attachmentsByMessage = new Map>(); for (const attachment of attachments.rows) { const current = attachmentsByMessage.get(attachment.message_id) ?? []; current.push({ id: attachment.id, name: attachment.original_name, contentType: attachment.content_type, size: Number(attachment.byte_size) }); attachmentsByMessage.set(attachment.message_id, current); } const messagesByConversation = new Map(); for (const message of messages.rows) { const collection = messagesByConversation.get(message.conversation_id) ?? []; collection.push({ ...mapMessage(message), attachments: attachmentsByMessage.get(message.id) ?? [] }); messagesByConversation.set(message.conversation_id, collection); } return conversations.rows.map((conversation) => ({ id: conversation.id, title: conversation.title, modelId: conversation.model_id, updatedAt: new Date(conversation.updated_at).toISOString(), messages: messagesByConversation.get(conversation.id) ?? [], })); } export async function createConversation( database: Queryable, userId: string, modelId: ModelId, title: string, ): Promise { const result = await database.query<{ id: string }>( `INSERT INTO conversations (user_id, model_id, title) VALUES ($1, $2, $3) RETURNING id`, [userId, modelId, title], ); return result.rows[0].id; } export async function ownedConversationExists( database: Queryable, userId: string, conversationId: string, ): Promise { const result = await database.query<{ id: string }>( "SELECT id FROM conversations WHERE id = $1 AND user_id = $2", [conversationId, userId], ); return Boolean(result.rows[0]); } export async function updateConversationModel( database: Queryable, userId: string, conversationId: string, modelId: ModelId, ): Promise { await database.query( `UPDATE conversations SET model_id = $1, updated_at = now() WHERE id = $2 AND user_id = $3`, [modelId, conversationId, userId], ); } export async function createMessage( database: Queryable, conversationId: string, role: ChatMessage["role"], content: string, status: MessageStatus, ): Promise { const result = await database.query<{ id: string }>( `INSERT INTO chat_messages (conversation_id, role, content, status) VALUES ($1, $2, $3, $4) RETURNING id`, [conversationId, role, content, status], ); await database.query( "UPDATE conversations SET updated_at = now() WHERE id = $1", [conversationId], ); return result.rows[0].id; } export async function listGatewayMessages( database: Queryable, userId: string, conversationId: string, ): Promise { const result = await database.query( `SELECT message.id, message.role, message.content FROM chat_messages AS message JOIN conversations AS conversation ON conversation.id = message.conversation_id WHERE conversation.id = $1 AND conversation.user_id = $2 AND message.status = 'complete' ORDER BY message.created_at ASC`, [conversationId, userId], ); const attachments = await database.query( `SELECT attachment.message_id, attachment.gate_file_id FROM chat_attachments AS attachment JOIN chat_messages AS message ON message.id = attachment.message_id JOIN conversations AS conversation ON conversation.id = message.conversation_id WHERE conversation.id = $1 AND conversation.user_id = $2 AND attachment.gate_file_id IS NOT NULL ORDER BY attachment.created_at ASC`, [conversationId, userId], ); const fileIdsByMessage = new Map(); for (const attachment of attachments.rows) { const fileIds = fileIdsByMessage.get(attachment.message_id) ?? []; fileIds.push(attachment.gate_file_id); fileIdsByMessage.set(attachment.message_id, fileIds); } const messages: ChatRequest["messages"] = []; for (const message of result.rows) { const fileIds = fileIdsByMessage.get(message.id) ?? []; if (fileIds.length === 0) { if (message.content) { messages.push({ role: message.role, content: message.content }); } continue; } messages.push({ role: message.role, content: [ ...(message.content ? [{ type: "text" as const, text: message.content }] : []), ...fileIds.map((fileId) => ({ type: "file" as const, file: { file_id: fileId }, })), ], }); } return messages; } export async function updateMessage( database: Queryable, messageId: string, content: string, status: MessageStatus, ): Promise { await database.query( "UPDATE chat_messages SET content = $1, status = $2 WHERE id = $3", [content, status, messageId], ); } export async function getRetryPrompt( database: Queryable, userId: string, conversationId: string, ): Promise { const result = await database.query<{ content: string }>( `SELECT previous.content FROM chat_messages AS final JOIN chat_messages AS previous ON previous.conversation_id = final.conversation_id WHERE final.conversation_id = $1 AND final.role = 'assistant' AND final.status IN ('error', 'stopped') AND previous.role = 'user' AND previous.status = 'complete' AND previous.created_at <= final.created_at AND NOT EXISTS ( SELECT 1 FROM chat_messages AS intervening WHERE intervening.conversation_id = final.conversation_id AND intervening.created_at > previous.created_at AND intervening.created_at < final.created_at ) AND EXISTS ( SELECT 1 FROM conversations WHERE id = final.conversation_id AND user_id = $2 ) ORDER BY final.created_at DESC, previous.created_at DESC LIMIT 1`, [conversationId, userId], ); return result.rows[0]?.content ?? null; } export async function deleteConversation( database: Queryable, userId: string, conversationId: string, ): Promise { const result = await database.query<{ id: string }>( "DELETE FROM conversations WHERE id = $1 AND user_id = $2 RETURNING id", [conversationId, userId], ); return Boolean(result.rows[0]); } function mapMessage(row: MessageRow): ChatMessage { return { id: row.id, role: row.role, content: row.content, status: row.status, createdAt: new Date(row.created_at).toISOString(), }; }