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(),
};
}