{"version":3,"file":"index.cjs","names":["LoggerTransport"],"sources":["../../src/http/index.ts"],"sourcesContent":["import { LoggerTransport } from '@mastra/core/logger';\nimport type { BaseLogMessage, LogLevel } from '@mastra/core/logger';\n\ninterface RetryOptions {\n  maxRetries?: number;\n  retryDelay?: number;\n  exponentialBackoff?: boolean;\n}\n\ninterface HttpTransportOptions {\n  url: string;\n  method?: 'POST' | 'PUT' | 'PATCH';\n  headers?: Record<string, string>;\n  batchSize?: number;\n  flushInterval?: number;\n  timeout?: number;\n  retryOptions?: RetryOptions;\n}\n\nexport class HttpTransport extends LoggerTransport {\n  private url: string;\n  private method: string;\n  private headers: Record<string, string>;\n  private batchSize: number;\n  private flushInterval: number;\n  private timeout: number;\n  private retryOptions: Required<RetryOptions>;\n  private logBuffer: BaseLogMessage[];\n  private lastFlush: number;\n  private flushIntervalId: NodeJS.Timeout;\n\n  constructor(options: HttpTransportOptions) {\n    super({ objectMode: true });\n\n    if (!options.url) {\n      throw new Error('HTTP URL is required');\n    }\n\n    this.url = options.url;\n    this.method = options.method || 'POST';\n    this.headers = {\n      'Content-Type': 'application/json',\n      ...options.headers,\n    };\n    this.batchSize = options.batchSize || 100;\n    this.flushInterval = options.flushInterval || 10000;\n    this.timeout = options.timeout || 30000;\n    this.retryOptions = {\n      maxRetries: options.retryOptions?.maxRetries || 3,\n      retryDelay: options.retryOptions?.retryDelay || 1000,\n      exponentialBackoff: options.retryOptions?.exponentialBackoff || true,\n    };\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 HTTP endpoint:', err);\n      });\n    }, this.flushInterval);\n  }\n\n  private async makeHttpRequest(data: any, retryCount = 0): Promise<Response> {\n    const controller = new AbortController();\n    const timeoutId = setTimeout(() => controller.abort(), this.timeout);\n\n    try {\n      const body = JSON.stringify({ logs: data });\n\n      const response = await fetch(this.url, {\n        method: this.method,\n        headers: this.headers,\n        body,\n        signal: controller.signal,\n      });\n\n      clearTimeout(timeoutId);\n\n      if (!response.ok) {\n        throw new Error(`HTTP ${response.status}: ${response.statusText}`);\n      }\n\n      return response;\n    } catch (error) {\n      clearTimeout(timeoutId);\n\n      if (retryCount < this.retryOptions.maxRetries) {\n        const delay = this.retryOptions.exponentialBackoff\n          ? this.retryOptions.retryDelay * Math.pow(2, retryCount)\n          : this.retryOptions.retryDelay;\n\n        await new Promise(resolve => setTimeout(resolve, delay));\n        return this.makeHttpRequest(data, retryCount + 1);\n      }\n\n      throw error;\n    }\n  }\n\n  async _flush(): Promise<void> {\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      await this.makeHttpRequest(logs);\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): void {\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 HTTP endpoint:', 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): void {\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;\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    // HttpTransport is write-only by default\n    // Subclasses can override this method to implement log retrieval\n    console.warn(\n      'HttpTransport.listLogs: This transport is write-only. Override this method to implement log retrieval.',\n    );\n\n    return {\n      logs: [],\n      total: 0,\n      page: params?.page ?? 1,\n      perPage: params?.perPage ?? 100,\n      hasMore: false,\n    };\n  }\n\n  async listLogsByRunId({\n    runId: _runId,\n    fromDate: _fromDate,\n    toDate: _toDate,\n    logLevel: _logLevel,\n    filters: _filters,\n    page,\n    perPage,\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    // HttpTransport is write-only by default\n    // Subclasses can override this method to implement log retrieval\n    console.warn(\n      'HttpTransport.listLogsByRunId: This transport is write-only. Override this method to implement log retrieval.',\n    );\n\n    return {\n      logs: [],\n      total: 0,\n      page: page ?? 1,\n      perPage: perPage ?? 100,\n      hasMore: false,\n    };\n  }\n\n  // Utility methods\n  public getBufferedLogs(): BaseLogMessage[] {\n    return [...this.logBuffer];\n  }\n\n  public clearBuffer(): void {\n    this.logBuffer = [];\n  }\n\n  public getLastFlushTime(): number {\n    return this.lastFlush;\n  }\n}\n"],"mappings":";;;AAmBA,IAAa,gBAAb,cAAmCA,oBAAAA,gBAAgB;CACjD;CACA;CACA;CACA;CACA;CACA;CACA;CACA;CACA;CACA;CAEA,YAAY,SAA+B;EACzC,MAAM,EAAE,YAAY,KAAK,CAAC;EAE1B,IAAI,CAAC,QAAQ,KACX,MAAM,IAAI,MAAM,sBAAsB;EAGxC,KAAK,MAAM,QAAQ;EACnB,KAAK,SAAS,QAAQ,UAAU;EAChC,KAAK,UAAU;GACb,gBAAgB;GAChB,GAAG,QAAQ;EACb;EACA,KAAK,YAAY,QAAQ,aAAa;EACtC,KAAK,gBAAgB,QAAQ,iBAAiB;EAC9C,KAAK,UAAU,QAAQ,WAAW;EAClC,KAAK,eAAe;GAClB,YAAY,QAAQ,cAAc,cAAc;GAChD,YAAY,QAAQ,cAAc,cAAc;GAChD,oBAAoB,QAAQ,cAAc,sBAAsB;EAClE;EAEA,KAAK,YAAY,CAAC;EAClB,KAAK,YAAY,KAAK,IAAI;EAG1B,KAAK,kBAAkB,kBAAkB;GACvC,KAAK,OAAO,CAAC,CAAC,OAAM,QAAO;IACzB,QAAQ,MAAM,yCAAyC,GAAG;GAC5D,CAAC;EACH,GAAG,KAAK,aAAa;CACvB;CAEA,MAAc,gBAAgB,MAAW,aAAa,GAAsB;EAC1E,MAAM,aAAa,IAAI,gBAAgB;EACvC,MAAM,YAAY,iBAAiB,WAAW,MAAM,GAAG,KAAK,OAAO;EAEnE,IAAI;GACF,MAAM,OAAO,KAAK,UAAU,EAAE,MAAM,KAAK,CAAC;GAE1C,MAAM,WAAW,MAAM,MAAM,KAAK,KAAK;IACrC,QAAQ,KAAK;IACb,SAAS,KAAK;IACd;IACA,QAAQ,WAAW;GACrB,CAAC;GAED,aAAa,SAAS;GAEtB,IAAI,CAAC,SAAS,IACZ,MAAM,IAAI,MAAM,QAAQ,SAAS,OAAO,IAAI,SAAS,YAAY;GAGnE,OAAO;EACT,SAAS,OAAO;GACd,aAAa,SAAS;GAEtB,IAAI,aAAa,KAAK,aAAa,YAAY;IAC7C,MAAM,QAAQ,KAAK,aAAa,qBAC5B,KAAK,aAAa,aAAa,KAAK,IAAI,GAAG,UAAU,IACrD,KAAK,aAAa;IAEtB,MAAM,IAAI,SAAQ,YAAW,WAAW,SAAS,KAAK,CAAC;IACvD,OAAO,KAAK,gBAAgB,MAAM,aAAa,CAAC;GAClD;GAEA,MAAM;EACR;CACF;CAEA,MAAM,SAAwB;EAC5B,IAAI,KAAK,UAAU,WAAW,GAC5B;EAGF,MAAM,MAAM,KAAK,IAAI;EACrB,MAAM,OAAO,KAAK,UAAU,OAAO,GAAG,KAAK,SAAS;EAEpD,IAAI;GACF,MAAM,KAAK,gBAAgB,IAAI;GAC/B,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,IAAoB;EAC1D,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,yCAAyC,GAAG;GAC5D,CAAC;GAIH,GAAG,MAAM,KAAK;EAChB,SAAS,OAAO;GACd,GAAG,KAAK;EACV;CACF;CAEA,SAAS,KAAY,IAAoB;EACvC,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;EAGD,QAAQ,KACN,wGACF;EAEA,OAAO;GACL,MAAM,CAAC;GACP,OAAO;GACP,MAAM,QAAQ,QAAQ;GACtB,SAAS,QAAQ,WAAW;GAC5B,SAAS;EACX;CACF;CAEA,MAAM,gBAAgB,EACpB,OAAO,QACP,UAAU,WACV,QAAQ,SACR,UAAU,WACV,SAAS,UACT,MACA,WAeC;EAGD,QAAQ,KACN,+GACF;EAEA,OAAO;GACL,MAAM,CAAC;GACP,OAAO;GACP,MAAM,QAAQ;GACd,SAAS,WAAW;GACpB,SAAS;EACX;CACF;CAGA,kBAA2C;EACzC,OAAO,CAAC,GAAG,KAAK,SAAS;CAC3B;CAEA,cAA2B;EACzB,KAAK,YAAY,CAAC;CACpB;CAEA,mBAAkC;EAChC,OAAO,KAAK;CACd;AACF"}