1 | |
2 |
|
3 | import * as wmill from "windmill-client" |
4 |
|
5 | |
6 | * New Email |
7 | * Emits messages received since the last poll that match a Gmail search query (default `in:inbox`), oldest first, with headers, labels and snippet. The first run sets the watermark to now and emits nothing. |
8 | */ |
9 | export async function main(auth: RT.Gmail, q: string = "in:inbox") { |
10 | const lastChecked: number | undefined = await wmill.getState() |
11 |
|
12 | |
13 | if (!lastChecked) { |
14 | await wmill.setState(Date.now()) |
15 | return [] |
16 | } |
17 |
|
18 | const base = "https://gmail.googleapis.com/gmail/v1/users/me/messages" |
19 | const headers = { |
20 | Authorization: `Bearer ${auth.token}`, |
21 | Accept: "application/json", |
22 | } |
23 |
|
24 | |
25 | const ids: string[] = [] |
26 | let pageToken: string | undefined |
27 | do { |
28 | const url = new URL(base) |
29 | url.searchParams.append("q", `${q} after:${Math.floor(lastChecked / 1000)}`) |
30 | url.searchParams.append("maxResults", "500") |
31 | if (pageToken) url.searchParams.append("pageToken", pageToken) |
32 | const response = await fetch(url, { headers }) |
33 | if (!response.ok) { |
34 | throw new Error(`${response.status} ${await response.text()}`) |
35 | } |
36 | const page = (await response.json()) as { |
37 | messages?: { id: string }[] |
38 | nextPageToken?: string |
39 | } |
40 | ids.push(...(page.messages ?? []).map((m) => m.id)) |
41 | pageToken = page.nextPageToken |
42 | } while (pageToken) |
43 |
|
44 | |
45 | |
46 | const messages: { internalDate: string }[] = [] |
47 | for (let i = 0; i < ids.length; i += 25) { |
48 | const batch = await Promise.all( |
49 | ids.slice(i, i + 25).map(async (id) => { |
50 | const url = new URL(`${base}/${id}`) |
51 | url.searchParams.append("format", "metadata") |
52 | for (const h of ["From", "To", "Cc", "Subject", "Date", "Message-ID"]) |
53 | url.searchParams.append("metadataHeaders", h) |
54 | const r = await fetch(url, { headers }) |
55 | if (!r.ok) throw new Error(`${r.status} ${await r.text()}`) |
56 | return (await r.json()) as { internalDate: string } |
57 | }) |
58 | ) |
59 | messages.push(...batch) |
60 | } |
61 |
|
62 | |
63 | const fresh = messages |
64 | .filter((m) => Number(m.internalDate) > lastChecked) |
65 | .sort((a, b) => Number(a.internalDate) - Number(b.internalDate)) |
66 |
|
67 | if (fresh.length > 0) { |
68 | await wmill.setState(Number(fresh[fresh.length - 1].internalDate)) |
69 | } |
70 |
|
71 | return fresh |
72 | } |
73 |
|