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
8 changes: 3 additions & 5 deletions .github/workflows/playwright.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -130,18 +130,16 @@ jobs:

- name: Start QStash dev server
run: |
docker run -d --name qstash-dev \
--network host \
public.ecr.aws/upstash/qstash:latest \
qstash dev
nohup npx --yes @upstash/qstash-cli@2.37.18 dev > /tmp/qstash.log 2>&1 &
for i in $(seq 1 30); do
if curl -s -o /dev/null -w "%{http_code}" http://127.0.0.1:8080 | grep -qE '^[0-9]{3}$'; then
echo "QStash is up"
exit 0
fi
sleep 1
done
docker logs qstash-dev
echo "QStash failed to start"
cat /tmp/qstash.log
exit 1

- name: Configure Tinybird Local
Expand Down
3 changes: 3 additions & 0 deletions apps/web/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,9 @@ NEXTAUTH_URL=http://localhost:8888 # (only needed for localhost)

# Secret for Vercel cron jobs + sync-embeddings
CRON_SECRET=
# Shared with the LoopWork demo. Mints geo-accurate clicks via POST /api/demo/click
# and backdated commissions via POST /api/demo/commission
DEMO_CLICK_SECRET=
# Encryption key (AES-256-GCM) for encrypting sensitive data in the database
ENCRYPTION_KEY=
# Email unsubscribe token secret (optional, falls back to NEXTAUTH_SECRET)
Expand Down
148 changes: 68 additions & 80 deletions apps/web/app/(ee)/api/cron/invoices/retry-failed/route.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { prisma } from "@/lib/prisma";
import { createPaymentIntent } from "@/lib/stripe/create-payment-intent";
import { ACME_WORKSPACE_ID, DUB_WORKSPACE_ID } from "@dub/utils";
Expand All @@ -12,93 +11,82 @@ const schema = z.object({
});

