Skip to content

Commit 512124e

Browse files
authored
Merge pull request #596 from bright5455/feature/dau-rollup-jobs
implemented nightly DAU rollup job
2 parents 336ce5a + 501ad52 commit 512124e

6 files changed

Lines changed: 239 additions & 2 deletions

File tree

backend/src/analytics/analytics.module.ts

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { Module } from '@nestjs/common';
22
import { TypeOrmModule } from '@nestjs/typeorm';
33
import { AnalyticsEvent } from './entities/analytics-event.entity';
44
import { RetentionCohort } from './entities/retention-cohort.entity';
5+
import { DailyActiveUser } from './entities/daily-active-user.entity';
56
import { UsersAnalyticsListener } from './listeners/users-analytics.listener';
67
import { AnalyticsController } from './controllers/analytics.controller';
78
import { AnalyticsService } from './analytics.service';
@@ -12,7 +13,13 @@ import { GetChurnRiskProvider } from './providers/get-churn-risk.provider';
1213
import { ExportCsvProvider } from './providers/export-csv.provider';
1314

1415
@Module({
15-
imports: [TypeOrmModule.forFeature([AnalyticsEvent, RetentionCohort])],
16+
imports: [
17+
TypeOrmModule.forFeature([
18+
AnalyticsEvent,
19+
RetentionCohort,
20+
DailyActiveUser,
21+
]),
22+
],
1623
controllers: [AnalyticsController],
1724
providers: [
1825
AnalyticsService,
Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,32 @@
1+
import {
2+
Column,
3+
CreateDateColumn,
4+
Entity,
5+
Index,
6+
PrimaryGeneratedColumn,
7+
Unique,
8+
} from 'typeorm';
9+
10+
/**
11+
* One row per user who was active on a given day.
12+
*
13+
* Populated by the nightly `DailyActiveUsersRollupJob` from raw
14+
* `AnalyticsEvent` rows, so `GetDauProvider` can compute DAU with a cheap
15+
* `COUNT` instead of a `SELECT DISTINCT` scan over raw events.
16+
*/
17+
@Entity('daily_active_users')
18+
@Unique(['date', 'userId'])
19+
export class DailyActiveUser {
20+
@PrimaryGeneratedColumn('uuid')
21+
id: string;
22+
23+
@Index()
24+
@Column({ type: 'date' })
25+
date: string;
26+
27+
@Column({ type: 'varchar', length: 100 })
28+
userId: string;
29+
30+
@CreateDateColumn({ type: 'timestamptz' })
31+
createdAt: Date;
32+
}

backend/src/analytics/entities/retention-cohort.entity.ts

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@ import {
1616
* on-the-fly from raw `AnalyticsEvent` rows.
1717
*/
1818
@Entity('retention_cohorts')
19-
@Index(['cohortDate'])
2019
export class RetentionCohort {
2120
@PrimaryGeneratedColumn('uuid')
2221
id: string;
Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,131 @@
1+
import { Logger } from '@nestjs/common';
2+
import { Test, TestingModule } from '@nestjs/testing';
3+
import { getRepositoryToken } from '@nestjs/typeorm';
4+
import { AnalyticsEvent } from '../entities/analytics-event.entity';
5+
import { DailyActiveUser } from '../entities/daily-active-user.entity';
6+
import { DailyActiveUsersRollupJob } from './daily-active-users-rollup.job';
7+
8+
describe('DailyActiveUsersRollupJob', () => {
9+
let job: DailyActiveUsersRollupJob;
10+
const mockAnalyticsEventRepo = {
11+
find: jest.fn<Promise<Pick<AnalyticsEvent, 'userId'>[]>, [unknown]>(),
12+
};
13+
const mockDailyActiveUserRepo = {
14+
delete: jest.fn<Promise<unknown>, [unknown]>(),
15+
create: jest.fn(
16+
(data: Partial<DailyActiveUser>) => data as DailyActiveUser,
17+
),
18+
save: jest.fn<Promise<DailyActiveUser[]>, [DailyActiveUser[]]>(),
19+
};
20+
21+
beforeEach(async () => {
22+
jest.clearAllMocks();
23+
24+
const module: TestingModule = await Test.createTestingModule({
25+
providers: [
26+
DailyActiveUsersRollupJob,
27+
{
28+
provide: getRepositoryToken(AnalyticsEvent),
29+
useValue: mockAnalyticsEventRepo,
30+
},
31+
{
32+
provide: getRepositoryToken(DailyActiveUser),
33+
useValue: mockDailyActiveUserRepo,
34+
},
35+
],
36+
}).compile();
37+
38+
job = module.get<DailyActiveUsersRollupJob>(DailyActiveUsersRollupJob);
39+
});
40+
41+
it('should be defined', () => {
42+
expect(job).toBeDefined();
43+
});
44+
45+
describe('rollupForDate', () => {
46+
it('dedupes users who fired multiple events on the same day', async () => {
47+
mockAnalyticsEventRepo.find.mockResolvedValue([
48+
{ userId: 'user-1' },
49+
{ userId: 'user-1' },
50+
{ userId: 'user-2' },
51+
{ userId: 'user-1' },
52+
{ userId: 'user-3' },
53+
{ userId: 'user-2' },
54+
] as AnalyticsEvent[]);
55+
56+
const result = await job.rollupForDate(
57+
new Date('2026-07-22T12:00:00.000Z'),
58+
);
59+
60+
expect(result).toEqual({ date: '2026-07-22', rowsWritten: 3 });
61+
expect(mockDailyActiveUserRepo.save).toHaveBeenCalledTimes(1);
62+
63+
const savedRows = mockDailyActiveUserRepo.save.mock.calls[0][0];
64+
const savedUserIds = savedRows.map((row) => row.userId).sort();
65+
expect(savedUserIds).toEqual(['user-1', 'user-2', 'user-3']);
66+
savedRows.forEach((row) => expect(row.date).toBe('2026-07-22'));
67+
});
68+
69+
it('queries raw events for the correct UTC day bounds', async () => {
70+
mockAnalyticsEventRepo.find.mockResolvedValue([]);
71+
72+
await job.rollupForDate(new Date('2026-07-22T12:00:00.000Z'));
73+
74+
expect(mockAnalyticsEventRepo.find).toHaveBeenCalledWith(
75+
expect.objectContaining({
76+
select: ['userId'],
77+
where: expect.objectContaining({ timestamp: expect.anything() }),
78+
}),
79+
);
80+
});
81+
82+
it('is idempotent: re-running for the same day replaces rather than duplicates rows', async () => {
83+
mockAnalyticsEventRepo.find.mockResolvedValue([
84+
{ userId: 'user-1' },
85+
{ userId: 'user-2' },
86+
] as AnalyticsEvent[]);
87+
88+
await job.rollupForDate(new Date('2026-07-22T00:00:00.000Z'));
89+
await job.rollupForDate(new Date('2026-07-22T00:00:00.000Z'));
90+
91+
expect(mockDailyActiveUserRepo.delete).toHaveBeenCalledTimes(2);
92+
expect(mockDailyActiveUserRepo.delete).toHaveBeenNthCalledWith(1, {
93+
date: '2026-07-22',
94+
});
95+
expect(mockDailyActiveUserRepo.delete).toHaveBeenNthCalledWith(2, {
96+
date: '2026-07-22',
97+
});
98+
expect(mockDailyActiveUserRepo.save).toHaveBeenCalledTimes(2);
99+
});
100+
101+
it('writes no rows and skips save when no users were active', async () => {
102+
mockAnalyticsEventRepo.find.mockResolvedValue([]);
103+
104+
const result = await job.rollupForDate(
105+
new Date('2026-07-22T00:00:00.000Z'),
106+
);
107+
108+
expect(result.rowsWritten).toBe(0);
109+
expect(mockDailyActiveUserRepo.delete).toHaveBeenCalledWith({
110+
date: '2026-07-22',
111+
});
112+
expect(mockDailyActiveUserRepo.save).not.toHaveBeenCalled();
113+
});
114+
115+
it('logs a summary of rows written on completion', async () => {
116+
mockAnalyticsEventRepo.find.mockResolvedValue([
117+
{ userId: 'user-1' },
118+
] as AnalyticsEvent[]);
119+
const loggerSpy = jest.spyOn(Logger.prototype, 'log');
120+
121+
await job.rollupForDate(new Date('2026-07-22T00:00:00.000Z'));
122+
123+
expect(loggerSpy).toHaveBeenCalledWith(
124+
expect.stringContaining('2026-07-22'),
125+
);
126+
expect(loggerSpy).toHaveBeenCalledWith(
127+
expect.stringContaining('1 rows written'),
128+
);
129+
});
130+
});
131+
});
Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,66 @@
1+
import { Injectable, Logger } from '@nestjs/common';
2+
import { Cron, CronExpression } from '@nestjs/schedule';
3+
import { InjectRepository } from '@nestjs/typeorm';
4+
import { Between, Repository } from 'typeorm';
5+
import { AnalyticsEvent } from '../entities/analytics-event.entity';
6+
import { DailyActiveUser } from '../entities/daily-active-user.entity';
7+
8+
/**
9+
* Nightly rollup that materializes `DailyActiveUser` rows from raw
10+
* `AnalyticsEvent` data for the previous day, so DAU can be read cheaply
11+
* without scanning raw events on every dashboard load.
12+
*/
13+
@Injectable()
14+
export class DailyActiveUsersRollupJob {
15+
private readonly logger = new Logger(DailyActiveUsersRollupJob.name);
16+
17+
constructor(
18+
@InjectRepository(AnalyticsEvent)
19+
private readonly analyticsEventRepository: Repository<AnalyticsEvent>,
20+
@InjectRepository(DailyActiveUser)
21+
private readonly dailyActiveUserRepository: Repository<DailyActiveUser>,
22+
) {}
23+
24+
@Cron(CronExpression.EVERY_DAY_AT_MIDNIGHT)
25+
async handleCron(): Promise<void> {
26+
const yesterday = new Date();
27+
yesterday.setDate(yesterday.getDate() - 1);
28+
await this.rollupForDate(yesterday);
29+
}
30+
31+
/**
32+
* Recomputes the `DailyActiveUser` rows for `targetDate`. Safe to re-run
33+
* for the same day: existing rows for that date are replaced rather than
34+
* appended to.
35+
*/
36+
async rollupForDate(
37+
targetDate: Date,
38+
): Promise<{ date: string; rowsWritten: number }> {
39+
const dateStr = targetDate.toISOString().split('T')[0];
40+
const startOfDay = new Date(`${dateStr}T00:00:00.000Z`);
41+
const endOfDay = new Date(`${dateStr}T23:59:59.999Z`);
42+
43+
const events = await this.analyticsEventRepository.find({
44+
select: ['userId'],
45+
where: { timestamp: Between(startOfDay, endOfDay) },
46+
});
47+
48+
const distinctUserIds = Array.from(
49+
new Set(events.map((event) => event.userId)),
50+
);
51+
52+
await this.dailyActiveUserRepository.delete({ date: dateStr });
53+
if (distinctUserIds.length > 0) {
54+
const rows = distinctUserIds.map((userId) =>
55+
this.dailyActiveUserRepository.create({ date: dateStr, userId }),
56+
);
57+
await this.dailyActiveUserRepository.save(rows);
58+
}
59+
60+
this.logger.log(
61+
`DAU rollup for ${dateStr}: ${distinctUserIds.length} rows written`,
62+
);
63+
64+
return { date: dateStr, rowsWritten: distinctUserIds.length };
65+
}
66+
}

backend/src/app.module.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { Module, NestModule, MiddlewareConsumer, RequestMethod } from '@nestjs/c
22
import { TypeOrmModule } from '@nestjs/typeorm';
33
import { ConfigModule, ConfigService } from '@nestjs/config';
44
import { EventEmitterModule } from '@nestjs/event-emitter';
5+
import { ScheduleModule } from '@nestjs/schedule';
56
import { RedisModule } from './redis/redis.module';
67
import { AuthModule } from './auth/auth.module';
78
import appConfig from './config/app.config';
@@ -36,6 +37,7 @@ import { AnalyticsModule } from './analytics/analytics.module';
3637
load: [appConfig, databaseConfig, jwtConfig],
3738
}),
3839
EventEmitterModule.forRoot(),
40+
ScheduleModule.forRoot(),
3941
TypeOrmModule.forRootAsync({
4042
imports: [ConfigModule],
4143
inject: [ConfigService],

0 commit comments

Comments
 (0)