All files / src/services/eventsub TwitchWebhookServer.ts

100% Statements 342/342
100% Branches 119/119
100% Functions 12/12
100% Lines 342/342

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 3433x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 3x 1x 1x 1x 6x 6x 1x 18x 18x 18x 18x 18x 18x 18x 18x 1x 554x 554x 3x 554x 554x 554x 552x 551x 554x 554x 554x 554x 554x 554x 554x 554x 554x 554x 549x 548x 547x 546x 545x 544x 554x 554x 554x 543x 543x 554x 1x 1x 554x 542x 542x 542x 542x 554x 554x 541x 541x 541x 554x 1x 1x 554x 540x 540x 1377x 538x 540x 554x 554x 554x 554x 4x 4x 3x 3x 3x 3x 3x 1x 4x 4x 1x 1x 1x 4x 4x 4x 4x 4x 4x 4x 554x 2x 2x 2x 2x 1x 1x 1x 1x 1x 554x 554x 528x 528x 1x 554x 554x 554x 526x 526x 526x 526x 257x 257x 257x 257x 526x 269x 269x 269x 269x 269x 269x 269x 269x 266x 266x 266x 266x 266x 266x 266x 269x 526x 526x 526x 526x 554x 4x 4x 2x 2x 4x 4x 4x 2x 2x 4x 4x 554x 1x 795x 795x 795x 795x 795x 795x 1x 269x 269x 269x 269x 269x 269x 269x 269x 269x 269x 269x 4x 4x 4x 269x 269x 269x 269x 269x 7x 7x 7x 7x 7x 7x 7x 7x 7x 7x 4x 4x 4x 4x 4x 4x 4x 4x 2x 2x 2x 2x 2x 2x 4x 1x 1x 1x 1x 1x 4x 7x 269x 1x 6x 6x 7x 7x 7x 3x 3x 3x 2x 2x 2x 2x 2x 2x 3x 1x 1x 1x 1x 1x 1x 1x 1x 1x 7x 7x 7x 6x 6x 1x 7x 5x 5x 5x 5x 5x 5x 5x 5x 2x 2x 2x 2x 2x 2x 2x 5x 3x 3x 3x 5x 5x 3x 3x 3x 5x 5x 5x 5x 2x 2x 2x 2x 2x 5x 1x 1x 7x 7x 7x 7x 7x 7x 7x 7x 7x 1x 24x 24x 24x 7x 7x 24x 24x 1x 544x 544x 544x 543x 543x 543x 544x 1085x 1085x 1085x 1085x 1x 1x 1x 1085x 542x 544x 543x 543x 544x 542x 542x 542x 542x 542x 542x 544x  
import { createServer, type Server } from "node:http"
import { createHmac, timingSafeEqual, randomUUID } from "node:crypto"
import { z } from "zod"
import { serializeLogError, type Logger } from "@workspace/logger"
import { routes, type EventSubRouter } from "./EventSubRouter"
import type { ConvexService } from "../ConvexService"
 
