{"version":3,"file":"message_histories.cjs","names":["BaseListChatMessageHistory"],"sources":["../src/message_histories.ts"],"sourcesContent":["import { v4 } from \"@langchain/core/utils/uuid\";\nimport type { D1Database } from \"@cloudflare/workers-types\";\nimport { BaseListChatMessageHistory } from \"@langchain/core/chat_history\";\nimport {\n  BaseMessage,\n  StoredMessage,\n  StoredMessageData,\n  mapChatMessagesToStoredMessages,\n  mapStoredMessagesToChatMessages,\n} from \"@langchain/core/messages\";\n/**\n * Type definition for the input parameters required when instantiating a\n * CloudflareD1MessageHistory object.\n */\nexport type CloudflareD1MessageHistoryInput = {\n  tableName?: string;\n  sessionId: string;\n  database?: D1Database;\n};\n\n/**\n * Interface for the data transfer object used when selecting stored\n * messages from the Cloudflare D1 database.\n */\ninterface selectStoredMessagesDTO {\n  id: string;\n  session_id: string;\n  type: string;\n  content: string;\n  role: string | null;\n  name: string | null;\n  additional_kwargs: string;\n}\n\n/**\n * Class for storing and retrieving chat message history from a\n * Cloudflare D1 database. Extends the BaseListChatMessageHistory class.\n * @example\n * ```typescript\n * const memory = new BufferMemory({\n *   returnMessages: true,\n *   chatHistory: new CloudflareD1MessageHistory({\n *     tableName: \"stored_message\",\n *     sessionId: \"example\",\n *     database: env.DB,\n *   }),\n * });\n *\n * const chainInput = { input };\n *\n * const res = await memory.chatHistory.invoke(chainInput);\n * await memory.saveContext(chainInput, {\n *   output: res,\n * });\n * ```\n */\nexport class CloudflareD1MessageHistory extends BaseListChatMessageHistory {\n  lc_namespace = [\"langchain\", \"stores\", \"message\", \"cloudflare_d1\"];\n\n  public database: D1Database;\n\n  private tableName: string;\n\n  private sessionId: string;\n\n  private tableInitialized: boolean;\n\n  constructor(fields: CloudflareD1MessageHistoryInput) {\n    super(fields);\n\n    const { sessionId, database, tableName } = fields;\n\n    if (database) {\n      this.database = database;\n    } else {\n      throw new Error(\n        \"Either a client or config must be provided to CloudflareD1MessageHistory\"\n      );\n    }\n\n    this.tableName = tableName || \"langchain_chat_histories\";\n    this.tableInitialized = false;\n    this.sessionId = sessionId;\n  }\n\n  /**\n   * Private method to ensure that the necessary table exists in the\n   * Cloudflare D1 database before performing any operations. If the table\n   * does not exist, it is created.\n   * @returns Promise that resolves to void.\n   */\n  private async ensureTable(): Promise<void> {\n    if (this.tableInitialized) {\n      return;\n    }\n\n    const query = `CREATE TABLE IF NOT EXISTS ${this.tableName} (id TEXT PRIMARY KEY, session_id TEXT, type TEXT, content TEXT, role TEXT, name TEXT, additional_kwargs TEXT);`;\n    await this.database.prepare(query).bind().all();\n\n    const idIndexQuery = `CREATE INDEX IF NOT EXISTS id_index ON ${this.tableName} (id);`;\n    await this.database.prepare(idIndexQuery).bind().all();\n\n    const sessionIdIndexQuery = `CREATE INDEX IF NOT EXISTS session_id_index ON ${this.tableName} (session_id);`;\n    await this.database.prepare(sessionIdIndexQuery).bind().all();\n\n    this.tableInitialized = true;\n  }\n\n  /**\n   * Method to retrieve all messages from the Cloudflare D1 database for the\n   * current session.\n   * @returns Promise that resolves to an array of BaseMessage objects.\n   */\n  async getMessages(): Promise<BaseMessage[]> {\n    await this.ensureTable();\n\n    const query = `SELECT * FROM ${this.tableName} WHERE session_id = ?`;\n    const rawStoredMessages = await this.database\n      .prepare(query)\n      .bind(this.sessionId)\n      .all();\n    const storedMessagesObject =\n      rawStoredMessages.results as unknown as selectStoredMessagesDTO[];\n\n    const orderedMessages: StoredMessage[] = storedMessagesObject.map(\n      (message) => {\n        const data = {\n          content: message.content,\n          additional_kwargs: JSON.parse(message.additional_kwargs),\n        } as StoredMessageData;\n\n        if (message.role) {\n          data.role = message.role;\n        }\n\n        if (message.name) {\n          data.name = message.name;\n        }\n\n        return {\n          type: message.type,\n          data,\n        };\n      }\n    );\n\n    return mapStoredMessagesToChatMessages(orderedMessages);\n  }\n\n  /**\n   * Method to add a new message to the Cloudflare D1 database for the current\n   * session.\n   * @param message The BaseMessage object to be added to the database.\n   * @returns Promise that resolves to void.\n   */\n  async addMessage(message: BaseMessage): Promise<void> {\n    await this.ensureTable();\n\n    const messageToAdd = mapChatMessagesToStoredMessages([message]);\n\n    const query = `INSERT INTO ${this.tableName} (id, session_id, type, content, role, name, additional_kwargs) VALUES(?, ?, ?, ?, ?, ?, ?)`;\n\n    const id = v4();\n\n    await this.database\n      .prepare(query)\n      .bind(\n        id,\n        this.sessionId,\n        messageToAdd[0].type || null,\n        messageToAdd[0].data.content || null,\n        messageToAdd[0].data.role || null,\n        messageToAdd[0].data.name || null,\n        JSON.stringify(messageToAdd[0].data.additional_kwargs)\n      )\n      .all();\n  }\n\n  /**\n   * Method to delete all messages from the Cloudflare D1 database for the\n   * current session.\n   * @returns Promise that resolves to void.\n   */\n  async clear(): Promise<void> {\n    await this.ensureTable();\n\n    const query = `DELETE FROM ? WHERE session_id = ? `;\n    await this.database\n      .prepare(query)\n      .bind(this.tableName, this.sessionId)\n      .all();\n  }\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;;;;;AAwDA,IAAa,6BAAb,cAAgDA,6BAAAA,2BAA2B;CACzE,eAAe;EAAC;EAAa;EAAU;EAAW;EAAgB;CAElE;CAEA;CAEA;CAEA;CAEA,YAAY,QAAyC;AACnD,QAAM,OAAO;EAEb,MAAM,EAAE,WAAW,UAAU,cAAc;AAE3C,MAAI,SACF,MAAK,WAAW;MAEhB,OAAM,IAAI,MACR,2EACD;AAGH,OAAK,YAAY,aAAa;AAC9B,OAAK,mBAAmB;AACxB,OAAK,YAAY;;;;;;;;CASnB,MAAc,cAA6B;AACzC,MAAI,KAAK,iBACP;EAGF,MAAM,QAAQ,8BAA8B,KAAK,UAAU;AAC3D,QAAM,KAAK,SAAS,QAAQ,MAAM,CAAC,MAAM,CAAC,KAAK;EAE/C,MAAM,eAAe,0CAA0C,KAAK,UAAU;AAC9E,QAAM,KAAK,SAAS,QAAQ,aAAa,CAAC,MAAM,CAAC,KAAK;EAEtD,MAAM,sBAAsB,kDAAkD,KAAK,UAAU;AAC7F,QAAM,KAAK,SAAS,QAAQ,oBAAoB,CAAC,MAAM,CAAC,KAAK;AAE7D,OAAK,mBAAmB;;;;;;;CAQ1B,MAAM,cAAsC;AAC1C,QAAM,KAAK,aAAa;EAExB,MAAM,QAAQ,iBAAiB,KAAK,UAAU;AA8B9C,UAAA,GAAA,yBAAA,kCAxBE,MAL8B,KAAK,SAClC,QAAQ,MAAM,CACd,KAAK,KAAK,UAAU,CACpB,KAAK,EAEY,QAE0C,KAC3D,YAAY;GACX,MAAM,OAAO;IACX,SAAS,QAAQ;IACjB,mBAAmB,KAAK,MAAM,QAAQ,kBAAkB;IACzD;AAED,OAAI,QAAQ,KACV,MAAK,OAAO,QAAQ;AAGtB,OAAI,QAAQ,KACV,MAAK,OAAO,QAAQ;AAGtB,UAAO;IACL,MAAM,QAAQ;IACd;IACD;IAIiD,CAAC;;;;;;;;CASzD,MAAM,WAAW,SAAqC;AACpD,QAAM,KAAK,aAAa;EAExB,MAAM,gBAAA,GAAA,yBAAA,iCAA+C,CAAC,QAAQ,CAAC;EAE/D,MAAM,QAAQ,eAAe,KAAK,UAAU;EAE5C,MAAM,MAAA,GAAA,2BAAA,KAAS;AAEf,QAAM,KAAK,SACR,QAAQ,MAAM,CACd,KACC,IACA,KAAK,WACL,aAAa,GAAG,QAAQ,MACxB,aAAa,GAAG,KAAK,WAAW,MAChC,aAAa,GAAG,KAAK,QAAQ,MAC7B,aAAa,GAAG,KAAK,QAAQ,MAC7B,KAAK,UAAU,aAAa,GAAG,KAAK,kBAAkB,CACvD,CACA,KAAK;;;;;;;CAQV,MAAM,QAAuB;AAC3B,QAAM,KAAK,aAAa;AAGxB,QAAM,KAAK,SACR,QAAQ,sCAAM,CACd,KAAK,KAAK,WAAW,KAAK,UAAU,CACpC,KAAK"}