{"version":3,"file":"index.cjs","names":["LoggerTransport"],"sources":["../../src/upstash/index.ts"],"sourcesContent":["import { LoggerTransport } from '@mastra/core/logger';\nimport type { BaseLogMessage, LogLevel } from '@mastra/core/logger';\n\nexport class UpstashTransport extends LoggerTransport {\n  upstashUrl: string;\n  upstashToken: string;\n  listName: string;\n  maxListLength: number;\n  batchSize: number;\n  flushInterval: number;\n  logBuffer: any[];\n  lastFlush: number;\n  flushIntervalId: NodeJS.Timeout;\n\n  constructor(opts: {\n    listName?: string;\n    maxListLength?: number;\n    batchSize?: number;\n    upstashUrl: string;\n    flushInterval?: number;\n    upstashToken: string;\n  }) {\n    super({ objectMode: true });\n\n    if (!opts.upstashUrl || !opts.upstashToken) {\n      throw new Error('Upstash URL and token are required');\n    }\n\n    this.upstashUrl = opts.upstashUrl;\n    this.upstashToken = opts.upstashToken;\n    this.listName = opts.listName || 'application-logs';\n    this.maxListLength = opts.maxListLength || 10000;\n    this.batchSize = opts.batchSize || 100;\n    this.flushInterval = opts.flushInterval || 10000;\n\n    this.logBuffer = [];\n    this.lastFlush = Date.now();\n\n    // Start flush interval\n    this.flushIntervalId = setInterval(() => {\n      this._flush().catch(err => {\n        console.error('Error flushing logs to Upstash:', err);\n      });\n    }, this.flushInterval);\n  }\n\n  private async executeUpstashCommand(command: any[]): Promise<any> {\n    const response = await fetch(`${this.upstashUrl}/pipeline`, {\n      method: 'POST',\n      headers: {\n        Authorization: `Bearer ${this.upstashToken}`,\n        'Content-Type': 'application/json',\n      },\n      body: JSON.stringify([command]),\n    });\n\n    if (!response.ok) {\n      throw new Error(`Failed to execute Upstash command: ${response.statusText}`);\n    }\n\n    return response.json();\n  }\n\n  async _flush() {\n    if (this.logBuffer.length === 0) {\n      return;\n    }\n\n    const now = Date.now();\n    const logs = this.logBuffer.splice(0, this.batchSize);\n\n    try {\n      // Prepare the Upstash Redis command\n      const command = ['LPUSH', this.listName, ...logs.map(log => JSON.stringify(log))];\n\n      // Trim the list if it exceeds maxListLength\n      if (this.maxListLength > 0) {\n        command.push('LTRIM', this.listName, 0 as any, (this.maxListLength - 1) as any);\n      }\n\n      // Send logs to Upstash Redis\n      await this.executeUpstashCommand(command);\n      this.lastFlush = now;\n    } catch (error) {\n      // On error, put logs back in the buffer\n      this.logBuffer.unshift(...logs);\n      throw error;\n    }\n  }\n\n  _write(chunk: any, encoding?: string, callback?: (error?: Error | null) => void): boolean {\n    if (typeof callback === 'function') {\n      this._transform(chunk, encoding || 'utf8', callback);\n      return true;\n    }\n\n    this._transform(chunk, encoding || 'utf8', (error: Error | null) => {\n      if (error) console.error('Transform error in write:', error);\n    });\n    return true;\n  }\n\n  _transform(chunk: string, _enc: string, cb: Function) {\n    try {\n      // Parse the log line if it's a string\n      const log = typeof chunk === 'string' ? JSON.parse(chunk) : chunk;\n\n      // Add timestamp if not present\n      if (!log.time) {\n        log.time = Date.now();\n      }\n\n      // Add to buffer\n      this.logBuffer.push(log);\n\n      // Flush if buffer reaches batch size\n      if (this.logBuffer.length >= this.batchSize) {\n        this._flush().catch(err => {\n          console.error('Error flushing logs to Upstash:', err);\n        });\n      }\n\n      // Pass through the log\n      cb(null, chunk);\n    } catch (error) {\n      cb(error);\n    }\n  }\n\n  _destroy(err: Error, cb: Function) {\n    clearInterval(this.flushIntervalId);\n\n    // Final flush\n    if (this.logBuffer.length > 0) {\n      this._flush()\n        .then(() => cb(err))\n        .catch(flushErr => {\n          console.error('Error in final flush:', flushErr);\n          cb(err || flushErr);\n        });\n    } else {\n      cb(err);\n    }\n  }\n\n  async listLogs(params?: {\n    fromDate?: Date;\n    toDate?: Date;\n    logLevel?: LogLevel;\n    filters?: Record<string, any>;\n    returnPaginationResults?: boolean; // default true\n    page?: number;\n    perPage?: number;\n  }): Promise<{\n    logs: BaseLogMessage[];\n    total: number;\n    page: number;\n    perPage: number;\n    hasMore: boolean;\n  }> {\n    try {\n      // Get all logs from the list\n      const command = ['LRANGE', this.listName, 0, -1];\n      const response = await this.executeUpstashCommand(command);\n\n      const logs =\n        (response?.[0]?.result?.map((log: string) => {\n          try {\n            // Parse the logs from JSON strings back to objects\n            return JSON.parse(log);\n          } catch {\n            return {};\n          }\n        }) as BaseLogMessage[]) || [];\n\n      let filteredLogs = logs.filter(record => record !== null && typeof record === 'object');\n\n      const {\n        fromDate,\n        toDate,\n        logLevel,\n        filters,\n        returnPaginationResults: returnPaginationResultsInput,\n        page: pageInput,\n        perPage: perPageInput,\n      } = params || {};\n\n      const page = pageInput === 0 ? 1 : (pageInput ?? 1);\n      const perPage = perPageInput ?? 100;\n      const returnPaginationResults = returnPaginationResultsInput ?? true;\n\n      if (filters) {\n        filteredLogs = filteredLogs.filter(log =>\n          Object.entries(filters || {}).every(([key, value]) => log[key as keyof BaseLogMessage] === value),\n        );\n      }\n\n      if (logLevel) {\n        filteredLogs = filteredLogs.filter(log => log.level === logLevel);\n      }\n\n      if (fromDate) {\n        filteredLogs = filteredLogs.filter(log => new Date(log.time)?.getTime() >= fromDate!.getTime());\n      }\n\n      if (toDate) {\n        filteredLogs = filteredLogs.filter(log => new Date(log.time)?.getTime() <= toDate!.getTime());\n      }\n\n      if (!returnPaginationResults) {\n        return {\n          logs: filteredLogs,\n          total: filteredLogs.length,\n          page,\n          perPage: filteredLogs.length,\n          hasMore: false,\n        };\n      }\n\n      const total = filteredLogs.length;\n      const resolvedPerPage = perPage || 100;\n      const start = (page - 1) * resolvedPerPage;\n      const end = start + resolvedPerPage;\n      const paginatedLogs = filteredLogs.slice(start, end);\n      const hasMore = end < total;\n\n      return {\n        logs: paginatedLogs,\n        total,\n        page,\n        perPage: resolvedPerPage,\n        hasMore,\n      };\n    } catch (error) {\n      console.error('Error getting logs from Upstash:', error);\n      return {\n        logs: [],\n        total: 0,\n        page: params?.page ?? 1,\n        perPage: params?.perPage ?? 100,\n        hasMore: false,\n      };\n    }\n  }\n\n  async listLogsByRunId({\n    runId,\n    fromDate,\n    toDate,\n    logLevel,\n    filters,\n    page: pageInput,\n    perPage: perPageInput,\n  }: {\n    runId: string;\n    fromDate?: Date;\n    toDate?: Date;\n    logLevel?: LogLevel;\n    filters?: Record<string, any>;\n    page?: number;\n    perPage?: number;\n  }): Promise<{\n    logs: BaseLogMessage[];\n    total: number;\n    page: number;\n    perPage: number;\n    hasMore: boolean;\n  }> {\n    try {\n      const page = pageInput === 0 ? 1 : (pageInput ?? 1);\n      const perPage = perPageInput ?? 100;\n      const allLogs = await this.listLogs({ fromDate, toDate, logLevel, filters, returnPaginationResults: false });\n      const logs = (allLogs?.logs?.filter((log: any) => log.runId === runId) || []) as BaseLogMessage[];\n      const total = logs.length;\n      const resolvedPerPage = perPage || 100;\n      const start = (page - 1) * resolvedPerPage;\n      const end = start + resolvedPerPage;\n      const paginatedLogs = logs.slice(start, end);\n      const hasMore = end < total;\n\n      return {\n        logs: paginatedLogs,\n        total,\n        page,\n        perPage: resolvedPerPage,\n        hasMore,\n      };\n    } catch (error) {\n      console.error('Error getting logs by runId from Upstash:', error);\n      return {\n        logs: [],\n        total: 0,\n        page: pageInput ?? 1,\n        perPage: perPageInput ?? 100,\n        hasMore: false,\n      };\n    }\n  }\n}\n"],"mappings":";;;AAGA,IAAa,mBAAb,cAAsCA,oBAAAA,gBAAgB;CACpD;CACA;CACA;CACA;CACA;CACA;CACA;CACA;CACA;CAEA,YAAY,MAOT;EACD,MAAM,EAAE,YAAY,KAAK,CAAC;EAE1B,IAAI,CAAC,KAAK,cAAc,CAAC,KAAK,cAC5B,MAAM,IAAI,MAAM,oCAAoC;EAGtD,KAAK,aAAa,KAAK;EACvB,KAAK,eAAe,KAAK;EACzB,KAAK,WAAW,KAAK,YAAY;EACjC,KAAK,gBAAgB,KAAK,iBAAiB;EAC3C,KAAK,YAAY,KAAK,aAAa;EACnC,KAAK,gBAAgB,KAAK,iBAAiB;EAE3C,KAAK,YAAY,CAAC;EAClB,KAAK,YAAY,KAAK,IAAI;EAG1B,KAAK,kBAAkB,kBAAkB;GACvC,KAAK,OAAO,CAAC,CAAC,OAAM,QAAO;IACzB,QAAQ,MAAM,mCAAmC,GAAG;GACtD,CAAC;EACH,GAAG,KAAK,aAAa;CACvB;CAEA,MAAc,sBAAsB,SAA8B;EAChE,MAAM,WAAW,MAAM,MAAM,GAAG,KAAK,WAAW,YAAY;GAC1D,QAAQ;GACR,SAAS;IACP,eAAe,UAAU,KAAK;IAC9B,gBAAgB;GAClB;GACA,MAAM,KAAK,UAAU,CAAC,OAAO,CAAC;EAChC,CAAC;EAED,IAAI,CAAC,SAAS,IACZ,MAAM,IAAI,MAAM,sCAAsC,SAAS,YAAY;EAG7E,OAAO,SAAS,KAAK;CACvB;CAEA,MAAM,SAAS;EACb,IAAI,KAAK,UAAU,WAAW,GAC5B;EAGF,MAAM,MAAM,KAAK,IAAI;EACrB,MAAM,OAAO,KAAK,UAAU,OAAO,GAAG,KAAK,SAAS;EAEpD,IAAI;GAEF,MAAM,UAAU;IAAC;IAAS,KAAK;IAAU,GAAG,KAAK,KAAI,QAAO,KAAK,UAAU,GAAG,CAAC;GAAC;GAGhF,IAAI,KAAK,gBAAgB,GACvB,QAAQ,KAAK,SAAS,KAAK,UAAU,GAAW,KAAK,gBAAgB,CAAS;GAIhF,MAAM,KAAK,sBAAsB,OAAO;GACxC,KAAK,YAAY;EACnB,SAAS,OAAO;GAEd,KAAK,UAAU,QAAQ,GAAG,IAAI;GAC9B,MAAM;EACR;CACF;CAEA,OAAO,OAAY,UAAmB,UAAoD;EACxF,IAAI,OAAO,aAAa,YAAY;GAClC,KAAK,WAAW,OAAO,YAAY,QAAQ,QAAQ;GACnD,OAAO;EACT;EAEA,KAAK,WAAW,OAAO,YAAY,SAAS,UAAwB;GAClE,IAAI,OAAO,QAAQ,MAAM,6BAA6B,KAAK;EAC7D,CAAC;EACD,OAAO;CACT;CAEA,WAAW,OAAe,MAAc,IAAc;EACpD,IAAI;GAEF,MAAM,MAAM,OAAO,UAAU,WAAW,KAAK,MAAM,KAAK,IAAI;GAG5D,IAAI,CAAC,IAAI,MACP,IAAI,OAAO,KAAK,IAAI;GAItB,KAAK,UAAU,KAAK,GAAG;GAGvB,IAAI,KAAK,UAAU,UAAU,KAAK,WAChC,KAAK,OAAO,CAAC,CAAC,OAAM,QAAO;IACzB,QAAQ,MAAM,mCAAmC,GAAG;GACtD,CAAC;GAIH,GAAG,MAAM,KAAK;EAChB,SAAS,OAAO;GACd,GAAG,KAAK;EACV;CACF;CAEA,SAAS,KAAY,IAAc;EACjC,cAAc,KAAK,eAAe;EAGlC,IAAI,KAAK,UAAU,SAAS,GAC1B,KAAK,OAAO,CAAC,CACV,WAAW,GAAG,GAAG,CAAC,CAAC,CACnB,OAAM,aAAY;GACjB,QAAQ,MAAM,yBAAyB,QAAQ;GAC/C,GAAG,OAAO,QAAQ;EACpB,CAAC;OAEH,GAAG,GAAG;CAEV;CAEA,MAAM,SAAS,QAcZ;EACD,IAAI;GAEF,MAAM,UAAU;IAAC;IAAU,KAAK;IAAU;IAAG;GAAE;GAa/C,IAAI,iBATD,MAHoB,KAAK,sBAAsB,OAAO,EAAA,GAG3C,EAAE,EAAE,QAAQ,KAAK,QAAgB;IAC3C,IAAI;KAEF,OAAO,KAAK,MAAM,GAAG;IACvB,QAAQ;KACN,OAAO,CAAC;IACV;GACF,CAAC,KAA0B,CAAC,EAAA,CAEN,QAAO,WAAU,WAAW,QAAQ,OAAO,WAAW,QAAQ;GAEtF,MAAM,EACJ,UACA,QACA,UACA,SACA,yBAAyB,8BACzB,MAAM,WACN,SAAS,iBACP,UAAU,CAAC;GAEf,MAAM,OAAO,cAAc,IAAI,IAAK,aAAa;GACjD,MAAM,UAAU,gBAAgB;GAChC,MAAM,0BAA0B,gCAAgC;GAEhE,IAAI,SACF,eAAe,aAAa,QAAO,QACjC,OAAO,QAAQ,WAAW,CAAC,CAAC,CAAC,CAAC,OAAO,CAAC,KAAK,WAAW,IAAI,SAAiC,KAAK,CAClG;GAGF,IAAI,UACF,eAAe,aAAa,QAAO,QAAO,IAAI,UAAU,QAAQ;GAGlE,IAAI,UACF,eAAe,aAAa,QAAO,QAAO,IAAI,KAAK,IAAI,IAAI,CAAC,CAAE,QAAQ,KAAK,SAAU,QAAQ,CAAC;GAGhG,IAAI,QACF,eAAe,aAAa,QAAO,QAAO,IAAI,KAAK,IAAI,IAAI,CAAC,CAAE,QAAQ,KAAK,OAAQ,QAAQ,CAAC;GAG9F,IAAI,CAAC,yBACH,OAAO;IACL,MAAM;IACN,OAAO,aAAa;IACpB;IACA,SAAS,aAAa;IACtB,SAAS;GACX;GAGF,MAAM,QAAQ,aAAa;GAC3B,MAAM,kBAAkB,WAAW;GACnC,MAAM,SAAS,OAAO,KAAK;GAC3B,MAAM,MAAM,QAAQ;GAIpB,OAAO;IACL,MAJoB,aAAa,MAAM,OAAO,GAI5B;IAClB;IACA;IACA,SAAS;IACT,SAPc,MAAM;GAQtB;EACF,SAAS,OAAO;GACd,QAAQ,MAAM,oCAAoC,KAAK;GACvD,OAAO;IACL,MAAM,CAAC;IACP,OAAO;IACP,MAAM,QAAQ,QAAQ;IACtB,SAAS,QAAQ,WAAW;IAC5B,SAAS;GACX;EACF;CACF;CAEA,MAAM,gBAAgB,EACpB,OACA,UACA,QACA,UACA,SACA,MAAM,WACN,SAAS,gBAeR;EACD,IAAI;GACF,MAAM,OAAO,cAAc,IAAI,IAAK,aAAa;GACjD,MAAM,UAAU,gBAAgB;GAEhC,MAAM,QAAQ,MADQ,KAAK,SAAS;IAAE;IAAU;IAAQ;IAAU;IAAS,yBAAyB;GAAM,CAAC,EAAA,EACpF,MAAM,QAAQ,QAAa,IAAI,UAAU,KAAK,KAAK,CAAC;GAC3E,MAAM,QAAQ,KAAK;GACnB,MAAM,kBAAkB,WAAW;GACnC,MAAM,SAAS,OAAO,KAAK;GAC3B,MAAM,MAAM,QAAQ;GAIpB,OAAO;IACL,MAJoB,KAAK,MAAM,OAAO,GAIpB;IAClB;IACA;IACA,SAAS;IACT,SAPc,MAAM;GAQtB;EACF,SAAS,OAAO;GACd,QAAQ,MAAM,6CAA6C,KAAK;GAChE,OAAO;IACL,MAAM,CAAC;IACP,OAAO;IACP,MAAM,aAAa;IACnB,SAAS,gBAAgB;IACzB,SAAS;GACX;EACF;CACF;AACF"}