0

New Email

by
Published today

Emits new Gmail messages matching a search query since the last poll.

Scriptยท trigger gmail Verified

The script

Submitted by hugo989 Typescript (fetch-only)
Verified 4 hours ago
1
//native
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
  // First run: set the watermark to now and don't emit a backlog.
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
  // `after:` takes epoch seconds and is inclusive at the boundary second.
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
  // Fetch in batches: Gmail allows ~50 messages.get per second per user, and a
45
  // long gap between polls can return hundreds of ids.
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
  // Drop what the inclusive boundary second already emitted.
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