aegida-console / lib / db / conversations.ts
conversations.ts
Raw
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<MessageRow, "id" | "role" | "content">;
type GatewayAttachmentRow = { message_id: string; gate_file_id: string };

export async function listConversations(
  database: Queryable,
  userId: string,
): Promise<Conversation[]> {
  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<ConversationRow>(
    `SELECT id, title, model_id, updated_at
     FROM conversations
     WHERE user_id = $1
     ORDER BY updated_at DESC`,
    [userId],
  );
  const messages = await database.query<MessageRow>(
    `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<AttachmentRow>(
    `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<string, NonNullable<ChatMessage["attachments"]>>();
  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<string, ChatMessage[]>();
  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<string> {
  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<boolean> {
  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<void> {
  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<string> {
  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<ChatRequest["messages"]> {
  const result = await database.query<GatewayMessageRow>(
    `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<GatewayAttachmentRow>(
    `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<string, string[]>();
  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<void> {
  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<string | null> {
  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<boolean> {
  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(),
  };
}