router.ts 2.6 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273
  1. import { Buffer } from "node:buffer"
  2. import { timingSafeEqual } from "node:crypto"
  3. import { Effect, Schema } from "effect"
  4. import * as Semaphore from "effect/Semaphore"
  5. import { HttpRouter, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
  6. import { Resource } from "sst/resource"
  7. import { Ingest } from "./ingest"
  8. import { isShuttingDown } from "./shutdown"
  9. const MAX_CONCURRENT_INGEST_REQUESTS = 8
  10. const IngestPayload = Schema.Struct({
  11. events: Schema.optional(Schema.Unknown),
  12. })
  13. export const Routes = HttpRouter.use((router) =>
  14. Effect.gen(function* () {
  15. const ingestService = yield* Ingest
  16. const ingestRequests = yield* Semaphore.make(MAX_CONCURRENT_INGEST_REQUESTS)
  17. yield* Effect.all(
  18. [
  19. router.add("GET", "/health", () => json(200, { ok: true })),
  20. router.add("GET", "/ready", () => json(isShuttingDown() ? 503 : 200, { ok: !isShuttingDown() })),
  21. router.add("POST", "/", ingestRequests.withPermit(ingest(ingestService))),
  22. ],
  23. { discard: true },
  24. )
  25. }),
  26. )
  27. const ingest = (ingestService: Ingest.Service) =>
  28. Effect.gen(function* () {
  29. const request = yield* HttpServerRequest.HttpServerRequest
  30. if (!isAuthorized(request.headers)) return yield* json(401, { ok: false, error: "Unauthorized" })
  31. const payload = yield* HttpServerRequest.schemaBodyJson(IngestPayload).pipe(
  32. Effect.match({
  33. onFailure: () => undefined,
  34. onSuccess: (value) => value,
  35. }),
  36. )
  37. if (!payload) return yield* json(400, { ok: false, error: "Invalid JSON body" })
  38. const events = Array.isArray(payload.events) ? payload.events : []
  39. if (events.length === 0) return yield* json(202, { ok: true, records: 0 })
  40. return yield* ingestService.write(events).pipe(
  41. Effect.flatMap((result) => json(202, { ok: true, records: result.records })),
  42. Effect.catchTag("IngestError", (error) =>
  43. json(502, { ok: false, records: countRecords(events), failed: error.failed }),
  44. ),
  45. )
  46. })
  47. function isAuthorized(headers: Record<string, string | undefined>) {
  48. const actual = Buffer.from(headers.authorization ?? headers.Authorization ?? "")
  49. const expected = Buffer.from(`Bearer ${Resource.LakeIngestConfig.secret}`)
  50. if (actual.length !== expected.length) return false
  51. return timingSafeEqual(actual, expected)
  52. }
  53. function countRecords(items: unknown[]) {
  54. let records = 0
  55. for (const item of items) {
  56. if (Boolean(item) && typeof item === "object" && !Array.isArray(item)) records++
  57. }
  58. return records
  59. }
  60. function json(status: number, body: Record<string, unknown>) {
  61. return HttpServerResponse.json(body, { status }).pipe(Effect.orDie)
  62. }