feat(cloud): replace novu for all notification triggers
This commit is contained in:
@@ -1,4 +1,9 @@
|
||||
import { describe, it, expect, vi, beforeEach } from "vitest";
|
||||
import { env } from "next-runtime-env";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
import * as memberRepo from "@kan/db/repository/member.repo";
|
||||
|
||||
import { createDatabaseHooks } from "./hooks";
|
||||
|
||||
vi.mock("next-runtime-env", () => ({
|
||||
env: vi.fn(),
|
||||
@@ -15,11 +20,11 @@ vi.mock("@kan/db/repository/user.repo", () => ({
|
||||
}));
|
||||
|
||||
vi.mock("@kan/email", () => ({
|
||||
notificationClient: null,
|
||||
createSubscriber: vi.fn(),
|
||||
triggerSubscriberWorkflow: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("@kan/shared", () => ({
|
||||
createEmailUnsubscribeLink: vi.fn(),
|
||||
createS3Client: vi.fn(),
|
||||
}));
|
||||
|
||||
@@ -27,17 +32,10 @@ vi.mock("@aws-sdk/client-s3", () => ({
|
||||
PutObjectCommand: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("@novu/api/models/components", () => ({
|
||||
ChatOrPushProviderEnum: { Discord: "discord" },
|
||||
}));
|
||||
|
||||
import { env } from "next-runtime-env";
|
||||
import * as memberRepo from "@kan/db/repository/member.repo";
|
||||
import { createDatabaseHooks } from "./hooks";
|
||||
|
||||
const mockEnv = env as ReturnType<typeof vi.fn>;
|
||||
const mockGetByEmailAndStatus =
|
||||
memberRepo.getByEmailAndStatus as ReturnType<typeof vi.fn>;
|
||||
const mockGetByEmailAndStatus = memberRepo.getByEmailAndStatus as ReturnType<
|
||||
typeof vi.fn
|
||||
>;
|
||||
|
||||
const db = {} as Parameters<typeof createDatabaseHooks>[0];
|
||||
|
||||
|
||||
@@ -1,14 +1,13 @@
|
||||
import { PutObjectCommand } from "@aws-sdk/client-s3";
|
||||
import { ChatOrPushProviderEnum } from "@novu/api/models/components";
|
||||
import { createAuthMiddleware } from "better-auth/api";
|
||||
import { env } from "next-runtime-env";
|
||||
|
||||
import type { dbClient } from "@kan/db/client";
|
||||
import * as memberRepo from "@kan/db/repository/member.repo";
|
||||
import * as userRepo from "@kan/db/repository/user.repo";
|
||||
import { createSubscriber, notificationClient } from "@kan/email";
|
||||
import { createSubscriber, triggerSubscriberWorkflow } from "@kan/email";
|
||||
import { createLogger } from "@kan/logger";
|
||||
import { createEmailUnsubscribeLink, createS3Client } from "@kan/shared";
|
||||
import { createS3Client } from "@kan/shared";
|
||||
|
||||
import { downloadImage } from "./utils";
|
||||
|
||||
@@ -94,80 +93,49 @@ export function createDatabaseHooks(db: dbClient) {
|
||||
}
|
||||
}
|
||||
|
||||
if (notificationClient) {
|
||||
const [firstName, ...rest] = (user.name || "")
|
||||
.split(" ")
|
||||
.filter(Boolean);
|
||||
const lastName = rest.length ? rest.join(" ") : undefined;
|
||||
const [firstName, ...rest] = (user.name || "")
|
||||
.split(" ")
|
||||
.filter(Boolean);
|
||||
const lastName = rest.length ? rest.join(" ") : undefined;
|
||||
|
||||
try {
|
||||
const avatarUrl = avatarKey
|
||||
? `${env("NEXT_PUBLIC_STORAGE_URL")}/${env("NEXT_PUBLIC_AVATAR_BUCKET_NAME")}/${avatarKey}`
|
||||
: undefined;
|
||||
try {
|
||||
const avatarUrl = avatarKey
|
||||
? `${env("NEXT_PUBLIC_STORAGE_URL")}/${env("NEXT_PUBLIC_AVATAR_BUCKET_NAME")}/${avatarKey}`
|
||||
: undefined;
|
||||
|
||||
const unsubscribeUrl = await createEmailUnsubscribeLink(user.id);
|
||||
await createSubscriber({
|
||||
publicId: user.id,
|
||||
email: user.email,
|
||||
externalId: user.id,
|
||||
firstName,
|
||||
lastName,
|
||||
name: user.name,
|
||||
attributes: {
|
||||
avatarUrl,
|
||||
emailVerified: user.emailVerified,
|
||||
stripeCustomerId: user.stripeCustomerId,
|
||||
createdAt: user.createdAt,
|
||||
updatedAt: user.updatedAt,
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
log.error({ err: error }, "Error creating subscriber");
|
||||
}
|
||||
|
||||
log.info(
|
||||
{
|
||||
workflowId: "user-signup",
|
||||
userId: user.id,
|
||||
email: user.email,
|
||||
},
|
||||
"Triggering Novu workflow",
|
||||
);
|
||||
await notificationClient.trigger({
|
||||
to: {
|
||||
subscriberId: user.id,
|
||||
firstName: firstName,
|
||||
lastName: lastName,
|
||||
email: user.email,
|
||||
avatar: avatarUrl,
|
||||
data: {
|
||||
emailVerified: user.emailVerified,
|
||||
stripeCustomerId: user.stripeCustomerId,
|
||||
createdAt: user.createdAt,
|
||||
updatedAt: user.updatedAt,
|
||||
},
|
||||
},
|
||||
payload: {
|
||||
emailUnsubscribeUrl: unsubscribeUrl,
|
||||
},
|
||||
workflowId: "user-signup",
|
||||
});
|
||||
log.info(
|
||||
{ workflowId: "user-signup", userId: user.id },
|
||||
"Novu workflow triggered",
|
||||
);
|
||||
|
||||
await notificationClient.subscribers.credentials.update(
|
||||
{
|
||||
providerId: ChatOrPushProviderEnum.Discord,
|
||||
credentials: {
|
||||
webhookUrl: env("DISCORD_WEBHOOK_URL"),
|
||||
},
|
||||
integrationIdentifier: "discord",
|
||||
},
|
||||
user.id,
|
||||
);
|
||||
} catch (error) {
|
||||
log.error(
|
||||
{ err: error },
|
||||
"Error adding user to notification client",
|
||||
);
|
||||
}
|
||||
|
||||
try {
|
||||
await createSubscriber({
|
||||
publicId: user.id,
|
||||
email: user.email,
|
||||
externalId: user.id,
|
||||
firstName,
|
||||
lastName,
|
||||
name: user.name,
|
||||
});
|
||||
} catch (error) {
|
||||
log.error({ err: error }, "Error creating subscriber");
|
||||
}
|
||||
try {
|
||||
log.info(
|
||||
{ workflowId: "user-signup", userId: user.id, email: user.email },
|
||||
"Triggering user-signup workflow",
|
||||
);
|
||||
await triggerSubscriberWorkflow("user-signup", {
|
||||
publicId: user.id,
|
||||
});
|
||||
log.info(
|
||||
{ workflowId: "user-signup", userId: user.id },
|
||||
"user-signup workflow triggered",
|
||||
);
|
||||
} catch (error) {
|
||||
log.error({ err: error }, "Error triggering user-signup workflow");
|
||||
}
|
||||
},
|
||||
},
|
||||
|
||||
@@ -3,9 +3,8 @@ import type Stripe from "stripe";
|
||||
|
||||
import type { dbClient } from "@kan/db/client";
|
||||
import * as userRepo from "@kan/db/repository/user.repo";
|
||||
import { notificationClient } from "@kan/email";
|
||||
import { triggerSubscriberWorkflow } from "@kan/email";
|
||||
import { createLogger } from "@kan/logger";
|
||||
import { createEmailUnsubscribeLink } from "@kan/shared";
|
||||
|
||||
const log = createLogger("auth");
|
||||
|
||||
@@ -24,30 +23,22 @@ export async function triggerWorkflow(
|
||||
cancellationDetails?: Stripe.Subscription.CancellationDetails | null,
|
||||
) {
|
||||
try {
|
||||
if (!subscription.stripeCustomerId || !notificationClient) return;
|
||||
if (!subscription.stripeCustomerId) return;
|
||||
|
||||
const user = await userRepo.getByStripeCustomerId(
|
||||
db,
|
||||
subscription.stripeCustomerId,
|
||||
);
|
||||
|
||||
if (!user || !notificationClient) return;
|
||||
if (!user) return;
|
||||
|
||||
const unsubscribeUrl = await createEmailUnsubscribeLink(user.id);
|
||||
|
||||
log.info({ workflowId, userId: user.id }, "Triggering Novu workflow");
|
||||
await notificationClient.trigger({
|
||||
to: {
|
||||
subscriberId: user.id,
|
||||
},
|
||||
payload: {
|
||||
...subscription,
|
||||
cancellationDetails,
|
||||
emailUnsubscribeUrl: unsubscribeUrl,
|
||||
},
|
||||
log.info({ workflowId, userId: user.id }, "Triggering workflow");
|
||||
await triggerSubscriberWorkflow(
|
||||
workflowId,
|
||||
});
|
||||
log.info({ workflowId, userId: user.id }, "Novu workflow triggered");
|
||||
{ publicId: user.id },
|
||||
{ ...subscription, cancellationDetails },
|
||||
);
|
||||
log.info({ workflowId, userId: user.id }, "Workflow triggered");
|
||||
} catch (error) {
|
||||
log.error({ err: error, workflowId }, "Error triggering workflow");
|
||||
}
|
||||
|
||||
@@ -24,7 +24,6 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@kan/logger": "workspace:^",
|
||||
"@novu/api": "^3.11.0",
|
||||
"@react-email/components": "^1.0.1",
|
||||
"nodemailer": "^7.0.3",
|
||||
"react-email": "^5.0.6"
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
export const name = "email";
|
||||
|
||||
export { sendEmail } from "./sendEmail";
|
||||
export { notificationClient } from "./notificationClient";
|
||||
export {
|
||||
createSubscriber,
|
||||
updateSubscriberPreferences,
|
||||
triggerSubscriberWorkflow,
|
||||
} from "./subscriberClient";
|
||||
|
||||
@@ -1,6 +0,0 @@
|
||||
import { Novu } from "@novu/api";
|
||||
|
||||
export const notificationClient =
|
||||
process.env.NEXT_PUBLIC_KAN_ENV === "cloud" && process.env.NOVU_API_KEY
|
||||
? new Novu({ secretKey: process.env.NOVU_API_KEY })
|
||||
: null;
|
||||
@@ -14,6 +14,46 @@ export const subscriberClient =
|
||||
}
|
||||
: null;
|
||||
|
||||
async function subscriberRequest(
|
||||
method: string,
|
||||
path: string,
|
||||
body: unknown,
|
||||
errorMessage: string,
|
||||
) {
|
||||
if (!subscriberClient) return;
|
||||
|
||||
const url = `${subscriberClient.apiUrl}${path}`;
|
||||
|
||||
log.debug({ method, url, body }, "subscriber.dev request");
|
||||
|
||||
try {
|
||||
const response = await fetch(url, {
|
||||
method,
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"X-API-Key": subscriberClient.apiKey,
|
||||
},
|
||||
body: JSON.stringify(body),
|
||||
});
|
||||
|
||||
const responseBody = await response.text().catch(() => undefined);
|
||||
|
||||
log.debug(
|
||||
{ method, url, status: response.status, body: responseBody },
|
||||
"subscriber.dev response",
|
||||
);
|
||||
|
||||
if (!response.ok) {
|
||||
log.error(
|
||||
{ status: response.status, body: responseBody },
|
||||
errorMessage,
|
||||
);
|
||||
}
|
||||
} catch (error) {
|
||||
log.error({ err: error }, errorMessage);
|
||||
}
|
||||
}
|
||||
|
||||
interface CreateSubscriberInput {
|
||||
publicId: string;
|
||||
email: string;
|
||||
@@ -21,36 +61,18 @@ interface CreateSubscriberInput {
|
||||
firstName?: string;
|
||||
lastName?: string;
|
||||
name?: string;
|
||||
attributes?: Record<string, unknown>;
|
||||
}
|
||||
|
||||
export async function createSubscriber(input: CreateSubscriberInput) {
|
||||
if (!subscriberClient) return;
|
||||
|
||||
try {
|
||||
const response = await fetch(
|
||||
`${subscriberClient.apiUrl}/environments/${subscriberClient.environmentId}/subscribers`,
|
||||
{
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"X-API-Key": subscriberClient.apiKey,
|
||||
},
|
||||
body: JSON.stringify(input),
|
||||
},
|
||||
);
|
||||
|
||||
if (!response.ok) {
|
||||
log.error(
|
||||
{
|
||||
status: response.status,
|
||||
body: await response.text().catch(() => undefined),
|
||||
},
|
||||
"Failed to create subscriber.dev subscriber",
|
||||
);
|
||||
}
|
||||
} catch (error) {
|
||||
log.error({ err: error }, "Error creating subscriber.dev subscriber");
|
||||
}
|
||||
await subscriberRequest(
|
||||
"POST",
|
||||
`/environments/${subscriberClient.environmentId}/subscribers`,
|
||||
input,
|
||||
"Failed to create subscriber.dev subscriber",
|
||||
);
|
||||
}
|
||||
|
||||
interface UpdateSubscriberPreferencesInput {
|
||||
@@ -63,29 +85,31 @@ export async function updateSubscriberPreferences(
|
||||
) {
|
||||
if (!subscriberClient) return;
|
||||
|
||||
try {
|
||||
const response = await fetch(
|
||||
`${subscriberClient.apiUrl}/environments/${subscriberClient.environmentId}/subscribers/${subscriberId}/preferences`,
|
||||
{
|
||||
method: "PATCH",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"X-API-Key": subscriberClient.apiKey,
|
||||
},
|
||||
body: JSON.stringify(input),
|
||||
},
|
||||
);
|
||||
|
||||
if (!response.ok) {
|
||||
log.error(
|
||||
{
|
||||
status: response.status,
|
||||
body: await response.text().catch(() => undefined),
|
||||
},
|
||||
"Failed to update subscriber preferences",
|
||||
);
|
||||
}
|
||||
} catch (error) {
|
||||
log.error({ err: error }, "Error updating subscriber.dev preferences");
|
||||
}
|
||||
await subscriberRequest(
|
||||
"PATCH",
|
||||
`/environments/${subscriberClient.environmentId}/subscribers/${subscriberId}/preferences`,
|
||||
input,
|
||||
"Failed to update subscriber preferences",
|
||||
);
|
||||
}
|
||||
|
||||
interface TriggerWorkflowSubscriberInput {
|
||||
publicId?: string;
|
||||
externalId?: string;
|
||||
email?: string;
|
||||
}
|
||||
|
||||
export async function triggerSubscriberWorkflow(
|
||||
key: string,
|
||||
subscriber: TriggerWorkflowSubscriberInput,
|
||||
payload?: Record<string, unknown>,
|
||||
) {
|
||||
if (!subscriberClient) return;
|
||||
|
||||
await subscriberRequest(
|
||||
"POST",
|
||||
`/environments/${subscriberClient.environmentId}/workflows/trigger`,
|
||||
{ key, subscriber, payload },
|
||||
"Failed to trigger subscriber.dev workflow",
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1,32 +0,0 @@
|
||||
import { SignJWT } from "jose";
|
||||
import { env } from "next-runtime-env";
|
||||
|
||||
const encoder = new TextEncoder();
|
||||
|
||||
/**
|
||||
* Creates a long‑lived unsubscribe link for a given user/subscriber.
|
||||
*
|
||||
* `${NEXT_PUBLIC_BASE_URL}/unsubscribe?token=<jwt>`
|
||||
*
|
||||
* The JWT payload only contains the subscriberId. There is no expiry
|
||||
* on purpose – unsubscribe links should remain valid indefinitely.
|
||||
*
|
||||
*/
|
||||
export async function createEmailUnsubscribeLink(
|
||||
userId: string,
|
||||
): Promise<string | null> {
|
||||
const baseUrl = env("NEXT_PUBLIC_BASE_URL");
|
||||
const secret = process.env.EMAIL_UNSUBSCRIBE_SECRET;
|
||||
|
||||
if (!baseUrl || !secret) {
|
||||
// Environment not configured for unsubscribe links.
|
||||
return null;
|
||||
}
|
||||
|
||||
const token = await new SignJWT({ subscriberId: userId })
|
||||
.setProtectedHeader({ alg: "HS256" })
|
||||
// No expiration on purpose; unsubscribe links are long‑lived.
|
||||
.sign(encoder.encode(secret));
|
||||
|
||||
return `${baseUrl}/unsubscribe?token=${encodeURIComponent(token)}`;
|
||||
}
|
||||
@@ -2,7 +2,6 @@ export * from "./generateUID";
|
||||
export * from "./generateSlug";
|
||||
export * from "./generateWorkspacePrefix";
|
||||
export * from "./subscriptions";
|
||||
export * from "./email";
|
||||
export * from "./dueDateFilters";
|
||||
export * from "./s3";
|
||||
export * from "./mentions";
|
||||
|
||||
Reference in New Issue
Block a user