feat: add gmail-hook monorepo with watcher, api, and web
Ship the receipt pipeline (Livin/Grab/BRI), PocketBase spendings, PM2 ecosystem, and a proper .gitignore that excludes secrets, SQLite data, and media dumps. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -0,0 +1,169 @@
|
||||
import { timingSafeEqual } from "node:crypto";
|
||||
import { PUBSUB_VERIFICATION_TOKEN } from "../config.js";
|
||||
import { getUserByEmail, getWatchState, saveWatchState } from "../db.js";
|
||||
import {
|
||||
gmailClientForUser,
|
||||
googleStatus,
|
||||
listHistory,
|
||||
startWatch,
|
||||
} from "../google.js";
|
||||
import { HttpError } from "../http.js";
|
||||
import { runPipeline } from "../pipeline.js";
|
||||
import { processFetchedMessage } from "../process-fetched.js";
|
||||
|
||||
const userLocks = new Map();
|
||||
|
||||
function withUserLock(userId, fn) {
|
||||
const prev = userLocks.get(userId) || Promise.resolve();
|
||||
const next = prev.then(fn, fn);
|
||||
userLocks.set(
|
||||
userId,
|
||||
next.then(
|
||||
() => {},
|
||||
() => {},
|
||||
),
|
||||
);
|
||||
return next;
|
||||
}
|
||||
|
||||
function tokenFromRequest(request) {
|
||||
const q = request.query?.token;
|
||||
if (typeof q === "string" && q) return q;
|
||||
const header =
|
||||
request.headers["x-pubsub-verification-token"] ||
|
||||
request.headers["x-goog-pubsub-verification-token"];
|
||||
if (typeof header === "string" && header) return header;
|
||||
const auth = request.headers.authorization;
|
||||
if (typeof auth === "string" && auth.startsWith("Bearer ")) {
|
||||
return auth.slice(7);
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
function tokensMatch(provided, expected) {
|
||||
const a = Buffer.from(provided);
|
||||
const b = Buffer.from(expected);
|
||||
if (a.length !== b.length) return false;
|
||||
return timingSafeEqual(a, b);
|
||||
}
|
||||
|
||||
function verifyPush(request) {
|
||||
if (!PUBSUB_VERIFICATION_TOKEN) return;
|
||||
if (!tokensMatch(tokenFromRequest(request), PUBSUB_VERIFICATION_TOKEN)) {
|
||||
throw new HttpError(403, "Invalid Pub/Sub verification token");
|
||||
}
|
||||
}
|
||||
|
||||
function decodeNotification(body) {
|
||||
const data = body?.message?.data;
|
||||
if (typeof data !== "string" || !data) return null;
|
||||
try {
|
||||
const json = Buffer.from(data, "base64").toString("utf8");
|
||||
const parsed = JSON.parse(json);
|
||||
if (!parsed?.emailAddress || parsed.historyId == null) return null;
|
||||
return {
|
||||
emailAddress: parsed.emailAddress,
|
||||
historyId: String(parsed.historyId),
|
||||
};
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async function processNotification(log, { fetchHandlers, processHandlers }, notification) {
|
||||
const user = getUserByEmail(notification.emailAddress);
|
||||
if (!user) {
|
||||
log.warn(
|
||||
{ email: notification.emailAddress },
|
||||
"gmail-watch: unknown mailbox, acking",
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
await withUserLock(user.id, () =>
|
||||
handleUserNotification(log, { fetchHandlers, processHandlers }, user, notification),
|
||||
);
|
||||
}
|
||||
|
||||
async function handleUserNotification(
|
||||
log,
|
||||
{ fetchHandlers, processHandlers },
|
||||
user,
|
||||
notification,
|
||||
) {
|
||||
const gmail = await gmailClientForUser(user.id);
|
||||
|
||||
const state = getWatchState(user.id);
|
||||
const startHistoryId = state?.history_id;
|
||||
if (!startHistoryId) {
|
||||
log.info(
|
||||
{ userId: user.id },
|
||||
"gmail-watch: no cursor yet, starting watch and skipping backlog",
|
||||
);
|
||||
await startWatch(user.id, { resetCursor: true });
|
||||
return;
|
||||
}
|
||||
|
||||
let result;
|
||||
try {
|
||||
result = await listHistory(gmail, startHistoryId);
|
||||
} catch (err) {
|
||||
const status = googleStatus(err);
|
||||
if (status === 404) {
|
||||
log.warn(
|
||||
{ userId: user.id, startHistoryId },
|
||||
"gmail-watch: history expired, re-watching without backfill",
|
||||
);
|
||||
await startWatch(user.id, { resetCursor: true });
|
||||
return;
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
|
||||
for (const messageId of result.messageIds) {
|
||||
const fetched = await runPipeline(fetchHandlers, {
|
||||
notification,
|
||||
userId: user.id,
|
||||
messageId,
|
||||
gmail,
|
||||
log,
|
||||
});
|
||||
await processFetchedMessage(fetched, processHandlers);
|
||||
}
|
||||
|
||||
saveWatchState(user.id, {
|
||||
historyId: result.historyId || notification.historyId,
|
||||
expiration: state.expiration || Date.now() + 7 * 24 * 3600 * 1000,
|
||||
});
|
||||
}
|
||||
|
||||
export async function registerPubsubRoutes(app, { fetchHandlers, processHandlers }) {
|
||||
app.post("/pubsub/gmail", async (request, reply) => {
|
||||
verifyPush(request);
|
||||
|
||||
const notification = decodeNotification(request.body);
|
||||
if (!notification) {
|
||||
request.log.warn("gmail-watch: undecodable Pub/Sub payload, acking");
|
||||
return reply.code(204).send();
|
||||
}
|
||||
|
||||
try {
|
||||
await processNotification(
|
||||
request.log,
|
||||
{ fetchHandlers, processHandlers },
|
||||
notification,
|
||||
);
|
||||
} catch (err) {
|
||||
const status = googleStatus(err) || err.statusCode || 500;
|
||||
if (status === 401 || status === 403) {
|
||||
throw new HttpError(
|
||||
500,
|
||||
"Gmail access expired or was revoked. Sign in again.",
|
||||
);
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
|
||||
return reply.code(204).send();
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user