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..1beb4e72552 --- /dev/null +++ b/apps/web/app/(ee)/api/cron/streams/update-click-stats/route.ts @@ -0,0 +1,355 @@ +import { withCron } from "@/lib/cron/with-cron"; +import { conn } from "@/lib/planetscale"; +import { redis } from "@/lib/upstash/redis"; +import { RedisStreamEntry } from "@/lib/upstash/redis-streams/client"; +import { + LinkClickEvent, + linkClickEventStream, +} from "@/lib/upstash/redis-streams/link-click-events"; +import { log } from "@dub/utils"; +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 = () => + linkClickEventStream.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 { + const lastClickedAt = new Date(update.lastClicked) + .toISOString() + .slice(0, 19) + .replace("T", " "); + + await conn.execute( + "UPDATE Link SET clicks = clicks + ?, lastClicked = GREATEST(COALESCE(lastClicked, ?), ?) WHERE id = ?", + [update.clicks, lastClickedAt, lastClickedAt, 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 p SET p.usage = p.usage + ?, p.totalClicks = p.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", + }); +}; + +// 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, + }); + + if (!acquired) { + return logAndRespond( + "[update-click-stats] Another run is in progress. Skipping...", + ); + } + + const { linkUpdates, errors, totalProcessed } = + await processClickStatsStreamBatch(); + + const streamInfo = await linkClickEventStream.getStreamInfo(); + await maybeAlertOnBacklog(streamInfo); + + if (!linkUpdates.length) { + return logAndRespond({ + success: true, + processed: 0, + streamInfo, + message: "No updates to process", + }); + } + + await redis.del(LOCK_KEY); + + return logAndRespond({ + success: true, + processed: totalProcessed, + errors: errors?.length || 0, + streamInfo, + message: `Successfully processed ${totalProcessed} click stats updates`, + }); +}); 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-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/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); + }); }); 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}`); })(), ); } 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/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/tinybird/record-click.ts b/apps/web/lib/tinybird/record-click.ts index a5deaca0129..38873c91833 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 { publishLinkClickEvent } from "../upstash/redis-streams/link-click-events"; 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, + }), + + publishLinkClickEvent({ + 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/link-click-events.ts b/apps/web/lib/upstash/redis-streams/link-click-events.ts new file mode 100644 index 00000000000..065d1902b6d --- /dev/null +++ b/apps/web/lib/upstash/redis-streams/link-click-events.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"; + +const STREAM_KEY = "link:click:events"; + +export const linkClickEventStream = new RedisStream(STREAM_KEY); + +export interface LinkClickEvent { + linkId: string; + timestamp: string; + workspaceId?: string; + programId?: string; + partnerId?: string; +} + +// 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, +}: LinkClickEvent) => { + const payload = { + linkId, + timestamp, + ...(workspaceId && { workspaceId }), + ...(programId && partnerId && { programId, partnerId }), + }; + + try { + return await redis.xadd(STREAM_KEY, "*", payload); + } catch (error) { + logger.error("stream.publish_failed", { + service: "upstash", + streamKey: 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 p SET p.usage = p.usage + 1, p.totalClicks = p.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/lib/upstash/redis-streams/partner-activity.ts b/apps/web/lib/upstash/redis-streams/partner-activity.ts index 3fb5761e062..ee4cf03517a 100644 --- a/apps/web/lib/upstash/redis-streams/partner-activity.ts +++ b/apps/web/lib/upstash/redis-streams/partner-activity.ts @@ -12,10 +12,10 @@ 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 +// 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 index 9f5e930fe6c..005bf418273 100644 --- a/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts +++ b/apps/web/lib/upstash/redis-streams/workspace-clicks-usage.ts @@ -1,7 +1,7 @@ import { redis } from "../redis"; import { RedisStream } from "./client"; -/* Workspace Clicks Usage Stream */ +// TODO: Remove once stream is drained const WORKSPACE_CLICKS_USAGE_UPDATES_STREAM_KEY = "workspace:usage:updates"; export const workspaceClicksUsageStream = new RedisStream( 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, ) => { 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(), }); }), 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": "* * * * *"