Skip to content

Commit 12a152a

Browse files
committed
feat(followers): implement CRDT-based follower count to eliminate lost updates (#757)
Replace monolithic follower count with a per-node G-Counter CRDT structure to prevent lost updates under high concurrency. - FollowerCounterShard and FollowEvent Prisma schema definitions - NODE_ID configuration option in config schema - Follow and unfollow operations using atomic shard increments and idempotent events - getFollowerCount summing across all node shards with 10s Redis caching - Nightly shard compaction merging per-node shards atomically - Unit tests verifying 100 concurrent follows, 50 follows + 30 unfollows, double-follow idempotency, multi-node resolution, and compaction
1 parent dea4063 commit 12a152a

8 files changed

Lines changed: 605 additions & 0 deletions

File tree

prisma/schema/follower.prisma

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
// prisma/schema/follower.prisma
2+
3+
enum FollowDirection {
4+
FOLLOW
5+
UNFOLLOW
6+
}
7+
8+
model FollowerCounterShard {
9+
id String @id @default(cuid())
10+
creatorWallet String
11+
nodeId String
12+
increments BigInt @default(0)
13+
decrements BigInt @default(0)
14+
updatedAt DateTime @updatedAt
15+
16+
@@unique([creatorWallet, nodeId])
17+
@@index([creatorWallet])
18+
@@map("follower_counter_shards")
19+
}
20+
21+
model FollowEvent {
22+
id String @id @default(cuid())
23+
followerWallet String
24+
creatorWallet String
25+
direction FollowDirection
26+
createdAt DateTime @default(now())
27+
28+
@@unique([followerWallet, creatorWallet])
29+
@@index([creatorWallet])
30+
@@map("follow_events")
31+
}

src/config.schema.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ export const envSchema = z
4646
DATABASE_URL: z
4747
.string()
4848
.min(1, 'DATABASE_URL is required in the environment variables'),
49+
NODE_ID: z.string().default('node-local'),
4950

5051
GMAIL_USER: z.string(),
5152
GMAIL_APP_PASSWORD: z.string(),
Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
import { prisma } from '../utils/prisma.utils';
2+
import { compactShardsForCreator } from '../modules/followers/follower.service';
3+
import { logger } from '../utils/logger.utils';
4+
5+
export async function runNightlyFollowerShardCompaction(): Promise<{
6+
creatorsCompacted: number;
7+
}> {
8+
logger.info('Starting nightly follower shard compaction job');
9+
10+
const creators = await prisma.followerCounterShard.findMany({
11+
distinct: ['creatorWallet'],
12+
select: { creatorWallet: true },
13+
});
14+
15+
let count = 0;
16+
for (const { creatorWallet } of creators) {
17+
const compacted = await compactShardsForCreator(creatorWallet);
18+
if (compacted) {
19+
count++;
20+
}
21+
}
22+
23+
logger.info(
24+
{ creators_compacted: count, total_creators_evaluated: creators.length },
25+
'Nightly follower shard compaction completed'
26+
);
27+
28+
return { creatorsCompacted: count };
29+
}
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
import { AsyncController } from '../../types/auth.types';
2+
import {
3+
follow,
4+
unfollow,
5+
getFollowerCount,
6+
} from './follower.service';
7+
import { sendSuccess, sendValidationError } from '../../utils/api-response.utils';
8+
9+
export const httpFollow: AsyncController = async (req, res, next) => {
10+
try {
11+
const { creatorWallet } = req.params;
12+
const followerWallet = req.body?.followerWallet || req.jwtPayload?.walletAddress;
13+
14+
if (!creatorWallet || !followerWallet) {
15+
sendValidationError(res, 'creatorWallet param and followerWallet body/token are required');
16+
return;
17+
}
18+
19+
const result = await follow(followerWallet, creatorWallet);
20+
sendSuccess(res, result);
21+
} catch (err) {
22+
next(err);
23+
}
24+
};
25+
26+
export const httpUnfollow: AsyncController = async (req, res, next) => {
27+
try {
28+
const { creatorWallet } = req.params;
29+
const followerWallet = req.body?.followerWallet || req.jwtPayload?.walletAddress;
30+
31+
if (!creatorWallet || !followerWallet) {
32+
sendValidationError(res, 'creatorWallet param and followerWallet body/token are required');
33+
return;
34+
}
35+
36+
const result = await unfollow(followerWallet, creatorWallet);
37+
sendSuccess(res, result);
38+
} catch (err) {
39+
next(err);
40+
}
41+
};
42+
43+
export const httpGetFollowerCount: AsyncController = async (req, res, next) => {
44+
try {
45+
const { creatorWallet } = req.params;
46+
if (!creatorWallet) {
47+
sendValidationError(res, 'creatorWallet parameter is required');
48+
return;
49+
}
50+
51+
const count = await getFollowerCount(creatorWallet);
52+
sendSuccess(res, { creatorWallet, count });
53+
} catch (err) {
54+
next(err);
55+
}
56+
};
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import { Router } from 'express';
2+
import {
3+
httpFollow,
4+
httpUnfollow,
5+
httpGetFollowerCount,
6+
} from './follower.controllers';
7+
8+
const followerRouter = Router();
9+
10+
followerRouter.post('/:creatorWallet/follow', httpFollow);
11+
followerRouter.post('/:creatorWallet/unfollow', httpUnfollow);
12+
followerRouter.get('/:creatorWallet/count', httpGetFollowerCount);
13+
14+
export default followerRouter;
Lines changed: 228 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,228 @@
1+
import {
2+
follow,
3+
unfollow,
4+
getFollowerCount,
5+
compactShardsForCreator,
6+
} from './follower.service';
7+
import { prisma } from '../../utils/prisma.utils';
8+
9+
jest.mock('../../utils/logger.utils', () => ({
10+
logger: {
11+
debug: jest.fn(),
12+
info: jest.fn(),
13+
warn: jest.fn(),
14+
error: jest.fn(),
15+
},
16+
}));
17+
18+
// In-memory mock for Prisma to test CRDT logic isolated from real DB
19+
jest.mock('../../utils/prisma.utils', () => {
20+
const followEvents = new Map<string, { followerWallet: string; creatorWallet: string; direction: string }>();
21+
const shards = new Map<string, { id: string; creatorWallet: string; nodeId: string; increments: bigint; decrements: bigint; updatedAt: Date }>();
22+
23+
return {
24+
prisma: {
25+
followEvent: {
26+
findUnique: jest.fn().mockImplementation(async ({ where }) => {
27+
const key = `${where.followerWallet_creatorWallet.followerWallet}:${where.followerWallet_creatorWallet.creatorWallet}`;
28+
return followEvents.get(key) || null;
29+
}),
30+
upsert: jest.fn().mockImplementation(async ({ where, update, create }) => {
31+
const key = `${where.followerWallet_creatorWallet.followerWallet}:${where.followerWallet_creatorWallet.creatorWallet}`;
32+
const existing = followEvents.get(key);
33+
const direction = existing ? update.direction : create.direction;
34+
const record = {
35+
followerWallet: where.followerWallet_creatorWallet.followerWallet,
36+
creatorWallet: where.followerWallet_creatorWallet.creatorWallet,
37+
direction,
38+
};
39+
followEvents.set(key, record);
40+
return record;
41+
}),
42+
update: jest.fn().mockImplementation(async ({ where, data }) => {
43+
const key = `${where.followerWallet_creatorWallet.followerWallet}:${where.followerWallet_creatorWallet.creatorWallet}`;
44+
const existing = followEvents.get(key);
45+
if (existing) {
46+
existing.direction = data.direction;
47+
}
48+
return existing;
49+
}),
50+
},
51+
followerCounterShard: {
52+
findUnique: jest.fn().mockImplementation(async ({ where }) => {
53+
const key = `${where.creatorWallet_nodeId.creatorWallet}:${where.creatorWallet_nodeId.nodeId}`;
54+
return shards.get(key) || null;
55+
}),
56+
findMany: jest.fn().mockImplementation(async ({ where }) => {
57+
const list = Array.from(shards.values()).filter(
58+
s => s.creatorWallet === where.creatorWallet
59+
);
60+
if (where.updatedAt?.gte) {
61+
return list.filter(s => s.updatedAt >= where.updatedAt.gte);
62+
}
63+
return list;
64+
}),
65+
findFirst: jest.fn().mockImplementation(async ({ where }) => {
66+
const list = Array.from(shards.values()).filter(
67+
s => s.creatorWallet === where.creatorWallet
68+
);
69+
if (where.updatedAt?.gte) {
70+
return list.find(s => s.updatedAt >= where.updatedAt.gte) || null;
71+
}
72+
return list[0] || null;
73+
}),
74+
create: jest.fn().mockImplementation(async ({ data }) => {
75+
const key = `${data.creatorWallet}:${data.nodeId}`;
76+
let record = shards.get(key);
77+
if (record) {
78+
record.increments += BigInt(data.increments ?? 0);
79+
record.decrements += BigInt(data.decrements ?? 0);
80+
} else {
81+
record = {
82+
id: key,
83+
creatorWallet: data.creatorWallet,
84+
nodeId: data.nodeId,
85+
increments: BigInt(data.increments ?? 0),
86+
decrements: BigInt(data.decrements ?? 0),
87+
updatedAt: new Date(),
88+
};
89+
shards.set(key, record);
90+
}
91+
return record;
92+
}),
93+
update: jest.fn().mockImplementation(async ({ where, data }) => {
94+
const key = `${where.creatorWallet_nodeId.creatorWallet}:${where.creatorWallet_nodeId.nodeId}`;
95+
let record = shards.get(key);
96+
if (!record) {
97+
record = {
98+
id: key,
99+
creatorWallet: where.creatorWallet_nodeId.creatorWallet,
100+
nodeId: where.creatorWallet_nodeId.nodeId,
101+
increments: 0n,
102+
decrements: 0n,
103+
updatedAt: new Date(),
104+
};
105+
shards.set(key, record);
106+
}
107+
if (data.increments?.increment) {
108+
record.increments += BigInt(data.increments.increment);
109+
}
110+
if (data.decrements?.increment) {
111+
record.decrements += BigInt(data.decrements.increment);
112+
}
113+
record.updatedAt = new Date();
114+
return record;
115+
}),
116+
deleteMany: jest.fn().mockImplementation(async ({ where }) => {
117+
for (const [key, shard] of shards.entries()) {
118+
if (shard.creatorWallet === where.creatorWallet) {
119+
shards.delete(key);
120+
}
121+
}
122+
return { count: 1 };
123+
}),
124+
},
125+
$executeRaw: jest.fn().mockRejectedValue(new Error('raw query fallback')),
126+
$transaction: jest.fn().mockImplementation(async (actions) => Promise.all(actions)),
127+
_reset: () => {
128+
followEvents.clear();
129+
shards.clear();
130+
},
131+
},
132+
};
133+
});
134+
135+
jest.mock('../../utils/redis.utils', () => {
136+
const redisStore = new Map<string, string>();
137+
return {
138+
getRedis: () => ({
139+
get: jest.fn().mockImplementation(async (key: string) => redisStore.get(key) ?? null),
140+
set: jest.fn().mockImplementation(async (key: string, val: string) => {
141+
redisStore.set(key, val);
142+
return 'OK';
143+
}),
144+
del: jest.fn().mockImplementation(async (key: string) => {
145+
redisStore.delete(key);
146+
return 1;
147+
}),
148+
}),
149+
};
150+
});
151+
152+
describe('CRDT Follower Counter (#757)', () => {
153+
const creatorWallet = 'GCREATOR_CRDT_TEST';
154+
155+
beforeEach(() => {
156+
(prisma as any)._reset();
157+
jest.clearAllMocks();
158+
});
159+
160+
it('100 concurrent follows from different wallets produces count of exactly 100', async () => {
161+
const promises = Array.from({ length: 100 }, (_, i) =>
162+
follow(`follower-wallet-${i}`, creatorWallet, 'node-1')
163+
);
164+
165+
await Promise.all(promises);
166+
167+
const count = await getFollowerCount(creatorWallet);
168+
expect(count).toBe(100);
169+
});
170+
171+
it('50 concurrent follows and 30 concurrent unfollows produces net count of 20', async () => {
172+
// First 50 follow
173+
const followPromises = Array.from({ length: 50 }, (_, i) =>
174+
follow(`follower-${i}`, creatorWallet, 'node-1')
175+
);
176+
await Promise.all(followPromises);
177+
178+
// 30 of them unfollow
179+
const unfollowPromises = Array.from({ length: 30 }, (_, i) =>
180+
unfollow(`follower-${i}`, creatorWallet, 'node-1')
181+
);
182+
await Promise.all(unfollowPromises);
183+
184+
const count = await getFollowerCount(creatorWallet);
185+
expect(count).toBe(20);
186+
});
187+
188+
it('double-follow from the same wallet is idempotent and increments count only once', async () => {
189+
const first = await follow('wallet-alice', creatorWallet, 'node-1');
190+
const second = await follow('wallet-alice', creatorWallet, 'node-1');
191+
192+
expect(first.followed).toBe(true);
193+
expect(second.followed).toBe(false);
194+
195+
const count = await getFollowerCount(creatorWallet);
196+
expect(count).toBe(1);
197+
});
198+
199+
it('sums counts across multiple node shards correctly', async () => {
200+
await follow('wallet-1', creatorWallet, 'node-alpha');
201+
await follow('wallet-2', creatorWallet, 'node-alpha');
202+
await follow('wallet-3', creatorWallet, 'node-beta');
203+
204+
const count = await getFollowerCount(creatorWallet);
205+
expect(count).toBe(3);
206+
});
207+
208+
it('nightly compaction merges shards atomically while maintaining identical resolved count', async () => {
209+
await follow('wallet-1', creatorWallet, 'node-1');
210+
await follow('wallet-2', creatorWallet, 'node-2');
211+
await follow('wallet-3', creatorWallet, 'node-3');
212+
213+
const countBefore = await getFollowerCount(creatorWallet);
214+
expect(countBefore).toBe(3);
215+
216+
// Mock Date.now to simulate 6 minutes passing so compaction is allowed
217+
const realNow = Date.now;
218+
jest.spyOn(Date, 'now').mockReturnValue(realNow() + 6 * 60 * 1000);
219+
220+
const compacted = await compactShardsForCreator(creatorWallet);
221+
expect(compacted).toBe(true);
222+
223+
const countAfter = await getFollowerCount(creatorWallet);
224+
expect(countAfter).toBe(3);
225+
226+
jest.restoreAllMocks();
227+
});
228+
});

0 commit comments

Comments
 (0)