pr_01m47d15m3e54sn21z27rpy5n9/services/events/src/index.ts
Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Initial g1t: services, event bus, intents and attempts | 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 | } |