// POST /api/cron/invoices/retry-failed
export async function POST(req: Request) {
try {
const rawBody = await req.text();

await verifyQstashSignature({
req,
rawBody,
});

const { invoiceId } = schema.parse(JSON.parse(rawBody));

const invoice = await prisma.invoice.findUnique({
where: {
id: invoiceId,
},
select: {
id: true,
type: true,
status: true,
total: true,
failedAttempts: true,
workspace: {
select: {
id: true,
stripeId: true,
},
export const POST = withCron(async ({ rawBody }) => {
const { invoiceId } = schema.parse(JSON.parse(rawBody));

const invoice = await prisma.invoice.findUnique({
where: {
id: invoiceId,
},
select: {
id: true,
type: true,
status: true,
total: true,
failedAttempts: true,
workspace: {
select: {
id: true,
stripeId: true,
},
},
});
},
});

if (!invoice) {
console.log(`Invoice ${invoiceId} not found.`);
return new Response(`Invoice ${invoiceId} not found.`);
}
if (!invoice) {
console.log(`Invoice ${invoiceId} not found.`);
return new Response(`Invoice ${invoiceId} not found.`);
}

if (invoice.status !== "failed") {
console.log(`Invoice ${invoiceId} is not failed.`);
return new Response(`Invoice ${invoiceId} is not failed.`);
}
if (invoice.status !== "failed") {
console.log(`Invoice ${invoiceId} is not failed.`);
return new Response(`Invoice ${invoiceId} is not failed.`);
}

if (invoice.failedAttempts >= 3) {
console.log(`Invoice ${invoiceId} has reached max failed attempts of 3.`);
return new Response(
`Invoice ${invoiceId} has reached max failed attempts of 3.`,
);
}
if (invoice.failedAttempts >= 3) {
console.log(`Invoice ${invoiceId} has reached max failed attempts of 3.`);
return new Response(
`Invoice ${invoiceId} has reached max failed attempts of 3.`,
);
}

if (invoice.type !== "domainRenewal") {
console.log(`Only domain renewals can be retried at this time.`);
return new Response(`Only domain renewals can be retried at this time.`);
}
if (invoice.type !== "domainRenewal") {
console.log(`Only domain renewals can be retried at this time.`);
return new Response(`Only domain renewals can be retried at this time.`);
}

let { workspace } = invoice;
let { workspace } = invoice;

// If Acme workspace, use Dub workspace stripeId
if (workspace.id === ACME_WORKSPACE_ID) {
const dubWorkspace = await prisma.project.findUniqueOrThrow({
where: {
id: DUB_WORKSPACE_ID,
},
select: {
stripeId: true,
},
});
// If Acme workspace, use Dub workspace stripeId
if (workspace.id === ACME_WORKSPACE_ID) {
const dubWorkspace = await prisma.project.findUniqueOrThrow({
where: {
id: DUB_WORKSPACE_ID,
},
select: {
stripeId: true,
},
});

workspace = {
...workspace,
stripeId: dubWorkspace.stripeId,
};
}
workspace = {
...workspace,
stripeId: dubWorkspace.stripeId,
};
}

if (!workspace.stripeId) {
console.log(`Workspace ${workspace.id} has no stripeId.`);
return new Response(`Workspace ${workspace.id} has no stripeId.`);
}
if (!workspace.stripeId) {
console.log(`Workspace ${workspace.id} has no stripeId.`);
return new Response(`Workspace ${workspace.id} has no stripeId.`);
}

await createPaymentIntent({
stripeId: workspace.stripeId,
amount: invoice.total,
invoiceId: invoice.id,
statementDescriptor: "DUB.CO DOMAIN RENEWAL",
description: `Domain renewal invoice (${invoice.id})`,
idempotencyKey: `${invoice.id}-${invoice.failedAttempts}`,
});
await createPaymentIntent({
stripeId: workspace.stripeId,
amount: invoice.total,
invoiceId: invoice.id,
statementDescriptor: "DUB.CO DOMAIN RENEWAL",
description: `Domain renewal invoice (${invoice.id})`,
idempotencyKey: `${invoice.id}-${invoice.failedAttempts}`,
});

return new Response(`Retrying invoice charge ${invoice.id}...`);
} catch (error) {
return handleAndReturnErrorResponse(error);
}
}
return new Response(`Retrying invoice charge ${invoice.id}...`);
});
12 changes: 4 additions & 8 deletions apps/web/app/(ee)/api/cron/payouts/balance-available/route.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { BANK_ACCOUNT_STATUS_DESCRIPTIONS } from "@/lib/constants/payouts";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { getPartnerBankAccount } from "@/lib/partners/get-partner-bank-account";
import { prisma } from "@/lib/prisma";
import { stripe } from "@/lib/stripe";
Expand All @@ -24,11 +23,8 @@ const payloadSchema = z.object({
});

// POST /api/cron/payouts/balance-available
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();
await verifyQstashSignature({ req, rawBody });

const { stripeAccount } = payloadSchema.parse(JSON.parse(rawBody));

const partner = await prisma.partner.findUnique({
Expand Down Expand Up @@ -222,6 +218,6 @@ export async function POST(req: Request) {
type: "errors",
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
12 changes: 4 additions & 8 deletions apps/web/app/(ee)/api/cron/payouts/charge-succeeded/route.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { prisma } from "@/lib/prisma";
import { log } from "@dub/utils";
import { PartnerPayoutMethod } from "@prisma/client";
Expand All @@ -21,11 +20,8 @@ const payloadSchema = z.object({
// POST /api/cron/payouts/charge-succeeded
// This route is used to process the charge-succeeded event from Stripe.
// We're intentionally offloading this to a cron job so we can return a 200 to Stripe immediately.
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();
await verifyQstashSignature({ req, rawBody });

const { invoiceId } = payloadSchema.parse(JSON.parse(rawBody));

const invoice = await prisma.invoice.findUnique({
Expand Down Expand Up @@ -131,6 +127,6 @@ export async function POST(req: Request) {
type: "cron",
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
12 changes: 4 additions & 8 deletions apps/web/app/(ee)/api/cron/payouts/payout-failed/route.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { getPartnerBankAccount } from "@/lib/partners/get-partner-bank-account";
import { prisma } from "@/lib/prisma";
import { sendEmail } from "@dub/email";
Expand All @@ -19,11 +18,8 @@ const payloadSchema = z.object({
});

// POST /api/cron/payouts/payout-failed
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();
await verifyQstashSignature({ req, rawBody });

const { stripeAccount, stripePayout } = payloadSchema.parse(
JSON.parse(rawBody),
);
Expand Down Expand Up @@ -86,6 +82,6 @@ export async function POST(req: Request) {
type: "errors",
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
12 changes: 4 additions & 8 deletions apps/web/app/(ee)/api/cron/payouts/payout-paid/route.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { prisma } from "@/lib/prisma";
import { sendEmail } from "@dub/email";
import PartnerPayoutWithdrawalCompleted from "@dub/email/templates/partner-payout-withdrawal-completed";
Expand All @@ -19,11 +18,8 @@ const payloadSchema = z.object({
});

// POST /api/cron/payouts/payout-paid
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();
await verifyQstashSignature({ req, rawBody });

const { stripeAccount, stripePayout } = payloadSchema.parse(
JSON.parse(rawBody),
);
Expand Down Expand Up @@ -83,6 +79,6 @@ export async function POST(req: Request) {
type: "errors",
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
13 changes: 4 additions & 9 deletions apps/web/app/(ee)/api/cron/payouts/process/route.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { CUTOFF_PERIOD_ENUM } from "@/lib/partners/cutoff-period";
import { prisma } from "@/lib/prisma";
import { log } from "@dub/utils";
Expand All @@ -24,12 +23,8 @@ const processPayoutsCronSchema = z.object({
// POST /api/cron/payouts/process
// This route is used to process payouts for a given invoice
// we're intentionally offloading this to a cron job to avoid blocking the main thread
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();

await verifyQstashSignature({ req, rawBody });

const {
workspaceId,
userId,
Expand Down Expand Up @@ -107,6 +102,6 @@ export async function POST(req: Request) {
mention: true,
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
16 changes: 4 additions & 12 deletions apps/web/app/(ee)/api/cron/payouts/process/updates/route.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import { recordAuditLog } from "@/lib/api/audit-logs/record-audit-log";
import { handleAndReturnErrorResponse } from "@/lib/api/errors";
import { qstash } from "@/lib/cron";
import { verifyQstashSignature } from "@/lib/cron/verify-qstash";
import { withCron } from "@/lib/cron/with-cron";
import { prisma } from "@/lib/prisma";
import { sendBatchEmail } from "@dub/email";
import PartnerPayoutConfirmed from "@dub/email/templates/partner-payout-confirmed";
Expand All @@ -20,15 +19,8 @@ const BATCH_SIZE = 100;

// POST /api/cron/payouts/process/updates
// Recursive cron job to handle side effects of the `cron/payouts/process` job (recordAuditLog, sendBatchEmails)
export async function POST(req: Request) {
export const POST = withCron(async ({ rawBody }) => {
try {
const rawBody = await req.text();

await verifyQstashSignature({
req,
rawBody,
});

const { invoiceId, startingAfter } = payloadSchema.parse(
JSON.parse(rawBody),
);
Expand Down Expand Up @@ -149,6 +141,6 @@ export async function POST(req: Request) {
mention: true,
});

return handleAndReturnErrorResponse(error);
throw error;
}
}
});
Loading
Loading