diff --git a/fluxer_api/src/api/admin/tests/AdminLastActiveIpSearch.test.ts b/fluxer_api/src/api/admin/tests/AdminLastActiveIpSearch.test.ts index ce59f33a3..ce30f0dd3 100644 --- a/fluxer_api/src/api/admin/tests/AdminLastActiveIpSearch.test.ts +++ b/fluxer_api/src/api/admin/tests/AdminLastActiveIpSearch.test.ts @@ -2,6 +2,7 @@ import {afterAll, beforeAll, beforeEach, describe, expect, test} from 'vitest'; import {createTestAccount, setUserACLs} from '../../auth/tests/AuthTestUtils'; +import {getUserActivityBuffer} from '../../middleware/ServiceSingletons'; import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness'; import {HTTP_STATUS} from '../../test/TestConstants'; import {createBuilder} from '../../test/TestRequestBuilder'; @@ -19,6 +20,7 @@ async function setLastActiveIp(harness: ApiTestHarness, token: string, ip: strin .header('x-forwarded-for', ip) .expect(HTTP_STATUS.OK) .execute(); + await getUserActivityBuffer().drainAndFlush(); } describe('Admin last active IP search', () => { diff --git a/fluxer_api/src/api/admin/tests/AdminSearchEndpoints.test.ts b/fluxer_api/src/api/admin/tests/AdminSearchEndpoints.test.ts index 5b8ff92d6..9f3be9ae4 100644 --- a/fluxer_api/src/api/admin/tests/AdminSearchEndpoints.test.ts +++ b/fluxer_api/src/api/admin/tests/AdminSearchEndpoints.test.ts @@ -3,6 +3,7 @@ import {afterEach, beforeEach, describe, expect, test} from 'vitest'; import {createTestAccount, setUserACLs} from '../../auth/tests/AuthTestUtils'; import {createDmChannel, createFriendship, createGuild} from '../../channel/tests/ChannelTestUtils'; +import {getUserActivityBuffer} from '../../middleware/ServiceSingletons'; import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness'; import {HTTP_STATUS} from '../../test/TestConstants'; import {createBuilder} from '../../test/TestRequestBuilder'; @@ -13,6 +14,7 @@ async function setLastActiveIp(harness: ApiTestHarness, token: string, ip: strin .header('x-forwarded-for', ip) .expect(HTTP_STATUS.OK) .execute(); + await getUserActivityBuffer().drainAndFlush(); } describe('Admin Search Endpoints', () => { diff --git a/fluxer_api/src/api/admin/tests/AdminSearchFieldCoverage.test.ts b/fluxer_api/src/api/admin/tests/AdminSearchFieldCoverage.test.ts index 958b0ff6a..376de5d38 100644 --- a/fluxer_api/src/api/admin/tests/AdminSearchFieldCoverage.test.ts +++ b/fluxer_api/src/api/admin/tests/AdminSearchFieldCoverage.test.ts @@ -5,6 +5,7 @@ import type {UserAdminResponse} from '@fluxer/schema/src/domains/admin/AdminUser import {afterEach, beforeEach, describe, expect, test} from 'vitest'; import {createTestAccount, setUserACLs} from '../../auth/tests/AuthTestUtils'; import {createGuild} from '../../channel/tests/ChannelTestUtils'; +import {getUserActivityBuffer} from '../../middleware/ServiceSingletons'; import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness'; import {HTTP_STATUS} from '../../test/TestConstants'; import {createBuilder, createBuilderWithoutAuth} from '../../test/TestRequestBuilder'; @@ -40,6 +41,7 @@ async function setLastActiveIp(harness: ApiTestHarness, token: string, ip: strin .header('x-forwarded-for', ip) .expect(HTTP_STATUS.OK) .execute(); + await getUserActivityBuffer().drainAndFlush(); } describe('Admin Search Field Coverage', () => { diff --git a/fluxer_api/src/api/auth/AuthSession.ts b/fluxer_api/src/api/auth/AuthSession.ts index 707f7e60e..5435bd885 100644 --- a/fluxer_api/src/api/auth/AuthSession.ts +++ b/fluxer_api/src/api/auth/AuthSession.ts @@ -35,15 +35,6 @@ interface LogoutAuthSessionsParams { sessionIdHashes: Array; } -interface UpdateUserActivityParams { - userId: UserID; - clientIp: string; - user?: User; - action?: 'session_authenticated' | 'bearer_fallback_session_authenticated' | 'unknown'; - tokenType?: 'session' | 'bearer'; - sessionId?: string; -} - interface DispatchAuthSessionChangeParams { userId: UserID; oldAuthSessionIdHash: string; @@ -153,11 +144,6 @@ export async function updateAuthSessionLastUsed(ctx: ApiContext, tokenHash: Uint await ctx.services.userActivityBuffer.recordAuthSessionActivity(Buffer.from(tokenHash), new Date()); } -export async function updateUserActivity(ctx: ApiContext, {userId, clientIp}: UpdateUserActivityParams): Promise { - const {users} = ctx.services; - await users.updateUserActivity(userId, clientIp); -} - export async function revokeToken(ctx: ApiContext, token: string): Promise { const {users, gateway} = ctx.services; const tokenHash = Buffer.from(AuthUtility.getTokenIdHash(ctx, token)); diff --git a/fluxer_api/src/api/middleware/UserMiddleware.ts b/fluxer_api/src/api/middleware/UserMiddleware.ts index 570520c4d..0814a6d27 100644 --- a/fluxer_api/src/api/middleware/UserMiddleware.ts +++ b/fluxer_api/src/api/middleware/UserMiddleware.ts @@ -112,17 +112,6 @@ export const UserMiddleware = createMiddleware(async (ctx, next) => { if (authSession) { void AuthSession.updateAuthSessionLastUsed(apiContext, authSession.sessionIdHash); const user = await apiContext.services.users.findUniqueAssert(authSession.userId); - const sessionId = Buffer.from(authSession.sessionIdHash).toString('base64url'); - void AuthSession.updateUserActivity(apiContext, { - userId: authSession.userId, - clientIp: resolvedClientIp, - user, - action: 'session_authenticated', - tokenType: 'session', - sessionId, - }).catch((error: unknown) => { - Logger.warn({error, userId: authSession.userId}, 'Failed to update user activity telemetry'); - }); ctx.set('authSession', authSession); ctx.set('authTokenType', 'session'); setUserInContext(ctx, user, true); diff --git a/fluxer_api/src/api/user/repositories/IUserAuthRepository.ts b/fluxer_api/src/api/user/repositories/IUserAuthRepository.ts index ea4d59e7e..59f24752e 100644 --- a/fluxer_api/src/api/user/repositories/IUserAuthRepository.ts +++ b/fluxer_api/src/api/user/repositories/IUserAuthRepository.ts @@ -44,7 +44,6 @@ export interface IUserAuthRepository { createPhoneToken(token: PhoneVerificationToken, phone: string, userId: UserID | null): Promise; getPhoneToken(token: PhoneVerificationToken): Promise; deletePhoneToken(token: PhoneVerificationToken): Promise; - updateUserActivity(userId: UserID, clientIp: string): Promise; checkIpAuthorized(userId: UserID, ip: string): Promise; createAuthorizedIp(userId: UserID, ip: string): Promise; createIpAuthorizationToken(userId: UserID, token: string, email: string): Promise; diff --git a/fluxer_api/src/api/user/repositories/UserAuthRepository.ts b/fluxer_api/src/api/user/repositories/UserAuthRepository.ts index c7b15ea06..26d84f3b5 100644 --- a/fluxer_api/src/api/user/repositories/UserAuthRepository.ts +++ b/fluxer_api/src/api/user/repositories/UserAuthRepository.ts @@ -153,10 +153,6 @@ export class UserAuthRepository implements IUserAuthRepository { return this.tokenRepository.deletePhoneToken(token); } - async updateUserActivity(userId: UserID, clientIp: string): Promise { - return this.ipAuthorizationRepository.updateUserActivity(userId, clientIp); - } - async checkIpAuthorized(userId: UserID, ip: string): Promise { return this.ipAuthorizationRepository.checkIpAuthorized(userId, ip); } diff --git a/fluxer_api/src/api/user/repositories/UserRepository.ts b/fluxer_api/src/api/user/repositories/UserRepository.ts index 67ee137c5..9cb3703ee 100644 --- a/fluxer_api/src/api/user/repositories/UserRepository.ts +++ b/fluxer_api/src/api/user/repositories/UserRepository.ts @@ -348,10 +348,6 @@ export class UserRepository implements IUserRepositoryAggregate { return this.authRepo.deletePhoneToken(token); } - async updateUserActivity(userId: UserID, clientIp: string): Promise { - return this.authRepo.updateUserActivity(userId, clientIp); - } - async checkIpAuthorized(userId: UserID, ip: string): Promise { return this.authRepo.checkIpAuthorized(userId, ip); } diff --git a/fluxer_api/src/api/user/repositories/account/crud/UserDataRepository.ts b/fluxer_api/src/api/user/repositories/account/crud/UserDataRepository.ts index fd0da4eee..34ca02e79 100644 --- a/fluxer_api/src/api/user/repositories/account/crud/UserDataRepository.ts +++ b/fluxer_api/src/api/user/repositories/account/crud/UserDataRepository.ts @@ -174,6 +174,7 @@ export class UserDataRepository { }; }> { const {userId, lastActiveAt, lastActiveIp} = params; + const previousData = (await this.getActivityTracking(userId)) ?? {last_active_at: null, last_active_ip: null}; await upsertOne( Users.patchByPk( {user_id: userId}, @@ -184,7 +185,7 @@ export class UserDataRepository { ), ); return { - previousData: {last_active_at: null, last_active_ip: null}, + previousData, updatedData: {last_active_at: lastActiveAt, last_active_ip: lastActiveIp ?? null}, }; } diff --git a/fluxer_api/src/api/user/repositories/account/crud/UserIndexRepository.ts b/fluxer_api/src/api/user/repositories/account/crud/UserIndexRepository.ts index 114c04557..832c035ae 100644 --- a/fluxer_api/src/api/user/repositories/account/crud/UserIndexRepository.ts +++ b/fluxer_api/src/api/user/repositories/account/crud/UserIndexRepository.ts @@ -163,7 +163,7 @@ export class UserIndexRepository { ); } } - await batch.execute(); + await batch.execute(false); } async deleteIndices( diff --git a/fluxer_api/src/api/user/repositories/auth/IpAuthorizationRepository.ts b/fluxer_api/src/api/user/repositories/auth/IpAuthorizationRepository.ts index b13e6d06e..9b7556c8a 100644 --- a/fluxer_api/src/api/user/repositories/auth/IpAuthorizationRepository.ts +++ b/fluxer_api/src/api/user/repositories/auth/IpAuthorizationRepository.ts @@ -111,15 +111,6 @@ export class IpAuthorizationRepository { return {userId: result.user_id, email: result.email}; } - async updateUserActivity(userId: UserID, clientIp: string): Promise { - const now = new Date(); - await this.userAccountRepository.updateLastActiveAt({ - userId, - lastActiveAt: now, - lastActiveIp: clientIp, - }); - } - async getAuthorizedIps(userId: UserID): Promise< Array<{ ip: string; diff --git a/fluxer_api/src/api/user/services/UserActivityBuffer.ts b/fluxer_api/src/api/user/services/UserActivityBuffer.ts index ba2112d72..6c240ca3c 100644 --- a/fluxer_api/src/api/user/services/UserActivityBuffer.ts +++ b/fluxer_api/src/api/user/services/UserActivityBuffer.ts @@ -6,8 +6,9 @@ import type {UserID} from '../../BrandedTypes'; import {upsertOne} from '../../database/CassandraQueryExecution'; import {Db} from '../../database/CassandraTypes'; import {Logger} from '../../Logger'; -import {AuthSessions, Users} from '../../Tables'; +import {AuthSessions} from '../../Tables'; import {isJsonRecord, parseJsonRecord} from '../../utils/JsonBoundaryUtils'; +import {UserAccountRepository} from '../repositories/account/UserAccountRepository'; const PENDING_HASH_KEY = 'user_activity:pending'; const PENDING_AUTH_SESSION_HASH_KEY = 'auth_session_activity:pending'; @@ -16,6 +17,10 @@ const WRITE_CONCURRENCY = 64; const AUTH_SESSION_TOUCH_DEBOUNCE_TTL_SECONDS = seconds('5 minutes'); type ActivityWriter = typeof upsertOne; +interface UserActivityAccountWriter { + updateLastActiveAt(params: {userId: UserID; lastActiveAt: Date; lastActiveIp?: string}): Promise; +} + interface PendingEntry { ts: number; ip: string | null; @@ -58,10 +63,16 @@ function isStringRecord(value: unknown): value is Record { export class UserActivityBuffer { private readonly kv: IKVProvider; private readonly writer: ActivityWriter; + private readonly accounts: UserActivityAccountWriter; - constructor(kv: IKVProvider, writer: ActivityWriter = upsertOne) { + constructor( + kv: IKVProvider, + writer: ActivityWriter = upsertOne, + accounts: UserActivityAccountWriter = new UserAccountRepository(kv), + ) { this.kv = kv; this.writer = writer; + this.accounts = accounts; } recordActivity(userId: UserID, timestamp: Date, ip: string | null): void { @@ -128,15 +139,11 @@ export class UserActivityBuffer { const chunk = drained.slice(i, i + WRITE_CONCURRENCY); const results = await Promise.allSettled( chunk.map(({userId, entry}) => - this.writer( - Users.patchByPk( - {user_id: userId}, - { - last_active_at: Db.set(new Date(entry.ts)), - last_active_ip: entry.ip !== null ? Db.set(entry.ip) : Db.clear(), - }, - ), - ), + this.accounts.updateLastActiveAt({ + userId, + lastActiveAt: new Date(entry.ts), + lastActiveIp: entry.ip ?? undefined, + }), ), ); for (const r of results) { diff --git a/fluxer_api/src/api/user/tests/UserActivityFlush.test.ts b/fluxer_api/src/api/user/tests/UserActivityFlush.test.ts new file mode 100644 index 000000000..0660068ae --- /dev/null +++ b/fluxer_api/src/api/user/tests/UserActivityFlush.test.ts @@ -0,0 +1,53 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {afterAll, beforeAll, beforeEach, describe, expect, test} from 'vitest'; +import {createTestAccount} from '../../auth/tests/AuthTestUtils'; +import {createUserID} from '../../BrandedTypes'; +import {getUserActivityBuffer, getUserRepository} from '../../middleware/ServiceSingletons'; +import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness'; +import {HTTP_STATUS} from '../../test/TestConstants'; +import {createBuilder} from '../../test/TestRequestBuilder'; + +const REGISTRATION_IP = '198.51.100.7'; +const REQUEST_IP = '203.0.113.7'; + +async function fetchMeFromIp(harness: ApiTestHarness, token: string, ip: string): Promise { + await createBuilder(harness, token).get('/users/@me').header('x-forwarded-for', ip).expect(HTTP_STATUS.OK).execute(); +} + +describe('User activity buffering', () => { + let harness: ApiTestHarness; + beforeAll(async () => { + harness = await createApiTestHarness(); + }); + beforeEach(async () => { + await harness.reset(); + }); + afterAll(async () => { + await harness?.shutdown(); + }); + test('an authenticated request buffers last active instead of writing it', async () => { + const account = await createTestAccount(harness, {ipAddress: REGISTRATION_IP}); + const userId = createUserID(BigInt(account.userId)); + await fetchMeFromIp(harness, account.token, REQUEST_IP); + const pending = await harness.kvProvider.hgetall('user_activity:pending'); + expect(pending[account.userId]).toBeDefined(); + const beforeFlush = await getUserRepository().getActivityTracking(userId); + expect(beforeFlush?.last_active_ip).toBe(REGISTRATION_IP); + await getUserActivityBuffer().drainAndFlush(); + const afterFlush = await getUserRepository().getActivityTracking(userId); + expect(afterFlush?.last_active_ip).toBe(REQUEST_IP); + }); + test('flushing moves the last active IP index and prunes the previous address', async () => { + const account = await createTestAccount(harness, {ipAddress: REGISTRATION_IP}); + const userRepository = getUserRepository(); + const seeded = await userRepository.listUserIdsByLastActiveIp(REGISTRATION_IP, 10, 0); + expect(seeded.userIds.map((id) => id.toString())).toContain(account.userId); + await fetchMeFromIp(harness, account.token, REQUEST_IP); + await getUserActivityBuffer().drainAndFlush(); + const previous = await userRepository.listUserIdsByLastActiveIp(REGISTRATION_IP, 10, 0); + expect(previous.userIds.map((id) => id.toString())).not.toContain(account.userId); + const current = await userRepository.listUserIdsByLastActiveIp(REQUEST_IP, 10, 0); + expect(current.userIds.map((id) => id.toString())).toContain(account.userId); + }); +});