log-processor.ts 7.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186
  1. import { Resource } from "@opencode-ai/console-resource"
  2. import type { TraceItem } from "@cloudflare/workers-types"
  3. export default {
  4. async tail(events: TraceItem[]) {
  5. for (const event of events) {
  6. if (!event.event) continue
  7. if (!("request" in event.event)) continue
  8. if (event.event.request.method !== "POST") continue
  9. const url = new URL(event.event.request.url)
  10. if (
  11. url.pathname !== "/zen/v1/chat/completions" &&
  12. url.pathname !== "/zen/v1/messages" &&
  13. url.pathname !== "/zen/v1/responses" &&
  14. !url.pathname.startsWith("/zen/v1/models/") &&
  15. url.pathname !== "/zen/go/v1/chat/completions" &&
  16. url.pathname !== "/zen/go/v1/messages" &&
  17. url.pathname !== "/zen/go/v1/responses" &&
  18. !url.pathname.startsWith("/zen/go/v1/models/")
  19. )
  20. continue
  21. let data: Record<string, unknown> = {
  22. "cf.continent": event.event.request.cf?.continent,
  23. "cf.country": event.event.request.cf?.country,
  24. "cf.city": event.event.request.cf?.city,
  25. "cf.region": event.event.request.cf?.region,
  26. "cf.latitude": event.event.request.cf?.latitude,
  27. "cf.longitude": event.event.request.cf?.longitude,
  28. "cf.timezone": event.event.request.cf?.timezone,
  29. duration: event.wallTime,
  30. request_length: parseInt(event.event.request.headers["content-length"] ?? "0"),
  31. status: event.event.response?.status ?? 0,
  32. ip: event.event.request.headers["x-real-ip"],
  33. }
  34. const time = new Date(event.eventTimestamp ?? Date.now()).toISOString()
  35. const events = [
  36. ...event.logs.flatMap((log) =>
  37. log.message.flatMap((message: string) => {
  38. if (!message.startsWith("_metric:")) return []
  39. const json = JSON.parse(message.slice(8)) as Record<string, unknown>
  40. data = { ...data, ...json }
  41. if ("llm.error.code" in json) {
  42. return [{ time, data: { ...data, event_type: "llm.error" } }]
  43. }
  44. return []
  45. }),
  46. ),
  47. { time, data: { ...data, event_type: "completions" } },
  48. ]
  49. console.log(JSON.stringify(data, null, 2))
  50. const lakeIngest = getLakeIngest()
  51. const [honeycomb, lake] = await Promise.all([
  52. fetch("https://api.honeycomb.io/1/batch/zen", {
  53. method: "POST",
  54. headers: {
  55. "Content-Type": "application/json",
  56. "X-Honeycomb-Team": Resource.HONEYCOMB_API_KEY.value,
  57. },
  58. body: JSON.stringify(events),
  59. }),
  60. ...(lakeIngest
  61. ? [
  62. fetch(lakeIngest.url, {
  63. method: "POST",
  64. headers: {
  65. "Content-Type": "application/json",
  66. Authorization: `Bearer ${lakeIngest.secret}`,
  67. },
  68. body: JSON.stringify({ events: events.map((event) => toLakeEvent(event.time, event.data)) }),
  69. }),
  70. ]
  71. : []),
  72. ])
  73. console.log(honeycomb.status)
  74. console.log(await honeycomb.text())
  75. if (lake) {
  76. console.log(lake.status)
  77. console.log(await lake.text())
  78. }
  79. }
  80. },
  81. }
  82. function getLakeIngest(): { url: string; secret: string } | undefined {
  83. try {
  84. return Resource.LakeIngest
  85. } catch {
  86. return undefined
  87. }
  88. }
  89. function toLakeEvent(time: string, data: Record<string, unknown>) {
  90. return {
  91. _datalake_key: "inference.event",
  92. event_timestamp: time,
  93. event_date: time.slice(0, 10),
  94. event_type: string(data, "event_type"),
  95. dataset: "zen",
  96. cf_continent: string(data, "cf.continent"),
  97. cf_country: string(data, "cf.country"),
  98. cf_city: string(data, "cf.city"),
  99. cf_region: string(data, "cf.region"),
  100. cf_latitude: number(data, "cf.latitude"),
  101. cf_longitude: number(data, "cf.longitude"),
  102. cf_timezone: string(data, "cf.timezone"),
  103. duration: number(data, "duration"),
  104. request_length: integer(data, "request_length"),
  105. status: integer(data, "status"),
  106. ip: string(data, "ip"),
  107. is_stream: boolean(data, "is_stream"),
  108. session: string(data, "session"),
  109. request: string(data, "request"),
  110. client: string(data, "client"),
  111. user_agent: string(data, "user_agent"),
  112. model_variant: string(data, "model.variant"),
  113. source: string(data, "source"),
  114. provider: string(data, "provider"),
  115. provider_model: string(data, "provider.model"),
  116. model: string(data, "model"),
  117. llm_error_code: integer(data, "llm.error.code"),
  118. llm_error_message: string(data, "llm.error.message"),
  119. error_response: string(data, "error.response"),
  120. error_type: string(data, "error.type"),
  121. error_message: string(data, "error.message"),
  122. error_cause: string(data, "error.cause"),
  123. error_cause2: string(data, "error.cause2"),
  124. api_key: string(data, "api_key"),
  125. workspace: string(data, "workspace"),
  126. is_subscription: boolean(data, "isSubscription"),
  127. subscription: string(data, "subscription"),
  128. response_length: integer(data, "response_length"),
  129. time_to_first_byte: integer(data, "time_to_first_byte"),
  130. timestamp_first_byte: integer(data, "timestamp.first_byte"),
  131. timestamp_last_byte: integer(data, "timestamp.last_byte"),
  132. tokens_input: integer(data, "tokens.input"),
  133. tokens_output: integer(data, "tokens.output"),
  134. tokens_reasoning: integer(data, "tokens.reasoning"),
  135. tokens_cache_read: integer(data, "tokens.cache_read"),
  136. tokens_cache_write_5m: integer(data, "tokens.cache_write_5m"),
  137. tokens_cache_write_1h: integer(data, "tokens.cache_write_1h"),
  138. cost_input_microcents: integer(data, "cost.input.microcents"),
  139. cost_output_microcents: integer(data, "cost.output.microcents"),
  140. cost_cache_read_microcents: integer(data, "cost.cache_read.microcents"),
  141. cost_cache_write_microcents: integer(data, "cost.cache_write.microcents"),
  142. cost_total_microcents: integer(data, "cost.total.microcents"),
  143. cost_input: integer(data, "cost.input"),
  144. cost_output: integer(data, "cost.output"),
  145. cost_cache_read: integer(data, "cost.cache_read"),
  146. cost_cache_write_5m: integer(data, "cost.cache_write_5m"),
  147. cost_cache_write_1h: integer(data, "cost.cache_write_1h"),
  148. cost_total: integer(data, "cost.total"),
  149. }
  150. }
  151. function string(data: Record<string, unknown>, key: string) {
  152. const value = data[key]
  153. if (typeof value === "string") return value
  154. if (typeof value === "number" || typeof value === "boolean") return String(value)
  155. return undefined
  156. }
  157. function boolean(data: Record<string, unknown>, key: string) {
  158. const value = data[key]
  159. if (typeof value === "boolean") return value
  160. if (typeof value === "string") return value === "true" ? true : value === "false" ? false : undefined
  161. return undefined
  162. }
  163. function integer(data: Record<string, unknown>, key: string) {
  164. const value = number(data, key)
  165. if (value === undefined) return undefined
  166. return Math.round(value)
  167. }
  168. function number(data: Record<string, unknown>, key: string) {
  169. const value = data[key]
  170. if (typeof value === "number") return Number.isFinite(value) ? value : undefined
  171. if (typeof value === "string") {
  172. const parsed = Number(value)
  173. return Number.isFinite(parsed) ? parsed : undefined
  174. }
  175. return undefined
  176. }