g1t/services/events/src/index.ts

117 lines3,141 bytesCodeBlame

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 attempts1import { WorkerEntrypoint } from "cloudflare:workers";
2
3import {
4 type EventQuery,
5 type EventsApi,
6 type G1tEvent,
7 type NewEvent,
8 newId,
9} from "@g1t/contracts";
10
11export 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
18type 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
28const DEFAULT_PAGE = 50;
29const MAX_PAGE = 200;
30
31function 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 */
47export 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}