| 1 | import { WorkerEntrypoint } from "cloudflare:workers"; |
| 2 | |
| 3 | import { |
| 4 | type EventQuery, |
| 5 | type EventsApi, |
| 6 | type G1tEvent, |
| 7 | type NewEvent, |
| 8 | newId, |
| 9 | } from "@g1t/contracts"; |
| 10 | |
| 11 | export interface EventsEnv { |
| 12 | DB: D1Database; |
| 13 | BUS: Queue<G1tEvent>; |
| 14 | /** Every `SUBSCRIBER_*` binding is a queue that receives all events. */ |
| 15 | [subscriber: `SUBSCRIBER_${string}`]: Queue<G1tEvent>; |
| 16 | } |
| 17 | |
| 18 | type EventRow = { |
| 19 | id: string; |
| 20 | type: string; |
| 21 | source: string; |
| 22 | time: number; |
| 23 | repo_id: string | null; |
| 24 | actor: string | null; |
| 25 | data: string; |
| 26 | }; |
| 27 | |
| 28 | const DEFAULT_PAGE = 50; |
| 29 | const MAX_PAGE = 200; |
| 30 | |
| 31 | function toEvent(row: EventRow): G1tEvent { |
| 32 | return { |
| 33 | id: row.id, |
| 34 | type: row.type, |
| 35 | source: row.source, |
| 36 | time: row.time, |
| 37 | repoId: row.repo_id, |
| 38 | actor: row.actor, |
| 39 | data: JSON.parse(row.data), |
| 40 | } as G1tEvent; |
| 41 | } |
| 42 | |
| 43 | /** |
| 44 | * The event bus. Publishing enqueues; the queue consumer below writes the |
| 45 | * durable log and fans each batch out to every subscriber's own queue. |
| 46 | */ |
| 47 | export default class EventsService |
| 48 | extends WorkerEntrypoint<EventsEnv> |
| 49 | implements EventsApi |
| 50 | { |
| 51 | async publish(events: NewEvent[]): Promise<void> { |
| 52 | if (events.length === 0) return; |
| 53 | const now = Date.now(); |
| 54 | await this.env.BUS.sendBatch( |
| 55 | events.map((event) => ({ |
| 56 | body: { ...event, id: newId("evt", now), time: now } as G1tEvent, |
| 57 | })), |
| 58 | ); |
| 59 | } |
| 60 | |
| 61 | /** Callers must have checked that the viewer may see `query.repoId`. */ |
| 62 | async list(query: EventQuery): Promise<G1tEvent[]> { |
| 63 | const where: string[] = []; |
| 64 | const params: unknown[] = []; |
| 65 | if (query.repoId) { |
| 66 | where.push("repo_id = ?"); |
| 67 | params.push(query.repoId); |
| 68 | } |
| 69 | if (query.types?.length) { |
| 70 | where.push(`type IN (${query.types.map(() => "?").join(", ")})`); |
| 71 | params.push(...query.types); |
| 72 | } |
| 73 | if (query.before) { |
| 74 | where.push("id < ?"); |
| 75 | params.push(query.before); |
| 76 | } |
| 77 | const limit = Math.min(query.limit ?? DEFAULT_PAGE, MAX_PAGE); |
| 78 | const { results } = await this.env.DB.prepare( |
| 79 | `SELECT * FROM events ${where.length ? `WHERE ${where.join(" AND ")}` : ""} |
| 80 | ORDER BY id DESC LIMIT ?`, |
| 81 | ) |
| 82 | .bind(...params, limit) |
| 83 | .all<EventRow>(); |
| 84 | return results.map(toEvent); |
| 85 | } |
| 86 | |
| 87 | async queue(batch: MessageBatch<G1tEvent>): Promise<void> { |
| 88 | const events = batch.messages.map((message) => message.body); |
| 89 | const insert = this.env.DB.prepare( |
| 90 | // Redelivered batches must not duplicate log rows. |
| 91 | `INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data) |
| 92 | VALUES (?, ?, ?, ?, ?, ?, ?)`, |
| 93 | ); |
| 94 | await this.env.DB.batch( |
| 95 | events.map((event) => |
| 96 | insert.bind( |
| 97 | event.id, |
| 98 | event.type, |
| 99 | event.source, |
| 100 | event.time, |
| 101 | event.repoId, |
| 102 | event.actor, |
| 103 | JSON.stringify(event.data), |
| 104 | ), |
| 105 | ), |
| 106 | ); |
| 107 | |
| 108 | const subscribers = Object.entries(this.env) |
| 109 | .filter(([name]) => name.startsWith("SUBSCRIBER_")) |
| 110 | .map(([, queue]) => queue as Queue<G1tEvent>); |
| 111 | await Promise.all( |
| 112 | subscribers.map((queue) => |
| 113 | queue.sendBatch(events.map((event) => ({ body: event }))), |
| 114 | ), |
| 115 | ); |
| 116 | } |
| 117 | } |