export const MAX_WEBHOOK_BYTES = 65536
const envelope = z.object({
  subscription: z.object({
    id: z.string().min(1).max(512),
    type: z.string(),
    version: z.string(),
    condition: z.record(z.string(), z.string()),
  }),
  challenge: z.string().min(1).max(4096).optional(),
  event: z.unknown().optional(),
})
export class TwitchWebhookServer {
  private server?: Server
  private readonly pending = new Set<Promise<unknown>>()
  get isListening(): boolean {
    return this.server?.listening === true
  }
  constructor(
    private readonly secret: string,
    private readonly convex: Pick<
      ConvexService,
      "reserve" | "subscriptionState" | "begin" | "finish" | "pendingEvents"
    >,
    private readonly router: Pick<EventSubRouter, "dispatch">,
    private readonly logger: Logger
  ) {}
  async handle(request: Request, now = Date.now()): Promise<Response> {
    if (
      request.method === "GET" &&
      new URL(request.url).pathname === "/healthz"
    )
      return new Response("ok")
    if (
      request.method !== "POST" ||
      new URL(request.url).pathname !== "/eventsub"
    )
      return new Response(null, { status: 404 })
    const messageId = request.headers.get("Twitch-Eventsub-Message-Id") ?? ""
    const timestamp =
      request.headers.get("Twitch-Eventsub-Message-Timestamp") ?? ""
    const signature =
      request.headers.get("Twitch-Eventsub-Message-Signature") ?? ""
    const time = Date.parse(timestamp)
    if (
      !messageId ||
      messageId.length > 512 ||
      !/^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,9})?Z$/.test(timestamp) ||
      !Number.isFinite(time) ||
      time < now - 600000 ||
      time > now + 60000 ||
      !/^sha256=[a-f0-9]{64}$/.test(signature)
    )
      return new Response(null, { status: 403 })
    let raw: Uint8Array
    try {
      raw = await boundedBody(request)
    } catch {
      return new Response(null, { status: 413 })
    }
    const expected = createHmac("sha256", this.secret)
      .update(messageId + timestamp)
      .update(raw)
      .digest()
    if (!timingSafeEqual(expected, Buffer.from(signature.slice(7), "hex")))
      return new Response(null, { status: 403 })
    try {
      let text: string
      try {
        text = new TextDecoder("utf-8", { fatal: true }).decode(raw)
      } catch {
        return new Response(null, { status: 400 })
      }
      const value = envelope.parse(JSON.parse(text))
      const route = routes.find(
        (route) =>
          route.type === value.subscription.type &&
          route.version === value.subscription.version
      )
      if (!route) return new Response(null, { status: 400 })
      const type = request.headers.get("Twitch-Eventsub-Message-Type")
      if (this.pending.size >= 256) return new Response(null, { status: 503 })
      if (type === "webhook_callback_verification") {
        if (!value.challenge) return new Response(null, { status: 400 })
        const task = this.convex
          .subscriptionState(
            value.subscription.id,
            false,
            route.key,
            value.subscription.condition.broadcaster_user_id ??
              value.subscription.condition.to_broadcaster_user_id
          )
          .catch(() => {
            this.logger.warn("EventSub verification status unavailable", {
              subscription: value.subscription.id,
            })
          })
          .finally(() => this.pending.delete(task))
        this.pending.add(task)
        return new Response(value.challenge, {
          headers: { "Content-Type": "text/plain; charset=utf-8" },
        })
      }
      if (type === "revocation") {
        const work = this.convex.subscriptionState(value.subscription.id, true)
        this.track(work)
        await work
        this.logger.warn("EventSub revoked", {
          type: route.type,
          subscription: value.subscription.id,
        })
        return new Response(null, { status: 204 })
      }
      if (type !== "notification") return new Response(null, { status: 400 })
      const parsed = route.parse(value.event)
      const conditionId =
        value.subscription.condition.broadcaster_user_id ??
        value.subscription.condition.to_broadcaster_user_id
      if (conditionId !== parsed.broadcasterId)
        return new Response(null, { status: 400 })
      this.logger.info("EventSub received", {
        type: route.type,
      })
      const work = (async () => {
        if (route.key === "streamOnline") {
          await this.router.dispatch(parsed, route.key, messageId)
          this.logger.info("Convex stream.online action accepted", {
            messageId,
          })
        } else {
          const receipt = await this.convex.reserve(
            messageId,
            route.key,
            parsed.broadcasterId,
            JSON.stringify(value.event)
          )
          if (receipt.kind === "pending")
            this.track(
              this.dispatchPending(
                parsed,
                route.key,
                messageId,
                receipt.template
              )
            )
        }
        return new Response(null, { status: 204 })
      })()
      this.track(work)
      return await work
    } catch (error) {
      if (!(error instanceof z.ZodError) && !(error instanceof SyntaxError))
        this.logger.error("EventSub backend action failed", {
          error: serializeLogError(error),
        })
      return new Response(null, {
        status:
          error instanceof z.ZodError || error instanceof SyntaxError
            ? 400
            : 503,
      })
    }
  }
  private track<T>(work: Promise<T>): void {
    this.pending.add(work)
    void work.then(
      () => this.pending.delete(work),
      () => this.pending.delete(work)
    )
  }
  private async dispatchPending(
    parsed: ReturnType<(typeof routes)[number]["parse"]>,
    key: (typeof routes)[number]["key"],
    messageId: string,
    template?: string,
    retry = 0
  ): Promise<void> {
    const attempt = randomUUID()
    let started = false
    let considered = false
    try {
      await this.router.dispatch(parsed, key, messageId, template, async () => {
        considered = true
        started = await this.convex.begin(messageId, attempt)
        return started
      })
      // Handlers that deliberately suppress output still settle their receipt.
      if (!considered) started = await this.convex.begin(messageId, attempt)
      if (started) await this.convex.finish(messageId, attempt, true)
    } catch (error) {
      if (started)
        await this.convex.finish(messageId, attempt, false).catch(() => {})
      this.logger.error("Event dispatch failed", {
        messageId,
        type: key,
        error: serializeLogError(error),
      })
      // Before send reservation the durable pending payload is recoverable.
      // After reservation an ambiguous POST is terminal and never replayed.
      if (!started && retry < 2) {
        try {
          const receipt = await this.convex.reserve(
            messageId,
            key,
            parsed.broadcasterId
          )
          if (receipt.kind === "pending")
            await this.dispatchPending(
              parsed,
              key,
              messageId,
              receipt.template,
              retry + 1
            )
        } catch (failure) {
          this.logger.error("Event retry deferred until recovery", {
            messageId,
            error: serializeLogError(failure),
          })
        }
      }
    }
  }
  async recoverPending(): Promise<void> {
    let cursor: string | undefined
    do {
      const page = await this.convex.pendingEvents(cursor)
      await Promise.all(
        page.events.map(async (event) => {
          const route = routes.find((route) => route.key === event.key)
          if (!route) return
          const parsed = route.parse(JSON.parse(event.eventJson))
          const receipt = await this.convex.reserve(
            event.messageId,
            route.key,
            parsed.broadcasterId,
            event.eventJson
          )
          if (receipt.kind === "pending") {
            const task = this.dispatchPending(
              parsed,
              route.key,
              event.messageId,
              receipt.template
            )
            this.track(task)
            await task
          }
        })
      )
      cursor = page.cursor ?? undefined
    } while (cursor)
  }
  async start(port: number, host = "127.0.0.1"): Promise<void> {
    this.server = createServer(async (incoming, outgoing) => {
      try {
        const chunks: Buffer[] = []
        let size = 0
        for await (const chunk of incoming.iterator({
          destroyOnReturn: false,
        })) {
          size += chunk.length
          if (size > MAX_WEBHOOK_BYTES) {
            outgoing.shouldKeepAlive = false
            incoming.pause()
            outgoing
              .writeHead(413, { Connection: "close" })
              .end(() => incoming.destroy())
            return
          }
          chunks.push(chunk)
        }
        const headers = new Headers()
        for (const [name, value] of Object.entries(incoming.headers))
          if (typeof value === "string") headers.set(name, value)
        const body = Buffer.concat(chunks)
        const request = new Request(`http://localhost${incoming.url}`, {
          method: incoming.method,
          headers,
          ...(incoming.method === "POST" ? { body } : {}),
        })
        const response = await this.handle(request)
        outgoing
          .writeHead(
            response.status,
            Object.fromEntries(response.headers.entries())
          )
          .end(Buffer.from(await response.arrayBuffer()))
      } catch {
        outgoing.writeHead(503).end()
      }
    })
    this.server.requestTimeout = 10000
    this.server.headersTimeout = 5000
    await new Promise<void>((resolve, reject) => {
      this.server!.once("error", reject)
      this.server!.listen(port, host, resolve)
    })
    this.logger.info("Twitch webhook listener ready", { host, port })
  }
  async stop(): Promise<void> {
    this.server?.closeAllConnections()
    if (this.server?.listening)
      await new Promise<void>((resolve, reject) =>
        this.server!.close((error) => (error ? reject(error) : resolve()))
      )
    await Promise.all(this.pending)
  }
}
export async function boundedBody(request: Request): Promise<Uint8Array> {
  if (!request.body) return new Uint8Array()
  const reader = request.body.getReader()
  const chunks: Uint8Array[] = []
  let size = 0
  try {
    while (true) {
      const result = await reader.read()
      if (result.done) break
      size += result.value.byteLength
      if (size > MAX_WEBHOOK_BYTES) {
        await reader.cancel()
        throw new Error("Body too large.")
      }
      chunks.push(result.value)
    }
  } finally {
    reader.releaseLock()
  }
  const body = new Uint8Array(size)
  let offset = 0
  for (const chunk of chunks) {
    body.set(chunk, offset)
    offset += chunk.length
  }
  return body
}