From 5083a67ff45966fe4f97d209d82600ca3d562f27 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Tue, 28 Jul 2026 22:57:09 +0530 Subject: [PATCH 1/8] Batch click stats updates through a Redis stream cron. --- .../cron/streams/update-click-stats/route.ts | 389 ++++++++++++++++++ apps/web/lib/tinybird/record-click.ts | 59 +-- .../lib/upstash/redis-streams/click-stats.ts | 74 ++++ apps/web/vercel.json | 4 + 4 files changed, 482 insertions(+), 44 deletions(-) create mode 100644 apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts create mode 100644 apps/web/lib/upstash/redis-streams/click-stats.ts diff --git a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts new file mode 100644 index 00000000000..32de415dd0a --- /dev/null +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -0,0 +1,389 @@ +import { qstash } from "@/lib/cron"; +import { withCron } from "@/lib/cron/with-cron"; +import { conn } from "@/lib/planetscale"; +import { redis } from "@/lib/upstash/redis"; +import { + ClickStatsEvent, + clickStatsStream, +} from "@/lib/upstash/redis-streams/click-stats"; +import { RedisStreamEntry } from "@/lib/upstash/redis-streams/client"; +import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils"; +import { format } from "date-fns"; +import { NextResponse } from "next/server"; +import { logAndRespond } from "../../utils"; + +export const dynamic = "force-dynamic"; + +const BATCH_SIZE = 10_000; // Max stream entries to consume per cron run +const SUB_BATCH_SIZE = 50; // DB updates to run in parallel within each batch +const BACKLOG_ALERT_THRESHOLD = 50_000; // Alert when stream length exceeds this +const BACKLOG_AGE_ALERT_MS = 5 * 60 * 1000; // Alert when oldest pending entry is older than 5 minutes +const LOCK_KEY = "lock:update-click-stats"; // Prevents concurrent GET/POST from double-counting +const LOCK_TTL_SECONDS = 600; // ≥ cron maxDuration (600s) so lock outlives a running invocation + +type LinkAggregate = { + linkId: string; + clicks: number; + lastClicked: number; + entryIds: string[]; +}; + +type WorkspaceAggregate = { + workspaceId: string; + clicks: number; +}; + +type EnrollmentAggregate = { + programId: string; + partnerId: string; + clicks: number; +}; + +const aggregateClickStats = ( + entries: RedisStreamEntry[], +): { + linkUpdates: LinkAggregate[]; + workspaceUpdates: WorkspaceAggregate[]; + enrollmentUpdates: EnrollmentAggregate[]; +} => { + const links = new Map(); + const workspaces = new Map(); + const enrollments = new Map(); + + for (const entry of entries) { + const { linkId, workspaceId, programId, partnerId, timestamp } = entry.data; + + if (!linkId) { + continue; + } + + const parsedTimestamp = Date.parse(timestamp); + const lastClicked = Number.isFinite(parsedTimestamp) + ? parsedTimestamp + : Date.now(); + + const existingLink = links.get(linkId); + if (existingLink) { + existingLink.clicks += 1; + existingLink.lastClicked = Math.max( + existingLink.lastClicked, + lastClicked, + ); + existingLink.entryIds.push(entry.id); + } else { + links.set(linkId, { + linkId, + clicks: 1, + lastClicked, + entryIds: [entry.id], + }); + } + + if (workspaceId) { + const existingWorkspace = workspaces.get(workspaceId); + if (existingWorkspace) { + existingWorkspace.clicks += 1; + } else { + workspaces.set(workspaceId, { + workspaceId, + clicks: 1, + }); + } + } + + if (programId && partnerId) { + const key = `${programId}:${partnerId}`; + const existingEnrollment = enrollments.get(key); + if (existingEnrollment) { + existingEnrollment.clicks += 1; + } else { + enrollments.set(key, { + programId, + partnerId, + clicks: 1, + }); + } + } + } + + return { + linkUpdates: Array.from(links.values()), + workspaceUpdates: Array.from(workspaces.values()), + enrollmentUpdates: Array.from(enrollments.values()), + }; +}; + +const processInSubBatches = async ( + items: T[], + handler: (item: T) => Promise<{ success: boolean; error?: unknown }>, +) => { + const errors: unknown[] = []; + let totalProcessed = 0; + + for (let i = 0; i < items.length; i += SUB_BATCH_SIZE) { + const batch = items.slice(i, i + SUB_BATCH_SIZE); + const results = await Promise.allSettled(batch.map(handler)); + + for (const result of results) { + if (result.status === "fulfilled" && result.value.success) { + totalProcessed++; + } else if (result.status === "fulfilled" && result.value.error) { + errors.push(result.value.error); + } else if (result.status === "rejected") { + errors.push(result.reason); + } + } + } + + return { totalProcessed, errors }; +}; + +const getStreamEntryAgeMs = (entryId: string | null) => { + if (!entryId) { + return null; + } + + const timestampMs = Number(entryId.split("-")[0]); + if (!Number.isFinite(timestampMs)) { + return null; + } + + return Date.now() - timestampMs; +}; + +const processClickStatsStreamBatch = () => + clickStatsStream.processBatch( + async (entries) => { + if (!entries || entries.length === 0) { + return { + success: true, + linkUpdates: [], + workspaceUpdates: [], + enrollmentUpdates: [], + processedEntryIds: [], + totalProcessed: 0, + errors: [], + }; + } + + console.log(`Aggregating ${entries.length} click stats events`); + + const { linkUpdates, workspaceUpdates, enrollmentUpdates } = + aggregateClickStats(entries); + + if (linkUpdates.length === 0) { + console.log("No click stats updates to process"); + return { + success: true, + linkUpdates: [], + workspaceUpdates: [], + enrollmentUpdates: [], + processedEntryIds: entries.map((entry) => entry.id), + totalProcessed: 0, + errors: [], + }; + } + + console.log( + `Processing ${linkUpdates.length} link, ${workspaceUpdates.length} workspace, ${enrollmentUpdates.length} enrollment click stats updates...`, + ); + + const processedEntryIds: string[] = []; + const errors: unknown[] = []; + + const linkResult = await processInSubBatches( + linkUpdates, + async (update) => { + try { + await conn.execute( + "UPDATE Link SET clicks = clicks + ?, lastClicked = ? WHERE id = ?", + [ + update.clicks, + format(new Date(update.lastClicked), "yyyy-MM-dd HH:mm:ss"), + update.linkId, + ], + ); + processedEntryIds.push(...update.entryIds); + return { success: true }; + } catch (error) { + console.error(`Failed to update link ${update.linkId}:`, error); + return { + success: false, + error: { linkId: update.linkId, error }, + }; + } + }, + ); + errors.push(...linkResult.errors); + + const workspaceResult = await processInSubBatches( + workspaceUpdates, + async (update) => { + try { + await conn.execute( + "UPDATE Project SET usage = usage + ?, totalClicks = totalClicks + ? WHERE id = ?", + [update.clicks, update.clicks, update.workspaceId], + ); + return { success: true }; + } catch (error) { + console.error( + `Failed to update workspace ${update.workspaceId}:`, + error, + ); + return { + success: false, + error: { workspaceId: update.workspaceId, error }, + }; + } + }, + ); + errors.push(...workspaceResult.errors); + + const enrollmentResult = await processInSubBatches( + enrollmentUpdates, + async (update) => { + try { + await conn.execute( + "UPDATE ProgramEnrollment SET totalClicks = totalClicks + ? WHERE programId = ? AND partnerId = ?", + [update.clicks, update.programId, update.partnerId], + ); + return { success: true }; + } catch (error) { + console.error( + `Failed to update program enrollment ${update.programId}:${update.partnerId}:`, + error, + ); + return { + success: false, + error: { + programId: update.programId, + partnerId: update.partnerId, + error, + }, + }; + } + }, + ); + errors.push(...enrollmentResult.errors); + + const totalProcessed = + linkResult.totalProcessed + + workspaceResult.totalProcessed + + enrollmentResult.totalProcessed; + + console.log( + `Processed ${linkResult.totalProcessed}/${linkUpdates.length} links, ${workspaceResult.totalProcessed}/${workspaceUpdates.length} workspaces, ${enrollmentResult.totalProcessed}/${enrollmentUpdates.length} enrollments`, + ); + + if (errors.length > 0) { + console.error( + `Encountered ${errors.length} errors while processing click stats:`, + errors.slice(0, 5), + ); + } + + return { + linkUpdates, + workspaceUpdates, + enrollmentUpdates, + errors, + totalProcessed, + processedEntryIds, + entriesProcessed: entries.length, + }; + }, + { + count: BATCH_SIZE, + deleteAfterRead: true, + }, + ); + +const maybeAlertOnBacklog = async (streamInfo: { + length: number; + firstEntryId: string | null; +}) => { + const ageMs = getStreamEntryAgeMs(streamInfo.firstEntryId); + const isBackloggedByLength = streamInfo.length > BACKLOG_ALERT_THRESHOLD; + const isBackloggedByAge = + ageMs !== null && ageMs > BACKLOG_AGE_ALERT_MS && streamInfo.length > 0; + + if (!isBackloggedByLength && !isBackloggedByAge) { + return; + } + + await log({ + message: `Click stats stream backlog alert: length=${streamInfo.length}, oldestAgeMs=${ageMs ?? "unknown"}, firstEntryId=${streamInfo.firstEntryId ?? "none"}`, + type: "alerts", + }); +}; + +const executeClickStatsCron = async () => { + const { + linkUpdates, + errors, + totalProcessed, + entriesProcessed = 0, + } = await processClickStatsStreamBatch(); + + const streamInfo = await clickStatsStream.getStreamInfo(); + await maybeAlertOnBacklog(streamInfo); + + const hasMore = + streamInfo.length > 0 || (entriesProcessed ?? 0) >= BATCH_SIZE; + + if (hasMore) { + await qstash.publishJSON({ + url: `${APP_DOMAIN_WITH_NGROK}/api/cron/streams/update-click-stats`, + method: "POST", + body: {}, + }); + } + + if (!linkUpdates.length) { + return NextResponse.json({ + success: true, + message: "No updates to process", + processed: 0, + streamInfo, + hasMore, + }); + } + + const response = { + success: true, + processed: totalProcessed, + errors: errors?.length || 0, + streamInfo, + hasMore, + message: `Successfully processed ${totalProcessed} click stats updates`, + }; + + console.log(response); + + return NextResponse.json(response); +}; + +const runWithLock = async () => { + const acquired = await redis.set(LOCK_KEY, "1", { + nx: true, + ex: LOCK_TTL_SECONDS, + }); + + if (!acquired) { + return logAndRespond( + "[update-click-stats] Another run is in progress. Skipping...", + ); + } + + try { + return await executeClickStatsCron(); + } finally { + await redis.del(LOCK_KEY); + } +}; + +// GET /api/cron/streams/update-click-stats +export const GET = withCron(async () => runWithLock()); + +// POST /api/cron/streams/update-click-stats (recursively called by QStash) +export const POST = withCron(async () => runWithLock()); diff --git a/apps/web/lib/tinybird/record-click.ts b/apps/web/lib/tinybird/record-click.ts index a5deaca0129..4080a7d3a57 100644 --- a/apps/web/lib/tinybird/record-click.ts +++ b/apps/web/lib/tinybird/record-click.ts @@ -12,11 +12,9 @@ import { recordClickCache } from "../api/links/record-click-cache"; import { detectBot } from "../middleware/utils/detect-bot"; import { detectQr } from "../middleware/utils/detect-qr"; import { getIdentityHash } from "../middleware/utils/get-identity-hash"; -import { conn } from "../planetscale"; import { redis } from "../upstash"; -import { publishPartnerActivityEvent } from "../upstash/redis-streams/partner-activity"; +import { publishClickStatsEvent } from "../upstash/redis-streams/click-stats"; import { publishWorkspaceClickEvent } from "../upstash/redis-streams/workspace-click-events"; -import { publishWorkspaceClicksUsageEvent } from "../upstash/redis-streams/workspace-clicks-usage"; /** * Recording clicks with geo, ua, referer and timestamp data @@ -183,44 +181,19 @@ export async function recordClick({ ).then((res) => res.json()), // cache the recorded click for the corresponding IP address in Redis for 1 hour - recordClickCache.set({ domain, key, identityHash, clickId }), - - // increment the click count for the link (based on their ID) - // we have to use planetscale connection directly (not prismaEdge) because of connection pooling - conn.execute( - "UPDATE Link SET clicks = clicks + 1, lastClicked = NOW() WHERE id = ?", - [linkId], - ), - // if the link is associated with a workspace + has a destination URL - // increment the usage count for the workspace - workspaceId && - url && - publishWorkspaceClicksUsageEvent({ - linkId, - workspaceId, - timestamp: clickData.timestamp, - }).catch(() => { - // Fallback on writing directly to the database - return conn.execute( - "UPDATE Project p JOIN Link l ON p.id = l.projectId SET p.usage = p.usage + 1, p.totalClicks = p.totalClicks + 1 WHERE l.id = ?", - [linkId], - ); - }), - - programId && - partnerId && - publishPartnerActivityEvent({ - programId, - partnerId, - eventType: "click", - timestamp: new Date().toISOString(), - }).catch(() => { - // Fallback on writing directly to the database - return conn.execute( - "UPDATE ProgramEnrollment SET totalClicks = totalClicks + 1 WHERE programId = ? AND partnerId = ?", - [programId, partnerId], - ); - }), + recordClickCache.set({ + domain, + key, + identityHash, + clickId, + }), + + publishClickStatsEvent({ + linkId, + timestamp: clickData.timestamp, + ...(workspaceId && url && { workspaceId }), + ...(programId && partnerId && { programId, partnerId }), + }), // Publish the click event publishWorkspaceClickEvent(clickData), @@ -234,9 +207,7 @@ export async function recordClick({ const operations = [ "Tinybird click event ingestion", "recordClickCache set", - "Link clicks increment", - "Workspace usage increment", - "Program enrollment totalClicks increment", + "Click stats stream publish", "Workspace click event publish", ]; return { diff --git a/apps/web/lib/upstash/redis-streams/click-stats.ts b/apps/web/lib/upstash/redis-streams/click-stats.ts new file mode 100644 index 00000000000..a1557ea38e7 --- /dev/null +++ b/apps/web/lib/upstash/redis-streams/click-stats.ts @@ -0,0 +1,74 @@ +import { logger, toErrorFields } from "@/lib/axiom/server"; +import { conn } from "@/lib/planetscale"; +import { redis } from "../redis"; +import { RedisStream } from "./client"; + +// TODO: clean up those legacy streams once remaining publishers +// are migrated (track-sale → workspace-clicks-usage; lead/sale/commission → +// partner-activity) or confirmed still needed. + +const CLICK_STATS_STREAM_KEY = "click:stats:updates"; + +export const clickStatsStream = new RedisStream(CLICK_STATS_STREAM_KEY); + +export interface ClickStatsEvent { + linkId: string; + timestamp: string; + workspaceId?: string; + programId?: string; + partnerId?: string; +} + +export const publishClickStatsEvent = async ({ + linkId, + timestamp, + workspaceId, + programId, + partnerId, +}: ClickStatsEvent) => { + const payload = { + linkId, + timestamp, + ...(workspaceId && { workspaceId }), + ...(programId && partnerId && { programId, partnerId }), + }; + + try { + return await redis.xadd(CLICK_STATS_STREAM_KEY, "*", payload); + } catch (error) { + logger.error("stream.publish_failed", { + service: "upstash", + streamKey: CLICK_STATS_STREAM_KEY, + error: toErrorFields(error), + correlation: { + linkId, + workspaceId, + programId, + partnerId, + }, + }); + + return await Promise.allSettled([ + conn.execute( + "UPDATE Link SET clicks = clicks + 1, lastClicked = NOW() WHERE id = ?", + [linkId], + ), + + workspaceId + ? conn.execute( + "UPDATE Project SET usage = usage + 1, totalClicks = totalClicks + 1 WHERE id = ?", + [workspaceId], + ) + : null, + + programId && partnerId + ? conn.execute( + "UPDATE ProgramEnrollment SET totalClicks = totalClicks + 1 WHERE programId = ? AND partnerId = ?", + [programId, partnerId], + ) + : null, + + logger.flush(), + ]); + } +}; diff --git a/apps/web/vercel.json b/apps/web/vercel.json index 91c88525c6b..8e7130cda40 100644 --- a/apps/web/vercel.json +++ b/apps/web/vercel.json @@ -24,6 +24,10 @@ "path": "/api/cron/streams/update-workspace-clicks", "schedule": "* * * * *" }, + { + "path": "/api/cron/streams/update-click-stats", + "schedule": "* * * * *" + }, { "path": "/api/cron/streams/update-workspace-links-usage", "schedule": "* * * * *" From 52ef12328dc6bab761cfc719d0c225c96b4c4b60 Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 10:38:52 -0700 Subject: [PATCH 2/8] remove publishWorkspaceClicksUsageEvent --- apps/web/lib/api/conversions/track-sale.ts | 14 ++++++--- .../redis-streams/workspace-clicks-usage.ts | 31 ------------------- 2 files changed, 9 insertions(+), 36 deletions(-) delete mode 100644 apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts diff --git a/apps/web/lib/api/conversions/track-sale.ts b/apps/web/lib/api/conversions/track-sale.ts index 725b6d26854..b22e04de093 100644 --- a/apps/web/lib/api/conversions/track-sale.ts +++ b/apps/web/lib/api/conversions/track-sale.ts @@ -16,7 +16,6 @@ import { } from "@/lib/tinybird"; import { CustomerSource, LeadEventTB, WorkspaceProps } from "@/lib/types"; import { redis } from "@/lib/upstash"; -import { publishWorkspaceClicksUsageEvent } from "@/lib/upstash/redis-streams/workspace-clicks-usage"; import { sendWorkspaceWebhook } from "@/lib/webhook/publish"; import { transformLeadEventData, @@ -658,10 +657,15 @@ const _trackSale = async ({ ] : []), - publishWorkspaceClicksUsageEvent({ - linkId: link.id, - workspaceId: workspace.id, - timestamp: new Date().toISOString(), + prisma.project.update({ + where: { + id: workspace.id, + }, + data: { + usage: { + increment: 1, + }, + }, }), ]); diff --git a/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts b/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts deleted file mode 100644 index 9f5e930fe6c..00000000000 --- a/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts +++ /dev/null @@ -1,31 +0,0 @@ -import { redis } from "../redis"; -import { RedisStream } from "./client"; - -/* Workspace Clicks Usage Stream */ -const WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY = "workspace:usage:updates"; - -export const workspaceClicksUsageStream = new RedisStream( - WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY, -); - -export interface WorkspaceClicksUsageEvent { - linkId: string; - workspaceId: string; - timestamp: string; -} -// Publishes a click event to any relevant streams in a single transaction -export const publishWorkspaceClicksUsageEvent = async ( - event: WorkspaceClicksUsageEvent, -) => { - const { linkId, workspaceId, timestamp } = event; - try { - return await redis.xadd(WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY, "*", { - linkId, - workspaceId, - timestamp, - }); - } catch (error) { - console.error("Failed to publish workspace clicks usage event:", error); - throw error; - } -}; From 9128638bec5477792b86acc45e7e5c5dd1ccdf9b Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 10:43:13 -0700 Subject: [PATCH 3/8] update types --- apps/web/lib/api/partners/sync-partner-links-stats.ts | 2 +- apps/web/lib/upstash/redis-streams/click-stats.ts | 4 ---- apps/web/lib/upstash/redis-streams/partner-activity.ts | 2 +- apps/web/scripts/partners/aggregate-stats-seeding.ts | 2 +- 4 files changed, 3 insertions(+), 7 deletions(-) diff --git a/apps/web/lib/api/partners/sync-partner-links-stats.ts b/apps/web/lib/api/partners/sync-partner-links-stats.ts index 9b8f577d119..bf29f88043f 100644 --- a/apps/web/lib/api/partners/sync-partner-links-stats.ts +++ b/apps/web/lib/api/partners/sync-partner-links-stats.ts @@ -9,7 +9,7 @@ export const syncPartnerLinksStats = async ({ }: { partnerId: string; programId: string; - eventType: "click" | "lead" | "sale"; + eventType: "lead" | "sale"; }) => { try { return await publishPartnerActivityEvent({ diff --git a/apps/web/lib/upstash/redis-streams/click-stats.ts b/apps/web/lib/upstash/redis-streams/click-stats.ts index a1557ea38e7..c0a1b5096ef 100644 --- a/apps/web/lib/upstash/redis-streams/click-stats.ts +++ b/apps/web/lib/upstash/redis-streams/click-stats.ts @@ -3,10 +3,6 @@ import { conn } from "@/lib/planetscale"; import { redis } from "../redis"; import { RedisStream } from "./client"; -// TODO: clean up those legacy streams once remaining publishers -// are migrated (track-sale → workspace-clicks-usage; lead/sale/commission → -// partner-activity) or confirmed still needed. - const CLICK_STATS_STREAM_KEY = "click:stats:updates"; export const clickStatsStream = new RedisStream(CLICK_STATS_STREAM_KEY); diff --git a/apps/web/lib/upstash/redis-streams/partner-activity.ts b/apps/web/lib/upstash/redis-streams/partner-activity.ts index 3fb5761e062..7af5848be7d 100644 --- a/apps/web/lib/upstash/redis-streams/partner-activity.ts +++ b/apps/web/lib/upstash/redis-streams/partner-activity.ts @@ -12,7 +12,7 @@ export interface PartnerActivityEvent { programId: string; partnerId: string; timestamp: string; - eventType: "click" | "lead" | "sale" | "commission"; + eventType: "lead" | "sale" | "commission"; } // Publishes a partner activity event to the stream diff --git a/apps/web/scripts/partners/aggregate-stats-seeding.ts b/apps/web/scripts/partners/aggregate-stats-seeding.ts index 22de25c2da5..70ce5330488 100644 --- a/apps/web/scripts/partners/aggregate-stats-seeding.ts +++ b/apps/web/scripts/partners/aggregate-stats-seeding.ts @@ -54,7 +54,7 @@ async function main() { await publishPartnerActivityEvent({ partnerId: partnerLink.partnerId!, programId: partnerLink.programId!, - eventType: "click", + eventType: "lead", timestamp: new Date().toISOString(), }); }), From da15e9de217316352768b5ac0213ff0b0e4c2369 Mon Sep 17 00:00:00 2001 From: Kiran K Date: Tue, 28 Jul 2026 23:15:10 +0530 Subject: [PATCH 4/8] Qualify Project columns in click stats usage updates. --- apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts | 2 +- apps/web/lib/upstash/redis-streams/click-stats.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts index 32de415dd0a..63512f46d52 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -221,7 +221,7 @@ const processClickStatsStreamBatch = () => async (update) => { try { await conn.execute( - "UPDATE Project SET usage = usage + ?, totalClicks = totalClicks + ? WHERE id = ?", + "UPDATE Project p SET p.usage = p.usage + ?, p.totalClicks = p.totalClicks + ? WHERE id = ?", [update.clicks, update.clicks, update.workspaceId], ); return { success: true }; diff --git a/apps/web/lib/upstash/redis-streams/click-stats.ts b/apps/web/lib/upstash/redis-streams/click-stats.ts index a1557ea38e7..1d3293766db 100644 --- a/apps/web/lib/upstash/redis-streams/click-stats.ts +++ b/apps/web/lib/upstash/redis-streams/click-stats.ts @@ -56,7 +56,7 @@ export const publishClickStatsEvent = async ({ workspaceId ? conn.execute( - "UPDATE Project SET usage = usage + 1, totalClicks = totalClicks + 1 WHERE id = ?", + "UPDATE Project p SET p.usage = p.usage + 1, p.totalClicks = p.totalClicks + 1 WHERE id = ?", [workspaceId], ) : null, From 76ea2347fe74afca9fba0eadefbd0973251e4d2a Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 10:52:07 -0700 Subject: [PATCH 5/8] rename to publishLinkClickEvent / linkClickEventStream --- .../cron/streams/update-click-stats/route.ts | 14 ++++----- .../streams/update-workspace-clicks/route.ts | 1 + apps/web/lib/tinybird/record-click.ts | 4 +-- .../{click-stats.ts => link-click-events.ts} | 18 ++++++----- .../upstash/redis-streams/partner-activity.ts | 2 +- .../redis-streams/workspace-click-events.ts | 1 + .../redis-streams/workspace-clicks-usage.ts | 31 +++++++++++++++++++ .../redis-streams/workspace-links-usage.ts | 1 + 8 files changed, 55 insertions(+), 17 deletions(-) rename apps/web/lib/upstash/redis-streams/{click-stats.ts => link-click-events.ts} (74%) create mode 100644 apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts diff --git a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts index 32de415dd0a..25ddabe3a8f 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -2,11 +2,11 @@ import { qstash } from "@/lib/cron"; import { withCron } from "@/lib/cron/with-cron"; import { conn } from "@/lib/planetscale"; import { redis } from "@/lib/upstash/redis"; -import { - ClickStatsEvent, - clickStatsStream, -} from "@/lib/upstash/redis-streams/click-stats"; import { RedisStreamEntry } from "@/lib/upstash/redis-streams/client"; +import { + LinkClickEvent, + linkClickEventStream, +} from "@/lib/upstash/redis-streams/link-click-events"; import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils"; import { format } from "date-fns"; import { NextResponse } from "next/server"; @@ -40,7 +40,7 @@ type EnrollmentAggregate = { }; const aggregateClickStats = ( - entries: RedisStreamEntry[], + entries: RedisStreamEntry[], ): { linkUpdates: LinkAggregate[]; workspaceUpdates: WorkspaceAggregate[]; @@ -152,7 +152,7 @@ const getStreamEntryAgeMs = (entryId: string | null) => { }; const processClickStatsStreamBatch = () => - clickStatsStream.processBatch( + linkClickEventStream.processBatch( async (entries) => { if (!entries || entries.length === 0) { return { @@ -325,7 +325,7 @@ const executeClickStatsCron = async () => { entriesProcessed = 0, } = await processClickStatsStreamBatch(); - const streamInfo = await clickStatsStream.getStreamInfo(); + const streamInfo = await linkClickEventStream.getStreamInfo(); await maybeAlertOnBacklog(streamInfo); const hasMore = diff --git a/apps/web/app/(ee)/api/cron/streams/update-workspace-clicks/route.ts b/apps/web/app/(ee)/api/cron/streams/update-workspace-clicks/route.ts index 4055d92da7a..e3f676c9d75 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-workspace-clicks/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-workspace-clicks/route.ts @@ -7,6 +7,7 @@ import { } from "@/lib/upstash/redis-streams/workspace-clicks-usage"; import { NextResponse } from "next/server"; +// TODO: Remove once stream is drained export const dynamic = "force-dynamic"; const BATCH_SIZE = 10000; diff --git a/apps/web/lib/tinybird/record-click.ts b/apps/web/lib/tinybird/record-click.ts index 4080a7d3a57..38873c91833 100644 --- a/apps/web/lib/tinybird/record-click.ts +++ b/apps/web/lib/tinybird/record-click.ts @@ -13,7 +13,7 @@ import { detectBot } from "../middleware/utils/detect-bot"; import { detectQr } from "../middleware/utils/detect-qr"; import { getIdentityHash } from "../middleware/utils/get-identity-hash"; import { redis } from "../upstash"; -import { publishClickStatsEvent } from "../upstash/redis-streams/click-stats"; +import { publishLinkClickEvent } from "../upstash/redis-streams/link-click-events"; import { publishWorkspaceClickEvent } from "../upstash/redis-streams/workspace-click-events"; /** @@ -188,7 +188,7 @@ export async function recordClick({ clickId, }), - publishClickStatsEvent({ + publishLinkClickEvent({ linkId, timestamp: clickData.timestamp, ...(workspaceId && url && { workspaceId }), diff --git a/apps/web/lib/upstash/redis-streams/click-stats.ts b/apps/web/lib/upstash/redis-streams/link-click-events.ts similarity index 74% rename from apps/web/lib/upstash/redis-streams/click-stats.ts rename to apps/web/lib/upstash/redis-streams/link-click-events.ts index c0a1b5096ef..48dc9e6d3b8 100644 --- a/apps/web/lib/upstash/redis-streams/click-stats.ts +++ b/apps/web/lib/upstash/redis-streams/link-click-events.ts @@ -3,11 +3,11 @@ import { conn } from "@/lib/planetscale"; import { redis } from "../redis"; import { RedisStream } from "./client"; -const CLICK_STATS_STREAM_KEY = "click:stats:updates"; +const STREAM_KEY = "link:click:events"; -export const clickStatsStream = new RedisStream(CLICK_STATS_STREAM_KEY); +export const linkClickEventStream = new RedisStream(STREAM_KEY); -export interface ClickStatsEvent { +export interface LinkClickEvent { linkId: string; timestamp: string; workspaceId?: string; @@ -15,13 +15,17 @@ export interface ClickStatsEvent { partnerId?: string; } -export const publishClickStatsEvent = async ({ +// Publishes a link click event to the stream to update: +// - link clicks count + lastClicked timestamp +// - workspace usage + totalClicks +// - program enrollment totalClicks +export const publishLinkClickEvent = async ({ linkId, timestamp, workspaceId, programId, partnerId, -}: ClickStatsEvent) => { +}: LinkClickEvent) => { const payload = { linkId, timestamp, @@ -30,11 +34,11 @@ export const publishClickStatsEvent = async ({ }; try { - return await redis.xadd(CLICK_STATS_STREAM_KEY, "*", payload); + return await redis.xadd(STREAM_KEY, "*", payload); } catch (error) { logger.error("stream.publish_failed", { service: "upstash", - streamKey: CLICK_STATS_STREAM_KEY, + streamKey: STREAM_KEY, error: toErrorFields(error), correlation: { linkId, diff --git a/apps/web/lib/upstash/redis-streams/partner-activity.ts b/apps/web/lib/upstash/redis-streams/partner-activity.ts index 7af5848be7d..ee4cf03517a 100644 --- a/apps/web/lib/upstash/redis-streams/partner-activity.ts +++ b/apps/web/lib/upstash/redis-streams/partner-activity.ts @@ -15,7 +15,7 @@ export interface PartnerActivityEvent { eventType: "lead" | "sale" | "commission"; } -// Publishes a partner activity event to the stream +// Publishes a partner activity event to the stream to update programEnrollment stats (e.g. totalClicks, totalLeads, totalSales, totalCommissions) export const publishPartnerActivityEvent = async ( event: PartnerActivityEvent, ) => { diff --git a/apps/web/lib/upstash/redis-streams/workspace-click-events.ts b/apps/web/lib/upstash/redis-streams/workspace-click-events.ts index 2868bd98470..c6659ca6a1d 100644 --- a/apps/web/lib/upstash/redis-streams/workspace-click-events.ts +++ b/apps/web/lib/upstash/redis-streams/workspace-click-events.ts @@ -8,6 +8,7 @@ const STREAM_KEY = "workspace:click:events"; export const workspaceClickEventStream = new RedisStream(STREAM_KEY); +// Publishes a workspace click event to the stream (only for workspaces with link.clicked webhooks) export const publishWorkspaceClickEvent = async (event) => { try { const parsedEvent = clickEventSchemaTB.parse({ diff --git a/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts b/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts new file mode 100644 index 00000000000..005bf418273 --- /dev/null +++ b/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts @@ -0,0 +1,31 @@ +import { redis } from "../redis"; +import { RedisStream } from "./client"; + +// TODO: Remove once stream is drained +const WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY = "workspace:usage:updates"; + +export const workspaceClicksUsageStream = new RedisStream( + WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY, +); + +export interface WorkspaceClicksUsageEvent { + linkId: string; + workspaceId: string; + timestamp: string; +} +// Publishes a click event to any relevant streams in a single transaction +export const publishWorkspaceClicksUsageEvent = async ( + event: WorkspaceClicksUsageEvent, +) => { + const { linkId, workspaceId, timestamp } = event; + try { + return await redis.xadd(WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY, "*", { + linkId, + workspaceId, + timestamp, + }); + } catch (error) { + console.error("Failed to publish workspace clicks usage event:", error); + throw error; + } +}; diff --git a/apps/web/lib/upstash/redis-streams/workspace-links-usage.ts b/apps/web/lib/upstash/redis-streams/workspace-links-usage.ts index d7c029d5f87..a0dfa568378 100644 --- a/apps/web/lib/upstash/redis-streams/workspace-links-usage.ts +++ b/apps/web/lib/upstash/redis-streams/workspace-links-usage.ts @@ -15,6 +15,7 @@ export interface WorkspaceLinksUsageEvent { timestamp: string; } +// Publishes a workspace links usage event to the stream to update workspace linksUsage and totalLinks export const publishWorkspaceLinksUsageEvent = async ( event: WorkspaceLinksUsageEvent, ) => { From 7ac706f207344125073fc723ad8f46b5a84264d2 Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 11:03:41 -0700 Subject: [PATCH 6/8] remove date-fns dep, address CR feedback --- .../api/cron/streams/update-click-stats/route.ts | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts index 545909a8a94..c3413adda4e 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -8,7 +8,6 @@ import { linkClickEventStream, } from "@/lib/upstash/redis-streams/link-click-events"; import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils"; -import { format } from "date-fns"; import { NextResponse } from "next/server"; import { logAndRespond } from "../../utils"; @@ -195,13 +194,14 @@ const processClickStatsStreamBatch = () => linkUpdates, async (update) => { try { + const lastClickedAt = new Date(update.lastClicked) + .toISOString() + .slice(0, 19) + .replace("T", " "); + await conn.execute( - "UPDATE Link SET clicks = clicks + ?, lastClicked = ? WHERE id = ?", - [ - update.clicks, - format(new Date(update.lastClicked), "yyyy-MM-dd HH:mm:ss"), - update.linkId, - ], + "UPDATE Link SET clicks = clicks + ?, lastClicked = GREATEST(COALESCE(lastClicked, ?), ?) WHERE id = ?", + [update.clicks, lastClickedAt, lastClickedAt, update.linkId], ); processedEntryIds.push(...update.entryIds); return { success: true }; From 914237cafff5343967f92e803484154e5dcc56e6 Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 11:06:12 -0700 Subject: [PATCH 7/8] fix ts error --- .../app/(ee)/api/track/application/route.ts | 26 +++++++------------ 1 file changed, 9 insertions(+), 17 deletions(-) diff --git a/apps/web/app/(ee)/api/track/application/route.ts b/apps/web/app/(ee)/api/track/application/route.ts index ac02785cb77..ed5b08b0788 100644 --- a/apps/web/app/(ee)/api/track/application/route.ts +++ b/apps/web/app/(ee)/api/track/application/route.ts @@ -1,7 +1,6 @@ import { COMMON_CORS_HEADERS } from "@/lib/api/cors"; import { createId } from "@/lib/api/create-id"; import { DubApiError, handleAndReturnErrorResponse } from "@/lib/api/errors"; -import { syncPartnerLinksStats } from "@/lib/api/partners/sync-partner-links-stats"; import { parseRequestBody } from "@/lib/api/utils"; import { getIP } from "@/lib/api/utils/get-ip"; import { @@ -19,6 +18,7 @@ import { recordClickZodSchema, } from "@/lib/tinybird/record-click-zod"; import { ratelimit } from "@/lib/upstash"; +import { publishLinkClickEvent } from "@/lib/upstash/redis-streams/link-click-events"; import { MARKETPLACE_RESERVED_SLUGS } from "@/ui/program-marketplace/utils/urls"; import { capitalize, @@ -281,25 +281,17 @@ async function trackVisitEvent({ `Tracked click event for network partner ${referredByPartner.id}`, ); - await prisma.link.update({ - where: { - id: networkReferralLink.id, - }, - data: { - clicks: { increment: 1 }, - lastClicked: new Date(), - }, + await publishLinkClickEvent({ + linkId: networkReferralLink.id, + timestamp: new Date().toISOString(), + workspaceId: NETWORK_WORKSPACE_ID, + programId: NETWORK_PROGRAM_ID, + partnerId: referredByPartner.id, }); + console.log( - `Updated link ${networkReferralLink.id} to ${networkReferralLink.clicks + 1} clicks`, + `Published link click event for network referral link ${networkReferralLink.id}`, ); - - await syncPartnerLinksStats({ - partnerId: referredByPartner.id, - programId: NETWORK_PROGRAM_ID, - eventType: "click", - }); - console.log(`Synced click stats for partner ${referredByPartner.id}`); })(), ); } From 7e1b5803f9f885df780a0535cd47e7497f23b37a Mon Sep 17 00:00:00 2001 From: Steven Tey Date: Tue, 28 Jul 2026 12:28:39 -0700 Subject: [PATCH 8/8] address coderabbit feedback --- .../cron/streams/update-click-stats/route.ts | 78 ++++++------------- .../streams/update-partner-stats/route.ts | 12 +-- .../update-workspace-links-usage/route.ts | 12 +-- 3 files changed, 30 insertions(+), 72 deletions(-) diff --git a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts index c3413adda4e..1beb4e72552 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -1,4 +1,3 @@ -import { qstash } from "@/lib/cron"; import { withCron } from "@/lib/cron/with-cron"; import { conn } from "@/lib/planetscale"; import { redis } from "@/lib/upstash/redis"; @@ -7,8 +6,7 @@ import { LinkClickEvent, linkClickEventStream, } from "@/lib/upstash/redis-streams/link-click-events"; -import { APP_DOMAIN_WITH_NGROK, log } from "@dub/utils"; -import { NextResponse } from "next/server"; +import { log } from "@dub/utils"; import { logAndRespond } from "../../utils"; export const dynamic = "force-dynamic"; @@ -317,73 +315,41 @@ const maybeAlertOnBacklog = async (streamInfo: { }); }; -const executeClickStatsCron = async () => { - const { - linkUpdates, - errors, - totalProcessed, - entriesProcessed = 0, - } = await processClickStatsStreamBatch(); +// GET /api/cron/streams/update-click-stats +export const GET = withCron(async () => { + const acquired = await redis.set(LOCK_KEY, "1", { + nx: true, + ex: LOCK_TTL_SECONDS, + }); - const streamInfo = await linkClickEventStream.getStreamInfo(); - await maybeAlertOnBacklog(streamInfo); + if (!acquired) { + return logAndRespond( + "[update-click-stats] Another run is in progress. Skipping...", + ); + } - const hasMore = - streamInfo.length > 0 || (entriesProcessed ?? 0) >= BATCH_SIZE; + const { linkUpdates, errors, totalProcessed } = + await processClickStatsStreamBatch(); - if (hasMore) { - await qstash.publishJSON({ - url: `${APP_DOMAIN_WITH_NGROK}/api/cron/streams/update-click-stats`, - method: "POST", - body: {}, - }); - } + const streamInfo = await linkClickEventStream.getStreamInfo(); + await maybeAlertOnBacklog(streamInfo); if (!linkUpdates.length) { - return NextResponse.json({ + return logAndRespond({ success: true, - message: "No updates to process", processed: 0, streamInfo, - hasMore, + message: "No updates to process", }); } - const response = { + await redis.del(LOCK_KEY); + + return logAndRespond({ success: true, processed: totalProcessed, errors: errors?.length || 0, streamInfo, - hasMore, message: `Successfully processed ${totalProcessed} click stats updates`, - }; - - console.log(response); - - return NextResponse.json(response); -}; - -const runWithLock = async () => { - const acquired = await redis.set(LOCK_KEY, "1", { - nx: true, - ex: LOCK_TTL_SECONDS, }); - - if (!acquired) { - return logAndRespond( - "[update-click-stats] Another run is in progress. Skipping...", - ); - } - - try { - return await executeClickStatsCron(); - } finally { - await redis.del(LOCK_KEY); - } -}; - -// GET /api/cron/streams/update-click-stats -export const GET = withCron(async () => runWithLock()); - -// POST /api/cron/streams/update-click-stats (recursively called by QStash) -export const POST = withCron(async () => runWithLock()); +}); diff --git a/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts b/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts index 9205796015c..7e083361fd0 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-partner-stats/route.ts @@ -8,7 +8,7 @@ import { import { toCentsNumber } from "@dub/utils"; import { ProgramEnrollment } from "@prisma/client"; import { differenceInDays, format } from "date-fns"; -import { NextResponse } from "next/server"; +import { logAndRespond } from "../../utils"; export const dynamic = "force-dynamic"; @@ -351,7 +351,7 @@ export const GET = withCron(async () => { await processPartnerActivityStreamBatch(); if (!updates.length) { - return NextResponse.json({ + return logAndRespond({ success: true, message: "No updates to process", processed: 0, @@ -360,15 +360,11 @@ export const GET = withCron(async () => { const streamInfo = await partnerActivityStream.getStreamInfo(); - const response = { + return logAndRespond({ success: true, processed: totalProcessed, errors: errors?.length || 0, streamInfo, message: `Successfully processed ${totalProcessed} partner activity updates`, - }; - - console.log(response); - - return NextResponse.json(response); + }); }); diff --git a/apps/web/app/(ee)/api/cron/streams/update-workspace-links-usage/route.ts b/apps/web/app/(ee)/api/cron/streams/update-workspace-links-usage/route.ts index ac77b6019bb..c08e7a2e75a 100644 --- a/apps/web/app/(ee)/api/cron/streams/update-workspace-links-usage/route.ts +++ b/apps/web/app/(ee)/api/cron/streams/update-workspace-links-usage/route.ts @@ -11,7 +11,7 @@ import { workspaceLinksUsageStream, } from "@/lib/upstash/redis-streams/workspace-links-usage"; import { log } from "@dub/utils"; -import { NextResponse } from "next/server"; +import { logAndRespond } from "../../utils"; export const dynamic = "force-dynamic"; @@ -317,7 +317,7 @@ export const GET = withCron(async () => { } = await processWorkspaceLinksUsageBatch(); if (!updates.length) { - return NextResponse.json({ + return logAndRespond({ success: true, message: "No updates to process", processed: 0, @@ -326,7 +326,7 @@ export const GET = withCron(async () => { const streamInfo = await workspaceLinksUsageStream.getStreamInfo(); - const response = { + return logAndRespond({ success: true, processed: totalProcessed, notificationsSent, @@ -334,9 +334,5 @@ export const GET = withCron(async () => { lastProcessedId, streamInfo, message: `Successfully processed ${totalProcessed} workspace links usage updates`, - }; - - console.log(response); - - return NextResponse.json(response); + }); });