Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions apps/web/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,8 @@ LLM_API_KEY=

# Tinybird
TINYBIRD_TOKEN=
# Token with DATASOURCES:CREATE, used to delete a user's data on account deletion
TINYBIRD_DELETE_TOKEN=
TINYBIRD_BASE_URL=https://api.us-east.tinybird.co/

# Stripe AI overage billing (optional; unset = unlimited/no overage charges)
Expand Down
1 change: 1 addition & 0 deletions apps/web/env.ts
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,7 @@ const parsedEnv = createEnv({
APNS_TRANSPORT: z.enum(["apns", "fake"]).optional(),

TINYBIRD_TOKEN: z.string().optional(),
TINYBIRD_DELETE_TOKEN: z.string().optional(),
TINYBIRD_BASE_URL: z.string().default("https://api.us-east.tinybird.co/"),

API_KEY_SALT: z.string().optional(),
Expand Down
104 changes: 104 additions & 0 deletions apps/web/scripts/purge-orphaned-tinybird-data.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
// Deletes Tinybird rows that identify mailboxes that no longer exist,
// left behind before account deletion cleaned up Tinybird.
//
// Run with TINYBIRD_DELETE_TOKEN set to a token that can also read the datasources
// (for example the workspace admin token).
// Dry run (counts only): `pnpm --filter inbox-zero-ai exec tsx scripts/purge-orphaned-tinybird-data.ts`
// Delete: `pnpm --filter inbox-zero-ai exec tsx scripts/purge-orphaned-tinybird-data.ts --apply`

import "dotenv/config";
import chunk from "lodash/chunk";
import { deleteTinybirdEmailData } from "@inboxzero/tinybird";
import prisma from "@/utils/prisma";

const EMAIL_DATASOURCES = [
"email_action",
"email",
"last_and_oldest_emails_mv",
];
const BATCH_SIZE = 100;

async function main() {
const apply = process.argv.includes("--apply");
if (!process.env.TINYBIRD_TOKEN || !process.env.TINYBIRD_DELETE_TOKEN) {
throw new Error("TINYBIRD_TOKEN and TINYBIRD_DELETE_TOKEN must be set");
}

const emailAccounts = await prisma.emailAccount.findMany({
select: { email: true },
});
const liveEmails = new Set(
emailAccounts.map(({ email }) => email.toLowerCase()),
);

const orphanedEmails = new Set<string>();
for (const datasource of EMAIL_DATASOURCES) {
const owners = await getDistinctValues(datasource, "ownerEmail");
if (!owners) {
console.log(`${datasource}: datasource not found, skipping`);
continue;
}
const orphaned = owners.filter(
(email) => !liveEmails.has(email.toLowerCase()),
);
for (const email of orphaned) orphanedEmails.add(email);
console.log(
`${datasource}: ${owners.length} mailboxes, ${orphaned.length} orphaned`,
);
}

// AI usage rows are kept; only rows that fell back to an email address as
// the user id identify a person.
const aiCallUsers = await getDistinctValues("aiCall", "userId");
const orphanedAiCallEmails =
aiCallUsers?.filter(
(value) => value.includes("@") && !liveEmails.has(value.toLowerCase()),
) ?? [];
for (const email of orphanedAiCallEmails) orphanedEmails.add(email);
console.log(
aiCallUsers
? `aiCall: ${orphanedAiCallEmails.length} orphaned email-keyed users`
: "aiCall: datasource not found, skipping",
);

if (!apply) {
console.log("Dry run. Re-run with --apply to delete.");
return;
}

for (const emails of chunk([...orphanedEmails], BATCH_SIZE)) {
await deleteTinybirdEmailData(emails);
}
console.log("Delete jobs submitted.");
}

async function getDistinctValues(datasource: string, column: string) {
const url = new URL(
"/v0/sql",
process.env.TINYBIRD_BASE_URL || "https://api.us-east.tinybird.co/",
);
url.searchParams.set(
"q",
`SELECT DISTINCT ${column} AS value FROM ${datasource} FORMAT JSON`,
);

const response = await fetch(url, {
headers: { Authorization: `Bearer ${process.env.TINYBIRD_DELETE_TOKEN}` },
});
if (response.status === 404) return null;
if (!response.ok) {
throw new Error(
`Tinybird query on ${datasource} failed: [${response.status}] ${await response.text()}`,
);
}

const body = (await response.json()) as { data: { value: string }[] };
return body.data.map(({ value }) => value);
}

main()
.catch((error) => {
console.error(error);
process.exitCode = 1;
})
.finally(() => prisma.$disconnect());
7 changes: 7 additions & 0 deletions apps/web/utils/actions/user.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { deleteAccountUploadDirectory } from "@/utils/mail-api/upload-blobs";
import { deleteUser } from "@/utils/user/delete";
import { clearLastEmailAccountCookie } from "@/utils/cookies.server";
import { LAST_EMAIL_ACCOUNT_COOKIE } from "@/utils/cookies";
import { deleteTinybirdEmailData } from "@inboxzero/tinybird";
import { deleteAccountAction, deleteEmailAccountAction } from "./user";

vi.mock("@/utils/prisma");
Expand Down Expand Up @@ -40,6 +41,9 @@ vi.mock("@/utils/cookies.server", () => ({
vi.mock("@/utils/mail-api/upload-blobs", () => ({
deleteAccountUploadDirectory: vi.fn(() => Promise.resolve()),
}));
vi.mock("@inboxzero/tinybird", () => ({
deleteTinybirdEmailData: vi.fn(() => Promise.resolve()),
}));
vi.mock("@/utils/user/delete", () => ({
deleteUser: vi.fn(),
}));
Expand Down Expand Up @@ -221,6 +225,9 @@ describe("deleteEmailAccountAction", () => {
expect(prisma.$transaction.mock.invocationCallOrder[0]).toBeLessThan(
vi.mocked(deleteAccountUploadDirectory).mock.invocationCallOrder[0],
);
expect(deleteTinybirdEmailData).toHaveBeenCalledWith([
"secondary@example.com",
]);
expect(clearLastEmailAccountCookie).not.toHaveBeenCalled();
});

Expand Down
8 changes: 8 additions & 0 deletions apps/web/utils/actions/user.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import {
} from "@/utils/actions/user.validation";
import { clearLastEmailAccountCookie } from "@/utils/cookies.server";
import { deleteAccountUploadDirectory } from "@/utils/mail-api/upload-blobs";
import { deleteTinybirdEmailData } from "@inboxzero/tinybird";
import { aliasPosthogUser } from "@/utils/posthog";
import {
cleanupAIDraftsForAccount,
Expand Down Expand Up @@ -242,6 +243,13 @@ export const deleteEmailAccountAction = actionClientUser
);
}

after(() =>
deleteTinybirdEmailData([emailAccount.email]).catch((error) => {
logger.error("Error deleting Tinybird data", { error });
captureException(error);
}),
);

await deleteAccountUploadDirectory(emailAccountId).catch((error) => {
logger.error("Failed to delete account mail uploads", {
error,
Expand Down
29 changes: 27 additions & 2 deletions apps/web/utils/user/delete.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { createTestLogger } from "@/__tests__/helpers";
import prisma from "@/utils/__mocks__/prisma";
import { deleteAccountUploadDirectory } from "@/utils/mail-api/upload-blobs";
import { deleteUser } from "@/utils/user/delete";
import { deleteTinybirdEmailData } from "@inboxzero/tinybird";
import { createEmailProvider } from "@/utils/email/provider";

vi.mock("@/utils/prisma");
Expand All @@ -20,8 +21,8 @@ vi.mock("@inboxzero/loops", () => ({
vi.mock("@inboxzero/transactional-email", () => ({
deleteContact: vi.fn(),
}));
vi.mock("@inboxzero/tinybird-ai-analytics", () => ({
deleteTinybirdAiCalls: vi.fn(() => Promise.resolve()),
vi.mock("@inboxzero/tinybird", () => ({
deleteTinybirdEmailData: vi.fn(() => Promise.resolve()),
}));
vi.mock("@/utils/posthog", () => ({
deletePosthogUser: vi.fn(() => Promise.resolve()),
Expand Down Expand Up @@ -83,6 +84,29 @@ describe("deleteUser", () => {
expect(deleteAccountUploadDirectory).not.toHaveBeenCalled();
});

it("keeps Tinybird data when deleting the user fails", async () => {
prisma.account.findMany.mockResolvedValue([
{
provider: "google",
access_token: null,
refresh_token: null,
expires_at: null,
emailAccount: {
id: "email-account-1",
email: "owner@example.com",
watchEmailsSubscriptionId: null,
},
},
] as Awaited<ReturnType<typeof prisma.account.findMany>>);
prisma.executedRule.findMany.mockResolvedValue([]);
prisma.user.deleteMany.mockRejectedValue(new Error("database unavailable"));

await expect(deleteUser({ userId: "user-1", logger })).rejects.toThrow(
"database unavailable",
);
expect(deleteTinybirdEmailData).not.toHaveBeenCalled();
});

it("deletes solo organizations before deleting the user", async () => {
prisma.account.findMany.mockResolvedValue([
{
Expand Down Expand Up @@ -111,6 +135,7 @@ describe("deleteUser", () => {
prisma.user.deleteMany.mockResolvedValue({ count: 1 } as any);

await deleteUser({ userId: "user-1", logger });
expect(deleteTinybirdEmailData).toHaveBeenCalledWith(["owner@example.com"]);
expect(withThreadPageBufferDeletion).toHaveBeenCalledWith(
["email-account-1"],
expect.any(Function),
Expand Down
21 changes: 12 additions & 9 deletions apps/web/utils/user/delete.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@ import { deleteContact as deleteLoopsContact } from "@inboxzero/loops";
import { deleteContact as deleteResendContact } from "@inboxzero/transactional-email";
import { withThreadPageBufferDeletion } from "@/utils/redis/thread-page-buffer";
import prisma from "@/utils/prisma";
import { deleteTinybirdAiCalls } from "@inboxzero/tinybird-ai-analytics";
import { deleteTinybirdEmailData } from "@inboxzero/tinybird";
import { after } from "next/server";
import {
deletePosthogUser,
trackUserDeleted,
Expand Down Expand Up @@ -69,14 +70,6 @@ export async function deleteUser({
captureException(error);
});

deleteTinybirdAiCalls({ userId }).catch((error) => {
logger.error("Error deleting Tinybird AI calls", {
error,
userId,
});
captureException(error);
});

clearCachedResearchForUser(userId).catch((error) => {
logger.error("Error clearing cached research", { error });
captureException(error);
Expand Down Expand Up @@ -140,6 +133,16 @@ export async function deleteUser({
throw originalError;
}
});

const emails = accounts
.map((account) => account.emailAccount?.email)
.filter((email): email is string => Boolean(email));
after(() =>
deleteTinybirdEmailData(emails).catch((error) => {
logger.error("Error deleting Tinybird data", { error });
captureException(error);
}),
);
} catch (error) {
logger.error("Error during user resources deletion process", {
error,
Expand Down
41 changes: 0 additions & 41 deletions packages/tinybird-ai-analytics/src/delete.ts

This file was deleted.

1 change: 0 additions & 1 deletion packages/tinybird-ai-analytics/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
// biome-ignore lint/performance/noBarrelFile: fix later
export * from "./client";
export * from "./publish";
export * from "./delete";
export * from "./query";
6 changes: 5 additions & 1 deletion packages/tinybird/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,9 @@
"name": "@inboxzero/tinybird",
"version": "0.0.0",
"main": "src/index.ts",
"scripts": {
"test": "vitest run"
},
"dependencies": {
"@chronark/zod-bird": "1.0.0",
"p-retry": "8.0.1",
Expand All @@ -10,6 +13,7 @@
"devDependencies": {
"@types/node": "24.10.1",
"tsconfig": "workspace:*",
"typescript": "6.0.3"
"typescript": "6.0.3",
"vitest": "4.1.10"
}
}
Loading
Loading